From bf49475921d075bf8f12b4195c06ea1bda76a049 Mon Sep 17 00:00:00 2001 From: Andrew Valleteau Date: Mon, 21 Sep 2026 16:31:51 +0000 Subject: [PATCH 1/5] fix(stack): surface failed lazy wakes through the proxy and status A wake failure reached clients only as a dead connection: the request path answered 502 and the upgrade path destroyed the socket, both discarding the cause, while preparation failures never reached the service observation. Any failed wake was therefore indistinguishable from a flake. Both proxy arms now log the route and the underlying cause, with a client that disconnects first carrying its own tagged error so ordinary aborts stay quiet, and a failed wake preparation is recorded on the observation so status names it. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01LcUaThuQfDKPwi9mCE9XeD --- .../stack/src/HttpProxy.integration.test.ts | 75 ++++++++++++++++--- packages/stack/src/HttpProxy.ts | 21 +++++- .../stack/src/Service.integration.test.ts | 20 +++++ packages/stack/src/Service.ts | 4 +- 4 files changed, 105 insertions(+), 15 deletions(-) diff --git a/packages/stack/src/HttpProxy.integration.test.ts b/packages/stack/src/HttpProxy.integration.test.ts index 88ab3eca80..df5c991e67 100644 --- a/packages/stack/src/HttpProxy.integration.test.ts +++ b/packages/stack/src/HttpProxy.integration.test.ts @@ -1,6 +1,6 @@ import { NodeHttpClient, NodeServices } from "@effect/platform-node"; import { expect, it } from "@effect/vitest"; -import { Data, Deferred, Effect, Fiber, Layer } from "effect"; +import { Data, Deferred, Effect, Fiber, Layer, Logger } from "effect"; import { HttpClient, HttpClientRequest } from "effect/unstable/http"; import { createServer, type Server, type ServerResponse } from "node:http"; // oxlint-disable-line effecttsgo/node-builtin-import -- raw server fixture. import { Socket } from "node:net"; // oxlint-disable-line effecttsgo/node-builtin-import -- raw disconnect fixture. @@ -35,6 +35,14 @@ class HttpProxyTestError extends Data.TaggedError("HttpProxyTestError")<{ readonly cause?: unknown; }> {} +const captureErrors = (lines: Array) => + Logger.layer([ + Logger.make(({ logLevel, message }) => { + if (logLevel === "Error") + lines.push((Array.isArray(message) ? message : [message]).map(String).join(" ")); + }), + ]); + const request = (port: number, path: string, body: Uint8Array) => Effect.gen(function* () { const client = yield* HttpClient.HttpClient; @@ -128,8 +136,9 @@ it.live("keeps the retained listener and remaining route after one route is remo ).pipe(Effect.provide(NodeServices.layer)), ); -it.live("interrupts target acquisition when a waiting client disconnects", () => - Effect.scoped( +it.live("interrupts target acquisition quietly when a waiting client disconnects", () => { + const logs: Array = []; + return Effect.scoped( Effect.gen(function* () { const backend = createServer((_request, response) => response.end("unused")); const backendAddress = yield* listen(backend); @@ -173,9 +182,10 @@ it.live("interrupts target acquisition when a waiting client disconnects", () => yield* Deferred.await(acquired).pipe(Effect.timeout("5 seconds")); yield* Effect.sync(() => client.destroy()); yield* Deferred.await(released); + expect(logs).toEqual([]); }), - ).pipe(Effect.provide(NodeServices.layer)), -); + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, captureErrors(logs)))); +}); it.live("forwards raw WebSocket upgrades, subprotocols, and echo frames", () => Effect.scoped( @@ -259,13 +269,14 @@ it.live("overrides the upstream host for HTTP routes when configured", () => ).pipe(Effect.provide(NodeServices.layer)), ); -it.live("returns a gateway error when a managed target cannot become ready", () => - Effect.scoped( +it.live("returns a gateway error naming the route and cause when a target cannot wake", () => { + const logs: Array = []; + return Effect.scoped( Effect.gen(function* () { const proxy = yield* makeHttpProxy({ host: "127.0.0.1", port: 0 }); yield* proxy.setRoutes([ { - id: "failed", + id: "rest", prefix: "/", target: Effect.fail(new ProxyError({ message: "readiness failed" })), }, @@ -275,9 +286,53 @@ it.live("returns a gateway error when a managed target cannot become ready", () expect(response.status).toBe(502); expect(response.headers["access-control-allow-origin"]).toBe("*"); expect(yield* response.text).toBe("Bad Gateway"); + expect(logs).toHaveLength(1); + expect(logs[0]).toContain("Route rest request failed"); + expect(logs[0]).toContain("readiness failed"); }), - ).pipe(Effect.provide(Layer.merge(NodeHttpClient.layerNodeHttp, NodeServices.layer))), -); + ).pipe( + Effect.provide( + Layer.mergeAll(NodeHttpClient.layerNodeHttp, NodeServices.layer, captureErrors(logs)), + ), + ); +}); + +it.live("closes an upgrade naming the route and cause when a target cannot wake", () => { + const logs: Array = []; + return Effect.scoped( + Effect.gen(function* () { + const proxy = yield* makeHttpProxy({ host: "127.0.0.1", port: 0 }); + yield* proxy.setRoutes([ + { + id: "realtime", + prefix: "/socket", + target: Effect.fail(new ProxyError({ message: "wake failed" })), + }, + ]); + const socket = yield* Effect.acquireRelease( + Effect.sync(() => new Socket()), + (value) => Effect.sync(() => value.destroy()), + ); + const received: Array = []; + yield* Effect.callback((resume) => { + socket.on("data", (chunk: Buffer) => received.push(chunk)); + // A destroyed upgrade reaches the client as a reset, which closes the socket either way. + socket.on("error", () => undefined); + socket.once("close", () => resume(Effect.void)); + socket.connect(proxy.port, "127.0.0.1", () => + socket.write( + "GET /socket HTTP/1.1\r\nHost: localhost\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\r\n", + ), + ); + return Effect.void; + }).pipe(Effect.timeout("5 seconds")); + expect(Buffer.concat(received)).toHaveLength(0); + expect(logs).toHaveLength(1); + expect(logs[0]).toContain("Route realtime upgrade failed"); + expect(logs[0]).toContain("wake failed"); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, captureErrors(logs)))); +}); it.live("disconnects a pending upstream response when its client closes", () => Effect.scoped( diff --git a/packages/stack/src/HttpProxy.ts b/packages/stack/src/HttpProxy.ts index a1a0e28bb5..832b63e993 100644 --- a/packages/stack/src/HttpProxy.ts +++ b/packages/stack/src/HttpProxy.ts @@ -15,6 +15,9 @@ class HttpProxyError extends Data.TaggedError("HttpProxyError")<{ readonly cause?: unknown; }> {} +/** Distinguishes a client that went away first from a genuine proxy failure. */ +class HttpProxyDisconnected extends Data.TaggedError("HttpProxyDisconnected") {} + export interface HttpRoute { readonly id: string; readonly prefix: string; @@ -85,8 +88,8 @@ const setCors = (response: ServerResponse, request: IncomingMessage) => { }; const disconnected = (request: IncomingMessage, response: ServerResponse) => - Effect.callback((resume) => { - const onAbort = () => resume(Effect.fail(errorFor("client disconnected"))); + Effect.callback((resume) => { + const onAbort = () => resume(Effect.fail(new HttpProxyDisconnected())); const onRequestClose = () => { if (!request.complete) onAbort(); }; @@ -201,8 +204,8 @@ const upgrade = Effect.fn("HttpProxy.upgrade")( Effect.gen(function* () { const backend = yield* Effect.raceFirst( route.target, - Effect.callback((resume) => { - const onClose = () => resume(Effect.fail(errorFor("client disconnected"))); + Effect.callback((resume) => { + const onClose = () => resume(Effect.fail(new HttpProxyDisconnected())); client.once("close", onClose); if (client.destroyed) onClose(); return Effect.sync(() => client.off("close", onClose)); @@ -286,6 +289,11 @@ export const makeHttpProxy = (options: { response.end("Not Found"); } else { yield* proxyRequest(request, response, route).pipe( + Effect.tapError((cause) => + cause._tag === "HttpProxyDisconnected" + ? Effect.void + : Effect.logError(`Route ${route.id} request failed`, cause), + ), Effect.catch(() => Effect.sync(() => { if (response.destroyed) return; @@ -318,6 +326,11 @@ export const makeHttpProxy = (options: { if (route === undefined) socket.destroy(); else yield* upgrade(request, socket, head, route).pipe( + Effect.tapError((cause) => + cause._tag === "HttpProxyDisconnected" + ? Effect.void + : Effect.logError(`Route ${route.id} upgrade failed`, cause), + ), Effect.catch(() => Effect.sync(() => socket.destroy())), ); }), diff --git a/packages/stack/src/Service.integration.test.ts b/packages/stack/src/Service.integration.test.ts index 403695fbf9..1860059d60 100644 --- a/packages/stack/src/Service.integration.test.ts +++ b/packages/stack/src/Service.integration.test.ts @@ -288,6 +288,26 @@ describe("service kernel", () => { ), ); + it.live("records a failed wake preparation on the observation", () => + Effect.scoped( + Effect.gen(function* () { + const plans = yield* Queue.unbounded(); + const invalid = yield* Ref.make(true); + const service = yield* makeService(makeDefinition(plans, undefined, undefined, invalid), { + id: "database-wake-prepare", + config: { version: 17 }, + }); + yield* service.arm; + const failure = yield* service.startAt(0, undefined, true).pipe(Effect.flip); + expect(failure).toBeInstanceOf(ServiceError); + const observation = yield* service.get; + expect(observation.lifecycle).toBe("stopped"); + expect(observation.wakeEnabled).toBe(true); + expect(observation.error?.message).toBe("invalid configuration"); + }), + ), + ); + it.live("does not prepare concurrent starts more than once", () => Effect.scoped( Effect.gen(function* () { diff --git a/packages/stack/src/Service.ts b/packages/stack/src/Service.ts index 5ee3f910e4..d2d8120122 100644 --- a/packages/stack/src/Service.ts +++ b/packages/stack/src/Service.ts @@ -428,7 +428,9 @@ export const makeService = ( if (existing.lifecycle === "running") return; } const nextConfig = candidate ?? (yield* Ref.get(config)); - if (definition.prepare !== undefined) yield* definition.prepare(nextConfig); + // A wake has no caller to receive a preparation failure, so observers record it instead. + if (definition.prepare !== undefined) + yield* definition.prepare(nextConfig).pipe(Effect.tapError((error) => update({ error }))); yield* run( Effect.gen(function* () { const observation = yield* SubscriptionRef.get(observations); From 52bfbd275749debec721a4121585893095f4ed49 Mon Sep 17 00:00:00 2001 From: Andrew Valleteau Date: Mon, 21 Sep 2026 16:31:57 +0000 Subject: [PATCH 2/5] fix(stack): redraw auto port probes so reserved ranges cannot exhaust them The auto scan walked contiguously from a random start and gave up after 64 consecutive bind failures, so a start inside a reserved block wider than 64 ports failed with "No public port is available" while tens of thousands of ports were free. Windows publishes such excluded ranges, which is where this surfaced. Each probe is now redrawn, and the exhausted error carries the last bind failure so a reserved range reads differently from an occupied port. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01LcUaThuQfDKPwi9mCE9XeD --- packages/stack/src/Ports.integration.test.ts | 58 ++++++++++++++++++++ packages/stack/src/Ports.ts | 13 ++++- 2 files changed, 68 insertions(+), 3 deletions(-) diff --git a/packages/stack/src/Ports.integration.test.ts b/packages/stack/src/Ports.integration.test.ts index 4e2428373a..e18571ded3 100644 --- a/packages/stack/src/Ports.integration.test.ts +++ b/packages/stack/src/Ports.integration.test.ts @@ -78,6 +78,64 @@ it.live("reports an occupied saved port without moving its assignment", () => ).pipe(Effect.provide(NodeServices.layer)), ); +it.live("allocates an auto port outside a contiguous range that refuses to bind", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped(); + const state = yield* makeTestState(root); + yield* state.save({ + id: "stack", + runtime: "native", + identity: { projectRoot: root, branchContext: "test", stackName: "ports" }, + instances: [], + composition: { members: [], dependencies: [] }, + ports: [], + }); + const ports = yield* makePorts(state); + const reservedBelow = 40000; + const acquired = yield* ports.acquire( + { stackId: "stack", key: "api", host: "127.0.0.1", port: "auto" }, + (host, port) => + port < reservedBelow + ? Effect.fail(new PortError({ key: "api", message: `bind EACCES ${host}:${port}` })) + : Effect.succeed(port), + ); + expect(acquired.port).toBeGreaterThanOrEqual(reservedBelow); + expect((yield* state.read("stack"))?.ports).toEqual([ + { key: "api", host: "127.0.0.1", port: acquired.port }, + ]); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); + +it.live("names the last bind failure when no public port is available", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped(); + const state = yield* makeTestState(root); + yield* state.save({ + id: "stack", + runtime: "native", + identity: { projectRoot: root, branchContext: "test", stackName: "ports" }, + instances: [], + composition: { members: [], dependencies: [] }, + ports: [], + }); + const ports = yield* makePorts(state); + const failure = yield* ports + .acquire({ stackId: "stack", key: "api", host: "127.0.0.1", port: "auto" }, (host, port) => + Effect.fail(new PortError({ key: "api", message: `bind EACCES ${host}:${port}` })), + ) + .pipe(Effect.flip); + expect(failure.message).toContain("No public port is available"); + expect(failure.message).toContain("bind EACCES 127.0.0.1:"); + expect((yield* state.read("stack"))?.ports).toEqual([]); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); + it.live("refuses allocation when another stack has unreadable claims", () => Effect.scoped( Effect.gen(function* () { diff --git a/packages/stack/src/Ports.ts b/packages/stack/src/Ports.ts index 821af3467d..d4de1176e6 100644 --- a/packages/stack/src/Ports.ts +++ b/packages/stack/src/Ports.ts @@ -58,10 +58,12 @@ export const makePorts = (state: State.Interface) => }); const owner = yield* Scope.Scope; - const first = yield* crypto.randomIntBetween(0, 29999); let failures = 0; + let lastFailure: PortError | undefined; for (let attempt = 0; attempt < 30000 && failures < 64; attempt++) { - const port = requested === "auto" ? 20000 + ((first + attempt) % 30000) : requested; + // Each probe is redrawn so a contiguous reserved range cannot exhaust the budget. + const port = + requested === "auto" ? yield* crypto.randomIntBetween(20000, 49999) : requested; if (claimed.has(port)) continue; const result = yield* Effect.uninterruptibleMask((restore) => Effect.gen(function* () { @@ -105,10 +107,15 @@ export const makePorts = (state: State.Interface) => ) return yield* Effect.failCause(result.cause); failures++; + lastFailure = error.value; } return yield* new PortError({ key: request.key, - message: "No public port is available", + message: + lastFailure === undefined + ? "No public port is available" + : `No public port is available: ${lastFailure.message}`, + cause: lastFailure, }); }), ), From 72ae7bd83f17b20fc818ad406b6b0c6987cc367f Mon Sep 17 00:00:00 2001 From: Andrew Valleteau Date: Mon, 21 Sep 2026 16:48:24 +0000 Subject: [PATCH 3/5] fix(stack): scan auto ports deterministically past reserved ranges Drawing a random port per probe cleared reserved ranges only by chance. A stride co-prime with the span visits every port once and keeps consecutive probes far apart instead, so a reserved range narrower than the stride costs at most one probe and can never exhaust the failure budget. The scan starts from an offset derived from the stack's checkout, id, and key, so separate checkouts and keys stay apart while one stack reassigns the same port across runs. Ports no longer needs the platform crypto service, and its callers stop carrying that requirement. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01LcUaThuQfDKPwi9mCE9XeD --- packages/stack/src/HostProcess.ts | 7 ++---- packages/stack/src/Ports.integration.test.ts | 25 +++++++++++++++++++ packages/stack/src/Ports.ts | 26 ++++++++++++++------ 3 files changed, 46 insertions(+), 12 deletions(-) diff --git a/packages/stack/src/HostProcess.ts b/packages/stack/src/HostProcess.ts index 77b49aff1b..51a3031395 100644 --- a/packages/stack/src/HostProcess.ts +++ b/packages/stack/src/HostProcess.ts @@ -56,7 +56,7 @@ export const acquireHost = Effect.fn("HostProcess.acquireHost")(function* ( readonly closeConnections: Effect.Effect; }, HostProcessError | PortError | State.StateError, - Scope.Scope | import("effect").Crypto.Crypto + Scope.Scope > { if ((yield* state.read(stackId)) === undefined) return yield* error("acquire", "Stack is not registered"); @@ -307,10 +307,7 @@ export const launchHost = Effect.fn("HostProcess.launchHost")(function* ( ): Effect.fn.Return< HostEndpoint, HostProcessError | State.StateError, - | Scope.Scope - | HttpClient.HttpClient - | import("effect").Crypto.Crypto - | ChildProcessSpawner.ChildProcessSpawner + Scope.Scope | HttpClient.HttpClient | ChildProcessSpawner.ChildProcessSpawner > { const existing = yield* connectHost(state, options.stackId).pipe( Effect.map(Option.some), diff --git a/packages/stack/src/Ports.integration.test.ts b/packages/stack/src/Ports.integration.test.ts index e18571ded3..e75b6cf89a 100644 --- a/packages/stack/src/Ports.integration.test.ts +++ b/packages/stack/src/Ports.integration.test.ts @@ -109,6 +109,31 @@ it.live("allocates an auto port outside a contiguous range that refuses to bind" ).pipe(Effect.provide(NodeServices.layer)), ); +it.live("reassigns the same auto port after its claim is released", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped(); + const state = yield* makeTestState(root); + yield* state.save({ + id: "stack", + runtime: "native", + identity: { projectRoot: root, branchContext: "test", stackName: "ports" }, + instances: [], + composition: { members: [], dependencies: [] }, + ports: [], + }); + const ports = yield* makePorts(state); + const request = { stackId: "stack", key: "api", host: "127.0.0.1", port: "auto" as const }; + const accept = (_host: string, port: number) => Effect.succeed(port); + const first = yield* ports.acquire(request, accept); + yield* ports.release("stack", "api"); + const again = yield* ports.acquire(request, accept); + expect(again.port).toBe(first.port); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); + it.live("names the last bind failure when no public port is available", () => Effect.scoped( Effect.gen(function* () { diff --git a/packages/stack/src/Ports.ts b/packages/stack/src/Ports.ts index d4de1176e6..3b9f62d273 100644 --- a/packages/stack/src/Ports.ts +++ b/packages/stack/src/Ports.ts @@ -1,6 +1,14 @@ -import { Cause, Crypto, Data, Effect, Exit, Option, Scope } from "effect"; +import { Cause, Data, Effect, Exit, Hash, Option, Scope } from "effect"; import type * as State from "./State.ts"; +const portBase = 20000; +const portSpan = 30000; +/** + * Co-prime with the span, so the scan visits every port once and a reserved range narrower than + * the stride cannot produce consecutive bind failures. + */ +const portStride = 7919; + export class PortError extends Data.TaggedError("PortError")<{ readonly key: string; readonly message: string; @@ -14,11 +22,13 @@ export interface PortRequest { readonly port: number | "auto"; } +/** Spreads the scan across the span so separate checkouts, stacks, and keys start apart. */ +const scanStart = (stack: State.SavedStack, key: string) => + Math.abs(Hash.string(`${stack.identity.projectRoot}:${stack.id}:${key}`)) % portSpan; + /** Coordinates durable public claims while retaining each successfully bound listener. */ export const makePorts = (state: State.Interface) => - Effect.gen(function* () { - const crypto = yield* Crypto.Crypto; - + Effect.sync(() => { const acquire = Effect.fn("Ports.acquire")( ( request: PortRequest, @@ -58,12 +68,14 @@ export const makePorts = (state: State.Interface) => }); const owner = yield* Scope.Scope; + const start = scanStart(stack, request.key); let failures = 0; let lastFailure: PortError | undefined; - for (let attempt = 0; attempt < 30000 && failures < 64; attempt++) { - // Each probe is redrawn so a contiguous reserved range cannot exhaust the budget. + for (let attempt = 0; attempt < portSpan && failures < 64; attempt++) { const port = - requested === "auto" ? yield* crypto.randomIntBetween(20000, 49999) : requested; + requested === "auto" + ? portBase + ((start + attempt * portStride) % portSpan) + : requested; if (claimed.has(port)) continue; const result = yield* Effect.uninterruptibleMask((restore) => Effect.gen(function* () { From e65b21b069295ce2ddd9e038ae80c1dfcbc60347 Mon Sep 17 00:00:00 2001 From: Andrew Valleteau Date: Mon, 21 Sep 2026 16:55:15 +0000 Subject: [PATCH 4/5] fix(stack): restore the documented port scan and keep client exits quiet The auto scan now follows the architecture ADR again: it walks `20000..32767` with stride 257, staying below the Linux ephemeral range and visiting every candidate once, so repeated probes can neither skip a port nor spend the bind allowance twice. Only the start differs from the ADR, deriving from the stack instead of a random draw; the ADR records that. A client that goes away mid-response or resets an established upgrade now fails with the disconnect error rather than a proxy error, so ordinary client exits stay out of the log, and a preparation failure reaches the observation only while the service is still the stopped, registered attempt that prepared it. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01LcUaThuQfDKPwi9mCE9XeD --- ...7-simplified-managed-stack-architecture.md | 3 +- .../stack/src/HttpProxy.integration.test.ts | 32 +++++++++++---- packages/stack/src/HttpProxy.ts | 34 +++++++-------- packages/stack/src/Ports.integration.test.ts | 3 +- packages/stack/src/Ports.ts | 10 ++--- .../stack/src/Service.integration.test.ts | 41 +++++++++++++++++++ packages/stack/src/Service.ts | 16 +++++++- 7 files changed, 103 insertions(+), 36 deletions(-) diff --git a/docs/adr/0017-simplified-managed-stack-architecture.md b/docs/adr/0017-simplified-managed-stack-architecture.md index d8fd09fea1..f9ee11209f 100644 --- a/docs/adr/0017-simplified-managed-stack-architecture.md +++ b/docs/adr/0017-simplified-managed-stack-architecture.md @@ -103,7 +103,8 @@ public and private claims before one state commit; successful public sockets remain held and are adopted directly. Temporary private TCP listeners remain held until commit and then close, so a private workload gap remains possible. Fresh automatic claims draw from -`20000..32767` with a random start and stride `257`, making up to 64 bounded +`20000..32767` with a start derived from the stack's project root, identifier, and +listener key, and stride `257`, making up to 64 bounded `EADDRINUSE`/`EACCES` attempts per newly selected binding while skipping durable sibling claims. Sticky values do not migrate. Failed acquisition preserves the diff --git a/packages/stack/src/HttpProxy.integration.test.ts b/packages/stack/src/HttpProxy.integration.test.ts index df5c991e67..814365921c 100644 --- a/packages/stack/src/HttpProxy.integration.test.ts +++ b/packages/stack/src/HttpProxy.integration.test.ts @@ -334,13 +334,23 @@ it.live("closes an upgrade naming the route and cause when a target cannot wake" ).pipe(Effect.provide(Layer.merge(NodeServices.layer, captureErrors(logs)))); }); -it.live("disconnects a pending upstream response when its client closes", () => - Effect.scoped( +it.live("disconnects a pending upstream response quietly when its client closes", () => { + const logs: Array = []; + return Effect.scoped( Effect.gen(function* () { const backend = createServer(); const address = yield* listen(backend); const proxy = yield* makeHttpProxy({ host: "127.0.0.1", port: 0 }); - yield* proxy.setRoutes([{ id: "pending", prefix: "/", target: Effect.succeed(address) }]); + const released = yield* Deferred.make(); + yield* proxy.setRoutes([ + { + id: "pending", + prefix: "/", + target: Effect.acquireRelease(Effect.succeed(address), () => + Deferred.succeed(released, undefined), + ), + }, + ]); yield* Effect.callback((resume) => { const client = new Socket(); client.on("error", (cause) => @@ -355,12 +365,15 @@ it.live("disconnects a pending upstream response when its client closes", () => ); return Effect.sync(() => client.destroy()); }).pipe(Effect.timeout("5 seconds")); + yield* Deferred.await(released).pipe(Effect.timeout("5 seconds")); + expect(logs).toEqual([]); }), - ).pipe(Effect.provide(NodeServices.layer)), -); + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, captureErrors(logs)))); +}); -it.live("releases a waiting WebSocket target when its client resets the connection", () => - Effect.scoped( +it.live("releases a waiting WebSocket target quietly when its client resets", () => { + const logs: Array = []; + return Effect.scoped( Effect.gen(function* () { const proxy = yield* makeHttpProxy({ host: "127.0.0.1", port: 0 }); const acquiring = yield* Deferred.make(); @@ -393,6 +406,7 @@ it.live("releases a waiting WebSocket target when its client resets the connecti yield* Deferred.await(acquiring); yield* Effect.sync(() => socket.resetAndDestroy()); yield* Deferred.await(released).pipe(Effect.timeout("5 seconds")); + expect(logs).toEqual([]); }), - ).pipe(Effect.provide(NodeServices.layer)), -); + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, captureErrors(logs)))); +}); diff --git a/packages/stack/src/HttpProxy.ts b/packages/stack/src/HttpProxy.ts index 832b63e993..c7ef933e63 100644 --- a/packages/stack/src/HttpProxy.ts +++ b/packages/stack/src/HttpProxy.ts @@ -134,33 +134,35 @@ const proxyRequest = Effect.fn("HttpProxy.proxyRequest")( (request: IncomingMessage, response: ServerResponse, route: HttpRoute) => Effect.gen(function* () { const backend = yield* Effect.raceFirst(route.target, disconnected(request, response)); - yield* Effect.callback((resume) => { + yield* Effect.callback((resume) => { let outgoing: ReturnType | undefined; let incoming: IncomingMessage | undefined; let settled = false; // Error listeners remain until collection because destroy may emit errors asynchronously. const cleanup = () => { - request.off("aborted", onError); + request.off("aborted", onClientGone); response.off("close", onResponseClose); response.off("finish", onFinish); incoming?.off("aborted", onError); }; - const finish = (result: Effect.Effect) => { + const finish = (result: Effect.Effect) => { if (settled) return; settled = true; cleanup(); resume(result); }; - const onError = (cause: Error) => { + const abandon = (result: Effect.Effect) => { if (settled) return; outgoing?.destroy(); incoming?.destroy(); - finish(Effect.fail(errorFor(cause))); + finish(result); }; + const onError = (cause: Error) => abandon(Effect.fail(errorFor(cause))); + const onClientGone = () => abandon(Effect.fail(new HttpProxyDisconnected())); const onFinish = () => finish(Effect.void); const onResponseClose = () => { - if (!response.writableEnded) onError(new Error("client response closed")); + if (!response.writableEnded) onClientGone(); }; outgoing = upstreamRequest( { @@ -186,7 +188,7 @@ const proxyRequest = Effect.fn("HttpProxy.proxyRequest")( }, ); outgoing.on("error", onError); - request.once("aborted", onError); + request.once("aborted", onClientGone); response.once("close", onResponseClose); request.pipe(outgoing); return Effect.sync(() => { @@ -212,7 +214,7 @@ const upgrade = Effect.fn("HttpProxy.upgrade")( }), ); const upstream = yield* connectInterruptibly(backend); - yield* Effect.callback((resume) => { + yield* Effect.callback((resume) => { let settled = false; const cleanup = () => { client.off("close", onClose); @@ -221,25 +223,23 @@ const upgrade = Effect.fn("HttpProxy.upgrade")( client.off("end", onClientEnd); upstream.off("end", onUpstreamEnd); }; - const finish = (result: Effect.Effect) => { + const finish = (result: Effect.Effect) => { if (settled) return; settled = true; cleanup(); resume(result); }; - const onError = (cause: Error) => { + const abandon = (result: Effect.Effect) => { client.destroy(); upstream.destroy(); - finish(Effect.fail(errorFor(cause))); - }; - const onClose = () => { - client.destroy(); - upstream.destroy(); - finish(Effect.void); + finish(result); }; + const onError = (cause: Error) => abandon(Effect.fail(errorFor(cause))); + const onClientGone = () => abandon(Effect.fail(new HttpProxyDisconnected())); + const onClose = () => abandon(Effect.void); const onClientEnd = () => upstream.end(); const onUpstreamEnd = () => client.end(); - client.on("error", onError); + client.on("error", onClientGone); client.once("close", onClose); upstream.on("error", onError); upstream.once("close", onClose); diff --git a/packages/stack/src/Ports.integration.test.ts b/packages/stack/src/Ports.integration.test.ts index e75b6cf89a..87fea5b706 100644 --- a/packages/stack/src/Ports.integration.test.ts +++ b/packages/stack/src/Ports.integration.test.ts @@ -93,7 +93,8 @@ it.live("allocates an auto port outside a contiguous range that refuses to bind" ports: [], }); const ports = yield* makePorts(state); - const reservedBelow = 40000; + // Reserves most of the span contiguously, as Windows excluded ranges do. + const reservedBelow = 30000; const acquired = yield* ports.acquire( { stackId: "stack", key: "api", host: "127.0.0.1", port: "auto" }, (host, port) => diff --git a/packages/stack/src/Ports.ts b/packages/stack/src/Ports.ts index 3b9f62d273..4109ed6424 100644 --- a/packages/stack/src/Ports.ts +++ b/packages/stack/src/Ports.ts @@ -2,12 +2,10 @@ import { Cause, Data, Effect, Exit, Hash, Option, Scope } from "effect"; import type * as State from "./State.ts"; const portBase = 20000; -const portSpan = 30000; -/** - * Co-prime with the span, so the scan visits every port once and a reserved range narrower than - * the stride cannot produce consecutive bind failures. - */ -const portStride = 7919; +/** Stays below the Linux ephemeral range, per the [architecture ADR](../../../docs/adr/0017-simplified-managed-stack-architecture.md). */ +const portSpan = 12768; +/** Co-prime with the span, so the scan visits every port once and steps past reserved ranges. */ +const portStride = 257; export class PortError extends Data.TaggedError("PortError")<{ readonly key: string; diff --git a/packages/stack/src/Service.integration.test.ts b/packages/stack/src/Service.integration.test.ts index 1860059d60..10d881b5f1 100644 --- a/packages/stack/src/Service.integration.test.ts +++ b/packages/stack/src/Service.integration.test.ts @@ -308,6 +308,47 @@ describe("service kernel", () => { ), ); + it.live("keeps a stale preparation failure off a relaunched observation", () => + Effect.scoped( + Effect.gen(function* () { + const plans = yield* Queue.unbounded(); + const staleStarted = yield* Deferred.make(); + const staleGate = yield* Deferred.make(); + const attempts = yield* Ref.make(0); + const service = yield* makeService( + { + ...makeDefinition(plans), + prepare: () => + Effect.gen(function* () { + if ((yield* Ref.updateAndGet(attempts, (value) => value + 1)) > 1) return; + yield* Deferred.succeed(staleStarted, undefined); + yield* Deferred.await(staleGate); + return yield* new ServiceError({ + operation: "prepare", + message: "stale preparation", + }); + }), + }, + { id: "database-stale-failure", config: { version: 17 } }, + ); + const plan = yield* makeRuntimePlan; + yield* Queue.offer(plans, plan); + yield* open(plan.launchGate); + const stale = yield* service.start.pipe(Effect.forkScoped); + yield* Deferred.await(staleStarted); + yield* service.start; + expect((yield* service.get).lifecycle).toBe("running"); + + yield* open(staleGate); + expect(Exit.isFailure(yield* Fiber.await(stale))).toBe(true); + const observation = yield* service.get; + expect(observation.lifecycle).toBe("running"); + expect(observation.error).toBeUndefined(); + yield* stopFixture(service, plan); + }), + ), + ); + it.live("does not prepare concurrent starts more than once", () => Effect.scoped( Effect.gen(function* () { diff --git a/packages/stack/src/Service.ts b/packages/stack/src/Service.ts index d2d8120122..7e64997c7e 100644 --- a/packages/stack/src/Service.ts +++ b/packages/stack/src/Service.ts @@ -412,6 +412,17 @@ export const makeService = ( }); }); + /** A wake has no caller to receive a preparation failure, so an idle observation records it. */ + const recordPreparationFailure = Effect.fn("Service.recordPreparationFailure")(function* ( + expectedRevision: number, + error: ServiceError, + ) { + if ((yield* Ref.get(revision)) !== expectedRevision) return; + yield* SubscriptionRef.update(observations, (value) => + value.registered && value.lifecycle === "stopped" ? { ...value, error } : value, + ); + }); + const startAt = Effect.fn("Service.startAt")(function* ( expectedRevision: number, candidate?: Config, @@ -428,9 +439,10 @@ export const makeService = ( if (existing.lifecycle === "running") return; } const nextConfig = candidate ?? (yield* Ref.get(config)); - // A wake has no caller to receive a preparation failure, so observers record it instead. if (definition.prepare !== undefined) - yield* definition.prepare(nextConfig).pipe(Effect.tapError((error) => update({ error }))); + yield* definition + .prepare(nextConfig) + .pipe(Effect.tapError((error) => recordPreparationFailure(expectedRevision, error))); yield* run( Effect.gen(function* () { const observation = yield* SubscriptionRef.get(observations); From 0756522086f5929161048c3902e3cb32ed5b5572 Mon Sep 17 00:00:00 2001 From: Andrew Valleteau Date: Mon, 21 Sep 2026 18:40:22 +0000 Subject: [PATCH 5/5] fix(stack): settle a proxied request before destroying its upstream Destroying a partially received upstream response emits `aborted` synchronously, and the abort listener re-entered the abandon path while the operation was still unsettled. A client that closed after receiving a body chunk therefore resettled as `HttpProxyError: undefined` and logged, defeating the quiet classification. Both proxy paths now settle and drop their listeners before destroying, and the mid-response disconnect is covered: the previous coverage closed the client before any upstream response arrived. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01LcUaThuQfDKPwi9mCE9XeD --- .../stack/src/HttpProxy.integration.test.ts | 45 +++++++++++++++++++ packages/stack/src/HttpProxy.ts | 6 ++- 2 files changed, 49 insertions(+), 2 deletions(-) diff --git a/packages/stack/src/HttpProxy.integration.test.ts b/packages/stack/src/HttpProxy.integration.test.ts index 814365921c..657a031bb9 100644 --- a/packages/stack/src/HttpProxy.integration.test.ts +++ b/packages/stack/src/HttpProxy.integration.test.ts @@ -371,6 +371,51 @@ it.live("disconnects a pending upstream response quietly when its client closes" ).pipe(Effect.provide(Layer.merge(NodeServices.layer, captureErrors(logs)))); }); +it.live("stays quiet when a client closes after receiving part of the response", () => { + const logs: Array = []; + return Effect.scoped( + Effect.gen(function* () { + const backend = createServer((incoming, outgoing) => { + incoming.resume(); + outgoing.writeHead(200, { "content-type": "application/octet-stream" }); + outgoing.write("first-chunk"); + }); + const address = yield* listen(backend); + const proxy = yield* makeHttpProxy({ host: "127.0.0.1", port: 0 }); + const released = yield* Deferred.make(); + yield* proxy.setRoutes([ + { + id: "streaming", + prefix: "/", + target: Effect.acquireRelease(Effect.succeed(address), () => + Deferred.succeed(released, undefined), + ), + }, + ]); + const client = yield* Effect.acquireRelease( + Effect.sync(() => new Socket()), + (socket) => Effect.sync(() => socket.destroy()), + ); + yield* Effect.callback((resume) => { + // Closing mid-response reaches the client as a reset, which is the disconnect under test. + client.on("error", () => undefined); + const onData = (chunk: Buffer) => { + if (!chunk.includes("first-chunk")) return; + client.destroy(); + resume(Effect.void); + }; + client.on("data", onData); + client.connect(proxy.port, "127.0.0.1", () => + client.write("GET / HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + return Effect.sync(() => client.off("data", onData)); + }).pipe(Effect.timeout("5 seconds")); + yield* Deferred.await(released).pipe(Effect.timeout("5 seconds")); + expect(logs).toEqual([]); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, captureErrors(logs)))); +}); + it.live("releases a waiting WebSocket target quietly when its client resets", () => { const logs: Array = []; return Effect.scoped( diff --git a/packages/stack/src/HttpProxy.ts b/packages/stack/src/HttpProxy.ts index c7ef933e63..1d091e22f5 100644 --- a/packages/stack/src/HttpProxy.ts +++ b/packages/stack/src/HttpProxy.ts @@ -152,11 +152,13 @@ const proxyRequest = Effect.fn("HttpProxy.proxyRequest")( cleanup(); resume(result); }; + // Settling first keeps the outcome: destroying a partial upstream response emits + // `aborted` synchronously, which would otherwise resettle as a proxy failure. const abandon = (result: Effect.Effect) => { if (settled) return; + finish(result); outgoing?.destroy(); incoming?.destroy(); - finish(result); }; const onError = (cause: Error) => abandon(Effect.fail(errorFor(cause))); const onClientGone = () => abandon(Effect.fail(new HttpProxyDisconnected())); @@ -230,9 +232,9 @@ const upgrade = Effect.fn("HttpProxy.upgrade")( resume(result); }; const abandon = (result: Effect.Effect) => { + finish(result); client.destroy(); upstream.destroy(); - finish(result); }; const onError = (cause: Error) => abandon(Effect.fail(errorFor(cause))); const onClientGone = () => abandon(Effect.fail(new HttpProxyDisconnected()));