Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002) | 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 | */ | |
| 19 | import { | |
| 20 | NOTIFY_VIEWER_HEADER, | |
| 21 | identityClient, | |
| 22 | type FeedDelivery, | |
| 23 | type ServiceBinding, | |
| 24 | type Viewer, | |
| 25 | } from "@g1t/contracts"; | |
| 26 | ||
| 27 | import { Feed, type FeedEnv } from "./feed.ts"; | |
| 28 | ||
| 29 | export { Feed } from "./feed.ts"; | |
| 30 | ||
| 31 | type Env = FeedEnv & { | |
| 32 | IDENTITY: ServiceBinding; | |
| 33 | FEEDS: DurableObjectNamespace<Feed>; | |
| 34 | }; | |
| 35 | ||
| 36 | function 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. */ | |
| 41 | async 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. */ | |
| 51 | const MAX_DELIVERIES = 1000; | |
| 52 | ||
| 53 | async 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 | */ | |
| 93 | async 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 | ||
| 109 | export 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 |
This file's history is long; its oldest lines are credited to the oldest commit read.