From d8210cf0113cf6aa782f856db683e9123cb27711 Mon Sep 17 00:00:00 2001 From: Tam Nhu Tran Date: Sat, 8 Aug 2026 21:17:13 -0400 Subject: [PATCH] fix(proxy): stabilize upstream fetch transport Refs #1686 --- package.json | 1 + scripts/ci-parity-gate.sh | 1 + src/ccs.ts | 8 +- src/proxy/proxy-daemon-entry.ts | 14 +- src/proxy/server/messages-route.ts | 93 ++++---- src/proxy/server/proxy-server.ts | 74 ++++++- src/proxy/server/upstream-transport.ts | 93 ++++++++ src/services/logging/log-redaction.ts | 12 +- src/services/logging/log-types.ts | 1 + src/utils/fetch-proxy-setup.ts | 23 +- .../proxy/fixtures/runtime-probe.cjs | 100 +++++++++ .../proxy/runtime-transport-matrix.test.ts | 200 ++++++++++++++++++ .../proxy/runtime-transport.test.ts | 195 +++++++++++++++++ tests/unit/proxy/messages-route.test.ts | 78 ++++++- .../logging/cli-entry-log-redaction.test.ts | 4 +- .../services/logging/log-redaction.test.ts | 33 ++- tests/unit/utils/fetch-proxy-setup.test.ts | 44 ++-- 17 files changed, 889 insertions(+), 85 deletions(-) create mode 100644 src/proxy/server/upstream-transport.ts create mode 100644 tests/integration/proxy/fixtures/runtime-probe.cjs create mode 100644 tests/integration/proxy/runtime-transport-matrix.test.ts create mode 100644 tests/integration/proxy/runtime-transport.test.ts diff --git a/package.json b/package.json index 4bfbfe9c..ef686cd3 100644 --- a/package.json +++ b/package.json @@ -81,6 +81,7 @@ "test:unit": "bun test tests/unit", "test:npm": "bun test tests/npm/", "test:native": "bash tests/native/unix/edge-cases.sh", + "test:runtime-matrix": "bun run build && CCS_RUNTIME_MATRIX=1 bun test tests/integration/proxy/runtime-transport-matrix.test.ts --timeout 60000", "test:e2e": "bun test tests/e2e/ --bail --timeout 60000", "report:hardening": "node scripts/hardening-inventory.js", "dev": "bun run build:server && node dist/ccs.js config --dev", diff --git a/scripts/ci-parity-gate.sh b/scripts/ci-parity-gate.sh index 796abd29..9616c92f 100755 --- a/scripts/ci-parity-gate.sh +++ b/scripts/ci-parity-gate.sh @@ -67,6 +67,7 @@ echo "[i] Running CI-parity local checks..." bun run typecheck bun run lint bun run format:check +bun run test:runtime-matrix bun run build:all bun run test:all CCS_E2E_SKIP_BUILD=1 bun run test:e2e diff --git a/src/ccs.ts b/src/ccs.ts index e8fe5996..7157997c 100644 --- a/src/ccs.ts +++ b/src/ccs.ts @@ -1,5 +1,3 @@ -import './utils/fetch-proxy-setup'; - import { ErrorManager } from './utils/error-manager'; import { fail } from './utils/ui'; // Import centralized error handling @@ -7,6 +5,7 @@ import { handleError, runCleanup } from './errors'; import { createLogger, runWithRequestId } from './services/logging'; import { redactArgv } from './services/logging/log-redaction'; +import { applyGlobalFetchProxy } from './utils/fetch-proxy-setup'; // Import target adapter system import { registerTarget, ClaudeAdapter, DroidAdapter, CodexAdapter } from './targets'; @@ -19,6 +18,11 @@ import { resolveProfileAndTarget } from './dispatcher/profile-resolver'; // ========== Main Execution ========== async function main(): Promise { + const fetchProxySetup = applyGlobalFetchProxy(); + if (fetchProxySetup.error) { + console.error(`[!] Skipping global fetch proxy setup: ${fetchProxySetup.error}`); + } + // Register target adapters (singleton wiring — stays in main) registerTarget(new ClaudeAdapter()); registerTarget(new DroidAdapter()); diff --git a/src/proxy/proxy-daemon-entry.ts b/src/proxy/proxy-daemon-entry.ts index 4db89231..29c34d88 100644 --- a/src/proxy/proxy-daemon-entry.ts +++ b/src/proxy/proxy-daemon-entry.ts @@ -1,7 +1,7 @@ import * as fs from 'fs'; import { resolveOpenAICompatProfileConfig } from './profile-router'; import { OPENAI_COMPAT_PROXY_DEFAULT_PORT } from './proxy-daemon-paths'; -import { startOpenAICompatProxyServer } from './server/proxy-server'; +import { closeOpenAICompatProxyServer, startOpenAICompatProxyServer } from './server/proxy-server'; import { loadSettings } from '../config/config-loader-facade'; interface RuntimeOptions { @@ -99,9 +99,17 @@ function startRuntime(options: RuntimeOptions): void { }); server.once('error', (error) => { process.stderr.write(String((error as Error).message) + '\n'); - process.exit(1); + void closeOpenAICompatProxyServer(server).finally(() => { + process.exitCode = 1; + }); }); - const shutdown = () => server.close(); + let shutdownPromise: Promise | undefined; + const shutdown = () => { + shutdownPromise ??= closeOpenAICompatProxyServer(server).catch((error) => { + process.stderr.write(`Proxy shutdown failed: ${String((error as Error).message)}\n`); + process.exitCode = 1; + }); + }; process.on('SIGTERM', shutdown); process.on('SIGINT', shutdown); } diff --git a/src/proxy/server/messages-route.ts b/src/proxy/server/messages-route.ts index 4a5f2f93..6cd08c15 100644 --- a/src/proxy/server/messages-route.ts +++ b/src/proxy/server/messages-route.ts @@ -1,5 +1,4 @@ import * as http from 'http'; -import { Agent } from 'undici'; import type { Dispatcher } from 'undici'; import type { OpenAICompatProfileConfig } from '../profile-router'; import { resolveProxyRequestRoute } from '../request-router'; @@ -10,11 +9,14 @@ import { import { ProxySseStreamTransformer } from '../transformers/sse-stream-transformer'; import { isAnthropicPassthroughProfile, resolveOpenAIChatCompletionsUrl } from '../upstream-url'; import { createLogger } from '../../services/logging'; -import { - createGlobalFetchProxyDispatcher, - type UpstreamAgentTimeoutOptions, -} from '../../utils/fetch-proxy-setup'; +import type { UpstreamAgentTimeoutOptions } from '../../utils/fetch-proxy-setup'; import { pipeWebResponseToNode, readRawBody, writeJson } from './http-helpers'; +import { + closeUpstreamDispatcher, + createUpstreamDispatcher, + fetchWithUpstreamTransport, + toLogErrorInfo, +} from './upstream-transport'; const REQUEST_TIMEOUT_MS = 600_000; // Keep undici's per-phase timeouts above the explicit request timeout so the @@ -421,17 +423,6 @@ export function buildUpstreamAgentTimeouts(): UpstreamAgentTimeoutOptions { return { headersTimeout: ceiling, bodyTimeout: ceiling }; } -let defaultUpstreamDispatcher: Dispatcher | null = null; - -function getDefaultUpstreamDispatcher(): Dispatcher { - if (!defaultUpstreamDispatcher) { - const timeouts = buildUpstreamAgentTimeouts(); - // Honor HTTP(S)_PROXY routing when configured; otherwise a plain Agent. - defaultUpstreamDispatcher = createGlobalFetchProxyDispatcher(timeouts) ?? new Agent(timeouts); - } - return defaultUpstreamDispatcher; -} - function formatTimeoutDuration(timeoutMs: number): string { return timeoutMs % 1000 === 0 ? `${timeoutMs / 1000} seconds` : `${timeoutMs}ms`; } @@ -489,7 +480,10 @@ export async function handleProxyMessagesRequest( res: http.ServerResponse, profile: OpenAICompatProfileConfig, expectedAuthToken: string, - insecureDispatcher?: Dispatcher + insecureDispatcher?: Dispatcher, + sharedUpstreamDispatcher?: Dispatcher, + upstreamFetch?: typeof globalThis.fetch, + forceInsecureTls = false ): Promise { const transformer = new ProxySseStreamTransformer(); const startedAt = Date.now(); @@ -563,20 +557,26 @@ export async function handleProxyMessagesRequest( } ); + const routeStaysOnActiveProfile = upstream.route.profile.profileName === profile.profileName; const useSharedInsecureDispatcher = - insecureDispatcher !== undefined && - upstream.route.profile.profileName === profile.profileName; + insecureDispatcher !== undefined && routeStaysOnActiveProfile; + const useServerInsecureTls = forceInsecureTls && routeStaysOnActiveProfile; const useProfileInsecureTls = upstream.route.profile.insecure === true; const ephemeralInsecureDispatcher = - useProfileInsecureTls && !useSharedInsecureDispatcher - ? new Agent({ connect: { rejectUnauthorized: false }, ...buildUpstreamAgentTimeouts() }) + useProfileInsecureTls && !useSharedInsecureDispatcher && !useServerInsecureTls + ? createUpstreamDispatcher(buildUpstreamAgentTimeouts(), true) : undefined; - const insecureTls = useSharedInsecureDispatcher || useProfileInsecureTls; + const insecureTls = + useSharedInsecureDispatcher || useServerInsecureTls || useProfileInsecureTls; + const ownedUpstreamDispatcher = + sharedUpstreamDispatcher || insecureTls + ? undefined + : createUpstreamDispatcher(buildUpstreamAgentTimeouts()); const dispatcher = useSharedInsecureDispatcher ? insecureDispatcher : useProfileInsecureTls ? ephemeralInsecureDispatcher - : getDefaultUpstreamDispatcher(); + : (sharedUpstreamDispatcher ?? ownedUpstreamDispatcher); try { logger.stage('dispatch', 'upstream.dispatch', 'Dispatching upstream fetch', { @@ -588,17 +588,20 @@ export async function handleProxyMessagesRequest( const upstreamUrl = resolveOpenAIChatCompletionsUrl(upstream.route.profile.baseUrl, { passthrough, }); - const upstreamResponse = await fetch( - upstreamUrl, - buildFetchInit( - upstream.route.profile, - upstream.body, - controller.signal, - req.headers, - passthrough, - dispatcher - ) + const fetchInit = buildFetchInit( + upstream.route.profile, + upstream.body, + controller.signal, + req.headers, + passthrough, + dispatcher ); + const upstreamResponse = upstreamFetch + ? await upstreamFetch(upstreamUrl, fetchInit) + : await fetchWithUpstreamTransport(upstreamUrl, fetchInit, { + dispatcher, + insecureTls, + }); logger.stage('upstream', 'upstream.response', 'Received upstream response', { profileName: profile.profileName, routedProfileName: upstream.route.profile.profileName, @@ -620,7 +623,7 @@ export async function handleProxyMessagesRequest( cleanupDisconnectHandlers(); if (ephemeralInsecureDispatcher) { try { - await ephemeralInsecureDispatcher.close(); + await closeUpstreamDispatcher(ephemeralInsecureDispatcher); } catch (closeError) { logger.stage( 'cleanup', @@ -635,13 +638,27 @@ export async function handleProxyMessagesRequest( ); } } + if (ownedUpstreamDispatcher) { + try { + await closeUpstreamDispatcher(ownedUpstreamDispatcher); + } catch (closeError) { + logger.stage( + 'cleanup', + 'request.dispatcher_close_failed', + 'Failed to close per-request upstream dispatcher', + { + profileName: profile.profileName, + routedProfileName: upstream.route.profile.profileName, + error: closeError instanceof Error ? closeError.message : String(closeError), + }, + { level: 'warn' } + ); + } + } } } catch (error) { const message = error instanceof Error ? error.message : 'Unknown proxy error'; - const errInfo = { - name: error instanceof Error ? error.name : 'Error', - message, - }; + const errInfo = toLogErrorInfo(error); logger.stage( 'cleanup', 'request.failed', diff --git a/src/proxy/server/proxy-server.ts b/src/proxy/server/proxy-server.ts index f9af6885..84dd4b2a 100644 --- a/src/proxy/server/proxy-server.ts +++ b/src/proxy/server/proxy-server.ts @@ -1,6 +1,5 @@ import * as http from 'http'; import { randomUUID } from 'crypto'; -import { Agent } from 'undici'; import type { OpenAICompatProfileConfig } from '../profile-router'; import { OPENAI_COMPAT_PROXY_SERVICE_NAME } from '../proxy-daemon-paths'; import { createLogger, withRequestContext } from '../../services/logging'; @@ -11,10 +10,19 @@ import { validateIncomingProxyAuth, } from './messages-route'; import { writeJson } from './http-helpers'; +import { closeUpstreamDispatcher, createUpstreamDispatcher } from './upstream-transport'; const REQUEST_ID_HEADER = 'x-ccs-request-id'; // Loose UUID-ish guard: accepts UUIDs and similar opaque ids; rejects empty / control chars. const REQUEST_ID_PATTERN = /^[A-Za-z0-9._-]{8,128}$/; +const DEFAULT_SHUTDOWN_FORCE_MS = 1_000; + +type ProxyServerLifecycle = { + cleanup: () => Promise; + shutdownPromise?: Promise; +}; + +const proxyServerLifecycles = new WeakMap(); function resolveInboundRequestId(headers: http.IncomingHttpHeaders): string { const raw = headers[REQUEST_ID_HEADER]; @@ -40,8 +48,17 @@ export function startOpenAICompatProxyServer(options: OpenAICompatProxyServerOpt port: options.port, }); const insecureDispatcher = options.insecure - ? new Agent({ connect: { rejectUnauthorized: false }, ...buildUpstreamAgentTimeouts() }) + ? createUpstreamDispatcher(buildUpstreamAgentTimeouts(), true) : undefined; + const upstreamDispatcher = createUpstreamDispatcher(buildUpstreamAgentTimeouts()); + let cleanupPromise: Promise | undefined; + const cleanup = () => { + cleanupPromise ??= Promise.all([ + closeUpstreamDispatcher(upstreamDispatcher), + closeUpstreamDispatcher(insecureDispatcher), + ]).then(() => undefined); + return cleanupPromise; + }; const server = http.createServer((req, res) => { const requestId = resolveInboundRequestId(req.headers); res.setHeader(REQUEST_ID_HEADER, requestId); @@ -121,7 +138,10 @@ export function startOpenAICompatProxyServer(options: OpenAICompatProxyServerOpt res, options.profile, options.authToken, - insecureDispatcher + insecureDispatcher, + upstreamDispatcher, + undefined, + options.insecure === true ); return; } @@ -138,9 +158,55 @@ export function startOpenAICompatProxyServer(options: OpenAICompatProxyServerOpt }); server.on('close', () => { logger.info('server.stop', 'OpenAI-compatible proxy server stopped'); - void insecureDispatcher?.close(); + void cleanup().catch((error) => { + logger.warn('server.dispatcher_close_failed', 'Failed to close upstream transport', { + error: error instanceof Error ? error.message : String(error), + }); + }); }); + proxyServerLifecycles.set(server, { cleanup }); server.listen(options.port, host); return server; } + +export function closeOpenAICompatProxyServer( + server: http.Server, + forceAfterMs = DEFAULT_SHUTDOWN_FORCE_MS +): Promise { + const lifecycle = proxyServerLifecycles.get(server); + if (lifecycle?.shutdownPromise) { + return lifecycle.shutdownPromise; + } + + const shutdownPromise = (async () => { + if (server.listening) { + await new Promise((resolve, reject) => { + let forceTimer: NodeJS.Timeout | undefined; + server.close((error) => { + if (forceTimer) { + clearTimeout(forceTimer); + } + if (error) { + reject(error); + } else { + resolve(); + } + }); + + // Node 18 does not reap idle keep-alive sockets as part of close(). + server.closeIdleConnections?.(); + if (forceAfterMs > 0) { + forceTimer = setTimeout(() => server.closeAllConnections?.(), forceAfterMs); + forceTimer.unref(); + } + }); + } + await lifecycle?.cleanup(); + })(); + + if (lifecycle) { + lifecycle.shutdownPromise = shutdownPromise; + } + return shutdownPromise; +} diff --git a/src/proxy/server/upstream-transport.ts b/src/proxy/server/upstream-transport.ts new file mode 100644 index 00000000..440decf2 --- /dev/null +++ b/src/proxy/server/upstream-transport.ts @@ -0,0 +1,93 @@ +import { Agent, fetch as undiciFetch, type Dispatcher } from 'undici'; +import type { LogErrorInfo } from '../../services/logging'; +import { + createGlobalFetchProxyDispatcher, + type UpstreamAgentTimeoutOptions, +} from '../../utils/fetch-proxy-setup'; +import { getProxyResolution, shouldBypassProxy } from '../../utils/proxy-env'; + +type BunRequestInit = RequestInit & { + proxy?: string; + tls?: { rejectUnauthorized: boolean }; +}; + +const dispatcherClosePromises = new WeakMap>(); + +export function isBunRuntime(): boolean { + return typeof process.versions.bun === 'string'; +} + +export function createUpstreamDispatcher( + options: UpstreamAgentTimeoutOptions, + insecure = false +): Dispatcher | undefined { + if (isBunRuntime()) { + return undefined; + } + if (insecure) { + return new Agent({ connect: { rejectUnauthorized: false }, ...options }); + } + return createGlobalFetchProxyDispatcher(options) ?? new Agent(options); +} + +export function closeUpstreamDispatcher(dispatcher?: Dispatcher): Promise { + if (!dispatcher) { + return Promise.resolve(); + } + const existing = dispatcherClosePromises.get(dispatcher); + if (existing) { + return existing; + } + + const close = (dispatcher as unknown as { close?: () => Promise }).close; + const closing = + typeof close === 'function' ? Promise.resolve(close.call(dispatcher)) : Promise.resolve(); + dispatcherClosePromises.set(dispatcher, closing); + return closing; +} + +export async function fetchWithUpstreamTransport( + input: string | URL | Request, + init: RequestInit, + options: { dispatcher?: Dispatcher; insecureTls?: boolean } = {} +): Promise { + if (!isBunRuntime()) { + return undiciFetch(input, init as Parameters[1]) as Promise; + } + + const requestUrl = new URL( + typeof input === 'string' || input instanceof URL ? input.toString() : input.url + ); + const bunInit: BunRequestInit = { ...init }; + delete (bunInit as Record).dispatcher; + + if (!shouldBypassProxy(requestUrl.hostname)) { + const proxy = getProxyResolution(requestUrl.protocol === 'https:', process.env, { + allowedProtocols: ['http:', 'https:'], + }); + if (proxy.url) { + bunInit.proxy = proxy.url; + } + } + if (options.insecureTls) { + bunInit.tls = { rejectUnauthorized: false }; + } + + return globalThis.fetch(input, bunInit); +} + +export function toLogErrorInfo(error: unknown, depth = 0): LogErrorInfo { + if (!(error instanceof Error)) { + return { name: 'Error', message: String(error) }; + } + + const errorWithDetails = error as Error & { code?: unknown; cause?: unknown }; + const info: LogErrorInfo = { name: error.name, message: error.message }; + if (typeof errorWithDetails.code === 'string') { + info.code = errorWithDetails.code; + } + if (depth < 3 && errorWithDetails.cause !== undefined && errorWithDetails.cause !== error) { + info.cause = toLogErrorInfo(errorWithDetails.cause, depth + 1); + } + return info; +} diff --git a/src/services/logging/log-redaction.ts b/src/services/logging/log-redaction.ts index 8218e4e4..b268e113 100644 --- a/src/services/logging/log-redaction.ts +++ b/src/services/logging/log-redaction.ts @@ -50,6 +50,8 @@ function maskAuthSchemeValue(value: string): string { const SECRET_TOKEN_PATTERN = /(?:Bearer|Basic|Token)\s+\S{8,}|sk-ant-[A-Za-z0-9_-]{16,}|sk-[A-Za-z0-9_-]{32,}|xox[bpoa]-[A-Za-z0-9-]{10,}|gh[opsu]_[A-Za-z0-9]{36,}|glpat-[A-Za-z0-9_-]{18,}|AIza[0-9A-Za-z_-]{35}|eyJ[A-Za-z0-9_-]{8,}\.eyJ[A-Za-z0-9_-]{8,}\.[A-Za-z0-9_-]*|(?:api[_-]?key|access[_-]?token|refresh[_-]?token|secret)(?:=|%3D)[A-Za-z0-9._~+/=-]{8,}/g; +const URL_USERINFO_PATTERN = /([a-z][a-z\d+.-]*:\/\/)([^\s/@]+)@/gi; + /** * Scrub known credential shapes from an arbitrary string. Preserves the * Bearer/Basic/Token scheme prefix when present so the entry stays readable. @@ -57,10 +59,12 @@ const SECRET_TOKEN_PATTERN = * (defense-in-depth for the hotpath console.error sweep). */ export function maskSecretTokens(value: string): string { - return value.replace(SECRET_TOKEN_PATTERN, (match) => { - const scheme = /^(Bearer|Basic|Token)\s+/.exec(match); - return scheme ? `${scheme[1]} [redacted]` : '[redacted]'; - }); + return value + .replace(URL_USERINFO_PATTERN, '$1[redacted]@') + .replace(SECRET_TOKEN_PATTERN, (match) => { + const scheme = /^(Bearer|Basic|Token)\s+/.exec(match); + return scheme ? `${scheme[1]} [redacted]` : '[redacted]'; + }); } function sanitizeValue(value: unknown, depth: number): unknown { diff --git a/src/services/logging/log-types.ts b/src/services/logging/log-types.ts index 10f29d9d..c48892e1 100644 --- a/src/services/logging/log-types.ts +++ b/src/services/logging/log-types.ts @@ -41,6 +41,7 @@ export interface LogErrorInfo { message: string; code?: string; stack?: string; + cause?: LogErrorInfo; } export interface LogEntry { diff --git a/src/utils/fetch-proxy-setup.ts b/src/utils/fetch-proxy-setup.ts index d71b4039..894da2eb 100644 --- a/src/utils/fetch-proxy-setup.ts +++ b/src/utils/fetch-proxy-setup.ts @@ -47,9 +47,12 @@ class RoutingProxyDispatcher extends Dispatcher { close(): Promise; close(callback: () => void): void; close(callback?: () => void): Promise | void { - const promise = Promise.all(this.getDispatchers().map((dispatcher) => dispatcher.close())).then( - () => undefined - ); + const promise = Promise.all( + this.getDispatchers().map((dispatcher) => { + const close = (dispatcher as unknown as { close?: () => Promise }).close; + return typeof close === 'function' ? close.call(dispatcher) : Promise.resolve(); + }) + ).then(() => undefined); if (callback) { promise.then( @@ -73,7 +76,14 @@ class RoutingProxyDispatcher extends Dispatcher { const error = typeof errOrCallback === 'function' ? null : errOrCallback; const done = typeof errOrCallback === 'function' ? errOrCallback : callback; const promise = Promise.all( - this.getDispatchers().map((dispatcher) => dispatcher.destroy(error ?? null)) + this.getDispatchers().map((dispatcher) => { + const destroy = ( + dispatcher as unknown as { destroy?: (error: Error | null) => Promise } + ).destroy; + return typeof destroy === 'function' + ? destroy.call(dispatcher, error ?? null) + : Promise.resolve(); + }) ).then(() => undefined); if (done) { @@ -158,11 +168,6 @@ export function applyGlobalFetchProxy(): { enabled: boolean; error?: string } { } } -const setupResult = applyGlobalFetchProxy(); -if (setupResult.error) { - console.error(`[!] Skipping global fetch proxy setup: ${setupResult.error}`); -} - function resolveGlobalFetchProxyConfig(): GlobalFetchProxyConfig { const httpProxy = getProxyResolution(false, process.env, { allowedProtocols: FETCH_PROXY_PROTOCOLS, diff --git a/tests/integration/proxy/fixtures/runtime-probe.cjs b/tests/integration/proxy/fixtures/runtime-probe.cjs new file mode 100644 index 00000000..a07dd929 --- /dev/null +++ b/tests/integration/proxy/fixtures/runtime-probe.cjs @@ -0,0 +1,100 @@ +const http = require('node:http'); +const path = require('node:path'); + +const repoRoot = process.env.CCS_RUNTIME_PROBE_ROOT || process.cwd(); +const { + closeOpenAICompatProxyServer, + startOpenAICompatProxyServer, +} = require(path.join(repoRoot, 'dist/proxy/server/proxy-server.js')); + +function waitForListening(server) { + if (server.listening) return Promise.resolve(); + return new Promise((resolve, reject) => { + server.once('error', reject); + server.once('listening', resolve); + }); +} + +function request(port, requestPath, options = {}) { + return new Promise((resolve, reject) => { + const req = http.request( + { + hostname: '127.0.0.1', + port, + path: requestPath, + method: options.method || 'GET', + headers: options.headers, + agent: options.agent, + }, + (res) => { + let body = ''; + res.setEncoding('utf8'); + res.on('data', (chunk) => (body += chunk)); + res.on('end', () => resolve({ status: res.statusCode || 0, body })); + } + ); + req.once('error', reject); + req.end(options.body); + }); +} + +async function startProbeServer() { + const server = startOpenAICompatProxyServer({ + profile: { + profileName: 'sentinel', + settingsPath: '/tmp/sentinel.settings.json', + baseUrl: process.env.CCS_RUNTIME_PROBE_UPSTREAM, + apiKey: 'sentinel-api-key-not-real', + provider: 'generic-chat-completion-api', + model: 'sentinel-model', + passthrough: false, + }, + port: 0, + authToken: 'sentinel-local-token', + }); + await waitForListening(server); + return server; +} + +async function runTransportProbe() { + const server = await startProbeServer(); + const port = server.address().port; + const response = await request(port, '/v1/messages', { + method: 'POST', + headers: { + 'content-type': 'application/json', + 'x-api-key': 'sentinel-local-token', + }, + body: JSON.stringify({ + model: 'sentinel-model', + messages: [{ role: 'user', content: 'hi' }], + }), + }); + await closeOpenAICompatProxyServer(server, 250); + const parsed = JSON.parse(response.body); + console.log( + JSON.stringify({ + status: response.status, + text: parsed.content?.[0]?.text, + error: parsed.error?.message, + }) + ); +} + +async function runShutdownProbe() { + const server = await startProbeServer(); + const port = server.address().port; + const agent = new http.Agent({ keepAlive: true }); + await request(port, '/health', { agent }); + const startedAt = Date.now(); + await closeOpenAICompatProxyServer(server, 250); + const elapsedMs = Date.now() - startedAt; + agent.destroy(); + console.log(JSON.stringify({ closed: !server.listening, elapsedMs })); +} + +const mode = process.env.CCS_RUNTIME_PROBE_MODE || 'transport'; +Promise.resolve(mode === 'shutdown' ? runShutdownProbe() : runTransportProbe()).catch((error) => { + console.error(JSON.stringify({ name: error.name, message: error.message })); + process.exitCode = 1; +}); diff --git a/tests/integration/proxy/runtime-transport-matrix.test.ts b/tests/integration/proxy/runtime-transport-matrix.test.ts new file mode 100644 index 00000000..3d0a2773 --- /dev/null +++ b/tests/integration/proxy/runtime-transport-matrix.test.ts @@ -0,0 +1,200 @@ +import { afterAll, beforeAll, describe, expect, it } from 'bun:test'; +import { spawn, spawnSync } from 'node:child_process'; +import * as fs from 'node:fs'; +import * as http from 'node:http'; +import * as os from 'node:os'; +import * as path from 'node:path'; + +const matrixEnabled = process.env.CCS_RUNTIME_MATRIX === '1'; +const fixturePath = path.join(process.cwd(), 'tests/integration/proxy/fixtures/runtime-probe.cjs'); +const proxyKeys = [ + 'HTTP_PROXY', + 'HTTPS_PROXY', + 'ALL_PROXY', + 'NO_PROXY', + 'http_proxy', + 'https_proxy', + 'all_proxy', + 'no_proxy', +]; + +type ProbeResult = { + status?: number; + text?: string; + error?: string; + closed?: boolean; + elapsedMs?: number; +}; + +function resolveNodeBinary(version: 18 | 22 | 26): string { + const result = spawnSync('npx', ['-y', `node@${version}`, '-p', 'process.execPath'], { + encoding: 'utf8', + env: Object.fromEntries( + Object.entries(process.env).filter(([key]) => !proxyKeys.includes(key)) + ), + }); + if (result.status !== 0) { + throw new Error(`Unable to resolve Node ${version}: ${result.stderr.trim()}`); + } + return result.stdout.trim(); +} + +function runProbe( + binary: string, + upstream: string, + env: Record = {}, + mode = 'transport' +): Promise { + const cleanEnv = Object.fromEntries( + Object.entries(process.env).filter(([key]) => !proxyKeys.includes(key)) + ); + const ccsHome = fs.mkdtempSync(path.join(os.tmpdir(), 'ccs-runtime-matrix-')); + return new Promise((resolve, reject) => { + const child = spawn(binary, [fixturePath], { + cwd: process.cwd(), + env: { + ...cleanEnv, + ...env, + CCS_RUNTIME_PROBE_ROOT: process.cwd(), + CCS_RUNTIME_PROBE_UPSTREAM: upstream, + CCS_RUNTIME_PROBE_MODE: mode, + CCS_HOME: ccsHome, + }, + stdio: ['ignore', 'pipe', 'pipe'], + }); + let stdout = ''; + let stderr = ''; + child.stdout.on('data', (chunk) => (stdout += chunk)); + child.stderr.on('data', (chunk) => (stderr += chunk)); + child.once('error', (error) => { + fs.rmSync(ccsHome, { recursive: true, force: true }); + reject(error); + }); + child.once('close', (code) => { + fs.rmSync(ccsHome, { recursive: true, force: true }); + if (code !== 0) { + reject(new Error(`Probe exited ${code}: ${stderr || stdout}`)); + return; + } + const lastLine = stdout.trim().split('\n').pop(); + if (!lastLine) { + reject(new Error(`Probe returned no output: ${stderr}`)); + return; + } + resolve(JSON.parse(lastLine) as ProbeResult); + }); + }); +} + +describe.skipIf(!matrixEnabled)('real runtime upstream transport matrix', () => { + let node18 = ''; + let node22 = ''; + let node26 = ''; + let upstreamServer: http.Server | undefined; + let upstreamPort = 0; + let corporateProxy: http.Server | undefined; + let corporateProxyPort = 0; + let corporateProxyHits = 0; + let sawProxyAuthorization = false; + + beforeAll(async () => { + node18 = resolveNodeBinary(18); + node22 = resolveNodeBinary(22); + node26 = resolveNodeBinary(26); + const payload = JSON.stringify({ + id: 'chatcmpl_matrix', + model: 'sentinel-model', + choices: [ + { + index: 0, + message: { role: 'assistant', content: 'matrix-ok' }, + finish_reason: 'stop', + }, + ], + usage: { prompt_tokens: 1, completion_tokens: 1 }, + }); + upstreamServer = http.createServer((_req, res) => { + res.writeHead(200, { 'content-type': 'application/json' }); + res.end(payload); + }); + corporateProxy = http.createServer((req, res) => { + corporateProxyHits += 1; + sawProxyAuthorization = String(req.headers['proxy-authorization'] || '').startsWith('Basic '); + res.writeHead(200, { + 'content-type': 'application/json', + 'content-length': Buffer.byteLength(payload), + }); + res.end(payload); + }); + corporateProxy.on('connect', (req, socket) => { + corporateProxyHits += 1; + sawProxyAuthorization = String(req.headers['proxy-authorization'] || '').startsWith('Basic '); + socket.write('HTTP/1.1 200 Connection Established\r\n\r\n'); + socket.once('data', () => { + socket.end( + `HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: ${Buffer.byteLength(payload)}\r\nConnection: close\r\n\r\n${payload}` + ); + }); + }); + await Promise.all([ + new Promise((resolve) => upstreamServer.listen(0, '127.0.0.1', resolve)), + new Promise((resolve) => corporateProxy.listen(0, '127.0.0.1', resolve)), + ]); + upstreamPort = (upstreamServer.address() as { port: number }).port; + corporateProxyPort = (corporateProxy.address() as { port: number }).port; + }, 60_000); + + afterAll(async () => { + const servers = [upstreamServer, corporateProxy].filter( + (server): server is http.Server => server !== undefined + ); + await Promise.all( + servers.map((server) => new Promise((resolve) => server.close(() => resolve()))) + ); + }, 10_000); + + it('uses the matching fetch transport on Node 18, 22, 26, and Bun', async () => { + const upstream = `http://127.0.0.1:${upstreamPort}/v1`; + const results = await Promise.all([ + runProbe(node18, upstream), + runProbe(node22, upstream), + runProbe(node26, upstream), + runProbe(process.execPath, upstream), + ]); + for (const result of results) { + expect(result).toMatchObject({ status: 200, text: 'matrix-ok' }); + } + }); + + it('routes Node 26 and Bun through a real credentialed proxy', async () => { + const proxyUrl = `http://sentinel-user:sentinel-password@127.0.0.1:${corporateProxyPort}`; + const nodeResult = await runProbe(node26, 'http://upstream.invalid/v1', { + HTTP_PROXY: proxyUrl, + NO_PROXY: '127.0.0.1,localhost', + }); + const bunResult = await runProbe(process.execPath, 'http://upstream.invalid/v1', { + HTTP_PROXY: proxyUrl, + NO_PROXY: '127.0.0.1,localhost', + }); + expect(nodeResult).toMatchObject({ status: 200, text: 'matrix-ok' }); + expect(bunResult).toMatchObject({ status: 200, text: 'matrix-ok' }); + expect(corporateProxyHits).toBeGreaterThanOrEqual(2); + expect(sawProxyAuthorization).toBe(true); + }); + + it('honors NO_PROXY without touching the proxy', async () => { + const hitsBefore = corporateProxyHits; + const result = await runProbe(node26, 'http://upstream.invalid/v1', { + HTTP_PROXY: `http://sentinel-user:sentinel-password@127.0.0.1:${corporateProxyPort}`, + NO_PROXY: 'upstream.invalid', + }); + expect(result.status).toBe(502); + expect(corporateProxyHits).toBe(hitsBefore); + }); + + it('shuts down Node 18 with an idle keep-alive connection', async () => { + const result = await runProbe(node18, `http://127.0.0.1:${upstreamPort}/v1`, {}, 'shutdown'); + expect(result.closed).toBe(true); + expect(result.elapsedMs).toBeLessThan(2_000); + }); +}); diff --git a/tests/integration/proxy/runtime-transport.test.ts b/tests/integration/proxy/runtime-transport.test.ts new file mode 100644 index 00000000..f3fad034 --- /dev/null +++ b/tests/integration/proxy/runtime-transport.test.ts @@ -0,0 +1,195 @@ +import { afterEach, beforeEach, describe, expect, it } from 'bun:test'; +import * as fs from 'node:fs'; +import * as http from 'node:http'; +import * as os from 'node:os'; +import * as path from 'node:path'; +import { getGlobalDispatcher } from 'undici'; +import { + closeOpenAICompatProxyServer, + startOpenAICompatProxyServer, +} from '../../../src/proxy/server/proxy-server'; +import type { OpenAICompatProfileConfig } from '../../../src/proxy/profile-router'; + +const PROXY_ENV_KEYS = ['HTTP_PROXY', 'HTTPS_PROXY', 'ALL_PROXY', 'NO_PROXY'] as const; + +function listen(server: http.Server): Promise { + return new Promise((resolve, reject) => { + server.once('error', reject); + server.listen(0, '127.0.0.1', () => { + const address = server.address(); + if (!address || typeof address === 'string') { + reject(new Error('Failed to resolve listening port')); + return; + } + resolve(address.port); + }); + }); +} + +function requestMessages(port: number): Promise<{ status: number; body: string }> { + return new Promise((resolve, reject) => { + const req = http.request( + { + hostname: '127.0.0.1', + port, + path: '/v1/messages', + method: 'POST', + headers: { + 'content-type': 'application/json', + 'x-api-key': 'sentinel-local-token', + }, + }, + (res) => { + let body = ''; + res.setEncoding('utf8'); + res.on('data', (chunk) => (body += chunk)); + res.on('end', () => resolve({ status: res.statusCode ?? 0, body })); + } + ); + req.once('error', reject); + req.end( + JSON.stringify({ + model: 'sentinel-model', + messages: [{ role: 'user', content: 'hi' }], + }) + ); + }); +} + +describe('runtime-aware upstream transport', () => { + const originalEnv = new Map(); + let ccsHome = ''; + let corporateProxy: http.Server | undefined; + let corporateProxyHits = 0; + let sawProxyAuthorization = false; + let proxyServer: http.Server | undefined; + + beforeEach(() => { + for (const key of PROXY_ENV_KEYS) { + originalEnv.set(key, process.env[key]); + delete process.env[key]; + } + originalEnv.set('CCS_HOME', process.env.CCS_HOME); + ccsHome = fs.mkdtempSync(path.join(os.tmpdir(), 'ccs-runtime-transport-')); + process.env.CCS_HOME = ccsHome; + corporateProxy = undefined; + proxyServer = undefined; + corporateProxyHits = 0; + sawProxyAuthorization = false; + }); + + afterEach(async () => { + if (proxyServer) { + await closeOpenAICompatProxyServer(proxyServer, 50); + } + if (corporateProxy) { + await new Promise((resolve) => corporateProxy.close(() => resolve())); + } + for (const key of [...PROXY_ENV_KEYS, 'CCS_HOME'] as const) { + const value = originalEnv.get(key); + if (value === undefined) delete process.env[key]; + else process.env[key] = value; + } + fs.rmSync(ccsHome, { recursive: true, force: true }); + }); + + it('uses a credentialed proxy and honors NO_PROXY without global fetch mutation', async () => { + const originalFetch = globalThis.fetch; + const upstreamPayload = JSON.stringify({ + id: 'chatcmpl_proxy', + model: 'sentinel-model', + choices: [ + { + index: 0, + message: { role: 'assistant', content: 'proxied-ok' }, + finish_reason: 'stop', + }, + ], + usage: { prompt_tokens: 1, completion_tokens: 1 }, + }); + corporateProxy = http.createServer((req, res) => { + corporateProxyHits += 1; + sawProxyAuthorization = String(req.headers['proxy-authorization'] || '').startsWith('Basic '); + res.writeHead(200, { + 'content-type': 'application/json', + 'content-length': Buffer.byteLength(upstreamPayload), + }); + res.end(upstreamPayload); + }); + const corporateProxyPort = await listen(corporateProxy); + process.env.HTTP_PROXY = `http://sentinel-user:sentinel-password@127.0.0.1:${corporateProxyPort}`; + + const profile: OpenAICompatProfileConfig = { + profileName: 'sentinel', + settingsPath: '/tmp/sentinel.settings.json', + baseUrl: 'http://upstream.invalid/v1', + apiKey: 'sentinel-api-key-not-real', + provider: 'generic-chat-completion-api', + model: 'sentinel-model', + }; + proxyServer = startOpenAICompatProxyServer({ + profile, + port: 0, + authToken: 'sentinel-local-token', + }); + const proxyPort = await new Promise((resolve, reject) => { + proxyServer.once('error', reject); + proxyServer.once('listening', () => { + const address = proxyServer.address(); + if (!address || typeof address === 'string') reject(new Error('Missing proxy port')); + else resolve(address.port); + }); + }); + + const proxied = await requestMessages(proxyPort); + expect(proxied.status).toBe(200); + expect(JSON.parse(proxied.body)).toMatchObject({ + content: [{ type: 'text', text: 'proxied-ok' }], + }); + expect(corporateProxyHits).toBe(1); + expect(sawProxyAuthorization).toBe(true); + + process.env.NO_PROXY = 'upstream.invalid'; + const hitsBeforeBypass = corporateProxyHits; + const bypassed = await requestMessages(proxyPort); + expect(bypassed.status).toBe(502); + expect(corporateProxyHits).toBe(hitsBeforeBypass); + expect(globalThis.fetch).toBe(originalFetch); + }); + + it('keeps global transport ownership unchanged across overlapping server lifecycles', async () => { + const originalFetch = globalThis.fetch; + const originalDispatcher = getGlobalDispatcher(); + const profile: OpenAICompatProfileConfig = { + profileName: 'sentinel', + settingsPath: '/tmp/sentinel.settings.json', + baseUrl: 'http://127.0.0.1:1/v1', + apiKey: 'sentinel-api-key-not-real', + provider: 'generic-chat-completion-api', + model: 'sentinel-model', + }; + const serverA = startOpenAICompatProxyServer({ + profile, + port: 0, + authToken: 'sentinel-local-token-a', + }); + const serverB = startOpenAICompatProxyServer({ + profile, + port: 0, + authToken: 'sentinel-local-token-b', + }); + + await Promise.all([ + new Promise((resolve) => serverA.once('listening', resolve)), + new Promise((resolve) => serverB.once('listening', resolve)), + ]); + await closeOpenAICompatProxyServer(serverA, 50); + expect(serverB.listening).toBe(true); + expect(globalThis.fetch).toBe(originalFetch); + expect(getGlobalDispatcher()).toBe(originalDispatcher); + + await closeOpenAICompatProxyServer(serverB, 50); + expect(globalThis.fetch).toBe(originalFetch); + expect(getGlobalDispatcher()).toBe(originalDispatcher); + }); +}); diff --git a/tests/unit/proxy/messages-route.test.ts b/tests/unit/proxy/messages-route.test.ts index eb66ee5a..9cca4a5a 100644 --- a/tests/unit/proxy/messages-route.test.ts +++ b/tests/unit/proxy/messages-route.test.ts @@ -9,8 +9,13 @@ import { resolveOpenAICompatProfileConfig } from '../../../src/proxy/profile-rou import { attachDisconnectAbortHandlers, buildUpstreamAgentTimeouts, - handleProxyMessagesRequest, + handleProxyMessagesRequest as handleProxyMessagesRequestImpl, } from '../../../src/proxy/server/messages-route'; +import { + closeUpstreamDispatcher, + isBunRuntime, + toLogErrorInfo, +} from '../../../src/proxy/server/upstream-transport'; import { loadSettings } from '../../../src/utils/config-manager'; class FakeSocket extends EventEmitter { @@ -70,6 +75,24 @@ function buildProfile(profileName: string) { return profile!; } +function handleProxyMessagesRequest( + req: Parameters[0], + res: Parameters[1], + profile: Parameters[2], + expectedAuthToken: string, + insecureDispatcher?: Parameters[4] +) { + return handleProxyMessagesRequestImpl( + req, + res, + profile, + expectedAuthToken, + insecureDispatcher, + undefined, + globalThis.fetch + ); +} + beforeEach(() => { originalCcsHome = process.env.CCS_HOME; originalFetch = globalThis.fetch; @@ -191,6 +214,45 @@ describe('buildUpstreamAgentTimeouts', () => { }); }); +describe('closeUpstreamDispatcher', () => { + it('closes an owned dispatcher once across repeated cleanup', async () => { + let closeCalls = 0; + const dispatcher = { + close: async () => { + closeCalls += 1; + }, + } as never; + + await Promise.all([ + closeUpstreamDispatcher(dispatcher), + closeUpstreamDispatcher(dispatcher), + closeUpstreamDispatcher(dispatcher), + ]); + + expect(closeCalls).toBe(1); + }); +}); + +describe('toLogErrorInfo', () => { + it('preserves nested transport cause metadata for diagnostics', () => { + const cause = Object.assign(new Error('invalid onError method'), { + name: 'InvalidArgumentError', + code: 'UND_ERR_INVALID_ARG', + }); + const error = new TypeError('fetch failed', { cause }); + + expect(toLogErrorInfo(error)).toEqual({ + name: 'TypeError', + message: 'fetch failed', + cause: { + name: 'InvalidArgumentError', + message: 'invalid onError method', + code: 'UND_ERR_INVALID_ARG', + }, + }); + }); +}); + describe('handleProxyMessagesRequest', () => { it('dispatches secure upstream fetches with an explicit dispatcher (not undici defaults)', async () => { const activeProfile = buildProfile('hf'); @@ -217,7 +279,11 @@ describe('handleProxyMessagesRequest', () => { // A bare global-dispatcher fallback would reapply undici's 300s // headersTimeout/bodyTimeout and undercut the proxy's request timeout. - expect(capturedDispatcher).toBeDefined(); + if (isBunRuntime()) { + expect(capturedDispatcher).toBeUndefined(); + } else { + expect(capturedDispatcher).toBeDefined(); + } expect(res.statusCode).toBe(502); }); @@ -403,9 +469,13 @@ describe('handleProxyMessagesRequest', () => { ); await pending; - expect(capturedDispatcher).toBeInstanceOf(Agent); + if (isBunRuntime()) { + expect(capturedDispatcher).toBeUndefined(); + } else { + expect(capturedDispatcher).toBeInstanceOf(Agent); + } expect(capturedDispatcher).not.toBe(sharedDispatcher); - expect(closeCalls).toBe(1); + expect(closeCalls).toBe(isBunRuntime() ? 0 : 1); expect(res.statusCode).toBe(502); }); diff --git a/tests/unit/services/logging/cli-entry-log-redaction.test.ts b/tests/unit/services/logging/cli-entry-log-redaction.test.ts index cf3e34c9..ebd1c11f 100644 --- a/tests/unit/services/logging/cli-entry-log-redaction.test.ts +++ b/tests/unit/services/logging/cli-entry-log-redaction.test.ts @@ -63,7 +63,9 @@ beforeEach(() => { baselineUncaughtExceptionListeners = process.listeners('uncaughtException'); baselineUnhandledRejectionListeners = process.listeners('unhandledRejection'); - mock.module('../../../../src/utils/fetch-proxy-setup', () => ({})); + mock.module('../../../../src/utils/fetch-proxy-setup', () => ({ + applyGlobalFetchProxy: () => ({ enabled: false }), + })); mock.module('../../../../src/utils/error-manager', () => ({ ErrorManager: class ErrorManager { static async showProfileNotFound(): Promise {} diff --git a/tests/unit/services/logging/log-redaction.test.ts b/tests/unit/services/logging/log-redaction.test.ts index 124165eb..eb223abc 100644 --- a/tests/unit/services/logging/log-redaction.test.ts +++ b/tests/unit/services/logging/log-redaction.test.ts @@ -1,5 +1,9 @@ import { describe, expect, it } from 'bun:test'; -import { redactContext } from '../../../../src/services/logging/log-redaction'; +import { + maskSecretTokens, + redactContext, + redactErrorInfo, +} from '../../../../src/services/logging/log-redaction'; describe('log redaction', () => { it('redacts sensitive keys and preserves non-sensitive values', () => { @@ -100,4 +104,31 @@ describe('log redaction', () => { }, }); }); + + it('redacts proxy URL userinfo from nested error causes', () => { + const redacted = redactErrorInfo({ + name: 'TypeError', + message: 'fetch failed', + cause: { + name: 'ProxyError', + code: 'ECONNREFUSED', + message: 'connect http://sentinel-user:sentinel-password@proxy.example:8080', + }, + }); + + expect(redacted?.cause).toEqual({ + name: 'ProxyError', + code: 'ECONNREFUSED', + message: 'connect http://[redacted]@proxy.example:8080', + }); + expect(JSON.stringify(redacted)).not.toContain('sentinel-user'); + expect(JSON.stringify(redacted)).not.toContain('sentinel-password'); + }); + + it('redacts URL userinfo without changing non-credential URLs', () => { + expect(maskSecretTokens('https://user:pass@example.com/path')).toBe( + 'https://[redacted]@example.com/path' + ); + expect(maskSecretTokens('https://example.com/path')).toBe('https://example.com/path'); + }); }); diff --git a/tests/unit/utils/fetch-proxy-setup.test.ts b/tests/unit/utils/fetch-proxy-setup.test.ts index 446d765c..22a02d08 100644 --- a/tests/unit/utils/fetch-proxy-setup.test.ts +++ b/tests/unit/utils/fetch-proxy-setup.test.ts @@ -4,7 +4,7 @@ import { Agent, Dispatcher, ProxyAgent, - fetch, + fetch as undiciFetch, getGlobalDispatcher, setGlobalDispatcher, } from 'undici'; @@ -26,18 +26,24 @@ const PROXY_ENV_KEYS = [ describe('global fetch proxy setup', () => { const originalEnv = new Map(); + const originalFetch = globalThis.fetch; const originalDispatcher = getGlobalDispatcher(); - beforeEach(() => { + globalThis.fetch = originalFetch; + setGlobalDispatcher(originalDispatcher); for (const key of PROXY_ENV_KEYS) { originalEnv.set(key, process.env[key]); delete process.env[key]; } - setGlobalDispatcher(originalDispatcher); }); - afterEach(() => { + afterEach(async () => { + const activeDispatcher = getGlobalDispatcher(); + globalThis.fetch = originalFetch; setGlobalDispatcher(originalDispatcher); + if (activeDispatcher !== originalDispatcher) { + await activeDispatcher.close(); + } for (const key of PROXY_ENV_KEYS) { const value = originalEnv.get(key); if (value === undefined) { @@ -112,18 +118,6 @@ describe('global fetch proxy setup', () => { }); }); - it('rebinds globalThis.fetch to undici fetch when proxying is enabled', () => { - const originalFetch = globalThis.fetch; - process.env.HTTP_PROXY = 'http://proxy.example:8080'; - - try { - expect(applyGlobalFetchProxy()).toEqual({ enabled: true }); - expect(globalThis.fetch).toBe(fetch); - } finally { - globalThis.fetch = originalFetch; - } - }); - it('bypasses loopback fetches even when HTTP_PROXY is set', async () => { process.env.HTTP_PROXY = 'http://127.0.0.1:9'; @@ -139,11 +133,14 @@ describe('global fetch proxy setup', () => { } try { - expect(applyGlobalFetchProxy()).toEqual({ enabled: true }); - - const response = await fetch(`http://127.0.0.1:${address.port}/health`); + const dispatcher = createGlobalFetchProxyDispatcher(); + expect(dispatcher).not.toBeNull(); + const response = await undiciFetch(`http://127.0.0.1:${address.port}/health`, { + dispatcher: dispatcher!, + }); expect(response.status).toBe(200); expect(await response.text()).toBe('ok'); + await dispatcher?.close(); } finally { await new Promise((resolve, reject) => { server.close((error) => (error ? reject(error) : resolve())); @@ -158,6 +155,15 @@ describe('global fetch proxy setup', () => { expect(result).toEqual({ directCalls: 0, proxyCalls: 1 }); }); + it('applies the proxy transport only when explicitly invoked', () => { + process.env.HTTP_PROXY = 'http://proxy.example:8080'; + + expect(globalThis.fetch).toBe(originalFetch); + expect(applyGlobalFetchProxy()).toEqual({ enabled: true }); + expect(globalThis.fetch).toBe(undiciFetch); + expect(getGlobalDispatcher()).not.toBe(originalDispatcher); + }); + it('routes HTTPS requests through HTTPS_PROXY', async () => { process.env.HTTPS_PROXY = 'http://proxy.example:8443';