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#[derive(Clone, Debug)]
31pub struct ResponseInput {
32 items: Vec<ResponseItem>,
33}
34
35impl ResponseInput {
36 #[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 #[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 #[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#[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 #[must_use]
131 pub fn output(&self) -> &[ResponseItem] {
132 &self.output
133 }
134
135 #[must_use]
137 pub fn output_text(&self) -> &str {
138 &self.output_text
139 }
140
141 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 #[must_use]
156 pub const fn usage(&self) -> Option<&Usage> {
157 self.usage.as_ref()
158 }
159
160 #[must_use]
166 pub const fn estimated_cost(&self) -> Option<&crate::EstimatedUsdCost> {
167 self.estimated_cost.as_ref()
168 }
169
170 #[must_use]
172 pub const fn cost_status(&self) -> crate::CostStatus {
173 self.cost_status
174 }
175
176 #[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#[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 #[must_use]
205 pub const fn usage(&self) -> Option<&Usage> {
206 self.usage.as_ref()
207 }
208
209 #[must_use]
211 pub const fn estimated_cost(&self) -> Option<&crate::EstimatedUsdCost> {
212 self.estimated_cost.as_ref()
213 }
214
215 #[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#[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#[derive(Clone)]
333pub struct ResponseError {
334 inner: Arc<ResponseErrorInner>,
335}
336
337#[derive(Clone, Copy, Debug, Eq, PartialEq)]
339#[non_exhaustive]
340pub enum ResponseErrorKind {
341 ContextWindowExceeded,
344 Service,
346 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 #[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 #[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 #[must_use]
395 pub fn responses_error(&self) -> Option<&crate::ResponsesError> {
396 self.source().and_then(error_chain_responses_error)
397 }
398
399 #[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}