mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-08-04 06:12:10 +03:00
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)
This commit is contained in:
@@ -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<Uint8Array, Uint8Array>();
|
||||
}
|
||||
|
||||
// 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<Uint8Array, Uint8Array>();
|
||||
}
|
||||
|
||||
let intervalId: ReturnType<typeof setInterval> | undefined;
|
||||
|
||||
const stop = () => {
|
||||
|
||||
75
tests/unit/sseHeartbeat.test.ts
Normal file
75
tests/unit/sseHeartbeat.test.ts
Normal file
@@ -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<Uint8Array>({
|
||||
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<Uint8Array, string>({
|
||||
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;
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user