File size: 2,936 Bytes
67d18ac | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 | /**
* SSE keepalive stream generator.
*
* Emits pure SSE comment frames (`: keepalive\n\n`) at a fixed interval.
* Used by the async bridge to keep the client connection alive during
* ticket-queue wait — comment frames are universally ignored by spec-compliant
* SSE clients (OpenAI / Anthropic SDKs) but reset their idle-timeout timers.
*
* Never emits `data:` frames — those carry semantic events and would confuse
* strict SDK parsers (e.g. `@ai-sdk/anthropic` zod schemas throw on unknown
* `type` values, killing the stream — see plan §3.6 anti-pattern #1).
*/
export interface KeepaliveOptions {
/** Interval between comment frames in ms. */
intervalMs: number;
/** Comment text. Default `"keepalive"`. Must not contain newlines. */
text?: string;
/** External abort — when fired, the stream closes gracefully. */
signal?: AbortSignal;
}
/**
* Build a `ReadableStream<Uint8Array>` that emits `: ${text}\n\n` every
* `intervalMs`. The stream closes when `signal` aborts. Text encoder is
* reused across frames; the underlying buffer is fresh per emit.
*
* Cadence guarantee: first emit happens after `intervalMs` (NOT immediately)
* so callers can compose with upstream-output streams without a leading comment.
*/
export function keepaliveStream(opts: KeepaliveOptions): ReadableStream<Uint8Array> {
const text = (opts.text ?? "keepalive").replace(/[\r\n]/g, " ");
const frame = new TextEncoder().encode(`: ${text}\n\n`);
let timer: ReturnType<typeof setTimeout> | undefined;
let aborted = false;
return new ReadableStream<Uint8Array>({
start(controller) {
if (opts.signal) {
if (opts.signal.aborted) {
aborted = true;
controller.close();
return;
}
opts.signal.addEventListener("abort", () => {
aborted = true;
if (timer) {
clearTimeout(timer);
timer = undefined;
}
try { controller.close(); } catch { /* already closed */ }
}, { once: true });
}
const tick = (): void => {
if (aborted) return;
try {
controller.enqueue(frame);
} catch {
// Controller closed by consumer; stop the timer.
if (timer) {
clearTimeout(timer);
timer = undefined;
}
return;
}
timer = setTimeout(tick, opts.intervalMs);
};
timer = setTimeout(tick, opts.intervalMs);
},
cancel() {
aborted = true;
if (timer) {
clearTimeout(timer);
timer = undefined;
}
},
});
}
/**
* Emit a single immediate keepalive frame as a `Uint8Array`.
* Useful for flushing one frame before pausing on a long upstream call.
*/
export function keepaliveFrame(text: string = "keepalive"): Uint8Array {
const clean = text.replace(/[\r\n]/g, " ");
return new TextEncoder().encode(`: ${clean}\n\n`);
}
|