jevstrudel.git / worker / src / listening-store.ts

Listening's storage, in the site's D1 database (the DB binding; the schema is worker/migrations/0003_listening.sql). listening.ts validates everything before it arrives here.

A performance is written in one batch (one transaction): its row, then every segment's row from one JSON parameter through json_each, so a performance of any length is two statements (D1 allows 100 bound parameters a statement and, on the free plan, 50 queries a request).

9import {
10  MOVES_WINDOW,
11  type ListeningStore,
12  type Moves,
13  type Performance,
14  type Player,
15  type PlayerTake,
16  type Reaction,
17  type Recording,
18  type Segment,
19  type Stored,
20  type Take,
21  type Tally,
22} from './listening';
24type Row = {
25  jev: number;
26  segment: number;
27  section: string | null;
28  after_section: string | null;
29  status: Segment['status'];
30  ms: number | null;
31  from_cycle: number | null;
32  answers: string;
33};

The rows of a performance: one per segment per jev, with the form's section for each segment and the one before it.

37export function segmentRows(p: Performance): Row[] {
38  const sectionAt = (n: number) =>
39    p.form ? (p.jevs[p.form.jev].segments[n]?.values[p.form.question]?.choice ?? null) : null;
40  return p.jevs.flatMap(({ segments }, jev) =>
41    segments.map((s, segment) => ({
42      jev,
43      segment,
44      section: sectionAt(segment),
45      after_section: segment > 0 ? sectionAt(segment - 1) : null,
46      status: s.status,
47      ms: s.ms ?? null,
48      from_cycle: s.from ?? null,
49      answers: JSON.stringify(s.values),
50    })),
51  );
52}

A performance stored without a time (the tests, and nothing else) is dated noon UTC of its day, as 0006 dates the takes from before it.

56const noonOf = (day: string) => Date.parse(`${day}T12:00:00Z`);
58const playerOf = (r: Record<string, unknown>): Player =>
59  typeof r.player_id === 'string' ? { id: r.player_id, name: r.player_name as string } : null;

The form's sections in order for each of ids that has a form (their form jev's rows), at most 50 ids a call (D1 allows 100 bound parameters).

63async function pathsOf(db: D1Database, ids: string[]): Promise<Map<string, string[]>> {
64  const paths = new Map<string, string[]>(ids.map((id) => [id, []]));
65  if (!ids.length) return paths;
66  const { results } = await db
67    .prepare(
68      `SELECT s.performance, s.section FROM performance_segments s JOIN performances p ON p.id = s.performance
69       WHERE s.performance IN (${ids.map(() => '?').join(', ')}) AND s.jev = p.form_jev
70       ORDER BY s.performance, s.segment`,
71    )
72    .bind(...ids)
73    .all<{ performance: string; section: string | null }>();
74  for (const r of results) if (r.section !== null) paths.get(r.performance)!.push(r.section);
75  return paths;
76}

The columns a take is listed with, and the take from them.

79const TAKE_COLUMNS = `p.id, p.song, p.day, p.recorded_at, p.ended, p.code_hash, p.form_jev,
80  u.id AS player_id, u.display_name AS player_name,
81  (SELECT count(*) FROM performance_segments s
82    WHERE s.performance = p.id AND s.jev = coalesce(p.form_jev, 0)) AS segments,
83  (SELECT count(DISTINCT s.segment) FROM performance_segments s
84    WHERE s.performance = p.id AND s.status = 'fallback') AS fallbacks`;
85const takeOf = (r: Record<string, unknown>, paths: Map<string, string[]>): Take => ({
86  id: r.id as string,
87  day: r.day as string,
88  at: r.recorded_at as number,
89  player: playerOf(r),
90  ended: r.ended as Take['ended'],
91  codeHash: r.code_hash as string,
92  segments: r.segments as number,
93  fallbacks: r.fallbacks as number,
94  path: r.form_jev === null ? null : (paths.get(r.id as string) ?? []),
95});
97export function d1Listening(db: D1Database): ListeningStore {
98  return {
99    async react({ song, section, reaction }: Reaction): Promise<void> {
100      await db
101        .prepare(
102          'INSERT INTO reactions (song, section, reaction, n) VALUES (?, ?, ?, 1) ON CONFLICT (song, section, reaction) DO UPDATE SET n = n + 1',
103        )
104        .bind(song, section, reaction)
105        .run();
106    },
107
108    async tally(song: string): Promise<Tally[]> {
109      const { results } = await db
110        .prepare(
111          `SELECT section,
112             COALESCE(SUM(CASE reaction WHEN 'fire' THEN n END), 0) AS fire,
113             COALESCE(SUM(CASE reaction WHEN 'sleep' THEN n END), 0) AS sleep
114           FROM reactions WHERE song = ? GROUP BY section ORDER BY section`,
115        )
116        .bind(song)
117        .all<Tally>();
118      return results;
119    },
120
121    async record(id: string, day: string, p: Performance, by: Recording = { at: noonOf(day), player: null }): Promise<void> {
122      await db.batch([
123        db
124          .prepare(
125            `INSERT INTO performances (id, song, code_hash, model, every, ended, jevs, form_jev, form_question, day, recorded_at, user_id)
126             VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
127          )
128          .bind(
129            id,
130            p.song,
131            p.codeHash,
132            p.model,
133            p.every,
134            p.ended,
135            p.jevs.length,
136            p.form?.jev ?? null,
137            p.form?.question ?? null,
138            day,
139            by.at,
140            by.player,
141          ),
142        db
143          .prepare(
144            `INSERT INTO performance_segments
145               (performance, jev, segment, section, after_section, status, ms, from_cycle, answers)
146             SELECT ?, r.value ->> '$.jev', r.value ->> '$.segment', r.value ->> '$.section',
147               r.value ->> '$.after_section', r.value ->> '$.status', r.value ->> '$.ms',
148               r.value ->> '$.from_cycle', r.value ->> '$.answers'
149             FROM json_each(?) AS r`,
150          )
151          .bind(id, JSON.stringify(segmentRows(p))),
152      ]);
153    },
154
155    async performance(id: string): Promise<Stored | null> {
156      const p = await db
157        .prepare(
158          `SELECT p.id, p.song, p.code_hash, p.model, p.every, p.ended, p.jevs, p.form_jev, p.form_question, p.day,
159             p.recorded_at, u.id AS player_id, u.display_name AS player_name
160           FROM performances p LEFT JOIN users u ON u.id = p.user_id WHERE p.id = ?`,
161        )
162        .bind(id)
163        .first<{
164          id: string;
165          song: string;
166          code_hash: string;
167          model: string;
168          every: number;
169          ended: Performance['ended'];
170          jevs: number;
171          form_jev: number | null;
172          form_question: string | null;
173          day: string;
174          recorded_at: number;
175          player_id: string | null;
176          player_name: string | null;
177        }>();
178      if (!p) return null;
179      const { results } = await db
180        .prepare(
181          'SELECT jev, segment, status, ms, from_cycle, answers FROM performance_segments WHERE performance = ? ORDER BY jev, segment',
182        )
183        .bind(id)
184        .all<Pick<Row, 'jev' | 'segment' | 'status' | 'ms' | 'from_cycle' | 'answers'>>();
185      const jevs: { segments: Segment[] }[] = Array.from({ length: p.jevs }, () => ({ segments: [] }));
186      for (const r of results) {
187        jevs[r.jev].segments[r.segment] = {
188          status: r.status,
189          ...(r.ms !== null ? { ms: r.ms } : {}),
190          ...(r.from_cycle !== null ? { from: r.from_cycle } : {}),
191          values: JSON.parse(r.answers),
192        };
193      }
194      return {
195        id: p.id,
196        song: p.song,
197        codeHash: p.code_hash,
198        model: p.model,
199        every: p.every,
200        ended: p.ended,
201        form: p.form_jev !== null && p.form_question !== null ? { jev: p.form_jev, question: p.form_question } : null,
202        jevs,
203        day: p.day,
204        at: p.recorded_at,
205        player: playerOf(p),
206      };
207    },

A song's performances, newest first, each with who played it. A segment's section is its form jev's row (every jev's row carries it; see segmentRows), so the path is read from that jev alone; without a form, from jev 0, and it is null.

213    async takes(song: string, limit: number, offset: number) {
214      const [total, page] = (await db.batch([
215        db.prepare('SELECT count(*) AS n FROM performances WHERE song = ?').bind(song),
216        db
217          .prepare(
218            `SELECT ${TAKE_COLUMNS} FROM performances p LEFT JOIN users u ON u.id = p.user_id WHERE p.song = ?
219             ORDER BY p.recorded_at DESC, p.rowid DESC LIMIT ? OFFSET ?`,
220          )
221          .bind(song, limit, offset),
222      ])) as D1Result<Record<string, unknown>>[];
223      const paths = await pathsOf(db, page.results.filter((r) => r.form_jev !== null).map((r) => r.id as string));
224      return { total: (total.results[0]?.n as number) ?? 0, takes: page.results.map((r) => takeOf(r, paths)) };
225    },

An account's own takes, newest first, with each one's song: its "recently played" and its profile (data.ts).

229    async playerTakes(userId: string, limit: number): Promise<PlayerTake[]> {
230      const { results } = await db
231        .prepare(
232          `SELECT ${TAKE_COLUMNS},
233             (SELECT r.title FROM song_revisions r WHERE p.song = 'listener:' || r.song AND r.screen = 'fine' ORDER BY r.rev DESC LIMIT 1) AS title
234           FROM performances p LEFT JOIN users u ON u.id = p.user_id WHERE p.user_id = ?
235           ORDER BY p.recorded_at DESC, p.rowid DESC LIMIT ?`,
236        )
237        .bind(userId, limit)
238        .all<Record<string, unknown>>();
239      const paths = await pathsOf(db, results.filter((r) => r.form_jev !== null).map((r) => r.id as string));
240      return results.map((r) => ({ ...takeOf(r, paths), song: r.song as string, title: (r.title as string | null) ?? null }));
241    },

How Jev moves through a song (listening.ts's Moves), over its latest MOVES_WINDOW performances. A segment is one request whichever jev()s it has: it fell back when any of them did, and took as long as the slowest. Answer times are aggregated here, from the sorted times.

247    async moves(song: string): Promise<Moves> {
248      const window = `SELECT id, coalesce(form_jev, 0) AS jev, form_jev FROM performances WHERE song = ?1
249                      ORDER BY recorded_at DESC, rowid DESC LIMIT ${MOVES_WINDOW}`;
250      const [counted, segments, next] = (await db.batch([
251        db.prepare(`SELECT (SELECT count(*) FROM (${window})) AS window, (SELECT count(*) FROM performances WHERE song = ?1) AS total`).bind(song),
252        db
253          .prepare(
254            `WITH w AS (${window}),
255             seg AS (
256               SELECT s.performance, s.segment, max(s.status = 'fallback') AS fell,
257                 max(s.status <> 'opening') AS asked, max(s.ms) AS ms
258               FROM performance_segments s JOIN w ON w.id = s.performance GROUP BY s.performance, s.segment
259             )
260             SELECT coalesce(f.section, CAST(f.segment + 1 AS TEXT)) AS section, seg.fell, seg.asked, seg.ms
261             FROM performance_segments f JOIN w ON w.id = f.performance AND f.jev = w.jev
262             JOIN seg ON seg.performance = f.performance AND seg.segment = f.segment
263             ORDER BY section, seg.ms`,
264          )
265          .bind(song),
266        db
267          .prepare(
268            `WITH w AS (${window})
269             SELECT f.after_section AS "from", f.section AS "to", count(*) AS n
270             FROM performance_segments f JOIN w ON w.id = f.performance AND f.jev = w.form_jev
271             WHERE f.after_section IS NOT NULL AND f.section IS NOT NULL
272             GROUP BY 1, 2 ORDER BY 1, n DESC, 2`,
273          )
274          .bind(song),
275      ])) as D1Result<Record<string, unknown>>[];
276      const bySection = new Map<string, { plays: number; asked: number; fallbacks: number; ms: number[] }>();
277      for (const r of segments.results) {
278        const name = r.section as string;
279        const s = bySection.get(name) ?? { plays: 0, asked: 0, fallbacks: 0, ms: [] };
280        s.plays += 1;
281        if (r.asked) s.asked += 1;
282        if (r.asked && r.fell) s.fallbacks += 1;
283        if (r.asked && typeof r.ms === 'number') s.ms.push(r.ms);
284        bySection.set(name, s);
285      }
286      const nexts = new Map<string, { section: string; n: number }[]>();
287      for (const r of next.results) {
288        const from = r.from as string;
289        nexts.set(from, [...(nexts.get(from) ?? []), { section: r.to as string, n: r.n as number }]);
290      }
291      const quantile = (sorted: number[], q: number) => sorted[Math.min(sorted.length - 1, Math.ceil(q * sorted.length) - 1)];
292      const sections = [...bySection].map(([section, s]) => {
293        const out = nexts.get(section) ?? [];
294        const all = out.reduce((a, x) => a + x.n, 0);
295        const ms = [...s.ms].sort((a, b) => a - b);
296        return {
297          section,
298          plays: s.plays,
299          asked: s.asked,
300          fallbacks: s.fallbacks,
301          ms: ms.length ? { median: quantile(ms, 0.5), p90: quantile(ms, 0.9) } : null,
302          next: out.map((x) => ({ ...x, share: x.n / all })),
303        };
304      });
305      const c = counted.results[0] ?? {};
306      return { song, performances: (c.window as number) ?? 0, of: (c.total as number) ?? 0, sections };
307    },
308  };
309}

What anyone may read of listening (data.ts): every reaction count, every performance, and a performance's rows as stored. Nothing here is about a listener; a performance id is its replay link.

314export function d1ListeningReads(db: D1Database) {
315  return {
316    async reactionsPage(limit: number, offset: number, song: string | null) {
317      const [total, page] = (await db.batch([
318        db.prepare('SELECT count(*) AS n FROM reactions WHERE ?1 IS NULL OR song = ?1').bind(song),
319        db
320          .prepare(
321            `SELECT song, section, reaction, n FROM reactions WHERE ?1 IS NULL OR song = ?1
322             ORDER BY song, section, reaction LIMIT ?2 OFFSET ?3`,
323          )
324          .bind(song, limit, offset),
325      ])) as D1Result<Record<string, unknown>>[];
326      return {
327        total: (total.results[0]?.n as number) ?? 0,
328        reactions: page.results as { song: string; section: string; reaction: 'fire' | 'sleep'; n: number }[],
329      };
330    },

Newest first, each with who played it (null when signed out).

333    async performancesPage(limit: number, offset: number, song: string | null) {
334      const [total, page] = (await db.batch([
335        db.prepare('SELECT count(*) AS n FROM performances WHERE ?1 IS NULL OR song = ?1').bind(song),
336        db
337          .prepare(
338            `SELECT p.id, p.song, p.model, p.every, p.ended, p.jevs, p.form_question AS form, p.day, p.recorded_at AS at,
339               u.id AS player_id, u.display_name AS player_name,
340               (SELECT count(*) FROM performance_segments s WHERE s.performance = p.id) AS segments,
341               (SELECT count(*) FROM performance_segments s WHERE s.performance = p.id AND s.status = 'fallback') AS fallbacks
342             FROM performances p LEFT JOIN users u ON u.id = p.user_id WHERE ?1 IS NULL OR p.song = ?1
343             ORDER BY p.recorded_at DESC, p.rowid DESC LIMIT ?2 OFFSET ?3`,
344          )
345          .bind(song, limit, offset),
346      ])) as D1Result<Record<string, unknown>>[];
347      return {
348        total: (total.results[0]?.n as number) ?? 0,
349        performances: page.results.map(({ player_id, player_name, ...r }) => ({
350          ...r,
351          player: playerOf({ player_id, player_name }),
352        })) as {
353          id: string;
354          song: string;
355          model: string;
356          every: number;
357          ended: string;
358          jevs: number;
359          form: string | null;
360          day: string;
361          at: number;
362          player: Player;
363          segments: number;
364          fallbacks: number;
365        }[],
366      };
367    },

A performance's rows exactly as stored: Jev's decision history.

370    async segments(id: string) {
371      const { results } = await db
372        .prepare(
373          `SELECT jev, segment, section, after_section, status, ms, from_cycle, answers FROM performance_segments
374           WHERE performance = ? ORDER BY segment, jev`,
375        )
376        .bind(id)
377        .all<Row>();
378      return results.map((r) => ({
379        jev: r.jev,
380        segment: r.segment,
381        section: r.section,
382        afterSection: r.after_section,
383        status: r.status,
384        ms: r.ms,
385        fromCycle: r.from_cycle,
386        answers: JSON.parse(r.answers) as Record<string, unknown>,
387      }));
388    },
390    async totals() {
391      const [reactions, performances, segments] = (await db.batch([
392        db.prepare('SELECT count(*) AS rows, coalesce(sum(n), 0) AS n FROM reactions'),
393        db.prepare('SELECT count(*) AS rows FROM performances'),
394        db.prepare(`SELECT count(*) AS rows, coalesce(sum(status = 'fallback'), 0) AS fallbacks FROM performance_segments`),
395      ])) as D1Result<Record<string, number>>[];
396      return {
397        reactions: { rows: reactions.results[0].rows, count: reactions.results[0].n },
398        performances: performances.results[0].rows,
399        segments: { rows: segments.results[0].rows, fallbacks: segments.results[0].fallbacks },
400      };
401    },
402  };
403}