Skip to content
147 linesCodeBlameRaw
1/**
2 * The notify service: live notifications, unread counts and browser push.
3 * One feed per person (src/feed.ts).
4 *
5 * Reached through service bindings only: `POST /rpc/<method>` with
6 * snake_case bodies (`notifyClient` in @g1t/contracts), and `GET /live`,
7 * the person's feed socket, which the site forwards after checking the
8 * session (NOTIFY_VIEWER_HEADER), with the counts it read (NOTIFY_SEED_HEADER).
9 *
10 * Who sends what:
11 * - chat: `deliver`, for every message (counts for each person in the
12 * conversation, a notification for those it is for) and every read;
13 * - events: `notify`, for every inbox item, and `set_inbox`, the count after
14 * items arrive or are marked, naming the person by username;
15 * - the site: `subscribe`, `unsubscribe`, `status`, `set_preferences`,
16 * `test`, `presence` and `set_presence`, for the person signed in;
17 * - agents and the site: `workspace_presence`, how a workspace's people
18 * show now (away, in Do Not Disturb, offline, and their statuses);
19 * - integrations, later: `set_presence` with a status of `source`
20 * `calendar` or `integration`, for the person they act for.
21 *
22 * Presence: each feed works out its person's from their tabs and tells
23 * one room per workspace (src/room.ts), which tells everyone there who is
24 * online. See docs.g1t.sh/guides/chat/, "Presence and status".
25 */
26import {
27 NOTIFY_VIEWER_HEADER,
28 identityClient,
29 type FeedDelivery,
30 type ServiceBinding,
31 type Viewer,
32} from "@g1t/contracts";
33
34import { FEED_USERNAME_HEADER, FEED_USER_ID_HEADER, Feed, type FeedEnv } from "./feed.ts";
35import { onlyPeople } from "./presence.ts";
36
37export { Feed } from "./feed.ts";
38export { Room } from "./room.ts";
39
40type Env = FeedEnv & {
41 IDENTITY: ServiceBinding;
42 FEEDS: DurableObjectNamespace<Feed>;
43};
44
45function feed(env: Env, userId: string) {
46 return env.FEEDS.get(env.FEEDS.idFromName(userId));
47}
48
49/** A user id from the arguments, or one looked up by username. */
50async function userIdOf(env: Env, given: { user_id?: unknown; username?: unknown; target?: unknown }): Promise<string | null> {
51 // Who it's for, at the top level or under `target` as the agents service sends it.
52 const target = given.target && typeof given.target === "object" ? (given.target as { user_id?: unknown; username?: unknown }) : null;
53 const args = { user_id: given.user_id ?? target?.user_id, username: given.username ?? target?.username };
54 if (typeof args.user_id === "string" && args.user_id) return args.user_id;
55 if (typeof args.username !== "string" || !args.username) return null;
56 const user = await identityClient(env.IDENTITY)
57 .userByUsername(args.username)
58 .catch(() => null);
59 return user?.id ?? null;
60}
61
62/** The most deliveries in one call; chat sends one per person in the conversation. */
63const MAX_DELIVERIES = 1000;
64
65async function answer(env: Env, method: string, args: any): Promise<Response> {
66 if (method === "workspace_presence") {
67 // For services: how a workspace's people show now, as its room keeps them.
68 const slug = typeof args?.workspace === "string" ? args.workspace.trim().toLowerCase() : "";
69 if (!slug || !env.ROOMS) return Response.json([]);
70 const everyone = await env.ROOMS.get(env.ROOMS.idFromName(slug)).people();
71 return Response.json(onlyPeople(everyone, Array.isArray(args.user_ids) ? args.user_ids : null));
72 }
73 if (method === "deliver") {
74 const items = (Array.isArray(args?.items) ? args.items : []).slice(0, MAX_DELIVERIES) as FeedDelivery[];
75 const byUser = new Map<string, FeedDelivery[]>();
76 for (const item of items) {
77 if (typeof item?.user_id !== "string" || !item.user_id) continue;
78 byUser.set(item.user_id, [...(byUser.get(item.user_id) ?? []), item]);
79 }
80 const results = await Promise.allSettled([...byUser].map(([id, mine]) => feed(env, id).deliver(mine)));
81 for (const result of results) if (result.status === "rejected") console.error("notify: a delivery failed", result.reason);
82 return Response.json({ ok: true });
83 }
84 const userId = await userIdOf(env, args ?? {});
85 if (!userId) return Response.json({ ok: false });
86 const stub = feed(env, userId);
87 switch (method) {
88 case "notify":
89 return Response.json(await stub.notify(args.notification));
90 case "set_inbox":
91 return Response.json(await stub.setInbox(Number(args.unread)));
92 case "subscribe":
93 return Response.json(await stub.subscribe(args.subscription, typeof args.user_agent === "string" ? args.user_agent : null));
94 case "unsubscribe":
95 return Response.json(await stub.unsubscribe(String(args.endpoint ?? "")));
96 case "status":
97 return Response.json(await stub.status(typeof args.endpoint === "string" ? args.endpoint : null));
98 case "set_preferences":
99 return Response.json(await stub.setPreferences(args.preferences));
100 case "test":
101 return Response.json(await stub.test(typeof args.username === "string" ? args.username : ""));
102 case "presence":
103 return Response.json(await stub.own({ user_id: userId, username: typeof args.username === "string" ? args.username : "" }));
104 case "set_presence":
105 return Response.json(
106 await stub.setPresence({ user_id: userId, username: typeof args.username === "string" ? args.username : "" }, args.change ?? {}),
107 );
108 default:
109 return new Response("Unknown method\n", { status: 404 });
110 }
111}
112
113/**
114 * `GET /live`, upgraded: the viewer comes in NOTIFY_VIEWER_HEADER, set by
115 * the site after checking the session; trusted only because this Worker
116 * is reachable through service bindings alone.
117 */
118async function live(request: Request, env: Env): Promise<Response> {
119 if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") {
120 return new Response("Expected a WebSocket upgrade\n", { status: 426 });
121 }
122 let viewer: Viewer = null;
123 try {
124 viewer = JSON.parse(request.headers.get(NOTIFY_VIEWER_HEADER) ?? "null") as Viewer;
125 } catch {
126 viewer = null;
127 }
128 if (!viewer?.id) return new Response("Sign in first\n", { status: 401 });
129 const headers = new Headers(request.headers);
130 headers.delete(NOTIFY_VIEWER_HEADER);
131 // Who the feed is for, so it can name them to the presence rooms.
132 headers.set(FEED_USER_ID_HEADER, viewer.id);
133 headers.set(FEED_USERNAME_HEADER, viewer.username);
134 return feed(env, viewer.id).fetch(new Request(request.url, { method: "GET", headers }));
135}
136
137export default {
138 async fetch(request: Request, env: Env): Promise<Response> {
139 const url = new URL(request.url);
140 if (request.method === "GET" && url.pathname === "/live") return live(request, env);
141 const match = url.pathname.match(/^\/rpc\/([a-z_]+)$/);
142 if (request.method !== "POST" || !match) return new Response("Not found\n", { status: 404 });
143 const args = (await request.json().catch(() => ({}))) as any;
144 return answer(env, match[1], args);
145 },
146} satisfies ExportedHandler<Env>;
147