diff --git a/mateclaw-server/src/main/java/vip/mate/channel/webchat/WebChatController.java b/mateclaw-server/src/main/java/vip/mate/channel/webchat/WebChatController.java index 58c3858a..b758a121 100644 --- a/mateclaw-server/src/main/java/vip/mate/channel/webchat/WebChatController.java +++ b/mateclaw-server/src/main/java/vip/mate/channel/webchat/WebChatController.java @@ -219,6 +219,17 @@ public class WebChatController { if (model != null) modelInfo[0] = model.toString(); if (provider != null) modelInfo[1] = provider.toString(); } + // Forward a curated subset of agent lifecycle events to the + // visitor SSE stream. The full event vocabulary (iteration_*, + // perf_summary, _routing_decision, feedback_event, ...) is + // internal — exposing it to 3rd-party websites would leak + // graph internals and complicate the SDK contract. The four + // types below are the ones that drive visible UX: typing + // indicator (phase), tool execution badges (tool_start/end), + // plan-execute checklist (plan). See docs/zh/webchat.md. + if (delta.isEvent()) { + forwardVisitorEvent(conversationId, delta.eventType(), delta.eventData()); + } if (delta.content() != null && !delta.content().isEmpty()) { assistantReply.append(delta.content()); if (!delta.persistenceOnly()) { @@ -1253,6 +1264,78 @@ public class WebChatController { } } + /** + * Forward a curated subset of agent lifecycle events to the visitor SSE + * stream as visitor-friendly {@code phase} / {@code tool_start} / + * {@code tool_end} / {@code plan} events. Internal event types + * ({@code _usage_final}, {@code _routing_decision}, {@code iteration_*}, + * {@code perf_summary}, {@code feedback_event}, {@code finish_reason}, + * {@code plan_step_*}) are dropped — they leak graph internals and have + * no visitor-facing value. + * + *

Tool arguments are deliberately not forwarded. The agent may + * invoke tools with PII / sensitive arguments (file paths, user queries, + * credentials); relaying those to a 3rd-party website frontend is a data + * leak. The frontend gets only the tool name and renders a localized + * label via its own lookup table. + * + *

Payloads are serialized via the injected {@link ObjectMapper} so + * nested maps/lists are encoded correctly (the hand-rolled {@link #escapeJson} + * helper is string-only). + * + *

Backward compat: visitors / SDKs that don't know these event types + * silently ignore them per the SSE spec. + */ + private void forwardVisitorEvent(String conversationId, String eventType, Map data) { + if (eventType == null || data == null) return; + Map payload; + String sseName; + switch (eventType) { + case "phase": + // Graph phase transition (planning / thinking / generating / + // summarizing / ...). Lets the SDK show a typing indicator + // before the first content_delta lands. + sseName = "phase"; + payload = Map.of( + "phase", String.valueOf(data.getOrDefault("phase", "")), + "timestamp", System.currentTimeMillis()); + break; + case "tool_call_started": + // Tool invocation started. Args intentionally omitted — see javadoc. + sseName = "tool_start"; + payload = Map.of( + "tool", String.valueOf(data.getOrDefault("toolName", + data.getOrDefault("tool", "")))); + break; + case "tool_call_completed": + // Tool invocation finished. Result content intentionally omitted. + sseName = "tool_end"; + payload = new java.util.LinkedHashMap<>(); + payload.put("tool", String.valueOf(data.getOrDefault("toolName", + data.getOrDefault("tool", "")))); + Object success = data.get("success"); + payload.put("success", success != null ? success : Boolean.TRUE); + break; + case "plan_created": + // Plan-Execute agents expose their step list. The SDK can render + // a checklist; subsequent plan_step_* events are dropped (too + // granular for a visitor view). + sseName = "plan"; + payload = Map.of("steps", data.getOrDefault("steps", List.of())); + break; + default: + // Curated allow-list: anything else is internal — silently drop. + return; + } + try { + String json = objectMapper.writeValueAsString(payload); + streamTracker.broadcast(conversationId, sseName, json); + } catch (Exception e) { + log.debug("[WebChat] Failed to serialize visitor event {} for {}: {}", + eventType, conversationId, e.getMessage()); + } + } + private String escapeJson(String value) { if (value == null) return "null"; return "\"" + value diff --git a/mateclaw-server/src/main/resources/docs/zh/webchat.md b/mateclaw-server/src/main/resources/docs/zh/webchat.md index ab25848d..9e900b65 100644 --- a/mateclaw-server/src/main/resources/docs/zh/webchat.md +++ b/mateclaw-server/src/main/resources/docs/zh/webchat.md @@ -74,6 +74,18 @@ MateClaw 的 WebChat 渠道让外部网站通过纯 HTTP / SSE 接入对话能 event: meta data: {"sessionId":"s1","conversationId":"webchat:abc123:v1:s1","visitorToken":"xxx.yyy"} +event: phase +data: {"phase":"planning","timestamp":1716700000000} + +event: tool_start +data: {"tool":"web_search"} + +event: tool_end +data: {"tool":"web_search","success":true} + +event: plan +data: {"steps":["search the web","summarize"]} + event: content_delta data: {"text":"你"} @@ -90,6 +102,24 @@ event: error data: {"message":"..."} (出错时) ``` +### 可选的实时进度事件 + +`phase` / `tool_start` / `tool_end` / `plan` 是**可选**事件 —— 用于在 +SDK 里展示"AI 正在打字..."气泡、工具执行徽章("正在搜索...")、Plan-Execute +步骤清单。SDK 可以全部忽略,只看 `content_delta` 也能完整渲染回复。 + +| 事件 | 触发时机 | 数据字段 | +|---|---|---| +| `phase` | agent 进入新的执行阶段(planning / generating / summarizing / ...) | `phase`, `timestamp` | +| `tool_start` | agent 调用工具 | `tool`(工具名) | +| `tool_end` | 工具调用完成 | `tool`, `success` | +| `plan` | Plan-Execute agent 拆解出步骤 | `steps`(字符串数组) | + +**注意**:`tool_start` / `tool_end` **只携带工具名**,不携带调用参数或返回 +结果 —— agent 工具调用可能涉及 PII(文件路径、用户查询、凭据),转发给 +第三方网站前端会有数据泄露风险。SDK 应基于工具名做本地化 label 映射 +(`web_search` → "正在搜索...")。 + ## 文件上传 / 下载 1. `POST /upload`(multipart):返回 `{fileId, fileName, contentType, size}`。 diff --git a/mateclaw-server/src/test/java/vip/mate/channel/webchat/WebChatStreamE2ETest.java b/mateclaw-server/src/test/java/vip/mate/channel/webchat/WebChatStreamE2ETest.java index a527c07f..9453abbc 100644 --- a/mateclaw-server/src/test/java/vip/mate/channel/webchat/WebChatStreamE2ETest.java +++ b/mateclaw-server/src/test/java/vip/mate/channel/webchat/WebChatStreamE2ETest.java @@ -373,4 +373,128 @@ class WebChatStreamE2ETest { assertThat(err.name).isEqualTo("error"); assertThat(err.data).contains("Invalid visitorId"); } + + // ------------------------------------------------------------------ + // Visitor-facing lifecycle events (phase / tool_start / tool_end / plan) + // ------------------------------------------------------------------ + + @Test + @DisplayName("phase event from agent → forwarded as SSE phase event (typing indicator)") + void forwardsPhaseEvent() throws Exception { + org.mockito.Mockito.when(agentService.chatStructuredStream( + eq(AGENT_ID), anyString(), anyString(), anyString(), isNull(), any())) + .thenReturn(Flux.just( + new AgentService.StreamDelta(null, null, "phase", + Map.of("phase", "planning", "timestamp", 1L), false), + new AgentService.StreamDelta(null, null, "phase", + Map.of("phase", "generating", "timestamp", 2L), false), + new AgentService.StreamDelta("ok", null))); + + List events = sendAndDrain( + streamPost(API_KEY, "{\"message\":\"hi\",\"visitorId\":\"vPhase\"}")); + + List phases = events.stream().filter(e -> "phase".equals(e.name)).toList(); + assertThat(phases).hasSize(2); + assertThat(phases.get(0).data).contains("\"phase\":\"planning\""); + assertThat(phases.get(1).data).contains("\"phase\":\"generating\""); + // Each carries a timestamp for client-side timeline rendering. + assertThat(phases.get(0).data).contains("\"timestamp\":"); + } + + @Test + @DisplayName("tool_call_started → tool_start SSE event; args are NOT leaked to the visitor") + void forwardsToolStartWithoutArgs() throws Exception { + org.mockito.Mockito.when(agentService.chatStructuredStream( + eq(AGENT_ID), anyString(), anyString(), anyString(), isNull(), any())) + .thenReturn(Flux.just( + new AgentService.StreamDelta(null, null, "tool_call_started", + Map.of("toolCallId", "call_1", + "toolName", "web_search", + "arguments", "secret query with PII", + "timestamp", 1L), false), + new AgentService.StreamDelta("done", null))); + + List events = sendAndDrain( + streamPost(API_KEY, "{\"message\":\"hi\",\"visitorId\":\"vTool\"}")); + + SseEvent toolStart = events.stream().filter(e -> "tool_start".equals(e.name)).findFirst().orElseThrow(); + assertThat(toolStart.data).contains("\"tool\":\"web_search\""); + // Critical: arguments must NOT be forwarded. + assertThat(toolStart.data).doesNotContain("secret query"); + assertThat(toolStart.data).doesNotContain("PII"); + assertThat(toolStart.data).doesNotContain("arguments"); + } + + @Test + @DisplayName("tool_call_completed → tool_end SSE event; result content is NOT leaked") + void forwardsToolEndWithoutResult() throws Exception { + org.mockito.Mockito.when(agentService.chatStructuredStream( + eq(AGENT_ID), anyString(), anyString(), anyString(), isNull(), any())) + .thenReturn(Flux.just( + new AgentService.StreamDelta(null, null, "tool_call_completed", + Map.of("toolCallId", "call_1", + "toolName", "web_search", + "result", "", + "success", true, + "timestamp", 1L), false), + new AgentService.StreamDelta("ack", null))); + + List events = sendAndDrain( + streamPost(API_KEY, "{\"message\":\"hi\",\"visitorId\":\"vToolEnd\"}")); + + SseEvent toolEnd = events.stream().filter(e -> "tool_end".equals(e.name)).findFirst().orElseThrow(); + assertThat(toolEnd.data).contains("\"tool\":\"web_search\""); + assertThat(toolEnd.data).contains("\"success\":true"); + // Result content is dropped. + assertThat(toolEnd.data).doesNotContain("huge internal result payload"); + assertThat(toolEnd.data).doesNotContain("\"result\""); + } + + @Test + @DisplayName("plan_created → plan SSE event with the step list") + void forwardsPlanEvent() throws Exception { + org.mockito.Mockito.when(agentService.chatStructuredStream( + eq(AGENT_ID), anyString(), anyString(), anyString(), isNull(), any())) + .thenReturn(Flux.just( + new AgentService.StreamDelta(null, null, "plan_created", + Map.of("planId", 42L, + "steps", List.of("search the web", "summarize"), + "timestamp", 1L), false), + new AgentService.StreamDelta("done", null))); + + List events = sendAndDrain( + streamPost(API_KEY, "{\"message\":\"hi\",\"visitorId\":\"vPlan\"}")); + + SseEvent plan = events.stream().filter(e -> "plan".equals(e.name)).findFirst().orElseThrow(); + assertThat(plan.data).contains("\"steps\":["); + assertThat(plan.data).contains("search the web"); + assertThat(plan.data).contains("summarize"); + } + + @Test + @DisplayName("internal event types (_routing_decision / perf_summary / iteration_* / ...) are silently dropped") + void dropsInternalEvents() throws Exception { + org.mockito.Mockito.when(agentService.chatStructuredStream( + eq(AGENT_ID), anyString(), anyString(), anyString(), isNull(), any())) + .thenReturn(Flux.just( + // Internal-only — must not produce an SSE event. + new AgentService.StreamDelta(null, null, "_routing_decision", + Map.of("sidecar", "vision"), false), + new AgentService.StreamDelta(null, null, "perf_summary", + Map.of("phase", "generate", "tokensPerSec", 42.0), false), + new AgentService.StreamDelta(null, null, "iteration_start", + Map.of("index", 0, "reason", "tool_call"), false), + new AgentService.StreamDelta(null, null, "finish_reason", + Map.of("reason", "STOP"), false), + new AgentService.StreamDelta("done", null))); + + List events = sendAndDrain( + streamPost(API_KEY, "{\"message\":\"hi\",\"visitorId\":\"vQuiet\"}")); + + // Only meta, content_delta (for "done"), and the terminal done event — + // none of the internal event types leaked. + List names = events.stream().map(e -> e.name).toList(); + assertThat(names).isNotEmpty(); + assertThat(names).doesNotContain("_routing_decision", "perf_summary", "iteration_start", "finish_reason"); + } }