diff --git a/packages/stack/src/Owner.ts b/packages/stack/src/Owner.ts index b36d9fcfa1..e25cc167c5 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.claims, + ).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..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 } from "./Ports.ts"; +import { makePorts, portBase, portSpan, PortError, reserveNativePort } from "./Ports.ts"; import * as State from "./State.ts"; const makeTestState = (root: string) => @@ -421,3 +421,91 @@ 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); + +// 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* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped(); + const state = yield* makeTestState(root); + + const probeScope = yield* Scope.make(); + const probe = yield* reserveNativePort([], "pooler", fixedStart).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(yield* state.claims, "pooler", fixedStart); + 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* () { + // A wildcard bind is reachable through loopback, so a loopback-only probe would miss it. + const blocked = yield* bindBlockingPort("0.0.0.0"); + const forcedStart = Effect.succeed(blocked.port - portBase); + + const reserved = yield* reserveNativePort([], "pooler", forcedStart); + expect(reserved.port).not.toBe(blocked.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([], "pooler", fixedStart).pipe( + Effect.provideService(Scope.Scope, firstScope), + ); + yield* Scope.close(firstScope, Exit.void); + + 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); + + // 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"); + + 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 d45216416b..8b55842335 100644 --- a/packages/stack/src/Ports.ts +++ b/packages/stack/src/Ports.ts @@ -1,10 +1,10 @@ -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"; -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; @@ -88,6 +88,100 @@ 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(); + +/** 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 }); + +/** + * 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: ReadonlyArray, + key: string, + randomStart: Effect.Effect, + excluded: ReadonlySet = emptyPortSet, + ): Effect.Effect => + Effect.acquireRelease( + Effect.gen(function* () { + 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++) { + 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(() => { @@ -174,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/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 64e7828df5..b029213213 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 { randomPortSpanStart, reserveNativePort } from "../Ports.ts"; import { ServiceError, type ServiceDefinition } from "../Service.ts"; +import type * as State from "../State.ts"; import { makeDatabase, DatabaseConfig, @@ -298,7 +300,11 @@ const databaseRecipe = ( }); export const makeServiceRecipe = Effect.fn("Catalog.makeServiceRecipe")( - (input: unknown, options: CatalogOptions) => + ( + input: unknown, + options: CatalogOptions, + readPortClaims: Effect.Effect, State.StateError>, + ) => Effect.gen(function* () { const endpointError = validateEndpointNames(input); if (endpointError !== undefined) return yield* endpointError; @@ -351,7 +357,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 deps: ProcessDependencies = { + fs, + path, + crypto, + client, + spawner, + container, + readPortClaims, + reserveNativePort: (key, claims, excluded) => + reserveNativePort(claims, key, randomPortSpanStart(crypto), 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..a2eea0fa60 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, @@ -17,10 +17,23 @@ import { import { HttpClient, HttpClientRequest } from "effect/unstable/http"; import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; import { ContainerError, type ContainerRuntime } from "../runtime/Container.ts"; +import { randomPortSpanStart, 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 testReadPortClaims = Effect.succeed([]); +const testReserveNativePort = ( + key: string, + claims: ReadonlyArray, + excluded: ReadonlySet, +) => + Effect.flatMap(Crypto.Crypto, (crypto) => + reserveNativePort(claims, key, randomPortSpanStart(crypto), excluded), + ).pipe(Effect.provide(NodeCrypto.layer)); + const options = (root: string) => ({ stackId: "catalog-functions", instanceId: "instance", @@ -92,6 +105,7 @@ describe("service catalog", () => { }, }, { ...dockerOptions(root), stackId, instanceId }, + Effect.succeed([]), ); const instance = yield* makeService(recipe.definition, { id: instanceId, @@ -193,6 +207,7 @@ describe("service catalog", () => { instanceId: "ancestor", cacheRoot: "/tmp/supabase-stack-artifacts", }, + Effect.succeed([]), ); const logs = yield* Ref.make(""); yield* recipe.logs.pipe( @@ -261,6 +276,7 @@ describe("service catalog", () => { instanceId: "deno-config", cacheRoot: "/tmp/supabase-stack-artifacts", }, + Effect.succeed([]), ); const logs = yield* Ref.make(""); yield* recipe.logs.pipe( @@ -332,6 +348,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(); @@ -442,6 +459,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( @@ -552,6 +570,8 @@ 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, }, ); 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 9f14cc70ce..de276d5051 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, @@ -25,6 +25,8 @@ 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 { randomPortSpanStart, reserveNativePort } from "../Ports.ts"; +import type * as State from "../State.ts"; import { makeArtifactStore, type ArtifactRequest, @@ -98,6 +100,17 @@ const isPortOccupied = (port: number): Effect.Effect => }); }); +// No saved stacks to consult. +const testReadPortClaims = Effect.succeed([]); +const testReserveNativePort = ( + key: string, + claims: ReadonlyArray, + excluded: ReadonlySet, +) => + Effect.flatMap(Crypto.Crypto, (crypto) => + reserveNativePort(claims, key, randomPortSpanStart(crypto), excluded), + ).pipe(Effect.provide(NodeCrypto.layer)); + describe("ProcessRecipe launch cleanup", () => { for (const scenario of [ "partial launch", @@ -188,6 +201,8 @@ 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); const service = yield* makeService(recipe.definition, { @@ -275,6 +290,8 @@ describe("ProcessRecipe launch cleanup", () => { client, spawner, container: undefined, + readPortClaims: testReadPortClaims, + reserveNativePort: testReserveNativePort, } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, nativeOptions, dependencies, nativeSpec); const service = yield* makeService(recipe.definition, { @@ -385,6 +402,8 @@ describe("ProcessRecipe launch cleanup", () => { client, spawner, container: undefined, + readPortClaims: testReadPortClaims, + reserveNativePort: testReserveNativePort, } satisfies ProcessDependencies; const recipe = yield* makeProcessRecipe(creation, nativeOptions, dependencies, nativeSpec); const service = yield* makeService(recipe.definition, { @@ -458,6 +477,8 @@ const realtimeService = Effect.fn(function* (container: ContainerRuntime) { client: yield* HttpClient.HttpClient, spawner: yield* ChildProcessSpawner.ChildProcessSpawner, container, + readPortClaims: testReadPortClaims, + reserveNativePort: testReserveNativePort, }, Realtime.makeSpec(), ); @@ -655,7 +676,16 @@ 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, + readPortClaims: testReadPortClaims, + reserveNativePort: testReserveNativePort, + }, { ...spec, env: (_creation, endpoints) => @@ -924,7 +954,16 @@ 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, + readPortClaims: testReadPortClaims, + reserveNativePort: testReserveNativePort, + }, Pooler.makeSpec(), ); if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); @@ -957,6 +996,79 @@ describe("process recipe startup", () => { ).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", () => Effect.scoped( Effect.gen(function* () { @@ -998,6 +1110,8 @@ describe("process recipe startup", () => { client, spawner: countingSpawner(spawner, mainLaunches, startupLaunches, true), container: undefined, + readPortClaims: testReadPortClaims, + reserveNativePort: testReserveNativePort, }, Pooler.makeSpec(), ); @@ -1125,7 +1239,16 @@ 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, + readPortClaims: testReadPortClaims, + reserveNativePort: testReserveNativePort, + }, Pooler.makeSpec(), ); if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); @@ -1187,6 +1310,8 @@ describe("process recipe startup", () => { client, spawner: countingSpawner(spawner, mainLaunches, startupLaunches), container: undefined, + readPortClaims: testReadPortClaims, + reserveNativePort: testReserveNativePort, }, Pooler.makeSpec(), ); @@ -1284,7 +1409,16 @@ 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, + readPortClaims: testReadPortClaims, + 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..8757974905 100644 --- a/packages/stack/src/services/ProcessRecipe.ts +++ b/packages/stack/src/services/ProcessRecipe.ts @@ -19,9 +19,9 @@ 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 * as State from "../State.ts"; import { type ContainerError, type ContainerProcess, @@ -161,6 +161,12 @@ export interface ProcessDependencies { readonly client: HttpClient.HttpClient; readonly spawner: ChildProcessSpawnerService["Service"]; readonly container: ContainerRuntime | undefined; + readonly readPortClaims: Effect.Effect, State.StateError>; + readonly reserveNativePort: ( + key: string, + claims: ReadonlyArray, + excluded: ReadonlySet, + ) => Effect.Effect; } const serviceError = mapToServiceError; @@ -189,50 +195,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 +441,26 @@ 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); + // 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(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)), ); @@ -701,6 +676,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", 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",