| 1 | /** |
| 2 | * One workspace's presence room, as a Durable Object named by its slug. |
| 3 | * |
| 4 | * Each member's feed (src/feed.ts) tells the room how they show whenever |
| 5 | * that changes: presence, status, Do Not Disturb. The room keeps the |
| 6 | * latest word from each, in its own SQLite storage, and passes every |
| 7 | * change on to the feeds of the members who are online now, which send it |
| 8 | * down their tabs open in this workspace. Members who are offline hear |
| 9 | * nothing; they read everyone afresh when a tab connects (`report` with |
| 10 | * `snapshot`). |
| 11 | * |
| 12 | * The room authorizes nothing: only feeds reach it, and a feed reports to |
| 13 | * the workspaces the site said its person belongs to. |
| 14 | */ |
| 15 | import { DurableObject } from "cloudflare:workers"; |
| 16 | |
| 17 | import type { PresenceEntry } from "@g1t/contracts"; |
| 18 | |
| 19 | import type { Feed } from "./feed.ts"; |
| 20 | import { liveEntry, sameEntry } from "./presence.ts"; |
| 21 | |
| 22 | export type RoomEnv = { FEEDS?: DurableObjectNamespace<Feed> }; |
| 23 | |
| 24 | /** The most members a room keeps. */ |
| 25 | const MAX_PEOPLE = 5000; |
| 26 | |
| 27 | export class Room extends DurableObject<RoomEnv> { |
| 28 | private readonly sql: SqlStorage; |
| 29 | |
| 30 | constructor(ctx: DurableObjectState, env: RoomEnv) { |
| 31 | super(ctx, env); |
| 32 | this.sql = ctx.storage.sql; |
| 33 | this.sql.exec(`CREATE TABLE IF NOT EXISTS people (user_id TEXT PRIMARY KEY, at INTEGER NOT NULL, presence TEXT NOT NULL, json TEXT NOT NULL)`); |
| 34 | } |
| 35 | |
| 36 | private all(): PresenceEntry[] { |
| 37 | const now = Date.now(); |
| 38 | return this.sql |
| 39 | .exec<{ json: string }>("SELECT json FROM people ORDER BY at DESC LIMIT ?", MAX_PEOPLE) |
| 40 | .toArray() |
| 41 | .map((row) => liveEntry(JSON.parse(row.json) as PresenceEntry, now)); |
| 42 | } |
| 43 | |
| 44 | private find(userId: string): PresenceEntry | null { |
| 45 | const row = this.sql.exec<{ json: string }>("SELECT json FROM people WHERE user_id = ?", userId).toArray()[0]; |
| 46 | return row ? (JSON.parse(row.json) as PresenceEntry) : null; |
| 47 | } |
| 48 | |
| 49 | /** |
| 50 | * A member's latest word: kept, and passed to everyone online when it |
| 51 | * changes how they show. With `snapshot`, everyone known comes back, for |
| 52 | * a tab that just connected. |
| 53 | */ |
| 54 | async report(workspace: string, entry: PresenceEntry, snapshot = false): Promise<PresenceEntry[] | null> { |
| 55 | const before = this.find(entry.user_id); |
| 56 | // An older word, arriving late, changes nothing. |
| 57 | if (!before || before.at <= entry.at) { |
| 58 | this.sql.exec( |
| 59 | "INSERT OR REPLACE INTO people (user_id, at, presence, json) VALUES (?, ?, ?, ?)", |
| 60 | entry.user_id, |
| 61 | entry.at, |
| 62 | entry.presence, |
| 63 | JSON.stringify(entry), |
| 64 | ); |
| 65 | if (!sameEntry(before, entry)) this.ctx.waitUntil(this.tell(workspace, entry)); |
| 66 | } |
| 67 | return snapshot ? this.all() : null; |
| 68 | } |
| 69 | |
| 70 | /** Everyone known, as they show now. */ |
| 71 | async people(): Promise<PresenceEntry[]> { |
| 72 | return this.all(); |
| 73 | } |
| 74 | |
| 75 | /** Passes a change to every member online now, the one it is about included (their other tabs). */ |
| 76 | private async tell(workspace: string, entry: PresenceEntry): Promise<void> { |
| 77 | const feeds = this.env.FEEDS; |
| 78 | if (!feeds) return; |
| 79 | const online = this.sql |
| 80 | .exec<{ user_id: string }>("SELECT user_id FROM people WHERE presence != 'offline' OR user_id = ?", entry.user_id) |
| 81 | .toArray() |
| 82 | .map((row) => row.user_id); |
| 83 | const results = await Promise.allSettled(online.map((id) => feeds.get(feeds.idFromName(id)).presence(workspace, [entry]))); |
| 84 | for (const result of results) if (result.status === "rejected") console.error("notify: telling a feed of presence failed", result.reason); |
| 85 | } |
| 86 | } |