This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new 7ce9a6821 feat(topic): support cloud test-message sending (#4496)
7ce9a6821 is described below
commit 7ce9a68215c4cd1c7025013d02135b43b6aa6ba1
Author: zmuxuny <[email protected]>
AuthorDate: Mon Sep 21 12:25:51 2026 +0800
feat(topic): support cloud test-message sending (#4496)
Sending a test message was Apache-only: MetadataService.sendMessage
rejected any
other vendor with 501 before reaching a provider. It now dispatches to the
resolved instance provider, and a new MESSAGE_SEND capability declares which
providers implement it — Aliyun and Tencent through their vendor SDKs,
Apache
through the existing admin client path.
On the topic page the test-message action is offered for cloud instances as
well,
restricted to NORMAL topics, which is what the vendor APIs accept here.
---
.../studio/instance/topic/MetadataService.java | 5 +--
.../studio/provider/InstanceCapability.java | 1 +
.../rocketmq/studio/provider/InstanceProvider.java | 6 +++
.../provider/alibaba/AliyunInstanceProvider.java | 30 +++++++++++++
.../provider/apache/ApacheInstanceProvider.java | 8 ++++
.../provider/tencent/TencentInstanceProvider.java | 29 ++++++++++++
.../studio/instance/topic/MetadataServiceTest.java | 18 ++++----
.../alibaba/AliyunInstanceProviderTest.java | 51 ++++++++++++++++++++++
.../apache/ApacheInstanceProviderTest.java | 1 +
.../tencent/TencentInstanceProviderTest.java | 44 +++++++++++++++++++
web/src/api/instance.ts | 1 +
.../pages/instance/__tests__/TopicPage.test.tsx | 26 +++++++++++
web/src/pages/instance/topic.tsx | 3 +-
13 files changed, 209 insertions(+), 14 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
index d45b4b513..e2c3e6394 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
@@ -280,10 +280,7 @@ public class MetadataService {
public SendMessageVO sendMessage(SendMessageDTO request) {
requireSendMessageRequest(request);
- if (resolve(request.getInstanceId()).vendor() !=
InstanceVendor.APACHE) {
- throw new BusinessException(501, "Sending messages is not
supported for cloud instances");
- }
- return adminClient.sendMessage(request);
+ return resolve(request.getInstanceId()).sendMessage(request);
}
/**
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceCapability.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceCapability.java
index 30fed36d2..a98a5433d 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceCapability.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceCapability.java
@@ -24,6 +24,7 @@ public enum InstanceCapability {
CONSUMER_GROUP_MANAGEMENT,
MESSAGE_QUERY,
MESSAGE_TRACE,
+ MESSAGE_SEND,
ACL_MANAGEMENT,
DLQ_MANAGEMENT
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
index e89882f0d..fddd62c1f 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
@@ -30,6 +30,8 @@ import
org.apache.rocketmq.studio.instance.message.DirectConsumeMessageDTO;
import
org.apache.rocketmq.studio.instance.message.DirectConsumeMessageResultVO;
import org.apache.rocketmq.studio.instance.message.MessageQueryResult;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageDTO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
@@ -175,6 +177,10 @@ public interface InstanceProvider {
TraceRecordVO getMessageTrace(String instanceId, String msgId, String
topic);
+ default SendMessageVO sendMessage(SendMessageDTO request) {
+ throw new UnsupportedOperationException("Message sending is not
supported");
+ }
+
default DirectConsumeMessageResultVO
consumeMessageDirectly(DirectConsumeMessageDTO request) {
throw new UnsupportedOperationException("Direct message consumption is
not supported");
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
index 876e42261..c82d21027 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
@@ -43,6 +43,9 @@ import
com.aliyun.sdk.service.rocketmq20220801.models.ListTopicsResponse;
import com.aliyun.sdk.service.rocketmq20220801.models.ListTopicsResponseBody;
import
com.aliyun.sdk.service.rocketmq20220801.models.ResetConsumeOffsetRequest;
import com.aliyun.sdk.service.rocketmq20220801.models.UpdateTopicRequest;
+import com.aliyun.sdk.service.rocketmq20220801.models.VerifySendMessageRequest;
+import
com.aliyun.sdk.service.rocketmq20220801.models.VerifySendMessageResponse;
+import
com.aliyun.sdk.service.rocketmq20220801.models.VerifySendMessageResponseBody;
import org.springframework.util.StringUtils;
import org.apache.rocketmq.studio.common.domain.PageResult;
@@ -57,6 +60,8 @@ import
org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
import org.apache.rocketmq.studio.instance.message.MessageQueryResult;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageDTO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.provider.InstanceProvider;
@@ -103,6 +108,7 @@ public class AliyunInstanceProvider implements
InstanceProvider {
InstanceCapability.CONSUMER_GROUP_MANAGEMENT,
InstanceCapability.MESSAGE_QUERY,
InstanceCapability.MESSAGE_TRACE,
+ InstanceCapability.MESSAGE_SEND,
InstanceCapability.ACL_MANAGEMENT);
}
@@ -453,6 +459,30 @@ public class AliyunInstanceProvider implements
InstanceProvider {
clientFactory.call(ctx.credentialId(), ctx.regionId(), client ->
client.resetConsumeOffset(request));
}
+ @Override
+ public SendMessageVO sendMessage(SendMessageDTO request) {
+ Context ctx = resolve(request.getInstanceId());
+ VerifySendMessageRequest apiRequest =
VerifySendMessageRequest.builder()
+ .instanceId(ctx.cloudInstanceId())
+ .topicName(request.getTopic())
+ .message(request.getBody())
+ .messageTag(request.getTag())
+ .messageKey(request.getKey())
+ .messageGroup(request.getMessageGroup())
+ .deliveryTimeStamp(request.getDeliveryTimestamp())
+ .userProperties(request.getProperties())
+ .build();
+ VerifySendMessageResponse response =
clientFactory.call(ctx.credentialId(), ctx.regionId(),
+ client -> client.verifySendMessage(apiRequest));
+ VerifySendMessageResponseBody body = response == null ? null :
response.getBody();
+ if (body == null || !Boolean.TRUE.equals(body.getSuccess()) ||
!StringUtils.hasText(body.getData())) {
+ String detail = body == null ? "empty response" :
(StringUtils.hasText(body.getMessage())
+ ? body.getMessage() : body.getCode());
+ throw new BusinessException(502, "Aliyun test message send failed:
" + detail);
+ }
+ return
SendMessageVO.builder().msgId(body.getData()).sendTime(System.currentTimeMillis()).build();
+ }
+
@Override
public List<MessageRecordVO> queryMessages(String instanceId, String
topic, String msgId,
String tag, String key, Long
startTime, Long endTime) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
index 213cb54f9..abcd75d13 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
@@ -29,6 +29,8 @@ import
org.apache.rocketmq.studio.instance.message.DirectConsumeMessageDTO;
import
org.apache.rocketmq.studio.instance.message.DirectConsumeMessageResultVO;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageDTO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
@@ -65,6 +67,7 @@ public class ApacheInstanceProvider implements
InstanceProvider {
InstanceCapability.CONSUMER_GROUP_MANAGEMENT,
InstanceCapability.MESSAGE_QUERY,
InstanceCapability.MESSAGE_TRACE,
+ InstanceCapability.MESSAGE_SEND,
InstanceCapability.ACL_MANAGEMENT,
InstanceCapability.DLQ_MANAGEMENT);
}
@@ -183,6 +186,11 @@ public class ApacheInstanceProvider implements
InstanceProvider {
return messageProvider.getMessageTrace(instanceId, msgId, topic);
}
+ @Override
+ public SendMessageVO sendMessage(SendMessageDTO request) {
+ return adminClient.sendMessage(request);
+ }
+
@Override
public DirectConsumeMessageResultVO
consumeMessageDirectly(DirectConsumeMessageDTO request) {
return messageProvider.consumeMessageDirectly(request);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
index 9730693a2..9685c0c99 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
@@ -42,6 +42,8 @@ import
com.tencentcloudapi.trocket.v20230308.models.DescribeTopicResponse;
import com.tencentcloudapi.trocket.v20230308.models.Filter;
import com.tencentcloudapi.trocket.v20230308.models.ModifyTopicRequest;
import
com.tencentcloudapi.trocket.v20230308.models.ResetConsumerGroupOffsetRequest;
+import com.tencentcloudapi.trocket.v20230308.models.SendMessageRequest;
+import com.tencentcloudapi.trocket.v20230308.models.SendMessageResponse;
import com.tencentcloudapi.trocket.v20230308.models.SubscriptionData;
import com.tencentcloudapi.trocket.v20230308.models.TopicItem;
import org.apache.rocketmq.studio.common.domain.PageResult;
@@ -63,6 +65,8 @@ import
org.apache.rocketmq.studio.instance.message.MessageRecordVO;
import org.apache.rocketmq.studio.instance.message.MessageQueryResult;
import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageDTO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.provider.InstanceProvider;
@@ -144,6 +148,7 @@ public class TencentInstanceProvider implements
InstanceProvider {
InstanceCapability.CONSUMER_GROUP_MANAGEMENT,
InstanceCapability.MESSAGE_QUERY,
InstanceCapability.MESSAGE_TRACE,
+ InstanceCapability.MESSAGE_SEND,
InstanceCapability.ACL_MANAGEMENT);
}
@@ -555,6 +560,30 @@ public class TencentInstanceProvider implements
InstanceProvider {
clientFactory.call(context.credentialId(), context.regionId(), client
-> client.ResetConsumerGroupOffset(request));
}
+ @Override
+ public SendMessageVO sendMessage(SendMessageDTO request) {
+ if (request.getProperties() != null &&
!request.getProperties().isEmpty()) {
+ throw new BusinessException(400, "Tencent Cloud test sending does
not support user properties");
+ }
+ if (StringUtils.hasText(request.getMessageGroup()) ||
request.getDeliveryTimestamp() != null) {
+ throw new BusinessException(400,
+ "Tencent Cloud test sending supports normal messages only;
FIFO and delayed fields are unsupported");
+ }
+ Context context = resolve(request.getInstanceId());
+ SendMessageRequest apiRequest = new SendMessageRequest();
+ apiRequest.setInstanceId(context.cloudInstanceId());
+ apiRequest.setTopic(request.getTopic());
+ apiRequest.setMsgBody(request.getBody());
+ apiRequest.setMsgKey(request.getKey());
+ apiRequest.setMsgTag(request.getTag());
+ SendMessageResponse response =
clientFactory.call(context.credentialId(), context.regionId(),
+ client -> client.SendMessage(apiRequest));
+ if (response == null || !StringUtils.hasText(response.getMsgId())) {
+ throw new BusinessException(502, "Tencent Cloud test message send
returned no message id");
+ }
+ return
SendMessageVO.builder().msgId(response.getMsgId()).sendTime(System.currentTimeMillis()).build();
+ }
+
@Override
public List<MessageRecordVO> queryMessages(String instanceId, String
topic, String msgId,
String tag, String key, Long
startTime, Long endTime) {
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
index a86498582..65f28b5fb 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
@@ -122,7 +122,7 @@ class MetadataServiceTest {
when(messageService.queryMessages(
"instance-a", "orders", "msg-original", null, null, null,
null))
.thenReturn(List.of(original));
- when(adminClient.sendMessage(any(SendMessageDTO.class)))
+ when(apacheProvider.sendMessage(any(SendMessageDTO.class)))
.thenReturn(SendMessageVO.builder().msgId("msg-new").build());
SendMessageVO result = metadataService.redeliverMessage(
@@ -130,7 +130,7 @@ class MetadataServiceTest {
assertThat(result.getMsgId()).isEqualTo("msg-new");
ArgumentCaptor<SendMessageDTO> request =
ArgumentCaptor.forClass(SendMessageDTO.class);
- verify(adminClient).sendMessage(request.capture());
+ verify(apacheProvider).sendMessage(request.capture());
assertThat(request.getValue().getTopic()).isEqualTo("orders-retry");
assertThat(request.getValue().getTag()).isEqualTo("paid");
assertThat(request.getValue().getKey()).isEqualTo("order-1");
@@ -148,13 +148,13 @@ class MetadataServiceTest {
when(messageService.queryMessages(
"instance-a", "orders", "msg-original", null, null, null,
null))
.thenReturn(List.of(original));
- when(adminClient.sendMessage(any(SendMessageDTO.class)))
+ when(apacheProvider.sendMessage(any(SendMessageDTO.class)))
.thenReturn(SendMessageVO.builder().msgId("msg-new").build());
metadataService.redeliverMessage("instance-a", "group-a", "orders",
"msg-original", null);
ArgumentCaptor<SendMessageDTO> request =
ArgumentCaptor.forClass(SendMessageDTO.class);
- verify(adminClient).sendMessage(request.capture());
+ verify(apacheProvider).sendMessage(request.capture());
assertThat(request.getValue().getTopic()).isEqualTo("%RETRY%group-a");
}
@@ -182,13 +182,13 @@ class MetadataServiceTest {
when(messageService.queryMessages(
"instance-a", "orders", "msg-original", null, null, null,
null))
.thenReturn(List.of(original));
- when(adminClient.sendMessage(any(SendMessageDTO.class)))
+ when(apacheProvider.sendMessage(any(SendMessageDTO.class)))
.thenReturn(SendMessageVO.builder().msgId("msg-new").build());
metadataService.redeliverMessage("instance-a", "group-a", "orders",
"msg-original", "orders-copy");
ArgumentCaptor<SendMessageDTO> request =
ArgumentCaptor.forClass(SendMessageDTO.class);
- verify(adminClient).sendMessage(request.capture());
+ verify(apacheProvider).sendMessage(request.capture());
assertThat(request.getValue().getProperties())
.containsExactlyInAnyOrderEntriesOf(Map.of("tenant", "alpha"));
assertThat(request.getValue().getTag()).isEqualTo("paid");
@@ -399,7 +399,7 @@ class MetadataServiceTest {
.key("order-1")
.body("hello")
.build();
-
when(adminClient.sendMessage(message)).thenReturn(SendMessageVO.builder().msgId("msg-1").build());
+
when(apacheProvider.sendMessage(message)).thenReturn(SendMessageVO.builder().msgId("msg-1").build());
metadataService.updateTopic(topic);
metadataService.deleteTopic("instance-a", " orders ");
@@ -708,13 +708,13 @@ class MetadataServiceTest {
.offsetMsgId("offset-001")
.build();
- when(adminClient.sendMessage(request)).thenReturn(expectedResult);
+ when(apacheProvider.sendMessage(request)).thenReturn(expectedResult);
SendMessageVO result = metadataService.sendMessage(request);
assertThat(result.getMsgId()).isEqualTo("msg-001");
assertThat(result.getOffsetMsgId()).isEqualTo("offset-001");
- verify(adminClient).sendMessage(request);
+ verify(apacheProvider).sendMessage(request);
verifyNoInteractions(operationAuditService);
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
index 59203a9df..b58cebd0c 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
@@ -37,6 +37,9 @@ import
com.aliyun.sdk.service.rocketmq20220801.models.ListTopicsResponseBody;
import
com.aliyun.sdk.service.rocketmq20220801.models.ResetConsumeOffsetRequest;
import
com.aliyun.sdk.service.rocketmq20220801.models.ResetConsumeOffsetResponse;
import
com.aliyun.sdk.service.rocketmq20220801.models.ResetConsumeOffsetResponseBody;
+import com.aliyun.sdk.service.rocketmq20220801.models.VerifySendMessageRequest;
+import
com.aliyun.sdk.service.rocketmq20220801.models.VerifySendMessageResponse;
+import
com.aliyun.sdk.service.rocketmq20220801.models.VerifySendMessageResponseBody;
import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
import org.apache.rocketmq.studio.common.domain.enums.SubscriptionMode;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
@@ -52,6 +55,8 @@ import
org.apache.rocketmq.studio.instance.message.MessageQueryResult;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageDTO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.provider.InstanceCapability;
import org.junit.jupiter.api.BeforeEach;
@@ -108,10 +113,56 @@ class AliyunInstanceProviderTest {
assertThat(provider.capabilities())
.contains(InstanceCapability.TOPIC_MANAGEMENT,
InstanceCapability.MESSAGE_QUERY,
+ InstanceCapability.MESSAGE_SEND,
InstanceCapability.ACL_MANAGEMENT)
.doesNotContain(InstanceCapability.DLQ_MANAGEMENT);
}
+ @Test
+ void sendMessageShouldMapAliyunTestSendFieldsTest() {
+ stubInstance();
+ stubCallThrough();
+ VerifySendMessageResponse response =
VerifySendMessageResponse.create().toBuilder()
+ .statusCode(200)
+ .body(VerifySendMessageResponseBody.builder()
+
.success(true).data("MSG-ALI-1").requestId("req-1").build())
+ .build();
+
when(asyncClient.verifySendMessage(any())).thenReturn(CompletableFuture.completedFuture(response));
+ SendMessageDTO request =
SendMessageDTO.builder().instanceId(STUDIO_INSTANCE_ID).topic("orders")
+
.tag("paid").key("order-1").body("payload").properties(Map.of("tenant",
"alpha"))
+
.messageGroup("order-1").deliveryTimestamp(1700000000000L).build();
+ SendMessageVO result = provider.sendMessage(request);
+ ArgumentCaptor<VerifySendMessageRequest> captor =
ArgumentCaptor.forClass(VerifySendMessageRequest.class);
+ verify(asyncClient).verifySendMessage(captor.capture());
+ VerifySendMessageRequest sent = captor.getValue();
+ assertThat(sent.getInstanceId()).isEqualTo(CLOUD_INSTANCE_ID);
+ assertThat(sent.getTopicName()).isEqualTo("orders");
+ assertThat(sent.getMessage()).isEqualTo("payload");
+ assertThat(sent.getMessageTag()).isEqualTo("paid");
+ assertThat(sent.getMessageKey()).isEqualTo("order-1");
+ assertThat(sent.getMessageGroup()).isEqualTo("order-1");
+ assertThat(sent.getDeliveryTimeStamp()).isEqualTo(1700000000000L);
+ assertThat(sent.getUserProperties()).containsKey("tenant");
+ assertThat(sent.getUserProperties().get("tenant")).isEqualTo("alpha");
+ assertThat(result.getMsgId()).isEqualTo("MSG-ALI-1");
+ }
+
+ @Test
+ void sendMessageShouldRejectUnsuccessfulAliyunResponseTest() {
+ stubInstance();
+ stubCallThrough();
+ VerifySendMessageResponse response =
VerifySendMessageResponse.create().toBuilder()
+ .statusCode(200)
+
.body(VerifySendMessageResponseBody.builder().success(false).message("send
denied").build())
+ .build();
+
when(asyncClient.verifySendMessage(any())).thenReturn(CompletableFuture.completedFuture(response));
+ SendMessageDTO request =
SendMessageDTO.builder().instanceId(STUDIO_INSTANCE_ID)
+ .topic("orders").body("payload").build();
+ assertThatThrownBy(() -> provider.sendMessage(request))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("send denied");
+ }
+
@Test
void listTopicsShouldMapMessageTypeAndFilterTest() {
stubInstance();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
index b562ec3f7..6968df8c5 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
@@ -81,6 +81,7 @@ class ApacheInstanceProviderTest {
InstanceCapability.CONSUMER_GROUP_MANAGEMENT,
InstanceCapability.MESSAGE_QUERY,
InstanceCapability.MESSAGE_TRACE,
+ InstanceCapability.MESSAGE_SEND,
InstanceCapability.ACL_MANAGEMENT,
InstanceCapability.DLQ_MANAGEMENT);
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
index 2f7d98b4a..64e8a0932 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
@@ -40,6 +40,8 @@ import
com.tencentcloudapi.trocket.v20230308.models.DescribeTopicResponse;
import com.tencentcloudapi.trocket.v20230308.models.Filter;
import com.tencentcloudapi.trocket.v20230308.models.ModifyTopicRequest;
import
com.tencentcloudapi.trocket.v20230308.models.ResetConsumerGroupOffsetRequest;
+import com.tencentcloudapi.trocket.v20230308.models.SendMessageRequest;
+import com.tencentcloudapi.trocket.v20230308.models.SendMessageResponse;
import com.tencentcloudapi.trocket.v20230308.models.SubscriptionData;
import com.tencentcloudapi.trocket.v20230308.models.TopicItem;
import com.tencentcloudapi.trocket.v20230308.TrocketClient;
@@ -61,6 +63,8 @@ import
org.apache.rocketmq.studio.instance.message.MessageQueryResult;
import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageDTO;
+import org.apache.rocketmq.studio.instance.topic.SendMessageVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.provider.InstanceCapability;
import org.junit.jupiter.api.BeforeEach;
@@ -126,10 +130,50 @@ class TencentInstanceProviderTest {
assertThat(provider.capabilities())
.contains(InstanceCapability.TOPIC_MANAGEMENT,
InstanceCapability.MESSAGE_QUERY,
+ InstanceCapability.MESSAGE_SEND,
InstanceCapability.ACL_MANAGEMENT)
.doesNotContain(InstanceCapability.DLQ_MANAGEMENT);
}
+ @Test
+ void sendMessageShouldMapTencentConsoleTestSendFieldsTest() throws
Exception {
+ SendMessageResponse response = new SendMessageResponse();
+ response.setMsgId("MSG-TENCENT-1");
+ response.setRequestId("req-1");
+ when(client.SendMessage(any())).thenReturn(response);
+ SendMessageDTO request =
SendMessageDTO.builder().instanceId(STUDIO_INSTANCE_ID).topic("orders")
+ .tag("paid").key("order-1").body("payload").build();
+ SendMessageVO result = provider.sendMessage(request);
+ ArgumentCaptor<SendMessageRequest> captor =
ArgumentCaptor.forClass(SendMessageRequest.class);
+ verify(client).SendMessage(captor.capture());
+ SendMessageRequest sent = captor.getValue();
+ assertThat(sent.getInstanceId()).isEqualTo(CLOUD_INSTANCE_ID);
+ assertThat(sent.getTopic()).isEqualTo("orders");
+ assertThat(sent.getMsgBody()).isEqualTo("payload");
+ assertThat(sent.getMsgKey()).isEqualTo("order-1");
+ assertThat(sent.getMsgTag()).isEqualTo("paid");
+ assertThat(result.getMsgId()).isEqualTo("MSG-TENCENT-1");
+ }
+
+ @Test
+ void sendMessageShouldRejectFieldsTencentConsoleApiCannotRepresentTest() {
+ SendMessageDTO request =
SendMessageDTO.builder().instanceId(STUDIO_INSTANCE_ID).topic("orders")
+ .body("payload").properties(java.util.Map.of("tenant",
"alpha")).build();
+ assertThatThrownBy(() -> provider.sendMessage(request))
+
.isInstanceOf(BusinessException.class).hasMessageContaining("user properties");
+ verifyNoInteractions(client);
+ }
+
+ @Test
+ void sendMessageShouldRejectTencentResponseWithoutMessageIdTest() throws
Exception {
+ when(client.SendMessage(any())).thenReturn(new SendMessageResponse());
+ SendMessageDTO request =
SendMessageDTO.builder().instanceId(STUDIO_INSTANCE_ID)
+ .topic("orders").body("payload").build();
+ assertThatThrownBy(() -> provider.sendMessage(request))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("no message id");
+ }
+
@Test
void countTopicsShouldClampOversizedTotals() throws Exception {
DescribeTopicListResponse response = new DescribeTopicListResponse();
diff --git a/web/src/api/instance.ts b/web/src/api/instance.ts
index 12c87bb13..eab47239b 100644
--- a/web/src/api/instance.ts
+++ b/web/src/api/instance.ts
@@ -25,6 +25,7 @@ export type InstanceCapability =
| 'CONSUMER_GROUP_MANAGEMENT'
| 'MESSAGE_QUERY'
| 'MESSAGE_TRACE'
+ | 'MESSAGE_SEND'
| 'ACL_MANAGEMENT'
| 'DLQ_MANAGEMENT';
diff --git a/web/src/pages/instance/__tests__/TopicPage.test.tsx
b/web/src/pages/instance/__tests__/TopicPage.test.tsx
index a4f27af07..752ac4f4b 100644
--- a/web/src/pages/instance/__tests__/TopicPage.test.tsx
+++ b/web/src/pages/instance/__tests__/TopicPage.test.tsx
@@ -382,6 +382,32 @@ describe('TopicPage', () => {
);
});
+ it('sends a normal test message from an Aliyun cloud topic', async () => {
+ const user = userEvent.setup();
+ instanceServiceMocks.listInstances.mockResolvedValue([
+ { ...selectedInstance, type: 'CLOUD', vendor: 'ALIYUN' },
+ ]);
+ mockTopicsList([buildTopics(1)[0]]);
+ renderWithProviders();
+
+ await user.click(await screen.findByRole('button', { name: /发送/ }));
+ const dialog = await getSendDialog();
+ fireEvent.change(within(dialog).getByLabelText('消息体 Body'), {
+ target: { value: 'cloud-test-payload' },
+ });
+ await user.click(within(dialog).getByRole('button', { name: /发\s*送/ }));
+
+ await waitFor(() =>
expect(topicServiceMocks.sendTopicMessage).toHaveBeenCalledTimes(1));
+ expect(topicServiceMocks.sendTopicMessage).toHaveBeenCalledWith({
+ topic: 'topic-01',
+ instanceId: 'instance-proxy-1',
+ tag: undefined,
+ key: undefined,
+ body: 'cloud-test-payload',
+ properties: {},
+ });
+ });
+
it('opens a clean create dialog after a cancelled edit', async () => {
const user = userEvent.setup();
renderWithProviders();
diff --git a/web/src/pages/instance/topic.tsx b/web/src/pages/instance/topic.tsx
index 4f6e8c4cb..823f049f0 100644
--- a/web/src/pages/instance/topic.tsx
+++ b/web/src/pages/instance/topic.tsx
@@ -350,6 +350,7 @@ const TopicPageContent = ({
const navigate = useNavigate();
const isCloudInstance =
selectedInstance?.vendor === 'ALIYUN' || selectedInstance?.vendor ===
'TENCENT';
+ const canSendTestMessage = (topic: Topic) => !isCloudInstance || topic.type
=== 'NORMAL';
const hasSelectedInstance = Boolean(selectedInstanceId);
// ─── State ─────────────────────────────────────────────────────
@@ -769,7 +770,7 @@ const TopicPageContent = ({
>
配置
</Button>
- {!isCloudInstance && (
+ {canSendTestMessage(record) && (
<Button
size="small"
icon={<SendOutlined />}