Skip to main content

toolkit_contract/runtime/
client.rs

1//! Per-call helpers used by the generated REST client.
2//!
3//! The macro keeps emitted code small by funnelling the common
4//! "send unary request and decode the response" path through
5//! [`send_unary`], and the streaming path through
6//! [`send_streaming`].
7
8use std::pin::Pin;
9use std::time::Duration;
10
11use futures_core::Stream;
12use serde::Serialize;
13use serde::de::DeserializeOwned;
14use toolkit_http::RequestBuilder;
15
16use crate::runtime::config::ReconnectConfig;
17use crate::runtime::http::{body_to_byte_stream, map_http_error, parse_retry_after};
18use crate::runtime::sse::{LastEventId, parse_sse_stream_with_id};
19use crate::runtime::transport_error::TransportError;
20
21/// Boxed byte stream produced from an HTTP response body for SSE parsing.
22type BoxByteStream = Pin<
23    Box<
24        dyn Stream<Item = Result<bytes::Bytes, Box<dyn std::error::Error + Send + Sync + 'static>>>
25            + Send,
26    >,
27>;
28
29/// Send a unary request and decode the JSON response.
30///
31/// `build` is a closure returning a `RequestBuilder` configured with method,
32/// URL, headers, and any auth state. It is invoked once per attempt.
33///
34/// `timeout` bounds the **whole** attempt — connection, response headers, and
35/// reading the response body — as a per-attempt deadline (mirroring the
36/// streaming path). Elapse maps to [`TransportError::Timeout`], which is
37/// transient, so a `#[retryable]` method retries a timed-out attempt. `None`
38/// leaves the deadline to the underlying transport.
39///
40/// # Errors
41/// Returns [`TransportError`] when the builder closure fails, the deadline elapses,
42/// the network call fails, the response body cannot be read, JSON deserialization of a
43/// success body fails, or the server returns a non-success HTTP status (mapped via
44/// [`map_http_error`]).
45pub async fn send_unary<F, R>(build: F, timeout: Option<Duration>) -> Result<R, TransportError>
46where
47    F: FnOnce() -> Result<RequestBuilder, TransportError>,
48    R: DeserializeOwned,
49{
50    let builder = build()?;
51    let op = async move {
52        let response = builder.send().await.map_err(TransportError::network)?;
53        let status = response.status();
54        if status.is_success() {
55            // `HttpResponse::bytes` does NOT enforce status check; we already did.
56            // Avoiding `.json()` here because it would re-check status and
57            // duplicate work on the success path.
58            let bytes = response.bytes().await.map_err(TransportError::network)?;
59            // A `204 No Content` (or any 2xx with an empty body) is valid for
60            // methods returning `Result<(), _>` and other unit-like `Ok` types.
61            // `serde_json::from_slice::<()>(b"")` would fail with an EOF error
62            // (and `()` deserializes only from JSON `null`), so map an empty
63            // success body to JSON `null` before decoding — `null` deserializes
64            // into `()` and `Option::None`, while any non-unit `R` still errors.
65            let bytes: &[u8] = if bytes.is_empty() { b"null" } else { &bytes };
66            serde_json::from_slice::<R>(bytes).map_err(TransportError::serialization)
67        } else {
68            let retry_after = parse_retry_after(response.headers());
69            let bytes = response.bytes().await.map_err(TransportError::network)?;
70            let body = String::from_utf8_lossy(&bytes).into_owned();
71            Err(map_http_error(status.as_u16(), body, retry_after))
72        }
73    };
74    match timeout {
75        Some(d) => tokio::time::timeout(d, op)
76            .await
77            .map_err(|_| TransportError::Timeout(d))?,
78        None => op.await,
79    }
80}
81
82/// Build the default `toolkit-http` client used by macro-generated REST clients.
83///
84/// Transport-layer retry is **disabled** — the SDK runs its own retry loop in
85/// [`retry_with_backoff`](crate::runtime::retry::retry_with_backoff).
86///
87/// `require_tls` selects the transport security mode: `false` (the default
88/// via [`ClientConfig::new`](crate::runtime::config::ClientConfig::new))
89/// allows plaintext `http://`, preserving the platform's existing in-mesh
90/// service-to-service convention; `true`
91/// ([`ClientConfig::with_require_tls`](crate::runtime::config::ClientConfig::with_require_tls))
92/// switches to `toolkit_http::TransportSecurity::TlsOnly`, rejecting
93/// plaintext. This matters because the tenant bearer token is forwarded via
94/// `Authorization` on whatever scheme the resolved endpoint uses — set
95/// `require_tls` when a client may talk to an endpoint outside a trusted
96/// network boundary.
97///
98/// `.with_otel()` is called in BOTH cfg variants below — it always installs
99/// `toolkit-http`'s `OtelLayer`, but the layer's actual W3C `traceparent`
100/// injection is itself `#[cfg(feature = "otel")]` **inside `toolkit-http`**
101/// (a no-op otherwise). Because Cargo unifies features workspace-wide, that
102/// gate tracks whether *anything* in the final binary enables
103/// `toolkit-http/otel` — which may not be this crate's own `otel` feature.
104/// The one thing genuinely gated by *this* crate's `otel` feature (which
105/// forwards `toolkit-http/otel`) is `.with_metrics(client_type)` below, since
106/// that builder method only exists when the dependency is compiled with the
107/// feature on. The generated method's per-method `tracing` span is always
108/// opened either way, independent of any of this.
109///
110/// # Errors
111/// Propagates [`toolkit_http::HttpError`] from the underlying builder (e.g. a
112/// TLS backend that cannot be constructed under FIPS).
113#[cfg(feature = "otel")]
114pub fn build_default_http_client(
115    client_type: &str,
116    require_tls: bool,
117) -> Result<toolkit_http::HttpClient, toolkit_http::HttpError> {
118    toolkit_http::HttpClient::builder()
119        .retry(None)
120        .transport(transport_security(require_tls))
121        .with_otel()
122        .with_metrics(client_type)
123        .build()
124}
125
126/// Non-`otel` build of [`build_default_http_client`]: RED metrics
127/// (`.with_metrics`) are NOT compiled in — that builder method only exists
128/// under `toolkit-http/otel`. `client_type` is unused in this configuration.
129/// `.with_otel()` is still called (see the doc on the `otel`-cfg sibling
130/// above): whether it actually propagates `traceparent` depends on whether
131/// `toolkit-http/otel` ends up enabled by *some* crate in the build, not on
132/// this crate's own `otel` feature.
133///
134/// # Errors
135/// Propagates [`toolkit_http::HttpError`] from the underlying builder (e.g. a
136/// TLS backend that cannot be constructed under FIPS).
137#[cfg(not(feature = "otel"))]
138pub fn build_default_http_client(
139    _client_type: &str,
140    require_tls: bool,
141) -> Result<toolkit_http::HttpClient, toolkit_http::HttpError> {
142    toolkit_http::HttpClient::builder()
143        .retry(None)
144        .transport(transport_security(require_tls))
145        .with_otel()
146        .build()
147}
148
149fn transport_security(require_tls: bool) -> toolkit_http::TransportSecurity {
150    if require_tls {
151        toolkit_http::TransportSecurity::TlsOnly
152    } else {
153        toolkit_http::TransportSecurity::AllowInsecureHttp
154    }
155}
156
157/// Add a JSON body to a request builder. Wraps `toolkit_http`'s fallible
158/// `.json()` (which can fail to serialize) in our [`TransportError`] surface
159/// so the macro emit path can `?` uniformly.
160///
161/// # Errors
162/// Returns [`TransportError::Serialization`] when `body` cannot be serialized to JSON.
163pub fn with_json_body<T: Serialize>(
164    builder: RequestBuilder,
165    body: &T,
166) -> Result<RequestBuilder, TransportError> {
167    builder.json(body).map_err(TransportError::serialization)
168}
169
170/// Builder for an SSE request that can be re-issued on reconnect.
171///
172/// `build` receives the latest seen `Last-Event-ID` (or `None` on the first
173/// attempt) and must return a fresh, configured `RequestBuilder`.
174/// Implementations should set the `Last-Event-ID` header from the parameter
175/// when present.
176pub trait StreamRequestFactory: Send + 'static {
177    /// Construct a `RequestBuilder` for the next stream attempt.
178    ///
179    /// # Errors
180    /// Returns [`TransportError`] when the factory cannot produce a builder
181    /// (e.g. URL composition or auth header attachment fails).
182    fn build(&self, last_event_id: Option<&str>) -> Result<RequestBuilder, TransportError>;
183}
184
185impl<F> StreamRequestFactory for F
186where
187    F: Fn(Option<&str>) -> Result<RequestBuilder, TransportError> + Send + 'static,
188{
189    fn build(&self, last: Option<&str>) -> Result<RequestBuilder, TransportError> {
190        (self)(last)
191    }
192}
193
194/// Send a streaming SSE request and adapt the response into a typed stream.
195///
196/// `factory` produces a fresh `RequestBuilder` per attempt; the
197/// `last_event_id` argument is `None` on the first attempt and contains the
198/// most recently observed SSE `id:` field on every reconnect.
199///
200/// The returned stream yields `Result<T, TransportError>` items per SSE
201/// event. With `reconnect.max_attempts == 0` (the default), a transient
202/// transport failure ends the stream immediately. With a non-zero limit,
203/// the client re-issues the request with a `Last-Event-ID` header up to
204/// `max_attempts` times, applying exponential backoff between attempts.
205pub fn send_streaming<F, T>(
206    factory: F,
207    reconnect: ReconnectConfig,
208    timeout: Option<Duration>,
209) -> Pin<Box<dyn Stream<Item = Result<T, TransportError>> + Send>>
210where
211    F: StreamRequestFactory,
212    T: DeserializeOwned + Send + 'static,
213{
214    use futures_util::StreamExt;
215
216    Box::pin(async_stream::try_stream! {
217        let last_id = LastEventId::empty();
218        let mut attempt = 0u32;
219
220        loop {
221            let snapshot = last_id.current();
222            let builder = match factory.build(snapshot.as_deref()) {
223                Ok(b) => b,
224                Err(e) => {
225                    // URL-build / serialization failure on the request side —
226                    // not eligible for reconnect (config-shape error).
227                    Err(e)?;
228                    return;
229                }
230            };
231            let send_fut = builder.send();
232            let response = match timeout {
233                Some(d) => if let Ok(r) = tokio::time::timeout(d, send_fut).await { r } else {
234                    let err = TransportError::Timeout(d);
235                    if attempt < reconnect.max_attempts && err.is_transient() {
236                        attempt += 1;
237                        sleep_backoff(&reconnect, attempt).await;
238                        continue;
239                    }
240                    Err(err)?;
241                    return;
242                },
243                None => send_fut.await,
244            };
245            let response = match response {
246                Ok(r) => r,
247                Err(e) => {
248                    // Pre-flight failures (DNS, connect refused) are
249                    // network-class — eligible for reconnect.
250                    let err = TransportError::network(e);
251                    if attempt < reconnect.max_attempts && err.is_transient() {
252                        attempt += 1;
253                        sleep_backoff(&reconnect, attempt).await;
254                        continue;
255                    }
256                    Err(err)?;
257                    return;
258                }
259            };
260            let status = response.status();
261            if !status.is_success() {
262                let status_code = status.as_u16();
263                let retry_after = parse_retry_after(response.headers());
264                // Bound the error-body read by the same per-attempt deadline as
265                // the initial send — otherwise a slow/oversized error body on
266                // the pre-stream path could stall the attempt indefinitely,
267                // defeating the timeout guarantee.
268                let body_fut = response.bytes();
269                let bytes_result = match timeout {
270                    Some(d) => if let Ok(r) = tokio::time::timeout(d, body_fut).await { r } else {
271                        let err = TransportError::Timeout(d);
272                        if attempt < reconnect.max_attempts && err.is_transient() {
273                            attempt += 1;
274                            sleep_backoff(&reconnect, attempt).await;
275                            continue;
276                        }
277                        Err(err)?;
278                        return;
279                    },
280                    None => body_fut.await,
281                };
282                let bytes = bytes_result.map_err(TransportError::network)?;
283                let body = String::from_utf8_lossy(&bytes).into_owned();
284                let err = map_http_error(status_code, body, retry_after);
285                // A transient status (e.g. 503 during a rolling deploy) on the
286                // initial connect is reconnect-eligible, consistent with the
287                // pre-flight failure branch above; other statuses are domain
288                // errors and bubble straight through.
289                if attempt < reconnect.max_attempts && err.is_transient() {
290                    attempt += 1;
291                    sleep_backoff(&reconnect, attempt).await;
292                    continue;
293                }
294                Err(err)?;
295                return;
296            }
297
298            // `parse_sse_stream_with_id` needs `Unpin + 'static`; pinning
299            // on the stack with `pin_mut!` would borrow `byte_stream` for
300            // less than `'static`, so move ownership behind `Box::pin` and
301            // hand the boxed stream to the parser.
302            let byte_stream: BoxByteStream = Box::pin(body_to_byte_stream(response.into_body()));
303            let mut inner = parse_sse_stream_with_id::<T, _, _>(
304                byte_stream,
305                last_id.clone(),
306            );
307            let activity = inner.activity_handle();
308            let mut stream_err: Option<TransportError> = None;
309            // Whether this connection ever delivered an event. A connection that
310            // produced data was healthy, so the failure that follows it starts a
311            // fresh reconnect budget rather than continuing the previous one.
312            let mut delivered_an_event = false;
313            loop {
314                // Idle timeout is *idle*: keepalive comments (and any other wire
315                // chunk that dispatches no item) still count as activity, so a
316                // quiet-but-alive stream is not torn down. On elapse we only
317                // error if no chunk arrived while we waited.
318                let item = match timeout {
319                    Some(d) => loop {
320                        let before = activity.generation();
321                        match tokio::time::timeout(d, inner.next()).await {
322                            Ok(v) => break v,
323                            // Activity advanced since the wait started (e.g. a
324                            // keepalive comment) — loop again rather than
325                            // treating this as a genuine idle timeout. Falling
326                            // off this arm already re-enters the loop; an
327                            // explicit `continue` here is redundant.
328                            Err(_) if activity.generation() != before => {}
329                            Err(_) => {
330                                stream_err = Some(TransportError::Timeout(d));
331                                break None;
332                            }
333                        }
334                    },
335                    None => inner.next().await,
336                };
337                if stream_err.is_some() {
338                    break;
339                }
340                match item {
341                    Some(Ok(v)) => {
342                        delivered_an_event = true;
343                        yield v;
344                    }
345                    Some(Err(e)) => {
346                        stream_err = Some(e);
347                        break;
348                    }
349                    None => {
350                        // The underlying byte stream ended with no error, but
351                        // that alone doesn't mean the peer is finished: an
352                        // `event: done` frame and a bare connection close both
353                        // surface as `None` here. Treat a close WITHOUT an
354                        // explicit `done` as an anomaly (reconnect-eligible)
355                        // rather than silently reporting success — otherwise a
356                        // server restart / LB idle-timeout / rolling deploy
357                        // that closes the connection mid-stream would look
358                        // identical to a fully-delivered stream.
359                        if !inner.saw_done_event() {
360                            stream_err = Some(TransportError::sse(
361                                "SSE stream ended without a terminal `done` event",
362                            ));
363                        }
364                        break;
365                    }
366                }
367            }
368
369            // A connection that delivered events was healthy, so the budget
370            // starts over. Without this, `max_attempts` is a lifetime cap rather
371            // than a burst cap: a long-lived subscription that survives N
372            // unrelated blips over days would die on the N+1st, despite the
373            // stream being advertised as indefinitely reconnecting. Resetting
374            // also restarts the backoff at its base delay, which is what you
375            // want after a healthy period rather than resuming a 30s ceiling.
376            if delivered_an_event {
377                attempt = 0;
378            }
379
380            match stream_err {
381                None => return, // Stream ended cleanly (`event: done`).
382                Some(e) if attempt < reconnect.max_attempts && e.is_transient() => {
383                    attempt += 1;
384                    sleep_backoff(&reconnect, attempt).await;
385                    // Fall through to next loop iteration to retry.
386                }
387                Some(e) => {
388                    Err(e)?;
389                    return;
390                }
391            }
392        }
393    })
394}
395
396/// Compute backoff for reconnect attempt #N (1-indexed). Doubles the base
397/// delay each attempt, capped at `max_delay`. Multiplied by a ±25% jitter
398/// factor — fleet-wide reconnect synchronization (many clients reconnecting
399/// in lockstep after a shared upstream blip) is the real concern.
400async fn sleep_backoff(config: &ReconnectConfig, attempt: u32) {
401    use rand::RngExt;
402    let exp = attempt.saturating_sub(1);
403    let multiplier = 2u32.saturating_pow(exp);
404    let base = config
405        .base_delay
406        .saturating_mul(multiplier)
407        .min(config.max_delay);
408    let jitter: f64 = rand::rng().random_range(0.75..=1.25);
409    let secs = base.as_secs_f64() * jitter;
410    let delay = if secs.is_finite() && secs >= 0.0 {
411        Duration::from_secs_f64(secs).min(config.max_delay)
412    } else {
413        base
414    };
415    tokio::time::sleep(delay).await;
416}
417
418// `with_json_body` and the streaming path are exercised end-to-end (real
419// serialized body observed by a real server) by the integration tests in
420// `tests/rest_client_codegen.rs` (`unary_post_round_trip` et al.) — a local
421// unit test here could only check "builds without panicking" without
422// access to `toolkit_http::RequestBuilder`'s private fields, which is
423// strictly weaker than the existing round-trip coverage.