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 | * One reply: an agent answering a message in chat, in this Worker, with no | |
| 3 | * sandbox (docs/WORKSPACE.md, "Two kinds of turn"). | |
| 4 | * | |
| 5 | * A reply goes through the same doors an agent run does, so there is one | |
| 6 | * way g1t meters and bills model work: | |
| 7 | * | |
| 8 | * 1. The agent's own monthly and daily caps (`budget.ts`). | |
| 9 | * 2. Whether it may use a model at all: g1t's hosted models are open as | |
| 10 | * the runner decides (`hostedOpen`, `HOSTED_AGENT_WORKSPACES` and | |
| 11 | * billing's status), or the workspace's own provider; the agent's | |
| 12 | * `providers` narrows that. | |
| 13 | * 3. A model session from integrations (`openModelSession`, task `reply`), | |
| 14 | * which picks g1t's gateway or the workspace's own provider by the | |
| 15 | * workspace's model routes, as for a run. | |
| 16 | * 4. The compute gate's reservation (`ComputeGate.admit`, kind `agent`): | |
| 17 | * the workspace's spend limit, AI credit, pauses and g1t's breaker. | |
| 18 | * 5. A billing run (`start_run`), so the reply is charged as Agent tokens: | |
| 19 | * the model at the provider's price on g1t's models, plus the agent rate | |
| 20 | * on every token; comped terms and discounts are billing's. | |
| 21 | * 6. The model call through the model proxy (the `MODELS` binding, or | |
| 22 | * `MODELS_URL` without one), with the session's token, exactly as a | |
| 23 | * sandbox makes it: the proxy holds the keys, caps the session, and | |
| 24 | * reports its tokens to billing. | |
| 25 | * 7. `finish_run` with the reply's cost and tokens, and the reservation | |
| 26 | * settled at cost. | |
| 27 | * | |
| 28 | * Whatever happens, the reply's row says so: replied, blocked (with one | |
| 29 | * short notice in the conversation, not repeated), skipped or failed (with | |
| 30 | * one short apology). | |
| 31 | */ | |
| 32 | import { | |
| 33 | type AgentDelivery, | |
| 34 | type ModelSession, | |
| 35 | type RunTicket, | |
| 36 | type ServiceBinding, | |
| 37 | ComputeGate, | |
| 38 | MODEL_ESTIMATE_MICROS, | |
| 39 | billingClient, | |
| 40 | integrationsClient, | |
| 41 | newId, | |
| 42 | } from "@g1t/contracts"; | |
| 43 | ||
| 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) | 44 | import { CHAT_MAX_HOPS } from "../../../packages/contracts/src/chat.ts"; |
| Chat and workspace agents: channels, DMs and named agents you talk to | 45 | import { hostedOpen } from "../../runner/src/hosted.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) | 46 | import { type AgentRouting as Policy, routingReader } from "../../runner/src/model-env.ts"; |
| Chat and workspace agents: channels, DMs and named agents you talk to | 47 | import { type Spent, type Tokens, budgetBlock, chargedMicros, costMicros, replyCapMicros, totalTokens } from "./budget.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) | 48 | import { HISTORY_LIMIT, fixedHello, helloAsk, systemPrompt, turns } from "./prompt.ts"; |
| 49 | import { BUILTIN_NO_MODEL, type Specialist, capMentions, orchestratorInstructions, orchestratorTier } from "./orchestrator.ts"; | |
| 50 | import { REPLY_TIER, allowedProviders, replyModel } from "./routing.ts"; | |
| 51 | import { type Row, definitionOf, periods, selectAgents, spendStatements, toAgent } from "./store.ts"; | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 52 | import { type SurfaceMessage, surfaceFor } from "./surface.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) | 53 | import { Audience } from "./audience.ts"; |
| 54 | import { audiencePorts, toolPorts } from "./ports.ts"; | |
| 55 | import { rosterLines } from "./orchestrator.ts"; | |
| 56 | import { type ToolCall, type ToolPorts, ToolBox } from "./tools.ts"; | |
| 57 | import type { Surface } from "./surface.ts"; | |
| 58 | import { type ModelAnswer, type ModelMessage, type Send, NO_TOKENS, addTokens, runTurn } from "./turn.ts"; | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 59 | |
| 60 | export type ReplyEnv = { | |
| 61 | DB: D1Database; | |
| 62 | CHAT: ServiceBinding; | |
| 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) | 63 | /** People and their access, for the audience; members and teams for the roster. */ |
| 64 | IDENTITY: ServiceBinding; | |
| 65 | /** Code, for read tools. */ | |
| 66 | REPOS: ServiceBinding; | |
| 67 | /** Issues and pull requests, for read tools. */ | |
| 68 | WORK: ServiceBinding; | |
| 69 | /** Code search, for read tools. */ | |
| 70 | SEARCH: ServiceBinding; | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 71 | BILLING: ServiceBinding; |
| 72 | INTEGRATIONS: ServiceBinding; | |
| 73 | /** The model proxy, by service binding: how replies reach a model. */ | |
| 74 | MODELS?: ServiceBinding; | |
| 75 | HOSTED_AGENT_WORKSPACES: string; | |
| 76 | AGENT_ROUTING: string; | |
| 77 | /** The model proxy by address, only where there is no `MODELS` binding (a self-hosted install pointing elsewhere). */ | |
| 78 | MODELS_URL?: string; | |
| 79 | }; | |
| 80 | ||
| 81 | /** The longest one model answer may take. */ | |
| 82 | const MODEL_TIMEOUT_MS = 90_000; | |
| 83 | /** A notice that the agent cannot reply is posted once per conversation in this long. */ | |
| 84 | const NOTICE_QUIET_MS = 6 * 60 * 60 * 1000; | |
| 85 | ||
| 86 | const APOLOGY = "Sorry, something went wrong on my side and I couldn't answer that. Try again in a moment."; | |
| 87 | ||
| 88 | /** Staff's model defaults on top of `AGENT_ROUTING`, read at most once a minute, as in the runner. */ | |
| 89 | const routingNow = routingReader(); | |
| 90 | let gate: ComputeGate | null = null; | |
| 91 | ||
| 92 | /** | |
| 93 | * Where a reply's spend shows in billing: the workspace, under the agent. | |
| 94 | * Billing keys runs and reservations by a repository; no repository is | |
| 95 | * named with an `@`, so the agent's line never mixes with a project's. | |
| 96 | */ | |
| 97 | export function billingRepo(workspace: string, handle: string): { namespace: string; name: string } { | |
| 98 | return { namespace: workspace.toLowerCase(), name: `@${handle}` }; | |
| 99 | } | |
| 100 | ||
| 101 | async function rpc<T>(service: ServiceBinding, method: string, args: object): Promise<T> { | |
| 102 | const response = await service.fetch(`https://service/rpc/${method}`, { | |
| 103 | method: "POST", | |
| 104 | headers: { "content-type": "application/json" }, | |
| 105 | body: JSON.stringify(args), | |
| 106 | }); | |
| 107 | if (!response.ok) throw new Error(`${method} failed with status ${response.status}`); | |
| 108 | return (await response.json()) as T; | |
| 109 | } | |
| 110 | ||
| 111 | type PriceTerms = { marginPercent: number; rate: number; rateOwn: number }; | |
| 112 | let terms: { value: PriceTerms; until: number } | null = null; | |
| 113 | ||
| 114 | /** The model margin and the agent rates from billing's price book, kept ten minutes. */ | |
| 115 | async function priceTerms(billing: ServiceBinding): Promise<PriceTerms> { | |
| 116 | if (terms && terms.until > Date.now()) return terms.value; | |
| 117 | type Book = { prices?: { meter: string; priceMicros?: number; price_micros?: number }[]; modelMarginPercent?: number; model_margin_percent?: number }; | |
| 118 | const book = await rpc<Book>(billing, "prices", {}).catch(() => null); | |
| 119 | const price = (meter: string) => { | |
| 120 | const found = book?.prices?.find((p) => p.meter === meter); | |
| 121 | return found?.priceMicros ?? found?.price_micros ?? 0; | |
| 122 | }; | |
| 123 | const value = { | |
| 124 | marginPercent: book?.modelMarginPercent ?? book?.model_margin_percent ?? 0, | |
| 125 | rate: price("agent_tokens"), | |
| 126 | rateOwn: price("agent_tokens_own"), | |
| 127 | }; | |
| 128 | // A failed read is tried again in a minute, not kept. | |
| 129 | terms = { value, until: Date.now() + (book ? 10 * 60_000 : 60_000) }; | |
| 130 | return value; | |
| 131 | } | |
| 132 | ||
| 133 | async function sha256Hex(text: string): Promise<string> { | |
| 134 | const digest = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(text)); | |
| 135 | return [...new Uint8Array(digest)].map((b) => b.toString(16).padStart(2, "0")).join(""); | |
| 136 | } | |
| 137 | ||
| 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) | 138 | /** |
| 139 | * How a reply asks the model: one Messages API request through the model | |
| 140 | * proxy, with the model session's token, non-streamed. | |
| 141 | */ | |
| 142 | function sendFor(env: ReplyEnv, token: string): Send { | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 143 | // The binding, not the proxy's public address: a Worker fetching another |
| 144 | // Worker's domain on the same zone can be refused or loop. The proxy reads | |
| 145 | // only the path and the token, so the host does not matter. | |
| 146 | const path = "/anthropic/v1/messages"; | |
| 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) | 147 | const fetcher = (url: string, init: RequestInit) => (env.MODELS ? env.MODELS.fetch(url, init) : fetch(url, init)); |
| 148 | return async (body) => { | |
| 149 | const response = await fetcher(env.MODELS ? `https://models${path}` : `${env.MODELS_URL!.replace(/\/+$/, "")}${path}`, { | |
| 150 | method: "POST", | |
| 151 | headers: { "content-type": "application/json", "x-api-key": token, "anthropic-version": "2023-06-01" }, | |
| 152 | body: JSON.stringify(body), | |
| 153 | signal: AbortSignal.timeout(MODEL_TIMEOUT_MS), | |
| 154 | }); | |
| 155 | const json = (await response.json().catch(() => null)) as (ModelAnswer & { error?: { message?: string } }) | null; | |
| 156 | if (!response.ok || !json) throw new Error(`the model answered ${response.status}: ${json?.error?.message ?? "no answer"}`); | |
| 157 | return json; | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 158 | }; |
| 159 | } | |
| 160 | ||
| 161 | /** Who asked, from the message that woke the agent (or their latest one). */ | |
| 162 | function askerIn(history: SurfaceMessage[], delivery: AgentDelivery): SurfaceMessage["author"] | null { | |
| 163 | const woken = history.find((m) => m.id === delivery.message_id); | |
| 164 | if (woken) return woken.author; | |
| 165 | return [...history].reverse().find((m) => m.author.kind === "user" && m.author.id === delivery.asked_by)?.author ?? null; | |
| 166 | } | |
| 167 | ||
| 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) | 168 | /** |
| 169 | * An agent's colleagues: every agent of the workspace but itself that is | |
| 170 | * not archived (docs/WORKSPACE.md, "Agents know each other"). | |
| 171 | */ | |
| 172 | async function team(db: D1Database, workspaceId: string, selfId: string, now: Date): Promise<Specialist[]> { | |
| 173 | const rows = await db | |
| 174 | .prepare(`${selectAgents("a.workspace_id = ?3 AND a.archived_at IS NULL AND a.id <> ?4")} ORDER BY a.builtin DESC, a.handle LIMIT 50`) | |
| 175 | .bind(...periods(now), workspaceId, selfId) | |
| 176 | .all<Row>(); | |
| 177 | return rows.results.map((row) => { | |
| 178 | const agent = toAgent(row, now); | |
| 179 | return { | |
| 180 | handle: agent.handle, | |
| 181 | display_name: agent.display_name, | |
| 182 | role: agent.role, | |
| 183 | title: agent.title, | |
| 184 | team: agent.team, | |
| 185 | department: agent.department, | |
| 186 | responsibilities: agent.responsibilities, | |
| 187 | status: agent.status, | |
| 188 | spent_month_micros: agent.spent_month_micros, | |
| 189 | monthly_micros: agent.budget.monthly_micros, | |
| 190 | }; | |
| 191 | }); | |
| 192 | } | |
| 193 | ||
| 194 | /** | |
| 195 | * Consulting a colleague (docs/WORKSPACE.md, "Agents know each other"): | |
| 196 | * the colleague answers in a nested turn that posts nothing, with the same | |
| 197 | * audience (so it can read no more than the conversation may), the same | |
| 198 | * asker, one hop further, its own routing limits, on the same model | |
| 199 | * session, billed to this reply. A compact card in the thread says who | |
| 200 | * asked whom. Consults don't nest: the colleague can't consult in turn. | |
| 201 | */ | |
| 202 | function consulting(input: { | |
| 203 | db: D1Database; | |
| 204 | row: Row; | |
| 205 | delivery: DeskWork; | |
| 206 | send: Send; | |
| 207 | policy: Policy; | |
| 208 | ownModel: boolean; | |
| 209 | sessionModel: string | null; | |
| 210 | own: string | null; | |
| 211 | ports: (consult: ToolPorts["consult"]) => ToolPorts; | |
| 212 | hops: number; | |
| 213 | now: Date; | |
| 214 | surface: Surface; | |
| 215 | asker: { name: string; display_name: string | null; access: NonNullable<AgentDelivery["asker"]> | null }; | |
| 216 | spent: { tokens: Tokens; cost: number }; | |
| 217 | }): { ask: ToolPorts["consult"]; attach(box: ToolBox): void } { | |
| 218 | let parent: ToolBox | null = null; | |
| 219 | const noNesting: ToolPorts["consult"] = async () => ({ ok: false, message: "Consults don't nest: answer with what you have." }); | |
| 220 | const ask: ToolPorts["consult"] = async (handle, question) => { | |
| 221 | const { row, delivery } = input; | |
| 222 | const colleague = await input.db | |
| 223 | .prepare("SELECT * FROM agents WHERE workspace_id = ? AND handle = ? AND archived_at IS NULL") | |
| 224 | .bind(row.workspace_id, handle) | |
| 225 | .first<Row>(); | |
| 226 | if (!colleague || colleague.id === row.id) return { ok: false, message: `There is no other agent called @${handle} here.` }; | |
| 227 | // No ping-pong: never back to the agent that sent this work. | |
| 228 | if ((delivery.chain ?? []).at(-1) === colleague.id) return { ok: false, message: `@${handle} sent you this work; answer with what you have.` }; | |
| 229 | if (input.hops + 1 > CHAT_MAX_HOPS) return { ok: false, message: "This request has been passed along too many times; answer with what you have." }; | |
| 230 | const definition = definitionOf(colleague); | |
| 231 | const allowed = allowedProviders(definition.routing, input.own); | |
| 232 | if (input.ownModel ? !allowed.own : !allowed.hosted) return { ok: false, message: `@${handle} can't use the model this conversation runs on.` }; | |
| 233 | const model = replyModel( | |
| 234 | input.policy, | |
| 235 | { ...definition.routing, pinned: input.ownModel ? definition.routing.pinned : null }, | |
| 236 | { named: input.ownModel ? input.sessionModel : null }, | |
| 237 | ); | |
| 238 | const tools = parent | |
| 239 | ? parent.forColleague(input.ports(noNesting), { agentId: colleague.id, notConsult: [colleague.handle, row.handle], hops: input.hops + 1, maxHops: input.hops + 1 }) | |
| 240 | : null; | |
| 241 | const system = systemPrompt({ | |
| 242 | agent: { ...definition, id: colleague.id }, | |
| 243 | workspace: delivery.workspace, | |
| 244 | channel: { kind: delivery.channel_kind, name: delivery.channel_name }, | |
| 245 | asker: input.asker, | |
| 246 | today: input.now, | |
| 247 | tools: tools ? { code: tools.definitions().some((tool) => tool.name === "read_file") } : null, | |
| 248 | consultedBy: row.handle, | |
| 249 | }); | |
| 250 | const result = await runTurn(input.send, { | |
| 251 | model: model.model, | |
| 252 | system, | |
| 253 | messages: [{ role: "user", content: `@${row.handle} (agent) asks you: ${question}` }], | |
| 254 | tools, | |
| 255 | price: input.ownModel ? null : model.price, | |
| 256 | }); | |
| 257 | input.spent.tokens = addTokens(input.spent.tokens, result.tokens); | |
| 258 | input.spent.cost += input.ownModel ? 0 : result.cost; | |
| 259 | const answer = result.text || "(no answer)"; | |
| 260 | const exchange = `Q: ${question}\nA: ${answer}`; | |
| 261 | // The exchange, collapsed, in the thread. No @ in its text: it wakes nobody. | |
| 262 | await input.surface | |
| 263 | .post(`${row.display_name} asked ${colleague.display_name}`, { | |
| 264 | kind: "consult", | |
| 265 | title: `${row.display_name} asked ${colleague.display_name}`, | |
| 266 | detail: exchange.length > 1000 ? `${exchange.slice(0, 1000)}…` : exchange, | |
| 267 | state: null, | |
| 268 | href: null, | |
| 269 | }) | |
| 270 | .catch((error: unknown) => console.error("agents: a consult card was not posted", String(error))); | |
| 271 | return { ok: true, colleague: colleague.handle, answer }; | |
| 272 | }; | |
| 273 | return { | |
| 274 | ask, | |
| 275 | attach(box) { | |
| 276 | parent = box; | |
| 277 | }, | |
| 278 | }; | |
| 279 | } | |
| 280 | ||
| Chat and workspace agents: channels, DMs and named agents you talk to | 281 | type Outcome = { |
| 282 | status: "replied" | "blocked" | "skipped" | "failed"; | |
| 283 | error?: string | null; | |
| 284 | reply_id?: string | null; | |
| 285 | model?: string | null; | |
| 286 | tier?: string | null; | |
| 287 | tokens?: Tokens; | |
| 288 | cost?: number; | |
| 289 | charged?: number; | |
| 290 | }; | |
| 291 | ||
| 292 | /** | |
| 293 | * Answers `delivery` as its agent. Never throws: every way it ends is | |
| 294 | * recorded on the reply's row. A message handed over twice is answered | |
| 295 | * once. | |
| 296 | */ | |
| 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) | 297 | /** |
| 298 | * What a desk is handed: a message to answer, or (`hello`) the agent's | |
| 299 | * first message to the person who made it, in the DM that just opened. | |
| 300 | */ | |
| 301 | export type DeskWork = AgentDelivery & { hello?: boolean }; | |
| 302 | ||
| 303 | export async function reply(env: ReplyEnv, delivery: DeskWork, now = new Date()): Promise<void> { | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 304 | const db = env.DB; |
| 305 | const row = await db.prepare("SELECT * FROM agents WHERE id = ?").bind(delivery.agent_id).first<Row>(); | |
| 306 | if (!row || row.archived_at || row.workspace_id !== delivery.workspace_id) return; | |
| 307 | const id = newId("arp", now.getTime()); | |
| 308 | const claimed = await db | |
| 309 | .prepare( | |
| 310 | `INSERT INTO agent_replies (id, agent_id, workspace_id, channel_id, message_id, asked_by, agent_version, status, created_at) | |
| 311 | VALUES (?, ?, ?, ?, ?, ?, ?, 'working', ?) | |
| 312 | ON CONFLICT (agent_id, message_id) DO NOTHING RETURNING id`, | |
| 313 | ) | |
| 314 | .bind(id, row.id, row.workspace_id, delivery.channel_id, delivery.message_id, delivery.asked_by, row.version, now.toISOString()) | |
| 315 | .first<{ id: string }>(); | |
| 316 | if (!claimed) return; | |
| 317 | ||
| 318 | const surface = surfaceFor(env.CHAT, delivery); | |
| 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) | 319 | // 👀 as soon as the desk has it, while everything else goes on. |
| 320 | const acknowledging = surface.acknowledge(); | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 321 | const slug = delivery.workspace.toLowerCase(); |
| 322 | const billing = billingClient(env.BILLING); | |
| 323 | const integrations = integrationsClient(env.INTEGRATIONS); | |
| 324 | gate ??= new ComputeGate(env.BILLING); | |
| 325 | ||
| 326 | let session: ModelSession | null = null; | |
| 327 | let reservation: string | null = null; | |
| 328 | let settled = false; | |
| 329 | // What the answer used, once there is one: counted however the reply ends. | |
| 330 | let usage: Partial<Outcome> = {}; | |
| 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) | 331 | // What the read tools did, for the audit table, and who the audience was. |
| 332 | let toolCalls: ToolCall[] = []; | |
| 333 | let audienceHash: string | null = null; | |
| 334 | // What colleagues consulted along the way used: billed to this reply. | |
| 335 | const consulted = { tokens: NO_TOKENS, cost: 0 }; | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 336 | |
| 337 | const finish = async (outcome: Outcome) => { | |
| 338 | const charged = outcome.charged ?? 0; | |
| 339 | const statements = [ | |
| 340 | db | |
| 341 | .prepare( | |
| 342 | `UPDATE agent_replies SET status = ?, error = ?, reply_id = ?, model = ?, tier = ?, input_tokens = ?, output_tokens = ?, | |
| 343 | cost_micros = ?, charged_micros = ?, finished_at = ? WHERE id = ?`, | |
| 344 | ) | |
| 345 | .bind( | |
| 346 | outcome.status, | |
| 347 | outcome.error?.slice(0, 1000) ?? null, | |
| 348 | outcome.reply_id ?? null, | |
| 349 | outcome.model ?? null, | |
| 350 | outcome.tier ?? null, | |
| 351 | (outcome.tokens?.input ?? 0) + (outcome.tokens?.cacheRead ?? 0) + (outcome.tokens?.cacheWrite ?? 0), | |
| 352 | outcome.tokens?.output ?? 0, | |
| 353 | outcome.cost ?? 0, | |
| 354 | charged, | |
| 355 | new Date().toISOString(), | |
| 356 | id, | |
| 357 | ), | |
| 358 | ...(charged > 0 ? spendStatements(db, row.id, charged, now) : []), | |
| 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) | 359 | // Every tool call: what was asked, for whom, and whether it was read or withheld. |
| 360 | ...toolCalls.map((call, n) => | |
| 361 | db | |
| 362 | .prepare( | |
| 363 | `INSERT INTO agent_tool_calls (id, reply_id, agent_id, workspace_id, asked_by, tool, args, audience_hash, outcome, bytes, created_at) | |
| 364 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, | |
| 365 | ) | |
| 366 | .bind(`${id}_${n}`, id, row.id, row.workspace_id, delivery.asked_by, call.tool, call.args, audienceHash ?? "", call.outcome, call.bytes, now.toISOString()), | |
| 367 | ), | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 368 | ]; |
| 369 | await db.batch(statements); | |
| 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) | 370 | // ✅ when it answered; 👀 taken back after a notice, an apology or nothing. |
| 371 | await acknowledging; | |
| 372 | await surface.settle(outcome.status === "replied" ? "done" : "withdrawn"); | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 373 | }; |
| 374 | ||
| 375 | /** Says once, in this conversation, why the agent cannot answer; a repeat within hours is kept back. */ | |
| 376 | const notice = async (message: string, reason: string) => { | |
| 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) | 377 | // An agent's first words are never an excuse: without a model, a fixed hello. |
| 378 | if (delivery.hello) { | |
| 379 | const posted = await surface.post(fixedHello(row, delivery.asker?.username ?? null)).catch(() => null); | |
| 380 | return await finish({ status: "blocked", error: reason, reply_id: posted }); | |
| 381 | } | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 382 | const since = new Date(now.getTime() - NOTICE_QUIET_MS).toISOString(); |
| 383 | const recent = await db | |
| 384 | .prepare( | |
| 385 | `SELECT 1 FROM agent_replies WHERE agent_id = ? AND channel_id = ? AND status = 'blocked' AND error = ? | |
| 386 | AND reply_id IS NOT NULL AND created_at > ? AND id <> ? LIMIT 1`, | |
| 387 | ) | |
| 388 | .bind(row.id, delivery.channel_id, reason, since, id) | |
| 389 | .first(); | |
| 390 | const posted = recent ? null : await surface.post(message).catch(() => null); | |
| 391 | await finish({ status: "blocked", error: reason, reply_id: posted }); | |
| 392 | }; | |
| 393 | ||
| 394 | try { | |
| 395 | if (!env.MODELS && !env.MODELS_URL) return await notice("I can't reply here: this installation has no model proxy set up.", "no_models_url"); | |
| 396 | const definition = definitionOf(row); | |
| 397 | const [month, day] = periods(now); | |
| 398 | const spentRows = await db | |
| 399 | .prepare("SELECT period, micros FROM agent_spend WHERE agent_id = ? AND period IN (?, ?)") | |
| 400 | .bind(row.id, month, day) | |
| 401 | .all<{ period: string; micros: number }>(); | |
| 402 | const spent: Spent = { | |
| 403 | month: spentRows.results.find((r) => r.period === month)?.micros ?? 0, | |
| 404 | day: spentRows.results.find((r) => r.period === day)?.micros ?? 0, | |
| 405 | }; | |
| 406 | ||
| 407 | // 1. The agent's own caps. | |
| 408 | const blocked = budgetBlock(definition.budget, spent, now); | |
| 409 | if (blocked) return await notice(blocked.message, `budget_${blocked.cap}`); | |
| 410 | ||
| 411 | // 2. Whether it may use a model at all. | |
| 412 | const [own, status] = await Promise.all([ | |
| 413 | integrations.modelProvider(slug).catch(() => null), | |
| 414 | billing.status().catch(() => ({ enabled: false, live: false })), | |
| 415 | ]); | |
| 416 | const allowed = allowedProviders(definition.routing, own?.id ?? null); | |
| 417 | const mayHosted = hostedOpen(slug, env.HOSTED_AGENT_WORKSPACES, status) && allowed.hosted; | |
| 418 | if (!mayHosted && !allowed.own) { | |
| 419 | const why = own | |
| 420 | ? "My settings don't let me use any model this workspace has. An owner can change my providers on my profile." | |
| 421 | : allowed.hosted | |
| 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) | 422 | ? row.builtin |
| 423 | ? BUILTIN_NO_MODEL | |
| 424 | : "I can't reply yet: g1t's hosted models aren't open to this workspace, and it has no model provider of its own. An owner can connect one under Integrations." | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 425 | : "My settings let me use only this workspace's own model providers, and it has none. An owner can connect one under Integrations, or change my providers on my profile."; |
| 426 | return await notice(why, "no_model"); | |
| 427 | } | |
| 428 | ||
| 429 | // Read the conversation while showing that the agent is on it. | |
| 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) | 430 | // A hello has no conversation yet: it is asked to introduce itself. |
| 431 | const [, history] = await Promise.all([surface.typing(), delivery.hello ? Promise.resolve([]) : surface.history(HISTORY_LIMIT)]); | |
| 432 | const conversation = delivery.hello ? [{ role: "user" as const, content: helloAsk(delivery.asker?.username ?? null) }] : turns(history, row.id); | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 433 | if (!conversation.length) return await finish({ status: "skipped", error: "nothing to answer" }); |
| 434 | const author = askerIn(history, delivery); | |
| 435 | const askerName = delivery.asker?.username ?? author?.name ?? 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) | 436 | // Every agent knows its colleagues; @g1t also steps up a tier to decide |
| 437 | // who gets the work in a long thread. | |
| 438 | const specialists = await team(db, row.workspace_id, row.id, now); | |
| 439 | const start = row.builtin ? orchestratorTier(history.length, specialists.filter((a) => a.handle !== "g1t").length) : REPLY_TIER; | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 440 | |
| 441 | // 3. A model session, routed by the workspace's model routes. | |
| 442 | const repo = billingRepo(slug, row.handle); | |
| 443 | const policy = await routingNow(env.AGENT_ROUTING, () => billing.modelDefaults()); | |
| 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) | 444 | const provisional = replyModel(policy, { ...definition.routing, pinned: null }, { start }); |
| Chat and workspace agents: channels, DMs and named agents you talk to | 445 | const opened = await integrations.openModelSession({ |
| 446 | workspace: slug, | |
| 447 | repo, | |
| 448 | number: 0, | |
| 449 | task: "reply", | |
| 450 | hostedOpen: mayHosted, | |
| 451 | tier: provisional.tier, | |
| 452 | requestedBy: askerName, | |
| 453 | }); | |
| 454 | if (!opened.ok) return await notice(`I can't reply right now: ${opened.error.message}`, "model_route"); | |
| 455 | session = opened.value; | |
| 456 | const ownModel = session.billedTo === "workspace"; | |
| 457 | if (ownModel ? !allowed.own : !mayHosted) { | |
| 458 | return await notice( | |
| 459 | "This workspace routes replies to a model my settings don't allow. An owner can change my providers on my profile, or the workspace's model routes under Integrations.", | |
| 460 | "provider_not_allowed", | |
| 461 | ); | |
| 462 | } | |
| 463 | // A pinned model is for the workspace's own endpoints; on g1t's models the tier decides. | |
| 464 | const model = replyModel( | |
| 465 | policy, | |
| 466 | { ...definition.routing, pinned: ownModel ? definition.routing.pinned : 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) | 467 | { chosen: session.tierChoice ?? null, named: ownModel ? session.model : null, start }, |
| Chat and workspace agents: channels, DMs and named agents you talk to | 468 | ); |
| 469 | ||
| 470 | // 4. The workspace's own limits, through the compute gate. | |
| 471 | const ent = await gate.entitlements(slug); | |
| 472 | const admission = await gate.admit( | |
| 473 | { workspace: slug, repo, public: false, kind: "agent", estimateMicros: ownModel ? 0 : MODEL_ESTIMATE_MICROS.reply, hostedModel: !ownModel }, | |
| 474 | ent, | |
| 475 | ); | |
| 476 | if (!admission.ok) return await notice(`I can't reply right now: ${admission.message}`, `workspace_${admission.code}`); | |
| 477 | reservation = admission.reservation?.id ?? null; | |
| 478 | ||
| 479 | // 5. The run the reply is billed as. | |
| 480 | const named = ownModel && session.model ? session.model : null; | |
| 481 | const started = await billing.startRun({ | |
| 482 | workspace: slug, | |
| 483 | repo, | |
| 484 | number: 0, | |
| 485 | task: "reply", | |
| 486 | model: ownModel ? `${model.modelName} (${session.providerName ?? "own provider"})` : model.modelName, | |
| 487 | billedTo: ownModel ? "workspace" : "g1t", | |
| 488 | session: session.id, | |
| 489 | tier: named ? null : model.tier, | |
| 490 | }); | |
| 491 | if (!started.ok) return await notice(`I can't reply right now: ${started.error.message}`, "billing"); | |
| 492 | const ticket: RunTicket | null = started.value; | |
| 493 | const cap = replyCapMicros(definition.budget, spent, ent && ent.runCapMicros > 0 ? ent.runCapMicros : null); | |
| 494 | if (cap) await integrations.capModelSessions([await sha256Hex(session.token)], cap).catch(() => 0); | |
| 495 | ||
| 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) | 496 | // 6. The answer, with read tools within the audience (a hello reads nothing). |
| 497 | const send = sendFor(env, session.token); | |
| 498 | const hops = Math.max(0, Math.floor(delivery.hops || 0)); | |
| 499 | const chain = delivery.chain ?? []; | |
| 500 | const sender = chain.length ? await db.prepare("SELECT handle FROM agents WHERE id = ?").bind(chain[chain.length - 1]).first<{ handle: string }>() : null; | |
| 501 | let toolbox: ToolBox | null = null; | |
| 502 | if (!delivery.hello) { | |
| 503 | try { | |
| 504 | const audience = await Audience.build(slug, delivery.asked_by, audiencePorts(env, slug, delivery.channel_id)); | |
| 505 | const ports = (consult: ToolPorts["consult"]) => toolPorts(env, slug, row.workspace_id, delivery.channel_id, consult); | |
| 506 | // A colleague's answer for this agent: a nested turn that posts nothing. | |
| 507 | audienceHash = audience.hash; | |
| 508 | const consult = consulting({ | |
| 509 | db, | |
| 510 | row, | |
| 511 | delivery, | |
| 512 | send, | |
| 513 | policy, | |
| 514 | ownModel, | |
| 515 | sessionModel: session.model, | |
| 516 | own: own?.id ?? null, | |
| 517 | ports, | |
| 518 | hops, | |
| 519 | now, | |
| 520 | surface, | |
| 521 | asker: { name: askerName ?? "someone", display_name: author?.display_name ?? null, access: delivery.asker ?? null }, | |
| 522 | spent: consulted, | |
| 523 | }); | |
| 524 | toolbox = new ToolBox(audience, ports(consult.ask), { | |
| 525 | agentId: row.id, | |
| 526 | notConsult: [row.handle, ...(sender ? [sender.handle] : [])], | |
| 527 | hops, | |
| 528 | maxHops: CHAT_MAX_HOPS, | |
| 529 | }); | |
| 530 | consult.attach(toolbox); | |
| 531 | } catch (error) { | |
| 532 | // Without an audience nothing may be read: the reply goes on with this conversation only. | |
| 533 | console.error("agents: no audience for a reply, so no tools", row.id, String(error)); | |
| 534 | } | |
| 535 | } | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 536 | const system = systemPrompt({ |
| 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) | 537 | agent: { |
| 538 | ...definition, | |
| 539 | id: row.id, | |
| 540 | // @g1t's job is fixed; what the workspace wrote is added to it. | |
| 541 | instructions: row.builtin ? orchestratorInstructions(specialists.filter((a) => a.handle !== "g1t"), definition.instructions) : definition.instructions, | |
| 542 | }, | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 543 | workspace: delivery.workspace, |
| 544 | channel: { kind: delivery.channel_kind, name: delivery.channel_name }, | |
| 545 | asker: { name: askerName ?? "someone", display_name: author?.display_name ?? null, access: delivery.asker ?? null }, | |
| 546 | today: now, | |
| 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) | 547 | tools: toolbox ? { code: toolbox.definitions().some((tool) => tool.name === "read_file") } : null, |
| 548 | // @g1t's team is in its job; everyone else is told who their colleagues are. | |
| 549 | colleagues: row.builtin ? null : rosterLines(specialists), | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 550 | }); |
| 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) | 551 | toolCalls = toolbox?.calls ?? []; |
| 552 | const answer = await runTurn(send, { model: model.model, system, messages: conversation as ModelMessage[], tools: toolbox, price: ownModel ? null : model.price }); | |
| 553 | // Colleagues consulted along the way were billed to this reply. | |
| 554 | const tokens = addTokens(answer.tokens, consulted.tokens); | |
| 555 | const cost = (ownModel ? 0 : answer.cost) + consulted.cost; | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 556 | const priced = await priceTerms(env.BILLING); |
| 557 | const charged = chargedMicros({ | |
| 558 | costMicros: cost, | |
| 559 | hosted: !ownModel, | |
| 560 | marginPercent: priced.marginPercent, | |
| 561 | ratePerMillionMicros: ownModel ? priced.rateOwn : priced.rate, | |
| 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) | 562 | tokens: totalTokens(tokens), |
| Chat and workspace agents: channels, DMs and named agents you talk to | 563 | }); |
| 564 | ||
| 565 | // 7. Bill it, whether or not the answer could be posted: the tokens were used. | |
| 566 | if (ticket) { | |
| 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) | 567 | await rpc(env.BILLING, "finish_run", { runId: ticket.runId, token: ticket.token, costUsd: cost / 1_000_000, turns: answer.rounds, tokens }).catch( |
| Chat and workspace agents: channels, DMs and named agents you talk to | 568 | (error: unknown) => console.error("agents: finish_run failed", ticket.runId, String(error)), |
| 569 | ); | |
| 570 | } | |
| 571 | if (reservation) { | |
| 572 | settled = true; | |
| 573 | await gate.settle(reservation, cost); | |
| 574 | } | |
| 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) | 575 | usage = { model: model.modelName, tier: named ? null : model.tier, tokens, cost, charged }; |
| Chat and workspace agents: channels, DMs and named agents you talk to | 576 | if (!answer.text) { |
| 577 | const posted = await surface.post(APOLOGY).catch(() => null); | |
| 578 | return await finish({ status: "failed", error: "the model gave no text", reply_id: posted, ...usage }); | |
| 579 | } | |
| 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) | 580 | // At most two specialists woken by one of @g1t's messages: a rail, not only a rule in its prompt. |
| 581 | const text = row.builtin ? capMentions(answer.text, specialists.map((agent) => agent.handle)) : answer.text; | |
| 582 | const posted = await surface.post(text); | |
| Chat and workspace agents: channels, DMs and named agents you talk to | 583 | await finish({ status: "replied", reply_id: posted, ...usage }); |
| 584 | } catch (error) { | |
| 585 | const message = error instanceof Error ? error.message : String(error); | |
| 586 | console.error("agents: a reply failed", row.id, delivery.message_id, message); | |
| 587 | const posted = await surface.post(APOLOGY).catch(() => null); | |
| 588 | await finish({ status: "failed", error: message, reply_id: posted, ...usage }).catch((failure: unknown) => | |
| 589 | console.error("agents: a failed reply was not recorded", id, String(failure)), | |
| 590 | ); | |
| 591 | } finally { | |
| 592 | // What was held is given back however the reply ended, and its model | |
| 593 | // session's token stops working. | |
| 594 | if (reservation && !settled) await gate.settle(reservation, 0); | |
| 595 | if (session) await integrations.closeModelSessions([await sha256Hex(session.token)]).catch(() => 0); | |
| 596 | } | |
| 597 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.