From 0169725b4e1fd797149af8a335a6c6ee9f88e0f5 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 16:21:00 -0400 Subject: [PATCH 1/2] fix(llm): request streaming usage on openai_compatible providers Chat Completions only emits the trailing usage chunk when the request sets `stream_options: {"include_usage": true}`. The openai_compatible codec never sent it, so providers that follow the spec strictly returned no usage at all on streamed responses. Every message came back with zero tokens, and the catalog cost estimate multiplied those zeros into $0. Kimi is the visible case: a run's kimi-k3 stages report 0 tokens and no dollars, while an openrouter stage in the same run bills normally because OpenRouter volunteers usage (and an in-band cost) without being asked. Send the opt-in whenever we stream. Providers that already volunteer usage accept the field and are unaffected. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/codec/openai_compatible/request.rs | 54 +++++++++++++++++-- .../src/codec/openai_compatible/wire.rs | 10 ++++ ...tible__stream_text_happy_path_request.snap | 5 +- 3 files changed, 65 insertions(+), 4 deletions(-) diff --git a/lib/components/fabro-llm/src/codec/openai_compatible/request.rs b/lib/components/fabro-llm/src/codec/openai_compatible/request.rs index 12be9e84b..1ec6dc8fb 100644 --- a/lib/components/fabro-llm/src/codec/openai_compatible/request.rs +++ b/lib/components/fabro-llm/src/codec/openai_compatible/request.rs @@ -1,7 +1,7 @@ //! Request encoding: canonical `Request` → Chat Completions body. use super::translate; -use super::wire::{ApiRequest, ChatMessage}; +use super::wire::{ApiRequest, ChatMessage, StreamOptions}; use crate::codec::{CodecCtx, EncodedRequest, cache, merge_named_provider_options}; use crate::error::Error; @@ -10,8 +10,10 @@ use crate::error::Error; const KNOWN_OPTION_KEYS: &[&str] = &["auto_cache"]; /// Build the Chat Completions request for `ctx.request`. `stream` toggles the -/// `stream` body field. The body is assembled as a `serde_json::Value` so -/// `provider_options.` fields can be merged in before sending. +/// `stream` body field and the `stream_options.include_usage` opt-in that makes +/// providers emit the trailing usage chunk. The body is assembled as a +/// `serde_json::Value` so `provider_options.` fields can be +/// merged in before sending. /// /// Returns an error when the request contains a custom tool definition, which /// the Chat Completions tool envelope cannot represent. @@ -55,6 +57,9 @@ pub(super) fn encode(ctx: &CodecCtx<'_>, stream: bool) -> Result` merge runs after the body is built, so a + /// caller pointed at a gateway that rejects the field can still turn it + /// off. + #[test] + fn provider_options_can_override_stream_options() { + let mut request = minimal_request(); + request.provider_options = Some(serde_json::json!({ + "kimi": { "stream_options": serde_json::Value::Null } + })); + + let body = encode_body(&request, "kimi", true); + + assert_eq!(body["stream_options"], serde_json::Value::Null); } #[test] diff --git a/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs b/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs index 5bc0a3b93..31019a230 100644 --- a/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs +++ b/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs @@ -26,6 +26,16 @@ pub(super) struct ApiRequest { pub response_format: Option, #[serde(skip_serializing_if = "Option::is_none")] pub stream: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub stream_options: Option, +} + +/// Streaming options. Chat Completions only emits the trailing usage chunk +/// when the request opts in, so without this a streamed response reports zero +/// tokens and costs are estimated at $0. +#[derive(serde::Serialize)] +pub(super) struct StreamOptions { + pub include_usage: bool, } #[derive(serde::Serialize)] diff --git a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_text_happy_path_request.snap b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_text_happy_path_request.snap index fa711fab9..297b156fa 100644 --- a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_text_happy_path_request.snap +++ b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_text_happy_path_request.snap @@ -11,5 +11,8 @@ expression: rendered } ], "max_tokens": 128, - "stream": true + "stream": true, + "stream_options": { + "include_usage": true + } } From 53a92f75f4d6d10fb115fd791859133dfcd08d6c Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 17:47:05 -0400 Subject: [PATCH 2/2] test(llm): model streamed usage in OpenAI twin --- .../src/codec/openai_compatible/request.rs | 43 +------------------ .../tests/it/wire/openai_compatible.rs | 3 +- .../openai/src/openai/chat_completions.rs | 7 ++- test/twin/openai/src/openai/models.rs | 14 ++++++ test/twin/openai/src/sse.rs | 16 ++++++- .../openai/tests/chat_completions_contract.rs | 39 +++++++++++++++++ 6 files changed, 77 insertions(+), 45 deletions(-) diff --git a/lib/components/fabro-llm/src/codec/openai_compatible/request.rs b/lib/components/fabro-llm/src/codec/openai_compatible/request.rs index 1ec6dc8fb..f6bbf3b16 100644 --- a/lib/components/fabro-llm/src/codec/openai_compatible/request.rs +++ b/lib/components/fabro-llm/src/codec/openai_compatible/request.rs @@ -177,13 +177,10 @@ mod tests { tool_choice: None, response_format: None, stream: Some(true), - stream_options: Some(StreamOptions { - include_usage: true, - }), + stream_options: None, }; let json = serde_json::to_value(&req).unwrap(); assert_eq!(json["stream"], true); - assert_eq!(json["stream_options"]["include_usage"], true); let req_no_stream = ApiRequest { model: "test".into(), @@ -201,44 +198,6 @@ mod tests { }; let json_no_stream = serde_json::to_value(&req_no_stream).unwrap(); assert!(json_no_stream.get("stream").is_none()); - assert!(json_no_stream.get("stream_options").is_none()); - } - - /// Chat Completions only emits the trailing usage chunk when the request - /// opts in; without it streamed responses report zero tokens and cost - /// estimation silently produces $0. - #[test] - fn encode_opts_into_streaming_usage_when_streaming() { - let request = minimal_request(); - - let body = encode_body(&request, "kimi", true); - - assert_eq!(body["stream"], true); - assert_eq!(body["stream_options"]["include_usage"], true); - } - - #[test] - fn encode_omits_stream_options_when_not_streaming() { - let request = minimal_request(); - - let body = encode_body(&request, "kimi", false); - - assert!(body.get("stream_options").is_none()); - } - - /// The `provider_options.` merge runs after the body is built, so a - /// caller pointed at a gateway that rejects the field can still turn it - /// off. - #[test] - fn provider_options_can_override_stream_options() { - let mut request = minimal_request(); - request.provider_options = Some(serde_json::json!({ - "kimi": { "stream_options": serde_json::Value::Null } - })); - - let body = encode_body(&request, "kimi", true); - - assert_eq!(body["stream_options"], serde_json::Value::Null); } #[test] diff --git a/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs b/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs index c8facd920..4327176fc 100644 --- a/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs +++ b/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs @@ -560,7 +560,8 @@ async fn stream_text_happy_path_capture() -> (WireCapture, Vec, #[serde(default)] pub stream: bool, + stream_options: Option, pub tools: Option>, pub tool_choice: Option, pub response_format: Option, pub stop: Option, } +#[derive(Clone, Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct ChatStreamOptions { + #[serde(default)] + include_usage: bool, +} + impl ChatCompletionsRequest { + pub fn include_stream_usage(&self) -> bool { + self.stream_options + .as_ref() + .is_some_and(|options| options.include_usage) + } + pub fn extract_user_text(&self) -> String { let pieces: Vec = self .messages diff --git a/test/twin/openai/src/sse.rs b/test/twin/openai/src/sse.rs index eb49c3f47..2d30a0f88 100644 --- a/test/twin/openai/src/sse.rs +++ b/test/twin/openai/src/sse.rs @@ -236,7 +236,11 @@ pub fn responses_sse_response(plan: &ResponsePlan, transport: TransportOptions) stream_response(events, transport) } -pub fn chat_sse_response(plan: &ResponsePlan, transport: TransportOptions) -> Response { +pub fn chat_sse_response( + plan: &ResponsePlan, + include_usage: bool, + transport: TransportOptions, +) -> Response { let mut events = Vec::new(); let content = plan.chat_content(); events.push(chat_chunk(&json!({ @@ -321,6 +325,16 @@ pub fn chat_sse_response(plan: &ResponsePlan, transport: TransportOptions) -> Re "finish_reason": if plan.tool_calls.is_empty() { "stop" } else { "tool_calls" }, }] }))); + if include_usage { + events.push(chat_chunk(&json!({ + "id": format!("chatcmpl_{}", plan.id), + "object": "chat.completion.chunk", + "created": plan.created, + "model": plan.model, + "choices": [], + "usage": plan.usage.chat_completions_json(), + }))); + } events.push("data: [DONE]\n\n".to_owned()); } diff --git a/test/twin/openai/tests/chat_completions_contract.rs b/test/twin/openai/tests/chat_completions_contract.rs index 6fe4c27fe..e17369c42 100644 --- a/test/twin/openai/tests/chat_completions_contract.rs +++ b/test/twin/openai/tests/chat_completions_contract.rs @@ -50,9 +50,48 @@ async fn chat_completions_stream_uses_same_canonical_plan() { assert_eq!(status, 200); assert!(joined.contains("\"content\":\"deterministic: stream same plan\"")); + assert!(!joined.contains("\"usage\"")); assert!(joined.contains("data: [DONE]")); } +#[tokio::test] +async fn chat_completions_stream_includes_usage_when_requested() { + let server = common::spawn_server().await.expect("server should start"); + + let (status, chunks) = server + .post_chat_stream(json!({ + "model": "gpt-test", + "messages": [{ "role": "user", "content": "stream with usage" }], + "stream": true, + "stream_options": { "include_usage": true } + })) + .await; + + assert_eq!(status, 200); + let transcript = + common::parse_sse_transcript(chunks.join("").as_bytes()).expect("valid SSE transcript"); + let usage_chunk = transcript + .events + .iter() + .filter(|event| event.data != "[DONE]") + .map(|event| { + serde_json::from_str::(&event.data).expect("valid JSON chunk") + }) + .find(|chunk| chunk.get("usage").is_some()) + .expect("trailing usage chunk"); + + assert_eq!(usage_chunk["choices"], json!([])); + assert_eq!( + usage_chunk["usage"], + json!({ + "prompt_tokens": 3, + "completion_tokens": 5, + "total_tokens": 8 + }) + ); + assert!(transcript.done); +} + #[tokio::test] async fn chat_completions_accepts_supported_openai_compatible_fields() { let server = common::spawn_server().await.expect("server should start");