Skip to content
1,920 linesCodeBlameRaw
1/**
2 * The chat service: a workspace's channels, direct messages, threads and
3 * messages. People and agents are members alike. Plan: docs/WORKSPACE.md.
4 *
5 * Reached through service bindings: `POST /rpc/<method>` with snake_case
6 * bodies (`chatClient` in @g1t/contracts), and `GET /live` for a channel's
7 * socket, which the site forwards after checking the session. Each channel
8 * has a room (src/room.ts) that delivers what happens in it live.
9 *
10 * Workspaces are kept by id, so renaming one changes nothing here; who is
11 * in a workspace comes from the viewer's memberships, as in every service.
12 */
13
14import {
15 CHAT_MAX_HOPS,
16 askerAccess,
17 CHAT_VIEWER_HEADER,
18 fail,
19 identityClient,
20 newId,
21 ok,
22 openD1,
23 parsePrincipalKey,
24 principalKey,
25 workspaceAgentsClient,
26 type AgentDelivery,
27 type AgentFoundMessage,
28 type AgentPostMessage,
29 type AskerAccess,
30 type Channel,
31 type ChannelChange,
32 type ChannelDetail,
33 type ChannelMember,
34 type ChatAudience,
35 type ChatLiveEvent,
36 type ChatMessage,
37 type ChatReaction,
38 type CustomEmoji,
39 type EmojiFile,
40 type EmojiList,
41 type EmojiUpload,
42 type ChatSettings,
43 type ChatSettingsView,
44 type ChatSidebar,
45 type ChatSidebarEntry,
46
47 type Member,
48 type MemberProfile,
49 type MessageCard,
50 type MessagePage,
51 type NewChannel,
52 type PostMessage,
53 type Principal,
54 type Result,
55 type ServiceBinding,
56 type User,
57 type Viewer,
58 type Workspace,
59 type WorkspaceAgent,
60} from "@g1t/contracts";
61
62import { audienceKind, isShared, likePattern, readableBy } from "./audience.ts";
63import { MAX_HOPS, addsOrchestrator, chainFor, deliveries, delivery, sender, type Chain } from "./delivery.ts";
64import {
65 MAX_REACTIONS_PER_MESSAGE,
66 emojiImage,
67 emojiName,
68 fromBase64,
69 mayRemove,
70 mayUpload,
71 reactionEmoji,
72 roomForReaction,
73 tallyReactions,
74 type ReactionRow,
75 type ReactionTally,
76} from "./emoji.ts";
77import { mentionedHandles, mentionsColumn } from "./mentions.ts";
78import { AGENT_TYPING_MS, historyOf, historySize, messageBody, meterDay, pageOf, pageSize } from "./messages.ts";
79import { GENERAL, MAX_DM_MEMBERS, channelName, dmKey, dmMembers } from "./names.ts";
80import { ROOM_MEMBER_HEADER, type ChannelRoom, type RoomMember } from "./room.ts";
81import { mayCreateChannel, mayManageChannel, permissionsFor, rowFor, settingsChange, settingsOf, type SettingsRow } from "./settings.ts";
82import { dmTitle, sidebarOrder, tally, type UnreadRow } from "./unread.ts";
83// Live notifications and counts (services/notify).
84import { notifyMessage, notifyMuted, notifyRead } from "./notify.ts";
85
86export { ChannelRoom } from "./room.ts";
87
88// The hop limit here is the one in the contract.
89const SAME_HOP_LIMIT: typeof CHAT_MAX_HOPS = MAX_HOPS;
90void SAME_HOP_LIMIT;
91
92type Env = {
93 DB: D1Database;
94 IDENTITY: ServiceBinding;
95 AGENTS: ServiceBinding;
96 ROOMS: DurableObjectNamespace<ChannelRoom>;
97 /** The avatars namespace: custom emoji images, under `emoji/<sha256>`, which the usercontent origin serves. */
98 AVATARS: KVNamespace;
99 /** Live notifications and unread counts (services/notify); absent, nobody is told. */
100 NOTIFY?: ServiceBinding;
101};
102
103type EmojiRow = {
104 workspace_id: string;
105 name: string;
106 alias_of: string | null;
107 file: string;
108 content_type: CustomEmoji["content_type"];
109 bytes: number;
110 created_by: string;
111 created_at: string;
112 deleted_at: string | null;
113};
114
115/** The viewer's role in a workspace, or null when they are not in it. */
116function roleOf(viewer: Viewer, workspace: string): "owner" | "member" | null {
117 return viewer?.workspaces?.find((m) => m.slug === workspace.toLowerCase())?.role ?? null;
118}
119
120async function sha256(bytes: Uint8Array): Promise<string> {
121 const digest = await crypto.subtle.digest("SHA-256", bytes);
122 return [...new Uint8Array(digest)].map((b) => b.toString(16).padStart(2, "0")).join("");
123}
124
125type ChannelRow = {
126 id: string;
127 workspace_id: string;
128 kind: "channel" | "dm";
129 name: string | null;
130 topic: string | null;
131 private: number;
132 dm_key: string | null;
133 created_by: string;
134 created_at: string;
135 archived_at: string | null;
136 last_message_at: string | null;
137};
138
139type MemberRow = {
140 channel_id: string;
141 principal: string;
142 role: "owner" | "member";
143 starred: number;
144 muted: number;
145 last_read_id: string | null;
146 joined_at: string;
147};
148
149type MessageRow = {
150 id: string;
151 channel_id: string;
152 author: string;
153 kind: "text" | "card";
154 body: string;
155 card: string | null;
156 mentions: string;
157 thread_root: string | null;
158 reply_count: number;
159 last_reply_at: string | null;
160 created_at: string;
161 edited_at: string | null;
162 deleted_at: string | null;
163};
164
165/** What a method found out about the channel it was asked about. */
166type Place = { slug: string; workspace: Workspace; channel: ChannelRow; member: MemberRow | null };
167
168/** The most unread messages one sidebar reads to count; past it, counts are "at least". */
169const MAX_UNREAD_ROWS = 5_000;
170/** The longest channel topic. */
171const MAX_TOPIC = 250;
172/** How many others a direct message's sidebar entry shows. */
173const DM_FACES = 4;
174
175const now = () => new Date().toISOString();
176
177function isMember(viewer: Viewer, workspace: string): boolean {
178 return !!viewer?.workspaces?.some((m) => m.slug === workspace.toLowerCase());
179}
180
181function isOwner(viewer: Viewer, workspace: string): boolean {
182 return !!viewer?.workspaces?.some((m) => m.slug === workspace.toLowerCase() && m.role === "owner");
183}
184
185function userKey(viewer: User): string {
186 return principalKey({ kind: "user", id: viewer.id });
187}
188
189function toChannel(row: ChannelRow): Channel {
190 return {
191 id: row.id,
192 workspace_id: row.workspace_id,
193 kind: row.kind,
194 name: row.kind === "dm" ? null : row.name,
195 topic: row.topic,
196 private: row.kind === "dm" || !!row.private,
197 created_by: parsePrincipalKey(row.created_by) ?? { kind: "user", id: row.created_by },
198 created_at: row.created_at,
199 archived_at: row.archived_at,
200 last_message_at: row.last_message_at,
201 };
202}
203
204/** A card as kept, or null when what was sent is not one. */
205function cleanCard(card: unknown): MessageCard | null {
206 if (!card || typeof card !== "object") return null;
207 const c = card as Record<string, unknown>;
208 const text = (value: unknown, max: number) => (typeof value === "string" && value.trim() ? value.trim().slice(0, max) : null);
209 const kind = text(c.kind, 40);
210 const title = text(c.title, 300);
211 if (!kind || !title) return null;
212 const href = text(c.href, 2000);
213 return {
214 kind,
215 title,
216 detail: text(c.detail, 500),
217 state: text(c.state, 80),
218 // Relative to the site only: a card never links somewhere else.
219 href: href && href.startsWith("/") && !href.startsWith("//") ? href : null,
220 };
221}
222
223/** An asker handed back by the agents service, or null when what was sent is not one. */
224function cleanAsker(asker: unknown): AskerAccess | null {
225 if (!asker || typeof asker !== "object") return null;
226 const a = asker as Record<string, unknown>;
227 if (typeof a.username !== "string" || !["owner", "member", "outside"].includes(String(a.role))) return null;
228 return { username: a.username, role: a.role as AskerAccess["role"], can_write: a.can_write === true };
229}
230
231function bytesOf(text: string): number {
232 return new TextEncoder().encode(text).length;
233}
234
235class Chat {
236 private readonly workspaces = new Map<string, Promise<Workspace | null>>();
237 private readonly people = new Map<string, Promise<Map<string, Member>>>();
238 private readonly usernames = new Map<string, string>();
239 private readonly agents = new Map<string, WorkspaceAgent | null>();
240 private readonly settingsRows = new Map<string, Promise<SettingsRow | null>>();
241
242 /** `defer` runs work after the answer is sent: the request's waitUntil. */
243 constructor(
244 private readonly env: Env,
245 private readonly defer: (work: Promise<unknown>) => void = () => {},
246 ) {}
247
248 private get db() {
249 return this.env.DB;
250 }
251
252 // ── Who and where ───────────────────────────────────────────────────────
253
254 private workspace(slug: string): Promise<Workspace | null> {
255 const key = slug.toLowerCase();
256 let found = this.workspaces.get(key);
257 if (!found) {
258 found = identityClient(this.env.IDENTITY).getWorkspace(key);
259 this.workspaces.set(key, found);
260 }
261 return found;
262 }
263
264 /** The workspace's people by username, with their names and avatars; asked once per request. */
265 private members(slug: string, workspace: Workspace): Promise<Map<string, Member>> {
266 let found = this.people.get(workspace.id);
267 if (!found) {
268 // Asked as the workspace itself, so it works for agents' calls too.
269 const actor: User = {
270 id: workspace.id,
271 username: workspace.slug,
272 kind: "workspace",
273 verified: true,
274 workspaces: [{ slug: workspace.slug, role: "member" }],
275 };
276 found = identityClient(this.env.IDENTITY)
277 .listMembers(slug, actor)
278 .then((result) => new Map(result.ok ? result.value.map((m) => [m.username, m]) : []))
279 .catch((error) => {
280 console.error("chat could not list members of", slug, error);
281 return new Map<string, Member>();
282 });
283 this.people.set(workspace.id, found);
284 }
285 return found;
286 }
287
288 /** Agents by id; ones the agents service does not know are null. Asked once per request. */
289 private async agentsById(ids: string[]): Promise<Map<string, WorkspaceAgent | null>> {
290 const wanted = [...new Set(ids)].filter((id) => !this.agents.has(id));
291 if (wanted.length) {
292 let found: WorkspaceAgent[] = [];
293 try {
294 found = await workspaceAgentsClient(this.env.AGENTS).byIds(wanted);
295 } catch (error) {
296 console.error("chat could not resolve agents", error);
297 }
298 for (const id of wanted) this.agents.set(id, found.find((a) => a.id === id) ?? null);
299 }
300 return new Map(ids.map((id) => [id, this.agents.get(id) ?? null]));
301 }
302
303 /** An agent of this workspace that is not archived, or null. */
304 private async liveAgent(workspace: Workspace, id: string): Promise<WorkspaceAgent | null> {
305 const agent = (await this.agentsById([id])).get(id) ?? null;
306 return agent && agent.workspace_id === workspace.id && !agent.archived_at ? agent : null;
307 }
308
309 /** How each member key shows, for one workspace. */
310 private async profiles(slug: string, workspace: Workspace, keys: string[]): Promise<Map<string, MemberProfile>> {
311 const principals = [...new Set(keys)].map((key) => parsePrincipalKey(key)).filter((p): p is Principal => !!p);
312 const userIds = principals.filter((p) => p.kind === "user").map((p) => p.id);
313 const agentIds = principals.filter((p) => p.kind === "agent").map((p) => p.id);
314 const unnamed = userIds.filter((id) => !this.usernames.has(id));
315 const [named, people, agents] = await Promise.all([
316 unnamed.length ? identityClient(this.env.IDENTITY).usernames(unnamed).catch(() => ({}) as Record<string, string>) : ({} as Record<string, string>),
317 userIds.length ? this.members(slug, workspace) : new Map<string, Member>(),
318 this.agentsById(agentIds),
319 ]);
320 for (const [id, username] of Object.entries(named)) this.usernames.set(id, username);
321 const out = new Map<string, MemberProfile>();
322 for (const p of principals) {
323 if (p.kind === "user") {
324 const username = this.usernames.get(p.id) ?? null;
325 const person = username ? people.get(username) : undefined;
326 out.set(principalKey(p), {
327 ...p,
328 name: username ?? "ghost",
329 display_name: person?.name || username || "Former member",
330 avatar: person?.avatar ?? null,
331 role: null,
332 title: null,
333 avatar_seed: null,
334 });
335 } else {
336 const agent = agents.get(p.id) ?? null;
337 out.set(principalKey(p), {
338 ...p,
339 name: agent?.handle ?? p.id,
340 display_name: agent?.display_name ?? "Former agent",
341 avatar: agent?.avatar ?? null,
342 role: agent?.role ?? null,
343 title: agent?.title || null,
344 avatar_seed: agent?.avatar_seed ?? null,
345 });
346 }
347 }
348 return out;
349 }
350
351 private async profile(slug: string, workspace: Workspace, key: string): Promise<MemberProfile> {
352 return (await this.profiles(slug, workspace, [key])).get(key)!;
353 }
354
355 /** Whether `principal` may be added to a conversation in this workspace. */
356 private async belongs(slug: string, workspace: Workspace, principal: Principal): Promise<boolean> {
357 if (principal.kind === "agent") return !!(await this.liveAgent(workspace, principal.id));
358 if (!this.usernames.has(principal.id)) {
359 const named = await identityClient(this.env.IDENTITY).usernames([principal.id]);
360 for (const [id, username] of Object.entries(named)) this.usernames.set(id, username);
361 }
362 const username = this.usernames.get(principal.id);
363 return !!username && (await this.members(slug, workspace)).has(username);
364 }
365
366 /**
367 * The viewer's workspace, checked: they must belong to it, as in every
368 * other service.
369 */
370 private async viewerWorkspace(slug: string, viewer: Viewer): Promise<Result<Workspace>> {
371 if (!viewer) return fail("unauthenticated", "Sign in to use chat.");
372 if (!slug || !isMember(viewer, slug)) return fail("forbidden", "Only members of a workspace can use its chat.");
373 const workspace = await this.workspace(slug);
374 return workspace ? ok(workspace) : fail("not_found", "No such workspace.");
375 }
376
377 /**
378 * A channel the viewer may read (`read`: any public one in their
379 * workspace, or one they are in) or write in (`member`: one they are in).
380 * A private channel or direct message they are not in is not found, so
381 * its existence does not leak.
382 */
383 private async place(
384 slug: string,
385 channelId: string,
386 viewer: Viewer,
387 need: "read" | "member",
388 ): Promise<Result<Place>> {
389 const found = await this.viewerWorkspace(slug, viewer);
390 if (!found.ok) return found;
391 const workspace = found.value;
392 const [channel, member] = await Promise.all([
393 this.db
394 .prepare("SELECT * FROM channels WHERE id = ? AND workspace_id = ?")
395 .bind(String(channelId ?? ""), workspace.id)
396 .first<ChannelRow>(),
397 this.db
398 .prepare("SELECT * FROM channel_members WHERE channel_id = ? AND principal = ?")
399 .bind(String(channelId ?? ""), userKey(viewer!))
400 .first<MemberRow>(),
401 ]);
402 if (!channel) return fail("not_found", "No such channel.");
403 const open = channel.kind === "channel" && !channel.private;
404 if (!member && !open) return fail("not_found", "No such channel.");
405 if (!member && need === "member") return fail("forbidden", `Join #${channel.name} first.`);
406 return ok({ slug: slug.toLowerCase(), workspace, channel, member });
407 }
408
409 private room(channelId: string) {
410 return this.env.ROOMS.get(this.env.ROOMS.idFromName(channelId));
411 }
412
413 /** Tells everyone looking at a channel, after the answer is sent. */
414 private broadcast(channelId: string, event: ChatLiveEvent, except: string | null = null): void {
415 this.defer(
416 this.room(channelId)
417 .broadcast(event, except)
418 .catch((error: unknown) => console.error("chat could not broadcast to", channelId, error)),
419 );
420 }
421
422 /**
423 * Messages as they go out, with their reactions. `me` (a member key) is
424 * the viewer an answer is for; null for what everyone in a room gets,
425 * where no reaction is anyone's own.
426 */
427 private async toMessages(slug: string, workspace: Workspace, rows: MessageRow[], me: string | null = null): Promise<ChatMessage[]> {
428 const reactions = await this.reactionsOf(rows.filter((r) => !r.deleted_at).map((r) => r.id), me);
429 const reactors = [...reactions.values()].flatMap((list) => list.flatMap((r) => r.by));
430 const profiles = await this.profiles(slug, workspace, [...rows.map((r) => r.author), ...reactors]);
431 return rows.map((row) => {
432 const gone = !!row.deleted_at;
433 return {
434 reactions: gone
435 ? []
436 : (reactions.get(row.id) ?? []).map((r) => ({ ...r, by: r.by.map((key) => profiles.get(key)!).filter(Boolean) })),
437 id: row.id,
438 channel_id: row.channel_id,
439 author: profiles.get(row.author)!,
440 kind: row.kind,
441 body: gone ? "" : row.body,
442 card: gone || !row.card ? null : (JSON.parse(row.card) as MessageCard),
443 thread_root: row.thread_root,
444 reply_count: row.reply_count,
445 last_reply_at: row.last_reply_at,
446 created_at: row.created_at,
447 edited_at: row.edited_at,
448 deleted_at: row.deleted_at,
449 };
450 });
451 }
452
453 /** Reactions on these messages, counted (src/emoji.ts). One read, by the reactions table's key. */
454 private async reactionsOf(ids: string[], me: string | null): Promise<Map<string, ReactionTally[]>> {
455 if (!ids.length) return new Map();
456 const rows = await this.db
457 .prepare(
458 "SELECT message_id, emoji, principal, created_at FROM reactions WHERE message_id IN (SELECT value FROM json_each(?))",
459 )
460 .bind(JSON.stringify(ids))
461 .all<ReactionRow>();
462 return tallyReactions(rows.results, me);
463 }
464
465 private async messageRow(channelId: string, id: string): Promise<MessageRow | null> {
466 return this.db.prepare("SELECT * FROM messages WHERE id = ? AND channel_id = ?").bind(String(id ?? ""), channelId).first<MessageRow>();
467 }
468
469 /** Sends a message as it now is to everyone looking at its channel. */
470 private rebroadcast(place: Place, id: string): void {
471 this.defer(
472 (async () => {
473 const row = await this.messageRow(place.channel.id, id);
474 if (!row) return;
475 const [message] = await this.toMessages(place.slug, place.workspace, [row]);
476 await this.room(place.channel.id).broadcast({ type: "message.updated", message });
477 })().catch((error) => console.error("chat could not rebroadcast", id, error)),
478 );
479 }
480
481 // ── The sidebar ─────────────────────────────────────────────────────────
482
483 /**
484 * Puts a person in the workspace's #general, once. A new workspace has
485 * no channels, so the first sidebar anyone in it asks for creates
486 * #general; and everyone who asks for the sidebar is put in it the first
487 * time, so a new workspace has somewhere to talk and a new member lands
488 * where everyone is. Someone who leaves it is not put back
489 * (`general_joined`). A private channel someone named `general` is never
490 * joined this way.
491 */
492 private async ensureGeneral(workspace: Workspace, me: string): Promise<void> {
493 const seen = await this.db
494 .prepare("SELECT 1 FROM general_joined WHERE workspace_id = ? AND principal = ?")
495 .bind(workspace.id, me)
496 .first();
497 if (seen) return;
498 const at = now();
499 await this.db
500 .prepare(
501 "INSERT OR IGNORE INTO channels (id, workspace_id, kind, name, topic, private, created_by, created_at) VALUES (?, ?, 'channel', ?, ?, 0, ?, ?)",
502 )
503 .bind(newId("chn"), workspace.id, GENERAL, "Anything and everything for the whole workspace.", me, at)
504 .run();
505 const general = await this.db
506 .prepare("SELECT * FROM channels WHERE workspace_id = ? AND name = ?")
507 .bind(workspace.id, GENERAL)
508 .first<ChannelRow>();
509 const statements = [
510 this.db
511 .prepare("INSERT OR IGNORE INTO general_joined (workspace_id, principal, joined_at) VALUES (?, ?, ?)")
512 .bind(workspace.id, me, at),
513 ];
514 // The workspace's default channels (#general unless its owners chose
515 // others): public and not archived only, whatever was kept.
516 const settings = await this.settings(workspace, general?.id ?? null);
517 const defaults = settings.default_channels.length
518 ? await this.db
519 .prepare(
520 "SELECT * FROM channels WHERE workspace_id = ? AND kind = 'channel' AND private = 0 AND archived_at IS NULL AND id IN (SELECT value FROM json_each(?))",
521 )
522 .bind(workspace.id, JSON.stringify(settings.default_channels))
523 .all<ChannelRow>()
524 : { results: [] as ChannelRow[] };
525 for (const channel of defaults.results) {
526 statements.push(this.joinStatement(channel.id, me, channel.created_by === me ? "owner" : "member", at));
527 }
528 await this.db.batch(statements);
529 }
530
531 /**
532 * Adds a member. Someone joining starts with everything already said
533 * read, so a long channel does not greet them with its whole history as
534 * unread.
535 */
536 private joinStatement(channelId: string, principal: string, role: "owner" | "member", at: string): D1PreparedStatement {
537 return this.db
538 .prepare(
539 "INSERT OR IGNORE INTO channel_members (channel_id, principal, role, last_read_id, joined_at) VALUES (?1, ?2, ?3, (SELECT MAX(id) FROM messages WHERE channel_id = ?1), ?4)",
540 )
541 .bind(channelId, principal, role, at);
542 }
543
544 async sidebar(a: { workspace: string; viewer: Viewer }): Promise<Result<ChatSidebar>> {
545 const found = await this.viewerWorkspace(a.workspace, a.viewer);
546 if (!found.ok) return found;
547 const workspace = found.value;
548 const slug = a.workspace.toLowerCase();
549 const me = userKey(a.viewer!);
550 await this.ensureGeneral(workspace, me);
551
552 const [joined, unread, dmOthers, browsable] = await Promise.all([
553 this.db
554 .prepare(
555 `SELECT c.*, m.starred, m.muted, m.last_read_id
556 FROM channel_members m JOIN channels c ON c.id = m.channel_id
557 WHERE m.principal = ? AND c.workspace_id = ? AND c.archived_at IS NULL`,
558 )
559 .bind(me, workspace.id)
560 .all<ChannelRow & { starred: number; muted: number; last_read_id: string | null }>(),
561 this.db
562 .prepare(
563 `SELECT msg.channel_id, msg.id, msg.author, msg.mentions
564 FROM channel_members m
565 JOIN channels c ON c.id = m.channel_id
566 JOIN messages msg ON msg.channel_id = m.channel_id AND msg.id > COALESCE(m.last_read_id, '')
567 WHERE m.principal = ?1 AND c.workspace_id = ?2 AND c.archived_at IS NULL
568 AND msg.deleted_at IS NULL AND msg.author != ?1
569 LIMIT ${MAX_UNREAD_ROWS}`,
570 )
571 .bind(me, workspace.id)
572 .all<UnreadRow>(),
573 this.db
574 .prepare(
575 `SELECT o.channel_id, o.principal
576 FROM channel_members m
577 JOIN channels c ON c.id = m.channel_id AND c.kind = 'dm'
578 JOIN channel_members o ON o.channel_id = m.channel_id AND o.principal != m.principal
579 WHERE m.principal = ? AND c.workspace_id = ? AND c.archived_at IS NULL
580 ORDER BY o.joined_at, o.principal`,
581 )
582 .bind(me, workspace.id)
583 .all<{ channel_id: string; principal: string }>(),
584 this.db
585 .prepare(
586 `SELECT COUNT(*) AS n FROM channels c
587 WHERE c.workspace_id = ? AND c.kind = 'channel' AND c.private = 0 AND c.archived_at IS NULL
588 AND NOT EXISTS (SELECT 1 FROM channel_members m WHERE m.channel_id = c.id AND m.principal = ?)`,
589 )
590 .bind(workspace.id, me)
591 .first<{ n: number }>(),
592 ]);
593
594 const others = new Map<string, string[]>();
595 for (const row of dmOthers.results) others.set(row.channel_id, [...(others.get(row.channel_id) ?? []), row.principal]);
596 const profiles = await this.profiles(slug, workspace, [me, ...dmOthers.results.map((r) => r.principal)]);
597 const self = profiles.get(me) ?? null;
598 const counts = tally(
599 unread.results,
600 me,
601 self?.name ?? a.viewer!.username,
602 new Map(joined.results.map((row) => [row.id, row.last_read_id])),
603 );
604
605 const entries: ChatSidebarEntry[] = joined.results.map((row) => {
606 const faces = (others.get(row.id) ?? []).map((key) => profiles.get(key)!).filter(Boolean);
607 const count = counts.get(row.id) ?? { unread: 0, mentions: 0 };
608 return {
609 channel: toChannel(row),
610 title: row.kind === "dm" ? dmTitle(faces, self) : (row.name ?? ""),
611 others: row.kind === "dm" ? faces.slice(0, DM_FACES) : [],
612 starred: !!row.starred,
613 muted: !!row.muted,
614 unread: count.unread,
615 mentions: count.mentions,
616 };
617 });
618 const settings = await this.settings(workspace);
619 return ok({ entries: sidebarOrder(entries), browsable: browsable?.n ?? 0, can: permissionsFor(settings, roleOf(a.viewer, slug)) });
620 }
621
622 // ── Channels ────────────────────────────────────────────────────────────
623
624 async channel(a: { workspace: string; channel_id: string; viewer: Viewer }): Promise<Result<ChannelDetail>> {
625 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
626 if (!found.ok) return found;
627 const { slug, workspace, channel, member } = found.value;
628 const [rows, settings] = await Promise.all([
629 this.db
630 .prepare("SELECT * FROM channel_members WHERE channel_id = ? ORDER BY joined_at, principal")
631 .bind(channel.id)
632 .all<MemberRow>(),
633 this.settings(workspace),
634 ]);
635 const profiles = await this.profiles(slug, workspace, rows.results.map((r) => r.principal));
636 return ok({
637 channel: toChannel(channel),
638 can_manage: channel.kind === "channel" && mayManageChannel(settings, roleOf(a.viewer, slug), member?.role ?? null),
639 members: rows.results.map((row) => ({
640 channel_id: row.channel_id,
641 member: profiles.get(row.principal)!,
642 role: row.role,
643 starred: !!row.starred,
644 muted: !!row.muted,
645 last_read_id: row.last_read_id,
646 joined_at: row.joined_at,
647 })),
648 });
649 }
650
651 /** A channel by its name, as the site's URLs name them; read like `channel`. */
652 async channelByName(a: { workspace: string; name: string; viewer: Viewer }): Promise<Result<ChannelDetail>> {
653 const found = await this.viewerWorkspace(a.workspace, a.viewer);
654 if (!found.ok) return found;
655 const named = channelName(a.name ?? "");
656 if (!named.ok) return fail("not_found", "No such channel.");
657 const row = await this.db
658 .prepare("SELECT id FROM channels WHERE workspace_id = ? AND name = ?")
659 .bind(found.value.id, named.name)
660 .first<{ id: string }>();
661 if (!row) return fail("not_found", "No such channel.");
662 // A private channel the viewer is not in stays not found there.
663 return this.channel({ workspace: a.workspace, channel_id: row.id, viewer: a.viewer });
664 }
665
666 /**
667 * Every public channel, and the private ones the viewer is in; a private
668 * one they are not in stays unseen. With `archived`, the archived ones.
669 */
670 async browse(a: { workspace: string; viewer: Viewer; archived?: boolean }): Promise<Result<Channel[]>> {
671 const found = await this.viewerWorkspace(a.workspace, a.viewer);
672 if (!found.ok) return found;
673 const rows = await this.db
674 .prepare(
675 `SELECT * FROM channels c
676 WHERE c.workspace_id = ?1 AND c.kind = 'channel'
677 AND (c.archived_at IS NULL) = (?3 = 0)
678 AND (c.private = 0 OR EXISTS (SELECT 1 FROM channel_members m WHERE m.channel_id = c.id AND m.principal = ?2))
679 ORDER BY c.name`,
680 )
681 .bind(found.value.id, userKey(a.viewer!), a.archived === true ? 1 : 0)
682 .all<ChannelRow>();
683 return ok(rows.results.map(toChannel));
684 }
685
686 async createChannel(a: { workspace: string; viewer: Viewer; input: NewChannel }): Promise<Result<Channel>> {
687 const found = await this.viewerWorkspace(a.workspace, a.viewer);
688 if (!found.ok) return found;
689 const workspace = found.value;
690 const named = channelName(a.input?.name ?? "");
691 if (!named.ok) return fail("invalid", named.message);
692 const topic = typeof a.input?.topic === "string" ? a.input.topic.trim() : "";
693 if (topic.length > MAX_TOPIC) return fail("invalid", `A topic is at most ${MAX_TOPIC} characters.`);
694 const isPrivate = !!a.input?.private;
695 if (!mayCreateChannel(await this.settings(workspace), roleOf(a.viewer, a.workspace), isPrivate)) {
696 return fail("forbidden", `Only owners can create ${isPrivate ? "private" : "public"} channels in this workspace.`);
697 }
698 const taken = await this.db
699 .prepare("SELECT 1 FROM channels WHERE workspace_id = ? AND name = ?")
700 .bind(workspace.id, named.name)
701 .first();
702 if (taken) return fail("conflict", `#${named.name} already exists.`);
703 const me = userKey(a.viewer!);
704 const row: ChannelRow = {
705 id: newId("chn"),
706 workspace_id: workspace.id,
707 kind: "channel",
708 name: named.name,
709 topic: topic || null,
710 private: isPrivate ? 1 : 0,
711 dm_key: null,
712 created_by: me,
713 created_at: now(),
714 archived_at: null,
715 last_message_at: null,
716 };
717 try {
718 await this.db.batch([
719 this.db
720 .prepare(
721 "INSERT INTO channels (id, workspace_id, kind, name, topic, private, created_by, created_at) VALUES (?, ?, 'channel', ?, ?, ?, ?, ?)",
722 )
723 .bind(row.id, row.workspace_id, row.name, row.topic, row.private, me, row.created_at),
724 this.joinStatement(row.id, me, "owner", row.created_at),
725 ]);
726 } catch (error) {
727 if (String(error).includes("UNIQUE")) return fail("conflict", `#${named.name} already exists.`);
728 throw error;
729 }
730 return ok(toChannel(row));
731 }
732
733 async openDm(a: { workspace: string; viewer: Viewer; members: Principal[] }): Promise<Result<Channel>> {
734 const found = await this.viewerWorkspace(a.workspace, a.viewer);
735 if (!found.ok) return found;
736 const workspace = found.value;
737 const slug = a.workspace.toLowerCase();
738 const me = userKey(a.viewer!);
739 const asked = Array.isArray(a.members) ? a.members : [];
740 const principals: Principal[] = [];
741 for (const member of asked) {
742 const p = member && parsePrincipalKey(`${member.kind}:${member.id}`);
743 if (!p) return fail("invalid", "Each member is a person or an agent, by id.");
744 principals.push(p);
745 }
746 const members = dmMembers(me, principals.map(principalKey));
747 if (members.length > MAX_DM_MEMBERS) {
748 return fail("invalid", `A direct message has at most ${MAX_DM_MEMBERS} people and agents. Make a private channel instead.`);
749 }
750 const key = dmKey(members);
751 const existing = await this.db
752 .prepare("SELECT * FROM channels WHERE workspace_id = ? AND dm_key = ?")
753 .bind(workspace.id, key)
754 .first<ChannelRow>();
755 if (existing) return ok(toChannel(existing));
756
757 for (const member of members) {
758 if (member === me) continue;
759 const p = parsePrincipalKey(member)!;
760 if (!(await this.belongs(slug, workspace, p))) {
761 return fail("not_found", p.kind === "agent" ? "No such agent in this workspace." : "That person is not in this workspace.");
762 }
763 }
764 const at = now();
765 await this.db
766 .prepare(
767 "INSERT OR IGNORE INTO channels (id, workspace_id, kind, private, dm_key, created_by, created_at) VALUES (?, ?, 'dm', 1, ?, ?, ?)",
768 )
769 .bind(newId("chn"), workspace.id, key, me, at)
770 .run();
771 // Read back by key: if two people opened it at once, both get the one that won.
772 const channel = await this.db
773 .prepare("SELECT * FROM channels WHERE workspace_id = ? AND dm_key = ?")
774 .bind(workspace.id, key)
775 .first<ChannelRow>();
776 if (!channel) return fail("conflict", "The direct message could not be opened. Try again.");
777 await this.db.batch(members.map((member) => this.joinStatement(channel.id, member, "member", at)));
778 return ok(toChannel(channel));
779 }
780
781 async join(a: { workspace: string; channel_id: string; viewer: Viewer }): Promise<Result<null>> {
782 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
783 if (!found.ok) return found;
784 const { channel, member } = found.value;
785 if (member) return ok(null);
786 if (channel.archived_at) return fail("invalid", "This channel is archived.");
787 await this.joinStatement(channel.id, userKey(a.viewer!), "member", now()).run();
788 return ok(null);
789 }
790
791 async leave(a: { workspace: string; channel_id: string; viewer: Viewer }): Promise<Result<null>> {
792 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
793 if (!found.ok) return found;
794 const { channel, member } = found.value;
795 if (!member) return ok(null);
796 if (channel.kind === "dm") return fail("invalid", "A direct message can't be left. Mute it instead.");
797 const me = userKey(a.viewer!);
798 await this.db.prepare("DELETE FROM channel_members WHERE channel_id = ? AND principal = ?").bind(channel.id, me).run();
799 // Out of a private channel, they may no longer read it, live either.
800 if (channel.private) {
801 this.defer(this.room(channel.id).drop(me).catch((error: unknown) => console.error("chat could not drop", me, error)));
802 }
803 return ok(null);
804 }
805
806 async invite(a: { workspace: string; channel_id: string; viewer: Viewer; member: Principal }): Promise<Result<null>> {
807 const found = await this.place(a.workspace, a.channel_id, a.viewer, "member");
808 if (!found.ok) return found;
809 const { slug, workspace, channel } = found.value;
810 if (channel.kind === "dm") return fail("invalid", "People can't be added to a direct message. Start a new one with everyone in it.");
811 if (channel.archived_at) return fail("invalid", "This channel is archived.");
812 const p = a.member && parsePrincipalKey(`${a.member.kind}:${a.member.id}`);
813 if (!p) return fail("invalid", "Invite a person or an agent, by id.");
814 if (!(await this.belongs(slug, workspace, p))) {
815 return fail("not_found", p.kind === "agent" ? "No such agent in this workspace." : "That person is not in this workspace.");
816 }
817 await this.joinStatement(channel.id, principalKey(p), "member", now()).run();
818 return ok(null);
819 }
820
821 async setPreferences(a: {
822 workspace: string;
823 channel_id: string;
824 viewer: Viewer;
825 prefs: { starred?: boolean; muted?: boolean };
826 }): Promise<Result<null>> {
827 const found = await this.place(a.workspace, a.channel_id, a.viewer, "member");
828 if (!found.ok) return found;
829 const starred = typeof a.prefs?.starred === "boolean" ? (a.prefs.starred ? 1 : 0) : null;
830 const muted = typeof a.prefs?.muted === "boolean" ? (a.prefs.muted ? 1 : 0) : null;
831 await this.db
832 .prepare(
833 "UPDATE channel_members SET starred = COALESCE(?, starred), muted = COALESCE(?, muted) WHERE channel_id = ? AND principal = ?",
834 )
835 .bind(starred, muted, found.value.channel.id, userKey(a.viewer!))
836 .run();
837 // Notify: the badge counts the conversation again, or leaves it out, in every tab.
838 if (muted !== null) {
839 this.defer(
840 notifyMuted(this.env.NOTIFY, { slug: found.value.slug, channel_id: found.value.channel.id, user_id: a.viewer!.id, muted: !!muted }).catch((error) =>
841 console.error("chat could not notify a mute", error),
842 ),
843 );
844 }
845 return ok(null);
846 }
847
848 // ── What owners decide (src/settings.ts) ────────────────────────────────
849
850 /** The row of a workspace's chat settings, read once per request. */
851 private settingsRow(workspace: Workspace): Promise<SettingsRow | null> {
852 let found = this.settingsRows.get(workspace.id);
853 if (!found) {
854 found = this.db
855 .prepare("SELECT emoji_upload, public_channels, private_channels, manage_channels, default_channels FROM chat_settings WHERE workspace_id = ?")
856 .bind(workspace.id)
857 .first<SettingsRow>();
858 this.settingsRows.set(workspace.id, found);
859 }
860 return found;
861 }
862
863 /**
864 * A workspace's chat settings. Default channels never chosen are
865 * #general, looked up unless `generalId` is given.
866 */
867 private async settings(workspace: Workspace, generalId?: string | null): Promise<ChatSettings> {
868 const row = await this.settingsRow(workspace);
869 if (row?.default_channels != null || generalId !== undefined) return settingsOf(row, generalId ?? null);
870 const general = await this.db
871 .prepare("SELECT id FROM channels WHERE workspace_id = ? AND name = ? AND kind = 'channel'")
872 .bind(workspace.id, GENERAL)
873 .first<{ id: string }>();
874 return settingsOf(row, general?.id ?? null);
875 }
876
877 async chatSettings(a: { workspace: string; viewer: Viewer }): Promise<Result<ChatSettingsView>> {
878 const found = await this.viewerWorkspace(a.workspace, a.viewer);
879 if (!found.ok) return found;
880 const workspace = found.value;
881 const [settings, channels] = await Promise.all([
882 this.settings(workspace),
883 this.db
884 .prepare("SELECT * FROM channels WHERE workspace_id = ? AND kind = 'channel' AND private = 0 AND archived_at IS NULL ORDER BY name")
885 .bind(workspace.id)
886 .all<ChannelRow>(),
887 ]);
888 return ok({ settings, can: permissionsFor(settings, roleOf(a.viewer, a.workspace)), channels: channels.results.map(toChannel) });
889 }
890
891 async setChatSettings(a: { workspace: string; viewer: Viewer; change: Partial<ChatSettings> }): Promise<Result<ChatSettings>> {
892 const found = await this.viewerWorkspace(a.workspace, a.viewer);
893 if (!found.ok) return found;
894 const workspace = found.value;
895 if (roleOf(a.viewer, a.workspace) !== "owner") return fail("forbidden", "Only owners can change the workspace's chat settings.");
896 const checked = settingsChange(a.change);
897 if (!checked.ok) return fail("invalid", checked.message);
898 const change = checked.change;
899 if (change.default_channels) {
900 // Only the workspace's public channels that are not archived: a private
901 // one would put people somewhere they were never invited.
902 const rows = await this.db
903 .prepare(
904 "SELECT id FROM channels WHERE workspace_id = ? AND kind = 'channel' AND private = 0 AND archived_at IS NULL AND id IN (SELECT value FROM json_each(?))",
905 )
906 .bind(workspace.id, JSON.stringify(change.default_channels))
907 .all<{ id: string }>();
908 const open = new Set(rows.results.map((r) => r.id));
909 if (change.default_channels.some((id) => !open.has(id))) return fail("invalid", "Default channels must be public channels that are not archived.");
910 }
911 const next: ChatSettings = { ...(await this.settings(workspace)), ...change };
912 const row = rowFor(next);
913 await this.db
914 .prepare(
915 `INSERT INTO chat_settings (workspace_id, emoji_upload, public_channels, private_channels, manage_channels, default_channels)
916 VALUES (?1, ?2, ?3, ?4, ?5, ?6)
917 ON CONFLICT (workspace_id) DO UPDATE SET emoji_upload = ?2, public_channels = ?3, private_channels = ?4, manage_channels = ?5, default_channels = ?6`,
918 )
919 .bind(workspace.id, row.emoji_upload, row.public_channels, row.private_channels, row.manage_channels, row.default_channels)
920 .run();
921 this.settingsRows.delete(workspace.id);
922 return ok(next);
923 }
924
925 /**
926 * Renames, archives or unarchives a channel, or changes its topic. The
927 * workspace's `manage_channels` setting says who may do the first three;
928 * any member may change the topic. #general stays #general, and stays.
929 */
930 async updateChannel(a: { workspace: string; channel_id: string; viewer: Viewer; change: ChannelChange }): Promise<Result<Channel>> {
931 // Archived or not, a channel is found the same way; a private one only by its members.
932 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
933 if (!found.ok) return found;
934 const { slug, workspace, channel, member } = found.value;
935 if (channel.kind === "dm") return fail("invalid", "A direct message has no name or topic to change.");
936 const change = a.change ?? {};
937 const next: ChannelRow = { ...channel };
938 if (change.name !== undefined || change.archived !== undefined) {
939 const settings = await this.settings(workspace);
940 if (!mayManageChannel(settings, roleOf(a.viewer, slug), member?.role ?? null)) {
941 return fail(
942 "forbidden",
943 settings.manage_channels === "owners"
944 ? "Only workspace owners can rename or archive channels here."
945 : "Only this channel's owners and workspace owners can rename or archive it.",
946 );
947 }
948 if (channel.name === GENERAL) return fail("invalid", "#general is where everyone is: it can't be renamed or archived.");
949 }
950 if (change.name !== undefined) {
951 const named = channelName(String(change.name ?? ""));
952 if (!named.ok) return fail("invalid", named.message);
953 if (named.name !== channel.name) {
954 const taken = await this.db
955 .prepare("SELECT 1 FROM channels WHERE workspace_id = ? AND name = ? AND id != ?")
956 .bind(workspace.id, named.name, channel.id)
957 .first();
958 if (taken || named.name === GENERAL) return fail("conflict", `#${named.name} already exists.`);
959 }
960 next.name = named.name;
961 }
962 if (change.topic !== undefined) {
963 if (!member) return fail("forbidden", `Join #${channel.name} first.`);
964 const topic = typeof change.topic === "string" ? change.topic.trim() : "";
965 if (topic.length > MAX_TOPIC) return fail("invalid", `A topic is at most ${MAX_TOPIC} characters.`);
966 next.topic = topic || null;
967 }
968 if (change.archived !== undefined) {
969 next.archived_at = change.archived ? (channel.archived_at ?? now()) : null;
970 } else if (channel.archived_at && (change.name !== undefined || change.topic !== undefined)) {
971 return fail("invalid", "This channel is archived. Unarchive it first.");
972 }
973 try {
974 await this.db
975 .prepare("UPDATE channels SET name = ?, topic = ?, archived_at = ? WHERE id = ?")
976 .bind(next.name, next.topic, next.archived_at, channel.id)
977 .run();
978 } catch (error) {
979 if (String(error).includes("UNIQUE")) return fail("conflict", `#${next.name} already exists.`);
980 throw error;
981 }
982 const updated = toChannel(next);
983 this.broadcast(channel.id, { type: "channel.updated", channel: updated });
984 return ok(updated);
985 }
986
987 // ── Messages ────────────────────────────────────────────────────────────
988
989 /**
990 * Newest first, a page at a time. With `thread_root`, that thread's
991 * replies, and the message they reply to as the oldest once the page
992 * reaches the start of the thread; without, the channel's top-level
993 * messages. A deleted message stays only while replies hang off it.
994 */
995 async messages(a: {
996 workspace: string;
997 channel_id: string;
998 viewer: Viewer;
999 before?: string | null;
1000 after?: string | null;
1001 limit?: number | null;
1002 thread_root?: string | null;
1003 }): Promise<Result<MessagePage>> {
1004 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
1005 if (!found.ok) return found;
1006 const { slug, workspace, channel } = found.value;
1007 const size = pageSize(a.limit);
1008 // Ids are lowercase letters, digits and `_`, so `~` sorts after every
1009 // one: the first page reads from the newest as a range of the index.
1010 const before = typeof a.before === "string" && a.before ? a.before : "~";
1011 const root = typeof a.thread_root === "string" && a.thread_root ? a.thread_root : null;
1012 if (typeof a.after === "string" && a.after) {
1013 // Catching up after a reconnect: what came after, oldest first. Deleted
1014 // ones too, so the client drops them; read one past the page to know
1015 // whether there is more.
1016 const newer = root
1017 ? await this.db
1018 .prepare("SELECT * FROM messages WHERE thread_root = ? AND channel_id = ? AND id > ? ORDER BY id LIMIT ?")
1019 .bind(root, channel.id, a.after, size + 1)
1020 .all<MessageRow>()
1021 : await this.db
1022 .prepare("SELECT * FROM messages WHERE channel_id = ? AND thread_root IS NULL AND id > ? ORDER BY id LIMIT ?")
1023 .bind(channel.id, a.after, size + 1)
1024 .all<MessageRow>();
1025 const page = pageOf(newer.results, size);
1026 return ok({ messages: await this.toMessages(slug, workspace, page.rows, userKey(a.viewer!)), older: null, newer: page.older });
1027 }
1028 const rows = root
1029 ? await this.db
1030 .prepare(
1031 `SELECT * FROM messages WHERE thread_root = ?1 AND channel_id = ?2 AND id < ?3
1032 ORDER BY id DESC LIMIT ?4`,
1033 )
1034 .bind(root, channel.id, before, size + 1)
1035 .all<MessageRow>()
1036 : await this.db
1037 .prepare(
1038 `SELECT * FROM messages
1039 WHERE channel_id = ?1 AND thread_root IS NULL AND id < ?2
1040 AND (deleted_at IS NULL OR reply_count > 0)
1041 ORDER BY id DESC LIMIT ?3`,
1042 )
1043 .bind(channel.id, before, size + 1)
1044 .all<MessageRow>();
1045 const page = pageOf(rows.results, size);
1046 let list = page.rows;
1047 if (root && page.older === null) {
1048 const first = await this.messageRow(channel.id, root);
1049 if (first) list = [...list, first];
1050 }
1051 return ok({ messages: await this.toMessages(slug, workspace, list, userKey(a.viewer!)), older: page.older });
1052 }
1053
1054 /**
1055 * Writes a message and everything that follows from it: the thread's
1056 * reply count, the channel's last activity, the author's own read mark,
1057 * the meter; then, after answering, tells the room and wakes the agents
1058 * it is for.
1059 */
1060 private async write(
1061 place: Place,
1062 author: string,
1063 input: { body: string; card: MessageCard | null; thread_root: string | null },
1064 chain: Chain<AskerAccess>,
1065 ): Promise<Result<ChatMessage>> {
1066 const { channel, workspace } = place;
1067 if (channel.archived_at) return fail("invalid", "This channel is archived.");
1068 let threadRoot: string | null = null;
1069 if (input.thread_root) {
1070 const root = await this.messageRow(channel.id, input.thread_root);
1071 if (!root || (root.deleted_at && !root.reply_count)) return fail("not_found", "No such message to reply to.");
1072 // A reply to a reply goes in the same thread.
1073 threadRoot = root.thread_root ?? root.id;
1074 }
1075 const at = now();
1076 const handles = mentionedHandles(input.body);
1077 const row: MessageRow = {
1078 id: newId("msg"),
1079 channel_id: channel.id,
1080 author,
1081 kind: input.card ? "card" : "text",
1082 body: input.body,
1083 card: input.card ? JSON.stringify(input.card) : null,
1084 mentions: mentionsColumn(handles),
1085 thread_root: threadRoot,
1086 reply_count: 0,
1087 last_reply_at: null,
1088 created_at: at,
1089 edited_at: null,
1090 deleted_at: null,
1091 };
1092 const statements = [
1093 this.db
1094 .prepare(
1095 "INSERT INTO messages (id, channel_id, author, kind, body, card, mentions, thread_root, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
1096 )
1097 .bind(row.id, row.channel_id, author, row.kind, row.body, row.card, row.mentions, threadRoot, at),
1098 this.db.prepare("UPDATE channels SET last_message_at = ? WHERE id = ?").bind(at, channel.id),
1099 // What you wrote, you have read.
1100 this.db
1101 .prepare("UPDATE channel_members SET last_read_id = ? WHERE channel_id = ? AND principal = ?")
1102 .bind(row.id, channel.id, author),
1103 this.db
1104 .prepare(
1105 `INSERT INTO chat_meter (workspace_id, day, messages, bytes) VALUES (?, ?, 1, ?)
1106 ON CONFLICT (workspace_id, day) DO UPDATE SET messages = messages + 1, bytes = bytes + excluded.bytes`,
1107 )
1108 .bind(workspace.id, meterDay(at), bytesOf(row.body) + bytesOf(row.card ?? "")),
1109 ];
1110 if (threadRoot) {
1111 statements.push(
1112 this.db
1113 .prepare("UPDATE messages SET reply_count = reply_count + 1, last_reply_at = ? WHERE id = ?")
1114 .bind(at, threadRoot),
1115 );
1116 }
1117 await this.db.batch(statements);
1118
1119 const [message] = await this.toMessages(place.slug, workspace, [row]);
1120 this.broadcast(channel.id, { type: "message.created", message });
1121 if (threadRoot) this.rebroadcast(place, threadRoot);
1122 this.defer(
1123 this.wake(place, row, handles, chain).catch((error) => console.error("chat could not hand", row.id, "to agents", error)),
1124 );
1125 // Notify: counts for everyone in the conversation, a notification for those it is for.
1126 this.defer(
1127 notifyMessage(this.db, this.env.NOTIFY, (keys) => this.profiles(place.slug, workspace, keys), { slug: place.slug, channel, row, handles }).catch(
1128 (error) => console.error("chat could not notify about", row.id, error),
1129 ),
1130 );
1131 return ok(message);
1132 }
1133
1134 /** Hands a new message to the agents it is for (src/delivery.ts). */
1135 private async wake(place: Place, row: MessageRow, handles: string[], chain: Chain<AskerAccess>): Promise<void> {
1136 const { channel, workspace } = place;
1137 // Only an agent's message mentioning someone, or a person's, can wake anyone.
1138 if (row.author.startsWith("agent:") && !handles.length) return;
1139 if (channel.kind === "channel" && !handles.length) return;
1140 const members = await this.db
1141 .prepare("SELECT principal FROM channel_members WHERE channel_id = ? AND principal LIKE 'agent:%'")
1142 .bind(channel.id)
1143 .all<{ principal: string }>();
1144 const ids = members.results.map((m) => m.principal.slice("agent:".length));
1145 let found = await this.agentsById(ids);
1146 const orchestratorIsMember = [...found.values()].some((agent) => !!agent?.builtin && agent.workspace_id === workspace.id && !agent.archived_at);
1147 if (addsOrchestrator({ channelKind: channel.kind, mentioned: handles, orchestratorIsMember })) {
1148 // Mentioning @g1t brings it in: every workspace has it, nobody invites it.
1149 const builtin = await workspaceAgentsClient(this.env.AGENTS)
1150 .builtin(place.slug, workspace.id)
1151 .catch((error: unknown) => {
1152 console.error("chat could not find @g1t for", place.slug, error);
1153 return null;
1154 });
1155 if (builtin?.ok) {
1156 await this.joinStatement(channel.id, `agent:${builtin.value.id}`, "member", now()).run();
1157 this.agents.set(builtin.value.id, builtin.value);
1158 ids.push(builtin.value.id);
1159 found = await this.agentsById(ids);
1160 }
1161 }
1162 if (!ids.length) return;
1163 const agents = [...found.values()].filter(
1164 (agent): agent is WorkspaceAgent => !!agent && agent.workspace_id === workspace.id && !agent.archived_at,
1165 );
1166 const wakes = deliveries({
1167 author: row.author,
1168 hops: chain.hops,
1169 channelKind: channel.kind,
1170 agents: agents.map((agent) => ({ id: agent.id, handle: agent.handle })),
1171 mentioned: handles,
1172 notTo: sender(chain.chain),
1173 });
1174 const client = workspaceAgentsClient(this.env.AGENTS);
1175 await Promise.all(
1176 wakes.map((wake) =>
1177 client
1178 .deliver(
1179 delivery(
1180 {
1181 workspace: place.slug,
1182 workspace_id: workspace.id,
1183 channel_id: channel.id,
1184 channel_kind: channel.kind,
1185 channel_name: channel.kind === "dm" ? null : channel.name,
1186 },
1187 wake,
1188 row,
1189 chain,
1190 ) satisfies AgentDelivery,
1191 )
1192 .then((result) => {
1193 if (!result.ok) console.error("agents refused delivery to", wake.agent_id, result.error.message);
1194 })
1195 .catch((error) => console.error("chat could not deliver to", wake.agent_id, error)),
1196 ),
1197 );
1198 }
1199
1200 async post(a: { workspace: string; channel_id: string; viewer: Viewer; message: PostMessage }): Promise<Result<ChatMessage>> {
1201 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
1202 if (!found.ok) return found;
1203 const body = messageBody(a.message?.body);
1204 if (!body.ok) return fail("invalid", body.message);
1205 const me = userKey(a.viewer!);
1206 let place = found.value;
1207 if (!place.member) {
1208 // Saying something in a public channel joins it, as reading does not.
1209 if (place.channel.archived_at) return fail("invalid", "This channel is archived.");
1210 await this.joinStatement(place.channel.id, me, "member", now()).run();
1211 place = { ...place, member: { channel_id: place.channel.id, principal: me } as MemberRow };
1212 }
1213 return this.write(
1214 place,
1215 me,
1216 { body: body.body, card: null, thread_root: a.message?.thread_root ?? null },
1217 // A person's message starts a chain.
1218 { hops: 0, asked_by: a.viewer!.id, asker: askerAccess(a.viewer!, a.workspace), chain: [] },
1219 );
1220 }
1221
1222 async edit(a: { workspace: string; channel_id: string; viewer: Viewer; id: string; body: string }): Promise<Result<ChatMessage>> {
1223 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
1224 if (!found.ok) return found;
1225 const place = found.value;
1226 const row = await this.messageRow(place.channel.id, a.id);
1227 if (!row || row.deleted_at) return fail("not_found", "No such message.");
1228 if (row.author !== userKey(a.viewer!)) return fail("forbidden", "Only its author can edit a message.");
1229 const body = messageBody(a.body, !!row.card);
1230 if (!body.ok) return fail("invalid", body.message);
1231 const at = now();
1232 const mentions = mentionsColumn(mentionedHandles(body.body));
1233 await this.db
1234 .prepare("UPDATE messages SET body = ?, mentions = ?, edited_at = ? WHERE id = ?")
1235 .bind(body.body, mentions, at, row.id)
1236 .run();
1237 const [message] = await this.toMessages(place.slug, place.workspace, [{ ...row, body: body.body, mentions, edited_at: at }]);
1238 this.broadcast(place.channel.id, { type: "message.updated", message });
1239 return ok(message);
1240 }
1241
1242 /** Deletes a message, keeping its place so its thread still hangs together. Its author or a workspace owner may. */
1243 async remove(a: { workspace: string; channel_id: string; viewer: Viewer; id: string }): Promise<Result<null>> {
1244 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
1245 if (!found.ok) return found;
1246 const place = found.value;
1247 const row = await this.messageRow(place.channel.id, a.id);
1248 if (!row || row.deleted_at) return fail("not_found", "No such message.");
1249 if (row.author !== userKey(a.viewer!) && !isOwner(a.viewer, a.workspace)) {
1250 return fail("forbidden", "Only its author or a workspace owner can delete a message.");
1251 }
1252 const statements = [
1253 this.db
1254 .prepare("UPDATE messages SET deleted_at = ?, body = '', card = NULL, mentions = '' WHERE id = ?")
1255 .bind(now(), row.id),
1256 ];
1257 if (row.thread_root) {
1258 statements.push(
1259 this.db.prepare("UPDATE messages SET reply_count = MAX(reply_count - 1, 0) WHERE id = ?").bind(row.thread_root),
1260 );
1261 }
1262 await this.db.batch(statements);
1263 this.broadcast(place.channel.id, { type: "message.deleted", channel_id: place.channel.id, id: row.id });
1264 if (row.thread_root) this.rebroadcast(place, row.thread_root);
1265 return ok(null);
1266 }
1267
1268 async markRead(a: { workspace: string; channel_id: string; viewer: Viewer; id: string }): Promise<Result<null>> {
1269 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
1270 if (!found.ok) return found;
1271 const { channel, member } = found.value;
1272 // Someone reading a public channel they have not joined keeps no read state.
1273 if (!member) return ok(null);
1274 const id = typeof a.id === "string" ? a.id : "";
1275 if (!id) return fail("invalid", "Say which message was read.");
1276 // Only forward: reading an old thread does not mark newer messages unread.
1277 const changed = await this.db
1278 .prepare(
1279 "UPDATE channel_members SET last_read_id = ?1 WHERE channel_id = ?2 AND principal = ?3 AND (last_read_id IS NULL OR last_read_id < ?1)",
1280 )
1281 .bind(id, channel.id, member.principal)
1282 .run();
1283 if (changed.meta.changes) {
1284 this.broadcast(channel.id, {
1285 type: "read",
1286 channel_id: channel.id,
1287 principal: { kind: "user", id: a.viewer!.id },
1288 last_read_id: id,
1289 });
1290 // Notify: the read drops the counts in every tab of theirs.
1291 this.defer(
1292 notifyRead(this.db, this.env.NOTIFY, { slug: found.value.slug, channel_id: channel.id, user_id: a.viewer!.id, username: a.viewer!.username, last_read_id: id }).catch(
1293 (error) => console.error("chat could not notify a read", error),
1294 ),
1295 );
1296 }
1297 return ok(null);
1298 }
1299
1300 // ── Agents ──────────────────────────────────────────────────────────────
1301
1302 /** The channel and agent for an agent's call: the agent must be of the workspace and in the channel. */
1303 private async agentPlace(slug: string, channelId: string, agentId: string): Promise<Result<{ place: Place; agent: WorkspaceAgent }>> {
1304 const workspace = await this.workspace(String(slug ?? ""));
1305 if (!workspace) return fail("not_found", "No such workspace.");
1306 const agent = await this.liveAgent(workspace, String(agentId ?? ""));
1307 if (!agent) return fail("not_found", "No such agent in this workspace.");
1308 const key = principalKey({ kind: "agent", id: agent.id });
1309 const [channel, member] = await Promise.all([
1310 this.db
1311 .prepare("SELECT * FROM channels WHERE id = ? AND workspace_id = ?")
1312 .bind(String(channelId ?? ""), workspace.id)
1313 .first<ChannelRow>(),
1314 this.db.prepare("SELECT * FROM channel_members WHERE channel_id = ? AND principal = ?").bind(String(channelId ?? ""), key).first<MemberRow>(),
1315 ]);
1316 if (!channel) return fail("not_found", "No such channel.");
1317 if (!member) return fail("forbidden", "The agent is not a member of this channel.");
1318 return ok({ place: { slug: slug.toLowerCase(), workspace, channel, member }, agent });
1319 }
1320
1321 async postAsAgent(a: { workspace: string; channel_id: string; agent_id: string; message: AgentPostMessage }): Promise<Result<ChatMessage>> {
1322 const found = await this.agentPlace(a.workspace, a.channel_id, a.agent_id);
1323 if (!found.ok) return found;
1324 const { place, agent } = found.value;
1325 const card = a.message?.card == null ? null : cleanCard(a.message.card);
1326 if (a.message?.card != null && !card) return fail("invalid", "A card needs a kind and a title.");
1327 const body = messageBody(a.message?.body ?? "", !!card);
1328 if (!body.ok) return fail("invalid", body.message);
1329 const hops = typeof a.message?.hops === "number" && a.message.hops >= 0 ? Math.floor(a.message.hops) : 0;
1330 const askedBy = typeof a.message?.asked_by === "string" && a.message.asked_by ? a.message.asked_by : agent.created_by;
1331 return this.write(
1332 place,
1333 principalKey({ kind: "agent", id: agent.id }),
1334 { body: body.body, card, thread_root: a.message?.thread_root ?? null },
1335 // The asker carries on from the delivery the agent is answering;
1336 // without one, agents it wakes treat the asker as unable to change code.
1337 { hops, asked_by: askedBy, asker: cleanAsker(a.message?.asker), chain: chainFor(a.message?.chain, agent.id) },
1338 );
1339 }
1340
1341 /**
1342 * Changes a message the agent posted (a session's live card, say): its
1343 * body, its card, or both. Only the agent's own messages; it wakes
1344 * nobody, and everyone in the conversation sees it change.
1345 */
1346 async updateAsAgent(a: {
1347 workspace: string;
1348 channel_id: string;
1349 agent_id: string;
1350 id: string;
1351 change: { body?: string; card?: MessageCard | null };
1352 }): Promise<Result<ChatMessage>> {
1353 const found = await this.agentPlace(a.workspace, a.channel_id, a.agent_id);
1354 if (!found.ok) return found;
1355 const { place, agent } = found.value;
1356 const row = await this.messageRow(place.channel.id, String(a.id ?? ""));
1357 if (!row || row.deleted_at) return fail("not_found", "No such message.");
1358 if (row.author !== principalKey({ kind: "agent", id: agent.id })) return fail("forbidden", "An agent can change only its own messages.");
1359 const change = a.change ?? {};
1360 const card = change.card === undefined ? (row.card ? (JSON.parse(row.card) as MessageCard) : null) : change.card === null ? null : cleanCard(change.card);
1361 if (change.card && !card) return fail("invalid", "A card needs a kind and a title.");
1362 const body = messageBody(change.body === undefined ? row.body : change.body, !!card);
1363 if (!body.ok) return fail("invalid", body.message);
1364 const cardJson = card ? JSON.stringify(card) : null;
1365 const at = now();
1366 // A card's live state is not an edit a person made: no "edited" mark for it.
1367 const edited = change.body !== undefined && change.body !== row.body ? at : row.edited_at;
1368 await this.db
1369 .prepare("UPDATE messages SET body = ?, card = ?, kind = ?, edited_at = ? WHERE id = ?")
1370 .bind(body.body, cardJson, card ? "card" : "text", edited, row.id)
1371 .run();
1372 const [message] = await this.toMessages(place.slug, place.workspace, [{ ...row, body: body.body, card: cardJson, kind: card ? "card" : "text", edited_at: edited }]);
1373 this.broadcast(place.channel.id, { type: "message.updated", message });
1374 return ok(message);
1375 }
1376
1377 /**
1378 * What an agent reads before replying, oldest first: a thread (its root,
1379 * then its latest replies), or the channel's latest top-level messages.
1380 * Only where the agent is a member, so it reads only what was said where
1381 * it was invited. Deleted messages are left out, save a thread's root.
1382 */
1383 async historyForAgent(a: {
1384 workspace: string;
1385 channel_id: string;
1386 agent_id: string;
1387 thread_root?: string | null;
1388 limit?: number | null;
1389 }): Promise<Result<ChatMessage[]>> {
1390 const found = await this.agentPlace(a.workspace, a.channel_id, a.agent_id);
1391 if (!found.ok) return found;
1392 const { place } = found.value;
1393 const size = historySize(a.limit);
1394 if (typeof a.thread_root === "string" && a.thread_root) {
1395 const asked = await this.messageRow(place.channel.id, a.thread_root);
1396 if (!asked) return fail("not_found", "No such thread.");
1397 // Asked from a reply: its whole thread.
1398 const root = asked.thread_root ? await this.messageRow(place.channel.id, asked.thread_root) : asked;
1399 if (!root) return fail("not_found", "No such thread.");
1400 const replies = await this.db
1401 .prepare(
1402 "SELECT * FROM messages WHERE thread_root = ? AND channel_id = ? AND deleted_at IS NULL ORDER BY id DESC LIMIT ?",
1403 )
1404 .bind(root.id, place.channel.id, Math.max(0, size - 1))
1405 .all<MessageRow>();
1406 return ok(await this.toMessages(place.slug, place.workspace, historyOf(replies.results, root)));
1407 }
1408 const rows = await this.db
1409 .prepare(
1410 "SELECT * FROM messages WHERE channel_id = ? AND thread_root IS NULL AND id < '~' AND deleted_at IS NULL ORDER BY id DESC LIMIT ?",
1411 )
1412 .bind(place.channel.id, size)
1413 .all<MessageRow>();
1414 return ok(await this.toMessages(place.slug, place.workspace, historyOf(rows.results)));
1415 }
1416
1417 // ── What an agent may read (src/audience.ts) ───────────────────────────
1418
1419 /** The people in a conversation, by user id. */
1420 private async peopleIn(channelId: string): Promise<string[]> {
1421 const rows = await this.db
1422 .prepare("SELECT principal FROM channel_members WHERE channel_id = ? AND principal LIKE 'user:%'")
1423 .bind(channelId)
1424 .all<{ principal: string }>();
1425 return rows.results.map((row) => row.principal.slice("user:".length));
1426 }
1427
1428 /** A conversation of this workspace and who reads it, worked out here, never taken from a caller. */
1429 private async audienceOf(slug: string, channelId: string): Promise<Result<{ workspace: Workspace; channel: ChannelRow; audience: ChatAudience }>> {
1430 const workspace = await this.workspace(String(slug ?? "").toLowerCase());
1431 if (!workspace) return fail("not_found", "No such workspace.");
1432 const channel = await this.db
1433 .prepare("SELECT * FROM channels WHERE id = ? AND workspace_id = ?")
1434 .bind(String(channelId ?? ""), workspace.id)
1435 .first<ChannelRow>();
1436 if (!channel) return fail("not_found", "No such conversation.");
1437 const people = await this.peopleIn(channel.id);
1438 return ok({ workspace, channel, audience: { kind: audienceKind(channel), member_user_ids: people, member_count: people.length } });
1439 }
1440
1441 async audience(a: { workspace: string; channel_id: string }): Promise<Result<ChatAudience>> {
1442 const found = await this.audienceOf(a.workspace, a.channel_id);
1443 return found.ok ? ok(found.value.audience) : found;
1444 }
1445
1446 /**
1447 * The conversations of the workspace the audience of `channelId` may all
1448 * read: public channels, and the ones every person in it is in (a DM only
1449 * with exactly them). At most 500, most recently active first.
1450 */
1451 private async readableFor(workspace: Workspace, audience: ChatAudience): Promise<Map<string, ChannelRow>> {
1452 const people = audience.member_user_ids;
1453 const rows = isShared({ kind: audience.kind, user_ids: people }) || !people.length
1454 ? await this.db
1455 .prepare("SELECT * FROM channels WHERE workspace_id = ? AND kind = 'channel' AND private = 0 ORDER BY last_message_at DESC LIMIT 500")
1456 .bind(workspace.id)
1457 .all<ChannelRow>()
1458 : await this.db
1459 .prepare(
1460 `SELECT * FROM channels WHERE workspace_id = ?1 AND (
1461 (kind = 'channel' AND private = 0)
1462 OR id IN (SELECT channel_id FROM channel_members WHERE principal IN (${people.map((_, i) => `?${i + 2}`).join(", ")})
1463 GROUP BY channel_id HAVING COUNT(DISTINCT principal) = ${people.length})
1464 ) ORDER BY last_message_at DESC LIMIT 500`,
1465 )
1466 .bind(workspace.id, ...people.map((id) => `user:${id}`))
1467 .all<ChannelRow>();
1468 const out = new Map<string, ChannelRow>();
1469 for (const channel of rows.results) {
1470 // The query finds candidates; the rule decides, a DM's people included.
1471 const target = { kind: channel.kind, private: channel.private, user_ids: channel.kind === "channel" && !channel.private ? [] : await this.peopleIn(channel.id) };
1472 if (readableBy(target, { kind: audience.kind, user_ids: people })) out.set(channel.id, channel);
1473 }
1474 return out;
1475 }
1476
1477 private async found(slug: string, workspace: Workspace, channels: Map<string, ChannelRow>, rows: MessageRow[]): Promise<AgentFoundMessage[]> {
1478 const messages = await this.toMessages(slug, workspace, rows);
1479 return messages.map((message) => {
1480 const channel = channels.get(message.channel_id)!;
1481 return { channel_id: channel.id, channel: channel.kind === "dm" ? null : channel.name, message };
1482 });
1483 }
1484
1485 async searchForAgent(a: { workspace: string; channel_id: string; query: string; limit?: number | null }): Promise<Result<AgentFoundMessage[]>> {
1486 const found = await this.audienceOf(a.workspace, a.channel_id);
1487 if (!found.ok) return found;
1488 const pattern = likePattern(a.query);
1489 if (!pattern) return fail("invalid", "Search for at least two characters.");
1490 const { workspace, audience } = found.value;
1491 const channels = await this.readableFor(workspace, audience);
1492 if (!channels.size) return ok([]);
1493 const ids = [...channels.keys()];
1494 const limit = Math.min(20, Math.max(1, Math.floor(Number(a.limit) || 20)));
1495 const rows = await this.db
1496 .prepare(
1497 `SELECT * FROM messages WHERE channel_id IN (${ids.map(() => "?").join(", ")}) AND deleted_at IS NULL AND body LIKE ? ESCAPE '\\'
1498 ORDER BY id DESC LIMIT ?`,
1499 )
1500 .bind(...ids, pattern, limit)
1501 .all<MessageRow>();
1502 return ok(await this.found(a.workspace.toLowerCase(), workspace, channels, rows.results));
1503 }
1504
1505 async threadForAgent(a: { workspace: string; channel_id: string; target_channel_id: string; id: string }): Promise<Result<AgentFoundMessage[]>> {
1506 const found = await this.audienceOf(a.workspace, a.channel_id);
1507 if (!found.ok) return found;
1508 const { workspace, audience } = found.value;
1509 // The same answer for a conversation that is not there and one the audience may not read.
1510 const hidden = fail("not_found", "Not available in this conversation.");
1511 const target = await this.db
1512 .prepare("SELECT * FROM channels WHERE id = ? AND workspace_id = ?")
1513 .bind(String(a.target_channel_id ?? ""), workspace.id)
1514 .first<ChannelRow>();
1515 if (!target) return hidden;
1516 const people = target.kind === "channel" && !target.private ? [] : await this.peopleIn(target.id);
1517 if (!readableBy({ kind: target.kind, private: target.private, user_ids: people }, { kind: audience.kind, user_ids: audience.member_user_ids })) return hidden;
1518 const asked = await this.messageRow(target.id, String(a.id ?? ""));
1519 if (!asked) return hidden;
1520 const root = asked.thread_root ? await this.messageRow(target.id, asked.thread_root) : asked;
1521 if (!root) return hidden;
1522 const replies = await this.db
1523 .prepare("SELECT * FROM messages WHERE thread_root = ? AND channel_id = ? AND deleted_at IS NULL ORDER BY id DESC LIMIT 49")
1524 .bind(root.id, target.id)
1525 .all<MessageRow>();
1526 const rows = historyOf(replies.results, root).filter((row) => !row.deleted_at);
1527 return ok(await this.found(a.workspace.toLowerCase(), workspace, new Map([[target.id, target]]), rows));
1528 }
1529
1530 async agentTyping(a: { workspace: string; channel_id: string; agent_id: string }): Promise<Result<null>> {
1531 const found = await this.agentPlace(a.workspace, a.channel_id, a.agent_id);
1532 if (!found.ok) return found;
1533 const { place, agent } = found.value;
1534 const key = principalKey({ kind: "agent", id: agent.id });
1535 const member = await this.profile(place.slug, place.workspace, key);
1536 this.broadcast(place.channel.id, {
1537 type: "typing",
1538 channel_id: place.channel.id,
1539 member,
1540 until: new Date(Date.now() + AGENT_TYPING_MS).toISOString(),
1541 });
1542 return ok(null);
1543 }
1544
1545 // ── Reactions (src/emoji.ts) ────────────────────────────────────────────
1546
1547 /** One message's reactions now, as `me` sees them. */
1548 private async reactionsFor(place: Place, messageId: string, me: string): Promise<ChatReaction[]> {
1549 const list = (await this.reactionsOf([messageId], me)).get(messageId) ?? [];
1550 const profiles = await this.profiles(place.slug, place.workspace, list.flatMap((r) => r.by));
1551 return list.map((r) => ({ ...r, by: r.by.map((key) => profiles.get(key)!).filter(Boolean) }));
1552 }
1553
1554 /**
1555 * Adds or takes back `who`'s reaction. Each member reacts once with each
1556 * emoji, a message holds at most 50 different ones, and a workspace's own
1557 * emoji must exist to be used (taking one back never needs it to). The
1558 * room hears of each change.
1559 */
1560 private async reactTo(place: Place, who: string, messageId: unknown, input: unknown, remove: boolean): Promise<Result<ChatReaction[]>> {
1561 const parsed = reactionEmoji(input);
1562 if (!parsed.ok) return fail("invalid", parsed.message);
1563 const row = await this.messageRow(place.channel.id, String(messageId ?? ""));
1564 if (!row || row.deleted_at) return fail("not_found", "No such message.");
1565 let changed = 0;
1566 if (remove) {
1567 const done = await this.db
1568 .prepare("DELETE FROM reactions WHERE message_id = ? AND principal = ? AND emoji = ?")
1569 .bind(row.id, who, parsed.emoji)
1570 .run();
1571 changed = done.meta.changes;
1572 } else {
1573 if (parsed.custom) {
1574 const known = await this.db
1575 .prepare("SELECT 1 FROM custom_emoji WHERE workspace_id = ? AND name = ? AND deleted_at IS NULL")
1576 .bind(place.workspace.id, parsed.custom)
1577 .first();
1578 if (!known) return fail("not_found", `This workspace has no :${parsed.custom}: emoji.`);
1579 }
1580 const kinds = await this.db.prepare("SELECT DISTINCT emoji FROM reactions WHERE message_id = ?").bind(row.id).all<{ emoji: string }>();
1581 if (!roomForReaction(new Set(kinds.results.map((k) => k.emoji)), parsed.emoji)) {
1582 return fail("invalid", `A message can have at most ${MAX_REACTIONS_PER_MESSAGE} different reactions.`);
1583 }
1584 const done = await this.db
1585 .prepare("INSERT OR IGNORE INTO reactions (message_id, principal, emoji, created_at) VALUES (?, ?, ?, ?)")
1586 .bind(row.id, who, parsed.emoji, now())
1587 .run();
1588 changed = done.meta.changes;
1589 }
1590 if (changed) {
1591 const member = await this.profile(place.slug, place.workspace, who);
1592 this.broadcast(place.channel.id, {
1593 type: remove ? "reaction.removed" : "reaction.added",
1594 channel_id: place.channel.id,
1595 message_id: row.id,
1596 emoji: parsed.emoji,
1597 member,
1598 });
1599 }
1600 return ok(await this.reactionsFor(place, row.id, who));
1601 }
1602
1603 async react(a: { workspace: string; channel_id: string; viewer: Viewer; message_id: string; emoji: string }): Promise<Result<ChatReaction[]>> {
1604 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
1605 if (!found.ok) return found;
1606 return this.reactTo(found.value, userKey(a.viewer!), a.message_id, a.emoji, false);
1607 }
1608
1609 async unreact(a: { workspace: string; channel_id: string; viewer: Viewer; message_id: string; emoji: string }): Promise<Result<ChatReaction[]>> {
1610 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
1611 if (!found.ok) return found;
1612 return this.reactTo(found.value, userKey(a.viewer!), a.message_id, a.emoji, true);
1613 }
1614
1615 /** An agent's reaction counts like anyone's; it must be in the channel. */
1616 async reactAsAgent(a: {
1617 workspace: string;
1618 channel_id: string;
1619 agent_id: string;
1620 message_id: string;
1621 emoji: string;
1622 remove?: boolean;
1623 }): Promise<Result<ChatReaction[]>> {
1624 const found = await this.agentPlace(a.workspace, a.channel_id, a.agent_id);
1625 if (!found.ok) return found;
1626 const { place, agent } = found.value;
1627 return this.reactTo(place, principalKey({ kind: "agent", id: agent.id }), a.message_id, a.emoji, a.remove === true);
1628 }
1629
1630 // ── A workspace's own emoji (src/emoji.ts) ──────────────────────────────
1631
1632 private async emojiUpload(workspace: Workspace): Promise<EmojiUpload> {
1633 return settingsOf(await this.settingsRow(workspace)).emoji_upload;
1634 }
1635
1636 private async toEmoji(slug: string, workspace: Workspace, rows: EmojiRow[]): Promise<CustomEmoji[]> {
1637 const profiles = await this.profiles(slug, workspace, rows.map((r) => r.created_by));
1638 return rows.map((row) => ({
1639 name: row.name,
1640 alias_of: row.alias_of,
1641 file: row.file,
1642 content_type: row.content_type,
1643 bytes: row.bytes,
1644 created_by: profiles.get(row.created_by)!,
1645 created_at: row.created_at,
1646 }));
1647 }
1648
1649 private liveEmoji(workspace: Workspace, name: string): Promise<EmojiRow | null> {
1650 return this.db
1651 .prepare("SELECT * FROM custom_emoji WHERE workspace_id = ? AND name = ? AND deleted_at IS NULL")
1652 .bind(workspace.id, name)
1653 .first<EmojiRow>();
1654 }
1655
1656 async listEmoji(a: { workspace: string; viewer: Viewer }): Promise<Result<EmojiList>> {
1657 const found = await this.viewerWorkspace(a.workspace, a.viewer);
1658 if (!found.ok) return found;
1659 const workspace = found.value;
1660 const [rows, setting] = await Promise.all([
1661 this.db
1662 .prepare("SELECT * FROM custom_emoji WHERE workspace_id = ? AND deleted_at IS NULL ORDER BY name")
1663 .bind(workspace.id)
1664 .all<EmojiRow>(),
1665 this.emojiUpload(workspace),
1666 ]);
1667 const role = roleOf(a.viewer, a.workspace);
1668 return ok({
1669 emoji: await this.toEmoji(a.workspace.toLowerCase(), workspace, rows.results),
1670 emoji_upload: setting,
1671 can_upload: mayUpload(setting, role),
1672 can_manage: role === "owner",
1673 });
1674 }
1675
1676 /** Who may add one: checked against the workspace's setting. */
1677 private async mayAdd(workspace: Workspace, viewer: Viewer, slug: string): Promise<Result<null>> {
1678 const setting = await this.emojiUpload(workspace);
1679 if (!mayUpload(setting, roleOf(viewer, slug))) return fail("forbidden", "Only owners can add emoji in this workspace.");
1680 return ok(null);
1681 }
1682
1683 async addEmoji(a: { workspace: string; viewer: Viewer; name: string; file: EmojiFile }): Promise<Result<CustomEmoji>> {
1684 const found = await this.viewerWorkspace(a.workspace, a.viewer);
1685 if (!found.ok) return found;
1686 const workspace = found.value;
1687 const allowed = await this.mayAdd(workspace, a.viewer, a.workspace);
1688 if (!allowed.ok) return allowed;
1689 const named = emojiName(a.name);
1690 if (!named.ok) return fail("invalid", named.message);
1691 if (await this.liveEmoji(workspace, named.name)) return fail("conflict", `:${named.name}: is already taken.`);
1692 const bytes = fromBase64(a.file?.data);
1693 if (!bytes) return fail("invalid", "Choose a PNG, GIF or WebP image of at most 256 KB.");
1694 const checked = emojiImage(bytes);
1695 if (!checked.ok) return fail("invalid", checked.message);
1696 const file = await sha256(bytes);
1697 // Kept by its hash, so the same image stored twice is one file; its
1698 // type is the one read from its bytes (usercontent serves only that).
1699 await this.env.AVATARS.put(`emoji/${file}`, bytes, { metadata: { contentType: checked.image.content_type } });
1700 const row: EmojiRow = {
1701 workspace_id: workspace.id,
1702 name: named.name,
1703 alias_of: null,
1704 file,
1705 content_type: checked.image.content_type,
1706 bytes: bytes.length,
1707 created_by: userKey(a.viewer!),
1708 created_at: now(),
1709 deleted_at: null,
1710 };
1711 const added = await this.insertEmoji(row);
1712 if (!added.ok) return added;
1713 const [emoji] = await this.toEmoji(a.workspace.toLowerCase(), workspace, [row]);
1714 return ok(emoji!);
1715 }
1716
1717 private async insertEmoji(row: EmojiRow): Promise<Result<null>> {
1718 try {
1719 await this.db
1720 .prepare(
1721 "INSERT INTO custom_emoji (workspace_id, name, alias_of, file, content_type, bytes, created_by, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
1722 )
1723 .bind(row.workspace_id, row.name, row.alias_of, row.file, row.content_type, row.bytes, row.created_by, row.created_at)
1724 .run();
1725 return ok(null);
1726 } catch (error) {
1727 if (String(error).includes("UNIQUE")) return fail("conflict", `:${row.name}: is already taken.`);
1728 throw error;
1729 }
1730 }
1731
1732 async aliasEmoji(a: { workspace: string; viewer: Viewer; name: string; target: string }): Promise<Result<CustomEmoji>> {
1733 const found = await this.viewerWorkspace(a.workspace, a.viewer);
1734 if (!found.ok) return found;
1735 const workspace = found.value;
1736 const allowed = await this.mayAdd(workspace, a.viewer, a.workspace);
1737 if (!allowed.ok) return allowed;
1738 const named = emojiName(a.name);
1739 if (!named.ok) return fail("invalid", named.message);
1740 const targetName = String(a.target ?? "").trim().replace(/^:+|:+$/g, "").toLowerCase();
1741 let target = await this.liveEmoji(workspace, targetName);
1742 // An alias of an alias names the emoji itself, so removing one never strands another.
1743 if (target?.alias_of) target = await this.liveEmoji(workspace, target.alias_of);
1744 if (!target) return fail("not_found", `This workspace has no :${targetName}: emoji.`);
1745 if (await this.liveEmoji(workspace, named.name)) return fail("conflict", `:${named.name}: is already taken.`);
1746 const row: EmojiRow = { ...target, name: named.name, alias_of: target.name, created_by: userKey(a.viewer!), created_at: now(), deleted_at: null };
1747 const added = await this.insertEmoji(row);
1748 if (!added.ok) return added;
1749 const [emoji] = await this.toEmoji(a.workspace.toLowerCase(), workspace, [row]);
1750 return ok(emoji!);
1751 }
1752
1753 async removeEmoji(a: { workspace: string; viewer: Viewer; name: string }): Promise<Result<null>> {
1754 const found = await this.viewerWorkspace(a.workspace, a.viewer);
1755 if (!found.ok) return found;
1756 const workspace = found.value;
1757 const name = String(a.name ?? "").trim().replace(/^:+|:+$/g, "").toLowerCase();
1758 const row = await this.liveEmoji(workspace, name);
1759 if (!row) return fail("not_found", `This workspace has no :${name}: emoji.`);
1760 if (!mayRemove(row.created_by, userKey(a.viewer!), roleOf(a.viewer, a.workspace))) {
1761 return fail("forbidden", "Only whoever added an emoji, or an owner, can remove it.");
1762 }
1763 // An emoji goes with its aliases; an alias goes alone.
1764 await this.db
1765 .prepare(
1766 "UPDATE custom_emoji SET deleted_at = ?1 WHERE workspace_id = ?2 AND deleted_at IS NULL AND (name = ?3 OR (?4 = 0 AND alias_of = ?3))",
1767 )
1768 .bind(now(), workspace.id, row.name, row.alias_of ? 1 : 0)
1769 .run();
1770 // Its image goes once nothing live shows it, in any workspace.
1771 this.defer(
1772 (async () => {
1773 const used = await this.db.prepare("SELECT 1 FROM custom_emoji WHERE file = ? AND deleted_at IS NULL LIMIT 1").bind(row.file).first();
1774 if (!used) await this.env.AVATARS.delete(`emoji/${row.file}`);
1775 })().catch((error) => console.error("chat could not forget emoji", row.file, error)),
1776 );
1777 return ok(null);
1778 }
1779
1780 async setEmojiUpload(a: { workspace: string; viewer: Viewer; value: EmojiUpload }): Promise<Result<EmojiUpload>> {
1781 const found = await this.viewerWorkspace(a.workspace, a.viewer);
1782 if (!found.ok) return found;
1783 if (roleOf(a.viewer, a.workspace) !== "owner") return fail("forbidden", "Only owners can change who adds emoji.");
1784 if (a.value !== "members" && a.value !== "admins") return fail("invalid", "Choose members or owners.");
1785 // The same setting as Settings → Chat (`setChatSettings`).
1786 const set = await this.setChatSettings({ workspace: a.workspace, viewer: a.viewer, change: { emoji_upload: a.value } });
1787 return set.ok ? ok(set.value.emoji_upload) : set;
1788 }
1789
1790 // ── The live socket ─────────────────────────────────────────────────────
1791
1792 /**
1793 * `GET /live?workspace=<slug>&channel=<id>`, upgraded to a WebSocket. The
1794 * viewer comes in CHAT_VIEWER_HEADER, set by the site after checking the
1795 * session; trusted only because this Worker is reachable through service
1796 * bindings alone (`workers_dev` is off and it has no routes). Checked
1797 * like any read, then handed to the channel's room.
1798 */
1799 async live(request: Request): Promise<Response> {
1800 if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") {
1801 return new Response("Expected a WebSocket upgrade\n", { status: 426 });
1802 }
1803 let viewer: Viewer = null;
1804 try {
1805 viewer = JSON.parse(request.headers.get(CHAT_VIEWER_HEADER) ?? "null") as Viewer;
1806 } catch {
1807 viewer = null;
1808 }
1809 if (!viewer?.id) return new Response("Sign in to use chat\n", { status: 401 });
1810 const url = new URL(request.url);
1811 const channelId = url.searchParams.get("channel") ?? "";
1812 let slug = (url.searchParams.get("workspace") ?? "").toLowerCase();
1813 if (!slug) {
1814 // Not named: whichever of the viewer's workspaces holds the channel.
1815 const row = await this.db.prepare("SELECT workspace_id FROM channels WHERE id = ?").bind(channelId).first<{ workspace_id: string }>();
1816 if (row) {
1817 const theirs = await Promise.all((viewer.workspaces ?? []).map((m) => this.workspace(m.slug)));
1818 slug = theirs.find((w) => w?.id === row.workspace_id)?.slug.toLowerCase() ?? "";
1819 }
1820 }
1821 const found = await this.place(slug, channelId, viewer, "read");
1822 if (!found.ok) return new Response(`${found.error.message}\n`, { status: found.error.code === "forbidden" ? 403 : 404 });
1823 const { workspace, channel } = found.value;
1824 const who: RoomMember = { channel_id: channel.id, member: await this.profile(slug, workspace, userKey(viewer)) };
1825 const headers = new Headers(request.headers);
1826 headers.delete(CHAT_VIEWER_HEADER);
1827 headers.set(ROOM_MEMBER_HEADER, JSON.stringify(who));
1828 return this.room(channel.id).fetch(new Request(request.url, { method: "GET", headers }));
1829 }
1830}
1831
1832/** One RPC method's answer. */
1833async function answer(service: Chat, method: string, args: any): Promise<Response> {
1834 switch (method) {
1835 case "sidebar":
1836 return Response.json(await service.sidebar(args));
1837 case "channel":
1838 return Response.json(await service.channel(args));
1839 case "channel_by_name":
1840 return Response.json(await service.channelByName(args));
1841 case "browse":
1842 return Response.json(await service.browse(args));
1843 case "create_channel":
1844 return Response.json(await service.createChannel(args));
1845 case "update_channel":
1846 return Response.json(await service.updateChannel(args));
1847 case "chat_settings":
1848 return Response.json(await service.chatSettings(args));
1849 case "set_chat_settings":
1850 return Response.json(await service.setChatSettings(args));
1851 case "open_dm":
1852 return Response.json(await service.openDm(args));
1853 case "join":
1854 return Response.json(await service.join(args));
1855 case "leave":
1856 return Response.json(await service.leave(args));
1857 case "invite":
1858 return Response.json(await service.invite(args));
1859 case "messages":
1860 return Response.json(await service.messages(args));
1861 case "post":
1862 return Response.json(await service.post(args));
1863 case "edit":
1864 return Response.json(await service.edit(args));
1865 case "remove":
1866 return Response.json(await service.remove(args));
1867 case "mark_read":
1868 return Response.json(await service.markRead(args));
1869 case "set_preferences":
1870 return Response.json(await service.setPreferences(args));
1871 case "post_as_agent":
1872 return Response.json(await service.postAsAgent(args));
1873 case "update_as_agent":
1874 return Response.json(await service.updateAsAgent(args));
1875 case "agent_typing":
1876 return Response.json(await service.agentTyping(args));
1877 case "history_for_agent":
1878 return Response.json(await service.historyForAgent(args));
1879 case "audience":
1880 return Response.json(await service.audience(args));
1881 case "search_for_agent":
1882 return Response.json(await service.searchForAgent(args));
1883 case "thread_for_agent":
1884 return Response.json(await service.threadForAgent(args));
1885 case "react":
1886 return Response.json(await service.react(args));
1887 case "unreact":
1888 return Response.json(await service.unreact(args));
1889 case "react_as_agent":
1890 return Response.json(await service.reactAsAgent(args));
1891 case "list_emoji":
1892 return Response.json(await service.listEmoji(args));
1893 case "add_emoji":
1894 return Response.json(await service.addEmoji(args));
1895 case "alias_emoji":
1896 return Response.json(await service.aliasEmoji(args));
1897 case "remove_emoji":
1898 return Response.json(await service.removeEmoji(args));
1899 case "set_emoji_upload":
1900 return Response.json(await service.setEmojiUpload(args));
1901 default:
1902 return new Response("Unknown method\n", { status: 404 });
1903 }
1904}
1905
1906export default {
1907 async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
1908 const url = new URL(request.url);
1909 if (request.method === "GET" && url.pathname === "/live") {
1910 return new Chat(env, (work) => ctx.waitUntil(work)).live(request);
1911 }
1912 const match = url.pathname.match(/^\/rpc\/([a-z_]+)$/);
1913 if (request.method !== "POST" || !match) return new Response("Not found\n", { status: 404 });
1914 // A replica near the caller when it asks for one (@g1t/contracts d1.ts).
1915 const opened = openD1(env.DB, request);
1916 const service = new Chat(Object.create(env, { DB: { value: opened.db } }) as Env, (work) => ctx.waitUntil(work));
1917 const args = (await request.json().catch(() => ({}))) as any;
1918 return opened.finish(await answer(service, match[1], args));
1919 },
1920} satisfies ExportedHandler<Env>;