fix: settle only after visible stream output

This commit is contained in:
Jesse_Chen
2026-07-21 13:01:37 +08:00
parent 53dddb8b43
commit 697aac2325
3 changed files with 71 additions and 7 deletions
+12 -6
View File
@@ -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<Uint8Array>,
value: string,
) {
await settleFirstOutput(value);
controller.enqueue(encoder.encode(value));
if (/\S/.test(value)) emitted = true;
}
const body = new ReadableStream<Uint8Array>({
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;
}
}
@@ -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,
@@ -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<void>((resolve) => {
markSecondReadStarted = resolve;
});
const never = new Promise<IteratorResult<string>>(() => {});
let reads = 0;
const reply: AsyncIterable<string> = {
[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;