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