From 90366903c4a673dadfb35cf59fbf3d7784dc58b1 Mon Sep 17 00:00:00 2001 From: opensource-elearning <159253500+opensource-elearning@users.noreply.github.com> Date: Mon, 31 Aug 2026 22:40:39 +0530 Subject: [PATCH] fix: prevent Claude Code session kills via liveness-aware readiness + auto model echo (#12189) - streamReadiness: reset deadline on each received chunk (keepalive = alive) with a hard maxTimeoutMs ceiling so truly-dead connections still fail fast. Preserves operator's 20s/100s intent for dead pulls while allowing slow-but-alive upstreams (reasoning warm-ups) to survive. - chatCore + codexIdentity: auto-detect Claude Code CLI via user-agent/originator headers and enable model echo for it. The response field now echoes the originally-requested alias/combo (e.g. ) instead of the resolved upstream id (e.g. ), so restores cleanly without 'could not be restored' errors. Refs: opensource-elearning/omniroute-fixes#1, diegosouzapw/OmniRoute#12185 --- open-sse/config/codexIdentity.ts | 34 +++++++++++++++++++++++++ open-sse/handlers/chatCore.ts | 11 ++++++-- open-sse/utils/streamReadiness.ts | 30 ++++++++++++++++++++-- open-sse/utils/streamReadinessPolicy.ts | 5 ++-- 4 files changed, 74 insertions(+), 6 deletions(-) diff --git a/open-sse/config/codexIdentity.ts b/open-sse/config/codexIdentity.ts index a081c5402b..8222490782 100644 --- a/open-sse/config/codexIdentity.ts +++ b/open-sse/config/codexIdentity.ts @@ -570,3 +570,37 @@ export function isVerifiedNativeCodexRequest( ): boolean { return isCodexOriginatedHeaders(headers) && hasNativeCodexTurnBinding(body); } + +/** + * Detect the Claude Code CLI as the request *client* from request headers. + * Used to auto-enable model echo so session restores work when the resolved + * upstream model (e.g. `oc/nemotron-3-ultra-free`) is not recognized by the + * Claude Code client on `--resume`. + */ +export function isClaudeCodeOriginatedHeaders( + headers: Headers | Record | null | undefined +): boolean { + const getHeader = (name: string): string => { + if (headers instanceof Headers) { + return headers.get(name)?.toLowerCase() ?? ""; + } + if (headers && typeof headers === "object") { + for (const [key, value] of Object.entries(headers as Record)) { + if (key.toLowerCase() === name && typeof value === "string") { + return value.toLowerCase(); + } + } + } + return ""; + }; + + // Claude Code identifies itself via the user-agent header + const userAgent = getHeader("user-agent"); + if (userAgent.includes("claude-code") || userAgent.includes("anthropic-ai/claude-code")) { + return true; + } + // Also check originator if present + const originator = getHeader("originator"); + if (originator.startsWith("claude-code")) return true; + return false; +} diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index 658b4084da..91f30bcbee 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -77,7 +77,7 @@ import { isStripReasoningRequested, } from "./chatCore/headers.ts"; import { markCodexScopeRateLimited } from "./chatCore/codexFailover.ts"; -import { getCodexClientSessionId, isCodexOriginatedHeaders } from "../config/codexIdentity.ts"; +import { getCodexClientSessionId, isCodexOriginatedHeaders, isClaudeCodeOriginatedHeaders } from "../config/codexIdentity.ts"; import { noteCodexTurnStateProvenance, readCodexTurnStateHeader, @@ -981,8 +981,14 @@ export async function handleChatCore({ const isCodexResponsesEcho = (isResponsesEndpoint || sourceFormat === FORMATS.OPENAI_RESPONSES) && isCodexOriginatedHeaders(clientRawRequest?.headers); + + // Detect Claude Code CLI so we can auto-enable model echo — this prevents + // session restore failures when the resolved upstream model (e.g. + // `oc/nemotron-3-ultra-free`) is not recognized by the client on `--resume`. + const isClaudeCodeClient = isClaudeCodeOriginatedHeaders(clientRawRequest?.headers); + let echoModel = - (settings.echoRequestedModelName === true || isCodexResponsesEcho) && + (settings.echoRequestedModelName === true || isCodexResponsesEcho || isClaudeCodeClient) && typeof requestedModel === "string" && requestedModel ? requestedModel @@ -5465,6 +5471,7 @@ export async function handleChatCore({ const streamReadiness = await ensureStreamReadiness(providerResponse, { timeoutMs: streamReadinessPolicy.timeoutMs, + maxTimeoutMs: streamReadinessPolicy.maxTimeoutMs, provider, model, log, diff --git a/open-sse/utils/streamReadiness.ts b/open-sse/utils/streamReadiness.ts index 1696a7a5c5..4bbeef1e0e 100644 --- a/open-sse/utils/streamReadiness.ts +++ b/open-sse/utils/streamReadiness.ts @@ -471,6 +471,9 @@ export async function ensureStreamReadiness( response: Response, options: { timeoutMs: number; + /** Hard ceiling for liveness-extended deadlines. When omitted, no hard ceiling + * is applied beyond `timeoutMs`. */ + maxTimeoutMs?: number; provider?: string | null; model?: string | null; log?: StreamReadinessLogger | null; @@ -489,7 +492,14 @@ export async function ensureStreamReadiness( }; const startedAt = Date.now(); const effectiveTimeoutMs = Math.max(0, Math.floor(options.timeoutMs)); - const deadline = startedAt + effectiveTimeoutMs; + // Hard ceiling: the deadline may extend on liveness signals (bytes arriving), + // but never past this absolute maximum. When maxTimeoutMs is omitted the + // initial timeoutMs itself acts as the ceiling (no extension). + const maxDeadline = + options.maxTimeoutMs != null + ? startedAt + Math.max(effectiveTimeoutMs, Math.floor(options.maxTimeoutMs)) + : startedAt + effectiveTimeoutMs; + let deadline = startedAt + effectiveTimeoutMs; let handedOffReader = false; const buildReadyResponse = () => @@ -500,7 +510,7 @@ export async function ensureStreamReadiness( }); const timeoutReason = () => - `Stream produced no non-ping SSE event within ${effectiveTimeoutMs}ms`; + `Stream produced no non-ping SSE event within ${deadline - startedAt}ms (max=${maxDeadline - startedAt}ms)`; try { while (true) { @@ -593,6 +603,22 @@ export async function ensureStreamReadiness( chunks.push(readResult.value); const decodedChunk = decoder.decode(readResult.value, { stream: true }); + // Liveness extension: bytes arrived → connection is alive, not dead. + // Reset the deadline so slow-but-alive upstreams (reasoning warm-ups, + // keepalive-only phases) are not aborted. The hard ceiling (maxDeadline) + // prevents unbounded waits and preserves the operator's fast-fail intent + // for truly dead connections. + const now = Date.now(); + if (deadline < maxDeadline) { + deadline = Math.min(now + effectiveTimeoutMs, maxDeadline); + if (now - startedAt > effectiveTimeoutMs) { + options.log?.debug?.( + "STREAM", + `readiness deadline extended to ${deadline - startedAt}ms (liveness signal) (${options.provider || "provider"}/${options.model || "unknown"})` + ); + } + } + if (appendStreamReadinessSignal(readinessState, decodedChunk)) { options.log?.debug?.( "STREAM", diff --git a/open-sse/utils/streamReadinessPolicy.ts b/open-sse/utils/streamReadinessPolicy.ts index 0507317bc3..0704c8ded9 100644 --- a/open-sse/utils/streamReadinessPolicy.ts +++ b/open-sse/utils/streamReadinessPolicy.ts @@ -14,6 +14,7 @@ export type StreamReadinessPolicyInput = { export type StreamReadinessPolicyResult = { timeoutMs: number; baseTimeoutMs: number; + maxTimeoutMs: number; reasons: string[]; }; @@ -121,7 +122,7 @@ export function resolveStreamReadinessTimeout( ): StreamReadinessPolicyResult { const baseTimeoutMs = Math.max(0, Math.floor(input.baseTimeoutMs || 0)); if (baseTimeoutMs <= 0) { - return { timeoutMs: baseTimeoutMs, baseTimeoutMs, reasons: ["disabled"] }; + return { timeoutMs: baseTimeoutMs, baseTimeoutMs, maxTimeoutMs: baseTimeoutMs, reasons: ["disabled"] }; } const maxTimeoutMs = Math.max(baseTimeoutMs, input.maxTimeoutMs ?? DEFAULT_MAX_TIMEOUT_MS); @@ -197,5 +198,5 @@ export function resolveStreamReadinessTimeout( timeoutMs = Math.min(timeoutMs, maxTimeoutMs); if (timeoutMs === baseTimeoutMs) reasons.push("base"); - return { timeoutMs, baseTimeoutMs, reasons }; + return { timeoutMs, baseTimeoutMs, maxTimeoutMs, reasons }; }