mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] For passing check style.
This commit is contained in:
@@ -33,7 +33,8 @@ public class HealthCheckServer implements StartAndShutdown {
|
||||
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
this.healthChecker = HttpServer.create(new InetSocketAddress(ConfigurationManager.getProxyConfig().getHealthCheckPort()), 0);
|
||||
this.healthChecker = HttpServer.create(
|
||||
new InetSocketAddress(ConfigurationManager.getProxyConfig().getHealthCheckPort()), 0);
|
||||
this.healthChecker.createContext("/status", new HealthCheckHandler());
|
||||
this.healthChecker.setExecutor(null);
|
||||
this.healthChecker.start();
|
||||
|
||||
@@ -39,11 +39,12 @@ import org.slf4j.LoggerFactory;
|
||||
|
||||
public class ProxyStartup {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(ProxyStartup.class);
|
||||
private static final ProxyStartAndShutdown proxyStartAndShutdown = new ProxyStartAndShutdown();
|
||||
private static final Logger LOGGER = LoggerFactory.getLogger(ProxyStartup.class);
|
||||
private static final ProxyStartAndShutdown PROXY_START_AND_SHUTDOWN = new ProxyStartAndShutdown();
|
||||
|
||||
private static class ProxyStartAndShutdown extends AbstractStartAndShutdown {
|
||||
@Override public void appendStartAndShutdown(StartAndShutdown startAndShutdown) {
|
||||
@Override
|
||||
public void appendStartAndShutdown(StartAndShutdown startAndShutdown) {
|
||||
super.appendStartAndShutdown(startAndShutdown);
|
||||
}
|
||||
}
|
||||
@@ -59,30 +60,29 @@ public class ProxyStartup {
|
||||
|
||||
// create and start grpcServer
|
||||
GrpcServer grpcServer = createGrpcServer();
|
||||
proxyStartAndShutdown.appendStartAndShutdown(grpcServer);
|
||||
PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(grpcServer);
|
||||
|
||||
// health check server
|
||||
final HealthCheckServer healthCheckServer = new HealthCheckServer();
|
||||
proxyStartAndShutdown.appendStartAndShutdown(healthCheckServer);
|
||||
PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(healthCheckServer);
|
||||
|
||||
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
|
||||
log.info("try to shutdown server");
|
||||
|
||||
LOGGER.info("try to shutdown server");
|
||||
try {
|
||||
proxyStartAndShutdown.shutdown();
|
||||
PROXY_START_AND_SHUTDOWN.shutdown();
|
||||
} catch (Exception e) {
|
||||
log.error("err when shutdown proxy", e);
|
||||
LOGGER.error("err when shutdown proxy", e);
|
||||
}
|
||||
}));
|
||||
} catch (Exception e) {
|
||||
System.err.println("find a unexpect err." + e);
|
||||
e.printStackTrace();
|
||||
log.error("find a unexpect err.", e);
|
||||
LOGGER.error("find a unexpect err.", e);
|
||||
System.exit(1);
|
||||
}
|
||||
|
||||
System.out.printf("%s%n", new Date() + " rmq-proxy startup successfully");
|
||||
log.info(new Date() + "rmq-proxy startup successfully");
|
||||
LOGGER.info(new Date() + "rmq-proxy startup successfully");
|
||||
}
|
||||
|
||||
private static GrpcServer createGrpcServer() throws Exception {
|
||||
@@ -93,15 +93,17 @@ public class ProxyStartup {
|
||||
} else if (ProxyMode.isLocalMode(proxyModeStr)) {
|
||||
BrokerController brokerController = createBrokerController();
|
||||
StartAndShutdown brokerControllerWrapper = new StartAndShutdown() {
|
||||
@Override public void start() throws Exception {
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
brokerController.start();
|
||||
}
|
||||
|
||||
@Override public void shutdown() throws Exception {
|
||||
@Override
|
||||
public void shutdown() throws Exception {
|
||||
brokerController.shutdown();
|
||||
}
|
||||
};
|
||||
proxyStartAndShutdown.appendStartAndShutdown(brokerControllerWrapper);
|
||||
PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(brokerControllerWrapper);
|
||||
grpcService = new LocalGrpcService(brokerController);
|
||||
} else {
|
||||
throw new IllegalArgumentException("try to start grpc server with wrong mode, use 'local' or 'cluster'");
|
||||
|
||||
@@ -21,9 +21,9 @@ import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
public class InvocationContext<R, W> {
|
||||
final private R request;
|
||||
final private CompletableFuture<W> response;
|
||||
final private long timestamp = System.currentTimeMillis();
|
||||
private final R request;
|
||||
private final CompletableFuture<W> response;
|
||||
private final long timestamp = System.currentTimeMillis();
|
||||
|
||||
public InvocationContext(R req, CompletableFuture<W> resp) {
|
||||
request = req;
|
||||
|
||||
+8
-2
@@ -52,11 +52,17 @@ public class GrpcClientChannel extends SimpleChannel {
|
||||
this.pollCommandResponseFutureRef.set(future);
|
||||
}
|
||||
|
||||
public static GrpcClientChannel create(ChannelManager channelManager, String group, String clientId, PollCommandResponseManager manager) {
|
||||
public static GrpcClientChannel create(
|
||||
ChannelManager channelManager,
|
||||
String group,
|
||||
String clientId,
|
||||
PollCommandResponseManager manager
|
||||
) {
|
||||
GrpcClientChannel channel = channelManager.createChannel(
|
||||
buildKey(group, clientId),
|
||||
() -> new GrpcClientChannel(group, clientId, manager),
|
||||
GrpcClientChannel.class);
|
||||
GrpcClientChannel.class
|
||||
);
|
||||
|
||||
channelManager.addGroupClientId(group, clientId);
|
||||
return channel;
|
||||
|
||||
+5
-3
@@ -46,9 +46,11 @@ public class HeaderInterceptor implements ServerInterceptor {
|
||||
private String parseSocketAddress(SocketAddress socketAddress) {
|
||||
if (socketAddress instanceof InetSocketAddress) {
|
||||
InetSocketAddress inetSocketAddress = (InetSocketAddress) socketAddress;
|
||||
return HostAndPort.fromParts(inetSocketAddress.getAddress()
|
||||
.getHostAddress(), inetSocketAddress.getPort())
|
||||
.toString();
|
||||
return HostAndPort.fromParts(
|
||||
inetSocketAddress.getAddress()
|
||||
.getHostAddress(),
|
||||
inetSocketAddress.getPort()
|
||||
).toString();
|
||||
}
|
||||
|
||||
return "";
|
||||
|
||||
@@ -181,8 +181,9 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
return this.clientService.pollCommand(ctx, request);
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<ReportThreadStackTraceResponse> reportThreadStackTrace(Context ctx,
|
||||
ReportThreadStackTraceRequest request) {
|
||||
@Override
|
||||
public CompletableFuture<ReportThreadStackTraceResponse> reportThreadStackTrace(Context ctx,
|
||||
ReportThreadStackTraceRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
@@ -136,7 +136,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
this.appendStartAndShutdown(new LocalGrpcServiceStartAndShutdown());
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request) {
|
||||
@Override
|
||||
public CompletableFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request) {
|
||||
return this.routeService.queryRoute(ctx, request);
|
||||
}
|
||||
|
||||
@@ -257,7 +258,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
return future;
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<AckMessageResponse> ackMessage(Context ctx, AckMessageRequest request) {
|
||||
@Override
|
||||
public CompletableFuture<AckMessageResponse> ackMessage(Context ctx, AckMessageRequest request) {
|
||||
Channel channel = channelManager.createChannel();
|
||||
SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
|
||||
AckMessageRequestHeader requestHeader = Converter.buildAckMessageRequestHeader(request);
|
||||
@@ -282,7 +284,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
return future;
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<NackMessageResponse> nackMessage(Context ctx, NackMessageRequest request) {
|
||||
@Override
|
||||
public CompletableFuture<NackMessageResponse> nackMessage(Context ctx, NackMessageRequest request) {
|
||||
Channel channel = channelManager.createChannel();
|
||||
SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
|
||||
|
||||
@@ -364,7 +367,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
return future;
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<QueryOffsetResponse> queryOffset(Context ctx, QueryOffsetRequest request) {
|
||||
@Override
|
||||
public CompletableFuture<QueryOffsetResponse> queryOffset(Context ctx, QueryOffsetRequest request) {
|
||||
Partition partition = request.getPartition();
|
||||
String topicName = Converter.getResourceNameWithNamespace(partition.getTopic());
|
||||
int queueId = partition.getId();
|
||||
@@ -386,7 +390,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
.build());
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<PullMessageResponse> pullMessage(Context ctx, PullMessageRequest request) {
|
||||
@Override
|
||||
public CompletableFuture<PullMessageResponse> pullMessage(Context ctx, PullMessageRequest request) {
|
||||
long timeRemaining = Context.current()
|
||||
.getDeadline()
|
||||
.timeRemaining(TimeUnit.MILLISECONDS);
|
||||
@@ -419,7 +424,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
return future;
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<PollCommandResponse> pollCommand(Context ctx, PollCommandRequest request) {
|
||||
@Override
|
||||
public CompletableFuture<PollCommandResponse> pollCommand(Context ctx, PollCommandRequest request) {
|
||||
String clientId = request.getClientId();
|
||||
CompletableFuture<PollCommandResponse> future = new CompletableFuture<>();
|
||||
switch (request.getGroupCase()) {
|
||||
@@ -479,6 +485,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
@Override
|
||||
public CompletableFuture<ReportMessageConsumptionResultResponse> reportMessageConsumptionResult(Context ctx,
|
||||
ReportMessageConsumptionResultRequest request) {
|
||||
|
||||
String commandId = request.getCommandId();
|
||||
PollCommandResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(commandId);
|
||||
if (pollCommandResponseFuture != null) {
|
||||
@@ -497,7 +504,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
.build());
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<NotifyClientTerminationResponse> notifyClientTermination(Context ctx,
|
||||
@Override
|
||||
public CompletableFuture<NotifyClientTerminationResponse> notifyClientTermination(Context ctx,
|
||||
NotifyClientTerminationRequest request) {
|
||||
Channel channel = channelManager.createChannel();
|
||||
SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel);
|
||||
@@ -512,7 +520,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
return new CompletableFuture<>();
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<ChangeInvisibleDurationResponse> changeInvisibleDuration(Context ctx,
|
||||
@Override
|
||||
public CompletableFuture<ChangeInvisibleDurationResponse> changeInvisibleDuration(Context ctx,
|
||||
ChangeInvisibleDurationRequest request) {
|
||||
Channel channel = channelManager.createChannel();
|
||||
SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
|
||||
|
||||
@@ -74,7 +74,6 @@ public class ClientService extends BaseService {
|
||||
if (request.hasProducerData()) {
|
||||
String producerGroup = Converter.getResourceNameWithNamespace(request.getProducerData().getGroup());
|
||||
GrpcClientChannel channel = GrpcClientChannel.create(channelManager, producerGroup, clientId, pollCommandResponseManager);
|
||||
//TODO: Use the faked MQ Version ?
|
||||
ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal());
|
||||
producerManager.registerProducer(producerGroup, clientChannelInfo);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user