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