Skip to main content

llm/providers/openai_responses/
streaming.rs

1use async_openai::types::responses::{OutputItem, ResponseUsage, Status};
2use futures::Stream;
3use serde::Deserialize;
4use tokio_stream::StreamExt;
5
6use crate::providers::tool_call_collector::ToolCallCollector;
7use crate::{LlmError, LlmResponse, Result, StopReason, TokenUsage};
8
9impl From<ResponseUsage> for TokenUsage {
10    fn from(usage: ResponseUsage) -> Self {
11        TokenUsage {
12            input_tokens: usage.input_tokens,
13            output_tokens: usage.output_tokens,
14            cache_read_tokens: Some(usage.input_tokens_details.cached_tokens),
15            reasoning_tokens: Some(usage.output_tokens_details.reasoning_tokens),
16            ..TokenUsage::default()
17        }
18    }
19}
20
21#[derive(Debug, Deserialize)]
22#[serde(tag = "type")]
23pub enum ResponsesStreamEvent {
24    #[serde(rename = "response.created")]
25    Created(ResponsesCreatedEvent),
26    #[serde(rename = "response.output_text.delta")]
27    OutputTextDelta(ResponsesTextDeltaEvent),
28    #[serde(rename = "response.output_item.added")]
29    OutputItemAdded(ResponsesOutputItemEvent),
30    #[serde(rename = "response.output_item.done")]
31    OutputItemDone(ResponsesOutputItemEvent),
32    #[serde(rename = "response.function_call_arguments.delta")]
33    FunctionCallArgumentsDelta(ResponsesFunctionCallArgumentsDeltaEvent),
34    #[serde(rename = "response.function_call_arguments.done")]
35    FunctionCallArgumentsDone(ResponsesFunctionCallArgumentsDoneEvent),
36    #[serde(rename = "response.reasoning_summary_text.delta")]
37    ReasoningSummaryTextDelta(ResponsesTextDeltaEvent),
38    #[serde(rename = "response.completed")]
39    Completed(ResponsesCompletedEvent),
40    #[serde(rename = "response.incomplete")]
41    Incomplete(ResponsesCompletedEvent),
42    #[serde(rename = "response.failed")]
43    Failed(ResponsesFailedEvent),
44    #[serde(rename = "error")]
45    Error(ResponsesErrorEvent),
46    #[serde(other)]
47    Ignored,
48}
49
50impl ResponsesStreamEvent {
51    /// Whether this event may legitimately arrive before `response.created`:
52    /// event types we ignore, and failures the endpoint reports *instead of*
53    /// opening a response. Rejecting those would replace the server's own
54    /// message with a generic interrupt.
55    fn may_precede_creation(&self) -> bool {
56        matches!(self, Self::Ignored | Self::Error(_) | Self::Failed(_))
57    }
58}
59
60#[derive(Debug, Deserialize)]
61pub struct ResponsesCreatedEvent {
62    pub response: ResponsesCreated,
63}
64
65#[derive(Debug, Deserialize)]
66pub struct ResponsesCreated {
67    pub id: String,
68}
69
70#[derive(Debug, Deserialize)]
71pub struct ResponsesFailedEvent {
72    pub response: ResponsesFailed,
73}
74
75#[derive(Debug, Deserialize)]
76pub struct ResponsesFailed {
77    #[serde(default)]
78    pub error: Option<ResponsesErrorEvent>,
79}
80
81#[derive(Debug, Deserialize)]
82pub struct ResponsesTextDeltaEvent {
83    pub delta: String,
84}
85
86#[derive(Debug, Deserialize)]
87pub struct ResponsesOutputItemEvent {
88    pub output_index: u32,
89    pub item: OutputItem,
90}
91
92#[derive(Debug, Deserialize)]
93pub struct ResponsesFunctionCallArgumentsDeltaEvent {
94    pub output_index: u32,
95    pub delta: String,
96}
97
98#[derive(Debug, Deserialize)]
99pub struct ResponsesFunctionCallArgumentsDoneEvent {
100    pub output_index: u32,
101}
102
103#[derive(Debug, Deserialize)]
104pub struct ResponsesCompletedEvent {
105    pub response: ResponsesCompleted,
106}
107
108#[derive(Debug, Deserialize)]
109pub struct ResponsesCompleted {
110    #[serde(default)]
111    pub usage: Option<ResponseUsage>,
112    #[serde(default)]
113    pub status: Option<Status>,
114}
115
116#[derive(Debug, Deserialize)]
117pub struct ResponsesErrorEvent {
118    pub message: String,
119}
120
121/// Process an `OpenAI` Responses event stream into `LlmResponse` items.
122pub fn process_response_stream<T>(stream: T) -> impl Stream<Item = Result<LlmResponse>> + Send
123where
124    T: Stream<Item = Result<ResponsesStreamEvent>> + Send + Unpin,
125{
126    async_stream::stream! {
127        let mut tool_collector = ToolCallCollector::<u32>::new();
128        let mut stream = Box::pin(stream);
129        let mut last_stop_reason: Option<StopReason> = None;
130        let mut started = false;
131        let mut terminal = false;
132        let mut failed = false;
133
134        while let Some(result) = stream.next().await {
135            let event = match result {
136                Ok(event) => event,
137                Err(e) => {
138                    yield Err(LlmError::StreamInterrupted(e.to_string()));
139                    failed = true;
140                    break;
141                }
142            };
143
144            if matches!(event, ResponsesStreamEvent::Created(_)) {
145                started = true;
146            } else if !started && !event.may_precede_creation() {
147                yield Err(LlmError::StreamInterrupted(
148                    "Responses stream emitted data before response.created".to_string(),
149                ));
150                failed = true;
151                break;
152            }
153
154            terminal = matches!(event, ResponsesStreamEvent::Completed(_) | ResponsesStreamEvent::Incomplete(_));
155            let responses = process_event(event, &mut tool_collector, &mut last_stop_reason);
156            let event_failed = responses.iter().any(Result::is_err);
157            for response in responses {
158                yield response;
159            }
160            if event_failed || terminal {
161                failed = event_failed;
162                break;
163            }
164        }
165
166        if !failed {
167            for tc in tool_collector.complete_all() {
168                yield Ok(LlmResponse::ToolRequestComplete { tool_call: tc });
169            }
170
171            if terminal {
172                yield Ok(LlmResponse::Done { stop_reason: last_stop_reason });
173            } else {
174                yield Err(LlmError::StreamInterrupted(
175                    "Responses stream ended before a terminal response event".to_string(),
176                ));
177            }
178        }
179    }
180}
181
182fn process_event(
183    event: ResponsesStreamEvent,
184    tool_collector: &mut ToolCallCollector<u32>,
185    last_stop_reason: &mut Option<StopReason>,
186) -> Vec<Result<LlmResponse>> {
187    let mut responses = Vec::new();
188    let incomplete = matches!(&event, ResponsesStreamEvent::Incomplete(_));
189
190    match event {
191        ResponsesStreamEvent::Created(e) => {
192            responses.push(Ok(LlmResponse::Start { message_id: e.response.id }));
193        }
194        ResponsesStreamEvent::OutputTextDelta(e) if !e.delta.is_empty() => {
195            responses.push(Ok(LlmResponse::Text { chunk: e.delta }));
196        }
197        ResponsesStreamEvent::OutputItemAdded(e) => {
198            if let OutputItem::FunctionCall(call) = e.item {
199                let tool_responses = tool_collector.handle_delta(e.output_index, call.id, Some(call.name), None);
200                responses.extend(tool_responses.into_iter().map(Ok));
201            }
202        }
203        ResponsesStreamEvent::FunctionCallArgumentsDelta(e) => {
204            let tool_responses = tool_collector.handle_delta(e.output_index, None, None, Some(e.delta));
205            responses.extend(tool_responses.into_iter().map(Ok));
206        }
207        ResponsesStreamEvent::FunctionCallArgumentsDone(e) => {
208            if let Some(tc) = tool_collector.complete_one(e.output_index) {
209                responses.push(Ok(LlmResponse::ToolRequestComplete { tool_call: tc }));
210            }
211        }
212        ResponsesStreamEvent::ReasoningSummaryTextDelta(e) if !e.delta.is_empty() => {
213            responses.push(Ok(LlmResponse::Reasoning { chunk: e.delta }));
214        }
215        ResponsesStreamEvent::OutputItemDone(e) => {
216            if let OutputItem::Reasoning(reasoning) = e.item
217                && let Some(id) = reasoning.id
218                && let Some(encrypted) = reasoning.encrypted_content
219            {
220                responses.push(Ok(LlmResponse::EncryptedReasoning { id, content: encrypted }));
221            }
222        }
223        ResponsesStreamEvent::Completed(e) | ResponsesStreamEvent::Incomplete(e) => {
224            if let Some(usage) = e.response.usage {
225                responses.push(Ok(LlmResponse::Usage { tokens: usage.into() }));
226            }
227            match e.response.status {
228                Some(Status::Completed) => *last_stop_reason = Some(StopReason::EndTurn),
229                Some(Status::Incomplete) => *last_stop_reason = Some(StopReason::Length),
230                _ if incomplete => {
231                    *last_stop_reason = Some(StopReason::Length);
232                }
233                _ => {}
234            }
235        }
236        ResponsesStreamEvent::Failed(e) => {
237            let message = e.response.error.map_or_else(|| "Unknown Responses API failure".to_string(), |e| e.message);
238            responses.push(Err(LlmError::ApiError(message)));
239        }
240        ResponsesStreamEvent::Error(e) => {
241            responses.push(Err(LlmError::ServerError {
242                status: None,
243                message: format!("Responses API error: {}", e.message),
244            }));
245        }
246        ResponsesStreamEvent::Ignored
247        | ResponsesStreamEvent::OutputTextDelta(_)
248        | ResponsesStreamEvent::ReasoningSummaryTextDelta(_) => {}
249    }
250
251    responses
252}
253
254#[cfg(test)]
255mod tests {
256    use super::*;
257    use crate::TokenUsage;
258    use async_openai::types::responses::{FunctionToolCall, ReasoningItem};
259    use serde_json::json;
260
261    async fn collect_responses(events: Vec<ResponsesStreamEvent>) -> Vec<LlmResponse> {
262        let stream = make_stream(events);
263        let mut response_stream = Box::pin(process_response_stream(stream));
264        let mut responses = Vec::new();
265        while let Some(result) = response_stream.next().await {
266            responses.push(result.unwrap());
267        }
268        responses
269    }
270
271    #[tokio::test]
272    async fn test_text_stream() {
273        let responses = collect_responses(vec![
274            text_delta("Hello"),
275            text_delta(" world"),
276            completed(Status::Completed, Some(make_usage(10, 5))),
277        ])
278        .await;
279
280        assert!(matches!(responses[0], LlmResponse::Start { .. }));
281        assert!(matches!(responses[1], LlmResponse::Text { ref chunk } if chunk == "Hello"));
282        assert!(matches!(responses[2], LlmResponse::Text { ref chunk } if chunk == " world"));
283        assert!(matches!(
284            responses[3],
285            LlmResponse::Usage { tokens: TokenUsage { input_tokens: 10, output_tokens: 5, .. } }
286        ));
287        assert!(matches!(responses[4], LlmResponse::Done { stop_reason: Some(StopReason::EndTurn) }));
288    }
289
290    #[tokio::test]
291    async fn test_tool_call_stream() {
292        let responses = collect_responses(vec![
293            ResponsesStreamEvent::OutputItemAdded(ResponsesOutputItemEvent {
294                output_index: 0,
295                item: OutputItem::FunctionCall(FunctionToolCall {
296                    id: Some("fc_1".to_string()),
297                    call_id: "call_1".to_string(),
298                    name: "read_file".to_string(),
299                    arguments: String::new(),
300                    status: None,
301                    namespace: None,
302                }),
303            }),
304            function_call_delta(r#"{"path":"#),
305            function_call_delta(r#""foo.rs"}"#),
306            ResponsesStreamEvent::FunctionCallArgumentsDone(ResponsesFunctionCallArgumentsDoneEvent {
307                output_index: 0,
308            }),
309            completed(Status::Completed, Some(make_usage(20, 10))),
310        ])
311        .await;
312
313        assert!(matches!(responses[0], LlmResponse::Start { .. }));
314        assert!(
315            matches!(&responses[1], LlmResponse::ToolRequestStart { id, name } if id == "fc_1" && name == "read_file")
316        );
317        assert!(matches!(responses[2], LlmResponse::ToolRequestArg { .. }));
318        assert!(matches!(responses[3], LlmResponse::ToolRequestArg { .. }));
319
320        let tc = responses.iter().find(|r| matches!(r, LlmResponse::ToolRequestComplete { .. }));
321        assert!(tc.is_some());
322        if let LlmResponse::ToolRequestComplete { tool_call } = tc.unwrap() {
323            assert_eq!(tool_call.id, "fc_1");
324            assert_eq!(tool_call.name, "read_file");
325            assert_eq!(tool_call.arguments, r#"{"path":"foo.rs"}"#);
326        }
327    }
328
329    #[tokio::test]
330    async fn test_error_event_is_retryable_server_error() {
331        let stream = make_stream(vec![ResponsesStreamEvent::Error(ResponsesErrorEvent {
332            message: "Rate limit exceeded".to_string(),
333        })]);
334        let mut response_stream = Box::pin(process_response_stream(stream));
335
336        let mut responses = Vec::new();
337        while let Some(result) = response_stream.next().await {
338            responses.push(result);
339        }
340
341        assert!(responses[0].is_ok());
342        let err = responses[1].as_ref().expect_err("expected error event to surface as Err");
343        assert!(matches!(err, LlmError::ServerError { status: None, .. }), "got {err:?}");
344        assert!(err.is_retryable(), "ResponseError must be retryable so the agent can recover");
345    }
346
347    #[tokio::test]
348    async fn test_reasoning_delta() {
349        let responses = collect_responses(vec![
350            reasoning_delta("Thinking about"),
351            reasoning_delta(" the problem"),
352            completed(Status::Completed, None),
353        ])
354        .await;
355
356        assert!(matches!(responses[1], LlmResponse::Reasoning { ref chunk } if chunk == "Thinking about"));
357        assert!(matches!(responses[2], LlmResponse::Reasoning { ref chunk } if chunk == " the problem"));
358    }
359
360    #[tokio::test]
361    async fn test_incomplete_status_gives_length_stop_reason() {
362        let responses = collect_responses(vec![completed(Status::Incomplete, None)]).await;
363
364        assert!(matches!(responses.last().unwrap(), LlmResponse::Done { stop_reason: Some(StopReason::Length) }));
365    }
366
367    #[tokio::test]
368    async fn test_stream_error_propagation_is_retryable() {
369        let events: Vec<Result<ResponsesStreamEvent>> =
370            vec![Err(LlmError::StreamInterrupted("connection lost".to_string()))];
371
372        let stream = tokio_stream::iter(events);
373        let mut response_stream = Box::pin(process_response_stream(stream));
374
375        let mut responses = Vec::new();
376        while let Some(result) = response_stream.next().await {
377            responses.push(result);
378        }
379
380        let err = responses[0].as_ref().expect_err("expected upstream Err to surface as Err");
381        assert!(matches!(err, LlmError::StreamInterrupted(_)), "got {err:?}");
382        assert_eq!(responses.len(), 1);
383        assert!(err.is_retryable(), "mid-stream interrupts must be retryable");
384    }
385
386    #[tokio::test]
387    async fn error_event_before_creation_keeps_the_servers_message() {
388        let events =
389            vec![Ok(ResponsesStreamEvent::Error(ResponsesErrorEvent { message: "Rate limit exceeded".to_string() }))];
390        let responses = process_response_stream(tokio_stream::iter(events)).collect::<Vec<_>>().await;
391
392        let err = responses[0].as_ref().expect_err("expected the error event to surface as Err");
393        assert!(matches!(err, LlmError::ServerError { .. }), "got {err:?}");
394        assert!(err.to_string().contains("Rate limit exceeded"), "server message was dropped: {err}");
395    }
396
397    #[tokio::test]
398    async fn failure_event_before_creation_keeps_the_servers_message() {
399        let events = vec![Ok(ResponsesStreamEvent::Failed(ResponsesFailedEvent {
400            response: ResponsesFailed { error: Some(ResponsesErrorEvent { message: "model overloaded".to_string() }) },
401        }))];
402        let responses = process_response_stream(tokio_stream::iter(events)).collect::<Vec<_>>().await;
403
404        let err = responses[0].as_ref().expect_err("expected the failure event to surface as Err");
405        assert!(matches!(err, LlmError::ApiError(_)), "got {err:?}");
406        assert!(err.to_string().contains("model overloaded"), "server message was dropped: {err}");
407    }
408
409    #[tokio::test]
410    async fn data_before_creation_is_interrupted() {
411        let events = vec![Ok(text_delta("leaked"))];
412        let responses = process_response_stream(tokio_stream::iter(events)).collect::<Vec<_>>().await;
413
414        assert!(matches!(responses[0], Err(LlmError::StreamInterrupted(_))), "{responses:?}");
415    }
416
417    #[tokio::test]
418    async fn stream_without_terminal_event_is_interrupted() {
419        let stream = make_stream(vec![text_delta("partial")]);
420        let responses = process_response_stream(stream).collect::<Vec<_>>().await;
421
422        assert!(matches!(responses[0], Ok(LlmResponse::Start { .. })));
423        assert!(matches!(responses[1], Ok(LlmResponse::Text { .. })));
424        assert!(matches!(responses[2], Err(LlmError::StreamInterrupted(_))));
425        assert!(!responses.iter().any(|response| matches!(response, Ok(LlmResponse::Done { .. }))));
426    }
427
428    #[tokio::test]
429    async fn captured_responses_fixture_uses_the_shared_processor() {
430        let responses = process_fixture(include_str!("../../../tests/fixtures/openai_responses/01_minimal.sse")).await;
431
432        assert!(responses.iter().all(Result::is_ok), "{responses:?}");
433        let usage = fixture_usage(&responses).expect("fixture should report usage");
434        assert!(usage.input_tokens > 0, "input_tokens should be > 0: {usage:?}");
435        assert!(usage.output_tokens > 0, "output_tokens should be > 0: {usage:?}");
436        assert!(matches!(responses.last(), Some(Ok(LlmResponse::Done { stop_reason: Some(StopReason::EndTurn) }))));
437    }
438
439    #[tokio::test]
440    async fn captured_reasoning_fixture_preserves_reasoning_usage() {
441        let responses =
442            process_fixture(include_str!("../../../tests/fixtures/openai_responses/02_reasoning.sse")).await;
443
444        assert!(responses.iter().all(Result::is_ok), "{responses:?}");
445        let usage = fixture_usage(&responses).expect("fixture should report usage");
446        assert!(usage.input_tokens > 0, "input_tokens should be > 0: {usage:?}");
447        assert!(usage.output_tokens > 0, "output_tokens should be > 0: {usage:?}");
448        assert!(usage.reasoning_tokens.is_some_and(|tokens| tokens > 0), "{usage:?}");
449    }
450
451    /// Decode a captured SSE body and run it through the shared processor.
452    async fn process_fixture(sse: &str) -> Vec<Result<LlmResponse>> {
453        let events = sse
454            .lines()
455            .filter_map(|line| line.strip_prefix("data: "))
456            .filter(|data| *data != "[DONE]")
457            .map(|data| serde_json::from_str::<ResponsesStreamEvent>(data).map_err(LlmError::from));
458        process_response_stream(tokio_stream::iter(events)).collect::<Vec<_>>().await
459    }
460
461    fn fixture_usage(responses: &[Result<LlmResponse>]) -> Option<TokenUsage> {
462        responses.iter().find_map(|response| match response {
463            Ok(LlmResponse::Usage { tokens }) => Some(*tokens),
464            _ => None,
465        })
466    }
467
468    #[test]
469    fn test_encrypted_reasoning_from_output_item_done() {
470        let event = ResponsesStreamEvent::OutputItemDone(ResponsesOutputItemEvent {
471            output_index: 0,
472            item: reasoning_item(Some("enc-blob-data")),
473        });
474
475        let mut tool_collector = ToolCallCollector::<u32>::new();
476        let mut stop_reason = None;
477        let responses = process_event(event, &mut tool_collector, &mut stop_reason);
478
479        assert_eq!(responses.len(), 1);
480        assert!(
481            matches!(&responses[0], Ok(LlmResponse::EncryptedReasoning { content, .. }) if content == "enc-blob-data")
482        );
483    }
484
485    #[tokio::test]
486    async fn test_usage_forwards_reasoning_and_cache_read() {
487        let responses =
488            collect_responses(vec![completed(Status::Completed, Some(make_usage_full(120, 80, 50, 30)))]).await;
489
490        let usage = responses.iter().find_map(|r| match r {
491            LlmResponse::Usage { tokens } => Some(*tokens),
492            _ => None,
493        });
494
495        assert_eq!(
496            usage,
497            Some(TokenUsage {
498                input_tokens: 120,
499                output_tokens: 80,
500                cache_read_tokens: Some(50),
501                reasoning_tokens: Some(30),
502                ..TokenUsage::default()
503            })
504        );
505    }
506
507    #[tokio::test]
508    async fn test_completed_without_output_deserializes_usage_and_stop_reason() {
509        let event: ResponsesStreamEvent = serde_json::from_value(json!({
510            "type": "response.completed",
511            "sequence_number": 1,
512            "response": {
513                "id": "resp_1",
514                "object": "response",
515                "created_at": 1_000_u64,
516                "status": "completed",
517                "background": false,
518                "completed_at": 2_000_u64,
519                "error": null,
520                "model": "test-model",
521                "usage": make_usage_json(100, 20, 0, 10)
522            }
523        }))
524        .unwrap();
525        let responses = collect_responses(vec![event]).await;
526
527        assert!(matches!(
528            responses.iter().find(|response| matches!(response, LlmResponse::Usage { .. })),
529            Some(LlmResponse::Usage {
530                tokens: TokenUsage { input_tokens: 100, output_tokens: 20, reasoning_tokens: Some(10), .. }
531            })
532        ));
533        assert!(matches!(responses.last().unwrap(), LlmResponse::Done { stop_reason: Some(StopReason::EndTurn) }));
534    }
535
536    #[test]
537    fn test_output_item_done_without_encrypted_content_is_ignored() {
538        let event = ResponsesStreamEvent::OutputItemDone(ResponsesOutputItemEvent {
539            output_index: 0,
540            item: reasoning_item(None),
541        });
542
543        let mut tool_collector = ToolCallCollector::<u32>::new();
544        let mut stop_reason = None;
545        let responses = process_event(event, &mut tool_collector, &mut stop_reason);
546
547        assert!(responses.is_empty());
548    }
549
550    fn text_delta(delta: &str) -> ResponsesStreamEvent {
551        ResponsesStreamEvent::OutputTextDelta(ResponsesTextDeltaEvent { delta: delta.to_string() })
552    }
553
554    fn reasoning_delta(delta: &str) -> ResponsesStreamEvent {
555        ResponsesStreamEvent::ReasoningSummaryTextDelta(ResponsesTextDeltaEvent { delta: delta.to_string() })
556    }
557
558    fn function_call_delta(delta: &str) -> ResponsesStreamEvent {
559        ResponsesStreamEvent::FunctionCallArgumentsDelta(ResponsesFunctionCallArgumentsDeltaEvent {
560            output_index: 0,
561            delta: delta.to_string(),
562        })
563    }
564
565    fn completed(status: Status, usage: Option<ResponseUsage>) -> ResponsesStreamEvent {
566        ResponsesStreamEvent::Completed(ResponsesCompletedEvent {
567            response: ResponsesCompleted { usage, status: Some(status) },
568        })
569    }
570
571    fn reasoning_item(encrypted_content: Option<&str>) -> OutputItem {
572        OutputItem::Reasoning(ReasoningItem {
573            id: Some("r_1".to_string()),
574            summary: vec![],
575            encrypted_content: encrypted_content.map(ToString::to_string),
576            content: None,
577            status: None,
578        })
579    }
580
581    fn make_stream(
582        events: Vec<ResponsesStreamEvent>,
583    ) -> impl Stream<Item = Result<ResponsesStreamEvent>> + Send + Unpin {
584        tokio_stream::iter(
585            std::iter::once(Ok(ResponsesStreamEvent::Created(ResponsesCreatedEvent {
586                response: ResponsesCreated { id: "resp_test".to_string() },
587            })))
588            .chain(events.into_iter().map(Ok))
589            .collect::<Vec<_>>(),
590        )
591    }
592
593    fn make_usage(input_tokens: u32, output_tokens: u32) -> ResponseUsage {
594        make_usage_full(input_tokens, output_tokens, 0, 0)
595    }
596
597    fn make_usage_full(
598        input_tokens: u32,
599        output_tokens: u32,
600        cached_tokens: u32,
601        reasoning_tokens: u32,
602    ) -> ResponseUsage {
603        serde_json::from_value(make_usage_json(input_tokens, output_tokens, cached_tokens, reasoning_tokens)).unwrap()
604    }
605
606    fn make_usage_json(
607        input_tokens: u32,
608        output_tokens: u32,
609        cached_tokens: u32,
610        reasoning_tokens: u32,
611    ) -> serde_json::Value {
612        json!({
613            "input_tokens": input_tokens,
614            "input_tokens_details": { "cached_tokens": cached_tokens },
615            "output_tokens": output_tokens,
616            "output_tokens_details": { "reasoning_tokens": reasoning_tokens },
617            "total_tokens": input_tokens + output_tokens
618        })
619    }
620}