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]
165 pub const fn estimated_cost(&self) -> Option<&crate::EstimatedUsdCost> {
166 self.estimated_cost.as_ref()
167 }
168
169 #[must_use]
171 pub const fn cost_status(&self) -> crate::CostStatus {
172 self.cost_status
173 }
174
175 #[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#[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 #[must_use]
204 pub const fn usage(&self) -> Option<&Usage> {
205 self.usage.as_ref()
206 }
207
208 #[must_use]
210 pub const fn estimated_cost(&self) -> Option<&crate::EstimatedUsdCost> {
211 self.estimated_cost.as_ref()
212 }
213
214 #[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#[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#[derive(Clone)]
332pub struct ResponseError {
333 inner: Arc<ResponseErrorInner>,
334}
335
336#[derive(Clone, Copy, Debug, Eq, PartialEq)]
338#[non_exhaustive]
339pub enum ResponseErrorKind {
340 ContextWindowExceeded,
343 Service,
345 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 #[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 #[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 #[must_use]
394 pub fn responses_error(&self) -> Option<&crate::ResponsesError> {
395 self.source().and_then(error_chain_responses_error)
396 }
397
398 #[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}