mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-09-17 12:22:34 +03:00
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
This commit is contained in:
committed by
GitHub
parent
668beed5b8
commit
90366903c4
@@ -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<string, unknown> | 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<string, unknown>)) {
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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 };
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user