Skip to content
1,211 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 AgentPostMessage,
28 type AskerAccess,
29 type Channel,
30 type ChannelMember,
31 type ChatLiveEvent,
32 type ChatMessage,
33 type ChatSidebar,
34 type ChatSidebarEntry,
35
36 type Member,
37 type MemberProfile,
38 type MessageCard,
39 type MessagePage,
40 type NewChannel,
41 type PostMessage,
42 type Principal,
43 type Result,
44 type ServiceBinding,
45 type User,
46 type Viewer,
47 type Workspace,
48 type WorkspaceAgent,
49} from "@g1t/contracts";
50
51import { MAX_HOPS, deliveries, delivery, type Chain } from "./delivery.ts";
52import { mentionedHandles, mentionsColumn } from "./mentions.ts";
53import { AGENT_TYPING_MS, historyOf, historySize, messageBody, meterDay, pageOf, pageSize } from "./messages.ts";
54import { GENERAL, MAX_DM_MEMBERS, channelName, dmKey, dmMembers } from "./names.ts";
55import { ROOM_MEMBER_HEADER, type ChannelRoom, type RoomMember } from "./room.ts";
56import { dmTitle, sidebarOrder, tally, type UnreadRow } from "./unread.ts";
57
58export { ChannelRoom } from "./room.ts";
59
60// The hop limit here is the one in the contract.
61const SAME_HOP_LIMIT: typeof CHAT_MAX_HOPS = MAX_HOPS;
62void SAME_HOP_LIMIT;
63
64type Env = {
65 DB: D1Database;
66 IDENTITY: ServiceBinding;
67 AGENTS: ServiceBinding;
68 ROOMS: DurableObjectNamespace<ChannelRoom>;
69};
70
71type ChannelRow = {
72 id: string;
73 workspace_id: string;
74 kind: "channel" | "dm";
75 name: string | null;
76 topic: string | null;
77 private: number;
78 dm_key: string | null;
79 created_by: string;
80 created_at: string;
81 archived_at: string | null;
82 last_message_at: string | null;
83};
84
85type MemberRow = {
86 channel_id: string;
87 principal: string;
88 role: "owner" | "member";
89 starred: number;
90 muted: number;
91 last_read_id: string | null;
92 joined_at: string;
93};
94
95type MessageRow = {
96 id: string;
97 channel_id: string;
98 author: string;
99 kind: "text" | "card";
100 body: string;
101 card: string | null;
102 mentions: string;
103 thread_root: string | null;
104 reply_count: number;
105 last_reply_at: string | null;
106 created_at: string;
107 edited_at: string | null;
108 deleted_at: string | null;
109};
110
111/** What a method found out about the channel it was asked about. */
112type Place = { slug: string; workspace: Workspace; channel: ChannelRow; member: MemberRow | null };
113
114/** The most unread messages one sidebar reads to count; past it, counts are "at least". */
115const MAX_UNREAD_ROWS = 5_000;
116/** The longest channel topic. */
117const MAX_TOPIC = 250;
118/** How many others a direct message's sidebar entry shows. */
119const DM_FACES = 4;
120
121const now = () => new Date().toISOString();
122
123function isMember(viewer: Viewer, workspace: string): boolean {
124 return !!viewer?.workspaces?.some((m) => m.slug === workspace.toLowerCase());
125}
126
127function isOwner(viewer: Viewer, workspace: string): boolean {
128 return !!viewer?.workspaces?.some((m) => m.slug === workspace.toLowerCase() && m.role === "owner");
129}
130
131function userKey(viewer: User): string {
132 return principalKey({ kind: "user", id: viewer.id });
133}
134
135function toChannel(row: ChannelRow): Channel {
136 return {
137 id: row.id,
138 workspace_id: row.workspace_id,
139 kind: row.kind,
140 name: row.kind === "dm" ? null : row.name,
141 topic: row.topic,
142 private: row.kind === "dm" || !!row.private,
143 created_by: parsePrincipalKey(row.created_by) ?? { kind: "user", id: row.created_by },
144 created_at: row.created_at,
145 archived_at: row.archived_at,
146 last_message_at: row.last_message_at,
147 };
148}
149
150/** A card as kept, or null when what was sent is not one. */
151function cleanCard(card: unknown): MessageCard | null {
152 if (!card || typeof card !== "object") return null;
153 const c = card as Record<string, unknown>;
154 const text = (value: unknown, max: number) => (typeof value === "string" && value.trim() ? value.trim().slice(0, max) : null);
155 const kind = text(c.kind, 40);
156 const title = text(c.title, 300);
157 if (!kind || !title) return null;
158 const href = text(c.href, 2000);
159 return {
160 kind,
161 title,
162 detail: text(c.detail, 500),
163 state: text(c.state, 80),
164 // Relative to the site only: a card never links somewhere else.
165 href: href && href.startsWith("/") && !href.startsWith("//") ? href : null,
166 };
167}
168
169/** An asker handed back by the agents service, or null when what was sent is not one. */
170function cleanAsker(asker: unknown): AskerAccess | null {
171 if (!asker || typeof asker !== "object") return null;
172 const a = asker as Record<string, unknown>;
173 if (typeof a.username !== "string" || !["owner", "member", "outside"].includes(String(a.role))) return null;
174 return { username: a.username, role: a.role as AskerAccess["role"], can_write: a.can_write === true };
175}
176
177function bytesOf(text: string): number {
178 return new TextEncoder().encode(text).length;
179}
180
181class Chat {
182 private readonly workspaces = new Map<string, Promise<Workspace | null>>();
183 private readonly people = new Map<string, Promise<Map<string, Member>>>();
184 private readonly usernames = new Map<string, string>();
185 private readonly agents = new Map<string, WorkspaceAgent | null>();
186
187 /** `defer` runs work after the answer is sent: the request's waitUntil. */
188 constructor(
189 private readonly env: Env,
190 private readonly defer: (work: Promise<unknown>) => void = () => {},
191 ) {}
192
193 private get db() {
194 return this.env.DB;
195 }
196
197 // ── Who and where ───────────────────────────────────────────────────────
198
199 private workspace(slug: string): Promise<Workspace | null> {
200 const key = slug.toLowerCase();
201 let found = this.workspaces.get(key);
202 if (!found) {
203 found = identityClient(this.env.IDENTITY).getWorkspace(key);
204 this.workspaces.set(key, found);
205 }
206 return found;
207 }
208
209 /** The workspace's people by username, with their names and avatars; asked once per request. */
210 private members(slug: string, workspace: Workspace): Promise<Map<string, Member>> {
211 let found = this.people.get(workspace.id);
212 if (!found) {
213 // Asked as the workspace itself, so it works for agents' calls too.
214 const actor: User = {
215 id: workspace.id,
216 username: workspace.slug,
217 kind: "workspace",
218 verified: true,
219 workspaces: [{ slug: workspace.slug, role: "member" }],
220 };
221 found = identityClient(this.env.IDENTITY)
222 .listMembers(slug, actor)
223 .then((result) => new Map(result.ok ? result.value.map((m) => [m.username, m]) : []))
224 .catch((error) => {
225 console.error("chat could not list members of", slug, error);
226 return new Map<string, Member>();
227 });
228 this.people.set(workspace.id, found);
229 }
230 return found;
231 }
232
233 /** Agents by id; ones the agents service does not know are null. Asked once per request. */
234 private async agentsById(ids: string[]): Promise<Map<string, WorkspaceAgent | null>> {
235 const wanted = [...new Set(ids)].filter((id) => !this.agents.has(id));
236 if (wanted.length) {
237 let found: WorkspaceAgent[] = [];
238 try {
239 found = await workspaceAgentsClient(this.env.AGENTS).byIds(wanted);
240 } catch (error) {
241 console.error("chat could not resolve agents", error);
242 }
243 for (const id of wanted) this.agents.set(id, found.find((a) => a.id === id) ?? null);
244 }
245 return new Map(ids.map((id) => [id, this.agents.get(id) ?? null]));
246 }
247
248 /** An agent of this workspace that is not archived, or null. */
249 private async liveAgent(workspace: Workspace, id: string): Promise<WorkspaceAgent | null> {
250 const agent = (await this.agentsById([id])).get(id) ?? null;
251 return agent && agent.workspace_id === workspace.id && !agent.archived_at ? agent : null;
252 }
253
254 /** How each member key shows, for one workspace. */
255 private async profiles(slug: string, workspace: Workspace, keys: string[]): Promise<Map<string, MemberProfile>> {
256 const principals = [...new Set(keys)].map((key) => parsePrincipalKey(key)).filter((p): p is Principal => !!p);
257 const userIds = principals.filter((p) => p.kind === "user").map((p) => p.id);
258 const agentIds = principals.filter((p) => p.kind === "agent").map((p) => p.id);
259 const unnamed = userIds.filter((id) => !this.usernames.has(id));
260 const [named, people, agents] = await Promise.all([
261 unnamed.length ? identityClient(this.env.IDENTITY).usernames(unnamed).catch(() => ({}) as Record<string, string>) : ({} as Record<string, string>),
262 userIds.length ? this.members(slug, workspace) : new Map<string, Member>(),
263 this.agentsById(agentIds),
264 ]);
265 for (const [id, username] of Object.entries(named)) this.usernames.set(id, username);
266 const out = new Map<string, MemberProfile>();
267 for (const p of principals) {
268 if (p.kind === "user") {
269 const username = this.usernames.get(p.id) ?? null;
270 const person = username ? people.get(username) : undefined;
271 out.set(principalKey(p), {
272 ...p,
273 name: username ?? "ghost",
274 display_name: person?.name || username || "Former member",
275 avatar: person?.avatar ?? null,
276 role: null,
277 });
278 } else {
279 const agent = agents.get(p.id) ?? null;
280 out.set(principalKey(p), {
281 ...p,
282 name: agent?.handle ?? p.id,
283 display_name: agent?.display_name ?? "Former agent",
284 avatar: agent?.avatar ?? null,
285 role: agent?.role ?? null,
286 });
287 }
288 }
289 return out;
290 }
291
292 private async profile(slug: string, workspace: Workspace, key: string): Promise<MemberProfile> {
293 return (await this.profiles(slug, workspace, [key])).get(key)!;
294 }
295
296 /** Whether `principal` may be added to a conversation in this workspace. */
297 private async belongs(slug: string, workspace: Workspace, principal: Principal): Promise<boolean> {
298 if (principal.kind === "agent") return !!(await this.liveAgent(workspace, principal.id));
299 if (!this.usernames.has(principal.id)) {
300 const named = await identityClient(this.env.IDENTITY).usernames([principal.id]);
301 for (const [id, username] of Object.entries(named)) this.usernames.set(id, username);
302 }
303 const username = this.usernames.get(principal.id);
304 return !!username && (await this.members(slug, workspace)).has(username);
305 }
306
307 /**
308 * The viewer's workspace, checked: they must belong to it, as in every
309 * other service.
310 */
311 private async viewerWorkspace(slug: string, viewer: Viewer): Promise<Result<Workspace>> {
312 if (!viewer) return fail("unauthenticated", "Sign in to use chat.");
313 if (!slug || !isMember(viewer, slug)) return fail("forbidden", "Only members of a workspace can use its chat.");
314 const workspace = await this.workspace(slug);
315 return workspace ? ok(workspace) : fail("not_found", "No such workspace.");
316 }
317
318 /**
319 * A channel the viewer may read (`read`: any public one in their
320 * workspace, or one they are in) or write in (`member`: one they are in).
321 * A private channel or direct message they are not in is not found, so
322 * its existence does not leak.
323 */
324 private async place(
325 slug: string,
326 channelId: string,
327 viewer: Viewer,
328 need: "read" | "member",
329 ): Promise<Result<Place>> {
330 const found = await this.viewerWorkspace(slug, viewer);
331 if (!found.ok) return found;
332 const workspace = found.value;
333 const [channel, member] = await Promise.all([
334 this.db
335 .prepare("SELECT * FROM channels WHERE id = ? AND workspace_id = ?")
336 .bind(String(channelId ?? ""), workspace.id)
337 .first<ChannelRow>(),
338 this.db
339 .prepare("SELECT * FROM channel_members WHERE channel_id = ? AND principal = ?")
340 .bind(String(channelId ?? ""), userKey(viewer!))
341 .first<MemberRow>(),
342 ]);
343 if (!channel) return fail("not_found", "No such channel.");
344 const open = channel.kind === "channel" && !channel.private;
345 if (!member && !open) return fail("not_found", "No such channel.");
346 if (!member && need === "member") return fail("forbidden", `Join #${channel.name} first.`);
347 return ok({ slug: slug.toLowerCase(), workspace, channel, member });
348 }
349
350 private room(channelId: string) {
351 return this.env.ROOMS.get(this.env.ROOMS.idFromName(channelId));
352 }
353
354 /** Tells everyone looking at a channel, after the answer is sent. */
355 private broadcast(channelId: string, event: ChatLiveEvent, except: string | null = null): void {
356 this.defer(
357 this.room(channelId)
358 .broadcast(event, except)
359 .catch((error: unknown) => console.error("chat could not broadcast to", channelId, error)),
360 );
361 }
362
363 private async toMessages(slug: string, workspace: Workspace, rows: MessageRow[]): Promise<ChatMessage[]> {
364 const profiles = await this.profiles(slug, workspace, rows.map((r) => r.author));
365 return rows.map((row) => {
366 const gone = !!row.deleted_at;
367 return {
368 id: row.id,
369 channel_id: row.channel_id,
370 author: profiles.get(row.author)!,
371 kind: row.kind,
372 body: gone ? "" : row.body,
373 card: gone || !row.card ? null : (JSON.parse(row.card) as MessageCard),
374 thread_root: row.thread_root,
375 reply_count: row.reply_count,
376 last_reply_at: row.last_reply_at,
377 created_at: row.created_at,
378 edited_at: row.edited_at,
379 deleted_at: row.deleted_at,
380 };
381 });
382 }
383
384 private async messageRow(channelId: string, id: string): Promise<MessageRow | null> {
385 return this.db.prepare("SELECT * FROM messages WHERE id = ? AND channel_id = ?").bind(String(id ?? ""), channelId).first<MessageRow>();
386 }
387
388 /** Sends a message as it now is to everyone looking at its channel. */
389 private rebroadcast(place: Place, id: string): void {
390 this.defer(
391 (async () => {
392 const row = await this.messageRow(place.channel.id, id);
393 if (!row) return;
394 const [message] = await this.toMessages(place.slug, place.workspace, [row]);
395 await this.room(place.channel.id).broadcast({ type: "message.updated", message });
396 })().catch((error) => console.error("chat could not rebroadcast", id, error)),
397 );
398 }
399
400 // ── The sidebar ─────────────────────────────────────────────────────────
401
402 /**
403 * Puts a person in the workspace's #general, once. A new workspace has
404 * no channels, so the first sidebar anyone in it asks for creates
405 * #general; and everyone who asks for the sidebar is put in it the first
406 * time, so a new workspace has somewhere to talk and a new member lands
407 * where everyone is. Someone who leaves it is not put back
408 * (`general_joined`). A private channel someone named `general` is never
409 * joined this way.
410 */
411 private async ensureGeneral(workspace: Workspace, me: string): Promise<void> {
412 const seen = await this.db
413 .prepare("SELECT 1 FROM general_joined WHERE workspace_id = ? AND principal = ?")
414 .bind(workspace.id, me)
415 .first();
416 if (seen) return;
417 const at = now();
418 await this.db
419 .prepare(
420 "INSERT OR IGNORE INTO channels (id, workspace_id, kind, name, topic, private, created_by, created_at) VALUES (?, ?, 'channel', ?, ?, 0, ?, ?)",
421 )
422 .bind(newId("chn"), workspace.id, GENERAL, "Anything and everything for the whole workspace.", me, at)
423 .run();
424 const general = await this.db
425 .prepare("SELECT * FROM channels WHERE workspace_id = ? AND name = ?")
426 .bind(workspace.id, GENERAL)
427 .first<ChannelRow>();
428 const statements = [
429 this.db
430 .prepare("INSERT OR IGNORE INTO general_joined (workspace_id, principal, joined_at) VALUES (?, ?, ?)")
431 .bind(workspace.id, me, at),
432 ];
433 if (general && !general.private && !general.archived_at) {
434 statements.push(this.joinStatement(general.id, me, general.created_by === me ? "owner" : "member", at));
435 }
436 await this.db.batch(statements);
437 }
438
439 /**
440 * Adds a member. Someone joining starts with everything already said
441 * read, so a long channel does not greet them with its whole history as
442 * unread.
443 */
444 private joinStatement(channelId: string, principal: string, role: "owner" | "member", at: string): D1PreparedStatement {
445 return this.db
446 .prepare(
447 "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)",
448 )
449 .bind(channelId, principal, role, at);
450 }
451
452 async sidebar(a: { workspace: string; viewer: Viewer }): Promise<Result<ChatSidebar>> {
453 const found = await this.viewerWorkspace(a.workspace, a.viewer);
454 if (!found.ok) return found;
455 const workspace = found.value;
456 const slug = a.workspace.toLowerCase();
457 const me = userKey(a.viewer!);
458 await this.ensureGeneral(workspace, me);
459
460 const [joined, unread, dmOthers, browsable] = await Promise.all([
461 this.db
462 .prepare(
463 `SELECT c.*, m.starred, m.muted, m.last_read_id
464 FROM channel_members m JOIN channels c ON c.id = m.channel_id
465 WHERE m.principal = ? AND c.workspace_id = ? AND c.archived_at IS NULL`,
466 )
467 .bind(me, workspace.id)
468 .all<ChannelRow & { starred: number; muted: number; last_read_id: string | null }>(),
469 this.db
470 .prepare(
471 `SELECT msg.channel_id, msg.id, msg.author, msg.mentions
472 FROM channel_members m
473 JOIN channels c ON c.id = m.channel_id
474 JOIN messages msg ON msg.channel_id = m.channel_id AND msg.id > COALESCE(m.last_read_id, '')
475 WHERE m.principal = ?1 AND c.workspace_id = ?2 AND c.archived_at IS NULL
476 AND msg.deleted_at IS NULL AND msg.author != ?1
477 LIMIT ${MAX_UNREAD_ROWS}`,
478 )
479 .bind(me, workspace.id)
480 .all<UnreadRow>(),
481 this.db
482 .prepare(
483 `SELECT o.channel_id, o.principal
484 FROM channel_members m
485 JOIN channels c ON c.id = m.channel_id AND c.kind = 'dm'
486 JOIN channel_members o ON o.channel_id = m.channel_id AND o.principal != m.principal
487 WHERE m.principal = ? AND c.workspace_id = ? AND c.archived_at IS NULL
488 ORDER BY o.joined_at, o.principal`,
489 )
490 .bind(me, workspace.id)
491 .all<{ channel_id: string; principal: string }>(),
492 this.db
493 .prepare(
494 `SELECT COUNT(*) AS n FROM channels c
495 WHERE c.workspace_id = ? AND c.kind = 'channel' AND c.private = 0 AND c.archived_at IS NULL
496 AND NOT EXISTS (SELECT 1 FROM channel_members m WHERE m.channel_id = c.id AND m.principal = ?)`,
497 )
498 .bind(workspace.id, me)
499 .first<{ n: number }>(),
500 ]);
501
502 const others = new Map<string, string[]>();
503 for (const row of dmOthers.results) others.set(row.channel_id, [...(others.get(row.channel_id) ?? []), row.principal]);
504 const profiles = await this.profiles(slug, workspace, [me, ...dmOthers.results.map((r) => r.principal)]);
505 const self = profiles.get(me) ?? null;
506 const counts = tally(
507 unread.results,
508 me,
509 self?.name ?? a.viewer!.username,
510 new Map(joined.results.map((row) => [row.id, row.last_read_id])),
511 );
512
513 const entries: ChatSidebarEntry[] = joined.results.map((row) => {
514 const faces = (others.get(row.id) ?? []).map((key) => profiles.get(key)!).filter(Boolean);
515 const count = counts.get(row.id) ?? { unread: 0, mentions: 0 };
516 return {
517 channel: toChannel(row),
518 title: row.kind === "dm" ? dmTitle(faces, self) : (row.name ?? ""),
519 others: row.kind === "dm" ? faces.slice(0, DM_FACES) : [],
520 starred: !!row.starred,
521 muted: !!row.muted,
522 unread: count.unread,
523 mentions: count.mentions,
524 };
525 });
526 return ok({ entries: sidebarOrder(entries), browsable: browsable?.n ?? 0 });
527 }
528
529 // ── Channels ────────────────────────────────────────────────────────────
530
531 async channel(a: { workspace: string; channel_id: string; viewer: Viewer }): Promise<Result<{ channel: Channel; members: ChannelMember[] }>> {
532 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
533 if (!found.ok) return found;
534 const { slug, workspace, channel } = found.value;
535 const rows = await this.db
536 .prepare("SELECT * FROM channel_members WHERE channel_id = ? ORDER BY joined_at, principal")
537 .bind(channel.id)
538 .all<MemberRow>();
539 const profiles = await this.profiles(slug, workspace, rows.results.map((r) => r.principal));
540 return ok({
541 channel: toChannel(channel),
542 members: rows.results.map((row) => ({
543 channel_id: row.channel_id,
544 member: profiles.get(row.principal)!,
545 role: row.role,
546 starred: !!row.starred,
547 muted: !!row.muted,
548 last_read_id: row.last_read_id,
549 joined_at: row.joined_at,
550 })),
551 });
552 }
553
554 /** A channel by its name, as the site's URLs name them; read like `channel`. */
555 async channelByName(a: { workspace: string; name: string; viewer: Viewer }): Promise<Result<{ channel: Channel; members: ChannelMember[] }>> {
556 const found = await this.viewerWorkspace(a.workspace, a.viewer);
557 if (!found.ok) return found;
558 const named = channelName(a.name ?? "");
559 if (!named.ok) return fail("not_found", "No such channel.");
560 const row = await this.db
561 .prepare("SELECT id FROM channels WHERE workspace_id = ? AND name = ?")
562 .bind(found.value.id, named.name)
563 .first<{ id: string }>();
564 if (!row) return fail("not_found", "No such channel.");
565 // A private channel the viewer is not in stays not found there.
566 return this.channel({ workspace: a.workspace, channel_id: row.id, viewer: a.viewer });
567 }
568
569 async browse(a: { workspace: string; viewer: Viewer }): Promise<Result<Channel[]>> {
570 const found = await this.viewerWorkspace(a.workspace, a.viewer);
571 if (!found.ok) return found;
572 const rows = await this.db
573 .prepare(
574 "SELECT * FROM channels WHERE workspace_id = ? AND kind = 'channel' AND private = 0 AND archived_at IS NULL ORDER BY name",
575 )
576 .bind(found.value.id)
577 .all<ChannelRow>();
578 return ok(rows.results.map(toChannel));
579 }
580
581 async createChannel(a: { workspace: string; viewer: Viewer; input: NewChannel }): Promise<Result<Channel>> {
582 const found = await this.viewerWorkspace(a.workspace, a.viewer);
583 if (!found.ok) return found;
584 const workspace = found.value;
585 const named = channelName(a.input?.name ?? "");
586 if (!named.ok) return fail("invalid", named.message);
587 const topic = typeof a.input?.topic === "string" ? a.input.topic.trim() : "";
588 if (topic.length > MAX_TOPIC) return fail("invalid", `A topic is at most ${MAX_TOPIC} characters.`);
589 const taken = await this.db
590 .prepare("SELECT 1 FROM channels WHERE workspace_id = ? AND name = ?")
591 .bind(workspace.id, named.name)
592 .first();
593 if (taken) return fail("conflict", `#${named.name} already exists.`);
594 const me = userKey(a.viewer!);
595 const row: ChannelRow = {
596 id: newId("chn"),
597 workspace_id: workspace.id,
598 kind: "channel",
599 name: named.name,
600 topic: topic || null,
601 private: a.input?.private ? 1 : 0,
602 dm_key: null,
603 created_by: me,
604 created_at: now(),
605 archived_at: null,
606 last_message_at: null,
607 };
608 try {
609 await this.db.batch([
610 this.db
611 .prepare(
612 "INSERT INTO channels (id, workspace_id, kind, name, topic, private, created_by, created_at) VALUES (?, ?, 'channel', ?, ?, ?, ?, ?)",
613 )
614 .bind(row.id, row.workspace_id, row.name, row.topic, row.private, me, row.created_at),
615 this.joinStatement(row.id, me, "owner", row.created_at),
616 ]);
617 } catch (error) {
618 if (String(error).includes("UNIQUE")) return fail("conflict", `#${named.name} already exists.`);
619 throw error;
620 }
621 return ok(toChannel(row));
622 }
623
624 async openDm(a: { workspace: string; viewer: Viewer; members: Principal[] }): Promise<Result<Channel>> {
625 const found = await this.viewerWorkspace(a.workspace, a.viewer);
626 if (!found.ok) return found;
627 const workspace = found.value;
628 const slug = a.workspace.toLowerCase();
629 const me = userKey(a.viewer!);
630 const asked = Array.isArray(a.members) ? a.members : [];
631 const principals: Principal[] = [];
632 for (const member of asked) {
633 const p = member && parsePrincipalKey(`${member.kind}:${member.id}`);
634 if (!p) return fail("invalid", "Each member is a person or an agent, by id.");
635 principals.push(p);
636 }
637 const members = dmMembers(me, principals.map(principalKey));
638 if (members.length > MAX_DM_MEMBERS) {
639 return fail("invalid", `A direct message has at most ${MAX_DM_MEMBERS} people and agents. Make a private channel instead.`);
640 }
641 const key = dmKey(members);
642 const existing = await this.db
643 .prepare("SELECT * FROM channels WHERE workspace_id = ? AND dm_key = ?")
644 .bind(workspace.id, key)
645 .first<ChannelRow>();
646 if (existing) return ok(toChannel(existing));
647
648 for (const member of members) {
649 if (member === me) continue;
650 const p = parsePrincipalKey(member)!;
651 if (!(await this.belongs(slug, workspace, p))) {
652 return fail("not_found", p.kind === "agent" ? "No such agent in this workspace." : "That person is not in this workspace.");
653 }
654 }
655 const at = now();
656 await this.db
657 .prepare(
658 "INSERT OR IGNORE INTO channels (id, workspace_id, kind, private, dm_key, created_by, created_at) VALUES (?, ?, 'dm', 1, ?, ?, ?)",
659 )
660 .bind(newId("chn"), workspace.id, key, me, at)
661 .run();
662 // Read back by key: if two people opened it at once, both get the one that won.
663 const channel = await this.db
664 .prepare("SELECT * FROM channels WHERE workspace_id = ? AND dm_key = ?")
665 .bind(workspace.id, key)
666 .first<ChannelRow>();
667 if (!channel) return fail("conflict", "The direct message could not be opened. Try again.");
668 await this.db.batch(members.map((member) => this.joinStatement(channel.id, member, "member", at)));
669 return ok(toChannel(channel));
670 }
671
672 async join(a: { workspace: string; channel_id: string; viewer: Viewer }): Promise<Result<null>> {
673 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
674 if (!found.ok) return found;
675 const { channel, member } = found.value;
676 if (member) return ok(null);
677 if (channel.archived_at) return fail("invalid", "This channel is archived.");
678 await this.joinStatement(channel.id, userKey(a.viewer!), "member", now()).run();
679 return ok(null);
680 }
681
682 async leave(a: { workspace: string; channel_id: string; viewer: Viewer }): Promise<Result<null>> {
683 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
684 if (!found.ok) return found;
685 const { channel, member } = found.value;
686 if (!member) return ok(null);
687 if (channel.kind === "dm") return fail("invalid", "A direct message can't be left. Mute it instead.");
688 const me = userKey(a.viewer!);
689 await this.db.prepare("DELETE FROM channel_members WHERE channel_id = ? AND principal = ?").bind(channel.id, me).run();
690 // Out of a private channel, they may no longer read it, live either.
691 if (channel.private) {
692 this.defer(this.room(channel.id).drop(me).catch((error: unknown) => console.error("chat could not drop", me, error)));
693 }
694 return ok(null);
695 }
696
697 async invite(a: { workspace: string; channel_id: string; viewer: Viewer; member: Principal }): Promise<Result<null>> {
698 const found = await this.place(a.workspace, a.channel_id, a.viewer, "member");
699 if (!found.ok) return found;
700 const { slug, workspace, channel } = found.value;
701 if (channel.kind === "dm") return fail("invalid", "People can't be added to a direct message. Start a new one with everyone in it.");
702 if (channel.archived_at) return fail("invalid", "This channel is archived.");
703 const p = a.member && parsePrincipalKey(`${a.member.kind}:${a.member.id}`);
704 if (!p) return fail("invalid", "Invite a person or an agent, by id.");
705 if (!(await this.belongs(slug, workspace, p))) {
706 return fail("not_found", p.kind === "agent" ? "No such agent in this workspace." : "That person is not in this workspace.");
707 }
708 await this.joinStatement(channel.id, principalKey(p), "member", now()).run();
709 return ok(null);
710 }
711
712 async setPreferences(a: {
713 workspace: string;
714 channel_id: string;
715 viewer: Viewer;
716 prefs: { starred?: boolean; muted?: boolean };
717 }): Promise<Result<null>> {
718 const found = await this.place(a.workspace, a.channel_id, a.viewer, "member");
719 if (!found.ok) return found;
720 const starred = typeof a.prefs?.starred === "boolean" ? (a.prefs.starred ? 1 : 0) : null;
721 const muted = typeof a.prefs?.muted === "boolean" ? (a.prefs.muted ? 1 : 0) : null;
722 await this.db
723 .prepare(
724 "UPDATE channel_members SET starred = COALESCE(?, starred), muted = COALESCE(?, muted) WHERE channel_id = ? AND principal = ?",
725 )
726 .bind(starred, muted, found.value.channel.id, userKey(a.viewer!))
727 .run();
728 return ok(null);
729 }
730
731 // ── Messages ────────────────────────────────────────────────────────────
732
733 /**
734 * Newest first, a page at a time. With `thread_root`, that thread's
735 * replies, and the message they reply to as the oldest once the page
736 * reaches the start of the thread; without, the channel's top-level
737 * messages. A deleted message stays only while replies hang off it.
738 */
739 async messages(a: {
740 workspace: string;
741 channel_id: string;
742 viewer: Viewer;
743 before?: string | null;
744 after?: string | null;
745 limit?: number | null;
746 thread_root?: string | null;
747 }): Promise<Result<MessagePage>> {
748 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
749 if (!found.ok) return found;
750 const { slug, workspace, channel } = found.value;
751 const size = pageSize(a.limit);
752 // Ids are lowercase letters, digits and `_`, so `~` sorts after every
753 // one: the first page reads from the newest as a range of the index.
754 const before = typeof a.before === "string" && a.before ? a.before : "~";
755 const root = typeof a.thread_root === "string" && a.thread_root ? a.thread_root : null;
756 if (typeof a.after === "string" && a.after) {
757 // Catching up after a reconnect: what came after, oldest first. Deleted
758 // ones too, so the client drops them; read one past the page to know
759 // whether there is more.
760 const newer = root
761 ? await this.db
762 .prepare("SELECT * FROM messages WHERE thread_root = ? AND channel_id = ? AND id > ? ORDER BY id LIMIT ?")
763 .bind(root, channel.id, a.after, size + 1)
764 .all<MessageRow>()
765 : await this.db
766 .prepare("SELECT * FROM messages WHERE channel_id = ? AND thread_root IS NULL AND id > ? ORDER BY id LIMIT ?")
767 .bind(channel.id, a.after, size + 1)
768 .all<MessageRow>();
769 const page = pageOf(newer.results, size);
770 return ok({ messages: await this.toMessages(slug, workspace, page.rows), older: null, newer: page.older });
771 }
772 const rows = root
773 ? await this.db
774 .prepare(
775 `SELECT * FROM messages WHERE thread_root = ?1 AND channel_id = ?2 AND id < ?3
776 ORDER BY id DESC LIMIT ?4`,
777 )
778 .bind(root, channel.id, before, size + 1)
779 .all<MessageRow>()
780 : await this.db
781 .prepare(
782 `SELECT * FROM messages
783 WHERE channel_id = ?1 AND thread_root IS NULL AND id < ?2
784 AND (deleted_at IS NULL OR reply_count > 0)
785 ORDER BY id DESC LIMIT ?3`,
786 )
787 .bind(channel.id, before, size + 1)
788 .all<MessageRow>();
789 const page = pageOf(rows.results, size);
790 let list = page.rows;
791 if (root && page.older === null) {
792 const first = await this.messageRow(channel.id, root);
793 if (first) list = [...list, first];
794 }
795 return ok({ messages: await this.toMessages(slug, workspace, list), older: page.older });
796 }
797
798 /**
799 * Writes a message and everything that follows from it: the thread's
800 * reply count, the channel's last activity, the author's own read mark,
801 * the meter; then, after answering, tells the room and wakes the agents
802 * it is for.
803 */
804 private async write(
805 place: Place,
806 author: string,
807 input: { body: string; card: MessageCard | null; thread_root: string | null },
808 chain: Chain<AskerAccess>,
809 ): Promise<Result<ChatMessage>> {
810 const { channel, workspace } = place;
811 if (channel.archived_at) return fail("invalid", "This channel is archived.");
812 let threadRoot: string | null = null;
813 if (input.thread_root) {
814 const root = await this.messageRow(channel.id, input.thread_root);
815 if (!root || (root.deleted_at && !root.reply_count)) return fail("not_found", "No such message to reply to.");
816 // A reply to a reply goes in the same thread.
817 threadRoot = root.thread_root ?? root.id;
818 }
819 const at = now();
820 const handles = mentionedHandles(input.body);
821 const row: MessageRow = {
822 id: newId("msg"),
823 channel_id: channel.id,
824 author,
825 kind: input.card ? "card" : "text",
826 body: input.body,
827 card: input.card ? JSON.stringify(input.card) : null,
828 mentions: mentionsColumn(handles),
829 thread_root: threadRoot,
830 reply_count: 0,
831 last_reply_at: null,
832 created_at: at,
833 edited_at: null,
834 deleted_at: null,
835 };
836 const statements = [
837 this.db
838 .prepare(
839 "INSERT INTO messages (id, channel_id, author, kind, body, card, mentions, thread_root, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
840 )
841 .bind(row.id, row.channel_id, author, row.kind, row.body, row.card, row.mentions, threadRoot, at),
842 this.db.prepare("UPDATE channels SET last_message_at = ? WHERE id = ?").bind(at, channel.id),
843 // What you wrote, you have read.
844 this.db
845 .prepare("UPDATE channel_members SET last_read_id = ? WHERE channel_id = ? AND principal = ?")
846 .bind(row.id, channel.id, author),
847 this.db
848 .prepare(
849 `INSERT INTO chat_meter (workspace_id, day, messages, bytes) VALUES (?, ?, 1, ?)
850 ON CONFLICT (workspace_id, day) DO UPDATE SET messages = messages + 1, bytes = bytes + excluded.bytes`,
851 )
852 .bind(workspace.id, meterDay(at), bytesOf(row.body) + bytesOf(row.card ?? "")),
853 ];
854 if (threadRoot) {
855 statements.push(
856 this.db
857 .prepare("UPDATE messages SET reply_count = reply_count + 1, last_reply_at = ? WHERE id = ?")
858 .bind(at, threadRoot),
859 );
860 }
861 await this.db.batch(statements);
862
863 const [message] = await this.toMessages(place.slug, workspace, [row]);
864 this.broadcast(channel.id, { type: "message.created", message });
865 if (threadRoot) this.rebroadcast(place, threadRoot);
866 this.defer(
867 this.wake(place, row, handles, chain).catch((error) => console.error("chat could not hand", row.id, "to agents", error)),
868 );
869 return ok(message);
870 }
871
872 /** Hands a new message to the agents it is for (src/delivery.ts). */
873 private async wake(place: Place, row: MessageRow, handles: string[], chain: Chain<AskerAccess>): Promise<void> {
874 const { channel, workspace } = place;
875 // Only an agent's message mentioning someone, or a person's, can wake anyone.
876 if (row.author.startsWith("agent:") && !handles.length) return;
877 if (channel.kind === "channel" && !handles.length) return;
878 const members = await this.db
879 .prepare("SELECT principal FROM channel_members WHERE channel_id = ? AND principal LIKE 'agent:%'")
880 .bind(channel.id)
881 .all<{ principal: string }>();
882 const ids = members.results.map((m) => m.principal.slice("agent:".length));
883 if (!ids.length) return;
884 const found = await this.agentsById(ids);
885 const agents = [...found.values()].filter(
886 (agent): agent is WorkspaceAgent => !!agent && agent.workspace_id === workspace.id && !agent.archived_at,
887 );
888 const wakes = deliveries({
889 author: row.author,
890 hops: chain.hops,
891 channelKind: channel.kind,
892 agents: agents.map((agent) => ({ id: agent.id, handle: agent.handle })),
893 mentioned: handles,
894 });
895 const client = workspaceAgentsClient(this.env.AGENTS);
896 await Promise.all(
897 wakes.map((wake) =>
898 client
899 .deliver(
900 delivery(
901 {
902 workspace: place.slug,
903 workspace_id: workspace.id,
904 channel_id: channel.id,
905 channel_kind: channel.kind,
906 channel_name: channel.kind === "dm" ? null : channel.name,
907 },
908 wake,
909 row,
910 chain,
911 ) satisfies AgentDelivery,
912 )
913 .then((result) => {
914 if (!result.ok) console.error("agents refused delivery to", wake.agent_id, result.error.message);
915 })
916 .catch((error) => console.error("chat could not deliver to", wake.agent_id, error)),
917 ),
918 );
919 }
920
921 async post(a: { workspace: string; channel_id: string; viewer: Viewer; message: PostMessage }): Promise<Result<ChatMessage>> {
922 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
923 if (!found.ok) return found;
924 const body = messageBody(a.message?.body);
925 if (!body.ok) return fail("invalid", body.message);
926 const me = userKey(a.viewer!);
927 let place = found.value;
928 if (!place.member) {
929 // Saying something in a public channel joins it, as reading does not.
930 if (place.channel.archived_at) return fail("invalid", "This channel is archived.");
931 await this.joinStatement(place.channel.id, me, "member", now()).run();
932 place = { ...place, member: { channel_id: place.channel.id, principal: me } as MemberRow };
933 }
934 return this.write(
935 place,
936 me,
937 { body: body.body, card: null, thread_root: a.message?.thread_root ?? null },
938 // A person's message starts a chain.
939 { hops: 0, asked_by: a.viewer!.id, asker: askerAccess(a.viewer!, a.workspace) },
940 );
941 }
942
943 async edit(a: { workspace: string; channel_id: string; viewer: Viewer; id: string; body: string }): Promise<Result<ChatMessage>> {
944 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
945 if (!found.ok) return found;
946 const place = found.value;
947 const row = await this.messageRow(place.channel.id, a.id);
948 if (!row || row.deleted_at) return fail("not_found", "No such message.");
949 if (row.author !== userKey(a.viewer!)) return fail("forbidden", "Only its author can edit a message.");
950 const body = messageBody(a.body, !!row.card);
951 if (!body.ok) return fail("invalid", body.message);
952 const at = now();
953 const mentions = mentionsColumn(mentionedHandles(body.body));
954 await this.db
955 .prepare("UPDATE messages SET body = ?, mentions = ?, edited_at = ? WHERE id = ?")
956 .bind(body.body, mentions, at, row.id)
957 .run();
958 const [message] = await this.toMessages(place.slug, place.workspace, [{ ...row, body: body.body, mentions, edited_at: at }]);
959 this.broadcast(place.channel.id, { type: "message.updated", message });
960 return ok(message);
961 }
962
963 /** Deletes a message, keeping its place so its thread still hangs together. Its author or a workspace owner may. */
964 async remove(a: { workspace: string; channel_id: string; viewer: Viewer; id: string }): Promise<Result<null>> {
965 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
966 if (!found.ok) return found;
967 const place = found.value;
968 const row = await this.messageRow(place.channel.id, a.id);
969 if (!row || row.deleted_at) return fail("not_found", "No such message.");
970 if (row.author !== userKey(a.viewer!) && !isOwner(a.viewer, a.workspace)) {
971 return fail("forbidden", "Only its author or a workspace owner can delete a message.");
972 }
973 const statements = [
974 this.db
975 .prepare("UPDATE messages SET deleted_at = ?, body = '', card = NULL, mentions = '' WHERE id = ?")
976 .bind(now(), row.id),
977 ];
978 if (row.thread_root) {
979 statements.push(
980 this.db.prepare("UPDATE messages SET reply_count = MAX(reply_count - 1, 0) WHERE id = ?").bind(row.thread_root),
981 );
982 }
983 await this.db.batch(statements);
984 this.broadcast(place.channel.id, { type: "message.deleted", channel_id: place.channel.id, id: row.id });
985 if (row.thread_root) this.rebroadcast(place, row.thread_root);
986 return ok(null);
987 }
988
989 async markRead(a: { workspace: string; channel_id: string; viewer: Viewer; id: string }): Promise<Result<null>> {
990 const found = await this.place(a.workspace, a.channel_id, a.viewer, "read");
991 if (!found.ok) return found;
992 const { channel, member } = found.value;
993 // Someone reading a public channel they have not joined keeps no read state.
994 if (!member) return ok(null);
995 const id = typeof a.id === "string" ? a.id : "";
996 if (!id) return fail("invalid", "Say which message was read.");
997 // Only forward: reading an old thread does not mark newer messages unread.
998 const changed = await this.db
999 .prepare(
1000 "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)",
1001 )
1002 .bind(id, channel.id, member.principal)
1003 .run();
1004 if (changed.meta.changes) {
1005 this.broadcast(channel.id, {
1006 type: "read",
1007 channel_id: channel.id,
1008 principal: { kind: "user", id: a.viewer!.id },
1009 last_read_id: id,
1010 });
1011 }
1012 return ok(null);
1013 }
1014
1015 // ── Agents ──────────────────────────────────────────────────────────────
1016
1017 /** The channel and agent for an agent's call: the agent must be of the workspace and in the channel. */
1018 private async agentPlace(slug: string, channelId: string, agentId: string): Promise<Result<{ place: Place; agent: WorkspaceAgent }>> {
1019 const workspace = await this.workspace(String(slug ?? ""));
1020 if (!workspace) return fail("not_found", "No such workspace.");
1021 const agent = await this.liveAgent(workspace, String(agentId ?? ""));
1022 if (!agent) return fail("not_found", "No such agent in this workspace.");
1023 const key = principalKey({ kind: "agent", id: agent.id });
1024 const [channel, member] = await Promise.all([
1025 this.db
1026 .prepare("SELECT * FROM channels WHERE id = ? AND workspace_id = ?")
1027 .bind(String(channelId ?? ""), workspace.id)
1028 .first<ChannelRow>(),
1029 this.db.prepare("SELECT * FROM channel_members WHERE channel_id = ? AND principal = ?").bind(String(channelId ?? ""), key).first<MemberRow>(),
1030 ]);
1031 if (!channel) return fail("not_found", "No such channel.");
1032 if (!member) return fail("forbidden", "The agent is not a member of this channel.");
1033 return ok({ place: { slug: slug.toLowerCase(), workspace, channel, member }, agent });
1034 }
1035
1036 async postAsAgent(a: { workspace: string; channel_id: string; agent_id: string; message: AgentPostMessage }): Promise<Result<ChatMessage>> {
1037 const found = await this.agentPlace(a.workspace, a.channel_id, a.agent_id);
1038 if (!found.ok) return found;
1039 const { place, agent } = found.value;
1040 const card = a.message?.card == null ? null : cleanCard(a.message.card);
1041 if (a.message?.card != null && !card) return fail("invalid", "A card needs a kind and a title.");
1042 const body = messageBody(a.message?.body ?? "", !!card);
1043 if (!body.ok) return fail("invalid", body.message);
1044 const hops = typeof a.message?.hops === "number" && a.message.hops >= 0 ? Math.floor(a.message.hops) : 0;
1045 const askedBy = typeof a.message?.asked_by === "string" && a.message.asked_by ? a.message.asked_by : agent.created_by;
1046 return this.write(
1047 place,
1048 principalKey({ kind: "agent", id: agent.id }),
1049 { body: body.body, card, thread_root: a.message?.thread_root ?? null },
1050 // The asker carries on from the delivery the agent is answering;
1051 // without one, agents it wakes treat the asker as unable to change code.
1052 { hops, asked_by: askedBy, asker: cleanAsker(a.message?.asker) },
1053 );
1054 }
1055
1056 /**
1057 * What an agent reads before replying, oldest first: a thread (its root,
1058 * then its latest replies), or the channel's latest top-level messages.
1059 * Only where the agent is a member, so it reads only what was said where
1060 * it was invited. Deleted messages are left out, save a thread's root.
1061 */
1062 async historyForAgent(a: {
1063 workspace: string;
1064 channel_id: string;
1065 agent_id: string;
1066 thread_root?: string | null;
1067 limit?: number | null;
1068 }): Promise<Result<ChatMessage[]>> {
1069 const found = await this.agentPlace(a.workspace, a.channel_id, a.agent_id);
1070 if (!found.ok) return found;
1071 const { place } = found.value;
1072 const size = historySize(a.limit);
1073 if (typeof a.thread_root === "string" && a.thread_root) {
1074 const asked = await this.messageRow(place.channel.id, a.thread_root);
1075 if (!asked) return fail("not_found", "No such thread.");
1076 // Asked from a reply: its whole thread.
1077 const root = asked.thread_root ? await this.messageRow(place.channel.id, asked.thread_root) : asked;
1078 if (!root) return fail("not_found", "No such thread.");
1079 const replies = await this.db
1080 .prepare(
1081 "SELECT * FROM messages WHERE thread_root = ? AND channel_id = ? AND deleted_at IS NULL ORDER BY id DESC LIMIT ?",
1082 )
1083 .bind(root.id, place.channel.id, Math.max(0, size - 1))
1084 .all<MessageRow>();
1085 return ok(await this.toMessages(place.slug, place.workspace, historyOf(replies.results, root)));
1086 }
1087 const rows = await this.db
1088 .prepare(
1089 "SELECT * FROM messages WHERE channel_id = ? AND thread_root IS NULL AND id < '~' AND deleted_at IS NULL ORDER BY id DESC LIMIT ?",
1090 )
1091 .bind(place.channel.id, size)
1092 .all<MessageRow>();
1093 return ok(await this.toMessages(place.slug, place.workspace, historyOf(rows.results)));
1094 }
1095
1096 async agentTyping(a: { workspace: string; channel_id: string; agent_id: string }): Promise<Result<null>> {
1097 const found = await this.agentPlace(a.workspace, a.channel_id, a.agent_id);
1098 if (!found.ok) return found;
1099 const { place, agent } = found.value;
1100 const key = principalKey({ kind: "agent", id: agent.id });
1101 const member = await this.profile(place.slug, place.workspace, key);
1102 this.broadcast(place.channel.id, {
1103 type: "typing",
1104 channel_id: place.channel.id,
1105 member,
1106 until: new Date(Date.now() + AGENT_TYPING_MS).toISOString(),
1107 });
1108 return ok(null);
1109 }
1110
1111 // ── The live socket ─────────────────────────────────────────────────────
1112
1113 /**
1114 * `GET /live?workspace=<slug>&channel=<id>`, upgraded to a WebSocket. The
1115 * viewer comes in CHAT_VIEWER_HEADER, set by the site after checking the
1116 * session; trusted only because this Worker is reachable through service
1117 * bindings alone (`workers_dev` is off and it has no routes). Checked
1118 * like any read, then handed to the channel's room.
1119 */
1120 async live(request: Request): Promise<Response> {
1121 if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") {
1122 return new Response("Expected a WebSocket upgrade\n", { status: 426 });
1123 }
1124 let viewer: Viewer = null;
1125 try {
1126 viewer = JSON.parse(request.headers.get(CHAT_VIEWER_HEADER) ?? "null") as Viewer;
1127 } catch {
1128 viewer = null;
1129 }
1130 if (!viewer?.id) return new Response("Sign in to use chat\n", { status: 401 });
1131 const url = new URL(request.url);
1132 const channelId = url.searchParams.get("channel") ?? "";
1133 let slug = (url.searchParams.get("workspace") ?? "").toLowerCase();
1134 if (!slug) {
1135 // Not named: whichever of the viewer's workspaces holds the channel.
1136 const row = await this.db.prepare("SELECT workspace_id FROM channels WHERE id = ?").bind(channelId).first<{ workspace_id: string }>();
1137 if (row) {
1138 const theirs = await Promise.all((viewer.workspaces ?? []).map((m) => this.workspace(m.slug)));
1139 slug = theirs.find((w) => w?.id === row.workspace_id)?.slug.toLowerCase() ?? "";
1140 }
1141 }
1142 const found = await this.place(slug, channelId, viewer, "read");
1143 if (!found.ok) return new Response(`${found.error.message}\n`, { status: found.error.code === "forbidden" ? 403 : 404 });
1144 const { workspace, channel } = found.value;
1145 const who: RoomMember = { channel_id: channel.id, member: await this.profile(slug, workspace, userKey(viewer)) };
1146 const headers = new Headers(request.headers);
1147 headers.delete(CHAT_VIEWER_HEADER);
1148 headers.set(ROOM_MEMBER_HEADER, JSON.stringify(who));
1149 return this.room(channel.id).fetch(new Request(request.url, { method: "GET", headers }));
1150 }
1151}
1152
1153/** One RPC method's answer. */
1154async function answer(service: Chat, method: string, args: any): Promise<Response> {
1155 switch (method) {
1156 case "sidebar":
1157 return Response.json(await service.sidebar(args));
1158 case "channel":
1159 return Response.json(await service.channel(args));
1160 case "channel_by_name":
1161 return Response.json(await service.channelByName(args));
1162 case "browse":
1163 return Response.json(await service.browse(args));
1164 case "create_channel":
1165 return Response.json(await service.createChannel(args));
1166 case "open_dm":
1167 return Response.json(await service.openDm(args));
1168 case "join":
1169 return Response.json(await service.join(args));
1170 case "leave":
1171 return Response.json(await service.leave(args));
1172 case "invite":
1173 return Response.json(await service.invite(args));
1174 case "messages":
1175 return Response.json(await service.messages(args));
1176 case "post":
1177 return Response.json(await service.post(args));
1178 case "edit":
1179 return Response.json(await service.edit(args));
1180 case "remove":
1181 return Response.json(await service.remove(args));
1182 case "mark_read":
1183 return Response.json(await service.markRead(args));
1184 case "set_preferences":
1185 return Response.json(await service.setPreferences(args));
1186 case "post_as_agent":
1187 return Response.json(await service.postAsAgent(args));
1188 case "agent_typing":
1189 return Response.json(await service.agentTyping(args));
1190 case "history_for_agent":
1191 return Response.json(await service.historyForAgent(args));
1192 default:
1193 return new Response("Unknown method\n", { status: 404 });
1194 }
1195}
1196
1197export default {
1198 async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
1199 const url = new URL(request.url);
1200 if (request.method === "GET" && url.pathname === "/live") {
1201 return new Chat(env, (work) => ctx.waitUntil(work)).live(request);
1202 }
1203 const match = url.pathname.match(/^\/rpc\/([a-z_]+)$/);
1204 if (request.method !== "POST" || !match) return new Response("Not found\n", { status: 404 });
1205 // A replica near the caller when it asks for one (@g1t/contracts d1.ts).
1206 const opened = openD1(env.DB, request);
1207 const service = new Chat(Object.create(env, { DB: { value: opened.db } }) as Env, (work) => ctx.waitUntil(work));
1208 const args = (await request.json().catch(() => ({}))) as any;
1209 return opened.finish(await answer(service, match[1], args));
1210 },
1211} satisfies ExportedHandler<Env>;