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 27fc97dd3 fix(studio): isolate MCP tool results from audit failures
(#4553)
27fc97dd3 is described below
commit 27fc97dd3084c7ccdf81a949aa0a204b5b4e06b5
Author: Yexi Xiang <[email protected]>
AuthorDate: Mon Sep 21 17:33:53 2026 +0800
fix(studio): isolate MCP tool results from audit failures (#4553)
ToolAuditFilter.record was the only one of the six auditService.record call
sites without a guard, so a persistence failure in the audit filter turned an
already-successful mutation into a reported failure. The record call is now
wrapped and the failure logged at warn level.
Fixes #4552
---
.../studio/ops/ai/tool/filter/ToolAuditFilter.java | 21 ++--
.../studio/ops/ai/mcp/McpToolRegistrarTest.java | 108 +++++++++++++++++++++
2 files changed, 122 insertions(+), 7 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/filter/ToolAuditFilter.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/filter/ToolAuditFilter.java
index e80647ca3..52ba8fe49 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/filter/ToolAuditFilter.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/filter/ToolAuditFilter.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.ops.ai.tool.filter;
import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.studio.ops.ai.tool.core.ToolExecutionContext;
import org.apache.rocketmq.studio.ops.ai.tool.core.ToolInvocation;
import org.apache.rocketmq.studio.ops.audit.AuditService;
@@ -24,6 +25,7 @@ import org.springframework.stereotype.Component;
@Component
@RequiredArgsConstructor
+@Slf4j
public class ToolAuditFilter implements ToolExecutionFilter {
private final AuditService auditService;
@@ -48,13 +50,18 @@ public class ToolAuditFilter implements ToolExecutionFilter
{
}
private void record(ToolExecutionContext context, String result, String
errorMessage) {
- auditService.record(context.operationType(),
- context.resourceType(),
- context.definition().name(),
- context.instanceId(),
- errorMessage,
- result
- );
+ try {
+ auditService.record(context.operationType(),
+ context.resourceType(),
+ context.definition().name(),
+ context.instanceId(),
+ errorMessage,
+ result
+ );
+ } catch (RuntimeException auditFailure) {
+ log.warn("Failed to record tool audit result={} tool={}
instance={}: {}",
+ result, context.definition().name(), context.instanceId(),
auditFailure.getMessage());
+ }
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/mcp/McpToolRegistrarTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/mcp/McpToolRegistrarTest.java
index 4344af62c..6d34f9d01 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/mcp/McpToolRegistrarTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/mcp/McpToolRegistrarTest.java
@@ -21,19 +21,37 @@ import io.modelcontextprotocol.common.McpTransportContext;
import io.modelcontextprotocol.server.McpServerFeatures;
import io.modelcontextprotocol.server.McpSyncServerExchange;
import io.modelcontextprotocol.spec.McpSchema;
+import org.apache.rocketmq.studio.instance.InstanceResolver;
import org.apache.rocketmq.studio.ops.ai.auth.McpAuthentication;
import org.apache.rocketmq.studio.ops.ai.tool.core.ToolDefinition;
+import org.apache.rocketmq.studio.ops.ai.tool.core.ToolError;
+import org.apache.rocketmq.studio.ops.ai.tool.core.ToolExecutionContext;
+import org.apache.rocketmq.studio.ops.ai.tool.core.ToolExecutionException;
import org.apache.rocketmq.studio.ops.ai.tool.core.ToolRiskLevel;
+import org.apache.rocketmq.studio.ops.ai.tool.catalog.ToolCatalog;
+import org.apache.rocketmq.studio.ops.ai.tool.contract.common.MutationOutput;
+import org.apache.rocketmq.studio.ops.ai.tool.contract.plan.ToolPlan;
+import org.apache.rocketmq.studio.ops.ai.tool.filter.ToolAuditFilter;
+import org.apache.rocketmq.studio.ops.ai.tool.filter.ToolFilterChain;
+import org.apache.rocketmq.studio.ops.ai.tool.filter.ToolMutationFilter;
+import org.apache.rocketmq.studio.ops.ai.tool.handler.MutationToolHandler;
import org.apache.rocketmq.studio.ops.ai.tool.service.ToolExecutionService;
+import org.apache.rocketmq.studio.ops.ai.tool.service.ToolTokenService;
+import org.apache.rocketmq.studio.ops.audit.AuditService;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import java.util.List;
import java.util.Map;
+import java.util.Optional;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.ArgumentMatchers.same;
+import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -74,6 +92,64 @@ class McpToolRegistrarTest {
assertThat(specification.tool().annotations().destructiveHint()).isFalse();
}
+ @Test
+ void
returnsSuccessfulMcpResultWhenMutationCompletesBeforeAuditPersistenceFailsTest()
{
+ ToolDefinition definition = toolDefinition();
+ CountingMutationHandler handler = new CountingMutationHandler();
+
+ AuditService audit = mock(AuditService.class);
+ doThrow(new IllegalStateException("audit database unavailable"))
+ .when(audit).record(anyString(), anyString(), anyString(),
anyString(), isNull(), eq("SUCCESS"));
+ McpSchema.CallToolResult result = call(definition,
toolExecutor(definition, handler, audit));
+
+ assertThat(handler.executions).isEqualTo(1);
+ assertThat(result.isError()).isFalse();
+ MutationOutput<?> output = (MutationOutput<?>)
result.structuredContent();
+ assertThat(output.status()).isEqualTo(MutationOutput.Status.EXECUTED);
+ assertThat(output.result()).isEqualTo(Map.of("topic", "orders"));
+ }
+
+ @Test
+ void preservesOriginalMcpFailureWhenFailedAuditPersistenceAlsoFailsTest() {
+ ToolDefinition definition = toolDefinition();
+ CountingMutationHandler handler = new CountingMutationHandler();
+ ToolExecutionException originalFailure =
ToolError.TOOL_CAPABILITY_UNSUPPORTED.exception(definition.name());
+ handler.failure = originalFailure;
+
+ AuditService audit = mock(AuditService.class);
+ doThrow(new IllegalStateException("audit database unavailable"))
+ .when(audit).record(anyString(), anyString(), anyString(),
anyString(), anyString(), eq("FAILED"));
+ McpSchema.CallToolResult result = call(definition,
toolExecutor(definition, handler, audit));
+
+ assertThat(handler.executions).isEqualTo(1);
+ assertThat(result.isError()).isTrue();
+ assertThat(result.structuredContent()).isEqualTo(Map.of(
+ "code", originalFailure.getErrorCode(),
+ "message", originalFailure.getMessage(),
+ "hint", originalFailure.getHint()));
+ }
+
+ private ToolExecutionService toolExecutor(
+ ToolDefinition definition, CountingMutationHandler handler,
AuditService audit) {
+ ToolCatalog catalog = mock(ToolCatalog.class);
+
when(catalog.find(definition.name())).thenReturn(Optional.of(definition));
+ when(catalog.getDefinition(definition.name())).thenReturn(definition);
+ when(catalog.list()).thenReturn(List.of(definition));
+ ToolFilterChain filters = new ToolFilterChain(List.of(
+ new ToolAuditFilter(audit),
+ new ToolMutationFilter(mock(ToolTokenService.class), true)));
+ return new ToolExecutionService(catalog, List.of(handler), filters,
mock(InstanceResolver.class));
+ }
+
+ private McpSchema.CallToolResult call(ToolDefinition definition,
ToolExecutionService toolExecutor) {
+ return McpToolRegistrar.toolSpecification(definition, toolExecutor,
objectMapper)
+ .callHandler().apply(exchange(AUTHENTICATION), new
McpSchema.CallToolRequest(
+ definition.name(), Map.of(
+ "instanceId", AUTHENTICATION.instanceId(),
+ "topic", "orders",
+ "confirm_token", "confirmed")));
+ }
+
private static McpSyncServerExchange exchange(McpAuthentication
authentication) {
McpSyncServerExchange exchange = mock(McpSyncServerExchange.class);
McpTransportContext context = authentication == null
@@ -102,4 +178,36 @@ class McpToolRegistrarTest {
false,
null);
}
+
+ private static final class CountingMutationHandler extends
MutationToolHandler<Map, Map<String, Object>> {
+
+ private int executions;
+ private RuntimeException failure;
+
+ private CountingMutationHandler() {
+ super(Map.class);
+ }
+
+ @Override
+ public String name() {
+ return "rmq.topic.update";
+ }
+
+ @Override
+ public ToolPlan preview(Map input, ToolExecutionContext context) {
+ return ToolPlan.builder("Update topic")
+ .after(Map.of("topic", input.get("topic")))
+ .impact("Updates the topic configuration.")
+ .build();
+ }
+
+ @Override
+ public Map<String, Object> execute(Map input, ToolExecutionContext
context) {
+ executions++;
+ if (failure != null) {
+ throw failure;
+ }
+ return Map.of("topic", input.get("topic"));
+ }
+ }
}