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