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