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 |
| The front page is Home: where you're needed comes first, however long ago it started, with what's at stake and how long it has waited; then what happened since you were last here, in words, with the agent work that settled, what landed, what deployed, what was decided and what was spent, switchable to the last 24 hours or 7 days; then what's running now. Your last visit is kept with your account, so it's the same on every device; Today's old address leads here. The home guide says how. | 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). | |
| 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 | 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. | |
| 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) | 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, | |
| 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 | 29 | type OwnPresence, |
| 30 | type PresenceChange, | |
| 31 | 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) | 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"; | |
| 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 | 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"; | |
| The front page is Home: where you're needed comes first, however long ago it started, with what's at stake and how long it has waited; then what happened since you were last here, in words, with the agent work that settled, what landed, what deployed, what was decided and what was spent, switchable to the last 24 hours or 7 days; then what's running now. Your last visit is kept with your account, so it's the same on every device; Today's old address leads here. The home guide says how. | 58 | import { type Visit, lastVisit, markVisit } from "./visits.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) | 59 | import { sendPush, type Vapid } from "./webpush.ts"; |
| 60 | ||
| 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 | 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"; | |
| 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) | 72 | |
| 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 | 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 }; | |
| 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) | 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)`); | |
| The front page is Home: where you're needed comes first, however long ago it started, with what's at stake and how long it has waited; then what happened since you were last here, in words, with the agent work that settled, what landed, what deployed, what was decided and what was spent, switchable to the last 24 hours or 7 days; then what's running now. Your last visit is kept with your account, so it's the same on every device; Today's old address leads here. The home guide says how. | 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)`); | |
| 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) | 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 | ||
| 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 | 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 | ||
| 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) | 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 | ||
| 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 | 261 | 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) | 262 | const text = JSON.stringify(event); |
| 263 | 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 | 264 | 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) | 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); | |
| 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 | 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; | |
| 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 pair = new WebSocketPair(); |
| 290 | const [client, server] = Object.values(pair) as [WebSocket, WebSocket]; | |
| 291 | 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 | 292 | 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) | 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))); | |
| 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 | 303 | // Presence after the socket is handed back: the rooms are not waited on. |
| 304 | 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) | 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") { | |
| 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 | 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(); | |
| 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) | 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 | } | |
| 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 | 363 | await this.left(socket); |
| 364 | } | |
| 365 | ||
| 366 | override async webSocketError(socket: WebSocket): Promise<void> { | |
| 367 | 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) | 368 | } |
| 369 | ||
| 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 | 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 | } | |
| 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) | 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(); | |
| 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 | 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 }); | |
| 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) | 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 | ||
| 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 | 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 | ||
| The front page is Home: where you're needed comes first, however long ago it started, with what's at stake and how long it has waited; then what happened since you were last here, in words, with the agent work that settled, what landed, what deployed, what was decided and what was spent, switchable to the last 24 hours or 7 days; then what's running now. Your last visit is kept with your account, so it's the same on every device; Today's old address leads here. The home guide says how. | 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 | ||
| 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) | 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 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.