Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 32 additions & 0 deletions src/__tests__/proxy.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -819,6 +819,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`, {
Expand Down
10 changes: 8 additions & 2 deletions src/middleware/errorHandler.contract.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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', () => {
Expand Down
152 changes: 139 additions & 13 deletions src/middleware/errorHandler.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,7 @@
import { Request, Response, NextFunction } from 'express';
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,
Expand Down Expand Up @@ -126,19 +129,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<ErrorEnvelope>,
mockNext
);
describe('when headers were already sent (mid-stream failure)', () => {
let destroy: jest.Mock;

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<ErrorEnvelope>, 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<ErrorEnvelope>, 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<ErrorEnvelope>, 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<ErrorEnvelope>, 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<ErrorEnvelope>, 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<ErrorEnvelope>, 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<ErrorEnvelope>, 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<ErrorEnvelope>, 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<ErrorEnvelope>, 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(mockRes.status).not.toHaveBeenCalled();
expect(mockRes.json).not.toHaveBeenCalled();
expect(outcome).toBe('aborted');
expect(logger.error).toHaveBeenCalledTimes(1);
} finally {
await new Promise<void>((resolve) => server.close(() => resolve()));
}
});

it('should include explicit catalog code when provided', () => {
Expand Down
44 changes: 42 additions & 2 deletions src/middleware/errorHandler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,13 +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<ErrorEnvelope>, 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
* - 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,
Expand Down Expand Up @@ -64,13 +95,16 @@ export function errorHandler(
});
const body = buildErrorEnvelope(normalized.code, normalized.message, requestId, normalized.details, normalized.retryAfterMs);

if (!res.headersSent) {
const headersAlreadySent = res.headersSent;

if (!headersAlreadySent) {
res.status(statusCode).json(body);
}

const logData = {
requestId,
statusCode,
...(headersAlreadySent ? { headersSent: true } : {}),
message: rawMessage,
...(isProduction ? {} : { err }),
};
Expand All @@ -84,4 +118,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);
}
}
21 changes: 20 additions & 1 deletion src/routes/proxyRoutes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -256,6 +256,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<void> => {
while (true) {
const { done, value } = await reader.read();
Expand All @@ -264,7 +272,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);
Expand Down