[ISSUE #3949] v2 test cases

This commit is contained in:
kaiyi.lk
2022-07-13 11:29:20 +08:00
committed by zhouxiang
parent f36fc44a3f
commit 6b877b5bdd
7 changed files with 206 additions and 145 deletions
@@ -42,8 +42,8 @@ public class TransactionalMessageCheckService extends ServiceThread {
@Override
public void run() {
log.info("Start transaction check service thread!");
long checkInterval = brokerController.getBrokerConfig().getTransactionCheckInterval();
while (!this.isStopped()) {
long checkInterval = brokerController.getBrokerConfig().getTransactionCheckInterval();
this.waitForRunning(checkInterval);
}
log.info("End transaction check service thread!");
@@ -29,7 +29,7 @@ public class ResponseBuilder {
t = t.getCause();
}
if (t instanceof ProxyException) {
ProxyException proxyException = (ProxyException) t.getCause();
ProxyException proxyException = (ProxyException) t;
return ResponseBuilder.buildStatus(proxyException.getCode(), proxyException.getMessage());
}
return ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "internal error");
@@ -96,12 +96,6 @@ public class ConsumerService extends BaseService {
public CompletableFuture<ReceiveMessageResponse> receiveMessage(Context ctx, ReceiveMessageRequest request) {
CompletableFuture<ReceiveMessageResponse> future = new CompletableFuture<>();
// register hook.
future.whenComplete((response, throwable) -> {
if (receiveMessageHook != null) {
receiveMessageHook.beforeResponse(ctx, request, response, throwable);
}
});
try {
PopMessageRequestHeader requestHeader = this.buildPopMessageRequestHeader(ctx, request);
@@ -111,26 +105,20 @@ public class ConsumerService extends BaseService {
throw new ProxyException(Code.FORBIDDEN, "no readable topic route for topic " + requestHeader.getTopic());
}
CompletableFuture<PopResult> popResultFuture = this.readConsumer.popMessage(
future = this.readConsumer.popMessage(
messageQueue.getBrokerAddr(),
messageQueue.getBrokerName(),
requestHeader,
requestHeader.getPollTime());
popResultFuture
.thenAccept(result -> {
try {
future.complete(convertToReceiveMessageResponse(ctx, request, result));
} catch (Throwable throwable) {
future.completeExceptionally(throwable);
}
})
.exceptionally(throwable -> {
future.completeExceptionally(throwable);
return null;
});
requestHeader.getPollTime())
.thenApply(result -> convertToReceiveMessageResponse(ctx, request, result));
} catch (Throwable t) {
future.completeExceptionally(t);
}
future.whenComplete((response, throwable) -> {
if (receiveMessageHook != null) {
receiveMessageHook.beforeResponse(ctx, request, response, throwable);
}
});
return future;
}
@@ -141,7 +129,8 @@ public class ConsumerService extends BaseService {
return GrpcConverter.buildPopMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx), fifo);
}
protected ReceiveMessageResponse convertToReceiveMessageResponse(Context ctx, ReceiveMessageRequest request, PopResult result) {
protected ReceiveMessageResponse convertToReceiveMessageResponse(Context ctx, ReceiveMessageRequest request,
PopResult result) {
PopStatus status = result.getPopStatus();
switch (status) {
case FOUND:
@@ -189,7 +178,8 @@ public class ConsumerService extends BaseService {
return resMessageList;
}
protected void forwardMessageToDLQ(Context ctx, ReceiveMessageRequest request, MessageExt messageExt, int maxReconsumeTimes) {
protected void forwardMessageToDLQ(Context ctx, ReceiveMessageRequest request, MessageExt messageExt,
int maxReconsumeTimes) {
CompletableFuture<RemotingCommand> future = new CompletableFuture<>();
ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader();
@@ -286,7 +276,8 @@ public class ConsumerService extends BaseService {
return future;
}
protected CompletableFuture<AckMessageResultEntry> processAckMessage(Context ctx, AckMessageRequest request, AckMessageEntry ackMessageEntry) {
protected CompletableFuture<AckMessageResultEntry> processAckMessage(Context ctx, AckMessageRequest request,
AckMessageEntry ackMessageEntry) {
CompletableFuture<AckMessageResultEntry> future = new CompletableFuture<>();
AckMessageResultEntry.Builder failResult = AckMessageResultEntry.newBuilder()
.setStatus(ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "ack message failed"))
@@ -311,11 +302,13 @@ public class ConsumerService extends BaseService {
return future;
}
protected AckMessageRequestHeader buildAckMessageRequestHeader(Context ctx, AckMessageRequest request, ReceiptHandle handle) {
protected AckMessageRequestHeader buildAckMessageRequestHeader(Context ctx, AckMessageRequest request,
ReceiptHandle handle) {
return GrpcConverter.buildAckMessageRequestHeader(request, handle);
}
protected AckMessageResultEntry convertToAckMessageResultEntry(Context ctx, AckMessageEntry ackMessageEntry, AckResult ackResult) {
protected AckMessageResultEntry convertToAckMessageResultEntry(Context ctx, AckMessageEntry ackMessageEntry,
AckResult ackResult) {
if (AckStatus.OK.equals(ackResult.getStatus())) {
return AckMessageResultEntry.newBuilder()
.setMessageId(ackMessageEntry.getMessageId())
@@ -332,11 +325,6 @@ public class ConsumerService extends BaseService {
public CompletableFuture<NackMessageResponse> nackMessage(Context ctx, NackMessageRequest request) {
CompletableFuture<NackMessageResponse> future = new CompletableFuture<>();
future.whenComplete((response, throwable) -> {
if (nackMessageHook != null) {
nackMessageHook.beforeResponse(ctx, request, response, throwable);
}
});
try {
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
@@ -344,51 +332,35 @@ public class ConsumerService extends BaseService {
Settings settings = grpcClientManager.getClientSettings(ctx);
int maxDeliveryAttempts = settings.getSubscription().getBackoffPolicy().getMaxAttempts();
if (request.getDeliveryAttempt() >= maxDeliveryAttempts) {
CompletableFuture<RemotingCommand> resultFuture = this.producer.sendMessageBack(
future = this.producer.sendMessageBack(
brokerAddr,
this.buildConsumerSendMsgBackToDLQRequestHeader(ctx, request, maxDeliveryAttempts)
);
resultFuture
.thenAccept(result -> {
try {
future.complete(convertToNackMessageResponse(ctx, request, result));
if (result.getCode() == ResponseCode.SUCCESS) {
writeConsumer.ackMessage(
brokerAddr,
this.buildAckMessageRequestHeader(ctx, request));
}
} catch (Throwable throwable) {
future.completeExceptionally(throwable);
}
})
.exceptionally(throwable -> {
future.completeExceptionally(throwable);
return null;
});
).thenApply(result -> {
if (result.getCode() == ResponseCode.SUCCESS) {
writeConsumer.ackMessage(
brokerAddr,
this.buildAckMessageRequestHeader(ctx, request));
}
return convertToNackMessageResponse(ctx, request, result);
});
} else {
ChangeInvisibleTimeRequestHeader requestHeader = this.buildChangeInvisibleTimeRequestHeader(ctx, request);
CompletableFuture<AckResult> resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader);
resultFuture
.thenAccept(result -> {
try {
future.complete(convertToNackMessageResponse(ctx, request, result));
} catch (Throwable throwable) {
future.completeExceptionally(throwable);
}
})
.exceptionally(throwable -> {
future.completeExceptionally(throwable);
return null;
});
future = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader)
.thenApply(result -> convertToNackMessageResponse(ctx, request, result));
}
} catch (Throwable t) {
future.completeExceptionally(t);
}
future.whenComplete((response, throwable) -> {
if (nackMessageHook != null) {
nackMessageHook.beforeResponse(ctx, request, response, throwable);
}
});
return future;
}
protected ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(Context ctx, NackMessageRequest request) {
protected ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(Context ctx,
NackMessageRequest request) {
return GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, this.delayPolicy);
}
@@ -396,12 +368,14 @@ public class ConsumerService extends BaseService {
return GrpcConverter.buildAckMessageRequestHeader(request);
}
protected ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackToDLQRequestHeader(Context ctx, NackMessageRequest request,
protected ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackToDLQRequestHeader(Context ctx,
NackMessageRequest request,
int maxReconsumeTimes) {
return GrpcConverter.buildConsumerSendMsgBackToDLQRequestHeader(request, maxReconsumeTimes);
}
protected NackMessageResponse convertToNackMessageResponse(Context ctx, NackMessageRequest request, AckResult ackResult) {
protected NackMessageResponse convertToNackMessageResponse(Context ctx, NackMessageRequest request,
AckResult ackResult) {
if (AckStatus.OK.equals(ackResult.getStatus())) {
return NackMessageResponse.newBuilder()
.setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name()))
@@ -412,7 +386,8 @@ public class ConsumerService extends BaseService {
.build();
}
protected NackMessageResponse convertToNackMessageResponse(Context ctx, NackMessageRequest request, RemotingCommand sendMsgBackToDLQResult) {
protected NackMessageResponse convertToNackMessageResponse(Context ctx, NackMessageRequest request,
RemotingCommand sendMsgBackToDLQResult) {
return NackMessageResponse.newBuilder()
.setStatus(ResponseBuilder.buildStatus(sendMsgBackToDLQResult.getCode(), sendMsgBackToDLQResult.getRemark()))
.build();
@@ -421,32 +396,22 @@ public class ConsumerService extends BaseService {
public CompletableFuture<ChangeInvisibleDurationResponse> changeInvisibleDuration(Context ctx,
ChangeInvisibleDurationRequest request) {
CompletableFuture<ChangeInvisibleDurationResponse> future = new CompletableFuture<>();
future.whenComplete((response, throwable) -> {
if (changeInvisibleDurationHook != null) {
changeInvisibleDurationHook.beforeResponse(ctx, request, response, throwable);
}
});
try {
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
ChangeInvisibleTimeRequestHeader requestHeader = convertToChangeInvisibleTimeRequestHeader(ctx, request);
CompletableFuture<AckResult> resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader);
resultFuture
.thenAccept(result -> {
try {
future.complete(convertToChangeInvisibleDurationResponse(ctx, request, result));
} catch (Throwable throwable) {
future.completeExceptionally(throwable);
}
})
.exceptionally(throwable -> {
future.completeExceptionally(throwable);
return null;
});
future = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader)
.thenApply(result -> convertToChangeInvisibleDurationResponse(ctx, request, result));
} catch (Throwable t) {
future.completeExceptionally(t);
}
future.whenComplete((response, throwable) -> {
if (changeInvisibleDurationHook != null) {
changeInvisibleDurationHook.beforeResponse(ctx, request, response, throwable);
}
});
return future;
}
@@ -71,11 +71,6 @@ public class ProducerService extends BaseService {
public CompletableFuture<SendMessageResponse> sendMessage(Context ctx, SendMessageRequest request) {
CompletableFuture<SendMessageResponse> future = new CompletableFuture<>();
future.whenComplete((response, throwable) -> {
if (sendMessageHook != null) {
sendMessageHook.beforeResponse(ctx, request, response, throwable);
}
});
try {
Pair<SendMessageRequestHeader, List<org.apache.rocketmq.common.message.Message>> requestPair = this.buildSendMessageRequest(ctx, request);
@@ -89,28 +84,21 @@ public class ProducerService extends BaseService {
}
// send message to broker.
CompletableFuture<SendResult> sendResultCompletableFuture = this.producer.sendMessage(
future = this.producer.sendMessage(
selectableMessageQueue.getBrokerAddr(),
selectableMessageQueue.getBrokerName(),
message,
requestHeader
);
sendResultCompletableFuture
.thenAccept(result -> {
try {
future.complete(convertToSendMessageResponse(ctx, request, result));
} catch (Throwable throwable) {
future.completeExceptionally(throwable);
}
})
.exceptionally(e -> {
future.completeExceptionally(e);
return null;
});
).thenApply(result -> convertToSendMessageResponse(ctx, request, result));
} catch (Throwable t) {
future.completeExceptionally(t);
}
future.whenComplete((response, throwable) -> {
if (sendMessageHook != null) {
sendMessageHook.beforeResponse(ctx, request, response, throwable);
}
});
return future;
}
@@ -145,11 +133,6 @@ public class ProducerService extends BaseService {
public CompletableFuture<ForwardMessageToDeadLetterQueueResponse> forwardMessageToDeadLetterQueue(Context ctx,
ForwardMessageToDeadLetterQueueRequest request) {
CompletableFuture<ForwardMessageToDeadLetterQueueResponse> future = new CompletableFuture<>();
future.whenComplete((response, throwable) -> {
if (forwardMessageToDLQHook != null) {
forwardMessageToDLQHook.beforeResponse(ctx, request, response, throwable);
}
});
try {
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
@@ -158,22 +141,16 @@ public class ProducerService extends BaseService {
AckMessageRequestHeader ackMessageRequestHeader = GrpcConverter.buildAckMessageRequestHeader(
request.getTopic(), request.getGroup(), receiptHandle);
CompletableFuture<RemotingCommand> resultFuture = this.producer.sendMessageBackThenAckOrg(brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader);
resultFuture
.thenAccept(result ->
future.complete(
ForwardMessageToDeadLetterQueueResponse.newBuilder()
.setStatus(ResponseBuilder.buildStatus(result.getCode(), result.getRemark()))
.build()
)
)
.exceptionally(throwable -> {
future.completeExceptionally(throwable);
return null;
});
future = this.producer.sendMessageBackThenAckOrg(brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader)
.thenApply(result -> convertToForwardMessageToDeadLetterQueueResponse(ctx, result));
} catch (Throwable t) {
future.completeExceptionally(t);
}
future.whenComplete((response, throwable) -> {
if (forwardMessageToDLQHook != null) {
forwardMessageToDLQHook.beforeResponse(ctx, request, response, throwable);
}
});
return future;
}
@@ -181,4 +158,11 @@ public class ProducerService extends BaseService {
ForwardMessageToDeadLetterQueueRequest request) {
return GrpcConverter.buildConsumerSendMsgBackRequestHeader(request);
}
protected ForwardMessageToDeadLetterQueueResponse convertToForwardMessageToDeadLetterQueueResponse(Context ctx,
RemotingCommand result) {
return ForwardMessageToDeadLetterQueueResponse.newBuilder()
.setStatus(ResponseBuilder.buildStatus(result.getCode(), result.getRemark()))
.build();
}
}
@@ -40,6 +40,7 @@ public class ClusterGrpcTest extends GrpcBaseTest {
@Before
public void setUp() throws Exception {
super.setUp();
ConfigurationManager.getProxyConfig().setTransactionHeartbeatPeriodSecond(3);
grpcForwardService = new ClusterGrpcService();
grpcForwardService.start();
GrpcMessagingProcessor processor = new GrpcMessagingProcessor(grpcForwardService);
@@ -92,6 +93,11 @@ public class ClusterGrpcTest extends GrpcBaseTest {
super.testSendReceiveMessageThenToDLQ();
}
@Test
public void testSimpleConsumerSendAndRecv() throws Exception {
super.testSimpleConsumerSendAndRecv();
}
@Test
public void testSimpleConsumerToDLQ() throws Exception {
super.testSimpleConsumerToDLQ();
@@ -20,8 +20,11 @@ package org.apache.rocketmq.test.grpc.v2;
import apache.rocketmq.v2.AckMessageEntry;
import apache.rocketmq.v2.AckMessageRequest;
import apache.rocketmq.v2.AckMessageResponse;
import apache.rocketmq.v2.AckMessageResultEntry;
import apache.rocketmq.v2.Address;
import apache.rocketmq.v2.AddressScheme;
import apache.rocketmq.v2.ChangeInvisibleDurationRequest;
import apache.rocketmq.v2.ChangeInvisibleDurationResponse;
import apache.rocketmq.v2.ClientType;
import apache.rocketmq.v2.Code;
import apache.rocketmq.v2.EndTransactionRequest;
@@ -54,6 +57,7 @@ import apache.rocketmq.v2.TransactionResolution;
import apache.rocketmq.v2.TransactionSource;
import com.google.protobuf.ByteString;
import com.google.protobuf.Duration;
import com.google.protobuf.util.Durations;
import com.google.protobuf.util.Timestamps;
import io.grpc.Channel;
import io.grpc.Metadata;
@@ -77,6 +81,7 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@@ -120,6 +125,10 @@ public class GrpcBaseTest extends BaseConf {
protected static final int defaultQueueNums = 8;
public void setUp() throws Exception {
brokerController1.getBrokerConfig().setTransactionCheckInterval(3 * 1000);
brokerController2.getBrokerConfig().setTransactionCheckInterval(3 * 1000);
brokerController3.getBrokerConfig().setTransactionCheckInterval(3 * 1000);
header.put(InterceptorConstants.CLIENT_ID, "client-id" + UUID.randomUUID());
header.put(InterceptorConstants.LANGUAGE, "JAVA");
@@ -148,7 +157,8 @@ public class GrpcBaseTest extends BaseConf {
return MetadataUtils.attachHeaders(stub, header);
}
protected CompletableFuture<Settings> sendClientSettings(MessagingServiceGrpc.MessagingServiceStub stub, Settings clientSettings) {
protected CompletableFuture<Settings> sendClientSettings(MessagingServiceGrpc.MessagingServiceStub stub,
Settings clientSettings) {
CompletableFuture<Settings> future = new CompletableFuture<>();
StreamObserver<TelemetryCommand> requestStreamObserver = stub.telemetry(new DefaultTelemetryCommandStreamObserver() {
@Override
@@ -217,8 +227,8 @@ public class GrpcBaseTest extends BaseConf {
ReceiveMessageResponse response = receiveMessage(blockingStub, topic, group).get(0);
assertReceiveMessage(response, messageId);
String receiptHandle = response.getMessages(0).getSystemProperties().getReceiptHandle();
AckMessageResponse ackMessageResponse = blockingStub.ackMessage(buildAckMessageRequest(group, topic, messageId, receiptHandle));
assertAck(ackMessageResponse);
AckMessageResponse ackMessageResponse = blockingStub.ackMessage(buildAckMessageRequest(topic, group, messageId, receiptHandle));
assertAllAckOk(ackMessageResponse);
}
public void testSendReceiveMessageThenToDLQ() throws Exception {
@@ -241,7 +251,7 @@ public class GrpcBaseTest extends BaseConf {
Message message = receiveResponse.getMessages(0);
NackMessageResponse nackMessageResponse = blockingStub.nackMessage(buildNackMessageRequest(
group, topic, messageId, message.getSystemProperties().getReceiptHandle(), 1
topic, group, messageId, message.getSystemProperties().getReceiptHandle(), 1
));
assertNackMessageResponse(nackMessageResponse);
@@ -258,7 +268,7 @@ public class GrpcBaseTest extends BaseConf {
message = receiveRetryResponseRef.get().getMessages(0);
nackMessageResponse = blockingStub.nackMessage(buildNackMessageRequest(
group, topic, messageId, message.getSystemProperties().getReceiptHandle(), 2
topic, group, messageId, message.getSystemProperties().getReceiptHandle(), 2
));
assertNackMessageResponse(nackMessageResponse);
@@ -331,7 +341,7 @@ public class GrpcBaseTest extends BaseConf {
SendMessageResponse sendResponse = blockingStub.sendMessage(buildTransactionSendMessageRequest(topic, messageId));
assertSendMessage(sendResponse, messageId);
await().atMost(java.time.Duration.ofSeconds(60)).until(() -> {
await().atMost(java.time.Duration.ofSeconds(90)).until(() -> {
if (telemetryCommandRef.get() == null) {
return false;
}
@@ -364,6 +374,64 @@ public class GrpcBaseTest extends BaseConf {
}
}
public void testSimpleConsumerSendAndRecv() throws Exception {
String topic = initTopicOnSampleTopicBroker(broker1Name);
String group = MQRandomUtils.getRandomConsumerGroup();
int maxDeliveryAttempts = 16;
boolean fifo = false;
// init consumer offset
this.sendClientSettings(stub, buildSimpleConsumerClientSettings(maxDeliveryAttempts, fifo)).get();
receiveMessage(blockingStub, topic, group);
this.sendClientSettings(stub, buildProducerClientSettings(topic)).get();
String messageId = createUniqID();
SendMessageResponse sendResponse = blockingStub.sendMessage(buildSendMessageRequest(topic, messageId));
assertSendMessage(sendResponse, messageId);
this.sendClientSettings(stub, buildSimpleConsumerClientSettings(maxDeliveryAttempts, fifo)).get();
ReceiveMessageResponse receiveResponse = receiveMessage(blockingStub, topic, group).get(0);
assertReceiveMessage(receiveResponse, messageId);
String receiptHandle = receiveResponse.getMessages(0).getSystemProperties().getReceiptHandle();
ChangeInvisibleDurationResponse changeResponse = blockingStub.changeInvisibleDuration(buildChangeInvisibleDurationRequest(topic, group, receiptHandle, 5));
assertChangeInvisibleDurationResponse(changeResponse, receiptHandle);
List<String> ackHandles = new ArrayList<>();
ackHandles.add(changeResponse.getReceiptHandle());
await().atMost(java.time.Duration.ofSeconds(20)).until(() -> {
ReceiveMessageResponse receiveRetryResponse = receiveMessage(blockingStub, topic, group).get(0);
if (receiveRetryResponse.getMessagesCount() <= 0) {
return false;
}
if (receiveRetryResponse.getMessages(0).getSystemProperties()
.getMessageId().equals(messageId)) {
ackHandles.add(receiveRetryResponse.getMessages(0).getSystemProperties().getReceiptHandle());
return true;
}
return false;
});
assertThat(ackHandles.size()).isEqualTo(2);
AckMessageResponse ackMessageResponse = blockingStub.ackMessage(buildAckMessageRequest(topic, group,
AckMessageEntry.newBuilder().setMessageId(messageId).setReceiptHandle(ackHandles.get(0)).build(),
AckMessageEntry.newBuilder().setMessageId(messageId).setReceiptHandle(ackHandles.get(1)).build()));
assertThat(ackMessageResponse.getStatus().getCode()).isEqualTo(Code.OK);
int okNum = 0;
int expireNum = 0;
for (AckMessageResultEntry entry : ackMessageResponse.getEntriesList()) {
if (entry.getStatus().getCode().equals(Code.OK)) {
okNum++;
} else if (entry.getStatus().getCode().equals(Code.RECEIPT_HANDLE_EXPIRED)) {
expireNum++;
}
}
assertThat(okNum).isEqualTo(1);
assertThat(expireNum).isEqualTo(1);
}
public void testSimpleConsumerToDLQ() throws Exception {
String topic = initTopicOnSampleTopicBroker(broker1Name);
String group = MQRandomUtils.getRandomConsumerGroup();
@@ -411,20 +479,22 @@ public class GrpcBaseTest extends BaseConf {
assertThat(receiveMessageCount.get()).isEqualTo(maxDeliveryAttempts);
}
public List<ReceiveMessageResponse> receiveMessage(MessagingServiceGrpc.MessagingServiceBlockingStub stub, String topic, String group) {
public List<ReceiveMessageResponse> receiveMessage(MessagingServiceGrpc.MessagingServiceBlockingStub stub,
String topic, String group) {
List<ReceiveMessageResponse> responseList = new ArrayList<>();
Iterator<ReceiveMessageResponse> responseIterator = stub.withDeadlineAfter(15, TimeUnit.SECONDS)
.receiveMessage(buildReceiveMessageRequest(group, topic));
.receiveMessage(buildReceiveMessageRequest(topic, group));
while (responseIterator.hasNext()) {
responseList.add(responseIterator.next());
}
return responseList;
}
public List<ReceiveMessageResponse> receiveMessage(MessagingServiceGrpc.MessagingServiceBlockingStub stub, String topic, String group, int timeSeconds) {
public List<ReceiveMessageResponse> receiveMessage(MessagingServiceGrpc.MessagingServiceBlockingStub stub,
String topic, String group, int timeSeconds) {
List<ReceiveMessageResponse> responseList = new ArrayList<>();
Iterator<ReceiveMessageResponse> responseIterator = stub.withDeadlineAfter(timeSeconds, TimeUnit.SECONDS)
.receiveMessage(buildReceiveMessageRequest(group, topic));
Iterator<ReceiveMessageResponse> responseIterator = stub.withDeadlineAfter(timeSeconds, TimeUnit.SECONDS)
.receiveMessage(buildReceiveMessageRequest(topic, group));
while (responseIterator.hasNext()) {
responseList.add(responseIterator.next());
}
@@ -493,7 +563,7 @@ public class GrpcBaseTest extends BaseConf {
.build();
}
public ReceiveMessageRequest buildReceiveMessageRequest(String group, String topic) {
public ReceiveMessageRequest buildReceiveMessageRequest(String topic, String group) {
return ReceiveMessageRequest.newBuilder()
.setGroup(Resource.newBuilder()
.setName(group)
@@ -511,7 +581,15 @@ public class GrpcBaseTest extends BaseConf {
.build();
}
public AckMessageRequest buildAckMessageRequest(String group, String topic, String messageId, String receiptHandle) {
public AckMessageRequest buildAckMessageRequest(String topic, String group, String messageId,
String receiptHandle) {
return buildAckMessageRequest(topic, group, AckMessageEntry.newBuilder()
.setMessageId(messageId)
.setReceiptHandle(receiptHandle)
.build());
}
public AckMessageRequest buildAckMessageRequest(String topic, String group, AckMessageEntry... entry) {
return AckMessageRequest.newBuilder()
.setGroup(Resource.newBuilder()
.setName(group)
@@ -519,14 +597,12 @@ public class GrpcBaseTest extends BaseConf {
.setTopic(Resource.newBuilder()
.setName(topic)
.build())
.addEntries(AckMessageEntry.newBuilder()
.setMessageId(messageId)
.setReceiptHandle(receiptHandle)
.build())
.addAllEntries(Arrays.stream(entry).collect(Collectors.toList()))
.build();
}
public NackMessageRequest buildNackMessageRequest(String group, String topic, String messageId, String receiptHandle,
public NackMessageRequest buildNackMessageRequest(String topic, String group, String messageId,
String receiptHandle,
int deliveryAttempt) {
return NackMessageRequest.newBuilder()
.setDeliveryAttempt(deliveryAttempt)
@@ -541,7 +617,8 @@ public class GrpcBaseTest extends BaseConf {
.build();
}
public EndTransactionRequest buildEndTransactionRequest(String topic, String messageId, String transactionId, TransactionResolution resolution) {
public EndTransactionRequest buildEndTransactionRequest(String topic, String messageId, String transactionId,
TransactionResolution resolution) {
return EndTransactionRequest.newBuilder()
.setMessageId(messageId)
.setTopic(Resource.newBuilder()
@@ -553,6 +630,16 @@ public class GrpcBaseTest extends BaseConf {
.build();
}
public ChangeInvisibleDurationRequest buildChangeInvisibleDurationRequest(String topic, String group,
String receiptHandle, int second) {
return ChangeInvisibleDurationRequest.newBuilder()
.setTopic(Resource.newBuilder().setName(topic).build())
.setGroup(Resource.newBuilder().setName(group).build())
.setInvisibleDuration(Durations.fromSeconds(second))
.setReceiptHandle(receiptHandle)
.build();
}
public void assertQueryRoute(QueryRouteResponse response, int messageQueueSize) {
assertThat(response.getStatus().getCode()).isEqualTo(Code.OK);
assertThat(response.getMessageQueuesList().size()).isEqualTo(messageQueueSize);
@@ -580,9 +667,13 @@ public class GrpcBaseTest extends BaseConf {
.getMessageId()).isEqualTo(messageId);
}
public void assertAck(AckMessageResponse response) {
public void assertAllAckOk(AckMessageResponse response) {
assertThat(response.getStatus()
.getCode()).isEqualTo(Code.OK);
for (AckMessageResultEntry entry : response.getEntriesList()) {
assertThat(entry.getStatus()
.getCode()).isEqualTo(Code.OK);
}
}
public void assertNackMessageResponse(NackMessageResponse response) {
@@ -599,6 +690,11 @@ public class GrpcBaseTest extends BaseConf {
assertThat(response.getStatus().getCode()).isEqualTo(Code.OK);
}
public void assertChangeInvisibleDurationResponse(ChangeInvisibleDurationResponse response, String prevHandle) {
assertThat(response.getStatus().getCode()).isEqualTo(Code.OK);
assertThat(response.getReceiptHandle()).isNotEqualTo(prevHandle);
}
public Settings buildAccessPointClientSettings(int port) {
return Settings.newBuilder()
.setAccessPoint(Endpoints.newBuilder()
@@ -66,4 +66,14 @@ public class LocalGrpcTest extends GrpcBaseTest {
public void testSendReceiveMessageThenToDLQ() throws Exception {
super.testSendReceiveMessageThenToDLQ();
}
@Test
public void testSimpleConsumerSendAndRecv() throws Exception {
super.testSimpleConsumerSendAndRecv();
}
@Test
public void testSimpleConsumerToDLQ() throws Exception {
super.testSimpleConsumerToDLQ();
}
}