diff --git a/apps/cli/docs/stack-commands.md b/apps/cli/docs/stack-commands.md index dd8cbdd896..1363c94bbe 100644 --- a/apps/cli/docs/stack-commands.md +++ b/apps/cli/docs/stack-commands.md @@ -88,6 +88,13 @@ owner are unreachable. Plain `status` JSON `env` comes from saved bindings and c stopped or sleeping members, while `--env` requires a reachable owner and a running primary database. +Starting a stopped stack whose `config.toml` moved an endpoint to a different port, or to and from +automatic, applies that change instead of failing: text output prints one line per changed endpoint, +for example `api: 54321 → 26199`, and JSON/stream-json output add an `endpoint_changes` array (each +entry naming the endpoint and its `from`/`to` port) to the success payload, present only when a +change applied. A changed artifact version or PostgreSQL major version still fails, naming the +changed setting and suggesting `supabase stack destroy`. + ## Exporting environment variables ```sh diff --git a/apps/cli/docs/supabase-home.md b/apps/cli/docs/supabase-home.md index 4e5507407a..982f2ceb0b 100644 --- a/apps/cli/docs/supabase-home.md +++ b/apps/cli/docs/supabase-home.md @@ -90,7 +90,11 @@ no separate command that rewrites it, and no second project-level pinned-version Raw `supabase/config.toml` values and their origins are loaded before defaults are applied. Explicit sticky values are persisted as `exact` intents in each managed document. Omitted values remain `automatic`; sibling worktrees and branches have independent stack identities and allocations. -Automatic allocation avoids ports saved by any stack, so stopped stacks keep their URLs. An exact +Automatic allocation avoids ports saved by any stack, so stopped stacks keep their URLs. A port +intent is stable: it does not migrate on its own, but starting a stopped stack whose configuration +moved an endpoint to a different exact port, or to and from automatic, re-plans that endpoint and +claims the new one instead of failing; a line in the command's output names the endpoint and its +old and new port. An exact port is rejected only when a listener already answers on it or the port cannot be bound; another stack's saved claim alone never blocks it, and a conflict names the stack that saved the port. Runtime-only service ports are selected by the managed supervisor and are not written to the document. diff --git a/apps/cli/src/command-internal/db-config.integration.test.ts b/apps/cli/src/command-internal/db-config.integration.test.ts index 9e8a09e3ae..927a11b252 100644 --- a/apps/cli/src/command-internal/db-config.integration.test.ts +++ b/apps/cli/src/command-internal/db-config.integration.test.ts @@ -487,6 +487,7 @@ describe("dbConfigResolver (db-url under the stack backend)", () => { stop: unused, restart: unused, }, + startupEndpointChanges: unused, stop: unused, destroy: unused, commands: { run: () => unused }, diff --git a/apps/cli/src/commands/db/dump/dump.integration.test.ts b/apps/cli/src/commands/db/dump/dump.integration.test.ts index af9f38fc3b..61f6d3789f 100644 --- a/apps/cli/src/commands/db/dump/dump.integration.test.ts +++ b/apps/cli/src/commands/db/dump/dump.integration.test.ts @@ -151,6 +151,7 @@ const managedDumpStackApi = (runtime: "native" | "docker") => { stop: Effect.succeed([]), restart: Effect.succeed([]), }, + startupEndpointChanges: Effect.succeed([]), stop: Effect.void, destroy: Effect.succeed({ runtimeCleanup: "complete" as const }), commands: { run: runCommand }, diff --git a/apps/cli/src/commands/db/reset/reset.integration.test.ts b/apps/cli/src/commands/db/reset/reset.integration.test.ts index 4e960d99b1..74677a57b4 100644 --- a/apps/cli/src/commands/db/reset/reset.integration.test.ts +++ b/apps/cli/src/commands/db/reset/reset.integration.test.ts @@ -806,6 +806,7 @@ function mockResetStackApi(opts: { }), restart: Effect.die("unused"), }, + startupEndpointChanges: Effect.die("unused"), stop: Effect.die("unused"), destroy: Effect.die("unused"), commands: { run: () => Effect.die("unused") }, diff --git a/apps/cli/src/commands/db/schema/declarative/generate/generate.integration.test.ts b/apps/cli/src/commands/db/schema/declarative/generate/generate.integration.test.ts index 7344c9711b..dd3bb185fe 100644 --- a/apps/cli/src/commands/db/schema/declarative/generate/generate.integration.test.ts +++ b/apps/cli/src/commands/db/schema/declarative/generate/generate.integration.test.ts @@ -162,6 +162,7 @@ function generateStackApi(workdir: string) { stop: unusedStack, restart: unusedStack, }, + startupEndpointChanges: unusedStack, stop: unusedStack, destroy: unusedStack, commands: { run: unusedStackFn }, diff --git a/apps/cli/src/commands/db/schema/declarative/sync/sync.integration.test.ts b/apps/cli/src/commands/db/schema/declarative/sync/sync.integration.test.ts index 85ebc7eaaa..5b6958433f 100644 --- a/apps/cli/src/commands/db/schema/declarative/sync/sync.integration.test.ts +++ b/apps/cli/src/commands/db/schema/declarative/sync/sync.integration.test.ts @@ -164,6 +164,7 @@ function syncStackApi(workdir: string, port: number) { stop: unusedSync, restart: unusedSync, }, + startupEndpointChanges: unusedSync, stop: unusedSync, destroy: unusedSync, commands: { run: unusedSyncFn }, diff --git a/apps/cli/src/commands/db/start/start.integration.test.ts b/apps/cli/src/commands/db/start/start.integration.test.ts index debb4eae72..79754724be 100644 --- a/apps/cli/src/commands/db/start/start.integration.test.ts +++ b/apps/cli/src/commands/db/start/start.integration.test.ts @@ -1731,6 +1731,7 @@ describe("db start stack backend", () => { stop: Effect.succeed([]), restart: Effect.succeed([]), }, + startupEndpointChanges: Effect.succeed([]), stop: Effect.void, destroy: Effect.succeed({ runtimeCleanup: "complete" as const }), commands: { run: () => Effect.die("unused") }, diff --git a/apps/cli/src/commands/experimental/stack/prepare/prepare.integration.test.ts b/apps/cli/src/commands/experimental/stack/prepare/prepare.integration.test.ts index 32d6862aa3..f73ad86e5a 100644 --- a/apps/cli/src/commands/experimental/stack/prepare/prepare.integration.test.ts +++ b/apps/cli/src/commands/experimental/stack/prepare/prepare.integration.test.ts @@ -155,6 +155,7 @@ const makeFixture = (root: string, options: FixtureOptions = {}) => { stop: Effect.succeed([]), restart: Effect.succeed([]), }, + startupEndpointChanges: Effect.succeed([]), stop: Effect.void, destroy: Effect.succeed({ runtimeCleanup: "complete" as const }), commands: { run: () => Effect.die("unused") }, diff --git a/apps/cli/src/commands/experimental/stack/start/SIDE_EFFECTS.md b/apps/cli/src/commands/experimental/stack/start/SIDE_EFFECTS.md index 4c6e60d465..e13ec46ae3 100644 --- a/apps/cli/src/commands/experimental/stack/start/SIDE_EFFECTS.md +++ b/apps/cli/src/commands/experimental/stack/start/SIDE_EFFECTS.md @@ -87,10 +87,35 @@ verification, replace the saved configuration of the existing instances; their i and ports are retained. Changed exclusions reuse existing service identities, data, and ports. Removed services remain saved and stopped so including them again can reuse them; a saved stopped instance of a newly included service is reused when its endpoints and versions still match. The -project configuration file is unchanged. A changed endpoint, artifact version, or PostgreSQL major -version fails before modifying the stopped composition, naming the `config.toml` key or -`SUPABASE_*` env var behind the change with its saved and requested values, and suggesting either -reverting it or running the stack's exact `supabase stack destroy` command to recreate it. +project configuration file is unchanged. The requested creations travel into the stack package's own +owner startup: when every incompatible path across the whole composition is a changed endpoint, each +is re-planned there while the new owner alone holds the stack's lease, before it registers endpoint +namespaces from the saved state: as late as practical, just before its own normal endpoint binding +claims the newly requested port, or a freshly chosen automatic one, it saves the updated endpoint +intent with the old port claim dropped, reusing every check a live composition bind already applies. +The rollback covers only this save-and-claim commit, which finishes before the owner serves RPC or +publishes its holder: a failure or interruption there, not only a claim conflict, restores the saved +state and claims as they read before the re-plan, except a claim whose old port another stack took in +the meantime, which is left unclaimed so the next start reports it as a normal port conflict instead +of overlapping that stack's claim; a hard process death in this window is an accepted limitation, and +the next successful start converges the saved state again. A later startup failure, once that commit +succeeds, keeps the committed (consistent) state instead of rolling it back, since an attached client +may already have persisted its own change by then; the next start reuses it. This includes a failure +during the CLI's own database preparation (see First startup and retries below). A concurrent start +attaches to whichever owner wins that race instead of re-planning again. If the +saved stack's owner exits between this command's liveness check and the moment it opens the stack, +the freshly spawned replacement owner boots without the requested creations and this start falls +back to today's rejection; every later start now sees that replacement owner as running and skips +the re-plan too. Recovering means: stop the stack, then start it again. Text +output prints one line per changed endpoint naming its old and new port; JSON and stream-json output +add the same changes to the success payload. Any other incompatible path blocks the re-plan for the +whole composition, even for a member whose own change is purely a changed endpoint: a changed +`config.toml`-backed or env-var-backed setting (such as a changed PostgreSQL major version) still +fails before modifying the stopped composition, naming the key or env var behind the change with its +saved and requested values and suggesting reverting it; a changed catalog-pinned artifact version or +a same-major PostgreSQL build mismatch, which no `config.toml` key or env var controls, instead uses +a plain label with no revert advice. Either way the failure suggests running the stack's exact +`supabase stack destroy` command to recreate it. ## First startup and retries @@ -134,10 +159,11 @@ and warnings written while the spinner is shown appear on their own rows. JSON output returns the stack `id`, its saved `runtime`, `endpoints` keyed by service and endpoint name (protocol, address, port, and URL, matching `stack status`, with no synthetic entries), `lazy_services` listing members that start on their first request (empty with `--eager`), `env` -(the same connection map `stack status --env` exports, present on every success path), and an -empty message. See [`docs/stack-commands.md`](../../../../../docs/stack-commands.md) for an -example. Failures retain typed command errors and package diagnostics. Telemetry state is flushed -after success or failure. +(the same connection map `stack status --env` exports, present on every success path), an +`endpoint_changes` array naming each re-planned endpoint with its old and new port when the start +applied any, and an empty message. See +[`docs/stack-commands.md`](../../../../../docs/stack-commands.md) for an example. Failures retain +typed command errors and package diagnostics. Telemetry state is flushed after success or failure. A rejected configuration change additionally carries `stack_changes` on the JSON/stream-json error envelope: one entry per affected service (a shared setting such as the API port appears once per diff --git a/apps/cli/src/commands/experimental/stack/start/start-export-pointer.integration.test.ts b/apps/cli/src/commands/experimental/stack/start/start-export-pointer.integration.test.ts index ef49fc9753..b867cb135d 100644 --- a/apps/cli/src/commands/experimental/stack/start/start-export-pointer.integration.test.ts +++ b/apps/cli/src/commands/experimental/stack/start/start-export-pointer.integration.test.ts @@ -110,6 +110,7 @@ function makeDatabaseStack(sqlPort: number, credentials: StackCredentials): Stac stop: Effect.die("unused"), restart: Effect.die("unused"), }, + startupEndpointChanges: Effect.die("unused"), stop: Effect.die("unused"), destroy: Effect.die("unused"), commands: { run: () => Effect.die("unused") }, diff --git a/apps/cli/src/commands/experimental/stack/start/start.handler.ts b/apps/cli/src/commands/experimental/stack/start/start.handler.ts index 80ba093d48..4d2ee8c148 100644 --- a/apps/cli/src/commands/experimental/stack/start/start.handler.ts +++ b/apps/cli/src/commands/experimental/stack/start/start.handler.ts @@ -18,7 +18,9 @@ import { import { RuntimeInfo } from "../../../../shared/runtime/runtime-info.service.ts"; import { Effect, FileSystem, Fiber, Option, Path, Redacted, Ref } from "effect"; import { + apiRoute, resolveNativePostgresUser, + type EndpointPortChange, type Observation, type PlannedInstance, type ServiceCreation, @@ -381,6 +383,12 @@ const selectedCreations = ( return !exclusions.includes(capability); }); +/** Names a changed endpoint the way the connection summary does: `api`, or `service.endpoint`. */ +const endpointLabel = (change: EndpointPortChange) => + change.endpoint === "http" && apiRoute(change.service) !== undefined + ? "api" + : `${change.service}.${change.endpoint}`; + const isServing = (status: Pick) => status.lifecycle === "running" && status.health === "healthy"; @@ -456,10 +464,51 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags ? output.info(postgresUser.message) : Effect.void; const configBeforeCreate = - target.id === undefined ? yield* loadStartConfig(target.projectRoot, fs, path) : undefined; + target.id === undefined || !target.hostRunning + ? yield* loadStartConfig(target.projectRoot, fs, path) + : undefined; if (target.id === undefined) yield* ensurePostgresUser; const stateRoot = path.join(settings.supabaseHome, "stacks"); const cacheRoot = path.join(settings.supabaseHome, "cache", "stack"); + const resolveRequested = ( + stackId: string, + config: Effect.Success>["config"], + ) => + Effect.gen(function* () { + const creations = yield* config.creations(stackId).pipe( + Effect.mapError( + (error) => + new StackCommandStartError({ + reason: "invalid-config", + message: error.message, + cause: error, + }), + ), + ); + return yield* Effect.forEach( + selectedCreations(creations, exclusions), + withProjectFunctionsEnv, + ).pipe( + Effect.mapError( + (cause) => + new StackCommandStartError({ + reason: "invalid-config", + message: cause.message, + cause, + }), + ), + ); + }); + // The saved stack's owner is not running, so its requested creations travel into the owner's + // own startup: it re-plans and commits a changed endpoint's port while it alone holds the + // stack's lease, before it registers endpoint namespaces from the saved state. A concurrent + // start attaches to whichever owner wins that race instead of re-planning again. A running + // owner already bound its endpoints at its own startup and keeps today's behavior of applying + // endpoint changes only after stop and start. + const requestedForReplan = + target.id !== undefined && !target.hostRunning && configBeforeCreate !== undefined + ? yield* resolveRequested(target.id, configBeforeCreate.config) + : undefined; const startupComplete = yield* Ref.make(false); const stack = yield* Effect.acquireRelease( target.id === undefined @@ -471,7 +520,13 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags startOwner: true, ...(target.name === undefined ? {} : { name: target.name }), }) - : stackApi.open({ id: target.id, stateRoot, cacheRoot, startOwner: true }), + : stackApi.open({ + id: target.id, + stateRoot, + cacheRoot, + startOwner: true, + ...(requestedForReplan === undefined ? {} : { requestedCreations: requestedForReplan }), + }), (stack) => Ref.get(startupComplete).pipe( Effect.flatMap((complete) => @@ -501,6 +556,9 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags }, currentShellPlatform(), ); + // Set once a fully stopped stack's owner reports the endpoint changes it applied at its own + // startup; a stack that was already running never re-plans, so this stays empty for it. + let endpointChanges: ReadonlyArray = []; const reportReady = (report: Effect.Success>, message: string) => Effect.gen(function* () { const credentials = yield* summaryCredentials(stack.credentials.get, output.warn); @@ -517,6 +575,15 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags .filter(({ activation }) => activation === "lazy") .map(({ service }) => service), env, + ...(endpointChanges.length === 0 + ? {} + : { + endpoint_changes: endpointChanges.map((change) => ({ + endpoint: endpointLabel(change), + from: change.from, + to: change.to, + })), + }), }); if (message.length > 0) yield* output.success(message); yield* output.raw( @@ -545,7 +612,6 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags : status.lifecycle !== "starting" && status.wakeEnabled, ); if (fullyStarted) { - yield* Effect.annotateCurrentSpan({ "stack.path": "already-running" }); yield* Ref.set(startupComplete, true); yield* reportReady( yield* startReport(stack, currentInstances), @@ -561,10 +627,6 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags lifecycle === "running" || lifecycle === "starting" || wakeEnabled, ); if (resumable) { - yield* Effect.annotateCurrentSpan({ - "stack.path": "resume", - "stack.service_count": currentInstances.length, - }); yield* output.info( "Resuming the saved stack services. Run `supabase stack stop`, then `supabase stack start` to apply configuration or service-selection changes.", ); @@ -608,6 +670,9 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags message: "The stack is in a partial lifecycle state", suggestion: "Run supabase stack stop, then supabase stack start to recover the stack.", }); + endpointChanges = yield* stack.startupEndpointChanges.pipe(Effect.mapError(stackError)); + for (const change of endpointChanges) + yield* output.info(`${endpointLabel(change)}: ${change.from} → ${change.to}`); if (target.id !== undefined) yield* ensurePostgresUser; const shadowDatabase = composition.members.length === 0 @@ -621,25 +686,9 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags }); const { config, keys, toml } = configBeforeCreate ?? (yield* loadStartConfig(target.projectRoot, fs, path)); - const creations = yield* config.creations(stack.id).pipe( - Effect.mapError( - (error) => - new StackCommandStartError({ - reason: "invalid-config", - message: error.message, - cause: error, - }), - ), - ); - const requested = yield* Effect.forEach( - selectedCreations(creations, exclusions), - withProjectFunctionsEnv, - ).pipe( - Effect.mapError( - (cause) => - new StackCommandStartError({ reason: "invalid-config", message: cause.message, cause }), - ), - ); + // Reuses the creations already resolved for the re-plan above instead of reading the + // Functions dotenv a second time; only a new or already-running stack has none yet. + const requested = requestedForReplan ?? (yield* resolveRequested(stack.id, config)); if ( requested.some(({ service }) => service === "studio") && !requested.some(({ service }) => service === "rest") @@ -705,7 +754,11 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags ); const initialComposition = composition.members.length === 0; const serviceKindsChanged = !sameKinds(currentInstances, requested); - const planned = yield* stack.composition.plan(requested).pipe(Effect.mapError(stackError)); + // `requested` is the whole desired composition (exclusions already applied), not a partial + // comparison, so an excluded sibling's saved port must not anchor a shared endpoint's port. + const planned = yield* stack.composition + .plan(requested, { requestKind: "complete" }) + .pipe(Effect.mapError(stackError)); const stackIdentity = { id: stack.id, ...(target.name === undefined ? {} : { name: target.name }), @@ -759,12 +812,6 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags const candidate = candidates[0]; if (candidate !== undefined) reuseIds.push(candidate.id); } - yield* Effect.annotateCurrentSpan({ - "stack.path": "start", - "stack.service_count": requested.length, - "stack.initial_composition": initialComposition, - "stack.service_kinds_changed": serviceKindsChanged, - }); const starting = yield* output.task("Starting local Supabase stack..."); const members = yield* stack.composition .supabase(requested, { @@ -838,7 +885,6 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags message: "The stack has no saved credentials", }); if (initialComposition || serviceKindsChanged) { - yield* Effect.annotateCurrentSpan({ "stack.migrations_applied": true }); const migrations = initialComposition ? { workdir: target.projectRoot, @@ -868,7 +914,6 @@ export const stackStart = Effect.fn("experimental.stack.start")(function* (flags (message) => new SeedConfigLoadError({ message }), ); if (hasConfiguredBuckets(context.config)) { - yield* Effect.annotateCurrentSpan({ "stack.storage_seeded": true }); yield* storage.start.pipe( Effect.tapError((error) => starting.fail(error.message)), Effect.mapError(stackError), diff --git a/apps/cli/src/commands/experimental/stack/start/start.integration.test.ts b/apps/cli/src/commands/experimental/stack/start/start.integration.test.ts index ab3f102e4c..8d4adc9af1 100644 --- a/apps/cli/src/commands/experimental/stack/start/start.integration.test.ts +++ b/apps/cli/src/commands/experimental/stack/start/start.integration.test.ts @@ -22,6 +22,8 @@ import { import { postgresVersion } from "@supabase/stack/internal/artifacts"; import { StackError, + type EndpointPortChange, + type PlanOptions, type ServiceCreation, type ServiceCreationInput, type ServiceInstance, @@ -223,6 +225,7 @@ interface MemberStatus { const fakeStack = (compositionStart?: Stack["composition"]["start"]) => { let members: Array = []; + const memberCreations = new Map(); let stopped = 0; let hostStopped = 0; let hostDestroyed = 0; @@ -233,6 +236,26 @@ const fakeStack = (compositionStart?: Stack["composition"]["start"]) => { let activations = new Map(); const memberStatuses = new Map(); const memberReadiness = new Map>(); + /** The endpoint changes a simulated owner startup re-planned, until the next `applyStartupReplan`. */ + let pendingStartupChanges: ReadonlyArray = []; + /** Builds a member from a creation, tracking it for a later `reassignEndpoint`. */ + const buildMember = (id: string, creation: ServiceCreation) => { + memberCreations.set(id, creation); + return instance( + creation, + id, + () => memberStatuses.get(id)?.lifecycle ?? lifecycle, + () => + memberStatuses.get(id)?.wakeEnabled ?? + (lifecycle === "running" && activations.get(id) === "lazy"), + () => { + const status = memberStatuses.get(id); + if (status?.health !== undefined) return status.health; + return (status?.lifecycle ?? lifecycle) === "running" ? "healthy" : undefined; + }, + () => memberReadiness.get(id) ?? Effect.void, + ); + }; let savedCredentials: StackCredentials = { jwtSecret: DEFAULT_LOCAL_JWT_SECRET, postgresRootKey: DEFAULT_POSTGRES_ROOT_KEY, @@ -302,20 +325,7 @@ const fakeStack = (compositionStart?: Stack["composition"]["start"]) => { previous !== undefined && options?.reuseIds?.includes(previous.id) ? previous.id : `${creation.service}-member-${composed}`; - return instance( - requireConcreteCreation(creation), - id, - () => memberStatuses.get(id)?.lifecycle ?? lifecycle, - () => - memberStatuses.get(id)?.wakeEnabled ?? - (lifecycle === "running" && activations.get(id) === "lazy"), - () => { - const status = memberStatuses.get(id); - if (status?.health !== undefined) return status.health; - return (status?.lifecycle ?? lifecycle) === "running" ? "healthy" : undefined; - }, - () => memberReadiness.get(id) ?? Effect.void, - ); + return buildMember(id, requireConcreteCreation(creation)); }); activations = new Map( members.map(({ id, service }) => [ @@ -327,7 +337,7 @@ const fakeStack = (compositionStart?: Stack["composition"]["start"]) => { }), // Delegates to the production planner so paths/shared-API-port normalization match what // `packages/stack` actually reports, instead of a hand-rolled approximation. - plan: (creations: ReadonlyArray) => + plan: (creations: ReadonlyArray, options?: PlanOptions) => Effect.forEach(members, (member) => member.status.pipe(Effect.map(({ config }) => ({ id: member.id, creation: config }))), ).pipe( @@ -344,6 +354,7 @@ const fakeStack = (compositionStart?: Stack["composition"]["start"]) => { }, }, creations, + options, ), ), ), @@ -366,6 +377,7 @@ const fakeStack = (compositionStart?: Stack["composition"]["start"]) => { }), restart: Effect.succeed([]), }, + startupEndpointChanges: Effect.sync(() => pendingStartupChanges), stop: Effect.sync(() => { hostStopped += 1; }), @@ -392,6 +404,10 @@ const fakeStack = (compositionStart?: Stack["composition"]["start"]) => { get composed() { return composed; }, + /** Whether a detached owner is still up, the way a live stack's would be between calls. */ + get hostRunning() { + return lifecycle === "running"; + }, get catalogApplied() { return catalogApplied; }, @@ -411,6 +427,74 @@ const fakeStack = (compositionStart?: Stack["composition"]["start"]) => { setMemberReadiness(id: string, ready: Effect.Effect) { memberReadiness.set(id, ready); }, + /** Simulates a pre-launch endpoint replan: a stopped stack's saved member adopts new endpoints. */ + reassignEndpoint(service: ServiceCreation["service"], endpoints: ServiceCreation["endpoints"]) { + const member = members.find((entry) => entry.service === service); + const current = member === undefined ? undefined : memberCreations.get(member.id); + if (member === undefined || current === undefined) return; + members = members.map((entry) => + entry.id === member.id + ? buildMember(member.id, { ...current, endpoints } as ServiceCreation) + : entry, + ); + }, + /** + * Simulates the owner's own startup re-plan: every saved member whose only incompatible path + * is its own endpoint adopts the requested endpoints, the way a real owner boot applies the + * saved intent to every such member, shared "api" key or dedicated. REST's own change is + * recorded for `startupEndpointChanges`, since that is the only one whose live port existing + * tests assert. + */ + applyStartupReplan(requested: ReadonlyArray) { + // Cleared on every simulated owner startup, not only when a port actually changes, so an + // unchanged restart after a reported change does not keep reporting it. + pendingStartupChanges = []; + // Reads each member's live creation the same way `plan()` does: a member's own `restart` + // can hold a newer creation (e.g. a changed artifact version) in its own closure, which + // `memberCreations` never sees again after construction. + const liveCreations = new Map( + members.map((member) => [member.id, Effect.runSync(member.status).config]), + ); + const planned = planSupabaseComposition( + { + instances: members.map(({ id }) => ({ id, creation: liveCreations.get(id)! })), + composition: { + members: members.map(({ id }) => ({ id, activation: activations.get(id) ?? "eager" })), + dependencies: [], + }, + }, + requested, + { requestKind: "complete" }, + ); + // A changed endpoint only re-plans when every other incompatible path on every member is + // also a pure endpoint reassignment; a changed database version (or any other non-endpoint + // path) blocks the whole re-plan, the way `planEndpointReplan` does. + const purelyEndpointChanges = planned.every( + (entry) => + !entry.member || + entry.change !== "incompatible" || + entry.paths.every((path) => path === "endpoints" || path.startsWith("endpoints.")), + ); + if (!purelyEndpointChanges) return; + // The fake never resolves an automatic port for real, so it stands in the same canned + // value `instance()`'s own status reports for REST, letting a change show up either way. + const resolvedPort = (port: number | "auto" | undefined) => (port === "auto" ? 23457 : port); + for (const entry of planned) { + if (!entry.member || entry.change !== "incompatible") continue; + const request = requested.find((creation) => creation.service === entry.service); + if (request === undefined) continue; + if (entry.service === "rest" && request.service === "rest") { + const current = liveCreations.get(entry.id); + const from = resolvedPort( + current?.service === "rest" ? current.endpoints?.http?.port : undefined, + ); + const to = resolvedPort(request.endpoints?.http?.port); + if (from !== undefined && to !== undefined && from !== to) + pendingStartupChanges = [{ service: "rest", endpoint: "http", from, to }]; + } + this.reassignEndpoint(request.service, request.endpoints); + } + }, }; }; @@ -430,12 +514,16 @@ const layers = ( projectRoot: root, ...(existing ? { id: fixture.stack.id } : {}), runtime: "native" as const, - hostRunning: false, + hostRunning: existing && fixture.hostRunning, }), }); const api = Layer.succeed(StackApi, { create: () => Effect.succeed(fixture.stack), - open: () => Effect.succeed(fixture.stack), + open: (options) => { + if (options.requestedCreations !== undefined) + fixture.applyStartupReplan(options.requestedCreations); + return Effect.succeed(fixture.stack); + }, discover: () => Effect.succeed([]), find: () => Effect.die("identity not used"), }); @@ -1223,122 +1311,212 @@ describe("experimental stack start", () => { }).pipe(Effect.provide(BunServices.layer)), ); - it.live("rejects a changed endpoint after the stack is stopped", () => - Effect.gen(function* () { - const fs = yield* FileSystem.FileSystem; - const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-changed-port-" }); - yield* fs.makeDirectory(`${root}/supabase`, { recursive: true }); - yield* fs.writeFileString(`${root}/supabase/config.toml`, 'project_id = "changed-port"\n'); - const fixture = fakeStack(); - yield* stackStart(flags()).pipe(Effect.provide(layers(root, fixture))); + it.live( + "re-plans a changed endpoint on a stopped stack instead of rejecting it, prints the change, and reports none on the next unchanged restart", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-changed-port-" }); + yield* fs.makeDirectory(`${root}/supabase`, { recursive: true }); + yield* fs.writeFileString(`${root}/supabase/config.toml`, 'project_id = "changed-port"\n'); + const onlyRest = [ + "auth", + "realtime", + "storage", + "functions", + "studio", + "mail", + "analytics", + "pooler", + ]; + const fixture = fakeStack(); + yield* stackStart(flags(onlyRest)).pipe(Effect.provide(layers(root, fixture))); + const restId = fixture.members.find(({ service }) => service === "rest")?.id; - yield* fs.writeFileString( - `${root}/supabase/config.toml`, - 'project_id = "changed-port"\n[api]\nport = 54999\n', - ); - yield* fixture.stack.composition.stop; - const error = yield* stackStart(flags()).pipe( - Effect.provide(layers(root, fixture)), - Effect.flip, - ); + yield* fs.writeFileString( + `${root}/supabase/config.toml`, + 'project_id = "changed-port"\n[api]\nport = 54999\n', + ); + yield* fixture.stack.composition.stop; + const output = mockOutput(); + yield* stackStart(flags(onlyRest)).pipe(Effect.provide(layers(root, fixture, output))); - expect(error).toMatchObject({ - reason: "invalid-config", - message: expect.stringContaining("[api] port: saved automatic, requested 54999"), - suggestion: expect.stringContaining( - `supabase stack destroy --stack-id ${fixture.stack.id}`, - ), - }); - expect(fixture.composed).toBe(1); - }).pipe(Effect.provide(BunServices.layer)), + const rest = fixture.members.find(({ service }) => service === "rest"); + expect(rest?.id).toBe(restId); + expect(fixture.composed).toBe(2); + expect(output.messages.map(({ message }) => message)).toContain("api: 23457 → 54999"); + + // A later stop/start with the now-saved config unchanged reports no endpoint change, + // rather than repeating the one applied by the previous re-plan. + yield* fixture.stack.composition.stop; + const secondOutput = mockOutput(); + yield* stackStart(flags(onlyRest)).pipe( + Effect.provide(layers(root, fixture, secondOutput)), + ); + expect(secondOutput.messages.some(({ message }) => message.includes("→"))).toBe(false); + }).pipe(Effect.provide(BunServices.layer)), ); - it.live("names the config key and both values when a saved port changes", () => + it.live("does not re-plan endpoints while the stack is already running", () => Effect.gen(function* () { const fs = yield* FileSystem.FileSystem; - const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-port-config-" }); + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-running-no-replan-" }); yield* fs.makeDirectory(`${root}/supabase`, { recursive: true }); yield* fs.writeFileString( `${root}/supabase/config.toml`, - 'project_id = "port-config"\n[api]\nport = 54321\n', + 'project_id = "running-no-replan"\n', ); const fixture = fakeStack(); yield* stackStart(flags()).pipe(Effect.provide(layers(root, fixture))); yield* fs.writeFileString( `${root}/supabase/config.toml`, - 'project_id = "port-config"\n[api]\nport = 54999\n', - ); - yield* fixture.stack.composition.stop; - const error = yield* stackStart(flags()).pipe( - Effect.provide(layers(root, fixture)), - Effect.flip, + 'project_id = "running-no-replan"\n[api]\nport = 54999\n', ); - - expect(error).toMatchObject({ - reason: "invalid-config", - message: expect.stringContaining("[api] port: saved 54321, requested 54999"), + // A fully running stack reports "already running" and returns before ever comparing the + // saved composition against the request, so a changed endpoint neither re-plans nor + // rejects; this layer dies if a running stack ever passed requested creations into `open`, + // which only a stopped stack's re-plan does. + const base = layers(root, fixture); + const api = Layer.succeed(StackApi, { + create: () => Effect.die("unused"), + open: (options) => + options.requestedCreations === undefined + ? Effect.succeed(fixture.stack) + : Effect.die("a running stack must not pass requested creations to re-plan"), + discover: () => Effect.succeed([]), + find: () => Effect.die("identity not used"), }); + const result = yield* stackStart(flags()).pipe(Effect.provide(Layer.merge(base, api))); + expect(result).toBe(fixture.stack.id); + expect(fixture.composed).toBe(1); }).pipe(Effect.provide(BunServices.layer)), ); - it.live("names a dedicated (non-shared) port's own config key", () => - Effect.gen(function* () { - const fs = yield* FileSystem.FileSystem; - const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-dedicated-port-" }); - yield* fs.makeDirectory(`${root}/supabase`, { recursive: true }); - yield* fs.writeFileString(`${root}/supabase/config.toml`, 'project_id = "dedicated-port"\n'); - const fixture = fakeStack(); - yield* stackStart(flags()).pipe(Effect.provide(layers(root, fixture))); + it.live( + "names the config key and both values when a saved port changes alongside a database version", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-port-config-" }); + yield* fs.makeDirectory(`${root}/supabase`, { recursive: true }); + yield* fs.writeFileString( + `${root}/supabase/config.toml`, + 'project_id = "port-config"\n[api]\nport = 54321\n', + ); + const fixture = fakeStack(); + yield* stackStart(flags()).pipe(Effect.provide(layers(root, fixture))); - yield* fs.writeFileString( - `${root}/supabase/config.toml`, - 'project_id = "dedicated-port"\n[studio]\nport = 12345\n', - ); - yield* fixture.stack.composition.stop; - const error = yield* stackStart(flags()).pipe( - Effect.provide(layers(root, fixture)), - Effect.flip, - ); + // A changed port alone on a stopped stack now re-plans instead of failing; pairing it + // with a database version change, which never re-plans, keeps this error path reachable. + yield* fs.writeFileString( + `${root}/supabase/config.toml`, + 'project_id = "port-config"\n[db]\nmajor_version = 15\n[api]\nport = 54999\n', + ); + yield* fixture.stack.composition.stop; + const error = yield* stackStart(flags()).pipe( + Effect.provide(layers(root, fixture)), + Effect.flip, + ); - expect(error).toMatchObject({ - reason: "invalid-config", - message: expect.stringContaining("[studio] port: saved automatic, requested 12345"), - }); - // A dedicated port only affects its own service, unlike the shared API port. - expect(error).toBeInstanceOf(StackCommandStartError); - if (error instanceof StackCommandStartError) - expect(error.message).not.toContain("[api] port"); - }).pipe(Effect.provide(BunServices.layer)), + expect(error).toMatchObject({ + reason: "invalid-config", + message: expect.stringContaining("[api] port: saved 54321, requested 54999"), + }); + }).pipe(Effect.provide(BunServices.layer)), ); - // The production planner (`fixedApiPorts`/`withSharedApiPort`) normalizes a requested - // automatic shared-API port to the composition's already-fixed value whenever one exists, so a - // saved fixed port going back to automatic in `config.toml` reuses the saved port rather than - // failing. This locks down that non-obvious compatible case: it is not an incompatible path. it.live( - "accepts a shared API port going from fixed back to automatic, reusing the saved port", + "names a dedicated (non-shared) port's own config key alongside a database version change", () => Effect.gen(function* () { const fs = yield* FileSystem.FileSystem; - const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-api-to-auto-" }); + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-dedicated-port-" }); yield* fs.makeDirectory(`${root}/supabase`, { recursive: true }); yield* fs.writeFileString( `${root}/supabase/config.toml`, - 'project_id = "api-to-auto"\n[api]\nport = 54321\n', + 'project_id = "dedicated-port"\n', ); const fixture = fakeStack(); yield* stackStart(flags()).pipe(Effect.provide(layers(root, fixture))); - yield* fs.writeFileString(`${root}/supabase/config.toml`, 'project_id = "api-to-auto"\n'); + // A changed port alone on a stopped stack now re-plans instead of failing; pairing it + // with a database version change, which never re-plans, keeps this error path reachable. + yield* fs.writeFileString( + `${root}/supabase/config.toml`, + 'project_id = "dedicated-port"\n[db]\nmajor_version = 15\n[studio]\nport = 12345\n', + ); yield* fixture.stack.composition.stop; - // Does not throw: the planner treats this as compatible (`change: "unchanged"`), not an - // incompatible path to report. Reusing the already-bound port for the resumed instance is - // `packages/stack`'s own concern, not asserted here. + const error = yield* stackStart(flags()).pipe( + Effect.provide(layers(root, fixture)), + Effect.flip, + ); + + expect(error).toMatchObject({ + reason: "invalid-config", + message: expect.stringContaining("[studio] port: saved automatic, requested 12345"), + }); + // A dedicated port only affects its own service, unlike the shared API port. + expect(error).toBeInstanceOf(StackCommandStartError); + if (error instanceof StackCommandStartError) + expect(error.message).not.toContain("[api] port"); + }).pipe(Effect.provide(BunServices.layer)), + ); + + it.live( + "reports a shared API port transitioning from automatic to fixed alongside a database version change", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-api-to-fixed-" }); + yield* fs.makeDirectory(`${root}/supabase`, { recursive: true }); + yield* fs.writeFileString(`${root}/supabase/config.toml`, 'project_id = "api-to-fixed"\n'); + const fixture = fakeStack(); yield* stackStart(flags()).pipe(Effect.provide(layers(root, fixture))); + + // A changed port alone on a stopped stack now re-plans instead of failing; pairing it + // with a database version change, which never re-plans, keeps this error path reachable. + yield* fs.writeFileString( + `${root}/supabase/config.toml`, + 'project_id = "api-to-fixed"\n[db]\nmajor_version = 15\n[api]\nport = 54999\n', + ); + yield* fixture.stack.composition.stop; + const error = yield* stackStart(flags()).pipe( + Effect.provide(layers(root, fixture)), + Effect.flip, + ); + + expect(error).toMatchObject({ + reason: "invalid-config", + message: expect.stringContaining("[api] port: saved automatic, requested 54999"), + }); }).pipe(Effect.provide(BunServices.layer)), ); + // A saved fixed shared-API port going back to automatic is an incompatible endpoint path, like + // any other changed port, since every sibling sharing that port reports it; this locks down + // that it re-plans successfully rather than rejecting the start, the way a single service's own + // port change does. + it.live("accepts a shared API port going from fixed back to automatic by re-planning it", () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-api-to-auto-" }); + yield* fs.makeDirectory(`${root}/supabase`, { recursive: true }); + yield* fs.writeFileString( + `${root}/supabase/config.toml`, + 'project_id = "api-to-auto"\n[api]\nport = 54321\n', + ); + const fixture = fakeStack(); + yield* stackStart(flags()).pipe(Effect.provide(layers(root, fixture))); + + yield* fs.writeFileString(`${root}/supabase/config.toml`, 'project_id = "api-to-auto"\n'); + yield* fixture.stack.composition.stop; + // Does not throw: every sibling sharing the "api" port reports a purely-endpoint + // incompatible path, so the whole composition re-plans instead of rejecting the start. + yield* stackStart(flags()).pipe(Effect.provide(layers(root, fixture))); + }).pipe(Effect.provide(BunServices.layer)), + ); + it.live( "collects simultaneous database-version and port changes into one error with plural revert wording", () => @@ -1375,31 +1553,39 @@ describe("experimental stack start", () => { }).pipe(Effect.provide(BunServices.layer)), ); - it.live("names the env var override when SUPABASE_*_PORT set the saved port", () => - Effect.gen(function* () { - const fs = yield* FileSystem.FileSystem; - const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-port-env-" }); - yield* fs.makeDirectory(`${root}/supabase`, { recursive: true }); - yield* fs.writeFileString(`${root}/supabase/config.toml`, 'project_id = "port-env"\n'); - const fixture = fakeStack(); - yield* withEnvVar( - "SUPABASE_API_PORT", - "54321", - stackStart(flags()).pipe(Effect.provide(layers(root, fixture))), - ); + it.live( + "names the env var override when SUPABASE_*_PORT set the saved port alongside a database version change", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-port-env-" }); + yield* fs.makeDirectory(`${root}/supabase`, { recursive: true }); + yield* fs.writeFileString(`${root}/supabase/config.toml`, 'project_id = "port-env"\n'); + const fixture = fakeStack(); + yield* withEnvVar( + "SUPABASE_API_PORT", + "54321", + stackStart(flags()).pipe(Effect.provide(layers(root, fixture))), + ); - yield* fixture.stack.composition.stop; - const error = yield* withEnvVar( - "SUPABASE_API_PORT", - "54999", - stackStart(flags()).pipe(Effect.provide(layers(root, fixture)), Effect.flip), - ); + // A changed port alone on a stopped stack now re-plans instead of failing; pairing it + // with a database version change, which never re-plans, keeps this error path reachable. + yield* fs.writeFileString( + `${root}/supabase/config.toml`, + 'project_id = "port-env"\n[db]\nmajor_version = 15\n', + ); + yield* fixture.stack.composition.stop; + const error = yield* withEnvVar( + "SUPABASE_API_PORT", + "54999", + stackStart(flags()).pipe(Effect.provide(layers(root, fixture)), Effect.flip), + ); - expect(error).toMatchObject({ - reason: "invalid-config", - message: expect.stringContaining("SUPABASE_API_PORT: saved 54321, requested 54999"), - }); - }).pipe(Effect.provide(BunServices.layer)), + expect(error).toMatchObject({ + reason: "invalid-config", + message: expect.stringContaining("SUPABASE_API_PORT: saved 54321, requested 54999"), + }); + }).pipe(Effect.provide(BunServices.layer)), ); it.live( diff --git a/apps/cli/src/commands/experimental/stack/start/start.native.integration.test.ts b/apps/cli/src/commands/experimental/stack/start/start.native.integration.test.ts index 85363964ec..13d359b766 100644 --- a/apps/cli/src/commands/experimental/stack/start/start.native.integration.test.ts +++ b/apps/cli/src/commands/experimental/stack/start/start.native.integration.test.ts @@ -145,6 +145,45 @@ enabled = false enabled = false `; +/** REST and Auth share the API listener; an absent \`[api] port\` requests an automatic one. */ +const n2ProjectConfig = (apiPort: number | undefined) => ` +project_id = "stack-start-n2-replan" +${apiPort === undefined ? "" : `\n[api]\nport = ${apiPort}\n`} +[auth] +enabled = true + +[realtime] +enabled = false + +[storage] +enabled = false + +[edge_runtime] +enabled = false + +[studio] +enabled = false + +[analytics] +enabled = false + +[db.pooler] +enabled = false + +[local_smtp] +enabled = false +`; +const onlyRestAndAuth = [ + "realtime", + "storage", + "functions", + "studio", + "mail", + "analytics", + "pooler", +]; +const onlyRest = [...onlyRestAndAuth, "auth"]; + const makeLayers = (root: string, apiLayer = liveStackApi, workdir = root) => { const settings = mockCommandSettings({ workdir, supabaseHome: root }); const resolver = stackTargetResolverLayer.pipe( @@ -326,7 +365,7 @@ describe("experimental stack start native lifecycle", () => { ); it.live( - "stops owners after bind and pre-compose config failures across retries", + "stops owners after bind failures across retries, and rejects an invalid config before spawning one", () => Effect.gen(function* () { const fs = yield* FileSystem.FileSystem; @@ -403,6 +442,9 @@ describe("experimental stack start native lifecycle", () => { expect(yield* stack.services.list).toHaveLength(0); expect((yield* stack.composition.describe).members).toHaveLength(0); } + // A stopped stack now reads its project config before reopening the stack, to carry a + // changed endpoint into the owner's own startup; an invalid config rejects here without + // spawning another owner, leaving `ownerPids[2]` unset. yield* fs.writeFileString(path.join(root, "supabase", "config.toml"), "project_id = ["); activeAttempt = 2; const invalidConfigStart = yield* Effect.scoped(Effect.exit(stackStart(flags([])))); @@ -412,11 +454,7 @@ describe("experimental stack start native lifecycle", () => { expect(Option.isSome(configError)).toBe(true); if (Option.isSome(configError)) expect(configError.value).toMatchObject({ reason: "invalid-config" }); - const configFailureOwnerPid = ownerPids[2]; - expect(configFailureOwnerPid).toBeDefined(); - if (configFailureOwnerPid === undefined) - return yield* Effect.die("owner PID missing after config failure"); - expect(yield* ownerHasExited(configFailureOwnerPid)).toBe(true); + expect(ownerPids[2]).toBeUndefined(); const definitions = yield* api.discover({ stateRoot: path.join(root, "stacks") }); expect(definitions).toHaveLength(1); const definition = definitions[0]; @@ -430,7 +468,7 @@ describe("experimental stack start native lifecycle", () => { }); expect(yield* stack.services.list).toHaveLength(0); expect((yield* stack.composition.describe).members).toHaveLength(0); - expect(ownerPids.filter((pid) => pid !== undefined)).toHaveLength(3); + expect(ownerPids.filter((pid) => pid !== undefined)).toHaveLength(2); }), ).pipe(Effect.provide(fixture.layer)); }).pipe(Effect.provide(BunServices.layer)), @@ -497,4 +535,71 @@ describe("experimental stack start native lifecycle", () => { }).pipe(Effect.provide(BunServices.layer)), { timeout: 180_000 }, ); + + it.live( + "re-plans the shared API port after excluding a sibling that previously shared it", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-start-n2-replan-" }); + yield* fs.makeDirectory(path.join(root, "supabase"), { recursive: true }); + // Below every OS ephemeral range, so another test's outbound socket cannot already hold it. + const fixedPort = 24_530; + const locations = (r: string) => ({ + stateRoot: path.join(r, "stacks"), + cacheRoot: path.join(r, "cache"), + }); + yield* fs.writeFileString( + path.join(root, "supabase", "config.toml"), + n2ProjectConfig(fixedPort), + ); + const fixture = makeLayers(root); + yield* Effect.ensuring( + Effect.gen(function* () { + const api = yield* StackApi; + const stackId = yield* stackStart(flags(onlyRestAndAuth)); + const stack = yield* api.open({ id: stackId, ...locations(root) }); + const restBefore = (yield* stack.services.list).find( + (instance) => instance.service === "rest", + ); + if (restBefore === undefined) return yield* Effect.die("REST missing"); + const statusBefore = yield* restBefore.status; + const portBefore = statusBefore.endpoints.find(({ name }) => name === "http")?.port; + expect(portBefore).toBe(fixedPort); + + yield* stack.stop; + // The config drops the fixed API port and now excludes Auth too: REST alone remains, + // so Auth's still-saved fixed port must not anchor REST's shared port back to it. + yield* fs.writeFileString( + path.join(root, "supabase", "config.toml"), + n2ProjectConfig(undefined), + ); + const restartedId = yield* stackStart(flags(onlyRest)); + expect(restartedId).toBe(stackId); + const restarted = yield* api.open({ id: restartedId, ...locations(root) }); + const restAfter = (yield* restarted.services.list).find( + (instance) => instance.service === "rest", + ); + if (restAfter === undefined) return yield* Effect.die("REST missing after restart"); + expect(restAfter.id).toBe(restBefore.id); + const statusAfter = yield* restAfter.status; + const portAfter = statusAfter.endpoints.find(({ name }) => name === "http")?.port; + expect(portAfter).toBeDefined(); + // Automatic selection can legitimately land back on the old fixed port, so the + // meaningful checks are that a replan to automatic actually ran and that its claim + // agrees with the live bind, not that the number differs. + const changes = yield* restarted.startupEndpointChanges; + const restChange = changes.find((change) => change.service === "rest"); + expect(restChange?.from).toBe(fixedPort); + expect(restChange?.to).toBe(portAfter); + }), + Effect.gen(function* () { + const api = yield* StackApi; + yield* destroyTestStacks(api, path.join(root, "stacks"), path.join(root, "cache")); + }), + ).pipe(Effect.provide(fixture.layer)); + }).pipe(Effect.provide(BunServices.layer)), + { timeout: 120_000 }, + ); }); diff --git a/apps/cli/src/commands/experimental/stack/status/status.integration.test.ts b/apps/cli/src/commands/experimental/stack/status/status.integration.test.ts index ffa4890839..09ea063368 100644 --- a/apps/cli/src/commands/experimental/stack/status/status.integration.test.ts +++ b/apps/cli/src/commands/experimental/stack/status/status.integration.test.ts @@ -174,6 +174,7 @@ const makeStack = ( stop: Effect.die("unused"), restart: Effect.die("unused"), }, + startupEndpointChanges: Effect.die("unused"), stop: Effect.die("unused"), destroy: Effect.die("unused"), commands: { diff --git a/apps/cli/src/commands/functions/serve/serve.stack.integration.test.ts b/apps/cli/src/commands/functions/serve/serve.stack.integration.test.ts index 45f486e1cd..5bc1f2667e 100644 --- a/apps/cli/src/commands/functions/serve/serve.stack.integration.test.ts +++ b/apps/cli/src/commands/functions/serve/serve.stack.integration.test.ts @@ -283,6 +283,7 @@ const fixture = ( stop: Effect.die("unused"), restart: Effect.die("unused"), }, + startupEndpointChanges: Effect.die("unused"), stop: Effect.die("unused"), destroy: Effect.die("unused"), commands: { run: () => Effect.die("unused") }, diff --git a/apps/cli/tests/helpers/storage.ts b/apps/cli/tests/helpers/storage.ts index 5a5a1f2479..15057da7e5 100644 --- a/apps/cli/tests/helpers/storage.ts +++ b/apps/cli/tests/helpers/storage.ts @@ -261,6 +261,7 @@ export function buildStorageStackApi( stop: Effect.die("unused"), restart: Effect.die("unused"), }, + startupEndpointChanges: Effect.die("unused"), stop: Effect.die("unused"), destroy: Effect.die("unused"), commands: { run: () => Effect.die("unused") }, diff --git a/docs/adr/0017-simplified-managed-stack-architecture.md b/docs/adr/0017-simplified-managed-stack-architecture.md index f9ee11209f..a8ac174819 100644 --- a/docs/adr/0017-simplified-managed-stack-architecture.md +++ b/docs/adr/0017-simplified-managed-stack-architecture.md @@ -106,7 +106,24 @@ Fresh automatic claims draw from `20000..32767` with a start derived from the stack's project root, identifier, and listener key, and stride `257`, making up to 64 bounded `EADDRINUSE`/`EACCES` attempts per newly selected binding while skipping -durable sibling claims. Sticky values do not migrate. Failed acquisition +durable sibling claims. Sticky values are stable: they do not migrate on their +own, but a stopped stack's own owner startup re-plans an endpoint whose saved +port differs from the current configuration while it alone holds the stack's +lease, before it registers endpoint namespaces from the saved state: as late +as practical, just before its own normal endpoint binding claims the new one, +reusing every check a live composition bind already applies, it saves the +updated endpoint intent with the changed endpoint's old claim dropped. The +rollback covers only this save-and-claim commit, which finishes before the +owner serves RPC or publishes its holder: a failure or interruption there +restores the exact document read before the re-plan, except a claim whose old +port another stack claimed in the meantime, which stays unclaimed so the next +start reports a normal port conflict instead of overlapping that stack's +claim; a hard process death in this window is an accepted limitation, left +for the next successful start to converge. A later startup failure, once that +commit succeeds, keeps the committed (consistent) state instead of rolling it +back, since an attached client may already have persisted its own change by +then; the next start reuses it. A concurrent start attaches to whichever +owner wins the lease instead of re-planning again. Failed acquisition preserves the previous successful arrays, including claims removed or reconfigured by the new definition, and releases attempted sockets; those arrays change only after diff --git a/packages/stack/README.md b/packages/stack/README.md index 5363725a53..1701a8975c 100644 --- a/packages/stack/README.md +++ b/packages/stack/README.md @@ -65,7 +65,7 @@ Functions use a package-provided, self-contained Edge Runtime main service unles On Linux, native Functions project files must be outside `/tmp`: Edge Runtime uses a private filesystem at that path. Docker and Podman mount project files at a separate runtime path. -`open({ id, stateRoot, cacheRoot })` reconnects to a saved stack. The package stores the stack document at `//state.json` and service data at `//data/`. `discover({ stateRoot })` lists saved definitions and port assignments separately from live-owner availability. It skips each entry that cannot be read or decoded and reports it to `onInvalidState(id, error)`; only a failure to read `stateRoot` itself fails discovery. Port allocation skips the same entries. Offline definitions are not live lifecycle observations. +`open({ id, stateRoot, cacheRoot })` reconnects to a saved stack. The package stores the stack document at `//state.json` and service data at `//data/`. `open({ ..., startOwner: true, requestedCreations })` passes those creations into a freshly spawned owner's own startup: while it alone holds the stack's lease, it re-plans a changed endpoint's port against the saved state, and, as late as practical, just before its own normal endpoint binding claims the new port, saves the updated endpoint intent with the old claim dropped, reusing every check a live composition bind already applies; `stack.startupEndpointChanges` reports what that startup applied. The rollback covers only this save-and-claim commit, which finishes before the owner serves RPC or publishes its holder: a failure or interruption there restores the exact document read before the re-plan, except a claim whose old port another stack claimed in the meantime, which stays unclaimed so the next start reports a normal port conflict instead of overlapping that stack's claim. A later startup failure, once that commit succeeds, keeps the committed (consistent) state instead of rolling it back, since an attached client may already have persisted its own change by then; the next start reuses it. An owner already running ignores `requestedCreations` and never re-plans; a concurrent `open` attaches to whichever owner wins the lease. `discover({ stateRoot })` lists saved definitions and port assignments separately from live-owner availability. It skips each entry that cannot be read or decoded and reports it to `onInvalidState(id, error)`; only a failure to read `stateRoot` itself fails discovery. Port allocation skips the same entries. Offline definitions are not live lifecycle observations. `find({ stateRoot, projectRoot, name })` derives the stack ID with the same identity rules as `create` and reads only that stack; `find({ stateRoot, id })` reads a known ID. It returns the saved definition with the live owner's endpoint, if any, or nothing when no such stack is saved. Unlike `discover`, an unreadable state document fails the call instead of being skipped. The Effect entrypoint's `StackId` schema validates an ID before lookup. diff --git a/packages/stack/src/HostProcess.integration.test.ts b/packages/stack/src/HostProcess.integration.test.ts index f6a89455ba..c90f218edf 100644 --- a/packages/stack/src/HostProcess.integration.test.ts +++ b/packages/stack/src/HostProcess.integration.test.ts @@ -10,6 +10,8 @@ import { Layer, Option, Path, + PlatformError, + Redacted, Schema, Stream, Tracer, @@ -18,6 +20,10 @@ import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; import * as HttpClient from "effect/unstable/http/HttpClient"; import * as HttpClientRequest from "effect/unstable/http/HttpClientRequest"; import * as FetchHttpClient from "effect/unstable/http/FetchHttpClient"; +// oxlint-disable-next-line effecttsgo/node-builtin-import -- reads the spawned owner's own argv. +import { execFileSync } from "node:child_process"; +// oxlint-disable-next-line effecttsgo/node-builtin-import -- checks for the launcher's startup payload file left behind. +import { existsSync, readdirSync } from "node:fs"; // oxlint-disable-next-line effecttsgo/node-builtin-import -- integration verifies exact-port reopening. import * as Net from "node:net"; import { fileURLToPath } from "node:url"; @@ -30,6 +36,7 @@ import { type HostAccess, } from "./HostProcess.ts"; import { discover } from "./effect.ts"; +import type { ServiceCreationInput } from "./services/Catalog.ts"; import { watchLeaseRelease } from "../tests/owner.ts"; import * as State from "./State.ts"; @@ -49,6 +56,20 @@ const savedStack = (root: string, stackName: string): State.SavedStack => ({ ports: [], }); +/** The live command line of a running process, read the POSIX or Windows way. */ +const processCommandLine = (pid: number) => + process.platform === "win32" + ? execFileSync( + "powershell", + [ + "-NoProfile", + "-Command", + `(Get-CimInstance Win32_Process -Filter "ProcessId=${pid}").CommandLine`, + ], + { encoding: "utf8" }, + ) + : execFileSync("ps", ["-ww", "-p", String(pid), "-o", "args="], { encoding: "utf8" }); + /** A loopback listener that accepts connections and never answers, counting each one. */ const silentListener = Effect.acquireRelease( Effect.callback< @@ -517,3 +538,149 @@ it.live("keeps the owner secret out of recorded HTTP span attributes", () => }), ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), ); + +it.live("keeps the spawned owner's requested creations out of its process argv", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "host-process-argv-secrecy-" }); + const state = yield* makeTestState(root); + yield* state.save(savedStack(root, "argv-secrecy")); + const marker = `argv-secrecy-password-${root.split("/").at(-1)}`; + const requestedCreations: ReadonlyArray = [ + { + service: "database", + config: { + version: "17", + databasePassword: Redacted.make(marker), + jwtSecret: Redacted.make("argv-secrecy-jwt-secret-with-32-characters"), + jwtExpiry: 3600, + }, + endpoints: { sql: { port: "auto" } }, + }, + ]; + const access = yield* launchHost(state, { + stateRoot: root, + cacheRoot: root, + stackId: "stack", + entrypoint: fixtureEntrypoint, + requestedCreations, + }); + yield* Effect.sync(() => { + expect(processCommandLine(access.endpoint.pid)).not.toContain(marker); + }).pipe(Effect.ensuring(bestEffortShutdown(root)(access))); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), +); + +it.live("deletes the spawned owner's startup payload file once it reports ready", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const root = yield* fs.makeTempDirectoryScoped({ + prefix: "host-process-payload-cleanup-", + }); + const state = yield* makeTestState(root); + const ownerDir = path.dirname(state.ownerLog("stack")); + // The real consumer (`internal/host-process.ts`), not the lightweight test fixture, reads + // and deletes this file; the default entrypoint exercises it. + const access = yield* launchHost(state, { + stateRoot: root, + cacheRoot: root, + stackId: "stack", + register: savedStack(root, "payload-cleanup"), + }); + yield* Effect.sync(() => { + const leftover = existsSync(ownerDir) + ? readdirSync(ownerDir).filter((name) => name.startsWith("startup-")) + : []; + expect(leftover).toEqual([]); + }).pipe(Effect.ensuring(bestEffortShutdown(root)(access))); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), +); + +it.live("deletes the spawned owner's startup payload file even when its startup fails", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const root = yield* fs.makeTempDirectoryScoped({ + prefix: "host-process-payload-cleanup-fail-", + }); + const state = yield* makeTestState(root); + const ownerDir = path.dirname(state.ownerLog("stack")); + // This fixture never reads argv past the stack id, so only the launcher's own release + // deletes the file here, proving the cleanup does not depend on the child consuming it. + const failure = yield* launchHost(state, { + stateRoot: root, + cacheRoot: root, + stackId: "stack", + entrypoint: fileURLToPath(new URL("../tests/failing-owner-fixture.ts", import.meta.url)), + register: savedStack(root, "payload-cleanup-fail"), + }).pipe(Effect.flip); + expect(failure.message).toContain("owner startup failed"); + const leftover = existsSync(ownerDir) + ? readdirSync(ownerDir).filter((name) => name.startsWith("startup-")) + : []; + expect(leftover).toEqual([]); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), +); + +it.live("deletes a startup payload file that fails to write after creating partial content", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const root = yield* fs.makeTempDirectoryScoped({ + prefix: "host-process-payload-write-failure-", + }); + const state = yield* makeTestState(root); + const ownerDir = path.dirname(state.ownerLog("stack")); + // Writes a few bytes of the real payload through the real filesystem, then fails, the + // way a disk running out of space would after creating the file. + const failingWrites = Layer.effect( + FileSystem.FileSystem, + Effect.map(FileSystem.FileSystem, (real) => ({ + ...real, + writeFileString: ( + file: string, + content: string, + options?: Parameters[2], + ) => + file.includes("startup-") + ? real.writeFileString(file, content.slice(0, 4), options).pipe( + Effect.andThen( + Effect.fail( + PlatformError.systemError({ + _tag: "Unknown", + module: "FileSystem", + method: "writeFile", + pathOrDescriptor: file, + description: "injected write failure", + cause: Object.assign(new Error("no space left on device"), { + code: "ENOSPC", + }), + }), + ), + ), + ) + : real.writeFileString(file, content, options), + })), + ); + const failure = yield* launchHost(state, { + stateRoot: root, + cacheRoot: root, + stackId: "stack", + register: savedStack(root, "payload-write-failure"), + }).pipe(Effect.provide(failingWrites), Effect.flip); + expect(failure.message).toContain("injected write failure"); + const leftover = existsSync(ownerDir) + ? readdirSync(ownerDir).filter((name) => name.startsWith("startup-")) + : []; + expect(leftover).toEqual([]); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), +); diff --git a/packages/stack/src/HostProcess.ts b/packages/stack/src/HostProcess.ts index 88f34aa4de..eb799ab948 100644 --- a/packages/stack/src/HostProcess.ts +++ b/packages/stack/src/HostProcess.ts @@ -31,6 +31,7 @@ import { HOST_PROCESS_DISPATCH_SENTINEL, isBunVirtualPath } from "./internal/dis import { failureMessage } from "./internal/failure-message.ts"; import { stackSourceDigest } from "./internal/release.ts"; import { StackRpc } from "./Rpc.ts"; +import { ServiceCreationInput } from "./services/Catalog.ts"; import { SavedStack, type Interface as StateInterface, type StateError } from "./State.ts"; declare const SUPABASE_STACK_BUILD_ID: string | undefined; @@ -337,6 +338,20 @@ export const shutdownHost = Effect.fn("HostProcess.shutdownHost")(function* ( ); }); +/** + * A spawned owner's registration and requested creations carry database, JWT, and root-key + * secrets, and Functions env values. They travel through a private, owner-only file instead of + * argv, which a live process list (`ps`) exposes to any local user for the owner's whole + * lifetime: the launcher writes this payload once, under the owner's own state directory with + * owner-only permissions, and the spawned child reads and deletes it immediately, before it does + * anything else. + */ +export const HostStartupPayload = Schema.Struct({ + register: Schema.optionalKey(SavedStack), + requestedCreations: Schema.optionalKey(Schema.Array(Schema.toCodecJson(ServiceCreationInput))), +}); +export interface HostStartupPayload extends Schema.Schema.Type {} + export interface LaunchOptions { readonly stateRoot: string; readonly cacheRoot: string; @@ -344,6 +359,13 @@ export interface LaunchOptions { readonly entrypoint?: string; /** Registers this definition once the spawned owner holds the lease. */ readonly register?: SavedStack; + /** + * Re-plans changed endpoint ports against the saved state once a spawned owner holds the lease, + * before it registers endpoint namespaces from that state. Only a spawned owner acts on this; a + * launcher that attaches to an already-running owner leaves its endpoints as that owner bound + * them at its own startup. + */ + readonly requestedCreations?: ReadonlyArray; /** Ties a spawned owner to the enclosing scope, whose closure destroys the stack. */ readonly lifeline?: boolean; } @@ -423,127 +445,175 @@ const spawnOwner = Effect.fn("HostProcess.spawnOwner")(function* ( state: StateInterface, options: LaunchOptions, entrypoint: string, -): Effect.fn.Return { +): Effect.fn.Return< + Spawned, + HostProcessError, + Scope.Scope | FileSystem.FileSystem | Path.Path | Crypto.Crypto +> { const fs = yield* FileSystem.FileSystem; const path = yield* Path.Path; + const crypto = yield* Crypto.Crypto; const log = state.ownerLog(options.stackId); + const ownerDir = path.dirname(log); yield* fs - .makeDirectory(path.dirname(log), { recursive: true, mode: 0o700 }) + .makeDirectory(ownerDir, { recursive: true, mode: 0o700 }) .pipe(Effect.mapError((cause) => error("startup", cause))); - const register = - options.register === undefined - ? [] - : [ - yield* Schema.encodeEffect(Schema.fromJsonString(SavedStack))(options.register).pipe( - Effect.mapError((cause) => error("startup", cause)), - ), - ]; + const payload: HostStartupPayload = { + ...(options.register === undefined ? {} : { register: options.register }), + ...(options.requestedCreations === undefined + ? {} + : { requestedCreations: options.requestedCreations }), + }; + // A secret-bearing payload travels through a private, owner-only file instead of argv, which a + // live process list exposes to any local user for the owner's whole lifetime. The child reads + // and deletes this file immediately, before doing anything else; this file's own lifetime spans + // the write below and the whole handshake that follows, so a log-open failure, a spawn failure, + // or this launcher's own cancellation also deletes it, idempotently alongside a child that + // already did. Only the path is computed in the acquire step below, so a write that creates the + // file and then fails still reaches the release. return yield* Effect.acquireUseRelease( - Effect.try({ - try: () => openSync(log, "a+", 0o600), - catch: (cause) => error("startup", `Cannot open the owner log ${log}: ${String(cause)}`), - }), - (descriptor) => + Object.keys(payload).length === 0 + ? Effect.void + : Effect.gen(function* () { + const name = Array.from(yield* crypto.randomBytes(16), (byte) => + byte.toString(16).padStart(2, "0"), + ).join(""); + return path.join(ownerDir, `startup-${name}.json`); + }).pipe(Effect.mapError((cause) => error("startup", cause))), + (payloadFile) => Effect.gen(function* () { - const failure = (cause: unknown, reason?: HostFailureReason) => { - const tail = logTail(descriptor); - return error( - "startup", - `Stack owner failed to start: ${error("startup", cause).message} (owner log: ${log})${tail.length === 0 ? "" : `\n${tail}`}`, - reason, - ); - }; - const spawnFailed = yield* Deferred.make(); + if (payloadFile !== undefined) { + const encoded = yield* Schema.encodeEffect(Schema.fromJsonString(HostStartupPayload))( + payload, + ).pipe(Effect.mapError((cause) => error("startup", cause))); + yield* fs + .writeFileString(payloadFile, encoded, { mode: 0o600 }) + .pipe(Effect.mapError((cause) => error("startup", cause))); + } + const extraArgs = payloadFile === undefined ? [] : [payloadFile]; return yield* Effect.acquireUseRelease( Effect.try({ - try: () => { - const started = spawn( - process.execPath, - [entrypoint, options.stateRoot, options.cacheRoot, options.stackId, ...register], - { - cwd: process.cwd(), - detached: true, - windowsHide: true, - stdio: [ - options.lifeline === true ? "pipe" : "ignore", - descriptor, - descriptor, - "pipe", - ], - }, - ); - started.on("error", (cause) => { - Deferred.doneUnsafe(spawnFailed, Exit.fail(failure(cause))); - }); - return started; - }, - catch: failure, + try: () => openSync(log, "a+", 0o600), + catch: (cause) => + error("startup", `Cannot open the owner log ${log}: ${String(cause)}`), }), - (child) => + (descriptor) => Effect.gen(function* () { - const readiness = child.stdio[3]; - const line = yield* ( - readiness === null || readiness === undefined || !("read" in readiness) - ? Effect.fail(failure("Owner readiness descriptor is unavailable")) - : NodeStream.fromReadable({ evaluate: () => readiness, onError: failure }).pipe( - Stream.decodeText, - Stream.splitLines, - Stream.runHead, - Effect.flatMap( - Option.match({ - onNone: () => - Effect.fail(failure("Owner exited before reporting readiness")), - onSome: Effect.succeed, - }), + const failure = (cause: unknown, reason?: HostFailureReason) => { + const tail = logTail(descriptor); + return error( + "startup", + `Stack owner failed to start: ${error("startup", cause).message} (owner log: ${log})${tail.length === 0 ? "" : `\n${tail}`}`, + reason, + ); + }; + const spawnFailed = yield* Deferred.make(); + return yield* Effect.acquireUseRelease( + Effect.try({ + try: () => { + const started = spawn( + process.execPath, + [ + entrypoint, + options.stateRoot, + options.cacheRoot, + options.stackId, + ...extraArgs, + ], + { + cwd: process.cwd(), + detached: true, + windowsHide: true, + stdio: [ + options.lifeline === true ? "pipe" : "ignore", + descriptor, + descriptor, + "pipe", + ], + }, + ); + started.on("error", (cause) => { + Deferred.doneUnsafe(spawnFailed, Exit.fail(failure(cause))); + }); + return started; + }, + catch: failure, + }), + (child) => + Effect.gen(function* () { + const readiness = child.stdio[3]; + const line = yield* ( + readiness === null || readiness === undefined || !("read" in readiness) + ? Effect.fail(failure("Owner readiness descriptor is unavailable")) + : NodeStream.fromReadable({ + evaluate: () => readiness, + onError: failure, + }).pipe( + Stream.decodeText, + Stream.splitLines, + Stream.runHead, + Effect.flatMap( + Option.match({ + onNone: () => + Effect.fail(failure("Owner exited before reporting readiness")), + onSome: Effect.succeed, + }), + ), + Effect.timeoutOrElse({ + duration: "30 seconds", + orElse: () => + Effect.fail(failure("Timed out waiting for owner readiness")), + }), + ) + ).pipe( + Effect.raceFirst(Deferred.await(spawnFailed)), + Effect.flatMap((text) => + Schema.decodeEffect(Schema.fromJsonString(readyLine))(text).pipe( + Effect.mapError(failure), + ), ), - Effect.timeoutOrElse({ - duration: "30 seconds", - orElse: () => Effect.fail(failure("Timed out waiting for owner readiness")), - }), - ) - ).pipe( - Effect.raceFirst(Deferred.await(spawnFailed)), - Effect.flatMap((text) => - Schema.decodeEffect(Schema.fromJsonString(readyLine))(text).pipe( - Effect.mapError(failure), - ), - ), - Effect.ensuring(Effect.sync(() => readiness?.destroy())), + Effect.ensuring(Effect.sync(() => readiness?.destroy())), + ); + if (line.type === "error") { + yield* awaitExit(child).pipe(Effect.timeout("5 seconds"), Effect.ignore); + if (line.reason === "lease-held" && options.register === undefined) + return { _tag: "LeaseHeld" } as const; + return yield* line.reason === "lease-held" || line.reason === "exists" + ? error("startup", "Stack already exists; use open") + : failure(line.message, line.reason); + } + if (line.endpoint.stackId !== options.stackId) + return yield* failure( + "Owner readiness identity does not match the requested stack", + ); + child.unref(); + if (options.lifeline === true) { + const stdin = child.stdin; + // Ending the lifeline of an owner that already exited reports EPIPE, which changes nothing. + stdin?.on("error", () => undefined); + if ( + stdin !== null && + Predicate.hasProperty(stdin, "unref") && + typeof stdin.unref === "function" + ) + stdin.unref(); + yield* Effect.addFinalizer(() => closeLifeline(child)); + } + return { + _tag: "Ready", + access: { endpoint: line.endpoint, secret: line.secret }, + } as const; + }), + (child, exit) => (Exit.isSuccess(exit) ? Effect.void : terminate(child)), ); - if (line.type === "error") { - yield* awaitExit(child).pipe(Effect.timeout("5 seconds"), Effect.ignore); - if (line.reason === "lease-held" && options.register === undefined) - return { _tag: "LeaseHeld" } as const; - return yield* line.reason === "lease-held" || line.reason === "exists" - ? error("startup", "Stack already exists; use open") - : failure(line.message, line.reason); - } - if (line.endpoint.stackId !== options.stackId) - return yield* failure( - "Owner readiness identity does not match the requested stack", - ); - child.unref(); - if (options.lifeline === true) { - const stdin = child.stdin; - // Ending the lifeline of an owner that already exited reports EPIPE, which changes nothing. - stdin?.on("error", () => undefined); - if ( - stdin !== null && - Predicate.hasProperty(stdin, "unref") && - typeof stdin.unref === "function" - ) - stdin.unref(); - yield* Effect.addFinalizer(() => closeLifeline(child)); - } - return { - _tag: "Ready", - access: { endpoint: line.endpoint, secret: line.secret }, - } as const; }), - (child, exit) => (Exit.isSuccess(exit) ? Effect.void : terminate(child)), + (descriptor) => Effect.sync(() => closeSync(descriptor)), ); }), - (descriptor) => Effect.sync(() => closeSync(descriptor)), + (payloadFile) => + payloadFile === undefined + ? Effect.void + : fs.remove(payloadFile, { force: true }).pipe(Effect.ignore), ); }); diff --git a/packages/stack/src/Owner.ts b/packages/stack/src/Owner.ts index e25cc167c5..472ac3850f 100644 --- a/packages/stack/src/Owner.ts +++ b/packages/stack/src/Owner.ts @@ -116,6 +116,15 @@ export interface Interface { }; readonly setDraining: (draining: boolean) => Effect.Effect; readonly getServing: Effect.Effect; + /** + * Binds each instance's configured endpoints, claiming any not yet bound through the same + * `Ports.acquire` path and checks a normal composition bind uses; already-bound endpoints are + * untouched. Lets a re-planned endpoint's new port claim happen at owner startup, before the + * saved definition's endpoint namespaces otherwise bind only when the composition starts. + */ + readonly claimEndpoints: ( + ids: ReadonlyArray, + ) => Effect.Effect; } export class Service extends Context.Service()("@supabase/stack/Owner") {} @@ -706,6 +715,14 @@ const makeOwner = Effect.fn("Owner.make")(function* (options: OwnerOptions) { }, setDraining: (value) => Ref.set(draining, value), getServing: Ref.get(draining).pipe(Effect.map((isDraining) => !isDraining)), + claimEndpoints: (ids) => + Effect.forEach( + ids, + (id) => orchestrator.get(id).pipe(Effect.flatMap((entry) => entry.bind)), + { + discard: true, + }, + ).pipe(Effect.withSpan("Owner.claimEndpoints")), } satisfies Interface; }); diff --git a/packages/stack/src/Ports.ts b/packages/stack/src/Ports.ts index 8b55842335..109a4f4359 100644 --- a/packages/stack/src/Ports.ts +++ b/packages/stack/src/Ports.ts @@ -52,7 +52,12 @@ const loopbackOccupied = (port: number) => concurrency: "unbounded", }).pipe(Effect.map((answers) => answers.some(Boolean))); -const claimantOf = (stacks: ReadonlyArray, stackId: string, port: number) => +/** The other stack, if any, whose registry claim already covers this port. */ +export const claimantOf = ( + stacks: ReadonlyArray, + stackId: string, + port: number, +) => stacks.find((other) => other.id !== stackId && other.ports.some((claim) => claim.port === port)) ?.id; diff --git a/packages/stack/src/PromiseClient.ts b/packages/stack/src/PromiseClient.ts index e834f6cb37..af6d796a9a 100644 --- a/packages/stack/src/PromiseClient.ts +++ b/packages/stack/src/PromiseClient.ts @@ -418,9 +418,11 @@ export const stackAdapter = (client: Client) => { ), options, ), - plan: (creations, options) => + plan: (creations, planOptions, options) => client.run( - decodeCreations("plan", creations).pipe(Effect.flatMap(handle.composition.plan)), + decodeCreations("plan", creations).pipe( + Effect.flatMap((decoded) => handle.composition.plan(decoded, planOptions)), + ), options, ), }, diff --git a/packages/stack/src/Rpc.ts b/packages/stack/src/Rpc.ts index 9b296509c4..60639e0f13 100644 --- a/packages/stack/src/Rpc.ts +++ b/packages/stack/src/Rpc.ts @@ -152,6 +152,15 @@ export const OwnerRpc = RpcGroup.make( Rpc.make("restartComposition", { success: Schema.Array(Observation), error: StackError }), ); +/** A changed endpoint the owner's own startup re-planned, with its bound port before and after. */ +export const EndpointPortChange = Schema.Struct({ + service: Schema.String, + endpoint: Schema.String, + from: Schema.Int, + to: Schema.Int, +}); +export interface EndpointPortChange extends Schema.Schema.Type {} + /** The private transport contract; lifecycle admission remains in the owner. */ export const StackRpc = OwnerRpc.add( Rpc.make("runCommand", { @@ -164,4 +173,8 @@ export const StackRpc = OwnerRpc.add( payload: { attachmentId: Schema.String, bytes: Schema.NullOr(Schema.Uint8ArrayFromBase64) }, error: StackError, }), + Rpc.make("startupEndpointChanges", { + success: Schema.Array(EndpointPortChange), + error: StackError, + }), ); diff --git a/packages/stack/src/StackHost.integration.test.ts b/packages/stack/src/StackHost.integration.test.ts index 914f877ed2..4620662807 100644 --- a/packages/stack/src/StackHost.integration.test.ts +++ b/packages/stack/src/StackHost.integration.test.ts @@ -32,10 +32,18 @@ import * as Owner from "./Owner.ts"; import { OrchestratorError } from "./Orchestrator.ts"; import { CommandEvent, StackError } from "./Rpc.ts"; import * as State from "./State.ts"; -import { bindControl, makeRuntime } from "./StackHost.ts"; +import { + bindControl, + commitOrRestoreEndpointReplan, + makeRuntime, + runStackHost, + StackHostError, +} from "./StackHost.ts"; import { shutdownOwner } from "../tests/owner.ts"; +import { captureLogs } from "../tests/logs.ts"; import { postgres } from "./Commands.ts"; import * as CommandRunner from "./host/CommandRunner.ts"; +import type { ServiceCreationInput } from "./services/Catalog.ts"; class HostTestError extends Data.TaggedError("HostTestError")<{ readonly message: string }> {} @@ -1056,3 +1064,239 @@ it.live("destroys a stack only after an abandoned composition settles", () => }), ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), ); + +it.live("keeps the re-planned document committed when a later startup step fails", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-host-replan-commit-" }); + const state = yield* stateFor(root); + const fixedPort = 24_613; + const saved: State.SavedStack = { + id: "stack", + lifetime: "detached", + identity: { projectRoot: root, branchContext: "main", stackName: "replan-commit" }, + runtime: "native", + instances: [ + { + id: "rest-1", + creation: { service: "rest", config: {}, endpoints: { http: { port: fixedPort } } }, + }, + ], + composition: { members: [{ id: "rest-1", activation: "eager" }], dependencies: [] }, + ports: [{ key: "api", host: "127.0.0.1", port: fixedPort }], + }; + yield* state.save(saved); + const requestedFixed: ServiceCreationInput = { + service: "rest", + config: {}, + endpoints: { http: { port: fixedPort + 1 } }, + }; + // The claim already committed, under the lease, before the owner ever served RPC or + // published its holder; a later, unrelated startup failure must not undo it, since an + // attached client could have persisted its own change in the meantime. + const failure = yield* runStackHost({ + stateRoot: root, + cacheRoot: root, + stackId: "stack", + requestedCreations: [requestedFixed], + onReady: () => + Effect.fail( + new StackHostError({ operation: "startup", message: "injected late failure" }), + ), + }).pipe(Effect.flip); + expect(failure.message).toContain("injected late failure"); + + const after = yield* state.read("stack"); + expect(after?.ports.find((claim) => claim.key === "api")?.port).toBe(fixedPort + 1); + expect(after?.instances.find(({ id }) => id === "rest-1")?.creation.endpoints).toEqual({ + http: { port: fixedPort + 1 }, + }); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), +); + +it.live( + "does not let a later startup failure overwrite a mutation an attached client made after the re-plan committed", + () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ + prefix: "stack-host-replan-no-overwrite-", + }); + const state = yield* stateFor(root); + const fixedPort = 24_623; + const saved: State.SavedStack = { + id: "stack", + lifetime: "detached", + identity: { projectRoot: root, branchContext: "main", stackName: "replan-no-overwrite" }, + runtime: "native", + instances: [ + { + id: "rest-1", + creation: { service: "rest", config: {}, endpoints: { http: { port: fixedPort } } }, + }, + ], + composition: { members: [{ id: "rest-1", activation: "eager" }], dependencies: [] }, + ports: [{ key: "api", host: "127.0.0.1", port: fixedPort }], + }; + yield* state.save(saved); + const requestedFixed: ServiceCreationInput = { + service: "rest", + config: {}, + endpoints: { http: { port: fixedPort + 1 } }, + }; + // `onReady` runs only after the owner serves RPC and publishes its holder, so building a + // client from the same `access` it receives stands in for a real client that attached in + // that window: its acknowledged mutation must survive the later injected failure, rather + // than being overwritten by a restore of the pre-re-plan document. + const failure = yield* runStackHost({ + stateRoot: root, + cacheRoot: root, + stackId: "stack", + requestedCreations: [requestedFixed], + onReady: (access) => + Effect.scoped( + Effect.gen(function* () { + const client = yield* ownerClient(access); + yield* client.configureComposition({ + members: [{ id: "rest-1", activation: "lazy" }], + dependencies: [], + }); + return yield* new StackHostError({ + operation: "startup", + message: "injected late failure", + }); + }), + ).pipe( + Effect.mapError((cause) => + cause instanceof StackHostError + ? cause + : new StackHostError({ operation: "startup", message: String(cause) }), + ), + Effect.provide(NodeHttpClient.layerNodeHttp), + ), + }).pipe(Effect.flip); + expect(failure.message).toContain("injected late failure"); + + const after = yield* state.read("stack"); + // The re-plan's new port claim stays committed... + expect(after?.instances.find(({ id }) => id === "rest-1")?.creation.endpoints).toEqual({ + http: { port: fixedPort + 1 }, + }); + // ...and so does the attached client's own acknowledged mutation. + expect(after?.composition.members).toEqual([{ id: "rest-1", activation: "lazy" }]); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), +); + +it.live("notes a restore failure in the startup error without losing the original message", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ + prefix: "stack-host-replan-restore-fails-", + }); + const state = yield* stateFor(root); + const fixedPort = 24_629; + const registered: State.SavedStack = { + id: "stack", + lifetime: "detached", + identity: { projectRoot: root, branchContext: "main", stackName: "replan-restore-fails" }, + runtime: "native", + instances: [ + { + id: "rest-1", + creation: { service: "rest", config: {}, endpoints: { http: { port: fixedPort } } }, + }, + ], + composition: { members: [{ id: "rest-1", activation: "eager" }], dependencies: [] }, + ports: [{ key: "api", host: "127.0.0.1", port: fixedPort }], + }; + yield* state.save(registered); + // The restore's own lock is unavailable, deterministically, after the prepared document + // has already committed, without a real port conflict or process crash. + const failingState: State.Interface = { + ...state, + withLock: ( + _effect: Effect.Effect, + ): Effect.Effect => + Effect.fail(new State.StateError({ operation: "withLock", message: "lock unavailable" })), + }; + const logs: Array = []; + const result = yield* commitOrRestoreEndpointReplan( + failingState, + registered, + [{ key: "api" }], + Effect.fail( + new StackHostError({ operation: "startup", message: "injected claim failure" }), + ), + ).pipe(Effect.flip, Effect.provide(captureLogs(["Error"])(logs))); + + expect(result.message).toContain("injected claim failure"); + expect(result.message).toContain("could not be restored"); + expect(logs.some((line) => line.includes("Restoring the saved endpoint state failed"))).toBe( + true, + ); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), +); + +it.live( + "restores the saved document when an external interrupt lands during the claim after the prepared document was saved", + () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-host-replan-interrupt-" }); + const state = yield* stateFor(root); + const fixedPort = 24_637; + const registered: State.SavedStack = { + id: "stack", + lifetime: "detached", + identity: { projectRoot: root, branchContext: "main", stackName: "replan-interrupt" }, + runtime: "native", + instances: [ + { + id: "rest-1", + creation: { service: "rest", config: {}, endpoints: { http: { port: fixedPort } } }, + }, + ], + composition: { members: [{ id: "rest-1", activation: "eager" }], dependencies: [] }, + ports: [{ key: "api", host: "127.0.0.1", port: fixedPort }], + }; + yield* state.save(registered); + const prepared: State.SavedStack = { + ...registered, + instances: [ + { + id: "rest-1", + creation: { service: "rest", config: {}, endpoints: { http: { port: "auto" } } }, + }, + ], + ports: [], + }; + const entered = yield* Deferred.make(); + // Mirrors the real sequence: the prepared document, with the old claim already dropped, + // is saved first, then the claim itself pauses, standing in for the window an external + // interrupt can land in before the owner's own endpoint binding ever completes. + const commit = state.save(prepared).pipe( + Effect.mapError( + (cause) => new StackHostError({ operation: "startup", message: cause.message }), + ), + Effect.andThen(Deferred.succeed(entered, undefined)), + Effect.andThen(Effect.never), + ); + const claiming = yield* Effect.forkScoped( + commitOrRestoreEndpointReplan(state, registered, [{ key: "api" }], commit), + ); + yield* Deferred.await(entered); + yield* Fiber.interrupt(claiming); + + const restored = yield* state.read("stack"); + expect(restored?.ports).toEqual(registered.ports); + expect(restored?.instances).toEqual(registered.instances); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))), +); diff --git a/packages/stack/src/StackHost.ts b/packages/stack/src/StackHost.ts index e2571409c3..30ad5c3aa4 100644 --- a/packages/stack/src/StackHost.ts +++ b/packages/stack/src/StackHost.ts @@ -40,12 +40,24 @@ import { } from "./HostProcess.ts"; import { projectSegmentFor } from "./identity/Identity.ts"; import * as Owner from "./Owner.ts"; -import { StackError, stackError, StackRpc, type RunCommandPayload } from "./Rpc.ts"; +import { + StackError, + stackError, + StackRpc, + type EndpointPortChange, + type RunCommandPayload, +} from "./Rpc.ts"; import { makeHostGateway } from "./runtime/Container.ts"; import * as State from "./State.ts"; import { sweepOrphans } from "./Sweep.ts"; import { makeCommandAttachments } from "./host/CommandAttachments.ts"; import * as CommandRunner from "./host/CommandRunner.ts"; +import { + prepareEndpointReplan, + reportedEndpointChanges, + restoreFailedEndpointReplan, +} from "./composition/EndpointReplan.ts"; +import type { ServiceCreationInput } from "./services/Catalog.ts"; export interface StackHostOptions { readonly stateRoot: string; @@ -53,6 +65,12 @@ export interface StackHostOptions { readonly stackId: string; /** Registers this definition once the owner holds the lease; the stack must not exist. */ readonly register?: State.SavedStack; + /** + * Re-plans changed endpoint ports against the saved state before the owner registers endpoint + * namespaces from it, while this process alone holds the lease; the owner's own normal endpoint + * binding then claims them. + */ + readonly requestedCreations?: ReadonlyArray; readonly release?: string; readonly onReady?: (access: HostAccess) => Effect.Effect; } @@ -166,6 +184,7 @@ export const makeRuntime = Effect.fn("StackHost.makeRuntime")( access: HostAccess, server: HttpServer.HttpServer["Service"], closeConnections: Effect.Effect, + startupEndpointChanges: ReadonlyArray = [], ): Effect.Effect => Effect.gen(function* () { const scope = yield* Scope.Scope; @@ -336,6 +355,7 @@ export const makeRuntime = Effect.fn("StackHost.makeRuntime")( readonly attachmentId: string; readonly bytes: Uint8Array | null; }) => attachments.input(attachmentId, bytes), + startupEndpointChanges: () => Effect.succeed(startupEndpointChanges), }); const rpc = yield* RpcServer.toHttpEffect(StackRpc, { streamBufferSize: 16 }).pipe( Effect.provide(Layer.merge(StackRpc.toLayer(handlers), RpcSerialization.layerNdjson)), @@ -381,6 +401,41 @@ export const makeRuntime = Effect.fn("StackHost.makeRuntime")( type HostEvent = "SIGTERM" | "SIGINT" | "creator-gone"; +/** + * Runs `commit` interruptibly; on any unsuccessful exit, including interruption, restores the + * saved endpoint state uninterruptibly before re-failing, so an external interrupt reaching this + * window still leaves the saved document consistent instead of holding new intents with missing + * or partial claims. If that restore itself fails, logs it and, only for a genuine typed failure + * (never for an interruption, which carries no message to extend), fails with that failure's + * message extended with recovery guidance, instead of letting the restore failure mask it. + */ +export const commitOrRestoreEndpointReplan = ( + state: State.Interface, + registered: State.SavedStack, + changedKeys: ReadonlyArray<{ readonly key: string }>, + commit: Effect.Effect, +): Effect.Effect => + Effect.uninterruptibleMask((restore) => + Effect.gen(function* () { + const exit = yield* restore(commit).pipe(Effect.exit); + if (Exit.isSuccess(exit)) return exit.value; + const restoreExit = yield* restoreFailedEndpointReplan(state, registered, changedKeys).pipe( + Effect.exit, + ); + if (Exit.isSuccess(restoreExit)) return yield* Effect.failCause(exit.cause); + yield* Effect.logError( + "Restoring the saved endpoint state failed after a failed re-plan", + restoreExit.cause, + ); + const failure = Cause.findErrorOption(exit.cause); + if (Option.isNone(failure)) return yield* Effect.failCause(exit.cause); + return yield* new StackHostError({ + ...failure.value, + message: `${failure.value.message} The saved endpoint state could not be restored either: stop the stack, then start it again, or destroy it to recreate it.`, + }); + }), + ); + export const runStackHost = Effect.fn("StackHost.run")( (options: StackHostOptions): Effect.Effect => Effect.scoped( @@ -421,8 +476,23 @@ export const runStackHost = Effect.fn("StackHost.run")( }), ); const started = yield* Effect.gen(function* () { - const saved = yield* state.read(id); - if (saved === undefined) return yield* hostError("startup", "Stack is not registered"); + const registered = yield* state.read(id); + if (registered === undefined) + return yield* hostError("startup", "Stack is not registered"); + // Re-plans a changed endpoint's port against the saved state while this process alone + // holds the lease, before any owner registers its endpoint namespaces from it; a + // concurrent launcher that loses the lease race never reaches here and simply attaches + // to whichever owner wins. `saved` already has the changed endpoints' old claims + // dropped; the document commits this below, as late as practical, and the owner's own + // normal endpoint binding claims the new ones through the same `Ports.acquire` path and + // checks a live composition bind already applies. + const preparation = + options.requestedCreations === undefined + ? undefined + : yield* prepareEndpointReplan(registered, options.requestedCreations).pipe( + Effect.mapError((cause) => hostError("startup", cause)), + ); + const saved = preparation?.saved ?? registered; const control = yield* bindControl(); const dataRootPath = path.join(options.stateRoot, saved.id, "data"); yield* fs.makeDirectory(dataRootPath, { recursive: true }); @@ -457,6 +527,42 @@ export const runStackHost = Effect.fn("StackHost.run")( ).pipe(Layer.provide(Layer.succeed(State.Service, state))), ); const owner = Context.get(services, Owner.Service); + // The document is saved here, as late as practical, because the owner's own endpoint + // binding below reads it back through the same `Ports.acquire` path a live composition + // bind uses. This commits before the owner serves RPC or publishes its holder below, so + // no attached client can race it: a failure here restores the exact document read before + // the re-plan, while this process still alone holds the lease; a hard process death in + // this window is an accepted limitation, and the next successful start converges the + // saved state again. A failure after this point must not trigger that restore, since by + // then a client may already have attached and persisted its own acknowledged change. + const commitReplan: Effect.Effect< + ReadonlyArray, + StackHostError + > = Effect.gen(function* () { + if (preparation !== undefined) + yield* state + .withLock(state.save(saved)) + .pipe(Effect.mapError((cause) => hostError("startup", cause))); + if (preparation === undefined) return []; + return yield* owner.claimEndpoints(preparation.changedInstanceIds).pipe( + Effect.mapError((cause) => hostError("startup", cause)), + Effect.andThen( + state.read(saved.id).pipe( + Effect.mapError((cause) => hostError("startup", cause)), + Effect.map((after) => reportedEndpointChanges(preparation.changedKeys, after)), + ), + ), + ); + }); + const startupEndpointChanges: ReadonlyArray = + preparation === undefined + ? yield* commitReplan + : yield* commitOrRestoreEndpointReplan( + state, + registered, + preparation.changedKeys, + commitReplan, + ); const endpoint: HostEndpoint = { stackId: saved.id, identity: saved.identity, @@ -474,6 +580,7 @@ export const runStackHost = Effect.fn("StackHost.run")( access, control.server, control.closeConnections, + startupEndpointChanges, ).pipe( Effect.provideService( CommandRunner.Service, diff --git a/packages/stack/src/composition/EndpointReplan.integration.test.ts b/packages/stack/src/composition/EndpointReplan.integration.test.ts new file mode 100644 index 0000000000..56f9219dd2 --- /dev/null +++ b/packages/stack/src/composition/EndpointReplan.integration.test.ts @@ -0,0 +1,102 @@ +import { NodeServices } from "@effect/platform-node"; +import { expect, it } from "@effect/vitest"; +import { Context, Deferred, Effect, Fiber, FileSystem, Layer } from "effect"; +import { claimantOf } from "../Ports.ts"; +import * as State from "../State.ts"; +import { restoreFailedEndpointReplan } from "./EndpointReplan.ts"; + +const stateFor = (root: string) => + Effect.gen(function* () { + const context = yield* Layer.build(State.layer({ root })); + return Context.get(context, State.Service); + }); + +const fixedPort = 24_718; + +const registeredStack = (root: string): State.SavedStack => ({ + id: "stack-a", + lifetime: "detached", + identity: { projectRoot: root, branchContext: "main", stackName: "replan-atomic" }, + runtime: "native", + instances: [ + { + id: "rest-1", + creation: { service: "rest", config: {}, endpoints: { http: { port: fixedPort } } }, + }, + ], + composition: { members: [{ id: "rest-1", activation: "eager" }], dependencies: [] }, + ports: [{ key: "api", host: "127.0.0.1", port: fixedPort }], +}); + +/** Claims a port for another stack, through the same registry check the restore uses, only if no other stack already holds it. */ +const claimIfFree = (state: State.Interface, stackId: string, port: number) => + state.withLock( + Effect.gen(function* () { + const current: State.SavedStack = (yield* state.read(stackId)) ?? { + id: stackId, + lifetime: "detached", + identity: { projectRoot: "/tmp", branchContext: "main", stackName: stackId }, + runtime: "native", + instances: [], + composition: { members: [], dependencies: [] }, + ports: [], + }; + const others = yield* state.claims; + if (claimantOf(others, stackId, port) !== undefined) return false; + yield* state.save({ + ...current, + ports: [...current.ports, { key: "other", host: "127.0.0.1", port }], + }); + return true; + }), + ); + +/** The other stack's claim commits before the restore takes the registry lock, so the restore's read of other stacks' claims inside that lock already sees it. */ +it.live( + "restores the saved document without overlapping a claim another stack committed first", + () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "endpoint-replan-atomic-" }); + const stateA = yield* stateFor(root); + const stateB = yield* stateFor(root); + + // `registered` is the pre-re-plan snapshot kept only for the restore to fall back to; the + // stack's own live document already has the touched "api" claim dropped, matching the + // real mid-re-plan state the owner's startup leaves behind before it fails. + const registered = registeredStack(root); + yield* stateA.save({ + ...registered, + ports: registered.ports.filter((claim) => claim.key !== "api"), + }); + + const otherClaimReady = yield* Deferred.make(); + const otherClaimDone = yield* Deferred.make(); + const otherClaim = yield* Deferred.await(otherClaimReady).pipe( + Effect.andThen(claimIfFree(stateB, "stack-b", fixedPort)), + Effect.tap((claimed) => Deferred.succeed(otherClaimDone, claimed)), + Effect.forkChild, + ); + const wrappedStateA: State.Interface = { + ...stateA, + withLock: (effect: Effect.Effect) => + Deferred.succeed(otherClaimReady, undefined).pipe( + Effect.andThen(Deferred.await(otherClaimDone)), + Effect.andThen(stateA.withLock(effect)), + ), + }; + + yield* restoreFailedEndpointReplan(wrappedStateA, registered, [{ key: "api" }]); + const claimedByOther = yield* Fiber.join(otherClaim); + expect(claimedByOther).toBe(true); + + const afterA = yield* stateA.read("stack-a"); + const afterB = yield* stateA.read("stack-b"); + expect(afterB?.ports).toEqual([{ key: "other", host: "127.0.0.1", port: fixedPort }]); + expect(afterA?.ports).toEqual([]); + expect(afterA?.instances).toEqual(registered.instances); + expect(afterA?.composition).toEqual(registered.composition); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); diff --git a/packages/stack/src/composition/EndpointReplan.ts b/packages/stack/src/composition/EndpointReplan.ts new file mode 100644 index 0000000000..d52f0a464f --- /dev/null +++ b/packages/stack/src/composition/EndpointReplan.ts @@ -0,0 +1,104 @@ +import { Effect, Schema } from "effect"; +import { claimantOf } from "../Ports.ts"; +import type { EndpointPortChange } from "../Rpc.ts"; +import { ServiceCreation, type ServiceCreationInput } from "../services/Catalog.ts"; +import type * as State from "../State.ts"; +import { planEndpointReplan } from "./Supabase.ts"; + +/** One changed endpoint's claim key and its port before the change, for reporting and restore. */ +export interface ChangedEndpointKey { + readonly key: string; + readonly service: string; + readonly endpoint: string; + readonly previousPort: number | undefined; +} + +interface EndpointReplanPreparation { + /** The document to register the owner from: new endpoint intents, changed claims dropped. */ + readonly saved: State.SavedStack; + /** The instances whose endpoints changed, to claim once the owner is built from `saved`. */ + readonly changedInstanceIds: ReadonlyArray; + readonly changedKeys: ReadonlyArray; +} + +/** + * Computes the saved document a stopped stack's owner registers from: endpoint intents updated to + * `requested`, and the changed endpoints' old port claims dropped. The owner's normal endpoint + * binding then claims them, reusing every check a live composition bind already applies; nothing + * here claims a port itself. Returns `undefined` when nothing changed. + */ +export const prepareEndpointReplan = Effect.fn("EndpointReplan.prepare")(function* ( + saved: State.SavedStack, + requested: ReadonlyArray, +) { + const plan = planEndpointReplan(saved, requested); + if (plan === undefined || plan.changes.length === 0) return undefined; + + const touchedKeys = new Set(plan.changes.map((change) => change.key)); + const instances = yield* Effect.forEach(saved.instances, (instance) => { + const endpoints = plan.endpointsByInstance.get(instance.id); + return endpoints === undefined + ? Effect.succeed(instance) + : Schema.decodeUnknownEffect(ServiceCreation)({ ...instance.creation, endpoints }).pipe( + Effect.map((creation) => ({ ...instance, creation })), + ); + }); + + const preparation: EndpointReplanPreparation = { + saved: { + ...saved, + instances, + ports: saved.ports.filter((claim) => !touchedKeys.has(claim.key)), + }, + changedInstanceIds: [...new Set(plan.changes.map((change) => change.id))], + changedKeys: plan.changes.map((change) => ({ + key: change.key, + service: change.service, + endpoint: change.endpoint, + previousPort: change.previousPort, + })), + }; + return preparation; +}); + +/** + * Restores the document read before a failed re-plan, dropping only a touched key whose old port + * another stack has since claimed; that key is left unclaimed, so the next start reports it as a + * normal port conflict instead of overlapping the other stack's claim. The read of other stacks' + * claims, the filtering, and the save run inside one lock, so a concurrent claim on the same port + * is either already visible here or still waiting for this lock, never both unseen and applied. + */ +export const restoreFailedEndpointReplan = Effect.fn("EndpointReplan.restore")(function* ( + state: State.Interface, + registered: State.SavedStack, + changedKeys: ReadonlyArray>, +) { + const touchedKeys = new Set(changedKeys.map((change) => change.key)); + yield* state.withLock( + Effect.gen(function* () { + const others = yield* state.claims; + const ports = registered.ports.filter( + (claim) => + !touchedKeys.has(claim.key) || + claimantOf(others, registered.id, claim.port) === undefined, + ); + yield* state.save({ ...registered, ports }); + }), + ); +}); + +/** + * The changed keys whose newly claimed port differs from its saved value, as reportable endpoint + * changes. An automatic re-plan that resolved back to its previous port is left out: persisting + * the updated intent still matters, but it is not a change worth reporting. + */ +export const reportedEndpointChanges = ( + changedKeys: ReadonlyArray, + after: State.SavedStack | undefined, +): ReadonlyArray => + changedKeys.flatMap((change) => { + const to = after?.ports.find((claim) => claim.key === change.key)?.port; + return change.previousPort === undefined || to === undefined || change.previousPort === to + ? [] + : [{ service: change.service, endpoint: change.endpoint, from: change.previousPort, to }]; + }); diff --git a/packages/stack/src/composition/EndpointReplan.unit.test.ts b/packages/stack/src/composition/EndpointReplan.unit.test.ts new file mode 100644 index 0000000000..e05d26a95f --- /dev/null +++ b/packages/stack/src/composition/EndpointReplan.unit.test.ts @@ -0,0 +1,110 @@ +import { expect, it } from "@effect/vitest"; +import { Effect } from "effect"; +import type { SavedStack } from "../State.ts"; +import type { ServiceCreation } from "../services/Catalog.ts"; +import { prepareEndpointReplan, reportedEndpointChanges } from "./EndpointReplan.ts"; + +const fixedPort = 54_321; + +/** REST and Auth share the "api" claim key at a fixed port, same shape the owner registers from. */ +const savedStack = (): SavedStack => ({ + id: "stack-1", + lifetime: "detached", + identity: { projectRoot: "/project", branchContext: "main", stackName: "stack-1" }, + runtime: "native", + instances: [ + { + id: "rest-1", + creation: { service: "rest", config: {}, endpoints: { http: { port: fixedPort } } }, + }, + { + id: "auth-1", + creation: { service: "auth", config: {}, endpoints: { http: { port: fixedPort } } }, + }, + ], + composition: { + members: [ + { id: "rest-1", activation: "eager" }, + { id: "auth-1", activation: "eager" }, + ], + dependencies: [], + }, + ports: [{ key: "api", host: "127.0.0.1", port: fixedPort }], +}); + +it.effect( + "drops the changed endpoint's old claim and updates the saved intents, without claiming a port", + () => + Effect.gen(function* () { + const saved = savedStack(); + const requested: ReadonlyArray = [ + { service: "rest", config: {}, endpoints: { http: { port: "auto" } } }, + { service: "auth", config: {}, endpoints: { http: { port: "auto" } } }, + ]; + const preparation = yield* prepareEndpointReplan(saved, requested); + if (preparation === undefined) return yield* Effect.die("Expected a replan"); + + // REST and Auth share the "api" claim key, so one change covers both; claiming it once + // (rest-1) is enough, since the later composition-wide bind is idempotent per key. + expect(preparation.changedInstanceIds).toEqual(["rest-1"]); + expect(preparation.changedKeys).toEqual([ + { key: "api", service: "rest", endpoint: "http", previousPort: fixedPort }, + ]); + + // The old "api" claim is dropped so the owner's normal binding claims it fresh. + expect(preparation.saved.ports).toEqual([]); + const rest = preparation.saved.instances.find((instance) => instance.id === "rest-1"); + const auth = preparation.saved.instances.find((instance) => instance.id === "auth-1"); + expect(rest?.creation.endpoints).toEqual({ http: { port: "auto" } }); + expect(auth?.creation.endpoints).toEqual({ http: { port: "auto" } }); + + // Unrelated saved state is left untouched. + expect(preparation.saved.id).toBe(saved.id); + expect(preparation.saved.composition).toEqual(saved.composition); + }), +); + +it.effect("returns undefined when the requested endpoints match the saved ports", () => + Effect.gen(function* () { + const saved = savedStack(); + const requested: ReadonlyArray = [ + { service: "rest", config: {}, endpoints: { http: { port: fixedPort } } }, + { service: "auth", config: {}, endpoints: { http: { port: fixedPort } } }, + ]; + const preparation = yield* prepareEndpointReplan(saved, requested); + expect(preparation).toBeUndefined(); + }), +); + +it.effect("returns undefined when the incompatibility is not a pure endpoint reassignment", () => + Effect.gen(function* () { + const saved = savedStack(); + // An artifact version bump alongside the port change is not a pure endpoint reassignment. + const requested: ReadonlyArray = [ + { service: "rest", version: "2", config: {}, endpoints: { http: { port: "auto" } } }, + { service: "auth", config: {}, endpoints: { http: { port: fixedPort } } }, + ]; + const preparation = yield* prepareEndpointReplan(saved, requested); + expect(preparation).toBeUndefined(); + }), +); + +it("leaves out a changed key whose automatic re-plan resolved back to its previous port", () => { + const changedKeys = [{ key: "api", service: "rest", endpoint: "http", previousPort: fixedPort }]; + const after: SavedStack = { + ...savedStack(), + ports: [{ key: "api", host: "127.0.0.1", port: fixedPort }], + }; + expect(reportedEndpointChanges(changedKeys, after)).toEqual([]); +}); + +it("reports a changed key whose claimed port differs from its previous one", () => { + const changedKeys = [{ key: "api", service: "rest", endpoint: "http", previousPort: fixedPort }]; + const after: SavedStack = { + ...savedStack(), + ports: [{ key: "api", host: "127.0.0.1", port: fixedPort + 1 }], + }; + expect(reportedEndpointChanges(changedKeys, after)).toEqual([ + { service: "rest", endpoint: "http", from: fixedPort, to: fixedPort + 1 }, + ]); +}); diff --git a/packages/stack/src/composition/Supabase.ts b/packages/stack/src/composition/Supabase.ts index 19fb3eeaf9..9b0b95f316 100644 --- a/packages/stack/src/composition/Supabase.ts +++ b/packages/stack/src/composition/Supabase.ts @@ -324,23 +324,51 @@ const compareCreation = ( : { change: "unchanged" }; }; +/** + * A `partial` request compares only the kinds it names, so an absent kind falls back to its saved + * sibling's port; a `complete` request replaces the whole composition, so an absent kind (for + * example one dropped with `--exclude`) must not constrain the shared port it would otherwise + * contribute. + */ +export interface PlanOptions { + readonly requestKind?: "partial" | "complete"; +} + +/** + * The fixed shared API port the request agrees on. A `partial` request falls back to a saved + * sibling's port for a member kind it leaves out; a `complete` request never does, so a request + * that moves every sharing member it carries to automatic is not masked by that kind's own stale + * saved intent, whether that kind is simply absent or deliberately excluded. + */ +const sharedApiPortFor = ( + saved: Pick, + requested: ReadonlyArray, + options: PlanOptions = {}, +): number | undefined => { + const members = new Set(saved.composition.members.map(({ id }) => id)); + const requestedKinds = new Set(requested.map(({ service }) => service)); + const ports = fixedApiPorts([ + ...requested, + ...(options.requestKind === "complete" + ? [] + : saved.instances + .filter(({ id, creation }) => members.has(id) && !requestedKinds.has(creation.service)) + .map(({ creation }) => creation)), + ]); + return ports.size === 1 ? [...ports][0] : undefined; +}; + /** Compares each saved instance of a requested kind with its request, without changing state. */ export const planSupabaseComposition = ( saved: Pick, requested: ReadonlyArray, + options: PlanOptions = {}, ): ReadonlyArray => { const members = new Set(saved.composition.members.map(({ id }) => id)); const memberKinds = new Set( saved.instances.filter(({ id }) => members.has(id)).map(({ creation }) => creation.service), ); - const requestedKinds = new Set(requested.map(({ service }) => service)); - const ports = fixedApiPorts([ - ...requested, - ...saved.instances - .filter(({ id, creation }) => members.has(id) && requestedKinds.has(creation.service)) - .map(({ creation }) => creation), - ]); - const sharedPort = ports.size === 1 ? [...ports][0] : undefined; + const sharedPort = sharedApiPortFor(saved, requested, options); return saved.instances.flatMap(({ id, creation }) => { const request = requested.find(({ service }) => service === creation.service); return request === undefined @@ -356,6 +384,97 @@ export const planSupabaseComposition = ( }); }; +/** The claim key an endpoint binds under: one shared "api" key, or a dedicated one per instance. */ +const endpointKey = (id: string, service: ServiceCreation["service"], name: string): string => + name === "http" && apiRoute(service) !== undefined ? "api" : `${id}:${name}`; + +/** A named endpoint whose port a stopped stack's start can migrate to the requested value. */ +interface EndpointPortChange { + readonly id: string; + readonly service: ServiceCreation["service"]; + readonly endpoint: string; + readonly key: string; + /** The port to request from the port registry: a specific number, or `"auto"`. */ + readonly port: number | "auto"; + /** The endpoint's last bound port, when the registry has a claim for its key. */ + readonly previousPort: number | undefined; +} + +interface EndpointReplan { + readonly changes: ReadonlyArray; + /** Each changed member's endpoint intents to persist once every change claims successfully. */ + readonly endpointsByInstance: ReadonlyMap; +} + +/** + * Plans the endpoint port changes a stopped stack's start can apply instead of failing. + * Returns `undefined` when a member is incompatible for a reason besides a pure endpoint port + * reassignment (an artifact or database version change, or an added or removed endpoint), so the + * caller leaves that case to the existing incompatible-change rejection. + */ +export const planEndpointReplan = ( + saved: Pick, + requested: ReadonlyArray, +): EndpointReplan | undefined => { + // The start re-plan always replaces the whole composition, so an excluded sibling's stale + // saved port must not anchor the shared port this request would otherwise move to automatic. + const planOptions: PlanOptions = { requestKind: "complete" }; + const planned = planSupabaseComposition(saved, requested, planOptions); + const incompatibleCount = planned.filter( + (entry) => entry.member && entry.change === "incompatible", + ).length; + if (incompatibleCount === 0) return undefined; + + const sharedPort = sharedApiPortFor(saved, requested, planOptions); + const changes: Array = []; + const endpointsByInstance = new Map(); + + for (const entry of planned) { + if (!entry.member || entry.change !== "incompatible") continue; + const savedInstance = saved.instances.find(({ id }) => id === entry.id); + const request = requested.find(({ service }) => service === entry.service); + if (savedInstance === undefined || request === undefined) return undefined; + const nonEndpointPaths = entry.paths.filter( + (path) => path !== "endpoints" && !path.startsWith("endpoints."), + ); + if (nonEndpointPaths.length > 0) return undefined; + + const requestedEndpoints = withSharedApiPort(request, sharedPort); + const requestedIntents = { service: request.service, endpoints: requestedEndpoints }; + const savedNames = endpointNames(savedInstance.creation); + const requestedNames = endpointNames(requestedIntents); + if ( + savedNames.length !== requestedNames.length || + !savedNames.every((name) => requestedNames.includes(name)) + ) + return undefined; + + for (const name of savedNames) { + const previous = endpointPort(savedInstance.creation, name); + const next = endpointPort(requestedIntents, name); + if (previous === next) continue; + const key = endpointKey(entry.id, request.service, name); + changes.push({ + id: entry.id, + service: entry.service, + endpoint: name, + key, + port: next, + previousPort: saved.ports.find((claim) => claim.key === key)?.port, + }); + } + endpointsByInstance.set(entry.id, requestedEndpoints); + } + + const seenKeys = new Set(); + const uniqueChanges = changes.filter((change) => { + if (seenKeys.has(change.key)) return false; + seenKeys.add(change.key); + return true; + }); + return { changes: uniqueChanges, endpointsByInstance }; +}; + const compositionError = (message: string, cause?: unknown) => new SupabaseCompositionError({ message, cause }); diff --git a/packages/stack/src/composition/Supabase.unit.test.ts b/packages/stack/src/composition/Supabase.unit.test.ts index 8d0e259c90..16cac6fe84 100644 --- a/packages/stack/src/composition/Supabase.unit.test.ts +++ b/packages/stack/src/composition/Supabase.unit.test.ts @@ -3,8 +3,13 @@ import { Deferred, Effect, Exit, Fiber, Redacted, Ref, Scope, Stream } from "eff import * as TestClock from "effect/testing/TestClock"; import * as Orchestrator from "../Orchestrator.ts"; import { makeService, ServiceError } from "../Service.ts"; +import type { SavedStack } from "../State.ts"; import type { ServiceCreation } from "../services/Catalog.ts"; -import { makeSupabaseComposition, type SupabaseCompositionOperations } from "./Supabase.ts"; +import { + makeSupabaseComposition, + planEndpointReplan, + type SupabaseCompositionOperations, +} from "./Supabase.ts"; /** A registered instance with no runtime behavior beyond an immediate healthy start and stop. */ const makeInstance = ( @@ -167,3 +172,42 @@ it.live( }), ).pipe(Effect.provide(TestClock.layer())), ); + +it("does not let an excluded sibling's stale fixed port mask a shared endpoint's own change", () => { + const fixedPort = 54_321; + const saved: Pick = { + instances: [ + { + id: "rest-1", + creation: { service: "rest", config: {}, endpoints: { http: { port: fixedPort } } }, + }, + { + id: "auth-1", + creation: { service: "auth", config: {}, endpoints: { http: { port: fixedPort } } }, + }, + ], + composition: { + members: [ + { id: "rest-1", activation: "eager" }, + { id: "auth-1", activation: "eager" }, + ], + dependencies: [], + }, + ports: [{ key: "api", host: "127.0.0.1", port: fixedPort }], + }; + // The config dropped [api] port and the start excludes auth, so only REST is requested. + const requested: ReadonlyArray = [ + { service: "rest", config: {}, endpoints: { http: { port: "auto" } } }, + ]; + const plan = planEndpointReplan(saved, requested); + expect(plan?.changes).toEqual([ + { + id: "rest-1", + service: "rest", + endpoint: "http", + key: "api", + port: "auto", + previousPort: fixedPort, + }, + ]); +}); diff --git a/packages/stack/src/effect.integration.test.ts b/packages/stack/src/effect.integration.test.ts index 723bb24fc1..f0c3ff4721 100644 --- a/packages/stack/src/effect.integration.test.ts +++ b/packages/stack/src/effect.integration.test.ts @@ -22,7 +22,9 @@ import { find, open, type DatabaseInstance, + type ServiceCreationInput, type ServiceInstance, + type StackError, } from "./effect.ts"; import { initialization, postgres } from "./Commands.ts"; import { fileURLToPath } from "node:url"; @@ -1065,3 +1067,475 @@ it.live("plans a project's own URL for an input whose supplying member is absent ); }).pipe(Effect.scoped, Effect.provide(layer)), ); + +const replanDatabase = (tag: string) => + ({ + service: "database", + config: { + version: "17", + databasePassword: Redacted.make(`${tag}-password`), + jwtSecret: Redacted.make(`${tag}-jwt-secret-at-least-32-characters-long`), + jwtExpiry: 3600, + }, + endpoints: { sql: { port: "auto" } }, + }) as const; + +const restPort = (status: { readonly endpoints: ReadonlyArray<{ name: string; port: number }> }) => + status.endpoints.find(({ name }) => name === "http")?.port; + +/** Stops the owner, then opens a fresh one that re-plans `requestedCreations` against the saved state. */ +const restartWithReplan = ( + stack: { readonly id: string; readonly stop: Effect.Effect }, + stateRoot: string, + cacheRoot: string, + requestedCreations: ReadonlyArray, +) => + stack.stop.pipe( + Effect.andThen( + open({ id: stack.id, stateRoot, cacheRoot, startOwner: true, requestedCreations }), + ), + ); + +it.live( + "re-plans a stopped stack's fixed endpoint to automatic, keeping data and other saved ports", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-api-replan-auto-" }); + const stateRoot = `${root}/state`; + const cacheRoot = `${root}/cache`; + const stack = yield* create({ projectRoot: root, stateRoot, cacheRoot, runtime: "native" }); + yield* Effect.ensuring( + Effect.gen(function* () { + const database = replanDatabase("replan-auto"); + const rest = { + service: "rest", + config: {}, + endpoints: { http: { port: FIXED_API_PORT } }, + } as const; + const members = yield* stack.composition.supabase([database, rest], { eager: true }); + const databaseId = members.find(({ service }) => service === "database")?.id; + const restId = members.find(({ service }) => service === "rest")?.id; + if (databaseId === undefined || restId === undefined) + return yield* Effect.die("Composition is missing a member"); + const sqlPortBefore = (yield* (yield* stack.services.get(databaseId)) + .status).endpoints.find(({ name }) => name === "sql")?.port; + + const requestedAuto = { ...rest, endpoints: { http: { port: "auto" } } } as const; + const reopened = yield* restartWithReplan(stack, stateRoot, cacheRoot, [ + database, + requestedAuto, + ]); + const changes = yield* reopened.startupEndpointChanges; + expect(changes).toHaveLength(1); + const change = changes[0]; + expect(change?.service).toBe("rest"); + expect(change?.endpoint).toBe("http"); + expect(change?.from).toBe(FIXED_API_PORT); + // Automatic selection can legitimately land back on the old number; what matters is + // that it claimed some port and the live bind agrees with that claim, checked below. + expect(change?.to).toEqual(expect.any(Number)); + + expect(yield* reopened.composition.plan([database, requestedAuto])).toEqual([ + { id: databaseId, service: "database", member: true, change: "unchanged" }, + { id: restId, service: "rest", member: true, change: "unchanged" }, + ]); + + yield* reopened.composition.supabase([database, requestedAuto], { + reuseIds: [databaseId, restId], + eager: true, + }); + const restStatus = yield* (yield* reopened.services.get(restId)).status; + expect(restPort(restStatus)).toBe(change?.to); + const sqlPortAfter = (yield* (yield* reopened.services.get(databaseId)) + .status).endpoints.find(({ name }) => name === "sql")?.port; + expect(sqlPortAfter).toBe(sqlPortBefore); + }), + destroyTestStack(stack), + ); + }).pipe(Effect.scoped, Effect.provide(layer)), +); + +it.live( + "re-plans a stopped stack's automatic endpoint to a fixed port, and a fixed port to a different one", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-api-replan-fixed-" }); + const stateRoot = `${root}/state`; + const cacheRoot = `${root}/cache`; + const stack = yield* create({ projectRoot: root, stateRoot, cacheRoot, runtime: "native" }); + yield* Effect.ensuring( + Effect.gen(function* () { + const database = replanDatabase("replan-fixed"); + const rest = { + service: "rest", + config: {}, + endpoints: { http: { port: "auto" } }, + } as const; + const members = yield* stack.composition.supabase([database, rest], { eager: true }); + const databaseId = members.find(({ service }) => service === "database")?.id; + const restId = members.find(({ service }) => service === "rest")?.id; + if (databaseId === undefined || restId === undefined) + return yield* Effect.die("Composition is missing a member"); + + const requestedFixed = { + ...rest, + endpoints: { http: { port: FIXED_API_PORT } }, + } as const; + const toFixed = yield* restartWithReplan(stack, stateRoot, cacheRoot, [ + database, + requestedFixed, + ]); + expect(yield* toFixed.startupEndpointChanges).toEqual([ + { service: "rest", endpoint: "http", from: expect.any(Number), to: FIXED_API_PORT }, + ]); + expect(yield* toFixed.composition.plan([database, requestedFixed])).toEqual([ + { id: databaseId, service: "database", member: true, change: "unchanged" }, + { id: restId, service: "rest", member: true, change: "unchanged" }, + ]); + + const requestedOtherFixed = { + ...rest, + endpoints: { http: { port: FIXED_API_PORT + 1 } }, + } as const; + const toOtherFixed = yield* restartWithReplan(toFixed, stateRoot, cacheRoot, [ + database, + requestedOtherFixed, + ]); + expect(yield* toOtherFixed.startupEndpointChanges).toEqual([ + { service: "rest", endpoint: "http", from: FIXED_API_PORT, to: FIXED_API_PORT + 1 }, + ]); + expect(yield* toOtherFixed.composition.plan([database, requestedOtherFixed])).toEqual([ + { id: databaseId, service: "database", member: true, change: "unchanged" }, + { id: restId, service: "rest", member: true, change: "unchanged" }, + ]); + }), + destroyTestStack(stack), + ); + }).pipe(Effect.scoped, Effect.provide(layer)), +); + +/** Binds an ephemeral port picked by the OS, so the conflict never collides with another test. */ +const holdPort = Effect.acquireRelease( + Effect.callback((resume) => { + const server = Net.createServer(); + server.once("error", (cause) => resume(Effect.fail(cause))); + server.listen(0, "127.0.0.1", () => resume(Effect.succeed(server))); + }), + (server) => + Effect.callback((resume) => { + if (!server.listening) return resume(Effect.void); + server.close(() => resume(Effect.void)); + }), +); + +it.live( + "fails to claim a conflicting requested port on a stopped stack, leaving its saved state unchanged", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-api-replan-conflict-" }); + const stateRoot = `${root}/state`; + const cacheRoot = `${root}/cache`; + const stack = yield* create({ projectRoot: root, stateRoot, cacheRoot, runtime: "native" }); + yield* Effect.ensuring( + Effect.gen(function* () { + const database = replanDatabase("replan-conflict"); + const rest = { + service: "rest", + config: {}, + endpoints: { http: { port: FIXED_API_PORT } }, + } as const; + const members = yield* stack.composition.supabase([database, rest], { eager: true }); + const databaseId = members.find(({ service }) => service === "database")?.id; + const restId = members.find(({ service }) => service === "rest")?.id; + if (databaseId === undefined || restId === undefined) + return yield* Effect.die("Composition is missing a member"); + + yield* stack.stop; + + const holder = yield* holdPort; + const address = holder.address(); + const conflictingPort = + typeof address === "object" && address !== null ? address.port : 0; + expect(conflictingPort).toBeGreaterThan(0); + + const requestedConflicting = { + ...rest, + endpoints: { http: { port: conflictingPort } }, + } as const; + const failure = yield* open({ + id: stack.id, + stateRoot, + cacheRoot, + startOwner: true, + requestedCreations: [database, requestedConflicting], + }).pipe(Effect.flip); + // Linux rejects the overlapping bind itself; other platforms fail the occupancy pre-check. + expect(failure.message).toMatch(new RegExp(`\\b${conflictingPort}\\b.*\\bin use\\b`)); + expect(holder.listening).toBe(true); + + const reopened = yield* open({ id: stack.id, stateRoot, cacheRoot }); + expect(yield* reopened.composition.plan([database, requestedConflicting])).toEqual([ + { id: databaseId, service: "database", member: true, change: "unchanged" }, + { + id: restId, + service: "rest", + member: true, + change: "incompatible", + paths: ["endpoints.http.port"], + }, + ]); + expect(yield* reopened.composition.plan([database, rest])).toEqual([ + { id: databaseId, service: "database", member: true, change: "unchanged" }, + { id: restId, service: "rest", member: true, change: "unchanged" }, + ]); + + // The restored claim is not just reported as unchanged: a start with the old config + // actually succeeds, re-plans nothing, and binds REST at its original fixed port. + const restarted = yield* restartWithReplan(reopened, stateRoot, cacheRoot, [ + database, + rest, + ]); + expect(yield* restarted.startupEndpointChanges).toEqual([]); + yield* restarted.composition.supabase([database, rest], { + reuseIds: [databaseId, restId], + eager: true, + }); + const restStatus = yield* (yield* restarted.services.get(restId)).status; + expect(restPort(restStatus)).toBe(FIXED_API_PORT); + }), + destroyTestStack(stack), + ); + }).pipe(Effect.scoped, Effect.provide(layer)), +); + +it.live( + "fails to claim two endpoints re-planned onto the same free port, leaving saved state unchanged", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-api-replan-collide-" }); + const stateRoot = `${root}/state`; + const cacheRoot = `${root}/cache`; + const stack = yield* create({ projectRoot: root, stateRoot, cacheRoot, runtime: "native" }); + yield* Effect.ensuring( + Effect.gen(function* () { + const database = replanDatabase("replan-collide"); + const rest = { + service: "rest", + config: {}, + endpoints: { http: { port: FIXED_API_PORT } }, + } as const; + const members = yield* stack.composition.supabase([database, rest], { eager: true }); + const databaseId = members.find(({ service }) => service === "database")?.id; + const restId = members.find(({ service }) => service === "rest")?.id; + if (databaseId === undefined || restId === undefined) + return yield* Effect.die("Composition is missing a member"); + + yield* stack.stop; + + // Database's sql endpoint and REST's http endpoint use distinct claim keys, so + // requesting the same literal port for both is a same-stack collision, not a + // same-key no-op. + const collidingPort = FIXED_API_PORT + 11; + const requestedDatabase = { + ...database, + endpoints: { sql: { port: collidingPort } }, + } as const; + const requestedRest = { ...rest, endpoints: { http: { port: collidingPort } } } as const; + const failure = yield* open({ + id: stack.id, + stateRoot, + cacheRoot, + startOwner: true, + requestedCreations: [requestedDatabase, requestedRest], + }).pipe(Effect.flip); + expect(failure.message).toContain("claimed by another listener of this stack"); + + const reopened = yield* open({ id: stack.id, stateRoot, cacheRoot }); + expect(yield* reopened.composition.plan([database, rest])).toEqual([ + { id: databaseId, service: "database", member: true, change: "unchanged" }, + { id: restId, service: "rest", member: true, change: "unchanged" }, + ]); + + const restarted = yield* restartWithReplan(reopened, stateRoot, cacheRoot, [ + database, + rest, + ]); + expect(yield* restarted.startupEndpointChanges).toEqual([]); + yield* restarted.composition.supabase([database, rest], { + reuseIds: [databaseId, restId], + eager: true, + }); + const restStatus = yield* (yield* restarted.services.get(restId)).status; + expect(restPort(restStatus)).toBe(FIXED_API_PORT); + }), + destroyTestStack(stack), + ); + }).pipe(Effect.scoped, Effect.provide(layer)), +); + +it.live( + "leaves a changed Postgres major version incompatible without re-planning any endpoint", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-api-replan-version-" }); + const stateRoot = `${root}/state`; + const cacheRoot = `${root}/cache`; + const stack = yield* create({ projectRoot: root, stateRoot, cacheRoot, runtime: "native" }); + yield* Effect.ensuring( + Effect.gen(function* () { + const database = replanDatabase("replan-version"); + const rest = { + service: "rest", + config: {}, + endpoints: { http: { port: FIXED_API_PORT } }, + } as const; + const members = yield* stack.composition.supabase([database, rest], { eager: true }); + const databaseId = members.find(({ service }) => service === "database")?.id; + const restId = members.find(({ service }) => service === "rest")?.id; + if (databaseId === undefined || restId === undefined) + return yield* Effect.die("Composition is missing a member"); + + const olderVersion = { ...database, config: { ...database.config, version: "15" } }; + const requestedAuto = { ...rest, endpoints: { http: { port: "auto" } } } as const; + const reopened = yield* restartWithReplan(stack, stateRoot, cacheRoot, [ + olderVersion, + requestedAuto, + ]); + expect(yield* reopened.startupEndpointChanges).toEqual([]); + expect(yield* reopened.composition.plan([olderVersion, requestedAuto])).toEqual([ + { + id: databaseId, + service: "database", + member: true, + change: "incompatible", + paths: ["config.version"], + }, + { + id: restId, + service: "rest", + member: true, + change: "incompatible", + paths: ["endpoints.http.port"], + }, + ]); + }), + destroyTestStack(stack), + ); + }).pipe(Effect.scoped, Effect.provide(layer)), +); + +it.live( + "does not let an excluded sibling's stale fixed port mask a shared endpoint moving to automatic", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-api-replan-excluded-" }); + const stateRoot = `${root}/state`; + const cacheRoot = `${root}/cache`; + const stack = yield* create({ projectRoot: root, stateRoot, cacheRoot, runtime: "native" }); + yield* Effect.ensuring( + Effect.gen(function* () { + const database = replanDatabase("replan-excluded"); + const rest = { + service: "rest", + config: {}, + endpoints: { http: { port: FIXED_API_PORT } }, + } as const; + const auth = { + service: "auth", + config: {}, + endpoints: { http: { port: FIXED_API_PORT } }, + } as const; + const members = yield* stack.composition.supabase([database, rest, auth], { + eager: true, + }); + const restId = members.find(({ service }) => service === "rest")?.id; + if (restId === undefined) return yield* Effect.die("Composition is missing REST"); + + // The config now omits [api] port and excludes auth, so only REST is requested. + const requestedAuto = { ...rest, endpoints: { http: { port: "auto" } } } as const; + const reopened = yield* restartWithReplan(stack, stateRoot, cacheRoot, [ + database, + requestedAuto, + ]); + const changes = yield* reopened.startupEndpointChanges; + // Automatic selection can legitimately land back on the old number; what matters is + // that it claimed some port and the live bind agrees with that claim, checked below. + expect(changes).toEqual([ + { service: "rest", endpoint: "http", from: FIXED_API_PORT, to: expect.any(Number) }, + ]); + + const databaseId = (yield* reopened.services.list).find( + ({ service }) => service === "database", + )?.id; + if (databaseId === undefined) return yield* Effect.die("Composition is missing database"); + yield* reopened.composition.supabase([database, requestedAuto], { + reuseIds: [databaseId, restId], + eager: true, + }); + const restStatus = yield* (yield* reopened.services.get(restId)).status; + expect(restPort(restStatus)).toBe(changes[0]?.to); + }), + destroyTestStack(stack), + ); + }).pipe(Effect.scoped, Effect.provide(layer)), +); + +it.live( + "lets only one of several concurrent starts of a stopped stack re-plan its changed endpoint", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "stack-api-replan-concurrent-" }); + const stateRoot = `${root}/state`; + const cacheRoot = `${root}/cache`; + const stack = yield* create({ projectRoot: root, stateRoot, cacheRoot, runtime: "native" }); + yield* Effect.ensuring( + Effect.gen(function* () { + const database = replanDatabase("replan-concurrent"); + const rest = { + service: "rest", + config: {}, + endpoints: { http: { port: FIXED_API_PORT } }, + } as const; + yield* stack.composition.supabase([database, rest], { eager: true }); + yield* stack.stop; + + const requestedAuto = { ...rest, endpoints: { http: { port: "auto" } } } as const; + // Several concurrent starts race on the same lease; only its winner re-plans, and the + // others attach to it instead, so every caller observes the identical applied change. + const opened = yield* Effect.all( + [1, 2, 3].map(() => + open({ + id: stack.id, + stateRoot, + cacheRoot, + startOwner: true, + requestedCreations: [database, requestedAuto], + }), + ), + { concurrency: "unbounded" }, + ); + const changeSets = yield* Effect.forEach( + opened, + (handle) => handle.startupEndpointChanges, + ); + for (const changes of changeSets) expect(changes).toEqual(changeSets[0]); + expect(changeSets[0]).toEqual([ + { service: "rest", endpoint: "http", from: FIXED_API_PORT, to: expect.any(Number) }, + ]); + + const entries = yield* discover({ stateRoot }); + const entry = entries.find(({ definition }) => definition.id === stack.id); + expect(entry?.definition.ports.filter(({ key }) => key === "api")).toHaveLength(1); + expect(new Set(opened.map((handle) => handle.id)).size).toBe(1); + }), + destroyTestStack(stack), + ); + }).pipe(Effect.scoped, Effect.provide(layer)), +); diff --git a/packages/stack/src/effect.ts b/packages/stack/src/effect.ts index 40e09273f6..67649f265a 100644 --- a/packages/stack/src/effect.ts +++ b/packages/stack/src/effect.ts @@ -36,6 +36,7 @@ import { import { planSupabaseComposition, type PlannedInstance, + type PlanOptions, type SupabaseCompositionOptions, } from "./composition/Supabase.ts"; import { removeStackContainersCommand } from "./runtime/Container.ts"; @@ -44,7 +45,7 @@ import { deriveStackId, resolveStackIdentity } from "./identity/Identity.ts"; import { failureMessage } from "./internal/failure-message.ts"; import * as State from "./State.ts"; import type { SavedStack, StackCredentials, StackKeysInput } from "./State.ts"; -import { StackError, type Definition, type Observation } from "./Rpc.ts"; +import { StackError, type Definition, type EndpointPortChange, type Observation } from "./Rpc.ts"; import { reclaimStack } from "./Sweep.ts"; import { ServiceCreationInput as ServiceCreationInputSchema, @@ -71,12 +72,13 @@ export type { CompositionConfig } from "./Orchestrator.ts"; export type { CreationChange, PlannedInstance, + PlanOptions, SupabaseCompositionOptions, } from "./composition/Supabase.ts"; export { StackIdSchema as StackId } from "./identity/StackId.ts"; export type { SavedStack } from "./State.ts"; export type { StackCredentials, StackKeysInput }; -export type { Observation } from "./Rpc.ts"; +export type { EndpointPortChange, Observation } from "./Rpc.ts"; export type { Command, InitializationCommand, @@ -111,6 +113,12 @@ export interface CreateOptions extends StackLocations { export interface OpenOptions extends StackLocations { readonly id: string; readonly startOwner?: boolean; + /** + * Endpoint intents a spawned owner re-plans against the saved state before it registers + * endpoint namespaces from it, applying a changed endpoint's port instead of keeping the saved + * one. Only a freshly spawned owner acts on this; an already-running owner is unaffected. + */ + readonly requestedCreations?: ReadonlyArray; } /** No owner serves the stack: nothing holds its lease, or a sweeper is cleaning it up. */ @@ -252,10 +260,15 @@ export interface Stack { ) => Effect.Effect, StackError>; /** * Compares the requested creations with every saved instance of the same kinds, ignoring - * inputs the composition supplies, without changing state or contacting the owner. + * inputs the composition supplies, without changing state or contacting the owner. A + * `complete` request (the default caller's whole desired composition) never falls back to a + * saved sibling's port for a shared endpoint kind the request leaves out; a `partial` request + * (comparing only the kinds it names) does, so an incremental check keeps matching an + * unrelated sibling's still-current port. */ readonly plan: ( services: ReadonlyArray, + options?: PlanOptions, ) => Effect.Effect, StackError>; readonly configure: (config: Orchestrator.CompositionConfig) => Effect.Effect; readonly describe: Effect.Effect; @@ -263,6 +276,14 @@ export interface Stack { readonly stop: Effect.Effect, StackError>; readonly restart: Effect.Effect, StackError>; }; + /** + * The endpoint port changes the owner this handle reaches applied at its own last startup, + * re-planned from `OpenOptions.requestedCreations`. An owner that was already running when this + * handle reached it reports whatever its own last startup applied (possibly none), the same as + * every other caller reaching that owner; it does not re-plan again for this call. Empty when + * that startup re-planned nothing at all. + */ + readonly startupEndpointChanges: Effect.Effect, StackError>; readonly stop: Effect.Effect; readonly destroy: Effect.Effect; readonly commands: { @@ -874,14 +895,14 @@ const makeHandle = Effect.fn("Stack.makeHandle")(function* ( ...(options?.eager === undefined ? {} : { eager: options.eager }), }), ).pipe(Effect.map((definitions) => definitions.map(instance))), - plan: (services: ReadonlyArray) => + plan: (services: ReadonlyArray, options?: PlanOptions) => Effect.forEach(services, (service) => Schema.decodeEffect(ServiceCreationInputSchema)(service), ).pipe( Effect.mapError((cause) => failure("plan", cause)), Effect.flatMap((requested) => savedDefinition.pipe( - Effect.map((current) => planSupabaseComposition(current, requested)), + Effect.map((current) => planSupabaseComposition(current, requested, options)), ), ), ), @@ -892,6 +913,11 @@ const makeHandle = Effect.fn("Stack.makeHandle")(function* ( stop: whileRunning("stopComposition", (rpc) => rpc.stopComposition(), []), restart: call("restartComposition", (rpc) => rpc.restartComposition()), }, + startupEndpointChanges: call( + "startupEndpointChanges", + (rpc) => rpc.startupEndpointChanges(), + "attach", + ), stop: shutdown(false).pipe(Effect.asVoid), destroy: shutdown(true), commands: { run }, @@ -963,7 +989,13 @@ export const open = Effect.fn("Stack.open")( ? undefined : saved.lifetime === "session" ? yield* connectHost(state, saved.id) - : yield* launchHost(state, { ...locations, stackId: saved.id }); + : yield* launchHost(state, { + ...locations, + stackId: saved.id, + ...(options.requestedCreations === undefined + ? {} + : { requestedCreations: options.requestedCreations }), + }); return yield* makeHandle(state, saved, locations, { access }); }, Effect.mapError((cause) => failure("open", cause)), diff --git a/packages/stack/src/host/Endpoints.ts b/packages/stack/src/host/Endpoints.ts index 036f245874..28558e6fa5 100644 --- a/packages/stack/src/host/Endpoints.ts +++ b/packages/stack/src/host/Endpoints.ts @@ -30,8 +30,11 @@ export const endpointPort = (creation: EndpointIntents, name: string): number | return endpoint.port === "auto" || typeof endpoint.port === "number" ? endpoint.port : "auto"; }; -/** Path prefix of a service on the shared API listener, or `undefined` when it has a dedicated one. */ -export const apiRoute = (service: ServiceCreation["service"]): string | undefined => { +/** + * Path prefix of a service on the shared API listener, or `undefined` when it has a dedicated one. + * Accepts a plain string too, for a service kind decoded off an RPC transport. + */ +export const apiRoute = (service: string): string | undefined => { switch (service) { case "rest": return "/rest/v1"; diff --git a/packages/stack/src/internal/host-process.ts b/packages/stack/src/internal/host-process.ts index eb459d0c37..f39bd7708d 100644 --- a/packages/stack/src/internal/host-process.ts +++ b/packages/stack/src/internal/host-process.ts @@ -1,7 +1,7 @@ import { Cause, Effect, Exit, Option, Schema } from "effect"; -// oxlint-disable-next-line effecttsgo/node-builtin-import -- readiness is an inherited launcher descriptor. -import { closeSync, writeSync } from "node:fs"; -import { SavedStack } from "../State.ts"; +// oxlint-disable-next-line effecttsgo/node-builtin-import -- readiness is an inherited launcher descriptor, and the startup payload file predates any service layer. +import { closeSync, readFileSync, unlinkSync, writeSync } from "node:fs"; +import { HostStartupPayload } from "../HostProcess.ts"; import { runStackHost, StackHostError, type StackHostOptions } from "../StackHost.ts"; const writeLine = (value: unknown) => @@ -34,7 +34,7 @@ const options = ( report: (value: unknown) => Effect.Effect, ): Effect.Effect => Effect.gen(function* () { - const [stateRoot, cacheRoot, stackId, register, ...rest] = args; + const [stateRoot, cacheRoot, stackId, payloadFile, ...rest] = args; if ( stateRoot === undefined || cacheRoot === undefined || @@ -43,21 +43,44 @@ const options = ( ) return yield* new StackHostError({ operation: "startup", - message: "Expected stateRoot, cacheRoot, stackId and an optional stack to register", + message: "Expected stateRoot, cacheRoot, stackId and an optional startup payload file", }); - const registered = - register === undefined + // The launcher writes this file once, under the owner's own state directory with owner-only + // permissions; reading and deleting it here, before anything else, keeps its secrets off argv + // and off this process's whole lifetime in a live process list. + const payload: HostStartupPayload | undefined = + payloadFile === undefined || payloadFile === "" ? undefined - : yield* Schema.decodeEffect(Schema.fromJsonString(SavedStack))(register).pipe( - Effect.mapError( - (cause) => new StackHostError({ operation: "startup", message: cause.message }), + : yield* Effect.gen(function* () { + const text = yield* Effect.try({ + try: () => readFileSync(payloadFile, "utf8"), + catch: (cause) => + new StackHostError({ operation: "startup", message: String(cause) }), + }); + return yield* Schema.decodeEffect(Schema.fromJsonString(HostStartupPayload))(text).pipe( + Effect.mapError( + (cause) => new StackHostError({ operation: "startup", message: cause.message }), + ), + ); + }).pipe( + Effect.ensuring( + Effect.sync(() => { + try { + unlinkSync(payloadFile); + } catch { + // Already removed, or the parent is cleaning up; the secret is gone either way. + } + }), ), ); return { stateRoot, cacheRoot, stackId, - ...(registered === undefined ? {} : { register: registered }), + ...(payload?.register === undefined ? {} : { register: payload.register }), + ...(payload?.requestedCreations === undefined + ? {} + : { requestedCreations: payload.requestedCreations }), onReady: ({ endpoint, secret }) => report({ type: "ready", endpoint, secret }), }; }); diff --git a/packages/stack/tests/host-process-fixture.ts b/packages/stack/tests/host-process-fixture.ts index f0210334c5..51d40130b4 100644 --- a/packages/stack/tests/host-process-fixture.ts +++ b/packages/stack/tests/host-process-fixture.ts @@ -14,6 +14,8 @@ import { import { HttpServerRequest, HttpServerResponse } from "effect/unstable/http"; // oxlint-disable-next-line effecttsgo/node-builtin-import -- test fixture owns inherited readiness and release descriptors. import { closeSync, createReadStream, writeSync } from "node:fs"; +// oxlint-disable-next-line effecttsgo/node-builtin-import -- consumes its own startup payload file the way the real consumer does. +import { existsSync, readFileSync, unlinkSync } from "node:fs"; import { currentRelease, launchHost, authorizes, type HostEndpoint } from "../src/HostProcess.ts"; import { bindControl } from "../src/StackHost.ts"; import * as State from "../src/State.ts"; @@ -44,9 +46,26 @@ const report = (value: unknown) => { } }; -const [stateRoot, cacheRoot, stackId, mode = "owner", ownerEntrypoint] = process.argv.slice(2); +const [stateRoot, cacheRoot, stackId, fourth, fifth] = process.argv.slice(2); if (stateRoot === undefined || stackId === undefined) throw new FixtureError({ message: "fixture arguments missing" }); +/** + * A spawned owner's startup payload travels as a file path in this same position instead of this + * fixture's own "mode"/entrypoint test convention; a real file path on disk, checked here, + * disambiguates it. Consuming it the way the real consumer does leaves nothing behind for a test + * that launches this fixture with `register` or `requestedCreations`. + */ +const payloadFile = fourth !== undefined && existsSync(fourth) ? fourth : undefined; +if (payloadFile !== undefined) { + readFileSync(payloadFile, "utf8"); + try { + unlinkSync(payloadFile); + } catch { + /* Already removed, or the parent is cleaning up; the secret is gone either way. */ + } +} +const mode = payloadFile === undefined ? (fourth ?? "owner") : "owner"; +const ownerEntrypoint = payloadFile === undefined ? fifth : undefined; /** Follows the owner protocol without services; `held` acknowledges shutdown but stays alive. */ const owner = Effect.scoped(