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') + }) +})