Skip to content

Commit 610377b

Browse files
authored
fix: make server listen/close idempotent and race-safe (#56)
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.
1 parent 404a558 commit 610377b

5 files changed

Lines changed: 250 additions & 29 deletions

File tree

package-lock.json

Lines changed: 16 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

package.json

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,7 @@
8282
"p-mutex": "^0.1.0",
8383
"secret-stream-http": "^1.0.1",
8484
"smp-noto-glyphs": "^1.0.0-pre.0",
85+
"start-stop-state-machine": "^2.0.0",
8586
"styled-map-package-api": "^5.0.0-pre.4",
8687
"typebox": "^1.0.61",
8788
"typed-event-target": "^3.4.0"

src/index.ts

Lines changed: 81 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -9,9 +9,12 @@ import {
99
Agent,
1010
createServer as createSecretStreamServer,
1111
} from 'secret-stream-http'
12+
import StartStopStateMachine from 'start-stop-state-machine'
1213

1314
import { Context } from './context.js'
15+
import { errors } from './lib/errors.js'
1416
import { fetchAPI } from './lib/fetch-api.js'
17+
import { noop } from './lib/utils.js'
1518
import { RootRouter } from './routes/root.js'
1619
import type { FetchContext } from './types.js'
1720

@@ -60,7 +63,7 @@ export function createServer(options: ServerOptions) {
6063
options.keyPair = Agent.keyPair()
6164
}
6265

63-
let deferredListen = pDefer<ListenResult>()
66+
let deferredListen = createDeferredListen()
6467
const context = new Context({
6568
...options,
6669
keyPair: options.keyPair,
@@ -100,41 +103,90 @@ export function createServer(options: ServerOptions) {
100103
localHttpServer.on('connection', onConnection)
101104
secretStreamServer.on('connection', onConnection)
102105

106+
const stateMachine = new StartStopStateMachine({
107+
start: listen,
108+
stop: close,
109+
})
110+
111+
async function listen(opts: ListenOptions = {}) {
112+
localHttpServer.listen(opts.localPort, '127.0.0.1')
113+
secretStreamServer.listen(opts.remotePort, '0.0.0.0')
114+
await Promise.all([
115+
once(localHttpServer, 'listening'),
116+
once(secretStreamServer, 'listening'),
117+
])
118+
const localPort = (localHttpServer.address() as AddressInfo).port
119+
const remotePort = (secretStreamServer.address() as AddressInfo).port
120+
deferredListen.resolve({ localPort, remotePort })
121+
return { localPort, remotePort }
122+
}
123+
124+
async function close() {
125+
// Throw away any pending listen promises since the server is closing
126+
deferredListen.reject(new errors.SERVER_CLOSED())
127+
// Remove connection listeners
128+
localHttpServer.off('connection', onConnection)
129+
secretStreamServer.off('connection', onConnection)
130+
localHttpServer.close()
131+
secretStreamServer.close()
132+
// Destroy all active connections to ensure clean shutdown
133+
for (const socket of connections) {
134+
socket.destroy()
135+
}
136+
connections.clear()
137+
await Promise.all([
138+
once(localHttpServer, 'close'),
139+
once(secretStreamServer, 'close'),
140+
])
141+
await context.close()
142+
// Reset deferred listen for potential restart with different ports
143+
deferredListen = createDeferredListen()
144+
}
145+
146+
function createDeferredListen() {
147+
const deferred = pDefer<ListenResult>()
148+
// close() rejects this to unblock getRemotePort(); swallow the rejection
149+
// so it isn't reported as unhandled when nothing is awaiting it.
150+
deferred.promise.catch(noop)
151+
return deferred
152+
}
153+
103154
return {
104155
async listen(opts: ListenOptions = {}) {
105-
localHttpServer.listen(opts.localPort, '127.0.0.1')
106-
secretStreamServer.listen(opts.remotePort, '0.0.0.0')
107-
await Promise.all([
108-
once(localHttpServer, 'listening'),
109-
once(secretStreamServer, 'listening'),
110-
])
111-
const localPort = (localHttpServer.address() as AddressInfo).port
112-
const remotePort = (secretStreamServer.address() as AddressInfo).port
113-
deferredListen.resolve({ localPort, remotePort })
114-
return { localPort, remotePort }
115-
},
116-
async close() {
117-
// Remove connection listeners
118-
localHttpServer.off('connection', onConnection)
119-
secretStreamServer.off('connection', onConnection)
120-
localHttpServer.close()
121-
secretStreamServer.close()
122-
// Destroy all active connections to ensure clean shutdown
123-
for (const socket of connections) {
124-
socket.destroy()
156+
const { value } = stateMachine.state
157+
// When already (or nearly) listening, listen() is idempotent — but
158+
// only for compatible ports. Requesting different explicit ports is a
159+
// programming error: close() first to listen on new ports.
160+
if (value === 'starting' || value === 'started') {
161+
const current = await stateMachine.started()
162+
assertCompatiblePorts(opts, current)
163+
return current
125164
}
126-
connections.clear()
127-
await Promise.all([
128-
once(localHttpServer, 'close'),
129-
once(secretStreamServer, 'close'),
130-
])
131-
await context.close()
132-
// Reset deferred listen for potential restart with different ports
133-
deferredListen = pDefer<ListenResult>()
165+
return stateMachine.start(opts)
166+
},
167+
close() {
168+
return stateMachine.stop()
134169
},
135170
}
136171
}
137172

173+
function assertCompatiblePorts(opts: ListenOptions, current: ListenResult) {
174+
// A falsy requested port (undefined or 0) means "any port", so it never
175+
// conflicts with the port already in use.
176+
assert(
177+
!opts.localPort || opts.localPort === current.localPort,
178+
new Error(
179+
`Server is already listening on local port ${current.localPort}; call close() before listening on port ${opts.localPort}`,
180+
),
181+
)
182+
assert(
183+
!opts.remotePort || opts.remotePort === current.remotePort,
184+
new Error(
185+
`Server is already listening on remote port ${current.remotePort}; call close() before listening on port ${opts.remotePort}`,
186+
),
187+
)
188+
}
189+
138190
function validateOptions(options: unknown): asserts options is ServerOptions {
139191
assert(
140192
typeof options === 'object' && options !== null,

src/lib/errors.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -120,6 +120,11 @@ const errorsList = [
120120
message: 'Invalid request',
121121
status: 400,
122122
},
123+
{
124+
code: 'SERVER_CLOSED',
125+
message: 'Server is closed',
126+
status: 503,
127+
},
123128
] as const satisfies Array<ErrorDefinition>
124129

125130
export const errors = {} as Record<

test/server-lifecycle.test.ts

Lines changed: 147 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,147 @@
1+
import fs from 'node:fs/promises'
2+
import os from 'node:os'
3+
import path from 'node:path'
4+
5+
import { Agent as SecretStreamAgent } from 'secret-stream-http'
6+
import { describe, it, expect } from 'vitest'
7+
8+
import { createServer, errors } from '../src/index.js'
9+
import { DEMOTILES_Z2, OSM_BRIGHT_Z6 } from './helpers.js'
10+
11+
/**
12+
* Create a server backed by a throwaway copy of the custom map fixture. The
13+
* server is closed and the temp dir removed when the test finishes.
14+
*/
15+
async function createTestServer(t: {
16+
onTestFinished: (fn: () => Promise<void> | void) => void
17+
}) {
18+
const tmpDir = await fs.mkdtemp(path.join(os.tmpdir(), 'map-server-test-'))
19+
const tmpCustomMapPath = path.join(tmpDir, 'custom-map.smp')
20+
await fs.copyFile(OSM_BRIGHT_Z6, tmpCustomMapPath)
21+
22+
const server = createServer({
23+
defaultOnlineStyleUrl: 'https://demotiles.maplibre.org/style.json',
24+
customMapPath: tmpCustomMapPath,
25+
fallbackMapPath: DEMOTILES_Z2,
26+
keyPair: SecretStreamAgent.keyPair(),
27+
})
28+
29+
t.onTestFinished(async () => {
30+
await server.close().catch(() => {})
31+
await fs.rm(tmpDir, { recursive: true, force: true }).catch(() => {})
32+
})
33+
34+
return server
35+
}
36+
37+
describe('Server lifecycle', () => {
38+
it('listen() is idempotent', async (t) => {
39+
const server = await createTestServer(t)
40+
41+
const first = await server.listen({ localPort: 0, remotePort: 0 })
42+
// A second listen() without a close() in between must not throw. The old
43+
// code called net.Server.listen() a second time, which throws
44+
// ERR_SERVER_ALREADY_LISTEN ("Listen method has been called more than
45+
// once without closing"). It should instead resolve to the same ports.
46+
const second = await server.listen({ localPort: 0, remotePort: 0 })
47+
48+
expect(second).toEqual(first)
49+
50+
const response = await fetch(
51+
`http://127.0.0.1:${first.localPort}/maps/custom/style.json`,
52+
)
53+
expect(response.status).toBe(200)
54+
})
55+
56+
it('concurrent listen() calls resolve to the same ports', async (t) => {
57+
const server = await createTestServer(t)
58+
59+
const [a, b] = await Promise.all([
60+
server.listen({ localPort: 0, remotePort: 0 }),
61+
server.listen({ localPort: 0, remotePort: 0 }),
62+
])
63+
64+
expect(a).toEqual(b)
65+
66+
const response = await fetch(
67+
`http://127.0.0.1:${a.localPort}/maps/custom/style.json`,
68+
)
69+
expect(response.status).toBe(200)
70+
})
71+
72+
it('listen() with different ports while already listening throws', async (t) => {
73+
const server = await createTestServer(t)
74+
75+
const { localPort } = await server.listen({ localPort: 0, remotePort: 0 })
76+
77+
// Requesting a different explicit port without closing first is a
78+
// programming error and must reject rather than silently keep the old port.
79+
await expect(
80+
server.listen({ localPort: localPort + 1 }),
81+
).rejects.toThrow(/already listening/)
82+
83+
// The original listener is untouched.
84+
const response = await fetch(
85+
`http://127.0.0.1:${localPort}/maps/custom/style.json`,
86+
)
87+
expect(response.status).toBe(200)
88+
})
89+
90+
it('close() is idempotent', async (t) => {
91+
const server = await createTestServer(t)
92+
93+
await server.listen({ localPort: 0, remotePort: 0 })
94+
95+
await server.close()
96+
// A second close() must resolve cleanly rather than calling
97+
// net.Server.close() on an already-closed server.
98+
await expect(server.close()).resolves.toBeUndefined()
99+
})
100+
101+
it('a close() racing a listen() on a started server leaves it working (last call wins)', async (t) => {
102+
const server = await createTestServer(t)
103+
104+
await server.listen({ localPort: 0, remotePort: 0 })
105+
106+
// Fire close() then listen() without awaiting close() in between. The
107+
// state machine serialises them so the last call (listen) wins.
108+
const closePromise = server.close()
109+
const listenPromise = server.listen({ localPort: 0, remotePort: 0 })
110+
111+
await closePromise
112+
const result = await listenPromise
113+
114+
const response = await fetch(
115+
`http://127.0.0.1:${result.localPort}/maps/custom/style.json`,
116+
)
117+
expect(response.status).toBe(200)
118+
})
119+
120+
it('close() called while listen() is in flight settles cleanly and the server can restart', async (t) => {
121+
const server = await createTestServer(t)
122+
123+
// Start listening but request a close before listen() has resolved. The
124+
// old code raced net.Server.close() against an in-progress listen() and
125+
// hung forever waiting on a 'close' event that never fired. The state
126+
// machine serialises the two so both promises settle (last call wins:
127+
// the server ends up stopped).
128+
const listenPromise = server.listen({ localPort: 0, remotePort: 0 })
129+
const closePromise = server.close()
130+
131+
await Promise.all([listenPromise, closePromise])
132+
133+
// The server must be restartable and functional after the race.
134+
const result = await server.listen({ localPort: 0, remotePort: 0 })
135+
const response = await fetch(
136+
`http://127.0.0.1:${result.localPort}/maps/custom/style.json`,
137+
)
138+
expect(response.status).toBe(200)
139+
}, 15_000)
140+
141+
it('exposes a typed SERVER_CLOSED error', () => {
142+
const err = new errors.SERVER_CLOSED()
143+
expect(err.code).toBe('SERVER_CLOSED')
144+
expect(err.status).toBe(503)
145+
expect(err.message).toBe('Server is closed')
146+
})
147+
})

0 commit comments

Comments
 (0)