Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix
- models: every /v1/messages answer is measured as it passes (input, output, cache reads and writes, from the answer or its stream) and reported to billing after it is sent, with the model and who the run was for. A failed report never touches the answer. - integrations: a model session knows who its run was for (requested_by, 0005); the runner passes it from each kind of work. - billing: token_usage (0030), a row per day, workspace, person, session and model; record_tokens adds to it, token_usage answers a window (42 days by default) for the workspace or one person: totals, the runs' cost at g1t's price whoever paid, active days and each day. People see their own; owners anyone's. - site: a Usage panel beside This week, Workspace or You: tokens, cost, active days and cache share; a two-row grid of days shaded in one hue by quartile of the busiest; and the token mix as one bar, each part labelled, in colours checked for every kind of colour vision. - docs: Token usage in BILLING_OPERATIONS; the models guide says tokens are counted per run for usage views. Billing still prices runs from AI Gateway.
| 112 | 112 | work is routed to, translates if the provider speaks OpenAI's API, and | |
| 113 | 113 | forwards the request. Answers stream straight back. | |
| 114 | 114 | ||
| 115 | + | As each answer passes, the proxy reads how many tokens it used (input, | |
| 116 | + | output, and cache reads and writes) and counts them for the run, under the | |
| 117 | + | person it was for. Those counts are for usage views; they never change what | |
| 118 | + | a run is charged. | |
| 119 | + | ||
| 115 | 120 | The token stops working within seconds of the run finishing, however it | |
| 116 | 121 | ends, and within seconds if you disconnect the provider. A run whose end | |
| 117 | 122 | g1t never hears about loses it three hours after it starts. Keys are sealed when you save them, and used |
| 32 | 32 | import type { ShellData } from "./shell"; | |
| 33 | 33 | import { useLiveRefresh } from "./agents"; | |
| 34 | 34 | import { Unavailable } from "./mission"; | |
| 35 | + | import { TokenUsagePanel } from "./token-usage"; | |
| 35 | 36 | import { Avatar, TimeAgo } from "./ui"; | |
| 36 | 37 | import { DropdownMenu, DropdownMenuContent, DropdownMenuItem, DropdownMenuLabel, DropdownMenuTrigger } from "./ui/dropdown-menu"; | |
| 37 | 38 | import type { Loaded } from "../routes/home"; | |
| 1093 | 1094 | )} | |
| 1094 | 1095 | </section> | |
| 1095 | 1096 | ||
| 1097 | + | {workspace && ( | |
| 1098 | + | <TokenUsagePanel | |
| 1099 | + | workspace={loaderData.tokens.workspace} | |
| 1100 | + | mine={loaderData.tokens.mine} | |
| 1101 | + | usageHref={`/${workspace}/-/usage`} | |
| 1102 | + | /> | |
| 1103 | + | )} | |
| 1104 | + | ||
| 1096 | 1105 | <section id="activity" className="scroll-mt-20 rounded-xl border border-line bg-surface p-5"> | |
| 1097 | 1106 | <div className="flex items-center justify-between"> | |
| 1098 | 1107 | <h2 className="text-base font-semibold tracking-tight">Activity</h2> |
| 1 | + | import { useState } from "react"; | |
| 2 | + | import { Link } from "react-router"; | |
| 3 | + | ||
| 4 | + | import { cn } from "../lib/cn"; | |
| 5 | + | import { usd } from "../lib/mission-control"; | |
| 6 | + | import { | |
| 7 | + | HEAT_LEVELS, | |
| 8 | + | type MixPart, | |
| 9 | + | type TokenUsageView, | |
| 10 | + | bestDay, | |
| 11 | + | cacheShare, | |
| 12 | + | compactTokens, | |
| 13 | + | heatLevel, | |
| 14 | + | shortDay, | |
| 15 | + | tokenMix, | |
| 16 | + | } from "../lib/token-usage"; | |
| 17 | + | ||
| 18 | + | /** | |
| 19 | + | * The token mix's colours, in its stacking order. Checked against the dark | |
| 20 | + | * surface (lightness band, chroma, contrast, and neighbours apart for every | |
| 21 | + | * kind of colour vision); change them only by re-running that check. | |
| 22 | + | */ | |
| 23 | + | const MIX_COLOR: Record<MixPart["key"], string> = { | |
| 24 | + | input: "#9085e9", | |
| 25 | + | cacheRead: "#199e70", | |
| 26 | + | cacheWrite: "#3987e5", | |
| 27 | + | output: "#d55181", | |
| 28 | + | }; | |
| 29 | + | ||
| 30 | + | /** The daily grid's single hue, stronger with more tokens. */ | |
| 31 | + | const HEAT_SHADE = ["var(--color-line)", ...Array.from({ length: HEAT_LEVELS }, (_, i) => `color-mix(in oklab, var(--color-merged) ${[28, 48, 72, 100][i]}%, var(--color-surface))`)]; | |
| 32 | + | ||
| 33 | + | type Scope = "workspace" | "mine"; | |
| 34 | + | ||
| 35 | + | /** | |
| 36 | + | * Model tokens over the last weeks, for the workspace or for you: what was | |
| 37 | + | * used and what it cost, how many days, how much came from the cache, each | |
| 38 | + | * day's intensity, and the mix of input, cache and output. | |
| 39 | + | */ | |
| 40 | + | export function TokenUsagePanel({ | |
| 41 | + | workspace, | |
| 42 | + | mine, | |
| 43 | + | usageHref, | |
| 44 | + | }: { | |
| 45 | + | workspace: TokenUsageView | null; | |
| 46 | + | mine: TokenUsageView | null; | |
| 47 | + | usageHref: string | null; | |
| 48 | + | }) { | |
| 49 | + | const [scope, setScope] = useState<Scope>("workspace"); | |
| 50 | + | const usage = scope === "mine" ? mine : workspace; | |
| 51 | + | return ( | |
| 52 | + | <section aria-labelledby="token-usage" className="rounded-xl border border-line bg-surface p-5"> | |
| 53 | + | <div className="flex items-center justify-between gap-3"> | |
| 54 | + | <h2 id="token-usage" className="text-base font-semibold tracking-tight"> | |
| 55 | + | Usage | |
| 56 | + | </h2> | |
| 57 | + | <div role="tablist" aria-label="Whose usage" className="flex rounded-md border border-line p-0.5 text-xs"> | |
| 58 | + | {(["workspace", "mine"] as const).map((option) => ( | |
| 59 | + | <button | |
| 60 | + | key={option} | |
| 61 | + | type="button" | |
| 62 | + | role="tab" | |
| 63 | + | aria-selected={scope === option} | |
| 64 | + | disabled={option === "mine" && !mine} | |
| 65 | + | onClick={() => setScope(option)} | |
| 66 | + | className={cn( | |
| 67 | + | "rounded px-2 py-0.5 transition-colors disabled:opacity-40", | |
| 68 | + | scope === option ? "bg-raised text-fg" : "text-muted hover:text-fg", | |
| 69 | + | )} | |
| 70 | + | > | |
| 71 | + | {option === "workspace" ? "Workspace" : "You"} | |
| 72 | + | </button> | |
| 73 | + | ))} | |
| 74 | + | </div> | |
| 75 | + | </div> | |
| 76 | + | {usage ? <Figures usage={usage} usageHref={usageHref} /> : <p className="mt-4 text-sm text-muted">Usage could not be loaded.</p>} | |
| 77 | + | </section> | |
| 78 | + | ); | |
| 79 | + | } | |
| 80 | + | ||
| 81 | + | function Figures({ usage, usageHref }: { usage: TokenUsageView; usageHref: string | null }) { | |
| 82 | + | const share = cacheShare(usage); | |
| 83 | + | const weeks = Math.round(usage.days / 7); | |
| 84 | + | if (usage.totalTokens === 0) { | |
| 85 | + | return ( | |
| 86 | + | <p className="mt-4 text-sm text-muted"> | |
| 87 | + | No model tokens in the last {weeks} weeks{usage.person ? " on work you asked for" : ""}. Agents' runs show here as they work. | |
| 88 | + | </p> | |
| 89 | + | ); | |
| 90 | + | } | |
| 91 | + | return ( | |
| 92 | + | <> | |
| 93 | + | <dl className="mt-4 grid grid-cols-2 gap-2"> | |
| 94 | + | <Tile label="Tokens" value={compactTokens(usage.totalTokens)} /> | |
| 95 | + | <Tile label="Cost" value={usd(usage.costMicros / 1e6)} title="What these runs were charged" /> | |
| 96 | + | <Tile label="Active days" value={String(usage.activeDays)} hint={`of ${usage.days}`} /> | |
| 97 | + | <Tile label="From cache" value={share == null ? "—" : `${Math.round(share * 100)}%`} title="Prompt tokens read from the model's cache" /> | |
| 98 | + | </dl> | |
| 99 | + | <DailyGrid byDay={usage.byDay} /> | |
| 100 | + | <Mix usage={usage} /> | |
| 101 | + | {usageHref && ( | |
| 102 | + | <Link to={usageHref} className="mt-4 block border-t border-line pt-3 text-xs text-muted hover:text-fg"> | |
| 103 | + | Every run, by repository and model | |
| 104 | + | </Link> | |
| 105 | + | )} | |
| 106 | + | </> | |
| 107 | + | ); | |
| 108 | + | } | |
| 109 | + | ||
| 110 | + | function Tile({ label, value, hint, title }: { label: string; value: string; hint?: string; title?: string }) { | |
| 111 | + | return ( | |
| 112 | + | <div className="rounded-lg border border-line bg-bg/40 px-3 py-2.5" title={title}> | |
| 113 | + | <dt className="text-[0.6875rem] text-muted">{label}</dt> | |
| 114 | + | <dd className="mt-0.5 text-lg font-semibold tracking-tight tabular-nums"> | |
| 115 | + | {value} | |
| 116 | + | {hint && <span className="ml-1 text-xs font-normal text-faint">{hint}</span>} | |
| 117 | + | </dd> | |
| 118 | + | </div> | |
| 119 | + | ); | |
| 120 | + | } | |
| 121 | + | ||
| 122 | + | /** | |
| 123 | + | * Each day as a cell, two rows, oldest first; shaded by quartile of the | |
| 124 | + | * busiest day. Hover or focus a day to read it. | |
| 125 | + | */ | |
| 126 | + | function DailyGrid({ byDay }: { byDay: TokenUsageView["byDay"] }) { | |
| 127 | + | const best = bestDay(byDay); | |
| 128 | + | const busiest = best?.tokens ?? 0; | |
| 129 | + | const [shown, setShown] = useState<{ day: string; tokens: number } | null>(null); | |
| 130 | + | const columns = Math.ceil(byDay.length / 2); | |
| 131 | + | const readout = shown ?? best; | |
| 132 | + | return ( | |
| 133 | + | <div className="mt-5"> | |
| 134 | + | <div className="flex items-baseline justify-between gap-2"> | |
| 135 | + | <h3 className="text-xs font-medium text-fg-soft">Each day</h3> | |
| 136 | + | <p className="text-[0.6875rem] text-muted tabular-nums" aria-live="polite"> | |
| 137 | + | {readout ? `${shown ? "" : "Busiest: "}${shortDay(readout.day)} · ${compactTokens(readout.tokens)}` : ""} | |
| 138 | + | </p> | |
| 139 | + | </div> | |
| 140 | + | <div | |
| 141 | + | className="mt-2 grid gap-[3px]" | |
| 142 | + | style={{ gridTemplateColumns: `repeat(${columns}, minmax(0, 1fr))` }} | |
| 143 | + | onMouseLeave={() => setShown(null)} | |
| 144 | + | > | |
| 145 | + | {byDay.map((entry) => ( | |
| 146 | + | <span | |
| 147 | + | key={entry.day} | |
| 148 | + | role="img" | |
| 149 | + | tabIndex={0} | |
| 150 | + | aria-label={`${shortDay(entry.day)}: ${compactTokens(entry.tokens)} tokens`} | |
| 151 | + | onMouseEnter={() => setShown(entry)} | |
| 152 | + | onFocus={() => setShown(entry)} | |
| 153 | + | onBlur={() => setShown(null)} | |
| 154 | + | className="aspect-square rounded-[3px] outline-offset-1 hover:ring-1 hover:ring-fg-soft focus-visible:outline-2 focus-visible:outline-accent" | |
| 155 | + | style={{ background: HEAT_SHADE[heatLevel(entry.tokens, busiest)] }} | |
| 156 | + | /> | |
| 157 | + | ))} | |
| 158 | + | </div> | |
| 159 | + | <div className="mt-1.5 flex items-center justify-between text-[0.6875rem] text-faint tabular-nums"> | |
| 160 | + | <span>{byDay[0] ? shortDay(byDay[0].day) : ""}</span> | |
| 161 | + | <span className="flex items-center gap-1" aria-hidden> | |
| 162 | + | Less | |
| 163 | + | {HEAT_SHADE.map((shade) => ( | |
| 164 | + | <span key={shade} className="size-2 rounded-[2px]" style={{ background: shade }} /> | |
| 165 | + | ))} | |
| 166 | + | More | |
| 167 | + | </span> | |
| 168 | + | <span>{byDay.length ? shortDay(byDay[byDay.length - 1].day) : ""}</span> | |
| 169 | + | </div> | |
| 170 | + | </div> | |
| 171 | + | ); | |
| 172 | + | } | |
| 173 | + | ||
| 174 | + | /** Input, cache reads, cache writes and output as one bar, each part labelled. */ | |
| 175 | + | function Mix({ usage }: { usage: TokenUsageView }) { | |
| 176 | + | const parts = tokenMix(usage).filter((part) => part.tokens > 0); | |
| 177 | + | return ( | |
| 178 | + | <div className="mt-5"> | |
| 179 | + | <h3 className="text-xs font-medium text-fg-soft">Token mix</h3> | |
| 180 | + | <div className="mt-2 flex h-2.5 gap-[2px] overflow-hidden rounded" role="img" aria-label={parts.map((p) => `${p.label} ${Math.round(p.share * 100)}%`).join(", ")}> | |
| 181 | + | {parts.map((part) => ( | |
| 182 | + | <span | |
| 183 | + | key={part.key} | |
| 184 | + | title={`${part.label}: ${compactTokens(part.tokens)} (${Math.round(part.share * 100)}%)`} | |
| 185 | + | className="h-full first:rounded-l last:rounded-r" | |
| 186 | + | style={{ width: `${part.share * 100}%`, minWidth: 3, background: MIX_COLOR[part.key] }} | |
| 187 | + | /> | |
| 188 | + | ))} | |
| 189 | + | </div> | |
| 190 | + | <ul className="mt-2.5 grid grid-cols-2 gap-x-3 gap-y-1.5 text-[0.6875rem]"> | |
| 191 | + | {tokenMix(usage).map((part) => ( | |
| 192 | + | <li key={part.key} className="flex min-w-0 items-center gap-1.5"> | |
| 193 | + | <span className="size-2 shrink-0 rounded-[2px]" style={{ background: MIX_COLOR[part.key] }} /> | |
| 194 | + | <span className="truncate text-muted">{part.label}</span> | |
| 195 | + | <span className="ml-auto text-fg-soft tabular-nums"> | |
| 196 | + | {compactTokens(part.tokens)} <span className="text-faint">{Math.round(part.share * 100)}%</span> | |
| 197 | + | </span> | |
| 198 | + | </li> | |
| 199 | + | ))} | |
| 200 | + | </ul> | |
| 201 | + | </div> | |
| 202 | + | ); | |
| 203 | + | } |
| 1 | + | import assert from "node:assert/strict"; | |
| 2 | + | import { test } from "node:test"; | |
| 3 | + | ||
| 4 | + | import { bestDay, cacheShare, compactTokens, heatLevel, shortDay, tokenMix, type TokenUsageView } from "./token-usage.ts"; | |
| 5 | + | ||
| 6 | + | const usage = (over: Partial<TokenUsageView> = {}): TokenUsageView => ({ | |
| 7 | + | since: "2026-08-26", | |
| 8 | + | days: 42, | |
| 9 | + | person: null, | |
| 10 | + | totalTokens: 0, | |
| 11 | + | inputTokens: 0, | |
| 12 | + | outputTokens: 0, | |
| 13 | + | cacheReadTokens: 0, | |
| 14 | + | cacheWriteTokens: 0, | |
| 15 | + | costMicros: 0, | |
| 16 | + | activeDays: 0, | |
| 17 | + | byDay: [], | |
| 18 | + | ...over, | |
| 19 | + | }); | |
| 20 | + | ||
| 21 | + | test("token counts read short", () => { | |
| 22 | + | assert.equal(compactTokens(8_500_000_000), "8.5B"); | |
| 23 | + | assert.equal(compactTokens(412_300_000), "412M"); | |
| 24 | + | assert.equal(compactTokens(12_340), "12.3K"); | |
| 25 | + | assert.equal(compactTokens(2_000_000), "2M"); | |
| 26 | + | assert.equal(compactTokens(940), "940"); | |
| 27 | + | assert.equal(compactTokens(0), "0"); | |
| 28 | + | }); | |
| 29 | + | ||
| 30 | + | test("cache share is reads over every prompt token", () => { | |
| 31 | + | assert.equal(cacheShare(usage({ inputTokens: 10, cacheReadTokens: 80, cacheWriteTokens: 10 })), 0.8); | |
| 32 | + | assert.equal(cacheShare(usage()), null); | |
| 33 | + | }); | |
| 34 | + | ||
| 35 | + | test("the grid shades by quartile of the busiest day", () => { | |
| 36 | + | assert.equal(heatLevel(0, 100), 0); | |
| 37 | + | assert.equal(heatLevel(1, 100), 1); | |
| 38 | + | assert.equal(heatLevel(26, 100), 2); | |
| 39 | + | assert.equal(heatLevel(100, 100), 4); | |
| 40 | + | assert.equal(heatLevel(5, 0), 0); | |
| 41 | + | }); | |
| 42 | + | ||
| 43 | + | test("the best day, and none on an empty window", () => { | |
| 44 | + | assert.deepEqual(bestDay([{ day: "a", tokens: 3 }, { day: "b", tokens: 9 }, { day: "c", tokens: 9 }]), { day: "b", tokens: 9 }); | |
| 45 | + | assert.equal(bestDay([{ day: "a", tokens: 0 }]), null); | |
| 46 | + | }); | |
| 47 | + | ||
| 48 | + | test("the mix is in its checked order and adds up to one", () => { | |
| 49 | + | const mix = tokenMix(usage({ inputTokens: 1, cacheReadTokens: 6, cacheWriteTokens: 1, outputTokens: 2 })); | |
| 50 | + | assert.deepEqual(mix.map((part) => part.key), ["input", "cacheRead", "cacheWrite", "output"]); | |
| 51 | + | assert.equal(mix.reduce((sum, part) => sum + part.share, 0), 1); | |
| 52 | + | assert.ok(tokenMix(usage()).every((part) => part.share === 0)); | |
| 53 | + | }); | |
| 54 | + | ||
| 55 | + | test("days are shown short, in UTC", () => { | |
| 56 | + | assert.equal(shortDay("2026-09-22"), "22 Sept"); | |
| 57 | + | assert.equal(shortDay("nope"), "nope"); | |
| 58 | + | }); |
| 1 | + | /** | |
| 2 | + | * Model tokens for mission control's usage panel: a person's or the | |
| 3 | + | * workspace's, over the last weeks (billing's `token_usage`). | |
| 4 | + | */ | |
| 5 | + | ||
| 6 | + | import type { TokenUsage } from "@g1t/contracts"; | |
| 7 | + | ||
| 8 | + | /** What billing answers. */ | |
| 9 | + | export type TokenUsageView = TokenUsage; | |
| 10 | + | ||
| 11 | + | /** "8.5B", "412M", "12.3K", "940". */ | |
| 12 | + | export function compactTokens(n: number): string { | |
| 13 | + | const units: [number, string][] = [ | |
| 14 | + | [1e12, "T"], | |
| 15 | + | [1e9, "B"], | |
| 16 | + | [1e6, "M"], | |
| 17 | + | [1e3, "K"], | |
| 18 | + | ]; | |
| 19 | + | for (const [size, unit] of units) { | |
| 20 | + | if (n >= size) { | |
| 21 | + | const value = n / size; | |
| 22 | + | return `${value >= 100 ? Math.round(value) : Number(value.toFixed(1))}${unit}`; | |
| 23 | + | } | |
| 24 | + | } | |
| 25 | + | return String(Math.max(0, Math.round(n))); | |
| 26 | + | } | |
| 27 | + | ||
| 28 | + | /** The share of prompt tokens read from the provider's cache, 0..1; null with no prompt. */ | |
| 29 | + | export function cacheShare(usage: Pick<TokenUsageView, "inputTokens" | "cacheReadTokens" | "cacheWriteTokens">): number | null { | |
| 30 | + | const prompt = usage.inputTokens + usage.cacheReadTokens + usage.cacheWriteTokens; | |
| 31 | + | return prompt > 0 ? usage.cacheReadTokens / prompt : null; | |
| 32 | + | } | |
| 33 | + | ||
| 34 | + | /** Shading levels for the daily grid: 0 for none, then 1..4 by quartile of the busiest day. */ | |
| 35 | + | export const HEAT_LEVELS = 4; | |
| 36 | + | ||
| 37 | + | export function heatLevel(tokens: number, busiest: number): number { | |
| 38 | + | if (tokens <= 0 || busiest <= 0) return 0; | |
| 39 | + | return Math.min(HEAT_LEVELS, Math.max(1, Math.ceil((tokens / busiest) * HEAT_LEVELS))); | |
| 40 | + | } | |
| 41 | + | ||
| 42 | + | /** The busiest day, or null when none had tokens. */ | |
| 43 | + | export function bestDay(byDay: TokenUsageView["byDay"]): { day: string; tokens: number } | null { | |
| 44 | + | let best: { day: string; tokens: number } | null = null; | |
| 45 | + | for (const entry of byDay) if (entry.tokens > 0 && (!best || entry.tokens > best.tokens)) best = entry; | |
| 46 | + | return best; | |
| 47 | + | } | |
| 48 | + | ||
| 49 | + | /** One part of the token mix, in stacking order. */ | |
| 50 | + | export type MixPart = { key: "input" | "cacheRead" | "cacheWrite" | "output"; label: string; tokens: number; share: number }; | |
| 51 | + | ||
| 52 | + | /** | |
| 53 | + | * New input, cache reads, cache writes and output, in that order: the | |
| 54 | + | * order the colours were checked in (neighbours stay apart for every kind | |
| 55 | + | * of colour vision). | |
| 56 | + | */ | |
| 57 | + | export function tokenMix(usage: TokenUsageView): MixPart[] { | |
| 58 | + | const parts: Omit<MixPart, "share">[] = [ | |
| 59 | + | { key: "input", label: "New input", tokens: usage.inputTokens }, | |
| 60 | + | { key: "cacheRead", label: "Cache reads", tokens: usage.cacheReadTokens }, | |
| 61 | + | { key: "cacheWrite", label: "Cache writes", tokens: usage.cacheWriteTokens }, | |
| 62 | + | { key: "output", label: "Output", tokens: usage.outputTokens }, | |
| 63 | + | ]; | |
| 64 | + | const total = parts.reduce((sum, part) => sum + part.tokens, 0); | |
| 65 | + | return parts.map((part) => ({ ...part, share: total > 0 ? part.tokens / total : 0 })); | |
| 66 | + | } | |
| 67 | + | ||
| 68 | + | /** "22 Sept"-style short day, read in UTC since days are UTC dates. */ | |
| 69 | + | export function shortDay(day: string): string { | |
| 70 | + | const date = new Date(`${day}T00:00:00Z`); | |
| 71 | + | return Number.isNaN(date.getTime()) | |
| 72 | + | ? day | |
| 73 | + | : date.toLocaleDateString("en-GB", { day: "numeric", month: "short", timeZone: "UTC" }); | |
| 74 | + | } |
| 202 | 202 | ); | |
| 203 | 203 | }); | |
| 204 | 204 | ||
| 205 | − | const [repos, perRepo, active, models, profile, runs, overview, usage, projectList, memories, invitations] = await Promise.all([ | |
| 205 | + | const [repos, perRepo, active, models, profile, runs, overview, usage, projectList, memories, invitations, tokens, myTokens] = await Promise.all([ | |
| 206 | 206 | reposP, | |
| 207 | 207 | soft("projects", perRepoP), | |
| 208 | 208 | soft("pulls", work.listActivePulls(viewer)), | |
| 215 | 215 | slug ? soft("memories", agents.listMemories(viewer, slug, null)) : null, | |
| 216 | 216 | // Repositories someone has invited the viewer to. | |
| 217 | 217 | soft("invitations", identity.myRepoInvitations(viewer)), | |
| 218 | + | // Model tokens over the last six weeks: the workspace's, and yours. | |
| 219 | + | slug ? soft("tokens", billing.tokenUsage(slug, viewer)) : null, | |
| 220 | + | slug ? soft("my_tokens", billing.tokenUsage(slug, viewer, { person: username })) : null, | |
| 218 | 221 | ]); | |
| 219 | 222 | ||
| 220 | 223 | const repoList = repos ?? []; | |
| 598 | 601 | weekCost, | |
| 599 | 602 | }, | |
| 600 | 603 | week, | |
| 604 | + | tokens: { workspace: okOr(tokens), mine: okOr(myTokens) }, | |
| 601 | 605 | groups, | |
| 602 | 606 | titles: shownTitles, | |
| 603 | 607 | // A run that is going makes the page worth refreshing on its own. |
| 350 | 350 | pub added_micros: i64, | |
| 351 | 351 | } | |
| 352 | 352 | ||
| 353 | + | /// `record_tokens`: what one model answer used, added to the day's count | |
| 354 | + | /// for its run. The model proxy sends it after each answer. For usage | |
| 355 | + | /// views only: runs are still priced from AI Gateway. Returns | |
| 356 | + | /// `Outcome<bool>`: false when there was nothing to count. | |
| 357 | + | #[derive(Debug, Serialize, Deserialize)] | |
| 358 | + | #[serde(rename_all = "camelCase")] | |
| 359 | + | pub struct RecordTokensArgs { | |
| 360 | + | pub workspace: String, | |
| 361 | + | /// The model session's id (`ModelSession::id`), one per run. | |
| 362 | + | pub session: String, | |
| 363 | + | /// The person the run is for, by username. Absent when nobody asked. | |
| 364 | + | #[serde(default)] | |
| 365 | + | pub person: Option<String>, | |
| 366 | + | pub model: String, | |
| 367 | + | /// On g1t's hosted models: `small` or `large`. | |
| 368 | + | #[serde(default)] | |
| 369 | + | pub tier: Option<String>, | |
| 370 | + | #[serde(default)] | |
| 371 | + | pub input: u64, | |
| 372 | + | #[serde(default)] | |
| 373 | + | pub output: u64, | |
| 374 | + | #[serde(default)] | |
| 375 | + | pub cache_read: u64, | |
| 376 | + | #[serde(default)] | |
| 377 | + | pub cache_write: u64, | |
| 378 | + | } | |
| 379 | + | ||
| 380 | + | /// `token_usage`: the model tokens a workspace's runs used, day by day, | |
| 381 | + | /// for the whole workspace or for one person. Members only; a member may | |
| 382 | + | /// ask only for themselves, an owner for anyone. Returns | |
| 383 | + | /// `Outcome<TokenUsage>`. | |
| 384 | + | #[derive(Debug, Serialize, Deserialize)] | |
| 385 | + | pub struct TokenUsageArgs { | |
| 386 | + | pub workspace: String, | |
| 387 | + | pub viewer: Viewer, | |
| 388 | + | /// A username: only the runs for them. | |
| 389 | + | #[serde(default)] | |
| 390 | + | pub person: Option<String>, | |
| 391 | + | /// How many days, to today: 42 when absent, 366 at most. | |
| 392 | + | #[serde(default)] | |
| 393 | + | pub days: Option<u32>, | |
| 394 | + | } | |
| 395 | + | ||
| 396 | + | /// One day's tokens. | |
| 397 | + | #[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] | |
| 398 | + | pub struct DayTokens { | |
| 399 | + | /// `YYYY-MM-DD`, UTC. | |
| 400 | + | pub day: String, | |
| 401 | + | pub tokens: u64, | |
| 402 | + | } | |
| 403 | + | ||
| 404 | + | /// The model tokens runs used over a window of days. | |
| 405 | + | #[derive(Clone, Debug, Serialize, Deserialize)] | |
| 406 | + | #[serde(rename_all = "camelCase")] | |
| 407 | + | pub struct TokenUsage { | |
| 408 | + | /// `YYYY-MM-DD`: the first day counted. | |
| 409 | + | pub since: String, | |
| 410 | + | pub days: u32, | |
| 411 | + | /// Null for the whole workspace. | |
| 412 | + | pub person: Option<String>, | |
| 413 | + | pub total_tokens: u64, | |
| 414 | + | pub input_tokens: u64, | |
| 415 | + | pub output_tokens: u64, | |
| 416 | + | pub cache_read_tokens: u64, | |
| 417 | + | pub cache_write_tokens: u64, | |
| 418 | + | /// What those runs were charged, as `usage` measures it. | |
| 419 | + | pub cost_micros: i64, | |
| 420 | + | /// Days in the window with any tokens. | |
| 421 | + | pub active_days: u32, | |
| 422 | + | /// Every day in the window, oldest first, zeros included. | |
| 423 | + | pub by_day: Vec<DayTokens>, | |
| 424 | + | } | |
| 425 | + | ||
| 353 | 426 | /// What a workspace pays a monthly price for. There is one plan, `plan` | |
| 354 | 427 | /// ("g1t"): a flat price per workspace, never per person, with included | |
| 355 | 428 | /// usage each month, more private storage, and deployments. Never free: |
| 387 | 387 | /// For `g1t`: the tier the run was routed to, `small` or `large`. | |
| 388 | 388 | #[serde(default)] | |
| 389 | 389 | pub tier: Option<String>, | |
| 390 | + | /// The person the run is for, by username: who asked g1t for the work. | |
| 391 | + | /// Null when nobody did. Never `g1t`, the agent itself. | |
| 392 | + | #[serde(default)] | |
| 393 | + | pub requested_by: Option<String>, | |
| 390 | 394 | /// For `endpoint`: where to send requests. | |
| 391 | 395 | pub base_url: Option<String>, | |
| 392 | 396 | /// For `anthropic` and `endpoint`: the workspace's key. | |
| 566 | 570 | /// when the run goes to g1t's models. | |
| 567 | 571 | #[serde(default)] | |
| 568 | 572 | pub tier: Option<String>, | |
| 573 | + | /// The person the run is for, by username, so usage can be shown per | |
| 574 | + | /// person. Null when nobody asked; `g1t`, the agent, is kept as null. | |
| 575 | + | #[serde(default)] | |
| 576 | + | pub requested_by: Option<String>, | |
| 569 | 577 | } | |
| 570 | 578 | ||
| 571 | 579 | /// `routes`: a workspace's model routes, one per kind of work that has its |
| 200 | 200 | emailed again weekly. The red bar on every sudo page shows margin, | |
| 201 | 201 | overall and leak alerts. | |
| 202 | 202 | ||
| 203 | + | ## Token usage | |
| 204 | + | ||
| 205 | + | The model proxy (`services/models`) reads Anthropic's `usage` from every | |
| 206 | + | `/v1/messages` answer, streamed or whole, on g1t's models and on a | |
| 207 | + | workspace's own provider alike (OpenAI-shaped providers are translated | |
| 208 | + | first). Count-tokens requests are not answers and are skipped. After the | |
| 209 | + | answer, it calls `record_tokens`, which adds input, output, cache reads and | |
| 210 | + | cache writes to one row per day, workspace, person, session and model in | |
| 211 | + | `token_usage` (migration `0030_token_usage.sql`). The person is who the run | |
| 212 | + | was for, from the model session's `requested_by`; never g1t's agent. A | |
| 213 | + | report that fails is dropped and never affects the answer. | |
| 214 | + | ||
| 215 | + | `token_usage` reads a window (42 days by default, 366 at most) for the | |
| 216 | + | workspace or one person: totals, every day's tokens and the active days, | |
| 217 | + | with `costMicros` the window's run charges from the ledger, measured as | |
| 218 | + | `usage` measures them. These counts are for views only: runs are still | |
| 219 | + | priced from AI Gateway's logs, never from `token_usage`. | |
| 220 | + | ||
| 203 | 221 | ## Tables (migration `0022_costs_and_margin.sql`) | |
| 204 | 222 | ||
| 205 | 223 | `cost_lines`, `cost_map`, `revenue_map`, `own_counts`, |
| 789 | 789 | /** What the workspace's agents cost since `since`, broken down. Members only. */ | |
| 790 | 790 | usage(workspace: string, viewer: Viewer, since: string): Promise<Result<Usage>>; | |
| 791 | 791 | /** | |
| 792 | + | * The model tokens the workspace's runs used, day by day over the last | |
| 793 | + | * `days` (42, at most 366), for everyone or for one `person`. Members | |
| 794 | + | * only; a member may ask only for themselves, an owner for anyone. | |
| 795 | + | */ | |
| 796 | + | tokenUsage(workspace: string, viewer: User, options?: { person?: string; days?: number }): Promise<Result<TokenUsage>>; | |
| 797 | + | /** | |
| 798 | + | * What one model answer used, added to its run's count for the day. The | |
| 799 | + | * model proxy sends it; for usage views only, as runs are priced from AI | |
| 800 | + | * Gateway. False when there was nothing to count. | |
| 801 | + | */ | |
| 802 | + | recordTokens(usage: { | |
| 803 | + | workspace: string; | |
| 804 | + | /** The model session's id, one per run. */ | |
| 805 | + | session: string; | |
| 806 | + | /** The person the run is for, by username. */ | |
| 807 | + | person?: string | null; | |
| 808 | + | model: string; | |
| 809 | + | tier?: "small" | "large" | null; | |
| 810 | + | input: number; | |
| 811 | + | output: number; | |
| 812 | + | cacheRead: number; | |
| 813 | + | cacheWrite: number; | |
| 814 | + | }): Promise<Result<boolean>>; | |
| 815 | + | /** | |
| 792 | 816 | * Prepays usage ($25 at least) and returns the page to send the person to: | |
| 793 | 817 | * by card with 3-D Secure, or by bank transfer from $1,000. Owners only. | |
| 794 | 818 | * The payment's id comes back to `returnUrl` as `session`. | |
| 959 | 983 | /** One slice of usage: what it was for, what it cost, how many runs. */ | |
| 960 | 984 | export type UsageSlice = { key: string; micros: number; runs: number }; | |
| 961 | 985 | ||
| 986 | + | /** The model tokens runs used over a window of days. */ | |
| 987 | + | export type TokenUsage = { | |
| 988 | + | /** `YYYY-MM-DD`, the first day counted. */ | |
| 989 | + | since: string; | |
| 990 | + | /** The window's length: 42 unless asked, 366 at most. */ | |
| 991 | + | days: number; | |
| 992 | + | /** Null for the whole workspace. */ | |
| 993 | + | person: string | null; | |
| 994 | + | totalTokens: number; | |
| 995 | + | inputTokens: number; | |
| 996 | + | outputTokens: number; | |
| 997 | + | cacheReadTokens: number; | |
| 998 | + | cacheWriteTokens: number; | |
| 999 | + | /** What those runs were charged, as `usage` measures it. */ | |
| 1000 | + | costMicros: number; | |
| 1001 | + | /** Days in the window with any tokens. */ | |
| 1002 | + | activeDays: number; | |
| 1003 | + | /** Every day in the window, oldest first, zeros included. */ | |
| 1004 | + | byDay: { day: string; tokens: number }[]; | |
| 1005 | + | }; | |
| 1006 | + | ||
| 962 | 1007 | /** What a workspace's agents cost over a period. */ | |
| 963 | 1008 | export type Usage = { | |
| 964 | 1009 | since: string; |
| 373 | 373 | before: filter.before ?? null, | |
| 374 | 374 | }), | |
| 375 | 375 | usage: (workspace, viewer, since) => call("usage", { workspace, viewer, since }), | |
| 376 | + | tokenUsage: (workspace, viewer, options = {}) => | |
| 377 | + | call("token_usage", { workspace, viewer, person: options.person ?? null, days: options.days ?? null }), | |
| 378 | + | recordTokens: (usage) => call("record_tokens", usage), | |
| 376 | 379 | checkout: (actor, workspace, amountCents, returnUrl, method = "card") => | |
| 377 | 380 | call("checkout", { actor, workspace, amountCents, returnUrl, method }), | |
| 378 | 381 | confirm: (workspace, viewer, session) => call("confirm", { workspace, viewer, session }), |
| 144 | 144 | session: string; | |
| 145 | 145 | /** For `g1t`: the tier the run was routed to. */ | |
| 146 | 146 | tier?: "small" | "large" | null; | |
| 147 | + | /** The person the run is for, by username. Null when nobody asked; never `g1t`. */ | |
| 148 | + | requestedBy: string | null; | |
| 147 | 149 | baseUrl: string | null; | |
| 148 | 150 | apiKey: string | null; | |
| 149 | 151 | authHeader: string | null; | |
| 187 | 189 | hostedOpen: boolean; | |
| 188 | 190 | /** The tier the run is routed to on g1t's hosted models, for the gateway's logs. */ | |
| 189 | 191 | tier?: "small" | "large" | null; | |
| 192 | + | /** The person the run is for, by username, so its tokens show under them. */ | |
| 193 | + | requestedBy?: string | null; | |
| 190 | 194 | }): Promise<Result<ModelSession>>; | |
| 191 | 195 | routes(workspace: string, viewer: Viewer): Promise<Result<ModelRoute[]>>; | |
| 192 | 196 | setRoutes(actor: User, workspace: string, routes: ModelRoute[]): Promise<Result<ModelRoute[]>>; |
| 1 | + | -- Model tokens, added up per day and run, for usage views: the model proxy | |
| 2 | + | -- reports what each answer used. Billing still prices runs from AI | |
| 3 | + | -- Gateway, never from these. | |
| 4 | + | CREATE TABLE token_usage ( | |
| 5 | + | -- YYYY-MM-DD, UTC. | |
| 6 | + | day TEXT NOT NULL, | |
| 7 | + | workspace TEXT NOT NULL, | |
| 8 | + | -- The person the run was for, by username; empty when nobody asked. | |
| 9 | + | person TEXT NOT NULL DEFAULT '', | |
| 10 | + | -- The model session's id, one per run (runs.session_id). | |
| 11 | + | session TEXT NOT NULL, | |
| 12 | + | model TEXT NOT NULL, | |
| 13 | + | -- On g1t's hosted models: small or large. | |
| 14 | + | tier TEXT, | |
| 15 | + | input INTEGER NOT NULL DEFAULT 0, | |
| 16 | + | output INTEGER NOT NULL DEFAULT 0, | |
| 17 | + | cache_read INTEGER NOT NULL DEFAULT 0, | |
| 18 | + | cache_write INTEGER NOT NULL DEFAULT 0, | |
| 19 | + | requests INTEGER NOT NULL DEFAULT 0, | |
| 20 | + | PRIMARY KEY (day, workspace, person, session, model) | |
| 21 | + | ); | |
| 22 | + | CREATE INDEX token_usage_by_workspace_day ON token_usage (workspace, day); | |
| 23 | + | CREATE INDEX token_usage_by_person_day ON token_usage (workspace, person, day); |
| 40 | 40 | mod retention; | |
| 41 | 41 | mod stripe; | |
| 42 | 42 | mod stripe_sync; | |
| 43 | + | mod tokens; | |
| 43 | 44 | ||
| 44 | 45 | use g1t_contracts::billing::*; | |
| 45 | 46 | use g1t_contracts::time::rfc3339; | |
| 1106 | 1107 | "account" => reply(&billing.account(args(body)?).await?), | |
| 1107 | 1108 | "ledger" => reply(&billing.ledger(args(body)?).await?), | |
| 1108 | 1109 | "usage" => reply(&billing.usage(args(body)?).await?), | |
| 1110 | + | "record_tokens" => reply(&billing.record_tokens(args(body)?).await?), | |
| 1111 | + | "token_usage" => reply(&billing.token_usage(args(body)?).await?), | |
| 1109 | 1112 | "checkout" => reply(&billing.checkout(args(body)?).await?), | |
| 1110 | 1113 | "confirm" => reply(&billing.confirm(args(body)?).await?), | |
| 1111 | 1114 | "can_start" => reply(&billing.can_start(args(body)?).await?), |
| 1 | + | //! Model tokens, counted per run for usage views. | |
| 2 | + | //! | |
| 3 | + | //! The model proxy reports what each answer used (`record_tokens`), and the | |
| 4 | + | //! day's row for that run adds it up: one row per day, workspace, person, | |
| 5 | + | //! session and model. `token_usage` reads a window of those back, for the | |
| 6 | + | //! whole workspace or for one person, with what the window's runs were | |
| 7 | + | //! charged. Billing still prices runs from AI Gateway, never from these. | |
| 8 | + | ||
| 9 | + | use futures_util::future::try_join; | |
| 10 | + | use g1t_contracts::billing::{DayTokens, RecordTokensArgs, TokenUsage, TokenUsageArgs}; | |
| 11 | + | use g1t_contracts::identity::AGENT_NAME; | |
| 12 | + | use g1t_contracts::time::rfc3339; | |
| 13 | + | use g1t_contracts::{FailureCode, Outcome, Role}; | |
| 14 | + | use g1t_kit::now_ms; | |
| 15 | + | use serde::Deserialize; | |
| 16 | + | use worker::wasm_bindgen::JsValue; | |
| 17 | + | use worker::Result; | |
| 18 | + | ||
| 19 | + | use crate::{Billing, members_only}; | |
| 20 | + | ||
| 21 | + | /// The window `token_usage` shows when asked for none, and the longest. | |
| 22 | + | pub(crate) const DEFAULT_DAYS: u32 = 42; | |
| 23 | + | pub(crate) const MAX_DAYS: u32 = 366; | |
| 24 | + | const DAY_MS: u64 = 86_400_000; | |
| 25 | + | ||
| 26 | + | /// One answer's tokens, ready to add to its day's row. | |
| 27 | + | #[derive(Debug, PartialEq)] | |
| 28 | + | pub(crate) struct TokenRow { | |
| 29 | + | pub day: String, | |
| 30 | + | pub workspace: String, | |
| 31 | + | pub person: String, | |
| 32 | + | pub session: String, | |
| 33 | + | pub model: String, | |
| 34 | + | pub tier: Option<String>, | |
| 35 | + | pub counts: [u64; 4], | |
| 36 | + | } | |
| 37 | + | ||
| 38 | + | /// The UTC day of a time, `YYYY-MM-DD`. | |
| 39 | + | fn day_of(ms: u64) -> String { | |
| 40 | + | rfc3339(ms)[..10].to_owned() | |
| 41 | + | } | |
| 42 | + | ||
| 43 | + | /// Who a run's tokens count for: a username, lowercased, or empty for | |
| 44 | + | /// nobody. The agent is not a person. | |
| 45 | + | pub(crate) fn person_of(username: Option<&str>) -> String { | |
| 46 | + | let name = username.unwrap_or_default().trim().to_lowercase(); | |
| 47 | + | if name == AGENT_NAME { String::new() } else { name } | |
| 48 | + | } | |
| 49 | + | ||
| 50 | + | /// What to add for one report, or None when there is nothing to count or | |
| 51 | + | /// nothing to count it under. | |
| 52 | + | pub(crate) fn token_row(a: &RecordTokensArgs, now: u64) -> Option<TokenRow> { | |
| 53 | + | let counts = [a.input, a.output, a.cache_read, a.cache_write]; | |
| 54 | + | let workspace = a.workspace.trim().to_lowercase(); | |
| 55 | + | let session = a.session.trim(); | |
| 56 | + | if counts.iter().all(|n| *n == 0) || workspace.is_empty() || session.is_empty() { | |
| 57 | + | return None; | |
| 58 | + | } | |
| 59 | + | let model = a.model.trim(); | |
| 60 | + | Some(TokenRow { | |
| 61 | + | day: day_of(now), | |
| 62 | + | workspace, | |
| 63 | + | person: person_of(a.person.as_deref()), | |
| 64 | + | session: session.chars().take(64).collect(), | |
| 65 | + | model: if model.is_empty() { "unknown".to_owned() } else { model.chars().take(200).collect() }, | |
| 66 | + | tier: a.tier.as_deref().filter(|tier| matches!(*tier, "small" | "large")).map(str::to_owned), | |
| 67 | + | counts, | |
| 68 | + | }) | |
| 69 | + | } | |
| 70 | + | ||
| 71 | + | /// The days of a window that ends today, oldest first, and how many. | |
| 72 | + | pub(crate) fn window(now: u64, days: Option<u32>) -> Vec<String> { | |
| 73 | + | let days = days.unwrap_or(DEFAULT_DAYS).clamp(1, MAX_DAYS); | |
| 74 | + | (0..u64::from(days)).rev().map(|back| day_of(now.saturating_sub(back * DAY_MS))).collect() | |
| 75 | + | } | |
| 76 | + | ||
| 77 | + | /// Every day of the window with its tokens, zeros included. | |
| 78 | + | pub(crate) fn fill(days: &[String], counted: &[(String, u64)]) -> Vec<DayTokens> { | |
| 79 | + | days.iter() | |
| 80 | + | .map(|day| DayTokens { | |
| 81 | + | day: day.clone(), | |
| 82 | + | tokens: counted.iter().filter(|(d, _)| d == day).map(|(_, n)| n).sum(), | |
| 83 | + | }) | |
| 84 | + | .collect() | |
| 85 | + | } | |
| 86 | + | ||
| 87 | + | /// D1 takes numbers as doubles; a count of tokens fits exactly. | |
| 88 | + | fn number(n: u64) -> JsValue { | |
| 89 | + | JsValue::from_f64(n as f64) | |
| 90 | + | } | |
| 91 | + | ||
| 92 | + | impl Billing { | |
| 93 | + | pub(crate) async fn record_tokens(&self, a: RecordTokensArgs) -> Result<Outcome<bool>> { | |
| 94 | + | let Some(row) = token_row(&a, now_ms()) else { | |
| 95 | + | return Ok(Outcome::Ok(false)); | |
| 96 | + | }; | |
| 97 | + | let [input, output, cache_read, cache_write] = row.counts; | |
| 98 | + | self.db | |
| 99 | + | .prepare( | |
| 100 | + | "INSERT INTO token_usage (day, workspace, person, session, model, tier, input, output, cache_read, cache_write, requests) | |
| 101 | + | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1) | |
| 102 | + | ON CONFLICT (day, workspace, person, session, model) DO UPDATE SET | |
| 103 | + | input = input + excluded.input, | |
| 104 | + | output = output + excluded.output, | |
| 105 | + | cache_read = cache_read + excluded.cache_read, | |
| 106 | + | cache_write = cache_write + excluded.cache_write, | |
| 107 | + | requests = requests + 1, | |
| 108 | + | tier = COALESCE(excluded.tier, tier)", | |
| 109 | + | ) | |
| 110 | + | .bind(&[ | |
| 111 | + | row.day.into(), | |
| 112 | + | row.workspace.into(), | |
| 113 | + | row.person.into(), | |
| 114 | + | row.session.into(), | |
| 115 | + | row.model.into(), | |
| 116 | + | row.tier.map_or(JsValue::NULL, JsValue::from), | |
| 117 | + | number(input), | |
| 118 | + | number(output), | |
| 119 | + | number(cache_read), | |
| 120 | + | number(cache_write), | |
| 121 | + | ])? | |
| 122 | + | .run() | |
| 123 | + | .await?; | |
| 124 | + | Ok(Outcome::Ok(true)) | |
| 125 | + | } | |
| 126 | + | ||
| 127 | + | pub(crate) async fn token_usage(&self, a: TokenUsageArgs) -> Result<Outcome<TokenUsage>> { | |
| 128 | + | let workspace = a.workspace.to_lowercase(); | |
| 129 | + | let Some(viewer) = a.viewer.filter(|viewer| viewer.is_member(&workspace)) else { | |
| 130 | + | return Ok(members_only()); | |
| 131 | + | }; | |
| 132 | + | let person = a.person.as_deref().map(|name| person_of(Some(name))).filter(|name| !name.is_empty()); | |
| 133 | + | if let Some(person) = &person | |
| 134 | + | && *person != viewer.username.to_lowercase() | |
| 135 | + | && viewer.role_in(&workspace) != Some(Role::Owner) | |
| 136 | + | { | |
| 137 | + | return Ok(Outcome::fail( | |
| 138 | + | FailureCode::Forbidden, | |
| 139 | + | "Only owners can see another person's usage.", | |
| 140 | + | )); | |
| 141 | + | } | |
| 142 | + | let days = window(now_ms(), a.days); | |
| 143 | + | let since = days[0].clone(); | |
| 144 | + | let mut binds: Vec<JsValue> = vec![workspace.as_str().into(), since.as_str().into()]; | |
| 145 | + | if let Some(person) = &person { | |
| 146 | + | binds.push(person.as_str().into()); | |
| 147 | + | } | |
| 148 | + | let for_person = if person.is_some() { " AND person = ?3" } else { "" }; | |
| 149 | + | ||
| 150 | + | #[derive(Deserialize)] | |
| 151 | + | struct DayRow { | |
| 152 | + | day: String, | |
| 153 | + | input: Option<f64>, | |
| 154 | + | output: Option<f64>, | |
| 155 | + | cache_read: Option<f64>, | |
| 156 | + | cache_write: Option<f64>, | |
| 157 | + | } | |
| 158 | + | #[derive(Deserialize)] | |
| 159 | + | struct Charged { | |
| 160 | + | micros: Option<f64>, | |
| 161 | + | } | |
| 162 | + | let by_day = async { | |
| 163 | + | self.db | |
| 164 | + | .prepare(format!( | |
| 165 | + | "SELECT day, SUM(input) AS input, SUM(output) AS output, SUM(cache_read) AS cache_read, SUM(cache_write) AS cache_write | |
| 166 | + | FROM token_usage WHERE workspace = ?1 AND day >= ?2{for_person} GROUP BY day" | |
| 167 | + | )) | |
| 168 | + | .bind(&binds)? | |
| 169 | + | .all() | |
| 170 | + | .await? | |
| 171 | + | .results::<DayRow>() | |
| 172 | + | }; | |
| 173 | + | // What the window's runs cost at g1t's price, whoever paid: what was | |
| 174 | + | // left to pay plus what included usage, credit, the trial, the | |
| 175 | + | // open-source pool or a comp covered, so a comped workspace's runs | |
| 176 | + | // still show their cost. While g1t charges nothing, the provider's | |
| 177 | + | // cost. For one person, the runs whose sessions counted their tokens. | |
| 178 | + | let measure = if self.free { | |
| 179 | + | "COALESCE(l.cost_micros, 0)" | |
| 180 | + | } else { | |
| 181 | + | "(-l.amount_micros + l.credit_micros + l.trial_micros + l.oss_micros + l.given_micros)" | |
| 182 | + | }; | |
| 183 | + | let sessions = if person.is_some() { | |
| 184 | + | " AND r.session_id IN (SELECT session FROM token_usage WHERE workspace = ?1 AND person = ?3 AND day >= ?2)" | |
| 185 | + | } else { | |
| 186 | + | "" | |
| 187 | + | }; | |
| 188 | + | let charged = async { | |
| 189 | + | self.db | |
| 190 | + | .prepare(format!( | |
| 191 | + | "SELECT SUM({measure}) AS micros FROM ledger l JOIN runs r ON r.id = l.reference | |
| 192 | + | WHERE l.workspace = ?1 AND l.kind = 'usage' AND l.created_at >= ?2{sessions}" | |
| 193 | + | )) | |
| 194 | + | .bind(&binds)? | |
| 195 | + | .first::<Charged>(None) | |
| 196 | + | .await | |
| 197 | + | }; | |
| 198 | + | // Both read independently, so they go to D1 at once. | |
| 199 | + | let (rows, charged) = try_join(by_day, charged).await?; | |
| 200 | + | ||
| 201 | + | let count = |n: Option<f64>| n.unwrap_or_default().max(0.0) as u64; | |
| 202 | + | let mut totals = [0u64; 4]; | |
| 203 | + | let mut counted = Vec::with_capacity(rows.len()); | |
| 204 | + | for row in &rows { | |
| 205 | + | let row_counts = [count(row.input), count(row.output), count(row.cache_read), count(row.cache_write)]; | |
| 206 | + | for (total, n) in totals.iter_mut().zip(row_counts) { | |
| 207 | + | *total += n; | |
| 208 | + | } | |
| 209 | + | counted.push((row.day.clone(), row_counts.iter().sum::<u64>())); | |
| 210 | + | } | |
| 211 | + | let by_day = fill(&days, &counted); | |
| 212 | + | Ok(Outcome::Ok(TokenUsage { | |
| 213 | + | since, | |
| 214 | + | days: days.len() as u32, | |
| 215 | + | person, | |
| 216 | + | total_tokens: totals.iter().sum(), | |
| 217 | + | input_tokens: totals[0], | |
| 218 | + | output_tokens: totals[1], | |
| 219 | + | cache_read_tokens: totals[2], | |
| 220 | + | cache_write_tokens: totals[3], | |
| 221 | + | cost_micros: charged.and_then(|c| c.micros).unwrap_or_default().round() as i64, | |
| 222 | + | active_days: by_day.iter().filter(|day| day.tokens > 0).count() as u32, | |
| 223 | + | by_day, | |
| 224 | + | })) | |
| 225 | + | } | |
| 226 | + | } | |
| 227 | + | ||
| 228 | + | #[cfg(test)] | |
| 229 | + | mod tests { | |
| 230 | + | use super::*; | |
| 231 | + | ||
| 232 | + | // 2026-10-06T12:00:00Z. | |
| 233 | + | const NOW: u64 = 1_791_288_000_000; | |
| 234 | + | ||
| 235 | + | fn args() -> RecordTokensArgs { | |
| 236 | + | RecordTokensArgs { | |
| 237 | + | workspace: " Acme ".into(), | |
| 238 | + | session: "ms_abc".into(), | |
| 239 | + | person: Some("Ada".into()), | |
| 240 | + | model: "claude-opus".into(), | |
| 241 | + | tier: Some("large".into()), | |
| 242 | + | input: 10, | |
| 243 | + | output: 5, | |
| 244 | + | cache_read: 0, | |
| 245 | + | cache_write: 0, | |
| 246 | + | } | |
| 247 | + | } | |
| 248 | + | ||
| 249 | + | #[test] | |
| 250 | + | fn a_report_is_added_to_today_under_its_person() { | |
| 251 | + | let row = token_row(&args(), NOW).unwrap(); | |
| 252 | + | assert_eq!(row.day, "2026-10-06"); | |
| 253 | + | assert_eq!(row.workspace, "acme"); | |
| 254 | + | assert_eq!(row.person, "ada"); | |
| 255 | + | assert_eq!(row.tier.as_deref(), Some("large")); | |
| 256 | + | assert_eq!(row.counts, [10, 5, 0, 0]); | |
| 257 | + | } | |
| 258 | + | ||
| 259 | + | #[test] | |
| 260 | + | fn nothing_used_or_nowhere_to_put_it_is_not_counted() { | |
| 261 | + | assert!(token_row(&RecordTokensArgs { input: 0, output: 0, ..args() }, NOW).is_none()); | |
| 262 | + | assert!(token_row(&RecordTokensArgs { session: " ".into(), ..args() }, NOW).is_none()); | |
| 263 | + | assert!(token_row(&RecordTokensArgs { workspace: String::new(), ..args() }, NOW).is_none()); | |
| 264 | + | } | |
| 265 | + | ||
| 266 | + | #[test] | |
| 267 | + | fn the_agent_is_nobody_and_odd_tiers_and_models_are_tidied() { | |
| 268 | + | let row = token_row( | |
| 269 | + | &RecordTokensArgs { person: Some("g1t".into()), tier: Some("huge".into()), model: " ".into(), ..args() }, | |
| 270 | + | NOW, | |
| 271 | + | ) | |
| 272 | + | .unwrap(); | |
| 273 | + | assert_eq!(row.person, ""); | |
| 274 | + | assert_eq!(row.tier, None); | |
| 275 | + | assert_eq!(row.model, "unknown"); | |
| 276 | + | assert_eq!(person_of(None), ""); | |
| 277 | + | } | |
| 278 | + | ||
| 279 | + | #[test] | |
| 280 | + | fn the_window_ends_today_oldest_first_and_is_bounded() { | |
| 281 | + | let days = window(NOW, Some(3)); | |
| 282 | + | assert_eq!(days, vec!["2026-10-04", "2026-10-05", "2026-10-06"]); | |
| 283 | + | assert_eq!(window(NOW, None).len(), DEFAULT_DAYS as usize); | |
| 284 | + | assert_eq!(window(NOW, Some(0)).len(), 1); | |
| 285 | + | assert_eq!(window(NOW, Some(5000)).len(), MAX_DAYS as usize); | |
| 286 | + | // Across a month's end. | |
| 287 | + | assert_eq!(window(NOW, Some(7))[0], "2026-09-30"); | |
| 288 | + | } | |
| 289 | + | ||
| 290 | + | #[test] | |
| 291 | + | fn every_day_is_filled_with_zeros_where_nothing_ran() { | |
| 292 | + | let days = window(NOW, Some(3)); | |
| 293 | + | let filled = fill(&days, &[("2026-10-05".into(), 40), ("2026-09-01".into(), 9)]); | |
| 294 | + | let tokens: Vec<u64> = filled.iter().map(|d| d.tokens).collect(); | |
| 295 | + | assert_eq!(tokens, vec![0, 40, 0]); | |
| 296 | + | assert_eq!(filled[2].day, "2026-10-06"); | |
| 297 | + | } | |
| 298 | + | } |
| 1 | + | -- The person a run is for, by username: who asked g1t for the work. The | |
| 2 | + | -- model proxy reports each run's tokens under them, for usage views. Null | |
| 3 | + | -- when nobody asked, and never g1t's own agent. | |
| 4 | + | ALTER TABLE model_sessions ADD COLUMN requested_by TEXT; |
| 123 | 123 | out | |
| 124 | 124 | } | |
| 125 | 125 | ||
| 126 | + | /// Who a run is for, as kept on its session: a username, lowercased. The | |
| 127 | + | /// agent's own name is not a person, so it is kept as nobody. | |
| 128 | + | fn requester(username: Option<&str>) -> Option<String> { | |
| 129 | + | let name = username?.trim().to_lowercase(); | |
| 130 | + | (!name.is_empty() && name != g1t_contracts::identity::AGENT_NAME).then_some(name) | |
| 131 | + | } | |
| 132 | + | ||
| 126 | 133 | /// A model session's public id: the start of its token's hash. | |
| 127 | 134 | fn session_id(token_hash: &str) -> String { | |
| 128 | 135 | format!("ms_{}", &token_hash[..token_hash.len().min(24)]) | |
| 139 | 146 | model: Option<String>, | |
| 140 | 147 | #[serde(default)] | |
| 141 | 148 | tier: Option<String>, | |
| 149 | + | #[serde(default)] | |
| 150 | + | requested_by: Option<String>, | |
| 142 | 151 | } | |
| 143 | 152 | ||
| 144 | 153 | #[derive(Deserialize)] | |
| 1229 | 1238 | .bind(&[rfc3339(now).into()])?, | |
| 1230 | 1239 | self.db | |
| 1231 | 1240 | .prepare( | |
| 1232 | − | "INSERT INTO model_sessions (token_hash, workspace, connection_id, repo, number, task, expires_at, model, tier) | |
| 1233 | − | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", | |
| 1241 | + | "INSERT INTO model_sessions (token_hash, workspace, connection_id, repo, number, task, expires_at, model, tier, requested_by) | |
| 1242 | + | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", | |
| 1234 | 1243 | ) | |
| 1235 | 1244 | .bind(&[ | |
| 1236 | 1245 | crypto::sha256_hex(&token).into(), | |
| 1247 | 1256 | .as_deref() | |
| 1248 | 1257 | .filter(|tier| connection.is_none() && matches!(*tier, "small" | "large")), | |
| 1249 | 1258 | ), | |
| 1259 | + | optional(requester(a.requested_by.as_deref()).as_deref()), | |
| 1250 | 1260 | ])?, | |
| 1251 | 1261 | ]) | |
| 1252 | 1262 | .await?; | |
| 1282 | 1292 | task: session.task, | |
| 1283 | 1293 | session: session_id(&session.token_hash), | |
| 1284 | 1294 | tier: session.tier, | |
| 1295 | + | requested_by: session.requested_by, | |
| 1285 | 1296 | base_url: None, | |
| 1286 | 1297 | api_key: None, | |
| 1287 | 1298 | auth_header: None, | |
| 1490 | 1501 | ||
| 1491 | 1502 | #[cfg(test)] | |
| 1492 | 1503 | mod close_tests { | |
| 1493 | − | use super::closable_hashes; | |
| 1504 | + | use super::{closable_hashes, requester}; | |
| 1505 | + | ||
| 1506 | + | #[test] | |
| 1507 | + | fn a_session_is_for_a_person_never_the_agent() { | |
| 1508 | + | assert_eq!(requester(Some(" Ada ")).as_deref(), Some("ada")); | |
| 1509 | + | assert_eq!(requester(Some("g1t")), None); | |
| 1510 | + | assert_eq!(requester(Some("G1T")), None); | |
| 1511 | + | assert_eq!(requester(Some(" ")), None); | |
| 1512 | + | assert_eq!(requester(None), None); | |
| 1513 | + | } | |
| 1494 | 1514 | ||
| 1495 | 1515 | #[test] | |
| 1496 | 1516 | fn only_token_hashes_are_closed() { |
| 11 | 11 | * the run ends (the runner closes its session then, and lookups are kept | |
| 12 | 12 | * only seconds), and nothing of the workspace's. | |
| 13 | 13 | * | |
| 14 | − | * Responses stream through. | |
| 14 | + | * Responses stream through. What each answer used is read from a copy as | |
| 15 | + | * it passes and reported to billing afterwards, counted per run for usage | |
| 16 | + | * views. | |
| 15 | 17 | */ | |
| 16 | − | import { type ModelUpstream, type ServiceBinding, integrationsClient } from "@g1t/contracts"; | |
| 18 | + | import { type ModelUpstream, type ServiceBinding, billingClient, integrationsClient } from "@g1t/contracts"; | |
| 17 | 19 | ||
| 18 | 20 | import { type AnthropicRequest, StreamTranslator, errorFromChat, estimateTokens, fromChat, toChat } from "./openai"; | |
| 21 | + | import { isAnswer, tokenReport } from "./report"; | |
| 19 | 22 | import { type HostedRouting, presentedToken, upstreamRequest } from "./route"; | |
| 23 | + | import { measure } from "./usage"; | |
| 20 | 24 | ||
| 21 | 25 | interface Env extends HostedRouting { | |
| 22 | 26 | INTEGRATIONS: ServiceBinding; | |
| 27 | + | BILLING: ServiceBinding; | |
| 23 | 28 | } | |
| 24 | 29 | ||
| 25 | 30 | /** | |
| 48 | 53 | ); | |
| 49 | 54 | } | |
| 50 | 55 | ||
| 56 | + | /** | |
| 57 | + | * Passes an answer through and, once it has all gone by, tells billing what | |
| 58 | + | * it used. Reporting happens after the answer, and a report that fails is | |
| 59 | + | * dropped: the answer never waits on it or breaks for it. | |
| 60 | + | */ | |
| 61 | + | function counted(answer: Response, upstream: ModelUpstream, env: Env, ctx: ExecutionContext): Response { | |
| 62 | + | const { response, tokens, model } = measure(answer); | |
| 63 | + | ctx.waitUntil( | |
| 64 | + | (async () => { | |
| 65 | + | const report = tokenReport(upstream, await model, await tokens); | |
| 66 | + | if (report) await billingClient(env.BILLING).recordTokens(report); | |
| 67 | + | })().catch(() => undefined), | |
| 68 | + | ); | |
| 69 | + | return response; | |
| 70 | + | } | |
| 71 | + | ||
| 51 | 72 | /** Sends an Anthropic request to a provider that speaks OpenAI's API. */ | |
| 52 | 73 | async function viaChat(upstream: ModelUpstream, path: string, request: Request): Promise<Response> { | |
| 53 | 74 | const body = (await request.json()) as AnthropicRequest; | |
| 95 | 116 | } | |
| 96 | 117 | ||
| 97 | 118 | export default { | |
| 98 | − | async fetch(request: Request, env: Env): Promise<Response> { | |
| 119 | + | async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> { | |
| 99 | 120 | const url = new URL(request.url); | |
| 100 | 121 | if (url.pathname === "/" || url.pathname === "") { | |
| 101 | 122 | return new Response("g1t's model proxy, for g1t's sandboxes. See https://docs.g1t.sh/guides/models/\n"); | |
| 107 | 128 | if (!upstream) return refuse(401, "This run's model token has expired, or its model connection was removed."); | |
| 108 | 129 | ||
| 109 | 130 | const path = url.pathname.slice("/anthropic".length) + url.search; | |
| 110 | − | if (upstream.api === "openai") return viaChat(upstream, path, request); | |
| 131 | + | // Both routes answer in Anthropic's shape, so one reading counts either. | |
| 132 | + | const answer = (response: Response) => (isAnswer(url.pathname.slice("/anthropic".length)) ? counted(response, upstream, env, ctx) : response); | |
| 133 | + | if (upstream.api === "openai") return answer(await viaChat(upstream, path, request)); | |
| 111 | 134 | ||
| 112 | 135 | const { url: target, headers } = upstreamRequest(upstream, env, path, request.headers); | |
| 113 | 136 | // A route that names a model gets it for every request of the run, | |
| 118 | 141 | body = JSON.stringify({ ...parsed, model: upstream.model }); | |
| 119 | 142 | headers.delete("content-length"); | |
| 120 | 143 | } | |
| 121 | − | return fetch(target, { method: request.method, headers, body }); | |
| 144 | + | return answer(await fetch(target, { method: request.method, headers, body })); | |
| 122 | 145 | }, | |
| 123 | 146 | } satisfies ExportedHandler<Env>; |
| 1 | + | import assert from "node:assert/strict"; | |
| 2 | + | import { test } from "node:test"; | |
| 3 | + | ||
| 4 | + | import type { ModelUpstream } from "@g1t/contracts"; | |
| 5 | + | ||
| 6 | + | import { isAnswer, tokenReport } from "./report.ts"; | |
| 7 | + | import { NO_TOKENS, measure } from "./usage.ts"; | |
| 8 | + | ||
| 9 | + | const upstream: ModelUpstream = { | |
| 10 | + | route: "g1t", | |
| 11 | + | api: "anthropic", | |
| 12 | + | model: null, | |
| 13 | + | official: false, | |
| 14 | + | provider: "g1t", | |
| 15 | + | workspace: "acme", | |
| 16 | + | repo: "acme/web", | |
| 17 | + | number: 7, | |
| 18 | + | task: "implement", | |
| 19 | + | session: "ms_abc", | |
| 20 | + | tier: "large", | |
| 21 | + | requestedBy: "ada", | |
| 22 | + | baseUrl: null, | |
| 23 | + | apiKey: null, | |
| 24 | + | authHeader: null, | |
| 25 | + | }; | |
| 26 | + | ||
| 27 | + | test("only messages are answers; counting tokens is not", () => { | |
| 28 | + | assert.equal(isAnswer("/v1/messages"), true); | |
| 29 | + | assert.equal(isAnswer("/v1/messages?beta=true"), true); | |
| 30 | + | assert.equal(isAnswer("/v1/messages/count_tokens?beta=true"), false); | |
| 31 | + | assert.equal(isAnswer("/v1/models"), false); | |
| 32 | + | }); | |
| 33 | + | ||
| 34 | + | test("an answer's tokens are reported under its run and person", () => { | |
| 35 | + | const report = tokenReport(upstream, "claude-opus-4", { input: 3, output: 9, cacheRead: 100, cacheWrite: 0 }); | |
| 36 | + | assert.deepEqual(report, { | |
| 37 | + | workspace: "acme", | |
| 38 | + | session: "ms_abc", | |
| 39 | + | person: "ada", | |
| 40 | + | model: "claude-opus-4", | |
| 41 | + | tier: "large", | |
| 42 | + | input: 3, | |
| 43 | + | output: 9, | |
| 44 | + | cacheRead: 100, | |
| 45 | + | cacheWrite: 0, | |
| 46 | + | }); | |
| 47 | + | // A route that names a model counts under it. | |
| 48 | + | assert.equal(tokenReport({ ...upstream, model: "gpt-x", tier: null }, "other", { ...NO_TOKENS, output: 1 })?.model, "gpt-x"); | |
| 49 | + | }); | |
| 50 | + | ||
| 51 | + | test("nothing used, or no session to count it under, is not reported", () => { | |
| 52 | + | assert.equal(tokenReport(upstream, null, { ...NO_TOKENS }), null); | |
| 53 | + | assert.equal(tokenReport({ ...upstream, session: "" }, null, { ...NO_TOKENS, input: 1 }), null); | |
| 54 | + | }); | |
| 55 | + | ||
| 56 | + | test("the model that answered is read from the answer", async () => { | |
| 57 | + | const body = `data: ${JSON.stringify({ type: "message_start", message: { model: "claude-x", usage: { input_tokens: 1 } } })}\n\n`; | |
| 58 | + | const streamed = measure(new Response(body, { headers: { "content-type": "text/event-stream" } })); | |
| 59 | + | await streamed.response.text(); | |
| 60 | + | assert.equal(await streamed.model, "claude-x"); | |
| 61 | + | const whole = measure(Response.json({ model: "claude-y", usage: { output_tokens: 2 } })); | |
| 62 | + | await whole.response.text(); | |
| 63 | + | assert.equal(await whole.model, "claude-y"); | |
| 64 | + | }); |
| 1 | + | /** | |
| 2 | + | * What the proxy tells billing about one answer: its tokens, under the | |
| 3 | + | * run's session and the person it is for, for usage views. Billing still | |
| 4 | + | * prices runs from AI Gateway, not from these. | |
| 5 | + | */ | |
| 6 | + | import type { ModelUpstream } from "@g1t/contracts"; | |
| 7 | + | ||
| 8 | + | import type { Tokens } from "./usage"; | |
| 9 | + | ||
| 10 | + | export type TokenReport = { | |
| 11 | + | workspace: string; | |
| 12 | + | session: string; | |
| 13 | + | person: string | null; | |
| 14 | + | model: string; | |
| 15 | + | tier: "small" | "large" | null; | |
| 16 | + | input: number; | |
| 17 | + | output: number; | |
| 18 | + | cacheRead: number; | |
| 19 | + | cacheWrite: number; | |
| 20 | + | }; | |
| 21 | + | ||
| 22 | + | /** | |
| 23 | + | * Whether a request's answer is a model's answer, and so used tokens. | |
| 24 | + | * Counting tokens is a question about a request, not an answer. | |
| 25 | + | */ | |
| 26 | + | export function isAnswer(path: string): boolean { | |
| 27 | + | return /^\/v1\/messages\/?(\?|$)/.test(path); | |
| 28 | + | } | |
| 29 | + | ||
| 30 | + | /** | |
| 31 | + | * The report for one answer, or null when it used nothing or its session | |
| 32 | + | * has no id to count it under. The model is the run's when its route names | |
| 33 | + | * one, else the one that answered. | |
| 34 | + | */ | |
| 35 | + | export function tokenReport(upstream: ModelUpstream, answeredBy: string | null, tokens: Tokens): TokenReport | null { | |
| 36 | + | const used = tokens.input + tokens.output + tokens.cacheRead + tokens.cacheWrite; | |
| 37 | + | if (used === 0 || !upstream.session) return null; | |
| 38 | + | return { | |
| 39 | + | workspace: upstream.workspace, | |
| 40 | + | session: upstream.session, | |
| 41 | + | person: upstream.requestedBy ?? null, | |
| 42 | + | model: upstream.model ?? answeredBy ?? "unknown", | |
| 43 | + | tier: upstream.tier ?? null, | |
| 44 | + | input: tokens.input, | |
| 45 | + | output: tokens.output, | |
| 46 | + | cacheRead: tokens.cacheRead, | |
| 47 | + | cacheWrite: tokens.cacheWrite, | |
| 48 | + | }; | |
| 49 | + | } |
| 5 | 5 | ||
| 6 | 6 | import { presentedToken, upstreamRequest } from "./route.ts"; | |
| 7 | 7 | ||
| 8 | − | const run = { workspace: "acme", repo: "acme/web", number: 7, task: "implement", session: "ms_abc", baseUrl: null, apiKey: null, authHeader: null, api: "anthropic" as const, model: null, official: false, provider: "g1t" }; | |
| 8 | + | const run = { workspace: "acme", repo: "acme/web", number: 7, task: "implement", session: "ms_abc", baseUrl: null, apiKey: null, authHeader: null, api: "anthropic" as const, model: null, official: false, provider: "g1t", requestedBy: null }; | |
| 9 | 9 | const hosted = { AI_GATEWAY_ID: "g1t", CLOUDFLARE_ACCOUNT_ID: "acct", AI_GATEWAY_TOKEN: "gw-token" }; | |
| 10 | 10 | ||
| 11 | 11 | function incoming(): Headers { |
| 47 | 47 | */ | |
| 48 | 48 | export class StreamUsage { | |
| 49 | 49 | tokens: Tokens = { ...NO_TOKENS }; | |
| 50 | + | /** The model that answered, as `message_start` names it. */ | |
| 51 | + | model: string | null = null; | |
| 50 | 52 | private pending = ""; | |
| 51 | 53 | ||
| 52 | 54 | push(text: string): void { | |
| 67 | 69 | ||
| 68 | 70 | private line(line: string): void { | |
| 69 | 71 | if (!line.startsWith("data:")) return; | |
| 70 | − | let event: { type?: string; message?: { usage?: Usage }; usage?: Usage }; | |
| 72 | + | let event: { type?: string; message?: { usage?: Usage; model?: unknown }; usage?: Usage }; | |
| 71 | 73 | try { | |
| 72 | 74 | event = JSON.parse(line.slice(5).trim()); | |
| 73 | 75 | } catch { | |
| 74 | 76 | return; | |
| 75 | 77 | } | |
| 76 | 78 | if (event.type === "message_start") { | |
| 79 | + | if (typeof event.message?.model === "string") this.model = event.message.model; | |
| 77 | 80 | const start = fromUsage(event.message?.usage); | |
| 78 | 81 | this.tokens = { ...start, output: Math.max(this.tokens.output, start.output) }; | |
| 79 | 82 | } else if (event.type === "message_delta" && event.usage) { | |
| 89 | 92 | ||
| 90 | 93 | /** | |
| 91 | 94 | * The answer as it was, and a promise of what it used, read from a copy of | |
| 92 | − | * its body. A failed answer, or one that is not a message, used nothing. | |
| 95 | + | * its body, with the model that answered when it says. A failed answer, or | |
| 96 | + | * one that is not a message, used nothing. | |
| 93 | 97 | */ | |
| 94 | − | export function measure(answer: Response): { response: Response; tokens: Promise<Tokens> } { | |
| 95 | − | if (!answer.ok || !answer.body) return { response: answer, tokens: Promise.resolve({ ...NO_TOKENS }) }; | |
| 98 | + | export function measure(answer: Response): { response: Response; tokens: Promise<Tokens>; model: Promise<string | null> } { | |
| 99 | + | if (!answer.ok || !answer.body) { | |
| 100 | + | return { response: answer, tokens: Promise.resolve({ ...NO_TOKENS }), model: Promise.resolve(null) }; | |
| 101 | + | } | |
| 96 | 102 | const [passed, copy] = answer.body.tee(); | |
| 97 | 103 | const response = new Response(passed, answer); | |
| 98 | 104 | const streaming = (answer.headers.get("content-type") ?? "").includes("text/event-stream"); | |
| 105 | + | let model: string | null = null; | |
| 99 | 106 | const tokens = (async () => { | |
| 100 | 107 | try { | |
| 101 | − | if (!streaming) return fromUsage(((await new Response(copy).json()) as { usage?: Usage }).usage); | |
| 108 | + | if (!streaming) { | |
| 109 | + | const whole = (await new Response(copy).json()) as { usage?: Usage; model?: unknown }; | |
| 110 | + | if (typeof whole.model === "string") model = whole.model; | |
| 111 | + | return fromUsage(whole.usage); | |
| 112 | + | } | |
| 102 | 113 | const reader = copy.pipeThrough(new TextDecoderStream()).getReader(); | |
| 103 | 114 | const usage = new StreamUsage(); | |
| 104 | 115 | for (;;) { | |
| 106 | 117 | if (done) break; | |
| 107 | 118 | usage.push(value); | |
| 108 | 119 | } | |
| 109 | − | return usage.finish(); | |
| 120 | + | const used = usage.finish(); | |
| 121 | + | model = usage.model; | |
| 122 | + | return used; | |
| 110 | 123 | } catch { | |
| 111 | 124 | return { ...NO_TOKENS }; | |
| 112 | 125 | } | |
| 113 | 126 | })(); | |
| 114 | − | return { response, tokens }; | |
| 127 | + | return { response, tokens, model: tokens.then(() => model) }; | |
| 115 | 128 | } |
| 7 | 7 | "workers_dev": false, | |
| 8 | 8 | // Sandboxes reach models only through here, with a token for their run. | |
| 9 | 9 | "routes": [{ "pattern": "models.g1t.sh", "custom_domain": true }], | |
| 10 | − | "services": [{ "binding": "INTEGRATIONS", "service": "g1t-integrations" }], | |
| 10 | + | "services": [ | |
| 11 | + | { "binding": "INTEGRATIONS", "service": "g1t-integrations" }, | |
| 12 | + | // Where each run's tokens are counted, for usage views. | |
| 13 | + | { "binding": "BILLING", "service": "g1t-billing" } | |
| 14 | + | ], | |
| 11 | 15 | "vars": { | |
| 12 | 16 | // g1t's own route, for runs g1t pays for: a Cloudflare AI Gateway, | |
| 13 | 17 | // which holds g1t's key. |
| 1334 | 1334 | * it (`chooseTier`), by what `route` says about it. Whether this is a | |
| 1335 | 1335 | * retry is asked only when it would change the answer: when the work | |
| 1336 | 1336 | * would otherwise go to the small tier. | |
| 1337 | + | * | |
| 1338 | + | * `requestedBy` is the person the run is for, by username, so the run's | |
| 1339 | + | * tokens are counted under them. | |
| 1337 | 1340 | */ | |
| 1338 | 1341 | private async modelEnv( | |
| 1339 | 1342 | task: AgentTask, | |
| 1340 | 1343 | repo: RepoPath, | |
| 1341 | 1344 | pull: number, | |
| 1345 | + | requestedBy: string | null, | |
| 1342 | 1346 | route: RouteInput = {}, | |
| 1343 | 1347 | ): Promise<Result<Record<string, string>>> { | |
| 1344 | 1348 | const routing = parseRouting(this.env.AGENT_ROUTING); | |
| 1358 | 1362 | task, | |
| 1359 | 1363 | hostedOpen: (await this.modelAccess(repo.namespace)).hosted, | |
| 1360 | 1364 | tier, | |
| 1365 | + | requestedBy, | |
| 1361 | 1366 | }); | |
| 1362 | 1367 | if (!opened.ok) return opened; | |
| 1363 | 1368 | session = opened.value; | |
| 1529 | 1534 | task: AgentTask, | |
| 1530 | 1535 | repo: RepoPath, | |
| 1531 | 1536 | pull: number, | |
| 1537 | + | requestedBy: string | null, | |
| 1532 | 1538 | route: RouteInput = {}, | |
| 1533 | 1539 | ): Promise<Record<string, string>> { | |
| 1534 | − | const vars = await this.modelEnv(task, repo, pull, route); | |
| 1540 | + | const vars = await this.modelEnv(task, repo, pull, requestedBy, route); | |
| 1535 | 1541 | if (!vars.ok) throw new Error(vars.error.message); | |
| 1536 | 1542 | return vars.value; | |
| 1537 | 1543 | } | |
| 2185 | 2191 | job.repo, | |
| 2186 | 2192 | job.author, | |
| 2187 | 2193 | ), | |
| 2188 | − | ...(await this.modelEnvOrThrow("implement", job.repo, job.number)), | |
| 2194 | + | ...(await this.modelEnvOrThrow("implement", job.repo, job.number, job.author.username)), | |
| 2189 | 2195 | }, | |
| 2190 | 2196 | }); | |
| 2191 | 2197 | } | |
| 2238 | 2244 | job.repo, | |
| 2239 | 2245 | job.author, | |
| 2240 | 2246 | ), | |
| 2241 | − | ...(await this.modelEnvOrThrow("implement", job.repo, job.number)), | |
| 2247 | + | ...(await this.modelEnvOrThrow("implement", job.repo, job.number, startedBy ?? job.author.username)), | |
| 2242 | 2248 | }, | |
| 2243 | 2249 | }); | |
| 2244 | 2250 | } | |
| 2492 | 2498 | repo, | |
| 2493 | 2499 | actor, | |
| 2494 | 2500 | ), | |
| 2495 | − | ...(await this.modelEnvOrThrow("update", repo, number, { | |
| 2501 | + | ...(await this.modelEnvOrThrow("update", repo, number, actor.username, { | |
| 2496 | 2502 | retried: () => this.failedBefore(actor, repo, "update", number), | |
| 2497 | 2503 | })), | |
| 2498 | 2504 | }, | |
| 2548 | 2554 | `It is for issue #${job.issue.number}: ${job.issue.title}\n\n${job.issue.body}`, | |
| 2549 | 2555 | await this.peopleSaid(job.author, repo, number), | |
| 2550 | 2556 | ]; | |
| 2551 | − | const model = await this.modelEnv("review", repo, number, { | |
| 2557 | + | const model = await this.modelEnv("review", repo, number, job.author.username, { | |
| 2552 | 2558 | change: job.files?.length ? changeSize(job.files, job.sensitive ?? []) : null, | |
| 2553 | 2559 | labels: job.issue?.labels ?? [], | |
| 2554 | 2560 | retried: () => this.failedBefore(job.author, repo, "review", number), | |
| 2612 | 2618 | const started = await work.startPlan(actor, repo, brief); | |
| 2613 | 2619 | if (!started.ok) return started; | |
| 2614 | 2620 | const job = started.value; | |
| 2615 | − | const model = await this.modelEnv("plan", repo, 0, { | |
| 2621 | + | const model = await this.modelEnv("plan", repo, 0, actor.username, { | |
| 2616 | 2622 | retried: () => this.failedBefore(actor, repo, "plan", null, job.brief), | |
| 2617 | 2623 | }); | |
| 2618 | 2624 | if (!model.ok) { | |
| 2767 | 2773 | // Opened without a branch, so it has a fork. | |
| 2768 | 2774 | const fork = pull.fork!; | |
| 2769 | 2775 | ||
| 2770 | − | const model = await this.modelEnv("implement", repo, pull.number); | |
| 2776 | + | const model = await this.modelEnv("implement", repo, pull.number, actor.username); | |
| 2771 | 2777 | if (!model.ok) { | |
| 2772 | 2778 | await work.closePull(actor, repo, pull.number); | |
| 2773 | 2779 | return model; | |
| 2914 | 2920 | ({ title, body } = found.value.issue); | |
| 2915 | 2921 | comments = found.value.comments; | |
| 2916 | 2922 | } | |
| 2917 | − | const model = await this.modelEnv("implement", job.repo, job.number); | |
| 2923 | + | const model = await this.modelEnv("implement", job.repo, job.number, job.actor.username); | |
| 2918 | 2924 | if (!model.ok) return model; | |
| 2919 | 2925 | const source = job.pull?.source ?? job.repo; | |
| 2920 | 2926 | // Reads the code; pushes nothing. Its answer is posted with its tools. |