diff --git a/CHANGELOG.md b/CHANGELOG.md index 786f7c0..eeed2c0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,12 @@ Release tags use the form `vX.Y.Z` and match `package.json`. GitHub Releases car ## [Unreleased] +## [0.5.9] - 2026-07-23 + +### Fixed + +- Long-running streamed responses now send protocol-safe keepalives so clients such as Claude Code do not cancel healthy requests during extended model silence. + ## [0.5.8] - 2026-07-23 ### Fixed @@ -94,7 +100,8 @@ Release tags use the form `vX.Y.Z` and match `package.json`. GitHub Releases car See [GitHub Releases](https://github.com/gitcommit90/rerouted/releases) for artifact digests and notes prior to the Keep a Changelog narrative. Notable themes in late 0.4.x included signed/notarized distribution, in-app updates, named routes, OAuth account pools, OpenAI chat completions and Responses routing, and launch hardening. -[Unreleased]: https://github.com/gitcommit90/rerouted/compare/v0.5.8...HEAD +[Unreleased]: https://github.com/gitcommit90/rerouted/compare/v0.5.9...HEAD +[0.5.9]: https://github.com/gitcommit90/rerouted/releases/tag/v0.5.9 [0.5.8]: https://github.com/gitcommit90/rerouted/releases/tag/v0.5.8 [0.5.7]: https://github.com/gitcommit90/rerouted/releases/tag/v0.5.7 [0.5.5]: https://github.com/gitcommit90/rerouted/releases/tag/v0.5.5 diff --git a/package-lock.json b/package-lock.json index 418ea14..a149eaf 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "@gitcommit90/rerouted", - "version": "0.5.8", + "version": "0.5.9", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@gitcommit90/rerouted", - "version": "0.5.8", + "version": "0.5.9", "license": "MIT", "bin": { "rerouted": "src/cli/index.js" diff --git a/package.json b/package.json index fe77720..3f0c426 100644 --- a/package.json +++ b/package.json @@ -1,7 +1,7 @@ { "name": "@gitcommit90/rerouted", "productName": "ReRouted", - "version": "0.5.8", + "version": "0.5.9", "description": "A local AI router for connected accounts, models, named routes, and automatic fallback.", "author": "gitcommit90", "license": "MIT", diff --git a/src/lib/anthropic-api.js b/src/lib/anthropic-api.js index ced7555..381d74f 100644 --- a/src/lib/anthropic-api.js +++ b/src/lib/anthropic-api.js @@ -378,6 +378,12 @@ function writeEvent(sink, type, data) { sink.write(`event: ${type}\ndata: ${JSON.stringify(data)}\n\n`); } +// Claude Code filters SSE comments and ping events before its stream-event +// watchdog. This zero-impact delta reaches that watchdog without adding content +// or completing the response; the real final delta remains authoritative. +const ANTHROPIC_SSE_HEARTBEAT = + 'event: message_delta\ndata: {"type":"message_delta","delta":{"stop_reason":null,"stop_sequence":null},"usage":{"output_tokens":0}}\n\n'; + async function pipeChatCompletionsSseToAnthropic(streamPipe, sink, requestedModel) { const parser = createSseParser(); const id = messageId(); @@ -571,6 +577,7 @@ module.exports = { toChatCompletionsBody, fromChatCompletion, pipeChatCompletionsSseToAnthropic, + ANTHROPIC_SSE_HEARTBEAT, toAnthropicError, estimateInputTokens, anthropicUsage, diff --git a/src/lib/gateway.js b/src/lib/gateway.js index cd98c32..5222fe7 100644 --- a/src/lib/gateway.js +++ b/src/lib/gateway.js @@ -3,6 +3,7 @@ const http = require("node:http"); const { DEFAULT_PORT } = require("./constants"); const logger = require("./logger"); +const { createSseHeartbeat, DEFAULT_SSE_HEARTBEAT_MS } = require("./sse"); const { toChatCompletionsBody, fromChatCompletion, @@ -15,10 +16,15 @@ const { pipeChatCompletionsSseToAnthropic, toAnthropicError, estimateInputTokens, + ANTHROPIC_SSE_HEARTBEAT, } = require("./anthropic-api"); const MAX_JSON_BODY_BYTES = 32 * 1024 * 1024; +function isClaudeCodeRequest(req) { + return /^claude-cli\//i.test(String(req.headers["user-agent"] || "")); +} + /** * OpenAI-compatible HTTP gateway. * Auth: Authorization: Bearer @@ -31,6 +37,7 @@ function createGateway({ port = DEFAULT_PORT, host = "127.0.0.1", maxBodyBytes = MAX_JSON_BODY_BYTES, + sseHeartbeatMs = DEFAULT_SSE_HEARTBEAT_MS, requestActivity, } = {}) { let server = null; @@ -329,13 +336,29 @@ function createGateway({ "Cache-Control": "no-cache", Connection: "keep-alive", }); + const heartbeat = anthropicRequest && isClaudeCodeRequest(req) + ? createSseHeartbeat(res, { + intervalMs: sseHeartbeatMs, + heartbeat: ANTHROPIC_SSE_HEARTBEAT, + }) + : null; + const streamSink = heartbeat?.sink || res; try { if (responsesRequest) { - await pipeChatCompletionsSseToResponses(result.streamPipe, res, body.model, body); + await pipeChatCompletionsSseToResponses( + result.streamPipe, + streamSink, + body.model, + body + ); } else if (anthropicRequest) { - await pipeChatCompletionsSseToAnthropic(result.streamPipe, res, body.model); + await pipeChatCompletionsSseToAnthropic( + result.streamPipe, + streamSink, + body.model + ); } else { - await result.streamPipe(res); + await result.streamPipe(streamSink); } activityStatus = 200; activityOutcome = "success"; @@ -349,6 +372,8 @@ function createGateway({ ); } } + } finally { + heartbeat?.stop(); } if (!res.writableEnded) res.end(); return; diff --git a/src/lib/sse.js b/src/lib/sse.js index b650338..c6fea3c 100644 --- a/src/lib/sse.js +++ b/src/lib/sse.js @@ -24,6 +24,57 @@ function formatSseData(obj) { } const SSE_DONE = "data: [DONE]\n\n"; +const SSE_KEEPALIVE = ": keepalive\n\n"; +const DEFAULT_SSE_HEARTBEAT_MS = 25_000; + +/** + * Wrap a client-facing SSE sink and emit the configured liveness frame whenever + * the stream has otherwise been silent for the configured interval. + */ +function createSseHeartbeat( + res, + { intervalMs = DEFAULT_SSE_HEARTBEAT_MS, heartbeat = SSE_KEEPALIVE } = {} +) { + let timer = null; + let stopped = false; + + function schedule() { + if (timer) clearTimeout(timer); + timer = null; + if (stopped || intervalMs <= 0) return; + timer = setTimeout(() => { + timer = null; + if (stopped || res.writableEnded || res.destroyed) return; + res.write(heartbeat); + schedule(); + }, intervalMs); + timer.unref?.(); + } + + const sink = { + write(...args) { + const written = res.write(...args); + schedule(); + return written; + }, + get writableEnded() { + return res.writableEnded; + }, + get destroyed() { + return res.destroyed; + }, + }; + + schedule(); + return { + sink, + stop() { + stopped = true; + if (timer) clearTimeout(timer); + timer = null; + }, + }; +} /** * Parse SSE text stream lines into { event, data } objects. @@ -126,6 +177,9 @@ module.exports = { openaiChunk, formatSseData, SSE_DONE, + SSE_KEEPALIVE, + DEFAULT_SSE_HEARTBEAT_MS, + createSseHeartbeat, createSseParser, chunkToString, pipeOpenAiSse, diff --git a/tests/gateway.test.js b/tests/gateway.test.js index be7e566..9ccec8e 100644 --- a/tests/gateway.test.js +++ b/tests/gateway.test.js @@ -1797,6 +1797,157 @@ describe("SSE chunk decoding", () => { assert.equal(chunkToString(u8), text); assert.notEqual(String(u8), text); }); + + it("keeps Anthropic Messages clients alive during downstream silence", async () => { + const store = createStore(tmpConfig()); + const apiKey = store.load().apiKey; + const gateway = createGateway({ + store, + sseHeartbeatMs: 10, + router: { + async chatCompletions() { + return { + ok: true, + stream: true, + streamPipe: async (sink) => { + await new Promise((resolve) => setTimeout(resolve, 35)); + sink.write( + `data: ${JSON.stringify({ + choices: [{ delta: { role: "assistant", content: "OK" } }], + })}\n\n` + ); + sink.write( + `data: ${JSON.stringify({ + choices: [{ delta: {}, finish_reason: "stop" }], + })}\n\n` + ); + sink.write("data: [DONE]\n\n"); + }, + }; + }, + }, + }); + const server = http.createServer((req, res) => gateway.handle(req, res)); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + + try { + const response = await fetch(`http://127.0.0.1:${server.address().port}/v1/messages`, { + method: "POST", + headers: { + Authorization: `Bearer ${apiKey}`, + "Content-Type": "application/json", + "User-Agent": "claude-cli/2.1.218 (external, sdk-cli)", + }, + body: JSON.stringify({ + model: "route", + max_tokens: 8, + messages: [], + stream: true, + }), + }); + const text = await response.text(); + assert.equal(response.status, 200); + assert.match( + text, + /event: message_delta\ndata: \{"type":"message_delta","delta":\{"stop_reason":null,"stop_sequence":null\},"usage":\{"output_tokens":0\}\}\n\n/ + ); + assert.match(text, /event: content_block_delta/); + const heartbeatBlock = text.split("\n\n").find( + (block) => block.startsWith("event: message_delta") && + block.includes('"stop_reason":null') + ); + const heartbeatData = JSON.parse( + heartbeatBlock.split("\n").find((line) => line.startsWith("data: ")).slice(6) + ); + assert.deepEqual(heartbeatData.usage, { output_tokens: 0 }); + + const genericResponse = await fetch( + `http://127.0.0.1:${server.address().port}/v1/messages`, + { + method: "POST", + headers: { Authorization: `Bearer ${apiKey}`, "Content-Type": "application/json" }, + body: JSON.stringify({ + model: "route", + max_tokens: 8, + messages: [], + stream: true, + }), + } + ); + const genericText = await genericResponse.text(); + assert.equal(genericResponse.status, 200); + assert.equal( + genericText.split("\n\n").filter( + (block) => block.startsWith("event: message_delta") && + block.includes('"stop_reason":null') + ).length, + 0 + ); + } finally { + await new Promise((resolve) => server.close(resolve)); + } + }); + + it("still cancels an active stream when the client disconnects", async () => { + const store = createStore(tmpConfig()); + const apiKey = store.load().apiKey; + let finishActivity; + const activityEnded = new Promise((resolve) => { + finishActivity = resolve; + }); + const gateway = createGateway({ + store, + sseHeartbeatMs: 10, + requestActivity: { + begin: () => "activity-1", + route() {}, + end: (_id, result) => finishActivity(result), + }, + router: { + async chatCompletions({ signal }) { + return { + ok: true, + stream: true, + streamPipe: async () => new Promise((resolve, reject) => { + if (signal.aborted) reject(signal.reason); + else signal.addEventListener("abort", () => reject(signal.reason), { once: true }); + }), + }; + }, + }, + }); + const server = http.createServer((req, res) => gateway.handle(req, res)); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + + try { + await new Promise((resolve, reject) => { + const request = http.request({ + host: "127.0.0.1", + port: server.address().port, + path: "/v1/messages", + method: "POST", + headers: { + Authorization: `Bearer ${apiKey}`, + "Content-Type": "application/json", + "User-Agent": "claude-cli/2.1.218 (external, sdk-cli)", + }, + }, (response) => { + let received = ""; + response.on("data", (chunk) => { + received += String(chunk); + if (!received.includes('"stop_reason":null')) return; + response.destroy(); + resolve(); + }); + }); + request.once("error", reject); + request.end(JSON.stringify({ model: "route", max_tokens: 8, messages: [], stream: true })); + }); + assert.deepEqual(await activityEnded, { status: 499, outcome: "canceled" }); + } finally { + await new Promise((resolve) => server.close(resolve)); + } + }); }); describe("OAuth → OpenAI SSE translation pipes", () => {