| 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; or OpenAI's, from |
| 5 | * a chat completion, a stream's chunks or an embeddings answer. The proxy |
| 6 | * adds these up per run for usage views, and the AI Gateway charges by |
| 7 | * them; billing still prices runs from AI Gateway. |
| 8 | */ |
| 9 | |
| 10 | export type Tokens = { |
| 11 | input: number; |
| 12 | output: number; |
| 13 | /** Prompt tokens read from the provider's cache. */ |
| 14 | cacheRead: number; |
| 15 | /** Prompt tokens written to the provider's cache, for either lifetime. */ |
| 16 | cacheWrite: number; |
| 17 | /** Of `cacheWrite`, those written to the hour-long cache. Absent when none. */ |
| 18 | cacheWrite1h?: number; |
| 19 | }; |
| 20 | |
| 21 | export const NO_TOKENS: Tokens = { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }; |
| 22 | |
| 23 | type Usage = { |
| 24 | input_tokens?: number; |
| 25 | output_tokens?: number; |
| 26 | cache_read_input_tokens?: number; |
| 27 | cache_creation_input_tokens?: number; |
| 28 | cache_creation?: { ephemeral_5m_input_tokens?: number; ephemeral_1h_input_tokens?: number } | null; |
| 29 | }; |
| 30 | |
| 31 | /** OpenAI's usage, from a chat completion or embeddings answer. */ |
| 32 | export type ChatUsage = { |
| 33 | prompt_tokens?: number; |
| 34 | completion_tokens?: number; |
| 35 | prompt_tokens_details?: { cached_tokens?: number } | null; |
| 36 | /** DeepSeek's name for cached prompt tokens. */ |
| 37 | prompt_cache_hit_tokens?: number; |
| 38 | }; |
| 39 | |
| 40 | const count = (value: unknown): number => (typeof value === "number" && Number.isFinite(value) && value > 0 ? value : 0); |
| 41 | |
| 42 | /** The tokens of one `usage` object. */ |
| 43 | export function fromUsage(usage: Usage | undefined | null): Tokens { |
| 44 | const tokens: Tokens = { |
| 45 | input: count(usage?.input_tokens), |
| 46 | output: count(usage?.output_tokens), |
| 47 | cacheRead: count(usage?.cache_read_input_tokens), |
| 48 | cacheWrite: count(usage?.cache_creation_input_tokens), |
| 49 | }; |
| 50 | const hour = count(usage?.cache_creation?.ephemeral_1h_input_tokens); |
| 51 | if (hour) tokens.cacheWrite1h = Math.min(hour, tokens.cacheWrite || hour); |
| 52 | if (hour && !tokens.cacheWrite) tokens.cacheWrite = hour + count(usage?.cache_creation?.ephemeral_5m_input_tokens); |
| 53 | return tokens; |
| 54 | } |
| 55 | |
| 56 | /** |
| 57 | * The tokens of one OpenAI `usage`: its prompt tokens less those read from |
| 58 | * the cache are input; the cached ones are cache reads. |
| 59 | */ |
| 60 | export function fromChatUsage(usage: ChatUsage | undefined | null): Tokens { |
| 61 | const prompt = count(usage?.prompt_tokens); |
| 62 | const cached = Math.min(prompt, count(usage?.prompt_tokens_details?.cached_tokens) || count(usage?.prompt_cache_hit_tokens)); |
| 63 | return { input: prompt - cached, output: count(usage?.completion_tokens), cacheRead: cached, cacheWrite: 0 }; |
| 64 | } |
| 65 | |
| 66 | export function total(tokens: Tokens): number { |
| 67 | return tokens.input + tokens.output + tokens.cacheRead + tokens.cacheWrite; |
| 68 | } |
| 69 | |
| 70 | /** |
| 71 | * Reads a server-sent event stream a chunk at a time. `message_start` |
| 72 | * carries the prompt's tokens; each `message_delta` carries the output so |
| 73 | * far (the last one is the total). |
| 74 | */ |
| 75 | export class StreamUsage { |
| 76 | tokens: Tokens = { ...NO_TOKENS }; |
| 77 | /** The model that answered, as `message_start` names it. */ |
| 78 | model: string | null = null; |
| 79 | private pending = ""; |
| 80 | |
| 81 | push(text: string): void { |
| 82 | this.pending += text; |
| 83 | let end = this.pending.indexOf("\n"); |
| 84 | while (end >= 0) { |
| 85 | this.line(this.pending.slice(0, end).trim()); |
| 86 | this.pending = this.pending.slice(end + 1); |
| 87 | end = this.pending.indexOf("\n"); |
| 88 | } |
| 89 | } |
| 90 | |
| 91 | finish(): Tokens { |
| 92 | this.line(this.pending.trim()); |
| 93 | this.pending = ""; |
| 94 | return this.tokens; |
| 95 | } |
| 96 | |
| 97 | private line(line: string): void { |
| 98 | if (!line.startsWith("data:")) return; |
| 99 | let event: { type?: string; message?: { usage?: Usage; model?: unknown }; usage?: Usage }; |
| 100 | try { |
| 101 | event = JSON.parse(line.slice(5).trim()); |
| 102 | } catch { |
| 103 | return; |
| 104 | } |
| 105 | if (event.type === "message_start") { |
| 106 | if (typeof event.message?.model === "string") this.model = event.message.model; |
| 107 | const start = fromUsage(event.message?.usage); |
| 108 | this.tokens = { ...start, output: Math.max(this.tokens.output, start.output) }; |
| 109 | } else if (event.type === "message_delta" && event.usage) { |
| 110 | const delta = fromUsage(event.usage); |
| 111 | this.tokens.output = Math.max(this.tokens.output, delta.output); |
| 112 | // Some providers repeat the prompt's tokens at the end. |
| 113 | if (delta.input) this.tokens.input = delta.input; |
| 114 | if (delta.cacheRead) this.tokens.cacheRead = delta.cacheRead; |
| 115 | if (delta.cacheWrite) this.tokens.cacheWrite = delta.cacheWrite; |
| 116 | if (delta.cacheWrite1h) this.tokens.cacheWrite1h = delta.cacheWrite1h; |
| 117 | } |
| 118 | } |
| 119 | } |
| 120 | |
| 121 | /** |
| 122 | * Reads a chat completion's stream a chunk at a time: the chunk that |
| 123 | * carries `usage` (the last, when the request asked for it) says it all. |
| 124 | */ |
| 125 | export class ChatStreamUsage { |
| 126 | tokens: Tokens = { ...NO_TOKENS }; |
| 127 | model: string | null = null; |
| 128 | private pending = ""; |
| 129 | |
| 130 | push(text: string): void { |
| 131 | this.pending += text; |
| 132 | let end = this.pending.indexOf("\n"); |
| 133 | while (end >= 0) { |
| 134 | this.line(this.pending.slice(0, end).trim()); |
| 135 | this.pending = this.pending.slice(end + 1); |
| 136 | end = this.pending.indexOf("\n"); |
| 137 | } |
| 138 | } |
| 139 | |
| 140 | finish(): Tokens { |
| 141 | this.line(this.pending.trim()); |
| 142 | this.pending = ""; |
| 143 | return this.tokens; |
| 144 | } |
| 145 | |
| 146 | private line(line: string): void { |
| 147 | if (!line.startsWith("data:")) return; |
| 148 | let chunk: { usage?: ChatUsage | null; model?: unknown }; |
| 149 | try { |
| 150 | chunk = JSON.parse(line.slice(5).trim()); |
| 151 | } catch { |
| 152 | return; |
| 153 | } |
| 154 | if (typeof chunk?.model === "string" && !this.model) this.model = chunk.model; |
| 155 | if (chunk?.usage) this.tokens = fromChatUsage(chunk.usage); |
| 156 | } |
| 157 | } |
| 158 | |
| 159 | /** |
| 160 | * The answer as it was, and a promise of what it used, read from a copy of |
| 161 | * its body, with the model that answered when it says. A failed answer, or |
| 162 | * one that is not a message, used nothing. |
| 163 | */ |
| 164 | export function measure( |
| 165 | answer: Response, |
| 166 | shape: "anthropic" | "openai" = "anthropic", |
| 167 | ): { response: Response; tokens: Promise<Tokens>; model: Promise<string | null> } { |
| 168 | if (!answer.ok || !answer.body) { |
| 169 | return { response: answer, tokens: Promise.resolve({ ...NO_TOKENS }), model: Promise.resolve(null) }; |
| 170 | } |
| 171 | const [passed, copy] = answer.body.tee(); |
| 172 | const response = new Response(passed, answer); |
| 173 | const streaming = (answer.headers.get("content-type") ?? "").includes("text/event-stream"); |
| 174 | let model: string | null = null; |
| 175 | const tokens = (async () => { |
| 176 | try { |
| 177 | if (!streaming) { |
| 178 | const whole = (await new Response(copy).json()) as { usage?: Usage & ChatUsage; model?: unknown }; |
| 179 | if (typeof whole.model === "string") model = whole.model; |
| 180 | return shape === "openai" ? fromChatUsage(whole.usage) : fromUsage(whole.usage); |
| 181 | } |
| 182 | const reader = copy.pipeThrough(new TextDecoderStream()).getReader(); |
| 183 | const usage = shape === "openai" ? new ChatStreamUsage() : new StreamUsage(); |
| 184 | for (;;) { |
| 185 | const { done, value } = await reader.read(); |
| 186 | if (done) break; |
| 187 | usage.push(value); |
| 188 | } |
| 189 | const used = usage.finish(); |
| 190 | model = usage.model; |
| 191 | return used; |
| 192 | } catch { |
| 193 | return { ...NO_TOKENS }; |
| 194 | } |
| 195 | })(); |
| 196 | return { response, tokens, model: tokens.then(() => model) }; |
| 197 | } |