A listening party's room: who is in it, and what they must all hear the same. The Durable Object in party.ts is a thin shell around this: it owns the WebSockets (the Hibernation API) and hands each event here, so the rules are plain code the tests run in Node.
One host and up to MAX_PEOPLE - 1 guests. The host's page is the only one
that asks Jev (website/src/jev/party.mjs): it announces each performance
(run: which song), when it started (clock, in the room's clock), and
every decision its jev()s settle (decision). The room keeps the current
performance so a guest who joins mid-song gets it all in one welcome,
and passes each message on to everyone else. Reactions (🔥/😴) from anyone
are counted here and the totals broadcast, so every page, the host's
included (whose Jev hears them), shows the same counts.
Live coordination only: nothing here outlives the room. The room's storage holds the current performance, and is deleted when the last person leaves (or found stale when someone arrives to an empty room).
22export const PARTY_PREFIX = '/jev/party/';
GET /jev/party/<room>: a WebSocket into that room's Durable Object, one
per room (idFromName). The per-visitor limit on joining is PARTY_LIMIT
(wrangler.json): a page joins a party once, and reconnects only when its
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';
A token bucket per connection: BURST messages at once, refilled at RATE a second. A section's decisions, a reaction or a clock ping are a message each; a song asks at most a few times a section, and a listener taps 🔥 by hand.
Bounds on what a host can make the room store: songs have at most 16 sections and a few jev()s.
What a connection remembers across hibernation (serializeAttachment).
64export type Attachment = { role: Role; id: string };
The socket as the room sees it: the DO's hibernatable WebSocket, or a test's fake.
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 };
Every write the room makes is unconfirmed. By default a Durable Object holds its outgoing messages until its writes are on disk (the output gate), so each decision's write delayed its own broadcast: measured locally, decisions reached a guest up to 10 s after the host sent them, and the song played its fallbacks meanwhile. The room's storage only saves a late joiner from waiting and survives hibernation; a write lost in a crash costs nothing, since a crash drops every socket and the host's page re-sends the whole performance when it reconnects.
92const LIVE: WriteOptions = { allowUnconfirmed: true };
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')}`;
Whether a connection may join: null when it may, else the HTTP status and
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}
Per-connection rate limit, kept in memory: a room that hibernates forgets it, which only ever refills a bucket.
139type Bucket = { tokens: number; at: number; warned: boolean };
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 }
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 }
A connection closed (or errored). Empty, the room forgets the party.
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}