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