mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] v2 support
This commit is contained in:
@@ -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<List<String>> 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<Long> getMaxOffset(String brokerAddr, String topic, int queueId) {
|
||||
return this.getMaxOffset(brokerAddr, topic, queueId, DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
public CompletableFuture<Long> getMaxOffset(Context ctx, String brokerAddr, String topic, int queueId) {
|
||||
return this.getMaxOffset(ctx, brokerAddr, topic, queueId, DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
|
||||
public CompletableFuture<Long> getMaxOffset(
|
||||
Context ctx,
|
||||
String brokerAddr,
|
||||
String topic,
|
||||
int queueId,
|
||||
@@ -78,15 +81,17 @@ public class DefaultForwardClient extends AbstractForwardClient {
|
||||
}
|
||||
|
||||
public CompletableFuture<Long> 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<Long> searchOffset(
|
||||
Context ctx,
|
||||
String brokerAddr,
|
||||
String topic,
|
||||
int queueId,
|
||||
|
||||
@@ -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<Integer> heartBeat(String brokerAddr, HeartbeatData heartbeatData) throws Exception {
|
||||
return this.heartBeat(brokerAddr, heartbeatData, DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
public CompletableFuture<Integer> heartBeat(Context ctx, String brokerAddr, HeartbeatData heartbeatData) throws Exception {
|
||||
return this.heartBeat(ctx, brokerAddr, heartbeatData, DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
public CompletableFuture<Integer> heartBeat(String brokerAddr, HeartbeatData heartbeatData, long timeout) throws Exception {
|
||||
public CompletableFuture<Integer> 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<SendResult> sendMessage(
|
||||
Context ctx,
|
||||
String address,
|
||||
String brokerName,
|
||||
List<Message> 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<SendResult> sendMessage(
|
||||
Context ctx,
|
||||
String address,
|
||||
String brokerName,
|
||||
List<Message> 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<SendResult> processSendMessageResponseFuture(
|
||||
protected CompletableFuture<SendResult> processSendMessageResponseFuture(
|
||||
Context ctx,
|
||||
String address,
|
||||
SendMessageRequestHeader requestHeader,
|
||||
CompletableFuture<SendResult> future) {
|
||||
@@ -108,14 +112,14 @@ public class ForwardProducer extends AbstractForwardClient {
|
||||
});
|
||||
}
|
||||
|
||||
public CompletableFuture<RemotingCommand> sendMessageBackThenAckOrg(String brokerAddr, ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader,
|
||||
public CompletableFuture<RemotingCommand> 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<RemotingCommand> sendMessageBackThenAckOrg(String brokerAddr, ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader,
|
||||
public CompletableFuture<RemotingCommand> 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<RemotingCommand> sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader) {
|
||||
return this.sendMessageBack(brokerAddr, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
public CompletableFuture<RemotingCommand> sendMessageBack(Context ctx, String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader) {
|
||||
return this.sendMessageBack(ctx, brokerAddr, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
|
||||
public CompletableFuture<RemotingCommand> sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) {
|
||||
public CompletableFuture<RemotingCommand> sendMessageBack(Context ctx, String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) {
|
||||
return this.getClient().sendMessageBackAsync(brokerAddr, requestHeader, timeoutMillis);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<PopResult> popMessage(String address, String brokerName,
|
||||
public CompletableFuture<PopResult> 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<PopResult> 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<PullResult> pullMessage(String address, PullMessageRequestHeader requestHeader) {
|
||||
return this.pullMessage(address, requestHeader, MAX_CONSUMER_TIMEOUT_MILLIS);
|
||||
public CompletableFuture<PullResult> pullMessage(Context ctx, String address, PullMessageRequestHeader requestHeader) {
|
||||
return this.pullMessage(ctx, address, requestHeader, MAX_CONSUMER_TIMEOUT_MILLIS);
|
||||
}
|
||||
|
||||
public CompletableFuture<PullResult> pullMessage(String address, PullMessageRequestHeader requestHeader,
|
||||
public CompletableFuture<PullResult> pullMessage(Context ctx, String address, PullMessageRequestHeader requestHeader,
|
||||
long timeoutMillis) {
|
||||
return this.getClient().pullMessageAsync(address, requestHeader, timeoutMillis);
|
||||
}
|
||||
|
||||
@@ -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<AckResult> ackMessage(String address, AckMessageRequestHeader requestHeader) {
|
||||
return this.ackMessage(address, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
public CompletableFuture<AckResult> ackMessage(Context ctx, String address, AckMessageRequestHeader requestHeader) {
|
||||
return this.ackMessage(ctx, address, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
|
||||
public CompletableFuture<AckResult> ackMessage(
|
||||
Context ctx,
|
||||
String address,
|
||||
AckMessageRequestHeader requestHeader,
|
||||
long timeoutMillis
|
||||
@@ -60,14 +62,16 @@ public class ForwardWriteConsumer extends AbstractForwardClient {
|
||||
}
|
||||
|
||||
public CompletableFuture<AckResult> 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<AckResult> 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
|
||||
|
||||
+3
-1
@@ -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<BrokerData> 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);
|
||||
}
|
||||
|
||||
+8
-5
@@ -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<AckResult> ackResultFuture = this.writeConsumer.ackMessage(brokerAddr, requestHeader);
|
||||
CompletableFuture<AckResult> 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);
|
||||
|
||||
+2
-1
@@ -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);
|
||||
|
||||
+1
-1
@@ -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());
|
||||
|
||||
+10
-10
@@ -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<String> 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<ReceiveMessageResponse> 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<ConsumerSendMsgBackRequestHeader> 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<ChangeInvisibleTimeRequestHeader> 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);
|
||||
|
||||
+2
-2
@@ -66,7 +66,7 @@ public class ProducerServiceTest extends BaseServiceTest {
|
||||
@Test
|
||||
public void testSendMessage() {
|
||||
CompletableFuture<SendResult> 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<SendResult> sendResultFuture = new CompletableFuture<>();
|
||||
when(producerClient.sendMessage(anyString(), anyString(), any(), any()))
|
||||
when(producerClient.sendMessage(any(), anyString(), anyString(), any(), any()))
|
||||
.thenReturn(sendResultFuture);
|
||||
sendResultFuture.completeExceptionally(ex);
|
||||
|
||||
|
||||
+3
-3
@@ -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())
|
||||
|
||||
Reference in New Issue
Block a user