1// A listening party's room: who is in it, and what they must all hear the 2// same. The Durable Object in party.ts is a thin shell around this: it owns 3// the WebSockets (the Hibernation API) and hands each event here, so the 4// rules are plain code the tests run in Node. 5// 6// One host and up to MAX_PEOPLE - 1 guests. The host's page is the only one 7// that asks Jev (website/src/jev/party.mjs): it announces each performance 8// (`run`: which song), when it started (`clock`, in the room's clock), and 9// every decision its jev()s settle (`decision`). The room keeps the current 10// performance so a guest who joins mid-song gets it all in one `welcome`, 11// and passes each message on to everyone else. Reactions (🔥/😴) from anyone 12// are counted here and the totals broadcast, so every page, the host's 13// included (whose Jev hears them), shows the same counts. 14// 15// Live coordination only: nothing here outlives the room. The room's 16// storage holds the current performance, and is deleted when the last 17// person leaves (or found stale when someone arrives to an empty room). 18 19import type { Env } from './env'; 20import type { Live } from './party-directory'; 21 22export const PARTY_PREFIX = '/jev/party/'; 23 24// GET /jev/party/<room>: a WebSocket into that room's Durable Object, one 25// per room (`idFromName`). The per-visitor limit on joining is PARTY_LIMIT 26// (wrangler.json): a page joins a party once, and reconnects only when its 27// network drops, so 20 a minute is far above any listener's need. 28export async function party(request: Request, env: Pick<Env, 'PARTY' | 'PARTY_LIMIT'>): Promise<Response> { 29 const room = new URL(request.url).pathname.slice(PARTY_PREFIX.length); 30 if (!ROOM.test(room)) return new Response('not found', { status: 404 }); 31 if (request.headers.get('Upgrade') !== 'websocket') return new Response('WebSocket only', { status: 426 }); 32 const visitor = request.headers.get('CF-Connecting-IP') ?? 'local'; 33 if (!(await env.PARTY_LIMIT.limit({ key: visitor })).success) { 34 return new Response('too many joins; try again in a minute', { status: 429 }); 35 } 36 // the room's name, for a listed room to give to someone joining from the 37 // list (party.ts); set here, never taken from the page 38 const headers = new Headers(request.headers); 39 headers.set(ROOM_HEADER, room); 40 return env.PARTY.get(env.PARTY.idFromName(room)).fetch(new Request(request, { headers })); 41} 42export const ROOM_HEADER = 'X-Jev-Party-Room'; 43 44export const MAX_PEOPLE = 50; 45export const MAX_MESSAGE_BYTES = 16 * 1024; 46// A token bucket per connection: BURST messages at once, refilled at RATE a 47// second. A section's decisions, a reaction or a clock ping are a message 48// each; a song asks at most a few times a section, and a listener taps 🔥 by 49// hand. 50export const BURST = 30; 51export const RATE = 5; 52// Bounds on what a host can make the room store: songs have at most 16 53// sections and a few jev()s. 54export const MAX_SEGMENT = 63; 55export const MAX_VIEW = 7; 56 57export const ROOM = /^[A-Za-z0-9_-]{16,43}$/; 58const KEY = /^[A-Za-z0-9_-]{16,64}$/; 59const SONG = /^[a-z0-9-]{1,64}\/[a-z0-9-]{1,64}$/; 60 61export type Role = 'host' | 'guest'; 62 63// What a connection remembers across hibernation (serializeAttachment). 64export type Attachment = { role: Role; id: string }; 65 66// The socket as the room sees it: the DO's hibernatable WebSocket, or a 67// test's fake. 68export interface Peer { 69 send(message: string): void; 70 close(code?: number, reason?: string): void; 71 deserializeAttachment(): unknown; 72} 73 74// The Durable Object's key-value storage, the part the room uses. 75export interface Storage { 76 get<T>(key: string): Promise<T | undefined>; 77 put(entries: Record<string, unknown>, options?: WriteOptions): Promise<void>; 78 delete(keys: string[], options?: WriteOptions): Promise<number>; 79 list<T>(options: { prefix: string }): Promise<Map<string, T>>; 80 deleteAll(options?: WriteOptions): Promise<void>; 81} 82type WriteOptions = { allowUnconfirmed?: boolean }; 83 84// Every write the room makes is unconfirmed. By default a Durable Object 85// holds its outgoing messages until its writes are on disk (the output 86// gate), so each decision's write delayed its own broadcast: measured 87// locally, decisions reached a guest up to 10 s after the host sent them, 88// and the song played its fallbacks meanwhile. The room's storage only 89// saves a late joiner from waiting and survives hibernation; a write lost 90// in a crash costs nothing, since a crash drops every socket and the 91// host's page re-sends the whole performance when it reconnects. 92const LIVE: WriteOptions = { allowUnconfirmed: true }; 93 94export type Run = { run: number; song: string; code: string }; 95export type Clock = { startedAt: number; cps: number }; 96export type Counts = { fire: number; sleep: number }; 97 98type Message = Record<string, unknown> & { t?: unknown }; 99 100const isInt = (x: unknown, max: number): x is number => Number.isInteger(x) && (x as number) >= 0 && (x as number) <= max; 101const isObject = (x: unknown): x is Record<string, unknown> => typeof x === 'object' && x !== null && !Array.isArray(x); 102 103async function sha256(text: string): Promise<string> { 104 const digest = await crypto.subtle.digest('SHA-256', new TextEncoder().encode(text)); 105 return [...new Uint8Array(digest)].map((b) => b.toString(16).padStart(2, '0')).join(''); 106} 107 108const decisionKey = (view: number, segment: number) => 109 `d:${String(view).padStart(2, '0')}:${String(segment).padStart(3, '0')}`; 110const reactionKey = (segment: number) => `r:${String(segment).padStart(3, '0')}`; 111 112// Whether a connection may join: null when it may, else the HTTP status and 113// why. `people` is how many are connected now. 114export async function admit( 115 storage: Storage, 116 people: Peer[], 117 { role, key }: { role: string | null; key: string | null }, 118): Promise<{ status: number; why: string } | null> { 119 if (role !== 'host' && role !== 'guest') return { status: 400, why: 'role must be host or guest' }; 120 if (people.length >= MAX_PEOPLE) return { status: 503, why: `the party is full (${MAX_PEOPLE})` }; 121 if (people.length === 0) { 122 // An empty room has no party, whatever its storage says (a party whose 123 // last connection dropped while the room was evicted): start clean. 124 await storage.deleteAll(LIVE); 125 if (role === 'guest') return { status: 404, why: 'no party here: the host has left' }; 126 } 127 if (role === 'host') { 128 if (!key || !KEY.test(key)) return { status: 400, why: 'a host needs its key' }; 129 const stored = await storage.get<string>('host'); 130 // The first host claims the room; later ones (the host's reload) must 131 // bring the same key, so a guest cannot take the party over. 132 if (stored && stored !== (await sha256(key))) return { status: 403, why: 'this party has another host' }; 133 } 134 return null; 135} 136 137// Per-connection rate limit, kept in memory: a room that hibernates forgets 138// it, which only ever refills a bucket. 139type Bucket = { tokens: number; at: number; warned: boolean }; 140 141export class Room { 142 private buckets = new WeakMap<Peer, Bucket>(); 143 144 constructor( 145 private storage: Storage, 146 // every open connection, the one an event is for included 147 private peers: () => Peer[], 148 private now: () => number = Date.now, 149 // told whenever the room's people or song change, null once it is 150 // empty: the party directory's report (party-directory.ts) 151 private changed: (live: Live | null) => void = () => {}, 152 // whether a room name is this room's own (party.ts: the name's object id 153 // is this object's): how a host's "show me here", ticked mid-party, 154 // gives the room its name without the room trusting the page for it 155 private ownName: (name: string) => boolean = () => false, 156 ) {} 157 158 private async report(except?: Peer) { 159 const people = this.people(except); 160 if (!people.n) return this.changed(null); 161 const [run, listed] = await Promise.all([this.storage.get<Run>('run'), this.storage.get<boolean>('listed')]); 162 this.changed({ people: people.n, host: people.host, song: run?.song ?? null, listed: listed === true }); 163 } 164 165 private role(peer: Peer): Role | null { 166 const a = peer.deserializeAttachment() as Attachment | null; 167 return a?.role ?? null; 168 } 169 170 private others(except?: Peer) { 171 return this.peers().filter((p) => p !== except); 172 } 173 174 private broadcast(message: unknown, except?: Peer) { 175 const text = JSON.stringify(message); 176 for (const p of this.others(except)) { 177 try { 178 p.send(text); 179 } catch { 180 // closing already; its close event tidies up 181 } 182 } 183 } 184 185 private people(except?: Peer) { 186 const ps = this.others(except); 187 return { t: 'people', n: ps.length, host: ps.some((p) => this.role(p) === 'host') }; 188 } 189 190 // A connection has been accepted (after admit()). 191 async joined(peer: Peer, role: Role, key: string | null) { 192 if (role === 'host') { 193 // one host: a reload's new connection replaces the old one 194 for (const p of this.others(peer)) if (this.role(p) === 'host') p.close(4000, 'the host reconnected elsewhere'); 195 if (!(await this.storage.get('host'))) await this.storage.put({ host: await sha256(key!) }, LIVE); 196 } 197 const [run, clock, decisions, reactions] = await Promise.all([ 198 this.storage.get<Run>('run'), 199 this.storage.get<Clock>('clock'), 200 this.storage.list<unknown>({ prefix: 'd:' }), 201 this.storage.list<Counts>({ prefix: 'r:' }), 202 ]); 203 const people = this.people(); 204 peer.send( 205 JSON.stringify({ 206 t: 'welcome', 207 role, 208 people: people.n, 209 host: people.host, 210 now: this.now(), 211 run: run ?? null, 212 clock: clock ?? null, 213 decisions: [...decisions.values()], 214 reactions: Object.fromEntries([...reactions].map(([k, v]) => [Number(k.slice(2)), v])), 215 }), 216 ); 217 this.broadcast(people, peer); 218 await this.report(); 219 } 220 221 // A connection closed (or errored). Empty, the room forgets the party. 222 async left(peer: Peer) { 223 const rest = this.others(peer); 224 if (!rest.length) { 225 await this.storage.deleteAll(LIVE); 226 this.changed(null); 227 return; 228 } 229 this.broadcast(this.people(peer), peer); 230 await this.report(peer); 231 } 232 233 private allow(peer: Peer): boolean { 234 const now = this.now(); 235 const b = this.buckets.get(peer) ?? { tokens: BURST, at: now, warned: false }; 236 b.tokens = Math.min(BURST, b.tokens + ((now - b.at) / 1000) * RATE); 237 b.at = now; 238 this.buckets.set(peer, b); 239 if (b.tokens < 1) { 240 if (!b.warned) peer.send(JSON.stringify({ t: 'error', why: 'too many messages; some were dropped' })); 241 b.warned = true; 242 return false; 243 } 244 b.tokens -= 1; 245 b.warned = false; 246 return true; 247 } 248 249 async message(peer: Peer, data: string | ArrayBuffer) { 250 const size = typeof data === 'string' ? new TextEncoder().encode(data).length : data.byteLength; 251 if (typeof data !== 'string' || size > MAX_MESSAGE_BYTES) { 252 peer.send(JSON.stringify({ t: 'error', why: `messages are JSON text of at most ${MAX_MESSAGE_BYTES} bytes` })); 253 return; 254 } 255 if (!this.allow(peer)) return; 256 let m: Message; 257 try { 258 m = JSON.parse(data); 259 } catch { 260 return; 261 } 262 if (!isObject(m)) return; 263 const host = this.role(peer) === 'host'; 264 265 switch (m.t) { 266 case 'ping': 267 // the room's clock, for the page's offset estimate (NTP-style) 268 if (typeof m.c === 'number') peer.send(JSON.stringify({ t: 'pong', c: m.c, s: this.now() })); 269 return; 270 271 case 'react': { 272 if (!isInt(m.segment, MAX_SEGMENT) || (m.kind !== 'fire' && m.kind !== 'sleep')) return; 273 if (!(await this.storage.get<Run>('run'))) return; 274 const key = reactionKey(m.segment); 275 const counts = (await this.storage.get<Counts>(key)) ?? { fire: 0, sleep: 0 }; 276 counts[m.kind] += 1; 277 await this.storage.put({ [key]: counts }, LIVE); 278 // to everyone, the reacting page included: it shows the totals 279 this.broadcast({ t: 'reactions', segment: m.segment, ...counts }); 280 return; 281 } 282 283 case 'run': { 284 if (!host || typeof m.song !== 'string' || !SONG.test(m.song) || typeof m.code !== 'string') return; 285 if (m.code.length > 128) return; 286 const previous = await this.storage.get<Run>('run'); 287 const run: Run = { run: (previous?.run ?? 0) + 1, song: m.song, code: m.code }; 288 // a new performance: the last one's decisions, clock and reactions go 289 const stale = [ 290 ...(await this.storage.list({ prefix: 'd:' })).keys(), 291 ...(await this.storage.list({ prefix: 'r:' })).keys(), 292 'clock', 293 ]; 294 await this.storage.delete(stale, LIVE); 295 await this.storage.put({ run }, LIVE); 296 this.broadcast({ t: 'run', ...run }, peer); 297 peer.send(JSON.stringify({ t: 'run', ...run })); 298 await this.report(); 299 return; 300 } 301 302 case 'clock': { 303 if (!host || typeof m.startedAt !== 'number' || !Number.isFinite(m.startedAt)) return; 304 if (typeof m.cps !== 'number' || !(m.cps > 0 && m.cps < 100)) return; 305 const run = await this.storage.get<Run>('run'); 306 if (!run || m.run !== run.run) return; 307 const clock: Clock = { startedAt: m.startedAt, cps: m.cps }; 308 await this.storage.put({ clock }, LIVE); 309 this.broadcast({ t: 'clock', run: run.run, ...clock }, peer); 310 return; 311 } 312 313 case 'decision': { 314 if (!host || !isInt(m.view, MAX_VIEW) || !isInt(m.segment, MAX_SEGMENT) || !isObject(m.d)) return; 315 const run = await this.storage.get<Run>('run'); 316 if (!run || m.run !== run.run) return; 317 // Stored and passed on as the host sent it: each guest checks every 318 // value against its own song before it plays (jevCore's follow mode). 319 const decision = { run: run.run, view: m.view, segment: m.segment, d: m.d }; 320 await this.storage.put({ [decisionKey(m.view, m.segment)]: decision }, LIVE); 321 this.broadcast({ t: 'decision', ...decision }, peer); 322 return; 323 } 324 325 case 'listed': { 326 // The host ticked or unticked "show me here" mid-party: the party 327 // joins or leaves the live list at once. Listed, the room keeps its 328 // name, proven its own (`ownName`), to give someone joining from the 329 // list; unlisted, it keeps none. Guests already in stay in. 330 if (!host || typeof m.listed !== 'boolean') return; 331 if (m.listed) { 332 if (typeof m.room !== 'string' || !ROOM.test(m.room) || !this.ownName(m.room)) return; 333 await this.storage.put({ listed: true, name: m.room }, LIVE); 334 } else { 335 await this.storage.put({ listed: false }, LIVE); 336 await this.storage.delete(['name'], LIVE); 337 } 338 await this.report(); 339 return; 340 } 341 342 case 'end': { 343 if (!host) return; 344 const run = await this.storage.get<Run>('run'); 345 if (!run || m.run !== run.run) return; 346 await this.storage.delete(['clock'], LIVE); 347 this.broadcast({ t: 'end', run: run.run }, peer); 348 return; 349 } 350 } 351 } 352}