From 5052aa676dad9dd2cb56dcbe9718902afb41b1b2 Mon Sep 17 00:00:00 2001 From: Julien Goux Date: Thu, 1 Oct 2026 13:15:37 +0200 Subject: [PATCH 1/5] fix(stack): reserve native backend ports below the ephemeral range Native backend ports were reserved by binding port 0, which the OS assigns from its ephemeral range; an unrelated outgoing connection can take that same port as its source port while it sits in TIME_WAIT, causing Bandit/Erlang to fail with eaddrinuse on reopen. Move backend reservation into Ports.ts so it shares the below-ephemeral scan with public auto allocation: it skips ports claimed by any saved stack and, reusing the same accepts-based occupancy check acquire uses for explicit public ports, ports held by a live listener on loopback or wildcard (a loopback-only bind can silently coexist with an existing wildcard listener on macOS, BSD and Windows, masking the real occupant). It never persists the backend port as a claim of its own. A launch attempt that loses its bind to another listener now excludes that port from the next retry's scan, so retries advance instead of repeating the same losing candidate. The release-before-child-bind window and the retry count are unchanged. --- packages/stack/src/Owner.ts | 24 ++-- packages/stack/src/Ports.integration.test.ts | 63 +++++++++- packages/stack/src/Ports.ts | 95 +++++++++++++++ packages/stack/src/services/Catalog.ts | 16 ++- .../services/Functions.integration.test.ts | 6 + .../ProcessRecipe.integration.test.ts | 114 +++++++++++++++++- packages/stack/src/services/ProcessRecipe.ts | 78 +++++------- 7 files changed, 330 insertions(+), 66 deletions(-) diff --git a/packages/stack/src/Owner.ts b/packages/stack/src/Owner.ts index b36d9fcfa1..ed692c05c6 100644 --- a/packages/stack/src/Owner.ts +++ b/packages/stack/src/Owner.ts @@ -288,16 +288,20 @@ const makeOwner = Effect.fn("Owner.make")(function* (options: OwnerOptions) { }); const recipeFor = (creation: ServiceCreation, id: string) => - makeServiceRecipe(creation, { - stackId, - instanceId: id, - project, - root: options.root, - cacheRoot: options.cacheRoot, - runtime, - helpers, - ...(options.hostGateway === undefined ? {} : { hostGateway: options.hostGateway }), - }).pipe(Effect.provideContext(services)); + makeServiceRecipe( + creation, + { + stackId, + instanceId: id, + project, + root: options.root, + cacheRoot: options.cacheRoot, + runtime, + helpers, + ...(options.hostGateway === undefined ? {} : { hostGateway: options.hostGateway }), + }, + options.state, + ).pipe(Effect.provideContext(services)); const persistCreation = (entry: Pick, creation: ServiceCreation) => updateState((current) => ({ diff --git a/packages/stack/src/Ports.integration.test.ts b/packages/stack/src/Ports.integration.test.ts index e5c4f9bb0d..341596977c 100644 --- a/packages/stack/src/Ports.integration.test.ts +++ b/packages/stack/src/Ports.integration.test.ts @@ -1,7 +1,7 @@ import { NodeServices, NodeSocketServer } from "@effect/platform-node"; import { expect, it } from "@effect/vitest"; import { Cause, Context, Effect, Exit, FileSystem, Layer, Option, Ref, Scope } from "effect"; -import { makePorts, PortError } from "./Ports.ts"; +import { makePorts, PortError, reserveNativePort } from "./Ports.ts"; import * as State from "./State.ts"; const makeTestState = (root: string) => @@ -421,3 +421,64 @@ it.live("lets exactly one of two stacks sharing a saved port bind it when both s }), ).pipe(Effect.provide(NodeServices.layer)), ); + +it.live("skips a native backend port claimed by another saved stack", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped(); + const state = yield* makeTestState(root); + + const probeScope = yield* Scope.make(); + const probe = yield* reserveNativePort(Effect.succeed([]), "backend", "pooler").pipe( + Effect.provideService(Scope.Scope, probeScope), + ); + yield* Scope.close(probeScope, Exit.void); + + yield* saveStack(state, root, "claimer", [ + { key: "db/sql", host: "127.0.0.1", port: probe.port }, + ]); + + const reserved = yield* reserveNativePort(state.claims, "backend", "pooler"); + expect(reserved.port).not.toBe(probe.port); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); + +it.live("skips a native backend port a wildcard listener holds", () => + Effect.scoped( + Effect.gen(function* () { + const probeScope = yield* Scope.make(); + const probe = yield* reserveNativePort(Effect.succeed([]), "wildcard-backend", "pooler").pipe( + Effect.provideService(Scope.Scope, probeScope), + ); + yield* Scope.close(probeScope, Exit.void); + + // A wildcard bind is reachable through loopback, so a loopback-only probe would miss it. + yield* bind("0.0.0.0", probe.port); + + const reserved = yield* reserveNativePort(Effect.succeed([]), "wildcard-backend", "pooler"); + expect(reserved.port).not.toBe(probe.port); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); + +it.live("excludes a port a previous attempt lost from the next reservation", () => + Effect.scoped( + Effect.gen(function* () { + const firstScope = yield* Scope.make(); + const first = yield* reserveNativePort(Effect.succeed([]), "retry-backend", "pooler").pipe( + Effect.provideService(Scope.Scope, firstScope), + ); + yield* Scope.close(firstScope, Exit.void); + + const second = yield* reserveNativePort( + Effect.succeed([]), + "retry-backend", + "pooler", + new Set([first.port]), + ); + expect(second.port).not.toBe(first.port); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); diff --git a/packages/stack/src/Ports.ts b/packages/stack/src/Ports.ts index d45216416b..244328ee45 100644 --- a/packages/stack/src/Ports.ts +++ b/packages/stack/src/Ports.ts @@ -88,6 +88,101 @@ const resolveRequest = ( return Effect.succeed({ saved, requested }); }; +export interface NativePortReservation { + readonly port: number; + readonly server: Net.Server; +} + +const closeNativePort = (server: Net.Server): Effect.Effect => + Effect.callback((resume) => { + if (!server.listening) { + resume(Effect.void); + return Effect.void; + } + server.close(() => resume(Effect.void)); + return Effect.void; + }); + +const bindNativePort = ( + key: string, + port: number, +): Effect.Effect => + Effect.callback((resume) => { + const server = Net.createServer((socket) => socket.destroy()); + const onError = (cause: Error) => + resume( + Effect.fail( + new PortError({ key, message: "Unable to reserve native service port", cause }), + ), + ); + server.once("error", onError); + server.listen({ host: "127.0.0.1", port }, () => { + const address = server.address(); + if (address === null || typeof address === "string") { + onError(new Error("Native service port reservation returned no address")); + } else { + resume(Effect.succeed({ port: address.port, server })); + } + }); + return Effect.sync(() => { + server.off("error", onError); + if (server.listening) server.close(); + }); + }); + +const emptyPortSet: ReadonlySet = new Set(); + +/** + * Scans the same below-ephemeral span as `acquire` for a backend port, skipping claimed and + * excluded ports; reuses `loopbackOccupied` so a wildcard listener a loopback-only bind would miss + * on macOS, BSD or Windows still rules out the candidate. Never persists one of its own. + */ +export const reserveNativePort = Effect.fn("Ports.reserveNativePort")( + ( + claims: Effect.Effect, State.StateError>, + stackId: string, + key: string, + excluded: ReadonlySet = emptyPortSet, + ): Effect.Effect => + Effect.acquireRelease( + Effect.gen(function* () { + const stacks = yield* claims.pipe( + Effect.mapError( + (cause) => new PortError({ key, message: "Unable to read port claims", cause }), + ), + ); + const claimed = new Set(stacks.flatMap((stack) => stack.ports.map((claim) => claim.port))); + const start = Math.abs(Hash.string(`${stackId}:${key}`)) % portSpan; + let failures = 0; + let lastFailure: PortError | undefined; + for (let attempt = 0; attempt < portSpan && failures < 64; attempt++) { + const port = portBase + ((start + attempt * portStride) % portSpan); + if (claimed.has(port) || excluded.has(port)) continue; + if (yield* loopbackOccupied(port)) { + failures++; + lastFailure = new PortError({ key, message: `Port ${port} is already in use` }); + continue; + } + const result = yield* Effect.exit(bindNativePort(key, port)); + if (Exit.isSuccess(result)) return result.value; + const error = Cause.findErrorOption(result.cause); + if (Option.isNone(error)) return yield* Effect.failCause(result.cause); + failures++; + lastFailure = error.value; + } + return yield* new PortError({ + key, + message: + lastFailure === undefined + ? "No native service port is available" + : `No native service port is available: ${lastFailure.message}`, + cause: lastFailure, + }); + }), + ({ server }) => closeNativePort(server), + ), +); + /** Claims steer auto allocation away from saved stacks; live listeners and binds decide conflicts for fixed ports. */ export const makePorts = (state: State.Interface, platform: NodeJS.Platform = process.platform) => Effect.sync(() => { diff --git a/packages/stack/src/services/Catalog.ts b/packages/stack/src/services/Catalog.ts index 64e7828df5..4afc4dd68f 100644 --- a/packages/stack/src/services/Catalog.ts +++ b/packages/stack/src/services/Catalog.ts @@ -2,7 +2,9 @@ import { Crypto, Effect, FileSystem, Path, Ref, Schema, Stream } from "effect"; import { ChildProcessSpawner } from "effect/unstable/process"; import { HttpClient } from "effect/unstable/http"; import { makeContainerRuntime } from "../runtime/Container.ts"; +import { reserveNativePort } from "../Ports.ts"; import { ServiceError, type ServiceDefinition } from "../Service.ts"; +import type * as State from "../State.ts"; import { makeDatabase, DatabaseConfig, @@ -298,7 +300,7 @@ const databaseRecipe = ( }); export const makeServiceRecipe = Effect.fn("Catalog.makeServiceRecipe")( - (input: unknown, options: CatalogOptions) => + (input: unknown, options: CatalogOptions, state?: State.Interface) => Effect.gen(function* () { const endpointError = validateEndpointNames(input); if (endpointError !== undefined) return yield* endpointError; @@ -351,7 +353,17 @@ export const makeServiceRecipe = Effect.fn("Catalog.makeServiceRecipe")( imageMirrors: slimImageMirrors, ...(options.hostGateway === undefined ? {} : { hostGateway: options.hostGateway }), }); - const deps: ProcessDependencies = { fs, path, crypto, client, spawner, container }; + const claims = state?.claims ?? Effect.succeed([]); + const deps: ProcessDependencies = { + fs, + path, + crypto, + client, + spawner, + container, + reserveNativePort: (stackId, key, excluded) => + reserveNativePort(claims, stackId, key, excluded), + }; switch (creation.service) { case "rest": return catalogRecipe( diff --git a/packages/stack/src/services/Functions.integration.test.ts b/packages/stack/src/services/Functions.integration.test.ts index 36042572d6..72b1f831b1 100644 --- a/packages/stack/src/services/Functions.integration.test.ts +++ b/packages/stack/src/services/Functions.integration.test.ts @@ -17,10 +17,15 @@ import { import { HttpClient, HttpClientRequest } from "effect/unstable/http"; import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; import { ContainerError, type ContainerRuntime } from "../runtime/Container.ts"; +import { reserveNativePort } from "../Ports.ts"; import { makeService } from "../Service.ts"; import { makeServiceRecipe } from "./Catalog.ts"; import * as Functions from "./Functions.ts"; +// No saved stacks to consult; this test never launches the native backend it configures. +const testReserveNativePort = (stackId: string, key: string, excluded: ReadonlySet) => + reserveNativePort(Effect.succeed([]), stackId, key, excluded); + const options = (root: string) => ({ stackId: "catalog-functions", instanceId: "instance", @@ -552,6 +557,7 @@ it.effect("passes POSIX project paths to a docker Functions container from a Win client: yield* HttpClient.HttpClient, spawner: yield* ChildProcessSpawner.ChildProcessSpawner, container, + reserveNativePort: testReserveNativePort, }, ); const scope = yield* Scope.make(); diff --git a/packages/stack/src/services/ProcessRecipe.integration.test.ts b/packages/stack/src/services/ProcessRecipe.integration.test.ts index 9f14cc70ce..5b1aa0f7c7 100644 --- a/packages/stack/src/services/ProcessRecipe.integration.test.ts +++ b/packages/stack/src/services/ProcessRecipe.integration.test.ts @@ -25,6 +25,7 @@ import * as Net from "node:net"; // oxlint-disable-next-line effecttsgo/node-builtin-import -- the collision fixture owns a local HTTP listener. import * as NodeHttp from "node:http"; import { catalogPins, resolveArtifact, type ServiceKind } from "../Artifacts.ts"; +import { reserveNativePort } from "../Ports.ts"; import { makeArtifactStore, type ArtifactRequest, @@ -98,6 +99,10 @@ const isPortOccupied = (port: number): Effect.Effect => }); }); +// No saved stacks to consult outside the claim-interaction test below. +const testReserveNativePort = (stackId: string, key: string, excluded: ReadonlySet) => + reserveNativePort(Effect.succeed([]), stackId, key, excluded); + describe("ProcessRecipe launch cleanup", () => { for (const scenario of [ "partial launch", @@ -188,6 +193,7 @@ describe("ProcessRecipe launch cleanup", () => { client: yield* HttpClient.HttpClient, spawner: yield* ChildProcessSpawner.ChildProcessSpawner, container, + reserveNativePort: testReserveNativePort, } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, options, dependencies, spec); const service = yield* makeService(recipe.definition, { @@ -275,6 +281,7 @@ describe("ProcessRecipe launch cleanup", () => { client, spawner, container: undefined, + reserveNativePort: testReserveNativePort, } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, nativeOptions, dependencies, nativeSpec); const service = yield* makeService(recipe.definition, { @@ -385,6 +392,7 @@ describe("ProcessRecipe launch cleanup", () => { client, spawner, container: undefined, + reserveNativePort: testReserveNativePort, } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, nativeOptions, dependencies, nativeSpec); const service = yield* makeService(recipe.definition, { @@ -458,6 +466,7 @@ const realtimeService = Effect.fn(function* (container: ContainerRuntime) { client: yield* HttpClient.HttpClient, spawner: yield* ChildProcessSpawner.ChildProcessSpawner, container, + reserveNativePort: testReserveNativePort, }, Realtime.makeSpec(), ); @@ -655,7 +664,15 @@ const nativeRestRecipe = Effect.fn(function* ( runtime: "native", platform: { os: process.platform, arch: process.arch }, }, - { fs, path, crypto, client, spawner, container: undefined }, + { + fs, + path, + crypto, + client, + spawner, + container: undefined, + reserveNativePort: testReserveNativePort, + }, { ...spec, env: (_creation, endpoints) => @@ -924,7 +941,15 @@ describe("process recipe startup", () => { runtime: "native", platform: { os: process.platform, arch: process.arch }, }, - { fs, path, crypto, client, spawner: interceptingSpawner, container: undefined }, + { + fs, + path, + crypto, + client, + spawner: interceptingSpawner, + container: undefined, + reserveNativePort: testReserveNativePort, + }, Pooler.makeSpec(), ); if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); @@ -957,6 +982,69 @@ describe("process recipe startup", () => { ).pipe(Effect.provide(platform)), ); + it.live( + "reserves a native backend port from the below-ephemeral range, never an OS-assigned one", + () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const crypto = yield* Crypto.Crypto; + const client = yield* HttpClient.HttpClient; + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "process-recipe-port-range-" }); + const cacheRoot = path.join(root, "cache"); + yield* nativePoolerArtifact(cacheRoot); + const creation: Pooler.Creation = { + service: "pooler", + config: { + databaseUrl: "postgresql://postgres:postgres@127.0.0.1:5432/postgres", + jwtSecret: "pooler-port-range-test-secret-with-more-than-32-characters", + tenant: "port-range-test", + poolMode: "transaction", + }, + }; + const recipe = yield* makeProcessRecipe( + creation, + { + stackId: "process-recipe-port-range", + instanceId: "instance", + root, + cacheRoot, + runtime: "native", + platform: { os: process.platform, arch: process.arch }, + }, + { + fs, + path, + crypto, + client, + spawner, + container: undefined, + reserveNativePort: testReserveNativePort, + }, + Pooler.makeSpec(), + ); + if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); + const scope = yield* Scope.fork(yield* Effect.scope, "sequential"); + const runtime = yield* recipe.definition.launch({ + id: "pooler", + config: creation, + scope, + }); + yield* runtime.health; + const endpoints = yield* Ref.get(recipe.endpoints); + const endpoint = endpoints.get("http"); + // A backend chosen from this range (Ports.ts's below-ephemeral 20000..32767 span) + // can never coincide with an OS-auto-assigned ephemeral port, so releasing the probe + // before the child binds cannot race an unrelated outgoing connection for the number. + expect(endpoint?.port).toBeGreaterThanOrEqual(20000); + expect(endpoint?.port).toBeLessThan(32768); + yield* runtime.stop; + }), + ).pipe(Effect.provide(platform)), + ); + it.live("stops after three consecutive native Pooler port collisions", () => Effect.scoped( Effect.gen(function* () { @@ -998,6 +1086,7 @@ describe("process recipe startup", () => { client, spawner: countingSpawner(spawner, mainLaunches, startupLaunches, true), container: undefined, + reserveNativePort: testReserveNativePort, }, Pooler.makeSpec(), ); @@ -1125,7 +1214,15 @@ describe("process recipe startup", () => { runtime: "native", platform: { os: process.platform, arch: process.arch }, }, - { fs, path, crypto, client, spawner: deadlineSpawner, container: undefined }, + { + fs, + path, + crypto, + client, + spawner: deadlineSpawner, + container: undefined, + reserveNativePort: testReserveNativePort, + }, Pooler.makeSpec(), ); if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); @@ -1187,6 +1284,7 @@ describe("process recipe startup", () => { client, spawner: countingSpawner(spawner, mainLaunches, startupLaunches), container: undefined, + reserveNativePort: testReserveNativePort, }, Pooler.makeSpec(), ); @@ -1284,7 +1382,15 @@ describe("process recipe startup", () => { runtime: "native", platform: { os: process.platform, arch: process.arch }, }, - { fs, path, crypto, client, spawner: failingSpawner, container: undefined }, + { + fs, + path, + crypto, + client, + spawner: failingSpawner, + container: undefined, + reserveNativePort: testReserveNativePort, + }, Pooler.makeSpec(), ); if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); diff --git a/packages/stack/src/services/ProcessRecipe.ts b/packages/stack/src/services/ProcessRecipe.ts index 73a81799b9..9abab73904 100644 --- a/packages/stack/src/services/ProcessRecipe.ts +++ b/packages/stack/src/services/ProcessRecipe.ts @@ -19,9 +19,8 @@ import { import { ChildProcessSpawner } from "effect/unstable/process"; import type { ChildProcessSpawner as ChildProcessSpawnerService } from "effect/unstable/process/ChildProcessSpawner"; import { HttpClient, HttpClientRequest } from "effect/unstable/http"; -import * as Net from "node:net"; import { prepareNativeArtifact, resolveArtifact, type ServiceKind } from "../Artifacts.ts"; -import { accepts } from "../Ports.ts"; +import { accepts, type NativePortReservation, type PortError } from "../Ports.ts"; import { type ContainerError, type ContainerProcess, @@ -161,6 +160,11 @@ export interface ProcessDependencies { readonly client: HttpClient.HttpClient; readonly spawner: ChildProcessSpawnerService["Service"]; readonly container: ContainerRuntime | undefined; + readonly reserveNativePort: ( + stackId: string, + key: string, + excluded: ReadonlySet, + ) => Effect.Effect; } const serviceError = mapToServiceError; @@ -189,50 +193,6 @@ const runtimeFromNative = (process: NativeProcess): RuntimeSession => ({ remove: Effect.void, }); -interface NativePortReservation { - readonly port: number; - readonly server: Net.Server; -} - -const closeNativePort = (server: Net.Server): Effect.Effect => - Effect.callback((resume) => { - if (!server.listening) { - resume(Effect.void); - return Effect.void; - } - server.close(() => resume(Effect.void)); - return Effect.void; - }); - -const reserveNativePort = Effect.fn("ProcessRecipe.reserveNativePort")( - (requested: number): Effect.Effect => - Effect.acquireRelease( - Effect.callback((resume) => { - const server = Net.createServer((socket) => socket.destroy()); - const onError = (cause: Error) => - resume( - Effect.fail( - catalogError("launch", "Unable to reserve native service port", undefined, cause), - ), - ); - server.once("error", onError); - server.listen({ host: "127.0.0.1", port: requested }, () => { - const address = server.address(); - if (address === null || typeof address === "string") { - onError(new Error("Native service port reservation returned no address")); - } else { - resume(Effect.succeed({ port: address.port, server })); - } - }); - return Effect.sync(() => { - server.off("error", onError); - if (server.listening) server.close(); - }); - }), - ({ server }) => closeNativePort(server), - ), -); - const startupTimeoutSeconds = 60; const nativeLaunchAttempts = 3; const probeTimeout = Duration.seconds(10); @@ -479,13 +439,24 @@ export const makeProcessRecipe = return yield* serviceError("launch", "Artifact was not prepared"); if (artifactRoot === undefined) return yield* serviceError("launch", "Artifact root was not prepared"); + const keyFor = (name: string) => `${options.instanceId}:${context.id}:${name}`; + // A port a prior attempt lost stays excluded so a retry advances instead of repeating it. + const excludedByKey = yield* Ref.make>>(new Map()); const reserveEndpoints = Effect.fn("ProcessRecipe.reserveEndpoints")(function* ( parent: Scope.Closeable, ) { const portScope = yield* Scope.fork(parent, "sequential"); - const reservations = yield* Effect.forEach(portNames, () => reserveNativePort(0), { - concurrency: 1, - }).pipe( + const excluded = yield* Ref.get(excludedByKey); + const reservations = yield* Effect.forEach( + portNames, + ([name]) => + deps.reserveNativePort( + options.stackId, + keyFor(name), + excluded.get(keyFor(name)) ?? new Set(), + ), + { concurrency: 1 }, + ).pipe( Scope.provide(portScope), Effect.mapError((cause) => serviceError("launch", cause)), ); @@ -701,6 +672,15 @@ export const makeProcessRecipe = (!(yield* Deferred.isDone(attempt.output.bindReady)) && (yield* anotherListenerHolds(attempt.selected)))); if (!collided) return yield* settleFailure(failure); + yield* Ref.update(excludedByKey, (current) => { + const next = new Map(current); + for (const [name, endpoint] of attempt.selected) { + if (endpoint.kind !== "tcp") continue; + const key = keyFor(name); + next.set(key, new Set([...(next.get(key) ?? []), endpoint.port])); + } + return next; + }); return yield* new NativePortCollision({ failure: serviceError( "launch", From 0ebeab4832d24f62807183b3be830e0d478ffe95 Mon Sep 17 00:00:00 2001 From: Julien Goux Date: Thu, 1 Oct 2026 15:08:00 +0200 Subject: [PATCH 2/5] fix(stack): check managed port occupancy for public auto allocation Review on the backend-port-reservation fix found a regression and two robustness gaps it introduced: - Public auto allocation (Ports.ts acquire) bound fresh candidates without checking occupancy first. On macOS, BSD and Windows a container stack's wildcard public bind can succeed on a port where another stack's native backend already listens on loopback, so loopback traffic reached the backend instead of the proxy. Auto allocation now reuses the same accepts-based occupancy check (loopbackOccupied) backend reservation uses, so one mechanism owns whether a managed port is free for both allocators. - Backend port reservation scanned from a start derived from hash(stackId:key). Backend ports aren't persisted, so stability bought nothing, and a reopened stack retried its previous incarnation's port, which can still be in server-side TIME_WAIT: the Node probe sets SO_REUSEADDR and passes, while a child process without it fails, consuming a retry. The start is now drawn from the injected Crypto service per reservation; claim and exclusion skipping are unchanged. - A reservation batch re-read state.claims once per endpoint instead of once per attempt. reserveEndpoints now reads it once and shares the snapshot across the batch; a later retry attempt still re-reads it. - A service's endpoints now reserve concurrently instead of serially, since each reservation holds its bound probe listener until the batch releases, so two endpoints can't settle on the same port. --- packages/stack/src/Ports.integration.test.ts | 47 ++++++++++++++----- packages/stack/src/Ports.ts | 33 ++++++++----- packages/stack/src/services/Catalog.ts | 8 ++-- .../services/Functions.integration.test.ts | 13 ++++- .../ProcessRecipe.integration.test.ts | 27 ++++++++++- packages/stack/src/services/ProcessRecipe.ts | 18 ++++--- 6 files changed, 107 insertions(+), 39 deletions(-) diff --git a/packages/stack/src/Ports.integration.test.ts b/packages/stack/src/Ports.integration.test.ts index 341596977c..7ee87c4cfb 100644 --- a/packages/stack/src/Ports.integration.test.ts +++ b/packages/stack/src/Ports.integration.test.ts @@ -422,6 +422,9 @@ it.live("lets exactly one of two stacks sharing a saved port bind it when both s ).pipe(Effect.provide(NodeServices.layer)), ); +// Fixed so a test can force two reservations to the same candidate; production uses randomPortSpanStart. +const fixedStart = Effect.succeed(0); + it.live("skips a native backend port claimed by another saved stack", () => Effect.scoped( Effect.gen(function* () { @@ -430,7 +433,7 @@ it.live("skips a native backend port claimed by another saved stack", () => const state = yield* makeTestState(root); const probeScope = yield* Scope.make(); - const probe = yield* reserveNativePort(Effect.succeed([]), "backend", "pooler").pipe( + const probe = yield* reserveNativePort([], "pooler", fixedStart).pipe( Effect.provideService(Scope.Scope, probeScope), ); yield* Scope.close(probeScope, Exit.void); @@ -439,7 +442,7 @@ it.live("skips a native backend port claimed by another saved stack", () => { key: "db/sql", host: "127.0.0.1", port: probe.port }, ]); - const reserved = yield* reserveNativePort(state.claims, "backend", "pooler"); + const reserved = yield* reserveNativePort(yield* state.claims, "pooler", fixedStart); expect(reserved.port).not.toBe(probe.port); }), ).pipe(Effect.provide(NodeServices.layer)), @@ -449,7 +452,7 @@ it.live("skips a native backend port a wildcard listener holds", () => Effect.scoped( Effect.gen(function* () { const probeScope = yield* Scope.make(); - const probe = yield* reserveNativePort(Effect.succeed([]), "wildcard-backend", "pooler").pipe( + const probe = yield* reserveNativePort([], "pooler", fixedStart).pipe( Effect.provideService(Scope.Scope, probeScope), ); yield* Scope.close(probeScope, Exit.void); @@ -457,7 +460,7 @@ it.live("skips a native backend port a wildcard listener holds", () => // A wildcard bind is reachable through loopback, so a loopback-only probe would miss it. yield* bind("0.0.0.0", probe.port); - const reserved = yield* reserveNativePort(Effect.succeed([]), "wildcard-backend", "pooler"); + const reserved = yield* reserveNativePort([], "pooler", fixedStart); expect(reserved.port).not.toBe(probe.port); }), ).pipe(Effect.provide(NodeServices.layer)), @@ -467,18 +470,40 @@ it.live("excludes a port a previous attempt lost from the next reservation", () Effect.scoped( Effect.gen(function* () { const firstScope = yield* Scope.make(); - const first = yield* reserveNativePort(Effect.succeed([]), "retry-backend", "pooler").pipe( + const first = yield* reserveNativePort([], "pooler", fixedStart).pipe( Effect.provideService(Scope.Scope, firstScope), ); yield* Scope.close(firstScope, Exit.void); - const second = yield* reserveNativePort( - Effect.succeed([]), - "retry-backend", - "pooler", - new Set([first.port]), - ); + const second = yield* reserveNativePort([], "pooler", fixedStart, new Set([first.port])); expect(second.port).not.toBe(first.port); }), ).pipe(Effect.provide(NodeServices.layer)), ); + +it.live("skips a public auto candidate a loopback listener already holds", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped(); + const state = yield* makeTestState(root); + yield* saveStack(state, root, "stack"); + const ports = yield* makePorts(state); + // A container-runtime stack's bind to the wildcard host would otherwise succeed here too. + const request = { stackId: "stack", key: "api", host: "0.0.0.0", port: "auto" as const }; + const accept = (_host: string, port: number) => Effect.succeed(port); + + const probeScope = yield* Scope.make(); + const probe = yield* ports + .acquire(request, accept) + .pipe(Effect.provideService(Scope.Scope, probeScope)); + yield* Scope.close(probeScope, Exit.void); + yield* ports.release("stack", "api"); + + yield* bind("127.0.0.1", probe.port); + + const acquired = yield* ports.acquire(request, accept); + expect(acquired.port).not.toBe(probe.port); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); diff --git a/packages/stack/src/Ports.ts b/packages/stack/src/Ports.ts index 244328ee45..6150b0f796 100644 --- a/packages/stack/src/Ports.ts +++ b/packages/stack/src/Ports.ts @@ -1,4 +1,4 @@ -import { Cause, Data, Effect, Exit, Hash, Option, Scope } from "effect"; +import { Cause, Crypto, Data, Effect, Exit, Hash, Option, Scope } from "effect"; import * as Net from "node:net"; import type * as State from "./State.ts"; @@ -132,27 +132,26 @@ const bindNativePort = ( const emptyPortSet: ReadonlySet = new Set(); +/** Backend ports aren't persisted, so a random start buys nothing by staying stable across reopens. */ +export const randomPortSpanStart = (crypto: Crypto.Crypto): Effect.Effect => + crypto.randomIntBetween(0, portSpan, { halfOpen: true }); + /** - * Scans the same below-ephemeral span as `acquire` for a backend port, skipping claimed and - * excluded ports; reuses `loopbackOccupied` so a wildcard listener a loopback-only bind would miss - * on macOS, BSD or Windows still rules out the candidate. Never persists one of its own. + * Scans the below-ephemeral span for a backend port, skipping claimed and excluded ports; reuses + * `loopbackOccupied` so a wildcard listener a loopback-only bind would miss on macOS, BSD or + * Windows still rules out the candidate. Never persists one of its own. */ export const reserveNativePort = Effect.fn("Ports.reserveNativePort")( ( - claims: Effect.Effect, State.StateError>, - stackId: string, + claims: ReadonlyArray, key: string, + randomStart: Effect.Effect, excluded: ReadonlySet = emptyPortSet, ): Effect.Effect => Effect.acquireRelease( Effect.gen(function* () { - const stacks = yield* claims.pipe( - Effect.mapError( - (cause) => new PortError({ key, message: "Unable to read port claims", cause }), - ), - ); - const claimed = new Set(stacks.flatMap((stack) => stack.ports.map((claim) => claim.port))); - const start = Math.abs(Hash.string(`${stackId}:${key}`)) % portSpan; + const claimed = new Set(claims.flatMap((stack) => stack.ports.map((claim) => claim.port))); + const start = yield* randomStart; let failures = 0; let lastFailure: PortError | undefined; for (let attempt = 0; attempt < portSpan && failures < 64; attempt++) { @@ -269,6 +268,14 @@ export const makePorts = (state: State.Interface, platform: NodeJS.Platform = pr ? portBase + ((start + attempt * portStride) % portSpan) : requested; if (requested === "auto" && claimed.has(port)) continue; + if (requested === "auto" && (yield* loopbackOccupied(port))) { + failures++; + lastFailure = new PortError({ + key: request.key, + message: `Port ${port} is already in use`, + }); + continue; + } const result = yield* Effect.uninterruptibleMask((restore) => Effect.gen(function* () { const scope = yield* Scope.fork(owner, "sequential"); diff --git a/packages/stack/src/services/Catalog.ts b/packages/stack/src/services/Catalog.ts index 4afc4dd68f..748f8ff349 100644 --- a/packages/stack/src/services/Catalog.ts +++ b/packages/stack/src/services/Catalog.ts @@ -2,7 +2,7 @@ import { Crypto, Effect, FileSystem, Path, Ref, Schema, Stream } from "effect"; import { ChildProcessSpawner } from "effect/unstable/process"; import { HttpClient } from "effect/unstable/http"; import { makeContainerRuntime } from "../runtime/Container.ts"; -import { reserveNativePort } from "../Ports.ts"; +import { randomPortSpanStart, reserveNativePort } from "../Ports.ts"; import { ServiceError, type ServiceDefinition } from "../Service.ts"; import type * as State from "../State.ts"; import { @@ -353,7 +353,6 @@ export const makeServiceRecipe = Effect.fn("Catalog.makeServiceRecipe")( imageMirrors: slimImageMirrors, ...(options.hostGateway === undefined ? {} : { hostGateway: options.hostGateway }), }); - const claims = state?.claims ?? Effect.succeed([]); const deps: ProcessDependencies = { fs, path, @@ -361,8 +360,9 @@ export const makeServiceRecipe = Effect.fn("Catalog.makeServiceRecipe")( client, spawner, container, - reserveNativePort: (stackId, key, excluded) => - reserveNativePort(claims, stackId, key, excluded), + readPortClaims: state?.claims ?? Effect.succeed([]), + reserveNativePort: (key, claims, excluded) => + reserveNativePort(claims, key, randomPortSpanStart(crypto), excluded), }; switch (creation.service) { case "rest": diff --git a/packages/stack/src/services/Functions.integration.test.ts b/packages/stack/src/services/Functions.integration.test.ts index 72b1f831b1..9460688d14 100644 --- a/packages/stack/src/services/Functions.integration.test.ts +++ b/packages/stack/src/services/Functions.integration.test.ts @@ -16,15 +16,23 @@ import { } from "effect"; import { HttpClient, HttpClientRequest } from "effect/unstable/http"; import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; +// oxlint-disable-next-line effecttsgo/node-builtin-import -- test-only candidate start, not the injected Crypto service. +import { randomInt as nodeRandomInt } from "node:crypto"; import { ContainerError, type ContainerRuntime } from "../runtime/Container.ts"; import { reserveNativePort } from "../Ports.ts"; import { makeService } from "../Service.ts"; +import type * as State from "../State.ts"; import { makeServiceRecipe } from "./Catalog.ts"; import * as Functions from "./Functions.ts"; // No saved stacks to consult; this test never launches the native backend it configures. -const testReserveNativePort = (stackId: string, key: string, excluded: ReadonlySet) => - reserveNativePort(Effect.succeed([]), stackId, key, excluded); +const testReadPortClaims = Effect.succeed([]); +const testRandomStart = Effect.sync(() => nodeRandomInt(0, 1_000_000)); +const testReserveNativePort = ( + key: string, + claims: ReadonlyArray, + excluded: ReadonlySet, +) => reserveNativePort(claims, key, testRandomStart, excluded); const options = (root: string) => ({ stackId: "catalog-functions", @@ -557,6 +565,7 @@ it.effect("passes POSIX project paths to a docker Functions container from a Win client: yield* HttpClient.HttpClient, spawner: yield* ChildProcessSpawner.ChildProcessSpawner, container, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, }, ); diff --git a/packages/stack/src/services/ProcessRecipe.integration.test.ts b/packages/stack/src/services/ProcessRecipe.integration.test.ts index 5b1aa0f7c7..7d0669eb28 100644 --- a/packages/stack/src/services/ProcessRecipe.integration.test.ts +++ b/packages/stack/src/services/ProcessRecipe.integration.test.ts @@ -24,8 +24,11 @@ import { systemError } from "effect/PlatformError"; import * as Net from "node:net"; // oxlint-disable-next-line effecttsgo/node-builtin-import -- the collision fixture owns a local HTTP listener. import * as NodeHttp from "node:http"; +// oxlint-disable-next-line effecttsgo/node-builtin-import -- test-only candidate start, not the injected Crypto service. +import { randomInt as nodeRandomInt } from "node:crypto"; import { catalogPins, resolveArtifact, type ServiceKind } from "../Artifacts.ts"; import { reserveNativePort } from "../Ports.ts"; +import type * as State from "../State.ts"; import { makeArtifactStore, type ArtifactRequest, @@ -100,8 +103,13 @@ const isPortOccupied = (port: number): Effect.Effect => }); // No saved stacks to consult outside the claim-interaction test below. -const testReserveNativePort = (stackId: string, key: string, excluded: ReadonlySet) => - reserveNativePort(Effect.succeed([]), stackId, key, excluded); +const testReadPortClaims = Effect.succeed([]); +const testRandomStart = Effect.sync(() => nodeRandomInt(0, 1_000_000)); +const testReserveNativePort = ( + key: string, + claims: ReadonlyArray, + excluded: ReadonlySet, +) => reserveNativePort(claims, key, testRandomStart, excluded); describe("ProcessRecipe launch cleanup", () => { for (const scenario of [ @@ -193,6 +201,7 @@ describe("ProcessRecipe launch cleanup", () => { client: yield* HttpClient.HttpClient, spawner: yield* ChildProcessSpawner.ChildProcessSpawner, container, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, options, dependencies, spec); @@ -281,6 +290,7 @@ describe("ProcessRecipe launch cleanup", () => { client, spawner, container: undefined, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, nativeOptions, dependencies, nativeSpec); @@ -392,6 +402,7 @@ describe("ProcessRecipe launch cleanup", () => { client, spawner, container: undefined, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, nativeOptions, dependencies, nativeSpec); @@ -466,6 +477,7 @@ const realtimeService = Effect.fn(function* (container: ContainerRuntime) { client: yield* HttpClient.HttpClient, spawner: yield* ChildProcessSpawner.ChildProcessSpawner, container, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, }, Realtime.makeSpec(), @@ -671,6 +683,7 @@ const nativeRestRecipe = Effect.fn(function* ( client, spawner, container: undefined, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, }, { @@ -948,6 +961,7 @@ describe("process recipe startup", () => { client, spawner: interceptingSpawner, container: undefined, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, }, Pooler.makeSpec(), @@ -1021,6 +1035,7 @@ describe("process recipe startup", () => { client, spawner, container: undefined, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, }, Pooler.makeSpec(), @@ -1040,6 +1055,10 @@ describe("process recipe startup", () => { // before the child binds cannot race an unrelated outgoing connection for the number. expect(endpoint?.port).toBeGreaterThanOrEqual(20000); expect(endpoint?.port).toBeLessThan(32768); + // Pooler reserves "http" and "sql" concurrently; they must never settle on the same port. + const sql = endpoints.get("sql"); + expect(sql?.port).toBeDefined(); + expect(sql?.port).not.toBe(endpoint?.port); yield* runtime.stop; }), ).pipe(Effect.provide(platform)), @@ -1086,6 +1105,7 @@ describe("process recipe startup", () => { client, spawner: countingSpawner(spawner, mainLaunches, startupLaunches, true), container: undefined, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, }, Pooler.makeSpec(), @@ -1221,6 +1241,7 @@ describe("process recipe startup", () => { client, spawner: deadlineSpawner, container: undefined, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, }, Pooler.makeSpec(), @@ -1284,6 +1305,7 @@ describe("process recipe startup", () => { client, spawner: countingSpawner(spawner, mainLaunches, startupLaunches), container: undefined, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, }, Pooler.makeSpec(), @@ -1389,6 +1411,7 @@ describe("process recipe startup", () => { client, spawner: failingSpawner, container: undefined, + readPortClaims: testReadPortClaims, reserveNativePort: testReserveNativePort, }, Pooler.makeSpec(), diff --git a/packages/stack/src/services/ProcessRecipe.ts b/packages/stack/src/services/ProcessRecipe.ts index 9abab73904..8757974905 100644 --- a/packages/stack/src/services/ProcessRecipe.ts +++ b/packages/stack/src/services/ProcessRecipe.ts @@ -21,6 +21,7 @@ import type { ChildProcessSpawner as ChildProcessSpawnerService } from "effect/u import { HttpClient, HttpClientRequest } from "effect/unstable/http"; import { prepareNativeArtifact, resolveArtifact, type ServiceKind } from "../Artifacts.ts"; import { accepts, type NativePortReservation, type PortError } from "../Ports.ts"; +import type * as State from "../State.ts"; import { type ContainerError, type ContainerProcess, @@ -160,9 +161,10 @@ export interface ProcessDependencies { readonly client: HttpClient.HttpClient; readonly spawner: ChildProcessSpawnerService["Service"]; readonly container: ContainerRuntime | undefined; + readonly readPortClaims: Effect.Effect, State.StateError>; readonly reserveNativePort: ( - stackId: string, key: string, + claims: ReadonlyArray, excluded: ReadonlySet, ) => Effect.Effect; } @@ -447,15 +449,17 @@ export const makeProcessRecipe = ) { const portScope = yield* Scope.fork(parent, "sequential"); const excluded = yield* Ref.get(excludedByKey); + // Read once and share across this batch; a later retry attempt re-reads it. + const claims = yield* deps.readPortClaims.pipe( + Effect.mapError((cause) => serviceError("launch", cause)), + ); const reservations = yield* Effect.forEach( portNames, ([name]) => - deps.reserveNativePort( - options.stackId, - keyFor(name), - excluded.get(keyFor(name)) ?? new Set(), - ), - { concurrency: 1 }, + deps.reserveNativePort(keyFor(name), claims, excluded.get(keyFor(name)) ?? new Set()), + // Each reservation holds its bound probe listener until the batch releases, so two + // endpoints never settle on the same port even when reserved concurrently. + { concurrency: "unbounded" }, ).pipe( Scope.provide(portScope), Effect.mapError((cause) => serviceError("launch", cause)), From b052de29520bb0d654a62decf4e9d5305ead5d93 Mon Sep 17 00:00:00 2001 From: Julien Goux Date: Fri, 2 Oct 2026 18:36:59 +0200 Subject: [PATCH 3/5] fix(stack): address port reservation review findings Review on #6941 flagged two flaky fixtures that released a reserved port before rebinding it, letting another process grab the number in the gap; both now keep a real listener held instead. The native backend port-contention test shared independent random scan starts across its concurrent reservations, so it never actually forced them to collide; it now forces one shared start. Both integration test files also duplicated the random scan-start generation that Ports.ts already exports as randomPortSpanStart. --- packages/stack/src/Ports.integration.test.ts | 36 ++++++++++--------- packages/stack/src/Ports.ts | 4 +-- .../services/Functions.integration.test.ts | 12 +++---- .../ProcessRecipe.integration.test.ts | 25 ++++++++----- 4 files changed, 43 insertions(+), 34 deletions(-) diff --git a/packages/stack/src/Ports.integration.test.ts b/packages/stack/src/Ports.integration.test.ts index 7ee87c4cfb..c7f4df0938 100644 --- a/packages/stack/src/Ports.integration.test.ts +++ b/packages/stack/src/Ports.integration.test.ts @@ -1,7 +1,7 @@ import { NodeServices, NodeSocketServer } from "@effect/platform-node"; import { expect, it } from "@effect/vitest"; import { Cause, Context, Effect, Exit, FileSystem, Layer, Option, Ref, Scope } from "effect"; -import { makePorts, PortError, reserveNativePort } from "./Ports.ts"; +import { makePorts, portBase, portSpan, PortError, reserveNativePort } from "./Ports.ts"; import * as State from "./State.ts"; const makeTestState = (root: string) => @@ -425,6 +425,17 @@ it.live("lets exactly one of two stacks sharing a saved port bind it when both s // Fixed so a test can force two reservations to the same candidate; production uses randomPortSpanStart. const fixedStart = Effect.succeed(0); +// Binds a real listener directly in the native-reservation span, retrying past occupied +// candidates, instead of reserving then releasing a port that something else could grab meanwhile. +const bindBlockingPort = (host: string) => + Effect.gen(function* () { + for (let offset = 0; offset < portSpan; offset++) { + const attempt = yield* Effect.exit(bind(host, portBase + offset)); + if (Exit.isSuccess(attempt)) return { port: portBase + offset, listener: attempt.value }; + } + return yield* Effect.die("No port in the native reservation span was free for the fixture"); + }); + it.live("skips a native backend port claimed by another saved stack", () => Effect.scoped( Effect.gen(function* () { @@ -451,17 +462,12 @@ it.live("skips a native backend port claimed by another saved stack", () => it.live("skips a native backend port a wildcard listener holds", () => Effect.scoped( Effect.gen(function* () { - const probeScope = yield* Scope.make(); - const probe = yield* reserveNativePort([], "pooler", fixedStart).pipe( - Effect.provideService(Scope.Scope, probeScope), - ); - yield* Scope.close(probeScope, Exit.void); - // A wildcard bind is reachable through loopback, so a loopback-only probe would miss it. - yield* bind("0.0.0.0", probe.port); + const blocked = yield* bindBlockingPort("0.0.0.0"); + const forcedStart = Effect.succeed(blocked.port - portBase); - const reserved = yield* reserveNativePort([], "pooler", fixedStart); - expect(reserved.port).not.toBe(probe.port); + const reserved = yield* reserveNativePort([], "pooler", forcedStart); + expect(reserved.port).not.toBe(blocked.port); }), ).pipe(Effect.provide(NodeServices.layer)), ); @@ -493,15 +499,11 @@ it.live("skips a public auto candidate a loopback listener already holds", () => const request = { stackId: "stack", key: "api", host: "0.0.0.0", port: "auto" as const }; const accept = (_host: string, port: number) => Effect.succeed(port); - const probeScope = yield* Scope.make(); - const probe = yield* ports - .acquire(request, accept) - .pipe(Effect.provideService(Scope.Scope, probeScope)); - yield* Scope.close(probeScope, Exit.void); + // A loopback-only listener keeps occupying the port while only its saved claim is released, + // so the next auto allocation has to skip it for real via loopbackOccupied, not a real bind. + const probe = yield* ports.acquire(request, (_host, port) => bind("127.0.0.1", port)); yield* ports.release("stack", "api"); - yield* bind("127.0.0.1", probe.port); - const acquired = yield* ports.acquire(request, accept); expect(acquired.port).not.toBe(probe.port); }), diff --git a/packages/stack/src/Ports.ts b/packages/stack/src/Ports.ts index 6150b0f796..8f0f76c461 100644 --- a/packages/stack/src/Ports.ts +++ b/packages/stack/src/Ports.ts @@ -2,9 +2,9 @@ import { Cause, Crypto, Data, Effect, Exit, Hash, Option, Scope } from "effect"; import * as Net from "node:net"; import type * as State from "./State.ts"; -const portBase = 20000; +export const portBase = 20000; /** Stays below the Linux ephemeral range, per the [architecture ADR](../../../docs/adr/0017-simplified-managed-stack-architecture.md). */ -const portSpan = 12768; +export const portSpan = 12768; /** Co-prime with the span, so the scan visits every port once and steps past reserved ranges. */ const portStride = 257; diff --git a/packages/stack/src/services/Functions.integration.test.ts b/packages/stack/src/services/Functions.integration.test.ts index 9460688d14..34eebfe083 100644 --- a/packages/stack/src/services/Functions.integration.test.ts +++ b/packages/stack/src/services/Functions.integration.test.ts @@ -1,4 +1,4 @@ -import { NodeHttpClient, NodePath, NodeServices } from "@effect/platform-node"; +import { NodeCrypto, NodeHttpClient, NodePath, NodeServices } from "@effect/platform-node"; import { describe, expect, it } from "@effect/vitest"; import { Cause, @@ -16,10 +16,8 @@ import { } from "effect"; import { HttpClient, HttpClientRequest } from "effect/unstable/http"; import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; -// oxlint-disable-next-line effecttsgo/node-builtin-import -- test-only candidate start, not the injected Crypto service. -import { randomInt as nodeRandomInt } from "node:crypto"; import { ContainerError, type ContainerRuntime } from "../runtime/Container.ts"; -import { reserveNativePort } from "../Ports.ts"; +import { randomPortSpanStart, reserveNativePort } from "../Ports.ts"; import { makeService } from "../Service.ts"; import type * as State from "../State.ts"; import { makeServiceRecipe } from "./Catalog.ts"; @@ -27,12 +25,14 @@ import * as Functions from "./Functions.ts"; // No saved stacks to consult; this test never launches the native backend it configures. const testReadPortClaims = Effect.succeed([]); -const testRandomStart = Effect.sync(() => nodeRandomInt(0, 1_000_000)); const testReserveNativePort = ( key: string, claims: ReadonlyArray, excluded: ReadonlySet, -) => reserveNativePort(claims, key, testRandomStart, excluded); +) => + Effect.flatMap(Crypto.Crypto, (crypto) => + reserveNativePort(claims, key, randomPortSpanStart(crypto), excluded), + ).pipe(Effect.provide(NodeCrypto.layer)); const options = (root: string) => ({ stackId: "catalog-functions", diff --git a/packages/stack/src/services/ProcessRecipe.integration.test.ts b/packages/stack/src/services/ProcessRecipe.integration.test.ts index 7d0669eb28..84221f606d 100644 --- a/packages/stack/src/services/ProcessRecipe.integration.test.ts +++ b/packages/stack/src/services/ProcessRecipe.integration.test.ts @@ -1,4 +1,4 @@ -import { NodeHttpClient, NodeServices } from "@effect/platform-node"; +import { NodeCrypto, NodeHttpClient, NodeServices } from "@effect/platform-node"; import { describe, expect, it } from "@effect/vitest"; import { PlatformError, @@ -24,10 +24,8 @@ import { systemError } from "effect/PlatformError"; import * as Net from "node:net"; // oxlint-disable-next-line effecttsgo/node-builtin-import -- the collision fixture owns a local HTTP listener. import * as NodeHttp from "node:http"; -// oxlint-disable-next-line effecttsgo/node-builtin-import -- test-only candidate start, not the injected Crypto service. -import { randomInt as nodeRandomInt } from "node:crypto"; import { catalogPins, resolveArtifact, type ServiceKind } from "../Artifacts.ts"; -import { reserveNativePort } from "../Ports.ts"; +import { randomPortSpanStart, reserveNativePort } from "../Ports.ts"; import type * as State from "../State.ts"; import { makeArtifactStore, @@ -104,12 +102,14 @@ const isPortOccupied = (port: number): Effect.Effect => // No saved stacks to consult outside the claim-interaction test below. const testReadPortClaims = Effect.succeed([]); -const testRandomStart = Effect.sync(() => nodeRandomInt(0, 1_000_000)); const testReserveNativePort = ( key: string, claims: ReadonlyArray, excluded: ReadonlySet, -) => reserveNativePort(claims, key, testRandomStart, excluded); +) => + Effect.flatMap(Crypto.Crypto, (crypto) => + reserveNativePort(claims, key, randomPortSpanStart(crypto), excluded), + ).pipe(Effect.provide(NodeCrypto.layer)); describe("ProcessRecipe launch cleanup", () => { for (const scenario of [ @@ -1018,6 +1018,14 @@ describe("process recipe startup", () => { poolMode: "transaction", }, }; + // Pooler reserves "http" and "sql" concurrently; sharing one scan start forces both + // to contend for the same first candidate so the reservation has to skip one of them. + const sharedStart = Effect.succeed(yield* randomPortSpanStart(crypto)); + const contendingReserveNativePort = ( + key: string, + claims: ReadonlyArray, + excluded: ReadonlySet, + ) => reserveNativePort(claims, key, sharedStart, excluded); const recipe = yield* makeProcessRecipe( creation, { @@ -1036,7 +1044,7 @@ describe("process recipe startup", () => { spawner, container: undefined, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: contendingReserveNativePort, }, Pooler.makeSpec(), ); @@ -1050,8 +1058,7 @@ describe("process recipe startup", () => { yield* runtime.health; const endpoints = yield* Ref.get(recipe.endpoints); const endpoint = endpoints.get("http"); - // A backend chosen from this range (Ports.ts's below-ephemeral 20000..32767 span) - // can never coincide with an OS-auto-assigned ephemeral port, so releasing the probe + // This range sits below the default Linux ephemeral range, so releasing the probe // before the child binds cannot race an unrelated outgoing connection for the number. expect(endpoint?.port).toBeGreaterThanOrEqual(20000); expect(endpoint?.port).toBeLessThan(32768); From 69fdb613da73297a5fe472a81558386c8378604f Mon Sep 17 00:00:00 2001 From: Julien Goux Date: Fri, 2 Oct 2026 18:46:12 +0200 Subject: [PATCH 4/5] test(stack): require port claims and make port collision fixtures race-free Second round of review on the backend-port-reservation fix: - makeServiceRecipe took an optional state and silently fell back to Effect.succeed([]) for port claims when it was omitted, so a caller could lose saved-stack claim protection without any signal. readPortClaims is now a required parameter; Owner.ts passes the stack's own state.claims, and every test passes an explicit claims reader. - Test helpers in Functions.integration.test.ts and ProcessRecipe.integration.test.ts drew their reservation start from node:crypto's randomInt behind a lint suppression. They now build the same randomPortSpanStart(crypto) production uses, sourced from each test's own Crypto service, or a fixed start where a test needs a deterministic candidate. - Two Ports.integration.test.ts collision fixtures released their discovered port before binding the listener meant to block it, leaving a window where another process could take the port and break the fixture. The native fixture now binds the blocking wildcard listener directly, retrying over candidates, and reserves with a fixed start equal to its port. The public fixture binds the initial loopback listener during acquisition and keeps it alive while releasing only the saved claim, so the port is never unowned. - The endpoint-uniqueness assertion let each reservation draw its own random start, so concurrent reservations rarely contended for the same candidate. It now shares one fixed start across the batch so the two endpoints must contend and skip each other's held listener. - Reworded a comment that overstated the below-ephemeral range's isolation: it sits below Linux's default ephemeral range, but a custom host dynamic-port range can still overlap it. --- packages/stack/src/Owner.ts | 2 +- packages/stack/src/Ports.integration.test.ts | 57 ++++++++++++------- ...nalyticsVectorImgproxy.integration.test.ts | 4 ++ .../services/AuthStorage.integration.test.ts | 4 ++ .../src/services/Catalog.integration.test.ts | 3 + packages/stack/src/services/Catalog.ts | 8 ++- .../services/Functions.integration.test.ts | 24 ++++---- .../src/services/Mail.integration.test.ts | 21 ++++++- .../src/services/Pooler.integration.test.ts | 2 + .../ProcessRecipe.integration.test.ts | 52 +++++++++-------- .../src/services/Realtime.integration.test.ts | 4 ++ .../src/services/Rest.integration.test.ts | 4 ++ .../src/services/Vector.integration.test.ts | 3 + 13 files changed, 129 insertions(+), 59 deletions(-) diff --git a/packages/stack/src/Owner.ts b/packages/stack/src/Owner.ts index ed692c05c6..e25cc167c5 100644 --- a/packages/stack/src/Owner.ts +++ b/packages/stack/src/Owner.ts @@ -300,7 +300,7 @@ const makeOwner = Effect.fn("Owner.make")(function* (options: OwnerOptions) { helpers, ...(options.hostGateway === undefined ? {} : { hostGateway: options.hostGateway }), }, - options.state, + options.state.claims, ).pipe(Effect.provideContext(services)); const persistCreation = (entry: Pick, creation: ServiceCreation) => diff --git a/packages/stack/src/Ports.integration.test.ts b/packages/stack/src/Ports.integration.test.ts index 7ee87c4cfb..14c0b456eb 100644 --- a/packages/stack/src/Ports.integration.test.ts +++ b/packages/stack/src/Ports.integration.test.ts @@ -1,7 +1,18 @@ import { NodeServices, NodeSocketServer } from "@effect/platform-node"; import { expect, it } from "@effect/vitest"; -import { Cause, Context, Effect, Exit, FileSystem, Layer, Option, Ref, Scope } from "effect"; -import { makePorts, PortError, reserveNativePort } from "./Ports.ts"; +import { + Cause, + Context, + Crypto, + Effect, + Exit, + FileSystem, + Layer, + Option, + Ref, + Scope, +} from "effect"; +import { makePorts, PortError, randomPortSpanStart, reserveNativePort } from "./Ports.ts"; import * as State from "./State.ts"; const makeTestState = (root: string) => @@ -425,6 +436,11 @@ it.live("lets exactly one of two stacks sharing a saved port bind it when both s // Fixed so a test can force two reservations to the same candidate; production uses randomPortSpanStart. const fixedStart = Effect.succeed(0); +// Mirrors Ports.ts's below-ephemeral span (see the architecture ADR), for fixtures that need a +// candidate in that range without depending on its private constants. +const belowEphemeralBase = 20000; +const belowEphemeralSpan = 12768; + it.live("skips a native backend port claimed by another saved stack", () => Effect.scoped( Effect.gen(function* () { @@ -451,17 +467,24 @@ it.live("skips a native backend port claimed by another saved stack", () => it.live("skips a native backend port a wildcard listener holds", () => Effect.scoped( Effect.gen(function* () { - const probeScope = yield* Scope.make(); - const probe = yield* reserveNativePort([], "pooler", fixedStart).pipe( - Effect.provideService(Scope.Scope, probeScope), - ); - yield* Scope.close(probeScope, Exit.void); + // Binds the blocking wildcard listener directly, retrying over candidates, instead of + // discovering a port with a probe and releasing it first: the port is never unowned. + const start = yield* randomPortSpanStart(yield* Crypto.Crypto); + let blocked: number | undefined; + for (let offset = 0; offset < 50 && blocked === undefined; offset++) { + const port = belowEphemeralBase + ((start + offset) % belowEphemeralSpan); + const result = yield* Effect.exit(bind("0.0.0.0", port)); + if (Exit.isSuccess(result)) blocked = port; + } + if (blocked === undefined) return yield* Effect.die("no candidate port available for test"); // A wildcard bind is reachable through loopback, so a loopback-only probe would miss it. - yield* bind("0.0.0.0", probe.port); - - const reserved = yield* reserveNativePort([], "pooler", fixedStart); - expect(reserved.port).not.toBe(probe.port); + const reserved = yield* reserveNativePort( + [], + "pooler", + Effect.succeed(blocked - belowEphemeralBase), + ); + expect(reserved.port).not.toBe(blocked); }), ).pipe(Effect.provide(NodeServices.layer)), ); @@ -493,17 +516,13 @@ it.live("skips a public auto candidate a loopback listener already holds", () => const request = { stackId: "stack", key: "api", host: "0.0.0.0", port: "auto" as const }; const accept = (_host: string, port: number) => Effect.succeed(port); - const probeScope = yield* Scope.make(); - const probe = yield* ports - .acquire(request, accept) - .pipe(Effect.provideService(Scope.Scope, probeScope)); - yield* Scope.close(probeScope, Exit.void); + // Binds the initial loopback listener during acquisition and keeps it alive while + // releasing the saved claim, so the port is never unowned afterward. + const held = yield* ports.acquire({ ...request, host: "127.0.0.1" }, bind); yield* ports.release("stack", "api"); - yield* bind("127.0.0.1", probe.port); - const acquired = yield* ports.acquire(request, accept); - expect(acquired.port).not.toBe(probe.port); + expect(acquired.port).not.toBe(held.port); }), ).pipe(Effect.provide(NodeServices.layer)), ); diff --git a/packages/stack/src/services/AnalyticsVectorImgproxy.integration.test.ts b/packages/stack/src/services/AnalyticsVectorImgproxy.integration.test.ts index de79a2f854..f7d43bdadb 100644 --- a/packages/stack/src/services/AnalyticsVectorImgproxy.integration.test.ts +++ b/packages/stack/src/services/AnalyticsVectorImgproxy.integration.test.ts @@ -41,6 +41,7 @@ describe("service catalog", () => { }, }, dockerOptions(root), + Effect.succeed([]), ); const database = yield* makeService(databaseRecipe.definition, { id: "database", @@ -57,6 +58,7 @@ describe("service catalog", () => { config: { databaseUrl, backend: "postgres", apiKey: "catalog-analytics" }, }, dockerOptions(root), + Effect.succeed([]), ); const analytics = yield* makeService(analyticsRecipe.definition, { id: "analytics", @@ -83,6 +85,7 @@ describe("service catalog", () => { }, }, dockerOptions(root), + Effect.succeed([]), ); yield* fs.writeFileString( `${root}/vector.yaml`, @@ -106,6 +109,7 @@ describe("service catalog", () => { const imgproxyRecipe = yield* makeServiceRecipe( { service: "imgproxy", config: { filePath: imageRoot } }, dockerOptions(root), + Effect.succeed([]), ); const imgproxy = yield* makeService(imgproxyRecipe.definition, { id: "imgproxy", diff --git a/packages/stack/src/services/AuthStorage.integration.test.ts b/packages/stack/src/services/AuthStorage.integration.test.ts index 437e7975e8..a28f9e6571 100644 --- a/packages/stack/src/services/AuthStorage.integration.test.ts +++ b/packages/stack/src/services/AuthStorage.integration.test.ts @@ -42,6 +42,7 @@ describe("service catalog", () => { }, }, dockerOptions(root), + Effect.succeed([]), ); const database = yield* makeService(databaseRecipe.definition, { id: "database", @@ -57,6 +58,7 @@ describe("service catalog", () => { config: { databaseUrl, jwtSecret: secret, jwtExpiry: 3600 }, }, dockerOptions(root), + Effect.succeed([]), ); const auth = yield* makeService(authRecipe.definition, { id: "auth", @@ -98,6 +100,7 @@ describe("service catalog", () => { const imgproxyRecipe = yield* makeServiceRecipe( { service: "imgproxy", config: { filePath: storageRoot } }, dockerOptions(root), + Effect.succeed([]), ); const imgproxy = yield* makeService(imgproxyRecipe.definition, { id: "imgproxy", @@ -117,6 +120,7 @@ describe("service catalog", () => { }, }, dockerOptions(root), + Effect.succeed([]), ); const storage = yield* makeService(storageRecipe.definition, { id: "storage", diff --git a/packages/stack/src/services/Catalog.integration.test.ts b/packages/stack/src/services/Catalog.integration.test.ts index be06f948a4..5cad63c214 100644 --- a/packages/stack/src/services/Catalog.integration.test.ts +++ b/packages/stack/src/services/Catalog.integration.test.ts @@ -24,6 +24,7 @@ describe("service catalog", () => { endpoints: { smtp: { port: "auto" } }, }, options(root), + Effect.succeed([]), ).pipe(Effect.exit); expect(Exit.isFailure(result)).toBe(true); const misplacedDatabaseVersion = yield* makeServiceRecipe( @@ -38,6 +39,7 @@ describe("service catalog", () => { }, }, options(root), + Effect.succeed([]), ).pipe(Effect.exit); expect(Exit.isFailure(misplacedDatabaseVersion)).toBe(true); }), @@ -56,6 +58,7 @@ describe("service catalog", () => { config: { databaseUrl: "postgres://db" }, }, options(root), + Effect.succeed([]), ), ); expect(error.operation).toBe("config"); diff --git a/packages/stack/src/services/Catalog.ts b/packages/stack/src/services/Catalog.ts index 748f8ff349..b029213213 100644 --- a/packages/stack/src/services/Catalog.ts +++ b/packages/stack/src/services/Catalog.ts @@ -300,7 +300,11 @@ const databaseRecipe = ( }); export const makeServiceRecipe = Effect.fn("Catalog.makeServiceRecipe")( - (input: unknown, options: CatalogOptions, state?: State.Interface) => + ( + input: unknown, + options: CatalogOptions, + readPortClaims: Effect.Effect, State.StateError>, + ) => Effect.gen(function* () { const endpointError = validateEndpointNames(input); if (endpointError !== undefined) return yield* endpointError; @@ -360,7 +364,7 @@ export const makeServiceRecipe = Effect.fn("Catalog.makeServiceRecipe")( client, spawner, container, - readPortClaims: state?.claims ?? Effect.succeed([]), + readPortClaims, reserveNativePort: (key, claims, excluded) => reserveNativePort(claims, key, randomPortSpanStart(crypto), excluded), }; diff --git a/packages/stack/src/services/Functions.integration.test.ts b/packages/stack/src/services/Functions.integration.test.ts index 9460688d14..2d5e920a8b 100644 --- a/packages/stack/src/services/Functions.integration.test.ts +++ b/packages/stack/src/services/Functions.integration.test.ts @@ -16,10 +16,8 @@ import { } from "effect"; import { HttpClient, HttpClientRequest } from "effect/unstable/http"; import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; -// oxlint-disable-next-line effecttsgo/node-builtin-import -- test-only candidate start, not the injected Crypto service. -import { randomInt as nodeRandomInt } from "node:crypto"; import { ContainerError, type ContainerRuntime } from "../runtime/Container.ts"; -import { reserveNativePort } from "../Ports.ts"; +import { randomPortSpanStart, reserveNativePort } from "../Ports.ts"; import { makeService } from "../Service.ts"; import type * as State from "../State.ts"; import { makeServiceRecipe } from "./Catalog.ts"; @@ -27,12 +25,10 @@ import * as Functions from "./Functions.ts"; // No saved stacks to consult; this test never launches the native backend it configures. const testReadPortClaims = Effect.succeed([]); -const testRandomStart = Effect.sync(() => nodeRandomInt(0, 1_000_000)); -const testReserveNativePort = ( - key: string, - claims: ReadonlyArray, - excluded: ReadonlySet, -) => reserveNativePort(claims, key, testRandomStart, excluded); +const testReserveNativePort = + (crypto: Crypto.Crypto) => + (key: string, claims: ReadonlyArray, excluded: ReadonlySet) => + reserveNativePort(claims, key, randomPortSpanStart(crypto), excluded); const options = (root: string) => ({ stackId: "catalog-functions", @@ -105,6 +101,7 @@ describe("service catalog", () => { }, }, { ...dockerOptions(root), stackId, instanceId }, + Effect.succeed([]), ); const instance = yield* makeService(recipe.definition, { id: instanceId, @@ -206,6 +203,7 @@ describe("service catalog", () => { instanceId: "ancestor", cacheRoot: "/tmp/supabase-stack-artifacts", }, + Effect.succeed([]), ); const logs = yield* Ref.make(""); yield* recipe.logs.pipe( @@ -274,6 +272,7 @@ describe("service catalog", () => { instanceId: "deno-config", cacheRoot: "/tmp/supabase-stack-artifacts", }, + Effect.succeed([]), ); const logs = yield* Ref.make(""); yield* recipe.logs.pipe( @@ -345,6 +344,7 @@ describe("service catalog", () => { instanceId: "plain-deno-config", cacheRoot: "/tmp/supabase-stack-artifacts", }, + Effect.succeed([]), ); const logs = yield* Ref.make(""); const warned = yield* Deferred.make(); @@ -455,6 +455,7 @@ for (const runtime of ["native", "docker"] as const) { runtime, cacheRoot: "/tmp/supabase-stack-artifacts", }, + Effect.succeed([]), ); const logs = yield* Ref.make(""); yield* recipe.logs.pipe( @@ -549,6 +550,7 @@ it.effect("passes POSIX project paths to a docker Functions container from a Win }, }, }; + const crypto = yield* Crypto.Crypto; const recipe = yield* Functions.makeRecipe( creation, { @@ -561,12 +563,12 @@ it.effect("passes POSIX project paths to a docker Functions container from a Win { fs: yield* FileSystem.FileSystem, path: yield* Path.Path.pipe(Effect.provide(NodePath.layerWin32)), - crypto: yield* Crypto.Crypto, + crypto, client: yield* HttpClient.HttpClient, spawner: yield* ChildProcessSpawner.ChildProcessSpawner, container, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), }, ); const scope = yield* Scope.make(); diff --git a/packages/stack/src/services/Mail.integration.test.ts b/packages/stack/src/services/Mail.integration.test.ts index f80d9a3073..1b6b262f88 100644 --- a/packages/stack/src/services/Mail.integration.test.ts +++ b/packages/stack/src/services/Mail.integration.test.ts @@ -64,6 +64,7 @@ describe("service catalog", () => { const recipe = yield* makeServiceRecipe( { service: "mail", config: {} }, dockerOptions(root), + Effect.succeed([]), ); const instance = yield* makeService(recipe.definition, { id: "mail", @@ -90,7 +91,11 @@ describe("service catalog", () => { const path = yield* Path.Path; const root = yield* fs.makeTempDirectoryScoped({ prefix: "catalog-mail-isolation-" }); const recipeFor = (instanceId: string) => - makeServiceRecipe({ service: "mail", config: {} }, { ...options(root), instanceId }); + makeServiceRecipe( + { service: "mail", config: {} }, + { ...options(root), instanceId }, + Effect.succeed([]), + ); const first = yield* recipeFor("mail-a"); const second = yield* recipeFor("mail-b"); const firstInstance = yield* makeService(first.definition, { @@ -121,7 +126,11 @@ describe("service catalog", () => { Effect.gen(function* () { const fs = yield* FileSystem.FileSystem; const root = yield* fs.makeTempDirectoryScoped({ prefix: "catalog-mail-persistence-" }); - const recipe = yield* makeServiceRecipe({ service: "mail", config: {} }, options(root)); + const recipe = yield* makeServiceRecipe( + { service: "mail", config: {} }, + options(root), + Effect.succeed([]), + ); const instance = yield* makeService(recipe.definition, { id: "mail", config: recipe.creation, @@ -155,6 +164,7 @@ describe("service catalog", () => { const recipe = yield* makeServiceRecipe( { service: "mail", config: {} }, dockerOptions(root), + Effect.succeed([]), ); const instance = yield* makeService(recipe.definition, { id: "mail", @@ -185,7 +195,11 @@ describe("service catalog", () => { const fs = yield* FileSystem.FileSystem; const path = yield* Path.Path; const root = yield* fs.makeTempDirectoryScoped({ prefix: "catalog-mail-destroy-" }); - const recipe = yield* makeServiceRecipe({ service: "mail", config: {} }, options(root)); + const recipe = yield* makeServiceRecipe( + { service: "mail", config: {} }, + options(root), + Effect.succeed([]), + ); const instance = yield* makeService(recipe.definition, { id: "mail", config: recipe.creation, @@ -209,6 +223,7 @@ describe("service catalog", () => { const recipe = yield* makeServiceRecipe( { service: "mail", config: {} }, { ...options(root), instanceId: "../escaped" }, + Effect.succeed([]), ); const instance = yield* makeService(recipe.definition, { id: "mail", diff --git a/packages/stack/src/services/Pooler.integration.test.ts b/packages/stack/src/services/Pooler.integration.test.ts index d6e0866bc4..e853e7f01f 100644 --- a/packages/stack/src/services/Pooler.integration.test.ts +++ b/packages/stack/src/services/Pooler.integration.test.ts @@ -41,6 +41,7 @@ describe("service catalog", () => { }, }, dockerOptions(root), + Effect.succeed([]), ); const database = yield* makeService(databaseRecipe.definition, { id: "database", @@ -70,6 +71,7 @@ describe("service catalog", () => { }, }, runtime === "native" ? options(root) : dockerOptions(root), + Effect.succeed([]), ); const pooler = yield* makeService(poolerRecipe.definition, { id: `pooler-${runtime}-${poolMode}`, diff --git a/packages/stack/src/services/ProcessRecipe.integration.test.ts b/packages/stack/src/services/ProcessRecipe.integration.test.ts index 7d0669eb28..f0d7bb1276 100644 --- a/packages/stack/src/services/ProcessRecipe.integration.test.ts +++ b/packages/stack/src/services/ProcessRecipe.integration.test.ts @@ -24,10 +24,8 @@ import { systemError } from "effect/PlatformError"; import * as Net from "node:net"; // oxlint-disable-next-line effecttsgo/node-builtin-import -- the collision fixture owns a local HTTP listener. import * as NodeHttp from "node:http"; -// oxlint-disable-next-line effecttsgo/node-builtin-import -- test-only candidate start, not the injected Crypto service. -import { randomInt as nodeRandomInt } from "node:crypto"; import { catalogPins, resolveArtifact, type ServiceKind } from "../Artifacts.ts"; -import { reserveNativePort } from "../Ports.ts"; +import { randomPortSpanStart, reserveNativePort } from "../Ports.ts"; import type * as State from "../State.ts"; import { makeArtifactStore, @@ -104,12 +102,17 @@ const isPortOccupied = (port: number): Effect.Effect => // No saved stacks to consult outside the claim-interaction test below. const testReadPortClaims = Effect.succeed([]); -const testRandomStart = Effect.sync(() => nodeRandomInt(0, 1_000_000)); -const testReserveNativePort = ( +const testReserveNativePort = + (crypto: Crypto.Crypto) => + (key: string, claims: ReadonlyArray, excluded: ReadonlySet) => + reserveNativePort(claims, key, randomPortSpanStart(crypto), excluded); +// Shares one start across every call so concurrent reservations must contend for the same +// candidate and skip each other's held listener, instead of each drawing its own random start. +const contendingReserveNativePort = ( key: string, claims: ReadonlyArray, excluded: ReadonlySet, -) => reserveNativePort(claims, key, testRandomStart, excluded); +) => reserveNativePort(claims, key, Effect.succeed(0), excluded); describe("ProcessRecipe launch cleanup", () => { for (const scenario of [ @@ -194,15 +197,16 @@ describe("ProcessRecipe launch cleanup", () => { return makeProcess("service", Effect.never); }), }; + const crypto = yield* Crypto.Crypto; const dependencies = { fs: yield* FileSystem.FileSystem, path: yield* Path.Path, - crypto: yield* Crypto.Crypto, + crypto, client: yield* HttpClient.HttpClient, spawner: yield* ChildProcessSpawner.ChildProcessSpawner, container, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, options, dependencies, spec); const service = yield* makeService(recipe.definition, { @@ -291,7 +295,7 @@ describe("ProcessRecipe launch cleanup", () => { spawner, container: undefined, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, nativeOptions, dependencies, nativeSpec); const service = yield* makeService(recipe.definition, { @@ -403,7 +407,7 @@ describe("ProcessRecipe launch cleanup", () => { spawner, container: undefined, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, nativeOptions, dependencies, nativeSpec); const service = yield* makeService(recipe.definition, { @@ -461,6 +465,7 @@ const realtimeService = Effect.fn(function* (container: ContainerRuntime) { service: "realtime", config: { databaseUrl: "postgresql://postgres:postgres@host.docker.internal:54322/postgres" }, }; + const crypto = yield* Crypto.Crypto; const recipe = yield* makeProcessRecipe( creation, { @@ -473,12 +478,12 @@ const realtimeService = Effect.fn(function* (container: ContainerRuntime) { { fs: yield* FileSystem.FileSystem, path: yield* Path.Path, - crypto: yield* Crypto.Crypto, + crypto, client: yield* HttpClient.HttpClient, spawner: yield* ChildProcessSpawner.ChildProcessSpawner, container, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), }, Realtime.makeSpec(), ); @@ -684,7 +689,7 @@ const nativeRestRecipe = Effect.fn(function* ( spawner, container: undefined, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), }, { ...spec, @@ -962,7 +967,7 @@ describe("process recipe startup", () => { spawner: interceptingSpawner, container: undefined, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), }, Pooler.makeSpec(), ); @@ -1036,7 +1041,7 @@ describe("process recipe startup", () => { spawner, container: undefined, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: contendingReserveNativePort, }, Pooler.makeSpec(), ); @@ -1050,12 +1055,13 @@ describe("process recipe startup", () => { yield* runtime.health; const endpoints = yield* Ref.get(recipe.endpoints); const endpoint = endpoints.get("http"); - // A backend chosen from this range (Ports.ts's below-ephemeral 20000..32767 span) - // can never coincide with an OS-auto-assigned ephemeral port, so releasing the probe - // before the child binds cannot race an unrelated outgoing connection for the number. + // This below-ephemeral range (Ports.ts's 20000..32767 span) sits below Linux's default + // ephemeral range, but a custom host dynamic-port range can still overlap it (the + // architecture ADR excludes relying on non-default ranges). expect(endpoint?.port).toBeGreaterThanOrEqual(20000); expect(endpoint?.port).toBeLessThan(32768); - // Pooler reserves "http" and "sql" concurrently; they must never settle on the same port. + // Pooler reserves "http" and "sql" concurrently from the same forced start; they must + // contend for the candidate and still never settle on the same port. const sql = endpoints.get("sql"); expect(sql?.port).toBeDefined(); expect(sql?.port).not.toBe(endpoint?.port); @@ -1106,7 +1112,7 @@ describe("process recipe startup", () => { spawner: countingSpawner(spawner, mainLaunches, startupLaunches, true), container: undefined, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), }, Pooler.makeSpec(), ); @@ -1242,7 +1248,7 @@ describe("process recipe startup", () => { spawner: deadlineSpawner, container: undefined, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), }, Pooler.makeSpec(), ); @@ -1306,7 +1312,7 @@ describe("process recipe startup", () => { spawner: countingSpawner(spawner, mainLaunches, startupLaunches), container: undefined, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), }, Pooler.makeSpec(), ); @@ -1412,7 +1418,7 @@ describe("process recipe startup", () => { spawner: failingSpawner, container: undefined, readPortClaims: testReadPortClaims, - reserveNativePort: testReserveNativePort, + reserveNativePort: testReserveNativePort(crypto), }, Pooler.makeSpec(), ); diff --git a/packages/stack/src/services/Realtime.integration.test.ts b/packages/stack/src/services/Realtime.integration.test.ts index 3d495a5ef9..06922425c0 100644 --- a/packages/stack/src/services/Realtime.integration.test.ts +++ b/packages/stack/src/services/Realtime.integration.test.ts @@ -40,6 +40,7 @@ describe("service catalog", () => { }, }, dockerOptions(root), + Effect.succeed([]), ); const database = yield* makeService(databaseRecipe.definition, { id: "database", @@ -53,6 +54,7 @@ describe("service catalog", () => { const realtimeRecipe = yield* makeServiceRecipe( { service: "realtime", config: { databaseUrl, jwtSecret: secret } }, dockerOptions(root), + Effect.succeed([]), ); const realtime = yield* makeService(realtimeRecipe.definition, { id: "realtime", @@ -71,6 +73,7 @@ describe("service catalog", () => { const pgmetaRecipe = yield* makeServiceRecipe( { service: "pgmeta", config: { databaseUrl } }, dockerOptions(root), + Effect.succeed([]), ); const pgmeta = yield* makeService(pgmetaRecipe.definition, { id: "pgmeta", @@ -99,6 +102,7 @@ describe("service catalog", () => { }, }, dockerOptions(root), + Effect.succeed([]), ); const studio = yield* makeService(studioRecipe.definition, { id: "studio", diff --git a/packages/stack/src/services/Rest.integration.test.ts b/packages/stack/src/services/Rest.integration.test.ts index a8d87db95f..0f902e8bbe 100644 --- a/packages/stack/src/services/Rest.integration.test.ts +++ b/packages/stack/src/services/Rest.integration.test.ts @@ -74,6 +74,7 @@ describe("service catalog", () => { }, }, dockerOptions(root), + Effect.succeed([]), ); const database = yield* makeService(databaseRecipe.definition, { id: "database", @@ -139,6 +140,7 @@ describe("service catalog", () => { }, }, dockerOptions(root), + Effect.succeed([]), ); const rest = yield* makeService(restRecipe.definition, { id: "rest", @@ -216,6 +218,7 @@ describe("service catalog", () => { }, }, { ...options(root), stackId, cacheRoot: `${tmpdir()}/supabase-stack-artifacts` }, + Effect.succeed([]), ); const database = yield* makeService(databaseRecipe.definition, { id: "database", @@ -256,6 +259,7 @@ describe("service catalog", () => { }, }, { ...options(root), stackId, cacheRoot: `${tmpdir()}/supabase-stack-artifacts` }, + Effect.succeed([]), ); const rest = yield* makeService(restRecipe.definition, { id: "rest", diff --git a/packages/stack/src/services/Vector.integration.test.ts b/packages/stack/src/services/Vector.integration.test.ts index 2eba9f07ec..557c37b6b0 100644 --- a/packages/stack/src/services/Vector.integration.test.ts +++ b/packages/stack/src/services/Vector.integration.test.ts @@ -33,6 +33,7 @@ describe("vector recipe", () => { endpoints: { http: { port: "auto" } }, }, options(root, runtime), + Effect.succeed([]), ); const vector = yield* makeService(recipe.definition, { id: "vector", @@ -74,6 +75,7 @@ describe("vector recipe", () => { endpoints: { http: { port: "auto" } }, }, options(root, "docker"), + Effect.succeed([]), ); const vector = yield* makeService(recipe.definition, { id: "vector", @@ -107,6 +109,7 @@ describe("vector recipe", () => { endpoints: { http: { port: "auto" } }, }, options(root, "docker"), + Effect.succeed([]), ); const vector = yield* makeService(recipe.definition, { id: "vector", From a86f9c1cd0282590a2f5629de3aebe092333cf32 Mon Sep 17 00:00:00 2001 From: Julien Goux Date: Fri, 2 Oct 2026 20:12:34 +0200 Subject: [PATCH 5/5] test(stack): align port reservation comments and test title with the review Say why backend scans start at a random offset, and drop the stale claim-interaction and OS-assigned wording. --- packages/stack/src/Ports.ts | 2 +- .../ProcessRecipe.integration.test.ts | 146 +++++++++--------- 2 files changed, 73 insertions(+), 75 deletions(-) diff --git a/packages/stack/src/Ports.ts b/packages/stack/src/Ports.ts index 8f0f76c461..8b55842335 100644 --- a/packages/stack/src/Ports.ts +++ b/packages/stack/src/Ports.ts @@ -132,7 +132,7 @@ const bindNativePort = ( const emptyPortSet: ReadonlySet = new Set(); -/** Backend ports aren't persisted, so a random start buys nothing by staying stable across reopens. */ +/** Random, so a reopened stack doesn't retry a backend port still in server-side `TIME_WAIT`. */ export const randomPortSpanStart = (crypto: Crypto.Crypto): Effect.Effect => crypto.randomIntBetween(0, portSpan, { halfOpen: true }); diff --git a/packages/stack/src/services/ProcessRecipe.integration.test.ts b/packages/stack/src/services/ProcessRecipe.integration.test.ts index 1d4d3c4a81..de276d5051 100644 --- a/packages/stack/src/services/ProcessRecipe.integration.test.ts +++ b/packages/stack/src/services/ProcessRecipe.integration.test.ts @@ -100,7 +100,7 @@ const isPortOccupied = (port: number): Effect.Effect => }); }); -// No saved stacks to consult outside the claim-interaction test below. +// No saved stacks to consult. const testReadPortClaims = Effect.succeed([]); const testReserveNativePort = ( key: string, @@ -996,79 +996,77 @@ describe("process recipe startup", () => { ).pipe(Effect.provide(platform)), ); - it.live( - "reserves a native backend port from the below-ephemeral range, never an OS-assigned one", - () => - Effect.scoped( - Effect.gen(function* () { - const fs = yield* FileSystem.FileSystem; - const path = yield* Path.Path; - const crypto = yield* Crypto.Crypto; - const client = yield* HttpClient.HttpClient; - const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; - const root = yield* fs.makeTempDirectoryScoped({ prefix: "process-recipe-port-range-" }); - const cacheRoot = path.join(root, "cache"); - yield* nativePoolerArtifact(cacheRoot); - const creation: Pooler.Creation = { - service: "pooler", - config: { - databaseUrl: "postgresql://postgres:postgres@127.0.0.1:5432/postgres", - jwtSecret: "pooler-port-range-test-secret-with-more-than-32-characters", - tenant: "port-range-test", - poolMode: "transaction", - }, - }; - // Pooler reserves "http" and "sql" concurrently; sharing one scan start forces both - // to contend for the same first candidate so the reservation has to skip one of them. - const sharedStart = Effect.succeed(yield* randomPortSpanStart(crypto)); - const contendingReserveNativePort = ( - key: string, - claims: ReadonlyArray, - excluded: ReadonlySet, - ) => reserveNativePort(claims, key, sharedStart, excluded); - const recipe = yield* makeProcessRecipe( - creation, - { - stackId: "process-recipe-port-range", - instanceId: "instance", - root, - cacheRoot, - runtime: "native", - platform: { os: process.platform, arch: process.arch }, - }, - { - fs, - path, - crypto, - client, - spawner, - container: undefined, - readPortClaims: testReadPortClaims, - reserveNativePort: contendingReserveNativePort, - }, - Pooler.makeSpec(), - ); - if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); - const scope = yield* Scope.fork(yield* Effect.scope, "sequential"); - const runtime = yield* recipe.definition.launch({ - id: "pooler", - config: creation, - scope, - }); - yield* runtime.health; - const endpoints = yield* Ref.get(recipe.endpoints); - const endpoint = endpoints.get("http"); - // Below Linux's default ephemeral range, so the released probe port isn't handed to an - // outgoing connection there; a custom host dynamic range can still overlap (ADR 0017). - expect(endpoint?.port).toBeGreaterThanOrEqual(20000); - expect(endpoint?.port).toBeLessThan(32768); - // Pooler reserves "http" and "sql" concurrently; they must never settle on the same port. - const sql = endpoints.get("sql"); - expect(sql?.port).toBeDefined(); - expect(sql?.port).not.toBe(endpoint?.port); - yield* runtime.stop; - }), - ).pipe(Effect.provide(platform)), + it.live("reserves a native backend port below the default ephemeral range", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const crypto = yield* Crypto.Crypto; + const client = yield* HttpClient.HttpClient; + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "process-recipe-port-range-" }); + const cacheRoot = path.join(root, "cache"); + yield* nativePoolerArtifact(cacheRoot); + const creation: Pooler.Creation = { + service: "pooler", + config: { + databaseUrl: "postgresql://postgres:postgres@127.0.0.1:5432/postgres", + jwtSecret: "pooler-port-range-test-secret-with-more-than-32-characters", + tenant: "port-range-test", + poolMode: "transaction", + }, + }; + // Pooler reserves "http" and "sql" concurrently; sharing one scan start forces both + // to contend for the same first candidate so the reservation has to skip one of them. + const sharedStart = Effect.succeed(yield* randomPortSpanStart(crypto)); + const contendingReserveNativePort = ( + key: string, + claims: ReadonlyArray, + excluded: ReadonlySet, + ) => reserveNativePort(claims, key, sharedStart, excluded); + const recipe = yield* makeProcessRecipe( + creation, + { + stackId: "process-recipe-port-range", + instanceId: "instance", + root, + cacheRoot, + runtime: "native", + platform: { os: process.platform, arch: process.arch }, + }, + { + fs, + path, + crypto, + client, + spawner, + container: undefined, + readPortClaims: testReadPortClaims, + reserveNativePort: contendingReserveNativePort, + }, + Pooler.makeSpec(), + ); + if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); + const scope = yield* Scope.fork(yield* Effect.scope, "sequential"); + const runtime = yield* recipe.definition.launch({ + id: "pooler", + config: creation, + scope, + }); + yield* runtime.health; + const endpoints = yield* Ref.get(recipe.endpoints); + const endpoint = endpoints.get("http"); + // Below Linux's default ephemeral range, so the released probe port isn't handed to an + // outgoing connection there; a custom host dynamic range can still overlap (ADR 0017). + expect(endpoint?.port).toBeGreaterThanOrEqual(20000); + expect(endpoint?.port).toBeLessThan(32768); + // Pooler reserves "http" and "sql" concurrently; they must never settle on the same port. + const sql = endpoints.get("sql"); + expect(sql?.port).toBeDefined(); + expect(sql?.port).not.toBe(endpoint?.port); + yield* runtime.stop; + }), + ).pipe(Effect.provide(platform)), ); it.live("stops after three consecutive native Pooler port collisions", () =>