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}