Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Chat and workspace agents: channels, DMs and named agents you talk to | 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 | ||
| 14 | import { | |
| 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 | ||
| 51 | import { MAX_HOPS, deliveries, delivery, type Chain } from "./delivery.ts"; | |
| 52 | import { mentionedHandles, mentionsColumn } from "./mentions.ts"; | |
| 53 | import { AGENT_TYPING_MS, historyOf, historySize, messageBody, meterDay, pageOf, pageSize } from "./messages.ts"; | |
| 54 | import { GENERAL, MAX_DM_MEMBERS, channelName, dmKey, dmMembers } from "./names.ts"; | |
| 55 | import { ROOM_MEMBER_HEADER, type ChannelRoom, type RoomMember } from "./room.ts"; | |
| 56 | import { dmTitle, sidebarOrder, tally, type UnreadRow } from "./unread.ts"; | |
| 57 | ||
| 58 | export { ChannelRoom } from "./room.ts"; | |
| 59 | ||
| 60 | // The hop limit here is the one in the contract. | |
| 61 | const SAME_HOP_LIMIT: typeof CHAT_MAX_HOPS = MAX_HOPS; | |
| 62 | void SAME_HOP_LIMIT; | |
| 63 | ||
| 64 | type Env = { | |
| 65 | DB: D1Database; | |
| 66 | IDENTITY: ServiceBinding; | |
| 67 | AGENTS: ServiceBinding; | |
| 68 | ROOMS: DurableObjectNamespace<ChannelRoom>; | |
| 69 | }; | |
| 70 | ||
| 71 | type 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 | ||
| 85 | type 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 | ||
| 95 | type 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. */ | |
| 112 | type 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". */ | |
| 115 | const MAX_UNREAD_ROWS = 5_000; | |
| 116 | /** The longest channel topic. */ | |
| 117 | const MAX_TOPIC = 250; | |
| 118 | /** How many others a direct message's sidebar entry shows. */ | |
| 119 | const DM_FACES = 4; | |
| 120 | ||
| 121 | const now = () => new Date().toISOString(); | |
| 122 | ||
| 123 | function isMember(viewer: Viewer, workspace: string): boolean { | |
| 124 | return !!viewer?.workspaces?.some((m) => m.slug === workspace.toLowerCase()); | |
| 125 | } | |
| 126 | ||
| 127 | function isOwner(viewer: Viewer, workspace: string): boolean { | |
| 128 | return !!viewer?.workspaces?.some((m) => m.slug === workspace.toLowerCase() && m.role === "owner"); | |
| 129 | } | |
| 130 | ||
| 131 | function userKey(viewer: User): string { | |
| 132 | return principalKey({ kind: "user", id: viewer.id }); | |
| 133 | } | |
| 134 | ||
| 135 | function 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. */ | |
| 151 | function 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. */ | |
| 170 | function 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 | ||
| 177 | function bytesOf(text: string): number { | |
| 178 | return new TextEncoder().encode(text).length; | |
| 179 | } | |
| 180 | ||
| 181 | class 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. */ | |
| 1154 | async 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 | ||
| 1197 | export 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>; |