| 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 (browsers) and |
| 13 | * devices (the phone app, pushed through Expo: src/expopush.ts) and |
| 14 | * preferences, and |
| 15 | * the person's status, Do Not Disturb and whether they set themselves away, |
| 16 | * and when they were last on each workspace's Home page (src/visits.ts). |
| 17 | * |
| 18 | * Presence is worked out here from the tabs (src/presence.ts) and told to |
| 19 | * the room of every workspace the person belongs to (src/room.ts), which |
| 20 | * passes it to everyone there who is online. A status or Do Not Disturb |
| 21 | * that runs out is cleared by an alarm and told the same way. |
| 22 | * |
| 23 | * The feed authorizes nothing: the Worker reaches it only for the person |
| 24 | * the site checked (src/index.ts). |
| 25 | */ |
| 26 | import { DurableObject } from "cloudflare:workers"; |
| 27 | |
| 28 | import { |
| 29 | NOTIFY_SEED_HEADER, |
| 30 | type ChannelCounts, |
| 31 | type OwnPresence, |
| 32 | type PresenceChange, |
| 33 | type PresenceEntry, |
| 34 | type FeedDelivery, |
| 35 | type FeedEvent, |
| 36 | type FeedNotification, |
| 37 | type FeedSeed, |
| 38 | type NotifyPreferences, |
| 39 | type NotifyStatus, |
| 40 | type PushDevice, |
| 41 | type PushSubscriptionJson, |
| 42 | } from "@g1t/contracts"; |
| 43 | |
| 44 | import { applyCounts, totals } from "./counts.ts"; |
| 45 | import { cleanDevice, sendExpo } from "./expopush.ts"; |
| 46 | import { cleanNotification, decide, mergePreferences, pushPayload, readPreferences, type PushPayload, type TabState } from "./prefs.ts"; |
| 47 | import { |
| 48 | OFFLINE_GRACE_MS, |
| 49 | applyChange, |
| 50 | cleanWorkspaces, |
| 51 | current, |
| 52 | dndOn, |
| 53 | entryOf, |
| 54 | nextExpiry, |
| 55 | ownOf, |
| 56 | presenceOf, |
| 57 | readKept, |
| 58 | sameEntry, |
| 59 | type Kept, |
| 60 | } from "./presence.ts"; |
| 61 | import type { Room } from "./room.ts"; |
| 62 | import { type Visit, lastVisit, markVisit } from "./visits.ts"; |
| 63 | import { sendPush, type Vapid } from "./webpush.ts"; |
| 64 | |
| 65 | export type FeedEnv = { |
| 66 | VAPID_PUBLIC_KEY?: string; |
| 67 | VAPID_PRIVATE_KEY?: string; |
| 68 | VAPID_SUBJECT?: string; |
| 69 | /** |
| 70 | * A secret, optional: the access token Expo's push service asks for once |
| 71 | * the Expo project turns on enhanced push security. Without it, pushes to |
| 72 | * phones go unauthenticated, which Expo takes otherwise. |
| 73 | */ |
| 74 | EXPO_ACCESS_TOKEN?: string; |
| 75 | /** |
| 76 | * `off` sends phones nothing: a self-hosted g1t that would rather not pass |
| 77 | * notifications through Expo's push service (and Google's and Apple's) |
| 78 | * sets PHONE_PUSH=off (deploy/self-host/configs.mjs). Phones then do not |
| 79 | * register, and the site's /.well-known/g1t.json says so to the app. |
| 80 | */ |
| 81 | PHONE_PUSH?: string; |
| 82 | /** One presence room per workspace (src/room.ts); without it, presence stays in the person's own tabs. */ |
| 83 | ROOMS?: DurableObjectNamespace<Room>; |
| 84 | }; |
| 85 | |
| 86 | /** The headers the Worker passes the person's id and username in, with a socket (src/index.ts). */ |
| 87 | export const FEED_USERNAME_HEADER = "x-g1t-notify-username"; |
| 88 | export const FEED_USER_ID_HEADER = "x-g1t-notify-user-id"; |
| 89 | |
| 90 | /** What each socket carries: focus, page, whether it has gone idle, and the workspace it is open in. */ |
| 91 | type Tab = { focused: boolean; path: string; at: number; idle?: boolean; workspace?: string | null }; |
| 92 | |
| 93 | /** Notifications kept for a tab that opens later. */ |
| 94 | const KEPT = 100; |
| 95 | /** Sent with `hello`: the latest few. */ |
| 96 | const HELLO = 20; |
| 97 | /** The most browsers one person gets pushes on. */ |
| 98 | const MAX_SUBSCRIPTIONS = 20; |
| 99 | /** The longest endpoint taken. */ |
| 100 | const MAX_ENDPOINT = 2000; |
| 101 | /** The most phones one person gets pushes on. */ |
| 102 | const MAX_DEVICES = 10; |
| 103 | |
| 104 | export class Feed extends DurableObject<FeedEnv> { |
| 105 | private readonly sql: SqlStorage; |
| 106 | |
| 107 | constructor(ctx: DurableObjectState, env: FeedEnv) { |
| 108 | super(ctx, env); |
| 109 | this.sql = ctx.storage.sql; |
| 110 | // Keepalives are answered without waking the object. |
| 111 | ctx.setWebSocketAutoResponse(new WebSocketRequestResponsePair("ping", "pong")); |
| 112 | this.sql.exec(`CREATE TABLE IF NOT EXISTS notifications (id TEXT PRIMARY KEY, at INTEGER NOT NULL, json TEXT NOT NULL)`); |
| 113 | this.sql.exec(`CREATE INDEX IF NOT EXISTS notifications_at ON notifications (at)`); |
| 114 | this.sql.exec( |
| 115 | `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))`, |
| 116 | ); |
| 117 | this.sql.exec( |
| 118 | `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)`, |
| 119 | ); |
| 120 | // The phone app's Expo push tokens (src/expopush.ts). |
| 121 | this.sql.exec(`CREATE TABLE IF NOT EXISTS devices (token TEXT PRIMARY KEY, platform TEXT NOT NULL, name TEXT, created_at INTEGER NOT NULL)`); |
| 122 | this.sql.exec(`CREATE TABLE IF NOT EXISTS meta (key TEXT PRIMARY KEY, value TEXT NOT NULL)`); |
| 123 | // When the person was last on each workspace's Home page (src/visits.ts). |
| 124 | this.sql.exec(`CREATE TABLE IF NOT EXISTS visits (workspace TEXT PRIMARY KEY, seen_at INTEGER NOT NULL, previous_at INTEGER)`); |
| 125 | } |
| 126 | |
| 127 | // ── Kept state ────────────────────────────────────────────────────────── |
| 128 | |
| 129 | private meta(key: string): string | null { |
| 130 | const row = this.sql.exec<{ value: string }>("SELECT value FROM meta WHERE key = ?", key).toArray()[0]; |
| 131 | return row?.value ?? null; |
| 132 | } |
| 133 | |
| 134 | private setMeta(key: string, value: string): void { |
| 135 | this.sql.exec("INSERT INTO meta (key, value) VALUES (?, ?) ON CONFLICT (key) DO UPDATE SET value = excluded.value", key, value); |
| 136 | } |
| 137 | |
| 138 | private preferences(): NotifyPreferences { |
| 139 | return readPreferences(this.meta("preferences")); |
| 140 | } |
| 141 | |
| 142 | private inboxUnread(): number { |
| 143 | return Number(this.meta("inbox_unread") ?? 0) || 0; |
| 144 | } |
| 145 | |
| 146 | private channelRows(workspace: string): ChannelCounts[] { |
| 147 | return this.sql |
| 148 | .exec<{ channel_id: string; unread: number; mentions: number; muted: number }>( |
| 149 | "SELECT channel_id, unread, mentions, muted FROM counts WHERE workspace = ?", |
| 150 | workspace, |
| 151 | ) |
| 152 | .toArray() |
| 153 | .map((r) => ({ channel_id: r.channel_id, unread: r.unread, mentions: r.mentions, muted: !!r.muted })); |
| 154 | } |
| 155 | |
| 156 | private counts(workspace: string): FeedEvent { |
| 157 | return { type: "counts", ...totals(workspace, this.channelRows(workspace), this.inboxUnread(), this.meta(`complete:${workspace}`) === "1") }; |
| 158 | } |
| 159 | |
| 160 | private workspaces(): string[] { |
| 161 | return this.sql.exec<{ workspace: string }>("SELECT DISTINCT workspace FROM counts").toArray().map((r) => r.workspace); |
| 162 | } |
| 163 | |
| 164 | private vapid(): Vapid | null { |
| 165 | const { VAPID_PUBLIC_KEY, VAPID_PRIVATE_KEY, VAPID_SUBJECT } = this.env; |
| 166 | return VAPID_PUBLIC_KEY && VAPID_PRIVATE_KEY ? { publicKey: VAPID_PUBLIC_KEY, privateKey: VAPID_PRIVATE_KEY, subject: VAPID_SUBJECT || "https://g1t.sh" } : null; |
| 167 | } |
| 168 | |
| 169 | private subscriptionCount(): number { |
| 170 | return this.sql.exec<{ n: number }>("SELECT COUNT(*) AS n FROM subscriptions").one().n; |
| 171 | } |
| 172 | |
| 173 | private deviceCount(): number { |
| 174 | if (!this.phonePush()) return 0; |
| 175 | return this.sql.exec<{ n: number }>("SELECT COUNT(*) AS n FROM devices").one().n; |
| 176 | } |
| 177 | |
| 178 | /** Whether this g1t sends phones anything (PHONE_PUSH). */ |
| 179 | private phonePush(): boolean { |
| 180 | return this.env.PHONE_PUSH !== "off"; |
| 181 | } |
| 182 | |
| 183 | // ── Presence ──────────────────────────────────────────────────────────── |
| 184 | |
| 185 | private kept(): Kept { |
| 186 | return readKept(this.meta("presence")); |
| 187 | } |
| 188 | |
| 189 | /** Who this feed is for, as the site last said. */ |
| 190 | private person(): { user_id: string; username: string } { |
| 191 | return { user_id: this.meta("user_id") ?? "", username: this.meta("username") ?? "" }; |
| 192 | } |
| 193 | |
| 194 | /** The open tabs as presence sees them, leaving out one that is closing. */ |
| 195 | private presenceTabs(closing?: WebSocket): { idle: boolean }[] { |
| 196 | return this.ctx |
| 197 | .getWebSockets() |
| 198 | .filter((socket) => socket !== closing && socket.readyState === WebSocket.OPEN) |
| 199 | .map((socket) => ({ idle: (socket.deserializeAttachment() as Tab | null)?.idle === true })); |
| 200 | } |
| 201 | |
| 202 | private entry(closing?: WebSocket): PresenceEntry { |
| 203 | const kept = this.kept(); |
| 204 | return entryOf(this.person(), presenceOf(this.presenceTabs(closing), kept.away_manual), kept, Date.now()); |
| 205 | } |
| 206 | |
| 207 | private room(workspace: string) { |
| 208 | const rooms = this.env.ROOMS; |
| 209 | return rooms ? rooms.get(rooms.idFromName(workspace)) : null; |
| 210 | } |
| 211 | |
| 212 | private memberOf(): string[] { |
| 213 | try { |
| 214 | return cleanWorkspaces(JSON.parse(this.meta("workspaces") ?? "[]")); |
| 215 | } catch { |
| 216 | return []; |
| 217 | } |
| 218 | } |
| 219 | |
| 220 | /** |
| 221 | * Works presence out again and, when how the person shows has changed, |
| 222 | * tells their own tabs and every workspace's room. `force` tells them |
| 223 | * anyway (a tab just connected, so a room may have missed the last word). |
| 224 | */ |
| 225 | private async refresh(options: { closing?: WebSocket; force?: boolean } = {}): Promise<PresenceEntry> { |
| 226 | const now = Date.now(); |
| 227 | // What ran out goes, so nobody is told of it again. |
| 228 | const kept = this.kept(); |
| 229 | const live = current(kept, now); |
| 230 | if (JSON.stringify(live) !== JSON.stringify(kept)) this.setMeta("presence", JSON.stringify(live)); |
| 231 | const entry = this.entry(options.closing); |
| 232 | let before: PresenceEntry | null = null; |
| 233 | try { |
| 234 | before = JSON.parse(this.meta("reported") ?? "null") as PresenceEntry | null; |
| 235 | } catch { |
| 236 | before = null; |
| 237 | } |
| 238 | const changed = !sameEntry(before, entry); |
| 239 | if (changed || options.force) { |
| 240 | this.setMeta("reported", JSON.stringify(entry)); |
| 241 | this.send({ type: "me", me: ownOf(entry, live) }, options.closing); |
| 242 | // Without an id there is nobody to name: the site always sends one with a socket. |
| 243 | if (entry.user_id) await Promise.allSettled(this.memberOf().map((workspace) => this.room(workspace)?.report(workspace, entry))); |
| 244 | } |
| 245 | await this.schedule(); |
| 246 | return entry; |
| 247 | } |
| 248 | |
| 249 | /** The next alarm: when a status or Do Not Disturb runs out, or when a person whose last tab closed shows offline. */ |
| 250 | private async schedule(): Promise<void> { |
| 251 | const now = Date.now(); |
| 252 | const times = [nextExpiry(this.kept(), now), Number(this.meta("offline_at") ?? 0) || null].filter((at): at is number => at != null && at > now); |
| 253 | if (times.length) await this.ctx.storage.setAlarm(Math.min(...times)); |
| 254 | else await this.ctx.storage.deleteAlarm(); |
| 255 | } |
| 256 | |
| 257 | override async alarm(): Promise<void> { |
| 258 | const offlineAt = Number(this.meta("offline_at") ?? 0); |
| 259 | if (offlineAt && offlineAt <= Date.now()) this.setMeta("offline_at", "0"); |
| 260 | await this.refresh(); |
| 261 | } |
| 262 | |
| 263 | /** A tab just connected: everyone hears it is here, and it hears how everyone in its workspace shows. */ |
| 264 | private async connected(socket: WebSocket, workspace: string | null): Promise<void> { |
| 265 | this.setMeta("offline_at", "0"); |
| 266 | const entry = await this.refresh({ force: true }); |
| 267 | if (!workspace || !entry.user_id) return; |
| 268 | const people = await this.room(workspace) |
| 269 | ?.report(workspace, entry, true) |
| 270 | .catch((error: unknown) => { |
| 271 | console.error("notify: reading a presence room failed", error); |
| 272 | return null; |
| 273 | }); |
| 274 | if (!people) return; |
| 275 | try { |
| 276 | socket.send(JSON.stringify({ type: "presence", workspace, people, full: true } satisfies FeedEvent)); |
| 277 | } catch { |
| 278 | // Gone already. |
| 279 | } |
| 280 | } |
| 281 | |
| 282 | // ── Sockets ───────────────────────────────────────────────────────────── |
| 283 | |
| 284 | private tabs(): TabState[] { |
| 285 | return this.ctx.getWebSockets().map((socket) => { |
| 286 | const tab = (socket.deserializeAttachment() as Tab | null) ?? { focused: false, path: "", at: 0 }; |
| 287 | const pinged = this.ctx.getWebSocketAutoResponseTimestamp(socket)?.getTime() ?? 0; |
| 288 | return { focused: tab.focused, seen_at: Math.max(tab.at, pinged) }; |
| 289 | }); |
| 290 | } |
| 291 | |
| 292 | private send(event: FeedEvent, except?: WebSocket): void { |
| 293 | const text = JSON.stringify(event); |
| 294 | for (const socket of this.ctx.getWebSockets()) { |
| 295 | if (socket === except) continue; |
| 296 | try { |
| 297 | socket.send(text); |
| 298 | } catch { |
| 299 | // Closing already. |
| 300 | } |
| 301 | } |
| 302 | } |
| 303 | |
| 304 | override async fetch(request: Request): Promise<Response> { |
| 305 | if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") { |
| 306 | return new Response("Expected a WebSocket upgrade\n", { status: 426 }); |
| 307 | } |
| 308 | let seed: FeedSeed | null = null; |
| 309 | try { |
| 310 | seed = JSON.parse(request.headers.get(NOTIFY_SEED_HEADER) ?? "null") as FeedSeed | null; |
| 311 | } catch { |
| 312 | seed = null; |
| 313 | } |
| 314 | if (seed) this.applySeed(seed); |
| 315 | const username = request.headers.get(FEED_USERNAME_HEADER); |
| 316 | const userId = request.headers.get(FEED_USER_ID_HEADER); |
| 317 | this.remember({ user_id: userId ?? "", username: username ?? "" }); |
| 318 | if (Array.isArray(seed?.workspaces)) this.setMeta("workspaces", JSON.stringify(cleanWorkspaces(seed.workspaces))); |
| 319 | const workspace = seed?.workspace?.toLowerCase() || null; |
| 320 | const pair = new WebSocketPair(); |
| 321 | const [client, server] = Object.values(pair) as [WebSocket, WebSocket]; |
| 322 | this.ctx.acceptWebSocket(server); |
| 323 | server.serializeAttachment({ focused: false, path: "", at: Date.now(), idle: false, workspace } satisfies Tab); |
| 324 | const hello: FeedEvent = { |
| 325 | type: "hello", |
| 326 | notifications: this.latest(HELLO), |
| 327 | vapid_public_key: this.env.VAPID_PUBLIC_KEY || null, |
| 328 | preferences: this.preferences(), |
| 329 | }; |
| 330 | server.send(JSON.stringify(hello)); |
| 331 | const shown = new Set(this.workspaces()); |
| 332 | if (seed?.workspace) shown.add(seed.workspace); |
| 333 | for (const workspace of shown) server.send(JSON.stringify(this.counts(workspace))); |
| 334 | // Presence after the socket is handed back: the rooms are not waited on. |
| 335 | this.ctx.waitUntil(this.connected(server, workspace)); |
| 336 | return new Response(null, { status: 101, webSocket: client }); |
| 337 | } |
| 338 | |
| 339 | /** Counts read from chat and the inbox just now replace what was kept for that workspace. */ |
| 340 | private applySeed(seed: FeedSeed): void { |
| 341 | if (typeof seed.inbox_unread === "number") this.setMeta("inbox_unread", String(Math.max(0, Math.floor(seed.inbox_unread)))); |
| 342 | const workspace = seed.workspace?.toLowerCase(); |
| 343 | if (!workspace || !Array.isArray(seed.per_channel)) return; |
| 344 | this.ctx.storage.transactionSync(() => { |
| 345 | this.sql.exec("DELETE FROM counts WHERE workspace = ?", workspace); |
| 346 | for (const row of seed.per_channel!.slice(0, 2000)) { |
| 347 | if (typeof row?.channel_id !== "string") continue; |
| 348 | const kept = applyCounts(null, { ...row, set: true }); |
| 349 | this.sql.exec( |
| 350 | "INSERT OR REPLACE INTO counts (workspace, channel_id, unread, mentions, muted) VALUES (?, ?, ?, ?, ?)", |
| 351 | workspace, |
| 352 | kept.channel_id, |
| 353 | kept.unread, |
| 354 | kept.mentions, |
| 355 | kept.muted ? 1 : 0, |
| 356 | ); |
| 357 | } |
| 358 | this.setMeta(`complete:${workspace}`, "1"); |
| 359 | }); |
| 360 | } |
| 361 | |
| 362 | override async webSocketMessage(socket: WebSocket, message: string | ArrayBuffer): Promise<void> { |
| 363 | if (typeof message !== "string" || message.length > 4096) return; |
| 364 | let frame: Record<string, unknown>; |
| 365 | try { |
| 366 | frame = JSON.parse(message); |
| 367 | } catch { |
| 368 | return; |
| 369 | } |
| 370 | if (frame.type === "state") { |
| 371 | const before = (socket.deserializeAttachment() as Tab | null) ?? { focused: false, path: "", at: 0 }; |
| 372 | const idle = frame.idle === true; |
| 373 | socket.serializeAttachment({ |
| 374 | focused: frame.focused === true, |
| 375 | path: String(frame.path ?? "").slice(0, 500), |
| 376 | at: Date.now(), |
| 377 | idle, |
| 378 | workspace: before.workspace ?? null, |
| 379 | } satisfies Tab); |
| 380 | // Gone idle, or back: away or active. |
| 381 | if (idle !== (before.idle === true)) await this.refresh(); |
| 382 | } else if (frame.type === "inbox" && Number.isFinite(frame.unread)) { |
| 383 | // The inbox count a page just read: every tab shows it at once. |
| 384 | await this.setInbox(Number(frame.unread)); |
| 385 | } |
| 386 | } |
| 387 | |
| 388 | override async webSocketClose(socket: WebSocket, code: number, reason: string): Promise<void> { |
| 389 | try { |
| 390 | socket.close(code, reason); |
| 391 | } catch { |
| 392 | // Already closed. |
| 393 | } |
| 394 | await this.left(socket); |
| 395 | } |
| 396 | |
| 397 | override async webSocketError(socket: WebSocket): Promise<void> { |
| 398 | await this.left(socket); |
| 399 | } |
| 400 | |
| 401 | /** |
| 402 | * A tab went. The last one leaves the person showing as they were for a |
| 403 | * moment, so a reload or a switch of workspace is not leaving; the alarm |
| 404 | * then says offline if no tab came back. |
| 405 | */ |
| 406 | private async left(socket: WebSocket): Promise<void> { |
| 407 | if (this.presenceTabs(socket).length === 0) { |
| 408 | this.setMeta("offline_at", String(Date.now() + OFFLINE_GRACE_MS)); |
| 409 | await this.schedule(); |
| 410 | return; |
| 411 | } |
| 412 | await this.refresh({ closing: socket }); |
| 413 | } |
| 414 | |
| 415 | // ── Notifications ─────────────────────────────────────────────────────── |
| 416 | |
| 417 | private latest(limit: number): FeedNotification[] { |
| 418 | return this.sql |
| 419 | .exec<{ json: string }>("SELECT json FROM notifications ORDER BY at DESC LIMIT ?", limit) |
| 420 | .toArray() |
| 421 | .map((r) => JSON.parse(r.json) as FeedNotification); |
| 422 | } |
| 423 | |
| 424 | /** Keeps it, unless it was told before. True when it is news. */ |
| 425 | private keep(notification: FeedNotification): boolean { |
| 426 | const had = this.sql.exec("SELECT 1 FROM notifications WHERE id = ?", notification.id).toArray().length > 0; |
| 427 | if (had) return false; |
| 428 | this.sql.exec("INSERT INTO notifications (id, at, json) VALUES (?, ?, ?)", notification.id, Date.now(), JSON.stringify(notification)); |
| 429 | this.sql.exec("DELETE FROM notifications WHERE id NOT IN (SELECT id FROM notifications ORDER BY at DESC LIMIT ?)", KEPT); |
| 430 | return true; |
| 431 | } |
| 432 | |
| 433 | /** Tells the person: open tabs at once, then browsers and phones when none is in front of them. */ |
| 434 | private async tell(notification: FeedNotification, test = false): Promise<number> { |
| 435 | if (!test && !this.keep(notification)) return 0; |
| 436 | // Somewhere to push: a browser or a phone. |
| 437 | const subscriptions = this.subscriptionCount() + this.deviceCount(); |
| 438 | const now = Date.now(); |
| 439 | // Do Not Disturb: kept and counted, neither toasted nor pushed. |
| 440 | const quiet = dndOn(this.kept().dnd_until, now); |
| 441 | const decision = decide({ prefs: this.preferences(), notification, tabs: this.tabs(), subscriptions, now, test, dnd: quiet }); |
| 442 | this.send({ type: "notification", notification, toast: decision.toast }); |
| 443 | return decision.push ? this.push(notification) : 0; |
| 444 | } |
| 445 | |
| 446 | /** Pushes to every browser and phone at once; one failing does not stop the other. */ |
| 447 | private async push(notification: FeedNotification): Promise<number> { |
| 448 | const payload = pushPayload(notification); |
| 449 | const results = await Promise.allSettled([this.pushBrowsers(payload), this.pushDevices(notification, payload)]); |
| 450 | let sent = 0; |
| 451 | for (const result of results) { |
| 452 | if (result.status === "fulfilled") sent += result.value; |
| 453 | else console.error("notify: pushing failed", result.reason); |
| 454 | } |
| 455 | return sent; |
| 456 | } |
| 457 | |
| 458 | /** Pushes to every browser subscribed; drops the ones the push service says are gone. */ |
| 459 | private async pushBrowsers(payload: PushPayload): Promise<number> { |
| 460 | const vapid = this.vapid(); |
| 461 | if (!vapid) return 0; |
| 462 | const subscriptions = this.sql.exec<{ endpoint: string; p256dh: string; auth: string }>("SELECT endpoint, p256dh, auth FROM subscriptions").toArray(); |
| 463 | const results = await Promise.allSettled( |
| 464 | subscriptions.map((s) => sendPush(s, payload, { vapid, urgency: payload.urgent ? "high" : "normal", topic: payload.tag, ttl: 12 * 3600 })), |
| 465 | ); |
| 466 | let sent = 0; |
| 467 | for (const result of results) { |
| 468 | if (result.status === "rejected") { |
| 469 | console.error("notify: a push failed", result.reason); |
| 470 | continue; |
| 471 | } |
| 472 | if (result.value.gone) this.sql.exec("DELETE FROM subscriptions WHERE endpoint = ?", result.value.endpoint); |
| 473 | else if (result.value.status < 300) sent++; |
| 474 | else console.error("notify: a push service answered", result.value.status); |
| 475 | } |
| 476 | return sent; |
| 477 | } |
| 478 | |
| 479 | /** Pushes to every phone through Expo; drops the ones Expo says the app is gone from. */ |
| 480 | private async pushDevices(notification: FeedNotification, payload: PushPayload): Promise<number> { |
| 481 | if (!this.phonePush()) return 0; |
| 482 | const tokens = this.sql.exec<{ token: string }>("SELECT token FROM devices").toArray().map((r) => r.token); |
| 483 | if (!tokens.length) return 0; |
| 484 | const result = await sendExpo(tokens, payload, notification, { accessToken: this.env.EXPO_ACCESS_TOKEN }); |
| 485 | for (const token of result.gone) this.sql.exec("DELETE FROM devices WHERE token = ?", token); |
| 486 | for (const error of result.errors) console.error("notify: a push to a phone failed", error); |
| 487 | return result.sent; |
| 488 | } |
| 489 | |
| 490 | // ── RPC (called by the Worker, src/index.ts) ──────────────────────────── |
| 491 | |
| 492 | async notify(value: unknown): Promise<{ ok: boolean }> { |
| 493 | const notification = cleanNotification(value); |
| 494 | if (!notification) return { ok: false }; |
| 495 | await this.tell(notification); |
| 496 | return { ok: true }; |
| 497 | } |
| 498 | |
| 499 | /** A batch for this person from chat: counts moved, notifications told. */ |
| 500 | async deliver(items: FeedDelivery[]): Promise<{ ok: boolean }> { |
| 501 | const changed = new Set<string>(); |
| 502 | const told: FeedNotification[] = []; |
| 503 | for (const item of items) { |
| 504 | const workspace = String(item.workspace ?? "").toLowerCase(); |
| 505 | if (item.counts?.channel_id && workspace) { |
| 506 | const before = this.sql |
| 507 | .exec<{ unread: number; mentions: number; muted: number }>( |
| 508 | "SELECT unread, mentions, muted FROM counts WHERE workspace = ? AND channel_id = ?", |
| 509 | workspace, |
| 510 | item.counts.channel_id, |
| 511 | ) |
| 512 | .toArray()[0]; |
| 513 | const after = applyCounts( |
| 514 | before ? { channel_id: item.counts.channel_id, unread: before.unread, mentions: before.mentions, muted: !!before.muted } : null, |
| 515 | item.counts, |
| 516 | ); |
| 517 | this.sql.exec( |
| 518 | "INSERT OR REPLACE INTO counts (workspace, channel_id, unread, mentions, muted) VALUES (?, ?, ?, ?, ?)", |
| 519 | workspace, |
| 520 | after.channel_id, |
| 521 | after.unread, |
| 522 | after.mentions, |
| 523 | after.muted ? 1 : 0, |
| 524 | ); |
| 525 | changed.add(workspace); |
| 526 | } |
| 527 | const notification = item.notification ? cleanNotification(item.notification) : null; |
| 528 | if (notification) told.push(notification); |
| 529 | } |
| 530 | for (const workspace of changed) this.send(this.counts(workspace)); |
| 531 | await Promise.all(told.map((n) => this.tell(n))); |
| 532 | return { ok: true }; |
| 533 | } |
| 534 | |
| 535 | /** The inbox count as it now is (events, after items arrive or are marked): every tab shows it at once. */ |
| 536 | async setInbox(value: number): Promise<{ ok: boolean }> { |
| 537 | const unread = Number.isFinite(value) ? Math.max(0, Math.floor(value)) : 0; |
| 538 | if (unread === this.inboxUnread()) return { ok: true }; |
| 539 | this.setMeta("inbox_unread", String(unread)); |
| 540 | this.send({ type: "inbox", unread }); |
| 541 | return { ok: true }; |
| 542 | } |
| 543 | |
| 544 | async subscribe(subscription: PushSubscriptionJson, userAgent: string | null): Promise<{ ok: boolean }> { |
| 545 | const endpoint = String(subscription?.endpoint ?? ""); |
| 546 | const { p256dh, auth } = subscription?.keys ?? ({} as PushSubscriptionJson["keys"]); |
| 547 | if (!/^https:\/\//.test(endpoint) || endpoint.length > MAX_ENDPOINT || typeof p256dh !== "string" || typeof auth !== "string") return { ok: false }; |
| 548 | this.sql.exec( |
| 549 | "INSERT OR REPLACE INTO subscriptions (endpoint, p256dh, auth, user_agent, created_at) VALUES (?, ?, ?, ?, ?)", |
| 550 | endpoint, |
| 551 | p256dh.slice(0, 200), |
| 552 | auth.slice(0, 100), |
| 553 | userAgent ? userAgent.slice(0, 300) : null, |
| 554 | Date.now(), |
| 555 | ); |
| 556 | // The oldest go first past the limit. |
| 557 | this.sql.exec("DELETE FROM subscriptions WHERE endpoint NOT IN (SELECT endpoint FROM subscriptions ORDER BY created_at DESC LIMIT ?)", MAX_SUBSCRIPTIONS); |
| 558 | return { ok: true }; |
| 559 | } |
| 560 | |
| 561 | async unsubscribe(endpoint: string): Promise<{ ok: boolean }> { |
| 562 | this.sql.exec("DELETE FROM subscriptions WHERE endpoint = ?", String(endpoint ?? "")); |
| 563 | return { ok: true }; |
| 564 | } |
| 565 | |
| 566 | /** A phone the app runs on, by its Expo push token; registering it again moves it to the newest. */ |
| 567 | async registerDevice(device: PushDevice): Promise<{ ok: boolean; off?: boolean }> { |
| 568 | if (!this.phonePush()) return { ok: false, off: true }; |
| 569 | const clean = cleanDevice(device); |
| 570 | if (!clean) return { ok: false }; |
| 571 | this.sql.exec( |
| 572 | "INSERT OR REPLACE INTO devices (token, platform, name, created_at) VALUES (?, ?, ?, ?)", |
| 573 | clean.token, |
| 574 | clean.platform, |
| 575 | clean.name, |
| 576 | Date.now(), |
| 577 | ); |
| 578 | // The oldest go first past the limit. |
| 579 | this.sql.exec("DELETE FROM devices WHERE token NOT IN (SELECT token FROM devices ORDER BY created_at DESC LIMIT ?)", MAX_DEVICES); |
| 580 | return { ok: true }; |
| 581 | } |
| 582 | |
| 583 | async unregisterDevice(token: string): Promise<{ ok: boolean }> { |
| 584 | this.sql.exec("DELETE FROM devices WHERE token = ?", String(token ?? "")); |
| 585 | return { ok: true }; |
| 586 | } |
| 587 | |
| 588 | async status(endpoint: string | null): Promise<NotifyStatus> { |
| 589 | const subscribed = !!endpoint && this.sql.exec("SELECT 1 FROM subscriptions WHERE endpoint = ?", endpoint).toArray().length > 0; |
| 590 | return { |
| 591 | preferences: this.preferences(), |
| 592 | subscriptions: this.subscriptionCount(), |
| 593 | devices: this.deviceCount(), |
| 594 | subscribed, |
| 595 | vapid_public_key: this.env.VAPID_PUBLIC_KEY || null, |
| 596 | }; |
| 597 | } |
| 598 | |
| 599 | async setPreferences(change: unknown): Promise<NotifyPreferences> { |
| 600 | const preferences = mergePreferences(this.preferences(), change); |
| 601 | this.setMeta("preferences", JSON.stringify(preferences)); |
| 602 | this.send({ type: "preferences", preferences }); |
| 603 | return preferences; |
| 604 | } |
| 605 | |
| 606 | /** Tabs open in `workspace` hear how these people now show (from its room). */ |
| 607 | async presence(workspace: string, people: PresenceEntry[]): Promise<void> { |
| 608 | const text = JSON.stringify({ type: "presence", workspace, people, full: false } satisfies FeedEvent); |
| 609 | for (const socket of this.ctx.getWebSockets()) { |
| 610 | const tab = socket.deserializeAttachment() as Tab | null; |
| 611 | if (tab?.workspace !== workspace) continue; |
| 612 | try { |
| 613 | socket.send(text); |
| 614 | } catch { |
| 615 | // Closing already. |
| 616 | } |
| 617 | } |
| 618 | } |
| 619 | |
| 620 | /** The person's own: presence, status, Do Not Disturb. */ |
| 621 | async own(person: { user_id: string; username: string }): Promise<OwnPresence> { |
| 622 | this.remember(person); |
| 623 | const kept = current(this.kept(), Date.now()); |
| 624 | return ownOf(this.entry(), kept); |
| 625 | } |
| 626 | |
| 627 | /** A change to the person's own, told to their tabs and to every workspace they are in. */ |
| 628 | async setPresence(person: { user_id: string; username: string }, change: PresenceChange): Promise<OwnPresence> { |
| 629 | this.remember(person); |
| 630 | const kept = applyChange(this.kept(), change, Date.now()); |
| 631 | this.setMeta("presence", JSON.stringify(kept)); |
| 632 | const entry = await this.refresh(); |
| 633 | return ownOf(entry, current(this.kept(), Date.now())); |
| 634 | } |
| 635 | |
| 636 | private remember(person: { user_id: string; username: string }): void { |
| 637 | if (person.user_id) this.setMeta("user_id", person.user_id.slice(0, 100)); |
| 638 | if (person.username) this.setMeta("username", person.username.slice(0, 100)); |
| 639 | } |
| 640 | |
| 641 | private visit(workspace: string): Visit | null { |
| 642 | return this.sql.exec<Visit>("SELECT seen_at, previous_at FROM visits WHERE workspace = ?", workspace).toArray()[0] ?? null; |
| 643 | } |
| 644 | |
| 645 | /** When the person was last on `workspace`'s Home page, as it counts from (src/visits.ts); null when never. */ |
| 646 | async lastVisit(workspace: string): Promise<{ seen_at: string | null }> { |
| 647 | const since = lastVisit(this.visit(workspace), Date.now()); |
| 648 | return { seen_at: since != null ? new Date(since).toISOString() : null }; |
| 649 | } |
| 650 | |
| 651 | /** Marks the person's visit to `workspace`'s Home page at `at`; it only moves forward. */ |
| 652 | async markVisit(workspace: string, at: unknown): Promise<{ seen_at: string | null }> { |
| 653 | const now = Date.now(); |
| 654 | const next = markVisit(this.visit(workspace), at, now); |
| 655 | if (next) { |
| 656 | this.sql.exec( |
| 657 | "INSERT INTO visits (workspace, seen_at, previous_at) VALUES (?, ?, ?) ON CONFLICT (workspace) DO UPDATE SET seen_at = excluded.seen_at, previous_at = excluded.previous_at", |
| 658 | workspace, |
| 659 | next.seen_at, |
| 660 | next.previous_at, |
| 661 | ); |
| 662 | } |
| 663 | return this.lastVisit(workspace); |
| 664 | } |
| 665 | |
| 666 | async test(username: string): Promise<{ ok: boolean; pushed: number }> { |
| 667 | const pushed = await this.tell( |
| 668 | { |
| 669 | id: `test:${crypto.randomUUID()}`, |
| 670 | kind: "dm", |
| 671 | workspace: "", |
| 672 | title: "g1t", |
| 673 | body: `Notifications are on, ${username || "there"}. This is what a message looks like.`, |
| 674 | href: "/settings/notifications", |
| 675 | actor: { kind: "system", id: "g1t", name: "g1t", avatar: null, avatar_seed: null }, |
| 676 | channel_id: null, |
| 677 | thread_root: null, |
| 678 | created_at: new Date().toISOString(), |
| 679 | }, |
| 680 | true, |
| 681 | ); |
| 682 | return { ok: true, pushed }; |
| 683 | } |
| 684 | } |