Skip to content
602 linesCodeBlameRaw

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 managers12 * 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 managers15 *
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 */
24import { DurableObject } from "cloudflare:workers";
25
26import {
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 managers29 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
41import { applyCounts, totals } from "./counts.ts";
42import { 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 managers43import {
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";
57import 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.58import { 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)59import { 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 managers61export 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). */
70export const FEED_USERNAME_HEADER = "x-g1t-notify-username";
71export 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 managers73/** What each socket carries: focus, page, whether it has gone idle, and the workspace it is open in. */
74type 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. */
77const KEPT = 100;
78/** Sent with `hello`: the latest few. */
79const HELLO = 20;
80/** The most browsers one person gets pushes on. */
81const MAX_SUBSCRIPTIONS = 20;
82/** The longest endpoint taken. */
83const MAX_ENDPOINT = 2000;
84
85export 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 managers152 // ── 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 managers261 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 managers264 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 managers284 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 managers292 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 managers303 // 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 managers340 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 managers363 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 managers370 /**
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 managers406 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 managers524 /** 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.