mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 13:49:50 +08:00
[ISSUE #3949] Complete LocalMessageService
This commit is contained in:
@@ -22,4 +22,5 @@ public enum ProxyExceptionCode {
|
||||
INVALID_BROKER_NAME,
|
||||
INVALID_RECEIPT_HANDLE,
|
||||
ILLEGAL_MESSAGE,
|
||||
INTERNAL_SERVER_ERROR,
|
||||
}
|
||||
|
||||
@@ -35,6 +35,7 @@ public class GrpcProxyException extends RuntimeException {
|
||||
CODE_MAPPING.put(ProxyExceptionCode.RECEIPT_HANDLE_EXPIRED, Code.RECEIPT_HANDLE_EXPIRED);
|
||||
CODE_MAPPING.put(ProxyExceptionCode.FORBIDDEN, Code.FORBIDDEN);
|
||||
CODE_MAPPING.put(ProxyExceptionCode.ILLEGAL_MESSAGE, Code.ILLEGAL_MESSAGE);
|
||||
CODE_MAPPING.put(ProxyExceptionCode.INTERNAL_SERVER_ERROR, Code.INTERNAL_SERVER_ERROR);
|
||||
}
|
||||
|
||||
public GrpcProxyException(Code code, String message) {
|
||||
|
||||
+187
-7
@@ -17,13 +17,18 @@
|
||||
package org.apache.rocketmq.proxy.service.message;
|
||||
|
||||
import io.netty.channel.ChannelHandlerContext;
|
||||
import java.util.Arrays;
|
||||
import java.nio.ByteBuffer;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.broker.BrokerController;
|
||||
import org.apache.rocketmq.client.consumer.AckResult;
|
||||
import org.apache.rocketmq.client.consumer.AckStatus;
|
||||
import org.apache.rocketmq.client.consumer.PopResult;
|
||||
import org.apache.rocketmq.client.consumer.PopStatus;
|
||||
import org.apache.rocketmq.client.producer.SendResult;
|
||||
import org.apache.rocketmq.client.producer.SendStatus;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
@@ -31,14 +36,20 @@ import org.apache.rocketmq.common.consumer.ReceiptHandle;
|
||||
import org.apache.rocketmq.common.message.Message;
|
||||
import org.apache.rocketmq.common.message.MessageBatch;
|
||||
import org.apache.rocketmq.common.message.MessageClientIDSetter;
|
||||
import org.apache.rocketmq.common.message.MessageConst;
|
||||
import org.apache.rocketmq.common.message.MessageDecoder;
|
||||
import org.apache.rocketmq.common.message.MessageExt;
|
||||
import org.apache.rocketmq.common.message.MessageQueue;
|
||||
import org.apache.rocketmq.common.protocol.RequestCode;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.EndTransactionRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ExtraInfoUtil;
|
||||
import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.PopMessageResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageResponseHeader;
|
||||
import org.apache.rocketmq.proxy.common.ProxyContext;
|
||||
@@ -94,7 +105,7 @@ public class LocalMessageService implements MessageService {
|
||||
}
|
||||
} catch (Exception e) {
|
||||
future.completeExceptionally(e);
|
||||
log.error("Failed to process send message command", e);
|
||||
log.error("Failed to process sendMessage command", e);
|
||||
} finally {
|
||||
channel.eraseInvocationContext(request.getOpaque());
|
||||
}
|
||||
@@ -136,27 +147,196 @@ public class LocalMessageService implements MessageService {
|
||||
@Override
|
||||
public CompletableFuture<RemotingCommand> sendMessageBack(ProxyContext ctx, ReceiptHandle handle, String messageId,
|
||||
ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) {
|
||||
return null;
|
||||
SimpleChannel channel = channelManager.createChannel(ctx);
|
||||
ChannelHandlerContext channelHandlerContext = channel.getChannelHandlerContext();
|
||||
RemotingCommand command = LocalRemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader);
|
||||
CompletableFuture<RemotingCommand> future = new CompletableFuture<>();
|
||||
try {
|
||||
RemotingCommand response = brokerController.getSendMessageProcessor()
|
||||
.processRequest(channelHandlerContext, command);
|
||||
future.complete(response);
|
||||
} catch (Exception e) {
|
||||
log.error("Fail to process sendMessageBack command", e);
|
||||
future.completeExceptionally(e);
|
||||
}
|
||||
return future;
|
||||
}
|
||||
|
||||
@Override public void endTransactionOneway(ProxyContext ctx, TransactionId transactionId,
|
||||
EndTransactionRequestHeader requestHeader, long timeoutMillis) {
|
||||
|
||||
SimpleChannel channel = channelManager.createChannel(ctx);
|
||||
ChannelHandlerContext channelHandlerContext = channel.getChannelHandlerContext();
|
||||
RemotingCommand command = LocalRemotingCommand.createRequestCommand(RequestCode.END_TRANSACTION, requestHeader);
|
||||
try {
|
||||
brokerController.getEndTransactionProcessor()
|
||||
.processRequest(channelHandlerContext, command);
|
||||
} catch (Exception e) {
|
||||
log.error("Fail to process endTransaction command", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<PopResult> popMessage(ProxyContext ctx, SelectableMessageQueue messageQueue,
|
||||
PopMessageRequestHeader requestHeader, long timeoutMillis) {
|
||||
return null;
|
||||
RemotingCommand request = LocalRemotingCommand.createRequestCommand(RequestCode.POP_MESSAGE, requestHeader);
|
||||
CompletableFuture<RemotingCommand> future = new CompletableFuture<>();
|
||||
SimpleChannel channel = channelManager.createInvocationChannel(ctx);
|
||||
InvocationContext invocationContext = new InvocationContext(future);
|
||||
channel.registerInvocationContext(request.getOpaque(), invocationContext);
|
||||
ChannelHandlerContext simpleChannelHandlerContext = channel.getChannelHandlerContext();
|
||||
try {
|
||||
RemotingCommand response = brokerController.getPopMessageProcessor().processRequest(simpleChannelHandlerContext, request);
|
||||
if (response != null) {
|
||||
invocationContext.handle(response);
|
||||
}
|
||||
} catch (Exception e) {
|
||||
future.completeExceptionally(e);
|
||||
log.error("Failed to process popMessage command", e);
|
||||
} finally {
|
||||
channel.eraseInvocationContext(request.getOpaque());
|
||||
}
|
||||
return future.thenApply(r -> {
|
||||
PopStatus popStatus;
|
||||
List<MessageExt> messageExtList = new ArrayList<>();
|
||||
switch (r.getCode()) {
|
||||
case ResponseCode.SUCCESS:
|
||||
popStatus = PopStatus.FOUND;
|
||||
ByteBuffer byteBuffer = ByteBuffer.wrap(r.getBody());
|
||||
messageExtList = MessageDecoder.decodes(byteBuffer);
|
||||
break;
|
||||
case ResponseCode.POLLING_FULL:
|
||||
popStatus = PopStatus.POLLING_FULL;
|
||||
break;
|
||||
case ResponseCode.POLLING_TIMEOUT:
|
||||
case ResponseCode.PULL_NOT_FOUND:
|
||||
popStatus = PopStatus.POLLING_NOT_FOUND;
|
||||
break;
|
||||
default:
|
||||
throw new ProxyException(ProxyExceptionCode.INTERNAL_SERVER_ERROR, r.getRemark());
|
||||
}
|
||||
PopResult popResult = new PopResult(popStatus, messageExtList);
|
||||
PopMessageResponseHeader responseHeader = (PopMessageResponseHeader) r.readCustomHeader();
|
||||
|
||||
if (popStatus == PopStatus.FOUND) {
|
||||
Map<String, Long> startOffsetInfo = null;
|
||||
Map<String, List<Long>> msgOffsetInfo = null;
|
||||
Map<String, Integer> orderCountInfo = null;
|
||||
if (requestHeader != null) {
|
||||
popResult.setInvisibleTime(responseHeader.getInvisibleTime());
|
||||
popResult.setPopTime(responseHeader.getPopTime());
|
||||
startOffsetInfo = ExtraInfoUtil.parseStartOffsetInfo(responseHeader.getStartOffsetInfo());
|
||||
msgOffsetInfo = ExtraInfoUtil.parseMsgOffsetInfo(responseHeader.getMsgOffsetInfo());
|
||||
orderCountInfo = ExtraInfoUtil.parseOrderCountInfo(responseHeader.getOrderCountInfo());
|
||||
}
|
||||
// <topicMark@queueId, msg queueOffset>
|
||||
Map<String, List<Long>> sortMap = new HashMap<>(16);
|
||||
for (MessageExt messageExt : messageExtList) {
|
||||
String key = ExtraInfoUtil.getStartOffsetInfoMapKey(messageExt.getTopic(), messageExt.getQueueId());
|
||||
if (!sortMap.containsKey(key)) {
|
||||
sortMap.put(key, new ArrayList<>(4));
|
||||
}
|
||||
sortMap.get(key).add(messageExt.getQueueOffset());
|
||||
}
|
||||
Map<String, String> map = new HashMap<>(5);
|
||||
for (MessageExt messageExt : messageExtList) {
|
||||
if (requestHeader != null) {
|
||||
if (startOffsetInfo == null) {
|
||||
// we should set the check point info to extraInfo field , if the command is popMsg
|
||||
// find pop ck offset
|
||||
String key = messageExt.getTopic() + messageExt.getQueueId();
|
||||
if (!map.containsKey(messageExt.getTopic() + messageExt.getQueueId())) {
|
||||
map.put(key, ExtraInfoUtil.buildExtraInfo(messageExt.getQueueOffset(), responseHeader.getPopTime(), responseHeader.getInvisibleTime(), responseHeader.getReviveQid(),
|
||||
messageExt.getTopic(), messageQueue.getBrokerName(), messageExt.getQueueId()));
|
||||
}
|
||||
messageExt.getProperties().put(MessageConst.PROPERTY_POP_CK, map.get(key) + MessageConst.KEY_SEPARATOR + messageExt.getQueueOffset());
|
||||
} else {
|
||||
String key = ExtraInfoUtil.getStartOffsetInfoMapKey(messageExt.getTopic(), messageExt.getQueueId());
|
||||
int index = sortMap.get(key).indexOf(messageExt.getQueueOffset());
|
||||
Long msgQueueOffset = msgOffsetInfo.get(key).get(index);
|
||||
if (msgQueueOffset != messageExt.getQueueOffset()) {
|
||||
log.warn("Queue offset [{}] of msg is strange, not equal to the stored in msg, {}", msgQueueOffset, messageExt);
|
||||
}
|
||||
|
||||
messageExt.getProperties().put(MessageConst.PROPERTY_POP_CK,
|
||||
ExtraInfoUtil.buildExtraInfo(startOffsetInfo.get(key), responseHeader.getPopTime(), responseHeader.getInvisibleTime(),
|
||||
responseHeader.getReviveQid(), messageExt.getTopic(), messageQueue.getBrokerName(), messageExt.getQueueId(), msgQueueOffset)
|
||||
);
|
||||
if (requestHeader.isOrder() && orderCountInfo != null) {
|
||||
Integer count = orderCountInfo.get(key);
|
||||
if (count != null && count > 0) {
|
||||
messageExt.setReconsumeTimes(count);
|
||||
}
|
||||
}
|
||||
}
|
||||
messageExt.getProperties().computeIfAbsent(MessageConst.PROPERTY_FIRST_POP_TIME, k -> String.valueOf(responseHeader.getPopTime()));
|
||||
}
|
||||
messageExt.setBrokerName(messageExt.getBrokerName());
|
||||
}
|
||||
}
|
||||
return popResult;
|
||||
});
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<AckResult> changeInvisibleTime(ProxyContext ctx, ReceiptHandle handle, String messageId,
|
||||
ChangeInvisibleTimeRequestHeader requestHeader, long timeoutMillis) {
|
||||
return null;
|
||||
SimpleChannel channel = channelManager.createChannel(ctx);
|
||||
ChannelHandlerContext channelHandlerContext = channel.getChannelHandlerContext();
|
||||
RemotingCommand command = LocalRemotingCommand.createRequestCommand(RequestCode.CHANGE_MESSAGE_INVISIBLETIME, requestHeader);
|
||||
CompletableFuture<RemotingCommand> future = new CompletableFuture<>();
|
||||
try {
|
||||
RemotingCommand response = brokerController.getChangeInvisibleTimeProcessor()
|
||||
.processRequest(channelHandlerContext, command);
|
||||
future.complete(response);
|
||||
} catch (Exception e) {
|
||||
log.error("Fail to process changeInvisibleTime command", e);
|
||||
future.completeExceptionally(e);
|
||||
}
|
||||
return future.thenApply(r -> {
|
||||
ChangeInvisibleTimeResponseHeader responseHeader = (ChangeInvisibleTimeResponseHeader) r.readCustomHeader();
|
||||
AckResult ackResult = new AckResult();
|
||||
if (ResponseCode.SUCCESS == r.getCode()) {
|
||||
ackResult.setStatus(AckStatus.OK);
|
||||
} else {
|
||||
ackResult.setStatus(AckStatus.NO_EXIST);
|
||||
}
|
||||
ackResult.setPopTime(responseHeader.getPopTime());
|
||||
ackResult.setExtraInfo(ReceiptHandle.builder()
|
||||
.startOffset(handle.getStartOffset())
|
||||
.retrieveTime(responseHeader.getPopTime())
|
||||
.invisibleTime(responseHeader.getInvisibleTime())
|
||||
.reviveQueueId(responseHeader.getReviveQid())
|
||||
.topicType(handle.getTopicType())
|
||||
.brokerName(handle.getBrokerName())
|
||||
.queueId(handle.getQueueId())
|
||||
.offset(handle.getOffset())
|
||||
.build()
|
||||
.encode());
|
||||
return ackResult;
|
||||
});
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<AckResult> ackMessage(ProxyContext ctx, ReceiptHandle handle, String messageId,
|
||||
AckMessageRequestHeader requestHeader, long timeoutMillis) {
|
||||
return null;
|
||||
SimpleChannel channel = channelManager.createChannel(ctx);
|
||||
ChannelHandlerContext channelHandlerContext = channel.getChannelHandlerContext();
|
||||
RemotingCommand command = LocalRemotingCommand.createRequestCommand(RequestCode.ACK_MESSAGE, requestHeader);
|
||||
CompletableFuture<RemotingCommand> future = new CompletableFuture<>();
|
||||
try {
|
||||
RemotingCommand response = brokerController.getAckMessageProcessor()
|
||||
.processRequest(channelHandlerContext, command);
|
||||
future.complete(response);
|
||||
} catch (Exception e) {
|
||||
log.error("Fail to process ackMessage command", e);
|
||||
future.completeExceptionally(e);
|
||||
}
|
||||
return future.thenApply(r -> {
|
||||
AckResult ackResult = new AckResult();
|
||||
if (ResponseCode.SUCCESS == r.getCode()) {
|
||||
ackResult.setStatus(AckStatus.OK);
|
||||
} else {
|
||||
ackResult.setStatus(AckStatus.NO_EXIST);
|
||||
}
|
||||
return ackResult;
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -26,7 +26,6 @@ import org.apache.rocketmq.client.producer.SendStatus;
|
||||
import org.apache.rocketmq.common.KeyBuilder;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.consumer.ReceiptHandle;
|
||||
import org.apache.rocketmq.common.message.Message;
|
||||
import org.apache.rocketmq.common.message.MessageAccessor;
|
||||
import org.apache.rocketmq.common.message.MessageClientIDSetter;
|
||||
import org.apache.rocketmq.common.message.MessageConst;
|
||||
@@ -34,8 +33,6 @@ import org.apache.rocketmq.common.message.MessageExt;
|
||||
import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.sysflag.MessageSysFlag;
|
||||
import org.apache.rocketmq.proxy.common.ProxyContext;
|
||||
import org.apache.rocketmq.proxy.service.route.MessageQueueView;
|
||||
import org.apache.rocketmq.proxy.service.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.service.transaction.TransactionId;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
@@ -44,7 +41,8 @@ import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.mockito.ArgumentCaptor;
|
||||
|
||||
import static org.junit.Assert.*;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.anyLong;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
|
||||
+243
@@ -17,23 +17,45 @@
|
||||
|
||||
package org.apache.rocketmq.proxy.service.message;
|
||||
|
||||
import java.net.InetSocketAddress;
|
||||
import java.nio.ByteBuffer;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import org.apache.rocketmq.broker.BrokerController;
|
||||
import org.apache.rocketmq.broker.processor.AckMessageProcessor;
|
||||
import org.apache.rocketmq.broker.processor.ChangeInvisibleTimeProcessor;
|
||||
import org.apache.rocketmq.broker.processor.EndTransactionProcessor;
|
||||
import org.apache.rocketmq.broker.processor.PopMessageProcessor;
|
||||
import org.apache.rocketmq.broker.processor.SendMessageProcessor;
|
||||
import org.apache.rocketmq.client.consumer.AckResult;
|
||||
import org.apache.rocketmq.client.consumer.AckStatus;
|
||||
import org.apache.rocketmq.client.consumer.PopResult;
|
||||
import org.apache.rocketmq.client.consumer.PopStatus;
|
||||
import org.apache.rocketmq.client.producer.SendResult;
|
||||
import org.apache.rocketmq.client.producer.SendStatus;
|
||||
import org.apache.rocketmq.common.BrokerConfig;
|
||||
import org.apache.rocketmq.common.consumer.ReceiptHandle;
|
||||
import org.apache.rocketmq.common.message.Message;
|
||||
import org.apache.rocketmq.common.message.MessageBatch;
|
||||
import org.apache.rocketmq.common.message.MessageClientIDSetter;
|
||||
import org.apache.rocketmq.common.message.MessageDecoder;
|
||||
import org.apache.rocketmq.common.message.MessageExt;
|
||||
import org.apache.rocketmq.common.message.MessageQueue;
|
||||
import org.apache.rocketmq.common.protocol.RequestCode;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.EndTransactionRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ExtraInfoUtil;
|
||||
import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.PopMessageResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageResponseHeader;
|
||||
import org.apache.rocketmq.proxy.common.ContextVariable;
|
||||
@@ -44,6 +66,7 @@ import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.config.InitConfigAndLoggerTest;
|
||||
import org.apache.rocketmq.proxy.service.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.service.channel.SimpleChannelHandlerContext;
|
||||
import org.apache.rocketmq.proxy.service.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.junit.Before;
|
||||
@@ -62,6 +85,14 @@ public class LocalMessageServiceTest extends InitConfigAndLoggerTest {
|
||||
@Mock
|
||||
private SendMessageProcessor sendMessageProcessorMock;
|
||||
@Mock
|
||||
private EndTransactionProcessor endTransactionProcessorMock;
|
||||
@Mock
|
||||
private PopMessageProcessor popMessageProcessorMock;
|
||||
@Mock
|
||||
private ChangeInvisibleTimeProcessor changeInvisibleTimeProcessorMock;
|
||||
@Mock
|
||||
private AckMessageProcessor ackMessageProcessorMock;
|
||||
@Mock
|
||||
private BrokerController brokerControllerMock;
|
||||
|
||||
private ProxyContext proxyContext;
|
||||
@@ -70,6 +101,8 @@ public class LocalMessageServiceTest extends InitConfigAndLoggerTest {
|
||||
|
||||
private String topic = "topic";
|
||||
|
||||
private String brokerName = "brokerName";
|
||||
|
||||
private int queueId = 0;
|
||||
|
||||
private long queueOffset = 0L;
|
||||
@@ -84,6 +117,10 @@ public class LocalMessageServiceTest extends InitConfigAndLoggerTest {
|
||||
ConfigurationManager.getProxyConfig().setNameSrvAddr("1.1.1.1");
|
||||
channelManager = new ChannelManager();
|
||||
Mockito.when(brokerControllerMock.getSendMessageProcessor()).thenReturn(sendMessageProcessorMock);
|
||||
Mockito.when(brokerControllerMock.getPopMessageProcessor()).thenReturn(popMessageProcessorMock);
|
||||
Mockito.when(brokerControllerMock.getChangeInvisibleTimeProcessor()).thenReturn(changeInvisibleTimeProcessorMock);
|
||||
Mockito.when(brokerControllerMock.getAckMessageProcessor()).thenReturn(ackMessageProcessorMock);
|
||||
Mockito.when(brokerControllerMock.getEndTransactionProcessor()).thenReturn(endTransactionProcessorMock);
|
||||
Mockito.when(brokerControllerMock.getBrokerConfig()).thenReturn(new BrokerConfig());
|
||||
localMessageService = new LocalMessageService(brokerControllerMock, channelManager, null);
|
||||
proxyContext = ProxyContext.create().withVal(ContextVariable.REMOTE_ADDRESS, "0.0.0.1")
|
||||
@@ -205,4 +242,210 @@ public class LocalMessageServiceTest extends InitConfigAndLoggerTest {
|
||||
ExecutionException exception = catchThrowableOfType(future::get, ExecutionException.class);
|
||||
assertThat(exception.getCause()).isInstanceOf(RemotingCommandException.class);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSendMessageBack() throws Exception {
|
||||
RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "");
|
||||
Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(SimpleChannelHandlerContext.class), Mockito.argThat(argument -> {
|
||||
boolean first = argument.getCode() == RequestCode.CONSUMER_SEND_MSG_BACK;
|
||||
boolean second = argument.readCustomHeader() instanceof ConsumerSendMsgBackRequestHeader;
|
||||
return first && second;
|
||||
}))).thenReturn(remotingCommand);
|
||||
ConsumerSendMsgBackRequestHeader requestHeader = new ConsumerSendMsgBackRequestHeader();
|
||||
CompletableFuture<RemotingCommand> future = localMessageService.sendMessageBack(proxyContext, null, null, requestHeader, 1000L);
|
||||
RemotingCommand response = future.get();
|
||||
assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEndTransaction() throws Exception {
|
||||
EndTransactionRequestHeader requestHeader = new EndTransactionRequestHeader();
|
||||
localMessageService.endTransactionOneway(proxyContext, null, requestHeader, 1000L);
|
||||
Mockito.verify(endTransactionProcessorMock, Mockito.times(1)).processRequest(Mockito.any(SimpleChannelHandlerContext.class), Mockito.argThat(argument -> {
|
||||
boolean first = argument.getCode() == RequestCode.END_TRANSACTION;
|
||||
boolean second = argument.readCustomHeader() instanceof EndTransactionRequestHeader;
|
||||
return first && second;
|
||||
}));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testPopMessageWriteAndFlush() throws Exception {
|
||||
int reviveQueueId = 1;
|
||||
long popTime = System.currentTimeMillis();
|
||||
long invisibleTime = 3000L;
|
||||
long startOffset = 100L;
|
||||
long restNum = 0L;
|
||||
StringBuilder startOffsetStringBuilder = new StringBuilder();
|
||||
StringBuilder messageOffsetStringBuilder = new StringBuilder();
|
||||
List<MessageExt> messageExtList = new ArrayList<>();
|
||||
List<Long> messageOffsetList = new ArrayList<>();
|
||||
MessageExt message1 = buildMessageExt(topic, 0, startOffset);
|
||||
messageExtList.add(message1);
|
||||
messageOffsetList.add(startOffset);
|
||||
byte[] body1 = MessageDecoder.encode(message1, false);
|
||||
MessageExt message2 = buildMessageExt(topic, 0, startOffset + 1);
|
||||
messageExtList.add(message2);
|
||||
messageOffsetList.add(startOffset + 1);
|
||||
ExtraInfoUtil.buildStartOffsetInfo(startOffsetStringBuilder, false, queueId, startOffset);
|
||||
ExtraInfoUtil.buildMsgOffsetInfo(messageOffsetStringBuilder, false, queueId, messageOffsetList);
|
||||
byte[] body2 = MessageDecoder.encode(message2, false);
|
||||
ByteBuffer byteBuffer1 = ByteBuffer.wrap(body1);
|
||||
ByteBuffer byteBuffer2 = ByteBuffer.wrap(body2);
|
||||
ByteBuffer b3 = ByteBuffer.allocate(byteBuffer1.limit() + byteBuffer2.limit());
|
||||
b3.put(byteBuffer1);
|
||||
b3.put(byteBuffer2);
|
||||
PopMessageRequestHeader requestHeader = new PopMessageRequestHeader();
|
||||
requestHeader.setInvisibleTime(invisibleTime);
|
||||
Mockito.when(popMessageProcessorMock.processRequest(Mockito.any(SimpleChannelHandlerContext.class), Mockito.argThat(argument -> {
|
||||
boolean first = argument.getCode() == RequestCode.POP_MESSAGE;
|
||||
boolean second = argument.readCustomHeader() instanceof PopMessageRequestHeader;
|
||||
return first && second;
|
||||
}))).thenAnswer(invocation -> {
|
||||
SimpleChannelHandlerContext simpleChannelHandlerContext = invocation.getArgument(0);
|
||||
RemotingCommand request = invocation.getArgument(1);
|
||||
RemotingCommand response = RemotingCommand.createResponseCommand(PopMessageResponseHeader.class);
|
||||
response.setOpaque(request.getOpaque());
|
||||
response.setCode(ResponseCode.SUCCESS);
|
||||
response.setBody(b3.array());
|
||||
PopMessageResponseHeader responseHeader = (PopMessageResponseHeader) response.readCustomHeader();
|
||||
responseHeader.setStartOffsetInfo(startOffsetStringBuilder.toString());
|
||||
responseHeader.setMsgOffsetInfo(messageOffsetStringBuilder.toString());
|
||||
responseHeader.setInvisibleTime(requestHeader.getInvisibleTime());
|
||||
responseHeader.setPopTime(popTime);
|
||||
responseHeader.setRestNum(restNum);
|
||||
responseHeader.setReviveQid(reviveQueueId);
|
||||
simpleChannelHandlerContext.writeAndFlush(response);
|
||||
return null;
|
||||
});
|
||||
MessageQueue messageQueue = new MessageQueue(topic, brokerName, queueId);
|
||||
CompletableFuture<PopResult> future = localMessageService.popMessage(proxyContext, new SelectableMessageQueue(messageQueue, ""), requestHeader, 1000L);
|
||||
PopResult popResult = future.get();
|
||||
assertThat(popResult.getPopTime()).isEqualTo(popTime);
|
||||
assertThat(popResult.getInvisibleTime()).isEqualTo(invisibleTime);
|
||||
assertThat(popResult.getPopStatus()).isEqualTo(PopStatus.FOUND);
|
||||
assertThat(popResult.getRestNum()).isEqualTo(restNum);
|
||||
assertThat(popResult.getMsgFoundList().size()).isEqualTo(messageExtList.size());
|
||||
for (int i = 0; i < popResult.getMsgFoundList().size(); i++) {
|
||||
assertMessageExt(popResult.getMsgFoundList().get(i), messageExtList.get(i));
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testPopMessagePollingTimeout() throws Exception {
|
||||
RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.POLLING_TIMEOUT, "");
|
||||
Mockito.when(popMessageProcessorMock.processRequest(Mockito.any(SimpleChannelHandlerContext.class), Mockito.argThat(argument -> {
|
||||
boolean first = argument.getCode() == RequestCode.POP_MESSAGE;
|
||||
boolean second = argument.readCustomHeader() instanceof PopMessageRequestHeader;
|
||||
return first && second;
|
||||
}))).thenReturn(remotingCommand);
|
||||
PopMessageRequestHeader requestHeader = new PopMessageRequestHeader();
|
||||
CompletableFuture<PopResult> future = localMessageService.popMessage(proxyContext, null, requestHeader, 1000L);
|
||||
PopResult popResult = future.get();
|
||||
assertThat(popResult.getPopStatus()).isEqualTo(PopStatus.POLLING_NOT_FOUND);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testChangeInvisibleTime() throws Exception {
|
||||
String messageId = "messageId";
|
||||
long popTime = System.currentTimeMillis();
|
||||
long invisibleTime = 3000L;
|
||||
int reviveQueueId = 1;
|
||||
ReceiptHandle handle = ReceiptHandle.builder()
|
||||
.startOffset(0L)
|
||||
.retrieveTime(popTime)
|
||||
.invisibleTime(invisibleTime)
|
||||
.reviveQueueId(reviveQueueId)
|
||||
.topicType(ReceiptHandle.NORMAL_TOPIC)
|
||||
.brokerName(brokerName)
|
||||
.queueId(queueId)
|
||||
.offset(queueOffset)
|
||||
.build();
|
||||
RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ChangeInvisibleTimeResponseHeader.class);
|
||||
remotingCommand.setCode(ResponseCode.SUCCESS);
|
||||
remotingCommand.setRemark("");
|
||||
long newPopTime = System.currentTimeMillis();
|
||||
long newInvisibleTime = 5000L;
|
||||
int newReviveQueueId = 2;
|
||||
ChangeInvisibleTimeResponseHeader responseHeader = (ChangeInvisibleTimeResponseHeader) remotingCommand.readCustomHeader();
|
||||
responseHeader.setReviveQid(newReviveQueueId);
|
||||
responseHeader.setInvisibleTime(newInvisibleTime);
|
||||
responseHeader.setPopTime(newPopTime);
|
||||
Mockito.when(changeInvisibleTimeProcessorMock.processRequest(Mockito.any(SimpleChannelHandlerContext.class), Mockito.argThat(argument -> {
|
||||
boolean first = argument.getCode() == RequestCode.CHANGE_MESSAGE_INVISIBLETIME;
|
||||
boolean second = argument.readCustomHeader() instanceof ChangeInvisibleTimeRequestHeader;
|
||||
return first && second;
|
||||
}))).thenReturn(remotingCommand);
|
||||
ChangeInvisibleTimeRequestHeader requestHeader = new ChangeInvisibleTimeRequestHeader();
|
||||
CompletableFuture<AckResult> future = localMessageService.changeInvisibleTime(proxyContext, handle, messageId,
|
||||
requestHeader, 1000L);
|
||||
AckResult ackResult = future.get();
|
||||
assertThat(ackResult.getStatus()).isEqualTo(AckStatus.OK);
|
||||
assertThat(ackResult.getPopTime()).isEqualTo(newPopTime);
|
||||
assertThat(ackResult.getExtraInfo()).isEqualTo(ReceiptHandle.builder()
|
||||
.startOffset(0L)
|
||||
.retrieveTime(newPopTime)
|
||||
.invisibleTime(newInvisibleTime)
|
||||
.reviveQueueId(newReviveQueueId)
|
||||
.topicType(ReceiptHandle.NORMAL_TOPIC)
|
||||
.brokerName(brokerName)
|
||||
.queueId(queueId)
|
||||
.offset(queueOffset)
|
||||
.build()
|
||||
.encode());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testAckMessage() throws Exception {
|
||||
String messageId = "messageId";
|
||||
long popTime = System.currentTimeMillis();
|
||||
long invisibleTime = 3000L;
|
||||
int reviveQueueId = 1;
|
||||
ReceiptHandle handle = ReceiptHandle.builder()
|
||||
.startOffset(0L)
|
||||
.retrieveTime(popTime)
|
||||
.invisibleTime(invisibleTime)
|
||||
.reviveQueueId(reviveQueueId)
|
||||
.topicType(ReceiptHandle.NORMAL_TOPIC)
|
||||
.brokerName(brokerName)
|
||||
.queueId(queueId)
|
||||
.offset(queueOffset)
|
||||
.build();
|
||||
RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, null);
|
||||
Mockito.when(ackMessageProcessorMock.processRequest(Mockito.any(SimpleChannelHandlerContext.class), Mockito.argThat(argument -> {
|
||||
boolean first = argument.getCode() == RequestCode.ACK_MESSAGE;
|
||||
boolean second = argument.readCustomHeader() instanceof AckMessageRequestHeader;
|
||||
return first && second;
|
||||
}))).thenReturn(remotingCommand);
|
||||
AckMessageRequestHeader requestHeader = new AckMessageRequestHeader();
|
||||
CompletableFuture<AckResult> future = localMessageService.ackMessage(proxyContext, handle, messageId,
|
||||
requestHeader, 1000L);
|
||||
AckResult ackResult = future.get();
|
||||
assertThat(ackResult.getStatus()).isEqualTo(AckStatus.OK);
|
||||
}
|
||||
|
||||
private MessageExt buildMessageExt(String topic, int queueId, long queueOffset) {
|
||||
MessageExt message1 = new MessageExt();
|
||||
message1.setTopic(topic);
|
||||
message1.setBody("body".getBytes(StandardCharsets.UTF_8));
|
||||
message1.setFlag(0);
|
||||
message1.setQueueId(queueId);
|
||||
message1.setQueueOffset(queueOffset);
|
||||
message1.setCommitLogOffset(1000L);
|
||||
message1.setSysFlag(0);
|
||||
message1.setBornTimestamp(0L);
|
||||
InetSocketAddress inetSocketAddress = new InetSocketAddress("127.0.0.1", 80);
|
||||
message1.setBornHost(inetSocketAddress);
|
||||
message1.setStoreHost(inetSocketAddress);
|
||||
message1.setReconsumeTimes(0);
|
||||
message1.setPreparedTransactionOffset(0L);
|
||||
message1.putUserProperty("K", "V");
|
||||
return message1;
|
||||
}
|
||||
|
||||
private void assertMessageExt(MessageExt messageExt1, MessageExt messageExt2) {
|
||||
assertThat(messageExt1.getBody()).isEqualTo(messageExt2.getBody());
|
||||
assertThat(messageExt1.getTopic()).isEqualTo(messageExt2.getTopic());
|
||||
assertThat(messageExt1.getQueueId()).isEqualTo(messageExt2.getQueueId());
|
||||
assertThat(messageExt1.getQueueOffset()).isEqualTo(messageExt2.getQueueOffset());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user