Skip to content
228 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 const deliveries: FeedDelivery[] = input.recipients.map((r) => {
80 const notification: FeedNotification | null = r.kind
81 ? {
82 id: message.id,
83 kind: r.kind,
84 workspace: slug,
85 title: r.kind === "thread_reply" ? `${author.display_name} replied${where}` : `${author.display_name}${where}`,
86 body: preview(message.card_title ?? message.body),
87 href: conversationHref(slug, channel, message.thread_root),
88 actor: {
89 kind: author.kind,
90 id: author.id,
91 name: author.display_name,
92 avatar: author.avatar,
93 avatar_seed: author.avatar_seed ?? null,
94 },
95 channel_id: channel.id,
96 thread_root: message.thread_root,
97 created_at: message.created_at,
98 }
99 : null;
100 return {
101 user_id: r.user_id,
102 workspace: slug,
103 counts: { channel_id: channel.id, unread: 1, mentions: r.mentioned ? 1 : 0, muted: r.muted },
104 notification,
105 };
106 });
107 // The author has read the conversation by writing in it.
108 if (message.author.startsWith("user:")) {
109 deliveries.push({
110 user_id: message.author.slice("user:".length),
111 workspace: slug,
112 counts: { channel_id: channel.id, unread: 0, mentions: 0, set: true },
113 notification: null,
114 });
115 }
116 return deliveries;
117}
118
119/** A person's counts for one conversation after they read up to `lastReadId`, from the rows after it. */
120export function countsAfterRead(rows: UnreadRow[], channelId: string, member: string, username: string, lastReadId: string): { unread: number; mentions: number } {
121 return tally(rows, member, username, new Map([[channelId, lastReadId]])).get(channelId) ?? { unread: 0, mentions: 0 };
122}
123
124
125// ── Wiring, used by src/index.ts ─────────────────────────────────────────
126
127type Db = D1Database;
128type Notify = { fetch(input: string, init?: RequestInit): Promise<Response> };
129
130async function deliver(notify: Notify, items: FeedDelivery[]): Promise<void> {
131 if (!items.length) return;
132 const response = await notify.fetch("https://service/rpc/deliver", {
133 method: "POST",
134 headers: { "content-type": "application/json" },
135 body: JSON.stringify({ items }),
136 });
137 if (!response.ok) throw new Error(`deliver failed with status ${response.status}`);
138}
139
140/** Tells notify of a new message: one batched call for everyone in the conversation. */
141export async function notifyMessage(
142 db: Db,
143 notify: Notify | undefined,
144 profiles: (keys: string[]) => Promise<Map<string, MemberProfile>>,
145 input: {
146 slug: string;
147 channel: { id: string; kind: "channel" | "dm"; name: string | null };
148 row: { id: string; author: string; body: string; card: string | null; thread_root: string | null; created_at: string };
149 handles: string[];
150 },
151): Promise<void> {
152 if (!notify) return;
153 const { channel, row } = input;
154 const [members, thread] = await Promise.all([
155 db
156 .prepare("SELECT principal, muted FROM channel_members WHERE channel_id = ? AND principal LIKE 'user:%'")
157 .bind(channel.id)
158 .all<{ principal: string; muted: number }>(),
159 row.thread_root
160 ? db
161 .prepare("SELECT DISTINCT author FROM messages WHERE channel_id = ?1 AND (id = ?2 OR thread_root = ?2) AND author LIKE 'user:%'")
162 .bind(channel.id, row.thread_root)
163 .all<{ author: string }>()
164 : Promise.resolve(null),
165 ]);
166 const keys = members.results.map((m) => m.principal);
167 const found = await profiles([...keys, row.author]);
168 const people: Person[] = members.results.map((m) => ({
169 key: m.principal,
170 user_id: m.principal.slice("user:".length),
171 username: found.get(m.principal)?.name ?? "",
172 muted: !!m.muted,
173 }));
174 const author = found.get(row.author);
175 if (!author) return;
176 let cardTitle: string | null = null;
177 if (row.card) {
178 try {
179 cardTitle = (JSON.parse(row.card) as { title?: string }).title ?? null;
180 } catch {
181 cardTitle = null;
182 }
183 }
184 await deliver(
185 notify,
186 messageDeliveries({
187 slug: input.slug,
188 channel,
189 message: { id: row.id, author: row.author, body: row.body, card_title: cardTitle, thread_root: row.thread_root, created_at: row.created_at },
190 author,
191 recipients: recipients({
192 author: row.author,
193 channelKind: channel.kind,
194 people,
195 mentioned: input.handles,
196 thread: thread ? new Set(thread.results.map((r) => r.author)) : null,
197 }),
198 }),
199 );
200}
201
202/** After a read: the person's counts for the conversation, as they now are, in every tab. */
203export async function notifyRead(
204 db: Db,
205 notify: Notify | undefined,
206 input: { slug: string; channel_id: string; user_id: string; username: string; last_read_id: string },
207): Promise<void> {
208 if (!notify) return;
209 const member = `user:${input.user_id}`;
210 const rows = await db
211 .prepare(
212 "SELECT channel_id, id, author, mentions FROM messages WHERE channel_id = ? AND id > ? AND deleted_at IS NULL AND author != ? LIMIT 5000",
213 )
214 .bind(input.channel_id, input.last_read_id, member)
215 .all<UnreadRow>();
216 const counts = countsAfterRead(rows.results, input.channel_id, member, input.username, input.last_read_id);
217 await deliver(notify, [
218 { user_id: input.user_id, workspace: input.slug, counts: { channel_id: input.channel_id, ...counts, set: true }, notification: null },
219 ]);
220}
221
222/** After muting or unmuting: the badge counts the conversation, or leaves it out, in every tab. */
223export async function notifyMuted(notify: Notify | undefined, input: { slug: string; channel_id: string; user_id: string; muted: boolean }): Promise<void> {
224 if (!notify) return;
225 await deliver(notify, [
226 { user_id: input.user_id, workspace: input.slug, counts: { channel_id: input.channel_id, unread: 0, mentions: 0, muted: input.muted }, notification: null },
227 ]);
228}