Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
110 changes: 81 additions & 29 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'

Expand Down Expand Up @@ -60,7 +63,7 @@ export function createServer(options: ServerOptions) {
options.keyPair = Agent.keyPair()
}

let deferredListen = pDefer<ListenResult>()
let deferredListen = createDeferredListen()
const context = new Context({
...options,
keyPair: options.keyPair,
Expand Down Expand Up @@ -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<ListenResult>()
// 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<ListenResult>()
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,
Expand Down
5 changes: 5 additions & 0 deletions src/lib/errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,11 @@ const errorsList = [
message: 'Invalid request',
status: 400,
},
{
code: 'SERVER_CLOSED',
message: 'Server is closed',
status: 503,
},
] as const satisfies Array<ErrorDefinition>

export const errors = {} as Record<
Expand Down
147 changes: 147 additions & 0 deletions test/server-lifecycle.test.ts
Original file line number Diff line number Diff line change
@@ -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) => 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')
})
})
Loading