Skip to main content

toolkit_contract/runtime/
config.rs

1//! Client, retry, and credential configuration consumed by generated clients.
2//!
3//! [`ClientConfig`] is shared by both the generated REST client and the
4//! generated gRPC client; [`InternalTokenProvider`] is the runtime source of the
5//! platform-plane credential those clients attach on `PlatformSecurityContext`
6//! methods.
7
8use std::borrow::Cow;
9use std::sync::Arc;
10use std::time::Duration;
11
12use secrecy::SecretString;
13
14/// Outcome of resolving the process's platform-plane credential on one outbound
15/// call.
16///
17/// Three-state (rather than `Option<SecretString>`) so an attach site can tell
18/// an intentionally unauthenticated deployment (Profile 1) apart from a broken
19/// credential source — e.g. the projected token file is transiently empty.
20/// Attach helpers stay silent on [`Self::NotConfigured`] and `warn!` on
21/// [`Self::Unavailable`], never emitting the token.
22#[derive(Debug)]
23pub enum CredentialState {
24    /// No credential configured (Profile 1 / `InternalCredential::None`). Attach
25    /// nothing, silently — a legitimate deployment.
26    NotConfigured,
27    /// A credential is configured and currently available; attach it.
28    Available(SecretString),
29    /// Configured but currently unavailable (empty token file, or the background
30    /// refresh has not run yet). Attach nothing but **warn** — a broken source,
31    /// not an intentional opt-out. Carries the reason (never the token).
32    Unavailable(Cow<'static, str>),
33}
34
35/// Source of the process's **platform-plane** internal credential.
36///
37/// Generated clients attach it — as the `X-ToolKit-Internal-Token` header /
38/// metadata, **never** `Authorization` — on methods whose plane marker is
39/// `PlatformSecurityContext` (`cpt-cf-adr-two-plane-auth`). The credential comes
40/// from the runtime (the bootstrap-selected `InternalCredential`), never the
41/// contract argument.
42///
43/// Invoked on every call so a rotating credential (e.g. a projected Kubernetes
44/// `ServiceAccount` token) is always attached in its current form; it returns a
45/// [`CredentialState`] to distinguish not-configured (silent) from unavailable
46/// (warn). Because it runs per-request on an async path, the closure **must not
47/// block, do I/O, or take a contended lock** (see [`InternalTokenProvider::new`]).
48#[derive(Clone)]
49pub struct InternalTokenProvider(Arc<dyn Fn() -> CredentialState + Send + Sync>);
50
51impl InternalTokenProvider {
52    /// Build a provider whose credential is resolved by `provider` on each call
53    /// (supports rotation).
54    ///
55    /// The closure **must not block, do I/O, or take a contended lock** — it is
56    /// called on every outbound platform-plane request from an async path. See
57    /// [`InternalTokenProvider`] and use `ServiceAccountTokenReader::token_provider`
58    /// as the reference pattern for a rotating credential.
59    #[must_use]
60    pub fn new(provider: impl Fn() -> CredentialState + Send + Sync + 'static) -> Self {
61        Self(Arc::new(provider))
62    }
63
64    /// Build a provider that always yields the given static `token`.
65    ///
66    /// Suitable for a non-rotating credential (e.g. a shared secret); prefer
67    /// [`InternalTokenProvider::new`] for rotating tokens.
68    #[must_use]
69    pub fn from_token(token: SecretString) -> Self {
70        Self::new(move || CredentialState::Available(token.clone()))
71    }
72
73    /// Resolve the current credential state.
74    #[must_use]
75    pub fn current(&self) -> CredentialState {
76        (self.0)()
77    }
78
79    /// Resolve the token to attach on an outbound platform-plane call, applying
80    /// the shared attach policy so the REST and gRPC helpers behave identically:
81    /// `None`/[`NotConfigured`](CredentialState::NotConfigured) → `None` (silent),
82    /// [`Available`](CredentialState::Available) → `Some`, and
83    /// [`Unavailable`](CredentialState::Unavailable) → `None` plus a `warn!`
84    /// naming the plane and `rpc` (never the token).
85    #[must_use]
86    pub fn resolve_for_attach(provider: Option<&Self>, rpc: &str) -> Option<SecretString> {
87        match provider.map(Self::current) {
88            None | Some(CredentialState::NotConfigured) => None,
89            Some(CredentialState::Available(token)) => Some(token),
90            Some(CredentialState::Unavailable(reason)) => {
91                tracing::warn!(
92                    plane = "platform",
93                    rpc,
94                    reason = %reason,
95                    "platform-plane credential configured but currently unavailable; \
96                     sending request without the internal token",
97                );
98                None
99            }
100        }
101    }
102}
103
104impl std::fmt::Debug for InternalTokenProvider {
105    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
106        // Never render the credential (or even hint at its presence beyond the
107        // opaque marker) so it cannot leak through a `{:?}` sink.
108        f.write_str("InternalTokenProvider(<fn>)")
109    }
110}
111
112/// Base configuration for a generated REST client.
113///
114/// `#[non_exhaustive]`: construct via [`ClientConfig::new`] and the `with_*`
115/// chain rather than a struct literal, so future transport knobs can be added
116/// without a breaking change.
117#[derive(Debug, Clone)]
118#[non_exhaustive]
119pub struct ClientConfig {
120    /// Base URL prefix (e.g., `https://billing.internal`).
121    /// Combined with the base path declared in the projection trait.
122    pub base_url: String,
123    /// Deadline applied to a **single** unary attempt — NOT to the whole logical
124    /// call. A `#[retryable]` method may make up to `retry.max_attempts` attempts,
125    /// so the worst-case wall-clock for a logical call is bounded by
126    /// `max_attempts × (timeout + retry.max_delay)` (the per-retry backoff is
127    /// itself clamped to [`RetryConfig::max_delay`], including a server-advised
128    /// `Retry-After`). There is deliberately no separate whole-call budget field.
129    pub timeout: Duration,
130    /// Per-**item** idle deadline for streams of any framing: the maximum gap
131    /// between two received wire chunks before the stream is treated as timed
132    /// out. A long-lived stream is NOT bounded by [`timeout`](Self::timeout)
133    /// (which would kill a healthy slow stream); it is bounded by this larger
134    /// idle deadline instead. Defaults to 60s (> the unary default).
135    ///
136    /// This is *idle*, not *quiet*: any wire chunk resets it, including ones
137    /// that dispatch no item (an SSE keepalive comment, a multipart part header
138    /// block arriving on its own).
139    pub stream_idle_timeout: Duration,
140    /// Retry policy applied to methods marked `#[retryable]`.
141    pub retry: RetryConfig,
142    /// Reconnect policy for streams of any framing. By default
143    /// `max_attempts: 0` — stream failures bubble up unchanged. Set explicitly
144    /// to opt into transparent re-open on a transient failure.
145    ///
146    /// Two limits are deliberate rather than accidental:
147    ///
148    /// - **It applies only to a method whose open is immediate**
149    ///   (`#[streaming] fn`). A fallible open (`#[streaming] async fn`) carries
150    ///   domain semantics the client must not blindly repeat — a re-open can
151    ///   collide with an exclusion lease the first open acquired, and its
152    ///   failure would land as a stream item, past the caller's open-time error
153    ///   handling. Generated code therefore passes
154    ///   [`ReconnectConfig::disabled`] for a fallible open regardless of this
155    ///   value.
156    /// - **Resume via `Last-Event-ID` is SSE-only.** A reconnected
157    ///   `multipart/mixed` stream re-issues the original request with no resume
158    ///   token, because the framing has none. The transport still reopens it,
159    ///   but that is a blind restart, not a resume: the server replays the body
160    ///   from its first part, so any items already delivered before the failure
161    ///   are **delivered again** (at-least-once, with no marker for the
162    ///   restart). Enable reconnect for a multipart stream only where the
163    ///   consumer tolerates duplicates; one needing exactly-once must instead
164    ///   leave reconnect disabled and run its own reopen loop with
165    ///   application-level dedup.
166    pub stream_reconnect: ReconnectConfig,
167    /// When `true`, the generated client refuses plaintext `http://` and
168    /// requires TLS (`toolkit_http::TransportSecurity::TlsOnly`) for every
169    /// request — including the bearer-carrying `Authorization` header, which
170    /// otherwise would ride whatever scheme `base_url` uses. Defaults to
171    /// `false`, preserving the platform's existing in-mesh service-to-service
172    /// convention where plaintext HTTP inside a secured network boundary is an
173    /// accepted, deliberate choice (see
174    /// [`build_default_http_client`](crate::runtime::client::build_default_http_client)).
175    /// Set this when a resolved endpoint may cross an untrusted network. Read by
176    /// both the REST and gRPC transports.
177    pub require_tls: bool,
178    /// Maximum *idle* keep-alive connections retained **per upstream host**
179    /// (active in-flight requests are not capped). Defaults to 128; raise via
180    /// [`ClientTuning`](crate::wiring::ClientTuning) for higher concurrency.
181    /// **REST transport only.**
182    ///
183    /// Keep it at or above the expected per-upstream concurrency: below that,
184    /// hyper closes excess connections as they idle and reopens them per
185    /// request, producing a `connect(2)` storm that dominates CPU. The 128
186    /// default clears the ~100-concurrent gear-to-gear traffic that motivated it
187    /// (the old `toolkit-http` default of 32 did not).
188    pub pool_max_idle_per_host: usize,
189    /// How long an idle keep-alive connection is retained before it is closed —
190    /// the companion of [`pool_max_idle_per_host`](Self::pool_max_idle_per_host)
191    /// (which bounds *how many*). Keep it above the gap between bursts to an
192    /// upstream so connections stay warm. Defaults to 90s. **REST transport only.**
193    ///
194    /// `None` does **not** mean "kept indefinitely": it leaves the hyper-util
195    /// setter unset, so hyper-util's own default (~90s) applies. The default is
196    /// an explicit `Some(90s)` for that reason.
197    pub pool_idle_timeout: Option<Duration>,
198    /// Maximum in-flight requests through this client at once (across all
199    /// upstream hosts). `None` disables the limiter; `Some(n)` caps at `n`, with
200    /// `Some(0)` clamped to 1 by the transport so the client can't wedge
201    /// shedding everything. Defaults to `Some(128)`, aligned with
202    /// [`pool_max_idle_per_host`](Self::pool_max_idle_per_host) so the idle pool
203    /// is fully reusable before load is shed. **REST transport only.**
204    ///
205    /// The cap bounds requests *waiting on response headers* — the tower permit
206    /// is released once headers arrive, so a long-lived SSE/multipart body holds
207    /// no slot while it streams; size it against in-flight requests, not open
208    /// streams. A shed request surfaces as
209    /// [`TransportError::Overloaded`](crate::runtime::transport_error::TransportError::Overloaded),
210    /// which is **not** transient, so a saturated client fails fast rather than
211    /// retrying into its own overload.
212    pub max_concurrent_requests: Option<usize>,
213    /// Source of the platform-plane internal credential attached to methods
214    /// whose plane marker is `PlatformSecurityContext` (carried as
215    /// `X-ToolKit-Internal-Token`). `None` (the default) attaches nothing —
216    /// legitimate for Profile 1 / in-process (`InternalCredential::None`);
217    /// the requirement is enforced server-side. The bootstrap layer populates
218    /// this from the process's selected `InternalCredential`. Tenant-plane
219    /// methods (`SecurityContext`) never consult it; they forward the caller's
220    /// bearer token from the argument.
221    pub internal_token_provider: Option<InternalTokenProvider>,
222}
223
224impl ClientConfig {
225    /// Create a new config with sensible defaults.
226    #[must_use]
227    pub fn new(base_url: impl Into<String>) -> Self {
228        Self {
229            base_url: base_url.into(),
230            timeout: Duration::from_secs(30),
231            stream_idle_timeout: Duration::from_mins(1),
232            retry: RetryConfig::default(),
233            stream_reconnect: ReconnectConfig::default(),
234            require_tls: false,
235            // `build_default_http_client` always sets these on the builder, so
236            // toolkit-http's own defaults (32/90s/100) never apply here.
237            pool_max_idle_per_host: 128,
238            pool_idle_timeout: Some(Duration::from_secs(90)),
239            max_concurrent_requests: Some(128),
240            internal_token_provider: None,
241        }
242    }
243
244    /// Override the per-call (unary) timeout.
245    #[must_use]
246    pub fn with_timeout(mut self, timeout: Duration) -> Self {
247        self.timeout = timeout;
248        self
249    }
250
251    /// Override the per-item stream idle deadline (max gap between wire
252    /// chunks). See [`Self::stream_idle_timeout`].
253    #[must_use]
254    pub fn with_stream_idle_timeout(mut self, idle: Duration) -> Self {
255        self.stream_idle_timeout = idle;
256        self
257    }
258
259    /// Override the retry policy.
260    #[must_use]
261    pub fn with_retry(mut self, retry: RetryConfig) -> Self {
262        self.retry = retry;
263        self
264    }
265
266    /// Override the stream reconnect policy. Use
267    /// [`ReconnectConfig::disabled()`] to disable (the default) or
268    /// [`ReconnectConfig::enabled()`] to opt in. See
269    /// [`Self::stream_reconnect`] for what it does and does not govern.
270    #[must_use]
271    pub fn with_stream_reconnect(mut self, stream_reconnect: ReconnectConfig) -> Self {
272        self.stream_reconnect = stream_reconnect;
273        self
274    }
275
276    /// Require TLS (reject plaintext `http://`) for this client. See
277    /// [`Self::require_tls`].
278    #[must_use]
279    pub fn with_require_tls(mut self, require_tls: bool) -> Self {
280        self.require_tls = require_tls;
281        self
282    }
283
284    /// Override the max idle keep-alive connections per upstream host. See
285    /// [`Self::pool_max_idle_per_host`].
286    #[must_use]
287    pub fn with_pool_max_idle_per_host(mut self, max: usize) -> Self {
288        self.pool_max_idle_per_host = max;
289        self
290    }
291
292    /// Override how long idle keep-alive connections are retained (`None` uses
293    /// hyper-util's default). See [`Self::pool_idle_timeout`].
294    #[must_use]
295    pub fn with_pool_idle_timeout(mut self, timeout: Option<Duration>) -> Self {
296        self.pool_idle_timeout = timeout;
297        self
298    }
299
300    /// Override the max concurrent in-flight requests (`None` disables the
301    /// limiter). See [`Self::max_concurrent_requests`].
302    #[must_use]
303    pub fn with_max_concurrent_requests(mut self, max: Option<usize>) -> Self {
304        self.max_concurrent_requests = max;
305        self
306    }
307
308    /// Set (or clear) the platform-plane internal-credential provider. See
309    /// [`Self::internal_token_provider`]. Accepts either an
310    /// [`InternalTokenProvider`] or an `Option<InternalTokenProvider>`, so the
311    /// bootstrap layer can pass through whatever the process selected without a
312    /// branch.
313    #[must_use]
314    pub fn with_internal_token_provider(
315        mut self,
316        provider: impl Into<Option<InternalTokenProvider>>,
317    ) -> Self {
318        self.internal_token_provider = provider.into();
319        self
320    }
321}
322
323/// Bounded exponential-backoff retry policy with full jitter.
324#[derive(Debug, Clone)]
325pub struct RetryConfig {
326    /// Maximum number of attempts (must be at least 1).
327    pub max_attempts: u32,
328    /// Base delay before the first retry.
329    pub base_delay: Duration,
330    /// Hard cap on the delay between retries.
331    pub max_delay: Duration,
332    /// Multiplier applied between consecutive retries.
333    pub multiplier: f64,
334}
335
336impl RetryConfig {
337    /// Disable retries entirely (single attempt).
338    #[must_use]
339    pub const fn off() -> Self {
340        Self {
341            max_attempts: 1,
342            base_delay: Duration::ZERO,
343            max_delay: Duration::ZERO,
344            multiplier: 1.0,
345        }
346    }
347}
348
349impl Default for RetryConfig {
350    fn default() -> Self {
351        Self {
352            max_attempts: 3,
353            base_delay: Duration::from_millis(100),
354            max_delay: Duration::from_secs(2),
355            multiplier: 2.0,
356        }
357    }
358}
359
360/// SSE reconnect policy. The streaming client tracks the latest `id:`
361/// field seen on the wire and, on transient stream failures, re-issues
362/// the request with a `Last-Event-ID: <stored>` header so the server can
363/// resume the event sequence (per HTML5 `EventSource` spec).
364///
365/// Default is **opt-in disabled** (`max_attempts: 0`) so existing SDKs see
366/// no behaviour change.
367#[derive(Debug, Clone)]
368pub struct ReconnectConfig {
369    /// Maximum number of *consecutive* reconnect attempts with no healthy
370    /// connection in between (the burst budget). `0` (default) disables
371    /// reconnect entirely — stream errors bubble up. The budget is reset by a
372    /// connection that both delivers an item and stays up at least
373    /// [`min_healthy_uptime`](Self::min_healthy_uptime).
374    pub max_attempts: u32,
375    /// Initial delay before the first reconnect attempt.
376    pub base_delay: Duration,
377    /// Hard cap on delay between reconnect attempts.
378    pub max_delay: Duration,
379    /// Minimum time a connection must stay up — *in addition to* delivering at
380    /// least one item — before its end resets the burst budget. Delivering a
381    /// single item is too weak a health signal on its own: a peer that emits
382    /// one item and immediately drops would reset the budget on every cycle and
383    /// reopen forever, re-sending the auth token each time (#4740). A connection
384    /// shorter than this counts against `max_attempts` like any other failed
385    /// attempt.
386    pub min_healthy_uptime: Duration,
387    /// Absolute lifetime ceiling on reopens, independent of budget resets. It
388    /// bounds the pathological peer that stays up *just past*
389    /// `min_healthy_uptime`, delivers an item, and drops on a loop — which would
390    /// otherwise reset the burst budget indefinitely. Set high enough that a
391    /// genuinely healthy long-lived subscription (which reconnects rarely) never
392    /// approaches it; `0` refuses reopen outright.
393    pub max_total_reopens: u32,
394}
395
396impl Default for ReconnectConfig {
397    fn default() -> Self {
398        Self {
399            max_attempts: 0,
400            base_delay: Duration::from_millis(500),
401            max_delay: Duration::from_secs(10),
402            min_healthy_uptime: Duration::from_secs(5),
403            max_total_reopens: 10_000,
404        }
405    }
406}
407
408impl ReconnectConfig {
409    /// Build a reconnect policy with up to `max_attempts` retries and the
410    /// supplied initial delay (capped by `max_delay`, default 10s).
411    #[must_use]
412    pub fn enabled(max_attempts: u32, base_delay: Duration) -> Self {
413        Self {
414            max_attempts,
415            base_delay,
416            ..Self::default()
417        }
418    }
419
420    /// A policy that never reconnects: the stream's first transport failure
421    /// ends it.
422    ///
423    /// The counterpart to [`ReconnectConfig::enabled`]. [`Default`] already
424    /// yields `max_attempts: 0`, so this is behaviourally the same value — it
425    /// exists so a call site that *must* not reconnect reads as a deliberate
426    /// choice rather than an accepted default.
427    ///
428    /// Generated clients pass this for a method whose open is fallible
429    /// (`#[streaming] async fn`). Such an open carries domain semantics the
430    /// client must not blindly repeat: any exclusion lease it acquired is owned
431    /// by the returned stream's lifetime, resume may be specified through the
432    /// contract's own cursor rather than `Last-Event-ID`, and a reconnect-time
433    /// failure would arrive as a stream *item* — past the caller's open-time
434    /// error handling, which is the whole reason the fallible shape exists.
435    #[must_use]
436    pub const fn disabled() -> Self {
437        // Every field is named explicitly (rather than `..Self::default()`) so a
438        // future field with an *enabling* default can't silently leak into a
439        // constructor documented as never reconnecting — matching
440        // [`RetryConfig::off`]. These are the same values as [`Default`].
441        Self {
442            max_attempts: 0,
443            base_delay: Duration::from_millis(500),
444            max_delay: Duration::from_secs(10),
445            min_healthy_uptime: Duration::from_secs(5),
446            max_total_reopens: 10_000,
447        }
448    }
449
450    /// Override the maximum delay between reconnect attempts.
451    #[must_use]
452    pub fn with_max_delay(mut self, max_delay: Duration) -> Self {
453        self.max_delay = max_delay;
454        self
455    }
456
457    /// Override the minimum healthy connection uptime that resets the burst
458    /// budget. See [`min_healthy_uptime`](Self::min_healthy_uptime).
459    #[must_use]
460    pub fn with_min_healthy_uptime(mut self, min_healthy_uptime: Duration) -> Self {
461        self.min_healthy_uptime = min_healthy_uptime;
462        self
463    }
464
465    /// Override the absolute lifetime cap on reopens. See
466    /// [`max_total_reopens`](Self::max_total_reopens).
467    #[must_use]
468    pub fn with_max_total_reopens(mut self, max_total_reopens: u32) -> Self {
469        self.max_total_reopens = max_total_reopens;
470        self
471    }
472}
473
474#[cfg(test)]
475#[cfg_attr(coverage_nightly, coverage(off))]
476mod tests {
477    use super::*;
478
479    #[test]
480    fn default_retry_has_three_attempts() {
481        let r = RetryConfig::default();
482        assert_eq!(r.max_attempts, 3);
483        assert!(r.base_delay > Duration::ZERO);
484    }
485
486    #[test]
487    fn off_yields_single_attempt() {
488        let r = RetryConfig::off();
489        assert_eq!(r.max_attempts, 1);
490    }
491
492    // The "never a second attempt" guarantee that `disabled()` carries is
493    // pinned behaviourally by `reconnect_is_derived_from_the_open_shape_not_from_client_config`
494    // (tests/rest_client_codegen.rs), which counts real server connections on
495    // the fallible-open path and asserts exactly one. A unit test that merely
496    // read back `disabled()`'s fields couldn't fail unless struct construction
497    // itself broke, so it isn't restated here.
498
499    #[test]
500    fn client_config_chains_overrides() {
501        let cfg = ClientConfig::new("https://x.example")
502            .with_timeout(Duration::from_secs(5))
503            .with_retry(RetryConfig::off());
504        assert_eq!(cfg.base_url, "https://x.example");
505        assert_eq!(cfg.timeout, Duration::from_secs(5));
506        assert_eq!(cfg.retry.max_attempts, 1);
507    }
508}