diff --git a/services/gateway/src/usage.ts b/services/gateway/src/usage.ts index 241a54c9..eb777579 100644 --- a/services/gateway/src/usage.ts +++ b/services/gateway/src/usage.ts @@ -112,12 +112,15 @@ export class UsageTee extends Transform { const parsed = JSON.parse(text); if (this.provider === "anthropic") { // Anthropic splits streamed usage across frames: `message_start` carries - // input and cache tokens under `message.usage`, `message_delta` carries - // output tokens under `delta.usage`, and a non-streamed body reports - // everything under a top-level `usage`. Merge each source separately - // rather than spreading them into one object: a spread copies absent - // keys as `undefined` and would drop an earlier frame's input tokens - // when the later frame reports only output. + // input and cache tokens under `message.usage`, and `message_delta` + // reports output tokens in a top-level `usage` — a sibling of `delta`, + // which itself carries only the stop fields. A non-streamed body also + // reports everything under a top-level `usage`. Merge each source + // separately rather than spreading them into one object: a spread + // copies absent keys as `undefined` and would drop an earlier frame's + // input tokens when the later frame reports only output. The + // `parsed.delta` source stays as tolerance for relays that nest usage + // inside the delta. for (const source of [parsed, parsed?.message, parsed?.delta]) { mergeUsage(this.usage, usageFromJson("anthropic", source)); } diff --git a/services/gateway/test/gateway.test.ts b/services/gateway/test/gateway.test.ts index 6a8574a1..cda30409 100644 --- a/services/gateway/test/gateway.test.ts +++ b/services/gateway/test/gateway.test.ts @@ -1310,7 +1310,7 @@ async function buildStub(state: StubState) { }); reply.raw.writeHead(200, { "content-type": "text/event-stream" }); reply.raw.write( - 'event: message_delta\ndata: {"type":"message_delta","delta":{"usage":{"output_tokens":333}}}\n\n', + 'event: message_delta\ndata: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":333}}\n\n', ); await new Promise((resolve) => setTimeout(resolve, 10_000)); return reply; @@ -1375,8 +1375,11 @@ async function buildStub(state: StubState) { function anthropicSseBody() { return [ - 'event: message_start\ndata: {"type":"message_start","message":{"usage":{"input_tokens":1000000}}}\n\n', + 'event: message_start\ndata: {"type":"message_start","message":{"id":"msg_stub","type":"message","role":"assistant","model":"claude-sonnet-5","content":[],"stop_reason":null,"usage":{"input_tokens":1000000,"output_tokens":1}}}\n\n', 'event: content_block_delta\ndata: {"type":"content_block_delta","delta":{"text":"ok"}}\n\n', - 'event: message_delta\ndata: {"type":"message_delta","delta":{"usage":{"output_tokens":1000000}}}\n\n', + // `usage` is a sibling of `delta`, and `delta` carries only the stop fields — + // see RawMessageDeltaEvent in @anthropic-ai/sdk. The null cache counters are + // what the API sends for an uncached request. + 'event: message_delta\ndata: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"input_tokens":null,"cache_creation_input_tokens":null,"cache_read_input_tokens":null,"output_tokens":1000000}}\n\n', ].join(""); } diff --git a/services/gateway/test/usage.test.ts b/services/gateway/test/usage.test.ts index 6cfed2e9..010ef13d 100644 --- a/services/gateway/test/usage.test.ts +++ b/services/gateway/test/usage.test.ts @@ -6,7 +6,7 @@ const MESSAGE_START = const CONTENT_DELTA = 'event: content_block_delta\ndata: {"type":"content_block_delta","delta":{"text":"ok"}}\n\n'; const MESSAGE_DELTA = - 'event: message_delta\ndata: {"type":"message_delta","delta":{"usage":{"output_tokens":700}}}\n\n'; + 'event: message_delta\ndata: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":700}}\n\n'; describe("UsageTee anthropic streaming usage", () => { it("meters input and cache tokens reported by message_start", async () => {