mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Add integration test
* Write response when POLLING_TIMEOUT * Fix Channel isWritable * Use TransactionId in GrpcClientChannel and SendMessageResponseHandler * Optimize sendMessage and nackMessage
This commit is contained in:
@@ -454,6 +454,8 @@ public class PopMessageProcessor implements NettyRequestProcessor {
|
||||
response = null;
|
||||
}
|
||||
break;
|
||||
case ResponseCode.POLLING_TIMEOUT:
|
||||
return response;
|
||||
default:
|
||||
assert false;
|
||||
}
|
||||
|
||||
@@ -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<R, W> extends SimpleChannel implements C
|
||||
return super.writeAndFlush(msg);
|
||||
}
|
||||
|
||||
public boolean isWritable(int opaque) {
|
||||
if (!inFlightRequestMap.containsKey(opaque)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
InvocationContext<R, W> 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<R, W> context) {
|
||||
|
||||
+9
-4
@@ -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) {
|
||||
|
||||
+9
-6
@@ -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) {
|
||||
|
||||
+4
-3
@@ -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<ReceiveMessageRequest, ReceiveMessageResponse> {
|
||||
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<ReceiveMes
|
||||
ReceiveMessageRequest request = context.getRequest();
|
||||
CompletableFuture<ReceiveMessageResponse> future = context.getResponse();
|
||||
|
||||
String brokerName = request.getMessageQueue().getBroker().getName();
|
||||
long currentTimeInMillis = System.currentTimeMillis();
|
||||
long popCosts = currentTimeInMillis - context.getTimestamp();
|
||||
try {
|
||||
|
||||
+26
-7
@@ -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<SendMessageRequest, SendMessageResponse> {
|
||||
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<SendMessageRe
|
||||
// org.apache.rocketmq.broker.processor.AbstractSendMessageProcessor#doResponse
|
||||
if (null != responseCommand) {
|
||||
SendMessageResponseHeader responseHeader = (SendMessageResponseHeader) responseCommand.readCustomHeader();
|
||||
String messageId = "";
|
||||
String transactionId = "";
|
||||
if (responseHeader != null) {
|
||||
messageId = responseHeader.getMsgId();
|
||||
transactionId = responseHeader.getTransactionId();
|
||||
int tranType = MessageSysFlag.getTransactionValue(sysFlag);
|
||||
String transactionIdString = "";
|
||||
if (responseCommand.getCode() == ResponseCode.SUCCESS && tranType == MessageSysFlag.TRANSACTION_PREPARED_TYPE) {
|
||||
long commitLogOffset = 0L;
|
||||
try {
|
||||
commitLogOffset = TransactionId.generateCommitLogOffset(responseHeader.getMsgId());
|
||||
} catch (UnknownHostException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
TransactionId transactionId = TransactionId.genByBrokerTransactionId(RemotingUtil.string2SocketAddress(localAddress),
|
||||
responseHeader.getTransactionId(), commitLogOffset, responseHeader.getQueueOffset());
|
||||
transactionIdString = transactionId.getProxyTransactionId();
|
||||
}
|
||||
SendMessageResponse response = SendMessageResponse.newBuilder()
|
||||
.setStatus(ResponseBuilder.buildStatus(responseCommand.getCode(), responseCommand.getRemark()))
|
||||
.addReceipts(SendReceipt.newBuilder()
|
||||
.setMessageId(StringUtils.defaultString(messageId))
|
||||
.setTransactionId(StringUtils.defaultString(transactionId))
|
||||
.setTransactionId(StringUtils.defaultString(transactionIdString))
|
||||
.build())
|
||||
.build();
|
||||
context.getResponse().complete(response);
|
||||
|
||||
+58
-29
@@ -86,16 +86,17 @@ import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.UnregisterClientRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.channel.InvocationContext;
|
||||
import org.apache.rocketmq.proxy.channel.SimpleChannel;
|
||||
import org.apache.rocketmq.proxy.channel.SimpleChannelHandlerContext;
|
||||
import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.common.DelayPolicy;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.channel.InvocationContext;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandRecord;
|
||||
import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandRecord;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel;
|
||||
@@ -105,7 +106,6 @@ import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.SendMessageChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.PullMessageResponseHandler;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.ReceiveMessageResponseHandler;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.SendMessageResponseHandler;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.RouteService;
|
||||
import org.apache.rocketmq.remoting.RemotingServer;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRemotingAbstract;
|
||||
@@ -214,13 +214,21 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
SendMessageRequestHeader requestHeader = GrpcConverter.buildSendMessageRequestHeader(request, topicName);
|
||||
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.SEND_MESSAGE, requestHeader);
|
||||
List<org.apache.rocketmq.common.message.Message> 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<SendMessageResponse> future = new CompletableFuture<>();
|
||||
@@ -249,7 +257,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
|
||||
@Override
|
||||
public CompletableFuture<ReceiveMessageResponse> 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<ReceiveMessageResponse> future = new CompletableFuture<>();
|
||||
@@ -304,22 +313,42 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
public CompletableFuture<NackMessageResponse> 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<NackMessageResponse> 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<PullMessageResponse> 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();
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<ReceiveMessageResponse> 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()
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user