mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-09-14 19:02:17 +03:00
* fix(sse): grace period before finalizing a client disconnect as 499 (#9653) A client that closes its connection right after reading a fully-completed SSE stream can race OmniRoute's own completion bookkeeping: the bytes already reached the client, but the transform stream's own completion callback (onStreamComplete, which flips streamCompletionRecorded) hasn't finished bubbling up when the disconnect handler fires, so the request gets persisted as a false 499 with zero token usage even though it delivered its full response. Confirmed live on real traffic before this fix: a request whose server log showed "disconnect: request_signal_aborted" at 18236ms was persisted with status 200 and full token usage (82814/1292) once the grace period let the real completion win the race, matching what the client actually received. createClientDisconnectGraceHandler (new leaf in streamFailureFinalization.ts) polls isStreamCompletionRecorded() for up to STREAM_DISCONNECT_GRACE_PERIOD_MS (default 10s, env-configurable, 0 disables) before finalizing as a failure. If a real completion lands within the window, handleStreamFailure's own guard is a no-op and the genuine 200 stands. Covered by tests/unit/stream-disconnect-grace-period-9653.test.ts (fake-timer driven: already-recorded completion short-circuits, disabled-grace-period finalizes immediately, a completion landing mid-window skips finalize entirely, and no completion ever landing finalizes once the deadline passes). (cherry picked from commit 5d0fe28c4246518f7c7b588795a6da1573b5df16) * chore(quality): rebaseline chatCore.ts for the disconnect grace-period fix Own growth from the disconnect grace-period fix: 5030->5039 (+9, the createClientDisconnectGraceHandler wiring at the existing onClientDisconnectFinalize call site). --------- Co-authored-by: Markus Hartung <mail@hartmark.se>
324 lines
9.1 KiB
TypeScript
324 lines
9.1 KiB
TypeScript
type EnvSource = Record<string, string | undefined>;
|
|
type TimeoutLogger = (message: string) => void;
|
|
|
|
type ReadTimeoutOptions = {
|
|
allowZero?: boolean;
|
|
logger?: TimeoutLogger;
|
|
};
|
|
|
|
export const DEFAULT_FETCH_TIMEOUT_MS = 600_000;
|
|
export const DEFAULT_STREAM_IDLE_TIMEOUT_MS = 600_000;
|
|
export const MAX_TIMER_TIMEOUT_MS = 2_147_483_647;
|
|
export const DEFAULT_SSE_HEARTBEAT_INTERVAL_MS = 15_000;
|
|
export const DEFAULT_STREAM_READINESS_TIMEOUT_MS = 80_000;
|
|
export const DEFAULT_STREAM_READINESS_MAX_TIMEOUT_MS = 180_000;
|
|
export const DEFAULT_FETCH_CONNECT_TIMEOUT_MS = 30_000;
|
|
export const DEFAULT_FETCH_KEEPALIVE_TIMEOUT_MS = 4_000;
|
|
export const DEFAULT_API_BRIDGE_PROXY_TIMEOUT_MS = 600_000;
|
|
export const DEFAULT_API_BRIDGE_SERVER_REQUEST_TIMEOUT_MS = 300_000;
|
|
export const DEFAULT_API_BRIDGE_SERVER_HEADERS_TIMEOUT_MS = 60_000;
|
|
export const DEFAULT_API_BRIDGE_SERVER_KEEPALIVE_TIMEOUT_MS = 5_000;
|
|
export const DEFAULT_API_BRIDGE_SERVER_SOCKET_TIMEOUT_MS = 0;
|
|
// Node's http.Server default keepAliveTimeout is 5_000ms with no Keep-Alive
|
|
// response header hint. Pooled keep-alive clients that don't race that exact
|
|
// window (e.g. the JVM java.net.http.HttpClient used by JetBrains AI
|
|
// Assistant) can reuse a socket the server has already torn down, getting 0
|
|
// response bytes back (#7003). Raise both well above any realistic client
|
|
// idle-pool window, mirroring the API bridge server's pattern.
|
|
export const DEFAULT_MAIN_SERVER_KEEPALIVE_TIMEOUT_MS = 65_000;
|
|
export const DEFAULT_MAIN_SERVER_HEADERS_TIMEOUT_MS = 66_000;
|
|
// A client that closes its connection right after reading a fully-completed
|
|
// SSE stream can race OmniRoute's own completion bookkeeping (#9653): the
|
|
// bytes already reached the client, but the disconnect handler can fire
|
|
// before the stream's own completion callback finishes recording it,
|
|
// persisting a false 499 with zero token usage. Before committing to that
|
|
// failure, wait this long for the real completion to land. Set to 0 to
|
|
// disable and restore the old immediate-fail behavior.
|
|
export const DEFAULT_STREAM_DISCONNECT_GRACE_PERIOD_MS = 10_000;
|
|
|
|
function hasEnvValue(env: EnvSource, name: string): boolean {
|
|
const raw = env[name];
|
|
return raw != null && raw.trim() !== "";
|
|
}
|
|
|
|
export type UpstreamTimeoutConfig = {
|
|
fetchTimeoutMs: number;
|
|
streamIdleTimeoutMs: number;
|
|
sseHeartbeatIntervalMs: number;
|
|
streamReadinessTimeoutMs: number;
|
|
streamReadinessMaxTimeoutMs: number;
|
|
fetchHeadersTimeoutMs: number;
|
|
fetchBodyTimeoutMs: number;
|
|
fetchConnectTimeoutMs: number;
|
|
fetchKeepAliveTimeoutMs: number;
|
|
streamDisconnectGracePeriodMs: number;
|
|
};
|
|
|
|
export type TlsClientTimeoutConfig = {
|
|
timeoutMs: number;
|
|
};
|
|
|
|
export type ApiBridgeTimeoutConfig = {
|
|
proxyTimeoutMs: number;
|
|
serverRequestTimeoutMs: number;
|
|
serverHeadersTimeoutMs: number;
|
|
serverKeepAliveTimeoutMs: number;
|
|
serverSocketTimeoutMs: number;
|
|
};
|
|
|
|
export type MainServerTimeoutConfig = {
|
|
keepAliveTimeoutMs: number;
|
|
headersTimeoutMs: number;
|
|
};
|
|
|
|
function readTimeoutMs(
|
|
env: EnvSource,
|
|
name: string,
|
|
defaultValue: number,
|
|
options: ReadTimeoutOptions = {}
|
|
): number {
|
|
const raw = env[name];
|
|
if (raw == null || raw.trim() === "") return defaultValue;
|
|
|
|
const parsed = Number(raw);
|
|
const isValid = Number.isFinite(parsed) && (options.allowZero ? parsed >= 0 : parsed > 0);
|
|
if (!isValid) {
|
|
options.logger?.(`Invalid ${name}="${raw}". Using default ${defaultValue}ms.`);
|
|
return defaultValue;
|
|
}
|
|
|
|
return Math.floor(parsed);
|
|
}
|
|
|
|
export function getUpstreamTimeoutConfig(
|
|
env: EnvSource = process.env,
|
|
logger?: TimeoutLogger
|
|
): UpstreamTimeoutConfig {
|
|
const sharedRequestTimeoutMs = hasEnvValue(env, "REQUEST_TIMEOUT_MS")
|
|
? readTimeoutMs(env, "REQUEST_TIMEOUT_MS", DEFAULT_FETCH_TIMEOUT_MS, {
|
|
allowZero: true,
|
|
logger,
|
|
})
|
|
: undefined;
|
|
const fetchTimeoutMs = readTimeoutMs(
|
|
env,
|
|
"FETCH_TIMEOUT_MS",
|
|
sharedRequestTimeoutMs ?? DEFAULT_FETCH_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
const streamIdleTimeoutMs = readTimeoutMs(
|
|
env,
|
|
"STREAM_IDLE_TIMEOUT_MS",
|
|
sharedRequestTimeoutMs ?? DEFAULT_STREAM_IDLE_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
const streamReadinessTimeoutMs = readTimeoutMs(
|
|
env,
|
|
"STREAM_READINESS_TIMEOUT_MS",
|
|
sharedRequestTimeoutMs ?? DEFAULT_STREAM_READINESS_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
const streamReadinessMaxTimeoutMs = readTimeoutMs(
|
|
env,
|
|
"STREAM_READINESS_MAX_TIMEOUT_MS",
|
|
DEFAULT_STREAM_READINESS_MAX_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
const sseHeartbeatIntervalMs = readTimeoutMs(
|
|
env,
|
|
"SSE_HEARTBEAT_INTERVAL_MS",
|
|
DEFAULT_SSE_HEARTBEAT_INTERVAL_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
const streamDisconnectGracePeriodMs = readTimeoutMs(
|
|
env,
|
|
"STREAM_DISCONNECT_GRACE_PERIOD_MS",
|
|
DEFAULT_STREAM_DISCONNECT_GRACE_PERIOD_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
|
|
return {
|
|
fetchTimeoutMs,
|
|
streamIdleTimeoutMs,
|
|
streamReadinessTimeoutMs,
|
|
streamReadinessMaxTimeoutMs,
|
|
sseHeartbeatIntervalMs,
|
|
streamDisconnectGracePeriodMs,
|
|
fetchHeadersTimeoutMs: readTimeoutMs(env, "FETCH_HEADERS_TIMEOUT_MS", fetchTimeoutMs, {
|
|
allowZero: true,
|
|
logger,
|
|
}),
|
|
fetchBodyTimeoutMs: readTimeoutMs(env, "FETCH_BODY_TIMEOUT_MS", fetchTimeoutMs, {
|
|
allowZero: true,
|
|
logger,
|
|
}),
|
|
fetchConnectTimeoutMs: readTimeoutMs(
|
|
env,
|
|
"FETCH_CONNECT_TIMEOUT_MS",
|
|
DEFAULT_FETCH_CONNECT_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
),
|
|
fetchKeepAliveTimeoutMs: readTimeoutMs(
|
|
env,
|
|
"FETCH_KEEPALIVE_TIMEOUT_MS",
|
|
DEFAULT_FETCH_KEEPALIVE_TIMEOUT_MS,
|
|
{
|
|
logger,
|
|
}
|
|
),
|
|
};
|
|
}
|
|
|
|
export function getStainlessTimeoutSeconds(
|
|
env: EnvSource = process.env,
|
|
logger?: TimeoutLogger
|
|
): number {
|
|
const { fetchTimeoutMs } = getUpstreamTimeoutConfig(env, logger);
|
|
return Math.max(1, Math.ceil(fetchTimeoutMs / 1_000));
|
|
}
|
|
|
|
export function getTlsClientTimeoutConfig(
|
|
env: EnvSource = process.env,
|
|
logger?: TimeoutLogger
|
|
): TlsClientTimeoutConfig {
|
|
const upstream = getUpstreamTimeoutConfig(env, logger);
|
|
|
|
return {
|
|
timeoutMs: readTimeoutMs(env, "TLS_CLIENT_TIMEOUT_MS", upstream.fetchTimeoutMs, {
|
|
allowZero: true,
|
|
logger,
|
|
}),
|
|
};
|
|
}
|
|
|
|
export function getApiBridgeTimeoutConfig(
|
|
env: EnvSource = process.env,
|
|
logger?: TimeoutLogger
|
|
): ApiBridgeTimeoutConfig {
|
|
const sharedRequestTimeoutMs = hasEnvValue(env, "REQUEST_TIMEOUT_MS")
|
|
? readTimeoutMs(env, "REQUEST_TIMEOUT_MS", DEFAULT_FETCH_TIMEOUT_MS, {
|
|
allowZero: true,
|
|
logger,
|
|
})
|
|
: undefined;
|
|
const proxyTimeoutMs = readTimeoutMs(
|
|
env,
|
|
"API_BRIDGE_PROXY_TIMEOUT_MS",
|
|
sharedRequestTimeoutMs ?? DEFAULT_API_BRIDGE_PROXY_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
const derivedRequestTimeoutMs =
|
|
proxyTimeoutMs > 0
|
|
? Math.max(proxyTimeoutMs, DEFAULT_API_BRIDGE_SERVER_REQUEST_TIMEOUT_MS)
|
|
: DEFAULT_API_BRIDGE_SERVER_REQUEST_TIMEOUT_MS;
|
|
const serverRequestDefaultMs =
|
|
sharedRequestTimeoutMs !== undefined
|
|
? sharedRequestTimeoutMs > 0
|
|
? Math.max(sharedRequestTimeoutMs, derivedRequestTimeoutMs)
|
|
: 0
|
|
: derivedRequestTimeoutMs;
|
|
const serverKeepAliveTimeoutMs = readTimeoutMs(
|
|
env,
|
|
"API_BRIDGE_SERVER_KEEPALIVE_TIMEOUT_MS",
|
|
DEFAULT_API_BRIDGE_SERVER_KEEPALIVE_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
const serverHeadersTimeoutMs = readTimeoutMs(
|
|
env,
|
|
"API_BRIDGE_SERVER_HEADERS_TIMEOUT_MS",
|
|
DEFAULT_API_BRIDGE_SERVER_HEADERS_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
|
|
return {
|
|
proxyTimeoutMs,
|
|
serverRequestTimeoutMs: readTimeoutMs(
|
|
env,
|
|
"API_BRIDGE_SERVER_REQUEST_TIMEOUT_MS",
|
|
serverRequestDefaultMs,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
),
|
|
serverHeadersTimeoutMs:
|
|
serverHeadersTimeoutMs > 0 && serverKeepAliveTimeoutMs > 0
|
|
? Math.max(serverHeadersTimeoutMs, serverKeepAliveTimeoutMs + 1_000)
|
|
: serverHeadersTimeoutMs,
|
|
serverKeepAliveTimeoutMs,
|
|
serverSocketTimeoutMs: readTimeoutMs(
|
|
env,
|
|
"API_BRIDGE_SERVER_SOCKET_TIMEOUT_MS",
|
|
DEFAULT_API_BRIDGE_SERVER_SOCKET_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
),
|
|
};
|
|
}
|
|
|
|
export function getMainServerTimeoutConfig(
|
|
env: EnvSource = process.env,
|
|
logger?: TimeoutLogger
|
|
): MainServerTimeoutConfig {
|
|
const keepAliveTimeoutMs = readTimeoutMs(
|
|
env,
|
|
"MAIN_SERVER_KEEPALIVE_TIMEOUT_MS",
|
|
DEFAULT_MAIN_SERVER_KEEPALIVE_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
const headersTimeoutMs = readTimeoutMs(
|
|
env,
|
|
"MAIN_SERVER_HEADERS_TIMEOUT_MS",
|
|
DEFAULT_MAIN_SERVER_HEADERS_TIMEOUT_MS,
|
|
{
|
|
allowZero: true,
|
|
logger,
|
|
}
|
|
);
|
|
|
|
return {
|
|
keepAliveTimeoutMs,
|
|
// Node requires headersTimeout > keepAliveTimeout to avoid its internal
|
|
// race-condition warning; keep both configurable but always coherent.
|
|
headersTimeoutMs:
|
|
headersTimeoutMs > 0 && keepAliveTimeoutMs > 0
|
|
? Math.max(headersTimeoutMs, keepAliveTimeoutMs + 1_000)
|
|
: headersTimeoutMs,
|
|
};
|
|
}
|