From 29084f87d4215683a5b41dee490da9ccdc2ccd5d Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Fri, 25 Mar 2022 20:29:22 +0800 Subject: [PATCH] [ISSUE #3949] Sort code for rebase --- .../rocketmq/broker/BrokerController.java | 4 +++ .../rocketmq/proxy/grpc/GrpcServer.java | 2 +- .../proxy/grpc/service/LocalGrpcService.java | 28 +++++++---------- .../service/cluster/ForwardClientService.java | 10 ++++-- .../grpc/service/LocalGrpcServiceTest.java | 31 +++++++++---------- 5 files changed, 39 insertions(+), 36 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java index 7bf9cdcd9c..ded926f898 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java @@ -2046,6 +2046,10 @@ public class BrokerController { return assignmentManager; } + public ClientManageProcessor getClientManageProcessor() { + return clientManageProcessor; + } + public SendMessageProcessor getSendMessageProcessor() { return sendMessageProcessor; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java index dc36a5e13f..0627e08fb1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java @@ -31,9 +31,9 @@ import java.util.List; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import org.apache.rocketmq.acl.AccessValidator; -import org.apache.rocketmq.broker.util.ServiceProvider; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.common.thread.ThreadPoolMonitor; +import org.apache.rocketmq.common.utils.ServiceProvider; import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.grpc.interceptor.AuthenticationInterceptor; 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 bf3d497edb..8b90940dac 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 @@ -204,12 +204,12 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo = new InvocationContext<>(request, future); channel.registerInvocationContext(command.getOpaque(), context); try { - CompletableFuture processorFuture = brokerController.getSendMessageProcessor() - .asyncProcessRequest(channelHandlerContext, command); - processorFuture.thenAccept(r -> { - handler.handle(r, context); + RemotingCommand response = brokerController.getSendMessageProcessor() + .processRequest(channelHandlerContext, command); + if (response != null) { + handler.handle(response, context); channel.eraseInvocationContext(command.getOpaque()); - }); + } } catch (final Exception e) { LOGGER.error("Failed to process send message command", e); channel.eraseInvocationContext(command.getOpaque()); @@ -313,18 +313,12 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo CompletableFuture future = new CompletableFuture<>(); try { - CompletableFuture processorFuture = brokerController.getSendMessageProcessor() - .asyncProcessRequest(channelHandlerContext, command); - processorFuture.thenAccept(r -> { - ForwardMessageToDeadLetterQueueResponse.Builder builder = ForwardMessageToDeadLetterQueueResponse.newBuilder(); - if (null != r) { - builder.setCommon(ResponseBuilder.buildCommon(r.getCode(), r.getRemark())); - } else { - builder.setCommon(ResponseBuilder.buildCommon(Code.INTERNAL, "Response command is null")); - } - ForwardMessageToDeadLetterQueueResponse response = builder.build(); - future.complete(response); - }); + RemotingCommand response = brokerController.getSendMessageProcessor() + .processRequest(channelHandlerContext, command); + + future.complete(ForwardMessageToDeadLetterQueueResponse.newBuilder() + .setCommon(ResponseBuilder.buildCommon(response.getCode(), response.getRemark())) + .build()); } catch (Exception e) { LOGGER.error("Exception raised when forwardMessageToDeadLetterQueue", e); future.completeExceptionally(e); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java index eb8144e34f..1b68c1c1ac 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java @@ -29,6 +29,8 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import org.apache.rocketmq.broker.client.ClientChannelInfo; +import org.apache.rocketmq.broker.client.ConsumerGroupEvent; +import org.apache.rocketmq.broker.client.ConsumerIdsChangeListener; import org.apache.rocketmq.broker.client.ConsumerManager; import org.apache.rocketmq.broker.client.ProducerManager; import org.apache.rocketmq.common.MQVersion; @@ -65,8 +67,12 @@ public class ForwardClientService extends BaseService { this.channelManager = channelManager; this.pollCommandResponseManager = pollCommandResponseManager; - this.consumerManager = new ConsumerManager((event, group, args) -> { - // nothing to do in handler. + this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListener() { + @Override public void handle(ConsumerGroupEvent event, String group, Object... args) { + } + + @Override public void shutdown() { + } }); this.producerManager = new ProducerManager(); this.producerManager.setProducerOfflineListener(connectorManager.getTransactionHeartbeatRegisterService()::onProducerGroupOffline); 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 899cc3cebc..6a252d2ca8 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 @@ -82,6 +82,7 @@ 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; +import org.apache.rocketmq.store.config.MessageStoreConfig; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -113,6 +114,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { Mockito.when(brokerControllerMock.getPopMessageProcessor()).thenReturn(popMessageProcessorMock); Mockito.when(brokerControllerMock.getPullMessageProcessor()).thenReturn(pullMessageProcessorMock); Mockito.when(brokerControllerMock.getBrokerConfig()).thenReturn(new BrokerConfig()); + Mockito.when(brokerControllerMock.getMessageStoreConfig()).thenReturn(new MessageStoreConfig()); localGrpcService = new LocalGrpcService(brokerControllerMock); metadata = new Metadata(); metadata.put(InterceptorConstants.REMOTE_ADDRESS, "1.1.1.1"); @@ -168,9 +170,8 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { public void testSendMessageError() throws Exception { String remark = "store putMessage return null"; RemotingCommand response = RemotingCommand.createResponseCommand(ResponseCode.SYSTEM_ERROR, remark); - CompletableFuture future = CompletableFuture.completedFuture(response); - Mockito.when(sendMessageProcessorMock.asyncProcessRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) - .thenReturn(future); + Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) + .thenReturn(response); SendMessageRequest request = SendMessageRequest.newBuilder() .setMessage(Message.newBuilder() .setSystemAttribute(SystemAttribute.newBuilder() @@ -188,9 +189,8 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { @Test public void testSendMessageWriteAndFlush() throws Exception { - CompletableFuture future = CompletableFuture.completedFuture(null); - Mockito.when(sendMessageProcessorMock.asyncProcessRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) - .thenReturn(future); + Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) + .thenReturn(null); SendMessageRequest request = SendMessageRequest.newBuilder() .setMessage(Message.newBuilder() .setSystemAttribute(SystemAttribute.newBuilder() @@ -206,7 +206,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { @Test public void testSendMessageWithException() throws Exception { - Mockito.when(sendMessageProcessorMock.asyncProcessRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) + Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) .thenThrow(new RemotingCommandException("test")); SendMessageRequest request = SendMessageRequest.newBuilder() .setMessage(Message.newBuilder() @@ -256,9 +256,9 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { CompletableFuture grpcFuture = localGrpcService.receiveMessage( Context.current() .withValue(InterceptorConstants.METADATA, metadata) + .attach() .withDeadlineAfter(20, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor( - new ThreadFactoryImpl("test"))) - .attach(), request); + new ThreadFactoryImpl("test"))), request); ReceiveMessageResponse r = grpcFuture.get(); assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); assertThat(r.getMessagesCount()).isEqualTo(1); @@ -275,9 +275,9 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { CompletableFuture grpcFuture = localGrpcService.receiveMessage( Context.current() .withValue(InterceptorConstants.METADATA, metadata) + .attach() .withDeadlineAfter(20, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor( - new ThreadFactoryImpl("test"))) - .attach(), request); + new ThreadFactoryImpl("test"))), request); assertThat(grpcFuture.isDone()).isFalse(); } @@ -345,11 +345,10 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { @Test public void testForwardMessageToDeadLetterQueue() throws Exception { RemotingCommand response = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, null); - CompletableFuture future = CompletableFuture.completedFuture(response); Mockito.when(brokerControllerMock.getSendMessageProcessor()).thenReturn(sendMessageProcessorMock); - Mockito.when(sendMessageProcessorMock.asyncProcessRequest(Mockito.any(ChannelHandlerContext.class), + Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.argThat(argument -> argument.getCode() == RequestCode.CONSUMER_SEND_MSG_BACK))) - .thenReturn(future); + .thenReturn(response); ForwardMessageToDeadLetterQueueRequest request = ForwardMessageToDeadLetterQueueRequest.newBuilder() .setReceiptHandle(ReceiptHandle.builder() .startOffset(0L) @@ -557,9 +556,9 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { CompletableFuture grpcFuture = localGrpcService.pullMessage( Context.current() .withValue(InterceptorConstants.METADATA, metadata) + .attach() .withDeadlineAfter(20, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor( - new ThreadFactoryImpl("test"))) - .attach(), request); + new ThreadFactoryImpl("test"))), request); PullMessageResponse r = grpcFuture.get(); assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); assertThat(r.getMessagesCount()).isEqualTo(1);