Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002) | 1 | /** |
| 2 | * One person's feed, as a Durable Object named by their user id. | |
| 3 | * | |
| 4 | * It holds a socket per open tab with the WebSocket Hibernation API, so a | |
| 5 | * person with g1t open all day costs nothing between notifications. Plain | |
| 6 | * `ping` keepalives are answered at the edge without waking it, and the | |
| 7 | * time of the last one says whether a tab is still there | |
| 8 | * (`getWebSocketAutoResponseTimestamp`). Each socket carries the tab's | |
| 9 | * focus and page (`serializeAttachment`), which survive hibernation. | |
| 10 | * | |
| 11 | * Kept in the object's own SQLite storage: the latest notifications, the | |
| 12 | * unread counts per conversation, push subscriptions and preferences. | |
| 13 | * | |
| 14 | * The feed authorizes nothing: the Worker reaches it only for the person | |
| 15 | * the site checked (src/index.ts). | |
| 16 | */ | |
| 17 | import { DurableObject } from "cloudflare:workers"; | |
| 18 | ||
| 19 | import { | |
| 20 | NOTIFY_SEED_HEADER, | |
| 21 | type ChannelCounts, | |
| 22 | type FeedDelivery, | |
| 23 | type FeedEvent, | |
| 24 | type FeedNotification, | |
| 25 | type FeedSeed, | |
| 26 | type NotifyPreferences, | |
| 27 | type NotifyStatus, | |
| 28 | type PushSubscriptionJson, | |
| 29 | } from "@g1t/contracts"; | |
| 30 | ||
| 31 | import { applyCounts, totals } from "./counts.ts"; | |
| 32 | import { cleanNotification, decide, mergePreferences, pushPayload, readPreferences, type TabState } from "./prefs.ts"; | |
| 33 | import { sendPush, type Vapid } from "./webpush.ts"; | |
| 34 | ||
| 35 | export type FeedEnv = { VAPID_PUBLIC_KEY?: string; VAPID_PRIVATE_KEY?: string; VAPID_SUBJECT?: string }; | |
| 36 | ||
| 37 | /** What each socket carries. */ | |
| 38 | type Tab = { focused: boolean; path: string; at: number }; | |
| 39 | ||
| 40 | /** Notifications kept for a tab that opens later. */ | |
| 41 | const KEPT = 100; | |
| 42 | /** Sent with `hello`: the latest few. */ | |
| 43 | const HELLO = 20; | |
| 44 | /** The most browsers one person gets pushes on. */ | |
| 45 | const MAX_SUBSCRIPTIONS = 20; | |
| 46 | /** The longest endpoint taken. */ | |
| 47 | const MAX_ENDPOINT = 2000; | |
| 48 | ||
| 49 | export class Feed extends DurableObject<FeedEnv> { | |
| 50 | private readonly sql: SqlStorage; | |
| 51 | ||
| 52 | constructor(ctx: DurableObjectState, env: FeedEnv) { | |
| 53 | super(ctx, env); | |
| 54 | this.sql = ctx.storage.sql; | |
| 55 | // Keepalives are answered without waking the object. | |
| 56 | ctx.setWebSocketAutoResponse(new WebSocketRequestResponsePair("ping", "pong")); | |
| 57 | this.sql.exec(`CREATE TABLE IF NOT EXISTS notifications (id TEXT PRIMARY KEY, at INTEGER NOT NULL, json TEXT NOT NULL)`); | |
| 58 | this.sql.exec(`CREATE INDEX IF NOT EXISTS notifications_at ON notifications (at)`); | |
| 59 | this.sql.exec( | |
| 60 | `CREATE TABLE IF NOT EXISTS counts (workspace TEXT NOT NULL, channel_id TEXT NOT NULL, unread INTEGER NOT NULL, mentions INTEGER NOT NULL, muted INTEGER NOT NULL DEFAULT 0, PRIMARY KEY (workspace, channel_id))`, | |
| 61 | ); | |
| 62 | this.sql.exec( | |
| 63 | `CREATE TABLE IF NOT EXISTS subscriptions (endpoint TEXT PRIMARY KEY, p256dh TEXT NOT NULL, auth TEXT NOT NULL, user_agent TEXT, created_at INTEGER NOT NULL)`, | |
| 64 | ); | |
| 65 | this.sql.exec(`CREATE TABLE IF NOT EXISTS meta (key TEXT PRIMARY KEY, value TEXT NOT NULL)`); | |
| 66 | } | |
| 67 | ||
| 68 | // ── Kept state ────────────────────────────────────────────────────────── | |
| 69 | ||
| 70 | private meta(key: string): string | null { | |
| 71 | const row = this.sql.exec<{ value: string }>("SELECT value FROM meta WHERE key = ?", key).toArray()[0]; | |
| 72 | return row?.value ?? null; | |
| 73 | } | |
| 74 | ||
| 75 | private setMeta(key: string, value: string): void { | |
| 76 | this.sql.exec("INSERT INTO meta (key, value) VALUES (?, ?) ON CONFLICT (key) DO UPDATE SET value = excluded.value", key, value); | |
| 77 | } | |
| 78 | ||
| 79 | private preferences(): NotifyPreferences { | |
| 80 | return readPreferences(this.meta("preferences")); | |
| 81 | } | |
| 82 | ||
| 83 | private inboxUnread(): number { | |
| 84 | return Number(this.meta("inbox_unread") ?? 0) || 0; | |
| 85 | } | |
| 86 | ||
| 87 | private channelRows(workspace: string): ChannelCounts[] { | |
| 88 | return this.sql | |
| 89 | .exec<{ channel_id: string; unread: number; mentions: number; muted: number }>( | |
| 90 | "SELECT channel_id, unread, mentions, muted FROM counts WHERE workspace = ?", | |
| 91 | workspace, | |
| 92 | ) | |
| 93 | .toArray() | |
| 94 | .map((r) => ({ channel_id: r.channel_id, unread: r.unread, mentions: r.mentions, muted: !!r.muted })); | |
| 95 | } | |
| 96 | ||
| 97 | private counts(workspace: string): FeedEvent { | |
| 98 | return { type: "counts", ...totals(workspace, this.channelRows(workspace), this.inboxUnread(), this.meta(`complete:${workspace}`) === "1") }; | |
| 99 | } | |
| 100 | ||
| 101 | private workspaces(): string[] { | |
| 102 | return this.sql.exec<{ workspace: string }>("SELECT DISTINCT workspace FROM counts").toArray().map((r) => r.workspace); | |
| 103 | } | |
| 104 | ||
| 105 | private vapid(): Vapid | null { | |
| 106 | const { VAPID_PUBLIC_KEY, VAPID_PRIVATE_KEY, VAPID_SUBJECT } = this.env; | |
| 107 | return VAPID_PUBLIC_KEY && VAPID_PRIVATE_KEY ? { publicKey: VAPID_PUBLIC_KEY, privateKey: VAPID_PRIVATE_KEY, subject: VAPID_SUBJECT || "https://g1t.sh" } : null; | |
| 108 | } | |
| 109 | ||
| 110 | private subscriptionCount(): number { | |
| 111 | return this.sql.exec<{ n: number }>("SELECT COUNT(*) AS n FROM subscriptions").one().n; | |
| 112 | } | |
| 113 | ||
| 114 | // ── Sockets ───────────────────────────────────────────────────────────── | |
| 115 | ||
| 116 | private tabs(): TabState[] { | |
| 117 | return this.ctx.getWebSockets().map((socket) => { | |
| 118 | const tab = (socket.deserializeAttachment() as Tab | null) ?? { focused: false, path: "", at: 0 }; | |
| 119 | const pinged = this.ctx.getWebSocketAutoResponseTimestamp(socket)?.getTime() ?? 0; | |
| 120 | return { focused: tab.focused, seen_at: Math.max(tab.at, pinged) }; | |
| 121 | }); | |
| 122 | } | |
| 123 | ||
| 124 | private send(event: FeedEvent): void { | |
| 125 | const text = JSON.stringify(event); | |
| 126 | for (const socket of this.ctx.getWebSockets()) { | |
| 127 | try { | |
| 128 | socket.send(text); | |
| 129 | } catch { | |
| 130 | // Closing already. | |
| 131 | } | |
| 132 | } | |
| 133 | } | |
| 134 | ||
| 135 | override async fetch(request: Request): Promise<Response> { | |
| 136 | if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") { | |
| 137 | return new Response("Expected a WebSocket upgrade\n", { status: 426 }); | |
| 138 | } | |
| 139 | let seed: FeedSeed | null = null; | |
| 140 | try { | |
| 141 | seed = JSON.parse(request.headers.get(NOTIFY_SEED_HEADER) ?? "null") as FeedSeed | null; | |
| 142 | } catch { | |
| 143 | seed = null; | |
| 144 | } | |
| 145 | if (seed) this.applySeed(seed); | |
| 146 | const pair = new WebSocketPair(); | |
| 147 | const [client, server] = Object.values(pair) as [WebSocket, WebSocket]; | |
| 148 | this.ctx.acceptWebSocket(server); | |
| 149 | server.serializeAttachment({ focused: false, path: "", at: Date.now() } satisfies Tab); | |
| 150 | const hello: FeedEvent = { | |
| 151 | type: "hello", | |
| 152 | notifications: this.latest(HELLO), | |
| 153 | vapid_public_key: this.env.VAPID_PUBLIC_KEY || null, | |
| 154 | preferences: this.preferences(), | |
| 155 | }; | |
| 156 | server.send(JSON.stringify(hello)); | |
| 157 | const shown = new Set(this.workspaces()); | |
| 158 | if (seed?.workspace) shown.add(seed.workspace); | |
| 159 | for (const workspace of shown) server.send(JSON.stringify(this.counts(workspace))); | |
| 160 | return new Response(null, { status: 101, webSocket: client }); | |
| 161 | } | |
| 162 | ||
| 163 | /** Counts read from chat and the inbox just now replace what was kept for that workspace. */ | |
| 164 | private applySeed(seed: FeedSeed): void { | |
| 165 | if (typeof seed.inbox_unread === "number") this.setMeta("inbox_unread", String(Math.max(0, Math.floor(seed.inbox_unread)))); | |
| 166 | const workspace = seed.workspace?.toLowerCase(); | |
| 167 | if (!workspace || !Array.isArray(seed.per_channel)) return; | |
| 168 | this.ctx.storage.transactionSync(() => { | |
| 169 | this.sql.exec("DELETE FROM counts WHERE workspace = ?", workspace); | |
| 170 | for (const row of seed.per_channel!.slice(0, 2000)) { | |
| 171 | if (typeof row?.channel_id !== "string") continue; | |
| 172 | const kept = applyCounts(null, { ...row, set: true }); | |
| 173 | this.sql.exec( | |
| 174 | "INSERT OR REPLACE INTO counts (workspace, channel_id, unread, mentions, muted) VALUES (?, ?, ?, ?, ?)", | |
| 175 | workspace, | |
| 176 | kept.channel_id, | |
| 177 | kept.unread, | |
| 178 | kept.mentions, | |
| 179 | kept.muted ? 1 : 0, | |
| 180 | ); | |
| 181 | } | |
| 182 | this.setMeta(`complete:${workspace}`, "1"); | |
| 183 | }); | |
| 184 | } | |
| 185 | ||
| 186 | override async webSocketMessage(socket: WebSocket, message: string | ArrayBuffer): Promise<void> { | |
| 187 | if (typeof message !== "string" || message.length > 4096) return; | |
| 188 | let frame: Record<string, unknown>; | |
| 189 | try { | |
| 190 | frame = JSON.parse(message); | |
| 191 | } catch { | |
| 192 | return; | |
| 193 | } | |
| 194 | if (frame.type === "state") { | |
| 195 | socket.serializeAttachment({ focused: frame.focused === true, path: String(frame.path ?? "").slice(0, 500), at: Date.now() } satisfies Tab); | |
| 196 | } else if (frame.type === "inbox" && Number.isFinite(frame.unread)) { | |
| 197 | // The inbox count a page just read: every tab shows it at once. | |
| 198 | await this.setInbox(Number(frame.unread)); | |
| 199 | } | |
| 200 | } | |
| 201 | ||
| 202 | override async webSocketClose(socket: WebSocket, code: number, reason: string): Promise<void> { | |
| 203 | try { | |
| 204 | socket.close(code, reason); | |
| 205 | } catch { | |
| 206 | // Already closed. | |
| 207 | } | |
| 208 | } | |
| 209 | ||
| 210 | override async webSocketError(): Promise<void> {} | |
| 211 | ||
| 212 | // ── Notifications ─────────────────────────────────────────────────────── | |
| 213 | ||
| 214 | private latest(limit: number): FeedNotification[] { | |
| 215 | return this.sql | |
| 216 | .exec<{ json: string }>("SELECT json FROM notifications ORDER BY at DESC LIMIT ?", limit) | |
| 217 | .toArray() | |
| 218 | .map((r) => JSON.parse(r.json) as FeedNotification); | |
| 219 | } | |
| 220 | ||
| 221 | /** Keeps it, unless it was told before. True when it is news. */ | |
| 222 | private keep(notification: FeedNotification): boolean { | |
| 223 | const had = this.sql.exec("SELECT 1 FROM notifications WHERE id = ?", notification.id).toArray().length > 0; | |
| 224 | if (had) return false; | |
| 225 | this.sql.exec("INSERT INTO notifications (id, at, json) VALUES (?, ?, ?)", notification.id, Date.now(), JSON.stringify(notification)); | |
| 226 | this.sql.exec("DELETE FROM notifications WHERE id NOT IN (SELECT id FROM notifications ORDER BY at DESC LIMIT ?)", KEPT); | |
| 227 | return true; | |
| 228 | } | |
| 229 | ||
| 230 | /** Tells the person: open tabs at once, then browsers when none is in front of them. */ | |
| 231 | private async tell(notification: FeedNotification, test = false): Promise<number> { | |
| 232 | if (!test && !this.keep(notification)) return 0; | |
| 233 | const subscriptions = this.subscriptionCount(); | |
| 234 | const decision = decide({ prefs: this.preferences(), notification, tabs: this.tabs(), subscriptions, now: Date.now(), test }); | |
| 235 | this.send({ type: "notification", notification, toast: decision.toast }); | |
| 236 | return decision.push ? this.push(notification) : 0; | |
| 237 | } | |
| 238 | ||
| 239 | /** Pushes to every browser subscribed; drops the ones the push service says are gone. */ | |
| 240 | private async push(notification: FeedNotification): Promise<number> { | |
| 241 | const vapid = this.vapid(); | |
| 242 | if (!vapid) return 0; | |
| 243 | const subscriptions = this.sql.exec<{ endpoint: string; p256dh: string; auth: string }>("SELECT endpoint, p256dh, auth FROM subscriptions").toArray(); | |
| 244 | const payload = pushPayload(notification); | |
| 245 | const results = await Promise.allSettled( | |
| 246 | subscriptions.map((s) => sendPush(s, payload, { vapid, urgency: payload.urgent ? "high" : "normal", topic: payload.tag, ttl: 12 * 3600 })), | |
| 247 | ); | |
| 248 | let sent = 0; | |
| 249 | for (const result of results) { | |
| 250 | if (result.status === "rejected") { | |
| 251 | console.error("notify: a push failed", result.reason); | |
| 252 | continue; | |
| 253 | } | |
| 254 | if (result.value.gone) this.sql.exec("DELETE FROM subscriptions WHERE endpoint = ?", result.value.endpoint); | |
| 255 | else if (result.value.status < 300) sent++; | |
| 256 | else console.error("notify: a push service answered", result.value.status); | |
| 257 | } | |
| 258 | return sent; | |
| 259 | } | |
| 260 | ||
| 261 | // ── RPC (called by the Worker, src/index.ts) ──────────────────────────── | |
| 262 | ||
| 263 | async notify(value: unknown): Promise<{ ok: boolean }> { | |
| 264 | const notification = cleanNotification(value); | |
| 265 | if (!notification) return { ok: false }; | |
| 266 | await this.tell(notification); | |
| 267 | return { ok: true }; | |
| 268 | } | |
| 269 | ||
| 270 | /** A batch for this person from chat: counts moved, notifications told. */ | |
| 271 | async deliver(items: FeedDelivery[]): Promise<{ ok: boolean }> { | |
| 272 | const changed = new Set<string>(); | |
| 273 | const told: FeedNotification[] = []; | |
| 274 | for (const item of items) { | |
| 275 | const workspace = String(item.workspace ?? "").toLowerCase(); | |
| 276 | if (item.counts?.channel_id && workspace) { | |
| 277 | const before = this.sql | |
| 278 | .exec<{ unread: number; mentions: number; muted: number }>( | |
| 279 | "SELECT unread, mentions, muted FROM counts WHERE workspace = ? AND channel_id = ?", | |
| 280 | workspace, | |
| 281 | item.counts.channel_id, | |
| 282 | ) | |
| 283 | .toArray()[0]; | |
| 284 | const after = applyCounts( | |
| 285 | before ? { channel_id: item.counts.channel_id, unread: before.unread, mentions: before.mentions, muted: !!before.muted } : null, | |
| 286 | item.counts, | |
| 287 | ); | |
| 288 | this.sql.exec( | |
| 289 | "INSERT OR REPLACE INTO counts (workspace, channel_id, unread, mentions, muted) VALUES (?, ?, ?, ?, ?)", | |
| 290 | workspace, | |
| 291 | after.channel_id, | |
| 292 | after.unread, | |
| 293 | after.mentions, | |
| 294 | after.muted ? 1 : 0, | |
| 295 | ); | |
| 296 | changed.add(workspace); | |
| 297 | } | |
| 298 | const notification = item.notification ? cleanNotification(item.notification) : null; | |
| 299 | if (notification) told.push(notification); | |
| 300 | } | |
| 301 | for (const workspace of changed) this.send(this.counts(workspace)); | |
| 302 | await Promise.all(told.map((n) => this.tell(n))); | |
| 303 | return { ok: true }; | |
| 304 | } | |
| 305 | ||
| 306 | /** The inbox count as it now is (events, after items arrive or are marked): every tab shows it at once. */ | |
| 307 | async setInbox(value: number): Promise<{ ok: boolean }> { | |
| 308 | const unread = Number.isFinite(value) ? Math.max(0, Math.floor(value)) : 0; | |
| 309 | if (unread === this.inboxUnread()) return { ok: true }; | |
| 310 | this.setMeta("inbox_unread", String(unread)); | |
| 311 | this.send({ type: "inbox", unread }); | |
| 312 | return { ok: true }; | |
| 313 | } | |
| 314 | ||
| 315 | async subscribe(subscription: PushSubscriptionJson, userAgent: string | null): Promise<{ ok: boolean }> { | |
| 316 | const endpoint = String(subscription?.endpoint ?? ""); | |
| 317 | const { p256dh, auth } = subscription?.keys ?? ({} as PushSubscriptionJson["keys"]); | |
| 318 | if (!/^https:\/\//.test(endpoint) || endpoint.length > MAX_ENDPOINT || typeof p256dh !== "string" || typeof auth !== "string") return { ok: false }; | |
| 319 | this.sql.exec( | |
| 320 | "INSERT OR REPLACE INTO subscriptions (endpoint, p256dh, auth, user_agent, created_at) VALUES (?, ?, ?, ?, ?)", | |
| 321 | endpoint, | |
| 322 | p256dh.slice(0, 200), | |
| 323 | auth.slice(0, 100), | |
| 324 | userAgent ? userAgent.slice(0, 300) : null, | |
| 325 | Date.now(), | |
| 326 | ); | |
| 327 | // The oldest go first past the limit. | |
| 328 | this.sql.exec("DELETE FROM subscriptions WHERE endpoint NOT IN (SELECT endpoint FROM subscriptions ORDER BY created_at DESC LIMIT ?)", MAX_SUBSCRIPTIONS); | |
| 329 | return { ok: true }; | |
| 330 | } | |
| 331 | ||
| 332 | async unsubscribe(endpoint: string): Promise<{ ok: boolean }> { | |
| 333 | this.sql.exec("DELETE FROM subscriptions WHERE endpoint = ?", String(endpoint ?? "")); | |
| 334 | return { ok: true }; | |
| 335 | } | |
| 336 | ||
| 337 | async status(endpoint: string | null): Promise<NotifyStatus> { | |
| 338 | const subscribed = !!endpoint && this.sql.exec("SELECT 1 FROM subscriptions WHERE endpoint = ?", endpoint).toArray().length > 0; | |
| 339 | return { preferences: this.preferences(), subscriptions: this.subscriptionCount(), subscribed, vapid_public_key: this.env.VAPID_PUBLIC_KEY || null }; | |
| 340 | } | |
| 341 | ||
| 342 | async setPreferences(change: unknown): Promise<NotifyPreferences> { | |
| 343 | const preferences = mergePreferences(this.preferences(), change); | |
| 344 | this.setMeta("preferences", JSON.stringify(preferences)); | |
| 345 | this.send({ type: "preferences", preferences }); | |
| 346 | return preferences; | |
| 347 | } | |
| 348 | ||
| 349 | async test(username: string): Promise<{ ok: boolean; pushed: number }> { | |
| 350 | const pushed = await this.tell( | |
| 351 | { | |
| 352 | id: `test:${crypto.randomUUID()}`, | |
| 353 | kind: "dm", | |
| 354 | workspace: "", | |
| 355 | title: "g1t", | |
| 356 | body: `Notifications are on, ${username || "there"}. This is what a message looks like.`, | |
| 357 | href: "/settings/notifications", | |
| 358 | actor: { kind: "system", id: "g1t", name: "g1t", avatar: null, avatar_seed: null }, | |
| 359 | channel_id: null, | |
| 360 | thread_root: null, | |
| 361 | created_at: new Date().toISOString(), | |
| 362 | }, | |
| 363 | true, | |
| 364 | ); | |
| 365 | return { ok: true, pushed }; | |
| 366 | } | |
| 367 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.