Skip to content
367 linesCodeBlameRaw
1/**
2 * One person's feed, as a Durable Object named by their user id.
3 *
4 * It holds a socket per open tab with the WebSocket Hibernation API, so a
5 * person with g1t open all day costs nothing between notifications. Plain
6 * `ping` keepalives are answered at the edge without waking it, and the
7 * time of the last one says whether a tab is still there
8 * (`getWebSocketAutoResponseTimestamp`). Each socket carries the tab's
9 * focus and page (`serializeAttachment`), which survive hibernation.
10 *
11 * Kept in the object's own SQLite storage: the latest notifications, the
12 * unread counts per conversation, push subscriptions and preferences.
13 *
14 * The feed authorizes nothing: the Worker reaches it only for the person
15 * the site checked (src/index.ts).
16 */
17import { DurableObject } from "cloudflare:workers";
18
19import {
20 NOTIFY_SEED_HEADER,
21 type ChannelCounts,
22 type FeedDelivery,
23 type FeedEvent,
24 type FeedNotification,
25 type FeedSeed,
26 type NotifyPreferences,
27 type NotifyStatus,
28 type PushSubscriptionJson,
29} from "@g1t/contracts";
30
31import { applyCounts, totals } from "./counts.ts";
32import { cleanNotification, decide, mergePreferences, pushPayload, readPreferences, type TabState } from "./prefs.ts";
33import { sendPush, type Vapid } from "./webpush.ts";
34
35export type FeedEnv = { VAPID_PUBLIC_KEY?: string; VAPID_PRIVATE_KEY?: string; VAPID_SUBJECT?: string };
36
37/** What each socket carries. */
38type Tab = { focused: boolean; path: string; at: number };
39
40/** Notifications kept for a tab that opens later. */
41const KEPT = 100;
42/** Sent with `hello`: the latest few. */
43const HELLO = 20;
44/** The most browsers one person gets pushes on. */
45const MAX_SUBSCRIPTIONS = 20;
46/** The longest endpoint taken. */
47const MAX_ENDPOINT = 2000;
48
49export class Feed extends DurableObject<FeedEnv> {
50 private readonly sql: SqlStorage;
51
52 constructor(ctx: DurableObjectState, env: FeedEnv) {
53 super(ctx, env);
54 this.sql = ctx.storage.sql;
55 // Keepalives are answered without waking the object.
56 ctx.setWebSocketAutoResponse(new WebSocketRequestResponsePair("ping", "pong"));
57 this.sql.exec(`CREATE TABLE IF NOT EXISTS notifications (id TEXT PRIMARY KEY, at INTEGER NOT NULL, json TEXT NOT NULL)`);
58 this.sql.exec(`CREATE INDEX IF NOT EXISTS notifications_at ON notifications (at)`);
59 this.sql.exec(
60 `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))`,
61 );
62 this.sql.exec(
63 `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)`,
64 );
65 this.sql.exec(`CREATE TABLE IF NOT EXISTS meta (key TEXT PRIMARY KEY, value TEXT NOT NULL)`);
66 }
67
68 // ── Kept state ──────────────────────────────────────────────────────────
69
70 private meta(key: string): string | null {
71 const row = this.sql.exec<{ value: string }>("SELECT value FROM meta WHERE key = ?", key).toArray()[0];
72 return row?.value ?? null;
73 }
74
75 private setMeta(key: string, value: string): void {
76 this.sql.exec("INSERT INTO meta (key, value) VALUES (?, ?) ON CONFLICT (key) DO UPDATE SET value = excluded.value", key, value);
77 }
78
79 private preferences(): NotifyPreferences {
80 return readPreferences(this.meta("preferences"));
81 }
82
83 private inboxUnread(): number {
84 return Number(this.meta("inbox_unread") ?? 0) || 0;
85 }
86
87 private channelRows(workspace: string): ChannelCounts[] {
88 return this.sql
89 .exec<{ channel_id: string; unread: number; mentions: number; muted: number }>(
90 "SELECT channel_id, unread, mentions, muted FROM counts WHERE workspace = ?",
91 workspace,
92 )
93 .toArray()
94 .map((r) => ({ channel_id: r.channel_id, unread: r.unread, mentions: r.mentions, muted: !!r.muted }));
95 }
96
97 private counts(workspace: string): FeedEvent {
98 return { type: "counts", ...totals(workspace, this.channelRows(workspace), this.inboxUnread(), this.meta(`complete:${workspace}`) === "1") };
99 }
100
101 private workspaces(): string[] {
102 return this.sql.exec<{ workspace: string }>("SELECT DISTINCT workspace FROM counts").toArray().map((r) => r.workspace);
103 }
104
105 private vapid(): Vapid | null {
106 const { VAPID_PUBLIC_KEY, VAPID_PRIVATE_KEY, VAPID_SUBJECT } = this.env;
107 return VAPID_PUBLIC_KEY && VAPID_PRIVATE_KEY ? { publicKey: VAPID_PUBLIC_KEY, privateKey: VAPID_PRIVATE_KEY, subject: VAPID_SUBJECT || "https://g1t.sh" } : null;
108 }
109
110 private subscriptionCount(): number {
111 return this.sql.exec<{ n: number }>("SELECT COUNT(*) AS n FROM subscriptions").one().n;
112 }
113
114 // ── Sockets ─────────────────────────────────────────────────────────────
115
116 private tabs(): TabState[] {
117 return this.ctx.getWebSockets().map((socket) => {
118 const tab = (socket.deserializeAttachment() as Tab | null) ?? { focused: false, path: "", at: 0 };
119 const pinged = this.ctx.getWebSocketAutoResponseTimestamp(socket)?.getTime() ?? 0;
120 return { focused: tab.focused, seen_at: Math.max(tab.at, pinged) };
121 });
122 }
123
124 private send(event: FeedEvent): void {
125 const text = JSON.stringify(event);
126 for (const socket of this.ctx.getWebSockets()) {
127 try {
128 socket.send(text);
129 } catch {
130 // Closing already.
131 }
132 }
133 }
134
135 override async fetch(request: Request): Promise<Response> {
136 if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") {
137 return new Response("Expected a WebSocket upgrade\n", { status: 426 });
138 }
139 let seed: FeedSeed | null = null;
140 try {
141 seed = JSON.parse(request.headers.get(NOTIFY_SEED_HEADER) ?? "null") as FeedSeed | null;
142 } catch {
143 seed = null;
144 }
145 if (seed) this.applySeed(seed);
146 const pair = new WebSocketPair();
147 const [client, server] = Object.values(pair) as [WebSocket, WebSocket];
148 this.ctx.acceptWebSocket(server);
149 server.serializeAttachment({ focused: false, path: "", at: Date.now() } satisfies Tab);
150 const hello: FeedEvent = {
151 type: "hello",
152 notifications: this.latest(HELLO),
153 vapid_public_key: this.env.VAPID_PUBLIC_KEY || null,
154 preferences: this.preferences(),
155 };
156 server.send(JSON.stringify(hello));
157 const shown = new Set(this.workspaces());
158 if (seed?.workspace) shown.add(seed.workspace);
159 for (const workspace of shown) server.send(JSON.stringify(this.counts(workspace)));
160 return new Response(null, { status: 101, webSocket: client });
161 }
162
163 /** Counts read from chat and the inbox just now replace what was kept for that workspace. */
164 private applySeed(seed: FeedSeed): void {
165 if (typeof seed.inbox_unread === "number") this.setMeta("inbox_unread", String(Math.max(0, Math.floor(seed.inbox_unread))));
166 const workspace = seed.workspace?.toLowerCase();
167 if (!workspace || !Array.isArray(seed.per_channel)) return;
168 this.ctx.storage.transactionSync(() => {
169 this.sql.exec("DELETE FROM counts WHERE workspace = ?", workspace);
170 for (const row of seed.per_channel!.slice(0, 2000)) {
171 if (typeof row?.channel_id !== "string") continue;
172 const kept = applyCounts(null, { ...row, set: true });
173 this.sql.exec(
174 "INSERT OR REPLACE INTO counts (workspace, channel_id, unread, mentions, muted) VALUES (?, ?, ?, ?, ?)",
175 workspace,
176 kept.channel_id,
177 kept.unread,
178 kept.mentions,
179 kept.muted ? 1 : 0,
180 );
181 }
182 this.setMeta(`complete:${workspace}`, "1");
183 });
184 }
185
186 override async webSocketMessage(socket: WebSocket, message: string | ArrayBuffer): Promise<void> {
187 if (typeof message !== "string" || message.length > 4096) return;
188 let frame: Record<string, unknown>;
189 try {
190 frame = JSON.parse(message);
191 } catch {
192 return;
193 }
194 if (frame.type === "state") {
195 socket.serializeAttachment({ focused: frame.focused === true, path: String(frame.path ?? "").slice(0, 500), at: Date.now() } satisfies Tab);
196 } else if (frame.type === "inbox" && Number.isFinite(frame.unread)) {
197 // The inbox count a page just read: every tab shows it at once.
198 await this.setInbox(Number(frame.unread));
199 }
200 }
201
202 override async webSocketClose(socket: WebSocket, code: number, reason: string): Promise<void> {
203 try {
204 socket.close(code, reason);
205 } catch {
206 // Already closed.
207 }
208 }
209
210 override async webSocketError(): Promise<void> {}
211
212 // ── Notifications ───────────────────────────────────────────────────────
213
214 private latest(limit: number): FeedNotification[] {
215 return this.sql
216 .exec<{ json: string }>("SELECT json FROM notifications ORDER BY at DESC LIMIT ?", limit)
217 .toArray()
218 .map((r) => JSON.parse(r.json) as FeedNotification);
219 }
220
221 /** Keeps it, unless it was told before. True when it is news. */
222 private keep(notification: FeedNotification): boolean {
223 const had = this.sql.exec("SELECT 1 FROM notifications WHERE id = ?", notification.id).toArray().length > 0;
224 if (had) return false;
225 this.sql.exec("INSERT INTO notifications (id, at, json) VALUES (?, ?, ?)", notification.id, Date.now(), JSON.stringify(notification));
226 this.sql.exec("DELETE FROM notifications WHERE id NOT IN (SELECT id FROM notifications ORDER BY at DESC LIMIT ?)", KEPT);
227 return true;
228 }
229
230 /** Tells the person: open tabs at once, then browsers when none is in front of them. */
231 private async tell(notification: FeedNotification, test = false): Promise<number> {
232 if (!test && !this.keep(notification)) return 0;
233 const subscriptions = this.subscriptionCount();
234 const decision = decide({ prefs: this.preferences(), notification, tabs: this.tabs(), subscriptions, now: Date.now(), test });
235 this.send({ type: "notification", notification, toast: decision.toast });
236 return decision.push ? this.push(notification) : 0;
237 }
238
239 /** Pushes to every browser subscribed; drops the ones the push service says are gone. */
240 private async push(notification: FeedNotification): Promise<number> {
241 const vapid = this.vapid();
242 if (!vapid) return 0;
243 const subscriptions = this.sql.exec<{ endpoint: string; p256dh: string; auth: string }>("SELECT endpoint, p256dh, auth FROM subscriptions").toArray();
244 const payload = pushPayload(notification);
245 const results = await Promise.allSettled(
246 subscriptions.map((s) => sendPush(s, payload, { vapid, urgency: payload.urgent ? "high" : "normal", topic: payload.tag, ttl: 12 * 3600 })),
247 );
248 let sent = 0;
249 for (const result of results) {
250 if (result.status === "rejected") {
251 console.error("notify: a push failed", result.reason);
252 continue;
253 }
254 if (result.value.gone) this.sql.exec("DELETE FROM subscriptions WHERE endpoint = ?", result.value.endpoint);
255 else if (result.value.status < 300) sent++;
256 else console.error("notify: a push service answered", result.value.status);
257 }
258 return sent;
259 }
260
261 // ── RPC (called by the Worker, src/index.ts) ────────────────────────────
262
263 async notify(value: unknown): Promise<{ ok: boolean }> {
264 const notification = cleanNotification(value);
265 if (!notification) return { ok: false };
266 await this.tell(notification);
267 return { ok: true };
268 }
269
270 /** A batch for this person from chat: counts moved, notifications told. */
271 async deliver(items: FeedDelivery[]): Promise<{ ok: boolean }> {
272 const changed = new Set<string>();
273 const told: FeedNotification[] = [];
274 for (const item of items) {
275 const workspace = String(item.workspace ?? "").toLowerCase();
276 if (item.counts?.channel_id && workspace) {
277 const before = this.sql
278 .exec<{ unread: number; mentions: number; muted: number }>(
279 "SELECT unread, mentions, muted FROM counts WHERE workspace = ? AND channel_id = ?",
280 workspace,
281 item.counts.channel_id,
282 )
283 .toArray()[0];
284 const after = applyCounts(
285 before ? { channel_id: item.counts.channel_id, unread: before.unread, mentions: before.mentions, muted: !!before.muted } : null,
286 item.counts,
287 );
288 this.sql.exec(
289 "INSERT OR REPLACE INTO counts (workspace, channel_id, unread, mentions, muted) VALUES (?, ?, ?, ?, ?)",
290 workspace,
291 after.channel_id,
292 after.unread,
293 after.mentions,
294 after.muted ? 1 : 0,
295 );
296 changed.add(workspace);
297 }
298 const notification = item.notification ? cleanNotification(item.notification) : null;
299 if (notification) told.push(notification);
300 }
301 for (const workspace of changed) this.send(this.counts(workspace));
302 await Promise.all(told.map((n) => this.tell(n)));
303 return { ok: true };
304 }
305
306 /** The inbox count as it now is (events, after items arrive or are marked): every tab shows it at once. */
307 async setInbox(value: number): Promise<{ ok: boolean }> {
308 const unread = Number.isFinite(value) ? Math.max(0, Math.floor(value)) : 0;
309 if (unread === this.inboxUnread()) return { ok: true };
310 this.setMeta("inbox_unread", String(unread));
311 this.send({ type: "inbox", unread });
312 return { ok: true };
313 }
314
315 async subscribe(subscription: PushSubscriptionJson, userAgent: string | null): Promise<{ ok: boolean }> {
316 const endpoint = String(subscription?.endpoint ?? "");
317 const { p256dh, auth } = subscription?.keys ?? ({} as PushSubscriptionJson["keys"]);
318 if (!/^https:\/\//.test(endpoint) || endpoint.length > MAX_ENDPOINT || typeof p256dh !== "string" || typeof auth !== "string") return { ok: false };
319 this.sql.exec(
320 "INSERT OR REPLACE INTO subscriptions (endpoint, p256dh, auth, user_agent, created_at) VALUES (?, ?, ?, ?, ?)",
321 endpoint,
322 p256dh.slice(0, 200),
323 auth.slice(0, 100),
324 userAgent ? userAgent.slice(0, 300) : null,
325 Date.now(),
326 );
327 // The oldest go first past the limit.
328 this.sql.exec("DELETE FROM subscriptions WHERE endpoint NOT IN (SELECT endpoint FROM subscriptions ORDER BY created_at DESC LIMIT ?)", MAX_SUBSCRIPTIONS);
329 return { ok: true };
330 }
331
332 async unsubscribe(endpoint: string): Promise<{ ok: boolean }> {
333 this.sql.exec("DELETE FROM subscriptions WHERE endpoint = ?", String(endpoint ?? ""));
334 return { ok: true };
335 }
336
337 async status(endpoint: string | null): Promise<NotifyStatus> {
338 const subscribed = !!endpoint && this.sql.exec("SELECT 1 FROM subscriptions WHERE endpoint = ?", endpoint).toArray().length > 0;
339 return { preferences: this.preferences(), subscriptions: this.subscriptionCount(), subscribed, vapid_public_key: this.env.VAPID_PUBLIC_KEY || null };
340 }
341
342 async setPreferences(change: unknown): Promise<NotifyPreferences> {
343 const preferences = mergePreferences(this.preferences(), change);
344 this.setMeta("preferences", JSON.stringify(preferences));
345 this.send({ type: "preferences", preferences });
346 return preferences;
347 }
348
349 async test(username: string): Promise<{ ok: boolean; pushed: number }> {
350 const pushed = await this.tell(
351 {
352 id: `test:${crypto.randomUUID()}`,
353 kind: "dm",
354 workspace: "",
355 title: "g1t",
356 body: `Notifications are on, ${username || "there"}. This is what a message looks like.`,
357 href: "/settings/notifications",
358 actor: { kind: "system", id: "g1t", name: "g1t", avatar: null, avatar_seed: null },
359 channel_id: null,
360 thread_root: null,
361 created_at: new Date().toISOString(),
362 },
363 true,
364 );
365 return { ok: true, pushed };
366 }
367}