mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Fix codes for passing style checking.
This commit is contained in:
@@ -40,10 +40,6 @@
|
||||
<groupId>org.apache.rocketmq</groupId>
|
||||
<artifactId>rocketmq-proto</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.apache.rocketmq</groupId>
|
||||
<artifactId>rocketmq-proxy</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.apache.rocketmq</groupId>
|
||||
<artifactId>rocketmq-broker</artifactId>
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
+5
-2
@@ -28,8 +28,11 @@ import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants;
|
||||
public class ContextInterceptor implements ServerInterceptor {
|
||||
|
||||
@Override
|
||||
public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(ServerCall<ReqT, RespT> call, Metadata headers,
|
||||
ServerCallHandler<ReqT, RespT> next) {
|
||||
public <R, W> ServerCall.Listener<R> interceptCall(
|
||||
ServerCall<R, W> call,
|
||||
Metadata headers,
|
||||
ServerCallHandler<R, W> next
|
||||
) {
|
||||
Context context = Context.current()
|
||||
.withValue(InterceptorConstants.METADATA, headers);
|
||||
return Contexts.interceptCall(context, call, headers, next);
|
||||
|
||||
@@ -100,7 +100,8 @@ public class LocalGrpcService implements GrpcForwardService {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<HeartbeatResponse> heartbeat(Context ctx, HeartbeatRequest request) {
|
||||
@Override
|
||||
public CompletableFuture<HeartbeatResponse> 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<HealthCheckResponse> healthCheck(Context ctx, HealthCheckRequest request) {
|
||||
@Override
|
||||
public CompletableFuture<HealthCheckResponse> 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;
|
||||
}
|
||||
|
||||
+2
-2
@@ -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 {
|
||||
+1
-1
@@ -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() {
|
||||
|
||||
+1
-1
@@ -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
|
||||
+3
-13
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user