From 697aac23257fd114fcafa8d2eb2c26cc37edc14f Mon Sep 17 00:00:00 2001 From: Jesse_Chen Date: Tue, 21 Jul 2026 13:01:37 +0800 Subject: [PATCH] fix: settle only after visible stream output --- frontend/src/lib/stream-text-response.ts | 18 ++++-- .../rectification-question-handoff.test.ts | 2 +- frontend/tests/stream-text-response.test.ts | 58 +++++++++++++++++++ 3 files changed, 71 insertions(+), 7 deletions(-) diff --git a/frontend/src/lib/stream-text-response.ts b/frontend/src/lib/stream-text-response.ts index d2d622d9..0566e4cf 100644 --- a/frontend/src/lib/stream-text-response.ts +++ b/frontend/src/lib/stream-text-response.ts @@ -161,11 +161,20 @@ export function streamTextResponse( let firstOutputSettled = false; async function settleFirstOutput(value: string) { - if (!value || firstOutputSettled) return; + if (!/\S/.test(value) || firstOutputSettled) return; firstOutputSettled = true; await options.onFirstOutput?.(); } + async function enqueueOutput( + controller: ReadableStreamDefaultController, + value: string, + ) { + await settleFirstOutput(value); + controller.enqueue(encoder.encode(value)); + if (/\S/.test(value)) emitted = true; + } + const body = new ReadableStream({ async pull(controller) { try { @@ -176,8 +185,7 @@ export function streamTextResponse( ? visibleTransformer.finish(pending) : pending; if (finalText) { - await settleFirstOutput(finalText); - controller.enqueue(encoder.encode(finalText)); + await enqueueOutput(controller, finalText); } settled = true; if (!emitted) { @@ -190,7 +198,6 @@ export function streamTextResponse( controller.close(); return; } - if (/\S/.test(value)) emitted = true; pending += value; if (pending.length <= guardTailLength) continue; @@ -201,8 +208,7 @@ export function streamTextResponse( ? visibleTransformer.push(stable) : stable; if (transformed) { - await settleFirstOutput(transformed); - controller.enqueue(encoder.encode(transformed)); + await enqueueOutput(controller, transformed); return; } } diff --git a/frontend/tests/rectification-question-handoff.test.ts b/frontend/tests/rectification-question-handoff.test.ts index 59ff52c8..54b05b7f 100644 --- a/frontend/tests/rectification-question-handoff.test.ts +++ b/frontend/tests/rectification-question-handoff.test.ts @@ -413,7 +413,7 @@ test("confirmed surface never revives an old local question after durable consum const consumed = { ...confirmedTurn(), pendingConsultationQuestion: null, - actions: [] as const, + actions: [], }; const markup = renderToStaticMarkup(React.createElement( ConversationalRectificationSurface, diff --git a/frontend/tests/stream-text-response.test.ts b/frontend/tests/stream-text-response.test.ts index b25e5dac..c7723b32 100644 --- a/frontend/tests/stream-text-response.test.ts +++ b/frontend/tests/stream-text-response.test.ts @@ -165,6 +165,64 @@ test("transformed empty streams still refund through the error settlement", asyn assert.equal(errors, 1); }); +test("a transformed short reply that fails while buffered reports no emitted output", async () => { + let observedEmitted: boolean | null = null; + async function* reply() { + yield "尚未冲出的短回答"; + throw new Error("upstream_failed"); + } + const response = streamTextResponse(reply(), { + mode: "mastra", + requestId: "00000000-0000-4000-8000-000000000007", + transformText: (text) => text, + onError: async (_error, emitted) => { observedEmitted = emitted; }, + }); + + await assert.rejects(response.text(), /upstream_failed/); + assert.equal(observedEmitted, false); +}); + +test("cancelling while transformed short output is buffered reports no emitted output", async () => { + let markSecondReadStarted = () => {}; + const secondReadStarted = new Promise((resolve) => { + markSecondReadStarted = resolve; + }); + const never = new Promise>(() => {}); + let reads = 0; + const reply: AsyncIterable = { + [Symbol.asyncIterator]() { + return { + next() { + reads += 1; + if (reads === 1) { + return Promise.resolve({ done: false, value: "尚未冲出的短回答" }); + } + markSecondReadStarted(); + return never; + }, + return() { + return Promise.resolve({ done: true, value: undefined }); + }, + }; + }, + }; + let observedEmitted: boolean | null = null; + const response = streamTextResponse(reply, { + mode: "mastra", + requestId: "00000000-0000-4000-8000-000000000008", + transformText: (text) => text, + onCancel: async (emitted) => { observedEmitted = emitted; }, + }); + const reader = response.body?.getReader(); + assert.ok(reader); + const pendingRead = reader.read(); + await secondReadStarted; + + await reader.cancel(); + await pendingRead; + assert.equal(observedEmitted, false); +}); + test("cancelling a transformed stream after visible output preserves emitted settlement", async () => { let charged = 0; let refunded = 0;