diff --git a/open-sse/utils/proxyDispatcher.ts b/open-sse/utils/proxyDispatcher.ts index aab0050a66..df120d11a6 100644 --- a/open-sse/utils/proxyDispatcher.ts +++ b/open-sse/utils/proxyDispatcher.ts @@ -108,10 +108,42 @@ function getProxyDispatcherOptions(env: Record = pro }; } +export function getDefaultDispatcherConnectionLimit( + env: Record = process.env +): number { + const raw = env.OMNIROUTE_DIRECT_DISPATCHER_CONNECTIONS; + if (raw == null || raw.trim() === "") return DEFAULT_PROXY_DISPATCHER_CONNECTIONS; + + const parsed = Number(raw); + if (!Number.isFinite(parsed) || parsed < 1) { + console.warn( + `[ProxyDispatcher] Invalid OMNIROUTE_DIRECT_DISPATCHER_CONNECTIONS="${raw}". Using default ${DEFAULT_PROXY_DISPATCHER_CONNECTIONS}.` + ); + return DEFAULT_PROXY_DISPATCHER_CONNECTIONS; + } + + return Math.min(Math.floor(parsed), MAX_PROXY_DISPATCHER_CONNECTIONS); +} + +function getDefaultDispatcherOptions(env: Record = process.env) { + const options = getDispatcherOptions(); + // #4580 — On the direct egress path, undici's default pipelining (1) lets a long + // SSE stream monopolize the single pooled socket per origin, serializing every + // other concurrent request to that same provider. Mirror the proxy fix (#4288): + // disable pipelining and keep several connections available. Unlike the proxy + // path we KEEP keep-alive — the 1ms TTL there is a cheap-proxy-socket workaround, + // not needed (and harmful to perf) for direct connections. + return { + ...options, + connections: getDefaultDispatcherConnectionLimit(env), + pipelining: 0, + }; +} + export function getDefaultDispatcher(): Dispatcher { const globalWithCache = globalThis as GlobalWithDispatcherCache; if (!globalWithCache[DEFAULT_DISPATCHER_KEY]) { - globalWithCache[DEFAULT_DISPATCHER_KEY] = new Agent(getDispatcherOptions()); + globalWithCache[DEFAULT_DISPATCHER_KEY] = new Agent(getDefaultDispatcherOptions()); } return globalWithCache[DEFAULT_DISPATCHER_KEY]; } @@ -362,6 +394,12 @@ export function __getProxyDispatcherOptionsForTest( return getProxyDispatcherOptions(env); } +export function __getDefaultDispatcherOptionsForTest( + env: Record = process.env +) { + return getDefaultDispatcherOptions(env); +} + export function createProxyDispatcher(proxyUrl: string): Dispatcher { const normalizedUrl = normalizeProxyUrl(proxyUrl, "proxy dispatcher"); const dispatcherCache = getDispatcherCache(); diff --git a/tests/unit/direct-dispatcher-pipelining-4580.test.ts b/tests/unit/direct-dispatcher-pipelining-4580.test.ts new file mode 100644 index 0000000000..a56ed90831 --- /dev/null +++ b/tests/unit/direct-dispatcher-pipelining-4580.test.ts @@ -0,0 +1,52 @@ +import { describe, it, afterEach } from "node:test"; +import assert from "node:assert/strict"; +import { + __getDefaultDispatcherOptionsForTest, + __getProxyDispatcherOptionsForTest, + getDefaultDispatcherConnectionLimit, + clearDispatcherCache, +} from "../../open-sse/utils/proxyDispatcher.ts"; + +afterEach(() => clearDispatcherCache()); + +// #4580 — On the DIRECT egress path, concurrent same-provider requests serialized +// behind a long/streaming request. The proxy dispatcher already got pipelining:0 + +// a connections cap in #4288, but the first-attempt direct dispatcher +// (getDispatcherOptions → new Agent) kept undici's default pipelining (1), so long +// SSE streams bottlenecked the single pooled socket. The direct dispatcher now +// mirrors that fix while KEEPING keep-alive (a proxy-only concern was the 1ms TTL). + +describe("#4580 direct dispatcher options", () => { + it("disables pipelining so concurrent streams open separate sockets", () => { + const opts = __getDefaultDispatcherOptionsForTest({}); + assert.equal(opts.pipelining, 0); + }); + + it("caps connections to a finite number (default 32)", () => { + const opts = __getDefaultDispatcherOptionsForTest({}); + assert.equal(typeof opts.connections, "number"); + assert.equal(opts.connections, 32); + }); + + it("preserves keep-alive (NOT the 1ms TTL the proxy path forces)", () => { + const direct = __getDefaultDispatcherOptionsForTest({}); + const proxy = __getProxyDispatcherOptionsForTest({}); + assert.equal(proxy.keepAliveTimeout, 1); + assert.ok( + (direct.keepAliveTimeout ?? 0) > 1, + `direct keepAliveTimeout should stay > 1 (got ${direct.keepAliveTimeout})` + ); + }); + + it("connection limit honors OMNIROUTE_DIRECT_DISPATCHER_CONNECTIONS", () => { + assert.equal(getDefaultDispatcherConnectionLimit({ OMNIROUTE_DIRECT_DISPATCHER_CONNECTIONS: "8" }), 8); + }); + + it("connection limit clamps invalid values to the default", () => { + assert.equal( + getDefaultDispatcherConnectionLimit({ OMNIROUTE_DIRECT_DISPATCHER_CONNECTIONS: "nonsense" }), + 32 + ); + assert.equal(getDefaultDispatcherConnectionLimit({}), 32); + }); +});