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.