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 49edbcaba8..249630a224 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 @@ -39,7 +39,6 @@ import java.util.Map; 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.message.MessageAccessor; import org.apache.rocketmq.common.message.MessageConst; import org.apache.rocketmq.common.message.MessageExt; @@ -192,21 +191,34 @@ public class SendMessageActivity extends AbstractMessingActivity { protected SendMessageResponse convertToSendMessageResponse(ProxyContext ctx, SendMessageRequest request, SendResult result) { - if (result.getSendStatus() != SendStatus.SEND_OK) { - return SendMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "send message failed, sendStatus=" + result.getSendStatus())) - .build(); + 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 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(); } protected static class SendMessageQueueSelector implements QueueSelector { 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 new file mode 100644 index 0000000000..682cf32b6b --- /dev/null +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivityTest.java @@ -0,0 +1,223 @@ +/* + * 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.grpc.v2.producer; + +import apache.rocketmq.v2.Code; +import apache.rocketmq.v2.Encoding; +import apache.rocketmq.v2.Message; +import apache.rocketmq.v2.MessageType; +import apache.rocketmq.v2.Resource; +import apache.rocketmq.v2.SendMessageRequest; +import apache.rocketmq.v2.SendMessageResponse; +import apache.rocketmq.v2.SystemProperties; +import com.google.protobuf.ByteString; +import com.google.protobuf.util.Durations; +import com.google.protobuf.util.Timestamps; +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.message.MessageClientIDSetter; +import org.apache.rocketmq.common.message.MessageConst; +import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.common.sysflag.MessageSysFlag; +import org.apache.rocketmq.proxy.common.ProxyContext; +import org.apache.rocketmq.proxy.grpc.v2.BaseActivityTest; +import org.apache.rocketmq.proxy.grpc.v2.common.GrpcProxyException; +import org.apache.rocketmq.remoting.common.RemotingUtil; +import org.assertj.core.util.Lists; +import org.junit.Before; +import org.junit.Test; + +import static org.junit.Assert.assertEquals; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.when; + +public class SendMessageActivityTest extends BaseActivityTest { + + private static final String TOPIC = "topic"; + private static final String CONSUMER_GROUP = "consumerGroup"; + + private SendMessageActivity sendMessageActivity; + + @Before + public void before() throws Throwable { + super.before(); + this.sendMessageActivity = new SendMessageActivity(this.messagingProcessor, this.grpcClientSettingsManager); + } + + @Test + public void sendMessage() throws Exception { + String msgId = MessageClientIDSetter.createUniqID(); + + SendResult sendResult = new SendResult(); + sendResult.setSendStatus(SendStatus.SEND_OK); + sendResult.setMsgId(msgId); + when(this.messagingProcessor.sendMessage(any(), any(), anyString(), any())) + .thenReturn(CompletableFuture.completedFuture(sendResult)); + + SendMessageResponse response = this.sendMessageActivity.sendMessage( + createContext(), + SendMessageRequest.newBuilder() + .addMessages(Message.newBuilder() + .setTopic(Resource.newBuilder() + .setName(TOPIC) + .build()) + .setSystemProperties(SystemProperties.newBuilder() + .setMessageId(msgId) + .setQueueId(0) + .setMessageType(MessageType.NORMAL) + .setBornTimestamp(Timestamps.fromMillis(System.currentTimeMillis())) + .setBornHost(StringUtils.defaultString(RemotingUtil.getLocalAddress(), "127.0.0.1:1234")) + .build()) + .setBody(ByteString.copyFromUtf8("123")) + .build()) + .build() + ).get(); + + assertEquals(Code.OK, response.getStatus().getCode()); + assertEquals(msgId, response.getReceipts(0).getMessageId()); + } + + @Test + public void testConvertToSendMessageResponse() { + assertEquals( + Code.MASTER_PERSISTENCE_TIMEOUT, + 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( + 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( + ProxyContext.create(), + SendMessageRequest.newBuilder().build(), + new SendResult(SendStatus.SLAVE_NOT_AVAILABLE, null, null, null, 0) + ).getStatus().getCode() + ); + assertEquals( + Code.OK, + this.sendMessageActivity.convertToSendMessageResponse( + ProxyContext.create(), + SendMessageRequest.newBuilder().build(), + new SendResult(SendStatus.SEND_OK, null, null, null, 0) + ).getStatus().getCode() + ); + } + + @Test(expected = GrpcProxyException.class) + public void testBuildErrorMessage() { + this.sendMessageActivity.buildMessage(null, + Lists.newArrayList( + Message.newBuilder() + .setTopic(Resource.newBuilder() + .setName(TOPIC) + .build()) + .setSystemProperties(SystemProperties.newBuilder() + .setMessageId(MessageClientIDSetter.createUniqID()) + .setQueueId(0) + .setMessageType(MessageType.NORMAL) + .setBornTimestamp(Timestamps.fromMillis(System.currentTimeMillis())) + .setBornHost(StringUtils.defaultString(RemotingUtil.getLocalAddress(), "127.0.0.1:1234")) + .build()) + .setBody(ByteString.copyFromUtf8("123")) + .build(), + Message.newBuilder() + .setTopic(Resource.newBuilder() + .setName(TOPIC + 2) + .build()) + .setSystemProperties(SystemProperties.newBuilder() + .setMessageId(MessageClientIDSetter.createUniqID()) + .setQueueId(0) + .setMessageType(MessageType.NORMAL) + .setBornTimestamp(Timestamps.fromMillis(System.currentTimeMillis())) + .setBornHost(StringUtils.defaultString(RemotingUtil.getLocalAddress(), "127.0.0.1:1234")) + .build()) + .setBody(ByteString.copyFromUtf8("123")) + .build() + ), + Resource.newBuilder().setName(TOPIC).build()); + } + + @Test + public void testBuildMessage() { + long deliveryTime = System.currentTimeMillis(); + String msgId = MessageClientIDSetter.createUniqID(); + + MessageExt messageExt = this.sendMessageActivity.buildMessage(null, + Lists.newArrayList( + Message.newBuilder() + .setTopic(Resource.newBuilder() + .setName(TOPIC) + .build()) + .setSystemProperties(SystemProperties.newBuilder() + .setMessageId(msgId) + .setQueueId(0) + .setMessageType(MessageType.DELAY) + .setDeliveryTimestamp(Timestamps.fromMillis(deliveryTime)) + .setBornTimestamp(Timestamps.fromMillis(System.currentTimeMillis())) + .setBornHost(StringUtils.defaultString(RemotingUtil.getLocalAddress(), "127.0.0.1:1234")) + .build()) + .setBody(ByteString.copyFromUtf8("123")) + .build() + ), + Resource.newBuilder().setName(TOPIC).build()).get(0); + + assertEquals(MessageClientIDSetter.getUniqID(messageExt), msgId); + assertEquals(String.valueOf(deliveryTime), messageExt.getProperty(MessageConst.PROPERTY_TIMER_DELIVER_MS)); + } + + @Test + public void testTxMessage() { + String msgId = MessageClientIDSetter.createUniqID(); + + MessageExt messageExt = this.sendMessageActivity.buildMessage(null, + Lists.newArrayList( + Message.newBuilder() + .setTopic(Resource.newBuilder() + .setName(TOPIC) + .build()) + .setSystemProperties(SystemProperties.newBuilder() + .setMessageId(msgId) + .setQueueId(0) + .setMessageType(MessageType.TRANSACTION) + .setOrphanedTransactionRecoveryDuration(Durations.fromSeconds(30)) + .setBodyEncoding(Encoding.GZIP) + .setBornTimestamp(Timestamps.fromMillis(System.currentTimeMillis())) + .setBornHost(StringUtils.defaultString(RemotingUtil.getLocalAddress(), "127.0.0.1:1234")) + .build()) + .setBody(ByteString.copyFromUtf8("123")) + .build() + ), + Resource.newBuilder().setName(TOPIC).build()).get(0); + + assertEquals(MessageClientIDSetter.getUniqID(messageExt), msgId); + assertEquals(MessageSysFlag.TRANSACTION_PREPARED_TYPE | MessageSysFlag.COMPRESSED_FLAG, messageExt.getSysFlag()); + } +} \ No newline at end of file