diff --git a/proxy/pom.xml b/proxy/pom.xml index 337e94c8c6..944de82261 100644 --- a/proxy/pom.xml +++ b/proxy/pom.xml @@ -40,10 +40,6 @@ org.apache.rocketmq rocketmq-proto - - org.apache.rocketmq - rocketmq-proxy - org.apache.rocketmq rocketmq-broker diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java b/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java index a3aff93c19..bb8805dccc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java @@ -17,7 +17,6 @@ package org.apache.rocketmq.proxy.configuration; -import java.util.concurrent.TimeUnit; import org.apache.rocketmq.proxy.grpc.common.ProxyMode; public class ProxyConfig { @@ -73,6 +72,8 @@ public class ProxyConfig { private int topicRouteThreadPoolNums = 36; private int topicRouteThreadPoolQueueCapacity = 50000; + private int longPollingReserveTimeInMillis = 10000; + public Integer getHealthCheckPort() { return healthCheckPort; } @@ -229,18 +230,6 @@ public class ProxyConfig { return forwardConsumerNum; } - public long getLongPollingReserveTimeMill() { - return longPollingReserveTimeMill; - } - - public void setLongPollingReserveTimeMill(long longPollingReserveTimeMill) { - this.longPollingReserveTimeMill = longPollingReserveTimeMill; - } - - public int getConsumerClientNum() { - return consumerClientNum; - } - public void setForwardConsumerNum(int forwardConsumerNum) { this.forwardConsumerNum = forwardConsumerNum; } @@ -332,4 +321,12 @@ public class ProxyConfig { public void setTopicRouteThreadPoolQueueCapacity(int topicRouteThreadPoolQueueCapacity) { this.topicRouteThreadPoolQueueCapacity = topicRouteThreadPoolQueueCapacity; } + + public int getLongPollingReserveTimeInMillis() { + return longPollingReserveTimeInMillis; + } + + public void setLongPollingReserveTimeInMillis(int longPollingReserveTimeInMillis) { + this.longPollingReserveTimeInMillis = longPollingReserveTimeInMillis; + } } 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 259f02d926..4c0a8f18ab 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 @@ -28,8 +28,11 @@ import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; public class ContextInterceptor implements ServerInterceptor { @Override - public ServerCall.Listener interceptCall(ServerCall call, Metadata headers, - ServerCallHandler next) { + public ServerCall.Listener interceptCall( + ServerCall call, + Metadata headers, + ServerCallHandler next + ) { Context context = Context.current() .withValue(InterceptorConstants.METADATA, headers); return Contexts.interceptCall(context, call, headers, next); 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 b806ec6e7a..63f0de555d 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 @@ -100,7 +100,8 @@ public class LocalGrpcService implements GrpcForwardService { return null; } - @Override public CompletableFuture heartbeat(Context ctx, HeartbeatRequest request) { + @Override + public CompletableFuture heartbeat(Context ctx, HeartbeatRequest request) { LanguageCode languageCode; String language = InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.LANGUAGE); languageCode = LanguageCode.valueOf(language); @@ -120,7 +121,8 @@ public class LocalGrpcService implements GrpcForwardService { return CompletableFuture.completedFuture(heartbeatResponse); } - @Override public CompletableFuture healthCheck(Context ctx, HealthCheckRequest request) { + @Override + public CompletableFuture healthCheck(Context ctx, HealthCheckRequest request) { LOGGER.trace("Received health check request from client: {}", request.getClientHost()); final HealthCheckResponse response = HealthCheckResponse.newBuilder() .setCommon(ResponseBuilder.buildCommon(Code.OK, "ok")) @@ -168,7 +170,7 @@ public class LocalGrpcService implements GrpcForwardService { long timeRemaining = Context.current() .getDeadline() .timeRemaining(TimeUnit.MILLISECONDS); - long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeMill(); + long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis(); if (pollTime <= 0) { pollTime = timeRemaining; } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/client/ClientManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/client/ForwardClientManagerTest.java similarity index 93% rename from proxy/src/test/java/org/apache/rocketmq/proxy/client/ClientManagerTest.java rename to proxy/src/test/java/org/apache/rocketmq/proxy/client/ForwardClientManagerTest.java index d3121f32e5..6c837555ed 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/client/ClientManagerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/client/ForwardClientManagerTest.java @@ -19,13 +19,13 @@ package org.apache.rocketmq.proxy.client; import org.apache.rocketmq.proxy.client.transaction.TransactionStateChecker; import org.apache.rocketmq.proxy.configuration.ConfigurationManager; -import org.apache.rocketmq.proxy.configuration.InitConfigurationTest; +import org.apache.rocketmq.proxy.configuration.InitConfigAndLoggerTest; import org.junit.Test; import org.mockito.Mockito; import static org.assertj.core.api.Assertions.assertThat; -public class ClientManagerTest extends InitConfigurationTest { +public class ForwardClientManagerTest extends InitConfigAndLoggerTest { @Test public void testClientManager() throws Exception { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/configuration/ConfigurationManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/configuration/ConfigurationManagerTest.java index 95a3a31aad..f0b37123d2 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/configuration/ConfigurationManagerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/configuration/ConfigurationManagerTest.java @@ -22,7 +22,7 @@ import org.junit.Test; import static org.assertj.core.api.Assertions.assertThat; -public class ConfigurationManagerTest extends InitConfigurationTest { +public class ConfigurationManagerTest extends InitConfigAndLoggerTest { @Test public void testInitEnv() { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/configuration/InitConfigurationTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/configuration/InitConfigAndLoggerTest.java similarity index 98% rename from proxy/src/test/java/org/apache/rocketmq/proxy/configuration/InitConfigurationTest.java rename to proxy/src/test/java/org/apache/rocketmq/proxy/configuration/InitConfigAndLoggerTest.java index 753abb923b..5cad87fac8 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/configuration/InitConfigurationTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/configuration/InitConfigAndLoggerTest.java @@ -28,7 +28,7 @@ import org.slf4j.LoggerFactory; import static org.apache.rocketmq.proxy.configuration.ConfigurationManager.RMQ_PROXY_HOME; -public class InitConfigurationTest { +public class InitConfigAndLoggerTest { public static String mockProxyHome = "/mock/rmq/proxy/home"; @Before 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 b78534a9fd..52c3e587fe 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 @@ -31,7 +31,6 @@ import io.grpc.Context; import io.grpc.Metadata; import io.netty.channel.ChannelHandlerContext; import java.net.InetSocketAddress; -import java.net.URL; import java.nio.charset.StandardCharsets; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executors; @@ -45,8 +44,7 @@ import org.apache.rocketmq.common.message.MessageDecoder; import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.ResponseCode; import org.apache.rocketmq.common.protocol.header.PopMessageResponseHeader; -import org.apache.rocketmq.proxy.configuration.ConfigurationManager; -import org.apache.rocketmq.proxy.configuration.InitConfigurationTest; +import org.apache.rocketmq.proxy.configuration.InitConfigAndLoggerTest; import org.apache.rocketmq.proxy.grpc.common.Converter; import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; import org.apache.rocketmq.remoting.exception.RemotingCommandException; @@ -58,11 +56,10 @@ import org.mockito.Mock; import org.mockito.Mockito; import org.mockito.junit.MockitoJUnitRunner; -import static org.apache.rocketmq.proxy.configuration.ConfigurationManager.RMQ_PROXY_HOME; import static org.assertj.core.api.Assertions.assertThat; @RunWith(MockitoJUnitRunner.class) -public class LocalGrpcServiceTest extends InitConfigurationTest { +public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { private LocalGrpcService localGrpcService; @Mock private SendMessageProcessor sendMessageProcessorMock; @@ -75,14 +72,7 @@ public class LocalGrpcServiceTest extends InitConfigurationTest { @Before public void setUp() throws Exception { - String mockProxyHome = "/mock/rmq/proxy/home"; - URL mockProxyHomeURL = getClass().getClassLoader().getResource("rmq-proxy-home"); - if (mockProxyHomeURL != null) { - mockProxyHome = mockProxyHomeURL.toURI().getPath(); - } - System.setProperty(RMQ_PROXY_HOME, mockProxyHome); - ConfigurationManager.initEnv(); - ConfigurationManager.intConfig(); + super.before(); Mockito.when(brokerControllerMock.getSendMessageProcessor()).thenReturn(sendMessageProcessorMock); Mockito.when(brokerControllerMock.getPopMessageProcessor()).thenReturn(popMessageProcessorMock); localGrpcService = new LocalGrpcService(brokerControllerMock);