Skip to main content

nanocodex_oai_api/session/
response.rs

1use std::{
2    convert::Infallible,
3    error::Error,
4    fmt,
5    future::{Future, IntoFuture},
6    pin::Pin,
7    sync::Arc,
8    task::{Context, Poll},
9};
10
11use ::tower::Service;
12use futures_util::{Stream, future::poll_fn};
13use tokio::sync::mpsc;
14use tracing::{Instrument, info_span};
15use web_time::Instant;
16
17use crate::{
18    ContentItem, EventSink, MessageRole, ResponseEvent, ResponseItem, ResponsesAttempt,
19    ResponsesAttemptFactory, ResponsesOutput, ResponsesServiceError, ResponsesServiceResponse,
20    Usage,
21};
22
23use super::{
24    builder::{ResponseTurn, Session},
25    compaction,
26    context::{assign_missing_response_item_ids, is_canonical_context_item},
27};
28
29/// Typed input accepted by `response.create`.
30#[derive(Clone, Debug)]
31pub struct ResponseInput {
32    items: Vec<ResponseItem>,
33}
34
35impl ResponseInput {
36    /// Creates one user message from ordered typed content.
37    ///
38    /// ```
39    /// use nanocodex_oai_api::{
40    ///     ImageDetail,
41    ///     responses::ContentItem,
42    ///     session::ResponseInput,
43    /// };
44    ///
45    /// let _input = ResponseInput::content([
46    ///     ContentItem::input_text("Describe the deployment diagram."),
47    ///     ContentItem::input_image_with_detail(
48    ///         "https://example.com/deployment-diagram.png",
49    ///         ImageDetail::High,
50    ///     ),
51    /// ]);
52    /// ```
53    #[must_use]
54    pub fn content(content: impl IntoIterator<Item = ContentItem>) -> Self {
55        Self {
56            items: vec![ResponseItem::message(MessageRole::User, content)],
57        }
58    }
59
60    /// Estimates model-visible tokens contributed by this input.
61    ///
62    /// This uses the same image, text, and tool-output accounting as managed
63    /// context. A caller can combine it with
64    /// [`ResponseTurn::active_context_tokens`] and choose its own compaction
65    /// margin before calling [`ResponseTurn::create`].
66    #[must_use]
67    pub fn estimated_tokens(&self) -> u64 {
68        self.items
69            .iter()
70            .map(compaction::estimate_item_tokens)
71            .fold(0, u64::saturating_add)
72    }
73
74    /// Creates ordered protocol-level input items.
75    ///
76    /// Use this after a completed response contains tool calls and the
77    /// application executes those calls itself. The output item must retain
78    /// the exact call ID returned by the model.
79    ///
80    /// ```
81    /// use nanocodex_oai_api::{
82    ///     responses::{FunctionOutputBody, ResponseItem},
83    ///     session::ResponseInput,
84    /// };
85    ///
86    /// let input = ResponseInput::items([ResponseItem::function_call_output(
87    ///     "call_region_01".to_owned(),
88    ///     FunctionOutputBody::Text(
89    ///         r#"{"region":"iad","status":"healthy"}"#.into(),
90    ///     ),
91    /// )]);
92    ///
93    /// assert!(input.estimated_tokens() > 0);
94    /// ```
95    #[must_use]
96    pub fn items(items: impl IntoIterator<Item = ResponseItem>) -> Self {
97        Self {
98            items: items.into_iter().collect(),
99        }
100    }
101}
102
103impl From<String> for ResponseInput {
104    fn from(text: String) -> Self {
105        Self::content([ContentItem::InputText {
106            text: text.into_boxed_str(),
107        }])
108    }
109}
110
111impl From<&str> for ResponseInput {
112    fn from(text: &str) -> Self {
113        Self::from(text.to_owned())
114    }
115}
116
117/// A completed and atomically committed Responses operation.
118#[derive(Clone)]
119pub struct CompletedResponse {
120    output: Arc<[ResponseItem]>,
121    output_text: Arc<str>,
122    usage: Option<Usage>,
123    estimated_cost: Option<crate::EstimatedUsdCost>,
124    cost_status: crate::CostStatus,
125    end_turn: Option<bool>,
126}
127
128impl CompletedResponse {
129    /// Returns every completed output item in provider order.
130    #[must_use]
131    pub fn output(&self) -> &[ResponseItem] {
132        &self.output
133    }
134
135    /// Returns concatenated assistant output text.
136    #[must_use]
137    pub fn output_text(&self) -> &str {
138        &self.output_text
139    }
140
141    /// Iterates over complete function and custom tool calls.
142    pub fn tool_calls(&self) -> impl Iterator<Item = &ResponseItem> {
143        self.output.iter().filter(|item| {
144            matches!(
145                item,
146                ResponseItem::FunctionCall { .. }
147                    | ResponseItem::CustomToolCall { .. }
148                    | ResponseItem::LocalShellCall { .. }
149                    | ResponseItem::ToolSearchCall { .. }
150            )
151        })
152    }
153
154    /// Returns token usage for this API operation.
155    #[must_use]
156    pub const fn usage(&self) -> Option<&Usage> {
157        self.usage.as_ref()
158    }
159
160    /// Returns the automatic local USD estimate.
161    ///
162    /// Nanocodex applies the built-in standard or priority
163    /// [`crate::MODEL`] rates. `None` means the provider omitted usage.
164    #[must_use]
165    pub const fn estimated_cost(&self) -> Option<&crate::EstimatedUsdCost> {
166        self.estimated_cost.as_ref()
167    }
168
169    /// Returns why an estimate is present or unavailable.
170    #[must_use]
171    pub const fn cost_status(&self) -> crate::CostStatus {
172        self.cost_status
173    }
174
175    /// Returns whether the model affirmatively ended its logical turn.
176    #[must_use]
177    pub const fn end_turn(&self) -> Option<bool> {
178        self.end_turn
179    }
180}
181
182impl fmt::Debug for CompletedResponse {
183    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
184        formatter
185            .debug_struct("CompletedResponse")
186            .field("output_items", &self.output.len())
187            .field("output_text", &self.output_text)
188            .field("end_turn", &self.end_turn)
189            .finish_non_exhaustive()
190    }
191}
192
193/// A completed and installed remote compaction.
194#[derive(Clone)]
195pub struct CompletedCompaction {
196    usage: Option<Usage>,
197    estimated_cost: Option<crate::EstimatedUsdCost>,
198    cost_status: crate::CostStatus,
199}
200
201impl CompletedCompaction {
202    /// Returns token usage reported by the compaction operation.
203    #[must_use]
204    pub const fn usage(&self) -> Option<&Usage> {
205        self.usage.as_ref()
206    }
207
208    /// Returns the automatic local USD estimate when usage was available.
209    #[must_use]
210    pub const fn estimated_cost(&self) -> Option<&crate::EstimatedUsdCost> {
211        self.estimated_cost.as_ref()
212    }
213
214    /// Returns why an estimate is present or unavailable.
215    #[must_use]
216    pub const fn cost_status(&self) -> crate::CostStatus {
217        self.cost_status
218    }
219}
220
221#[cfg(not(target_family = "wasm"))]
222type ResponseRun<'a> =
223    Pin<Box<dyn Future<Output = Result<CompletedResponse, ResponseError>> + Send + 'a>>;
224#[cfg(target_family = "wasm")]
225type ResponseRun<'a> = Pin<Box<dyn Future<Output = Result<CompletedResponse, ResponseError>> + 'a>>;
226
227/// A single Responses operation that is both a typed stream and an awaitable
228/// completed aggregate.
229#[must_use = "a response does no work unless it is streamed or awaited"]
230pub struct Response<'a> {
231    events: mpsc::Receiver<ResponseEvent>,
232    run: ResponseRun<'a>,
233    result: Option<Result<CompletedResponse, ResponseError>>,
234    run_finished: bool,
235    completed_event_seen: bool,
236    stream_error_emitted: bool,
237}
238
239impl<'a> Response<'a> {
240    pub(super) fn new(events: mpsc::Receiver<ResponseEvent>, run: ResponseRun<'a>) -> Self {
241        Self {
242            events,
243            run,
244            result: None,
245            run_finished: false,
246            completed_event_seen: false,
247            stream_error_emitted: false,
248        }
249    }
250}
251
252impl Stream for Response<'_> {
253    type Item = Result<ResponseEvent, ResponseError>;
254
255    fn poll_next(mut self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<Self::Item>> {
256        loop {
257            match self.events.poll_recv(context) {
258                Poll::Ready(Some(event)) => {
259                    if matches!(event, ResponseEvent::Completed { .. }) {
260                        self.completed_event_seen = true;
261                    }
262                    return Poll::Ready(Some(Ok(event)));
263                }
264                Poll::Ready(None) | Poll::Pending => {}
265            }
266
267            if !self.run_finished {
268                match self.run.as_mut().poll(context) {
269                    Poll::Ready(result) => {
270                        self.result = Some(result);
271                        self.run_finished = true;
272                        continue;
273                    }
274                    Poll::Pending => return Poll::Pending,
275                }
276            }
277
278            if !self.completed_event_seen
279                && let Some(Ok(result)) = self.result.as_ref()
280            {
281                let event = ResponseEvent::Completed {
282                    usage: result.usage.clone(),
283                    end_turn: result.end_turn,
284                };
285                self.completed_event_seen = true;
286                return Poll::Ready(Some(Ok(event)));
287            }
288
289            let stream_error = (!self.stream_error_emitted)
290                .then(|| self.result.as_ref())
291                .flatten()
292                .and_then(|result| result.as_ref().err())
293                .cloned();
294            if let Some(error) = stream_error {
295                self.stream_error_emitted = true;
296                return Poll::Ready(Some(Err(error)));
297            }
298            return Poll::Ready(None);
299        }
300    }
301}
302
303#[cfg(not(target_family = "wasm"))]
304type ResponseIntoFuture<'a> =
305    Pin<Box<dyn Future<Output = Result<CompletedResponse, ResponseError>> + Send + 'a>>;
306#[cfg(target_family = "wasm")]
307type ResponseIntoFuture<'a> =
308    Pin<Box<dyn Future<Output = Result<CompletedResponse, ResponseError>> + 'a>>;
309
310impl<'a> IntoFuture for Response<'a> {
311    type Output = Result<CompletedResponse, ResponseError>;
312    type IntoFuture = ResponseIntoFuture<'a>;
313
314    fn into_future(mut self) -> Self::IntoFuture {
315        Box::pin(async move {
316            while poll_fn(|context| Pin::new(&mut self).poll_next(context))
317                .await
318                .is_some()
319            {}
320            self.result.take().unwrap_or_else(|| {
321                Err(ResponseError::protocol(
322                    "response stream ended without a terminal service result",
323                ))
324            })
325        })
326    }
327}
328
329/// Cloneable typed failure returned by a response stream and its completed
330/// future.
331#[derive(Clone)]
332pub struct ResponseError {
333    inner: Arc<ResponseErrorInner>,
334}
335
336/// Stable classification for one failed Responses operation.
337#[derive(Clone, Copy, Debug, Eq, PartialEq)]
338#[non_exhaustive]
339pub enum ResponseErrorKind {
340    /// The provider rejected the request because its context window was
341    /// exceeded.
342    ContextWindowExceeded,
343    /// The configured Tower service, transport, or provider API failed.
344    Service,
345    /// A completed operation violated a Responses lifecycle invariant.
346    Protocol,
347}
348
349enum ResponseErrorInner {
350    Source {
351        kind: ResponseErrorKind,
352        error: Arc<dyn Error + Send + Sync>,
353    },
354    Protocol(Arc<str>),
355}
356
357impl ResponseError {
358    /// Wraps a caller-composed Tower service error without erasing its source.
359    #[must_use]
360    pub fn service(error: impl Error + Send + Sync + 'static) -> Self {
361        let kind = if error_chain_responses_error(&error)
362            .is_some_and(crate::ResponsesError::is_context_window_exceeded)
363        {
364            ResponseErrorKind::ContextWindowExceeded
365        } else {
366            ResponseErrorKind::Service
367        };
368        let error: Arc<dyn Error + Send + Sync> = Arc::new(error);
369        Self {
370            inner: Arc::new(ResponseErrorInner::Source { kind, error }),
371        }
372    }
373
374    fn protocol(detail: impl Into<Arc<str>>) -> Self {
375        Self {
376            inner: Arc::new(ResponseErrorInner::Protocol(detail.into())),
377        }
378    }
379
380    /// Returns the stable operation-level error class.
381    #[must_use]
382    pub fn kind(&self) -> ResponseErrorKind {
383        match self.inner.as_ref() {
384            ResponseErrorInner::Source { kind, .. } => *kind,
385            ResponseErrorInner::Protocol(_) => ResponseErrorKind::Protocol,
386        }
387    }
388
389    /// Returns the underlying provider or transport error when one exists.
390    ///
391    /// This traverses caller middleware and the standard Tower service error,
392    /// so callers do not need to downcast each possible service boundary.
393    #[must_use]
394    pub fn responses_error(&self) -> Option<&crate::ResponsesError> {
395        self.source().and_then(error_chain_responses_error)
396    }
397
398    /// Returns whether the provider rejected the request for context-window
399    /// exhaustion.
400    #[must_use]
401    pub fn is_context_window_exceeded(&self) -> bool {
402        matches!(self.kind(), ResponseErrorKind::ContextWindowExceeded)
403    }
404}
405
406impl From<ResponsesServiceError> for ResponseError {
407    fn from(error: ResponsesServiceError) -> Self {
408        Self::service(error)
409    }
410}
411
412impl From<crate::ResponsesError> for ResponseError {
413    fn from(error: crate::ResponsesError) -> Self {
414        Self::service(error)
415    }
416}
417
418impl From<::tower::BoxError> for ResponseError {
419    fn from(error: ::tower::BoxError) -> Self {
420        let kind = if error_chain_responses_error(error.as_ref())
421            .is_some_and(crate::ResponsesError::is_context_window_exceeded)
422        {
423            ResponseErrorKind::ContextWindowExceeded
424        } else {
425            ResponseErrorKind::Service
426        };
427        let error: Arc<dyn Error + Send + Sync> = Arc::from(error);
428        Self {
429            inner: Arc::new(ResponseErrorInner::Source { kind, error }),
430        }
431    }
432}
433
434impl From<Infallible> for ResponseError {
435    fn from(error: Infallible) -> Self {
436        match error {}
437    }
438}
439
440impl fmt::Display for ResponseError {
441    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
442        match self.inner.as_ref() {
443            ResponseErrorInner::Source {
444                kind: ResponseErrorKind::ContextWindowExceeded,
445                ..
446            } => formatter.write_str("Responses input exceeded the model context window"),
447            ResponseErrorInner::Source { error, .. } => error.fmt(formatter),
448            ResponseErrorInner::Protocol(detail) => {
449                write!(formatter, "invalid Responses state: {detail}")
450            }
451        }
452    }
453}
454
455impl fmt::Debug for ResponseError {
456    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
457        formatter
458            .debug_struct("ResponseError")
459            .field("message", &self.to_string())
460            .finish_non_exhaustive()
461    }
462}
463
464impl Error for ResponseError {
465    fn source(&self) -> Option<&(dyn Error + 'static)> {
466        match self.inner.as_ref() {
467            ResponseErrorInner::Source { error, .. } => Some(error.as_ref()),
468            ResponseErrorInner::Protocol(_) => None,
469        }
470    }
471}
472
473fn error_chain_responses_error<'a>(
474    mut error: &'a (dyn Error + 'static),
475) -> Option<&'a crate::ResponsesError> {
476    loop {
477        if let Some(service) = error.downcast_ref::<ResponsesServiceError>()
478            && let Some(error) = service.responses_error()
479        {
480            return Some(error);
481        }
482        if let Some(error) = error.downcast_ref::<crate::ResponsesError>() {
483            return Some(error);
484        }
485        let source = error.source()?;
486        error = source;
487    }
488}
489
490pub(super) async fn run_create<S>(
491    turn: &mut ResponseTurn<'_, S>,
492    input: ResponseInput,
493    sink: EventSink,
494    response_events: mpsc::Sender<ResponseEvent>,
495) -> Result<CompletedResponse, ResponseError>
496where
497    S: Service<ResponsesAttempt, Response = ResponsesServiceResponse>,
498    S::Error: Into<ResponseError>,
499{
500    let span = response_call_span(
501        turn.session,
502        turn.logical_turn,
503        "response.create",
504        input.items.len(),
505    );
506    let started_at = Instant::now();
507    let mut cancellation = ResponseCallCancellation::new(span.clone(), started_at);
508    let result = run_create_inner(turn, input, sink, response_events)
509        .instrument(span.clone())
510        .await;
511    cancellation.complete();
512    finish_create_span(&span, started_at, &result);
513    result
514}
515
516async fn run_create_inner<S>(
517    turn: &mut ResponseTurn<'_, S>,
518    mut input: ResponseInput,
519    sink: EventSink,
520    response_events: mpsc::Sender<ResponseEvent>,
521) -> Result<CompletedResponse, ResponseError>
522where
523    S: Service<ResponsesAttempt, Response = ResponsesServiceResponse>,
524    S::Error: Into<ResponseError>,
525{
526    if input.items.is_empty() {
527        return Err(ResponseError::protocol(
528            "response.create input must contain at least one item",
529        ));
530    }
531
532    let session = &mut *turn.session;
533    let call_index = session.next_call_index;
534    session.next_call_index = session.next_call_index.saturating_add(1);
535    assign_missing_response_item_ids(&mut input.items);
536    let observed_canonical_context = input
537        .items
538        .iter()
539        .filter(|item| is_canonical_context_item(item))
540        .cloned()
541        .collect::<Vec<_>>();
542    let reinject_canonical_context =
543        session.canonical_context_reinjection_pending && observed_canonical_context.is_empty();
544    let mut candidate = session.state.clone();
545    if reinject_canonical_context {
546        candidate.append(session.canonical_context.iter().cloned());
547    }
548    candidate.append(input.items);
549    let (prompt_history, prompt_repaired) = candidate.prompt_history_with_repair();
550    let previous_response_id = if prompt_repaired {
551        None
552    } else {
553        candidate.previous_response_id().map(str::to_owned)
554    };
555
556    let factory = ResponsesAttemptFactory::new(
557        session.profile.clone(),
558        sink,
559        Arc::clone(&session.transport_stats),
560    )
561    .with_response_events(response_events)
562    .for_logical_turn(turn.logical_turn);
563    let request = factory.generation(
564        call_index,
565        prompt_history.clone(),
566        candidate.shared_history(),
567        candidate.delta_start(),
568        previous_response_id.as_deref(),
569        session.thinking,
570        session.fast_mode,
571    );
572    let success = session.client.execute(request).await.map_err(Into::into)?;
573    candidate.observe_server_reasoning(success.server_reasoning_included());
574    let ResponsesOutput::Generation(response) = success.into_output() else {
575        return Err(ResponseError::protocol(
576            "response.create returned a non-generation output",
577        ));
578    };
579
580    if prompt_repaired {
581        candidate.adopt_prompt_history(prompt_history);
582    }
583    candidate.append(response.output_items.clone());
584    candidate.update_token_info(response.usage.as_ref());
585    candidate.set_previous_response_id(response.id);
586    candidate
587        .commit()
588        .map_err(|error| ResponseError::protocol(error.to_string()))?;
589    session.state = candidate;
590    if !observed_canonical_context.is_empty() {
591        session.canonical_context = observed_canonical_context;
592    }
593    session.canonical_context_reinjection_pending = false;
594
595    let output: Arc<[ResponseItem]> = response.output_items.into();
596    let output_text = response
597        .final_message
598        .unwrap_or_else(|| output_text(&output))
599        .into();
600    let (estimated_cost, cost_status) = estimate_cost(response.usage.as_ref(), session.fast_mode);
601    let completed = CompletedResponse {
602        output,
603        output_text,
604        usage: response.usage,
605        estimated_cost,
606        cost_status,
607        end_turn: response.end_turn,
608    };
609    turn.completed_generation = true;
610    Ok(completed)
611}
612
613pub(super) async fn run_compact<S>(
614    turn: &mut ResponseTurn<'_, S>,
615) -> Result<CompletedCompaction, ResponseError>
616where
617    S: Service<ResponsesAttempt, Response = ResponsesServiceResponse>,
618    S::Error: Into<ResponseError>,
619{
620    let span = response_call_span(
621        turn.session,
622        turn.logical_turn,
623        "response.compact",
624        turn.session.history_len(),
625    );
626    let started_at = Instant::now();
627    let mut cancellation = ResponseCallCancellation::new(span.clone(), started_at);
628    let result = run_compact_inner(turn).instrument(span.clone()).await;
629    cancellation.complete();
630    finish_compaction_span(&span, started_at, &result);
631    result
632}
633
634struct ResponseCallCancellation {
635    span: tracing::Span,
636    started_at: Instant,
637    completed: bool,
638}
639
640impl ResponseCallCancellation {
641    const fn new(span: tracing::Span, started_at: Instant) -> Self {
642        Self {
643            span,
644            started_at,
645            completed: false,
646        }
647    }
648
649    const fn complete(&mut self) {
650        self.completed = true;
651    }
652}
653
654impl Drop for ResponseCallCancellation {
655    fn drop(&mut self) {
656        if self.completed {
657            return;
658        }
659        self.span.record("error.class", "cancelled");
660        self.span.record("status", "cancelled");
661        self.span.record("otel.status_code", "ERROR");
662        self.span.record(
663            "duration_ns",
664            u64::try_from(self.started_at.elapsed().as_nanos()).unwrap_or(u64::MAX),
665        );
666        tracing::warn!(
667            target: "nanocodex_oai_api",
668            parent: &self.span,
669            "Responses call cancelled before a provider terminal event"
670        );
671    }
672}
673
674async fn run_compact_inner<S>(
675    turn: &mut ResponseTurn<'_, S>,
676) -> Result<CompletedCompaction, ResponseError>
677where
678    S: Service<ResponsesAttempt, Response = ResponsesServiceResponse>,
679    S::Error: Into<ResponseError>,
680{
681    let session = &mut *turn.session;
682    let call_index = session.next_call_index;
683    session.next_call_index = session.next_call_index.saturating_add(1);
684    let (sink, events) = EventSink::channel(session.profile.session_id().to_owned());
685    drop(events);
686    let factory = ResponsesAttemptFactory::new(
687        session.profile.clone(),
688        sink,
689        Arc::clone(&session.transport_stats),
690    )
691    .for_logical_turn(turn.logical_turn);
692    let mut history = session.state.prompt_history();
693    compaction::trim_tool_outputs_to_fit_context_window(&mut history, session.profile.prefix());
694    let request = factory.compaction(
695        call_index,
696        history.clone(),
697        history,
698        session.state.delta_start(),
699        session.state.previous_response_id(),
700        compaction::trigger(),
701        session.thinking,
702        session.fast_mode,
703    );
704    let success = session.client.execute(request).await.map_err(Into::into)?;
705    let server_reasoning_included = success.server_reasoning_included();
706    let ResponsesOutput::Compaction(response) = success.into_output() else {
707        return Err(ResponseError::protocol(
708            "response.compact returned a non-compaction output",
709        ));
710    };
711
712    let mut candidate = session.state.clone();
713    candidate.observe_server_reasoning(server_reasoning_included);
714    let mid_turn = turn.completed_generation;
715    let canonical_context = if mid_turn {
716        session.canonical_context.clone()
717    } else {
718        Vec::new()
719    };
720    candidate.install_compaction(response.item, canonical_context, session.profile.prefix());
721    session.state = candidate;
722    session.canonical_context_reinjection_pending = !mid_turn;
723
724    let (estimated_cost, cost_status) = estimate_cost(response.usage.as_ref(), session.fast_mode);
725    Ok(CompletedCompaction {
726        usage: response.usage,
727        estimated_cost,
728        cost_status,
729    })
730}
731
732fn response_call_span<S>(
733    session: &Session<S>,
734    logical_turn: u64,
735    method: &'static str,
736    input_item_count: usize,
737) -> tracing::Span {
738    info_span!(
739        target: "nanocodex_oai_api",
740        "responses.call",
741        otel.kind = "client",
742        otel.status_code = tracing::field::Empty,
743        session.id = %session.id,
744        response.method = method,
745        turn.index = logical_turn,
746        model.call_index = session.next_call_index,
747        model.input.item_count = input_item_count,
748        model.output.item_count = tracing::field::Empty,
749        response.end_turn = tracing::field::Empty,
750        usage.input_tokens = tracing::field::Empty,
751        usage.cached_input_tokens = tracing::field::Empty,
752        usage.cache_write_input_tokens = tracing::field::Empty,
753        usage.output_tokens = tracing::field::Empty,
754        usage.reasoning_output_tokens = tracing::field::Empty,
755        usage.total_tokens = tracing::field::Empty,
756        cost.usd = tracing::field::Empty,
757        cost.status = tracing::field::Empty,
758        cost.service_tier = tracing::field::Empty,
759        error.class = tracing::field::Empty,
760        status = tracing::field::Empty,
761        duration_ns = tracing::field::Empty,
762    )
763}
764
765fn finish_create_span(
766    span: &tracing::Span,
767    started_at: Instant,
768    result: &Result<CompletedResponse, ResponseError>,
769) {
770    match result {
771        Ok(response) => {
772            span.record("model.output.item_count", response.output.len());
773            if let Some(end_turn) = response.end_turn {
774                span.record("response.end_turn", end_turn);
775            }
776            record_response_usage_and_cost(
777                span,
778                response.usage.as_ref(),
779                response.cost_status,
780                response.estimated_cost.as_ref(),
781            );
782            finish_response_span(span, started_at, None);
783        }
784        Err(error) => finish_response_span(span, started_at, Some(error)),
785    }
786}
787
788fn finish_compaction_span(
789    span: &tracing::Span,
790    started_at: Instant,
791    result: &Result<CompletedCompaction, ResponseError>,
792) {
793    match result {
794        Ok(response) => {
795            record_response_usage_and_cost(
796                span,
797                response.usage.as_ref(),
798                response.cost_status,
799                response.estimated_cost.as_ref(),
800            );
801            finish_response_span(span, started_at, None);
802        }
803        Err(error) => finish_response_span(span, started_at, Some(error)),
804    }
805}
806
807fn record_response_usage_and_cost(
808    span: &tracing::Span,
809    usage: Option<&Usage>,
810    cost_status: crate::CostStatus,
811    estimated_cost: Option<&crate::EstimatedUsdCost>,
812) {
813    if let Some(usage) = usage {
814        span.record("usage.input_tokens", usage.input_tokens);
815        span.record(
816            "usage.cached_input_tokens",
817            usage
818                .input_tokens_details
819                .as_ref()
820                .map_or(0, |details| details.cached_tokens),
821        );
822        span.record(
823            "usage.cache_write_input_tokens",
824            usage
825                .input_tokens_details
826                .as_ref()
827                .map_or(0, |details| details.cache_write_tokens),
828        );
829        span.record("usage.output_tokens", usage.output_tokens);
830        span.record(
831            "usage.reasoning_output_tokens",
832            usage
833                .output_tokens_details
834                .as_ref()
835                .map_or(0, |details| details.reasoning_tokens),
836        );
837        span.record("usage.total_tokens", usage.total_tokens);
838    }
839    span.record("cost.status", cost_status.as_str());
840    if let Some(estimated_cost) = estimated_cost {
841        span.record("cost.usd", tracing::field::display(estimated_cost.amount()));
842        span.record("cost.service_tier", estimated_cost.service_tier().as_str());
843    }
844}
845
846fn finish_response_span(span: &tracing::Span, started_at: Instant, error: Option<&ResponseError>) {
847    span.record(
848        "duration_ns",
849        u64::try_from(started_at.elapsed().as_nanos()).unwrap_or(u64::MAX),
850    );
851    if let Some(error) = error {
852        let error_class = match error.kind() {
853            ResponseErrorKind::ContextWindowExceeded => "context_window_exceeded",
854            ResponseErrorKind::Service => "service",
855            ResponseErrorKind::Protocol => "protocol",
856        };
857        span.record("error.class", error_class);
858        span.record("status", "failed");
859        span.record("otel.status_code", "ERROR");
860        tracing::error!(
861            target: "nanocodex_oai_api",
862            parent: span,
863            error = %error,
864            error.class = error_class,
865            "Responses call failed"
866        );
867    } else {
868        span.record("status", "completed");
869        span.record("otel.status_code", "OK");
870    }
871}
872
873pub(super) fn estimate_cost(
874    usage: Option<&Usage>,
875    fast_mode: bool,
876) -> (Option<crate::EstimatedUsdCost>, crate::CostStatus) {
877    match usage {
878        Some(usage) => (
879            Some(crate::pricing::estimate(
880                usage,
881                if fast_mode {
882                    crate::pricing::ServiceTier::Priority
883                } else {
884                    crate::pricing::ServiceTier::Standard
885                },
886            )),
887            crate::CostStatus::EstimatedFromUsage,
888        ),
889        None => (None, crate::CostStatus::UsageNotReported),
890    }
891}
892
893fn output_text(items: &[ResponseItem]) -> String {
894    items
895        .iter()
896        .filter_map(|item| {
897            let ResponseItem::Message { content, .. } = item else {
898                return None;
899            };
900            Some(content.iter().filter_map(|content| {
901                let ContentItem::OutputText { text, .. } = content else {
902                    return None;
903                };
904                Some(text.as_ref())
905            }))
906        })
907        .flatten()
908        .collect()
909}