Skip to main content

codex_api/
telemetry.rs

1use crate::error::ApiError;
2use codex_client::Request;
3use codex_client::RequestTelemetry;
4use codex_client::Response;
5use codex_client::RetryPolicy;
6use codex_client::StreamResponse;
7use codex_client::TransportError;
8use codex_client::run_with_retry;
9use http::StatusCode;
10use std::future::Future;
11use std::sync::Arc;
12use std::time::Duration;
13use tokio::time::Instant;
14use tokio_tungstenite::tungstenite::Error;
15use tokio_tungstenite::tungstenite::Message;
16
17/// Generic telemetry.
18pub trait SseTelemetry: Send + Sync {
19    fn on_sse_poll(
20        &self,
21        result: &Result<
22            Option<
23                Result<
24                    eventsource_stream::Event,
25                    eventsource_stream::EventStreamError<TransportError>,
26                >,
27            >,
28            tokio::time::error::Elapsed,
29        >,
30        duration: Duration,
31    );
32}
33
34/// Telemetry for Responses WebSocket transport.
35pub trait WebsocketTelemetry: Send + Sync {
36    fn on_ws_request(&self, duration: Duration, error: Option<&ApiError>, connection_reused: bool);
37
38    fn on_ws_event(
39        &self,
40        result: &Result<Option<Result<Message, Error>>, ApiError>,
41        duration: Duration,
42    );
43}
44
45pub(crate) trait WithStatus {
46    fn status(&self) -> StatusCode;
47}
48
49fn http_status(err: &TransportError) -> Option<StatusCode> {
50    match err {
51        TransportError::Http { status, .. } => Some(*status),
52        _ => None,
53    }
54}
55
56impl WithStatus for Response {
57    fn status(&self) -> StatusCode {
58        self.status
59    }
60}
61
62impl WithStatus for StreamResponse {
63    fn status(&self) -> StatusCode {
64        self.status
65    }
66}
67
68pub(crate) async fn run_with_request_telemetry<T, F, Fut>(
69    policy: RetryPolicy,
70    telemetry: Option<Arc<dyn RequestTelemetry>>,
71    make_request: impl FnMut() -> Request,
72    send: F,
73) -> Result<T, TransportError>
74where
75    T: WithStatus,
76    F: Clone + Fn(Request) -> Fut,
77    Fut: Future<Output = Result<T, TransportError>>,
78{
79    // Wraps `run_with_retry` to attach per-attempt request telemetry for both
80    // unary and streaming HTTP calls.
81    run_with_retry(policy, make_request, move |req, attempt| {
82        let telemetry = telemetry.clone();
83        let send = send.clone();
84        async move {
85            let start = Instant::now();
86            let result = send(req).await;
87            if let Some(t) = telemetry.as_ref() {
88                let (status, err) = match &result {
89                    Ok(resp) => (Some(resp.status()), None),
90                    Err(err) => (http_status(err), Some(err)),
91                };
92                t.on_request(attempt, status, err, start.elapsed());
93            }
94            result
95        }
96    })
97    .await
98}