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);

Reply via email to