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 | ||
| 44 | import { hostedOpen } from "../../runner/src/hosted.ts"; | |
| 45 | import { routingReader } from "../../runner/src/model-env.ts"; | |
| 46 | import { type Spent, type Tokens, budgetBlock, chargedMicros, costMicros, replyCapMicros, totalTokens } from "./budget.ts"; | |
| 47 | import { HISTORY_LIMIT, audienceFor, systemPrompt, turns } from "./prompt.ts"; | |
| 48 | import { allowedProviders, replyModel } from "./routing.ts"; | |
| 49 | import { type Row, definitionOf, periods, spendStatements } from "./store.ts"; | |
| 50 | import { type SurfaceMessage, surfaceFor } from "./surface.ts"; | |
| 51 | ||
| 52 | export type ReplyEnv = { | |
| 53 | DB: D1Database; | |
| 54 | CHAT: ServiceBinding; | |
| 55 | BILLING: ServiceBinding; | |
| 56 | INTEGRATIONS: ServiceBinding; | |
| 57 | /** The model proxy, by service binding: how replies reach a model. */ | |
| 58 | MODELS?: ServiceBinding; | |
| 59 | HOSTED_AGENT_WORKSPACES: string; | |
| 60 | AGENT_ROUTING: string; | |
| 61 | /** The model proxy by address, only where there is no `MODELS` binding (a self-hosted install pointing elsewhere). */ | |
| 62 | MODELS_URL?: string; | |
| 63 | }; | |
| 64 | ||
| 65 | /** The longest one model answer may take. */ | |
| 66 | const MODEL_TIMEOUT_MS = 90_000; | |
| 67 | /** A chat answer is short; this bounds the cost of one that is not. */ | |
| 68 | const MAX_OUTPUT_TOKENS = 2048; | |
| 69 | /** A notice that the agent cannot reply is posted once per conversation in this long. */ | |
| 70 | const NOTICE_QUIET_MS = 6 * 60 * 60 * 1000; | |
| 71 | ||
| 72 | const APOLOGY = "Sorry, something went wrong on my side and I couldn't answer that. Try again in a moment."; | |
| 73 | ||
| 74 | /** Staff's model defaults on top of `AGENT_ROUTING`, read at most once a minute, as in the runner. */ | |
| 75 | const routingNow = routingReader(); | |
| 76 | let gate: ComputeGate | null = null; | |
| 77 | ||
| 78 | /** | |
| 79 | * Where a reply's spend shows in billing: the workspace, under the agent. | |
| 80 | * Billing keys runs and reservations by a repository; no repository is | |
| 81 | * named with an `@`, so the agent's line never mixes with a project's. | |
| 82 | */ | |
| 83 | export function billingRepo(workspace: string, handle: string): { namespace: string; name: string } { | |
| 84 | return { namespace: workspace.toLowerCase(), name: `@${handle}` }; | |
| 85 | } | |
| 86 | ||
| 87 | async function rpc<T>(service: ServiceBinding, method: string, args: object): Promise<T> { | |
| 88 | const response = await service.fetch(`https://service/rpc/${method}`, { | |
| 89 | method: "POST", | |
| 90 | headers: { "content-type": "application/json" }, | |
| 91 | body: JSON.stringify(args), | |
| 92 | }); | |
| 93 | if (!response.ok) throw new Error(`${method} failed with status ${response.status}`); | |
| 94 | return (await response.json()) as T; | |
| 95 | } | |
| 96 | ||
| 97 | type PriceTerms = { marginPercent: number; rate: number; rateOwn: number }; | |
| 98 | let terms: { value: PriceTerms; until: number } | null = null; | |
| 99 | ||
| 100 | /** The model margin and the agent rates from billing's price book, kept ten minutes. */ | |
| 101 | async function priceTerms(billing: ServiceBinding): Promise<PriceTerms> { | |
| 102 | if (terms && terms.until > Date.now()) return terms.value; | |
| 103 | type Book = { prices?: { meter: string; priceMicros?: number; price_micros?: number }[]; modelMarginPercent?: number; model_margin_percent?: number }; | |
| 104 | const book = await rpc<Book>(billing, "prices", {}).catch(() => null); | |
| 105 | const price = (meter: string) => { | |
| 106 | const found = book?.prices?.find((p) => p.meter === meter); | |
| 107 | return found?.priceMicros ?? found?.price_micros ?? 0; | |
| 108 | }; | |
| 109 | const value = { | |
| 110 | marginPercent: book?.modelMarginPercent ?? book?.model_margin_percent ?? 0, | |
| 111 | rate: price("agent_tokens"), | |
| 112 | rateOwn: price("agent_tokens_own"), | |
| 113 | }; | |
| 114 | // A failed read is tried again in a minute, not kept. | |
| 115 | terms = { value, until: Date.now() + (book ? 10 * 60_000 : 60_000) }; | |
| 116 | return value; | |
| 117 | } | |
| 118 | ||
| 119 | async function sha256Hex(text: string): Promise<string> { | |
| 120 | const digest = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(text)); | |
| 121 | return [...new Uint8Array(digest)].map((b) => b.toString(16).padStart(2, "0")).join(""); | |
| 122 | } | |
| 123 | ||
| 124 | type Answer = { text: string; tokens: Tokens }; | |
| 125 | ||
| 126 | /** One non-streamed answer from the model proxy, in Anthropic's Messages format. */ | |
| 127 | async function askModel( | |
| 128 | env: ReplyEnv, | |
| 129 | token: string, | |
| 130 | body: { model: string; system: string; messages: { role: string; content: string }[] }, | |
| 131 | ): Promise<Answer> { | |
| 132 | // The binding, not the proxy's public address: a Worker fetching another | |
| 133 | // Worker's domain on the same zone can be refused or loop. The proxy reads | |
| 134 | // only the path and the token, so the host does not matter. | |
| 135 | const path = "/anthropic/v1/messages"; | |
| 136 | const send = (url: string, init: RequestInit) => (env.MODELS ? env.MODELS.fetch(url, init) : fetch(url, init)); | |
| 137 | const response = await send(env.MODELS ? `https://models${path}` : `${env.MODELS_URL!.replace(/\/+$/, "")}${path}`, { | |
| 138 | method: "POST", | |
| 139 | headers: { "content-type": "application/json", "x-api-key": token, "anthropic-version": "2023-06-01" }, | |
| 140 | body: JSON.stringify({ ...body, max_tokens: MAX_OUTPUT_TOKENS }), | |
| 141 | signal: AbortSignal.timeout(MODEL_TIMEOUT_MS), | |
| 142 | }); | |
| 143 | const json = (await response.json().catch(() => null)) as { | |
| 144 | content?: { type: string; text?: string }[]; | |
| 145 | usage?: { input_tokens?: number; output_tokens?: number; cache_read_input_tokens?: number; cache_creation_input_tokens?: number }; | |
| 146 | error?: { message?: string }; | |
| 147 | } | null; | |
| 148 | if (!response.ok || !json) throw new Error(`the model answered ${response.status}: ${json?.error?.message ?? "no answer"}`); | |
| 149 | const text = (json.content ?? []) | |
| 150 | .filter((block) => block.type === "text" && block.text) | |
| 151 | .map((block) => block.text!.trim()) | |
| 152 | .join("\n\n") | |
| 153 | .trim(); | |
| 154 | const usage = json.usage ?? {}; | |
| 155 | return { | |
| 156 | text, | |
| 157 | tokens: { | |
| 158 | input: usage.input_tokens ?? 0, | |
| 159 | output: usage.output_tokens ?? 0, | |
| 160 | cacheRead: usage.cache_read_input_tokens ?? 0, | |
| 161 | cacheWrite: usage.cache_creation_input_tokens ?? 0, | |
| 162 | }, | |
| 163 | }; | |
| 164 | } | |
| 165 | ||
| 166 | /** Who asked, from the message that woke the agent (or their latest one). */ | |
| 167 | function askerIn(history: SurfaceMessage[], delivery: AgentDelivery): SurfaceMessage["author"] | null { | |
| 168 | const woken = history.find((m) => m.id === delivery.message_id); | |
| 169 | if (woken) return woken.author; | |
| 170 | return [...history].reverse().find((m) => m.author.kind === "user" && m.author.id === delivery.asked_by)?.author ?? null; | |
| 171 | } | |
| 172 | ||
| 173 | type Outcome = { | |
| 174 | status: "replied" | "blocked" | "skipped" | "failed"; | |
| 175 | error?: string | null; | |
| 176 | reply_id?: string | null; | |
| 177 | model?: string | null; | |
| 178 | tier?: string | null; | |
| 179 | tokens?: Tokens; | |
| 180 | cost?: number; | |
| 181 | charged?: number; | |
| 182 | }; | |
| 183 | ||
| 184 | /** | |
| 185 | * Answers `delivery` as its agent. Never throws: every way it ends is | |
| 186 | * recorded on the reply's row. A message handed over twice is answered | |
| 187 | * once. | |
| 188 | */ | |
| 189 | export async function reply(env: ReplyEnv, delivery: AgentDelivery, now = new Date()): Promise<void> { | |
| 190 | const db = env.DB; | |
| 191 | const row = await db.prepare("SELECT * FROM agents WHERE id = ?").bind(delivery.agent_id).first<Row>(); | |
| 192 | if (!row || row.archived_at || row.workspace_id !== delivery.workspace_id) return; | |
| 193 | const id = newId("arp", now.getTime()); | |
| 194 | const claimed = await db | |
| 195 | .prepare( | |
| 196 | `INSERT INTO agent_replies (id, agent_id, workspace_id, channel_id, message_id, asked_by, agent_version, status, created_at) | |
| 197 | VALUES (?, ?, ?, ?, ?, ?, ?, 'working', ?) | |
| 198 | ON CONFLICT (agent_id, message_id) DO NOTHING RETURNING id`, | |
| 199 | ) | |
| 200 | .bind(id, row.id, row.workspace_id, delivery.channel_id, delivery.message_id, delivery.asked_by, row.version, now.toISOString()) | |
| 201 | .first<{ id: string }>(); | |
| 202 | if (!claimed) return; | |
| 203 | ||
| 204 | const surface = surfaceFor(env.CHAT, delivery); | |
| 205 | const slug = delivery.workspace.toLowerCase(); | |
| 206 | const billing = billingClient(env.BILLING); | |
| 207 | const integrations = integrationsClient(env.INTEGRATIONS); | |
| 208 | gate ??= new ComputeGate(env.BILLING); | |
| 209 | ||
| 210 | let session: ModelSession | null = null; | |
| 211 | let reservation: string | null = null; | |
| 212 | let settled = false; | |
| 213 | // What the answer used, once there is one: counted however the reply ends. | |
| 214 | let usage: Partial<Outcome> = {}; | |
| 215 | ||
| 216 | const finish = async (outcome: Outcome) => { | |
| 217 | const charged = outcome.charged ?? 0; | |
| 218 | const statements = [ | |
| 219 | db | |
| 220 | .prepare( | |
| 221 | `UPDATE agent_replies SET status = ?, error = ?, reply_id = ?, model = ?, tier = ?, input_tokens = ?, output_tokens = ?, | |
| 222 | cost_micros = ?, charged_micros = ?, finished_at = ? WHERE id = ?`, | |
| 223 | ) | |
| 224 | .bind( | |
| 225 | outcome.status, | |
| 226 | outcome.error?.slice(0, 1000) ?? null, | |
| 227 | outcome.reply_id ?? null, | |
| 228 | outcome.model ?? null, | |
| 229 | outcome.tier ?? null, | |
| 230 | (outcome.tokens?.input ?? 0) + (outcome.tokens?.cacheRead ?? 0) + (outcome.tokens?.cacheWrite ?? 0), | |
| 231 | outcome.tokens?.output ?? 0, | |
| 232 | outcome.cost ?? 0, | |
| 233 | charged, | |
| 234 | new Date().toISOString(), | |
| 235 | id, | |
| 236 | ), | |
| 237 | ...(charged > 0 ? spendStatements(db, row.id, charged, now) : []), | |
| 238 | ]; | |
| 239 | await db.batch(statements); | |
| 240 | }; | |
| 241 | ||
| 242 | /** Says once, in this conversation, why the agent cannot answer; a repeat within hours is kept back. */ | |
| 243 | const notice = async (message: string, reason: string) => { | |
| 244 | const since = new Date(now.getTime() - NOTICE_QUIET_MS).toISOString(); | |
| 245 | const recent = await db | |
| 246 | .prepare( | |
| 247 | `SELECT 1 FROM agent_replies WHERE agent_id = ? AND channel_id = ? AND status = 'blocked' AND error = ? | |
| 248 | AND reply_id IS NOT NULL AND created_at > ? AND id <> ? LIMIT 1`, | |
| 249 | ) | |
| 250 | .bind(row.id, delivery.channel_id, reason, since, id) | |
| 251 | .first(); | |
| 252 | const posted = recent ? null : await surface.post(message).catch(() => null); | |
| 253 | await finish({ status: "blocked", error: reason, reply_id: posted }); | |
| 254 | }; | |
| 255 | ||
| 256 | try { | |
| 257 | 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"); | |
| 258 | const definition = definitionOf(row); | |
| 259 | const [month, day] = periods(now); | |
| 260 | const spentRows = await db | |
| 261 | .prepare("SELECT period, micros FROM agent_spend WHERE agent_id = ? AND period IN (?, ?)") | |
| 262 | .bind(row.id, month, day) | |
| 263 | .all<{ period: string; micros: number }>(); | |
| 264 | const spent: Spent = { | |
| 265 | month: spentRows.results.find((r) => r.period === month)?.micros ?? 0, | |
| 266 | day: spentRows.results.find((r) => r.period === day)?.micros ?? 0, | |
| 267 | }; | |
| 268 | ||
| 269 | // 1. The agent's own caps. | |
| 270 | const blocked = budgetBlock(definition.budget, spent, now); | |
| 271 | if (blocked) return await notice(blocked.message, `budget_${blocked.cap}`); | |
| 272 | ||
| 273 | // 2. Whether it may use a model at all. | |
| 274 | const [own, status] = await Promise.all([ | |
| 275 | integrations.modelProvider(slug).catch(() => null), | |
| 276 | billing.status().catch(() => ({ enabled: false, live: false })), | |
| 277 | ]); | |
| 278 | const allowed = allowedProviders(definition.routing, own?.id ?? null); | |
| 279 | const mayHosted = hostedOpen(slug, env.HOSTED_AGENT_WORKSPACES, status) && allowed.hosted; | |
| 280 | if (!mayHosted && !allowed.own) { | |
| 281 | const why = own | |
| 282 | ? "My settings don't let me use any model this workspace has. An owner can change my providers on my profile." | |
| 283 | : allowed.hosted | |
| 284 | ? "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." | |
| 285 | : "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."; | |
| 286 | return await notice(why, "no_model"); | |
| 287 | } | |
| 288 | ||
| 289 | // Read the conversation while showing that the agent is on it. | |
| 290 | const [, history] = await Promise.all([surface.typing(), surface.history(HISTORY_LIMIT)]); | |
| 291 | const conversation = turns(history, row.id); | |
| 292 | if (!conversation.length) return await finish({ status: "skipped", error: "nothing to answer" }); | |
| 293 | const author = askerIn(history, delivery); | |
| 294 | const askerName = delivery.asker?.username ?? author?.name ?? null; | |
| 295 | // v1 reads only this conversation, which its whole audience can read. | |
| 296 | // Tools, when replies get them, filter every result by this. | |
| 297 | void audienceFor(delivery); | |
| 298 | ||
| 299 | // 3. A model session, routed by the workspace's model routes. | |
| 300 | const repo = billingRepo(slug, row.handle); | |
| 301 | const policy = await routingNow(env.AGENT_ROUTING, () => billing.modelDefaults()); | |
| 302 | const provisional = replyModel(policy, { ...definition.routing, pinned: null }); | |
| 303 | const opened = await integrations.openModelSession({ | |
| 304 | workspace: slug, | |
| 305 | repo, | |
| 306 | number: 0, | |
| 307 | task: "reply", | |
| 308 | hostedOpen: mayHosted, | |
| 309 | tier: provisional.tier, | |
| 310 | requestedBy: askerName, | |
| 311 | }); | |
| 312 | if (!opened.ok) return await notice(`I can't reply right now: ${opened.error.message}`, "model_route"); | |
| 313 | session = opened.value; | |
| 314 | const ownModel = session.billedTo === "workspace"; | |
| 315 | if (ownModel ? !allowed.own : !mayHosted) { | |
| 316 | return await notice( | |
| 317 | "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.", | |
| 318 | "provider_not_allowed", | |
| 319 | ); | |
| 320 | } | |
| 321 | // A pinned model is for the workspace's own endpoints; on g1t's models the tier decides. | |
| 322 | const model = replyModel( | |
| 323 | policy, | |
| 324 | { ...definition.routing, pinned: ownModel ? definition.routing.pinned : null }, | |
| 325 | { chosen: session.tierChoice ?? null, named: ownModel ? session.model : null }, | |
| 326 | ); | |
| 327 | ||
| 328 | // 4. The workspace's own limits, through the compute gate. | |
| 329 | const ent = await gate.entitlements(slug); | |
| 330 | const admission = await gate.admit( | |
| 331 | { workspace: slug, repo, public: false, kind: "agent", estimateMicros: ownModel ? 0 : MODEL_ESTIMATE_MICROS.reply, hostedModel: !ownModel }, | |
| 332 | ent, | |
| 333 | ); | |
| 334 | if (!admission.ok) return await notice(`I can't reply right now: ${admission.message}`, `workspace_${admission.code}`); | |
| 335 | reservation = admission.reservation?.id ?? null; | |
| 336 | ||
| 337 | // 5. The run the reply is billed as. | |
| 338 | const named = ownModel && session.model ? session.model : null; | |
| 339 | const started = await billing.startRun({ | |
| 340 | workspace: slug, | |
| 341 | repo, | |
| 342 | number: 0, | |
| 343 | task: "reply", | |
| 344 | model: ownModel ? `${model.modelName} (${session.providerName ?? "own provider"})` : model.modelName, | |
| 345 | billedTo: ownModel ? "workspace" : "g1t", | |
| 346 | session: session.id, | |
| 347 | tier: named ? null : model.tier, | |
| 348 | }); | |
| 349 | if (!started.ok) return await notice(`I can't reply right now: ${started.error.message}`, "billing"); | |
| 350 | const ticket: RunTicket | null = started.value; | |
| 351 | const cap = replyCapMicros(definition.budget, spent, ent && ent.runCapMicros > 0 ? ent.runCapMicros : null); | |
| 352 | if (cap) await integrations.capModelSessions([await sha256Hex(session.token)], cap).catch(() => 0); | |
| 353 | ||
| 354 | // 6. The answer. | |
| 355 | const system = systemPrompt({ | |
| 356 | agent: { ...definition, id: row.id }, | |
| 357 | workspace: delivery.workspace, | |
| 358 | channel: { kind: delivery.channel_kind, name: delivery.channel_name }, | |
| 359 | asker: { name: askerName ?? "someone", display_name: author?.display_name ?? null, access: delivery.asker ?? null }, | |
| 360 | today: now, | |
| 361 | }); | |
| 362 | const answer = await askModel(env, session.token, { model: model.model, system, messages: conversation }); | |
| 363 | const cost = ownModel ? 0 : costMicros(answer.tokens, model.price); | |
| 364 | const priced = await priceTerms(env.BILLING); | |
| 365 | const charged = chargedMicros({ | |
| 366 | costMicros: cost, | |
| 367 | hosted: !ownModel, | |
| 368 | marginPercent: priced.marginPercent, | |
| 369 | ratePerMillionMicros: ownModel ? priced.rateOwn : priced.rate, | |
| 370 | tokens: totalTokens(answer.tokens), | |
| 371 | }); | |
| 372 | ||
| 373 | // 7. Bill it, whether or not the answer could be posted: the tokens were used. | |
| 374 | if (ticket) { | |
| 375 | await rpc(env.BILLING, "finish_run", { runId: ticket.runId, token: ticket.token, costUsd: cost / 1_000_000, turns: 1, tokens: answer.tokens }).catch( | |
| 376 | (error: unknown) => console.error("agents: finish_run failed", ticket.runId, String(error)), | |
| 377 | ); | |
| 378 | } | |
| 379 | if (reservation) { | |
| 380 | settled = true; | |
| 381 | await gate.settle(reservation, cost); | |
| 382 | } | |
| 383 | usage = { model: model.modelName, tier: named ? null : model.tier, tokens: answer.tokens, cost, charged }; | |
| 384 | if (!answer.text) { | |
| 385 | const posted = await surface.post(APOLOGY).catch(() => null); | |
| 386 | return await finish({ status: "failed", error: "the model gave no text", reply_id: posted, ...usage }); | |
| 387 | } | |
| 388 | const posted = await surface.post(answer.text); | |
| 389 | await finish({ status: "replied", reply_id: posted, ...usage }); | |
| 390 | } catch (error) { | |
| 391 | const message = error instanceof Error ? error.message : String(error); | |
| 392 | console.error("agents: a reply failed", row.id, delivery.message_id, message); | |
| 393 | const posted = await surface.post(APOLOGY).catch(() => null); | |
| 394 | await finish({ status: "failed", error: message, reply_id: posted, ...usage }).catch((failure: unknown) => | |
| 395 | console.error("agents: a failed reply was not recorded", id, String(failure)), | |
| 396 | ); | |
| 397 | } finally { | |
| 398 | // What was held is given back however the reply ended, and its model | |
| 399 | // session's token stops working. | |
| 400 | if (reservation && !settled) await gate.settle(reservation, 0); | |
| 401 | if (session) await integrations.closeModelSessions([await sha256Hex(session.token)]).catch(() => 0); | |
| 402 | } | |
| 403 | } |