diff --git a/broker/src/main/java/org/apache/rocketmq/broker/transaction/TransactionalMessageCheckService.java b/broker/src/main/java/org/apache/rocketmq/broker/transaction/TransactionalMessageCheckService.java index 143889a19c..6a3c2d2b29 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/transaction/TransactionalMessageCheckService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/transaction/TransactionalMessageCheckService.java @@ -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!"); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java index 27be2f4333..5ae5abde19 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java @@ -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"); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java index 1c8cea0125..819349ede5 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java @@ -96,12 +96,6 @@ public class ConsumerService extends BaseService { public CompletableFuture receiveMessage(Context ctx, ReceiveMessageRequest request) { CompletableFuture 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 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 future = new CompletableFuture<>(); ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader(); @@ -286,7 +276,8 @@ public class ConsumerService extends BaseService { return future; } - protected CompletableFuture processAckMessage(Context ctx, AckMessageRequest request, AckMessageEntry ackMessageEntry) { + protected CompletableFuture processAckMessage(Context ctx, AckMessageRequest request, + AckMessageEntry ackMessageEntry) { CompletableFuture 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 nackMessage(Context ctx, NackMessageRequest request) { CompletableFuture 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 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 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 changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request) { CompletableFuture 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 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; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java index d40691b3f6..55d6fe5bf7 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java @@ -71,11 +71,6 @@ public class ProducerService extends BaseService { public CompletableFuture sendMessage(Context ctx, SendMessageRequest request) { CompletableFuture future = new CompletableFuture<>(); - future.whenComplete((response, throwable) -> { - if (sendMessageHook != null) { - sendMessageHook.beforeResponse(ctx, request, response, throwable); - } - }); try { Pair> requestPair = this.buildSendMessageRequest(ctx, request); @@ -89,28 +84,21 @@ public class ProducerService extends BaseService { } // send message to broker. - CompletableFuture 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 forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request) { CompletableFuture 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 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(); + } } diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcTest.java index 3a7ff341d0..84fa336a3e 100644 --- a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcTest.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcTest.java @@ -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(); diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java index 75c881039b..97e714f49e 100644 --- a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java @@ -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 sendClientSettings(MessagingServiceGrpc.MessagingServiceStub stub, Settings clientSettings) { + protected CompletableFuture sendClientSettings(MessagingServiceGrpc.MessagingServiceStub stub, + Settings clientSettings) { CompletableFuture future = new CompletableFuture<>(); StreamObserver 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 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 receiveMessage(MessagingServiceGrpc.MessagingServiceBlockingStub stub, String topic, String group) { + public List receiveMessage(MessagingServiceGrpc.MessagingServiceBlockingStub stub, + String topic, String group) { List responseList = new ArrayList<>(); Iterator 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 receiveMessage(MessagingServiceGrpc.MessagingServiceBlockingStub stub, String topic, String group, int timeSeconds) { + public List receiveMessage(MessagingServiceGrpc.MessagingServiceBlockingStub stub, + String topic, String group, int timeSeconds) { List responseList = new ArrayList<>(); - Iterator responseIterator = stub.withDeadlineAfter(timeSeconds, TimeUnit.SECONDS) - .receiveMessage(buildReceiveMessageRequest(group, topic)); + Iterator 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() diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java index 282fff469b..6ddefba797 100644 --- a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java @@ -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(); + } }