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 287e352c4 feat(message): add direct message consumption for Apache
groups (#2508)
287e352c4 is described below
commit 287e352c4a80a4398d55d88ec78066e52870108c
Author: aias00 <[email protected]>
AuthorDate: Tue Aug 25 18:06:01 2026 +0800
feat(message): add direct message consumption for Apache groups (#2508)
* feat(studio): add direct message consumption
* test(message): update direct consume service fixture
Signed-off-by: liuhy <[email protected]>
* test: align test method naming
Signed-off-by: liuhy <[email protected]>
---------
Signed-off-by: liuhy <[email protected]>
---
.../instance/message/DirectConsumeMessageDTO.java | 26 +++++++
.../message/DirectConsumeMessageResultVO.java | 22 ++++++
.../studio/instance/message/MessageController.java | 8 ++
.../studio/instance/message/MessageProvider.java | 4 +
.../studio/instance/message/MessageService.java | 13 ++++
.../rocketmq/studio/provider/InstanceProvider.java | 6 ++
.../provider/apache/ApacheInstanceProvider.java | 7 ++
.../provider/apache/RocketMQMessageProvider.java | 15 ++++
.../instance/message/MessageControllerTest.java | 14 ++++
.../instance/message/MessageServiceTest.java | 45 ++++++++---
.../apache/RocketMQMessageProviderTest.java | 27 +++++++
web/src/api/message.test.ts | 21 +++++-
web/src/api/message.ts | 24 ++++++
.../__tests__/MessagePageAsyncState.test.tsx | 19 ++++-
web/src/pages/instance/message.tsx | 87 ++++++++++++++++++++--
web/src/services/messageService.test.ts | 22 +++++-
web/src/services/messageService.ts | 17 +++++
17 files changed, 356 insertions(+), 21 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/DirectConsumeMessageDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/DirectConsumeMessageDTO.java
new file mode 100644
index 000000000..c6de345e3
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/DirectConsumeMessageDTO.java
@@ -0,0 +1,26 @@
+/*
+ * 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.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ */
+package org.apache.rocketmq.studio.instance.message;
+
+import jakarta.validation.constraints.NotBlank;
+import lombok.Data;
+
+@Data
+public class DirectConsumeMessageDTO {
+ @NotBlank(message = "instanceId is required")
+ private String instanceId;
+ @NotBlank(message = "topic is required")
+ private String topic;
+ @NotBlank(message = "msgId is required")
+ private String msgId;
+ @NotBlank(message = "consumerGroup is required")
+ private String consumerGroup;
+ @NotBlank(message = "clientId is required")
+ private String clientId;
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/DirectConsumeMessageResultVO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/DirectConsumeMessageResultVO.java
new file mode 100644
index 000000000..ec4bafc27
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/DirectConsumeMessageResultVO.java
@@ -0,0 +1,22 @@
+/*
+ * 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.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ */
+package org.apache.rocketmq.studio.instance.message;
+
+import lombok.Builder;
+import lombok.Value;
+
+@Value
+@Builder
+public class DirectConsumeMessageResultVO {
+ String consumeResult;
+ String remark;
+ long spentTimeMillis;
+ boolean order;
+ boolean autoCommit;
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
index 314692cc8..506d4eee1 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
@@ -19,6 +19,8 @@ package org.apache.rocketmq.studio.instance.message;
import org.apache.rocketmq.studio.common.domain.Result;
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
@@ -77,4 +79,10 @@ public class MessageController {
@RequestParam long
offset) {
return Result.ok(messageService.pullMessageAtOffset(instanceId, topic,
brokerName, queueId, offset));
}
+
+ @PostMapping("/direct-consume")
+ public Result<DirectConsumeMessageResultVO> consumeMessageDirectly(
+ @jakarta.validation.Valid @RequestBody DirectConsumeMessageDTO
request) {
+ return Result.ok(messageService.consumeMessageDirectly(request));
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
index 09242f9fd..7f11e66f3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
@@ -28,4 +28,8 @@ public interface MessageProvider {
List<QueueOffsetVO> getQueueOffsets(String instanceId, String topic);
MessageRecordVO pullMessageAtOffset(String instanceId, String topic,
String brokerName, int queueId, long offset);
+
+ 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/instance/message/MessageService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
index 5a185795d..0b939aa1e 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
@@ -20,6 +20,7 @@ import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.provider.InstanceProviderRegistry;
+import org.apache.rocketmq.studio.audit.OperationAuditService;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
@@ -37,6 +38,7 @@ public class MessageService {
private final MessageProvider messageProvider;
private final InstanceProviderRegistry providerRegistry;
private final QueryHistoryService queryHistoryService;
+ private final OperationAuditService operationAuditService;
public List<MessageRecordVO> queryMessages(
String instanceId, String topic, String msgId, String tag, String
key, Long startTime, Long endTime) {
@@ -100,6 +102,17 @@ public class MessageService {
return messageProvider.pullMessageAtOffset(instanceId, topic,
brokerName, queueId, offset);
}
+ public DirectConsumeMessageResultVO
consumeMessageDirectly(DirectConsumeMessageDTO request) {
+ DirectConsumeMessageResultVO result =
providerRegistry.byInstanceId(request.getInstanceId())
+ .map(provider -> provider.consumeMessageDirectly(request))
+ .orElseGet(() ->
messageProvider.consumeMessageDirectly(request));
+ operationAuditService.record("DIRECT_CONSUME_MESSAGE", "MESSAGE",
request.getMsgId(), request.getInstanceId(),
+ "topic=" + request.getTopic() + ", consumerGroup=" +
request.getConsumerGroup()
+ + ", clientId=" + request.getClientId() + ", result="
+ result.getConsumeResult(),
+ "SUCCESS", null);
+ return result;
+ }
+
private void recordMessageQuery(String instanceId, String topic, String
msgId, String tag,
String key, Long startTime, Long endTime,
int resultCount) {
String queryType = StringUtils.hasText(msgId) ? "MSG_ID" :
StringUtils.hasText(key) ? "KEY" : "TOPIC";
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 13ee1133a..d83c8e377 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
@@ -23,6 +23,8 @@ import
org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
+import org.apache.rocketmq.studio.instance.message.DirectConsumeMessageDTO;
+import
org.apache.rocketmq.studio.instance.message.DirectConsumeMessageResultVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
@@ -98,4 +100,8 @@ public interface InstanceProvider {
String tag, String key, Long
startTime, Long endTime);
TraceRecordVO getMessageTrace(String instanceId, String msgId, String
topic);
+
+ 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/apache/ApacheInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
index 27b1fc909..a9e94e7dc 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
@@ -23,6 +23,8 @@ import
org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
import org.apache.rocketmq.studio.instance.message.MessageProvider;
+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.TopicConsumerVO;
@@ -155,4 +157,9 @@ public class ApacheInstanceProvider implements
InstanceProvider {
public TraceRecordVO getMessageTrace(String instanceId, String msgId,
String topic) {
return messageProvider.getMessageTrace(instanceId, msgId, topic);
}
+
+ @Override
+ public DirectConsumeMessageResultVO
consumeMessageDirectly(DirectConsumeMessageDTO request) {
+ return messageProvider.consumeMessageDirectly(request);
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
index b1b3c00de..009c0a84b 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
@@ -36,6 +36,8 @@ import
org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
import org.apache.rocketmq.studio.instance.message.ConsumerStatusVO;
import org.apache.rocketmq.studio.instance.message.MessageProvider;
+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.QueueOffsetVO;
import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
@@ -425,6 +427,19 @@ public class RocketMQMessageProvider implements
MessageProvider {
adminExt -> getMessageTrace(instanceId, (DefaultMQAdminExt)
adminExt, msgId, topic));
}
+ @Override
+ public DirectConsumeMessageResultVO
consumeMessageDirectly(DirectConsumeMessageDTO request) {
+ return runtimeAdminClientResolver.execute(request.getInstanceId(),
admin -> {
+
org.apache.rocketmq.remoting.protocol.body.ConsumeMessageDirectlyResult result =
+ ((DefaultMQAdminExt)
admin).consumeMessageDirectly(request.getConsumerGroup(), request.getClientId(),
+ request.getTopic(), request.getMsgId());
+ return DirectConsumeMessageResultVO.builder()
+ .consumeResult(result.getConsumeResult() == null ?
"UNKNOWN" : result.getConsumeResult().name())
+
.remark(result.getRemark()).spentTimeMillis(result.getSpentTimeMills())
+
.order(result.isOrder()).autoCommit(result.isAutoCommit()).build();
+ });
+ }
+
private TraceRecordVO getMessageTrace(String instanceId, DefaultMQAdminExt
adminExt, String msgId, String topic) {
long now = System.currentTimeMillis();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
index 5dba50b20..001dec2be 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
@@ -31,6 +31,7 @@ import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static
org.springframework.test.web.servlet.request.MockMvcRequestBuilders.get;
+import static
org.springframework.test.web.servlet.request.MockMvcRequestBuilders.post;
import static
org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath;
import static
org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@@ -99,4 +100,17 @@ class MessageControllerTest {
verify(messageService).getMessageTrace("instance-a", "msg-001",
"orders");
}
+
+ @Test
+ void directConsumeShouldPassValidatedRequestTest() throws Exception {
+ DirectConsumeMessageResultVO result =
DirectConsumeMessageResultVO.builder()
+
.consumeResult("CR_SUCCESS").remark("ok").spentTimeMillis(8).build();
+
when(messageService.consumeMessageDirectly(org.mockito.ArgumentMatchers.any())).thenReturn(result);
+
+
mockMvc.perform(post("/api/messages/direct-consume").contentType("application/json")
+
.content("{\"instanceId\":\"instance-a\",\"topic\":\"orders\",\"msgId\":\"msg-001\","
+ +
"\"consumerGroup\":\"billing\",\"clientId\":\"client-a\"}"))
+ .andExpect(status().isOk())
+
.andExpect(jsonPath("$.data.consumeResult").value("CR_SUCCESS"));
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
index 7c3ef6c11..a409242d9 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
@@ -13,6 +13,7 @@ package org.apache.rocketmq.studio.instance.message;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.provider.InstanceProviderRegistry;
import org.apache.rocketmq.studio.provider.InstanceProvider;
+import org.apache.rocketmq.studio.audit.OperationAuditService;
import org.junit.jupiter.api.Test;
import java.util.List;
@@ -31,7 +32,7 @@ class MessageServiceTest {
void rejectsNegativeQueueCoordinatesBeforeCallingProvider() {
MessageProvider provider = mock(MessageProvider.class);
MessageService service = new MessageService(provider,
mock(InstanceProviderRegistry.class),
- mock(QueryHistoryService.class));
+ mock(QueryHistoryService.class),
mock(OperationAuditService.class));
assertThatThrownBy(() -> service.pullMessageAtOffset("instance-a",
"TopicA", "broker-a", -1, 0))
.isInstanceOf(BusinessException.class)
@@ -47,7 +48,7 @@ class MessageServiceTest {
void rejectsKeyQueryWithoutTopicBeforeCallingProvider() {
MessageProvider provider = mock(MessageProvider.class);
InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
- MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class));
+ MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class), mock(OperationAuditService.class));
assertThatThrownBy(() -> service.queryMessages(null, null, null, null,
"order-1", null, null))
.isInstanceOf(BusinessException.class)
@@ -60,7 +61,7 @@ class MessageServiceTest {
void rejectsMessageIdQueryWithoutTopicBeforeCallingProvider() {
MessageProvider provider = mock(MessageProvider.class);
InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
- MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class));
+ MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class), mock(OperationAuditService.class));
assertThatThrownBy(() -> service.queryMessages(null, null, "msg-001",
null, null, null, null))
.isInstanceOf(BusinessException.class)
@@ -73,7 +74,7 @@ class MessageServiceTest {
void rejectsBlankMessageTraceIdBeforeCallingProvider() {
MessageProvider provider = mock(MessageProvider.class);
InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
- MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class));
+ MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class), mock(OperationAuditService.class));
assertThatThrownBy(() -> service.getMessageTrace("instance-a", " ",
null))
.isInstanceOf(BusinessException.class)
@@ -86,7 +87,7 @@ class MessageServiceTest {
void rejectsReversedTopicQueryWindowBeforeCallingProvider() {
MessageProvider provider = mock(MessageProvider.class);
InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
- MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class));
+ MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class), mock(OperationAuditService.class));
assertThatThrownBy(() -> service.queryMessages("instance-a", "TopicA",
null, null, null, 200L, 100L))
.isInstanceOf(BusinessException.class)
@@ -99,7 +100,7 @@ class MessageServiceTest {
void rejectsTopicQueryWindowLongerThanSevenDaysBeforeCallingProvider() {
MessageProvider provider = mock(MessageProvider.class);
InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
- MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class));
+ MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class), mock(OperationAuditService.class));
assertThatThrownBy(() -> service.queryMessages("instance-a", "TopicA",
null, null, null, 0L,
8L * 24 * 60 * 60 * 1000))
@@ -115,7 +116,7 @@ class MessageServiceTest {
InstanceProvider provider = mock(InstanceProvider.class);
InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
QueryHistoryService history = mock(QueryHistoryService.class);
- MessageService service = new MessageService(fallback, registry,
history);
+ MessageService service = new MessageService(fallback, registry,
history, mock(OperationAuditService.class));
when(registry.byInstanceId("cloud-instance")).thenReturn(Optional.of(provider));
when(provider.queryMessages("cloud-instance", "orders", null, null,
"ORDER-1", null, null))
.thenReturn(List.of(MessageRecordVO.builder().msgId("msg-1").build()));
@@ -127,11 +128,37 @@ class MessageServiceTest {
verifyNoInteractions(fallback);
}
+ @Test
+ void directsMessageConsumptionThroughSelectedProviderAndAuditsItTest() {
+ MessageProvider fallback = mock(MessageProvider.class);
+ InstanceProvider provider = mock(InstanceProvider.class);
+ InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
+ OperationAuditService audit = mock(OperationAuditService.class);
+ DirectConsumeMessageDTO request = new DirectConsumeMessageDTO();
+ request.setInstanceId("instance-a");
+ request.setTopic("orders");
+ request.setMsgId("msg-1");
+ request.setConsumerGroup("billing");
+ request.setClientId("client-a");
+ DirectConsumeMessageResultVO expected =
DirectConsumeMessageResultVO.builder()
+ .consumeResult("CR_SUCCESS").build();
+
when(registry.byInstanceId("instance-a")).thenReturn(Optional.of(provider));
+ when(provider.consumeMessageDirectly(request)).thenReturn(expected);
+ MessageService service = new MessageService(fallback, registry,
mock(QueryHistoryService.class), audit);
+
+
org.assertj.core.api.Assertions.assertThat(service.consumeMessageDirectly(request)).isSameAs(expected);
+
+
verify(audit).record(org.mockito.ArgumentMatchers.eq("DIRECT_CONSUME_MESSAGE"),
+ org.mockito.ArgumentMatchers.eq("MESSAGE"),
org.mockito.ArgumentMatchers.eq("msg-1"),
+ org.mockito.ArgumentMatchers.eq("instance-a"),
org.mockito.ArgumentMatchers.contains("billing"),
+ org.mockito.ArgumentMatchers.eq("SUCCESS"),
org.mockito.ArgumentMatchers.isNull());
+ }
+
@Test
void rejectsOverflowingTopicQueryWindowBeforeCallingProvider() {
MessageProvider provider = mock(MessageProvider.class);
InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
- MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class));
+ MessageService service = new MessageService(provider, registry,
mock(QueryHistoryService.class), mock(OperationAuditService.class));
assertThatThrownBy(() -> service.queryMessages("instance-a", "TopicA",
null, null, null,
0L, Long.MAX_VALUE))
@@ -146,7 +173,7 @@ class MessageServiceTest {
MessageProvider provider = mock(MessageProvider.class);
InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
QueryHistoryService history = mock(QueryHistoryService.class);
- MessageService service = new MessageService(provider, registry,
history);
+ MessageService service = new MessageService(provider, registry,
history, mock(OperationAuditService.class));
when(registry.byInstanceId("instance-a")).thenReturn(Optional.empty());
when(provider.queryMessages("instance-a", "TopicA", null, null, null,
1000L, 2000L))
.thenReturn(java.util.stream.IntStream.range(0, 200)
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
index 2501eda4b..5c8f2756f 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
@@ -28,6 +28,8 @@ import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageId;
import org.apache.rocketmq.common.message.MessageQueue;
import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
+import org.apache.rocketmq.remoting.protocol.body.ConsumeMessageDirectlyResult;
+import org.apache.rocketmq.remoting.protocol.body.CMResult;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
@@ -35,6 +37,7 @@ import
org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
+import org.apache.rocketmq.studio.instance.message.DirectConsumeMessageDTO;
import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
@@ -136,6 +139,30 @@ class RocketMQMessageProviderTest {
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
}
+ @Test
+ void directlyConsumesMessageForTheExplicitGroupAndClientTest() throws
Exception {
+ ConsumeMessageDirectlyResult brokerResult = new
ConsumeMessageDirectlyResult();
+ brokerResult.setConsumeResult(CMResult.CR_SUCCESS);
+ brokerResult.setRemark("consumed");
+ brokerResult.setSpentTimeMills(12);
+ when(adminExt.consumeMessageDirectly("billing", "client-a", "orders",
"msg-1"))
+ .thenReturn(brokerResult);
+ DirectConsumeMessageDTO request = new DirectConsumeMessageDTO();
+ request.setInstanceId("instance-a");
+ request.setTopic("orders");
+ request.setMsgId("msg-1");
+ request.setConsumerGroup("billing");
+ request.setClientId("client-a");
+
+
org.apache.rocketmq.studio.instance.message.DirectConsumeMessageResultVO result
=
+ provider.consumeMessageDirectly(request);
+
+ assertThat(result.getConsumeResult()).isEqualTo("CR_SUCCESS");
+ assertThat(result.getRemark()).isEqualTo("consumed");
+ assertThat(result.getSpentTimeMillis()).isEqualTo(12);
+ verify(adminExt).consumeMessageDirectly("billing", "client-a",
"orders", "msg-1");
+ }
+
@Test
void queryMessagesShouldRejectInvertedTimeRangeBeforeAdminLookup() throws
Exception {
assertThatThrownBy(() -> provider.queryMessages(
diff --git a/web/src/api/message.test.ts b/web/src/api/message.test.ts
index 7c7a651af..caa55b1b3 100644
--- a/web/src/api/message.test.ts
+++ b/web/src/api/message.test.ts
@@ -18,7 +18,7 @@
import MockAdapter from 'axios-mock-adapter';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import client from './client';
-import { getMessageTrace, queryMessagePage, queryMessages } from './message';
+import { consumeMessageDirectly, getMessageTrace, queryMessagePage,
queryMessages } from './message';
const mock = new MockAdapter(client);
@@ -140,4 +140,23 @@ describe('message API', () => {
trace,
);
});
+
+ it('posts direct consumption to the message API', async () => {
+ const request = {
+ instanceId: 'instance-1',
+ topic: 'orders',
+ msgId: 'msg-1',
+ consumerGroup: 'billing',
+ clientId: 'client-a',
+ };
+ const result = {
+ consumeResult: 'CR_SUCCESS',
+ spentTimeMillis: 8,
+ order: false,
+ autoCommit: true,
+ };
+ mock.onPost('/messages/direct-consume', request).reply(200, { code: 200,
data: result });
+
+ await expect(consumeMessageDirectly(request)).resolves.toEqual(result);
+ });
});
diff --git a/web/src/api/message.ts b/web/src/api/message.ts
index 81df854a3..10b9c478d 100644
--- a/web/src/api/message.ts
+++ b/web/src/api/message.ts
@@ -55,6 +55,22 @@ export interface MessageQueryPage {
resultMayBeTruncated: boolean;
}
+export interface DirectConsumeMessageRequest {
+ instanceId: string;
+ topic: string;
+ msgId: string;
+ consumerGroup: string;
+ clientId: string;
+}
+
+export interface DirectConsumeMessageResult {
+ consumeResult: string;
+ remark?: string;
+ spentTimeMillis: number;
+ order: boolean;
+ autoCommit: boolean;
+}
+
const toStoreTimestamp = (storeTime: MessageRecord['storeTime']): number => {
if (typeof storeTime === 'number') return storeTime;
@@ -134,6 +150,14 @@ export async function getMessageTrace(msgId: string,
instanceId?: string, topic?
return res.data.data;
}
+export async function consumeMessageDirectly(data:
DirectConsumeMessageRequest) {
+ const res = await client.post<{ data: DirectConsumeMessageResult }>(
+ '/messages/direct-consume',
+ data,
+ );
+ return res.data.data;
+}
+
// ─── DLQ ────────────────────────────────────────────────────────
export async function listDLQGroups(instanceId: string, search?: string, page
= 1, pageSize = 20) {
const res = await client.get<{ data: DLQGroupPage }>('/dlq', {
diff --git a/web/src/pages/instance/__tests__/MessagePageAsyncState.test.tsx
b/web/src/pages/instance/__tests__/MessagePageAsyncState.test.tsx
index 8f7ae5192..005120d9c 100644
--- a/web/src/pages/instance/__tests__/MessagePageAsyncState.test.tsx
+++ b/web/src/pages/instance/__tests__/MessagePageAsyncState.test.tsx
@@ -27,6 +27,7 @@ import MessagePage from '../message';
const serviceMocks = vi.hoisted(() => ({
getMessageTrace: vi.fn(),
queryMessages: vi.fn(),
+ consumeMessageDirectly: vi.fn(),
}));
const instanceFilterMocks = vi.hoisted(() => ({
useInstanceFilter: vi.fn(),
@@ -229,7 +230,7 @@ describe('MessagePage async request ownership', () => {
).toBeInTheDocument();
});
- it('keeps normal message resend disabled until a real API is wired', async
() => {
+ it('requiresGroupAndClientBeforeDirectConsumeTest', async () => {
serviceMocks.queryMessages.mockResolvedValue([createMessage('message-a')]);
const user = userEvent.setup();
renderPage();
@@ -240,9 +241,19 @@ describe('MessagePage async request ownership', () => {
await user.click(within(row).getByRole('button', { name: /详情/ }));
const dialog = await screen.findByRole('dialog', { name: '消息详情' });
- const resendButton = within(dialog).getByRole('button', { name: /重新发送/ });
- expect(resendButton).toBeDisabled();
- expect(resendButton).toHaveAttribute('title', '当前版本尚未接入普通消息重新发送接口');
+ await user.click(within(dialog).getByRole('button', { name: /直接消费/ }));
+ const consumeDialogTitle = await screen.findByText('直接消费消息');
+ const consumeDialog = consumeDialogTitle.closest('[role="dialog"]');
+ expect(consumeDialog).not.toBeNull();
+ expect(
+ within(consumeDialog as HTMLElement).getByPlaceholderText('目标消费者组'),
+ ).toBeInTheDocument();
+ expect(
+ within(consumeDialog as HTMLElement).getByPlaceholderText('在线客户端 ID'),
+ ).toBeInTheDocument();
+
+ await user.click(within(consumeDialog as HTMLElement).getByRole('button',
{ name: /执\s*行/ }));
+ expect(serviceMocks.consumeMessageDirectly).not.toHaveBeenCalled();
});
it('keeps the latest query loading and ignores an earlier query result',
async () => {
diff --git a/web/src/pages/instance/message.tsx
b/web/src/pages/instance/message.tsx
index d74cb2292..01130a01f 100644
--- a/web/src/pages/instance/message.tsx
+++ b/web/src/pages/instance/message.tsx
@@ -59,7 +59,11 @@ import {
import type { MessageQueryHistory, TraceQueryHistory } from
'../../api/messageHistory';
import { useLang } from '../../i18n/LangContext';
import type { MessageQuery, MessageRecord, TraceRecord } from
'../../api/message';
-import { getMessageTrace, queryMessagePage } from
'../../services/messageService';
+import {
+ consumeMessageDirectly,
+ getMessageTrace,
+ queryMessagePage,
+} from '../../services/messageService';
import { listTopics } from '../../services/topicService';
import { useInstanceFilter } from '../../hooks/useInstanceFilter';
import { downloadBlob } from '../../utils/download';
@@ -83,8 +87,6 @@ type ApiErrorLike = {
};
};
-const RESEND_UNAVAILABLE_MESSAGE = '当前版本尚未接入普通消息重新发送接口';
-
const QUERY_OPTIONS = [
{ value: 'topic' as const, label: '按 Topic 查询' },
{ value: 'key' as const, label: '按 Message Key' },
@@ -265,6 +267,10 @@ const MessagePageContent = ({
const [queryError, setQueryError] = useState<string | null>(null);
const [traceError, setTraceError] = useState<string | null>(null);
const [historyDrawerOpen, setHistoryDrawerOpen] = useState(false);
+ const [directConsumeOpen, setDirectConsumeOpen] = useState(false);
+ const [directConsumeGroup, setDirectConsumeGroup] = useState('');
+ const [directConsumeClientId, setDirectConsumeClientId] = useState('');
+ const [directConsumeSubmitting, setDirectConsumeSubmitting] =
useState(false);
const queryGenerationRef = useRef(0);
const traceGenerationRef = useRef(0);
@@ -419,6 +425,41 @@ const MessagePageContent = ({
setTraceError(null);
};
+ const openDirectConsume = () => {
+ setDirectConsumeGroup('');
+ setDirectConsumeClientId('');
+ setDirectConsumeOpen(true);
+ };
+
+ const handleDirectConsume = async () => {
+ if (
+ !selectedInstanceId ||
+ !selectedMsg ||
+ !directConsumeGroup.trim() ||
+ !directConsumeClientId.trim()
+ ) {
+ message.warning('请填写目标消费组和在线客户端 ID');
+ return;
+ }
+ setDirectConsumeSubmitting(true);
+ try {
+ const result = await consumeMessageDirectly({
+ instanceId: selectedInstanceId,
+ topic: selectedMsg.topic,
+ msgId: selectedMsg.msgId,
+ consumerGroup: directConsumeGroup.trim(),
+ clientId: directConsumeClientId.trim(),
+ });
+ const detail = [result.consumeResult,
result.remark].filter(Boolean).join(':');
+ message.info(`Broker 返回 ${detail || 'UNKNOWN'},耗时
${result.spentTimeMillis} ms`);
+ setDirectConsumeOpen(false);
+ } catch (error) {
+ message.error(getErrorMessage(error, '直接消费请求失败,请检查消费组和客户端是否在线'));
+ } finally {
+ setDirectConsumeSubmitting(false);
+ }
+ };
+
const handleDownload = (record: MessageRecord) => {
const blob = new Blob([formatBody(record.body)], { type:
'application/json' });
downloadBlob(blob, `${record.msgId}.json`);
@@ -890,16 +931,50 @@ const MessagePageContent = ({
<Button
type="primary"
icon={<SendOutlined />}
- disabled
- title={RESEND_UNAVAILABLE_MESSAGE}
+ disabled={!selectedInstanceId || !selectedMsg}
+ onClick={openDirectConsume}
>
- 重新发送
+ 直接消费
</Button>
</Flex>
}
>
<Tabs activeKey={modalTab} onChange={setModalTab} items={modalTabs} />
</Modal>
+
+ <Modal
+ title="直接消费消息"
+ open={directConsumeOpen}
+ onCancel={() => setDirectConsumeOpen(false)}
+ onOk={() => void handleDirectConsume()}
+ confirmLoading={directConsumeSubmitting}
+ okText="执行"
+ destroyOnHidden
+ >
+ <Alert
+ showIcon
+ type="warning"
+ message="Broker 会请求指定在线客户端立即消费该消息。"
+ description="这不是向 Topic 重新发送消息;Broker 返回的消费结果会原样显示。"
+ style={{ marginBottom: 16 }}
+ />
+ <Space direction="vertical" size={12} style={{ width: '100%' }}>
+ <Input value={selectedMsg?.topic} disabled addonBefore="Topic" />
+ <Input value={selectedMsg?.msgId} disabled addonBefore="Message ID"
/>
+ <Input
+ value={directConsumeGroup}
+ onChange={(event) => setDirectConsumeGroup(event.target.value)}
+ placeholder="目标消费者组"
+ addonBefore="Consumer group"
+ />
+ <Input
+ value={directConsumeClientId}
+ onChange={(event) => setDirectConsumeClientId(event.target.value)}
+ placeholder="在线客户端 ID"
+ addonBefore="Client ID"
+ />
+ </Space>
+ </Modal>
</div>
);
};
diff --git a/web/src/services/messageService.test.ts
b/web/src/services/messageService.test.ts
index 7521764ab..ff1113d6a 100644
--- a/web/src/services/messageService.test.ts
+++ b/web/src/services/messageService.test.ts
@@ -16,7 +16,12 @@
*/
import { describe, expect, it, vi } from 'vitest';
-import { getMessageTrace, listDLQGroups, queryMessages } from
'./messageService';
+import {
+ consumeMessageDirectly,
+ getMessageTrace,
+ listDLQGroups,
+ queryMessages,
+} from './messageService';
vi.mock('./dataMode', () => ({ isMockMode: () => true }));
vi.mock('../config', () => ({
@@ -76,6 +81,21 @@ describe('message service mock data', () => {
expect(second.items[0]).not.toBe(first.items[0]);
});
+ it('marks direct-consume results as mock-only in mock mode', async () => {
+ await expect(
+ consumeMessageDirectly({
+ instanceId: 'instance-1',
+ topic: 'orders',
+ msgId: 'msg-1',
+ consumerGroup: 'billing',
+ clientId: 'client-a',
+ }),
+ ).resolves.toMatchObject({
+ consumeResult: 'CR_SUCCESS',
+ remark: expect.stringContaining('Mock mode'),
+ });
+ });
+
it('filters and pages mock DLQ groups', async () => {
const filtered = await listDLQGroups('instance-1', 'order', 1, 20);
expect(filtered.items.every((group) =>
group.groupName.includes('order'))).toBe(true);
diff --git a/web/src/services/messageService.ts
b/web/src/services/messageService.ts
index a65aac08e..5c744570f 100644
--- a/web/src/services/messageService.ts
+++ b/web/src/services/messageService.ts
@@ -4,6 +4,8 @@ import { sortMessagesByStoreTimeDesc } from '../api/message';
import type {
MessageQuery,
MessageQueryPage,
+ DirectConsumeMessageRequest,
+ DirectConsumeMessageResult,
MessageRecord,
TraceRecord,
DLQGroup,
@@ -83,6 +85,21 @@ export async function getMessageTrace(
return messageApi.getMessageTrace(msgId, instanceId, topic);
}
+export async function consumeMessageDirectly(
+ request: DirectConsumeMessageRequest,
+): Promise<DirectConsumeMessageResult> {
+ if (isMockMode()) {
+ return {
+ consumeResult: 'CR_SUCCESS',
+ remark: 'Mock mode: request was not sent to a broker',
+ spentTimeMillis: 0,
+ order: false,
+ autoCommit: true,
+ };
+ }
+ return messageApi.consumeMessageDirectly(request);
+}
+
export async function listDLQGroups(
instanceId: string,
search?: string,