Skip to content
231 linesCodeBlameRaw
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 */
17import type { FeedDelivery, FeedNotification, MemberProfile, NotificationKind } from "@g1t/contracts";
18
19import { tally, type UnreadRow } from "./unread.ts";
20
21export type Person = { key: string; user_id: string; username: string; muted: boolean };
22
23export type Recipient = { user_id: string; kind: NotificationKind | null; mentioned: boolean; muted: boolean };
24
25export 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. */
51export 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. */
64export 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. */
70export 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. */
123export 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
130type Db = D1Database;
131type Notify = { fetch(input: string, init?: RequestInit): Promise<Response> };
132
133async 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. */
144export 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. */
206export 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. */
226export 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}