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.