diff --git a/packages/stack/src/HttpProxy.integration.test.ts b/packages/stack/src/HttpProxy.integration.test.ts index bec4715c09..8e6398d412 100644 --- a/packages/stack/src/HttpProxy.integration.test.ts +++ b/packages/stack/src/HttpProxy.integration.test.ts @@ -38,6 +38,7 @@ class HttpProxyTestError extends Data.TaggedError("HttpProxyTestError")<{ }> {} const captureErrors = captureLogs(["Error"]); +const captureWarnings = captureLogs(["Error", "Warn"]); /** Upstream that resets its first `drops` accepted connections without responding. */ const droppingBackend = (drops: number) => { @@ -581,11 +582,7 @@ it.live("retries a bodyless request once when the upstream drops the connection }), ).pipe( Effect.provide( - Layer.mergeAll( - NodeHttpClient.layerNodeHttp, - NodeServices.layer, - captureLogs(["Error", "Warn"])(logs), - ), + Layer.mergeAll(NodeHttpClient.layerNodeHttp, NodeServices.layer, captureWarnings(logs)), ), ); }); @@ -617,50 +614,63 @@ it.live("does not replay a request with a body when the upstream drops the conne ); }); +/** + * Keep-alive upstream that answers once per connection and resets any reused connection; + * connections after `answered` are dropped unanswered. + */ +const oneRequestPerConnectionBackend = (answered = Number.POSITIVE_INFINITY) => { + let connections = 0; + const sockets = new Set(); + const server = createTcpServer((socket) => { + connections += 1; + sockets.add(socket); + socket.on("error", () => {}); + socket.on("close", () => sockets.delete(socket)); + if (connections > answered) { + socket.destroy(); + return; + } + let buffered = Buffer.alloc(0); + let bodyEnd: number | undefined; + let replied = false; + socket.on("data", (chunk: Buffer) => { + if (replied) { + socket.resetAndDestroy(); + return; + } + buffered = Buffer.concat([buffered, chunk]); + if (bodyEnd === undefined) { + const headerEnd = buffered.indexOf("\r\n\r\n"); + if (headerEnd === -1) return; + const contentLength = Number( + /content-length:\s*(\d+)/iu.exec( + buffered.subarray(0, headerEnd).toString("latin1"), + )?.[1] ?? 0, + ); + bodyEnd = headerEnd + 4 + contentLength; + } + if (buffered.length < bodyEnd) return; + replied = true; + socket.write("HTTP/1.1 200 OK\r\nConnection: keep-alive\r\nContent-Length: 2\r\n\r\nok"); + }); + }); + return { + connections: () => connections, + listen: listen(server, { + beforeClose: () => { + for (const socket of sockets) socket.destroy(); + }, + }), + }; +}; + it.live( "succeeds a second POST when a keep-alive backend only answers the first request per connection", () => Effect.scoped( Effect.gen(function* () { - let connections = 0; - const sockets = new Set(); - // Keeps each connection open after answering, then resets it if a second request reuses it. - const backend = createTcpServer((socket) => { - connections += 1; - sockets.add(socket); - socket.on("error", () => {}); - socket.on("close", () => sockets.delete(socket)); - let buffered = Buffer.alloc(0); - let bodyEnd: number | undefined; - let answered = false; - socket.on("data", (chunk: Buffer) => { - if (answered) { - socket.resetAndDestroy(); - return; - } - buffered = Buffer.concat([buffered, chunk]); - if (bodyEnd === undefined) { - const headerEnd = buffered.indexOf("\r\n\r\n"); - if (headerEnd === -1) return; - const contentLength = Number( - /content-length:\s*(\d+)/iu.exec( - buffered.subarray(0, headerEnd).toString("latin1"), - )?.[1] ?? 0, - ); - bodyEnd = headerEnd + 4 + contentLength; - } - if (buffered.length < bodyEnd) return; - answered = true; - socket.write( - "HTTP/1.1 200 OK\r\nConnection: keep-alive\r\nContent-Length: 2\r\n\r\nok", - ); - }); - }); - const address = yield* listen(backend, { - beforeClose: () => { - for (const socket of sockets) socket.destroy(); - }, - }); + const backend = oneRequestPerConnectionBackend(); + const address = yield* backend.listen; const proxy = yield* makeHttpProxy({ host: "127.0.0.1", port: 0 }); yield* proxy.setRoutes([{ id: "mcp", prefix: "/", target: Effect.succeed(address) }]); const first = yield* request(proxy.port, "/mcp", new TextEncoder().encode("first-body")); @@ -669,11 +679,149 @@ it.live( const second = yield* request(proxy.port, "/mcp", new TextEncoder().encode("second-body")); expect(second.status).toBe(200); expect(new TextDecoder().decode(second.body)).toBe("ok"); - expect(connections).toBe(2); + expect(backend.connections()).toBe(2); }), ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), ); +it.live( + "retries a GET once on a fresh connection when its pooled connection resets unanswered", + () => { + const logs: Array = []; + return Effect.scoped( + Effect.gen(function* () { + const backend = oneRequestPerConnectionBackend(); + const address = yield* backend.listen; + const proxy = yield* makeHttpProxy({ host: "127.0.0.1", port: 0 }); + yield* proxy.setRoutes([{ id: "studio", prefix: "/", target: Effect.succeed(address) }]); + const first = yield* request(proxy.port, "/", new Uint8Array(), {}, "GET"); + expect(first.status).toBe(200); + const second = yield* request(proxy.port, "/", new Uint8Array(), {}, "GET"); + expect(second.status).toBe(200); + expect(new TextDecoder().decode(second.body)).toBe("ok"); + expect(backend.connections()).toBe(2); + expect(logs).toHaveLength(1); + expect(logs[0]).toContain("Route studio GET upstream failed before responding, retrying"); + }), + ).pipe( + Effect.provide( + Layer.mergeAll(NodeHttpClient.layerNodeHttp, NodeServices.layer, captureWarnings(logs)), + ), + ); + }, +); + +it.live("returns a gateway error when the fresh retry after a pooled reset also fails", () => { + const logs: Array = []; + return Effect.scoped( + Effect.gen(function* () { + const backend = oneRequestPerConnectionBackend(1); + const address = yield* backend.listen; + const proxy = yield* makeHttpProxy({ host: "127.0.0.1", port: 0 }); + yield* proxy.setRoutes([{ id: "studio", prefix: "/", target: Effect.succeed(address) }]); + const first = yield* request(proxy.port, "/", new Uint8Array(), {}, "GET"); + expect(first.status).toBe(200); + const second = yield* request(proxy.port, "/", new Uint8Array(), {}, "GET"); + expect(second.status).toBe(502); + expect(backend.connections()).toBe(2); + expect(logs).toHaveLength(2); + expect(logs[0]).toContain("Route studio GET upstream failed before responding, retrying"); + expect(logs[1]).toContain("Route studio request failed"); + }), + ).pipe( + Effect.provide( + Layer.mergeAll(NodeHttpClient.layerNodeHttp, NodeServices.layer, captureWarnings(logs)), + ), + ); +}); + +it.live("delivers a keyed POST once when the upstream resets after reading its body", () => { + const logs: Array = []; + return Effect.scoped( + Effect.gen(function* () { + const bodies: Array = []; + let connections = 0; + // Answers GETs on a keep-alive connection; resets after reading a POST body in full. + const backend = createTcpServer((socket) => { + connections += 1; + socket.on("error", () => {}); + let buffered = Buffer.alloc(0); + socket.on("data", (chunk: Buffer) => { + buffered = Buffer.concat([buffered, chunk]); + const headerEnd = buffered.indexOf("\r\n\r\n"); + if (headerEnd === -1) return; + const head = buffered.subarray(0, headerEnd).toString("latin1"); + const bodyEnd = headerEnd + 4 + Number(/content-length:\s*(\d+)/iu.exec(head)?.[1] ?? 0); + if (buffered.length < bodyEnd) return; + const body = buffered.subarray(headerEnd + 4, bodyEnd).toString(); + buffered = buffered.subarray(bodyEnd); + if (head.startsWith("GET ")) { + socket.write( + "HTTP/1.1 200 OK\r\nConnection: keep-alive\r\nContent-Length: 2\r\n\r\nok", + ); + return; + } + bodies.push(body); + socket.resetAndDestroy(); + }); + }); + const address = yield* listen(backend); + const proxy = yield* makeHttpProxy({ host: "127.0.0.1", port: 0 }); + yield* proxy.setRoutes([{ id: "rest", prefix: "/", target: Effect.succeed(address) }]); + const pooled = yield* request(proxy.port, "/rest/v1/", new Uint8Array(), {}, "GET"); + expect(pooled.status).toBe(200); + const post = yield* request(proxy.port, "/rest/v1/rpc", new TextEncoder().encode("insert"), { + "idempotency-key": "insert", + }); + expect(post.status).toBe(502); + expect(bodies).toEqual(["insert"]); + expect(connections).toBe(2); + expect(logs).toHaveLength(1); + expect(logs[0]).toContain("Route rest request failed"); + }), + ).pipe( + Effect.provide( + Layer.mergeAll(NodeHttpClient.layerNodeHttp, NodeServices.layer, captureWarnings(logs)), + ), + ); +}); + +it.live("reuses pooled upstream connections across thousands of concurrent requests", () => + Effect.scoped( + Effect.gen(function* () { + let connections = 0; + const backend = createServer((incoming, outgoing) => { + incoming.resume(); + incoming.once("end", () => outgoing.end(incoming.url)); + }); + backend.on("connection", () => { + connections += 1; + }); + const address = yield* listen(backend); + const proxy = yield* makeHttpProxy({ host: "127.0.0.1", port: 0 }); + yield* proxy.setRoutes([{ id: "rest", prefix: "/", target: Effect.succeed(address) }]); + const responses = yield* Effect.forEach( + Array.from({ length: 3000 }, (_, index) => index), + (index) => + request(proxy.port, `/rest/v1/items?id=eq.${index}`, new Uint8Array(), {}, "GET").pipe( + Effect.map((response) => ({ + index, + status: response.status, + body: new TextDecoder().decode(response.body), + })), + ), + { concurrency: 16 }, + ); + expect( + responses.filter( + ({ index, status, body }) => status !== 200 || body !== `/rest/v1/items?id=eq.${index}`, + ), + ).toEqual([]); + expect(connections).toBeLessThan(100); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), +); + it.live( "does not replay a bodyless non-idempotent request when the upstream drops the connection", () => { diff --git a/packages/stack/src/HttpProxy.ts b/packages/stack/src/HttpProxy.ts index 4b4e5f0485..d2ae57d119 100644 --- a/packages/stack/src/HttpProxy.ts +++ b/packages/stack/src/HttpProxy.ts @@ -1,7 +1,8 @@ -import { Data, Effect, FiberSet, Ref, Schedule, Scope } from "effect"; +import { Data, Effect, FiberSet, Ref, Scope } from "effect"; import { PortError } from "./Ports.ts"; import type { BackendAddress, ProxyError } from "./Proxy.ts"; import { + Agent, createServer, request as upstreamRequest, type IncomingMessage, @@ -74,6 +75,10 @@ const hasBody = (request: IncomingMessage) => request.headers["transfer-encoding"] !== undefined || Number(request.headers["content-length"] ?? 0) > 0; +// A write can commit before its connection resets, so only safe, bodyless requests are resent. +const isReplayable = (request: IncomingMessage) => + safeMethods.has(request.method ?? "GET") && !hasBody(request); + const headersFor = (headers: IncomingMessage["headers"]) => Object.fromEntries( Object.entries(headers).filter( @@ -225,30 +230,24 @@ const connectInterruptibly = Effect.fn("HttpProxy.connect")((address: BackendAdd }), ); -// Mirrors nginx `proxy_next_upstream error`: a backend that drops a fresh connection before -// answering gets one more attempt, but only when nothing sent to the client or upstream would -// need replaying. -const retryOnce = (request: IncomingMessage, response: ServerResponse, route: HttpRoute) => - Schedule.recurs(1).pipe( - Schedule.setInputType(), - Schedule.while( - ({ input }) => - input._tag === "HttpProxyError" && - input.responded === false && - !response.destroyed && - safeMethods.has(request.method ?? "GET") && - !hasBody(request), - ), - Schedule.tap(({ input }) => - Effect.logWarning( - `Route ${route.id} ${request.method ?? "GET"} upstream failed before responding, retrying`, - input, - ), - ), - ); +// Mirrors nginx `proxy_next_upstream error`: a backend that drops a connection before answering +// gets one more attempt, but only when nothing sent to the client or upstream would need replaying. +const isRetryable = + (request: IncomingMessage, response: ServerResponse) => + (error: HttpProxyError | HttpProxyDisconnected) => + error._tag === "HttpProxyError" && + error.responded === false && + !response.destroyed && + isReplayable(request); const forward = Effect.fn("HttpProxy.forward")( - (request: IncomingMessage, response: ServerResponse, route: HttpRoute, backend: BackendAddress) => + ( + request: IncomingMessage, + response: ServerResponse, + route: HttpRoute, + backend: BackendAddress, + agent: Agent | false, + ) => Effect.callback((resume) => { let outgoing: ReturnType | undefined; let incoming: IncomingMessage | undefined; @@ -287,9 +286,7 @@ const forward = Effect.fn("HttpProxy.forward")( host: "path" in backend ? undefined : backend.host, port: "path" in backend ? undefined : backend.port, socketPath: "path" in backend ? backend.path : undefined, - // A reused upstream connection can be reset by a just-woken backend, and a body - // cannot be replayed. - agent: false, + agent, method: request.method, path: pathFor(request, route), // Bun ends a streamed upstream response early when the request says Connection: close. @@ -326,11 +323,16 @@ const forward = Effect.fn("HttpProxy.forward")( ); const proxyRequest = Effect.fn("HttpProxy.proxyRequest")( - (request: IncomingMessage, response: ServerResponse, route: HttpRoute) => + (request: IncomingMessage, response: ServerResponse, route: HttpRoute, agent: Agent) => Effect.gen(function* () { const backend = yield* Effect.raceFirst(route.target, disconnected(request, response)); - yield* forward(request, response, route, backend).pipe( - Effect.retry(retryOnce(request, response, route)), + yield* forward(request, response, route, backend, isReplayable(request) ? agent : false).pipe( + Effect.catchIf(isRetryable(request, response), (error) => + Effect.logWarning( + `Route ${route.id} ${request.method ?? "GET"} upstream failed before responding, retrying`, + error, + ).pipe(Effect.andThen(forward(request, response, route, backend, false))), + ), ); }), ); @@ -405,6 +407,12 @@ export const makeHttpProxy = (options: { }): Effect.Effect => Effect.gen(function* () { const routes = yield* Ref.make>([]); + // Sockets per upstream stay unbounded so long-lived streamed responses never queue requests. + // Idle sockets close before Node upstreams' default 5 s keep-alive timeout can race a reuse. + const agent = yield* Effect.acquireRelease( + Effect.sync(() => new Agent({ keepAlive: true, timeout: 4_000 })), + (value) => Effect.sync(() => value.destroy()), + ); const runRequest = yield* FiberSet.makeRuntime(); const sockets = new Set(); const server = createServer((request, response) => { @@ -422,7 +430,7 @@ export const makeHttpProxy = (options: { response.statusCode = 404; response.end("Not Found"); } else { - yield* proxyRequest(request, response, route).pipe( + yield* proxyRequest(request, response, route, agent).pipe( Effect.tapError((cause) => cause._tag === "HttpProxyDisconnected" ? Effect.void