mirror of
https://gitee.com/mateos/mateclaw.git
synced 2026-09-16 04:18:17 +08:00
fix(conversation): capture runtime model/provider with token usage in IM and webchat paths
This commit is contained in:
parent
727373f67c
commit
8d01396130
@ -732,8 +732,9 @@ public class ChannelMessageRouter {
|
|||||||
// for any Web SSE viewer of the same conversationId.
|
// for any Web SSE viewer of the same conversationId.
|
||||||
StringBuilder replyAccumulator = new StringBuilder();
|
StringBuilder replyAccumulator = new StringBuilder();
|
||||||
final String channelType = adapter.getChannelType();
|
final String channelType = adapter.getChannelType();
|
||||||
// Token usage: capture _usage_final event emitted at stream end
|
// Token usage + model attribution: capture _usage_final event emitted at stream end
|
||||||
final int[] usage = {0, 0}; // [promptTokens, completionTokens]
|
final int[] usage = {0, 0}; // [promptTokens, completionTokens]
|
||||||
|
final String[] modelInfo = {null, null}; // [runtimeModel, runtimeProvider]
|
||||||
agentService.chatStructuredStream(agentId, promptText, conversationId,
|
agentService.chatStructuredStream(agentId, promptText, conversationId,
|
||||||
message.getSenderId(), chatOrigin)
|
message.getSenderId(), chatOrigin)
|
||||||
.doOnNext(delta -> {
|
.doOnNext(delta -> {
|
||||||
@ -742,6 +743,10 @@ public class ChannelMessageRouter {
|
|||||||
Map<String, Object> data = delta.eventData();
|
Map<String, Object> data = delta.eventData();
|
||||||
usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue();
|
usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue();
|
||||||
usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue();
|
usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue();
|
||||||
|
Object model = data.get("runtimeModelName");
|
||||||
|
Object provider = data.get("runtimeProviderId");
|
||||||
|
if (model != null) modelInfo[0] = model.toString();
|
||||||
|
if (provider != null) modelInfo[1] = provider.toString();
|
||||||
}
|
}
|
||||||
mirrorPlanEventToTracker(conversationId, delta, channelType);
|
mirrorPlanEventToTracker(conversationId, delta, channelType);
|
||||||
} else if (delta.content() != null) {
|
} else if (delta.content() != null) {
|
||||||
@ -778,7 +783,7 @@ public class ChannelMessageRouter {
|
|||||||
String status = isError ? "error" : "completed";
|
String status = isError ? "error" : "completed";
|
||||||
MessageEntity saved = conversationService.saveMessage(
|
MessageEntity saved = conversationService.saveMessage(
|
||||||
conversationId, "assistant", reply, null, status,
|
conversationId, "assistant", reply, null, status,
|
||||||
usage[0], usage[1], null, null);
|
usage[0], usage[1], modelInfo[0], modelInfo[1]);
|
||||||
savedAssistantId = saved != null ? saved.getId() : null;
|
savedAssistantId = saved != null ? saved.getId() : null;
|
||||||
if (!isError) {
|
if (!isError) {
|
||||||
publishConversationCompletedEvent(agentId, conversationId, message.getContent(), reply);
|
publishConversationCompletedEvent(agentId, conversationId, message.getContent(), reply);
|
||||||
@ -892,13 +897,18 @@ public class ChannelMessageRouter {
|
|||||||
// only reads `delta.content()` and would otherwise eat plan_created /
|
// only reads `delta.content()` and would otherwise eat plan_created /
|
||||||
// plan_step_* events, leaving the Web Console mirror with no
|
// plan_step_* events, leaving the Web Console mirror with no
|
||||||
// PlanStepsPanel for IM-routed conversations.
|
// PlanStepsPanel for IM-routed conversations.
|
||||||
// Token usage: capture _usage_final event emitted at stream end
|
// Token usage + model attribution: capture _usage_final event emitted at stream end
|
||||||
final int[] usage = {0, 0}; // [promptTokens, completionTokens]
|
final int[] usage = {0, 0}; // [promptTokens, completionTokens]
|
||||||
|
final String[] modelInfo = {null, null}; // [runtimeModel, runtimeProvider]
|
||||||
Flux<AgentService.StreamDelta> mirroredStream = stream.doOnNext(delta -> {
|
Flux<AgentService.StreamDelta> mirroredStream = stream.doOnNext(delta -> {
|
||||||
if (delta.isEvent() && "_usage_final".equals(delta.eventType())) {
|
if (delta.isEvent() && "_usage_final".equals(delta.eventType())) {
|
||||||
Map<String, Object> data = delta.eventData();
|
Map<String, Object> data = delta.eventData();
|
||||||
usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue();
|
usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue();
|
||||||
usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue();
|
usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue();
|
||||||
|
Object model = data.get("runtimeModelName");
|
||||||
|
Object provider = data.get("runtimeProviderId");
|
||||||
|
if (model != null) modelInfo[0] = model.toString();
|
||||||
|
if (provider != null) modelInfo[1] = provider.toString();
|
||||||
}
|
}
|
||||||
mirrorPlanEventToTracker(conversationId, delta, channelType);
|
mirrorPlanEventToTracker(conversationId, delta, channelType);
|
||||||
});
|
});
|
||||||
@ -922,7 +932,7 @@ public class ChannelMessageRouter {
|
|||||||
String status = isError ? "error" : "completed";
|
String status = isError ? "error" : "completed";
|
||||||
MessageEntity saved = conversationService.saveMessage(
|
MessageEntity saved = conversationService.saveMessage(
|
||||||
conversationId, "assistant", finalContent, null, status,
|
conversationId, "assistant", finalContent, null, status,
|
||||||
usage[0], usage[1], null, null);
|
usage[0], usage[1], modelInfo[0], modelInfo[1]);
|
||||||
if (!isError) {
|
if (!isError) {
|
||||||
publishConversationCompletedEvent(agentId, conversationId, promptText, finalContent);
|
publishConversationCompletedEvent(agentId, conversationId, promptText, finalContent);
|
||||||
}
|
}
|
||||||
|
|||||||
@ -118,8 +118,9 @@ public class WebChatController {
|
|||||||
// Pattern mirrors ChatController: always accumulate, only broadcast when the
|
// Pattern mirrors ChatController: always accumulate, only broadcast when the
|
||||||
// delta is not a persistence-only echo of content already streamed by inner nodes.
|
// delta is not a persistence-only echo of content already streamed by inner nodes.
|
||||||
StringBuilder assistantReply = new StringBuilder();
|
StringBuilder assistantReply = new StringBuilder();
|
||||||
// Token usage: capture _usage_final event emitted at stream end
|
// Token usage + model attribution: capture _usage_final event emitted at stream end
|
||||||
final int[] usage = {0, 0}; // [promptTokens, completionTokens]
|
final int[] usage = {0, 0}; // [promptTokens, completionTokens]
|
||||||
|
final String[] modelInfo = {null, null}; // [runtimeModel, runtimeProvider]
|
||||||
|
|
||||||
agentService.chatStructuredStream(agentId, message, conversationId, visitorId)
|
agentService.chatStructuredStream(agentId, message, conversationId, visitorId)
|
||||||
.doOnNext(delta -> {
|
.doOnNext(delta -> {
|
||||||
@ -127,6 +128,10 @@ public class WebChatController {
|
|||||||
Map<String, Object> data = delta.eventData();
|
Map<String, Object> data = delta.eventData();
|
||||||
usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue();
|
usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue();
|
||||||
usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue();
|
usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue();
|
||||||
|
Object model = data.get("runtimeModelName");
|
||||||
|
Object provider = data.get("runtimeProviderId");
|
||||||
|
if (model != null) modelInfo[0] = model.toString();
|
||||||
|
if (provider != null) modelInfo[1] = provider.toString();
|
||||||
}
|
}
|
||||||
if (delta.content() != null && !delta.content().isEmpty()) {
|
if (delta.content() != null && !delta.content().isEmpty()) {
|
||||||
assistantReply.append(delta.content());
|
assistantReply.append(delta.content());
|
||||||
@ -147,7 +152,7 @@ public class WebChatController {
|
|||||||
if (!reply.isBlank()) {
|
if (!reply.isBlank()) {
|
||||||
conversationService.saveMessage(
|
conversationService.saveMessage(
|
||||||
conversationId, "assistant", reply, List.of(),
|
conversationId, "assistant", reply, List.of(),
|
||||||
"completed", usage[0], usage[1], null, null);
|
"completed", usage[0], usage[1], modelInfo[0], modelInfo[1]);
|
||||||
}
|
}
|
||||||
completionPublisher.publish(
|
completionPublisher.publish(
|
||||||
agentId, conversationId, message, reply, "webchat");
|
agentId, conversationId, message, reply, "webchat");
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user