mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 13:49:50 +08:00
[ISSUE #3949] change sendResult to list of sendResult; return MULTIPLE_RESULTS when has multiple response code
This commit is contained in:
+16
-2
@@ -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<Code> responseCodes = new HashSet<>();
|
||||
List<AckMessageResultEntry> entryList = new ArrayList<>();
|
||||
for (CompletableFuture<AckMessageResultEntry> 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) {
|
||||
|
||||
+50
-30
@@ -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<SendReceipt> 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<SendResult> resultList) {
|
||||
SendMessageResponse.Builder builder = SendMessageResponse.newBuilder();
|
||||
|
||||
Set<Code> 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 {
|
||||
|
||||
+1
-1
@@ -125,7 +125,7 @@ public class DefaultMessagingProcessor extends AbstractStartAndShutdown implemen
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<SendResult> sendMessage(ProxyContext ctx, QueueSelector queueSelector,
|
||||
public CompletableFuture<List<SendResult>> sendMessage(ProxyContext ctx, QueueSelector queueSelector,
|
||||
String producerGroup, List<MessageExt> msg, long timeoutMillis) {
|
||||
return this.producerProcessor.sendMessage(ctx, queueSelector, producerGroup, msg, timeoutMillis);
|
||||
}
|
||||
|
||||
@@ -59,7 +59,7 @@ public interface MessagingProcessor extends StartAndShutdown {
|
||||
String topicName
|
||||
) throws Exception;
|
||||
|
||||
default CompletableFuture<SendResult> sendMessage(
|
||||
default CompletableFuture<List<SendResult>> 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<SendResult> sendMessage(
|
||||
CompletableFuture<List<SendResult>> sendMessage(
|
||||
ProxyContext ctx,
|
||||
QueueSelector queueSelector,
|
||||
String producerGroup,
|
||||
|
||||
@@ -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<SendResult> sendMessage(ProxyContext ctx, QueueSelector queueSelector,
|
||||
public CompletableFuture<List<SendResult>> sendMessage(ProxyContext ctx, QueueSelector queueSelector,
|
||||
String producerGroup, List<MessageExt> messageExtList, long timeoutMillis) {
|
||||
CompletableFuture<SendResult> future = new CompletableFuture<>();
|
||||
CompletableFuture<List<SendResult>> 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);
|
||||
}
|
||||
|
||||
-44
@@ -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<SendResult> processSendMessageResponseFuture(
|
||||
String brokerName,
|
||||
SendMessageRequestHeader requestHeader,
|
||||
CompletableFuture<SendResult> 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;
|
||||
});
|
||||
}
|
||||
}
|
||||
+9
-6
@@ -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<SendResult> sendMessage(ProxyContext ctx, SelectableMessageQueue messageQueue,
|
||||
public CompletableFuture<List<SendResult>> sendMessage(ProxyContext ctx, SelectableMessageQueue messageQueue,
|
||||
List<? extends Message> msgList, SendMessageRequestHeader requestHeader, long timeoutMillis) {
|
||||
CompletableFuture<SendResult> future;
|
||||
CompletableFuture<List<SendResult>> 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
|
||||
|
||||
+2
-2
@@ -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<SendResult> sendMessage(ProxyContext ctx, SelectableMessageQueue messageQueue,
|
||||
@Override public CompletableFuture<List<SendResult>> sendMessage(ProxyContext ctx, SelectableMessageQueue messageQueue,
|
||||
List<? extends Message> msgList, SendMessageRequestHeader requestHeader, long timeoutMillis) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -38,7 +38,7 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
|
||||
public interface MessageService {
|
||||
|
||||
CompletableFuture<SendResult> sendMessage(
|
||||
CompletableFuture<List<SendResult>> sendMessage(
|
||||
ProxyContext ctx,
|
||||
SelectableMessageQueue messageQueue,
|
||||
List<? extends Message> msgList,
|
||||
|
||||
+1
-1
@@ -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());
|
||||
|
||||
+45
-26
@@ -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)
|
||||
|
||||
@@ -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<SendMessageRequestHeader> 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<MessageExt> 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<SendResult> 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());
|
||||
|
||||
@@ -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<ReceiveMessageResponse> response, String messageId) {
|
||||
|
||||
Reference in New Issue
Block a user