diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/ProxyExceptionCode.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/ProxyExceptionCode.java index a297e50aed..2fbe49ef43 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/ProxyExceptionCode.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/ProxyExceptionCode.java @@ -22,4 +22,5 @@ public enum ProxyExceptionCode { INVALID_BROKER_NAME, INVALID_RECEIPT_HANDLE, ILLEGAL_MESSAGE, + INTERNAL_SERVER_ERROR, } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcProxyException.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcProxyException.java index 7cd5f5a43d..ed1eaf1197 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcProxyException.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcProxyException.java @@ -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) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java index c18e8fe4d1..4e893c384f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java @@ -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 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 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 popMessage(ProxyContext ctx, SelectableMessageQueue messageQueue, PopMessageRequestHeader requestHeader, long timeoutMillis) { - return null; + RemotingCommand request = LocalRemotingCommand.createRequestCommand(RequestCode.POP_MESSAGE, requestHeader); + CompletableFuture 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 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 startOffsetInfo = null; + Map> msgOffsetInfo = null; + Map 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()); + } + // + Map> 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 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 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 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 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 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; + }); } } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java index 4bac51a5d8..eac0abe553 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java @@ -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; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/message/LocalMessageServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/message/LocalMessageServiceTest.java index 7053e08e26..64ba5cba71 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/message/LocalMessageServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/message/LocalMessageServiceTest.java @@ -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 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 messageExtList = new ArrayList<>(); + List 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 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 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 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 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()); + } } \ No newline at end of file