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