From c2fd34872b4728d2ac1be6ede80704cfcaf5ae3c Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Mon, 25 Apr 2022 14:32:33 +0800 Subject: [PATCH] [ISSUE #3949] v2 support --- .../proxy/connector/DefaultForwardClient.java | 11 ++++-- .../proxy/connector/ForwardProducer.java | 36 ++++++++++--------- .../proxy/connector/ForwardReadConsumer.java | 12 ++++--- .../proxy/connector/ForwardWriteConsumer.java | 11 ++++-- .../TransactionHeartbeatRegisterService.java | 4 ++- .../v2/service/cluster/ConsumerService.java | 13 ++++--- .../v2/service/cluster/ProducerService.java | 3 +- .../service/cluster/TransactionService.java | 2 +- .../service/cluster/ConsumerServiceTest.java | 20 +++++------ .../service/cluster/ProducerServiceTest.java | 4 +-- .../cluster/TransactionServiceTest.java | 6 ++-- 11 files changed, 72 insertions(+), 50 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java index 3cf8bd9334..ee7b1e45c7 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java @@ -16,6 +16,7 @@ */ package org.apache.rocketmq.proxy.connector; +import io.grpc.Context; import java.util.List; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.exception.MQClientException; @@ -47,6 +48,7 @@ public class DefaultForwardClient extends AbstractForwardClient { } public CompletableFuture> getConsumerListByGroup( + Context ctx, String brokerAddr, GetConsumerListByGroupRequestHeader requestHeader, long timeoutMillis @@ -64,11 +66,12 @@ public class DefaultForwardClient extends AbstractForwardClient { return this.getClient().getTopicRouteInfoFromNameServer(topic, timeoutMillis); } - public CompletableFuture getMaxOffset(String brokerAddr, String topic, int queueId) { - return this.getMaxOffset(brokerAddr, topic, queueId, DEFAULT_MQ_CLIENT_TIMEOUT); + public CompletableFuture getMaxOffset(Context ctx, String brokerAddr, String topic, int queueId) { + return this.getMaxOffset(ctx, brokerAddr, topic, queueId, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture getMaxOffset( + Context ctx, String brokerAddr, String topic, int queueId, @@ -78,15 +81,17 @@ public class DefaultForwardClient extends AbstractForwardClient { } public CompletableFuture searchOffset( + Context ctx, String brokerAddr, String topic, int queueId, long timestamp ) { - return this.searchOffset(brokerAddr, topic, queueId, timestamp, DEFAULT_MQ_CLIENT_TIMEOUT); + return this.searchOffset(ctx, brokerAddr, topic, queueId, timestamp, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture searchOffset( + Context ctx, String brokerAddr, String topic, int queueId, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java index 355449ba7c..10202aa5b9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java @@ -16,6 +16,7 @@ */ package org.apache.rocketmq.proxy.connector; +import io.grpc.Context; import java.util.List; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.common.protocol.ResponseCode; @@ -54,31 +55,33 @@ public class ForwardProducer extends AbstractForwardClient { return clientFactory.getTransactionalProducer(name, threadCount); } - public CompletableFuture heartBeat(String brokerAddr, HeartbeatData heartbeatData) throws Exception { - return this.heartBeat(brokerAddr, heartbeatData, DEFAULT_MQ_CLIENT_TIMEOUT); + public CompletableFuture heartBeat(Context ctx, String brokerAddr, HeartbeatData heartbeatData) throws Exception { + return this.heartBeat(ctx, brokerAddr, heartbeatData, DEFAULT_MQ_CLIENT_TIMEOUT); } - public CompletableFuture heartBeat(String brokerAddr, HeartbeatData heartbeatData, long timeout) throws Exception { + public CompletableFuture heartBeat(Context ctx, String brokerAddr, HeartbeatData heartbeatData, long timeout) throws Exception { return this.getClient().sendHeartbeatAsync(brokerAddr, heartbeatData, timeout); } - public void endTransaction(String brokerAddr, EndTransactionRequestHeader requestHeader) throws Exception { - this.endTransaction(brokerAddr, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); + public void endTransaction(Context ctx, String brokerAddr, EndTransactionRequestHeader requestHeader) throws Exception { + this.endTransaction(ctx, brokerAddr, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } - public void endTransaction(String brokerAddr, EndTransactionRequestHeader requestHeader, long timeoutMillis) throws Exception { + public void endTransaction(Context ctx, String brokerAddr, EndTransactionRequestHeader requestHeader, long timeoutMillis) throws Exception { this.getClient().endTransactionOneway(brokerAddr, requestHeader, "end transaction from rmq proxy", timeoutMillis); } public CompletableFuture sendMessage( + Context ctx, String address, String brokerName, List msg, SendMessageRequestHeader requestHeader ) { - return this.sendMessage(address, brokerName, msg, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); + return this.sendMessage(ctx, address, brokerName, msg, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture sendMessage( + Context ctx, String address, String brokerName, List msg, @@ -91,10 +94,11 @@ public class ForwardProducer extends AbstractForwardClient { } else { future = this.getClient().sendMessageAsync(address, brokerName, msg, requestHeader, timeoutMillis); } - return processSendMessageResponseFuture(address, requestHeader, future); + return processSendMessageResponseFuture(ctx, address, requestHeader, future); } - private CompletableFuture processSendMessageResponseFuture( + protected CompletableFuture processSendMessageResponseFuture( + Context ctx, String address, SendMessageRequestHeader requestHeader, CompletableFuture future) { @@ -108,14 +112,14 @@ public class ForwardProducer extends AbstractForwardClient { }); } - public CompletableFuture sendMessageBackThenAckOrg(String brokerAddr, ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader, + public CompletableFuture sendMessageBackThenAckOrg(Context ctx, String brokerAddr, ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader, AckMessageRequestHeader ackMessageRequestHeader) { - return sendMessageBackThenAckOrg(brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader,DEFAULT_MQ_CLIENT_TIMEOUT); + return sendMessageBackThenAckOrg(ctx, brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader,DEFAULT_MQ_CLIENT_TIMEOUT); } - public CompletableFuture sendMessageBackThenAckOrg(String brokerAddr, ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader, + public CompletableFuture sendMessageBackThenAckOrg(Context ctx, String brokerAddr, ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader, AckMessageRequestHeader ackMessageRequestHeader, long timeoutMillis) { - return this.sendMessageBack(brokerAddr, sendMsgBackRequestHeader, timeoutMillis).whenComplete((result, throwable) -> { + return this.sendMessageBack(ctx, brokerAddr, sendMsgBackRequestHeader, timeoutMillis).whenComplete((result, throwable) -> { if (throwable != null || ResponseCode.SUCCESS != result.getCode()) { return; } @@ -123,11 +127,11 @@ public class ForwardProducer extends AbstractForwardClient { }); } - public CompletableFuture sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader) { - return this.sendMessageBack(brokerAddr, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); + public CompletableFuture sendMessageBack(Context ctx, String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader) { + return this.sendMessageBack(ctx, brokerAddr, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } - public CompletableFuture sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) { + public CompletableFuture sendMessageBack(Context ctx, String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) { return this.getClient().sendMessageBackAsync(brokerAddr, requestHeader, timeoutMillis); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java index e45db1cbf3..13973e8882 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java @@ -16,6 +16,7 @@ */ package org.apache.rocketmq.proxy.connector; +import io.grpc.Context; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.consumer.PopResult; import org.apache.rocketmq.client.consumer.PullResult; @@ -46,12 +47,13 @@ public class ForwardReadConsumer extends AbstractForwardClient { return clientFactory.getMQClient(name, threadCount); } - public CompletableFuture popMessage(String address, String brokerName, + public CompletableFuture popMessage(Context ctx, String address, String brokerName, PopMessageRequestHeader requestHeader) { - return this.popMessage(address, brokerName, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); + return this.popMessage(ctx, address, brokerName, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture popMessage( + Context ctx, String address, String brokerName, PopMessageRequestHeader requestHeader, @@ -60,11 +62,11 @@ public class ForwardReadConsumer extends AbstractForwardClient { return this.getClient().popMessageAsync(address, brokerName, requestHeader, timeoutMillis); } - public CompletableFuture pullMessage(String address, PullMessageRequestHeader requestHeader) { - return this.pullMessage(address, requestHeader, MAX_CONSUMER_TIMEOUT_MILLIS); + public CompletableFuture pullMessage(Context ctx, String address, PullMessageRequestHeader requestHeader) { + return this.pullMessage(ctx, address, requestHeader, MAX_CONSUMER_TIMEOUT_MILLIS); } - public CompletableFuture pullMessage(String address, PullMessageRequestHeader requestHeader, + public CompletableFuture pullMessage(Context ctx, String address, PullMessageRequestHeader requestHeader, long timeoutMillis) { return this.getClient().pullMessageAsync(address, requestHeader, timeoutMillis); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java index 60332c063b..6be934cbb6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java @@ -16,6 +16,7 @@ */ package org.apache.rocketmq.proxy.connector; +import io.grpc.Context; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.consumer.AckResult; import org.apache.rocketmq.proxy.connector.client.MQClientAPIExt; @@ -47,11 +48,12 @@ public class ForwardWriteConsumer extends AbstractForwardClient { return clientFactory.getMQClient(name, threadCount); } - public CompletableFuture ackMessage(String address, AckMessageRequestHeader requestHeader) { - return this.ackMessage(address, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); + public CompletableFuture ackMessage(Context ctx, String address, AckMessageRequestHeader requestHeader) { + return this.ackMessage(ctx, address, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture ackMessage( + Context ctx, String address, AckMessageRequestHeader requestHeader, long timeoutMillis @@ -60,14 +62,16 @@ public class ForwardWriteConsumer extends AbstractForwardClient { } public CompletableFuture changeInvisibleTimeAsync( + Context ctx, String address, String brokerName, ChangeInvisibleTimeRequestHeader requestHeader ) { - return this.changeInvisibleTimeAsync(address, brokerName, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); + return this.changeInvisibleTimeAsync(ctx, address, brokerName, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture changeInvisibleTimeAsync( + Context ctx, String address, String brokerName, ChangeInvisibleTimeRequestHeader requestHeader, @@ -77,6 +81,7 @@ public class ForwardWriteConsumer extends AbstractForwardClient { } public void updateConsumerOffsetOneWay( + Context ctx, String brokerAddr, UpdateConsumerOffsetRequestHeader header, long timeoutMillis diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionHeartbeatRegisterService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionHeartbeatRegisterService.java index 830bd29b69..92c7d63964 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionHeartbeatRegisterService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionHeartbeatRegisterService.java @@ -17,6 +17,7 @@ package org.apache.rocketmq.proxy.connector.transaction; import com.google.common.collect.Sets; +import io.grpc.Context; import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; @@ -174,6 +175,7 @@ public class TransactionHeartbeatRegisterService implements StartAndShutdown { protected void sendHeartBeatToCluster(String clusterName, HeartbeatData heartbeatData) { try { + Context ctx = Context.current(); MessageQueueWrapper messageQueue = this.topicRouteCache.getMessageQueue(clusterName); List brokerDataList = messageQueue.getTopicRouteData().getBrokerDatas(); if (brokerDataList == null) { @@ -183,7 +185,7 @@ public class TransactionHeartbeatRegisterService implements StartAndShutdown { heartbeatExecutors.submit(() -> { String brokerAddr = brokerData.selectBrokerAddr(); try { - this.forwardProducer.heartBeat(brokerAddr, heartbeatData); + this.forwardProducer.heartBeat(ctx, brokerAddr, heartbeatData); } catch (Exception e) { log.error("Send transactionHeartbeat to broker err. brokerAddr: {}", brokerAddr, e); } 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 3c9ac57364..dcbbf600c4 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 @@ -102,6 +102,7 @@ public class ConsumerService extends BaseService { } future = this.readConsumer.popMessage( + ctx, messageQueue.getBrokerAddr(), messageQueue.getBrokerName(), requestHeader, @@ -210,7 +211,7 @@ public class ConsumerService extends BaseService { group, handle); - future = this.producer.sendMessageBackThenAckOrg(brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader); + future = this.producer.sendMessageBackThenAckOrg(ctx, brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader); } catch (Throwable t) { future.completeExceptionally(t); } @@ -238,7 +239,7 @@ public class ConsumerService extends BaseService { ackMessageRequestHeader.setExtraInfo(handle.getReceiptHandle()); ackMessageRequestHeader.setOffset(handle.getOffset()); - future = this.writeConsumer.ackMessage(brokerAddr, ackMessageRequestHeader); + future = this.writeConsumer.ackMessage(ctx, brokerAddr, ackMessageRequestHeader); } catch (Throwable t) { future.completeExceptionally(t); } @@ -297,7 +298,7 @@ public class ConsumerService extends BaseService { String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); AckMessageRequestHeader requestHeader = this.buildAckMessageRequestHeader(ctx, request, receiptHandle); - CompletableFuture ackResultFuture = this.writeConsumer.ackMessage(brokerAddr, requestHeader); + CompletableFuture ackResultFuture = this.writeConsumer.ackMessage(ctx, brokerAddr, requestHeader); ackResultFuture .thenAccept(result -> future.complete(convertToAckMessageResultEntry(ctx, ackMessageEntry, result))) .exceptionally(throwable -> { @@ -341,11 +342,13 @@ public class ConsumerService extends BaseService { int maxDeliveryAttempts = settings.getSubscription().getBackoffPolicy().getMaxAttempts(); if (request.getDeliveryAttempt() >= maxDeliveryAttempts) { future = this.producer.sendMessageBack( + ctx, brokerAddr, this.buildConsumerSendMsgBackToDLQRequestHeader(ctx, request, maxDeliveryAttempts) ).thenApply(result -> { if (result.getCode() == ResponseCode.SUCCESS) { writeConsumer.ackMessage( + ctx, brokerAddr, this.buildAckMessageRequestHeader(ctx, request)); } @@ -353,7 +356,7 @@ public class ConsumerService extends BaseService { }); } else { ChangeInvisibleTimeRequestHeader requestHeader = this.buildChangeInvisibleTimeRequestHeader(ctx, request); - future = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader) + future = this.writeConsumer.changeInvisibleTimeAsync(ctx, brokerAddr, receiptHandle.getBrokerName(), requestHeader) .thenApply(result -> convertToNackMessageResponse(ctx, request, result)); } } catch (Throwable t) { @@ -411,7 +414,7 @@ public class ConsumerService extends BaseService { String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); ChangeInvisibleTimeRequestHeader requestHeader = convertToChangeInvisibleTimeRequestHeader(ctx, request); - future = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader) + future = this.writeConsumer.changeInvisibleTimeAsync(ctx, brokerAddr, receiptHandle.getBrokerName(), requestHeader) .thenApply(result -> convertToChangeInvisibleDurationResponse(ctx, request, result)); } catch (Throwable t) { future.completeExceptionally(t); 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 d3f0e1c812..38fb8adec8 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 @@ -67,6 +67,7 @@ public class ProducerService extends BaseService { // send message to broker. future = this.producer.sendMessage( + ctx, selectableMessageQueue.getBrokerAddr(), selectableMessageQueue.getBrokerName(), convertToMessageList(ctx, request), @@ -128,7 +129,7 @@ public class ProducerService extends BaseService { AckMessageRequestHeader ackMessageRequestHeader = GrpcConverter.buildAckMessageRequestHeader( request.getTopic(), request.getGroup(), receiptHandle); - future = this.producer.sendMessageBackThenAckOrg(brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader) + future = this.producer.sendMessageBackThenAckOrg(ctx, brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader) .thenApply(result -> convertToForwardMessageToDeadLetterQueueResponse(ctx, result)); } catch (Throwable t) { future.completeExceptionally(t); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java index eea4722e0a..eb71fa848e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java @@ -99,7 +99,7 @@ public class TransactionService extends BaseService implements TransactionStateC TransactionId handle = TransactionId.decode(request.getTransactionId()); String brokerAddr = RemotingHelper.parseSocketAddressAddr(handle.getBrokerAddr()); EndTransactionRequestHeader requestHeader = this.toEndTransactionRequestHeader(ctx, request); - this.forwardProducer.endTransaction(brokerAddr, requestHeader); + this.forwardProducer.endTransaction(ctx, brokerAddr, requestHeader); future.complete(EndTransactionResponse.newBuilder() .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) .build()); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java index 5d993ee4f0..ee57c45e7e 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java @@ -76,10 +76,10 @@ public class ConsumerServiceTest extends BaseServiceTest { createMessageExt("msg2", "msg2") ); PopResult popResult = new PopResult(PopStatus.FOUND, messageExtList); - when(readConsumerClient.popMessage(anyString(), anyString(), any(), anyLong())) + when(readConsumerClient.popMessage(any(), anyString(), anyString(), any(), anyLong())) .thenReturn(CompletableFuture.completedFuture(popResult)); when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr"); - when(writeConsumerClient.ackMessage(anyString(), any())) + when(writeConsumerClient.ackMessage(any(), anyString(), any())) .thenReturn(CompletableFuture.completedFuture(new AckResult())); Context ctx = Context.current().withDeadlineAfter(3, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor()); @@ -127,15 +127,15 @@ public class ConsumerServiceTest extends BaseServiceTest { createMessageExt("msg2", "msg2") ); PopResult popResult = new PopResult(PopStatus.FOUND, messageExtList); - when(readConsumerClient.popMessage(anyString(), anyString(), any(), anyLong())) + when(readConsumerClient.popMessage(any(), anyString(), anyString(), any(), anyLong())) .thenReturn(CompletableFuture.completedFuture(popResult)); when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr"); List toDLQMsgId = new ArrayList<>(); doAnswer(mock -> { - ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader = mock.getArgument(1); + ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader = mock.getArgument(2); toDLQMsgId.add(sendMsgBackRequestHeader.getOriginMsgId()); return CompletableFuture.completedFuture(RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "")); - }).when(producerClient).sendMessageBackThenAckOrg(anyString(), any(), any()); + }).when(producerClient).sendMessageBackThenAckOrg(any(), anyString(), any(), any()); Context ctx = Context.current().withDeadlineAfter(3, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor()); List responseList = consumerService.receiveMessage(ctx, @@ -164,7 +164,7 @@ public class ConsumerServiceTest extends BaseServiceTest { when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr"); AckResult ackResult = new AckResult(); ackResult.setStatus(AckStatus.OK); - when(writeConsumerClient.ackMessage(anyString(), any())).thenReturn(CompletableFuture.completedFuture(ackResult)); + when(writeConsumerClient.ackMessage(any(), anyString(), any())).thenReturn(CompletableFuture.completedFuture(ackResult)); AckMessageResponse response = consumerService.ackMessage(Context.current(), AckMessageRequest.newBuilder() .setTopic(Resource.newBuilder() @@ -187,9 +187,9 @@ public class ConsumerServiceTest extends BaseServiceTest { ReceiptHandle receiptHandle = createReceiptHandle(); AtomicReference headerRef = new AtomicReference<>(); doAnswer(mock -> { - headerRef.set(mock.getArgument(1)); + headerRef.set(mock.getArgument(2)); return CompletableFuture.completedFuture(RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "")); - }).when(producerClient).sendMessageBack(anyString(), any()); + }).when(producerClient).sendMessageBack(any(), anyString(), any()); when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr"); Settings clientSettings = createClientSettings(3); @@ -216,11 +216,11 @@ public class ConsumerServiceTest extends BaseServiceTest { ReceiptHandle receiptHandle = createReceiptHandle(); AtomicReference headerRef = new AtomicReference<>(); doAnswer(mock -> { - headerRef.set(mock.getArgument(2)); + headerRef.set(mock.getArgument(3)); AckResult ackResult = new AckResult(); ackResult.setStatus(AckStatus.OK); return CompletableFuture.completedFuture(ackResult); - }).when(writeConsumerClient).changeInvisibleTimeAsync(anyString(), anyString(), any()); + }).when(writeConsumerClient).changeInvisibleTimeAsync(any(), anyString(), anyString(), any()); when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr"); Settings clientSettings = createClientSettings(3); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerServiceTest.java index 566f0c763b..e350682276 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerServiceTest.java @@ -66,7 +66,7 @@ public class ProducerServiceTest extends BaseServiceTest { @Test public void testSendMessage() { CompletableFuture sendResultFuture = new CompletableFuture<>(); - when(producerClient.sendMessage(anyString(), anyString(), any(), any())) + when(producerClient.sendMessage(any(), anyString(), anyString(), any(), any())) .thenReturn(sendResultFuture); sendResultFuture.complete(new SendResult(SendStatus.SEND_OK, "msgId", new MessageQueue(), 1L, "txId", "offsetMsgId", "regionId")); @@ -121,7 +121,7 @@ public class ProducerServiceTest extends BaseServiceTest { RuntimeException ex = new RuntimeException(); CompletableFuture sendResultFuture = new CompletableFuture<>(); - when(producerClient.sendMessage(anyString(), anyString(), any(), any())) + when(producerClient.sendMessage(any(), anyString(), anyString(), any(), any())) .thenReturn(sendResultFuture); sendResultFuture.completeExceptionally(ex); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionServiceTest.java index 5909ffac16..e9344d62c7 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionServiceTest.java @@ -72,10 +72,10 @@ public class TransactionServiceTest extends BaseServiceTest { RemotingHelper.string2SocketAddress("127.0.0.1:8080"), "71F99B78B6E261357FA259CCA6456118", 1234, 5678); doAnswer(mock -> { - brokerAddrRef.set(mock.getArgument(0)); - headerRef.set(mock.getArgument(1)); + brokerAddrRef.set(mock.getArgument(1)); + headerRef.set(mock.getArgument(2)); return null; - }).when(producerClient).endTransaction(anyString(), any()); + }).when(producerClient).endTransaction(any(), anyString(), any()); EndTransactionResponse response = transactionService.endTransaction(Context.current(), EndTransactionRequest.newBuilder() .setTransactionId(transactionId.getProxyTransactionId())