jevstrudel.git / worker / src / party-room.ts
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}