g1t/services/models/src/usage.ts

115 lines3,943 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

The model proxy can read what an answer used: Anthropic's usage from a whole answer or a stream1/**
2 * What a model answer used, read as it passes through: Anthropic's
3 * `usage`, from a whole JSON answer or from a stream's `message_start`
4 * (input and cache) and `message_delta` (output) events. The proxy adds
5 * these up per run for usage views; billing still prices runs from AI
6 * Gateway, not from these.
7 */
8
9export type Tokens = {
10 input: number;
11 output: number;
12 /** Prompt tokens read from the provider's cache. */
13 cacheRead: number;
14 /** Prompt tokens written to the provider's cache. */
15 cacheWrite: number;
16};
17
18export const NO_TOKENS: Tokens = { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 };
19
20type Usage = {
21 input_tokens?: number;
22 output_tokens?: number;
23 cache_read_input_tokens?: number;
24 cache_creation_input_tokens?: number;
25};
26
27const count = (value: unknown): number => (typeof value === "number" && Number.isFinite(value) && value > 0 ? value : 0);
28
29/** The tokens of one `usage` object. */
30export function fromUsage(usage: Usage | undefined | null): Tokens {
31 return {
32 input: count(usage?.input_tokens),
33 output: count(usage?.output_tokens),
34 cacheRead: count(usage?.cache_read_input_tokens),
35 cacheWrite: count(usage?.cache_creation_input_tokens),
36 };
37}
38
39export function total(tokens: Tokens): number {
40 return tokens.input + tokens.output + tokens.cacheRead + tokens.cacheWrite;
41}
42
43/**
44 * Reads a server-sent event stream a chunk at a time. `message_start`
45 * carries the prompt's tokens; each `message_delta` carries the output so
46 * far (the last one is the total).
47 */
48export class StreamUsage {
49 tokens: Tokens = { ...NO_TOKENS };
50 private pending = "";
51
52 push(text: string): void {
53 this.pending += text;
54 let end = this.pending.indexOf("\n");
55 while (end >= 0) {
56 this.line(this.pending.slice(0, end).trim());
57 this.pending = this.pending.slice(end + 1);
58 end = this.pending.indexOf("\n");
59 }
60 }
61
62 finish(): Tokens {
63 this.line(this.pending.trim());
64 this.pending = "";
65 return this.tokens;
66 }
67
68 private line(line: string): void {
69 if (!line.startsWith("data:")) return;
70 let event: { type?: string; message?: { usage?: Usage }; usage?: Usage };
71 try {
72 event = JSON.parse(line.slice(5).trim());
73 } catch {
74 return;
75 }
76 if (event.type === "message_start") {
77 const start = fromUsage(event.message?.usage);
78 this.tokens = { ...start, output: Math.max(this.tokens.output, start.output) };
79 } else if (event.type === "message_delta" && event.usage) {
80 const delta = fromUsage(event.usage);
81 this.tokens.output = Math.max(this.tokens.output, delta.output);
82 // Some providers repeat the prompt's tokens at the end.
83 if (delta.input) this.tokens.input = delta.input;
84 if (delta.cacheRead) this.tokens.cacheRead = delta.cacheRead;
85 if (delta.cacheWrite) this.tokens.cacheWrite = delta.cacheWrite;
86 }
87 }
88}
89
90/**
91 * 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.
93 */
94export function measure(answer: Response): { response: Response; tokens: Promise<Tokens> } {
95 if (!answer.ok || !answer.body) return { response: answer, tokens: Promise.resolve({ ...NO_TOKENS }) };
96 const [passed, copy] = answer.body.tee();
97 const response = new Response(passed, answer);
98 const streaming = (answer.headers.get("content-type") ?? "").includes("text/event-stream");
99 const tokens = (async () => {
100 try {
101 if (!streaming) return fromUsage(((await new Response(copy).json()) as { usage?: Usage }).usage);
102 const reader = copy.pipeThrough(new TextDecoderStream()).getReader();
103 const usage = new StreamUsage();
104 for (;;) {
105 const { done, value } = await reader.read();
106 if (done) break;
107 usage.push(value);
108 }
109 return usage.finish();
110 } catch {
111 return { ...NO_TOKENS };
112 }
113 })();
114 return { response, tokens };
115}