Skip to content
119 linesCodeBlameRaw
1/**
2 * The notify service: live notifications, unread counts and browser push.
3 * One feed per person (src/feed.ts). Plan: docs/WORKSPACE.md, "Live
4 * notifications".
5 *
6 * Reached through service bindings only: `POST /rpc/<method>` with
7 * snake_case bodies (`notifyClient` in @g1t/contracts), and `GET /live`,
8 * the person's feed socket, which the site forwards after checking the
9 * session (NOTIFY_VIEWER_HEADER), with the counts it read (NOTIFY_SEED_HEADER).
10 *
11 * Who sends what:
12 * - chat: `deliver`, for every message (counts for each person in the
13 * conversation, a notification for those it is for) and every read;
14 * - events: `notify`, for every inbox item, and `set_inbox`, the count after
15 * items arrive or are marked, naming the person by username;
16 * - the site: `subscribe`, `unsubscribe`, `status`, `set_preferences` and
17 * `test`, for the person signed in.
18 */
19import {
20 NOTIFY_VIEWER_HEADER,
21 identityClient,
22 type FeedDelivery,
23 type ServiceBinding,
24 type Viewer,
25} from "@g1t/contracts";
26
27import { Feed, type FeedEnv } from "./feed.ts";
28
29export { Feed } from "./feed.ts";
30
31type Env = FeedEnv & {
32 IDENTITY: ServiceBinding;
33 FEEDS: DurableObjectNamespace<Feed>;
34};
35
36function feed(env: Env, userId: string) {
37 return env.FEEDS.get(env.FEEDS.idFromName(userId));
38}
39
40/** A user id from the arguments, or one looked up by username. */
41async function userIdOf(env: Env, args: { user_id?: unknown; username?: unknown }): Promise<string | null> {
42 if (typeof args.user_id === "string" && args.user_id) return args.user_id;
43 if (typeof args.username !== "string" || !args.username) return null;
44 const user = await identityClient(env.IDENTITY)
45 .userByUsername(args.username)
46 .catch(() => null);
47 return user?.id ?? null;
48}
49
50/** The most deliveries in one call; chat sends one per person in the conversation. */
51const MAX_DELIVERIES = 1000;
52
53async function answer(env: Env, method: string, args: any): Promise<Response> {
54 if (method === "deliver") {
55 const items = (Array.isArray(args?.items) ? args.items : []).slice(0, MAX_DELIVERIES) as FeedDelivery[];
56 const byUser = new Map<string, FeedDelivery[]>();
57 for (const item of items) {
58 if (typeof item?.user_id !== "string" || !item.user_id) continue;
59 byUser.set(item.user_id, [...(byUser.get(item.user_id) ?? []), item]);
60 }
61 const results = await Promise.allSettled([...byUser].map(([id, mine]) => feed(env, id).deliver(mine)));
62 for (const result of results) if (result.status === "rejected") console.error("notify: a delivery failed", result.reason);
63 return Response.json({ ok: true });
64 }
65 const userId = await userIdOf(env, args ?? {});
66 if (!userId) return Response.json({ ok: false });
67 const stub = feed(env, userId);
68 switch (method) {
69 case "notify":
70 return Response.json(await stub.notify(args.notification));
71 case "set_inbox":
72 return Response.json(await stub.setInbox(Number(args.unread)));
73 case "subscribe":
74 return Response.json(await stub.subscribe(args.subscription, typeof args.user_agent === "string" ? args.user_agent : null));
75 case "unsubscribe":
76 return Response.json(await stub.unsubscribe(String(args.endpoint ?? "")));
77 case "status":
78 return Response.json(await stub.status(typeof args.endpoint === "string" ? args.endpoint : null));
79 case "set_preferences":
80 return Response.json(await stub.setPreferences(args.preferences));
81 case "test":
82 return Response.json(await stub.test(typeof args.username === "string" ? args.username : ""));
83 default:
84 return new Response("Unknown method\n", { status: 404 });
85 }
86}
87
88/**
89 * `GET /live`, upgraded: the viewer comes in NOTIFY_VIEWER_HEADER, set by
90 * the site after checking the session; trusted only because this Worker
91 * is reachable through service bindings alone.
92 */
93async function live(request: Request, env: Env): Promise<Response> {
94 if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") {
95 return new Response("Expected a WebSocket upgrade\n", { status: 426 });
96 }
97 let viewer: Viewer = null;
98 try {
99 viewer = JSON.parse(request.headers.get(NOTIFY_VIEWER_HEADER) ?? "null") as Viewer;
100 } catch {
101 viewer = null;
102 }
103 if (!viewer?.id) return new Response("Sign in first\n", { status: 401 });
104 const headers = new Headers(request.headers);
105 headers.delete(NOTIFY_VIEWER_HEADER);
106 return feed(env, viewer.id).fetch(new Request(request.url, { method: "GET", headers }));
107}
108
109export default {
110 async fetch(request: Request, env: Env): Promise<Response> {
111 const url = new URL(request.url);
112 if (request.method === "GET" && url.pathname === "/live") return live(request, env);
113 const match = url.pathname.match(/^\/rpc\/([a-z_]+)$/);
114 if (request.method !== "POST" || !match) return new Response("Not found\n", { status: 404 });
115 const args = (await request.json().catch(() => ({}))) as any;
116 return answer(env, match[1], args);
117 },
118} satisfies ExportedHandler<Env>;
119