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
53 changes: 52 additions & 1 deletion open-sse/utils/stream.js
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,16 @@ export { COLORS, formatSSE };
// sharedEncoder is stateless — safe to share across streams
const sharedEncoder = new TextEncoder();

// Idle keepalive. Claude with extended thinking can stream ZERO bytes for 60s+
// while reasoning before the first token. A proxy in front of 9router
// (Cloudflare Tunnel, nginx, etc.) sees that silence as an idle connection and
// kills it at ~55-60s, so the request aborts before any token arrives (ttft is
// null, 0 tokens). Emitting an SSE comment line on a timer keeps the connection
// active. Comment lines (": ...") are ignored by every spec-compliant SSE
// client (Anthropic/OpenAI SDKs included), so this is invisible to callers.
const DEFAULT_HEARTBEAT_INTERVAL_MS = 10_000;
const HEARTBEAT_COMMENT = ": 9router-keepalive\n\n";

/**
* Stream modes
*/
Expand Down Expand Up @@ -44,7 +54,8 @@ export function createSSEStream(options = {}) {
connectionId = null,
body = null,
onStreamComplete = null,
apiKey = null
apiKey = null,
heartbeatIntervalMs = DEFAULT_HEARTBEAT_INTERVAL_MS
} = options;

let buffer = "";
Expand All @@ -67,10 +78,40 @@ export function createSSEStream(options = {}) {
// upstream EOF; cancel() runs on client disconnect, stall-abort, or upstream
// error. Without finalizing on cancel(), the streaming request detail row
// stays stuck at "[Streaming in progress...]" with 0 tokens forever.
// Idle-keepalive timer. Reset on every upstream chunk; fires only when the
// upstream has been silent for heartbeatIntervalMs. Cleared once the stream
// terminates (flush or cancel) so it never outlives the connection.
let heartbeatTimer = null;
let heartbeatController = null;
const stopHeartbeat = () => {
if (heartbeatTimer) {
clearInterval(heartbeatTimer);
heartbeatTimer = null;
}
};
const startHeartbeat = (controller) => {
heartbeatController = controller;
if (heartbeatIntervalMs <= 0 || heartbeatTimer) return;
heartbeatTimer = setInterval(() => {
if (finalized) {
stopHeartbeat();
return;
}
try {
heartbeatController.enqueue(sharedEncoder.encode(HEARTBEAT_COMMENT));
reqLogger?.appendConvertedChunk?.(HEARTBEAT_COMMENT);
} catch {
// Controller already closed/errored — nothing left to keep alive.
stopHeartbeat();
}
}, heartbeatIntervalMs);
};

let finalized = false;
const finalizeOnce = (interruptedReason = null) => {
if (finalized) return;
finalized = true;
stopHeartbeat();
if (!onStreamComplete) return;
const finalUsage = (mode === STREAM_MODE.TRANSLATE ? state?.usage : usage) || null;
onStreamComplete(
Expand All @@ -82,8 +123,18 @@ export function createSSEStream(options = {}) {
};

return new TransformStream({
start(controller) {
// Begin the idle keepalive immediately — the silent gap that kills the
// connection happens BEFORE the first upstream chunk arrives.
startHeartbeat(controller);
},

transform(chunk, controller) {
if (!ttftAt) ttftAt = Date.now();
// First real token arrived — the connection is no longer idle, so the
// keepalive has done its job. Stop it so we never interleave comment
// lines with genuine SSE data.
stopHeartbeat();
const text = decoder.decode(chunk, { stream: true });
buffer += text;
reqLogger?.appendProviderChunk?.(text);
Expand Down
129 changes: 129 additions & 0 deletions tests/unit/streaming-heartbeat-keepalive.test.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
import { beforeEach, describe, expect, it, vi } from "vitest";

// Regression guard for "streaming request aborts at ~55-60s with 0 tokens / null
// ttft when 9router is behind Cloudflare Tunnel" (follow-up to the interrupted-
// finalize fix).
//
// Root cause: Claude with extended thinking can stream ZERO bytes for 60s+ while
// reasoning before the first token. A proxy in front of 9router treats that
// silence as an idle connection and kills it before any token arrives. The fix
// emits an SSE comment line (": ...") on an idle timer so the connection stays
// active until the real first token shows up.
//
// Contract locked in here:
// 1. While upstream is silent, the transform emits heartbeat comment lines.
// 2. Comment lines are SSE comments (start with ":") so clients ignore them.
// 3. Once the first real chunk arrives, heartbeats stop (no interleaving).
// 4. The heartbeat timer is cleared on termination (no leak past finalize).

vi.mock("@/lib/usageDb.js", () => ({
trackPendingRequest: vi.fn(),
appendRequestLog: vi.fn(() => Promise.resolve()),
}));

const { createSSEStream } = await import("../../open-sse/utils/stream.js");

const HEARTBEAT_INTERVAL_MS = 20;

// Drive a source whose first chunk is delayed, so the transform sits idle long
// enough to emit heartbeats before any real data flows.
function pipeWithDelay({ firstChunkDelayMs, chunks, heartbeatIntervalMs, onStreamComplete }) {
const encoder = new TextEncoder();
let index = 0;
const source = new ReadableStream({
async pull(controller) {
if (index === 0 && firstChunkDelayMs > 0) {
await new Promise((resolve) => setTimeout(resolve, firstChunkDelayMs));
}
if (index >= chunks.length) {
controller.close();
return;
}
controller.enqueue(encoder.encode(chunks[index++]));
},
});

const transform = createSSEStream({
mode: "passthrough",
provider: "claude",
model: "claude-opus-4-8",
connectionId: "conn-1234",
body: { messages: [{ role: "user", content: "hi" }] },
onStreamComplete,
apiKey: "sk-test",
heartbeatIntervalMs,
});

const readable = source.pipeThrough(transform);
const reader = readable.getReader();
const decoder = new TextDecoder();

return (async () => {
const received = [];
while (true) {
const { value, done } = await reader.read();
if (done) break;
received.push(decoder.decode(value));
}
return received.join("");
})();
}

describe("streaming idle keepalive", () => {
beforeEach(() => vi.clearAllMocks());

it("emits SSE comment heartbeats while upstream is silent before the first token", async () => {
const onStreamComplete = vi.fn();
const output = await pipeWithDelay({
firstChunkDelayMs: HEARTBEAT_INTERVAL_MS * 3, // stay idle across several heartbeat ticks
chunks: [
'data: {"choices":[{"delta":{"content":"hello"}}]}\n\n',
"data: [DONE]\n\n",
],
heartbeatIntervalMs: HEARTBEAT_INTERVAL_MS,
onStreamComplete,
});

const heartbeatCount = (output.match(/: 9router-keepalive/g) || []).length;
expect(heartbeatCount).toBeGreaterThanOrEqual(1);
// Every heartbeat is an SSE comment line — clients ignore lines starting with ":".
for (const line of output.split("\n\n")) {
if (line.includes("9router-keepalive")) {
expect(line.trimStart().startsWith(":")).toBe(true);
}
}
// Real content still made it through after the idle gap.
expect(output).toContain('"content":"hello"');
expect(onStreamComplete).toHaveBeenCalledTimes(1);
});

it("does NOT emit heartbeats once tokens are flowing", async () => {
const onStreamComplete = vi.fn();
const output = await pipeWithDelay({
firstChunkDelayMs: 0, // first token immediately, never idle
chunks: [
'data: {"choices":[{"delta":{"content":"a"}}]}\n\n',
'data: {"choices":[{"delta":{"content":"b"}}]}\n\n',
"data: [DONE]\n\n",
],
heartbeatIntervalMs: HEARTBEAT_INTERVAL_MS,
onStreamComplete,
});

expect(output).not.toContain("9router-keepalive");
expect(onStreamComplete).toHaveBeenCalledTimes(1);
});

it("heartbeat disabled when interval <= 0", async () => {
const onStreamComplete = vi.fn();
const output = await pipeWithDelay({
firstChunkDelayMs: HEARTBEAT_INTERVAL_MS * 3,
chunks: ['data: {"choices":[{"delta":{"content":"x"}}]}\n\n', "data: [DONE]\n\n"],
heartbeatIntervalMs: 0,
onStreamComplete,
});

expect(output).not.toContain("9router-keepalive");
expect(onStreamComplete).toHaveBeenCalledTimes(1);
});
});