From 2dc265dfbfdb82a1eb9d3b27da6507153e5e988a Mon Sep 17 00:00:00 2001 From: avallete Date: Wed, 30 Sep 2026 20:43:30 +0200 Subject: [PATCH 1/4] fix(stack): reuse keep-alive upstream connections in the HTTP gateway The stack gateway opened a new upstream connection for every proxied request. Under sustained load this exhausts ephemeral ports natively (EADDRNOTAVAIL) and makes Docker Desktop's forwarder reset connections, so the gateway answers 502 in waves. Safe requests and requests carrying an Idempotency-Key now go through a keep-alive agent owned by the proxy, with bodies up to 1 MiB buffered so a request that meets a stale pooled connection is replayed once on a fresh one. Other requests keep a fresh connection and are never replayed. Idle pooled sockets close after 4 s, before Node upstreams' keep-alive timeout. Closes #6922 Co-Authored-By: Claude Opus 5.5 --- .../stack/src/HttpProxy.integration.test.ts | 300 +++++++++++++++--- packages/stack/src/HttpProxy.ts | 122 ++++--- 2 files changed, 336 insertions(+), 86 deletions(-) diff --git a/packages/stack/src/HttpProxy.integration.test.ts b/packages/stack/src/HttpProxy.integration.test.ts index bec4715c09..d3d622a1b9 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,211 @@ 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 POST once and fails it when the upstream resets after reading its body", () => { + const logs: Array = []; + return Effect.scoped( + Effect.gen(function* () { + const bodies: Array = []; + // Answers GETs on a keep-alive connection; resets after reading a POST body in full. + const backend = createTcpServer((socket) => { + 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")); + expect(post.status).toBe(502); + expect(bodies).toEqual(["insert"]); + 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( + "replays an idempotency-keyed POST once 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: "mcp", prefix: "/", target: Effect.succeed(address) }]); + const first = yield* request(proxy.port, "/mcp", new TextEncoder().encode("first-body"), { + "idempotency-key": "first", + }); + expect(first.status).toBe(200); + const second = yield* request(proxy.port, "/mcp", new TextEncoder().encode("second-body"), { + "idempotency-key": "second", + }); + 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 mcp POST upstream failed before responding, retrying"); + }), + ).pipe( + Effect.provide( + Layer.mergeAll(NodeHttpClient.layerNodeHttp, NodeServices.layer, captureWarnings(logs)), + ), + ); + }, +); + +it.live("sends a body too large to replay on a fresh upstream connection", () => { + 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: "storage", prefix: "/", target: Effect.succeed(address) }]); + const small = yield* request(proxy.port, "/object", new TextEncoder().encode("small"), { + "idempotency-key": "small", + }); + expect(small.status).toBe(200); + const large = yield* request(proxy.port, "/object", new Uint8Array(2 * 1024 * 1024).fill(1), { + "idempotency-key": "large", + }); + expect(large.status).toBe(200); + expect(new TextDecoder().decode(large.body)).toBe("ok"); + expect(backend.connections()).toBe(2); + expect(logs).toEqual([]); + }), + ).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) => { + let body = ""; + incoming.on("data", (chunk: Buffer) => { + body += chunk.toString(); + }); + incoming.once("end", () => outgoing.end(`${incoming.method} ${body}`)); + }); + 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) => + (index % 2 === 0 + ? request(proxy.port, "/rest/v1/", new Uint8Array(), {}, "GET") + : request(proxy.port, "/rest/v1/", new TextEncoder().encode(`row-${index}`), { + "idempotency-key": `key-${index}`, + }) + ).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 !== (index % 2 === 0 ? "GET " : `POST row-${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..158c35746f 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, @@ -15,6 +16,8 @@ class HttpProxyError extends Data.TaggedError("HttpProxyError")<{ readonly cause?: unknown; /** Whether upstream response headers had arrived when a proxied request failed. */ readonly responded?: boolean; + /** Whether the failed request was sent on a pooled keep-alive connection. */ + readonly reused?: boolean; }> {} /** Distinguishes a client that went away first from a genuine proxy failure. */ @@ -57,22 +60,47 @@ const hopByHop = new Set([ "upgrade", ]); -const errorFor = (cause: unknown, responded?: boolean) => +const errorFor = (cause: unknown, responded?: boolean, reused?: boolean) => new HttpProxyError({ message: cause instanceof Error ? cause.message : String(cause), cause, ...(responded === undefined ? {} : { responded }), + ...(reused === undefined ? {} : { reused }), }); // RFC 9110 section 9.2.1 safe methods only: user functions behind the proxy need not honor // PUT or DELETE idempotency. const safeMethods = new Set(["GET", "HEAD", "OPTIONS", "TRACE"]); -// RFC 9112 section 6: a request carries a body only when Content-Length or -// Transfer-Encoding is present, regardless of method. -const hasBody = (request: IncomingMessage) => - request.headers["transfer-encoding"] !== undefined || - Number(request.headers["content-length"] ?? 0) > 0; +// Bodies up to this size are buffered so a request that meets a stale pooled connection can be +// replayed on a fresh one. +const replayableBodyLimit = 1024 * 1024; + +// A POST can commit before its connection resets, so only safe or `Idempotency-Key` requests replay. +const isReplayable = (request: IncomingMessage) => + (safeMethods.has(request.method ?? "GET") || request.headers["idempotency-key"] !== undefined) && + request.headers["transfer-encoding"] === undefined && + Number(request.headers["content-length"] ?? 0) <= replayableBodyLimit; + +const bufferBody = ( + request: IncomingMessage, +): Effect.Effect => + !isReplayable(request) + ? Effect.undefined + : Effect.callback((resume) => { + const chunks: Array = []; + const onData = (chunk: Buffer) => chunks.push(chunk); + const onEnd = () => resume(Effect.succeed(Buffer.concat(chunks))); + const onError = () => resume(Effect.fail(new HttpProxyDisconnected())); + request.on("data", onData); + request.once("end", onEnd); + request.once("error", onError); + return Effect.sync(() => { + request.off("data", onData); + request.off("end", onEnd); + request.off("error", onError); + }); + }); const headersFor = (headers: IncomingMessage["headers"]) => Object.fromEntries( @@ -225,30 +253,26 @@ 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, - ), - ), - ); +// A reused connection failing before any response is most likely the keep-alive close race; a fresh +// one is retried only when nothing needs replaying, like nginx `proxy_next_upstream error`. +const isRetryable = + (request: IncomingMessage, response: ServerResponse, body: Buffer | undefined) => + (error: HttpProxyError | HttpProxyDisconnected) => + error._tag === "HttpProxyError" && + error.responded === false && + !response.destroyed && + body !== undefined && + (error.reused === true || (safeMethods.has(request.method ?? "GET") && body.length === 0)); const forward = Effect.fn("HttpProxy.forward")( - (request: IncomingMessage, response: ServerResponse, route: HttpRoute, backend: BackendAddress) => + ( + request: IncomingMessage, + response: ServerResponse, + route: HttpRoute, + backend: BackendAddress, + body: Buffer | undefined, + agent: Agent | false, + ) => Effect.callback((resume) => { let outgoing: ReturnType | undefined; let incoming: IncomingMessage | undefined; @@ -276,7 +300,7 @@ const forward = Effect.fn("HttpProxy.forward")( incoming?.destroy(); }; const onError = (cause: Error) => - abandon(Effect.fail(errorFor(cause, incoming !== undefined))); + abandon(Effect.fail(errorFor(cause, incoming !== undefined, outgoing?.reusedSocket))); const onClientGone = () => abandon(Effect.fail(new HttpProxyDisconnected())); const onFinish = () => finish(Effect.void); const onResponseClose = () => { @@ -287,9 +311,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. @@ -312,10 +334,8 @@ const forward = Effect.fn("HttpProxy.forward")( outgoing.on("error", onError); request.once("aborted", onClientGone); response.once("close", onResponseClose); - // A retried bodyless request was already drained by the first attempt and emits no - // further `end`, so pipe would never finish the upstream request. - if (request.readableEnded) outgoing.end(); - else request.pipe(outgoing); + if (body === undefined) request.pipe(outgoing); + else outgoing.end(body); return Effect.sync(() => { settled = true; cleanup(); @@ -326,11 +346,25 @@ 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)), + const body = yield* Effect.raceFirst(bufferBody(request), disconnected(request, response)); + // Only a replayable request may meet a stale pooled connection; a retry always opens a fresh one. + yield* forward( + request, + response, + route, + backend, + body, + body === undefined ? false : agent, + ).pipe( + Effect.catchIf(isRetryable(request, response, body), (error) => + Effect.logWarning( + `Route ${route.id} ${request.method ?? "GET"} upstream failed before responding, retrying`, + error, + ).pipe(Effect.andThen(forward(request, response, route, backend, body, false))), + ), ); }), ); @@ -405,6 +439,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 +462,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 From 958ba458f28f03e020f08e34c643b4fb48f1dc53 Mon Sep 17 00:00:00 2001 From: avallete Date: Wed, 30 Sep 2026 20:52:11 +0200 Subject: [PATCH 2/4] fix(stack): pool and replay only safe, bodyless upstream requests Stack upstreams do not deduplicate on Idempotency-Key, so replaying a keyed write after a pooled connection reset could run it twice. Only safe, bodyless requests now share pooled connections and get the single fresh-connection retry; every other request keeps its own connection and is never replayed. Co-Authored-By: Claude Opus 5.5 --- .../stack/src/HttpProxy.integration.test.ts | 81 ++----------------- packages/stack/src/HttpProxy.ts | 73 +++++------------ 2 files changed, 29 insertions(+), 125 deletions(-) diff --git a/packages/stack/src/HttpProxy.integration.test.ts b/packages/stack/src/HttpProxy.integration.test.ts index d3d622a1b9..a844a6f152 100644 --- a/packages/stack/src/HttpProxy.integration.test.ts +++ b/packages/stack/src/HttpProxy.integration.test.ts @@ -735,7 +735,7 @@ it.live("returns a gateway error when the fresh retry after a pooled reset also ); }); -it.live("delivers a POST once and fails it when the upstream resets after reading its body", () => { +it.live("delivers a keyed POST once when the upstream resets after reading its body", () => { const logs: Array = []; return Effect.scoped( Effect.gen(function* () { @@ -768,7 +768,9 @@ it.live("delivers a POST once and fails it when the upstream resets after readin 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")); + 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(logs).toHaveLength(1); @@ -781,74 +783,13 @@ it.live("delivers a POST once and fails it when the upstream resets after readin ); }); -it.live( - "replays an idempotency-keyed POST once 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: "mcp", prefix: "/", target: Effect.succeed(address) }]); - const first = yield* request(proxy.port, "/mcp", new TextEncoder().encode("first-body"), { - "idempotency-key": "first", - }); - expect(first.status).toBe(200); - const second = yield* request(proxy.port, "/mcp", new TextEncoder().encode("second-body"), { - "idempotency-key": "second", - }); - 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 mcp POST upstream failed before responding, retrying"); - }), - ).pipe( - Effect.provide( - Layer.mergeAll(NodeHttpClient.layerNodeHttp, NodeServices.layer, captureWarnings(logs)), - ), - ); - }, -); - -it.live("sends a body too large to replay on a fresh upstream connection", () => { - 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: "storage", prefix: "/", target: Effect.succeed(address) }]); - const small = yield* request(proxy.port, "/object", new TextEncoder().encode("small"), { - "idempotency-key": "small", - }); - expect(small.status).toBe(200); - const large = yield* request(proxy.port, "/object", new Uint8Array(2 * 1024 * 1024).fill(1), { - "idempotency-key": "large", - }); - expect(large.status).toBe(200); - expect(new TextDecoder().decode(large.body)).toBe("ok"); - expect(backend.connections()).toBe(2); - expect(logs).toEqual([]); - }), - ).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) => { - let body = ""; - incoming.on("data", (chunk: Buffer) => { - body += chunk.toString(); - }); - incoming.once("end", () => outgoing.end(`${incoming.method} ${body}`)); + incoming.resume(); + incoming.once("end", () => outgoing.end(incoming.url)); }); backend.on("connection", () => { connections += 1; @@ -859,12 +800,7 @@ it.live("reuses pooled upstream connections across thousands of concurrent reque const responses = yield* Effect.forEach( Array.from({ length: 3000 }, (_, index) => index), (index) => - (index % 2 === 0 - ? request(proxy.port, "/rest/v1/", new Uint8Array(), {}, "GET") - : request(proxy.port, "/rest/v1/", new TextEncoder().encode(`row-${index}`), { - "idempotency-key": `key-${index}`, - }) - ).pipe( + request(proxy.port, `/rest/v1/items?id=eq.${index}`, new Uint8Array(), {}, "GET").pipe( Effect.map((response) => ({ index, status: response.status, @@ -875,8 +811,7 @@ it.live("reuses pooled upstream connections across thousands of concurrent reque ); expect( responses.filter( - ({ index, status, body }) => - status !== 200 || body !== (index % 2 === 0 ? "GET " : `POST row-${index}`), + ({ index, status, body }) => status !== 200 || body !== `/rest/v1/items?id=eq.${index}`, ), ).toEqual([]); expect(connections).toBeLessThan(100); diff --git a/packages/stack/src/HttpProxy.ts b/packages/stack/src/HttpProxy.ts index 158c35746f..97a0814678 100644 --- a/packages/stack/src/HttpProxy.ts +++ b/packages/stack/src/HttpProxy.ts @@ -16,8 +16,6 @@ class HttpProxyError extends Data.TaggedError("HttpProxyError")<{ readonly cause?: unknown; /** Whether upstream response headers had arrived when a proxied request failed. */ readonly responded?: boolean; - /** Whether the failed request was sent on a pooled keep-alive connection. */ - readonly reused?: boolean; }> {} /** Distinguishes a client that went away first from a genuine proxy failure. */ @@ -60,47 +58,26 @@ const hopByHop = new Set([ "upgrade", ]); -const errorFor = (cause: unknown, responded?: boolean, reused?: boolean) => +const errorFor = (cause: unknown, responded?: boolean) => new HttpProxyError({ message: cause instanceof Error ? cause.message : String(cause), cause, ...(responded === undefined ? {} : { responded }), - ...(reused === undefined ? {} : { reused }), }); // RFC 9110 section 9.2.1 safe methods only: user functions behind the proxy need not honor // PUT or DELETE idempotency. const safeMethods = new Set(["GET", "HEAD", "OPTIONS", "TRACE"]); -// Bodies up to this size are buffered so a request that meets a stale pooled connection can be -// replayed on a fresh one. -const replayableBodyLimit = 1024 * 1024; +// RFC 9112 section 6: a request carries a body only when Content-Length or +// Transfer-Encoding is present, regardless of method. +const hasBody = (request: IncomingMessage) => + request.headers["transfer-encoding"] !== undefined || + Number(request.headers["content-length"] ?? 0) > 0; -// A POST can commit before its connection resets, so only safe or `Idempotency-Key` requests replay. +// A write can commit before its connection resets, so only safe, bodyless requests are resent. const isReplayable = (request: IncomingMessage) => - (safeMethods.has(request.method ?? "GET") || request.headers["idempotency-key"] !== undefined) && - request.headers["transfer-encoding"] === undefined && - Number(request.headers["content-length"] ?? 0) <= replayableBodyLimit; - -const bufferBody = ( - request: IncomingMessage, -): Effect.Effect => - !isReplayable(request) - ? Effect.undefined - : Effect.callback((resume) => { - const chunks: Array = []; - const onData = (chunk: Buffer) => chunks.push(chunk); - const onEnd = () => resume(Effect.succeed(Buffer.concat(chunks))); - const onError = () => resume(Effect.fail(new HttpProxyDisconnected())); - request.on("data", onData); - request.once("end", onEnd); - request.once("error", onError); - return Effect.sync(() => { - request.off("data", onData); - request.off("end", onEnd); - request.off("error", onError); - }); - }); + safeMethods.has(request.method ?? "GET") && !hasBody(request); const headersFor = (headers: IncomingMessage["headers"]) => Object.fromEntries( @@ -253,16 +230,15 @@ const connectInterruptibly = Effect.fn("HttpProxy.connect")((address: BackendAdd }), ); -// A reused connection failing before any response is most likely the keep-alive close race; a fresh -// one is retried only when nothing needs replaying, like nginx `proxy_next_upstream error`. +// 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, body: Buffer | undefined) => + (request: IncomingMessage, response: ServerResponse) => (error: HttpProxyError | HttpProxyDisconnected) => error._tag === "HttpProxyError" && error.responded === false && !response.destroyed && - body !== undefined && - (error.reused === true || (safeMethods.has(request.method ?? "GET") && body.length === 0)); + isReplayable(request); const forward = Effect.fn("HttpProxy.forward")( ( @@ -270,7 +246,6 @@ const forward = Effect.fn("HttpProxy.forward")( response: ServerResponse, route: HttpRoute, backend: BackendAddress, - body: Buffer | undefined, agent: Agent | false, ) => Effect.callback((resume) => { @@ -300,7 +275,7 @@ const forward = Effect.fn("HttpProxy.forward")( incoming?.destroy(); }; const onError = (cause: Error) => - abandon(Effect.fail(errorFor(cause, incoming !== undefined, outgoing?.reusedSocket))); + abandon(Effect.fail(errorFor(cause, incoming !== undefined))); const onClientGone = () => abandon(Effect.fail(new HttpProxyDisconnected())); const onFinish = () => finish(Effect.void); const onResponseClose = () => { @@ -334,8 +309,10 @@ const forward = Effect.fn("HttpProxy.forward")( outgoing.on("error", onError); request.once("aborted", onClientGone); response.once("close", onResponseClose); - if (body === undefined) request.pipe(outgoing); - else outgoing.end(body); + // A retried bodyless request was already drained by the first attempt and emits no + // further `end`, so pipe would never finish the upstream request. + if (request.readableEnded) outgoing.end(); + else request.pipe(outgoing); return Effect.sync(() => { settled = true; cleanup(); @@ -349,21 +326,13 @@ const proxyRequest = Effect.fn("HttpProxy.proxyRequest")( (request: IncomingMessage, response: ServerResponse, route: HttpRoute, agent: Agent) => Effect.gen(function* () { const backend = yield* Effect.raceFirst(route.target, disconnected(request, response)); - const body = yield* Effect.raceFirst(bufferBody(request), disconnected(request, response)); - // Only a replayable request may meet a stale pooled connection; a retry always opens a fresh one. - yield* forward( - request, - response, - route, - backend, - body, - body === undefined ? false : agent, - ).pipe( - Effect.catchIf(isRetryable(request, response, body), (error) => + // Only replayable requests share pooled connections; the retry always opens a fresh one. + 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, body, false))), + ).pipe(Effect.andThen(forward(request, response, route, backend, false))), ), ); }), From 23a63212b6f54f644ca7db1bd3a231d820551016 Mon Sep 17 00:00:00 2001 From: avallete Date: Wed, 30 Sep 2026 21:34:17 +0200 Subject: [PATCH 3/4] chore(stack): drop a comment that restates the pooling expression Co-Authored-By: Claude Opus 5.5 --- packages/stack/src/HttpProxy.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/packages/stack/src/HttpProxy.ts b/packages/stack/src/HttpProxy.ts index 97a0814678..d2ae57d119 100644 --- a/packages/stack/src/HttpProxy.ts +++ b/packages/stack/src/HttpProxy.ts @@ -326,7 +326,6 @@ const proxyRequest = Effect.fn("HttpProxy.proxyRequest")( (request: IncomingMessage, response: ServerResponse, route: HttpRoute, agent: Agent) => Effect.gen(function* () { const backend = yield* Effect.raceFirst(route.target, disconnected(request, response)); - // Only replayable requests share pooled connections; the retry always opens a fresh one. yield* forward(request, response, route, backend, isReplayable(request) ? agent : false).pipe( Effect.catchIf(isRetryable(request, response), (error) => Effect.logWarning( From 71d1c6c1947e807bc407ab1a059c6574cdb9d3d8 Mon Sep 17 00:00:00 2001 From: avallete Date: Wed, 30 Sep 2026 22:58:26 +0200 Subject: [PATCH 4/4] test(stack): assert a keyed POST uses a fresh upstream connection Co-Authored-By: Claude Opus 5.5 --- packages/stack/src/HttpProxy.integration.test.ts | 3 +++ 1 file changed, 3 insertions(+) diff --git a/packages/stack/src/HttpProxy.integration.test.ts b/packages/stack/src/HttpProxy.integration.test.ts index a844a6f152..8e6398d412 100644 --- a/packages/stack/src/HttpProxy.integration.test.ts +++ b/packages/stack/src/HttpProxy.integration.test.ts @@ -740,8 +740,10 @@ it.live("delivers a keyed POST once when the upstream resets after reading its b 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) => { @@ -773,6 +775,7 @@ it.live("delivers a keyed POST once when the upstream resets after reading its b }); 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"); }),