From dfb144ebb5e3fd781e75750ac92ca9d9c2ebf51c Mon Sep 17 00:00:00 2001 From: dbarr5 Date: Sun, 27 Sep 2026 16:46:15 -0400 Subject: [PATCH] Reject undelimited SSE frames at end of stream --- src/core/stream.ts | 13 ++++++++++--- test/stream.test.ts | 33 ++++++++++++++++++++++++++++++++- test/turn_lifecycle.test.ts | 17 +++++++++++++++++ 3 files changed, 59 insertions(+), 4 deletions(-) diff --git a/src/core/stream.ts b/src/core/stream.ts index bdf8261b..ade870f9 100644 --- a/src/core/stream.ts +++ b/src/core/stream.ts @@ -9,7 +9,7 @@ // be ignored. The CF-flush preamble (`:<4096 spaces>`) and `ping` heartbeat are // handled here (comment lines skipped; ping surfaced as a typed liveness frame). -import { StreamEventTooLargeError } from "./errors.js"; +import { StreamEventTooLargeError, StreamIncompleteError } from "./errors.js"; /** One event may contain model text, but never an unbounded delimiter-free body. */ export const MAX_SSE_EVENT_BYTES = 1_048_576; @@ -349,8 +349,15 @@ export async function* decodeSse( throw new StreamEventTooLargeError(MAX_SSE_EVENT_BYTES); } } - const tail = parseEvent(buf); - if (tail) yield tail; + // SSE dispatches an event only after a blank-line delimiter. A connection + // that closes after `data: {"type":"done"}` has not delivered a terminal + // frame: accepting that tail would turn a truncated response into success. + // A trailing comment/whitespace is harmless, but any unfinished data line + // leaves the turn uncertain (and lets dev-session consumers reconnect). + buf = (buf + td.decode()).replace(/\r\n/g, "\n"); + if (buf.split("\n").some((line) => line.startsWith("data:"))) { + throw new StreamIncompleteError(); + } } function numOrUndef(v: unknown): number | undefined { diff --git a/test/stream.test.ts b/test/stream.test.ts index 32246fcc..0d2708d8 100644 --- a/test/stream.test.ts +++ b/test/stream.test.ts @@ -5,8 +5,9 @@ import { decodeSse, normalizeFrame, parseEvent, + type StreamFrame, } from "../src/core/stream.js"; -import { StreamEventTooLargeError } from "../src/core/errors.js"; +import { StreamEventTooLargeError, StreamIncompleteError } from "../src/core/errors.js"; test("normalizeFrame done maps token fields and carries no signature", () => { const f = normalizeFrame({ @@ -112,6 +113,36 @@ test("decodeSse handles a frame split across chunks", async () => { assert.deepEqual(frames, [{ type: "delta", text: "split" }]); }); +test("decodeSse refuses an undelimited terminal frame at EOF", async () => { + async function* bytes(): AsyncGenerator { + const enc = new TextEncoder(); + yield enc.encode('data: {"type":"delta","text":"partial"}\n\n'); + yield enc.encode('data: {"type":"done","uvt":1,"cents":0}'); + } + const frames: StreamFrame[] = []; + await assert.rejects(async () => { + for await (const frame of decodeSse(bytes())) frames.push(frame); + }, StreamIncompleteError); + assert.deepEqual(frames, [{ type: "delta", text: "partial" }]); +}); + +test("decodeSse refuses a truncated data frame but ignores a trailing comment", async () => { + const enc = new TextEncoder(); + async function* truncated(): AsyncGenerator { + yield enc.encode('data: {"type":"done"'); + } + await assert.rejects(async () => { + for await (const _frame of decodeSse(truncated())) { /* drain */ } + }, StreamIncompleteError); + + async function* comment(): AsyncGenerator { + yield enc.encode(': keepalive'); + } + const frames = []; + for await (const frame of decodeSse(comment())) frames.push(frame); + assert.deepEqual(frames, []); +}); + test("decodeSse rejects an oversized delimiter-free event and closes its source", async () => { let closed = false; async function* bytes(): AsyncGenerator { diff --git a/test/turn_lifecycle.test.ts b/test/turn_lifecycle.test.ts index f2b4cf5d..9b3849fc 100644 --- a/test/turn_lifecycle.test.ts +++ b/test/turn_lifecycle.test.ts @@ -304,6 +304,23 @@ test("EOF without a terminal frame preserves partial output and exits nonzero as } }); +test("an undelimited done at EOF cannot make a headless chat turn succeed", async () => { + resetRegistry(); + const realFetch = globalThis.fetch; + globalThis.fetch = (async () => new Response( + 'data: {"type":"delta","text":"partial answer"}\n\ndata: {"type":"done","uvt":1,"cents":0}', + { status: 200, headers: { "content-type": "text/event-stream" } }, + )) as typeof globalThis.fetch; + try { + const result = await captureWrites(() => cmdChat(cloudContext(), "ship this")); + assert.equal(result.value, 1); + assert.equal(result.stdout, "partial answer"); + assert.match(result.stderr, /connection ended before the server finished/i); + } finally { + globalThis.fetch = realFetch; + } +}); + test("a successful stream returns its stable succeeded outcome", async () => { resetRegistry(); const realFetch = globalThis.fetch;