fix(agent): stabilize stream token output
This commit is contained in:
+8
-4
@@ -4,7 +4,7 @@
|
||||
"workspaces": {
|
||||
"": {
|
||||
"dependencies": {
|
||||
"@opencode-ai/plugin": "^1.16.2",
|
||||
"@opencode-ai/plugin": "^1.17.13",
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^24.7.2",
|
||||
@@ -13,6 +13,8 @@
|
||||
},
|
||||
},
|
||||
"packages": {
|
||||
"@ai-sdk/provider": ["@ai-sdk/provider@3.0.8", "", { "dependencies": { "json-schema": "^0.4.0" } }, "sha512-oGMAgGoQdBXbZqNG0Ze56CHjDZ1IDYOwGYxYjO5KLSlz5HiNQ9udIXsPZ61VWaHGZ5XW/jyjmr6t2xz2jGVwbQ=="],
|
||||
|
||||
"@msgpackr-extract/msgpackr-extract-darwin-arm64": ["@msgpackr-extract/msgpackr-extract-darwin-arm64@3.0.4", "", { "os": "darwin", "cpu": "arm64" }, "sha512-LCkGo6JDfaBhgST7UpPWgNgLINpcpabaHfyz5OBx75nUYxBsaEPxjnyNjWpeb/xBup/682QnBfRBy2/LvPutZQ=="],
|
||||
|
||||
"@msgpackr-extract/msgpackr-extract-darwin-x64": ["@msgpackr-extract/msgpackr-extract-darwin-x64@3.0.4", "", { "os": "darwin", "cpu": "x64" }, "sha512-zExlW9zUJKZH/tOtVMttwjKa4Xm/3KcNjnE3dPN92uCktwavMxpgCA3MoJK/DOnTWsQgo224OaST27/mPNAf+w=="],
|
||||
@@ -25,9 +27,9 @@
|
||||
|
||||
"@msgpackr-extract/msgpackr-extract-win32-x64": ["@msgpackr-extract/msgpackr-extract-win32-x64@3.0.4", "", { "os": "win32", "cpu": "x64" }, "sha512-CmCXPQrkbwExx3j946/PtHWHbYJiCRBRDl4BlkRQcJB/YOwQxJRTpoo7aTsortjgoJ1x7opzTSxn7C+ASSLVjQ=="],
|
||||
|
||||
"@opencode-ai/plugin": ["@opencode-ai/plugin@1.16.2", "", { "dependencies": { "@opencode-ai/sdk": "1.16.2", "effect": "4.0.0-beta.74", "zod": "4.1.8" }, "peerDependencies": { "@opentui/core": ">=0.3.2", "@opentui/keymap": ">=0.3.2", "@opentui/solid": ">=0.3.2" }, "optionalPeers": ["@opentui/core", "@opentui/keymap", "@opentui/solid"] }, "sha512-FaZhVXrbz93xsdGLCtarRDTeqFt8AkLfh8B34tFBj6G4HXVmKSgBwVXmtELKKC+08xMtawBC9hshiMbXryv6cg=="],
|
||||
"@opencode-ai/plugin": ["@opencode-ai/plugin@1.17.13", "", { "dependencies": { "@ai-sdk/provider": "3.0.8", "@opencode-ai/sdk": "1.17.13", "effect": "4.0.0-beta.83", "zod": "4.1.8" }, "peerDependencies": { "@opentui/core": ">=0.3.4", "@opentui/keymap": ">=0.3.4", "@opentui/solid": ">=0.3.4" }, "optionalPeers": ["@opentui/core", "@opentui/keymap", "@opentui/solid"] }, "sha512-JjoIs6qTPl0IrXo+J1GTqX9lcb2h2dMoW6BWMR2IbTT3MY2FQ/lvBXDzkx0d3zVRaYN+NTmV4u8Nth396xq+sQ=="],
|
||||
|
||||
"@opencode-ai/sdk": ["@opencode-ai/sdk@1.16.2", "", { "dependencies": { "cross-spawn": "7.0.6" } }, "sha512-Z/xZ7q79dYeE0afqIk/yFEcRNGEQFcE+H8ssYivUiy+xGZ1mGwT72jpaQZKBwPn3JH4sRCu4KA2lcktBQfcOjg=="],
|
||||
"@opencode-ai/sdk": ["@opencode-ai/sdk@1.17.13", "", { "dependencies": { "cross-spawn": "7.0.6" } }, "sha512-VItOGjMzRQx3zypwmeFLNhCiIx32kxS7FqzIJvVZLfyNGCifs3rfGC9qzNKWcxQo4SjNvAw++v4gWWU6Inv+JQ=="],
|
||||
|
||||
"@standard-schema/spec": ["@standard-schema/spec@1.1.0", "", {}, "sha512-l2aFy5jALhniG5HgqrD6jXLi/rUWrKvqN/qJx6yoJsgKhblVd+iqqU4RCXavm/jPityDo5TCvKMnpjKnOriy0w=="],
|
||||
|
||||
@@ -37,7 +39,7 @@
|
||||
|
||||
"detect-libc": ["detect-libc@2.1.2", "", {}, "sha512-Btj2BOOO83o3WyH59e8MgXsxEQVcarkUOpEYrubB0urwnN10yQ364rsiByU11nZlqWYZm05i/of7io4mzihBtQ=="],
|
||||
|
||||
"effect": ["effect@4.0.0-beta.74", "", { "dependencies": { "@standard-schema/spec": "^1.1.0", "fast-check": "^4.8.0", "find-my-way-ts": "^0.1.6", "ini": "^7.0.0", "kubernetes-types": "^1.30.0", "msgpackr": "^2.0.1", "multipasta": "^0.2.7", "toml": "^4.1.1", "uuid": "^14.0.0", "yaml": "^2.9.0" } }, "sha512-Yx+Kh12U+i2FmjwEfKs+ePFmpMd43RPD1oGqc/VraSS9bYzvF0Ff3PojwEFEVEewp8xc92Uxu28gTspU4qyvHA=="],
|
||||
"effect": ["effect@4.0.0-beta.83", "", { "dependencies": { "@standard-schema/spec": "^1.1.0", "fast-check": "^4.8.0", "find-my-way-ts": "^0.1.6", "ini": "^7.0.0", "kubernetes-types": "^1.30.0", "msgpackr": "^2.0.1", "multipasta": "^0.2.7", "toml": "^4.1.1", "uuid": "^14.0.0", "yaml": "^2.9.0" } }, "sha512-0wsak8RtgGAr9UWSbVDgJHZcUqMSvicHcvaZv1MbMM7MCGgW4Rn/137J1MHQbwYPcwYGxT/IqehFd+UbYuj78w=="],
|
||||
|
||||
"fast-check": ["fast-check@4.8.0", "", { "dependencies": { "pure-rand": "^8.0.0" } }, "sha512-GOJ158CUMnN6cSahsv4+ExARvIDuzzinFjkp0E9WtiBa5zcVeLozVkWaE4IzFcc+Y48Wp1EDlUZsXRyAztQcSg=="],
|
||||
|
||||
@@ -47,6 +49,8 @@
|
||||
|
||||
"isexe": ["isexe@2.0.0", "", {}, "sha512-RHxMLp9lnKHGHRng9QFhRCMbYAcVpn69smSGcq3f36xjgVVWThj4qqLbTLlq7Ssj8B+fIQ1EuCEGI2lKsyQeIw=="],
|
||||
|
||||
"json-schema": ["json-schema@0.4.0", "", {}, "sha512-es94M3nTIfsEPisRafak+HDLfHXnKBhV3vU5eqPcS3flIWqcxJWgXHXiey3YrpaNsanY5ei1VoYEbOzijuq9BA=="],
|
||||
|
||||
"kubernetes-types": ["kubernetes-types@1.30.0", "", {}, "sha512-Dew1okvhM/SQcIa2rcgujNndZwU8VnSapDgdxlYoB84ZlpAD43U6KLAFqYo17ykSFGHNPrg0qry0bP+GJd9v7Q=="],
|
||||
|
||||
"msgpackr": ["msgpackr@2.0.2", "", { "optionalDependencies": { "msgpackr-extract": "^3.0.4" } }, "sha512-c5hYOXFbP79Slh6Dzd2wzk+jnV7mX1UxfMYtilnY1NmalXPqG8DGb5cYCMBrW4AsH3zekBBZd4QrKz9NhtvYLQ=="],
|
||||
|
||||
@@ -4,7 +4,7 @@
|
||||
"typecheck": "tsc --noEmit -p tsconfig.json"
|
||||
},
|
||||
"dependencies": {
|
||||
"@opencode-ai/plugin": "^1.16.2"
|
||||
"@opencode-ai/plugin": "^1.17.13"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^24.7.2",
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
"": {
|
||||
"name": "tjwater-agent",
|
||||
"dependencies": {
|
||||
"@opencode-ai/sdk": "^1.16.2",
|
||||
"@opencode-ai/sdk": "^1.17.13",
|
||||
"ai": "7.0.8",
|
||||
"cors": "^2.8.5",
|
||||
"dotenv": "^17.2.3",
|
||||
@@ -30,7 +30,7 @@
|
||||
|
||||
"@ai-sdk/provider-utils": ["@ai-sdk/provider-utils@5.0.2", "", { "dependencies": { "@ai-sdk/provider": "4.0.1", "@standard-schema/spec": "^1.1.0", "@workflow/serde": "4.1.0", "eventsource-parser": "^3.0.8" }, "peerDependencies": { "zod": "^3.25.76 || ^4.1.8" } }, "sha512-EcmdjJb7yggsZPCbS3MFBpvAUnKaPW+QvanU5GzF00XCq0bqqAmvJ3MN19ejlmOETbW8sJNiq6qam48wTcbUNw=="],
|
||||
|
||||
"@opencode-ai/sdk": ["@opencode-ai/sdk@1.16.2", "", { "dependencies": { "cross-spawn": "7.0.6" } }, "sha512-Z/xZ7q79dYeE0afqIk/yFEcRNGEQFcE+H8ssYivUiy+xGZ1mGwT72jpaQZKBwPn3JH4sRCu4KA2lcktBQfcOjg=="],
|
||||
"@opencode-ai/sdk": ["@opencode-ai/sdk@1.17.13", "", { "dependencies": { "cross-spawn": "7.0.6" } }, "sha512-VItOGjMzRQx3zypwmeFLNhCiIx32kxS7FqzIJvVZLfyNGCifs3rfGC9qzNKWcxQo4SjNvAw++v4gWWU6Inv+JQ=="],
|
||||
|
||||
"@pinojs/redact": ["@pinojs/redact@0.4.0", "", {}, "sha512-k2ENnmBugE/rzQfEcdWHcCY+/FM3VLzH9cYEsbdsoqrvzAKRhUZeRNhAZvB8OitQJ1TBed3yqWtdjzS6wJKBwg=="],
|
||||
|
||||
|
||||
+1
-1
@@ -15,7 +15,7 @@
|
||||
"start:prod": "bun run check && bun src/server.ts"
|
||||
},
|
||||
"dependencies": {
|
||||
"@opencode-ai/sdk": "^1.16.2",
|
||||
"@opencode-ai/sdk": "^1.17.13",
|
||||
"ai": "7.0.8",
|
||||
"cors": "^2.8.5",
|
||||
"dotenv": "^17.2.3",
|
||||
|
||||
+20
-1
@@ -764,6 +764,15 @@ export const buildChatRouter = (
|
||||
textPart.state = "streaming";
|
||||
}
|
||||
writeUiChunk({ type: "text-delta", id: textPartId, delta });
|
||||
writeUiChunk({
|
||||
type: "data-stream_token",
|
||||
data: {
|
||||
session_id: clientSessionId,
|
||||
message_id: assistantUiMessageId,
|
||||
content: delta,
|
||||
},
|
||||
transient: true,
|
||||
} as UIMessageChunk);
|
||||
syncLatestUiMessages();
|
||||
};
|
||||
const writeDataPart = (event: string, data: Record<string, unknown>) => {
|
||||
@@ -879,6 +888,16 @@ export const buildChatRouter = (
|
||||
res.on("close", handleClientClose);
|
||||
|
||||
const publish = (event: string, data: Record<string, unknown>) => {
|
||||
const subscriberData =
|
||||
event === "token"
|
||||
? {
|
||||
...data,
|
||||
message_id:
|
||||
typeof data.message_id === "string"
|
||||
? data.message_id
|
||||
: assistantUiMessageId,
|
||||
}
|
||||
: data;
|
||||
if (event === "token") {
|
||||
activeRun.messages = updateLastAssistantMessage(activeRun.messages, (message) => ({
|
||||
...message,
|
||||
@@ -1011,7 +1030,7 @@ export const buildChatRouter = (
|
||||
}
|
||||
|
||||
for (const subscriber of activeRun.subscribers) {
|
||||
subscriber.write(event, data);
|
||||
subscriber.write(event, subscriberData);
|
||||
}
|
||||
void queueSessionUiStatePersist().catch((error) => {
|
||||
logger.warn({ err: error, sessionId: clientSessionId }, "failed to persist chat stream state");
|
||||
|
||||
@@ -99,8 +99,6 @@ type SmoothTextEndPart = Extract<SmoothTextStreamPart, { type: "text-end" }>;
|
||||
const segmenters = new Map<string, Intl.Segmenter | null>();
|
||||
const CJK_SCRIPT_PATTERN = /[\p{Script=Han}\p{Script=Hiragana}\p{Script=Katakana}\p{Script=Hangul}]/u;
|
||||
const CJK_PUNCTUATION_PATTERN = /[。!?!?,、;:,;:]/u;
|
||||
const CJK_SENTENCE_END_PATTERN = /[。!?!?]+(?:["'”’』」))\]\s]+)?/u;
|
||||
const CJK_PAUSE_PATTERN = /[,、;:,;:]+(?:\s+)?/u;
|
||||
const CJK_FALLBACK_TRIGGER_CHARS = 24;
|
||||
const CJK_FALLBACK_MIN_CHARS = 12;
|
||||
const CJK_FALLBACK_TARGET_CHARS = 18;
|
||||
@@ -133,20 +131,6 @@ function isCjkSmoothingContent(content: string, locale: string) {
|
||||
);
|
||||
}
|
||||
|
||||
function detectCjkPunctuationChunk(content: string) {
|
||||
const sentenceMatch = CJK_SENTENCE_END_PATTERN.exec(content);
|
||||
if (sentenceMatch) {
|
||||
return content.slice(0, sentenceMatch.index + sentenceMatch[0].length);
|
||||
}
|
||||
|
||||
const pauseMatch = CJK_PAUSE_PATTERN.exec(content);
|
||||
if (pauseMatch) {
|
||||
return content.slice(0, pauseMatch.index + pauseMatch[0].length);
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
function detectCjkFallbackChunk(content: string, locale: string) {
|
||||
if (codePointLength(content) < CJK_FALLBACK_TRIGGER_CHARS) {
|
||||
return null;
|
||||
@@ -179,10 +163,7 @@ export function detectTokenSmoothingChunk(content: string, locale = "zh") {
|
||||
}
|
||||
|
||||
if (isCjkSmoothingContent(content, locale)) {
|
||||
return (
|
||||
detectCjkPunctuationChunk(content) ??
|
||||
detectCjkFallbackChunk(content, locale)
|
||||
);
|
||||
return detectCjkFallbackChunk(content, locale);
|
||||
}
|
||||
|
||||
const wordChunk = content.match(/^\S+\s*/u)?.[0];
|
||||
|
||||
@@ -19,14 +19,9 @@ const createEventStream = (events: unknown[]) => ({
|
||||
const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));
|
||||
|
||||
describe("streamPromptResponse", () => {
|
||||
it("detects Chinese punctuation chunks for smooth token output", () => {
|
||||
expect(detectTokenSmoothingChunk("管网压力异常,需要继续分析", "zh")).toBe(
|
||||
"管网压力异常,",
|
||||
);
|
||||
|
||||
expect(detectTokenSmoothingChunk("管网压力异常需要继续分析。后续排查阀门", "zh")).toBe(
|
||||
"管网压力异常需要继续分析。",
|
||||
);
|
||||
it("does not split short Chinese text only because it has punctuation", () => {
|
||||
expect(detectTokenSmoothingChunk("管网压力异常,需要继续分析", "zh")).toBeNull();
|
||||
expect(detectTokenSmoothingChunk("管网压力异常需要继续分析。", "zh")).toBeNull();
|
||||
});
|
||||
|
||||
it("does not immediately release a single Chinese character", () => {
|
||||
@@ -67,7 +62,7 @@ describe("streamPromptResponse", () => {
|
||||
expect(events.map((item) => item.data.content)).toEqual(["第一段", "第二段"]);
|
||||
});
|
||||
|
||||
it("buffers Chinese single-character deltas until punctuation", async () => {
|
||||
it("buffers Chinese single-character deltas until a bounded chunk is available", async () => {
|
||||
const events: Array<{ event: string; data: Record<string, unknown> }> = [];
|
||||
const smoother = createTokenSmoother({
|
||||
enabled: true,
|
||||
@@ -81,13 +76,15 @@ describe("streamPromptResponse", () => {
|
||||
await sleep(8);
|
||||
expect(events).toHaveLength(0);
|
||||
|
||||
for (const char of Array.from("网压力异常,")) {
|
||||
const content = "网压力异常需要继续分析东部主干供水走廊和相关阀门状态";
|
||||
for (const char of Array.from(content)) {
|
||||
smoother.writeToken(char);
|
||||
}
|
||||
await sleep(20);
|
||||
|
||||
expect(events.map((item) => item.data.content).join("")).toBe("管网压力异常,");
|
||||
expect(events.length).toBeGreaterThan(0);
|
||||
await smoother.flush();
|
||||
expect(events.map((item) => item.data.content).join("")).toBe(`管${content}`);
|
||||
});
|
||||
|
||||
it("smooths long Chinese text without waiting for paragraph-sized buffers", async () => {
|
||||
|
||||
Reference in New Issue
Block a user