| 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 | */ |
| 26 | import { |
| 27 | NOTIFY_VIEWER_HEADER, |
| 28 | identityClient, |
| 29 | type FeedDelivery, |
| 30 | type ServiceBinding, |
| 31 | type Viewer, |
| 32 | } from "@g1t/contracts"; |
| 33 | |
| 34 | import { FEED_USERNAME_HEADER, FEED_USER_ID_HEADER, Feed, type FeedEnv } from "./feed.ts"; |
| 35 | import { onlyPeople } from "./presence.ts"; |
| 36 | |
| 37 | export { Feed } from "./feed.ts"; |
| 38 | export { Room } from "./room.ts"; |
| 39 | |
| 40 | type Env = FeedEnv & { |
| 41 | IDENTITY: ServiceBinding; |
| 42 | FEEDS: DurableObjectNamespace<Feed>; |
| 43 | }; |
| 44 | |
| 45 | function 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. */ |
| 50 | async 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. */ |
| 63 | const MAX_DELIVERIES = 1000; |
| 64 | |
| 65 | async 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 | */ |
| 118 | async 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 | |
| 137 | export 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 | |