diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java index 7708f415ea..b4007d413d 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java @@ -454,6 +454,8 @@ public class PopMessageProcessor implements NettyRequestProcessor { response = null; } break; + case ResponseCode.POLLING_TIMEOUT: + return response; default: assert false; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/InvocationChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/InvocationChannel.java index 4e04086797..83dc2428cc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/InvocationChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/InvocationChannel.java @@ -20,7 +20,6 @@ package org.apache.rocketmq.proxy.channel; import io.netty.channel.ChannelFuture; import java.util.Iterator; import java.util.Map; -import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import org.apache.rocketmq.proxy.common.Cleaner; @@ -50,17 +49,9 @@ public abstract class InvocationChannel extends SimpleChannel implements C return super.writeAndFlush(msg); } - public boolean isWritable(int opaque) { - if (!inFlightRequestMap.containsKey(opaque)) { - return false; - } - - InvocationContext invocationContext = inFlightRequestMap.get(opaque); - if (null != invocationContext) { - CompletableFuture future = invocationContext.getResponse(); - return null != future && !future.isCancelled() && !future.isCompletedExceptionally() && !future.isDone(); - } - return false; + @Override + public boolean isWritable() { + return inFlightRequestMap.size() > 0; } public void registerInvocationContext(int opaque, InvocationContext context) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java index e680c7ad90..4a4b176cbc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java @@ -60,18 +60,18 @@ public class TransactionId { } public static TransactionId genByBrokerTransactionId(String brokerAddr, SendResult sendResult) { - MessageId id = new MessageId(null, 0); + long commitLogOffset = 0L; try { if (sendResult.getOffsetMsgId() != null) { - id = MessageDecoder.decodeMessageId(sendResult.getOffsetMsgId()); + commitLogOffset = generateCommitLogOffset(sendResult.getOffsetMsgId()); } else { - id = MessageDecoder.decodeMessageId(sendResult.getMsgId()); + commitLogOffset = generateCommitLogOffset(sendResult.getMsgId()); } } catch (Exception e) { log.warn("genFromBrokerTransactionId failed. brokerAddr: {}, sendResult: {}", brokerAddr, sendResult, e); } return genByBrokerTransactionId(RemotingUtil.string2SocketAddress(brokerAddr), sendResult.getTransactionId(), - id.getOffset(), sendResult.getQueueOffset()); + commitLogOffset, sendResult.getQueueOffset()); } public static TransactionId genByBrokerTransactionId(SocketAddress brokerAddr, String orgTransactionId, @@ -127,6 +127,11 @@ public class TransactionId { .build(); } + public static long generateCommitLogOffset(String messageId) throws UnknownHostException { + MessageId id = MessageDecoder.decodeMessageId(messageId); + return id.getOffset(); + } + @Override public boolean equals(Object o) { if (this == o) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java index 2bbf6b84db..6b7ae4ce20 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java @@ -31,8 +31,10 @@ import org.apache.rocketmq.common.protocol.header.CheckTransactionStateRequestHe import org.apache.rocketmq.common.protocol.header.GetConsumerRunningInfoRequestHeader; import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.channel.SimpleChannel; -import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.common.TelemetryCommandManager; +import org.apache.rocketmq.proxy.connector.transaction.TransactionId; +import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; +import org.apache.rocketmq.remoting.common.RemotingUtil; import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class GrpcClientChannel extends SimpleChannel { @@ -116,19 +118,21 @@ public class GrpcClientChannel extends SimpleChannel { try { switch (command.getCode()) { case RequestCode.CHECK_TRANSACTION_STATE: { - final CheckTransactionStateRequestHeader requestHeader = command.decodeCommandCustomHeader(CheckTransactionStateRequestHeader.class); + final CheckTransactionStateRequestHeader header = (CheckTransactionStateRequestHeader) command.readCustomHeader(); MessageExt messageExt = MessageDecoder.decode(ByteBuffer.wrap(command.getBody()), true, false, false); + TransactionId transactionId = TransactionId.genByBrokerTransactionId(RemotingUtil.string2SocketAddress(localAddress), + header.getTransactionId(), messageExt.getCommitLogOffset(), messageExt.getQueueOffset()); streamObserver.onNext(TelemetryCommand.newBuilder() .setRecoverOrphanedTransactionCommand(RecoverOrphanedTransactionCommand.newBuilder() - .setTransactionId(requestHeader.getTransactionId()) + .setTransactionId(transactionId.getProxyTransactionId()) .setOrphanedTransactionalMessage(GrpcConverter.buildMessage(messageExt)) .build()) .build()); break; } case RequestCode.GET_CONSUMER_RUNNING_INFO: { - final GetConsumerRunningInfoRequestHeader requestHeader = command.decodeCommandCustomHeader(GetConsumerRunningInfoRequestHeader.class); - if (!requestHeader.isJstackEnable()) { + final GetConsumerRunningInfoRequestHeader header = (GetConsumerRunningInfoRequestHeader) command.readCustomHeader(); + if (!header.isJstackEnable()) { break; } String nonce = manager.putCommand(command.getOpaque()); @@ -141,7 +145,6 @@ public class GrpcClientChannel extends SimpleChannel { } } } catch (Exception ignore) { - } } if (msg instanceof TelemetryCommand) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/ReceiveMessageResponseHandler.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/ReceiveMessageResponseHandler.java index 88066c25d1..6d076490f4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/ReceiveMessageResponseHandler.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/ReceiveMessageResponseHandler.java @@ -36,8 +36,8 @@ import org.apache.rocketmq.common.message.MessageDecoder; import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.header.ExtraInfoUtil; import org.apache.rocketmq.common.protocol.header.PopMessageResponseHeader; -import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.channel.InvocationContext; +import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.apache.rocketmq.remoting.protocol.RemotingSysResponseCode; @@ -46,9 +46,11 @@ import org.slf4j.LoggerFactory; public class ReceiveMessageResponseHandler implements ResponseHandler { private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME); + private final String brokerName; private final boolean fifo; - public ReceiveMessageResponseHandler(boolean fifo) { + public ReceiveMessageResponseHandler(String brokerName, boolean fifo) { + this.brokerName = brokerName; this.fifo = fifo; } @@ -58,7 +60,6 @@ public class ReceiveMessageResponseHandler implements ResponseHandler future = context.getResponse(); - String brokerName = request.getMessageQueue().getBroker().getName(); long currentTimeInMillis = System.currentTimeMillis(); long popCosts = currentTimeInMillis - context.getTimestamp(); try { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/SendMessageResponseHandler.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/SendMessageResponseHandler.java index 0716c404f5..6f842df674 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/SendMessageResponseHandler.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/SendMessageResponseHandler.java @@ -20,14 +20,26 @@ package org.apache.rocketmq.proxy.grpc.v2.adapter.handler; import apache.rocketmq.v2.SendMessageRequest; import apache.rocketmq.v2.SendMessageResponse; import apache.rocketmq.v2.SendReceipt; +import java.net.UnknownHostException; import org.apache.commons.lang3.StringUtils; +import org.apache.rocketmq.common.protocol.ResponseCode; import org.apache.rocketmq.common.protocol.header.SendMessageResponseHeader; +import org.apache.rocketmq.common.sysflag.MessageSysFlag; import org.apache.rocketmq.proxy.channel.InvocationContext; +import org.apache.rocketmq.proxy.connector.transaction.TransactionId; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; +import org.apache.rocketmq.remoting.common.RemotingUtil; import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class SendMessageResponseHandler implements ResponseHandler { - public SendMessageResponseHandler() { + private final String messageId; + private final int sysFlag; + private final String localAddress; + + public SendMessageResponseHandler(String messageId, int sysFlag, String localAddress) { + this.messageId = messageId; + this.sysFlag = sysFlag; + this.localAddress = localAddress; } @Override public void handle(RemotingCommand responseCommand, @@ -37,17 +49,24 @@ public class SendMessageResponseHandler implements ResponseHandler messageList = GrpcConverter.buildMessage(request.getMessagesList(), topicName); - MessageBatch messageBatch = MessageBatch.generateFromList(messageList); - MessageClientIDSetter.setUniqID(messageBatch); - messageBatch.setBody(messageBatch.encode()); - command.setBody(messageBatch.encode()); + String messageId; + if (messageList.size() == 1) { + org.apache.rocketmq.common.message.Message message = messageList.get(0); + command.setBody(message.getBody()); + messageId = MessageClientIDSetter.getUniqID(message); + } else { + MessageBatch messageBatch = MessageBatch.generateFromList(messageList); + MessageClientIDSetter.setUniqID(messageBatch); + messageBatch.setBody(messageBatch.encode()); + command.setBody(messageBatch.encode()); + messageId = MessageClientIDSetter.getUniqID(messageBatch); + } command.makeCustomHeaderToNet(); - SendMessageResponseHandler handler = new SendMessageResponseHandler(); + SendMessageResponseHandler handler = new SendMessageResponseHandler(messageId, requestHeader.getSysFlag(), brokerController.getBrokerAddr()); SendMessageChannel channel = channelManager.createChannel(() -> new SendMessageChannel(handler), SendMessageChannel.class); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); CompletableFuture future = new CompletableFuture<>(); @@ -249,7 +257,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo @Override public CompletableFuture receiveMessage(Context ctx, ReceiveMessageRequest request) { - long pollTime = GrpcConverter.buildPollTimeFromContext(ctx); + long pollTime = ctx.getDeadline().timeRemaining(TimeUnit.MILLISECONDS); String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID); ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId); PopMessageRequestHeader requestHeader = GrpcConverter.buildPopMessageRequestHeader(request, pollTime, @@ -257,7 +265,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.POP_MESSAGE, requestHeader); command.makeCustomHeaderToNet(); - ReceiveMessageResponseHandler handler = new ReceiveMessageResponseHandler(clientSettings.getSettings().getSubscription().getFifo()); + ReceiveMessageResponseHandler handler = new ReceiveMessageResponseHandler(brokerController.getBrokerConfig().getBrokerName(), + clientSettings.getSettings().getSubscription().getFifo()); ReceiveMessageChannel channel = channelManager.createChannel(() -> new ReceiveMessageChannel(handler), ReceiveMessageChannel.class); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); CompletableFuture future = new CompletableFuture<>(); @@ -304,22 +313,42 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo public CompletableFuture nackMessage(Context ctx, NackMessageRequest request) { Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); - - ChangeInvisibleTimeRequestHeader requestHeader = GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy); - RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CHANGE_MESSAGE_INVISIBLETIME, requestHeader); - command.makeCustomHeaderToNet(); - CompletableFuture future = new CompletableFuture<>(); - try { - RemotingCommand responseCommand = brokerController.getChangeInvisibleTimeProcessor() - .processRequest(channelHandlerContext, command); - NackMessageResponse response = NackMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(responseCommand.getCode(), responseCommand.getRemark())) - .build(); - future.complete(response); - } catch (Exception e) { - log.error("Exception raised while nackMessage", e); - future.completeExceptionally(e); + + ClientSettings clientSettings = grpcClientManager.getClientSettings(ctx); + int maxReconsumeTimes = clientSettings.getSettings().getSubscription().getDeadLetterPolicy().getMaxDeliveryAttempts(); + if (request.getDeliveryAttempt() >= maxReconsumeTimes) { + ConsumerSendMsgBackRequestHeader requestHeader = GrpcConverter.buildConsumerSendMsgBackToDLQRequestHeader(request, maxReconsumeTimes); + RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader); + command.makeCustomHeaderToNet(); + + try { + RemotingCommand responseCommand = brokerController.getSendMessageProcessor() + .processRequest(channelHandlerContext, command); + NackMessageResponse response = NackMessageResponse.newBuilder() + .setStatus(ResponseBuilder.buildStatus(responseCommand.getCode(), responseCommand.getRemark())) + .build(); + future.complete(response); + } catch (Exception e) { + log.error("Exception raised while nackMessage", e); + future.completeExceptionally(e); + } + } else { + ChangeInvisibleTimeRequestHeader requestHeader = GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy); + RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CHANGE_MESSAGE_INVISIBLETIME, requestHeader); + command.makeCustomHeaderToNet(); + + try { + RemotingCommand responseCommand = brokerController.getChangeInvisibleTimeProcessor() + .processRequest(channelHandlerContext, command); + NackMessageResponse response = NackMessageResponse.newBuilder() + .setStatus(ResponseBuilder.buildStatus(responseCommand.getCode(), responseCommand.getRemark())) + .build(); + future.complete(response); + } catch (Exception e) { + log.error("Exception raised while nackMessage", e); + future.completeExceptionally(e); + } } return future; } @@ -402,7 +431,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo @Override public CompletableFuture pullMessage(Context ctx, PullMessageRequest request) { - long pollTime = org.apache.rocketmq.proxy.grpc.v1.adapter.GrpcConverter.buildPollTimeFromContext(ctx); + long pollTime = ctx.getDeadline().timeRemaining(TimeUnit.MILLISECONDS); PullMessageRequestHeader requestHeader = GrpcConverter.buildPullMessageRequestHeader(request, pollTime); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.PULL_MESSAGE, requestHeader); command.makeCustomHeaderToNet(); 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 168d011795..680f8a1e2c 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 @@ -76,4 +76,24 @@ public class ClusterGrpcTest extends GrpcBaseTest { assertQueryAssignment(response, brokerNum); } + + @Test + public void testSendReceiveMessage() throws Exception { + super.testSendReceiveMessage(); + } + + @Test + public void testTransactionCheckThenCommit() { + super.testTransactionCheckThenCommit(); + } + + @Test + public void testSendReceiveMessageThenToDLQ() throws Exception { + super.testSendReceiveMessageThenToDLQ(); + } + + @Test + public void testPullMessage() throws Exception { + super.testPullMessage(); + } } 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 5bb2471cb7..0415d6dbd1 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 @@ -30,6 +30,7 @@ import apache.rocketmq.v2.DeadLetterPolicy; import apache.rocketmq.v2.EndTransactionRequest; import apache.rocketmq.v2.EndTransactionResponse; import apache.rocketmq.v2.Endpoints; +import apache.rocketmq.v2.HeartbeatRequest; import apache.rocketmq.v2.Message; import apache.rocketmq.v2.MessageQueue; import apache.rocketmq.v2.MessageType; @@ -60,7 +61,6 @@ import apache.rocketmq.v2.TransactionResolution; import apache.rocketmq.v2.TransactionSource; import com.google.protobuf.ByteString; import com.google.protobuf.Duration; -import com.google.protobuf.Timestamp; import com.google.protobuf.util.Timestamps; import io.grpc.Channel; import io.grpc.Metadata; @@ -125,6 +125,7 @@ public class GrpcBaseTest extends BaseConf { public void setUp() throws Exception { header.put(InterceptorConstants.CLIENT_ID, "client-id" + UUID.randomUUID()); + header.put(InterceptorConstants.LANGUAGE, "JAVA"); String mockProxyHome = "/mock/rmq/proxy/home"; URL mockProxyHomeURL = getClass().getClassLoader().getResource("rmq-proxy-home"); @@ -202,7 +203,6 @@ public class GrpcBaseTest extends BaseConf { .build()); } - @Test public void testSendReceiveMessage() throws Exception { String topic = initTopicOnSampleTopicBroker(broker1Name); String group = "group"; @@ -229,7 +229,6 @@ public class GrpcBaseTest extends BaseConf { assertAck(ackMessageResponse); } - @Test public void testSendReceiveMessageThenToDLQ() throws Exception { String topic = initTopicOnSampleTopicBroker(broker1Name); this.sendClientSettings(stub, ClientSettings.newBuilder() @@ -260,7 +259,7 @@ public class GrpcBaseTest extends BaseConf { AtomicReference receiveRetryResponseRef = new AtomicReference<>(); await().atMost(java.time.Duration.ofSeconds(30)).until(() -> { - ReceiveMessageResponse receiveRetryResponse = receiveMessage(blockingStub, topic, group); + ReceiveMessageResponse receiveRetryResponse = receiveMessage(blockingStub, topic, group, 1); if (receiveRetryResponse.getMessagesCount() <= 0) { return false; } @@ -307,17 +306,48 @@ public class GrpcBaseTest extends BaseConf { try { requestStreamObserver.onNext(TelemetryCommand.newBuilder() - .setClientSettings(buildProducerClientSettings(topic)) + .setClientSettings(buildPushConsumerClientSettings()) .build()); - + await().atMost(java.time.Duration.ofSeconds(3)).until(() -> { + if (telemetryCommandRef.get() == null) { + return false; + } + if (telemetryCommandRef.get().getCommandCase() != TelemetryCommand.CommandCase.CLIENT_OVERWRITTEN_SETTINGS) { + return false; + } + return telemetryCommandRef.get() != null; + }); + telemetryCommandRef.set(null); // init consumer offset receiveMessage(blockingStub, topic, group); + requestStreamObserver.onNext(TelemetryCommand.newBuilder() + .setClientSettings(buildProducerClientSettings(topic)) + .build()); + blockingStub.heartbeat(HeartbeatRequest.newBuilder() + .setGroup(Resource.newBuilder() + .setName(group) + .build()) + .build()); + await().atMost(java.time.Duration.ofSeconds(3)).until(() -> { + if (telemetryCommandRef.get() == null) { + return false; + } + if (telemetryCommandRef.get().getCommandCase() != TelemetryCommand.CommandCase.CLIENT_OVERWRITTEN_SETTINGS) { + return false; + } + return telemetryCommandRef.get() != null; + }); + telemetryCommandRef.set(null); + String messageId = createUniqID(); SendMessageResponse sendResponse = blockingStub.sendMessage(buildTransactionSendMessageRequest(topic, messageId)); assertSendMessage(sendResponse, messageId); await().atMost(java.time.Duration.ofSeconds(60)).until(() -> { + if (telemetryCommandRef.get() == null) { + return false; + } if (telemetryCommandRef.get().getCommandCase() != TelemetryCommand.CommandCase.RECOVER_ORPHANED_TRANSACTION_COMMAND) { return false; } @@ -347,7 +377,6 @@ public class GrpcBaseTest extends BaseConf { } } - @Test public void testPullMessage() throws Exception { String topic = initTopicOnSampleTopicBroker(broker1Name); String group = "group"; @@ -379,6 +408,11 @@ public class GrpcBaseTest extends BaseConf { .receiveMessage(buildReceiveMessageRequest(group, topic)); } + public ReceiveMessageResponse receiveMessage(MessagingServiceGrpc.MessagingServiceBlockingStub stub, String topic, String group, int timeSeconds) { + return stub.withDeadlineAfter(timeSeconds, TimeUnit.SECONDS) + .receiveMessage(buildReceiveMessageRequest(group, topic)); + } + public QueryRouteRequest buildQueryRouteRequest(String topic) { return QueryRouteRequest.newBuilder() .setTopic(Resource.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 6fefd4337c..cf5023755a 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 @@ -51,4 +51,24 @@ public class LocalGrpcTest extends GrpcBaseTest { QueryRouteResponse response = blockingStub.queryRoute(buildQueryRouteRequest(topic)); assertQueryRoute(response, brokerControllerList.size() * defaultQueueNums); } + + @Test + public void testSendReceiveMessage() throws Exception { + super.testSendReceiveMessage(); + } + + @Test + public void testTransactionCheckThenCommit() { + super.testTransactionCheckThenCommit(); + } + + @Test + public void testSendReceiveMessageThenToDLQ() throws Exception { + super.testSendReceiveMessageThenToDLQ(); + } + + @Test + public void testPullMessage() throws Exception { + super.testPullMessage(); + } }