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