Files
windmill/multiplayer/test/cold_start.test.mjs
T
Alexander PetricandClaude Opus 5.5 e8c3f9514e fix(multiplayer): don't drop client messages during cold-start token verification (#11343)
* fix(multiplayer): don't drop client messages during cold-start token verification

`wss.on('connection')` awaits `verifyToken()` before `setupWSConnection()`
attaches the 'message' listener. On a cold process that await includes the
first `/api/debug/jwks` fetch (~30ms on ECS). A y-websocket client sends sync
step 1 the instant the socket opens, and `ws` drops messages emitted with no
listener attached, so that step 1 was lost and never answered with step 2 —
the client's provider never became `synced`.

Buffer messages from the moment the connection is accepted and replay them, in
order, once `setupWSConnection()` has installed its handlers. Rejected
connections drop the buffer and close with the same 4401/4403 codes as before.

Also prefetch the public key at startup when WINDMILL_BASE_URL is set. That is
insurance, not the fix: a connection arriving before the prefetch resolves
still relies on the buffer.

Adds `npm test` in multiplayer/ (node:test, no docker or backend needed) with a
fake JWKS endpoint that answers with a delay, which holds the cold window open
and makes the race deterministic; wired into the existing test_extra CI job.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix(multiplayer): cap what an unauthenticated peer can buffer pre-auth

Review follow-up.

The pre-auth buffer was unbounded: `ws` sets no `maxPayload` here and the JWKS
fetch has no timeout, so a peer that never authenticates could stream frames
into memory for as long as `verifyToken` was stalled. Cap it at 32 frames /
1 MiB — a real client only has sync step 1 and its first awareness update in
flight there — and close 1009 past that, dropping what was buffered.

A socket closed during verification (by the peer, or by that cap) is no longer
handed to setupWSConnection: it would be added to `doc.conns` with a 'close'
listener that can never fire.

The startup prefetch's .catch was dead code — getPublicKey() logs its own
failures and resolves to null rather than rejecting.

Test helper: pin REQUIRE_SIGNED_MULTIPLAYER_REQUESTS and BASE_INTERNAL_URL so an
ambient value cannot turn the rejection tests into false passes; bind the JWKS
server on port 0 instead of a released probe port, and retry the spawned server
on EADDRINUSE; destroy still-delayed JWKS responses on teardown, since
server.close() waits for in-flight requests.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* test(multiplayer): gate the JWKS response instead of delaying it

Review follow-up.

The cold window was held open by a 1500 ms delay on the fake JWKS response, but
that timer started when the startup prefetch reached the fake server, not when
the client sent its first frame. A slow enough machine could load the key before
the client connected, and the race test would then pass without ever exercising
the buffer — a false pass.

The fake JWKS server now parks every response until the test calls release(), so
the server provably holds no key while the client is sending. The race test
releases only after both frames are written to the socket, and asserts the
server has not logged the key as loaded at that point; the flood test never
releases until after the cap has closed the connection.

What is left to wall-clock time is 250 ms for bytes already written to the socket
to cross loopback into an otherwise idle server, rather than a window that had to
cover process startup, connect and handshake.

Also drops the prefetch precondition from the forged-token and flood tests so
each test still maps to one behaviour. Suite runs in ~1.1s instead of ~5.3s.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-09-25 15:54:11 +02:00

231 lines
9.0 KiB
JavaScript

/**
* Regression tests for the cold-start window of the multiplayer server.
*
* On a fresh process the first `wss.on('connection')` handler awaits the JWKS
* fetch before `setupWSConnection` attaches a 'message' listener. A y-websocket
* client sends sync step 1 the instant the socket opens, so any message landing
* inside that window used to be dropped by `ws` and never answered.
*
* The fake JWKS endpoint below parks the response until the test releases it,
* so the server provably has no key while the client is sending and the race is
* reproduced on every run rather than on a lucky schedule.
*/
import assert from 'node:assert/strict'
import test from 'node:test'
import crypto from 'node:crypto'
import { setTimeout as delay } from 'node:timers/promises'
import { WebSocket } from 'ws'
import * as Y from 'yjs'
import * as syncProtocol from 'y-protocols/sync'
import * as encoding from 'lib0/encoding'
import * as decoding from 'lib0/decoding'
import { mintToken, startJwksServer, startMultiplayerServer, waitFor } from './helpers.mjs'
const WORKSPACE = 'test_workspace'
const DOC_PATH = `${WORKSPACE}/f/foo/bar`
// Grace for frames the client has already written to the socket to be delivered
// over loopback and read by the (otherwise idle) server, before the key is
// released. Only these already-written bytes have to land in this time.
const FLIGHT_MARGIN_MS = 250
// Logged by server.mjs once the JWKS fetch resolves. While it is absent the
// server demonstrably has no key, so anything it receives is in the pre-auth window.
const KEY_LOADED = 'Successfully loaded Ed25519 public key'
const messageSync = 0
const syncStep2 = 1
/** `messageSync` + sync step 1 for an empty document, as y-websocket sends on open. */
function syncStep1Message() {
const encoder = encoding.createEncoder()
encoding.writeVarUint(encoder, messageSync)
syncProtocol.writeSyncStep1(encoder, new Y.Doc())
return encoding.toUint8Array(encoder)
}
/** `messageSync` + a Yjs update inserting `text` into the 'content' text type. */
function updateMessage(text) {
const doc = new Y.Doc()
doc.getText('content').insert(0, text)
const encoder = encoding.createEncoder()
encoding.writeVarUint(encoder, messageSync)
syncProtocol.writeUpdate(encoder, Y.encodeStateAsUpdate(doc))
return encoding.toUint8Array(encoder)
}
/** Read the (messageType, syncType) pair that prefixes a sync message. */
function messageKind(data) {
const decoder = decoding.createDecoder(data)
const messageType = decoding.readVarUint(decoder)
if (messageType !== messageSync) return { messageType }
return { messageType, syncType: decoding.readVarUint(decoder) }
}
/** Open a socket and collect every frame it receives, plus its close code. */
function openClient(url, { onOpen } = {}) {
const ws = new WebSocket(url)
const received = []
const client = { ws, received, closeCode: undefined }
ws.on('message', (data) => received.push(new Uint8Array(data)))
ws.on('close', (code) => {
client.closeCode = code
})
ws.on('error', () => {})
if (onOpen) ws.on('open', () => onOpen(ws))
return client
}
function hasKind(client, syncType) {
return client.received.some((data) => messageKind(data).syncType === syncType)
}
test('cold start: a sync step 1 sent while the key is still being fetched is answered', { timeout: 60000 }, async (t) => {
const jwks = await startJwksServer({ hold: true })
const server = await startMultiplayerServer({ WINDMILL_BASE_URL: jwks.baseUrl })
t.after(async () => {
await server.close()
await jwks.close()
})
const token = mintToken(jwks.privateKey, { workspaceId: WORKSPACE })
// The server has asked for the key and is parked on the reply, so it cannot
// authenticate anyone until this test lets it.
await waitFor(() => jwks.requests >= 1, { message: 'the server to request the JWKS' })
assert.ok(!server.output.includes(KEY_LOADED))
// Cold client: sends sync step 1 and an update the instant the socket opens,
// i.e. while the server is still awaiting the JWKS response.
let framesWritten = 0
const cold = openClient(`${server.url}/${DOC_PATH}?token=${token}`, {
onOpen: (ws) => {
ws.send(syncStep1Message(), () => framesWritten++)
ws.send(updateMessage('hello'), () => framesWritten++)
}
})
await waitFor(() => framesWritten === 2, { message: 'the cold client to write both frames' })
// Both frames are on the wire while the server still has no key, so they can
// only reach it inside the pre-auth window.
assert.ok(!server.output.includes(KEY_LOADED))
await delay(FLIGHT_MARGIN_MS)
jwks.release()
await waitFor(() => hasKind(cold, syncStep2), {
message: 'the server to answer the cold client with sync step 2'
})
// The echo of our own update proves it was applied to the server-side doc.
await waitFor(() => hasKind(cold, 2), {
message: 'the server to broadcast back the update sent during the cold window'
})
cold.ws.close()
// A second client must now see the update the cold client sent.
const warm = openClient(`${server.url}/${DOC_PATH}?token=${token}`, {
onOpen: (ws) => ws.send(syncStep1Message())
})
t.after(() => warm.ws.close())
await waitFor(() => hasKind(warm, syncStep2), {
message: 'the server to answer the second client with sync step 2'
})
const doc = new Y.Doc()
const step2 = warm.received.find((data) => messageKind(data).syncType === syncStep2)
const decoder = decoding.createDecoder(step2)
decoding.readVarUint(decoder) // messageSync
syncProtocol.readSyncMessage(decoder, encoding.createEncoder(), doc, null)
assert.equal(doc.getText('content').toString(), 'hello')
})
test('cold start: a forged token is rejected with 4403 and its messages are dropped', { timeout: 60000 }, async (t) => {
const jwks = await startJwksServer({ hold: true })
const server = await startMultiplayerServer({ WINDMILL_BASE_URL: jwks.baseUrl })
t.after(async () => {
await server.close()
await jwks.close()
})
const { privateKey: otherKey } = crypto.generateKeyPairSync('ed25519')
const forged = mintToken(otherKey, { workspaceId: WORKSPACE })
let framesWritten = 0
const client = openClient(`${server.url}/${DOC_PATH}?token=${forged}`, {
onOpen: (ws) => ws.send(syncStep1Message(), () => framesWritten++)
})
// Buffer the frame first, then let verification run and reject.
await waitFor(() => framesWritten === 1, { message: 'the forged client to write its frame' })
await delay(FLIGHT_MARGIN_MS)
jwks.release()
await waitFor(() => client.closeCode !== undefined, { message: 'the forged connection to be closed' })
assert.equal(client.closeCode, 4403)
assert.deepEqual(client.received, [])
})
test('cold start: a peer that floods before authenticating is closed and the server keeps serving', { timeout: 60000 }, async (t) => {
const jwks = await startJwksServer({ hold: true })
const server = await startMultiplayerServer({ WINDMILL_BASE_URL: jwks.baseUrl })
t.after(async () => {
await server.close()
await jwks.close()
})
const token = mintToken(jwks.privateKey, { workspaceId: WORKSPACE })
// Four 600 KiB frames sent while the key is still parked: the second one takes
// the buffer past MAX_PREAUTH_BYTES (1 MiB). The token is valid and the key is
// not released until after the close, so only the cap can close this socket.
const flooder = openClient(`${server.url}/${DOC_PATH}?token=${token}`, {
onOpen: (ws) => {
for (let i = 0; i < 4; i++) ws.send(Buffer.alloc(600 * 1024))
}
})
await waitFor(() => flooder.closeCode !== undefined, { message: 'the flooding connection to be closed' })
assert.equal(flooder.closeCode, 1009)
assert.deepEqual(flooder.received, [])
assert.ok(!server.output.includes(KEY_LOADED))
// The server is unharmed and still syncs a well-behaved client.
jwks.release()
const healthy = openClient(`${server.url}/${DOC_PATH}?token=${token}`, {
onOpen: (ws) => ws.send(syncStep1Message())
})
t.after(() => healthy.ws.close())
await waitFor(() => hasKind(healthy, syncStep2), {
message: 'the server to still answer a well-behaved client with sync step 2'
})
})
test('a connection without a token is rejected with 4401', { timeout: 60000 }, async (t) => {
const jwks = await startJwksServer()
const server = await startMultiplayerServer({ WINDMILL_BASE_URL: jwks.baseUrl })
t.after(async () => {
await server.close()
await jwks.close()
})
const client = openClient(`${server.url}/${DOC_PATH}`, {
onOpen: (ws) => ws.send(syncStep1Message())
})
await waitFor(() => client.closeCode !== undefined, { message: 'the unauthenticated connection to be closed' })
assert.equal(client.closeCode, 4401)
assert.deepEqual(client.received, [])
})
test('the public key is fetched at startup, before any client connects', { timeout: 60000 }, async (t) => {
const jwks = await startJwksServer()
const server = await startMultiplayerServer({ WINDMILL_BASE_URL: jwks.baseUrl })
t.after(async () => {
await server.close()
await jwks.close()
})
await waitFor(() => jwks.requests >= 1, { message: 'the server to prefetch the JWKS at startup' })
})