diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java index 30567fdb74..3f217294f4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java @@ -39,10 +39,8 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor { } @Override - public RemotingCommand processRequest( - ChannelHandlerContext ctx, - RemotingCommand request - ) throws RemotingCommandException { + public RemotingCommand processRequest(ChannelHandlerContext ctx, RemotingCommand request) + throws RemotingCommandException { if (request.getCode() == RequestCode.CHECK_TRANSACTION_STATE) { return this.checkTransactionState(ctx, request); } @@ -50,10 +48,8 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor { } @Override - public RemotingCommand checkTransactionState( - ChannelHandlerContext ctx, - RemotingCommand request - ) throws RemotingCommandException { + public RemotingCommand checkTransactionState(ChannelHandlerContext ctx, RemotingCommand request) + throws RemotingCommandException { final CheckTransactionStateRequestHeader requestHeader = (CheckTransactionStateRequestHeader) request.decodeCommandCustomHeader(CheckTransactionStateRequestHeader.class); final ByteBuffer byteBuffer = ByteBuffer.wrap(request.getBody()); @@ -61,18 +57,20 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor { if (messageExt != null) { final String group = messageExt.getProperty(MessageConst.PROPERTY_PRODUCER_GROUP); if (group != null) { - transactionStateChecker.checkTransactionState(new TransactionStateCheckRequest( - group, - requestHeader.getTranStateTableOffset(), - requestHeader.getCommitLogOffset(), - requestHeader.getMsgId(), - TransactionId.genFromBrokerTransactionId( - ctx.channel().remoteAddress(), - requestHeader.getTransactionId(), + transactionStateChecker.checkTransactionState( + new TransactionStateCheckRequest( + group, + requestHeader.getTranStateTableOffset(), requestHeader.getCommitLogOffset(), - requestHeader.getTranStateTableOffset()), - messageExt - )); + requestHeader.getMsgId(), + TransactionId.genFromBrokerTransactionId( + ctx.channel().remoteAddress(), + requestHeader.getTransactionId(), + requestHeader.getCommitLogOffset(), + requestHeader.getTranStateTableOffset()), + messageExt + ) + ); } } return null; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java index 263542b9b4..f60fac135f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java @@ -259,17 +259,15 @@ public class Converter { return endTransactionRequestHeader; } - public static PullMessageRequestHeader buildPullMessageRequestHeader(PullMessageRequest request, long pollTime) { + public static PullMessageRequestHeader buildPullMessageRequestHeader(PullMessageRequest request, long pollTimeoutInMillis) { Partition partition = request.getPartition(); String groupName = Converter.getResourceNameWithNamespace(request.getGroup()); String topicName = Converter.getResourceNameWithNamespace(partition.getTopic()); int queueId = partition.getId(); int sysFlag = PullSysFlag.buildSysFlag(false, true, true, false, false); - String expression = request.getFilterExpression() - .getExpression(); - String expressionType = Converter.buildExpressionType(request.getFilterExpression() - .getType()); + String expression = request.getFilterExpression().getExpression(); + String expressionType = Converter.buildExpressionType(request.getFilterExpression().getType()); PullMessageRequestHeader requestHeader = new PullMessageRequestHeader(); requestHeader.setConsumerGroup(groupName); @@ -279,7 +277,7 @@ public class Converter { requestHeader.setMaxMsgNums(request.getBatchSize()); requestHeader.setSysFlag(sysFlag); requestHeader.setCommitOffset(0L); - requestHeader.setSuspendTimeoutMillis(pollTime); + requestHeader.setSuspendTimeoutMillis(pollTimeoutInMillis); requestHeader.setSubscription(expression); requestHeader.setSubVersion(0L); requestHeader.setExpressionType(expressionType); @@ -296,22 +294,26 @@ public class Converter { } } MessageAccessor.setProperties(messageWithHeader, Maps.newHashMap(userProperties)); + // set tag String tag = message.getSystemAttribute().getTag(); if (!"".equals(tag)) { messageWithHeader.setTags(tag); } + // set keys List keysList = message.getSystemAttribute().getKeysList(); if (keysList.size() > 0) { messageWithHeader.setKeys(keysList); } + // set message id String messageId = message.getSystemAttribute().getMessageId(); if ("".equals(messageId)) { throw new IllegalArgumentException("message id is empty"); } MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX, messageId); + // set transaction property MessageType messageType = message.getSystemAttribute().getMessageType(); if (messageType.equals(MessageType.TRANSACTION)) { @@ -328,8 +330,7 @@ public class Converter { case DELAY_LEVEL: int delayLevel = message.getSystemAttribute().getDelayLevel(); if (delayLevel > 0) { - MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_DELAY_TIME_LEVEL, - String.valueOf(delayLevel)); + MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_DELAY_TIME_LEVEL, String.valueOf(delayLevel)); } break; case DELIVERY_TIMESTAMP: diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java index 4366b80c0b..889da0e2b4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java @@ -68,9 +68,9 @@ import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager; import org.apache.rocketmq.proxy.grpc.common.ProxyMode; import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.service.cluster.ClientService; +import org.apache.rocketmq.proxy.grpc.service.cluster.ConsumerService; import org.apache.rocketmq.proxy.grpc.service.cluster.ProducerService; import org.apache.rocketmq.proxy.grpc.service.cluster.PullMessageService; -import org.apache.rocketmq.proxy.grpc.service.cluster.ConsumerService; import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService; import org.apache.rocketmq.proxy.grpc.service.cluster.TransactionService; import org.slf4j.Logger; @@ -115,9 +115,11 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc @Override public CompletableFuture heartbeat(Context ctx, HeartbeatRequest request) { this.clientService.heartbeat(ctx, request, channelManager); - return CompletableFuture.completedFuture(HeartbeatResponse.newBuilder() - .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) - .build()); + return CompletableFuture.completedFuture( + HeartbeatResponse.newBuilder() + .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) + .build() + ); } @Override @@ -194,13 +196,16 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc public CompletableFuture notifyClientTermination(Context ctx, NotifyClientTerminationRequest request) { this.clientService.unregister(ctx, request, channelManager); - return CompletableFuture.completedFuture(NotifyClientTerminationResponse.newBuilder() - .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) - .build()); + return CompletableFuture.completedFuture( + NotifyClientTerminationResponse.newBuilder() + .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) + .build() + ); } - @Override public CompletableFuture changeInvisibleDuration(Context ctx, - ChangeInvisibleDurationRequest request) { + @Override + public CompletableFuture changeInvisibleDuration(Context ctx, + ChangeInvisibleDurationRequest request) { return null; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java index 5e6060a36d..3484ddfa01 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java @@ -40,7 +40,7 @@ import org.apache.rocketmq.proxy.grpc.common.ResponseHook; public class TransactionService extends BaseService implements TransactionStateChecker { - private ChannelManager channelManager; + private final ChannelManager channelManager; private final ForwardProducer forwardProducer; private volatile ResponseHook checkTransactionStateHook = null; @@ -56,23 +56,23 @@ public class TransactionService extends BaseService implements TransactionStateC public void checkTransactionState(TransactionStateCheckRequest checkData) { try { List clientIdList = this.channelManager.getClientIdList(checkData.getGroupId()); + // if clientIdList's size is 0, here will throw: java.lang.IllegalArgumentException: bound must be positive String clientId = clientIdList.get(ThreadLocalRandom.current().nextInt(clientIdList.size())); GrpcClientChannel channel = GrpcClientChannel.getChannel(this.channelManager, checkData.getGroupId(), clientId); String transactionId = checkData.getTransactionId().getProxyTransactionId(); Message message = Converter.buildMessage(checkData.getMessageExt()); - - PollCommandResponse commandResponse = PollCommandResponse.newBuilder() + PollCommandResponse response = PollCommandResponse.newBuilder() .setRecoverOrphanedTransactionCommand( RecoverOrphanedTransactionCommand.newBuilder() .setOrphanedTransactionalMessage(message) .setTransactionId(transactionId) .build() ).build(); - channel.writeAndFlush(commandResponse); + channel.writeAndFlush(response); if (this.checkTransactionStateHook != null) { - this.checkTransactionStateHook.beforeResponse(checkData, commandResponse, null); + this.checkTransactionStateHook.beforeResponse(checkData, response, null); } } catch (Throwable t) { if (this.checkTransactionStateHook != null) { @@ -88,8 +88,9 @@ public class TransactionService extends BaseService implements TransactionStateC endTransactionHook.beforeResponse(request, response, throwable); } }); + try { - EndTransactionRequestHeader requestHeader = this.convertToEndTransactionRequestHeader(ctx, request); + EndTransactionRequestHeader requestHeader = this.toEndTransactionRequestHeader(ctx, request); this.forwardProducer.endTransaction(requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); future.complete(EndTransactionResponse.newBuilder() .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) @@ -100,8 +101,7 @@ public class TransactionService extends BaseService implements TransactionStateC return future; } - protected EndTransactionRequestHeader convertToEndTransactionRequestHeader(Context ctx, - EndTransactionRequest request) { + protected EndTransactionRequestHeader toEndTransactionRequestHeader(Context ctx, EndTransactionRequest request) { return Converter.buildEndTransactionRequestHeader(request); } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/common/utils/FilterUtilTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/common/utils/FilterUtilTest.java index 2586060019..92ad3a362a 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/common/utils/FilterUtilTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/common/utils/FilterUtilTest.java @@ -17,6 +17,7 @@ package org.apache.rocketmq.proxy.common.utils; +import java.util.concurrent.ThreadLocalRandom; import org.apache.rocketmq.common.filter.FilterAPI; import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; import org.junit.Test;