mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Sort code for rebase
This commit is contained in:
@@ -2046,6 +2046,10 @@ public class BrokerController {
|
||||
return assignmentManager;
|
||||
}
|
||||
|
||||
public ClientManageProcessor getClientManageProcessor() {
|
||||
return clientManageProcessor;
|
||||
}
|
||||
|
||||
public SendMessageProcessor getSendMessageProcessor() {
|
||||
return sendMessageProcessor;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -204,12 +204,12 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
= new InvocationContext<>(request, future);
|
||||
channel.registerInvocationContext(command.getOpaque(), context);
|
||||
try {
|
||||
CompletableFuture<RemotingCommand> 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<ForwardMessageToDeadLetterQueueResponse> future = new CompletableFuture<>();
|
||||
try {
|
||||
CompletableFuture<RemotingCommand> 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);
|
||||
|
||||
+8
-2
@@ -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);
|
||||
|
||||
+15
-16
@@ -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<RemotingCommand> 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<RemotingCommand> 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<ReceiveMessageResponse> 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<ReceiveMessageResponse> 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<RemotingCommand> 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<PullMessageResponse> 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);
|
||||
|
||||
Reference in New Issue
Block a user