diff --git a/src/__tests__/proxy.integration.test.ts b/src/__tests__/proxy.integration.test.ts index 30b16a66..79a4580f 100644 --- a/src/__tests__/proxy.integration.test.ts +++ b/src/__tests__/proxy.integration.test.ts @@ -820,6 +820,38 @@ describe('Proxy usage recording – finish vs premature close', () => { void fetchError; }); + it('terminates the client stream when upstream fails after headers were flushed', async () => { + // Upstream promises 1000 bytes, sends a few, then resets the socket. + // Headers have already been flushed to the caller by then, so the error + // handler cannot send a 502; it must destroy the response so the caller + // sees an aborted stream instead of hanging or accepting a partial body. + setUpstreamHandler((_req, res) => { + res.writeHead(200, { 'content-type': 'application/json', 'content-length': '1000' }); + res.write('{"partial":'); + setTimeout(() => res.socket!.destroy(), 20); + }); + + const outcome = await Promise.race([ + (async () => { + const res = await fetch(`${proxyUrl}/v1/call/${TEST_API_SLUG}/mid-stream-terminate`, { + method: 'GET', + headers: { 'x-api-key': TEST_API_KEY }, + }); + return res.text().then( + () => 'completed' as const, + () => 'terminated' as const, + ); + })().catch(() => 'terminated' as const), + new Promise<'hung'>((resolve) => setTimeout(() => resolve('hung'), 3000)), + ]); + + expect(outcome).toBe('terminated'); + + await new Promise((resolve) => setTimeout(resolve, 50)); + expect(usageStore.getEvents()).toHaveLength(0); + expect(billing.getBalance(TEST_DEVELOPER_ID)).toBe(1000); + }); + it('does NOT double-count when the same requestId is seen twice', async () => { // Simulates a retry or duplicate delivery of the same logical request. const res1 = await fetch(`${proxyUrl}/v1/call/${TEST_API_SLUG}/idempotency-test`, { diff --git a/src/middleware/errorHandler.contract.test.ts b/src/middleware/errorHandler.contract.test.ts index 41f99521..cc6a236c 100644 --- a/src/middleware/errorHandler.contract.test.ts +++ b/src/middleware/errorHandler.contract.test.ts @@ -5,9 +5,12 @@ import { ValidationError } from './validate.js'; import { errorHandler } from './errorHandler.js'; function response() { - const result = { statusCode: 200, body: undefined as unknown, sent: false }; + const result = { statusCode: 200, body: undefined as unknown, sent: false, destroyedWith: undefined as unknown }; const value = { headersSent: false, + writableEnded: false, + destroyed: false, + destroy(err?: unknown) { result.destroyedWith = err; value.destroyed = true; return value; }, status(code: number) { result.statusCode = code; return value; }, json(body: unknown) { result.body = body; result.sent = true; return value; }, } as unknown as Response; @@ -113,8 +116,11 @@ describe('errorHandler contract', () => { it('does not overwrite a response that already sent headers', () => { const output = response(); (output.value as Response).headersSent = true; - errorHandler(new Error('already handled'), request(), output.value, jest.fn() as unknown as NextFunction); + const failure = new Error('already handled'); + errorHandler(failure, request(), output.value, jest.fn() as unknown as NextFunction); expect(output.result.sent).toBe(false); + // Nothing can be written after headers, so the stream is terminated instead. + expect(output.result.destroyedWith).toBe(failure); }); it('keeps error responses JSON-compatible for null and non-Error throws', () => { diff --git a/src/middleware/errorHandler.test.ts b/src/middleware/errorHandler.test.ts index 6409e1a9..20a9ddf5 100644 --- a/src/middleware/errorHandler.test.ts +++ b/src/middleware/errorHandler.test.ts @@ -1,4 +1,8 @@ -import { Request, Response, NextFunction } from 'express';import { errorHandler } from '../middleware/errorHandler.js'; +import express, { Request, Response, NextFunction } from 'express'; +import http from 'node:http'; +import type { Server } from 'node:http'; +import type { AddressInfo } from 'node:net'; +import { errorHandler } from '../middleware/errorHandler.js'; import { BadRequestError, UnauthorizedError, @@ -8,7 +12,8 @@ import { TooManyRequestsError, AppError, } from '../errors/index.js'; -import { ValidationError } from '../middleware/validate.js';import { logger } from '../logger.js'; +import { ValidationError } from '../middleware/validate.js'; +import { logger } from '../logger.js'; import type { ErrorEnvelope } from '../types/ResponseEnvelope.js'; jest.mock('../logger.js', () => ({ @@ -62,7 +67,7 @@ describe('Error Handler', () => { code: 'BAD_REQUEST', message: 'Test bad request', }); - expect(typeof call.timestamp).toBe(typeof 'string'); + expect(typeof call.timestamp).toBe('string'); expect(logger.error).toHaveBeenCalledWith( '[errorHandler]', @@ -125,19 +130,142 @@ describe('Error Handler', () => { expect(call.requestId).toBe('unknown'); }); - it('should not send response if headers already sent', () => { - mockRes.headersSent = true; - const error = new BadRequestError('Test error'); - - errorHandler( - error, - mockReq as Request, - mockRes as Response, - mockNext - ); + describe('when headers were already sent (mid-stream failure)', () => { + let destroy: jest.Mock; - expect(mockRes.status).not.toHaveBeenCalled(); - expect(mockRes.json).not.toHaveBeenCalled(); + beforeEach(() => { + destroy = jest.fn(); + Object.assign(mockRes, { + headersSent: true, + writableEnded: false, + destroyed: false, + destroy, + }); + }); + + it('does not write an envelope and destroys the response with the error', () => { + const error = new Error('upstream reset mid-stream'); + + errorHandler(error, mockReq as Request, mockRes as Response, mockNext); + + expect(mockRes.status).not.toHaveBeenCalled(); + expect(mockRes.json).not.toHaveBeenCalled(); + expect(destroy).toHaveBeenCalledTimes(1); + expect(destroy).toHaveBeenCalledWith(error); + }); + + it('logs the error exactly once with the requestId', () => { + errorHandler(new Error('boom'), mockReq as Request, mockRes as Response, mockNext); + + expect(logger.error).toHaveBeenCalledTimes(1); + expect(logger.error).toHaveBeenCalledWith( + '[errorHandler]', + expect.objectContaining({ requestId: 'test-request-id', headersSent: true }), + ); + }); + + it('does not delegate to next(err), which would log a second time', () => { + errorHandler(new Error('boom'), mockReq as Request, mockRes as Response, mockNext); + + expect(mockNext).not.toHaveBeenCalled(); + }); + + it('wraps non-Error throws so destroy always receives an Error', () => { + errorHandler('string failure', mockReq as Request, mockRes as Response, mockNext); + + expect(destroy).toHaveBeenCalledWith(expect.any(Error)); + expect((destroy.mock.calls[0][0] as Error).message).toBe('string failure'); + }); + + it('maps AppErrors thrown mid-stream to a destroy, not a status change', () => { + errorHandler(new BadRequestError('late'), mockReq as Request, mockRes as Response, mockNext); + + expect(mockRes.status).not.toHaveBeenCalled(); + expect(destroy).toHaveBeenCalledTimes(1); + }); + + it('leaves a cleanly ended response alone but still logs', () => { + (mockRes as { writableEnded: boolean }).writableEnded = true; + + errorHandler(new Error('after end'), mockReq as Request, mockRes as Response, mockNext); + + expect(destroy).not.toHaveBeenCalled(); + expect(logger.error).toHaveBeenCalledTimes(1); + }); + + it('does not destroy twice when the client already disconnected', () => { + (mockRes as { destroyed: boolean }).destroyed = true; + + errorHandler(new Error('client gone'), mockReq as Request, mockRes as Response, mockNext); + + expect(destroy).not.toHaveBeenCalled(); + expect(logger.error).toHaveBeenCalledTimes(1); + }); + + it('never throws if destroy itself fails', () => { + destroy.mockImplementation(() => { + throw new Error('destroy failed'); + }); + + expect(() => + errorHandler(new Error('boom'), mockReq as Request, mockRes as Response, mockNext), + ).not.toThrow(); + expect(logger.error).toHaveBeenCalledTimes(1); + expect(logger.warn).toHaveBeenCalledWith( + '[errorHandler] failed to destroy response after headers sent', + expect.objectContaining({ requestId: 'test-request-id' }), + ); + }); + }); + + it('leaves responses without headers sent unchanged (no destroy)', () => { + const destroy = jest.fn(); + Object.assign(mockRes, { destroy, writableEnded: false, destroyed: false }); + + errorHandler(new BadRequestError('Test error'), mockReq as Request, mockRes as Response, mockNext); + + expect(mockRes.status).toHaveBeenCalledWith(400); + expect(mockRes.json).toHaveBeenCalledTimes(1); + expect(destroy).not.toHaveBeenCalled(); + }); + + it('terminates a real socket when a handler throws after res.write', async () => { + const app = express(); + app.use((req, _res, next) => { + (req as Request & { id?: string }).id = 'real-socket-id'; + next(); + }); + app.get('/stream', (_req, res, next) => { + res.status(200).set('content-length', '1000'); + res.write('{"partial":'); + setImmediate(() => next(new Error('failed after first chunk'))); + }); + app.use(errorHandler); + + const server: Server = await new Promise((resolve) => { + const s = app.listen(0, () => resolve(s)); + }); + + try { + const { port } = server.address() as AddressInfo; + const outcome = await new Promise<'aborted' | 'completed'>((resolve) => { + const timeout = setTimeout(() => resolve('completed'), 2000); + http + .get({ port, path: '/stream', agent: false }, (res) => { + res.on('data', () => undefined); + res.on('end', () => { clearTimeout(timeout); resolve('completed'); }); + res.on('aborted', () => { clearTimeout(timeout); resolve('aborted'); }); + res.on('error', () => { clearTimeout(timeout); resolve('aborted'); }); + res.on('close', () => { clearTimeout(timeout); resolve(res.complete ? 'completed' : 'aborted'); }); + }) + .on('error', () => { clearTimeout(timeout); resolve('aborted'); }); + }); + + expect(outcome).toBe('aborted'); + expect(logger.error).toHaveBeenCalledTimes(1); + } finally { + await new Promise((resolve) => server.close(() => resolve())); + } }); it('should destroy the socket when headers are already sent', () => { diff --git a/src/middleware/errorHandler.ts b/src/middleware/errorHandler.ts index 23257ef3..83b07187 100644 --- a/src/middleware/errorHandler.ts +++ b/src/middleware/errorHandler.ts @@ -25,14 +25,44 @@ function extractValidationDetails(err: unknown): ValidationErrorDetail[] | undef return undefined; } +/** + * Terminate a response whose headers were already flushed. + * + * Once headers are out we can no longer send a JSON error envelope: writing one + * would corrupt the in-flight body, and simply returning leaves the client with + * a truncated stream that looks complete (or hangs until its own timeout) while + * the socket leaks. Destroying the response makes the client observe an aborted + * stream instead. + * + * - No-op if the response already ended cleanly or was already destroyed + * (e.g. the client disconnected first), so we never double-destroy or kill a + * keep-alive socket that completed successfully. + * - Never throws: an error handler must not raise a second error. + */ +function terminateAfterHeadersSent(res: Response, err: unknown, requestId: string): void { + if (res.writableEnded || res.destroyed) { + return; + } + + try { + res.destroy(err instanceof Error ? err : new Error(typeof err === 'string' ? err : 'Response aborted after headers were sent')); + } catch (destroyErr) { + logger.warn('[errorHandler] failed to destroy response after headers sent', { + requestId, + message: destroyErr instanceof Error ? destroyErr.message : String(destroyErr), + }); + } +} + /** * Global error-handling middleware (4-arg form). * - Catches errors thrown in routes/services * - Maps known AppError subclasses to HTTP status codes * - Returns consistent JSON envelope: { success: false, error: { code, message }, requestId, timestamp } * - Never sends stack traces to the client in production - * - Logs full error server-side - * - When headers are already sent, destroys the socket so the client sees a terminated stream + * - Logs full error server-side (exactly once, with requestId) + * - If headers were already sent (mid-stream failure), the error is logged and + * the response/socket is destroyed so the client sees a terminated stream */ export function errorHandler( err: unknown, @@ -77,20 +107,16 @@ export function errorHandler( normalized.simulationDetails, ); - if (!res.headersSent) { + const headersAlreadySent = res.headersSent; + + if (!headersAlreadySent) { res.status(statusCode).json(body); - } else { - // Headers already flushed: we cannot write a JSON envelope. - // Terminate the socket so the client observes a truncated stream - // instead of hanging until its own timeout. - if (typeof res.destroy === 'function') { - res.destroy(err instanceof Error ? err : undefined); - } } const logData = { requestId, statusCode, + ...(headersAlreadySent ? { headersSent: true } : {}), message: rawMessage, ...(isProduction ? {} : { err }), }; @@ -104,4 +130,10 @@ export function errorHandler( } else { logger.error("[errorHandler]", logData); } + + // Log first, then terminate, so the error is recorded exactly once even if + // teardown fails. + if (headersAlreadySent) { + terminateAfterHeadersSent(res, err, requestId); + } } diff --git a/src/routes/proxyRoutes.ts b/src/routes/proxyRoutes.ts index 6916059b..cafcf43d 100644 --- a/src/routes/proxyRoutes.ts +++ b/src/routes/proxyRoutes.ts @@ -268,6 +268,14 @@ export function createProxyRouter(deps: ProxyDeps): Router { res.status(upstreamStatus); if (upstreamRes.body) { const reader = upstreamRes.body.getReader(); + // Release the upstream connection if the client goes away before the + // body has been fully delivered. + const onClientClose = (): void => { + if (!res.writableFinished) { + void reader.cancel().catch(() => undefined); + } + }; + res.once('close', onClientClose); const pump = async (): Promise => { while (true) { const { done, value } = await reader.read(); @@ -276,7 +284,18 @@ export function createProxyRouter(deps: ProxyDeps): Router { } res.end(); }; - await pump(); + try { + await pump(); + } catch (streamErr) { + // Mid-stream failure: headers are already flushed, so the error + // handler cannot send an envelope. Cancel the upstream body so its + // socket is freed; errorHandler then destroys the client response + // so the caller sees a terminated stream instead of a hang. + void reader.cancel(streamErr).catch(() => undefined); + throw streamErr; + } finally { + res.removeListener('close', onClientClose); + } } else { const text = await upstreamRes.text(); res.send(text);