Skip to content
384 linesCodeBlameRaw
1/**
2 * Metered model work: the one door every reply and every session step goes
3 * through, so g1t meters and bills an agent's model work one way
4 * (docs.g1t.sh/guides/agent-budgets/).
5 *
6 * 1. The paying agent's own monthly and daily caps (`budget.ts`), the
7 * workspace's budget for all its agents together (`policy.ts`), and the
8 * budget of the person who asked (`person-budget.ts`).
9 * 2. Whether the agent may use a model at all: g1t's hosted models are open
10 * as 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`), routed by the
14 * workspace's model routes, as for a run.
15 * 4. The compute gate's reservation (`ComputeGate.admit`, kind `agent`):
16 * the workspace's spend limit, AI credit, pauses and g1t's breaker.
17 * 5. A billing run (`start_run`), so the work is charged as Agent tokens on
18 * the workspace's bill, under the paying agent, attributed to it and to
19 * who asked: Spend's "by agent" and "by person" read the ledger.
20 * 6. The work itself, through the model proxy with the session's token.
21 * 7. `finish_run` with its cost and tokens, the reservation settled at
22 * cost, and the charge added to the paying agent's and the workspace's
23 * agent spend.
24 *
25 * Who does the work and who pays can differ: a colleague brought into a
26 * session, or a subagent, works on the budget of the agent at the root of
27 * the session's tree, so a chain never escapes the budget that started it.
28 */
29import { type EffortLevel, type ModelSession, type ModelTier, type RunTicket, ComputeGate, MODEL_ESTIMATE_MICROS, billingClient, integrationsClient } from "@g1t/contracts";
30
31import { hostedOpen } from "../../runner/src/hosted.ts";
32import { type AgentRouting as Policy, routingReader } from "../../runner/src/model-env.ts";
33import { type Spent, type Tokens, budgetBlock, chargedMicros, personBlock, personLimit, replyCapMicros, totalTokens } from "./budget.ts";
34import { monthStart, ownBudget, personSpent } from "./person-budget.ts";
35import { BUILTIN_NO_MODEL } from "./orchestrator.ts";
36import { type PolicyRow, alertDue, markAlerted, policyBlock, readPolicy, workspaceSpendStatements } from "./policy.ts";
37import { dollars } from "./money.ts";
38import { type TeamsHere, teamBudgetBlock, teamSpends } from "./teammates.ts";
39import { type ReplyModel, allowedProviders, reasoningEffort, replyModel } from "./routing.ts";
40import { type Row, definitionOf, periods, spendStatements } from "./store.ts";
41import type { ModelAnswer, Send } from "./turn.ts";
42import type { ServiceBinding } from "@g1t/contracts";
43
44export type MeterEnv = {
45 DB: D1Database;
46 BILLING: ServiceBinding;
47 INTEGRATIONS: ServiceBinding;
48 MODELS?: ServiceBinding;
49 MODELS_URL?: string;
50 HOSTED_AGENT_WORKSPACES: string;
51 AGENT_ROUTING: string;
52 /** Notifications: the workspace's agent budget crossing 75, 90 or 100%. */
53 NOTIFY?: ServiceBinding;
54};
55
56/** The longest one model answer may take. */
57const MODEL_TIMEOUT_MS = 120_000;
58
59/** Staff's model defaults on top of `AGENT_ROUTING`, read at most once a minute, as in the runner. */
60const routingNow = routingReader();
61let gate: ComputeGate | null = null;
62
63/**
64 * Where an agent's spend shows in billing: the workspace, under the agent
65 * that pays. Billing keys runs and reservations by a repository; no
66 * repository is named with an `@`, so the agent's line never mixes with a
67 * project's.
68 */
69export function billingRepo(workspace: string, handle: string): { namespace: string; name: string } {
70 return { namespace: workspace.toLowerCase(), name: `@${handle}` };
71}
72
73async function rpc<T>(service: ServiceBinding, method: string, args: object): Promise<T> {
74 const response = await service.fetch(`https://service/rpc/${method}`, {
75 method: "POST",
76 headers: { "content-type": "application/json" },
77 body: JSON.stringify(args),
78 });
79 if (!response.ok) throw new Error(`${method} failed with status ${response.status}`);
80 return (await response.json()) as T;
81}
82
83type PriceTerms = { marginPercent: number; rate: number; rateOwn: number };
84let terms: { value: PriceTerms; until: number } | null = null;
85
86/** The model margin and the agent rates from billing's price book, kept ten minutes. */
87async function priceTerms(billing: ServiceBinding): Promise<PriceTerms> {
88 if (terms && terms.until > Date.now()) return terms.value;
89 type Book = { prices?: { meter: string; priceMicros?: number; price_micros?: number }[]; modelMarginPercent?: number; model_margin_percent?: number };
90 const book = await rpc<Book>(billing, "prices", {}).catch(() => null);
91 const price = (meter: string) => {
92 const found = book?.prices?.find((p) => p.meter === meter);
93 return found?.priceMicros ?? found?.price_micros ?? 0;
94 };
95 const value = {
96 marginPercent: book?.modelMarginPercent ?? book?.model_margin_percent ?? 0,
97 rate: price("agent_tokens"),
98 rateOwn: price("agent_tokens_own"),
99 };
100 // A failed read is tried again in a minute, not kept.
101 terms = { value, until: Date.now() + (book ? 10 * 60_000 : 60_000) };
102 return value;
103}
104
105export async function sha256Hex(text: string): Promise<string> {
106 const digest = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(text));
107 return [...new Uint8Array(digest)].map((b) => b.toString(16).padStart(2, "0")).join("");
108}
109
110/**
111 * How work asks the model: one Messages API request through the model
112 * proxy, with the model session's token, non-streamed.
113 */
114function sendFor(env: MeterEnv, token: string): Send {
115 // The binding, not the proxy's public address: a Worker fetching another
116 // Worker's domain on the same zone can be refused or loop. The proxy reads
117 // only the path and the token, so the host does not matter.
118 const path = "/anthropic/v1/messages";
119 const fetcher = (url: string, init: RequestInit) => (env.MODELS ? env.MODELS.fetch(url, init) : fetch(url, init));
120 return async (body) => {
121 const response = await fetcher(env.MODELS ? `https://models${path}` : `${env.MODELS_URL!.replace(/\/+$/, "")}${path}`, {
122 method: "POST",
123 headers: { "content-type": "application/json", "x-api-key": token, "anthropic-version": "2023-06-01" },
124 body: JSON.stringify(body),
125 signal: AbortSignal.timeout(MODEL_TIMEOUT_MS),
126 });
127 const json = (await response.json().catch(() => null)) as (ModelAnswer & { error?: { message?: string } }) | null;
128 if (!response.ok || !json) throw new Error(`the model answered ${response.status}: ${json?.error?.message ?? "no answer"}`);
129 return json;
130 };
131}
132
133/** What the work is given: how to ask the model, and on what. */
134export type Model = {
135 send: Send;
136 model: ReplyModel;
137 /** The workspace's own provider pays for the model (g1t charges only the agent rate). */
138 ownModel: boolean;
139 /** The reasoning effort to send with each request, or null to leave it to the model. */
140 effort: EffortLevel | null;
141 policy: Policy;
142 /** The model the workspace's route names, on its own provider. */
143 sessionModel: string | null;
144 /** The workspace's own model connection, if any. */
145 own: string | null;
146};
147
148/** What the work used: its own tokens and cost, and any it spent for others (consults). */
149export type WorkUsage = { tokens: Tokens; cost: number; rounds: number };
150
151export type MeterInput = {
152 /** The agent doing the work. */
153 row: Row;
154 /** The agent whose budget pays; the same agent unless this is part of another's session. */
155 payer: Row;
156 /** The workspace's slug. */
157 slug: string;
158 task: "reply" | "session";
159 /** The tier the work starts on, before the agent's limits: from its effort plan (routing.ts `effortPlan`). */
160 start: ModelTier;
161 /** The reasoning level its effort plan asks for; sent where the model takes it. */
162 effort?: EffortLevel | null;
163 /** Who asked, by username, for billing's record. */
164 askerName: string | null;
165 /** The person the work is for, by username: their budget counts it. Null for work no person asked for. */
166 person?: string | null;
167 /** What is left of a session's cap, so one step never overruns it. */
168 leftMicros?: number | null;
169 /** Agent tier limits narrower than the agent's own (a subagent's). */
170 limits?: { floor: ModelTier | null; ceiling: ModelTier | null } | null;
171 /** The paying agent's teams (teammates.ts): a team's budget caps what its agents spend together. */
172 teams?: TeamsHere | null;
173};
174
175export type MeterBlock = { ok: false; reason: string; message: string };
176export type MeterDone<T> = { ok: true; value: T; model: string; tier: ModelTier | null; tokens: Tokens; cost: number; charged: number };
177
178/**
179 * Runs `work` as metered model work for `input.row`, paid by
180 * `input.payer`: checks every limit, opens and closes the model session,
181 * bills it, and records the spend. A limit that says no comes back as a
182 * block with what to tell people, before anything was spent. Whatever
183 * `work` throws is thrown on, after what was held is given back.
184 */
185export async function metered<T extends WorkUsage>(env: MeterEnv, input: MeterInput, work: (model: Model) => Promise<T>, now = new Date()): Promise<MeterBlock | MeterDone<T>> {
186 const db = env.DB;
187 const { row, payer, slug } = input;
188 if (!env.MODELS && !env.MODELS_URL) return { ok: false, reason: "no_models_url", message: "This installation has no model proxy set up." };
189 const billing = billingClient(env.BILLING);
190 const integrations = integrationsClient(env.INTEGRATIONS);
191 gate ??= new ComputeGate(env.BILLING);
192
193 // 1. The paying agent's caps, and the workspace's budget for every agent.
194 const payerDefinition = definitionOf(payer);
195 const [month, day] = periods(now);
196 const asker = input.person ? input.person.toLowerCase() : null;
197 const [spentRows, policy, ownCap] = await Promise.all([
198 db
199 .prepare("SELECT period, micros FROM agent_spend WHERE agent_id = ? AND period IN (?, ?)")
200 .bind(payer.id, month, day)
201 .all<{ period: string; micros: number }>(),
202 readPolicy(db, row.workspace_id, month),
203 asker ? ownBudget(db, row.workspace_id, asker) : Promise.resolve(null),
204 ]);
205 const spent: Spent = {
206 month: spentRows.results.find((r) => r.period === month)?.micros ?? 0,
207 day: spentRows.results.find((r) => r.period === day)?.micros ?? 0,
208 };
209 const blocked = budgetBlock(payerDefinition.budget, spent, now);
210 if (blocked) {
211 const message = payer.id === row.id ? blocked.message : `@${payer.handle}, who this work is for, is out of budget.`;
212 return { ok: false, reason: `budget_${blocked.cap}`, message };
213 }
214 const pool = policyBlock(policy);
215 if (pool) return { ok: false, reason: "workspace_agent_budget", message: pool };
216 // The budget of the person the work is for, summed only when one applies.
217 const personCap = asker ? personLimit(policy.person_monthly_micros, ownCap) : null;
218 const personUsed = asker && personCap != null ? await personSpent(db, row.workspace_id, asker, monthStart(now)) : 0;
219 const personStop = asker ? personBlock(asker, personCap, personUsed, now) : null;
220 if (personStop) return { ok: false, reason: "person_budget", message: personStop };
221 // The budgets of the teams the paying agent is on, read only when one has one.
222 const teamStop = input.teams?.teams.some((team) => team.budget_micros) ? teamBudgetBlock(await teamSpends(db, input.teams, month)) : null;
223 if (teamStop) return { ok: false, reason: "team_budget", message: teamStop };
224
225 // 2. Whether it may use a model at all.
226 const definition = definitionOf(row);
227 const [own, status] = await Promise.all([
228 integrations.modelProvider(slug).catch(() => null),
229 billing.status().catch(() => ({ enabled: false, live: false })),
230 ]);
231 const allowed = allowedProviders(definition.routing, own?.id ?? null);
232 const mayHosted = hostedOpen(slug, env.HOSTED_AGENT_WORKSPACES, status) && allowed.hosted;
233 if (!mayHosted && !allowed.own) {
234 const message = own
235 ? "My settings don't let me use any model this workspace has. An owner can change my providers on my profile."
236 : allowed.hosted
237 ? row.builtin
238 ? BUILTIN_NO_MODEL
239 : "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."
240 : "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.";
241 return { ok: false, reason: "no_model", message };
242 }
243
244 // 3. A model session, routed by the workspace's model routes, billed under the payer.
245 const repo = billingRepo(slug, payer.handle);
246 const routing = await routingNow(env.AGENT_ROUTING, () => billing.modelDefaults());
247 const limits = { ...definition.routing, ...(input.limits ?? {}) };
248 const provisional = replyModel(routing, { ...limits, pinned: null }, { start: input.start });
249 const opened = await integrations.openModelSession({
250 workspace: slug,
251 repo,
252 number: 0,
253 task: input.task,
254 hostedOpen: mayHosted,
255 tier: provisional.tier,
256 requestedBy: input.askerName,
257 });
258 if (!opened.ok) return { ok: false, reason: "model_route", message: opened.error.message };
259 const session: ModelSession = opened.value;
260 let reservation: string | null = null;
261 let settled = false;
262 try {
263 const ownModel = session.billedTo === "workspace";
264 if (ownModel ? !allowed.own : !mayHosted) {
265 return {
266 ok: false,
267 reason: "provider_not_allowed",
268 message:
269 "This workspace routes agents 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.",
270 };
271 }
272 // A pinned model is for the workspace's own endpoints; on g1t's models the tier decides.
273 const model = replyModel(
274 routing,
275 { ...limits, pinned: ownModel ? definition.routing.pinned : null },
276 { chosen: session.tierChoice ?? null, named: ownModel ? session.model : null, start: input.start },
277 );
278
279 // 4. The workspace's own limits, through the compute gate.
280 const ent = await gate.entitlements(slug);
281 const estimate = ownModel ? 0 : input.task === "reply" ? MODEL_ESTIMATE_MICROS.reply : MODEL_ESTIMATE_MICROS.reply * 4;
282 const admission = await gate.admit({ workspace: slug, repo, public: false, kind: "agent", estimateMicros: estimate, hostedModel: !ownModel }, ent);
283 if (!admission.ok) return { ok: false, reason: `workspace_${admission.code}`, message: admission.message };
284 reservation = admission.reservation?.id ?? null;
285
286 // 5. The run it is billed as.
287 const named = ownModel && session.model ? session.model : null;
288 const started = await billing.startRun({
289 workspace: slug,
290 repo,
291 number: 0,
292 task: input.task,
293 model: ownModel ? `${model.modelName} (${session.providerName ?? "own provider"})` : model.modelName,
294 billedTo: ownModel ? "workspace" : "g1t",
295 session: session.id,
296 tier: named ? null : model.tier,
297 // Whose work it is, on every line the run puts on the ledger: the
298 // agent that pays, and who asked. Spend reads both from the ledger.
299 agent: payer.handle,
300 askedBy: input.askerName,
301 });
302 if (!started.ok) return { ok: false, reason: "billing", message: started.error.message };
303 const ticket: RunTicket | null = started.value;
304 const caps = [replyCapMicros(payerDefinition.budget, spent, ent && ent.runCapMicros > 0 ? ent.runCapMicros : null)];
305 if (input.leftMicros != null) caps.push(Math.max(1, Math.floor(input.leftMicros)));
306 const left = policy.monthly_micros ? policy.monthly_micros - policy.spent : null;
307 if (left != null) caps.push(Math.max(1, left));
308 if (personCap != null) caps.push(Math.max(1, personCap - personUsed));
309 const cap = caps.filter((c): c is number => c != null);
310 if (cap.length) await integrations.capModelSessions([await sha256Hex(session.token)], Math.min(...cap)).catch(() => 0);
311
312 // 6. The work.
313 const effort = input.effort ? reasoningEffort(input.effort, routing.tiers[model.tier] ?? null, !!named || (ownModel && model.price == null)) : null;
314 const result = await work({ send: sendFor(env, session.token), model, ownModel, effort, policy: routing, sessionModel: session.model, own: own?.id ?? null });
315 const priced = await priceTerms(env.BILLING);
316 const charged = chargedMicros({
317 costMicros: result.cost,
318 hosted: !ownModel,
319 marginPercent: priced.marginPercent,
320 ratePerMillionMicros: ownModel ? priced.rateOwn : priced.rate,
321 tokens: totalTokens(result.tokens),
322 });
323
324 // 7. Bill it, settle, and count it against the payer and the workspace.
325 if (ticket) {
326 await rpc(env.BILLING, "finish_run", { runId: ticket.runId, token: ticket.token, costUsd: result.cost / 1_000_000, turns: result.rounds, tokens: result.tokens }).catch(
327 (error: unknown) => console.error("agents: finish_run failed", ticket.runId, String(error)),
328 );
329 }
330 if (reservation) {
331 settled = true;
332 await gate.settle(reservation, result.cost);
333 }
334 if (charged > 0) {
335 await db.batch([...spendStatements(db, payer.id, charged, now, input.task), ...workspaceSpendStatements(db, row.workspace_id, charged, now)]);
336 await budgetAlert(env, slug, row.workspace_id, { ...policy, spent: policy.spent + charged }, month).catch((error: unknown) =>
337 console.error("agents: a budget alert was not sent", slug, String(error)),
338 );
339 }
340 return { ok: true, value: result, model: model.modelName, tier: named ? null : model.tier, tokens: result.tokens, cost: result.cost, charged };
341 } finally {
342 // What was held is given back however the work ended, and the model
343 // session's token stops working.
344 if (reservation && !settled) await gate.settle(reservation, 0);
345 await integrations.closeModelSessions([await sha256Hex(session.token)]).catch(() => 0);
346 }
347}
348
349/**
350 * Tells the owner who set the workspace's agent budget, once per level a
351 * month, when every agent's spend together crosses 75, 90 or 100% of it.
352 */
353async function budgetAlert(env: MeterEnv, slug: string, workspaceId: string, policy: PolicyRow, month: string): Promise<void> {
354 const level = alertDue(policy);
355 if (!level || !env.NOTIFY || !policy.monthly_micros) return;
356 if (!(await markAlerted(env.DB, workspaceId, month, level))) return;
357 const setBy = await env.DB.prepare("SELECT updated_by FROM agent_policies WHERE workspace_id = ?").bind(workspaceId).first<{ updated_by: string | null }>();
358 if (!setBy?.updated_by) return;
359 const title = level >= 100 ? "Your agents have used this month's budget" : `Your agents have used ${level}% of this month's budget`;
360 const body =
361 level >= 100
362 ? `${dollars(policy.spent)} of ${dollars(policy.monthly_micros)}. They won't start new work until it's raised or the month turns.`
363 : `${dollars(policy.spent)} of ${dollars(policy.monthly_micros)} so far this month.`;
364 await env.NOTIFY.fetch("https://service/rpc/notify", {
365 method: "POST",
366 headers: { "content-type": "application/json" },
367 body: JSON.stringify({
368 target: { username: setBy.updated_by },
369 notification: {
370 id: `agent-budget:${workspaceId}:${month}:${level}`,
371 kind: "approval",
372 workspace: slug,
373 title,
374 body,
375 href: `/${slug}/-/agents`,
376 actor: { kind: "system", id: "g1t", name: "g1t" },
377 channel_id: null,
378 created_at: new Date().toISOString(),
379 },
380 }),
381 });
382}
383
384export type { PolicyRow };