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