From f9c3a5e49ba0a3b9e6483ce18b71252bb686416e Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Fri, 18 Mar 2022 16:12:08 +0800 Subject: [PATCH] [ISSUE #3949] Improve readability. --- .../rocketmq/proxy/channel/ChannelManager.java | 2 +- .../utils/{FilterUtil.java => FilterUtils.java} | 2 +- .../proxy/connector/DefaultForwardClient.java | 16 +++++++++++++--- .../rocketmq/proxy/grpc/common/Converter.java | 3 ++- .../grpc/interceptor/ContextInterceptor.java | 1 - .../grpc/interceptor/HeaderInterceptor.java | 1 - .../InterceptorConstants.java | 2 +- .../proxy/grpc/service/LocalGrpcService.java | 2 +- .../grpc/service/cluster/ClientService.java | 2 +- .../grpc/service/cluster/PullMessageService.java | 4 ++-- .../proxy/common/utils/FilterUtilTest.java | 8 ++++---- .../proxy/grpc/service/LocalGrpcServiceTest.java | 2 +- 12 files changed, 27 insertions(+), 18 deletions(-) rename proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/{FilterUtil.java => FilterUtils.java} (98%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common => interceptor}/InterceptorConstants.java (97%) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java index 1d27b1633c..8fa53b3d7b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java @@ -32,7 +32,7 @@ import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.common.Cleaner; import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; -import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; +import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FilterUtil.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FilterUtils.java similarity index 98% rename from proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FilterUtil.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FilterUtils.java index a33c8c1762..2c9b663a9a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FilterUtil.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FilterUtils.java @@ -19,7 +19,7 @@ package org.apache.rocketmq.proxy.common.utils; import java.util.Set; import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; -public class FilterUtil { +public class FilterUtils { /** * Whether the message's tag matches consumerGroup's SubscriptionData * 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 88b9b4a9cd..b1e032c2f4 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 @@ -64,12 +64,22 @@ public class DefaultForwardClient extends AbstractForwardClient { return getClient().getTopicRouteInfoFromNameServer(topic, timeoutMillis); } - public CompletableFuture getMaxOffset(String brokerAddr, String topic, int queueId, long timeoutMillis) { + public CompletableFuture getMaxOffset( + String brokerAddr, + String topic, + int queueId, + long timeoutMillis + ) { return getClient().getMaxOffset(brokerAddr, topic, queueId, timeoutMillis); } - public CompletableFuture searchOffset(String brokerAddr, String topic, int queueId, long timestamp, - long timeoutMillis) { + public CompletableFuture searchOffset( + String brokerAddr, + String topic, + int queueId, + long timestamp, + long timeoutMillis + ) { return getClient().searchOffset(brokerAddr, topic, queueId, timestamp, timeoutMillis); } } 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 03f4877313..263542b9b4 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 @@ -193,7 +193,8 @@ public class Converter { changeInvisibleTimeRequestHeader.setQueueId(handle.getQueueId()); changeInvisibleTimeRequestHeader.setExtraInfo(handle.getReceiptHandle()); changeInvisibleTimeRequestHeader.setOffset(handle.getOffset()); - changeInvisibleTimeRequestHeader.setInvisibleTime(delayPolicy.getDelayInterval(ConfigurationManager.getProxyConfig().getRetryDelayLevelDelta() + request.getDeliveryAttempt())); + changeInvisibleTimeRequestHeader.setInvisibleTime( + delayPolicy.getDelayInterval(ConfigurationManager.getProxyConfig().getRetryDelayLevelDelta() + request.getDeliveryAttempt())); return changeInvisibleTimeRequestHeader; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/ContextInterceptor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/ContextInterceptor.java index d2a1bcd1b2..07d7ab9bf3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/ContextInterceptor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/ContextInterceptor.java @@ -23,7 +23,6 @@ import io.grpc.Metadata; import io.grpc.ServerCall; import io.grpc.ServerCallHandler; import io.grpc.ServerInterceptor; -import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; public class ContextInterceptor implements ServerInterceptor { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java index 698f75d00e..d106f8f0d2 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java @@ -25,7 +25,6 @@ import io.grpc.ServerCallHandler; import io.grpc.ServerInterceptor; import java.net.InetSocketAddress; import java.net.SocketAddress; -import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; public class HeaderInterceptor implements ServerInterceptor { @Override diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/InterceptorConstants.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/InterceptorConstants.java similarity index 97% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/InterceptorConstants.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/InterceptorConstants.java index cb175d2ef9..5a672f43ea 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/InterceptorConstants.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/InterceptorConstants.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.rocketmq.proxy.grpc.common; +package org.apache.rocketmq.proxy.grpc.interceptor; import io.grpc.Context; import io.grpc.Metadata; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java index 2072a21eae..0fa164afe4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java @@ -98,7 +98,7 @@ import org.apache.rocketmq.proxy.grpc.adapter.handler.ReceiveMessageResponseHand import org.apache.rocketmq.proxy.grpc.adapter.handler.SendMessageResponseHandler; import org.apache.rocketmq.proxy.grpc.common.Converter; import org.apache.rocketmq.proxy.grpc.common.DelayPolicy; -import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; +import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseFuture; import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager; import org.apache.rocketmq.proxy.grpc.common.ProxyMode; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java index 02c2003fb2..f2f20c517c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java @@ -32,7 +32,7 @@ import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; +import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager; import org.apache.rocketmq.remoting.protocol.LanguageCode; import org.slf4j.Logger; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java index f9f3a24986..6d7adab3da 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java @@ -34,7 +34,7 @@ import org.apache.rocketmq.client.consumer.PullResult; import org.apache.rocketmq.client.consumer.PullStatus; import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader; import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; -import org.apache.rocketmq.proxy.common.utils.FilterUtil; +import org.apache.rocketmq.proxy.common.utils.FilterUtils; import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; @@ -156,7 +156,7 @@ public class PullMessageService extends BaseService { PullStatus status = result.getPullStatus(); if (status.equals(PullStatus.FOUND)) { List messageList = result.getMsgFoundList().stream() - .filter(msg -> FilterUtil.isTagMatched(subscriptionData.getTagsSet(), msg.getTags())) // only return tag matched messages. + .filter(msg -> FilterUtils.isTagMatched(subscriptionData.getTagsSet(), msg.getTags())) // only return tag matched messages. .map(Converter::buildMessage) .collect(Collectors.toList()); 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 b6a5198a96..2586060019 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 @@ -27,25 +27,25 @@ public class FilterUtilTest { @Test public void testIsTagMatched() throws Exception { SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "tagA"); - assertThat(FilterUtil.isTagMatched(subscriptionData.getTagsSet(), "tagA")).isTrue(); + assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), "tagA")).isTrue(); } @Test public void testIsTagNotMatched() throws Exception { SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "tagA"); - assertThat(FilterUtil.isTagMatched(subscriptionData.getTagsSet(), "tagB")).isFalse(); + assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), "tagB")).isFalse(); } @Test public void testIsTagMatchedStar() throws Exception { SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "*"); - assertThat(FilterUtil.isTagMatched(subscriptionData.getTagsSet(), "tagA")).isTrue(); + assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), "tagA")).isTrue(); } @Test public void testIsTagNotMatchedNull() throws Exception { SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "tagA"); - assertThat(FilterUtil.isTagMatched(subscriptionData.getTagsSet(), null)).isFalse(); + assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), null)).isFalse(); } } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java index 093b0a47a6..6731ed8459 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java @@ -78,7 +78,7 @@ import org.apache.rocketmq.common.protocol.header.PullMessageResponseHeader; import org.apache.rocketmq.proxy.config.InitConfigAndLoggerTest; import org.apache.rocketmq.proxy.connector.transaction.TransactionId; import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; +import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.remoting.exception.RemotingCommandException; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.apache.rocketmq.store.MessageStore;