fix(proxy): stabilize upstream fetch transport

Refs #1686
This commit is contained in:
Tam Nhu Tran committed 2026-08-08 21:17:13 -04:00
1 parent 7fb56d62e1
commit d8210cf011
17 files changed
+889 -85

No files matched your search

+1
View File
@@ -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",
+1
View File
@@ -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
+6 -2
View File
@@ -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<void> {
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());
+11 -3
View File
@@ -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<void> | 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);
}
+55 -38
View File
@@ -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<void> {
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',
+70 -4
View File
@@ -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<void>;
shutdownPromise?: Promise<void>;
};
const proxyServerLifecycles = new WeakMap<http.Server, ProxyServerLifecycle>();
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<void> | 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<void> {
const lifecycle = proxyServerLifecycles.get(server);
if (lifecycle?.shutdownPromise) {
return lifecycle.shutdownPromise;
}
const shutdownPromise = (async () => {
if (server.listening) {
await new Promise<void>((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;
}
+93
View File
@@ -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<object, Promise<void>>();
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<void> {
if (!dispatcher) {
return Promise.resolve();
}
const existing = dispatcherClosePromises.get(dispatcher);
if (existing) {
return existing;
}
const close = (dispatcher as unknown as { close?: () => Promise<void> }).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<Response> {
if (!isBunRuntime()) {
return undiciFetch(input, init as Parameters<typeof undiciFetch>[1]) as Promise<Response>;
}
const requestUrl = new URL(
typeof input === 'string' || input instanceof URL ? input.toString() : input.url
);
const bunInit: BunRequestInit = { ...init };
delete (bunInit as Record<string, unknown>).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;
}
+8 -4
View File
@@ -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 {
+1
View File
@@ -41,6 +41,7 @@ export interface LogErrorInfo {
message: string;
code?: string;
stack?: string;
cause?: LogErrorInfo;
}
export interface LogEntry {
+14 -9
View File
@@ -47,9 +47,12 @@ class RoutingProxyDispatcher extends Dispatcher {
close(): Promise<void>;
close(callback: () => void): void;
close(callback?: () => void): Promise<void> | 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<void> }).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<void> }
).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,
@@ -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;
});
@@ -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<string, string> = {},
mode = 'transport'
): Promise<ProbeResult> {
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<void>((resolve) => upstreamServer.listen(0, '127.0.0.1', resolve)),
new Promise<void>((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<void>((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);
});
});
@@ -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<number> {
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<string, string | undefined>();
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<void>((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<number>((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<void>((resolve) => serverA.once('listening', resolve)),
new Promise<void>((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);
});
});
+74 -4
View File
@@ -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<typeof handleProxyMessagesRequestImpl>[0],
res: Parameters<typeof handleProxyMessagesRequestImpl>[1],
profile: Parameters<typeof handleProxyMessagesRequestImpl>[2],
expectedAuthToken: string,
insecureDispatcher?: Parameters<typeof handleProxyMessagesRequestImpl>[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);
});
@@ -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<void> {}
@@ -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');
});
});
+25 -19
View File
@@ -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<string, string | undefined>();
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<void>((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';