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