From 43912beb43902fd6487edf37c2a694853f22d885 Mon Sep 17 00:00:00 2001 From: Antony Rizzitelli Date: Fri, 25 Sep 2026 16:24:07 +0200 Subject: [PATCH 1/2] fix(port-reservation): stop the port placeholder from hanging startup The port held between `provide` and `start` was a bare `net.createServer()` without a connection handler. A client connecting during that window, such as a Kubernetes startup probe, got a socket that was never read nor closed, so `holder.close()` never called back, the real `listen()` was never reached and the HTTP server never started. The placeholder now answers every accepted connection with a minimal `503 Service Unavailable` carrying `Connection: close` and `Retry-After`, ends it, and tracks it; releasing the reservation destroys any connection still open so the close always completes. Releasing the reservations is additionally bounded by a timeout, and an error is logged when the servers are still not listening five seconds after `start`. --- README.md | 2 + src/index.ts | 32 +++- src/port-reservation.ts | 52 +++++- src/startup-deadline.ts | 37 +++++ src/test/reservation-connections.test.ts | 199 +++++++++++++++++++++++ 5 files changed, 310 insertions(+), 12 deletions(-) create mode 100644 src/startup-deadline.ts create mode 100644 src/test/reservation-connections.test.ts diff --git a/README.md b/README.md index 06108fe..0493c40 100644 --- a/README.md +++ b/README.md @@ -124,6 +124,8 @@ The reservation honours the existing port rules: - In development, the reservation falls back to the next free port (up to 20 above the requested one, then an OS-assigned port), exactly as `listen()` did before. - `port: 0` reserves an OS-assigned port and publishes it. +A client that connects while the port is only reserved — a health probe, an eager consumer — receives an immediate `503 Service Unavailable` with `Retry-After: 1` and a closed connection, so it retries once the real server listens. Releasing the reservation never waits on such a connection. If the servers are still not listening five seconds after `start`, the module logs an error. + ### One advertised origin The module names itself in exactly one way. The host recorded in `.antelope/dev.json` and the host in `API_LOCAL_BASE_URL` go through the same derivation, so a browser is never handed the same server under two spellings — `localhost` and `127.0.0.1` are distinct origins to it, worth a second preflight and a second cookie jar. diff --git a/src/index.ts b/src/index.ts index cebc98a..74d6c41 100644 --- a/src/index.ts +++ b/src/index.ts @@ -9,6 +9,11 @@ import { listenServer } from "./port-binding"; import type { Config } from "./server-config"; import { buildConfigVars } from "./config-vars"; import { createConfiguredServer } from "./server-factory"; +import { + armListenDeadline, + disarmListenDeadline, + settlesWithin, +} from "./startup-deadline"; import { configure, getConfig, setCorsConfig } from "./module-config"; export { configure, getConfig, setCorsConfig }; @@ -24,14 +29,31 @@ import { } from "./dev-registry"; import "./middlewares/cors"; +const RESERVATION_RELEASE_TIMEOUT_MS = 1000; +const LISTEN_DEADLINE_MS = 5000; + let servers: net.Server[] = []; let listening = false; let reservations: ReservedPort[] = []; -function releaseReservations(): Promise { +async function releaseReservations(): Promise { const pending = reservations; reservations = []; - return releaseReservedPorts(pending); + const isReleased = await settlesWithin( + releaseReservedPorts(pending), + RESERVATION_RELEASE_TIMEOUT_MS, + ); + if (!isReleased) { + Logging.Warn( + `Port reservations still closing after ${RESERVATION_RELEASE_TIMEOUT_MS} ms, binding the servers anyway`, + ); + } +} + +function reportMissedListenDeadline(): void { + Logging.Error( + `Servers are still not listening ${LISTEN_DEADLINE_MS} ms after start`, + ); } async function reserveConfiguredPorts(): Promise { @@ -113,6 +135,7 @@ function closeServers(): Promise { ); servers = []; listening = false; + disarmListenDeadline(); return Promise.all(closing).then(() => undefined); } @@ -122,7 +145,8 @@ export function start(): void { createConfiguredServer(serverConfig), ); - if (getConfig().autoListen !== false) { + if (getConfig().autoListen !== false && servers.length > 0) { + armListenDeadline(LISTEN_DEADLINE_MS, reportMissedListenDeadline); void serversClosed .then(() => listenServers()) .catch((error: unknown) => { @@ -154,6 +178,8 @@ export async function listenServers(): Promise { } catch (error) { listening = false; throw error; + } finally { + disarmListenDeadline(); } await registerDevServerEndpoints(getListeningEndpoints()); diff --git a/src/port-reservation.ts b/src/port-reservation.ts index 644c3b1..57fde9b 100644 --- a/src/port-reservation.ts +++ b/src/port-reservation.ts @@ -14,7 +14,9 @@ import { * Holding the socket is what keeps the published `API_PORT` and the port * the server eventually binds in sync: the probe-then-use race shrinks to * the instant between `release()` and the real `listen()`, instead of - * spanning the construction of every other module. + * spanning the construction of every other module. Any client reaching + * the port meanwhile gets an immediate `503` and a closed connection, so + * an early probe never keeps the holding socket from closing. */ export interface ReservedPort { port: number; @@ -30,26 +32,58 @@ export class PortReservationError extends Error { readonly code = "EADDRINUSE"; } +interface PortHolder { + server: net.Server; + sockets: Set; +} + interface PortHold { - holder?: net.Server; + holder?: PortHolder; error?: unknown; } const STRICT_PORT_HINT = " Port fallback is disabled; set strictPort to false in development to accept the next free port."; +const RESERVATION_RETRY_AFTER_SECONDS = 1; + +const RESERVATION_RESPONSE = [ + "HTTP/1.1 503 Service Unavailable", + "Connection: close", + "Content-Length: 0", + `Retry-After: ${RESERVATION_RETRY_AFTER_SECONDS}`, + "", + "", +].join("\r\n"); + +function rejectHeldConnection( + sockets: Set, + socket: net.Socket, +): void { + sockets.add(socket); + socket.once("close", () => sockets.delete(socket)); + socket.on("error", () => socket.destroy()); + socket.end(RESERVATION_RESPONSE); +} + function holdPort(port: number, host?: string): Promise { return new Promise((resolve) => { - const holder = net.createServer(); - holder.unref(); - holder.once("error", (error: Error) => resolve({ error })); - holder.listen({ port, host }, () => resolve({ holder })); + const sockets = new Set(); + const server = net.createServer((socket) => + rejectHeldConnection(sockets, socket), + ); + server.unref(); + server.once("error", (error: Error) => resolve({ error })); + server.listen({ port, host }, () => + resolve({ holder: { server, sockets } }), + ); }); } -function closeHolder(holder: net.Server): Promise { +function closeHolder(holder: PortHolder): Promise { return new Promise((resolve) => { - holder.close(() => resolve()); + holder.server.close(() => resolve()); + holder.sockets.forEach((socket) => socket.destroy()); }); } @@ -81,7 +115,7 @@ async function reservePort( const { holder, error } = await holdPort(candidatePort, config.host); if (holder) { return { - port: resolveBoundPort(holder, candidatePort), + port: resolveBoundPort(holder.server, candidatePort), release: () => closeHolder(holder), }; } diff --git a/src/startup-deadline.ts b/src/startup-deadline.ts new file mode 100644 index 0000000..d98fa32 --- /dev/null +++ b/src/startup-deadline.ts @@ -0,0 +1,37 @@ +let listenDeadline: NodeJS.Timeout | undefined; + +/** + * Resolves `true` once `task` settles, or `false` if `timeoutMs` elapses + * first, so a stuck step can never block the startup sequence for good. + */ +export function settlesWithin( + task: Promise, + timeoutMs: number, +): Promise { + let timer: NodeJS.Timeout | undefined; + const timeout = new Promise((resolve) => { + timer = setTimeout(() => resolve(false), timeoutMs); + timer.unref(); + }); + + return Promise.race([task.then(() => true), timeout]).finally(() => + clearTimeout(timer), + ); +} + +/** + * Calls `onMissed` unless `disarmListenDeadline` runs within `timeoutMs`. + */ +export function armListenDeadline( + timeoutMs: number, + onMissed: () => void, +): void { + disarmListenDeadline(); + listenDeadline = setTimeout(onMissed, timeoutMs); + listenDeadline.unref(); +} + +export function disarmListenDeadline(): void { + clearTimeout(listenDeadline); + listenDeadline = undefined; +} diff --git a/src/test/reservation-connections.test.ts b/src/test/reservation-connections.test.ts new file mode 100644 index 0000000..f78cf3f --- /dev/null +++ b/src/test/reservation-connections.test.ts @@ -0,0 +1,199 @@ +import sinon from "sinon"; +import * as net from "node:net"; +import * as http from "node:http"; +import assert from "node:assert"; +import { Logging } from "@antelopejs/interface-core/logging"; +import * as coreRuntime from "@antelopejs/interface-core/runtime"; + +import { setDevMode } from "../dev-mode"; +import type { Config } from "../server-config"; +import { + configure, + getConfig, + getListeningEndpoints, + listenServers, + provide, + start, + stop, +} from "../index"; + +const TEST_HOST = "127.0.0.1"; +const PROBED_RESERVATION_PORT = 25300; +const SILENT_RESERVATION_PORT = 25310; +const STUCK_RESERVATION_PORT = 25320; +const UNLISTENED_PORT = 25330; +const STUCK_RELEASE_TEST_TIMEOUT_MS = 4000; +const START_LISTEN_DEADLINE_MS = 5000; +const PUBLIC_BASE_URL = "https://api.example.com"; +const LISTEN_DEADLINE_MS = 1500; +const SERVICE_UNAVAILABLE_STATUS = 503; +const PROBE_REQUEST = `GET /health HTTP/1.1\r\nHost: ${TEST_HOST}\r\n\r\n`; + +interface ProbeResult { + socket: net.Socket; + response: Promise; +} + +function stubRuntime(): void { + sinon + .stub(coreRuntime, "GetRuntimeInfo") + .resolves({ dev: false, projectPath: "", env: "test" }); + sinon.stub(coreRuntime, "RegisterDevServer").resolves(); +} + +function stubCloseNeverCallingBack(): sinon.SinonStub { + const closeStub = sinon.stub(net.Server.prototype, "close"); + return closeStub.callsFake(function (this: net.Server) { + return closeStub.wrappedMethod.call(this); + }); +} + +function singleServerConfig(port: number): Config { + return { + autoListen: false, + publicBaseUrl: PUBLIC_BASE_URL, + servers: [{ protocol: "http", host: TEST_HOST, port }], + }; +} + +function connect(port: number): Promise { + return new Promise((resolve, reject) => { + const socket = net.connect(port, TEST_HOST, () => resolve(socket)); + socket.once("error", reject); + }); +} + +function collectResponse(socket: net.Socket): Promise { + return new Promise((resolve) => { + const chunks: Buffer[] = []; + socket.on("data", (chunk: Buffer) => chunks.push(chunk)); + socket.on("error", () => undefined); + socket.once("close", () => resolve(Buffer.concat(chunks).toString())); + }); +} + +async function probeWithRequest(port: number): Promise { + const socket = await connect(port); + const response = collectResponse(socket); + socket.end(PROBE_REQUEST); + return { socket, response }; +} + +function rejectAfter(ms: number, message: string): Promise { + return new Promise((_, reject) => { + setTimeout(() => reject(new Error(message)), ms).unref(); + }); +} + +function listenWithinDeadline(): Promise { + return Promise.race([ + listenServers(), + rejectAfter(LISTEN_DEADLINE_MS, "The real server never started listening"), + ]); +} + +function requestStatus(port: number): Promise { + return new Promise((resolve, reject) => { + http + .get({ host: TEST_HOST, port, path: "/health", agent: false }, (res) => { + res.resume(); + resolve(res.statusCode); + }) + .once("error", reject); + }); +} + +describe("Connections accepted during the port reservation", () => { + const openSockets: net.Socket[] = []; + let originalConfig: Config; + + before(() => { + originalConfig = getConfig(); + }); + + after(() => { + if (!originalConfig.servers?.length) { + return; + } + + configure(originalConfig); + start(); + }); + + afterEach(async () => { + openSockets.forEach((socket) => socket.destroy()); + openSockets.length = 0; + await stop(); + sinon.restore(); + setDevMode(false); + }); + + it("starts the real server after a probe hit the reservation", async () => { + stubRuntime(); + await provide(singleServerConfig(PROBED_RESERVATION_PORT)); + + const probe = await probeWithRequest(PROBED_RESERVATION_PORT); + openSockets.push(probe.socket); + + start(); + await listenWithinDeadline(); + + const reservationResponse = await probe.response; + assert.match( + reservationResponse, + new RegExp(`^HTTP/1\\.1 ${SERVICE_UNAVAILABLE_STATUS} `), + ); + assert.equal(getListeningEndpoints()[0].port, PROBED_RESERVATION_PORT); + assert.ok(await requestStatus(PROBED_RESERVATION_PORT)); + }); + + it("starts the real server while a silent client holds a connection", async () => { + stubRuntime(); + await provide(singleServerConfig(SILENT_RESERVATION_PORT)); + + const silentSocket = await connect(SILENT_RESERVATION_PORT); + silentSocket.on("error", () => undefined); + openSockets.push(silentSocket); + + start(); + await listenWithinDeadline(); + + assert.equal(getListeningEndpoints()[0].port, SILENT_RESERVATION_PORT); + assert.ok(await requestStatus(SILENT_RESERVATION_PORT)); + }); + + it("binds the servers even when a reservation never finishes closing", async function () { + this.timeout(STUCK_RELEASE_TEST_TIMEOUT_MS); + stubRuntime(); + const warnStub = sinon.stub(Logging, "Warn"); + await provide(singleServerConfig(STUCK_RESERVATION_PORT)); + + const closeStub = stubCloseNeverCallingBack(); + start(); + await listenServers(); + closeStub.restore(); + + sinon.assert.calledOnce(warnStub); + assert.equal(getListeningEndpoints()[0].port, STUCK_RESERVATION_PORT); + assert.ok(await requestStatus(STUCK_RESERVATION_PORT)); + }); + + it("logs an error when the servers miss the listen deadline", () => { + sinon + .stub(coreRuntime, "GetRuntimeInfo") + .returns(new Promise(() => undefined)); + const errorStub = sinon.stub(Logging, "Error"); + const clock = sinon.useFakeTimers({ + toFake: ["setTimeout", "clearTimeout"], + }); + + configure({ + publicBaseUrl: PUBLIC_BASE_URL, + servers: [{ protocol: "http", host: TEST_HOST, port: UNLISTENED_PORT }], + }); + start(); + clock.tick(START_LISTEN_DEADLINE_MS); + + sinon.assert.calledOnce(errorStub); + }); +}); From 4228ac14b3a4afe603b7418b34c106a269351c72 Mon Sep 17 00:00:00 2001 From: Antony Rizzitelli Date: Sat, 26 Sep 2026 13:36:34 +0200 Subject: [PATCH 2/2] refactor(port-reservation): bind the listening socket once and serve it from start Replace the reserve, close and re-listen handoff with a single listening socket bound during provide and served from start, like Go's net.Listen followed by http.Serve. The socket is created with pauseOnConnect, so connections accepted before start wait unread and are handed to the http or https server, with every later connection, through its connection event. Nothing is closed and bound again between provide and start, so an early client can no longer keep the server from listening. The port fallback rules now live in one place, the listener binding, used by provide and lazily by listenServers after a stop or when the configuration no longer matches the bound address. This removes the 503 placeholder answer, the release timeout and the listen deadline log, which only papered over the handoff. --- README.md | 22 +-- src/index.ts | 134 +++++++------- src/port-binding.ts | 78 +-------- src/port-listener.ts | 211 +++++++++++++++++++++++ src/port-reservation.ts | 164 ------------------ src/startup-deadline.ts | 37 ---- src/test/config-vars.test.ts | 120 ++++++------- src/test/early-connections.test.ts | 202 ++++++++++++++++++++++ src/test/reservation-connections.test.ts | 199 --------------------- 9 files changed, 553 insertions(+), 614 deletions(-) create mode 100644 src/port-listener.ts delete mode 100644 src/port-reservation.ts delete mode 100644 src/startup-deadline.ts create mode 100644 src/test/early-connections.test.ts delete mode 100644 src/test/reservation-connections.test.ts diff --git a/README.md b/README.md index 0493c40..7d6f6e8 100644 --- a/README.md +++ b/README.md @@ -92,11 +92,11 @@ In development (when the runtime reports `dev`), loopback origins — `localhost The module publishes three [config variables](https://antelopejs.com/docs/concepts/configuration#module-config-variables) other modules reference from their own configuration with `${@api.}`: -| Variable | Type | Description | -| --------------------- | ------ | --------------------------------------------------------------------------- | -| `API_PORT` | number | The port the server reserved during `provide`, and the port it later binds. | -| `API_LOCAL_BASE_URL` | string | The same-host origin, always on loopback: `http://127.0.0.1:`. | -| `API_PUBLIC_BASE_URL` | string | The origin external clients must use, from the `publicBaseUrl` key. | +| Variable | Type | Description | +| --------------------- | ------ | ------------------------------------------------------------------------ | +| `API_PORT` | number | The port the server bound during `provide`, and the port it serves. | +| `API_LOCAL_BASE_URL` | string | The same-host origin, always on loopback: `http://127.0.0.1:`. | +| `API_PUBLIC_BASE_URL` | string | The origin external clients must use, from the `publicBaseUrl` key. | All three derive from the **first entry of `servers[]`**, the same entry the dev registry and the frontend discovery already treat as the project's canonical api endpoint. Additional servers are still started, but they are not advertised through config variables. @@ -116,15 +116,15 @@ export default defineConfig({ The variables are published from the `provide` callback, which the core runs before any module constructs. -To publish a port it can guarantee, the module reserves it right there: it binds a throwaway socket on the configured port, holds it while every other module constructs, and releases it immediately before the real `listen()` in `start`. The value other modules receive is therefore the port the server actually binds — never a stale one. +To publish a port it can guarantee, the module binds the listening socket right there, once, and publishes the port read back from it. The socket stays open while every other module constructs, and `start` serves the HTTP or HTTPS server from it — the server never binds a port of its own. The value other modules receive is therefore the port the server actually serves — never a stale one. -The reservation honours the existing port rules: +Connections that arrive before `start` wait unread in the socket's queue and are served, like every later connection, once the server starts. -- `strictPort: true`, or any non-development runtime, reserves exactly the requested port or fails the boot with a `PortReservationError` naming the port. -- In development, the reservation falls back to the next free port (up to 20 above the requested one, then an OS-assigned port), exactly as `listen()` did before. -- `port: 0` reserves an OS-assigned port and publishes it. +Binding honours the existing port rules: -A client that connects while the port is only reserved — a health probe, an eager consumer — receives an immediate `503 Service Unavailable` with `Retry-After: 1` and a closed connection, so it retries once the real server listens. Releasing the reservation never waits on such a connection. If the servers are still not listening five seconds after `start`, the module logs an error. +- `strictPort: true`, or any non-development runtime, binds exactly the requested port or fails the boot with a `PortReservationError` naming the port. +- In development, binding falls back to the next free port (up to 20 above the requested one, then an OS-assigned port). +- `port: 0` binds an OS-assigned port and publishes it. ### One advertised origin diff --git a/src/index.ts b/src/index.ts index 74d6c41..36731e7 100644 --- a/src/index.ts +++ b/src/index.ts @@ -5,23 +5,19 @@ import type { ConfigVars } from "@antelopejs/interface-core/config"; import type { DevServerEndpoint } from "@antelopejs/interface-core/runtime"; import { resolveDevMode } from "./dev-mode"; -import { listenServer } from "./port-binding"; +import { logServerStarted } from "./port-binding"; import type { Config } from "./server-config"; import { buildConfigVars } from "./config-vars"; import { createConfiguredServer } from "./server-factory"; -import { - armListenDeadline, - disarmListenDeadline, - settlesWithin, -} from "./startup-deadline"; import { configure, getConfig, setCorsConfig } from "./module-config"; export { configure, getConfig, setCorsConfig }; import { - releaseReservedPorts, - type ReservedPort, - reserveServerPorts, -} from "./port-reservation"; + type BoundListener, + bindServerPorts, + closeListeners, + isBoundFor, +} from "./port-listener"; import { collectListeningEndpoints, registerDevServerEndpoints, @@ -29,71 +25,54 @@ import { } from "./dev-registry"; import "./middlewares/cors"; -const RESERVATION_RELEASE_TIMEOUT_MS = 1000; -const LISTEN_DEADLINE_MS = 5000; - let servers: net.Server[] = []; let listening = false; -let reservations: ReservedPort[] = []; - -async function releaseReservations(): Promise { - const pending = reservations; - reservations = []; - const isReleased = await settlesWithin( - releaseReservedPorts(pending), - RESERVATION_RELEASE_TIMEOUT_MS, - ); - if (!isReleased) { - Logging.Warn( - `Port reservations still closing after ${RESERVATION_RELEASE_TIMEOUT_MS} ms, binding the servers anyway`, - ); - } -} +let listeners: BoundListener[] = []; -function reportMissedListenDeadline(): void { - Logging.Error( - `Servers are still not listening ${LISTEN_DEADLINE_MS} ms after start`, - ); +function releaseListeners(): Promise { + const pending = listeners; + listeners = []; + return closeListeners(pending); } -async function reserveConfiguredPorts(): Promise { - await releaseReservations(); +async function bindConfiguredPorts(): Promise { + await releaseListeners(); const config = getConfig(); - reservations = await reserveServerPorts( + listeners = await bindServerPorts( config.servers ?? [], await shouldAllowPortFallback(config), ); } /** - * Copies the reserved ports onto the current configuration. Exposed for + * Copies the bound ports onto the current configuration. Exposed for * tests, which reproduce the rebuilt configuration the core may hand to * `construct`. * * `provide` and `construct` receive the configuration through separate * substitution passes, so the object `construct` sees may be a rebuilt * copy carrying the originally requested ports again. Re-applying the - * reservation is what keeps the published `API_PORT` and the port - * `start` binds identical, whatever the core hands over. + * bound ports is what keeps the published `API_PORT` and the port the + * module reports and logs identical, whatever the core hands over. */ export function applyReservedPorts(): void { const servers = getConfig().servers ?? []; - reservations.forEach((reservation, index) => { + listeners.forEach((listener, index) => { const serverConfig = servers[index]; if (serverConfig) { - serverConfig.port = reservation.port; + serverConfig.port = listener.port; } }); } async function publishConfigVars(): Promise { - await reserveConfiguredPorts(); + await bindConfiguredPorts(); try { return buildConfigVars(getConfig()); } catch (error) { - await releaseReservations(); + await releaseListeners(); throw error; } } @@ -103,8 +82,8 @@ async function publishConfigVars(): Promise { * * Nothing has constructed at this point, so this path awaits no other * module's interface: it only reads the runtime information the core - * registers before the module lifecycle starts, and holds a socket on - * the port the server will bind. + * registers before the module lifecycle starts, and binds the listening + * socket the server is served from once it starts. */ export async function provide(config: Config): Promise { configure(config); @@ -135,31 +114,62 @@ function closeServers(): Promise { ); servers = []; listening = false; - disarmListenDeadline(); return Promise.all(closing).then(() => undefined); } +/** + * Creates the configured servers and, unless `autoListen` is `false`, + * serves them from the listening sockets. Calling it again replaces the + * servers: the sockets stay bound and hand their connections to the new + * ones. + */ export function start(): void { - const serversClosed = closeServers(); + void closeServers(); servers = (getConfig().servers ?? []).map((serverConfig) => createConfiguredServer(serverConfig), ); - if (getConfig().autoListen !== false && servers.length > 0) { - armListenDeadline(LISTEN_DEADLINE_MS, reportMissedListenDeadline); - void serversClosed - .then(() => listenServers()) - .catch((error: unknown) => { - const message = error instanceof Error ? error.message : String(error); - Logging.Error(`Unable to start listening servers: ${message}`); - }); + if (getConfig().autoListen !== false) { + listenServers().catch((error: unknown) => { + const message = error instanceof Error ? error.message : String(error); + Logging.Error(`Unable to start listening servers: ${message}`); + }); } } export function getListeningEndpoints(): DevServerEndpoint[] { - return collectListeningEndpoints(servers, getConfig().servers); + if (!listening) { + return []; + } + + return collectListeningEndpoints( + listeners.map((listener) => listener.socket), + getConfig().servers, + ); } +async function ensureListeners(): Promise { + if (isBoundFor(listeners, getConfig().servers ?? [])) { + return; + } + + await bindConfiguredPorts(); +} + +function serveListeners(): void { + const configs = getConfig().servers ?? []; + listeners.forEach((listener, index) => { + listener.serve(servers[index]); + logServerStarted(configs[index], listener.requestedPort, listener.port); + }); +} + +/** + * Serves the created servers from their listening sockets. The sockets + * bound during `provide` are kept when the configuration still names + * their address; otherwise — after a `stop`, when `provide` never ran, or + * when the configuration changed — they are bound again. + */ export async function listenServers(): Promise { if (listening || servers.length === 0) { return; @@ -168,24 +178,16 @@ export async function listenServers(): Promise { listening = true; try { - const allowPortFallback = await shouldAllowPortFallback(getConfig()); - await releaseReservations(); - await Promise.all( - (getConfig().servers ?? []).map((serverConfig, index) => - listenServer(servers[index], serverConfig, allowPortFallback), - ), - ); + await ensureListeners(); } catch (error) { listening = false; throw error; - } finally { - disarmListenDeadline(); } + serveListeners(); await registerDevServerEndpoints(getListeningEndpoints()); } export async function stop(): Promise { - await releaseReservations(); - await closeServers(); + await Promise.all([closeServers(), releaseListeners()]); } diff --git a/src/port-binding.ts b/src/port-binding.ts index ef7c5ba..4b1543a 100644 --- a/src/port-binding.ts +++ b/src/port-binding.ts @@ -10,33 +10,6 @@ import { const MAX_PORT_FALLBACK_OFFSET = 20; -function listenOnce( - server: net.Server, - port: number, - host?: string, -): Promise { - return new Promise((resolve, reject) => { - const onListening = () => { - cleanup(); - resolve(); - }; - - const onError = (error: Error) => { - cleanup(); - reject(error); - }; - - const cleanup = () => { - server.off("listening", onListening); - server.off("error", onError); - }; - - server.once("listening", onListening); - server.once("error", onError); - server.listen(port, host); - }); -} - export function isPortInUseError(error: unknown): boolean { return (error as NodeJS.ErrnoException)?.code === "EADDRINUSE"; } @@ -81,33 +54,10 @@ export function buildCandidatePorts( return [...sequentialPorts, RANDOM_PORT]; } -async function listenServerWithFallback( - server: net.Server, - config: ServerConfig, - requestedPort: number, - allowPortFallback: boolean, -): Promise { - const candidatePorts = buildCandidatePorts(requestedPort, allowPortFallback); - let portInUseError: unknown = new Error( - `Unable to bind ${config.protocol} server on port ${requestedPort}`, - ); - - for (const candidatePort of candidatePorts) { - try { - await listenOnce(server, candidatePort, config.host); - return resolveBoundPort(server, candidatePort); - } catch (error) { - if (!isPortInUseError(error)) { - throw error; - } - portInUseError = error; - } - } - - throw portInUseError; -} - -function logServerStarted( +/** + * Logs where a server listens, naming the fallback when it moved. + */ +export function logServerStarted( config: ServerConfig, requestedPort: number, boundPort: number, @@ -122,23 +72,3 @@ function logServerStarted( `Port ${requestedPort} in use, listening on ${serverUrl} instead`, ); } - -export async function listenServer( - server: net.Server, - config: ServerConfig, - allowPortFallback: boolean, -): Promise { - if (server.listening) { - return; - } - - const requestedPort = resolveRequestedPort(config); - const boundPort = await listenServerWithFallback( - server, - config, - requestedPort, - allowPortFallback, - ); - config.port = boundPort; - logServerStarted(config, requestedPort, boundPort); -} diff --git a/src/port-listener.ts b/src/port-listener.ts new file mode 100644 index 0000000..fda1f34 --- /dev/null +++ b/src/port-listener.ts @@ -0,0 +1,211 @@ +import * as net from "node:net"; + +import type { ServerConfig } from "./server-config"; +import { + buildCandidatePorts, + isPortInUseError, + resolveBoundPort, + resolveRequestedPort, +} from "./port-binding"; + +interface ListenerAddress { + host?: string; + /** The port read from the bound socket. */ + port: number; + /** The port the configuration asked for, before any fallback. */ + requestedPort: number; +} + +/** + * The listening socket of a configured server, bound once and served later. + * + * The socket is bound during `provide`, so the published `API_PORT` is + * the port read back from it, and it stays open until `stop`: the http or + * https server never binds a port of its own. Connections accepted before + * `serve` are queued paused and unread, then handed to the server with + * everything that follows, like Go's `net.Listen` followed by + * `http.Serve`. + */ +export interface BoundListener extends ListenerAddress { + /** The bound socket, as the dev registry reads it. */ + socket: net.Server; + /** Hands every queued and future connection to `server`. */ + serve: (server: net.Server) => void; + /** Stops accepting, drops unserved connections and frees the port. */ + close: () => Promise; +} + +/** + * Raised when no port could be bound for a configured server. + */ +export class PortReservationError extends Error { + override readonly name = "PortReservationError"; + readonly code = "EADDRINUSE"; +} + +interface ListenAttempt { + socket?: net.Server; + error?: unknown; +} + +const STRICT_PORT_HINT = + " Port fallback is disabled; set strictPort to false in development to accept the next free port."; + +function listenOnce(port: number, host?: string): Promise { + return new Promise((resolve) => { + const socket = net.createServer({ pauseOnConnect: true }); + socket.once("error", (error: Error) => resolve({ error })); + socket.listen({ port, host }, () => resolve({ socket })); + }); +} + +/** + * Hands a paused connection to `server`. Emitting `connection` is how a + * server takes over a socket it did not accept; https wraps it in TLS. + */ +function handOver(server: net.Server, connection: net.Socket): void { + server.emit("connection", connection); + connection.resume(); +} + +/** + * Makes `server` treat `connection` events as its own. The server never + * calls `listen()`, and `listening` is the event that arms its connection + * tracking: request and headers timeouts, and closing idle keep-alive + * connections on `close()`. + */ +function adoptAsListening(server: net.Server): void { + server.emit("listening"); +} + +/** + * Turns a bound socket into a listener. Until served it does not keep the + * process alive, so a boot that fails after `provide` still exits. + */ +function createListener( + socket: net.Server, + address: ListenerAddress, +): BoundListener { + const queued = new Set(); + let target: net.Server | undefined; + + socket.unref(); + socket.on("connection", (connection: net.Socket) => { + if (target) { + handOver(target, connection); + return; + } + queued.add(connection); + connection.once("close", () => queued.delete(connection)); + }); + + return { + ...address, + socket, + serve: (server) => { + target = server; + socket.ref(); + adoptAsListening(server); + queued.forEach((connection) => handOver(server, connection)); + queued.clear(); + }, + close: () => + new Promise((resolve) => { + socket.close(() => resolve()); + queued.forEach((connection) => connection.destroy()); + }), + }; +} + +function buildBindError( + config: ServerConfig, + requestedPort: number, + allowPortFallback: boolean, +): PortReservationError { + const reason = `Unable to bind port ${requestedPort} for the ${config.protocol} server: the port is already in use.`; + return new PortReservationError( + allowPortFallback ? reason : `${reason}${STRICT_PORT_HINT}`, + ); +} + +/** + * Binds the port a server is configured for, falling back to the next + * free ports then to an OS-assigned one when `allowPortFallback` is set. + */ +async function bindListener( + config: ServerConfig, + allowPortFallback: boolean, +): Promise { + const requestedPort = resolveRequestedPort(config); + + for (const candidatePort of buildCandidatePorts( + requestedPort, + allowPortFallback, + )) { + const { socket, error } = await listenOnce(candidatePort, config.host); + if (socket) { + return createListener(socket, { + host: config.host, + port: resolveBoundPort(socket, candidatePort), + requestedPort, + }); + } + if (!isPortInUseError(error)) { + throw error; + } + } + + throw buildBindError(config, requestedPort, allowPortFallback); +} + +/** + * Binds a listener for every configured server and writes the bound port + * back into its configuration, so the published config variables and the + * served socket agree on a single value. + */ +export async function bindServerPorts( + configs: ServerConfig[], + allowPortFallback: boolean, +): Promise { + const listeners: BoundListener[] = []; + + try { + for (const config of configs) { + const listener = await bindListener(config, allowPortFallback); + config.port = listener.port; + listeners.push(listener); + } + } catch (error) { + await closeListeners(listeners); + throw error; + } + + return listeners; +} + +/** + * Tells whether `listeners` are bound exactly where `configs` ask, so + * they can keep serving instead of being bound again. + */ +export function isBoundFor( + listeners: BoundListener[], + configs: ServerConfig[], +): boolean { + return ( + listeners.length === configs.length && + listeners.every( + (listener, index) => + listener.host === configs[index].host && + listener.port === resolveRequestedPort(configs[index]), + ) + ); +} + +/** + * Closes every listener, freeing their ports. + */ +export function closeListeners(listeners: BoundListener[]): Promise { + return Promise.all(listeners.map((listener) => listener.close())).then( + () => undefined, + ); +} diff --git a/src/port-reservation.ts b/src/port-reservation.ts deleted file mode 100644 index 57fde9b..0000000 --- a/src/port-reservation.ts +++ /dev/null @@ -1,164 +0,0 @@ -import * as net from "node:net"; - -import type { ServerConfig } from "./server-config"; -import { - buildCandidatePorts, - isPortInUseError, - resolveBoundPort, - resolveRequestedPort, -} from "./port-binding"; - -/** - * A port held open by a throwaway socket between `construct` and `start`. - * - * Holding the socket is what keeps the published `API_PORT` and the port - * the server eventually binds in sync: the probe-then-use race shrinks to - * the instant between `release()` and the real `listen()`, instead of - * spanning the construction of every other module. Any client reaching - * the port meanwhile gets an immediate `503` and a closed connection, so - * an early probe never keeps the holding socket from closing. - */ -export interface ReservedPort { - port: number; - /** Closes the holding socket; call it right before the real `listen()`. */ - release: () => Promise; -} - -/** - * Raised when no port could be held for a configured server. - */ -export class PortReservationError extends Error { - override readonly name = "PortReservationError"; - readonly code = "EADDRINUSE"; -} - -interface PortHolder { - server: net.Server; - sockets: Set; -} - -interface PortHold { - holder?: PortHolder; - error?: unknown; -} - -const STRICT_PORT_HINT = - " Port fallback is disabled; set strictPort to false in development to accept the next free port."; - -const RESERVATION_RETRY_AFTER_SECONDS = 1; - -const RESERVATION_RESPONSE = [ - "HTTP/1.1 503 Service Unavailable", - "Connection: close", - "Content-Length: 0", - `Retry-After: ${RESERVATION_RETRY_AFTER_SECONDS}`, - "", - "", -].join("\r\n"); - -function rejectHeldConnection( - sockets: Set, - socket: net.Socket, -): void { - sockets.add(socket); - socket.once("close", () => sockets.delete(socket)); - socket.on("error", () => socket.destroy()); - socket.end(RESERVATION_RESPONSE); -} - -function holdPort(port: number, host?: string): Promise { - return new Promise((resolve) => { - const sockets = new Set(); - const server = net.createServer((socket) => - rejectHeldConnection(sockets, socket), - ); - server.unref(); - server.once("error", (error: Error) => resolve({ error })); - server.listen({ port, host }, () => - resolve({ holder: { server, sockets } }), - ); - }); -} - -function closeHolder(holder: PortHolder): Promise { - return new Promise((resolve) => { - holder.server.close(() => resolve()); - holder.sockets.forEach((socket) => socket.destroy()); - }); -} - -function buildReservationError( - config: ServerConfig, - requestedPort: number, - allowPortFallback: boolean, -): PortReservationError { - const reason = `Unable to reserve port ${requestedPort} for the ${config.protocol} server: the port is already in use.`; - return new PortReservationError( - allowPortFallback ? reason : `${reason}${STRICT_PORT_HINT}`, - ); -} - -/** - * Holds the port a server will later bind, honouring the same requested - * port, fallback range and random-port semantics as the real `listen()`. - */ -async function reservePort( - config: ServerConfig, - allowPortFallback: boolean, -): Promise { - const requestedPort = resolveRequestedPort(config); - - for (const candidatePort of buildCandidatePorts( - requestedPort, - allowPortFallback, - )) { - const { holder, error } = await holdPort(candidatePort, config.host); - if (holder) { - return { - port: resolveBoundPort(holder.server, candidatePort), - release: () => closeHolder(holder), - }; - } - if (!isPortInUseError(error)) { - throw error; - } - } - - throw buildReservationError(config, requestedPort, allowPortFallback); -} - -/** - * Reserves a port for every configured server and writes the reserved - * port back into its configuration, so the published config variables and - * the later `listen()` agree on a single value. - */ -export async function reserveServerPorts( - configs: ServerConfig[], - allowPortFallback: boolean, -): Promise { - const reservations: ReservedPort[] = []; - - try { - for (const config of configs) { - const reservation = await reservePort(config, allowPortFallback); - config.port = reservation.port; - reservations.push(reservation); - } - } catch (error) { - await releaseReservedPorts(reservations); - throw error; - } - - return reservations; -} - -/** - * Closes every holding socket, freeing the ports for their real servers. - */ -export function releaseReservedPorts( - reservations: ReservedPort[], -): Promise { - return Promise.all( - reservations.map((reservation) => reservation.release()), - ).then(() => undefined); -} diff --git a/src/startup-deadline.ts b/src/startup-deadline.ts deleted file mode 100644 index d98fa32..0000000 --- a/src/startup-deadline.ts +++ /dev/null @@ -1,37 +0,0 @@ -let listenDeadline: NodeJS.Timeout | undefined; - -/** - * Resolves `true` once `task` settles, or `false` if `timeoutMs` elapses - * first, so a stuck step can never block the startup sequence for good. - */ -export function settlesWithin( - task: Promise, - timeoutMs: number, -): Promise { - let timer: NodeJS.Timeout | undefined; - const timeout = new Promise((resolve) => { - timer = setTimeout(() => resolve(false), timeoutMs); - timer.unref(); - }); - - return Promise.race([task.then(() => true), timeout]).finally(() => - clearTimeout(timer), - ); -} - -/** - * Calls `onMissed` unless `disarmListenDeadline` runs within `timeoutMs`. - */ -export function armListenDeadline( - timeoutMs: number, - onMissed: () => void, -): void { - disarmListenDeadline(); - listenDeadline = setTimeout(onMissed, timeoutMs); - listenDeadline.unref(); -} - -export function disarmListenDeadline(): void { - clearTimeout(listenDeadline); - listenDeadline = undefined; -} diff --git a/src/test/config-vars.test.ts b/src/test/config-vars.test.ts index 15cefdb..7ee50b7 100644 --- a/src/test/config-vars.test.ts +++ b/src/test/config-vars.test.ts @@ -9,11 +9,11 @@ import { LOOPBACK_HOST } from "../server-origin"; import { collectListeningEndpoints } from "../dev-registry"; import type { Config, ServerConfig } from "../server-config"; import { + type BoundListener, + bindServerPorts, + closeListeners, PortReservationError, - releaseReservedPorts, - type ReservedPort, - reserveServerPorts, -} from "../port-reservation"; +} from "../port-listener"; import { API_LOCAL_BASE_URL, API_PORT, @@ -33,17 +33,17 @@ import { } from "../index"; const TEST_HOST = "127.0.0.1"; -const RESERVED_FREE_PORT = 25200; -const RESERVED_TAKEN_PORT = 25210; -const RESERVED_STRICT_PORT = 25220; +const BOUND_FREE_PORT = 25200; +const BOUND_TAKEN_PORT = 25210; +const BOUND_STRICT_PORT = 25220; const CONSTRUCT_FREE_PORT = 25230; const CONSTRUCT_TAKEN_PORT = 25240; const CONSTRUCT_STRICT_PORT = 25250; const SECONDARY_PORT = 25260; const REBUILT_CONFIG_PORT = 25270; const AGREEMENT_PORT = 25290; -const HELD_RESERVATION_PORT = 25280; -const RESERVATION_HOLD_MS = 250; +const HELD_LISTENER_PORT = 25280; +const LISTENER_HOLD_MS = 250; const RANDOM_PORT = 0; const PUBLIC_BASE_URL = "https://api.example.com"; @@ -108,15 +108,15 @@ describe("Published config variables", () => { const vars = buildConfigVars({ publicBaseUrl: PUBLIC_BASE_URL, servers: [ - { protocol: "http", host: TEST_HOST, port: RESERVED_FREE_PORT }, + { protocol: "http", host: TEST_HOST, port: BOUND_FREE_PORT }, { protocol: "http", host: TEST_HOST, port: SECONDARY_PORT }, ], }); - assert.equal(vars[API_PORT], RESERVED_FREE_PORT); + assert.equal(vars[API_PORT], BOUND_FREE_PORT); assert.equal( vars[API_LOCAL_BASE_URL], - `http://${TEST_HOST}:${RESERVED_FREE_PORT}`, + `http://${TEST_HOST}:${BOUND_FREE_PORT}`, ); assert.equal(vars[API_PUBLIC_BASE_URL], PUBLIC_BASE_URL); }); @@ -136,14 +136,12 @@ describe("Published config variables", () => { for (const { bindHost, urlHost } of hostCases) { const vars = buildConfigVars({ publicBaseUrl: PUBLIC_BASE_URL, - servers: [ - { protocol: "http", host: bindHost, port: RESERVED_FREE_PORT }, - ], + servers: [{ protocol: "http", host: bindHost, port: BOUND_FREE_PORT }], }); assert.equal( vars[API_LOCAL_BASE_URL], - `http://${urlHost}:${RESERVED_FREE_PORT}`, + `http://${urlHost}:${BOUND_FREE_PORT}`, `bind host ${String(bindHost)}`, ); } @@ -165,9 +163,7 @@ describe("Published config variables", () => { setDevMode(true); const vars = buildConfigVars({ - servers: [ - { protocol: "http", host: TEST_HOST, port: RESERVED_FREE_PORT }, - ], + servers: [{ protocol: "http", host: TEST_HOST, port: BOUND_FREE_PORT }], }); assert.equal(vars[API_PUBLIC_BASE_URL], vars[API_LOCAL_BASE_URL]); @@ -180,7 +176,7 @@ describe("Published config variables", () => { () => buildConfigVars({ servers: [ - { protocol: "http", host: TEST_HOST, port: RESERVED_FREE_PORT }, + { protocol: "http", host: TEST_HOST, port: BOUND_FREE_PORT }, ], }), (error: unknown) => @@ -192,79 +188,77 @@ describe("Published config variables", () => { it("strips trailing slashes from the configured public base url", () => { const vars = buildConfigVars({ publicBaseUrl: `${PUBLIC_BASE_URL}//`, - servers: [ - { protocol: "http", host: TEST_HOST, port: RESERVED_FREE_PORT }, - ], + servers: [{ protocol: "http", host: TEST_HOST, port: BOUND_FREE_PORT }], }); assert.equal(vars[API_PUBLIC_BASE_URL], PUBLIC_BASE_URL); }); }); -describe("Port reservation", () => { +describe("Port binding", () => { const blockers: net.Server[] = []; - let reservations: ReservedPort[] = []; + let listeners: BoundListener[] = []; async function blockPort(port: number): Promise { blockers.push(await occupyPort(port)); } afterEach(async () => { - await releaseReservedPorts(reservations); - reservations = []; + await closeListeners(listeners); + listeners = []; await Promise.all(blockers.map((blocker) => closeServer(blocker))); blockers.length = 0; }); - it("reserves the requested port and writes it back to the config", async () => { - const servers = singleServerConfig(RESERVED_FREE_PORT).servers ?? []; + it("binds the requested port and writes it back to the config", async () => { + const servers = singleServerConfig(BOUND_FREE_PORT).servers ?? []; - reservations = await reserveServerPorts(servers, false); + listeners = await bindServerPorts(servers, false); - assert.equal(reservations[0].port, RESERVED_FREE_PORT); - assert.equal(servers[0].port, RESERVED_FREE_PORT); + assert.equal(listeners[0].port, BOUND_FREE_PORT); + assert.equal(servers[0].port, BOUND_FREE_PORT); }); - it("reserves the next free port when fallback is allowed", async () => { - await blockPort(RESERVED_TAKEN_PORT); - const servers = singleServerConfig(RESERVED_TAKEN_PORT).servers ?? []; + it("binds the next free port when fallback is allowed", async () => { + await blockPort(BOUND_TAKEN_PORT); + const servers = singleServerConfig(BOUND_TAKEN_PORT).servers ?? []; - reservations = await reserveServerPorts(servers, true); + listeners = await bindServerPorts(servers, true); - assert.equal(reservations[0].port, RESERVED_TAKEN_PORT + 1); - assert.equal(servers[0].port, RESERVED_TAKEN_PORT + 1); + assert.equal(listeners[0].port, BOUND_TAKEN_PORT + 1); + assert.equal(servers[0].port, BOUND_TAKEN_PORT + 1); }); it("fails with a named error when fallback is not allowed", async () => { - await blockPort(RESERVED_STRICT_PORT); - const servers = singleServerConfig(RESERVED_STRICT_PORT).servers ?? []; + await blockPort(BOUND_STRICT_PORT); + const servers = singleServerConfig(BOUND_STRICT_PORT).servers ?? []; await assert.rejects( - reserveServerPorts(servers, false), - (error: unknown) => + bindServerPorts(servers, false), + (error: Error) => error instanceof PortReservationError && error.name === "PortReservationError" && - error.message.includes(String(RESERVED_STRICT_PORT)), + error.message.includes(String(BOUND_STRICT_PORT)), ); }); - it("reserves an operating system assigned port for port 0", async () => { + it("binds an operating system assigned port for port 0", async () => { const servers = singleServerConfig(RANDOM_PORT).servers ?? []; - reservations = await reserveServerPorts(servers, false); + listeners = await bindServerPorts(servers, false); - assert.ok(reservations[0].port > 0); - assert.equal(servers[0].port, reservations[0].port); + assert.ok(listeners[0].port > 0); + assert.equal(servers[0].port, listeners[0].port); }); - it("releases every reservation when one of them fails", async () => { - await blockPort(RESERVED_STRICT_PORT); + it("closes every listener when one of them fails", async () => { + await blockPort(BOUND_STRICT_PORT); await assert.rejects( - reserveServerPorts( + bindServerPorts( [ { protocol: "http", host: TEST_HOST, port: SECONDARY_PORT }, - { protocol: "http", host: TEST_HOST, port: RESERVED_STRICT_PORT }, + { protocol: "http", host: TEST_HOST, port: BOUND_STRICT_PORT }, ], false, ), @@ -318,7 +312,7 @@ describe("Config variables publication", () => { assert.equal(vars[API_PUBLIC_BASE_URL], PUBLIC_BASE_URL); }); - it("publishes the reserved port and binds it without drift", async () => { + it("publishes the bound port and serves it without drift", async () => { stubRuntime(true); await blockPort(CONSTRUCT_TAKEN_PORT); @@ -333,7 +327,7 @@ describe("Config variables publication", () => { assert.equal(endpoints[0].port, vars[API_PORT]); }); - it("re-applies the reservation to the configuration construct receives", async () => { + it("re-applies the bound port to the configuration construct receives", async () => { stubRuntime(true); await blockPort(REBUILT_CONFIG_PORT); @@ -352,23 +346,23 @@ describe("Config variables publication", () => { assert.equal(getListeningEndpoints()[0].port, vars[API_PORT]); }); - it("holds the reservation across a long provide to start window", async () => { + it("holds the bound port across a long provide to start window", async () => { stubRuntime(false); - const vars = await publish(singleServerConfig(HELD_RESERVATION_PORT)); - assert.equal(vars[API_PORT], HELD_RESERVATION_PORT); + const vars = await publish(singleServerConfig(HELD_LISTENER_PORT)); + assert.equal(vars[API_PORT], HELD_LISTENER_PORT); - await assert.rejects(occupyPort(HELD_RESERVATION_PORT), isPortInUseError); - await delay(RESERVATION_HOLD_MS); - await assert.rejects(occupyPort(HELD_RESERVATION_PORT), isPortInUseError); + await assert.rejects(occupyPort(HELD_LISTENER_PORT), isPortInUseError); + await delay(LISTENER_HOLD_MS); + await assert.rejects(occupyPort(HELD_LISTENER_PORT), isPortInUseError); start(); await listenServers(); - assert.equal(getListeningEndpoints()[0].port, HELD_RESERVATION_PORT); + assert.equal(getListeningEndpoints()[0].port, HELD_LISTENER_PORT); }); - it("fails the boot when strictPort cannot reserve the requested port", async () => { + it("fails the boot when strictPort cannot bind the requested port", async () => { stubRuntime(true); await blockPort(CONSTRUCT_STRICT_PORT); @@ -453,12 +447,12 @@ describe("Advertised host agreement", () => { for (const host of loopbackHosts) { const vars = buildConfigVars({ publicBaseUrl: PUBLIC_BASE_URL, - servers: [{ protocol: "http", host, port: RESERVED_FREE_PORT }], + servers: [{ protocol: "http", host, port: BOUND_FREE_PORT }], }); assert.equal( vars[API_LOCAL_BASE_URL], - `http://${LOOPBACK_HOST}:${RESERVED_FREE_PORT}`, + `http://${LOOPBACK_HOST}:${BOUND_FREE_PORT}`, `bind host ${String(host)}`, ); } diff --git a/src/test/early-connections.test.ts b/src/test/early-connections.test.ts new file mode 100644 index 0000000..f9252db --- /dev/null +++ b/src/test/early-connections.test.ts @@ -0,0 +1,202 @@ +import sinon from "sinon"; +import * as net from "node:net"; +import * as http from "node:http"; +import assert from "node:assert"; +import * as coreRuntime from "@antelopejs/interface-core/runtime"; + +import { setDevMode } from "../dev-mode"; +import type { Config } from "../server-config"; +import { API_PORT } from "../config-vars"; +import { + configure, + getConfig, + getListeningEndpoints, + listenServers, + provide, + start, + stop, +} from "../index"; + +const TEST_HOST = "127.0.0.1"; +const EARLY_REQUEST_PORT = 25300; +const SILENT_CLIENT_PORT = 25310; +const RESTART_PORT = 25320; +const REPLACED_SERVER_PORT = 25330; +const KEEP_ALIVE_PORT = 25340; +const RANDOM_PORT = 0; +const PUBLIC_BASE_URL = "https://api.example.com"; +const UNSERVED_WINDOW_MS = 100; +const DEADLINE_MS = 1500; +const SERVICE_UNAVAILABLE_STATUS = 503; + +function stubRuntime(): void { + sinon + .stub(coreRuntime, "GetRuntimeInfo") + .resolves({ dev: false, projectPath: "", env: "test" }); + sinon.stub(coreRuntime, "RegisterDevServer").resolves(); +} + +function singleServerConfig(port: number): Config { + return { + autoListen: false, + publicBaseUrl: PUBLIC_BASE_URL, + servers: [{ protocol: "http", host: TEST_HOST, port }], + }; +} + +function connect(port: number): Promise { + return new Promise((resolve, reject) => { + const socket = net.connect(port, TEST_HOST, () => resolve(socket)); + socket.once("error", reject); + }); +} + +function delay(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)); +} + +function rejectAfter(ms: number, message: string): Promise { + return new Promise((_, reject) => { + setTimeout(() => reject(new Error(message)), ms).unref(); + }); +} + +function withinDeadline(task: Promise, message: string): Promise { + return Promise.race([task, rejectAfter(DEADLINE_MS, message)]); +} + +function requestStatus( + port: number, + agent: http.Agent | false = false, +): Promise { + return new Promise((resolve, reject) => { + http + .get({ host: TEST_HOST, port, path: "/health", agent }, (res) => { + res.resume(); + res.once("end", () => resolve(res.statusCode)); + }) + .once("error", reject); + }); +} + +async function startServing(): Promise { + start(); + await withinDeadline(listenServers(), "The server never started serving"); +} + +async function isSettledWithin( + task: Promise, + ms: number, +): Promise { + let isSettled = false; + const markSettled = () => { + isSettled = true; + }; + task.then(markSettled, markSettled); + await delay(ms); + return isSettled; +} + +describe("Connections accepted before the server starts", () => { + const openSockets: net.Socket[] = []; + let originalConfig: Config; + + before(() => { + originalConfig = getConfig(); + }); + + after(() => { + if (!originalConfig.servers?.length) { + return; + } + + configure(originalConfig); + start(); + }); + + afterEach(async () => { + openSockets.forEach((socket) => socket.destroy()); + openSockets.length = 0; + await stop(); + sinon.restore(); + setDevMode(false); + }); + + it("answers a request sent during init once the server starts", async () => { + stubRuntime(); + await provide(singleServerConfig(EARLY_REQUEST_PORT)); + + const earlyStatus = requestStatus(EARLY_REQUEST_PORT); + assert.equal(await isSettledWithin(earlyStatus, UNSERVED_WINDOW_MS), false); + + await startServing(); + + const status = await withinDeadline( + earlyStatus, + "The early request was never answered", + ); + assert.notEqual(status, SERVICE_UNAVAILABLE_STATUS); + assert.equal(status, await requestStatus(EARLY_REQUEST_PORT)); + }); + + it("starts while a silent client holds a connection", async () => { + stubRuntime(); + await provide(singleServerConfig(SILENT_CLIENT_PORT)); + + const silentSocket = await connect(SILENT_CLIENT_PORT); + silentSocket.on("error", () => undefined); + openSockets.push(silentSocket); + + await startServing(); + + assert.equal(getListeningEndpoints()[0].port, SILENT_CLIENT_PORT); + assert.ok(await requestStatus(SILENT_CLIENT_PORT)); + }); + + it("publishes and serves the port bound for port 0", async () => { + stubRuntime(); + const vars = await provide(singleServerConfig(RANDOM_PORT)); + assert.ok(Number(vars[API_PORT]) > 0); + + await startServing(); + + assert.equal(getListeningEndpoints()[0].port, vars[API_PORT]); + assert.ok(await requestStatus(Number(vars[API_PORT]))); + }); + + it("binds the port again when starting after a stop", async () => { + stubRuntime(); + await provide(singleServerConfig(RESTART_PORT)); + await startServing(); + await stop(); + + assert.deepEqual(getListeningEndpoints(), []); + await assert.rejects(requestStatus(RESTART_PORT)); + + await startServing(); + + assert.equal(getListeningEndpoints()[0].port, RESTART_PORT); + assert.ok(await requestStatus(RESTART_PORT)); + }); + + it("hands the bound port to the new servers when started twice", async () => { + stubRuntime(); + await provide(singleServerConfig(REPLACED_SERVER_PORT)); + await startServing(); + await startServing(); + + assert.equal(getListeningEndpoints()[0].port, REPLACED_SERVER_PORT); + assert.ok(await requestStatus(REPLACED_SERVER_PORT)); + }); + + it("closes idle keep-alive connections on stop", async () => { + stubRuntime(); + const agent = new http.Agent({ keepAlive: true }); + await provide(singleServerConfig(KEEP_ALIVE_PORT)); + await startServing(); + await requestStatus(KEEP_ALIVE_PORT, agent); + + await withinDeadline(stop(), "Stop waited on an idle connection"); + agent.destroy(); + }); +}); diff --git a/src/test/reservation-connections.test.ts b/src/test/reservation-connections.test.ts deleted file mode 100644 index f78cf3f..0000000 --- a/src/test/reservation-connections.test.ts +++ /dev/null @@ -1,199 +0,0 @@ -import sinon from "sinon"; -import * as net from "node:net"; -import * as http from "node:http"; -import assert from "node:assert"; -import { Logging } from "@antelopejs/interface-core/logging"; -import * as coreRuntime from "@antelopejs/interface-core/runtime"; - -import { setDevMode } from "../dev-mode"; -import type { Config } from "../server-config"; -import { - configure, - getConfig, - getListeningEndpoints, - listenServers, - provide, - start, - stop, -} from "../index"; - -const TEST_HOST = "127.0.0.1"; -const PROBED_RESERVATION_PORT = 25300; -const SILENT_RESERVATION_PORT = 25310; -const STUCK_RESERVATION_PORT = 25320; -const UNLISTENED_PORT = 25330; -const STUCK_RELEASE_TEST_TIMEOUT_MS = 4000; -const START_LISTEN_DEADLINE_MS = 5000; -const PUBLIC_BASE_URL = "https://api.example.com"; -const LISTEN_DEADLINE_MS = 1500; -const SERVICE_UNAVAILABLE_STATUS = 503; -const PROBE_REQUEST = `GET /health HTTP/1.1\r\nHost: ${TEST_HOST}\r\n\r\n`; - -interface ProbeResult { - socket: net.Socket; - response: Promise; -} - -function stubRuntime(): void { - sinon - .stub(coreRuntime, "GetRuntimeInfo") - .resolves({ dev: false, projectPath: "", env: "test" }); - sinon.stub(coreRuntime, "RegisterDevServer").resolves(); -} - -function stubCloseNeverCallingBack(): sinon.SinonStub { - const closeStub = sinon.stub(net.Server.prototype, "close"); - return closeStub.callsFake(function (this: net.Server) { - return closeStub.wrappedMethod.call(this); - }); -} - -function singleServerConfig(port: number): Config { - return { - autoListen: false, - publicBaseUrl: PUBLIC_BASE_URL, - servers: [{ protocol: "http", host: TEST_HOST, port }], - }; -} - -function connect(port: number): Promise { - return new Promise((resolve, reject) => { - const socket = net.connect(port, TEST_HOST, () => resolve(socket)); - socket.once("error", reject); - }); -} - -function collectResponse(socket: net.Socket): Promise { - return new Promise((resolve) => { - const chunks: Buffer[] = []; - socket.on("data", (chunk: Buffer) => chunks.push(chunk)); - socket.on("error", () => undefined); - socket.once("close", () => resolve(Buffer.concat(chunks).toString())); - }); -} - -async function probeWithRequest(port: number): Promise { - const socket = await connect(port); - const response = collectResponse(socket); - socket.end(PROBE_REQUEST); - return { socket, response }; -} - -function rejectAfter(ms: number, message: string): Promise { - return new Promise((_, reject) => { - setTimeout(() => reject(new Error(message)), ms).unref(); - }); -} - -function listenWithinDeadline(): Promise { - return Promise.race([ - listenServers(), - rejectAfter(LISTEN_DEADLINE_MS, "The real server never started listening"), - ]); -} - -function requestStatus(port: number): Promise { - return new Promise((resolve, reject) => { - http - .get({ host: TEST_HOST, port, path: "/health", agent: false }, (res) => { - res.resume(); - resolve(res.statusCode); - }) - .once("error", reject); - }); -} - -describe("Connections accepted during the port reservation", () => { - const openSockets: net.Socket[] = []; - let originalConfig: Config; - - before(() => { - originalConfig = getConfig(); - }); - - after(() => { - if (!originalConfig.servers?.length) { - return; - } - - configure(originalConfig); - start(); - }); - - afterEach(async () => { - openSockets.forEach((socket) => socket.destroy()); - openSockets.length = 0; - await stop(); - sinon.restore(); - setDevMode(false); - }); - - it("starts the real server after a probe hit the reservation", async () => { - stubRuntime(); - await provide(singleServerConfig(PROBED_RESERVATION_PORT)); - - const probe = await probeWithRequest(PROBED_RESERVATION_PORT); - openSockets.push(probe.socket); - - start(); - await listenWithinDeadline(); - - const reservationResponse = await probe.response; - assert.match( - reservationResponse, - new RegExp(`^HTTP/1\\.1 ${SERVICE_UNAVAILABLE_STATUS} `), - ); - assert.equal(getListeningEndpoints()[0].port, PROBED_RESERVATION_PORT); - assert.ok(await requestStatus(PROBED_RESERVATION_PORT)); - }); - - it("starts the real server while a silent client holds a connection", async () => { - stubRuntime(); - await provide(singleServerConfig(SILENT_RESERVATION_PORT)); - - const silentSocket = await connect(SILENT_RESERVATION_PORT); - silentSocket.on("error", () => undefined); - openSockets.push(silentSocket); - - start(); - await listenWithinDeadline(); - - assert.equal(getListeningEndpoints()[0].port, SILENT_RESERVATION_PORT); - assert.ok(await requestStatus(SILENT_RESERVATION_PORT)); - }); - - it("binds the servers even when a reservation never finishes closing", async function () { - this.timeout(STUCK_RELEASE_TEST_TIMEOUT_MS); - stubRuntime(); - const warnStub = sinon.stub(Logging, "Warn"); - await provide(singleServerConfig(STUCK_RESERVATION_PORT)); - - const closeStub = stubCloseNeverCallingBack(); - start(); - await listenServers(); - closeStub.restore(); - - sinon.assert.calledOnce(warnStub); - assert.equal(getListeningEndpoints()[0].port, STUCK_RESERVATION_PORT); - assert.ok(await requestStatus(STUCK_RESERVATION_PORT)); - }); - - it("logs an error when the servers miss the listen deadline", () => { - sinon - .stub(coreRuntime, "GetRuntimeInfo") - .returns(new Promise(() => undefined)); - const errorStub = sinon.stub(Logging, "Error"); - const clock = sinon.useFakeTimers({ - toFake: ["setTimeout", "clearTimeout"], - }); - - configure({ - publicBaseUrl: PUBLIC_BASE_URL, - servers: [{ protocol: "http", host: TEST_HOST, port: UNLISTENED_PORT }], - }); - start(); - clock.tick(START_LISTEN_DEADLINE_MS); - - sinon.assert.calledOnce(errorStub); - }); -});