Commit

The model proxy can read what an answer used: Anthropic's usage from a whole answer or a stream

usage.ts reads input, output and cache tokens from a JSON answer's usage or a stream's message_start and message_delta, from a copy of the body, so the answer passes through untouched. Not wired in yet: next it is reported per run, for usage views.

syntaqxcommitted Parent7956b44Browse files
2 files+157−00/2 viewed
+42−0
1+import assert from "node:assert/strict";
2+import { test } from "node:test";
3+
4+import { StreamUsage, fromUsage, measure, total } from "./usage.ts";
5+
6+const sse = (events: object[]) => events.map((e) => `event: x\ndata: ${JSON.stringify(e)}\n\n`).join("");
7+
8+test("a stream's tokens come from message_start and the last message_delta", () => {
9+ const stream = sse([
10+ { type: "message_start", message: { usage: { input_tokens: 12, cache_read_input_tokens: 4000, cache_creation_input_tokens: 300, output_tokens: 1 } } },
11+ { type: "content_block_delta", delta: { text: "hi" } },
12+ { type: "message_delta", usage: { output_tokens: 40 } },
13+ { type: "message_delta", usage: { output_tokens: 95 } },
14+ { type: "message_stop" },
15+ ]);
16+ const usage = new StreamUsage();
17+ // Split mid-line, as chunks arrive.
18+ usage.push(stream.slice(0, 37));
19+ usage.push(stream.slice(37));
20+ assert.deepEqual(usage.finish(), { input: 12, output: 95, cacheRead: 4000, cacheWrite: 300 });
21+});
22+
23+test("a whole answer's usage, and nonsense counted as nothing", () => {
24+ assert.deepEqual(fromUsage({ input_tokens: 5, output_tokens: 7 }), { input: 5, output: 7, cacheRead: 0, cacheWrite: 0 });
25+ assert.equal(total(fromUsage({ input_tokens: -3, output_tokens: Number.NaN })), 0);
26+ const usage = new StreamUsage();
27+ usage.push("data: {not json\n\ndata: [DONE]\n");
28+ assert.equal(total(usage.finish()), 0);
29+});
30+
31+test("measuring passes the answer through unchanged", async () => {
32+ const body = sse([{ type: "message_start", message: { usage: { input_tokens: 3 } } }, { type: "message_delta", usage: { output_tokens: 9 } }]);
33+ const { response, tokens } = measure(new Response(body, { headers: { "content-type": "text/event-stream" } }));
34+ assert.equal(await response.text(), body);
35+ assert.deepEqual(await tokens, { input: 3, output: 9, cacheRead: 0, cacheWrite: 0 });
36+ const json = measure(new Response(JSON.stringify({ usage: { input_tokens: 2, output_tokens: 1 } }), { headers: { "content-type": "application/json" } }));
37+ assert.equal(JSON.parse(await json.response.text()).usage.input_tokens, 2);
38+ assert.equal(total(await json.tokens), 3);
39+ const failed = measure(new Response("no", { status: 500 }));
40+ assert.equal(failed.response.status, 500);
41+ assert.equal(total(await failed.tokens), 0);
42+});
+115−0
1+/**
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+
9+export 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+
18+export const NO_TOKENS: Tokens = { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 };
19+
20+type Usage = {
21+ input_tokens?: number;
22+ output_tokens?: number;
23+ cache_read_input_tokens?: number;
24+ cache_creation_input_tokens?: number;
25+};
26+
27+const count = (value: unknown): number => (typeof value === "number" && Number.isFinite(value) && value > 0 ? value : 0);
28+
29+/** The tokens of one `usage` object. */
30+export 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+
39+export 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+ */
48+export 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+ */
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 }) };
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+}