Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 10 additions & 3 deletions src/core/stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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 {
Expand Down
33 changes: 32 additions & 1 deletion test/stream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand Down Expand Up @@ -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<Uint8Array> {
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<Uint8Array> {
yield enc.encode('data: {"type":"done"');
}
await assert.rejects(async () => {
for await (const _frame of decodeSse(truncated())) { /* drain */ }
}, StreamIncompleteError);

async function* comment(): AsyncGenerator<Uint8Array> {
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<Uint8Array> {
Expand Down
17 changes: 17 additions & 0 deletions test/turn_lifecycle.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading