Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 0 additions & 5 deletions config/quality/eslint-suppressions.json
Original file line number Diff line number Diff line change
Expand Up @@ -713,11 +713,6 @@
"count": 2
}
},
"open-sse/utils/proxyDispatcher.ts": {
"@typescript-eslint/no-unused-vars": {
"count": 1
}
},
"open-sse/utils/setupPolyfill.ts": {
"@typescript-eslint/no-explicit-any": {
"count": 5
Expand Down
80 changes: 72 additions & 8 deletions open-sse/utils/proxyDispatcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,16 @@ import { getUpstreamTimeoutConfig } from "@/shared/utils/runtimeTimeouts";
import { stripIpv6Brackets, detectIpLiteralFamily, parseProxyFamily } from "./proxyFamily.ts";
import { createSocksDispatcherWithFamily } from "./socksConnectorWithFamily.ts";
import {
clearDispatcherCache,
createRoundRobinDispatcher,
getDefaultCachedDispatcher,
getDispatcherCache,
getLocalDefaultCachedDispatcher,
getLocalRetryCachedDispatcher,
getRetryCachedDispatcher,
setDefaultCachedDispatcher,
setDispatcherCacheEntry,
setLocalDefaultCachedDispatcher,
setLocalRetryCachedDispatcher,
setRetryCachedDispatcher,
} from "./proxyDispatcherCache.ts";

Expand All @@ -27,6 +30,22 @@ export const RELAY_TYPES: ReadonlySet<string> = new Set(["vercel", "deno", "clou
export function isRelayType(type: string | undefined | null): boolean {
return typeof type === "string" && RELAY_TYPES.has(type);
}

// Local-egress hostnames: host.docker.internal, *.internal, *.local.
// Match the same shape the proxyFetch.ts isLocalAddress() helper uses for
// PROXY bypass, but narrower on purpose: we only switch dispatcher options
// for mDNS-style hostnames, not RFC1918 IPs (those may still be cloud
// upstreams via a private tunnel). IPv6 brackets are stripped defensively.
const LOCAL_EGRESS_HOSTNAME_REGEX = /(?:^|\.)(?:internal|local)$/i;
const LOCAL_KEEPALIVE_MAX_TIMEOUT_MS = 1000;
const LOCAL_AUTO_SELECT_FAMILY_ATTEMPT_TIMEOUT_MS = 200;

export function isLocalEgressHostname(hostname: string | null | undefined): boolean {
if (!hostname) return false;
// Tolerate both bare hostname and URL.host (host:port), and IPv6 brackets.
const host = hostname.replace(/^\[/, "").replace(/\]$/, "").replace(/:\d+$/, "");
return LOCAL_EGRESS_HOSTNAME_REGEX.test(host);
}
const DEFAULT_PROXY_DISPATCHER_CONNECTIONS = 32;
const MAX_PROXY_DISPATCHER_CONNECTIONS = 256;

Expand All @@ -46,10 +65,11 @@ type ProxyConfigObject = {
family?: string;
};

function getDispatcherOptions() {
function getDispatcherOptions(hostname?: string) {
const timeouts = getUpstreamTimeoutConfig(process.env, (message) => {
console.warn(`[ProxyDispatcher] ${message}`);
});
const localEgress = isLocalEgressHostname(hostname);

return {
headersTimeout: timeouts.fetchHeadersTimeoutMs,
Expand All @@ -59,7 +79,13 @@ function getDispatcherOptions() {
// Without this, an upstream Keep-Alive: timeout=N header clamps
// keepAliveTimeout UP to undici's default keepAliveMaxTimeout (600 s),
// completely overriding the configured 1 s and restoring zombie-socket risk.
keepAliveMaxTimeout: timeouts.fetchKeepAliveTimeoutMs,
// For local-egress hostnames (host.docker.internal / *.internal / *.local)
// Docker Desktop's NAT silently drops idle keep-alive sockets well inside
// the default window, so cap keep-alive at 1 s on that path to force fresh
// sockets before the next request lands on a stale one.
keepAliveMaxTimeout: localEgress
? LOCAL_KEEPALIVE_MAX_TIMEOUT_MS
: timeouts.fetchKeepAliveTimeoutMs,
// 9router#1237: RFC 8305 Happy Eyeballs. undici does not
// enable it by default, so when DNS returns both AAAA (IPv6) and A (IPv4)
// and the IPv6 route is broken (e.g. NAT64 `64:ff9b::` without routing),
Expand All @@ -71,9 +97,14 @@ function getDispatcherOptions() {
// requires `port`; at runtime undici merges these into net.connect (the origin
// already carries host:port), so the partial pin is valid — cast to suppress
// the spurious missing-`port` error, mirroring the `proxyTls` cast below.
// Local-egress path shortens the per-family attempt to 200 ms because the
// IPv6 route to host.docker.internal is dead inside the container (verified
// 2026-09-21) and the default 1 s wait is pure latency on every healthy request.
connect: {
autoSelectFamily: true,
autoSelectFamilyAttemptTimeout: 1000,
autoSelectFamilyAttemptTimeout: localEgress
? LOCAL_AUTO_SELECT_FAMILY_ATTEMPT_TIMEOUT_MS
: 1000,
} as ProxyAgent.Options["proxyTls"],
};
}
Expand Down Expand Up @@ -152,8 +183,8 @@ function getDefaultDispatcherOptions(env: Record<string, string | undefined> = p
};
}

function createRoundRobinDirectDispatcher(connectionLimit: number): Dispatcher {
const baseOptions = getDispatcherOptions();
function createRoundRobinDirectDispatcher(connectionLimit: number, hostname?: string): Dispatcher {
const baseOptions = getDispatcherOptions(hostname);
const perAgentOptions = {
...baseOptions,
connections: 1,
Expand All @@ -163,7 +194,18 @@ function createRoundRobinDirectDispatcher(connectionLimit: number): Dispatcher {
return createRoundRobinDispatcher(dispatchers);
}

export function getDefaultDispatcher(): Dispatcher {
export function getDefaultDispatcher(hostname?: string): Dispatcher {
if (isLocalEgressHostname(hostname)) {
let dispatcher = getLocalDefaultCachedDispatcher();
if (!dispatcher) {
dispatcher = createRoundRobinDirectDispatcher(
getDefaultDispatcherConnectionLimit(),
hostname
);
setLocalDefaultCachedDispatcher(dispatcher);
}
return dispatcher;
}
let dispatcher = getDefaultCachedDispatcher();
if (!dispatcher) {
dispatcher = createRoundRobinDirectDispatcher(getDefaultDispatcherConnectionLimit());
Expand All @@ -184,8 +226,25 @@ export function getDefaultDispatcher(): Dispatcher {
* retry uses this no-keep-alive / no-pipelining dispatcher (mirroring the proxy
* dispatcher mitigation) to force a fresh socket. Healthy keep-alive reuse on
* the first attempt is preserved — only the retry pays the fresh-socket cost.
*
* Local-egress hostnames (host.docker.internal / *.internal / *.local) route
* to a parallel retry cache so a fresh-socket retry cannot pick up a stale
* socket from the cloud-upstream pool.
*/
export function getRetryDispatcher(): Dispatcher {
export function getRetryDispatcher(hostname?: string): Dispatcher {
if (isLocalEgressHostname(hostname)) {
let dispatcher = getLocalRetryCachedDispatcher();
if (!dispatcher) {
dispatcher = new Agent({
...getDispatcherOptions(hostname),
keepAliveTimeout: 1,
keepAliveMaxTimeout: 1,
pipelining: 0,
});
setLocalRetryCachedDispatcher(dispatcher);
}
return dispatcher;
}
let dispatcher = getRetryCachedDispatcher();
if (!dispatcher) {
dispatcher = new Agent({
Expand Down Expand Up @@ -432,6 +491,11 @@ export function __getDefaultDispatcherOptionsForTest(
return getDefaultDispatcherOptions(env);
}

/** Test-only accessor for the hostname-branched dispatcher options (local-egress shortening). */
export function __getDispatcherOptionsForTest(hostname?: string) {
return getDispatcherOptions(hostname);
}

export function __createRoundRobinDispatcherForTest(dispatchers: Dispatcher[]): Dispatcher {
return createRoundRobinDispatcher(dispatchers);
}
Expand Down
29 changes: 29 additions & 0 deletions open-sse/utils/proxyDispatcherCache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,13 @@ import type { Dispatcher } from "undici";
const DISPATCHER_CACHE_KEY = Symbol.for("omniroute.proxyDispatcher.cache");
const DEFAULT_DISPATCHER_KEY = Symbol.for("omniroute.proxyDispatcher.default");
const RETRY_DISPATCHER_KEY = Symbol.for("omniroute.proxyDispatcher.retry");
// Local-egress dispatchers: separate cache for hostnames like
// host.docker.internal / *.internal / *.local, where Docker Desktop's NAT
// silently drops idle keep-alive sockets within the global pool's
// keepAliveMaxTimeout window. Kept on their own cache so a wider keep-alive
// for cloud upstreams cannot pull the .internal sockets down with it.
const LOCAL_DEFAULT_DISPATCHER_KEY = Symbol.for("omniroute.proxyDispatcher.localDefault");
const LOCAL_RETRY_DISPATCHER_KEY = Symbol.for("omniroute.proxyDispatcher.localRetry");

/** Upper bound on cached per-URL proxy dispatchers; oldest entries are evicted first. */
const MAX_DISPATCHER_CACHE_ENTRIES = 512;
Expand All @@ -12,6 +19,8 @@ type GlobalWithDispatcherCache = typeof globalThis & {
[DISPATCHER_CACHE_KEY]?: DispatcherCache;
[DEFAULT_DISPATCHER_KEY]?: Dispatcher;
[RETRY_DISPATCHER_KEY]?: Dispatcher;
[LOCAL_DEFAULT_DISPATCHER_KEY]?: Dispatcher;
[LOCAL_RETRY_DISPATCHER_KEY]?: Dispatcher;
};

/**
Expand Down Expand Up @@ -94,6 +103,22 @@ export function setRetryCachedDispatcher(dispatcher: Dispatcher): void {
(globalThis as GlobalWithDispatcherCache)[RETRY_DISPATCHER_KEY] = dispatcher;
}

export function getLocalDefaultCachedDispatcher(): Dispatcher | undefined {
return (globalThis as GlobalWithDispatcherCache)[LOCAL_DEFAULT_DISPATCHER_KEY];
}

export function setLocalDefaultCachedDispatcher(dispatcher: Dispatcher): void {
(globalThis as GlobalWithDispatcherCache)[LOCAL_DEFAULT_DISPATCHER_KEY] = dispatcher;
}

export function getLocalRetryCachedDispatcher(): Dispatcher | undefined {
return (globalThis as GlobalWithDispatcherCache)[LOCAL_RETRY_DISPATCHER_KEY];
}

export function setLocalRetryCachedDispatcher(dispatcher: Dispatcher): void {
(globalThis as GlobalWithDispatcherCache)[LOCAL_RETRY_DISPATCHER_KEY] = dispatcher;
}

function closeDispatcher(dispatcher: Dispatcher | undefined): void {
if (!dispatcher) return;
try {
Expand All @@ -118,8 +143,12 @@ export function clearDispatcherCache(): void {
const globalWithCache = globalThis as GlobalWithDispatcherCache;
closeDispatcher(globalWithCache[DEFAULT_DISPATCHER_KEY]);
closeDispatcher(globalWithCache[RETRY_DISPATCHER_KEY]);
closeDispatcher(globalWithCache[LOCAL_DEFAULT_DISPATCHER_KEY]);
closeDispatcher(globalWithCache[LOCAL_RETRY_DISPATCHER_KEY]);
delete globalWithCache[DEFAULT_DISPATCHER_KEY];
delete globalWithCache[RETRY_DISPATCHER_KEY];
delete globalWithCache[LOCAL_DEFAULT_DISPATCHER_KEY];
delete globalWithCache[LOCAL_RETRY_DISPATCHER_KEY];
}

export function __cacheProxyDispatcherForTest(key: string, dispatcher: Dispatcher): void {
Expand Down
24 changes: 23 additions & 1 deletion open-sse/utils/proxyFetch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,12 @@ import { AsyncLocalStorage } from "node:async_hooks";
import { fetch as undiciFetch, Agent } from "undici";
import {
buildVercelRelayHeaders,
clearDispatcherCache,
createProxyDispatcher,
getDefaultDispatcher,
getProxyRetryDispatcher,
getRetryDispatcher,
isLocalEgressHostname,
isRelayType,
normalizeProxyUrl,
proxyConfigToUrl,
Expand Down Expand Up @@ -112,6 +114,7 @@ const TLS_PROVIDER_PROFILE: Record<string, { browser: string; os: string }> = {

type TlsProfileResult = { browserProfile?: string; os?: string };
function tlsProfileForProvider(provider: string | null | undefined): TlsProfileResult {

if (!provider) return {};
const p = TLS_PROVIDER_PROFILE[provider.trim().toLowerCase()];
return p ? { browserProfile: p.browser, os: p.os } : {};
Expand Down Expand Up @@ -857,11 +860,18 @@ async function patchedFetchUnrecorded(
}
for (let attempt = 0; attempt < maxAttempts; attempt++) {
try {
let hostnameForDispatcher: string | undefined;
try {
hostnameForDispatcher = new URL(targetUrl).hostname;
} catch {}
return await directFetchWithBoundedResponseStart(
input,
{
...options,
dispatcher: attempt === 0 ? getDefaultDispatcher() : getRetryDispatcher(),
dispatcher:
attempt === 0
? getDefaultDispatcher(hostnameForDispatcher)
: getRetryDispatcher(hostnameForDispatcher),
},
_undiciDirect,
resolveDirectHeadersTimeoutMs(undefined, directBodyForTimeout, attempt, !!options.signal)
Expand Down Expand Up @@ -948,6 +958,18 @@ async function patchedFetchUnrecorded(
console.warn(
`[ProxyFetch] Undici dispatcher failed, falling back to native fetch (after retry): ${describeFetchCause(dispatcherError)}`
);
// On PROXY_UNREACHABLE for local-egress hostnames (host.docker.internal,
// *.internal, *.local), drop the cached dispatcher pool: Docker
// Desktop's NAT silently drops idle keep-alive sockets inside the
// round-robin pool's keepAliveMaxTimeout window, and the pool never
// reaps them on PROXY_UNREACHABLE, so the next request must rebuild
// with fresh sockets (#4252-style stale-socket burst mitigation).
if (
isLocalEgressHostname(targetHostForLogs) &&
isProxyUnreachableError(dispatcherError)
) {
clearDispatcherCache();
}
try {
return await _nativeFallback(input, options);
} catch (nativeError) {
Expand Down
Loading