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}