From 727373f67cdee060b59a9ea91eec424011b4c2a1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=80=AA=E7=A8=8B=E4=BC=9F?= Date: Mon, 25 May 2026 20:25:22 +0800 Subject: [PATCH] fix(conversation): capture token usage in IM channel and webchat paths (#217) IM channels (Feishu, DingTalk, WeCom, etc.) and the WebChat widget were calling saveMessage without token usage parameters, causing promptTokens and completionTokens to default to 0. This made the Token Statistics module report significantly lower numbers than actual usage. Root cause: the _usage_final event (containing promptTokens / completionTokens) emitted by the agent graph at stream end was not being captured in these paths, unlike ChatController's StreamAccumulator which already handles it correctly. Fix: capture _usage_final events in doOnNext handlers for: - ChannelMessageRouter sync path (non-streaming IM adapters) - ChannelMessageRouter streaming path (DingTalk, etc.) - WebChatController SSE stream Refs #214 (remaining String-API paths covered by follow-up). --- .../mate/channel/ChannelMessageRouter.java | 25 ++++++++++++++++--- .../channel/webchat/WebChatController.java | 10 +++++++- 2 files changed, 30 insertions(+), 5 deletions(-) diff --git a/mateclaw-server/src/main/java/vip/mate/channel/ChannelMessageRouter.java b/mateclaw-server/src/main/java/vip/mate/channel/ChannelMessageRouter.java index ef0d33f0..5abd02a5 100644 --- a/mateclaw-server/src/main/java/vip/mate/channel/ChannelMessageRouter.java +++ b/mateclaw-server/src/main/java/vip/mate/channel/ChannelMessageRouter.java @@ -732,10 +732,17 @@ public class ChannelMessageRouter { // for any Web SSE viewer of the same conversationId. StringBuilder replyAccumulator = new StringBuilder(); final String channelType = adapter.getChannelType(); + // Token usage: capture _usage_final event emitted at stream end + final int[] usage = {0, 0}; // [promptTokens, completionTokens] agentService.chatStructuredStream(agentId, promptText, conversationId, message.getSenderId(), chatOrigin) .doOnNext(delta -> { if (delta.isEvent()) { + if ("_usage_final".equals(delta.eventType())) { + Map data = delta.eventData(); + usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue(); + usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue(); + } mirrorPlanEventToTracker(conversationId, delta, channelType); } else if (delta.content() != null) { // Match the legacy agentService.chat() behavior: include @@ -770,7 +777,8 @@ public class ChannelMessageRouter { boolean isError = errorClassifier.isErrorReply(reply); String status = isError ? "error" : "completed"; MessageEntity saved = conversationService.saveMessage( - conversationId, "assistant", reply, null, status); + conversationId, "assistant", reply, null, status, + usage[0], usage[1], null, null); savedAssistantId = saved != null ? saved.getId() : null; if (!isError) { publishConversationCompletedEvent(agentId, conversationId, message.getContent(), reply); @@ -884,8 +892,16 @@ public class ChannelMessageRouter { // only reads `delta.content()` and would otherwise eat plan_created / // plan_step_* events, leaving the Web Console mirror with no // PlanStepsPanel for IM-routed conversations. - Flux mirroredStream = stream.doOnNext(delta -> - mirrorPlanEventToTracker(conversationId, delta, channelType)); + // Token usage: capture _usage_final event emitted at stream end + final int[] usage = {0, 0}; // [promptTokens, completionTokens] + Flux mirroredStream = stream.doOnNext(delta -> { + if (delta.isEvent() && "_usage_final".equals(delta.eventType())) { + Map data = delta.eventData(); + usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue(); + usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue(); + } + mirrorPlanEventToTracker(conversationId, delta, channelType); + }); // Step 2: 委托渠道渲染(渠道内部消费 Flux 并处理 UI 更新) String finalContent = streamingAdapter.processStream(mirroredStream, message, conversationId); @@ -905,7 +921,8 @@ public class ChannelMessageRouter { boolean isError = errorClassifier.isErrorReply(finalContent); String status = isError ? "error" : "completed"; MessageEntity saved = conversationService.saveMessage( - conversationId, "assistant", finalContent, null, status); + conversationId, "assistant", finalContent, null, status, + usage[0], usage[1], null, null); if (!isError) { publishConversationCompletedEvent(agentId, conversationId, promptText, finalContent); } 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 4cd0f572..4f3ea9ff 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 @@ -118,9 +118,16 @@ public class WebChatController { // Pattern mirrors ChatController: always accumulate, only broadcast when the // delta is not a persistence-only echo of content already streamed by inner nodes. StringBuilder assistantReply = new StringBuilder(); + // Token usage: capture _usage_final event emitted at stream end + final int[] usage = {0, 0}; // [promptTokens, completionTokens] agentService.chatStructuredStream(agentId, message, conversationId, visitorId) .doOnNext(delta -> { + if (delta.isEvent() && "_usage_final".equals(delta.eventType())) { + Map data = delta.eventData(); + usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue(); + usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue(); + } if (delta.content() != null && !delta.content().isEmpty()) { assistantReply.append(delta.content()); if (!delta.persistenceOnly()) { @@ -139,7 +146,8 @@ public class WebChatController { try { if (!reply.isBlank()) { conversationService.saveMessage( - conversationId, "assistant", reply, List.of()); + conversationId, "assistant", reply, List.of(), + "completed", usage[0], usage[1], null, null); } completionPublisher.publish( agentId, conversationId, message, reply, "webchat");