| 1 | /** |
| 2 | * What the notify service is told about a new message, or a read: every |
| 3 | * person in the conversation has its counts moved, and those it is for are |
| 4 | * notified (docs/WORKSPACE.md, "Live notifications"). |
| 5 | * |
| 6 | * - A direct message notifies everyone else in it, muted or not. |
| 7 | * - An @mention notifies the person named, muted or not. |
| 8 | * - A reply notifies the people already in its thread (who started it or |
| 9 | * replied), unless they muted the conversation. |
| 10 | * - Everyone else in it only has their counts moved. |
| 11 | * - The author is never told of their own message; their own count for |
| 12 | * the conversation goes to nothing (what you wrote, you have read). |
| 13 | * - Agents have no feed: only people are told. |
| 14 | * |
| 15 | * `recipients` is pure, so the rules are tested apart from the service. |
| 16 | */ |
| 17 | import type { FeedDelivery, FeedNotification, MemberProfile, NotificationKind } from "@g1t/contracts"; |
| 18 | |
| 19 | import { tally, type UnreadRow } from "./unread.ts"; |
| 20 | |
| 21 | export type Person = { key: string; user_id: string; username: string; muted: boolean }; |
| 22 | |
| 23 | export type Recipient = { user_id: string; kind: NotificationKind | null; mentioned: boolean; muted: boolean }; |
| 24 | |
| 25 | export function recipients(input: { |
| 26 | /** `user:<id>` or `agent:<id>`. */ |
| 27 | author: string; |
| 28 | channelKind: "channel" | "dm"; |
| 29 | /** The people in the conversation (agents left out). */ |
| 30 | people: Person[]; |
| 31 | /** Handles the message mentions, lowercased. */ |
| 32 | mentioned: string[]; |
| 33 | /** For a reply: the people in its thread, by key; null for a top-level message. */ |
| 34 | thread: ReadonlySet<string> | null; |
| 35 | }): Recipient[] { |
| 36 | const named = new Set(input.mentioned.map((h) => h.toLowerCase())); |
| 37 | const out: Recipient[] = []; |
| 38 | for (const person of input.people) { |
| 39 | if (person.key === input.author) continue; |
| 40 | const mentioned = named.has(person.username.toLowerCase()); |
| 41 | let kind: NotificationKind | null = null; |
| 42 | if (input.channelKind === "dm") kind = "dm"; |
| 43 | else if (mentioned) kind = "mention"; |
| 44 | else if (input.thread?.has(person.key) && !person.muted) kind = "thread_reply"; |
| 45 | out.push({ user_id: person.user_id, kind, mentioned, muted: person.muted }); |
| 46 | } |
| 47 | return out; |
| 48 | } |
| 49 | |
| 50 | /** A message's text as a notification shows it: one line, no markup, short. */ |
| 51 | export function preview(body: string, max = 140): string { |
| 52 | const text = String(body ?? "") |
| 53 | .replace(/```[\s\S]*?```/g, " [code] ") |
| 54 | .replace(/`([^`]*)`/g, "$1") |
| 55 | .replace(/!\[[^\]]*\]\([^)]*\)/g, "[image]") |
| 56 | .replace(/\[([^\]]*)\]\([^)]*\)/g, "$1") |
| 57 | .replace(/[*_~>#]+/g, "") |
| 58 | .replace(/\s+/g, " ") |
| 59 | .trim(); |
| 60 | return text.length > max ? `${text.slice(0, max - 1).trimEnd()}…` : text; |
| 61 | } |
| 62 | |
| 63 | /** Where a conversation, or a thread in it, is on the site. */ |
| 64 | export function conversationHref(slug: string, channel: { id: string; kind: "channel" | "dm"; name: string | null }, threadRoot: string | null): string { |
| 65 | const base = channel.kind === "dm" || !channel.name ? `/${slug}/-/chat/dm/${channel.id}` : `/${slug}/-/chat/${channel.name}`; |
| 66 | return threadRoot ? `${base}?thread=${encodeURIComponent(threadRoot)}` : base; |
| 67 | } |
| 68 | |
| 69 | /** The deliveries for one new message. */ |
| 70 | export function messageDeliveries(input: { |
| 71 | slug: string; |
| 72 | channel: { id: string; kind: "channel" | "dm"; name: string | null }; |
| 73 | message: { id: string; author: string; body: string; card_title: string | null; thread_root: string | null; created_at: string }; |
| 74 | author: MemberProfile; |
| 75 | recipients: Recipient[]; |
| 76 | }): FeedDelivery[] { |
| 77 | const { slug, channel, message, author } = input; |
| 78 | const where = channel.kind === "dm" ? "" : ` in #${channel.name}`; |
| 79 | // The service sets `display_name` by `memberName`'s rule (display name, |
| 80 | // else the username in its chosen case), so pushes match the chat. |
| 81 | const shown = author.display_name.trim() || author.display_username || author.name; |
| 82 | const deliveries: FeedDelivery[] = input.recipients.map((r) => { |
| 83 | const notification: FeedNotification | null = r.kind |
| 84 | ? { |
| 85 | id: message.id, |
| 86 | kind: r.kind, |
| 87 | workspace: slug, |
| 88 | title: r.kind === "thread_reply" ? `${shown} replied${where}` : `${shown}${where}`, |
| 89 | body: preview(message.card_title ?? message.body), |
| 90 | href: conversationHref(slug, channel, message.thread_root), |
| 91 | actor: { |
| 92 | kind: author.kind, |
| 93 | id: author.id, |
| 94 | name: shown, |
| 95 | avatar: author.avatar, |
| 96 | avatar_seed: author.avatar_seed ?? null, |
| 97 | }, |
| 98 | channel_id: channel.id, |
| 99 | thread_root: message.thread_root, |
| 100 | created_at: message.created_at, |
| 101 | } |
| 102 | : null; |
| 103 | return { |
| 104 | user_id: r.user_id, |
| 105 | workspace: slug, |
| 106 | counts: { channel_id: channel.id, unread: 1, mentions: r.mentioned ? 1 : 0, muted: r.muted }, |
| 107 | notification, |
| 108 | }; |
| 109 | }); |
| 110 | // The author has read the conversation by writing in it. |
| 111 | if (message.author.startsWith("user:")) { |
| 112 | deliveries.push({ |
| 113 | user_id: message.author.slice("user:".length), |
| 114 | workspace: slug, |
| 115 | counts: { channel_id: channel.id, unread: 0, mentions: 0, set: true }, |
| 116 | notification: null, |
| 117 | }); |
| 118 | } |
| 119 | return deliveries; |
| 120 | } |
| 121 | |
| 122 | /** A person's counts for one conversation after they read up to `lastReadId`, from the rows after it. */ |
| 123 | export function countsAfterRead(rows: UnreadRow[], channelId: string, member: string, username: string, lastReadId: string): { unread: number; mentions: number } { |
| 124 | return tally(rows, member, username, new Map([[channelId, lastReadId]])).get(channelId) ?? { unread: 0, mentions: 0 }; |
| 125 | } |
| 126 | |
| 127 | |
| 128 | // ── Wiring, used by src/index.ts ───────────────────────────────────────── |
| 129 | |
| 130 | type Db = D1Database; |
| 131 | type Notify = { fetch(input: string, init?: RequestInit): Promise<Response> }; |
| 132 | |
| 133 | async function deliver(notify: Notify, items: FeedDelivery[]): Promise<void> { |
| 134 | if (!items.length) return; |
| 135 | const response = await notify.fetch("https://service/rpc/deliver", { |
| 136 | method: "POST", |
| 137 | headers: { "content-type": "application/json" }, |
| 138 | body: JSON.stringify({ items }), |
| 139 | }); |
| 140 | if (!response.ok) throw new Error(`deliver failed with status ${response.status}`); |
| 141 | } |
| 142 | |
| 143 | /** Tells notify of a new message: one batched call for everyone in the conversation. */ |
| 144 | export async function notifyMessage( |
| 145 | db: Db, |
| 146 | notify: Notify | undefined, |
| 147 | profiles: (keys: string[]) => Promise<Map<string, MemberProfile>>, |
| 148 | input: { |
| 149 | slug: string; |
| 150 | channel: { id: string; kind: "channel" | "dm"; name: string | null }; |
| 151 | row: { id: string; author: string; body: string; card: string | null; thread_root: string | null; created_at: string }; |
| 152 | handles: string[]; |
| 153 | }, |
| 154 | ): Promise<void> { |
| 155 | if (!notify) return; |
| 156 | const { channel, row } = input; |
| 157 | const [members, thread] = await Promise.all([ |
| 158 | db |
| 159 | .prepare("SELECT principal, muted FROM channel_members WHERE channel_id = ? AND principal LIKE 'user:%'") |
| 160 | .bind(channel.id) |
| 161 | .all<{ principal: string; muted: number }>(), |
| 162 | row.thread_root |
| 163 | ? db |
| 164 | .prepare("SELECT DISTINCT author FROM messages WHERE channel_id = ?1 AND (id = ?2 OR thread_root = ?2) AND author LIKE 'user:%'") |
| 165 | .bind(channel.id, row.thread_root) |
| 166 | .all<{ author: string }>() |
| 167 | : Promise.resolve(null), |
| 168 | ]); |
| 169 | const keys = members.results.map((m) => m.principal); |
| 170 | const found = await profiles([...keys, row.author]); |
| 171 | const people: Person[] = members.results.map((m) => ({ |
| 172 | key: m.principal, |
| 173 | user_id: m.principal.slice("user:".length), |
| 174 | username: found.get(m.principal)?.name ?? "", |
| 175 | muted: !!m.muted, |
| 176 | })); |
| 177 | const author = found.get(row.author); |
| 178 | if (!author) return; |
| 179 | let cardTitle: string | null = null; |
| 180 | if (row.card) { |
| 181 | try { |
| 182 | cardTitle = (JSON.parse(row.card) as { title?: string }).title ?? null; |
| 183 | } catch { |
| 184 | cardTitle = null; |
| 185 | } |
| 186 | } |
| 187 | await deliver( |
| 188 | notify, |
| 189 | messageDeliveries({ |
| 190 | slug: input.slug, |
| 191 | channel, |
| 192 | message: { id: row.id, author: row.author, body: row.body, card_title: cardTitle, thread_root: row.thread_root, created_at: row.created_at }, |
| 193 | author, |
| 194 | recipients: recipients({ |
| 195 | author: row.author, |
| 196 | channelKind: channel.kind, |
| 197 | people, |
| 198 | mentioned: input.handles, |
| 199 | thread: thread ? new Set(thread.results.map((r) => r.author)) : null, |
| 200 | }), |
| 201 | }), |
| 202 | ); |
| 203 | } |
| 204 | |
| 205 | /** After a read: the person's counts for the conversation, as they now are, in every tab. */ |
| 206 | export async function notifyRead( |
| 207 | db: Db, |
| 208 | notify: Notify | undefined, |
| 209 | input: { slug: string; channel_id: string; user_id: string; username: string; last_read_id: string }, |
| 210 | ): Promise<void> { |
| 211 | if (!notify) return; |
| 212 | const member = `user:${input.user_id}`; |
| 213 | const rows = await db |
| 214 | .prepare( |
| 215 | "SELECT channel_id, id, author, mentions FROM messages WHERE channel_id = ? AND id > ? AND deleted_at IS NULL AND author != ? LIMIT 5000", |
| 216 | ) |
| 217 | .bind(input.channel_id, input.last_read_id, member) |
| 218 | .all<UnreadRow>(); |
| 219 | const counts = countsAfterRead(rows.results, input.channel_id, member, input.username, input.last_read_id); |
| 220 | await deliver(notify, [ |
| 221 | { user_id: input.user_id, workspace: input.slug, counts: { channel_id: input.channel_id, ...counts, set: true }, notification: null }, |
| 222 | ]); |
| 223 | } |
| 224 | |
| 225 | /** After muting or unmuting: the badge counts the conversation, or leaves it out, in every tab. */ |
| 226 | export async function notifyMuted(notify: Notify | undefined, input: { slug: string; channel_id: string; user_id: string; muted: boolean }): Promise<void> { |
| 227 | if (!notify) return; |
| 228 | await deliver(notify, [ |
| 229 | { user_id: input.user_id, workspace: input.slug, counts: { channel_id: input.channel_id, unread: 0, mentions: 0, muted: input.muted }, notification: null }, |
| 230 | ]); |
| 231 | } |