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 />}

Reply via email to