import http, { type IncomingHttpHeaders, type IncomingMessage, type Server, type ServerResponse, } from 'node:http'; import { afterEach, describe, expect, it } from 'vitest'; import { startLlmTransientFaultProxy, withLoopbackNoProxy, } from '../scripts/llm-transient-fault-proxy.mjs'; interface ProxyStats { requestCount: number; faultInjectedCount: number; heldRequestCount: number; forwardedRequestCount: number; forwardingReleased: boolean; stopped: boolean; } interface ProxyHandle { url: string; baseUrl: string; port: number; stats: ProxyStats; getStats(): ProxyStats; getRequestLog(): ReadonlyArray<{ sequence: number; acceptedAtMs: number; faultInjectedAtMs: number | null; heldAtMs: number | null; forwardingStartedAtMs: number | null; }>; waitForFault(timeoutMs?: number): Promise; waitForHeldRequest(timeoutMs?: number): Promise; releaseForwarding(): ProxyStats; stop(): Promise; } interface CapturedRequest { method: string | undefined; url: string | undefined; headers: IncomingHttpHeaders; body: string; } interface HttpResult { statusCode: number | undefined; headers: IncomingHttpHeaders; body: string; } const proxies = new Set(); const servers = new Set(); const TEST_UPSTREAM_FALLBACK_PORTS = Array.from( { length: 128 }, (_, index) => 62_300 + index, ); afterEach(async () => { await Promise.allSettled([...proxies].map((proxy) => proxy.stop())); proxies.clear(); for (const server of servers) server.closeAllConnections?.(); await Promise.allSettled([...servers].map((server) => closeServer(server))); servers.clear(); }); function deferred() { let resolve!: (value: T | PromiseLike) => void; let reject!: (reason?: unknown) => void; const promise = new Promise((resolvePromise, rejectPromise) => { resolve = resolvePromise; reject = rejectPromise; }); return { promise, resolve, reject }; } function listen( server: Server, { host, port }: { host: string; port: number }, ) { return new Promise((resolve, reject) => { const cleanup = () => { server.off('error', onError); server.off('listening', onListening); }; const onError = (error: Error) => { cleanup(); reject(error); }; const onListening = () => { cleanup(); resolve(); }; server.once('error', onError); server.once('listening', onListening); server.listen({ host, port, exclusive: true }); }); } function closeServer(server: Server) { if (!server.listening) return Promise.resolve(); return new Promise((resolve) => server.close(() => resolve())); } function errorCode(error: unknown) { return error && typeof error === 'object' && 'code' in error ? error.code : undefined; } async function listenTestServer(server: Server) { try { await listen(server, { host: '127.0.0.1', port: 0 }); return; } catch (error) { if (errorCode(error) !== 'EADDRINUSE') throw error; } for (const port of TEST_UPSTREAM_FALLBACK_PORTS) { try { await listen(server, { host: '127.0.0.1', port }); return; } catch (error) { if (errorCode(error) !== 'EADDRINUSE') throw error; } } throw new Error('test upstream fallback port pool is exhausted'); } async function startServer( handler: ( request: IncomingMessage, response: ServerResponse, ) => void | Promise, ) { const server = http.createServer((request, response) => { void Promise.resolve(handler(request, response)).catch(() => { response.destroy(); }); }); await listenTestServer(server); servers.add(server); const address = server.address(); if (!address || typeof address === 'string') { throw new Error('test server did not expose a TCP address'); } return { server, url: `http://127.0.0.1:${address.port}` }; } async function startProxy(options: Record) { const proxy = (await startLlmTransientFaultProxy(options)) as ProxyHandle; proxies.add(proxy); return proxy; } function readBody(request: IncomingMessage) { return new Promise((resolve, reject) => { const chunks: Buffer[] = []; request.on('data', (chunk) => chunks.push(Buffer.from(chunk))); request.once('end', () => resolve(Buffer.concat(chunks).toString('utf8'))); request.once('error', reject); }); } function request( baseUrl: string, { method = 'POST', path = '/', headers = {}, body = '', onChunk, }: { method?: string; path?: string; headers?: Record; body?: string; onChunk?: (chunk: string) => void; } = {}, ) { const base = new URL(baseUrl); return new Promise((resolve, reject) => { const outgoing = http.request( { hostname: base.hostname, port: base.port, method, path, headers, agent: false, }, (response) => { const chunks: Buffer[] = []; response.on('data', (chunk) => { const buffer = Buffer.from(chunk); chunks.push(buffer); onChunk?.(buffer.toString('utf8')); }); response.once('end', () => { resolve({ statusCode: response.statusCode, headers: response.headers, body: Buffer.concat(chunks).toString('utf8'), }); }); response.once('aborted', () => reject(new Error('response aborted'))); response.once('error', reject); }, ); outgoing.once('error', reject); outgoing.end(body); }); } describe('LLM transient fault proxy', () => { it('pins loopback in both no-proxy variants without dropping existing entries', () => { const source = { HTTP_PROXY: 'http://proxy.example', NO_PROXY: 'internal.example, 127.0.0.1', no_proxy: 'legacy.example', }; const childEnvironment = withLoopbackNoProxy(source); expect(childEnvironment).toEqual({ HTTP_PROXY: 'http://proxy.example', NO_PROXY: 'internal.example,127.0.0.1,legacy.example,localhost,::1', no_proxy: 'internal.example,127.0.0.1,legacy.example,localhost,::1', }); expect(source).toEqual({ HTTP_PROXY: 'http://proxy.example', NO_PROXY: 'internal.example, 127.0.0.1', no_proxy: 'legacy.example', }); }); it('resets the first POST before upstream and then streams an intact request and response', async () => { const captured: CapturedRequest[] = []; const releaseResponse = deferred(); const firstResponseChunk = deferred(); const upstream = await startServer(async (incoming, response) => { captured.push({ method: incoming.method, url: incoming.url, headers: incoming.headers, body: await readBody(incoming), }); response.writeHead(201, { 'content-type': 'text/plain', 'x-upstream-stream': 'yes', }); response.write('first-'); await releaseResponse.promise; response.end('second'); }); const proxy = await startProxy({ upstreamBaseUrl: upstream.url, faultCount: 1, }); const failedRequest = request(proxy.url, { path: '/v1/chat/completions?mode=fault', headers: { 'content-type': 'application/json' }, body: '{"attempt":1}', }); await expect(failedRequest).rejects.toBeInstanceOf(Error); await expect(proxy.waitForFault()).resolves.toMatchObject({ faultInjectedCount: 1, }); expect(captured).toHaveLength(0); const forwardedRequest = request(proxy.url, { path: '/v1/chat/completions?mode=stream', headers: { authorization: 'Bearer fixture-secret', 'content-type': 'application/json', 'x-request-marker': 'preserved', }, body: '{"prompt":"sensitive-body-marker"}', onChunk: (chunk) => firstResponseChunk.resolve(chunk), }); await expect(firstResponseChunk.promise).resolves.toBe('first-'); expect(captured).toEqual([ expect.objectContaining({ method: 'POST', url: '/v1/chat/completions?mode=stream', headers: expect.objectContaining({ authorization: 'Bearer fixture-secret', 'content-type': 'application/json', 'x-request-marker': 'preserved', }), body: '{"prompt":"sensitive-body-marker"}', }), ]); releaseResponse.resolve(); await expect(forwardedRequest).resolves.toMatchObject({ statusCode: 201, headers: expect.objectContaining({ 'x-upstream-stream': 'yes' }), body: 'first-second', }); expect(proxy.stats).toEqual({ requestCount: 2, faultInjectedCount: 1, heldRequestCount: 0, forwardedRequestCount: 1, forwardingReleased: true, stopped: false, }); }); it('uses only frozen timing metadata to select a later fault', async () => { const capturedBodies: string[] = []; const selectorCalls: Array<{ metadata: { sequence: number; acceptedAtMs: number }; frozen: boolean; }> = []; const upstream = await startServer(async (incoming, response) => { capturedBodies.push(await readBody(incoming)); response.end('forwarded'); }); const proxy = await startProxy({ upstreamBaseUrl: upstream.url, faultCount: 1, async shouldInjectFault(metadata: { sequence: number; acceptedAtMs: number; }) { selectorCalls.push({ metadata: { ...metadata }, frozen: Object.isFrozen(metadata), }); await Promise.resolve(); return metadata.sequence === 2; }, }); await expect( request(proxy.url, { headers: { authorization: 'Bearer a' }, body: 'first-sensitive-body', }), ).resolves.toMatchObject({ statusCode: 200, body: 'forwarded' }); await expect( request(proxy.url, { headers: { authorization: 'Bearer b' }, body: 'selected-sensitive-body', }), ).rejects.toBeInstanceOf(Error); await expect(proxy.waitForFault()).resolves.toMatchObject({ faultInjectedCount: 1, }); await expect( request(proxy.url, { body: 'after-fault-body' }), ).resolves.toMatchObject({ statusCode: 200, body: 'forwarded' }); expect(selectorCalls).toEqual([ { metadata: { sequence: 1, acceptedAtMs: expect.any(Number), }, frozen: true, }, { metadata: { sequence: 2, acceptedAtMs: expect.any(Number), }, frozen: true, }, ]); expect(Object.keys(selectorCalls[0].metadata).sort()).toEqual([ 'acceptedAtMs', 'sequence', ]); expect(JSON.stringify(selectorCalls)).not.toMatch( /secret|sensitive|authorization|body/iu, ); expect(capturedBodies).toEqual([ 'first-sensitive-body', 'after-fault-body', ]); expect(proxy.stats).toEqual({ requestCount: 3, faultInjectedCount: 1, heldRequestCount: 0, forwardedRequestCount: 2, forwardingReleased: true, stopped: false, }); }); it('does not count a delayed selector decision after the proxy stops', async () => { let upstreamRequestCount = 0; const selectorStarted = deferred(); const selectorDecision = deferred(); const upstream = await startServer(async (incoming, response) => { upstreamRequestCount += 1; await readBody(incoming); response.end('unexpected'); }); const proxy = await startProxy({ upstreamBaseUrl: upstream.url, faultCount: 1, async shouldInjectFault() { selectorStarted.resolve(); return selectorDecision.promise; }, }); const pendingRequest = request(proxy.url, { body: 'delayed-selector' }); await selectorStarted.promise; const stopPromise = proxy.stop(); selectorDecision.resolve(true); await expect(pendingRequest).rejects.toBeInstanceOf(Error); await expect(stopPromise).resolves.toBeUndefined(); expect(upstreamRequestCount).toBe(0); expect(proxy.stats).toMatchObject({ requestCount: 1, faultInjectedCount: 0, heldRequestCount: 0, forwardedRequestCount: 0, stopped: true, }); }); it('holds the post-fault request before upstream until forwarding is released', async () => { const captured: CapturedRequest[] = []; const upstream = await startServer(async (incoming, response) => { captured.push({ method: incoming.method, url: incoming.url, headers: incoming.headers, body: await readBody(incoming), }); response.writeHead(200, { 'x-held-request': 'released' }); response.end('forwarded-after-release'); }); const proxy = await startProxy({ upstreamBaseUrl: upstream.url, faultCount: 1, holdAfterFault: true, }); await expect( request(proxy.url, { path: '/fault', body: 'first-attempt' }), ).rejects.toBeInstanceOf(Error); await proxy.waitForFault(); const heldRequest = request(proxy.url, { path: '/v1/responses?held=1', headers: { authorization: 'Bearer held-secret', 'content-type': 'application/json', }, body: '{"held":"request-body"}', }); await expect(proxy.waitForHeldRequest()).resolves.toMatchObject({ heldRequestCount: 1, forwardedRequestCount: 0, }); const heldLog = proxy.getRequestLog(); expect(heldLog).toEqual([ { sequence: 1, acceptedAtMs: expect.any(Number), faultInjectedAtMs: expect.any(Number), heldAtMs: null, forwardingStartedAtMs: null, }, { sequence: 2, acceptedAtMs: expect.any(Number), faultInjectedAtMs: null, heldAtMs: expect.any(Number), forwardingStartedAtMs: null, }, ]); expect(heldLog[1].acceptedAtMs).toBeGreaterThanOrEqual( heldLog[0].acceptedAtMs, ); expect(heldLog[1].heldAtMs).toBeGreaterThanOrEqual(heldLog[1].acceptedAtMs); expect(JSON.stringify(heldLog)).not.toMatch( /held-secret|request-body|responses|authorization/iu, ); expect(captured).toHaveLength(0); expect(proxy.releaseForwarding()).toMatchObject({ forwardingReleased: true, heldRequestCount: 1, }); await expect(heldRequest).resolves.toMatchObject({ statusCode: 200, headers: expect.objectContaining({ 'x-held-request': 'released' }), body: 'forwarded-after-release', }); expect(captured).toEqual([ expect.objectContaining({ method: 'POST', url: '/v1/responses?held=1', headers: expect.objectContaining({ authorization: 'Bearer held-secret', }), body: '{"held":"request-body"}', }), ]); expect(proxy.stats.forwardedRequestCount).toBe(1); expect( proxy.getRequestLog()[1].forwardingStartedAtMs, ).toBeGreaterThanOrEqual(heldLog[1].heldAtMs); }); it('holds before injecting a remaining configured fault', async () => { let upstreamRequestCount = 0; const upstream = await startServer(async (incoming, response) => { upstreamRequestCount += 1; await readBody(incoming); response.end('after-two-faults'); }); const proxy = await startProxy({ upstreamBaseUrl: upstream.url, faultCount: 2, holdAfterFault: true, }); await expect( request(proxy.url, { body: 'fault-one' }), ).rejects.toBeInstanceOf(Error); const secondFault = request(proxy.url, { body: 'fault-two' }); await proxy.waitForHeldRequest(); expect(proxy.stats).toMatchObject({ faultInjectedCount: 1, heldRequestCount: 1, forwardedRequestCount: 0, }); proxy.releaseForwarding(); await expect(secondFault).rejects.toBeInstanceOf(Error); expect(proxy.stats.faultInjectedCount).toBe(2); expect(upstreamRequestCount).toBe(0); await expect( request(proxy.url, { body: 'forward-third' }), ).resolves.toMatchObject({ statusCode: 200, body: 'after-two-faults' }); expect(upstreamRequestCount).toBe(1); }); it('keeps sensitive upstream, header, and body values out of stats and errors', async () => { const sensitiveQuery = 'upstream-sensitive-query-marker'; const upstream = await startServer(async (incoming, response) => { await readBody(incoming); response.end('sensitive-response-marker'); }); const upstreamBaseUrl = `${upstream.url}/v1///?token=${sensitiveQuery}`; const proxy = await startProxy({ upstreamBaseUrl, faultCount: 0, }); await request(proxy.url, { path: '/v1/chat/completions', headers: { authorization: 'Bearer sensitive-authorization-marker' }, body: 'sensitive-request-body-marker', }); let timeoutError: unknown; try { await proxy.waitForHeldRequest(20); } catch (error) { timeoutError = error; } const publicSurface = JSON.stringify({ proxy, stats: proxy.getStats(), error: String(timeoutError), }); for (const sensitiveValue of [ upstream.url, upstreamBaseUrl, sensitiveQuery, 'sensitive-authorization-marker', 'sensitive-request-body-marker', 'sensitive-response-marker', ]) { expect(publicSurface).not.toContain(sensitiveValue); } expect(String(timeoutError)).toBe( 'LlmTransientFaultProxyError: timed out waiting for proxy held', ); }); it('returns a normalized local baseUrl and preserves the upstream base path', async () => { const captured: CapturedRequest[] = []; const upstream = await startServer(async (incoming, response) => { captured.push({ method: incoming.method, url: incoming.url, headers: incoming.headers, body: await readBody(incoming), }); response.end('base-path-preserved'); }); const proxy = await startProxy({ upstreamBaseUrl: `${upstream.url}/v1///`, faultCount: 0, }); expect(proxy.baseUrl).toBe(`${proxy.url}/v1`); await expect( request(proxy.url, { path: '/v10/chat/completions', }), ).resolves.toMatchObject({ statusCode: 400 }); await expect( request(proxy.url, { path: '/v1/../chat/completions', }), ).resolves.toMatchObject({ statusCode: 400 }); expect(captured).toHaveLength(0); const runtimeUrl = `${proxy.baseUrl}/chat/completions`; const runtimePath = new URL(runtimeUrl).pathname; await expect( request(proxy.url, { path: runtimePath, headers: { 'content-type': 'application/json' }, body: '{"basePath":"/v1"}', }), ).resolves.toMatchObject({ statusCode: 200, body: 'base-path-preserved', }); expect(captured).toEqual([ expect.objectContaining({ method: 'POST', url: '/v1/chat/completions', body: '{"basePath":"/v1"}', }), ]); }); it('stops idempotently, unblocks held requests, and leaves no forwarded request', async () => { let upstreamRequestCount = 0; const upstream = await startServer(async (incoming, response) => { upstreamRequestCount += 1; await readBody(incoming); response.end('unexpected'); }); const proxy = await startProxy({ upstreamBaseUrl: upstream.url, holdAfterFault: true, }); await expect(request(proxy.url, { body: 'fault' })).rejects.toBeInstanceOf( Error, ); const heldRequest = request(proxy.url, { body: 'held-until-stop' }); await proxy.waitForHeldRequest(); const firstStop = proxy.stop(); const secondStop = proxy.stop(); expect(secondStop).toBe(firstStop); await Promise.all([firstStop, secondStop]); await expect(proxy.stop()).resolves.toBeUndefined(); await expect(heldRequest).rejects.toBeInstanceOf(Error); expect(upstreamRequestCount).toBe(0); expect(proxy.stats).toMatchObject({ faultInjectedCount: 1, heldRequestCount: 1, forwardedRequestCount: 0, stopped: true, }); }); it('stops an active upstream request and its downstream connection', async () => { const upstreamReceived = deferred(); const upstreamSocketClosed = deferred(); const upstream = await startServer(async (incoming) => { incoming.socket.once('close', () => upstreamSocketClosed.resolve()); await readBody(incoming); upstreamReceived.resolve(); }); const proxy = await startProxy({ upstreamBaseUrl: upstream.url, faultCount: 0, }); const activeRequest = request(proxy.url, { body: 'active-forward' }); await upstreamReceived.promise; expect(proxy.stats.forwardedRequestCount).toBe(1); await proxy.stop(); await expect(activeRequest).rejects.toBeInstanceOf(Error); await expect(upstreamSocketClosed.promise).resolves.toBeUndefined(); expect(proxy.stats.stopped).toBe(true); }); it('falls back to an injected high loopback port pool after port zero exhaustion', async () => { const upstream = await startServer(async (incoming, response) => { await readBody(incoming); response.end('ok'); }); const attempts: Array<{ host: string; port: number }> = []; const fallbackPorts = Array.from( { length: 64 }, (_, index) => 62_128 + index, ); const injectedListen = async ( server: Server, options: { host: string; port: number }, ) => { attempts.push(options); if (options.port === 0) { throw Object.assign(new Error('simulated ephemeral port exhaustion'), { code: 'EADDRINUSE', }); } await listen(server, options); }; const proxy = await startProxy({ upstreamBaseUrl: upstream.url, faultCount: 0, fallbackPorts, listen: injectedListen, }); expect(attempts[0]).toEqual({ host: '127.0.0.1', port: 0 }); expect(attempts.slice(1).every(({ host }) => host === '127.0.0.1')).toBe( true, ); expect(fallbackPorts).toContain(proxy.port); await expect( request(proxy.url, { body: 'fallback-body' }), ).resolves.toMatchObject({ statusCode: 200, body: 'ok' }); }); it('fails closed for non-POST and absolute-form requests without consuming faults', async () => { let upstreamRequestCount = 0; const upstream = await startServer(async (incoming, response) => { upstreamRequestCount += 1; await readBody(incoming); response.end('unexpected'); }); const proxy = await startProxy({ upstreamBaseUrl: upstream.url }); await expect( request(proxy.url, { method: 'GET', path: '/models' }), ).resolves.toMatchObject({ statusCode: 405 }); await expect( request(proxy.url, { path: 'http://second-origin.invalid/v1/chat/completions', body: 'absolute-form', }), ).resolves.toMatchObject({ statusCode: 400 }); expect(upstreamRequestCount).toBe(0); expect(proxy.stats).toMatchObject({ requestCount: 2, faultInjectedCount: 0, forwardedRequestCount: 0, }); await expect( request(proxy.url, { body: 'valid-post' }), ).rejects.toBeInstanceOf(Error); expect(proxy.stats.faultInjectedCount).toBe(1); expect(upstreamRequestCount).toBe(0); }); it('returns upstream redirects without following them to a second origin', async () => { let secondOriginRequestCount = 0; const secondOrigin = await startServer(async (incoming, response) => { secondOriginRequestCount += 1; await readBody(incoming); response.end('must-not-be-requested'); }); const upstream = await startServer(async (incoming, response) => { await readBody(incoming); response.writeHead(307, { location: `${secondOrigin.url}/redirect-target`, }); response.end('redirect-not-followed'); }); const proxy = await startProxy({ upstreamBaseUrl: upstream.url, faultCount: 0, }); await expect( request(proxy.url, { body: 'redirect-source' }), ).resolves.toMatchObject({ statusCode: 307, headers: expect.objectContaining({ location: `${secondOrigin.url}/redirect-target`, }), body: 'redirect-not-followed', }); expect(secondOriginRequestCount).toBe(0); expect(proxy.stats.forwardedRequestCount).toBe(1); }); it('rejects unsafe upstream URL forms without echoing them', async () => { const unsafeUrls = [ 'ftp://unsafe-url-marker.invalid/v1', 'https://user:password@unsafe-url-marker.invalid/v1', 'https://unsafe-url-marker.invalid/v1#fragment', ]; for (const upstreamBaseUrl of unsafeUrls) { let startupError: unknown; try { await startLlmTransientFaultProxy({ upstreamBaseUrl }); } catch (error) { startupError = error; } expect(startupError).toBeInstanceOf(Error); expect(String(startupError)).not.toContain('unsafe-url-marker'); } }); });