jevstrudel.git / worker / src / party-room.ts

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).

19import type { Env } from './env';
20import type { Live } from './party-directory';
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';
44export const MAX_PEOPLE = 50;
45export const MAX_MESSAGE_BYTES = 16 * 1024;

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.

50export const BURST = 30;
51export const RATE = 5;

Bounds on what a host can make the room store: songs have at most 16 sections and a few jev()s.

54export const MAX_SEGMENT = 63;
55export const MAX_VIEW = 7;
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';

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.

68export interface Peer {
69  send(message: string): void;
70  close(code?: number, reason?: string): void;
71  deserializeAttachment(): unknown;
72}

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.

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  }
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}