mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-08-02 21:32:10 +03:00
fix(proxy): close dispatchers when clearing cache (#5202)
Integrated into release/v3.8.39
This commit is contained in:
@@ -4,10 +4,18 @@ import { socksDispatcher } from "fetch-socks";
|
||||
import { getUpstreamTimeoutConfig } from "@/shared/utils/runtimeTimeouts";
|
||||
import { stripIpv6Brackets, detectIpLiteralFamily, parseProxyFamily } from "./proxyFamily.ts";
|
||||
import { createSocksDispatcherWithFamily } from "./socksConnectorWithFamily.ts";
|
||||
import {
|
||||
clearDispatcherCache,
|
||||
createRoundRobinDispatcher,
|
||||
getDefaultCachedDispatcher,
|
||||
getDispatcherCache,
|
||||
getRetryCachedDispatcher,
|
||||
setDefaultCachedDispatcher,
|
||||
setRetryCachedDispatcher,
|
||||
} from "./proxyDispatcherCache.ts";
|
||||
|
||||
export { __cacheProxyDispatcherForTest, clearDispatcherCache } from "./proxyDispatcherCache.ts";
|
||||
|
||||
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");
|
||||
const SUPPORTED_PROTOCOLS = new Set(["http:", "https:", "socks5:"]);
|
||||
// Edge-relay proxy types. These do NOT go through an HTTP/SOCKS dispatcher —
|
||||
// the caller wraps the upstream URL with buildRelayHeaders() and fetches the
|
||||
@@ -21,12 +29,6 @@ export function isRelayType(type: string | undefined | null): boolean {
|
||||
const DEFAULT_PROXY_DISPATCHER_CONNECTIONS = 32;
|
||||
const MAX_PROXY_DISPATCHER_CONNECTIONS = 256;
|
||||
|
||||
type DispatcherCache = Map<string, Dispatcher>;
|
||||
type GlobalWithDispatcherCache = typeof globalThis & {
|
||||
[DISPATCHER_CACHE_KEY]?: DispatcherCache;
|
||||
[DEFAULT_DISPATCHER_KEY]?: Dispatcher;
|
||||
[RETRY_DISPATCHER_KEY]?: Dispatcher;
|
||||
};
|
||||
type SocksDispatcherOptions = {
|
||||
type: number;
|
||||
host: string;
|
||||
@@ -43,79 +45,6 @@ type ProxyConfigObject = {
|
||||
family?: string;
|
||||
};
|
||||
|
||||
/**
|
||||
* Direct upstream fan-out dispatcher.
|
||||
*
|
||||
* A single Undici Agent configured with `connections > 1` should be enough in
|
||||
* theory, but real Codex `/backend-api/codex/responses` streams on Node 24 have
|
||||
* still been observed queuing every subsequent same-origin request until the
|
||||
* previous stream emits trailers. Using several one-connection Agents gives
|
||||
* each long SSE stream an independent pool/client and prevents one stream from
|
||||
* monopolizing the effective queue while keeping pipelining disabled.
|
||||
*/
|
||||
class RoundRobinDispatcher {
|
||||
private readonly dispatchers: Dispatcher[];
|
||||
private nextIndex = 0;
|
||||
|
||||
constructor(dispatchers: Dispatcher[]) {
|
||||
this.dispatchers = dispatchers;
|
||||
}
|
||||
|
||||
dispatch(options: Dispatcher.DispatchOptions, handler: Dispatcher.DispatchHandler): boolean {
|
||||
const dispatcher = this.dispatchers[this.nextIndex % this.dispatchers.length];
|
||||
this.nextIndex = (this.nextIndex + 1) % this.dispatchers.length;
|
||||
return dispatcher.dispatch(options, handler);
|
||||
}
|
||||
|
||||
close(callback?: () => void): Promise<void> | void {
|
||||
const done = Promise.all(this.dispatchers.map((dispatcher) => dispatcher.close())).then(
|
||||
() => undefined
|
||||
);
|
||||
if (callback) {
|
||||
done.then(callback);
|
||||
return;
|
||||
}
|
||||
return done;
|
||||
}
|
||||
|
||||
destroy(
|
||||
errorOrCallback?: Error | null | (() => void),
|
||||
callback?: () => void
|
||||
): Promise<void> | void {
|
||||
const callbackFn = typeof errorOrCallback === "function" ? errorOrCallback : callback;
|
||||
const error = typeof errorOrCallback === "function" ? null : (errorOrCallback ?? null);
|
||||
const done = Promise.all(this.dispatchers.map((dispatcher) => dispatcher.destroy(error))).then(
|
||||
() => undefined
|
||||
);
|
||||
if (callbackFn) {
|
||||
done.then(callbackFn);
|
||||
return;
|
||||
}
|
||||
return done;
|
||||
}
|
||||
}
|
||||
|
||||
function getDispatcherCache(): DispatcherCache {
|
||||
const globalWithCache = globalThis as GlobalWithDispatcherCache;
|
||||
if (!globalWithCache[DISPATCHER_CACHE_KEY]) {
|
||||
globalWithCache[DISPATCHER_CACHE_KEY] = new Map();
|
||||
}
|
||||
return globalWithCache[DISPATCHER_CACHE_KEY];
|
||||
}
|
||||
|
||||
/**
|
||||
* Clear all cached proxy dispatchers.
|
||||
* Call this when proxy configuration changes to avoid stale connections.
|
||||
*/
|
||||
export function clearDispatcherCache() {
|
||||
const cache = getDispatcherCache();
|
||||
cache.clear();
|
||||
|
||||
const globalWithCache = globalThis as GlobalWithDispatcherCache;
|
||||
delete globalWithCache[DEFAULT_DISPATCHER_KEY];
|
||||
delete globalWithCache[RETRY_DISPATCHER_KEY];
|
||||
}
|
||||
|
||||
function getDispatcherOptions() {
|
||||
const timeouts = getUpstreamTimeoutConfig(process.env, (message) => {
|
||||
console.warn(`[ProxyDispatcher] ${message}`);
|
||||
@@ -208,17 +137,16 @@ function createRoundRobinDirectDispatcher(connectionLimit: number): Dispatcher {
|
||||
pipelining: 0,
|
||||
};
|
||||
const dispatchers = Array.from({ length: connectionLimit }, () => new Agent(perAgentOptions));
|
||||
return new RoundRobinDispatcher(dispatchers) as unknown as Dispatcher;
|
||||
return createRoundRobinDispatcher(dispatchers);
|
||||
}
|
||||
|
||||
export function getDefaultDispatcher(): Dispatcher {
|
||||
const globalWithCache = globalThis as GlobalWithDispatcherCache;
|
||||
if (!globalWithCache[DEFAULT_DISPATCHER_KEY]) {
|
||||
globalWithCache[DEFAULT_DISPATCHER_KEY] = createRoundRobinDirectDispatcher(
|
||||
getDefaultDispatcherConnectionLimit()
|
||||
);
|
||||
let dispatcher = getDefaultCachedDispatcher();
|
||||
if (!dispatcher) {
|
||||
dispatcher = createRoundRobinDirectDispatcher(getDefaultDispatcherConnectionLimit());
|
||||
setDefaultCachedDispatcher(dispatcher);
|
||||
}
|
||||
return globalWithCache[DEFAULT_DISPATCHER_KEY];
|
||||
return dispatcher;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -235,16 +163,17 @@ export function getDefaultDispatcher(): Dispatcher {
|
||||
* the first attempt is preserved — only the retry pays the fresh-socket cost.
|
||||
*/
|
||||
export function getRetryDispatcher(): Dispatcher {
|
||||
const globalWithCache = globalThis as GlobalWithDispatcherCache;
|
||||
if (!globalWithCache[RETRY_DISPATCHER_KEY]) {
|
||||
globalWithCache[RETRY_DISPATCHER_KEY] = new Agent({
|
||||
let dispatcher = getRetryCachedDispatcher();
|
||||
if (!dispatcher) {
|
||||
dispatcher = new Agent({
|
||||
...getDispatcherOptions(),
|
||||
keepAliveTimeout: 1,
|
||||
keepAliveMaxTimeout: 1,
|
||||
pipelining: 0,
|
||||
});
|
||||
setRetryCachedDispatcher(dispatcher);
|
||||
}
|
||||
return globalWithCache[RETRY_DISPATCHER_KEY];
|
||||
return dispatcher;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -482,7 +411,7 @@ export function __getDefaultDispatcherOptionsForTest(
|
||||
}
|
||||
|
||||
export function __createRoundRobinDispatcherForTest(dispatchers: Dispatcher[]): Dispatcher {
|
||||
return new RoundRobinDispatcher(dispatchers) as unknown as Dispatcher;
|
||||
return createRoundRobinDispatcher(dispatchers);
|
||||
}
|
||||
|
||||
export function createProxyDispatcher(proxyUrl: string): Dispatcher {
|
||||
|
||||
124
open-sse/utils/proxyDispatcherCache.ts
Normal file
124
open-sse/utils/proxyDispatcherCache.ts
Normal file
@@ -0,0 +1,124 @@
|
||||
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");
|
||||
|
||||
type DispatcherCache = Map<string, Dispatcher>;
|
||||
type GlobalWithDispatcherCache = typeof globalThis & {
|
||||
[DISPATCHER_CACHE_KEY]?: DispatcherCache;
|
||||
[DEFAULT_DISPATCHER_KEY]?: Dispatcher;
|
||||
[RETRY_DISPATCHER_KEY]?: Dispatcher;
|
||||
};
|
||||
|
||||
/**
|
||||
* Direct upstream fan-out dispatcher.
|
||||
*
|
||||
* A single Undici Agent configured with `connections > 1` should be enough in
|
||||
* theory, but real Codex `/backend-api/codex/responses` streams on Node 24 have
|
||||
* still been observed queuing every subsequent same-origin request until the
|
||||
* previous stream emits trailers. Using several one-connection Agents gives
|
||||
* each long SSE stream an independent pool/client and prevents one stream from
|
||||
* monopolizing the effective queue while keeping pipelining disabled.
|
||||
*/
|
||||
class RoundRobinDispatcher {
|
||||
private readonly dispatchers: Dispatcher[];
|
||||
private nextIndex = 0;
|
||||
|
||||
constructor(dispatchers: Dispatcher[]) {
|
||||
this.dispatchers = dispatchers;
|
||||
}
|
||||
|
||||
dispatch(options: Dispatcher.DispatchOptions, handler: Dispatcher.DispatchHandler): boolean {
|
||||
const dispatcher = this.dispatchers[this.nextIndex % this.dispatchers.length];
|
||||
this.nextIndex = (this.nextIndex + 1) % this.dispatchers.length;
|
||||
return dispatcher.dispatch(options, handler);
|
||||
}
|
||||
|
||||
close(callback?: () => void): Promise<void> | void {
|
||||
const done = Promise.all(this.dispatchers.map((dispatcher) => dispatcher.close())).then(
|
||||
() => undefined
|
||||
);
|
||||
if (callback) {
|
||||
done.then(callback);
|
||||
return;
|
||||
}
|
||||
return done;
|
||||
}
|
||||
|
||||
destroy(
|
||||
errorOrCallback?: Error | null | (() => void),
|
||||
callback?: () => void
|
||||
): Promise<void> | void {
|
||||
const callbackFn = typeof errorOrCallback === "function" ? errorOrCallback : callback;
|
||||
const error = typeof errorOrCallback === "function" ? null : (errorOrCallback ?? null);
|
||||
const done = Promise.all(this.dispatchers.map((dispatcher) => dispatcher.destroy(error))).then(
|
||||
() => undefined
|
||||
);
|
||||
if (callbackFn) {
|
||||
done.then(callbackFn);
|
||||
return;
|
||||
}
|
||||
return done;
|
||||
}
|
||||
}
|
||||
|
||||
export function createRoundRobinDispatcher(dispatchers: Dispatcher[]): Dispatcher {
|
||||
return new RoundRobinDispatcher(dispatchers) as unknown as Dispatcher;
|
||||
}
|
||||
|
||||
export function getDispatcherCache(): DispatcherCache {
|
||||
const globalWithCache = globalThis as GlobalWithDispatcherCache;
|
||||
if (!globalWithCache[DISPATCHER_CACHE_KEY]) {
|
||||
globalWithCache[DISPATCHER_CACHE_KEY] = new Map();
|
||||
}
|
||||
return globalWithCache[DISPATCHER_CACHE_KEY];
|
||||
}
|
||||
|
||||
export function getDefaultCachedDispatcher(): Dispatcher | undefined {
|
||||
return (globalThis as GlobalWithDispatcherCache)[DEFAULT_DISPATCHER_KEY];
|
||||
}
|
||||
|
||||
export function setDefaultCachedDispatcher(dispatcher: Dispatcher): void {
|
||||
(globalThis as GlobalWithDispatcherCache)[DEFAULT_DISPATCHER_KEY] = dispatcher;
|
||||
}
|
||||
|
||||
export function getRetryCachedDispatcher(): Dispatcher | undefined {
|
||||
return (globalThis as GlobalWithDispatcherCache)[RETRY_DISPATCHER_KEY];
|
||||
}
|
||||
|
||||
export function setRetryCachedDispatcher(dispatcher: Dispatcher): void {
|
||||
(globalThis as GlobalWithDispatcherCache)[RETRY_DISPATCHER_KEY] = dispatcher;
|
||||
}
|
||||
|
||||
function closeDispatcher(dispatcher: Dispatcher | undefined): void {
|
||||
if (!dispatcher) return;
|
||||
try {
|
||||
const result = dispatcher.close();
|
||||
if (result && typeof (result as Promise<void>).catch === "function") {
|
||||
void (result as Promise<void>).catch(() => {});
|
||||
}
|
||||
} catch {}
|
||||
}
|
||||
|
||||
/**
|
||||
* Clear all cached proxy dispatchers.
|
||||
* Call this when proxy configuration changes to avoid stale connections.
|
||||
*/
|
||||
export function clearDispatcherCache(): void {
|
||||
const cache = getDispatcherCache();
|
||||
for (const dispatcher of cache.values()) {
|
||||
closeDispatcher(dispatcher);
|
||||
}
|
||||
cache.clear();
|
||||
|
||||
const globalWithCache = globalThis as GlobalWithDispatcherCache;
|
||||
closeDispatcher(globalWithCache[DEFAULT_DISPATCHER_KEY]);
|
||||
closeDispatcher(globalWithCache[RETRY_DISPATCHER_KEY]);
|
||||
delete globalWithCache[DEFAULT_DISPATCHER_KEY];
|
||||
delete globalWithCache[RETRY_DISPATCHER_KEY];
|
||||
}
|
||||
|
||||
export function __cacheProxyDispatcherForTest(key: string, dispatcher: Dispatcher): void {
|
||||
getDispatcherCache().set(key, dispatcher);
|
||||
}
|
||||
@@ -7,6 +7,7 @@ import {
|
||||
proxyConfigToUrl,
|
||||
normalizeProxyUrl,
|
||||
clearDispatcherCache,
|
||||
__cacheProxyDispatcherForTest,
|
||||
} from "../../open-sse/utils/proxyDispatcher.ts";
|
||||
|
||||
afterEach(() => clearDispatcherCache());
|
||||
@@ -77,4 +78,22 @@ describe("proxyDispatcher connection pool", () => {
|
||||
});
|
||||
assert.equal(options.connections, 256);
|
||||
});
|
||||
|
||||
it("closes cached proxy dispatchers when the proxy cache is cleared", () => {
|
||||
let closeCount = 0;
|
||||
const dispatcher = {
|
||||
dispatch() {
|
||||
return true;
|
||||
},
|
||||
close() {
|
||||
closeCount += 1;
|
||||
},
|
||||
destroy() {},
|
||||
};
|
||||
|
||||
__cacheProxyDispatcherForTest("http://proxy.example.com:8080", dispatcher as never);
|
||||
clearDispatcherCache();
|
||||
|
||||
assert.equal(closeCount, 1);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user