mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-03 16:02:12 +00:00
* 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>
* fix(multiplayer): survive a malformed frame instead of exiting the process
`setupWSConnection`'s message handler decoded whatever an authenticated peer
put on the wire with no guard: `decoding.readVarUint`,
`syncProtocol.readSyncMessage` and `awarenessProtocol.applyAwarenessUpdate` all
throw on input they cannot parse, `ws` re-emits a listener's exception on the
process, and server.mjs installs no `uncaughtException` handler. One bad frame
from one client therefore killed the whole multiplayer server, taking every
other document and every other client with it.
Catch decode/apply failures, log the document, the client address and the error
message (never the payload), and close only the offending connection with 1007
"invalid frame payload data". Frames that arrive once a connection is no longer
OPEN are ignored, so the replay of the pre-auth buffer stops at the first
refusal instead of applying the rest.
docker/entrypoint-extra.sh made that outage permanent: on a service exit it
logged a bare PID and then `wait`ed on the rest, so the container stayed up with
a dead service and the health checks in front of it — which probe the LSP — saw
nothing wrong. It now names the service that died, stops the others through the
same shutdown path SIGTERM uses, and exits non-zero so the orchestrator replaces
the container. The "no services enabled" branch still sleeps.
Tests: multiplayer/test/malformed_frame.test.mjs covers four malformed payloads
from an authenticated client and one replayed out of the pre-auth buffer,
asserting the 1007 close, a live server process, an undisturbed bystander and a
real edit still propagating. docker/test_entrypoint_extra.sh runs the real
entrypoint in a container with stub services.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* test(multiplayer): prove the replayed malformed frame really is buffered pre-auth
Assert the server has not yet logged the loaded key when the frame is written,
and give it the same in-flight margin as the cold-start tests before releasing
the JWKS response, so the frame provably goes through the replay path rather
than landing on an already-authenticated connection.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* ci(extra): run the entrypoint supervision tests in publish_extra
The multiplayer unit tests already run there; the entrypoint test needs only
docker and the checkout, so run it in the same job, before the image build.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* fix(multiplayer): handle WebSocket protocol errors and bound the shutdown
Two crash paths of the same class as the malformed-frame one, from review.
`ws` fails a frame it cannot parse at the protocol level — an unmasked frame
from a client, a reserved opcode, a bad RSV bit — inside its Receiver, before
the application 'message' handler ever sees it, and `receiverOnError` ends with
`websocket.emit('error', err)`. With no 'error' listener that is an unhandled
EventEmitter error, so it exited the process just as a malformed payload did.
(A raw socket error such as ECONNRESET does not: ws 8.21.3's `socketOnError`
swallows those.) Add the listener on the accepted socket, before authentication
so the pre-auth window is covered too, and one on the server.
Log messages now go through `describeError`, which collapses whitespace and
truncates, so nothing that reaches an error message can forge or flood a log
line.
`stop_services` ended in a bare `wait`. On the `docker stop` path dockerd
provides the deadline; the "a service died" path signals itself, so a service
that is wedged or slow to honour SIGTERM would hold the container open
indefinitely — the state that path exists to prevent. Bound it: SIGTERM, wait
SHUTDOWN_GRACE_SECS (10 by default), then SIGKILL the stragglers by name.
Tests: multiplayer/test/socket_error.test.mjs (authenticated and pre-auth
illegal frames, asserting a live process and continued service), a
SIGTERM-ignoring stub scenario in docker/test_entrypoint_extra.sh, and that
harness is now bounded throughout — `timeout -k` on foreground runs, a watchdog
around the backgrounded ones, and an optional outer timeout on `docker run`.
`--entrypoint bash` so the documented windmill-extra:test override runs the
harness instead of the image's real entrypoint. The stubs publish a readiness
marker and the dying one waits for them, removing a startup race in the harness.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* fix(multiplayer): never log peer bytes, and keep a fatal server error fatal
Two review findings on the previous commit, both mine to answer for.
`describeError` collapsed whitespace, which is not enough. An error message is
not always a fixed string: `applyAwarenessUpdate` runs `JSON.parse` on the
peer's bytes and V8 quotes ~30 bytes of the offending input back verbatim, ESC
included, so a peer could put terminal escapes and forged content into a log
line. Strip everything outside printable ASCII instead, and say so where the
comment previously claimed the messages were fixed strings.
`wss.on('error')` was worse than the crash it replaced for one case: `ws`
forwards the HTTP server's errors there, so a failed listen (EADDRINUSE) was
logged and the process then exited 0 — a clean shutdown as far as anything
upstream could tell. It now sets a non-zero exit code. Setting `process.exitCode`
rather than calling `process.exit()` keeps the log line from being truncated.
`openClient` in the test helpers now records the socket error it was already
swallowing, so a failed connection reports its cause instead of surfacing as a
bare `waitFor` timeout.
Tests: a malformed awareness frame whose state is `x\x1b[2J OWNED THE LOG` added
to the payload table, with every case now asserting exactly one refusal line and
no control characters in it (1 fail before, 0 after, 3 runs); and a server that
cannot listen must exit non-zero (1 fail before, 0 after, 3 runs). 14/14 on 5
consecutive runs.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* test(multiplayer): make the exit-status helper robust to spawn and stdio races
Review follow-ups on runMultiplayerServerUntilExit.
Wait for 'close', not 'exit': 'exit' fires when the child terminates, which can
be before its stdio pipes are drained, and the caller reads the output. On the
EADDRINUSE path the child writes one line and exits immediately after, which is
exactly the shape that loses it.
Listen for 'error' too. A child that fails to spawn emits neither 'exit' nor
'close', so the promise would never settle and the SIGKILL guard could not help.
Report whether the guard fired, rather than leaving the caller to infer it from
the exit signal: `signal` is null for every child exit on Windows, so a server
that hung after the listen error would have looked like one that exited on its
own. The test asserts on that flag instead.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* test(multiplayer): assert the refusal line itself, and name hasSyncType for what it takes
Two review nits on the test helpers.
The `doc="..."` assertion searched the whole server log, where CONNECT and
DISCONNECT also name the document, so it would have passed even if the refusal
stopped naming anything. Every assertion about the refusal is now made against
the refusal line, which the test already isolates, and it also checks the peer
is named.
`hasKind` took a sync sub-type but was named as if it took any message kind, and
the two families overlap numerically (`syncStep1 === messageSync === 0`), so a
caller passing the wrong one got a silently wrong answer. No runtime check can
tell aliased numbers apart, so the fix is the name: `hasSyncType`, with the
overlap spelled out where the constants are declared.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* test(multiplayer): make the pre-auth tests prove which path they took
The replay test could not establish that its frame went through the pre-auth
buffer: a frame delivered after setup, into the live handler, produces the same
close code and the same refusal line, so if the in-flight margin were ever
missed the test would quietly become a duplicate of the main-loop cases rather
than fail. server.mjs now logs REPLAY when, and only when, it replays a buffered
pre-auth message — worth having on its own, since that path only runs when a
client beat the JWKS fetch on a slow-starting instance — and the test asserts on
it. Removing that log line turns the test red, which is the point.
The socket-error pre-auth test gated on `jwks.requests >= 1`, which the startup
warm-up already satisfies, so it proved nothing about the offender. What makes
it the pre-auth case is that the JWKS response stays parked for the whole test;
it now asserts the server never logged CONNECT, which is exact.
`killedByTimeout` was set before the kill, so a child that exited on its own just
before the timeout — with 'close' still pending on the stdio drain, the very
window this helper waits for — would have been reported as killed. It now claims
the rescue only when there was a live process to signal.
14/14 on eight consecutive runs.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* test(multiplayer): document the frame-recording contract in openClient
ws hands every frame over as a Buffer under the default binaryType, text frames
included, so recording them as Uint8Array is lossless for both. Worth stating:
ws 7 delivered text frames as strings, where new Uint8Array(string) would have
been a silent zero-fill, and the difference is not visible at the call site.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* test(multiplayer): pin the close code for an illegal frame, and trim comments to the 4-line rule
Both socket-error tests waited for a close and never checked what it was, so an
abrupt 1006 teardown would have passed while the comment beside the payload
claimed 1002. `ws` sends 1002 for an unmasked frame in both the authenticated
and pre-auth cases, confirmed over repeated runs; that is now a named constant
asserted in each test, mirroring malformed_frame.test.mjs. Changing the expected
value turns both red.
The comments added by this branch also ran past the four lines AGENTS.md allows,
and several justified the change to a reader rather than stating the invariant.
Condensed to the invariant, at the site that would break it.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
428 lines
15 KiB
JavaScript
428 lines
15 KiB
JavaScript
#!/usr/bin/env node
|
|
/**
|
|
* Simple y-websocket server with connection logging and JWT authentication.
|
|
* Run with: node server.mjs
|
|
*/
|
|
|
|
import http from 'http'
|
|
import crypto from 'node:crypto'
|
|
import { WebSocketServer } from 'ws'
|
|
import * as Y from 'yjs'
|
|
import * as syncProtocol from 'y-protocols/sync'
|
|
import * as awarenessProtocol from 'y-protocols/awareness'
|
|
import * as encoding from 'lib0/encoding'
|
|
import * as decoding from 'lib0/decoding'
|
|
|
|
const PORT = process.env.PORT || 3002
|
|
const HOST = process.env.HOST || '0.0.0.0'
|
|
const WINDMILL_BASE_URL = process.env.WINDMILL_BASE_URL || process.env.BASE_INTERNAL_URL
|
|
const REQUIRE_SIGNED_REQUESTS = process.env.REQUIRE_SIGNED_MULTIPLAYER_REQUESTS !== 'false'
|
|
|
|
const messageSync = 0
|
|
const messageAwareness = 1
|
|
|
|
// Caps on what an unauthenticated peer may buffer while its token is verified.
|
|
// A real client only has sync step 1 and its first awareness update in flight
|
|
// there, so both limits are far above legitimate use.
|
|
const MAX_PREAUTH_MESSAGES = 32
|
|
const MAX_PREAUTH_BYTES = 1024 * 1024
|
|
|
|
// --- JWT verification ---
|
|
|
|
let cachedPublicKey = null
|
|
let publicKeyFetchPromise = null
|
|
|
|
function base64urlDecode(str) {
|
|
const padding = '='.repeat((4 - str.length % 4) % 4)
|
|
const base64 = str.replace(/-/g, '+').replace(/_/g, '/') + padding
|
|
const binary = atob(base64)
|
|
return Uint8Array.from(binary, c => c.charCodeAt(0))
|
|
}
|
|
|
|
async function getPublicKey() {
|
|
if (cachedPublicKey) return cachedPublicKey
|
|
if (publicKeyFetchPromise) return publicKeyFetchPromise
|
|
|
|
if (!WINDMILL_BASE_URL) {
|
|
console.warn(`[${new Date().toISOString()}] WINDMILL_BASE_URL not set - cannot fetch public key`)
|
|
return null
|
|
}
|
|
|
|
publicKeyFetchPromise = (async () => {
|
|
try {
|
|
const jwksUrl = `${WINDMILL_BASE_URL.replace(/\/$/, '')}/api/debug/jwks`
|
|
console.log(`[${new Date().toISOString()}] Fetching JWKS from ${jwksUrl}`)
|
|
const response = await fetch(jwksUrl)
|
|
if (!response.ok) {
|
|
throw new Error(`Failed to fetch JWKS: ${response.status} ${response.statusText}`)
|
|
}
|
|
|
|
const jwks = await response.json()
|
|
if (!jwks.keys || jwks.keys.length === 0) {
|
|
throw new Error('No keys in JWKS')
|
|
}
|
|
|
|
const jwk = jwks.keys[0]
|
|
if (jwk.kty !== 'OKP' || jwk.crv !== 'Ed25519') {
|
|
throw new Error(`Unsupported key type: ${jwk.kty}/${jwk.crv}`)
|
|
}
|
|
|
|
const publicKeyBytes = base64urlDecode(jwk.x)
|
|
const key = await crypto.subtle.importKey(
|
|
'raw',
|
|
publicKeyBytes,
|
|
{ name: 'Ed25519' },
|
|
true,
|
|
['verify']
|
|
)
|
|
|
|
cachedPublicKey = key
|
|
console.log(`[${new Date().toISOString()}] Successfully loaded Ed25519 public key from JWKS`)
|
|
return key
|
|
} catch (error) {
|
|
console.error(`[${new Date().toISOString()}] Failed to fetch/parse JWKS: ${error}`)
|
|
return null
|
|
} finally {
|
|
publicKeyFetchPromise = null
|
|
}
|
|
})()
|
|
|
|
return publicKeyFetchPromise
|
|
}
|
|
|
|
/**
|
|
* Verify a JWT multiplayer token.
|
|
* Returns null if valid, or an error message if invalid.
|
|
*/
|
|
async function verifyToken(token, docName) {
|
|
const publicKey = await getPublicKey()
|
|
|
|
if (!publicKey) {
|
|
if (REQUIRE_SIGNED_REQUESTS) {
|
|
return 'Public key not available but signed requests are required. Set WINDMILL_BASE_URL.'
|
|
}
|
|
console.warn(`[${new Date().toISOString()}] Public key not available - signature verification disabled`)
|
|
return null
|
|
}
|
|
|
|
const parts = token.split('.')
|
|
if (parts.length !== 3) {
|
|
return 'Invalid JWT format'
|
|
}
|
|
|
|
const [headerB64, claimsB64, signatureB64] = parts
|
|
|
|
try {
|
|
const message = new TextEncoder().encode(`${headerB64}.${claimsB64}`)
|
|
const signature = base64urlDecode(signatureB64)
|
|
|
|
const isValid = await crypto.subtle.verify(
|
|
{ name: 'Ed25519' },
|
|
publicKey,
|
|
signature,
|
|
message
|
|
)
|
|
|
|
if (!isValid) {
|
|
return 'Invalid JWT signature'
|
|
}
|
|
|
|
const claimsJson = new TextDecoder().decode(base64urlDecode(claimsB64))
|
|
const claims = JSON.parse(claimsJson)
|
|
|
|
// Check expiration
|
|
const now = Math.floor(Date.now() / 1000)
|
|
if (now > claims.exp) {
|
|
return `Token expired: ${now - claims.exp} seconds ago`
|
|
}
|
|
|
|
// Check purpose
|
|
if (claims.purpose !== 'multiplayer') {
|
|
return `Invalid token purpose: ${claims.purpose}`
|
|
}
|
|
|
|
// Check workspace matches the doc room (format: "{workspace}/{path}")
|
|
const docWorkspace = docName.split('/')[0]
|
|
if (docWorkspace && claims.workspace_id !== docWorkspace) {
|
|
return `Token workspace "${claims.workspace_id}" does not match room workspace "${docWorkspace}"`
|
|
}
|
|
|
|
return null
|
|
} catch (error) {
|
|
return `JWT verification error: ${error}`
|
|
}
|
|
}
|
|
|
|
// --- Y.js document management ---
|
|
|
|
// Store docs in memory
|
|
const docs = new Map()
|
|
|
|
const getYDoc = (docname) => {
|
|
let doc = docs.get(docname)
|
|
if (!doc) {
|
|
doc = new Y.Doc()
|
|
doc.name = docname
|
|
docs.set(docname, doc)
|
|
}
|
|
return doc
|
|
}
|
|
|
|
/**
|
|
* Every error message is peer-controlled: `applyAwarenessUpdate` runs
|
|
* `JSON.parse` on the peer's bytes and V8 quotes the offending input back. Only
|
|
* printable ASCII survives, so nothing in it can forge or flood a log line.
|
|
*/
|
|
const describeError = (error) => String(error?.message ?? error).replace(/[^\x20-\x7e]+/g, ' ').slice(0, 200)
|
|
|
|
const send = (conn, message) => {
|
|
if (conn.readyState === 1) { // WebSocket.OPEN
|
|
conn.send(message, err => { if (err) console.error(err) })
|
|
}
|
|
}
|
|
|
|
const setupWSConnection = (conn, req, docName, bufferedMessages = []) => {
|
|
const doc = getYDoc(docName)
|
|
const clientIp = req.socket.remoteAddress
|
|
|
|
// Initialize awareness
|
|
if (!doc.awareness) {
|
|
doc.awareness = new awarenessProtocol.Awareness(doc)
|
|
}
|
|
|
|
const awareness = doc.awareness
|
|
|
|
// Track connections per doc
|
|
if (!doc.conns) doc.conns = new Set()
|
|
doc.conns.add(conn)
|
|
|
|
const messageHandler = (message) => {
|
|
// A frame that arrived before we decided to close this connection (or that a
|
|
// replay below already made moot) must not still be applied to the document.
|
|
if (conn.readyState !== 1) return // WebSocket.OPEN
|
|
|
|
try {
|
|
const data = new Uint8Array(message)
|
|
const decoder = decoding.createDecoder(data)
|
|
const messageType = decoding.readVarUint(decoder)
|
|
|
|
switch (messageType) {
|
|
case messageSync: {
|
|
const encoder = encoding.createEncoder()
|
|
encoding.writeVarUint(encoder, messageSync)
|
|
syncProtocol.readSyncMessage(decoder, encoder, doc, null)
|
|
if (encoding.length(encoder) > 1) {
|
|
send(conn, encoding.toUint8Array(encoder))
|
|
}
|
|
break
|
|
}
|
|
case messageAwareness:
|
|
awarenessProtocol.applyAwarenessUpdate(awareness, decoding.readVarUint8Array(decoder), conn)
|
|
break
|
|
}
|
|
} catch (error) {
|
|
// The decoders throw on unparseable input and `ws` re-emits a listener's
|
|
// exception on the process, so one bad frame from one peer would end the
|
|
// server for every document and client. Drop only the offender, with RFC
|
|
// 6455's 1007; the frame itself is untrusted and is never logged.
|
|
console.warn(`[${new Date().toISOString()}] MALFORMED MESSAGE: doc="${docName}" from=${clientIp} error="${describeError(error)}"`)
|
|
conn.close(1007, 'Invalid message')
|
|
}
|
|
}
|
|
conn.on('message', messageHandler)
|
|
|
|
// Send initial sync step 1
|
|
{
|
|
const encoder = encoding.createEncoder()
|
|
encoding.writeVarUint(encoder, messageSync)
|
|
syncProtocol.writeSyncStep1(encoder, doc)
|
|
send(conn, encoding.toUint8Array(encoder))
|
|
}
|
|
|
|
// Send awareness states
|
|
const awarenessStates = awareness.getStates()
|
|
if (awarenessStates.size > 0) {
|
|
const encoder = encoding.createEncoder()
|
|
encoding.writeVarUint(encoder, messageAwareness)
|
|
encoding.writeVarUint8Array(encoder, awarenessProtocol.encodeAwarenessUpdate(awareness, Array.from(awarenessStates.keys())))
|
|
send(conn, encoding.toUint8Array(encoder))
|
|
}
|
|
|
|
// Broadcast awareness changes
|
|
const awarenessChangeHandler = ({ added, updated, removed }, origin) => {
|
|
const changedClients = added.concat(updated).concat(removed)
|
|
const encoder = encoding.createEncoder()
|
|
encoding.writeVarUint(encoder, messageAwareness)
|
|
encoding.writeVarUint8Array(encoder, awarenessProtocol.encodeAwarenessUpdate(awareness, changedClients))
|
|
const message = encoding.toUint8Array(encoder)
|
|
doc.conns.forEach(c => send(c, message))
|
|
}
|
|
awareness.on('update', awarenessChangeHandler)
|
|
|
|
// Broadcast doc updates
|
|
const updateHandler = (update, origin) => {
|
|
const encoder = encoding.createEncoder()
|
|
encoding.writeVarUint(encoder, messageSync)
|
|
syncProtocol.writeUpdate(encoder, update)
|
|
const message = encoding.toUint8Array(encoder)
|
|
doc.conns.forEach(c => {
|
|
if (c !== origin) send(c, message)
|
|
})
|
|
}
|
|
doc.on('update', updateHandler)
|
|
|
|
conn.on('close', () => {
|
|
doc.conns.delete(conn)
|
|
awareness.off('update', awarenessChangeHandler)
|
|
doc.off('update', updateHandler)
|
|
|
|
// Clean up awareness for this connection
|
|
awarenessProtocol.removeAwarenessStates(awareness, [doc.clientID], null)
|
|
})
|
|
|
|
// Replay, in order, the messages that arrived while the token was being
|
|
// verified, now that the handlers above are in place. The log line is the only
|
|
// outside evidence that a frame took this path rather than the live handler,
|
|
// and marks a client that beat the JWKS fetch on a slow-starting instance.
|
|
if (bufferedMessages.length > 0) {
|
|
console.log(`[${new Date().toISOString()}] REPLAY: doc="${docName}" from=${clientIp} messages=${bufferedMessages.length}`)
|
|
}
|
|
bufferedMessages.forEach(messageHandler)
|
|
}
|
|
|
|
// --- HTTP + WebSocket server ---
|
|
|
|
const server = http.createServer((req, res) => {
|
|
// Strip /ws_mp/ prefix if present (when accessed without reverse proxy path stripping)
|
|
if (req.url?.startsWith('/ws_mp/')) {
|
|
req.url = req.url.slice('/ws_mp'.length)
|
|
} else if (req.url === '/ws_mp') {
|
|
req.url = '/'
|
|
}
|
|
console.log(`[${new Date().toISOString()}] HTTP ${req.method} ${req.url} from=${req.socket.remoteAddress}`)
|
|
if (req.url === '/' || req.url === '/health') {
|
|
res.writeHead(200, {
|
|
'Content-Type': 'application/json',
|
|
'Access-Control-Allow-Origin': '*'
|
|
})
|
|
res.end(JSON.stringify({ status: 'ok', service: 'multiplayer' }))
|
|
} else {
|
|
res.writeHead(404)
|
|
res.end('not found')
|
|
}
|
|
})
|
|
|
|
const wss = new WebSocketServer({ server })
|
|
|
|
// `ws` forwards the HTTP server's errors here, and those are fatal: a failed
|
|
// listen leaves nothing to serve, and exiting 0 would read as a clean shutdown.
|
|
// `process.exitCode` rather than `process.exit()`, which truncates this line.
|
|
wss.on('error', (error) => {
|
|
console.error(`[${new Date().toISOString()}] WEBSOCKET SERVER ERROR: ${describeError(error)}`)
|
|
process.exitCode = 1
|
|
})
|
|
|
|
wss.on('connection', async (ws, req) => {
|
|
let docName = req.url?.slice(1).split('?')[0] || 'unknown'
|
|
|
|
// Strip ws_mp/ prefix if present (when accessed without reverse proxy path stripping)
|
|
if (docName.startsWith('ws_mp/')) {
|
|
docName = docName.slice('ws_mp/'.length)
|
|
}
|
|
|
|
const clientIp = req.socket.remoteAddress
|
|
|
|
// A frame `ws` cannot parse at the protocol level fails in its Receiver, never
|
|
// reaching the handler below, and an unhandled 'error' on an EventEmitter ends
|
|
// the process. Attached before authentication, since the pre-auth window is
|
|
// exposed too; `ws` has already closed the connection by the time this runs.
|
|
ws.on('error', (error) => {
|
|
console.warn(`[${new Date().toISOString()}] SOCKET ERROR: doc="${docName}" from=${clientIp} error="${describeError(error)}"`)
|
|
})
|
|
|
|
// Handle ping test — respond and close immediately
|
|
if (docName === '__ping__') {
|
|
console.log(`[${new Date().toISOString()}] WS ping from=${clientIp}`)
|
|
ws.send(JSON.stringify({ type: 'pong', service: 'multiplayer' }))
|
|
ws.close()
|
|
return
|
|
}
|
|
|
|
// A y-websocket client sends sync step 1 as soon as the socket opens, which can
|
|
// be before token verification resolves (the first connection after startup has
|
|
// to wait for the JWKS fetch). `ws` drops messages emitted with no listener
|
|
// attached, so buffer them here and replay them once the connection is accepted.
|
|
// The peer is not authenticated yet and the JWKS fetch has no timeout, so what
|
|
// it may buffer is capped.
|
|
const bufferedMessages = []
|
|
let bufferedBytes = 0
|
|
const bufferMessage = (message) => {
|
|
bufferedBytes += message.length
|
|
if (bufferedMessages.length >= MAX_PREAUTH_MESSAGES || bufferedBytes > MAX_PREAUTH_BYTES) {
|
|
console.warn(`[${new Date().toISOString()}] REJECTED: doc="${docName}" from=${clientIp} reason="too much data before authentication"`)
|
|
bufferedMessages.length = 0
|
|
ws.off('message', bufferMessage)
|
|
ws.close(1009, 'Too much data before authentication')
|
|
return
|
|
}
|
|
bufferedMessages.push(message)
|
|
}
|
|
ws.on('message', bufferMessage)
|
|
|
|
// Verify JWT token
|
|
const urlParams = new URLSearchParams(req.url?.split('?')[1] || '')
|
|
const token = urlParams.get('token')
|
|
|
|
if (!token) {
|
|
if (REQUIRE_SIGNED_REQUESTS) {
|
|
console.warn(`[${new Date().toISOString()}] REJECTED: doc="${docName}" from=${clientIp} reason="no token"`)
|
|
ws.off('message', bufferMessage)
|
|
ws.close(4401, 'Authentication required')
|
|
return
|
|
}
|
|
console.warn(`[${new Date().toISOString()}] WARN: no token for doc="${docName}" from=${clientIp} (signed requests not required)`)
|
|
} else {
|
|
const error = await verifyToken(token, docName)
|
|
if (error) {
|
|
console.warn(`[${new Date().toISOString()}] REJECTED: doc="${docName}" from=${clientIp} reason="${error}"`)
|
|
ws.off('message', bufferMessage)
|
|
ws.close(4403, 'Token verification failed')
|
|
return
|
|
}
|
|
}
|
|
|
|
ws.off('message', bufferMessage)
|
|
|
|
// The socket may already be gone: closed by the peer while the token was being
|
|
// verified, or by the buffer cap above. Setting a doc connection up on it would
|
|
// add it to `doc.conns` with a 'close' listener that can no longer fire.
|
|
if (ws.readyState !== 1) { // WebSocket.OPEN
|
|
return
|
|
}
|
|
|
|
console.log(`[${new Date().toISOString()}] CONNECT: doc="${docName}" from=${clientIp}`)
|
|
|
|
ws.on('close', () => {
|
|
console.log(`[${new Date().toISOString()}] DISCONNECT: doc="${docName}" from=${clientIp}`)
|
|
})
|
|
|
|
setupWSConnection(ws, req, docName, bufferedMessages)
|
|
})
|
|
|
|
server.listen(PORT, HOST, () => {
|
|
console.log(`[${new Date().toISOString()}] Multiplayer server running at ${HOST}:${PORT}`)
|
|
if (REQUIRE_SIGNED_REQUESTS) {
|
|
console.log(`[${new Date().toISOString()}] Signed requests REQUIRED (set REQUIRE_SIGNED_MULTIPLAYER_REQUESTS=false to disable)`)
|
|
} else {
|
|
console.log(`[${new Date().toISOString()}] Signed requests DISABLED`)
|
|
}
|
|
|
|
// Warm the JWKS cache so the first connection does not have to wait for it.
|
|
// Best-effort insurance only: getPublicKey() stays the lazy fallback, logs its
|
|
// own failures and resolves to null rather than rejecting, and a connection
|
|
// arriving before this resolves is handled by the message buffer.
|
|
if (WINDMILL_BASE_URL) {
|
|
getPublicKey()
|
|
}
|
|
})
|