diff --git a/README.md b/README.md index 06108fe..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,13 +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: + +- `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 cebc98a..36731e7 100644 --- a/src/index.ts +++ b/src/index.ts @@ -5,7 +5,7 @@ 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"; @@ -13,10 +13,11 @@ 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, @@ -26,52 +27,52 @@ import "./middlewares/cors"; let servers: net.Server[] = []; let listening = false; -let reservations: ReservedPort[] = []; +let listeners: BoundListener[] = []; -function releaseReservations(): Promise { - const pending = reservations; - reservations = []; - return releaseReservedPorts(pending); +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; } } @@ -81,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); @@ -116,26 +117,59 @@ function closeServers(): Promise { 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) { - void serversClosed - .then(() => listenServers()) - .catch((error: unknown) => { - const message = error instanceof Error ? error.message : String(error); - Logging.Error(`Unable to start listening servers: ${message}`); - }); + 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; @@ -144,22 +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; } + 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 644c3b1..0000000 --- a/src/port-reservation.ts +++ /dev/null @@ -1,130 +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. - */ -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 PortHold { - holder?: 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 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 })); - }); -} - -function closeHolder(holder: net.Server): Promise { - return new Promise((resolve) => { - holder.close(() => resolve()); - }); -} - -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, 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/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(); + }); +});