Merge main into model-registry
This commit is contained in:
@@ -27,6 +27,7 @@ import type {
|
||||
AssistantMessage,
|
||||
Context,
|
||||
Model,
|
||||
ProviderEnv,
|
||||
SimpleStreamOptions,
|
||||
StreamFunction,
|
||||
StreamOptions,
|
||||
@@ -40,6 +41,7 @@ import {
|
||||
} from "../utils/diagnostics.ts";
|
||||
import { AssistantMessageEventStream } from "../utils/event-stream.ts";
|
||||
import { headersToRecord } from "../utils/headers.ts";
|
||||
import { resolveHttpProxyUrlForTarget } from "../utils/node-http-proxy.ts";
|
||||
import { clampOpenAIPromptCacheKey } from "./openai-prompt-cache.ts";
|
||||
import { convertResponsesMessages, convertResponsesTools, processResponsesStream } from "./openai-responses-shared.ts";
|
||||
import { buildBaseOptions } from "./simple-options.ts";
|
||||
@@ -53,7 +55,9 @@ 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;
|
||||
// Keep a bounded pre-header timeout so zero-event Codex SSE stalls fail instead of
|
||||
// leaving callers stuck on "Working..." indefinitely. See #4945.
|
||||
const DEFAULT_SSE_HEADER_TIMEOUT_MS = 20_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;
|
||||
@@ -812,19 +816,13 @@ type WebSocketConstructor = new (
|
||||
) => WebSocketLike;
|
||||
|
||||
let _cachedWebsocket: WebSocketConstructor | null = null;
|
||||
async function getWebSocketConstructor(): Promise<WebSocketConstructor | null> {
|
||||
if (_cachedWebsocket) return _cachedWebsocket;
|
||||
async function getWebSocketConstructor(env?: ProviderEnv): Promise<WebSocketConstructor | null> {
|
||||
if (!env && _cachedWebsocket) return _cachedWebsocket;
|
||||
|
||||
// bun doesn't respect http proxy envs, ref: https://github.com/oven-sh/bun/issues/15489
|
||||
// TODO: remove this when bun supports proxy envs in websocket.
|
||||
if (
|
||||
process?.versions?.bun &&
|
||||
(process.env.HTTP_PROXY || process.env.HTTPS_PROXY || process.env.http_proxy || process.env.https_proxy)
|
||||
) {
|
||||
const m = await dynamicImport("proxy-from-env");
|
||||
const getProxyForUrl = (m as { getProxyForUrl: (url: string | object | URL) => string }).getProxyForUrl;
|
||||
|
||||
_cachedWebsocket = class extends WebSocket {
|
||||
if (typeof process !== "undefined" && process.versions?.bun) {
|
||||
const WebSocketWithProxy = class extends WebSocket {
|
||||
constructor(url: string | URL, options?: string | string[] | Record<string, unknown>) {
|
||||
let _opts: Record<string, unknown> = {};
|
||||
if (Array.isArray(options) || typeof options === "string") {
|
||||
@@ -833,11 +831,17 @@ async function getWebSocketConstructor(): Promise<WebSocketConstructor | null> {
|
||||
_opts = { ...options };
|
||||
}
|
||||
|
||||
const proxy = getProxyForUrl(url.toString().replace(/^wss:/, "https:").replace(/^ws:/, "http:"));
|
||||
super(url, { ..._opts, ...(proxy ? { proxy } : {}) } as any);
|
||||
const proxyUrl = resolveHttpProxyUrlForTarget(
|
||||
url.toString().replace(/^wss:/, "https:").replace(/^ws:/, "http:"),
|
||||
env,
|
||||
);
|
||||
super(url, { ..._opts, ...(proxyUrl ? { proxy: proxyUrl.toString() } : {}) } as any);
|
||||
}
|
||||
};
|
||||
return _cachedWebsocket;
|
||||
if (!env) {
|
||||
_cachedWebsocket = WebSocketWithProxy;
|
||||
}
|
||||
return WebSocketWithProxy;
|
||||
}
|
||||
|
||||
const ctor = (globalThis as { WebSocket?: unknown }).WebSocket;
|
||||
@@ -892,8 +896,9 @@ async function connectWebSocket(
|
||||
headers: Headers,
|
||||
signal?: AbortSignal,
|
||||
connectTimeoutMs = DEFAULT_WEBSOCKET_CONNECT_TIMEOUT_MS,
|
||||
env?: ProviderEnv,
|
||||
): Promise<WebSocketLike> {
|
||||
const WebSocketCtor = await getWebSocketConstructor();
|
||||
const WebSocketCtor = await getWebSocketConstructor(env);
|
||||
if (!WebSocketCtor) {
|
||||
throw new Error("WebSocket transport is not available in this runtime");
|
||||
}
|
||||
@@ -970,6 +975,7 @@ async function acquireWebSocket(
|
||||
sessionId: string | undefined,
|
||||
signal?: AbortSignal,
|
||||
connectTimeoutMs?: number,
|
||||
env?: ProviderEnv,
|
||||
): Promise<{
|
||||
socket: WebSocketLike;
|
||||
entry?: CachedWebSocketConnection;
|
||||
@@ -977,7 +983,7 @@ async function acquireWebSocket(
|
||||
release: (options?: { keep?: boolean }) => void;
|
||||
}> {
|
||||
if (!sessionId) {
|
||||
const socket = await connectWebSocket(url, headers, signal, connectTimeoutMs);
|
||||
const socket = await connectWebSocket(url, headers, signal, connectTimeoutMs, env);
|
||||
return {
|
||||
socket,
|
||||
reused: false,
|
||||
@@ -1009,7 +1015,7 @@ async function acquireWebSocket(
|
||||
};
|
||||
}
|
||||
if (cached.busy) {
|
||||
const socket = await connectWebSocket(url, headers, signal, connectTimeoutMs);
|
||||
const socket = await connectWebSocket(url, headers, signal, connectTimeoutMs, env);
|
||||
return {
|
||||
socket,
|
||||
reused: false,
|
||||
@@ -1024,7 +1030,7 @@ async function acquireWebSocket(
|
||||
}
|
||||
}
|
||||
|
||||
const socket = await connectWebSocket(url, headers, signal, connectTimeoutMs);
|
||||
const socket = await connectWebSocket(url, headers, signal, connectTimeoutMs, env);
|
||||
const entry: CachedWebSocketConnection = { socket, busy: true };
|
||||
websocketSessionCache.set(sessionId, entry);
|
||||
return {
|
||||
@@ -1310,6 +1316,7 @@ async function processWebSocketStream(
|
||||
options?.sessionId,
|
||||
options?.signal,
|
||||
websocketConnectTimeoutMs,
|
||||
options?.env,
|
||||
);
|
||||
let keepConnection = true;
|
||||
const useCachedContext = options?.transport === "websocket-cached" || options?.transport === "auto";
|
||||
|
||||
Reference in New Issue
Block a user