From 5db2e20c3bddb9affd25706d008444675c618ddd Mon Sep 17 00:00:00 2001 From: Adam <835217250@qq.com> Date: Sat, 18 Jul 2026 22:34:14 +0800 Subject: [PATCH] feat(sse): allow disabling `:` comment heartbeats via OMNIROUTE_SSE_COMMENTS=off (#7036) * feat(sse): allow disabling `:` comment heartbeats via OMNIROUTE_SSE_COMMENTS=off * fix(sse): guard process access for edge/Workers + export helper * test(sse): cover sseCommentsEnabled + heartbeat suppression (#7036) --- open-sse/utils/sseHeartbeat.ts | 21 +++++++++ tests/unit/sseHeartbeat.test.ts | 75 +++++++++++++++++++++++++++++++++ 2 files changed, 96 insertions(+) create mode 100644 tests/unit/sseHeartbeat.test.ts diff --git a/open-sse/utils/sseHeartbeat.ts b/open-sse/utils/sseHeartbeat.ts index 9a12214145..005a19b172 100644 --- a/open-sse/utils/sseHeartbeat.ts +++ b/open-sse/utils/sseHeartbeat.ts @@ -61,6 +61,20 @@ type SseHeartbeatTransformOptions = { const HEARTBEAT_ENCODER = new TextEncoder(); +/** + * Whether OmniRoute may emit SSE `:` comment lines (e.g. the `: keepalive` heartbeat). + * Some strict OpenAI-compatible clients parse every SSE line as JSON and crash on `:` comments. + * Set OMNIROUTE_SSE_COMMENTS=off to suppress comment-shaped heartbeats (they become a no-op). + * Defaults to enabled for backward compatibility. + */ +export function sseCommentsEnabled(): boolean { + // SSR/edge safety: `process` is not defined in Workers/Deno/edge runtimes. + if (typeof process === "undefined") return true; + const v = process.env.OMNIROUTE_SSE_COMMENTS; + if (v === undefined || v === "") return true; + return v.trim().toLowerCase() !== "off"; +} + export function createSseHeartbeatTransform({ intervalMs = DEFAULT_SSE_HEARTBEAT_INTERVAL_MS, signal, @@ -72,6 +86,13 @@ export function createSseHeartbeatTransform({ return new TransformStream(); } + // Opt-out for strict OpenAI-compatible clients that JSON.parse every SSE line and + // crash on `:` comment heartbeats. OMNIROUTE_SSE_COMMENTS=off disables comment-shaped + // heartbeats (they become a no-op); valid `data:` heartbeats are unaffected. + if (!sseCommentsEnabled() && shape === HEARTBEAT_SHAPES.COMMENT) { + return new TransformStream(); + } + let intervalId: ReturnType | undefined; const stop = () => { diff --git a/tests/unit/sseHeartbeat.test.ts b/tests/unit/sseHeartbeat.test.ts new file mode 100644 index 0000000000..36d51b1e79 --- /dev/null +++ b/tests/unit/sseHeartbeat.test.ts @@ -0,0 +1,75 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { + sseCommentsEnabled, + shapeForClientFormat, + createSseHeartbeatTransform, + HEARTBEAT_SHAPES, +} from "../../open-sse/utils/sseHeartbeat.ts"; + +function withEnv(value: string | undefined, fn: () => void) { + const prev = process.env.OMNIROUTE_SSE_COMMENTS; + try { + if (value === undefined) delete process.env.OMNIROUTE_SSE_COMMENTS; + else process.env.OMNIROUTE_SSE_COMMENTS = value; + fn(); + } finally { + if (prev === undefined) delete process.env.OMNIROUTE_SSE_COMMENTS; + else process.env.OMNIROUTE_SSE_COMMENTS = prev; + } +} + +test("sseCommentsEnabled defaults to true when the env var is unset", () => { + withEnv(undefined, () => assert.equal(sseCommentsEnabled(), true)); +}); + +test("sseCommentsEnabled is false only when set to 'off' (case-insensitive)", () => { + withEnv("off", () => assert.equal(sseCommentsEnabled(), false)); + withEnv("OFF", () => assert.equal(sseCommentsEnabled(), false)); + withEnv("on", () => assert.equal(sseCommentsEnabled(), true)); + withEnv("false", () => assert.equal(sseCommentsEnabled(), true)); +}); + +test("shapeForClientFormat maps known client formats", () => { + assert.equal(shapeForClientFormat("claude"), HEARTBEAT_SHAPES.ANTHROPIC_PING); + assert.equal(shapeForClientFormat("openai"), HEARTBEAT_SHAPES.OPENAI_CHUNK); + assert.equal(shapeForClientFormat("openai-responses"), HEARTBEAT_SHAPES.OPENAI_RESPONSES_IN_PROGRESS); + assert.equal(shapeForClientFormat(undefined), HEARTBEAT_SHAPES.COMMENT); +}); + +test("createSseHeartbeatTransform suppresses COMMENT heartbeats when OMNIROUTE_SSE_COMMENTS=off", async () => { + const prev = process.env.OMNIROUTE_SSE_COMMENTS; + process.env.OMNIROUTE_SSE_COMMENTS = "off"; + try { + const enc = new TextEncoder(); + const dec = new TextDecoder(); + const input = new ReadableStream({ + start(c) { + c.enqueue(enc.encode("data: hello\n\n")); + c.close(); + }, + }); + const reader = input + .pipeThrough(createSseHeartbeatTransform({ shape: HEARTBEAT_SHAPES.COMMENT, intervalMs: 20 })) + .pipeThrough( + new TransformStream({ + transform(chunk, ctrl) { + ctrl.enqueue(dec.decode(chunk)); + }, + }) + ) + .getReader(); + const chunks: string[] = []; + let res = await reader.read(); + while (!res.done) { + chunks.push(res.value as string); + res = await reader.read(); + } + const out = chunks.join(""); + assert.ok(!out.includes(": keepalive"), "no comment heartbeat should be emitted when disabled"); + assert.ok(out.includes("data: hello"), "original chunk passes through unchanged"); + } finally { + if (prev === undefined) delete process.env.OMNIROUTE_SSE_COMMENTS; + else process.env.OMNIROUTE_SSE_COMMENTS = prev; + } +});