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