diff --git a/mateclaw-server/src/main/java/vip/mate/agent/graph/executor/ToolExecutionExecutor.java b/mateclaw-server/src/main/java/vip/mate/agent/graph/executor/ToolExecutionExecutor.java index 6334b077..28730ddf 100644 --- a/mateclaw-server/src/main/java/vip/mate/agent/graph/executor/ToolExecutionExecutor.java +++ b/mateclaw-server/src/main/java/vip/mate/agent/graph/executor/ToolExecutionExecutor.java @@ -172,7 +172,7 @@ public class ToolExecutionExecutor { // 4. 分类: concurrencySafe boolean safe = isConcurrencySafe(toolName); - preparedCalls.add(new PreparedToolCall(toolCall, callback, arguments, safe, allResponses.size())); + preparedCalls.add(new PreparedToolCall(toolCall, callback, arguments, safe, allResponses.size(), conversationId)); // 占位,Phase 2 填充 allResponses.add(null); } @@ -235,6 +235,14 @@ public class ToolExecutionExecutor { private void executePreparedCalls(List preparedCalls, List allResponses, List events) { + if (!preparedCalls.isEmpty() && streamTracker != null) { + String conversationId = preparedCalls.get(0).conversationId; + String phase = classifyBatchPhase(preparedCalls); + streamTracker.updatePhase(conversationId, phase); + streamTracker.broadcastObject(conversationId, "phase", GraphEventPublisher.phase(phase, Map.of( + "toolCount", preparedCalls.size() + )).data()); + } // 分组: 连续的 safe 工具可以并行,遇到 unsafe 工具则先等待所有 safe 完成再独占执行 List> batches = buildExecutionBatches(preparedCalls); @@ -320,6 +328,11 @@ public class ToolExecutionExecutor { List events) { String toolName = pc.toolCall.name(); try { + if (streamTracker != null) { + streamTracker.updateRunningTool(pc.conversationId, toolName); + streamTracker.broadcastObject(pc.conversationId, GraphEventPublisher.EVENT_TOOL_START, + GraphEventPublisher.toolStart(toolName, pc.arguments).data()); + } log.info("[ToolExecutor] Executing tool: {} with args: {}", toolName, pc.arguments != null && pc.arguments.length() > 200 ? pc.arguments.substring(0, 200) + "..." : pc.arguments); @@ -338,12 +351,22 @@ public class ToolExecutionExecutor { log.info("[ToolExecutor] Tool {} returned {} chars", toolName, rawLen); } events.add(GraphEventPublisher.toolComplete(toolName, result, true)); + if (streamTracker != null) { + streamTracker.broadcastObject(pc.conversationId, GraphEventPublisher.EVENT_TOOL_COMPLETE, + GraphEventPublisher.toolComplete(toolName, result, true).data()); + streamTracker.updateRunningTool(pc.conversationId, null); + } return new ToolResponseMessage.ToolResponse( pc.toolCall.id(), toolName, result != null ? result : ""); } catch (Exception e) { log.error("[ToolExecutor] Tool {} execution failed: {}", toolName, e.getMessage(), e); String normalizedError = normalizeToolExecutionError(e); events.add(GraphEventPublisher.toolComplete(toolName, normalizedError, false)); + if (streamTracker != null) { + streamTracker.broadcastObject(pc.conversationId, GraphEventPublisher.EVENT_TOOL_COMPLETE, + GraphEventPublisher.toolComplete(toolName, normalizedError, false).data()); + streamTracker.updateRunningTool(pc.conversationId, null); + } return new ToolResponseMessage.ToolResponse( pc.toolCall.id(), toolName, normalizedError); } @@ -408,6 +431,13 @@ public class ToolExecutionExecutor { return !DEFAULT_UNSAFE_TOOLS.contains(toolName); } + private String classifyBatchPhase(List preparedCalls) { + boolean memoryOnly = preparedCalls.stream().allMatch(pc -> + "read_workspace_memory_file".equals(pc.toolCall.name()) + || "list_workspace_memory_files".equals(pc.toolCall.name())); + return memoryOnly ? "reading_memory" : "executing_tool"; + } + private String normalizeToolExecutionError(Exception e) { String message = e != null && e.getMessage() != null ? e.getMessage() : "未知错误"; String lower = message.toLowerCase(Locale.ROOT); @@ -444,7 +474,8 @@ public class ToolExecutionExecutor { ToolCallback callback, String arguments, boolean concurrencySafe, - int resultIndex + int resultIndex, + String conversationId ) {} private record ApprovalBarrier(String pendingId, String toolName) {} diff --git a/mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java b/mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java index d948ea21..f63ab15e 100644 --- a/mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java +++ b/mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java @@ -206,6 +206,10 @@ public class ReasoningNode implements NodeAction { GraphEventPublisher.GraphEvent phaseEvent = GraphEventPublisher.phase("reasoning", Map.of("iteration", accessor.iterationCount())); + pushPhase(conversationId, "reasoning", Map.of( + "iteration", accessor.iterationCount(), + "llmCallCount", nextLlmCallCount + )); NodeStreamingChatHelper.StreamResult result; try { @@ -225,6 +229,11 @@ public class ReasoningNode implements NodeAction { messages.size(), compactedMessages.size()); // compact retry 是第 2 次 LLM 调用,先递增再调用 nextLlmCallCount++; + pushPhase(conversationId, "reasoning", Map.of( + "iteration", accessor.iterationCount(), + "llmCallCount", nextLlmCallCount, + "compacted", true + )); result = streamingHelper.streamCall(chatModel, retryPrompt, conversationId, "reasoning_compact_retry"); } else { log.warn("[ReasoningNode] Compaction did not reduce messages, cannot retry"); @@ -296,6 +305,10 @@ public class ReasoningNode implements NodeAction { log.info("[ReasoningNode] LLM requested {} tool call(s): {}", result.toolCalls().size(), result.toolCalls().stream().map(AssistantMessage.ToolCall::name).toList()); + pushPhase(conversationId, "executing_tool", Map.of( + "iteration", accessor.iterationCount(), + "toolCount", result.toolCalls().size() + )); return MateClawStateAccessor.output() .needsToolCall(true) @@ -315,6 +328,10 @@ public class ReasoningNode implements NodeAction { } else { String content = result.text(); log.info("[ReasoningNode] LLM produced final answer ({} chars)", content != null ? content.length() : 0); + pushPhase(conversationId, "drafting_answer", Map.of( + "iteration", accessor.iterationCount(), + "answerChars", content != null ? content.length() : 0 + )); return MateClawStateAccessor.output() .needsToolCall(false) @@ -347,4 +364,12 @@ public class ReasoningNode implements NodeAction { throw new RuntimeException("无法反序列化 forced_tool_call: " + e.getMessage(), e); } } + + private void pushPhase(String conversationId, String phase, Map extra) { + if (streamTracker == null || !StringUtils.hasText(conversationId)) { + return; + } + streamTracker.updatePhase(conversationId, phase); + streamTracker.broadcastObject(conversationId, "phase", GraphEventPublisher.phase(phase, extra).data()); + } } diff --git a/mateclaw-server/src/main/java/vip/mate/agent/graph/node/SummarizingNode.java b/mateclaw-server/src/main/java/vip/mate/agent/graph/node/SummarizingNode.java index 6095af3c..018c9589 100644 --- a/mateclaw-server/src/main/java/vip/mate/agent/graph/node/SummarizingNode.java +++ b/mateclaw-server/src/main/java/vip/mate/agent/graph/node/SummarizingNode.java @@ -82,6 +82,10 @@ public class SummarizingNode implements NodeAction { log.info("[SummarizingNode] Summarizing {} observations ({} total chars) for user query", observations.size(), accessor.totalObservationChars()); + pushPhase(conversationId, "summarizing_observations", Map.of( + "observationCount", observations.size(), + "summaryChars", accessor.totalObservationChars() + )); // 构建 summarize prompt StringBuilder observationText = new StringBuilder(); @@ -176,4 +180,12 @@ public class SummarizingNode implements NodeAction { "summaryChars", summaryContent.length())))) .build(); } + + private void pushPhase(String conversationId, String phase, Map extra) { + if (streamTracker == null || conversationId == null || conversationId.isEmpty()) { + return; + } + streamTracker.updatePhase(conversationId, phase); + streamTracker.broadcastObject(conversationId, "phase", GraphEventPublisher.phase(phase, extra).data()); + } } diff --git a/mateclaw-server/src/main/java/vip/mate/channel/web/ChatController.java b/mateclaw-server/src/main/java/vip/mate/channel/web/ChatController.java index d6566c63..53afa1f1 100644 --- a/mateclaw-server/src/main/java/vip/mate/channel/web/ChatController.java +++ b/mateclaw-server/src/main/java/vip/mate/channel/web/ChatController.java @@ -1218,8 +1218,14 @@ public class ChatController { runtimeProviderId = String.valueOf(data.getOrDefault("runtimeProviderId", "")); return; } + if ("phase".equals(delta.eventType())) { + String phase = String.valueOf(delta.eventData().getOrDefault("phase", "")); + if (!phase.isBlank()) { + streamTracker.updatePhase(conversationId, phase); + } + } // 累积工具调用事件,用于持久化到消息历史 - accumulateToolEvent(delta.eventType(), delta.eventData()); + accumulateToolEvent(delta.eventType(), delta.eventData(), conversationId); try { broadcastEvent(conversationId, delta.eventType(), delta.eventData()); } catch (Exception e) { @@ -1229,6 +1235,7 @@ public class ChatController { } if (delta.content() != null && !delta.content().isBlank()) { content.append(delta.content()); + streamTracker.updatePhase(conversationId, "drafting_answer"); if (!delta.persistenceOnly()) { broadcastEvent(conversationId, "content_delta", Map.of("delta", delta.content())); } @@ -1243,9 +1250,10 @@ public class ChatController { boolean isAwaitingApproval() { return awaitingApproval; } - private void accumulateToolEvent(String eventType, Map data) { + private void accumulateToolEvent(String eventType, Map data, String conversationId) { if ("tool_approval_requested".equals(eventType)) { awaitingApproval = true; + streamTracker.updatePhase(conversationId, "awaiting_approval"); } else if ("tool_call_started".equals(eventType)) { Map tc = new LinkedHashMap<>(); tc.put("name", data.getOrDefault("toolName", "")); diff --git a/mateclaw-ui/src/components/chat/StreamLoadingBar.vue b/mateclaw-ui/src/components/chat/StreamLoadingBar.vue index 3feedf0f..857871d8 100644 --- a/mateclaw-ui/src/components/chat/StreamLoadingBar.vue +++ b/mateclaw-ui/src/components/chat/StreamLoadingBar.vue @@ -2,8 +2,12 @@
{{ phaseIcon }} - {{ statusText }} - {{ runningToolName }} +
+ {{ statusText }} + {{ runningToolName }} + {{ statusDetail }} + {{ slowHint }} +
@@ -25,7 +29,7 @@