From bc86a2ea9a5ec59628f753c373bd586fa8339c15 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Mon, 23 May 2022 16:12:42 +0800 Subject: [PATCH] [ISSUE #3949] change sendResult to list of sendResult; return MULTIPLE_RESULTS when has multiple response code --- .../grpc/v2/consumer/AckMessageActivity.java | 18 ++++- .../grpc/v2/producer/SendMessageActivity.java | 80 ++++++++++++------- .../processor/DefaultMessagingProcessor.java | 2 +- .../proxy/processor/MessagingProcessor.java | 4 +- .../proxy/processor/ProducerProcessor.java | 22 ++++- .../message/AbstractMessageService.java | 44 ---------- .../message/ClusterMessageService.java | 15 ++-- .../service/message/LocalMessageService.java | 4 +- .../proxy/service/message/MessageService.java | 2 +- .../v2/consumer/AckMessageActivityTest.java | 2 +- .../v2/producer/SendMessageActivityTest.java | 71 ++++++++++------ .../processor/ProducerProcessorTest.java | 25 +++++- .../rocketmq/test/grpc/v2/GrpcBaseIT.java | 4 +- 13 files changed, 170 insertions(+), 123 deletions(-) delete mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/service/message/AbstractMessageService.java diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivity.java index 7226717a73..65c4f4fb3d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivity.java @@ -23,7 +23,9 @@ import apache.rocketmq.v2.AckMessageResultEntry; import apache.rocketmq.v2.Code; import io.grpc.Context; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; +import java.util.Set; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.consumer.AckResult; import org.apache.rocketmq.client.consumer.AckStatus; @@ -34,6 +36,7 @@ import org.apache.rocketmq.proxy.grpc.v2.common.GrpcClientSettingsManager; import org.apache.rocketmq.proxy.grpc.v2.common.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.common.ResponseBuilder; import org.apache.rocketmq.proxy.processor.MessagingProcessor; +import org.checkerframework.checker.units.qual.C; public class AckMessageActivity extends AbstractMessingActivity { @@ -56,13 +59,24 @@ public class AckMessageActivity extends AbstractMessingActivity { future.completeExceptionally(throwable); return; } + + Set responseCodes = new HashSet<>(); List entryList = new ArrayList<>(); for (CompletableFuture entryFuture : futures) { - entryFuture.thenAccept(entryList::add); + AckMessageResultEntry entryResult = entryFuture.join(); + responseCodes.add(entryResult.getStatus().getCode()); + entryList.add(entryResult); } AckMessageResponse.Builder responseBuilder = AckMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) .addAllEntries(entryList); + if (responseCodes.size() > 1) { + responseBuilder.setStatus(ResponseBuilder.buildStatus(Code.MULTIPLE_RESULTS, Code.MULTIPLE_RESULTS.name())); + } else if (responseCodes.size() == 1) { + Code code = responseCodes.stream().findAny().get(); + responseBuilder.setStatus(ResponseBuilder.buildStatus(code, code.name())); + } else { + responseBuilder.setStatus(ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "ack message result is empty")); + } future.complete(responseBuilder.build()); }); } catch (Throwable t) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java index 249630a224..998c0276fb 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java @@ -23,9 +23,8 @@ import apache.rocketmq.v2.MessageType; import apache.rocketmq.v2.Resource; import apache.rocketmq.v2.SendMessageRequest; import apache.rocketmq.v2.SendMessageResponse; -import apache.rocketmq.v2.SendReceipt; +import apache.rocketmq.v2.SendResultEntry; import apache.rocketmq.v2.SystemProperties; -import com.beust.jcommander.internal.Lists; import com.google.common.collect.Maps; import com.google.common.hash.Hashing; import com.google.protobuf.Duration; @@ -34,8 +33,10 @@ import com.google.protobuf.util.Durations; import com.google.protobuf.util.Timestamps; import io.grpc.Context; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.CompletableFuture; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.client.producer.SendResult; @@ -190,35 +191,54 @@ public class SendMessageActivity extends AbstractMessingActivity { } protected SendMessageResponse convertToSendMessageResponse(ProxyContext ctx, SendMessageRequest request, - SendResult result) { - switch (result.getSendStatus()) { - case FLUSH_DISK_TIMEOUT: - return SendMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.MASTER_PERSISTENCE_TIMEOUT, "send message failed, sendStatus=" + result.getSendStatus())) - .build(); - case FLUSH_SLAVE_TIMEOUT: - return SendMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.SLAVE_PERSISTENCE_TIMEOUT, "send message failed, sendStatus=" + result.getSendStatus())) - .build(); - case SLAVE_NOT_AVAILABLE: - return SendMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.HA_NOT_AVAILABLE, "send message failed, sendStatus=" + result.getSendStatus())) - .build(); - case SEND_OK: - List sendReceiptList = Lists.newArrayList(); - sendReceiptList.add(SendReceipt.newBuilder() - .setMessageId(StringUtils.defaultString(result.getMsgId())) - .setTransactionId(StringUtils.defaultString(result.getTransactionId())) - .build()); - return SendMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) - .addAllReceipts(sendReceiptList) - .build(); - default: - return SendMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "send message failed, sendStatus=" + result.getSendStatus())) - .build(); + List resultList) { + SendMessageResponse.Builder builder = SendMessageResponse.newBuilder(); + + Set responseCodes = new HashSet<>(); + for (SendResult result : resultList) { + SendResultEntry resultEntry; + switch (result.getSendStatus()) { + case FLUSH_DISK_TIMEOUT: + resultEntry = SendResultEntry.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.MASTER_PERSISTENCE_TIMEOUT, "send message failed, sendStatus=" + result.getSendStatus())) + .build(); + break; + case FLUSH_SLAVE_TIMEOUT: + resultEntry = SendResultEntry.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.SLAVE_PERSISTENCE_TIMEOUT, "send message failed, sendStatus=" + result.getSendStatus())) + .build(); + break; + case SLAVE_NOT_AVAILABLE: + resultEntry = SendResultEntry.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.HA_NOT_AVAILABLE, "send message failed, sendStatus=" + result.getSendStatus())) + .build(); + break; + case SEND_OK: + resultEntry = SendResultEntry.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) + .setOffset(result.getQueueOffset()) + .setMessageId(StringUtils.defaultString(result.getMsgId())) + .setTransactionId(StringUtils.defaultString(result.getTransactionId())) + .build(); + break; + default: + resultEntry = SendResultEntry.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "send message failed, sendStatus=" + result.getSendStatus())) + .build(); + break; + } + builder.addEntries(resultEntry); + responseCodes.add(resultEntry.getStatus().getCode()); } + if (responseCodes.size() > 1) { + builder.setStatus(ResponseBuilder.buildStatus(Code.MULTIPLE_RESULTS, Code.MULTIPLE_RESULTS.name())); + } else if (responseCodes.size() == 1) { + Code code = responseCodes.stream().findAny().get(); + builder.setStatus(ResponseBuilder.buildStatus(code, code.name())); + } else { + builder.setStatus(ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "send status is empty")); + } + return builder.build(); } protected static class SendMessageQueueSelector implements QueueSelector { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java index c159dec05d..44816f1adc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java @@ -125,7 +125,7 @@ public class DefaultMessagingProcessor extends AbstractStartAndShutdown implemen } @Override - public CompletableFuture sendMessage(ProxyContext ctx, QueueSelector queueSelector, + public CompletableFuture> sendMessage(ProxyContext ctx, QueueSelector queueSelector, String producerGroup, List msg, long timeoutMillis) { return this.producerProcessor.sendMessage(ctx, queueSelector, producerGroup, msg, timeoutMillis); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java index eee9f98731..1932c3f83b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java @@ -59,7 +59,7 @@ public interface MessagingProcessor extends StartAndShutdown { String topicName ) throws Exception; - default CompletableFuture sendMessage( + default CompletableFuture> sendMessage( ProxyContext ctx, QueueSelector queueSelector, String producerGroup, @@ -68,7 +68,7 @@ public interface MessagingProcessor extends StartAndShutdown { return sendMessage(ctx, queueSelector, producerGroup, msg, DEFAULT_TIMEOUT_MILLS); } - CompletableFuture sendMessage( + CompletableFuture> sendMessage( ProxyContext ctx, QueueSelector queueSelector, String producerGroup, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java index 086151e96b..a5baf4f9a5 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java @@ -19,7 +19,9 @@ package org.apache.rocketmq.proxy.processor; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; +import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.client.producer.SendResult; +import org.apache.rocketmq.client.producer.SendStatus; import org.apache.rocketmq.common.MixAll; import org.apache.rocketmq.common.consumer.ReceiptHandle; import org.apache.rocketmq.common.message.MessageAccessor; @@ -29,12 +31,14 @@ import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.ResponseCode; 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.common.ProxyException; import org.apache.rocketmq.proxy.common.ProxyExceptionCode; import org.apache.rocketmq.proxy.common.utils.FutureUtils; import org.apache.rocketmq.proxy.service.ServiceManager; import org.apache.rocketmq.proxy.service.route.SelectableMessageQueue; +import org.apache.rocketmq.proxy.service.transaction.TransactionId; import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class ProducerProcessor extends AbstractProcessor { @@ -47,9 +51,9 @@ public class ProducerProcessor extends AbstractProcessor { this.executor = executor; } - public CompletableFuture sendMessage(ProxyContext ctx, QueueSelector queueSelector, + public CompletableFuture> sendMessage(ProxyContext ctx, QueueSelector queueSelector, String producerGroup, List messageExtList, long timeoutMillis) { - CompletableFuture future = new CompletableFuture<>(); + CompletableFuture> future = new CompletableFuture<>(); try { String topic = messageExtList.get(0).getTopic(); SelectableMessageQueue messageQueue = queueSelector.select(ctx, @@ -65,7 +69,19 @@ public class ProducerProcessor extends AbstractProcessor { messageQueue, messageExtList, requestHeader, - timeoutMillis); + timeoutMillis) + .thenApplyAsync(sendResultList -> { + for (SendResult sendResult : sendResultList) { + int tranType = MessageSysFlag.getTransactionValue(requestHeader.getSysFlag()); + if (SendStatus.SEND_OK.equals(sendResult.getSendStatus()) && + tranType == MessageSysFlag.TRANSACTION_PREPARED_TYPE && + StringUtils.isNotBlank(sendResult.getTransactionId())) { + TransactionId transactionId = TransactionId.genByBrokerTransactionId(messageQueue.getBrokerName(), sendResult); + sendResult.setTransactionId(transactionId.getProxyTransactionId()); + } + } + return sendResultList; + }, this.executor); } catch (Throwable t) { future.completeExceptionally(t); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/AbstractMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/AbstractMessageService.java deleted file mode 100644 index b7aa35ae84..0000000000 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/AbstractMessageService.java +++ /dev/null @@ -1,44 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.apache.rocketmq.proxy.service.message; - -import java.util.concurrent.CompletableFuture; -import org.apache.commons.lang3.StringUtils; -import org.apache.rocketmq.client.producer.SendResult; -import org.apache.rocketmq.client.producer.SendStatus; -import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; -import org.apache.rocketmq.common.sysflag.MessageSysFlag; -import org.apache.rocketmq.proxy.service.transaction.TransactionId; - -public abstract class AbstractMessageService implements MessageService { - - protected CompletableFuture processSendMessageResponseFuture( - String brokerName, - SendMessageRequestHeader requestHeader, - CompletableFuture future) { - return future.thenApply(sendResult -> { - int tranType = MessageSysFlag.getTransactionValue(requestHeader.getSysFlag()); - if (SendStatus.SEND_OK.equals(sendResult.getSendStatus()) && - tranType == MessageSysFlag.TRANSACTION_PREPARED_TYPE && - StringUtils.isNotBlank(sendResult.getTransactionId())) { - TransactionId transactionId = TransactionId.genByBrokerTransactionId(brokerName, sendResult); - sendResult.setTransactionId(transactionId.getProxyTransactionId()); - } - return sendResult; - }); - } -} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/ClusterMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/ClusterMessageService.java index b48cade6da..4fc8fdad9e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/ClusterMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/ClusterMessageService.java @@ -16,6 +16,7 @@ */ package org.apache.rocketmq.proxy.service.message; +import com.google.common.collect.Lists; import java.util.List; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.consumer.AckResult; @@ -40,7 +41,7 @@ import org.apache.rocketmq.proxy.service.transaction.TransactionId; import org.apache.rocketmq.remoting.exception.RemotingException; import org.apache.rocketmq.remoting.protocol.RemotingCommand; -public class ClusterMessageService extends AbstractMessageService { +public class ClusterMessageService implements MessageService { private final TopicRouteService topicRouteService; private final MQClientAPIFactory mqClientAPIFactory; @@ -51,19 +52,21 @@ public class ClusterMessageService extends AbstractMessageService { } @Override - public CompletableFuture sendMessage(ProxyContext ctx, SelectableMessageQueue messageQueue, + public CompletableFuture> sendMessage(ProxyContext ctx, SelectableMessageQueue messageQueue, List msgList, SendMessageRequestHeader requestHeader, long timeoutMillis) { - CompletableFuture future; + CompletableFuture> future; if (msgList.size() == 1) { future = this.mqClientAPIFactory.getClient().sendMessageAsync( messageQueue.getBrokerAddr(), - messageQueue.getBrokerName(), msgList.get(0), requestHeader, timeoutMillis); + messageQueue.getBrokerName(), msgList.get(0), requestHeader, timeoutMillis) + .thenApply(Lists::newArrayList); } else { future = this.mqClientAPIFactory.getClient().sendMessageAsync( messageQueue.getBrokerAddr(), - messageQueue.getBrokerName(), msgList, requestHeader, timeoutMillis); + messageQueue.getBrokerName(), msgList, requestHeader, timeoutMillis) + .thenApply(Lists::newArrayList); } - return processSendMessageResponseFuture(messageQueue.getBrokerName(), requestHeader, future); + return future; } @Override 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 e8ff4a353c..055f9e85e1 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 @@ -36,7 +36,7 @@ import org.apache.rocketmq.proxy.service.transaction.TransactionId; import org.apache.rocketmq.remoting.RPCHook; import org.apache.rocketmq.remoting.protocol.RemotingCommand; -public class LocalMessageService extends AbstractMessageService { +public class LocalMessageService implements MessageService { private BrokerController brokerController; @@ -44,7 +44,7 @@ public class LocalMessageService extends AbstractMessageService { this.brokerController = brokerController; } - @Override public CompletableFuture sendMessage(ProxyContext ctx, SelectableMessageQueue messageQueue, + @Override public CompletableFuture> sendMessage(ProxyContext ctx, SelectableMessageQueue messageQueue, List msgList, SendMessageRequestHeader requestHeader, long timeoutMillis) { return null; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/MessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/MessageService.java index 4d42c317cb..35f0ea147b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/MessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/MessageService.java @@ -38,7 +38,7 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand; public interface MessageService { - CompletableFuture sendMessage( + CompletableFuture> sendMessage( ProxyContext ctx, SelectableMessageQueue messageQueue, List msgList, diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivityTest.java index 353b145fba..6ca311d6ff 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivityTest.java @@ -82,7 +82,7 @@ public class AckMessageActivityTest extends BaseActivityTest { .build() ).get(); - assertEquals(Code.OK, response.getStatus().getCode()); + assertEquals(Code.MULTIPLE_RESULTS, response.getStatus().getCode()); assertEquals(3, response.getEntriesCount()); assertEquals(Code.RECEIPT_HANDLE_EXPIRED, response.getEntries(0).getStatus().getCode()); assertEquals(Code.OK, response.getEntries(1).getStatus().getCode()); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivityTest.java index 682cf32b6b..a23fdd5612 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivityTest.java @@ -70,7 +70,7 @@ public class SendMessageActivityTest extends BaseActivityTest { sendResult.setSendStatus(SendStatus.SEND_OK); sendResult.setMsgId(msgId); when(this.messagingProcessor.sendMessage(any(), any(), anyString(), any())) - .thenReturn(CompletableFuture.completedFuture(sendResult)); + .thenReturn(CompletableFuture.completedFuture(Lists.newArrayList(sendResult))); SendMessageResponse response = this.sendMessageActivity.sendMessage( createContext(), @@ -92,43 +92,62 @@ public class SendMessageActivityTest extends BaseActivityTest { ).get(); assertEquals(Code.OK, response.getStatus().getCode()); - assertEquals(msgId, response.getReceipts(0).getMessageId()); + assertEquals(msgId, response.getEntries(0).getMessageId()); } @Test public void testConvertToSendMessageResponse() { - assertEquals( - Code.MASTER_PERSISTENCE_TIMEOUT, - this.sendMessageActivity.convertToSendMessageResponse( + { + SendMessageResponse response = this.sendMessageActivity.convertToSendMessageResponse( ProxyContext.create(), SendMessageRequest.newBuilder().build(), - new SendResult(SendStatus.FLUSH_DISK_TIMEOUT, null, null, null, 0) - ).getStatus().getCode() - ); - assertEquals( - Code.SLAVE_PERSISTENCE_TIMEOUT, - this.sendMessageActivity.convertToSendMessageResponse( + Lists.newArrayList(new SendResult(SendStatus.FLUSH_DISK_TIMEOUT, null, null, null, 0)) + ); + assertEquals(Code.MASTER_PERSISTENCE_TIMEOUT, response.getStatus().getCode()); + assertEquals(Code.MASTER_PERSISTENCE_TIMEOUT, response.getEntries(0).getStatus().getCode()); + } + + { + SendMessageResponse response = this.sendMessageActivity.convertToSendMessageResponse( ProxyContext.create(), SendMessageRequest.newBuilder().build(), - new SendResult(SendStatus.FLUSH_SLAVE_TIMEOUT, null, null, null, 0) - ).getStatus().getCode() - ); - assertEquals( - Code.HA_NOT_AVAILABLE, - this.sendMessageActivity.convertToSendMessageResponse( + Lists.newArrayList(new SendResult(SendStatus.FLUSH_SLAVE_TIMEOUT, null, null, null, 0)) + ); + assertEquals(Code.SLAVE_PERSISTENCE_TIMEOUT, response.getStatus().getCode()); + assertEquals(Code.SLAVE_PERSISTENCE_TIMEOUT, response.getEntries(0).getStatus().getCode()); + } + + { + SendMessageResponse response = this.sendMessageActivity.convertToSendMessageResponse( ProxyContext.create(), SendMessageRequest.newBuilder().build(), - new SendResult(SendStatus.SLAVE_NOT_AVAILABLE, null, null, null, 0) - ).getStatus().getCode() - ); - assertEquals( - Code.OK, - this.sendMessageActivity.convertToSendMessageResponse( + Lists.newArrayList(new SendResult(SendStatus.SLAVE_NOT_AVAILABLE, null, null, null, 0)) + ); + assertEquals(Code.HA_NOT_AVAILABLE, response.getStatus().getCode()); + assertEquals(Code.HA_NOT_AVAILABLE, response.getEntries(0).getStatus().getCode()); + } + + { + SendMessageResponse response = this.sendMessageActivity.convertToSendMessageResponse( ProxyContext.create(), SendMessageRequest.newBuilder().build(), - new SendResult(SendStatus.SEND_OK, null, null, null, 0) - ).getStatus().getCode() - ); + Lists.newArrayList(new SendResult(SendStatus.SEND_OK, null, null, null, 0)) + ); + assertEquals(Code.OK, response.getStatus().getCode()); + assertEquals(Code.OK, response.getEntries(0).getStatus().getCode()); + } + + { + SendMessageResponse response = this.sendMessageActivity.convertToSendMessageResponse( + ProxyContext.create(), + SendMessageRequest.newBuilder().build(), + Lists.newArrayList( + new SendResult(SendStatus.SEND_OK, null, null, null, 0), + new SendResult(SendStatus.SLAVE_NOT_AVAILABLE, null, null, null, 0) + ) + ); + assertEquals(Code.MULTIPLE_RESULTS, response.getStatus().getCode()); + } } @Test(expected = GrpcProxyException.class) 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 9f138fd562..4bac51a5d8 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 @@ -22,19 +22,24 @@ import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executors; import org.apache.rocketmq.client.producer.SendResult; +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; 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; +import org.assertj.core.util.Lists; import org.junit.Before; import org.junit.Test; import org.mockito.ArgumentCaptor; @@ -62,18 +67,27 @@ public class ProducerProcessorTest extends BaseProcessorTest { @Test public void testSendMessage() throws Throwable { + String txId = MessageClientIDSetter.createUniqID(); + String msgId = MessageClientIDSetter.createUniqID(); + + SendResult sendResult = new SendResult(); + sendResult.setSendStatus(SendStatus.SEND_OK); + sendResult.setTransactionId(txId); + sendResult.setMsgId(msgId); ArgumentCaptor requestHeaderArgumentCaptor = ArgumentCaptor.forClass(SendMessageRequestHeader.class); when(this.messageService.sendMessage(any(), any(), any(), requestHeaderArgumentCaptor.capture(), anyLong())) - .thenReturn(CompletableFuture.completedFuture(mock(SendResult.class))); + .thenReturn(CompletableFuture.completedFuture(Lists.newArrayList(sendResult))); List messageExtList = new ArrayList<>(); MessageExt messageExt = createMessageExt(MixAll.getRetryTopic(CONSUMER_GROUP), "tag", 0, 0); + messageExt.setSysFlag(MessageSysFlag.TRANSACTION_PREPARED_TYPE); MessageAccessor.putProperty(messageExt, MessageConst.PROPERTY_RECONSUME_TIME, "1"); MessageAccessor.putProperty(messageExt, MessageConst.PROPERTY_MAX_RECONSUME_TIMES, "16"); messageExtList.add(messageExt); SelectableMessageQueue messageQueue = mock(SelectableMessageQueue.class); + when(messageQueue.getBrokerName()).thenReturn("mockBroker"); - SendResult sendResult = this.producerProcessor.sendMessage( + List sendResultList = this.producerProcessor.sendMessage( createContext(), (ctx, messageQueueView) -> messageQueue, PRODUCER_GROUP, @@ -81,7 +95,12 @@ public class ProducerProcessorTest extends BaseProcessorTest { 3000 ).get(); - assertNotNull(sendResult); + assertNotNull(sendResultList); + TransactionId transactionId = TransactionId.decode(sendResultList.get(0).getTransactionId()); + assertNotNull(transactionId); + assertEquals(txId, transactionId.getBrokerTransactionId()); + assertEquals("mockBroker", transactionId.getBrokerName()); + SendMessageRequestHeader requestHeader = requestHeaderArgumentCaptor.getValue(); assertEquals(PRODUCER_GROUP, requestHeader.getProducerGroup()); assertEquals(MixAll.getRetryTopic(CONSUMER_GROUP), requestHeader.getTopic()); diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseIT.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseIT.java index efbda29d54..38873fa9e8 100644 --- a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseIT.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseIT.java @@ -336,7 +336,7 @@ public class GrpcBaseIT extends BaseConf { AckMessageResponse ackMessageResponse = blockingStub.ackMessage(buildAckMessageRequest(topic, group, AckMessageEntry.newBuilder().setMessageId(messageId).setReceiptHandle(ackHandles.get(0)).build(), AckMessageEntry.newBuilder().setMessageId(messageId).setReceiptHandle(ackHandles.get(1)).build())); - assertThat(ackMessageResponse.getStatus().getCode()).isEqualTo(Code.OK); + assertThat(ackMessageResponse.getStatus().getCode()).isEqualTo(Code.MULTIPLE_RESULTS); int okNum = 0; int expireNum = 0; for (AckMessageResultEntry entry : ackMessageResponse.getEntriesList()) { @@ -550,7 +550,7 @@ public class GrpcBaseIT extends BaseConf { public void assertSendMessage(SendMessageResponse response, String messageId) { assertThat(response.getStatus() .getCode()).isEqualTo(Code.OK); - assertThat(response.getReceipts(0).getMessageId()).isEqualTo(messageId); + assertThat(response.getEntries(0).getMessageId()).isEqualTo(messageId); } public Message assertAndGetReceiveMessage(List response, String messageId) {