Skip to content
403 linesCodeBlameRaw

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 to1/**
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 */
32import {
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
44import { hostedOpen } from "../../runner/src/hosted.ts";
45import { routingReader } from "../../runner/src/model-env.ts";
46import { type Spent, type Tokens, budgetBlock, chargedMicros, costMicros, replyCapMicros, totalTokens } from "./budget.ts";
47import { HISTORY_LIMIT, audienceFor, systemPrompt, turns } from "./prompt.ts";
48import { allowedProviders, replyModel } from "./routing.ts";
49import { type Row, definitionOf, periods, spendStatements } from "./store.ts";
50import { type SurfaceMessage, surfaceFor } from "./surface.ts";
51
52export 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. */
66const MODEL_TIMEOUT_MS = 90_000;
67/** A chat answer is short; this bounds the cost of one that is not. */
68const MAX_OUTPUT_TOKENS = 2048;
69/** A notice that the agent cannot reply is posted once per conversation in this long. */
70const NOTICE_QUIET_MS = 6 * 60 * 60 * 1000;
71
72const 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. */
75const routingNow = routingReader();
76let 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 */
83export function billingRepo(workspace: string, handle: string): { namespace: string; name: string } {
84 return { namespace: workspace.toLowerCase(), name: `@${handle}` };
85}
86
87async 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
97type PriceTerms = { marginPercent: number; rate: number; rateOwn: number };
98let terms: { value: PriceTerms; until: number } | null = null;
99
100/** The model margin and the agent rates from billing's price book, kept ten minutes. */
101async 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
119async 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
124type Answer = { text: string; tokens: Tokens };
125
126/** One non-streamed answer from the model proxy, in Anthropic's Messages format. */
127async 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). */
167function 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
173type 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 */
189export 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}