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