fix(ai): timeout Codex SSE header stalls
This commit is contained in:
@@ -33,6 +33,7 @@ import type {
|
||||
StreamOptions,
|
||||
Usage,
|
||||
} from "../types.ts";
|
||||
import { combineAbortSignals } from "../utils/abort-signals.ts";
|
||||
import {
|
||||
appendAssistantMessageDiagnostic,
|
||||
createAssistantMessageDiagnostic,
|
||||
@@ -53,6 +54,7 @@ const JWT_CLAIM_PATH = "https://api.openai.com/auth" as const;
|
||||
const DEFAULT_MAX_RETRIES = 0;
|
||||
const BASE_DELAY_MS = 1000;
|
||||
const DEFAULT_MAX_RETRY_DELAY_MS = 60_000;
|
||||
const DEFAULT_SSE_HEADER_TIMEOUT_MS = 10_000;
|
||||
const DEFAULT_WEBSOCKET_CONNECT_TIMEOUT_MS = 15_000;
|
||||
const CODEX_TOOL_CALL_PROVIDERS = new Set(["openai", "openai-codex", "opencode"]);
|
||||
const WEBSOCKET_MESSAGE_TOO_BIG_CLOSE_CODE = 1009;
|
||||
@@ -172,6 +174,20 @@ function normalizeTimeoutMs(value: number | undefined): number | undefined {
|
||||
return Math.floor(value);
|
||||
}
|
||||
|
||||
function createSSEHeaderTimeout(): { signal: AbortSignal; clear: () => void; error: () => Error | undefined } {
|
||||
const controller = new AbortController();
|
||||
let error: Error | undefined;
|
||||
const timeout = setTimeout(() => {
|
||||
error = new Error(`Codex SSE response headers timed out after ${DEFAULT_SSE_HEADER_TIMEOUT_MS}ms`);
|
||||
controller.abort(error);
|
||||
}, DEFAULT_SSE_HEADER_TIMEOUT_MS);
|
||||
return {
|
||||
signal: controller.signal,
|
||||
clear: () => clearTimeout(timeout),
|
||||
error: () => error,
|
||||
};
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Main Stream Function
|
||||
// ============================================================================
|
||||
@@ -294,12 +310,22 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses"
|
||||
}
|
||||
|
||||
try {
|
||||
response = await fetch(resolveCodexUrl(model.baseUrl), {
|
||||
method: "POST",
|
||||
headers: sseHeaders,
|
||||
body: bodyJson,
|
||||
signal: options?.signal,
|
||||
});
|
||||
const headerTimeout = createSSEHeaderTimeout();
|
||||
const combinedSignal = combineAbortSignals([options?.signal, headerTimeout.signal]);
|
||||
try {
|
||||
response = await fetch(resolveCodexUrl(model.baseUrl), {
|
||||
method: "POST",
|
||||
headers: sseHeaders,
|
||||
body: bodyJson,
|
||||
signal: combinedSignal.signal,
|
||||
});
|
||||
} catch (error) {
|
||||
const timeoutError = headerTimeout.error();
|
||||
throw timeoutError && !options?.signal?.aborted ? timeoutError : error;
|
||||
} finally {
|
||||
combinedSignal.cleanup();
|
||||
headerTimeout.clear();
|
||||
}
|
||||
await options?.onResponse?.(
|
||||
{ status: response.status, headers: headersToRecord(response.headers) },
|
||||
model,
|
||||
|
||||
41
packages/ai/src/utils/abort-signals.ts
Normal file
41
packages/ai/src/utils/abort-signals.ts
Normal file
@@ -0,0 +1,41 @@
|
||||
export interface CombinedAbortSignal {
|
||||
signal?: AbortSignal;
|
||||
cleanup: () => void;
|
||||
}
|
||||
|
||||
export function combineAbortSignals(signals: readonly (AbortSignal | undefined)[]): CombinedAbortSignal {
|
||||
const activeSignals = signals.filter((signal): signal is AbortSignal => signal !== undefined);
|
||||
if (activeSignals.length === 0) {
|
||||
return { cleanup: () => {} };
|
||||
}
|
||||
if (activeSignals.length === 1) {
|
||||
return { signal: activeSignals[0], cleanup: () => {} };
|
||||
}
|
||||
|
||||
const controller = new AbortController();
|
||||
const listeners: Array<{ signal: AbortSignal; listener: () => void }> = [];
|
||||
const abort = (signal: AbortSignal) => {
|
||||
if (!controller.signal.aborted) {
|
||||
controller.abort(signal.reason);
|
||||
}
|
||||
};
|
||||
|
||||
for (const signal of activeSignals) {
|
||||
if (signal.aborted) {
|
||||
abort(signal);
|
||||
break;
|
||||
}
|
||||
const listener = () => abort(signal);
|
||||
signal.addEventListener("abort", listener, { once: true });
|
||||
listeners.push({ signal, listener });
|
||||
}
|
||||
|
||||
return {
|
||||
signal: controller.signal,
|
||||
cleanup: () => {
|
||||
for (const { signal, listener } of listeners) {
|
||||
signal.removeEventListener("abort", listener);
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user