mirror of
https://gitee.com/mateos/mateclaw.git
synced 2026-09-13 03:13:41 +08:00
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).
This commit is contained in:
parent
07eb625d11
commit
727373f67c
@ -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<String, Object> 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<AgentService.StreamDelta> 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<AgentService.StreamDelta> mirroredStream = stream.doOnNext(delta -> {
|
||||
if (delta.isEvent() && "_usage_final".equals(delta.eventType())) {
|
||||
Map<String, Object> 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);
|
||||
}
|
||||
|
||||
@ -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<String, Object> 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");
|
||||
|
||||
Loading…
Reference in New Issue
Block a user