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