Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 14 additions & 10 deletions packages/stack/src/Owner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -288,16 +288,20 @@ const makeOwner = Effect.fn("Owner.make")(function* (options: OwnerOptions) {
});

const recipeFor = (creation: ServiceCreation, id: string) =>
makeServiceRecipe(creation, {
stackId,
instanceId: id,
project,
root: options.root,
cacheRoot: options.cacheRoot,
runtime,
helpers,
...(options.hostGateway === undefined ? {} : { hostGateway: options.hostGateway }),
}).pipe(Effect.provideContext(services));
makeServiceRecipe(
creation,
{
stackId,
instanceId: id,
project,
root: options.root,
cacheRoot: options.cacheRoot,
runtime,
helpers,
...(options.hostGateway === undefined ? {} : { hostGateway: options.hostGateway }),
},
options.state.claims,
).pipe(Effect.provideContext(services));

const persistCreation = (entry: Pick<Entry, "id" | "creation">, creation: ServiceCreation) =>
updateState((current) => ({
Expand Down
90 changes: 89 additions & 1 deletion packages/stack/src/Ports.integration.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { NodeServices, NodeSocketServer } from "@effect/platform-node";
import { expect, it } from "@effect/vitest";
import { Cause, Context, Effect, Exit, FileSystem, Layer, Option, Ref, Scope } from "effect";
import { makePorts, PortError } from "./Ports.ts";
import { makePorts, portBase, portSpan, PortError, reserveNativePort } from "./Ports.ts";
import * as State from "./State.ts";

const makeTestState = (root: string) =>
Expand Down Expand Up @@ -421,3 +421,91 @@ it.live("lets exactly one of two stacks sharing a saved port bind it when both s
}),
).pipe(Effect.provide(NodeServices.layer)),
);

// Fixed so a test can force two reservations to the same candidate; production uses randomPortSpanStart.
const fixedStart = Effect.succeed(0);

// Binds a real listener directly in the native-reservation span, retrying past occupied
// candidates, instead of reserving then releasing a port that something else could grab meanwhile.
const bindBlockingPort = (host: string) =>
Effect.gen(function* () {
for (let offset = 0; offset < portSpan; offset++) {
const attempt = yield* Effect.exit(bind(host, portBase + offset));
if (Exit.isSuccess(attempt)) return { port: portBase + offset, listener: attempt.value };
}
return yield* Effect.die("No port in the native reservation span was free for the fixture");
});

it.live("skips a native backend port claimed by another saved stack", () =>
Effect.scoped(
Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const root = yield* fs.makeTempDirectoryScoped();
const state = yield* makeTestState(root);

const probeScope = yield* Scope.make();
const probe = yield* reserveNativePort([], "pooler", fixedStart).pipe(
Effect.provideService(Scope.Scope, probeScope),
);
yield* Scope.close(probeScope, Exit.void);

yield* saveStack(state, root, "claimer", [
{ key: "db/sql", host: "127.0.0.1", port: probe.port },
]);

const reserved = yield* reserveNativePort(yield* state.claims, "pooler", fixedStart);
expect(reserved.port).not.toBe(probe.port);
}),
).pipe(Effect.provide(NodeServices.layer)),
);

it.live("skips a native backend port a wildcard listener holds", () =>
Effect.scoped(
Effect.gen(function* () {
// A wildcard bind is reachable through loopback, so a loopback-only probe would miss it.
const blocked = yield* bindBlockingPort("0.0.0.0");
const forcedStart = Effect.succeed(blocked.port - portBase);

const reserved = yield* reserveNativePort([], "pooler", forcedStart);
expect(reserved.port).not.toBe(blocked.port);
}),
).pipe(Effect.provide(NodeServices.layer)),
);

it.live("excludes a port a previous attempt lost from the next reservation", () =>
Effect.scoped(
Effect.gen(function* () {
const firstScope = yield* Scope.make();
const first = yield* reserveNativePort([], "pooler", fixedStart).pipe(
Effect.provideService(Scope.Scope, firstScope),
);
yield* Scope.close(firstScope, Exit.void);

const second = yield* reserveNativePort([], "pooler", fixedStart, new Set([first.port]));
expect(second.port).not.toBe(first.port);
}),
).pipe(Effect.provide(NodeServices.layer)),
);

it.live("skips a public auto candidate a loopback listener already holds", () =>
Effect.scoped(
Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const root = yield* fs.makeTempDirectoryScoped();
const state = yield* makeTestState(root);
yield* saveStack(state, root, "stack");
const ports = yield* makePorts(state);
// A container-runtime stack's bind to the wildcard host would otherwise succeed here too.
const request = { stackId: "stack", key: "api", host: "0.0.0.0", port: "auto" as const };
const accept = (_host: string, port: number) => Effect.succeed(port);

// A loopback-only listener keeps occupying the port while only its saved claim is released,
// so the next auto allocation has to skip it for real via loopbackOccupied, not a real bind.
const probe = yield* ports.acquire(request, (_host, port) => bind("127.0.0.1", port));
yield* ports.release("stack", "api");

const acquired = yield* ports.acquire(request, accept);
expect(acquired.port).not.toBe(probe.port);
}),
).pipe(Effect.provide(NodeServices.layer)),
);
108 changes: 105 additions & 3 deletions packages/stack/src/Ports.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
import { Cause, Data, Effect, Exit, Hash, Option, Scope } from "effect";
import { Cause, Crypto, Data, Effect, Exit, Hash, Option, Scope } from "effect";
import * as Net from "node:net";
import type * as State from "./State.ts";

const portBase = 20000;
export const portBase = 20000;
/** Stays below the Linux ephemeral range, per the [architecture ADR](../../../docs/adr/0017-simplified-managed-stack-architecture.md). */
const portSpan = 12768;
export const portSpan = 12768;
/** Co-prime with the span, so the scan visits every port once and steps past reserved ranges. */
const portStride = 257;

Expand Down Expand Up @@ -88,6 +88,100 @@ const resolveRequest = (
return Effect.succeed({ saved, requested });
};

export interface NativePortReservation {
readonly port: number;
readonly server: Net.Server;
}

const closeNativePort = (server: Net.Server): Effect.Effect<void> =>
Effect.callback<void, never>((resume) => {
if (!server.listening) {
resume(Effect.void);
return Effect.void;
}
server.close(() => resume(Effect.void));
return Effect.void;
});

const bindNativePort = (
key: string,
port: number,
): Effect.Effect<NativePortReservation, PortError> =>
Effect.callback<NativePortReservation, PortError>((resume) => {
const server = Net.createServer((socket) => socket.destroy());
const onError = (cause: Error) =>
resume(
Effect.fail(
new PortError({ key, message: "Unable to reserve native service port", cause }),
),
);
server.once("error", onError);
server.listen({ host: "127.0.0.1", port }, () => {
const address = server.address();
if (address === null || typeof address === "string") {
onError(new Error("Native service port reservation returned no address"));
} else {
resume(Effect.succeed({ port: address.port, server }));
}
});
return Effect.sync(() => {
server.off("error", onError);
if (server.listening) server.close();
});
});

const emptyPortSet: ReadonlySet<number> = new Set();

/** Random, so a reopened stack doesn't retry a backend port still in server-side `TIME_WAIT`. */
export const randomPortSpanStart = (crypto: Crypto.Crypto): Effect.Effect<number> =>
crypto.randomIntBetween(0, portSpan, { halfOpen: true });

/**
* Scans the below-ephemeral span for a backend port, skipping claimed and excluded ports; reuses
* `loopbackOccupied` so a wildcard listener a loopback-only bind would miss on macOS, BSD or
* Windows still rules out the candidate. Never persists one of its own.
*/
export const reserveNativePort = Effect.fn("Ports.reserveNativePort")(
(
claims: ReadonlyArray<State.StackClaims>,
key: string,
randomStart: Effect.Effect<number>,
excluded: ReadonlySet<number> = emptyPortSet,
): Effect.Effect<NativePortReservation, PortError, Scope.Scope> =>
Effect.acquireRelease(
Effect.gen(function* () {
const claimed = new Set(claims.flatMap((stack) => stack.ports.map((claim) => claim.port)));
const start = yield* randomStart;
let failures = 0;
let lastFailure: PortError | undefined;
for (let attempt = 0; attempt < portSpan && failures < 64; attempt++) {
const port = portBase + ((start + attempt * portStride) % portSpan);
if (claimed.has(port) || excluded.has(port)) continue;
if (yield* loopbackOccupied(port)) {
failures++;
lastFailure = new PortError({ key, message: `Port ${port} is already in use` });
continue;
}
Comment thread
jgoux marked this conversation as resolved.
const result = yield* Effect.exit(bindNativePort(key, port));
if (Exit.isSuccess(result)) return result.value;
const error = Cause.findErrorOption(result.cause);
if (Option.isNone(error)) return yield* Effect.failCause(result.cause);
failures++;
lastFailure = error.value;
}
return yield* new PortError({
key,
message:
lastFailure === undefined
? "No native service port is available"
: `No native service port is available: ${lastFailure.message}`,
cause: lastFailure,
});
}),
({ server }) => closeNativePort(server),
),
);

/** Claims steer auto allocation away from saved stacks; live listeners and binds decide conflicts for fixed ports. */
export const makePorts = (state: State.Interface, platform: NodeJS.Platform = process.platform) =>
Effect.sync(() => {
Expand Down Expand Up @@ -174,6 +268,14 @@ export const makePorts = (state: State.Interface, platform: NodeJS.Platform = pr
? portBase + ((start + attempt * portStride) % portSpan)
: requested;
if (requested === "auto" && claimed.has(port)) continue;
if (requested === "auto" && (yield* loopbackOccupied(port))) {
failures++;
lastFailure = new PortError({
key: request.key,
message: `Port ${port} is already in use`,
});
continue;
}
Comment thread
jgoux marked this conversation as resolved.
const result = yield* Effect.uninterruptibleMask((restore) =>
Effect.gen(function* () {
const scope = yield* Scope.fork(owner, "sequential");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ describe("service catalog", () => {
},
},
dockerOptions(root),
Effect.succeed([]),
);
const database = yield* makeService(databaseRecipe.definition, {
id: "database",
Expand All @@ -57,6 +58,7 @@ describe("service catalog", () => {
config: { databaseUrl, backend: "postgres", apiKey: "catalog-analytics" },
},
dockerOptions(root),
Effect.succeed([]),
);
const analytics = yield* makeService(analyticsRecipe.definition, {
id: "analytics",
Expand All @@ -83,6 +85,7 @@ describe("service catalog", () => {
},
},
dockerOptions(root),
Effect.succeed([]),
);
yield* fs.writeFileString(
`${root}/vector.yaml`,
Expand All @@ -106,6 +109,7 @@ describe("service catalog", () => {
const imgproxyRecipe = yield* makeServiceRecipe(
{ service: "imgproxy", config: { filePath: imageRoot } },
dockerOptions(root),
Effect.succeed([]),
);
const imgproxy = yield* makeService(imgproxyRecipe.definition, {
id: "imgproxy",
Expand Down
4 changes: 4 additions & 0 deletions packages/stack/src/services/AuthStorage.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ describe("service catalog", () => {
},
},
dockerOptions(root),
Effect.succeed([]),
);
const database = yield* makeService(databaseRecipe.definition, {
id: "database",
Expand All @@ -57,6 +58,7 @@ describe("service catalog", () => {
config: { databaseUrl, jwtSecret: secret, jwtExpiry: 3600 },
},
dockerOptions(root),
Effect.succeed([]),
);
const auth = yield* makeService(authRecipe.definition, {
id: "auth",
Expand Down Expand Up @@ -98,6 +100,7 @@ describe("service catalog", () => {
const imgproxyRecipe = yield* makeServiceRecipe(
{ service: "imgproxy", config: { filePath: storageRoot } },
dockerOptions(root),
Effect.succeed([]),
);
const imgproxy = yield* makeService(imgproxyRecipe.definition, {
id: "imgproxy",
Expand All @@ -117,6 +120,7 @@ describe("service catalog", () => {
},
},
dockerOptions(root),
Effect.succeed([]),
);
const storage = yield* makeService(storageRecipe.definition, {
id: "storage",
Expand Down
3 changes: 3 additions & 0 deletions packages/stack/src/services/Catalog.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ describe("service catalog", () => {
endpoints: { smtp: { port: "auto" } },
},
options(root),
Effect.succeed([]),
).pipe(Effect.exit);
expect(Exit.isFailure(result)).toBe(true);
const misplacedDatabaseVersion = yield* makeServiceRecipe(
Expand All @@ -38,6 +39,7 @@ describe("service catalog", () => {
},
},
options(root),
Effect.succeed([]),
).pipe(Effect.exit);
expect(Exit.isFailure(misplacedDatabaseVersion)).toBe(true);
}),
Expand All @@ -56,6 +58,7 @@ describe("service catalog", () => {
config: { databaseUrl: "postgres://db" },
},
options(root),
Effect.succeed([]),
),
);
expect(error.operation).toBe("config");
Expand Down
20 changes: 18 additions & 2 deletions packages/stack/src/services/Catalog.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,9 @@ import { Crypto, Effect, FileSystem, Path, Ref, Schema, Stream } from "effect";
import { ChildProcessSpawner } from "effect/unstable/process";
import { HttpClient } from "effect/unstable/http";
import { makeContainerRuntime } from "../runtime/Container.ts";
import { randomPortSpanStart, reserveNativePort } from "../Ports.ts";
import { ServiceError, type ServiceDefinition } from "../Service.ts";
import type * as State from "../State.ts";
import {
makeDatabase,
DatabaseConfig,
Expand Down Expand Up @@ -298,7 +300,11 @@ const databaseRecipe = (
});

export const makeServiceRecipe = Effect.fn("Catalog.makeServiceRecipe")(
(input: unknown, options: CatalogOptions) =>
(
input: unknown,
options: CatalogOptions,
readPortClaims: Effect.Effect<ReadonlyArray<State.StackClaims>, State.StateError>,
) =>
Effect.gen(function* () {
const endpointError = validateEndpointNames(input);
if (endpointError !== undefined) return yield* endpointError;
Expand Down Expand Up @@ -351,7 +357,17 @@ export const makeServiceRecipe = Effect.fn("Catalog.makeServiceRecipe")(
imageMirrors: slimImageMirrors,
...(options.hostGateway === undefined ? {} : { hostGateway: options.hostGateway }),
});
const deps: ProcessDependencies = { fs, path, crypto, client, spawner, container };
const deps: ProcessDependencies = {
fs,
path,
crypto,
client,
spawner,
container,
readPortClaims,
reserveNativePort: (key, claims, excluded) =>
reserveNativePort(claims, key, randomPortSpanStart(crypto), excluded),
};
switch (creation.service) {
case "rest":
return catalogRecipe(
Expand Down
Loading
Loading