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#[derive(Debug, Clone)]
114pub struct ClientConfig {
115    /// Base URL prefix (e.g., `https://billing.internal`).
116    /// Combined with the base path declared in the projection trait.
117    pub base_url: String,
118    /// Deadline applied to a **single** unary attempt — NOT to the whole logical
119    /// call. A `#[retryable]` method may make up to `retry.max_attempts` attempts,
120    /// so the worst-case wall-clock for a logical call is bounded by
121    /// `max_attempts × (timeout + retry.max_delay)` (the per-retry backoff is
122    /// itself clamped to [`RetryConfig::max_delay`], including a server-advised
123    /// `Retry-After`). There is deliberately no separate whole-call budget field.
124    pub timeout: Duration,
125    /// Per-**item** idle deadline for streams of any framing: the maximum gap
126    /// between two received wire chunks before the stream is treated as timed
127    /// out. A long-lived stream is NOT bounded by [`timeout`](Self::timeout)
128    /// (which would kill a healthy slow stream); it is bounded by this larger
129    /// idle deadline instead. Defaults to 60s (> the unary default).
130    ///
131    /// This is *idle*, not *quiet*: any wire chunk resets it, including ones
132    /// that dispatch no item (an SSE keepalive comment, a multipart part header
133    /// block arriving on its own).
134    pub stream_idle_timeout: Duration,
135    /// Retry policy applied to methods marked `#[retryable]`.
136    pub retry: RetryConfig,
137    /// Reconnect policy for streams of any framing. By default
138    /// `max_attempts: 0` — stream failures bubble up unchanged. Set explicitly
139    /// to opt into transparent re-open on a transient failure.
140    ///
141    /// Two limits are deliberate rather than accidental:
142    ///
143    /// - **It applies only to a method whose open is immediate**
144    ///   (`#[streaming] fn`). A fallible open (`#[streaming] async fn`) carries
145    ///   domain semantics the client must not blindly repeat — a re-open can
146    ///   collide with an exclusion lease the first open acquired, and its
147    ///   failure would land as a stream item, past the caller's open-time error
148    ///   handling. Generated code therefore passes
149    ///   [`ReconnectConfig::disabled`] for a fallible open regardless of this
150    ///   value.
151    /// - **Resume via `Last-Event-ID` is SSE-only.** A reconnected
152    ///   `multipart/mixed` stream re-issues the original request with no resume
153    ///   token, because the framing has none. The transport still reopens it,
154    ///   but that is a blind restart, not a resume: the server replays the body
155    ///   from its first part, so any items already delivered before the failure
156    ///   are **delivered again** (at-least-once, with no marker for the
157    ///   restart). Enable reconnect for a multipart stream only where the
158    ///   consumer tolerates duplicates; one needing exactly-once must instead
159    ///   leave reconnect disabled and run its own reopen loop with
160    ///   application-level dedup.
161    pub stream_reconnect: ReconnectConfig,
162    /// When `true`, the generated client refuses plaintext `http://` and
163    /// requires TLS (`toolkit_http::TransportSecurity::TlsOnly`) for every
164    /// request — including the bearer-carrying `Authorization` header, which
165    /// otherwise would ride whatever scheme `base_url` uses. Defaults to
166    /// `false`, preserving the platform's existing in-mesh
167    /// service-to-service convention where plaintext HTTP inside a secured
168    /// network boundary is an accepted, deliberate choice (see
169    /// [`build_default_http_client`](crate::runtime::client::build_default_http_client)).
170    /// Set this when a resolved endpoint may cross an untrusted network.
171    pub require_tls: bool,
172    /// Source of the platform-plane internal credential attached to methods
173    /// whose plane marker is `PlatformSecurityContext` (carried as
174    /// `X-ToolKit-Internal-Token`). `None` (the default) attaches nothing —
175    /// legitimate for Profile 1 / in-process (`InternalCredential::None`);
176    /// the requirement is enforced server-side. The bootstrap layer populates
177    /// this from the process's selected `InternalCredential`. Tenant-plane
178    /// methods (`SecurityContext`) never consult it; they forward the caller's
179    /// bearer token from the argument.
180    pub internal_token_provider: Option<InternalTokenProvider>,
181}
182
183impl ClientConfig {
184    /// Create a new config with sensible defaults.
185    #[must_use]
186    pub fn new(base_url: impl Into<String>) -> Self {
187        Self {
188            base_url: base_url.into(),
189            timeout: Duration::from_secs(30),
190            stream_idle_timeout: Duration::from_mins(1),
191            retry: RetryConfig::default(),
192            stream_reconnect: ReconnectConfig::default(),
193            require_tls: false,
194            internal_token_provider: None,
195        }
196    }
197
198    /// Override the per-call (unary) timeout.
199    #[must_use]
200    pub fn with_timeout(mut self, timeout: Duration) -> Self {
201        self.timeout = timeout;
202        self
203    }
204
205    /// Override the per-item stream idle deadline (max gap between wire
206    /// chunks). See [`Self::stream_idle_timeout`].
207    #[must_use]
208    pub fn with_stream_idle_timeout(mut self, idle: Duration) -> Self {
209        self.stream_idle_timeout = idle;
210        self
211    }
212
213    /// Override the retry policy.
214    #[must_use]
215    pub fn with_retry(mut self, retry: RetryConfig) -> Self {
216        self.retry = retry;
217        self
218    }
219
220    /// Override the stream reconnect policy. Use
221    /// [`ReconnectConfig::disabled()`] to disable (the default) or
222    /// [`ReconnectConfig::enabled()`] to opt in. See
223    /// [`Self::stream_reconnect`] for what it does and does not govern.
224    #[must_use]
225    pub fn with_stream_reconnect(mut self, stream_reconnect: ReconnectConfig) -> Self {
226        self.stream_reconnect = stream_reconnect;
227        self
228    }
229
230    /// Require TLS (reject plaintext `http://`) for this client. See
231    /// [`Self::require_tls`].
232    #[must_use]
233    pub fn with_require_tls(mut self, require_tls: bool) -> Self {
234        self.require_tls = require_tls;
235        self
236    }
237
238    /// Set (or clear) the platform-plane internal-credential provider. See
239    /// [`Self::internal_token_provider`]. Accepts either an
240    /// [`InternalTokenProvider`] or an `Option<InternalTokenProvider>`, so the
241    /// bootstrap layer can pass through whatever the process selected without a
242    /// branch.
243    #[must_use]
244    pub fn with_internal_token_provider(
245        mut self,
246        provider: impl Into<Option<InternalTokenProvider>>,
247    ) -> Self {
248        self.internal_token_provider = provider.into();
249        self
250    }
251}
252
253/// Bounded exponential-backoff retry policy with full jitter.
254#[derive(Debug, Clone)]
255pub struct RetryConfig {
256    /// Maximum number of attempts (must be at least 1).
257    pub max_attempts: u32,
258    /// Base delay before the first retry.
259    pub base_delay: Duration,
260    /// Hard cap on the delay between retries.
261    pub max_delay: Duration,
262    /// Multiplier applied between consecutive retries.
263    pub multiplier: f64,
264}
265
266impl RetryConfig {
267    /// Disable retries entirely (single attempt).
268    #[must_use]
269    pub const fn off() -> Self {
270        Self {
271            max_attempts: 1,
272            base_delay: Duration::ZERO,
273            max_delay: Duration::ZERO,
274            multiplier: 1.0,
275        }
276    }
277}
278
279impl Default for RetryConfig {
280    fn default() -> Self {
281        Self {
282            max_attempts: 3,
283            base_delay: Duration::from_millis(100),
284            max_delay: Duration::from_secs(2),
285            multiplier: 2.0,
286        }
287    }
288}
289
290/// SSE reconnect policy. The streaming client tracks the latest `id:`
291/// field seen on the wire and, on transient stream failures, re-issues
292/// the request with a `Last-Event-ID: <stored>` header so the server can
293/// resume the event sequence (per HTML5 `EventSource` spec).
294///
295/// Default is **opt-in disabled** (`max_attempts: 0`) so existing SDKs see
296/// no behaviour change.
297#[derive(Debug, Clone)]
298pub struct ReconnectConfig {
299    /// Maximum number of *consecutive* reconnect attempts with no healthy
300    /// connection in between (the burst budget). `0` (default) disables
301    /// reconnect entirely — stream errors bubble up. The budget is reset by a
302    /// connection that both delivers an item and stays up at least
303    /// [`min_healthy_uptime`](Self::min_healthy_uptime).
304    pub max_attempts: u32,
305    /// Initial delay before the first reconnect attempt.
306    pub base_delay: Duration,
307    /// Hard cap on delay between reconnect attempts.
308    pub max_delay: Duration,
309    /// Minimum time a connection must stay up — *in addition to* delivering at
310    /// least one item — before its end resets the burst budget. Delivering a
311    /// single item is too weak a health signal on its own: a peer that emits
312    /// one item and immediately drops would reset the budget on every cycle and
313    /// reopen forever, re-sending the auth token each time (#4740). A connection
314    /// shorter than this counts against `max_attempts` like any other failed
315    /// attempt.
316    pub min_healthy_uptime: Duration,
317    /// Absolute lifetime ceiling on reopens, independent of budget resets. It
318    /// bounds the pathological peer that stays up *just past*
319    /// `min_healthy_uptime`, delivers an item, and drops on a loop — which would
320    /// otherwise reset the burst budget indefinitely. Set high enough that a
321    /// genuinely healthy long-lived subscription (which reconnects rarely) never
322    /// approaches it; `0` refuses reopen outright.
323    pub max_total_reopens: u32,
324}
325
326impl Default for ReconnectConfig {
327    fn default() -> Self {
328        Self {
329            max_attempts: 0,
330            base_delay: Duration::from_millis(500),
331            max_delay: Duration::from_secs(10),
332            min_healthy_uptime: Duration::from_secs(5),
333            max_total_reopens: 10_000,
334        }
335    }
336}
337
338impl ReconnectConfig {
339    /// Build a reconnect policy with up to `max_attempts` retries and the
340    /// supplied initial delay (capped by `max_delay`, default 10s).
341    #[must_use]
342    pub fn enabled(max_attempts: u32, base_delay: Duration) -> Self {
343        Self {
344            max_attempts,
345            base_delay,
346            ..Self::default()
347        }
348    }
349
350    /// A policy that never reconnects: the stream's first transport failure
351    /// ends it.
352    ///
353    /// The counterpart to [`ReconnectConfig::enabled`]. [`Default`] already
354    /// yields `max_attempts: 0`, so this is behaviourally the same value — it
355    /// exists so a call site that *must* not reconnect reads as a deliberate
356    /// choice rather than an accepted default.
357    ///
358    /// Generated clients pass this for a method whose open is fallible
359    /// (`#[streaming] async fn`). Such an open carries domain semantics the
360    /// client must not blindly repeat: any exclusion lease it acquired is owned
361    /// by the returned stream's lifetime, resume may be specified through the
362    /// contract's own cursor rather than `Last-Event-ID`, and a reconnect-time
363    /// failure would arrive as a stream *item* — past the caller's open-time
364    /// error handling, which is the whole reason the fallible shape exists.
365    #[must_use]
366    pub const fn disabled() -> Self {
367        // Every field is named explicitly (rather than `..Self::default()`) so a
368        // future field with an *enabling* default can't silently leak into a
369        // constructor documented as never reconnecting — matching
370        // [`RetryConfig::off`]. These are the same values as [`Default`].
371        Self {
372            max_attempts: 0,
373            base_delay: Duration::from_millis(500),
374            max_delay: Duration::from_secs(10),
375            min_healthy_uptime: Duration::from_secs(5),
376            max_total_reopens: 10_000,
377        }
378    }
379
380    /// Override the maximum delay between reconnect attempts.
381    #[must_use]
382    pub fn with_max_delay(mut self, max_delay: Duration) -> Self {
383        self.max_delay = max_delay;
384        self
385    }
386
387    /// Override the minimum healthy connection uptime that resets the burst
388    /// budget. See [`min_healthy_uptime`](Self::min_healthy_uptime).
389    #[must_use]
390    pub fn with_min_healthy_uptime(mut self, min_healthy_uptime: Duration) -> Self {
391        self.min_healthy_uptime = min_healthy_uptime;
392        self
393    }
394
395    /// Override the absolute lifetime cap on reopens. See
396    /// [`max_total_reopens`](Self::max_total_reopens).
397    #[must_use]
398    pub fn with_max_total_reopens(mut self, max_total_reopens: u32) -> Self {
399        self.max_total_reopens = max_total_reopens;
400        self
401    }
402}
403
404#[cfg(test)]
405#[cfg_attr(coverage_nightly, coverage(off))]
406mod tests {
407    use super::*;
408
409    #[test]
410    fn default_retry_has_three_attempts() {
411        let r = RetryConfig::default();
412        assert_eq!(r.max_attempts, 3);
413        assert!(r.base_delay > Duration::ZERO);
414    }
415
416    #[test]
417    fn off_yields_single_attempt() {
418        let r = RetryConfig::off();
419        assert_eq!(r.max_attempts, 1);
420    }
421
422    // The "never a second attempt" guarantee that `disabled()` carries is
423    // pinned behaviourally by `reconnect_is_derived_from_the_open_shape_not_from_client_config`
424    // (tests/rest_client_codegen.rs), which counts real server connections on
425    // the fallible-open path and asserts exactly one. A unit test that merely
426    // read back `disabled()`'s fields couldn't fail unless struct construction
427    // itself broke, so it isn't restated here.
428
429    #[test]
430    fn client_config_chains_overrides() {
431        let cfg = ClientConfig::new("https://x.example")
432            .with_timeout(Duration::from_secs(5))
433            .with_retry(RetryConfig::off());
434        assert_eq!(cfg.base_url, "https://x.example");
435        assert_eq!(cfg.timeout, Duration::from_secs(5));
436        assert_eq!(cfg.retry.max_attempts, 1);
437    }
438}