From 541c53a61a1369ef0a8232d2da0a6d6eaa068e45 Mon Sep 17 00:00:00 2001 From: Gregor MacLennan Date: Tue, 16 Jun 2026 17:59:44 +0100 Subject: [PATCH] fix: make server listen/close idempotent and race-safe Wrap listen/close in a start-stop state machine so overlapping or repeated calls are serialised and the server ends in the state of the last call. listen() is idempotent for compatible ports and throws when called with different explicit ports while already listening. In-flight getRemotePort() is rejected with a typed SERVER_CLOSED error. --- package-lock.json | 16 ++++ package.json | 1 + src/index.ts | 110 ++++++++++++++++++------- src/lib/errors.ts | 5 ++ test/server-lifecycle.test.ts | 147 ++++++++++++++++++++++++++++++++++ 5 files changed, 250 insertions(+), 29 deletions(-) create mode 100644 test/server-lifecycle.test.ts diff --git a/package-lock.json b/package-lock.json index 1a104ce..e94371a 100644 --- a/package-lock.json +++ b/package-lock.json @@ -20,6 +20,7 @@ "p-mutex": "^0.1.0", "secret-stream-http": "^1.0.1", "smp-noto-glyphs": "^1.0.0-pre.0", + "start-stop-state-machine": "^2.0.0", "styled-map-package-api": "^5.0.0-pre.4", "typebox": "^1.0.61", "typed-event-target": "^3.4.0" @@ -5677,6 +5678,15 @@ "dev": true, "license": "MIT" }, + "node_modules/start-stop-state-machine": { + "version": "2.0.0", + "resolved": "https://registry.npmjs.org/start-stop-state-machine/-/start-stop-state-machine-2.0.0.tgz", + "integrity": "sha512-PLwGZc4FsCb9S9MTRwdTTgKjNPILQ3fzS6QpCdeDyrI+XFTZnhgd1Aj9uko5SuCs5+8Knx8XPYPa9brMDsJLmg==", + "license": "MIT", + "dependencies": { + "tiny-typed-emitter": "^2.1.0" + } + }, "node_modules/statuses": { "version": "2.0.2", "resolved": "https://registry.npmjs.org/statuses/-/statuses-2.0.2.tgz", @@ -5976,6 +5986,12 @@ "integrity": "sha512-SVqEcMZBsZF9mA78rjzCrYrUs37LMJk3ShZ851ygZYW1cMeIjs9mL57KO6Iv5mmjSQnOe/29/VAfGXo+oRCiVw==", "license": "MIT" }, + "node_modules/tiny-typed-emitter": { + "version": "2.1.0", + "resolved": "https://registry.npmjs.org/tiny-typed-emitter/-/tiny-typed-emitter-2.1.0.tgz", + "integrity": "sha512-qVtvMxeXbVej0cQWKqVSSAHmKZEHAvxdF8HEUBFWts8h+xEo5m/lEiPakuyZ3BnCBjOD8i24kzNOiOLLgsSxhA==", + "license": "MIT" + }, "node_modules/tinybench": { "version": "2.9.0", "resolved": "https://registry.npmjs.org/tinybench/-/tinybench-2.9.0.tgz", diff --git a/package.json b/package.json index 27db500..a207d97 100644 --- a/package.json +++ b/package.json @@ -82,6 +82,7 @@ "p-mutex": "^0.1.0", "secret-stream-http": "^1.0.1", "smp-noto-glyphs": "^1.0.0-pre.0", + "start-stop-state-machine": "^2.0.0", "styled-map-package-api": "^5.0.0-pre.4", "typebox": "^1.0.61", "typed-event-target": "^3.4.0" diff --git a/src/index.ts b/src/index.ts index f93b97b..9380a6d 100644 --- a/src/index.ts +++ b/src/index.ts @@ -9,9 +9,12 @@ import { Agent, createServer as createSecretStreamServer, } from 'secret-stream-http' +import StartStopStateMachine from 'start-stop-state-machine' import { Context } from './context.js' +import { errors } from './lib/errors.js' import { fetchAPI } from './lib/fetch-api.js' +import { noop } from './lib/utils.js' import { RootRouter } from './routes/root.js' import type { FetchContext } from './types.js' @@ -60,7 +63,7 @@ export function createServer(options: ServerOptions) { options.keyPair = Agent.keyPair() } - let deferredListen = pDefer() + let deferredListen = createDeferredListen() const context = new Context({ ...options, keyPair: options.keyPair, @@ -100,41 +103,90 @@ export function createServer(options: ServerOptions) { localHttpServer.on('connection', onConnection) secretStreamServer.on('connection', onConnection) + const stateMachine = new StartStopStateMachine({ + start: listen, + stop: close, + }) + + async function listen(opts: ListenOptions = {}) { + localHttpServer.listen(opts.localPort, '127.0.0.1') + secretStreamServer.listen(opts.remotePort, '0.0.0.0') + await Promise.all([ + once(localHttpServer, 'listening'), + once(secretStreamServer, 'listening'), + ]) + const localPort = (localHttpServer.address() as AddressInfo).port + const remotePort = (secretStreamServer.address() as AddressInfo).port + deferredListen.resolve({ localPort, remotePort }) + return { localPort, remotePort } + } + + async function close() { + // Throw away any pending listen promises since the server is closing + deferredListen.reject(new errors.SERVER_CLOSED()) + // Remove connection listeners + localHttpServer.off('connection', onConnection) + secretStreamServer.off('connection', onConnection) + localHttpServer.close() + secretStreamServer.close() + // Destroy all active connections to ensure clean shutdown + for (const socket of connections) { + socket.destroy() + } + connections.clear() + await Promise.all([ + once(localHttpServer, 'close'), + once(secretStreamServer, 'close'), + ]) + await context.close() + // Reset deferred listen for potential restart with different ports + deferredListen = createDeferredListen() + } + + function createDeferredListen() { + const deferred = pDefer() + // close() rejects this to unblock getRemotePort(); swallow the rejection + // so it isn't reported as unhandled when nothing is awaiting it. + deferred.promise.catch(noop) + return deferred + } + return { async listen(opts: ListenOptions = {}) { - localHttpServer.listen(opts.localPort, '127.0.0.1') - secretStreamServer.listen(opts.remotePort, '0.0.0.0') - await Promise.all([ - once(localHttpServer, 'listening'), - once(secretStreamServer, 'listening'), - ]) - const localPort = (localHttpServer.address() as AddressInfo).port - const remotePort = (secretStreamServer.address() as AddressInfo).port - deferredListen.resolve({ localPort, remotePort }) - return { localPort, remotePort } - }, - async close() { - // Remove connection listeners - localHttpServer.off('connection', onConnection) - secretStreamServer.off('connection', onConnection) - localHttpServer.close() - secretStreamServer.close() - // Destroy all active connections to ensure clean shutdown - for (const socket of connections) { - socket.destroy() + const { value } = stateMachine.state + // When already (or nearly) listening, listen() is idempotent — but + // only for compatible ports. Requesting different explicit ports is a + // programming error: close() first to listen on new ports. + if (value === 'starting' || value === 'started') { + const current = await stateMachine.started() + assertCompatiblePorts(opts, current) + return current } - connections.clear() - await Promise.all([ - once(localHttpServer, 'close'), - once(secretStreamServer, 'close'), - ]) - await context.close() - // Reset deferred listen for potential restart with different ports - deferredListen = pDefer() + return stateMachine.start(opts) + }, + close() { + return stateMachine.stop() }, } } +function assertCompatiblePorts(opts: ListenOptions, current: ListenResult) { + // A falsy requested port (undefined or 0) means "any port", so it never + // conflicts with the port already in use. + assert( + !opts.localPort || opts.localPort === current.localPort, + new Error( + `Server is already listening on local port ${current.localPort}; call close() before listening on port ${opts.localPort}`, + ), + ) + assert( + !opts.remotePort || opts.remotePort === current.remotePort, + new Error( + `Server is already listening on remote port ${current.remotePort}; call close() before listening on port ${opts.remotePort}`, + ), + ) +} + function validateOptions(options: unknown): asserts options is ServerOptions { assert( typeof options === 'object' && options !== null, diff --git a/src/lib/errors.ts b/src/lib/errors.ts index 7b646ae..92e142d 100644 --- a/src/lib/errors.ts +++ b/src/lib/errors.ts @@ -120,6 +120,11 @@ const errorsList = [ message: 'Invalid request', status: 400, }, + { + code: 'SERVER_CLOSED', + message: 'Server is closed', + status: 503, + }, ] as const satisfies Array export const errors = {} as Record< diff --git a/test/server-lifecycle.test.ts b/test/server-lifecycle.test.ts new file mode 100644 index 0000000..f4cefad --- /dev/null +++ b/test/server-lifecycle.test.ts @@ -0,0 +1,147 @@ +import fs from 'node:fs/promises' +import os from 'node:os' +import path from 'node:path' + +import { Agent as SecretStreamAgent } from 'secret-stream-http' +import { describe, it, expect } from 'vitest' + +import { createServer, errors } from '../src/index.js' +import { DEMOTILES_Z2, OSM_BRIGHT_Z6 } from './helpers.js' + +/** + * Create a server backed by a throwaway copy of the custom map fixture. The + * server is closed and the temp dir removed when the test finishes. + */ +async function createTestServer(t: { + onTestFinished: (fn: () => Promise | void) => void +}) { + const tmpDir = await fs.mkdtemp(path.join(os.tmpdir(), 'map-server-test-')) + const tmpCustomMapPath = path.join(tmpDir, 'custom-map.smp') + await fs.copyFile(OSM_BRIGHT_Z6, tmpCustomMapPath) + + const server = createServer({ + defaultOnlineStyleUrl: 'https://demotiles.maplibre.org/style.json', + customMapPath: tmpCustomMapPath, + fallbackMapPath: DEMOTILES_Z2, + keyPair: SecretStreamAgent.keyPair(), + }) + + t.onTestFinished(async () => { + await server.close().catch(() => {}) + await fs.rm(tmpDir, { recursive: true, force: true }).catch(() => {}) + }) + + return server +} + +describe('Server lifecycle', () => { + it('listen() is idempotent', async (t) => { + const server = await createTestServer(t) + + const first = await server.listen({ localPort: 0, remotePort: 0 }) + // A second listen() without a close() in between must not throw. The old + // code called net.Server.listen() a second time, which throws + // ERR_SERVER_ALREADY_LISTEN ("Listen method has been called more than + // once without closing"). It should instead resolve to the same ports. + const second = await server.listen({ localPort: 0, remotePort: 0 }) + + expect(second).toEqual(first) + + const response = await fetch( + `http://127.0.0.1:${first.localPort}/maps/custom/style.json`, + ) + expect(response.status).toBe(200) + }) + + it('concurrent listen() calls resolve to the same ports', async (t) => { + const server = await createTestServer(t) + + const [a, b] = await Promise.all([ + server.listen({ localPort: 0, remotePort: 0 }), + server.listen({ localPort: 0, remotePort: 0 }), + ]) + + expect(a).toEqual(b) + + const response = await fetch( + `http://127.0.0.1:${a.localPort}/maps/custom/style.json`, + ) + expect(response.status).toBe(200) + }) + + it('listen() with different ports while already listening throws', async (t) => { + const server = await createTestServer(t) + + const { localPort } = await server.listen({ localPort: 0, remotePort: 0 }) + + // Requesting a different explicit port without closing first is a + // programming error and must reject rather than silently keep the old port. + await expect( + server.listen({ localPort: localPort + 1 }), + ).rejects.toThrow(/already listening/) + + // The original listener is untouched. + const response = await fetch( + `http://127.0.0.1:${localPort}/maps/custom/style.json`, + ) + expect(response.status).toBe(200) + }) + + it('close() is idempotent', async (t) => { + const server = await createTestServer(t) + + await server.listen({ localPort: 0, remotePort: 0 }) + + await server.close() + // A second close() must resolve cleanly rather than calling + // net.Server.close() on an already-closed server. + await expect(server.close()).resolves.toBeUndefined() + }) + + it('a close() racing a listen() on a started server leaves it working (last call wins)', async (t) => { + const server = await createTestServer(t) + + await server.listen({ localPort: 0, remotePort: 0 }) + + // Fire close() then listen() without awaiting close() in between. The + // state machine serialises them so the last call (listen) wins. + const closePromise = server.close() + const listenPromise = server.listen({ localPort: 0, remotePort: 0 }) + + await closePromise + const result = await listenPromise + + const response = await fetch( + `http://127.0.0.1:${result.localPort}/maps/custom/style.json`, + ) + expect(response.status).toBe(200) + }) + + it('close() called while listen() is in flight settles cleanly and the server can restart', async (t) => { + const server = await createTestServer(t) + + // Start listening but request a close before listen() has resolved. The + // old code raced net.Server.close() against an in-progress listen() and + // hung forever waiting on a 'close' event that never fired. The state + // machine serialises the two so both promises settle (last call wins: + // the server ends up stopped). + const listenPromise = server.listen({ localPort: 0, remotePort: 0 }) + const closePromise = server.close() + + await Promise.all([listenPromise, closePromise]) + + // The server must be restartable and functional after the race. + const result = await server.listen({ localPort: 0, remotePort: 0 }) + const response = await fetch( + `http://127.0.0.1:${result.localPort}/maps/custom/style.json`, + ) + expect(response.status).toBe(200) + }, 15_000) + + it('exposes a typed SERVER_CLOSED error', () => { + const err = new errors.SERVER_CLOSED() + expect(err.code).toBe('SERVER_CLOSED') + expect(err.status).toBe(503) + expect(err.message).toBe('Server is closed') + }) +})