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::ir::binding::StreamFraming;
17use crate::runtime::config::{ClientConfig, ReconnectConfig};
18use crate::runtime::http::{
19    body_to_byte_stream, map_http_error, parse_retry_after, read_error_body_prefix,
20};
21use crate::runtime::multipart::{
22    MultipartStream, boundary_from_content_type, parse_multipart_stream,
23};
24use crate::runtime::sse::{LastEventId, SseStream, StreamActivity, parse_sse_stream_with_id};
25use crate::runtime::transport_error::TransportError;
26
27/// Boxed byte stream produced from an HTTP response body for framing parsers.
28type BoxByteStream = Pin<
29    Box<
30        dyn Stream<Item = Result<bytes::Bytes, Box<dyn std::error::Error + Send + Sync + 'static>>>
31            + Send,
32    >,
33>;
34
35/// Send a unary request and decode the JSON response.
36///
37/// `build` is a closure returning a `RequestBuilder` configured with method,
38/// URL, headers, and any auth state. It is invoked once per attempt.
39///
40/// `timeout` bounds the **whole** attempt — connection, response headers, and
41/// reading the response body — as a per-attempt deadline (mirroring the
42/// streaming path). Elapse maps to [`TransportError::Timeout`], which is
43/// transient, so a `#[retryable]` method retries a timed-out attempt. `None`
44/// leaves the deadline to the underlying transport.
45///
46/// # Errors
47/// Returns [`TransportError`] when the builder closure fails, the deadline elapses,
48/// the network call fails, the response body cannot be read, JSON deserialization of a
49/// success body fails, or the server returns a non-success HTTP status (mapped via
50/// [`map_http_error`]).
51pub async fn send_unary<F, R>(build: F, timeout: Option<Duration>) -> Result<R, TransportError>
52where
53    F: FnOnce() -> Result<RequestBuilder, TransportError>,
54    R: DeserializeOwned,
55{
56    let builder = build()?;
57    let op = async move {
58        let response = builder.send().await.map_err(map_send_error)?;
59        let status = response.status();
60        if status.is_success() {
61            // `HttpResponse::bytes` does NOT enforce status check; we already did.
62            // Avoiding `.json()` here because it would re-check status and
63            // duplicate work on the success path.
64            let bytes = response.bytes().await.map_err(TransportError::network)?;
65            // A `204 No Content` (or any 2xx with an empty body) is valid for
66            // methods returning `Result<(), _>` and other unit-like `Ok` types.
67            // `serde_json::from_slice::<()>(b"")` would fail with an EOF error
68            // (and `()` deserializes only from JSON `null`), so map an empty
69            // success body to JSON `null` before decoding — `null` deserializes
70            // into `()` and `Option::None`, while any non-unit `R` still errors.
71            let bytes: &[u8] = if bytes.is_empty() { b"null" } else { &bytes };
72            serde_json::from_slice::<R>(bytes).map_err(TransportError::serialization)
73        } else {
74            let retry_after = parse_retry_after(response.headers());
75            let bytes = response.bytes().await.map_err(TransportError::network)?;
76            let body = String::from_utf8_lossy(&bytes).into_owned();
77            Err(map_http_error(status.as_u16(), body, retry_after))
78        }
79    };
80    match timeout {
81        Some(d) => tokio::time::timeout(d, op)
82            .await
83            .map_err(|_| TransportError::Timeout(d))?,
84        None => op.await,
85    }
86}
87
88/// Shared transport-construction chain for both cfg variants of
89/// [`build_default_http_client`] — in one place so a knob can't be added to one
90/// arm and forgotten in the other.
91///
92/// Reads only transport-construction fields; call-time fields (`retry`,
93/// streaming, `internal_token_provider`) are applied per call and must not be
94/// wired here. `timeout` is wired on purpose: the SDK also enforces it per call,
95/// and matching the transport's `TimeoutLayer` to it stops a `config.timeout`
96/// above toolkit-http's 30s default being silently capped at 30s.
97fn transport_builder(config: &ClientConfig) -> toolkit_http::HttpClientBuilder {
98    toolkit_http::HttpClient::builder()
99        .retry(None)
100        .transport(transport_security(config.require_tls))
101        .timeout(config.timeout)
102        .pool_max_idle_per_host(config.pool_max_idle_per_host)
103        .pool_idle_timeout(config.pool_idle_timeout)
104        .concurrency_limit(
105            config
106                .max_concurrent_requests
107                .map(|max_concurrent_requests| toolkit_http::RateLimitConfig {
108                    max_concurrent_requests,
109                }),
110        )
111}
112
113/// Build the default `toolkit-http` client used by macro-generated REST clients.
114///
115/// Transport-layer retry is **disabled** — the SDK runs its own retry loop in
116/// [`retry_with_backoff`](crate::runtime::retry::retry_with_backoff).
117///
118/// Only the *transport-construction* fields of `config` are read — TLS mode,
119/// the connection-pool and concurrency knobs, and the per-attempt `timeout`;
120/// call-time fields (`retry`, the streaming knobs, `internal_token_provider`)
121/// are applied per call by the generated client, not on the transport. In
122/// particular `config.require_tls` selects the transport security mode: `false`
123/// (the default via
124/// [`ClientConfig::new`](crate::runtime::config::ClientConfig::new)) allows
125/// plaintext `http://`, preserving the platform's existing in-mesh
126/// service-to-service convention; `true`
127/// ([`ClientConfig::with_require_tls`](crate::runtime::config::ClientConfig::with_require_tls))
128/// switches to `toolkit_http::TransportSecurity::TlsOnly`, rejecting
129/// plaintext. This matters because the tenant bearer token is forwarded via
130/// `Authorization` on whatever scheme the resolved endpoint uses — set
131/// `require_tls` when a client may talk to an endpoint outside a trusted
132/// network boundary.
133///
134/// `.with_otel()` is called in BOTH cfg variants below — it always installs
135/// `toolkit-http`'s `OtelLayer`, but the layer's actual W3C `traceparent`
136/// injection is itself `#[cfg(feature = "otel")]` **inside `toolkit-http`**
137/// (a no-op otherwise). Because Cargo unifies features workspace-wide, that
138/// gate tracks whether *anything* in the final binary enables
139/// `toolkit-http/otel` — which may not be this crate's own `otel` feature.
140/// The one thing genuinely gated by *this* crate's `otel` feature (which
141/// forwards `toolkit-http/otel`) is `.with_metrics(client_type)` below, since
142/// that builder method only exists when the dependency is compiled with the
143/// feature on. The generated method's per-method `tracing` span is always
144/// opened either way, independent of any of this.
145///
146/// # Errors
147/// Propagates [`toolkit_http::HttpError`] from the underlying builder (e.g. a
148/// TLS backend that cannot be constructed under FIPS).
149#[cfg(feature = "otel")]
150pub fn build_default_http_client(
151    client_type: &str,
152    config: &ClientConfig,
153) -> Result<toolkit_http::HttpClient, toolkit_http::HttpError> {
154    transport_builder(config)
155        .with_otel()
156        .with_metrics(client_type)
157        .build()
158}
159
160/// Non-`otel` build of [`build_default_http_client`]: RED metrics
161/// (`.with_metrics`) are NOT compiled in — that builder method only exists
162/// under `toolkit-http/otel`. `client_type` is unused in this configuration.
163/// `.with_otel()` is still called (see the doc on the `otel`-cfg sibling
164/// above): whether it actually propagates `traceparent` depends on whether
165/// `toolkit-http/otel` ends up enabled by *some* crate in the build, not on
166/// this crate's own `otel` feature.
167///
168/// # Errors
169/// Propagates [`toolkit_http::HttpError`] from the underlying builder (e.g. a
170/// TLS backend that cannot be constructed under FIPS).
171#[cfg(not(feature = "otel"))]
172pub fn build_default_http_client(
173    _client_type: &str,
174    config: &ClientConfig,
175) -> Result<toolkit_http::HttpClient, toolkit_http::HttpError> {
176    transport_builder(config).with_otel().build()
177}
178
179fn transport_security(require_tls: bool) -> toolkit_http::TransportSecurity {
180    if require_tls {
181        toolkit_http::TransportSecurity::TlsOnly
182    } else {
183        toolkit_http::TransportSecurity::AllowInsecureHttp
184    }
185}
186
187/// Map a `toolkit-http` send error to a [`TransportError`]. Two kinds get a
188/// dedicated variant instead of the generic [`TransportError::Network`]:
189///
190/// - `Overloaded`: the concurrency limiter shed the request before it was sent
191///   → [`TransportError::Overloaded`] (non-transient), so the retry loop can
192///   tell a locally-shed request from a network failure.
193/// - `Timeout` / `DeadlineExceeded`: the transport timeout layer fired. It is
194///   set to `config.timeout`, which the SDK also enforces per call, so either
195///   timer can win; both map to [`TransportError::Timeout`] to keep the
196///   classification stable.
197fn map_send_error(err: toolkit_http::HttpError) -> TransportError {
198    match err {
199        toolkit_http::HttpError::Overloaded => TransportError::Overloaded,
200        toolkit_http::HttpError::Timeout(d) | toolkit_http::HttpError::DeadlineExceeded(d) => {
201            TransportError::Timeout(d)
202        }
203        other => TransportError::network(other),
204    }
205}
206
207/// Add a JSON body to a request builder. Wraps `toolkit_http`'s fallible
208/// `.json()` (which can fail to serialize) in our [`TransportError`] surface
209/// so the macro emit path can `?` uniformly.
210///
211/// # Errors
212/// Returns [`TransportError::Serialization`] when `body` cannot be serialized to JSON.
213pub fn with_json_body<T: Serialize>(
214    builder: RequestBuilder,
215    body: &T,
216) -> Result<RequestBuilder, TransportError> {
217    builder.json(body).map_err(TransportError::serialization)
218}
219
220/// Builder for a streaming request that can be re-issued on reconnect.
221///
222/// `build` receives the latest seen `Last-Event-ID` (or `None` on the first
223/// attempt) and must return a fresh, configured `RequestBuilder`.
224/// Implementations should set the `Last-Event-ID` header from the parameter
225/// when present.
226///
227/// The parameter is permanently `None` under any framing other than SSE:
228/// `Last-Event-ID` is an SSE mechanism and no other framing has a resume
229/// token, so the signature stays framing-free rather than growing a resume-token
230/// abstraction with exactly one inhabitant.
231pub trait StreamRequestFactory: Send + 'static {
232    /// Construct a `RequestBuilder` for the next stream attempt.
233    ///
234    /// # Errors
235    /// Returns [`TransportError`] when the factory cannot produce a builder
236    /// (e.g. URL composition or auth header attachment fails).
237    fn build(&self, last_event_id: Option<&str>) -> Result<RequestBuilder, TransportError>;
238}
239
240impl<F> StreamRequestFactory for F
241where
242    F: Fn(Option<&str>) -> Result<RequestBuilder, TransportError> + Send + 'static,
243{
244    fn build(&self, last: Option<&str>) -> Result<RequestBuilder, TransportError> {
245        (self)(last)
246    }
247}
248
249/// A configured streaming request, ready to be opened.
250///
251/// Replaces what was a positional argument list on [`send_streaming`]. The
252/// knobs are independent and all but the factory have a meaningful default, so
253/// a builder keeps call sites readable as more of them arrive.
254///
255/// Defaults: [`StreamFraming::ServerSentEvents`], no reconnect
256/// ([`ReconnectConfig::disabled`]), no open deadline and no idle deadline.
257///
258/// # Stability
259/// **Not settled surface.** This exists to be driven by generated code; it has
260/// no hand-written consumer yet. Expect it to change until a real consumer has
261/// exercised it.
262pub struct StreamRequest<F> {
263    factory: F,
264    framing: StreamFraming,
265    reconnect: ReconnectConfig,
266    open_timeout: Option<Duration>,
267    idle_timeout: Option<Duration>,
268}
269
270impl<F: StreamRequestFactory> StreamRequest<F> {
271    /// Start from a request factory, with SSE framing, reconnect disabled and
272    /// no open or idle deadline.
273    #[must_use]
274    pub fn new(factory: F) -> Self {
275        Self {
276            factory,
277            framing: StreamFraming::default(),
278            reconnect: ReconnectConfig::disabled(),
279            open_timeout: None,
280            idle_timeout: None,
281        }
282    }
283
284    /// Select the wire framing the response body is parsed as.
285    ///
286    /// This does **not** set the request's `Accept` header — the factory owns
287    /// the request, so advertising the matching media type is its job. The two
288    /// come from one declaration in generated code.
289    #[must_use]
290    pub fn framing(mut self, framing: StreamFraming) -> Self {
291        self.framing = framing;
292        self
293    }
294
295    /// Set the reconnect policy applied to transient open and stream failures.
296    #[must_use]
297    pub fn reconnect(mut self, reconnect: ReconnectConfig) -> Self {
298        self.reconnect = reconnect;
299        self
300    }
301
302    /// Set the deadline for **opening** the stream: the connect, the response
303    /// headers, and (on a non-success status) the error-body read. This is a
304    /// bound on the open handshake only, distinct from
305    /// [`idle_timeout`](Self::idle_timeout), which bounds the gap between items
306    /// once the stream is live. Without it, a peer that accepts the socket and
307    /// never answers hangs the open forever. Generated code defaults this from
308    /// the client's unary [`timeout`](crate::runtime::config::ClientConfig::timeout).
309    #[must_use]
310    pub fn open_timeout(mut self, open_timeout: Duration) -> Self {
311        self.open_timeout = Some(open_timeout);
312        self
313    }
314
315    /// Set the per-item idle deadline. This is *idle*, not total: any wire
316    /// chunk resets it, including ones that dispatch no item.
317    #[must_use]
318    pub fn idle_timeout(mut self, idle_timeout: Duration) -> Self {
319        self.idle_timeout = Some(idle_timeout);
320        self
321    }
322}
323
324/// Why one open attempt failed, and whether the reconnect budget applies.
325///
326/// The distinction is not derivable from [`TransportError::is_transient`]: a
327/// request-shape failure is a configuration bug that will fail identically on
328/// every attempt, so it is fatal regardless of how its error class is
329/// otherwise classified.
330enum OpenFailure {
331    /// Request-shape failure (URL composition, auth attach), or an error-body
332    /// read that itself failed. Never retried.
333    Fatal(TransportError),
334    /// Connect, send, or non-success status. Retried when the error is
335    /// transient and budget remains.
336    Retryable(TransportError),
337}
338
339/// Send one request and await response headers, bounded by `timeout`.
340// cancel-safe: holds only the in-flight send; being dropped (the `timeout`
341// below wraps `send`, so this is the future cancelled on elapse) abandons the
342// attempt with nothing committed and no buffered state to lose.
343async fn send_open_request(
344    builder: RequestBuilder,
345    timeout: Option<Duration>,
346) -> Result<toolkit_http::HttpResponse, OpenFailure> {
347    let send_fut = builder.send();
348    let sent = match timeout {
349        Some(d) => match tokio::time::timeout(d, send_fut).await {
350            Ok(r) => r,
351            Err(_) => return Err(OpenFailure::Retryable(TransportError::Timeout(d))),
352        },
353        None => send_fut.await,
354    };
355    // Pre-flight failures are reconnect-eligible, except a concurrency-limiter
356    // shed (`Overloaded`), which is fail-fast: bubble it up as `Fatal`.
357    sent.map_err(|e| match map_send_error(e) {
358        overloaded @ TransportError::Overloaded => OpenFailure::Fatal(overloaded),
359        other => OpenFailure::Retryable(other),
360    })
361}
362
363/// Read and classify a non-success open response.
364// cancel-safe: reads and discards the error body only to build a message;
365// dropping it abandons that read with nothing committed. Awaited to completion
366// inside `open`.
367async fn open_status_failure(
368    response: toolkit_http::HttpResponse,
369    timeout: Option<Duration>,
370) -> OpenFailure {
371    let status_code = response.status().as_u16();
372    let retry_after = parse_retry_after(response.headers());
373    // Bound the error-body read two ways. In *time*, by the same per-attempt
374    // deadline as the initial send — otherwise a slow error body on the
375    // pre-stream path could stall the attempt indefinitely, defeating the
376    // timeout guarantee. In *size*, by reading only a bounded prefix rather
377    // than `response.bytes()` (which buffers up to the client's `max_body_size`,
378    // megabytes by default) — the body only feeds a diagnostic, so a large or
379    // hostile error body cannot force an outsized allocation here.
380    let body_fut = read_error_body_prefix(response.into_body());
381    let bytes_result = match timeout {
382        Some(d) => match tokio::time::timeout(d, body_fut).await {
383            Ok(r) => r,
384            Err(_) => return OpenFailure::Retryable(TransportError::Timeout(d)),
385        },
386        None => body_fut.await,
387    };
388    let bytes = match bytes_result {
389        Ok(b) => b,
390        // A failed error-body read ends the stream rather than consuming a
391        // reconnect attempt: we no longer know what the peer said, so
392        // re-issuing would be guessing.
393        Err(e) => return OpenFailure::Fatal(TransportError::network(e)),
394    };
395    let body = String::from_utf8_lossy(&bytes).into_owned();
396    // A transient status (e.g. 503 during a rolling deploy) is
397    // reconnect-eligible, consistent with the pre-flight failure branch; other
398    // statuses are domain errors and bubble straight through.
399    OpenFailure::Retryable(map_http_error(status_code, body, retry_after))
400}
401
402/// One open attempt against an already-built request: send it and check the
403/// status.
404///
405/// Steps 1-3 of the streaming lifecycle, shared by [`send_streaming`] and
406/// [`open_streaming`]. Performs no retry of its own.
407///
408/// Takes the built `RequestBuilder` rather than the factory deliberately: an
409/// `async fn` holding a `&F` would put that reference in the returned future,
410/// which is `Send` only if `F: Sync` — a bound [`StreamRequestFactory`] does
411/// not require and should not have to.
412async fn attempt_open(
413    builder: RequestBuilder,
414    timeout: Option<Duration>,
415) -> Result<toolkit_http::HttpResponse, OpenFailure> {
416    let response = send_open_request(builder, timeout).await?;
417    if response.status().is_success() {
418        return Ok(response);
419    }
420    Err(open_status_failure(response, timeout).await)
421}
422
423/// A successfully opened streaming response, together with whatever
424/// framing-specific setup its body needs before it can be parsed.
425///
426/// An enum rather than a struct carrying an `Option<String>` boundary, so
427/// "multipart with no boundary" is not a representable state. The multipart arm
428/// is why this type exists at all: its boundary lives in the response's
429/// `Content-Type`, so a multipart open can still fail *after* a `200` — and
430/// resolving it here, at open time, is what makes that failure an `Err` from
431/// [`open_streaming`] rather than the returned stream's first item.
432enum OpenedResponse {
433    /// SSE needs no setup beyond the body itself.
434    ServerSentEvents(toolkit_http::HttpResponse),
435    /// `multipart/mixed`, with the boundary read from the response's
436    /// `Content-Type`.
437    MultipartMixed {
438        response: toolkit_http::HttpResponse,
439        boundary: String,
440    },
441}
442
443/// Perform the framing-specific setup a success response needs.
444///
445/// # Errors
446/// [`TransportError::Framing`] when a `multipart/mixed` response carries no
447/// `Content-Type`, or one from which no boundary can be read.
448fn open_response(
449    framing: StreamFraming,
450    response: toolkit_http::HttpResponse,
451) -> Result<OpenedResponse, TransportError> {
452    match framing {
453        StreamFraming::ServerSentEvents => Ok(OpenedResponse::ServerSentEvents(response)),
454        StreamFraming::MultipartMixed => {
455            let content_type = response
456                .headers()
457                .get(http::header::CONTENT_TYPE)
458                .and_then(|v| v.to_str().ok())
459                .ok_or_else(|| {
460                    TransportError::framing(
461                        StreamFraming::MultipartMixed,
462                        "success response carries no readable `Content-Type` header",
463                    )
464                })?;
465            let boundary = boundary_from_content_type(content_type)?;
466            Ok(OpenedResponse::MultipartMixed { response, boundary })
467        }
468    }
469}
470
471/// Owns everything that persists across open attempts: the request factory,
472/// the reconnect budget, and the `Last-Event-ID` cell the SSE parser advances.
473///
474/// Exists so the retry loop can be written once and shared by the eager and
475/// lazy entry points. Its methods take `&mut self` rather than `&F` — `&mut T`
476/// is `Send` when `T: Send`, so this keeps the futures `Send` without
477/// demanding `F: Sync`.
478struct StreamOpener<F> {
479    factory: F,
480    framing: StreamFraming,
481    reconnect: ReconnectConfig,
482    /// Deadline for the open handshake (connect, headers, error-body read),
483    /// applied per attempt. Distinct from `idle_timeout`, which bounds only the
484    /// item loop once the stream is live.
485    open_timeout: Option<Duration>,
486    idle_timeout: Option<Duration>,
487    last_id: LastEventId,
488    /// Reconnect attempts consumed so far. Shared by the open and the stream
489    /// so a burst of failures with no delivered item is capped as one budget
490    /// rather than one per phase. Reset by a healthy connection.
491    attempt: u32,
492    /// Reopens performed over the whole lifetime of this stream, never reset.
493    /// Bounds a peer that keeps resetting the burst budget from reopening
494    /// forever (#4740), against `reconnect.max_total_reopens`.
495    total_reopens: u32,
496}
497
498impl<F: StreamRequestFactory> StreamOpener<F> {
499    fn new(request: StreamRequest<F>) -> Self {
500        Self {
501            factory: request.factory,
502            framing: request.framing,
503            reconnect: request.reconnect,
504            open_timeout: request.open_timeout,
505            idle_timeout: request.idle_timeout,
506            last_id: LastEventId::empty(),
507            attempt: 0,
508            total_reopens: 0,
509        }
510    }
511
512    /// Open the stream, retrying transient failures against the budget.
513    // cancel-safe: the only carried state is the budget counters in `self`,
514    // which are monotonic, so a drop mid-attempt at worst counts an attempt
515    // that wasn't retried (conservative) — never a double-open or a lost item.
516    // The driver awaits it to completion rather than racing it under a timeout.
517    async fn open(&mut self) -> Result<OpenedResponse, TransportError> {
518        loop {
519            let snapshot = match self.framing {
520                // Re-read per attempt: a reconnect replays the most recently
521                // observed `id:` field, which the parser may have advanced
522                // since the previous open.
523                StreamFraming::ServerSentEvents => self.last_id.current(),
524                // `Last-Event-ID` is an SSE mechanism and `multipart/mixed`
525                // has no resume token of its own, so the factory is handed
526                // `None` on every attempt. Stated as a branch rather than left
527                // to the cell simply never being advanced, so the reason is
528                // visible at the one place it decides anything: resuming a
529                // multipart stream is the caller's own reopen loop, not the
530                // transport's.
531                StreamFraming::MultipartMixed => None,
532            };
533            // URL-build / serialization failure on the request side — not
534            // eligible for reconnect (config-shape error).
535            let built = self.factory.build(snapshot.as_deref());
536            let attempted = match built {
537                Ok(builder) => attempt_open(builder, self.open_timeout).await,
538                Err(e) => Err(OpenFailure::Fatal(e)),
539            };
540            match attempted {
541                // A framing-setup failure (a `200` whose `Content-Type` names
542                // no usable boundary) is not retried: it is a response-shape
543                // error that repeats identically, the same argument as the
544                // request-shape branch above. Returning it from here is also
545                // what makes it an `Err` from `open_streaming` rather than the
546                // stream's first item.
547                Ok(response) => return open_response(self.framing, response),
548                Err(OpenFailure::Fatal(e)) => return Err(e),
549                Err(OpenFailure::Retryable(err)) => {
550                    if !self.consume_attempt(&err).await {
551                        return Err(err);
552                    }
553                }
554            }
555        }
556    }
557
558    /// Consume one reconnect attempt for `err`, backing off before returning
559    /// `true`. Returns `false` when the budget is spent or the error class is
560    /// not retryable, meaning the caller must surface `err`.
561    // cancel-safe: the counter increments precede the backoff `sleep`, so a drop
562    // during the sleep leaves the budget consistently decremented (conservative)
563    // rather than in a half-updated state. No wire data flows through here.
564    async fn consume_attempt(&mut self, err: &TransportError) -> bool {
565        // Absolute lifetime ceiling, checked before the burst budget so it also
566        // stops a peer that keeps resetting that budget. It survives
567        // `reset_budget` (which only clears `attempt`), so it is the one bound a
568        // healthy-looking flapping peer cannot escape.
569        if self.total_reopens >= self.reconnect.max_total_reopens {
570            tracing::warn!(
571                total_reopens = self.total_reopens,
572                max_total_reopens = self.reconnect.max_total_reopens,
573                framing = ?self.framing,
574                error = %err,
575                "stream lifetime reopen cap reached; surfacing error",
576            );
577            return false;
578        }
579        if self.attempt < self.reconnect.max_attempts && err.is_transient() {
580            self.attempt += 1;
581            self.total_reopens += 1;
582            // A server-advised `Retry-After` (from a 429/503 open failure) wins
583            // over computed backoff, clamped to `max_delay` so a hostile or
584            // misconfigured peer cannot stall the reopen indefinitely — the same
585            // preference the unary retry loop applies in `retry::next_delay`.
586            let delay = match err.retry_after() {
587                Some(advised) => advised.min(self.reconnect.max_delay),
588                None => backoff_delay(&self.reconnect, self.attempt),
589            };
590            // The single choke point for every reconnect: an operator otherwise
591            // sees only the final error and nothing when a reopen succeeds.
592            tracing::warn!(
593                attempt = self.attempt,
594                max_attempts = self.reconnect.max_attempts,
595                delay_ms = u64::try_from(delay.as_millis()).unwrap_or(u64::MAX),
596                framing = ?self.framing,
597                error = %err,
598                "stream failed; reconnecting after backoff",
599            );
600            tokio::time::sleep(delay).await;
601            true
602        } else {
603            // A non-transient error, or the budget is spent: no more reopens.
604            if self.attempt >= self.reconnect.max_attempts && err.is_transient() {
605                tracing::warn!(
606                    attempt = self.attempt,
607                    max_attempts = self.reconnect.max_attempts,
608                    framing = ?self.framing,
609                    error = %err,
610                    "stream reconnect budget exhausted; surfacing error",
611                );
612            }
613            false
614        }
615    }
616
617    /// Reset the budget after a connection that delivered items.
618    fn reset_budget(&mut self) {
619        self.attempt = 0;
620    }
621}
622
623/// The framing parser driving one opened response body, behind the single
624/// surface [`drive`] needs.
625///
626/// A `dyn` decoder trait is not available here — decoding is generic in `T`, so
627/// the trait could not be object-safe — and the two framings genuinely differ
628/// in only three places, so an enum keeps each difference visible at the point
629/// it is decided rather than dispersing it behind a vtable.
630enum FramedItems<T> {
631    Sse(SseStream<T, BoxByteStream>),
632    Multipart(MultipartStream<T, BoxByteStream>),
633}
634
635impl<T> FramedItems<T>
636where
637    T: DeserializeOwned + 'static,
638{
639    /// Build the parser for an opened response.
640    ///
641    /// `last_id` is consumed **only** by the SSE arm: the cell is
642    /// SSE-parser-owned and multipart never reads or writes it.
643    fn new(opened: OpenedResponse, last_id: LastEventId) -> Self {
644        match opened {
645            OpenedResponse::ServerSentEvents(response) => {
646                // The parsers need `Unpin + 'static`; pinning on the stack with
647                // `pin_mut!` would borrow the byte stream for less than
648                // `'static`, so move ownership behind `Box::pin` and hand the
649                // boxed stream over.
650                let byte_stream: BoxByteStream =
651                    Box::pin(body_to_byte_stream(response.into_body()));
652                Self::Sse(parse_sse_stream_with_id::<T, _, _>(byte_stream, last_id))
653            }
654            OpenedResponse::MultipartMixed { response, boundary } => {
655                ::std::mem::drop(last_id);
656                let byte_stream: BoxByteStream =
657                    Box::pin(body_to_byte_stream(response.into_body()));
658                Self::Multipart(parse_multipart_stream::<T, _, _>(byte_stream, &boundary))
659            }
660        }
661    }
662
663    /// Wire-activity counter, so the idle deadline stays *idle* rather than
664    /// merely *quiet* under either framing.
665    fn activity_handle(&self) -> StreamActivity {
666        match self {
667            Self::Sse(s) => s.activity_handle(),
668            Self::Multipart(s) => s.activity_handle(),
669        }
670    }
671
672    // cancel-safe: the parser state (buffered wire bytes, partial frame) lives
673    // in `self`, not this future — `drive` drops and re-creates this future
674    // around its idle `timeout` on every quiet window, and resumes without
675    // losing bytes. Both arms delegate to a `Stream::next` that only reads into
676    // the stream's own buffer. This is the load-bearing one: if it regressed to
677    // holding partial state in the future, every idle elapse would drop bytes.
678    async fn next(&mut self) -> Option<Result<T, TransportError>> {
679        use futures_util::StreamExt;
680        match self {
681            Self::Sse(s) => s.next().await,
682            Self::Multipart(s) => s.next().await,
683        }
684    }
685
686    /// What a **graceful** end of the underlying byte stream means for this
687    /// framing. `None` is a clean end; `Some` ends the stream with that error,
688    /// which is transient and so reconnect-eligible.
689    ///
690    /// This is the one place the two framings disagree about success, and
691    /// getting it backwards is silent in both directions — every completed
692    /// multipart stream would become an error, or every truncated SSE stream a
693    /// success.
694    fn end_of_stream_error(&self) -> Option<TransportError> {
695        match self {
696            // SSE has no application-level terminator, so the transport must
697            // supply one: an `event: done` frame and a bare connection close
698            // both surface as an end of stream. Treat a close WITHOUT an
699            // explicit `done` as an anomaly rather than silently reporting
700            // success — otherwise a server restart / LB idle-timeout / rolling
701            // deploy that closes the connection mid-stream would look identical
702            // to a fully-delivered stream.
703            Self::Sse(s) => (!s.saw_done_event()).then(|| {
704                TransportError::framing(
705                    StreamFraming::ServerSentEvents,
706                    "SSE stream ended without a terminal `done` event",
707                )
708            }),
709            // A complete multipart body ends with the `--<boundary>--` close
710            // delimiter (RFC 2046). A graceful EOF that did NOT see it is a
711            // truncation — a proxy idle-timeout, an LB half-close, a rolling
712            // deploy closing the connection mid-body — so report it, mirroring
713            // the SSE arm (#4740). On the public path the caller has no other
714            // way to tell truncation from completion: `send_streaming` /
715            // `open_streaming` hand back a boxed stream that erases
716            // `MultipartStream::saw_close_delimiter`, and a generic item type
717            // carries no terminal marker of its own. An *aborted* body is a
718            // different thing and still arrives as a `Network` error item, not
719            // as an end of stream.
720            Self::Multipart(s) => (!s.saw_close_delimiter()).then(|| {
721                TransportError::framing(
722                    StreamFraming::MultipartMixed,
723                    "multipart body ended without a closing `--<boundary>--` delimiter",
724                )
725            }),
726        }
727    }
728}
729
730/// Parse one already-opened response body and yield its items.
731///
732/// Steps 4-5 of the streaming lifecycle. Ends after yielding at most one
733/// `Err`; a clean end yields nothing further. Reconnect is the caller's
734/// concern.
735fn drive<T>(
736    opened: OpenedResponse,
737    last_id: LastEventId,
738    idle_timeout: Option<Duration>,
739) -> impl Stream<Item = Result<T, TransportError>> + Send
740where
741    T: DeserializeOwned + Send + 'static,
742{
743    async_stream::stream! {
744        let mut inner = FramedItems::<T>::new(opened, last_id);
745        let activity = inner.activity_handle();
746        loop {
747            // Idle timeout is *idle*: keepalive comments (and any other wire
748            // chunk that dispatches no item) still count as activity, so a
749            // quiet-but-alive stream is not torn down. On elapse we only
750            // error if no chunk arrived while we waited.
751            let item = match idle_timeout {
752                Some(d) => loop {
753                    let before = activity.generation();
754                    match tokio::time::timeout(d, inner.next()).await {
755                        Ok(v) => break v,
756                        // Activity advanced since the wait started (e.g. a
757                        // keepalive comment) — loop again rather than
758                        // treating this as a genuine idle timeout. Falling
759                        // off this arm already re-enters the loop; an
760                        // explicit `continue` here is redundant.
761                        Err(_) if activity.generation() != before => {}
762                        Err(_) => {
763                            yield Err(TransportError::Timeout(d));
764                            return;
765                        }
766                    }
767                },
768                None => inner.next().await,
769            };
770            match item {
771                Some(Ok(v)) => yield Ok(v),
772                Some(Err(e)) => {
773                    yield Err(e);
774                    return;
775                }
776                None => {
777                    // The underlying byte stream ended with no error. Whether
778                    // that is success depends on the framing — see
779                    // `FramedItems::end_of_stream_error`.
780                    if let Some(e) = inner.end_of_stream_error() {
781                        yield Err(e);
782                    }
783                    return;
784                }
785            }
786        }
787    }
788}
789
790/// Drive an opened response to completion, re-opening on transient failures.
791///
792/// `first` is the already-opened response; subsequent attempts re-open through
793/// [`StreamOpener::open`].
794fn drive_with_reconnect<F, T>(
795    mut opener: StreamOpener<F>,
796    first: OpenedResponse,
797) -> impl Stream<Item = Result<T, TransportError>> + Send
798where
799    F: StreamRequestFactory,
800    T: DeserializeOwned + Send + 'static,
801{
802    use futures_util::StreamExt;
803
804    async_stream::try_stream! {
805        let mut pending = Some(first);
806        loop {
807            let response = match pending.take() {
808                Some(r) => r,
809                // Only reached on a reopen (`pending` is `Some` on the first
810                // iteration), so a success here is a recovered reconnect.
811                None => match opener.open().await {
812                    Ok(r) => {
813                        tracing::debug!(
814                            attempt = opener.attempt,
815                            framing = ?opener.framing,
816                            "stream reconnect succeeded",
817                        );
818                        r
819                    }
820                    Err(e) => {
821                        Err(e)?;
822                        return;
823                    }
824                },
825            };
826
827            // Measured from the moment the connection is live (its response is
828            // open) to when its byte stream ends, so a one-item-then-drop peer
829            // registers a near-zero uptime.
830            let connection_started = tokio::time::Instant::now();
831            let mut inner = Box::pin(
832                drive::<T>(response, opener.last_id.clone(), opener.idle_timeout),
833            );
834            let mut stream_err: Option<TransportError> = None;
835            // Whether this connection ever delivered an event. Necessary but not
836            // sufficient for "healthy" — see the reset gate below.
837            let mut delivered_an_event = false;
838            while let Some(item) = inner.next().await {
839                match item {
840                    Ok(v) => {
841                        delivered_an_event = true;
842                        yield v;
843                    }
844                    Err(e) => {
845                        stream_err = Some(e);
846                        break;
847                    }
848                }
849            }
850            let uptime = connection_started.elapsed();
851            // Drop the parser before re-opening: it owns the previous
852            // connection's body, and the next attempt gets a fresh one.
853            ::std::mem::drop(inner);
854
855            // A *healthy* connection resets the burst budget so `max_attempts`
856            // is a burst cap, not a lifetime cap: a long-lived subscription that
857            // survives N unrelated blips over days should not die on the N+1st.
858            // "Healthy" needs both a delivered item AND a minimum uptime —
859            // delivering one item then dropping immediately is the #4740 peer
860            // that would otherwise reset the budget forever and reopen at
861            // `base_delay` indefinitely, re-sending the auth token each time. A
862            // too-brief connection instead counts against the budget, and the
863            // absolute `max_total_reopens` cap (in `consume_attempt`) backstops
864            // a peer that games the uptime threshold.
865            if delivered_an_event && uptime >= opener.reconnect.min_healthy_uptime {
866                opener.reset_budget();
867            }
868
869            match stream_err {
870                // Stream ended cleanly, by whatever its framing calls clean.
871                None => return,
872                Some(e) => {
873                    if !opener.consume_attempt(&e).await {
874                        Err(e)?;
875                        return;
876                    }
877                    // Budget consumed and backoff applied — fall through to
878                    // the next iteration, which re-opens.
879                }
880            }
881        }
882    }
883}
884
885/// Open a stream **eagerly**: connect, check the response status, and perform
886/// the framing's own setup before returning, so an open-time failure is an
887/// `Err` from this call rather than the returned stream's first item.
888///
889/// This is the entry point for a contract method declared
890/// `#[streaming] async fn` — a stream whose open is a distinct, fallible
891/// operation. Open-time failures (a `404`, a `409` carrying a domain
892/// `Problem`, a `410`) reach the caller before any item exists, where its
893/// open-time error handling can act on them.
894///
895/// The request's reconnect policy governs the *open* as well as the stream:
896/// with [`ReconnectConfig::disabled`] (the [`StreamRequest::new`] default, and
897/// what generated code passes for a fallible open per D6) exactly one open
898/// attempt is made.
899///
900/// # Errors
901/// Returns [`TransportError`] when the request cannot be built, the connect
902/// fails or times out, the response status is not a success, or the framing's
903/// setup fails — for [`StreamFraming::MultipartMixed`] that last case is a
904/// `200` whose `Content-Type` names no usable boundary, which is why a
905/// multipart open can fail after a success status. Failures after a successful
906/// open arrive as items of the returned stream.
907///
908/// # Stability
909/// **Not settled surface** — see [`StreamRequest`].
910pub async fn open_streaming<F, T>(
911    request: StreamRequest<F>,
912) -> Result<Pin<Box<dyn Stream<Item = Result<T, TransportError>> + Send>>, TransportError>
913where
914    F: StreamRequestFactory,
915    T: DeserializeOwned + Send + 'static,
916{
917    let mut opener = StreamOpener::new(request);
918    let response = opener.open().await?;
919    Ok(Box::pin(drive_with_reconnect(opener, response)))
920}
921
922/// Send a streaming request and adapt the response into a typed stream, with
923/// the open deferred until the stream is first polled.
924///
925/// The returned stream yields `Result<T, TransportError>` items, one per SSE
926/// event or per `multipart/mixed` part. With `reconnect.max_attempts == 0` (the
927/// default), a transient transport failure ends the stream immediately. With a
928/// non-zero limit, the client re-issues the request up to `max_attempts` times,
929/// applying exponential backoff between attempts — replaying `Last-Event-ID`
930/// under SSE framing, and with no resume token under any other. A
931/// `multipart/mixed` reopen therefore replays the body from its first part,
932/// redelivering any items already yielded (at-least-once); see
933/// [`ClientConfig::stream_reconnect`](crate::runtime::config::ClientConfig::stream_reconnect).
934///
935/// The open is **lazy**: nothing is sent until the returned stream is first
936/// polled, so a connect failure, a non-success status or a framing-setup
937/// failure arrives as the stream's first item. That is the shape a
938/// `#[streaming] fn` method needs, since it has nowhere else to put an error.
939/// Use [`open_streaming`] when the open must be able to fail on its own.
940///
941/// # Stability
942/// The [`StreamRequest`] parameter is **not settled surface** — see
943/// [`StreamRequest`]. This entry point took a positional argument list before
944/// a framing selector existed.
945pub fn send_streaming<F, T>(
946    request: StreamRequest<F>,
947) -> Pin<Box<dyn Stream<Item = Result<T, TransportError>> + Send>>
948where
949    F: StreamRequestFactory,
950    T: DeserializeOwned + Send + 'static,
951{
952    use futures_util::StreamExt;
953
954    Box::pin(async_stream::try_stream! {
955        // Deferred into the stream so this entry point stays lazy.
956        let mut inner = match open_streaming::<F, T>(request).await {
957            Ok(s) => s,
958            Err(e) => {
959                Err(e)?;
960                return;
961            }
962        };
963        while let Some(item) = inner.next().await {
964            yield item?;
965        }
966    })
967}
968
969/// Compute the (jittered) backoff delay for reconnect attempt #N (1-indexed).
970/// Doubles the base delay each attempt, capped at `max_delay`, then multiplies
971/// by a ±25% jitter factor — fleet-wide reconnect synchronization (many clients
972/// reconnecting in lockstep after a shared upstream blip) is the real concern.
973///
974/// Pure so the caller can both log the chosen delay and sleep the same value —
975/// the jitter is random, so it must be drawn exactly once.
976fn backoff_delay(config: &ReconnectConfig, attempt: u32) -> Duration {
977    use rand::RngExt;
978    let exp = attempt.saturating_sub(1);
979    let multiplier = 2u32.saturating_pow(exp);
980    let base = config
981        .base_delay
982        .saturating_mul(multiplier)
983        .min(config.max_delay);
984    let jitter: f64 = rand::rng().random_range(0.75..=1.25);
985    let secs = base.as_secs_f64() * jitter;
986    if secs.is_finite() && secs >= 0.0 {
987        Duration::from_secs_f64(secs).min(config.max_delay)
988    } else {
989        base
990    }
991}
992
993// `with_json_body` and the streaming path are exercised end-to-end (real
994// serialized body observed by a real server) by the integration tests in
995// `tests/rest_client_codegen.rs` (`unary_post_round_trip` et al.) — a local
996// unit test here could only check "builds without panicking" without
997// access to `toolkit_http::RequestBuilder`'s private fields, which is
998// strictly weaker than the existing round-trip coverage.