Skip to content
Merged
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
238 changes: 193 additions & 45 deletions packages/stack/src/HttpProxy.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) => {
Expand Down Expand Up @@ -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)),
),
);
});
Expand Down Expand Up @@ -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<Socket>();
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<Socket>();
// 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"));
Expand All @@ -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<string> = [];
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<string> = [];
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<string> = [];
return Effect.scoped(
Effect.gen(function* () {
const bodies: Array<string> = [];
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"]);
Comment thread
avallete marked this conversation as resolved.
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",
() => {
Expand Down
68 changes: 38 additions & 30 deletions packages/stack/src/HttpProxy.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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<HttpProxyError | HttpProxyDisconnected>(),
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<void, HttpProxyError | HttpProxyDisconnected>((resume) => {
let outgoing: ReturnType<typeof upstreamRequest> | undefined;
let incoming: IncomingMessage | undefined;
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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))),
),
);
}),
);
Expand Down Expand Up @@ -405,6 +407,12 @@ export const makeHttpProxy = (options: {
}): Effect.Effect<HttpProxy, PortError, Scope.Scope> =>
Effect.gen(function* () {
const routes = yield* Ref.make<ReadonlyArray<HttpRoute>>([]);
// 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<Socket>();
const server = createServer((request, response) => {
Expand All @@ -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
Expand Down
Loading