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 042b50d98 fix(message): record direct-consume audit failures as FAILED
(#4607)
042b50d98 is described below
commit 042b50d985a6b5d65ec6213bc5d633f0c8a6794f
Author: Apulupie <[email protected]>
AuthorDate: Mon Sep 21 21:05:15 2026 +0800
fix(message): record direct-consume audit failures as FAILED (#4607)
`MessageService.consumeMessageDirectly` wrote its `DIRECT_CONSUME_MESSAGE`
audit row with a hard-coded `"SUCCESS"`, and when the provider threw it wrote
no row at all. The REST entry point in `MessageController` is reachable in
production, so a rollback result and an unsupported provider both landed in the
audit log as a success — or as nothing.
The result is now judged from `consumeResult`: only `CR_SUCCESS` audits
`SUCCESS`, everything else (including the `UNKNOWN` the provider writes for a
null result) audits `FAILED`, and the `RuntimeException` path records `FAILED`
with the exception message before rethrowing. The audit write itself is wrapped
so an audit failure cannot mask the consume outcome, the same way
`recordTraceQuery` does next to it. This is the only direct-consume audit point
in the tree, so both result paths [...]
Fixes #4608
---
.../studio/instance/message/MessageService.java | 36 +++++++++++----
.../instance/message/MessageServiceTest.java | 52 ++++++++++++++++++++++
2 files changed, 80 insertions(+), 8 deletions(-)
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 250f2220c..ce62474e3 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
@@ -134,14 +134,34 @@ public class MessageService {
}
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;
+ String detail = "topic=" + request.getTopic() + ", consumerGroup=" +
request.getConsumerGroup()
+ + ", clientId=" + request.getClientId();
+ try {
+ DirectConsumeMessageResultVO result =
providerRegistry.byInstanceId(request.getInstanceId())
+ .map(provider -> provider.consumeMessageDirectly(request))
+ .orElseGet(() ->
messageProvider.consumeMessageDirectly(request));
+ recordDirectConsumeAudit(request, detail + ", result=" +
result.getConsumeResult(),
+ auditResult(result.getConsumeResult()), null);
+ return result;
+ } catch (RuntimeException e) {
+ recordDirectConsumeAudit(request, detail, "FAILED",
e.getMessage());
+ throw e;
+ }
+ }
+
+ /** consumeResult mirrors the broker-side CMResult enum, where only
CR_SUCCESS means consumed. */
+ private static String auditResult(String consumeResult) {
+ return "CR_SUCCESS".equals(consumeResult) ? "SUCCESS" : "FAILED";
+ }
+
+ private void recordDirectConsumeAudit(DirectConsumeMessageDTO request,
String detail,
+ String result, String errorInfo) {
+ try {
+ operationAuditService.record("DIRECT_CONSUME_MESSAGE", "MESSAGE",
request.getMsgId(),
+ request.getInstanceId(), detail, result, errorInfo);
+ } catch (RuntimeException auditFailure) {
+ log.warn("Failed to record direct-consume audit: {}",
auditFailure.getMessage());
+ }
}
public TraceRecordVO getMessageTrace(String instanceId, String msgId,
String topic, String traceTopic) {
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 e92108109..58d5267e9 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
@@ -155,6 +155,58 @@ class MessageServiceTest {
org.mockito.ArgumentMatchers.eq("SUCCESS"),
org.mockito.ArgumentMatchers.isNull());
}
+ @Test
+ void auditsDirectConsumeFailureWhenConsumeResultIsNotSuccessTest() {
+ 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");
+
when(registry.byInstanceId("instance-a")).thenReturn(Optional.of(provider));
+
when(provider.consumeMessageDirectly(request)).thenReturn(DirectConsumeMessageResultVO.builder()
+ .consumeResult("CR_ROLLBACK").build());
+ MessageService service = new MessageService(fallback, registry,
mock(QueryHistoryService.class), audit);
+
+ service.consumeMessageDirectly(request);
+
+
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("FAILED"),
org.mockito.ArgumentMatchers.isNull());
+ }
+
+ @Test
+ void auditsDirectConsumeFailureWhenProviderThrowsTest() {
+ 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");
+
when(registry.byInstanceId("instance-a")).thenReturn(Optional.of(provider));
+ when(provider.consumeMessageDirectly(request)).thenThrow(new
IllegalStateException("client offline"));
+ MessageService service = new MessageService(fallback, registry,
mock(QueryHistoryService.class), audit);
+
+ assertThatThrownBy(() -> service.consumeMessageDirectly(request))
+ .isInstanceOf(IllegalStateException.class)
+ .hasMessage("client offline");
+
+
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("FAILED"),
org.mockito.ArgumentMatchers.eq("client offline"));
+ verifyNoInteractions(fallback);
+ }
+
@Test
void rejectsOverflowingTopicQueryWindowBeforeCallingProvider() {
MessageProvider provider = mock(MessageProvider.class);