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 | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 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. | |
| 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) | 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, | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 28 | type OwnPresence, |
| 29 | type PresenceChange, | |
| 30 | type PresenceEntry, | |
| 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) | 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"; | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 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"; | |
| 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) | 57 | import { sendPush, type Vapid } from "./webpush.ts"; |
| 58 | ||
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 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"; | |
| 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) | 70 | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 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 }; | |
| 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) | 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 | ||
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 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 | ||
| 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) | 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 | ||
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 257 | private send(event: FeedEvent, except?: WebSocket): void { |
| 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) | 258 | const text = JSON.stringify(event); |
| 259 | for (const socket of this.ctx.getWebSockets()) { | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 260 | if (socket === except) continue; |
| 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) | 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); | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 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; | |
| 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) | 285 | const pair = new WebSocketPair(); |
| 286 | const [client, server] = Object.values(pair) as [WebSocket, WebSocket]; | |
| 287 | this.ctx.acceptWebSocket(server); | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 288 | server.serializeAttachment({ focused: false, path: "", at: Date.now(), idle: false, workspace } satisfies Tab); |
| 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) | 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))); | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 299 | // Presence after the socket is handed back: the rooms are not waited on. |
| 300 | this.ctx.waitUntil(this.connected(server, workspace)); | |
| 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) | 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") { | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 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(); | |
| 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) | 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 | } | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 359 | await this.left(socket); |
| 360 | } | |
| 361 | ||
| 362 | override async webSocketError(socket: WebSocket): Promise<void> { | |
| 363 | await this.left(socket); | |
| 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) | 364 | } |
| 365 | ||
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 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 | } | |
| 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) | 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(); | |
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 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 }); | |
| 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) | 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 | ||
| One kind of access token; presence and status; usernames keep their case; the tour is a miniature of the real app; icons for password managers | 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 | ||
| 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) | 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 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.