Skip to main content

nmbrs_metrics/
controls.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Dynamic controls — see SRD 23.
5//!
6//! A [`Control<T>`] is a named, typed, observable cell that
7//! coordinates a **confirmed-apply** write protocol across a
8//! set of registered [`ControlApplier<T>`] consumers. A caller
9//! that [`Control::set`]s a new value does not proceed until
10//! every applier has acknowledged the change is in effect; if
11//! any applier fails or times out, the committed value is not
12//! advanced and the writer sees the aggregated error.
13//!
14//! Controls are structural runtime properties of a component
15//! (SRD 24): they are declared when the owning component /
16//! wrapper / dispenser is instantiated, enumerable through the
17//! component tree, and looked up — never implicitly created —
18//! by writers.
19//!
20//! Follow-up integrations:
21//! - **Gauge reification.** Each control's current value is
22//!   also published as a numeric gauge under the metric name
23//!   `control.<name>`, visible through every metric sink that
24//!   reads the component tree (SRD 23 §"Reification as
25//!   metrics"). Opt in via [`ControlBuilder::reify_as_gauge`].
26//! - **Adapter-side integrations** (rate limiter, fiber
27//!   executor): the control primitive works end-to-end with
28//!   its tests; wiring it into `nmbrs-rate` and the fiber pool
29//!   lands in a follow-up pass.
30
31use std::collections::HashMap;
32use std::fmt;
33use std::future::Future;
34use std::pin::Pin;
35use std::sync::atomic::{AtomicU64, Ordering};
36use std::sync::{Arc, Mutex, OnceLock};
37use std::time::{Duration, Instant};
38
39use futures::future::join_all;
40use tokio::sync::Mutex as AsyncMutex;
41
42use crate::instruments::gauge::ValueGauge;
43use crate::labels::Labels;
44use crate::snapshot::MetricSet;
45
46// =========================================================================
47// Rev allocator (session-wide monotonic revision counter)
48// =========================================================================
49
50/// A session-wide revision counter. Every successful control
51/// write allocates a fresh `rev` via [`Self::next`].
52///
53/// Callers who build a session share one `RevAllocator` across
54/// every `Control` in that session; the counter is strictly
55/// increasing and usable for log ordering or reverse lookups.
56///
57/// When no explicit allocator is wired (e.g. tests, single-
58/// purpose uses), a process-global allocator is available via
59/// [`RevAllocator::global`].
60#[derive(Debug, Default)]
61pub struct RevAllocator {
62    counter: AtomicU64,
63}
64
65impl RevAllocator {
66    pub fn new() -> Self {
67        Self::default()
68    }
69
70    pub fn next(&self) -> u64 {
71        self.counter.fetch_add(1, Ordering::AcqRel).wrapping_add(1)
72    }
73
74    /// Shared allocator used when no session-scoped one is
75    /// explicitly supplied. Useful for tests and for short-lived
76    /// single-session processes — production runs should give
77    /// each session its own allocator.
78    pub fn global() -> &'static Arc<RevAllocator> {
79        static GLOBAL: OnceLock<Arc<RevAllocator>> = OnceLock::new();
80        GLOBAL.get_or_init(|| Arc::new(RevAllocator::new()))
81    }
82}
83
84// =========================================================================
85// Origin, Versioned<T>, error types
86// =========================================================================
87
88/// Attribution for a control write. Every committed `Versioned`
89/// remembers who caused the change so log / replay / summary
90/// tools can explain "where did this come from?" after the fact.
91#[derive(Clone, Debug, PartialEq, Eq)]
92pub enum ControlOrigin {
93    /// Initial seed from `params:` at scenario start.
94    Launch,
95    /// Future `nmbrs ctl` CLI mutation.
96    Cli,
97    /// Keybind / input in the TUI.
98    Tui,
99    /// Scripted feedback loop (GK `control_set(...)` node).
100    Polydat { binding: String },
101    /// Runtime governor (e.g. the phase `throttle:` backpressure
102    /// loop) — `source` names the governing phase.
103    Governor { source: String },
104    /// External API — `source` identifies the caller (endpoint, auth id, etc.)
105    Api { source: String },
106    /// Test harness — never seen in production.
107    Test,
108}
109
110/// A committed control value bundled with its metadata.
111#[derive(Clone, Debug)]
112pub struct Versioned<T: Clone> {
113    pub value: T,
114    /// Session-wide monotonic revision. Strictly increasing
115    /// across every successful write to every control sharing
116    /// a [`RevAllocator`].
117    pub rev: u64,
118    /// Wall-clock time at which the commit happened.
119    pub updated_at: Instant,
120    pub origin: ControlOrigin,
121}
122
123/// Error outcomes returned by [`Control::set`]. A write can
124/// reject at validation (before fan-out), at apply (during
125/// fan-out), or up front because the control is `final` at its
126/// declaring scope. In every case the committed value is
127/// unchanged.
128#[derive(Clone, Debug, PartialEq, Eq)]
129pub enum SetError {
130    /// The pre-fanout validator rejected the value. No applier
131    /// was called. Message explains why.
132    ValidationFailed(String),
133    /// One or more appliers returned `Err` (or timed out). The
134    /// list carries the index + error message for each failure.
135    ApplyFailed(Vec<ApplyFailure>),
136    /// The control was declared `final` at `scope` and is
137    /// pinned to its Launch value. Runtime writes (any origin
138    /// other than `Launch`) are rejected.
139    FinalViolation { scope: String },
140}
141
142/// Declared visibility scope of a [`Control`]. Governs how
143/// descendant components resolve the control and how `set`
144/// writes propagate to newly-instantiated components below the
145/// declaring node. See SRD 23 §"Branch-scoped and final
146/// controls".
147#[derive(Clone, Copy, Debug, PartialEq, Eq)]
148pub enum BranchScope {
149    /// Control applies only to the component that declares it.
150    /// Descendants resolve the same name only if they declare
151    /// it themselves (or if an ancestor declares it with
152    /// [`BranchScope::Subtree`]).
153    Local,
154    /// Control applies to the declaring component and every
155    /// descendant in the subtree. Descendants resolving the
156    /// control name walk up and find this declaration;
157    /// descendant components constructed after a successful
158    /// write read the new committed value at construction.
159    Subtree,
160}
161
162#[derive(Clone, Debug, PartialEq, Eq)]
163pub struct ApplyFailure {
164    /// Registration index of the failing applier (0-based, in
165    /// insertion order). The applier itself is opaque to the
166    /// writer — the index plus the component the control lives
167    /// on is enough to locate the offending subscriber.
168    pub applier_index: usize,
169    /// Why the apply failed — either the `Err` the applier
170    /// returned or a timeout notice.
171    pub message: String,
172}
173
174impl fmt::Display for SetError {
175    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
176        match self {
177            Self::ValidationFailed(msg) => write!(f, "validation failed: {msg}"),
178            Self::ApplyFailed(failures) => {
179                write!(f, "{} applier(s) failed:", failures.len())?;
180                for fail in failures {
181                    write!(f, " [#{}: {}]", fail.applier_index, fail.message)?;
182                }
183                Ok(())
184            }
185            Self::FinalViolation { scope } => {
186                write!(
187                    f,
188                    "control is declared final at scope '{scope}'; runtime writes are rejected"
189                )
190            }
191        }
192    }
193}
194
195impl std::error::Error for SetError {}
196
197// =========================================================================
198// ControlApplier trait
199// =========================================================================
200
201/// Async callback registered on a [`Control`]. Each call attempts
202/// to apply a new value and returns `Ok(())` once the new value
203/// is in effect for this subscriber, or `Err(msg)` if the apply
204/// could not be completed.
205///
206/// Implementors should keep their own state consistent with
207/// either the old or the new value; partial / corrupted apply
208/// states are outside the protocol.
209pub trait ControlApplier<T>: Send + Sync + 'static
210where
211    T: Clone + Send + Sync + 'static,
212{
213    fn apply(&self, value: T) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send + '_>>;
214}
215
216/// Convenience wrapper for synchronous appliers. Most appliers
217/// just set an atomic or reconfigure a rate limiter and have no
218/// genuine async work to do — this wrapper lifts a plain
219/// `Fn(T) -> Result<(), String>` into a [`ControlApplier<T>`].
220pub struct SyncApplier<T, F>
221where
222    T: Clone + Send + Sync + 'static,
223    F: Fn(T) -> Result<(), String> + Send + Sync + 'static,
224{
225    f: F,
226    _marker: std::marker::PhantomData<fn(T)>,
227}
228
229impl<T, F> SyncApplier<T, F>
230where
231    T: Clone + Send + Sync + 'static,
232    F: Fn(T) -> Result<(), String> + Send + Sync + 'static,
233{
234    pub fn new(f: F) -> Self {
235        Self {
236            f,
237            _marker: std::marker::PhantomData,
238        }
239    }
240}
241
242impl<T, F> ControlApplier<T> for SyncApplier<T, F>
243where
244    T: Clone + Send + Sync + 'static,
245    F: Fn(T) -> Result<(), String> + Send + Sync + 'static,
246{
247    fn apply(&self, value: T) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send + '_>> {
248        let out = (self.f)(value);
249        Box::pin(async move { out })
250    }
251}
252
253// =========================================================================
254// Control<T>
255// =========================================================================
256
257type Validator<T> = Box<dyn Fn(&T) -> Result<(), String> + Send + Sync>;
258type ToF64<T> = Box<dyn Fn(&T) -> Option<f64> + Send + Sync>;
259type FromF64<T> = Box<dyn Fn(f64) -> Result<T, String> + Send + Sync>;
260
261pub struct Control<T: Clone + Send + Sync + 'static> {
262    inner: Arc<ControlInner<T>>,
263}
264
265impl<T: Clone + Send + Sync + 'static> Clone for Control<T> {
266    fn clone(&self) -> Self {
267        Self {
268            inner: self.inner.clone(),
269        }
270    }
271}
272
273struct ControlInner<T: Clone + Send + Sync + 'static> {
274    name: String,
275    /// Serializes writers across validate → fan-out → ack →
276    /// commit. Held by `set()` for its entire duration so two
277    /// writers can't overlap each other's fan-outs.
278    write_lock: AsyncMutex<()>,
279    /// Current committed value. Readers take a snapshot under
280    /// this RwLock; writers update it only on a fully-successful
281    /// apply.
282    committed: std::sync::RwLock<Versioned<T>>,
283    appliers: Mutex<Vec<Arc<dyn ControlApplier<T>>>>,
284    validator: OnceLock<Validator<T>>,
285    apply_timeout: Duration,
286    rev_allocator: Arc<RevAllocator>,
287    /// Optional gauge reification — when declared via
288    /// [`ControlBuilder::reify_as_gauge`], the control publishes
289    /// its current value as a numeric gauge. The gauge is
290    /// `Arc`-shared so readers that want to sample without going
291    /// through the control machinery can do so directly.
292    gauge: Option<GaugeReification<T>>,
293    /// When `Some(scope)`, the control is pinned: only the
294    /// Launch-origin seeding writes succeed; every subsequent
295    /// runtime write returns [`SetError::FinalViolation`] with
296    /// this scope name so the operator knows where the pin is
297    /// declared.
298    final_at_scope: Option<String>,
299    /// Declared visibility of this control within the component
300    /// tree.
301    branch_scope: BranchScope,
302    /// Optional converter from `f64` to `T`. Enables type-erased
303    /// writes from Polydat `control_set(name, value)` nodes and the
304    /// web API's JSON-number body. Controls that don't declare a
305    /// converter reject `f64` writes with `ValidationFailed`.
306    from_f64: Option<FromF64<T>>,
307}
308
309struct GaugeReification<T> {
310    gauge: Arc<ValueGauge>,
311    to_f64: ToF64<T>,
312}
313
314impl<T: Clone + Send + Sync + 'static> Control<T> {
315    /// Name of this control within its owning component's
316    /// registry.
317    pub fn name(&self) -> &str {
318        &self.inner.name
319    }
320
321    /// Borrow the current committed value.
322    pub fn get(&self) -> Versioned<T> {
323        self.inner
324            .committed
325            .read()
326            .unwrap_or_else(|e| e.into_inner())
327            .clone()
328    }
329
330    /// Convenience shortcut for `self.get().value` when the
331    /// caller doesn't need the metadata.
332    pub fn value(&self) -> T {
333        self.get().value
334    }
335
336    /// Register an applier. Every subsequent [`Self::set`] call
337    /// will include this applier in its fanout. There is no
338    /// unregister — appliers live for the control's lifetime.
339    /// Returns the insertion index so callers can correlate
340    /// later failure reports.
341    pub fn register_applier<A>(&self, applier: A) -> usize
342    where
343        A: ControlApplier<T>,
344    {
345        let arc: Arc<dyn ControlApplier<T>> = Arc::new(applier);
346        let mut g = self
347            .inner
348            .appliers
349            .lock()
350            .unwrap_or_else(|e| e.into_inner());
351        g.push(arc);
352        g.len() - 1
353    }
354
355    /// Number of currently-registered appliers.
356    pub fn applier_count(&self) -> usize {
357        self.inner
358            .appliers
359            .lock()
360            .unwrap_or_else(|e| e.into_inner())
361            .len()
362    }
363
364    /// Atomic write. Returns `Ok(rev)` with the new revision
365    /// only if every registered applier acknowledged the new
366    /// value within the control's apply timeout. On any failure
367    /// the committed value is not advanced.
368    ///
369    /// A control declared `final` via
370    /// [`ControlBuilder::final_at_scope`] rejects every write
371    /// whose origin is not [`ControlOrigin::Launch`] with
372    /// [`SetError::FinalViolation`]. The Launch-seed write
373    /// still goes through so the initial value can be
374    /// installed; runtime writers (CLI, TUI, GK, API) are
375    /// rejected against the pin.
376    pub async fn set(&self, value: T, origin: ControlOrigin) -> Result<u64, SetError> {
377        let _guard = self.inner.write_lock.lock().await;
378
379        // 0. Final-declaration pin. Checked before validation
380        //    so operators see the pin message rather than a
381        //    validation error that happens to also reject the
382        //    value.
383        if let Some(ref scope) = self.inner.final_at_scope
384            && origin != ControlOrigin::Launch
385        {
386            return Err(SetError::FinalViolation {
387                scope: scope.clone(),
388            });
389        }
390
391        // 1. Validation runs before any applier is called.
392        if let Some(validator) = self.inner.validator.get() {
393            validator(&value).map_err(SetError::ValidationFailed)?;
394        }
395
396        // 2. Fan out to all appliers concurrently, each with the
397        //    per-control timeout. A missing apply (timeout) is
398        //    treated as a failure just like an explicit Err.
399        let appliers: Vec<Arc<dyn ControlApplier<T>>> = {
400            let g = self
401                .inner
402                .appliers
403                .lock()
404                .unwrap_or_else(|e| e.into_inner());
405            g.clone()
406        };
407
408        let timeout = self.inner.apply_timeout;
409        let futures = appliers.iter().enumerate().map(|(idx, applier)| {
410            let v = value.clone();
411            let applier = applier.clone();
412            async move {
413                let fut = applier.apply(v);
414                match tokio::time::timeout(timeout, fut).await {
415                    Ok(Ok(())) => Ok(idx),
416                    Ok(Err(msg)) => Err(ApplyFailure {
417                        applier_index: idx,
418                        message: msg,
419                    }),
420                    Err(_) => Err(ApplyFailure {
421                        applier_index: idx,
422                        message: format!("apply timed out after {:?}", timeout),
423                    }),
424                }
425            }
426        });
427
428        let results = join_all(futures).await;
429        let failures: Vec<ApplyFailure> = results.into_iter().filter_map(|r| r.err()).collect();
430
431        if !failures.is_empty() {
432            return Err(SetError::ApplyFailed(failures));
433        }
434
435        // 3. All appliers succeeded — allocate a rev and commit.
436        let rev = self.inner.rev_allocator.next();
437        let versioned = Versioned {
438            value: value.clone(),
439            rev,
440            updated_at: Instant::now(),
441            origin,
442        };
443        *self
444            .inner
445            .committed
446            .write()
447            .unwrap_or_else(|e| e.into_inner()) = versioned;
448        // 4. Publish the new value to the reified gauge (if any).
449        //    This lands after commit so a reader that samples the
450        //    gauge concurrently always sees a value at least as
451        //    new as the committed `Versioned` — never newer.
452        self.publish_gauge(&value);
453
454        Ok(rev)
455    }
456
457    /// Current reified-gauge value, if the control opted in via
458    /// [`ControlBuilder::reify_as_gauge`] and the conversion
459    /// returns `Some`. The registry snapshots this into a
460    /// `MetricSet` so control values flow through the normal
461    /// metric sinks.
462    pub fn gauge_f64(&self) -> Option<f64> {
463        let g = self.inner.gauge.as_ref()?;
464        let v = self.get().value;
465        (g.to_f64)(&v)
466    }
467
468    /// Raw handle to the reified [`ValueGauge`] for callers that
469    /// want to sample the gauge directly without going through
470    /// the control machinery (e.g. registering it on a
471    /// [`crate::component::Component`] via
472    /// `register_instrument`). `None` if the control wasn't
473    /// declared with gauge reification.
474    pub fn reified_gauge(&self) -> Option<Arc<ValueGauge>> {
475        self.inner.gauge.as_ref().map(|g| g.gauge.clone())
476    }
477
478    /// `true` if the control was declared `final` via
479    /// [`ControlBuilder::final_at_scope`]. Final controls
480    /// reject every non-Launch write.
481    pub fn is_final(&self) -> bool {
482        self.inner.final_at_scope.is_some()
483    }
484
485    /// Name of the scope at which the control was declared
486    /// `final`, if any. `None` for non-final controls.
487    pub fn final_scope(&self) -> Option<&str> {
488        self.inner.final_at_scope.as_deref()
489    }
490
491    /// Declared branch scope of this control.
492    pub fn branch_scope(&self) -> BranchScope {
493        self.inner.branch_scope
494    }
495
496    fn publish_gauge(&self, value: &T) {
497        if let Some(ref g) = self.inner.gauge
498            && let Some(f) = (g.to_f64)(value)
499        {
500            g.gauge.set(f);
501        }
502    }
503}
504
505// =========================================================================
506// ControlBuilder
507// =========================================================================
508
509/// Construct a [`Control<T>`]. Every control needs an initial
510/// value, a name, and a rev allocator. Apply timeout and
511/// validator are optional.
512pub struct ControlBuilder<T: Clone + Send + Sync + 'static> {
513    name: String,
514    initial: T,
515    rev_allocator: Arc<RevAllocator>,
516    apply_timeout: Duration,
517    validator: Option<Validator<T>>,
518    reify: Option<ToF64<T>>,
519    final_at_scope: Option<String>,
520    branch_scope: BranchScope,
521    from_f64: Option<FromF64<T>>,
522}
523
524impl<T: Clone + Send + Sync + 'static> ControlBuilder<T> {
525    pub fn new(name: &str, initial: T) -> Self {
526        Self {
527            name: name.to_string(),
528            initial,
529            rev_allocator: RevAllocator::global().clone(),
530            // Five seconds is long enough for a sluggish reconfigure
531            // (e.g. restarting a rate limiter) but short enough that
532            // a stuck subscriber surfaces as an error well within a
533            // human operator's attention span.
534            apply_timeout: Duration::from_secs(5),
535            validator: None,
536            reify: None,
537            final_at_scope: None,
538            branch_scope: BranchScope::Local,
539            from_f64: None,
540        }
541    }
542
543    pub fn rev_allocator(mut self, allocator: Arc<RevAllocator>) -> Self {
544        self.rev_allocator = allocator;
545        self
546    }
547
548    pub fn apply_timeout(mut self, d: Duration) -> Self {
549        self.apply_timeout = d;
550        self
551    }
552
553    /// Register a pre-fanout validator. Bad values reject the
554    /// write before any applier is called. Typical use: bounds
555    /// checking (`0 < concurrency <= 10_000`), format parsing.
556    pub fn validator<F>(mut self, f: F) -> Self
557    where
558        F: Fn(&T) -> Result<(), String> + Send + Sync + 'static,
559    {
560        self.validator = Some(Box::new(f));
561        self
562    }
563
564    /// Publish the control's current value as a numeric gauge.
565    /// `to_f64` converts each committed value to an `f64` (or
566    /// `None` to suppress that sample — useful for enum-valued
567    /// controls where some variants don't have a meaningful
568    /// numeric projection). The reified gauge is captured
569    /// alongside instrument-emitted metrics by the tree walk —
570    /// consumers (SQLite, VictoriaMetrics push, summary report,
571    /// TUI) pick it up without any extra wiring.
572    ///
573    /// Metric name: `control.<control_name>`. Labels: the
574    /// effective-labels of the component the control lives on.
575    pub fn reify_as_gauge<F>(mut self, to_f64: F) -> Self
576    where
577        F: Fn(&T) -> Option<f64> + Send + Sync + 'static,
578    {
579        self.reify = Some(Box::new(to_f64));
580        self
581    }
582
583    /// Pin the control's value at the declaring scope. The
584    /// Launch-origin seed write still goes through so the
585    /// initial value can be installed; every subsequent runtime
586    /// write returns [`SetError::FinalViolation`] with
587    /// `scope_name` so the operator can see where the pin lives.
588    ///
589    /// Use when a workload (or a parent scope) needs to
590    /// guarantee that a runtime writer cannot override a value
591    /// this scope has chosen — e.g. a DDL phase that
592    /// deliberately runs at `concurrency=1` regardless of what
593    /// the operator nudges the session-wide concurrency to.
594    pub fn final_at_scope(mut self, scope_name: impl Into<String>) -> Self {
595        self.final_at_scope = Some(scope_name.into());
596        self
597    }
598
599    /// Declare the visibility of this control within the
600    /// component tree. Defaults to [`BranchScope::Local`].
601    ///
602    /// [`BranchScope::Subtree`] means every descendant
603    /// component resolves the control name through walk-up;
604    /// descendants constructed after a successful write read
605    /// the new committed value at construction time. See
606    /// SRD 23 §"Branch-scoped and final controls".
607    pub fn branch_scope(mut self, scope: BranchScope) -> Self {
608        self.branch_scope = scope;
609        self
610    }
611
612    /// Register a converter that lets type-erased writers
613    /// (GK `control_set`, web API JSON bodies) push `f64`
614    /// values into this control. Without a converter,
615    /// f64 writes return [`SetError::ValidationFailed`] with
616    /// a "no f64 setter" message.
617    ///
618    /// For `Control<f64>` a trivial converter is `|v| Ok(v)`;
619    /// for `Control<u32>` a reasonable converter validates the
620    /// range (`if v < 0 || v > u32::MAX as f64 { Err(...) }`)
621    /// and casts. More complex types (e.g. a rate spec) can
622    /// interpret `f64` as ops/sec and synthesize a fresh
623    /// domain value.
624    pub fn from_f64<F>(mut self, f: F) -> Self
625    where
626        F: Fn(f64) -> Result<T, String> + Send + Sync + 'static,
627    {
628        self.from_f64 = Some(Box::new(f));
629        self
630    }
631
632    pub fn build(self) -> Control<T> {
633        let versioned = Versioned {
634            value: self.initial.clone(),
635            rev: 0,
636            updated_at: Instant::now(),
637            origin: ControlOrigin::Launch,
638        };
639        let validator_slot: OnceLock<Validator<T>> = OnceLock::new();
640        if let Some(v) = self.validator {
641            let _ = validator_slot.set(v);
642        }
643        // If the caller opted into gauge reification, seed the
644        // ValueGauge with the initial value (skipping when the
645        // conversion returns None).
646        let gauge = self.reify.map(|to_f64| {
647            let gauge = Arc::new(ValueGauge::new(Labels::of("control", &self.name)));
648            if let Some(f) = to_f64(&self.initial) {
649                gauge.set(f);
650            }
651            GaugeReification { gauge, to_f64 }
652        });
653        Control {
654            inner: Arc::new(ControlInner {
655                name: self.name,
656                write_lock: AsyncMutex::new(()),
657                committed: std::sync::RwLock::new(versioned),
658                appliers: Mutex::new(Vec::new()),
659                validator: validator_slot,
660                apply_timeout: self.apply_timeout,
661                rev_allocator: self.rev_allocator,
662                gauge,
663                final_at_scope: self.final_at_scope,
664                branch_scope: self.branch_scope,
665                from_f64: self.from_f64,
666            }),
667        }
668    }
669}
670
671// =========================================================================
672// Type-erased control handle
673// =========================================================================
674
675/// Erased view of a [`Control<T>`] for the registry. Lets the
676/// registry hold controls of different value types in one
677/// homogeneous map, and lets enumeration / discovery ask
678/// questions that don't need the value's static type.
679pub trait ErasedControl: Send + Sync {
680    fn name(&self) -> &str;
681    fn rev(&self) -> u64;
682    fn origin(&self) -> ControlOrigin;
683    fn applier_count(&self) -> usize;
684    /// Human-readable rendering of the current value, for
685    /// diagnostics / `dryrun=controls` output.
686    fn value_string(&self) -> String;
687    /// Static type name of `T` — stable enough to hint the
688    /// caller what concrete type they'd need to downcast to.
689    fn value_type_name(&self) -> &'static str;
690    /// The reified gauge value (if the control opted in), for
691    /// the registry's MetricSet snapshot. `None` means either
692    /// no reification or the current value has no numeric
693    /// projection right now.
694    fn gauge_f64(&self) -> Option<f64>;
695    /// True if the control was declared with gauge reification.
696    fn has_reified_gauge(&self) -> bool;
697    /// True if the control was declared `final`. Final
698    /// controls reject every non-Launch write.
699    fn is_final(&self) -> bool;
700    /// Scope name carried by a final declaration, if any.
701    fn final_scope(&self) -> Option<String>;
702    /// Declared branch scope of this control.
703    fn branch_scope(&self) -> BranchScope;
704    /// `true` if the control opts into `f64`-routed writes via
705    /// [`ControlBuilder::from_f64`].
706    fn accepts_f64_writes(&self) -> bool;
707    /// Type-erased write through the `f64` converter. Returns
708    /// the boxed async future so the caller can await from any
709    /// runtime. If the control didn't register a converter the
710    /// future resolves to [`SetError::ValidationFailed`] with a
711    /// "no f64 setter registered" message.
712    fn set_f64(
713        &self,
714        value: f64,
715        origin: ControlOrigin,
716    ) -> Pin<Box<dyn Future<Output = Result<u64, SetError>> + Send>>;
717}
718
719impl<T> ErasedControl for Control<T>
720where
721    T: Clone + Send + Sync + fmt::Debug + 'static,
722{
723    fn name(&self) -> &str {
724        Control::name(self)
725    }
726
727    fn rev(&self) -> u64 {
728        self.get().rev
729    }
730
731    fn origin(&self) -> ControlOrigin {
732        self.get().origin
733    }
734
735    fn applier_count(&self) -> usize {
736        Control::applier_count(self)
737    }
738
739    fn value_string(&self) -> String {
740        format!("{:?}", self.get().value)
741    }
742
743    fn value_type_name(&self) -> &'static str {
744        std::any::type_name::<T>()
745    }
746
747    fn gauge_f64(&self) -> Option<f64> {
748        Control::gauge_f64(self)
749    }
750
751    fn has_reified_gauge(&self) -> bool {
752        self.inner.gauge.is_some()
753    }
754
755    fn is_final(&self) -> bool {
756        Control::is_final(self)
757    }
758
759    fn final_scope(&self) -> Option<String> {
760        Control::final_scope(self).map(|s| s.to_string())
761    }
762
763    fn branch_scope(&self) -> BranchScope {
764        Control::branch_scope(self)
765    }
766
767    fn accepts_f64_writes(&self) -> bool {
768        self.inner.from_f64.is_some()
769    }
770
771    fn set_f64(
772        &self,
773        value: f64,
774        origin: ControlOrigin,
775    ) -> Pin<Box<dyn Future<Output = Result<u64, SetError>> + Send>> {
776        let inner = self.inner.clone();
777        let self_clone = Control { inner };
778        Box::pin(async move {
779            let converter = match self_clone.inner.from_f64.as_ref() {
780                Some(c) => c,
781                None => {
782                    return Err(SetError::ValidationFailed(format!(
783                        "control '{}' has no f64 setter registered — \
784                         declare one via ControlBuilder::from_f64",
785                        self_clone.inner.name,
786                    )));
787                }
788            };
789            let typed = match converter(value) {
790                Ok(t) => t,
791                Err(msg) => return Err(SetError::ValidationFailed(msg)),
792            };
793            self_clone.set(typed, origin).await
794        })
795    }
796}
797
798// =========================================================================
799// ControlRegistry — per-component store of declared controls
800// =========================================================================
801
802/// Holds every control declared on one component. Lookup by
803/// name is typed via [`Self::get::<T>`]; enumeration returns
804/// erased handles so callers can render / inspect without
805/// knowing `T`.
806#[derive(Default)]
807pub struct ControlRegistry {
808    entries: std::sync::RwLock<HashMap<String, Arc<dyn std::any::Any + Send + Sync>>>,
809    erased: std::sync::RwLock<HashMap<String, Arc<dyn ErasedControl>>>,
810}
811
812impl ControlRegistry {
813    pub fn new() -> Self {
814        Self::default()
815    }
816
817    /// Register a control. Panics if a control by the same name
818    /// already exists on this component — declaration is
819    /// structural and name collisions are bugs.
820    pub fn declare<T>(&self, control: Control<T>)
821    where
822        T: Clone + Send + Sync + fmt::Debug + 'static,
823    {
824        let name = control.name().to_string();
825        let erased: Arc<dyn ErasedControl> = Arc::new(control.clone());
826        let typed: Arc<dyn std::any::Any + Send + Sync> = Arc::new(control);
827
828        let mut entries = self.entries.write().unwrap_or_else(|e| e.into_inner());
829        let mut erased_map = self.erased.write().unwrap_or_else(|e| e.into_inner());
830        assert!(
831            !entries.contains_key(&name),
832            "ControlRegistry: duplicate control declaration for '{name}'",
833        );
834        entries.insert(name.clone(), typed);
835        erased_map.insert(name, erased);
836    }
837
838    /// Look up a control by name with a specific value type.
839    /// Returns `None` if the name doesn't exist OR if the
840    /// caller's `T` doesn't match the registered type.
841    pub fn get<T>(&self, name: &str) -> Option<Control<T>>
842    where
843        T: Clone + Send + Sync + 'static,
844    {
845        let entries = self.entries.read().unwrap_or_else(|e| e.into_inner());
846        let any = entries.get(name)?.clone();
847        drop(entries);
848        any.downcast::<Control<T>>().ok().map(|arc| (*arc).clone())
849    }
850
851    /// Erased access — returns the control's metadata without
852    /// needing to know `T`. Used by discovery / enumeration.
853    pub fn get_erased(&self, name: &str) -> Option<Arc<dyn ErasedControl>> {
854        self.erased
855            .read()
856            .unwrap_or_else(|e| e.into_inner())
857            .get(name)
858            .cloned()
859    }
860
861    /// Enumerate every declared control on this component.
862    /// Order is unspecified.
863    pub fn list(&self) -> Vec<Arc<dyn ErasedControl>> {
864        self.erased
865            .read()
866            .unwrap_or_else(|e| e.into_inner())
867            .values()
868            .cloned()
869            .collect()
870    }
871
872    pub fn len(&self) -> usize {
873        self.entries.read().unwrap_or_else(|e| e.into_inner()).len()
874    }
875
876    pub fn is_empty(&self) -> bool {
877        self.len() == 0
878    }
879
880    /// Produce a `MetricSet` with two families per declared
881    /// control:
882    ///
883    /// 1. **Numeric gauge** `control.<name>` — emitted only when
884    ///    the control has a reified-gauge projection *and* the
885    ///    current value projects to a numeric `f64`. Missing or
886    ///    `None` projections skip this family.
887    /// 2. **Info family** `control_info.<name>` — emitted for
888    ///    every control regardless of numeric projection.
889    ///    Constant value `1.0`; the current value is carried on
890    ///    a `value="..."` label so sinks that filter / group by
891    ///    dimensional labels can read the symbolic value
892    ///    directly. Follows the OpenMetrics "info" pattern (SRD
893    ///    23 §"Non-numeric controls: info family").
894    ///
895    /// Enum / bool / string controls (no numeric projection)
896    /// flow through the metrics pipeline via (2). Numeric
897    /// controls get both families so reporters can query the
898    /// running total as a number AND group by its textual form.
899    ///
900    /// Called from [`crate::component::capture_tree`] at every
901    /// scheduler tick.
902    pub fn snapshot_gauges(&self, base_labels: &Labels, captured_at: Instant) -> MetricSet {
903        let mut set = MetricSet::at(captured_at, Duration::ZERO);
904        let erased = self.erased.read().unwrap_or_else(|e| e.into_inner());
905        for ctl in erased.values() {
906            // Numeric family: only when a projection exists.
907            if let Some(value) = ctl.gauge_f64() {
908                let family_name = format!("control_{}", ctl.name());
909                let labels = base_labels.with("control", ctl.name());
910                set.insert_gauge(&family_name, labels, value, captured_at);
911            }
912            // Info family: always. Carries the symbolic value
913            // via the `value` label; the gauge itself is the
914            // OpenMetrics-style constant `1.0`. Downstream sinks
915            // filter with `control="name" value="running"`.
916            let info_family = format!("control_info_{}", ctl.name());
917            let info_labels = base_labels
918                .with("control", ctl.name())
919                .with("value", ctl.value_string());
920            set.insert_gauge(&info_family, info_labels, 1.0, captured_at);
921        }
922        set
923    }
924}
925
926// =========================================================================
927// Tests
928// =========================================================================
929
930#[cfg(test)]
931mod tests {
932    use super::*;
933    use std::sync::atomic::{AtomicU32, AtomicUsize};
934
935    // ---- RevAllocator -----------------------------------------
936
937    #[test]
938    fn rev_allocator_is_monotonic() {
939        let a = RevAllocator::new();
940        let r1 = a.next();
941        let r2 = a.next();
942        let r3 = a.next();
943        assert!(r1 < r2 && r2 < r3);
944        // No zero — callers use 0 as the "never written" sentinel
945        // so allocate strictly > 0.
946        assert!(r1 >= 1);
947    }
948
949    // ---- Control basics ---------------------------------------
950
951    fn build_u32(name: &str, initial: u32) -> Control<u32> {
952        ControlBuilder::new(name, initial)
953            .apply_timeout(Duration::from_secs(1))
954            .build()
955    }
956
957    #[tokio::test]
958    async fn initial_value_is_seeded() {
959        let c = build_u32("concurrency", 16);
960        let got = c.get();
961        assert_eq!(got.value, 16);
962        assert_eq!(got.rev, 0);
963        assert_eq!(got.origin, ControlOrigin::Launch);
964    }
965
966    #[tokio::test]
967    async fn set_with_no_appliers_commits_and_advances_rev() {
968        let c = build_u32("concurrency", 4);
969        let rev = c.set(32, ControlOrigin::Test).await.unwrap();
970        assert!(rev >= 1);
971        let v = c.get();
972        assert_eq!(v.value, 32);
973        assert_eq!(v.rev, rev);
974        assert_eq!(v.origin, ControlOrigin::Test);
975    }
976
977    #[tokio::test]
978    async fn two_successive_sets_produce_strictly_increasing_revs() {
979        let c = build_u32("c", 1);
980        let r1 = c.set(2, ControlOrigin::Test).await.unwrap();
981        let r2 = c.set(3, ControlOrigin::Test).await.unwrap();
982        assert!(r1 < r2);
983    }
984
985    // ---- Validator --------------------------------------------
986
987    #[tokio::test]
988    async fn validator_rejects_bad_values() {
989        let c: Control<u32> = ControlBuilder::new("concurrency", 4)
990            .validator(|v| {
991                if *v == 0 {
992                    Err("must be > 0".into())
993                } else if *v > 10_000 {
994                    Err("too large".into())
995                } else {
996                    Ok(())
997                }
998            })
999            .build();
1000
1001        match c.set(0, ControlOrigin::Test).await {
1002            Err(SetError::ValidationFailed(msg)) => assert!(msg.contains("must be > 0")),
1003            other => panic!("expected ValidationFailed, got {other:?}"),
1004        }
1005        // Old value untouched.
1006        assert_eq!(c.value(), 4);
1007
1008        // Valid value still commits.
1009        c.set(32, ControlOrigin::Test).await.unwrap();
1010        assert_eq!(c.value(), 32);
1011    }
1012
1013    // ---- Appliers: success path -------------------------------
1014
1015    #[tokio::test]
1016    async fn single_applier_sees_new_value() {
1017        let seen = Arc::new(AtomicU32::new(0));
1018        let c = build_u32("c", 0);
1019        let seen_clone = seen.clone();
1020        c.register_applier(SyncApplier::new(move |v: u32| {
1021            seen_clone.store(v, Ordering::SeqCst);
1022            Ok(())
1023        }));
1024        c.set(42, ControlOrigin::Test).await.unwrap();
1025        assert_eq!(seen.load(Ordering::SeqCst), 42);
1026    }
1027
1028    #[tokio::test]
1029    async fn multiple_appliers_all_see_same_value_on_success() {
1030        let c = build_u32("c", 0);
1031        let counts: Vec<Arc<AtomicU32>> = (0..5).map(|_| Arc::new(AtomicU32::new(0))).collect();
1032        for counter in &counts {
1033            let c2 = counter.clone();
1034            c.register_applier(SyncApplier::new(move |v: u32| {
1035                c2.store(v, Ordering::SeqCst);
1036                Ok(())
1037            }));
1038        }
1039        c.set(99, ControlOrigin::Test).await.unwrap();
1040        for counter in &counts {
1041            assert_eq!(counter.load(Ordering::SeqCst), 99);
1042        }
1043    }
1044
1045    // ---- Appliers: failure paths ------------------------------
1046
1047    #[tokio::test]
1048    async fn any_applier_error_fails_the_set_and_reports_index() {
1049        let c = build_u32("c", 0);
1050        // Applier 0 succeeds, 1 fails, 2 succeeds.
1051        c.register_applier(SyncApplier::new(|_| Ok(())));
1052        c.register_applier(SyncApplier::new(|_| Err("subsystem X offline".into())));
1053        c.register_applier(SyncApplier::new(|_| Ok(())));
1054
1055        match c.set(1, ControlOrigin::Test).await {
1056            Err(SetError::ApplyFailed(failures)) => {
1057                assert_eq!(failures.len(), 1);
1058                assert_eq!(failures[0].applier_index, 1);
1059                assert!(failures[0].message.contains("subsystem X offline"));
1060            }
1061            other => panic!("expected ApplyFailed, got {other:?}"),
1062        }
1063        // Committed value unchanged on failure.
1064        assert_eq!(c.value(), 0);
1065        assert_eq!(c.get().rev, 0);
1066    }
1067
1068    #[tokio::test]
1069    async fn multiple_applier_failures_all_reported() {
1070        let c = build_u32("c", 0);
1071        c.register_applier(SyncApplier::new(|_| Err("first".into())));
1072        c.register_applier(SyncApplier::new(|_| Ok(())));
1073        c.register_applier(SyncApplier::new(|_| Err("third".into())));
1074
1075        match c.set(1, ControlOrigin::Test).await {
1076            Err(SetError::ApplyFailed(failures)) => {
1077                assert_eq!(failures.len(), 2);
1078                let indices: Vec<usize> = failures.iter().map(|f| f.applier_index).collect();
1079                assert!(indices.contains(&0));
1080                assert!(indices.contains(&2));
1081            }
1082            other => panic!("expected ApplyFailed with 2 failures, got {other:?}"),
1083        }
1084        assert_eq!(c.value(), 0);
1085    }
1086
1087    #[tokio::test]
1088    async fn applier_timeout_fails_the_set() {
1089        let c: Control<u32> = ControlBuilder::new("c", 0)
1090            .apply_timeout(Duration::from_millis(50))
1091            .build();
1092
1093        struct SlowApplier;
1094        impl ControlApplier<u32> for SlowApplier {
1095            fn apply(
1096                &self,
1097                _v: u32,
1098            ) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send + '_>> {
1099                Box::pin(async {
1100                    tokio::time::sleep(Duration::from_secs(10)).await;
1101                    Ok(())
1102                })
1103            }
1104        }
1105        c.register_applier(SlowApplier);
1106
1107        match c.set(1, ControlOrigin::Test).await {
1108            Err(SetError::ApplyFailed(failures)) => {
1109                assert_eq!(failures.len(), 1);
1110                assert!(
1111                    failures[0].message.contains("timed out"),
1112                    "message = {}",
1113                    failures[0].message,
1114                );
1115            }
1116            other => panic!("expected timeout as ApplyFailed, got {other:?}"),
1117        }
1118        assert_eq!(c.value(), 0);
1119    }
1120
1121    // ---- Serialization / concurrency --------------------------
1122
1123    #[tokio::test]
1124    async fn concurrent_writers_serialize_through_write_lock() {
1125        // Two writers hit the same control concurrently. The
1126        // write lock serializes them; both commit (no drops).
1127        // Their revisions must differ and come from the same
1128        // allocator's monotonic sequence.
1129        let c = Arc::new(build_u32("c", 0));
1130        let apply_count = Arc::new(AtomicUsize::new(0));
1131        let ac = apply_count.clone();
1132        c.register_applier(SyncApplier::new(move |_: u32| {
1133            ac.fetch_add(1, Ordering::SeqCst);
1134            // Small synthetic delay so the concurrent writers
1135            // meaningfully contend on the lock.
1136            std::thread::sleep(Duration::from_millis(5));
1137            Ok(())
1138        }));
1139
1140        let c1 = c.clone();
1141        let c2 = c.clone();
1142        let h1 = tokio::spawn(async move { c1.set(10, ControlOrigin::Test).await });
1143        let h2 = tokio::spawn(async move { c2.set(20, ControlOrigin::Test).await });
1144
1145        let r1 = h1.await.unwrap().unwrap();
1146        let r2 = h2.await.unwrap().unwrap();
1147        assert_ne!(r1, r2, "revs must be distinct across concurrent writes");
1148        assert_eq!(apply_count.load(Ordering::SeqCst), 2);
1149        // Final value is one of the two, matching whichever won the race.
1150        let final_v = c.value();
1151        assert!(final_v == 10 || final_v == 20);
1152    }
1153
1154    // ---- Registry ---------------------------------------------
1155
1156    #[test]
1157    fn registry_declare_and_typed_lookup() {
1158        let reg = ControlRegistry::new();
1159        let c = build_u32("concurrency", 16);
1160        reg.declare(c);
1161        let looked_up: Control<u32> = reg.get("concurrency").unwrap();
1162        assert_eq!(looked_up.value(), 16);
1163        assert_eq!(reg.len(), 1);
1164    }
1165
1166    #[test]
1167    fn registry_get_with_wrong_type_returns_none() {
1168        let reg = ControlRegistry::new();
1169        reg.declare(build_u32("concurrency", 8));
1170        // Asking with a different T should miss cleanly.
1171        let wrong: Option<Control<u64>> = reg.get("concurrency");
1172        assert!(wrong.is_none());
1173    }
1174
1175    #[test]
1176    fn registry_missing_control_returns_none() {
1177        let reg = ControlRegistry::new();
1178        assert!(reg.get::<u32>("nope").is_none());
1179        assert!(reg.get_erased("nope").is_none());
1180    }
1181
1182    #[test]
1183    #[should_panic(expected = "duplicate control declaration")]
1184    fn registry_duplicate_declaration_panics() {
1185        let reg = ControlRegistry::new();
1186        reg.declare(build_u32("c", 1));
1187        reg.declare(build_u32("c", 2));
1188    }
1189
1190    #[tokio::test]
1191    async fn erased_view_surfaces_name_rev_and_value_string() {
1192        let reg = ControlRegistry::new();
1193        let c = build_u32("concurrency", 8);
1194        reg.declare(c.clone());
1195        c.set(16, ControlOrigin::Test).await.unwrap();
1196
1197        let erased = reg.get_erased("concurrency").unwrap();
1198        assert_eq!(erased.name(), "concurrency");
1199        assert!(erased.rev() >= 1);
1200        assert_eq!(erased.value_string(), "16");
1201        assert_eq!(erased.origin(), ControlOrigin::Test);
1202        assert_eq!(erased.applier_count(), 0);
1203    }
1204
1205    #[test]
1206    fn registry_list_returns_all_controls() {
1207        let reg = ControlRegistry::new();
1208        reg.declare(build_u32("a", 1));
1209        reg.declare(build_u32("b", 2));
1210        let listed = reg.list();
1211        assert_eq!(listed.len(), 2);
1212        let mut names: Vec<&str> = listed.iter().map(|c| c.name()).collect();
1213        names.sort();
1214        assert_eq!(names, vec!["a", "b"]);
1215    }
1216
1217    // ---- Complex value type ----------------------------------
1218
1219    // ---- Gauge reification -----------------------------------
1220
1221    #[tokio::test]
1222    async fn reified_gauge_seeds_with_initial_value() {
1223        let c: Control<u32> = ControlBuilder::new("concurrency", 32u32)
1224            .reify_as_gauge(|v| Some(*v as f64))
1225            .build();
1226        // Before any set, the gauge already holds the initial
1227        // value so a capture-tree tick right after declaration
1228        // shows the control with a sensible reading.
1229        assert_eq!(c.gauge_f64(), Some(32.0));
1230    }
1231
1232    #[tokio::test]
1233    async fn reified_gauge_updates_on_commit() {
1234        let c: Control<u32> = ControlBuilder::new("concurrency", 8u32)
1235            .reify_as_gauge(|v| Some(*v as f64))
1236            .build();
1237        c.set(64, ControlOrigin::Test).await.unwrap();
1238        assert_eq!(c.gauge_f64(), Some(64.0));
1239    }
1240
1241    #[tokio::test]
1242    async fn reified_gauge_unchanged_on_apply_failure() {
1243        let c: Control<u32> = ControlBuilder::new("concurrency", 8u32)
1244            .reify_as_gauge(|v| Some(*v as f64))
1245            .build();
1246        c.register_applier(SyncApplier::new(|_: u32| Err("no".into())));
1247        let _ = c.set(999, ControlOrigin::Test).await;
1248        // Commit never happened → gauge stays at the initial.
1249        assert_eq!(c.gauge_f64(), Some(8.0));
1250    }
1251
1252    #[tokio::test]
1253    async fn unreified_control_has_no_gauge() {
1254        let c = build_u32("c", 1);
1255        assert!(c.gauge_f64().is_none());
1256        assert!(c.reified_gauge().is_none());
1257    }
1258
1259    #[tokio::test]
1260    async fn to_f64_can_suppress_samples() {
1261        // Enum-valued controls may not have a meaningful numeric
1262        // projection for every variant — `to_f64` returning None
1263        // skips the sample rather than forcing a sentinel.
1264        #[derive(Clone, Debug, PartialEq)]
1265        enum Mode {
1266            Off,
1267            On(u32),
1268        }
1269        let c: Control<Mode> = ControlBuilder::new("explain", Mode::Off)
1270            .reify_as_gauge(|m| match m {
1271                Mode::Off => None,
1272                Mode::On(n) => Some(*n as f64),
1273            })
1274            .build();
1275        assert!(c.gauge_f64().is_none());
1276        c.set(Mode::On(7), ControlOrigin::Test).await.unwrap();
1277        assert_eq!(c.gauge_f64(), Some(7.0));
1278        c.set(Mode::Off, ControlOrigin::Test).await.unwrap();
1279        assert!(c.gauge_f64().is_none());
1280    }
1281
1282    #[tokio::test]
1283    async fn registry_snapshot_emits_numeric_gauge_for_reified_controls() {
1284        let reg = ControlRegistry::new();
1285        reg.declare(
1286            ControlBuilder::new("concurrency", 4u32)
1287                .reify_as_gauge(|v| Some(*v as f64))
1288                .build(),
1289        );
1290        reg.declare(ControlBuilder::new("non_reified", 99u32).build());
1291        let base = crate::labels::Labels::of("phase", "rampup");
1292        let now = std::time::Instant::now();
1293        let snap = reg.snapshot_gauges(&base, now);
1294
1295        // Numeric gauge for the reified control.
1296        let family = snap
1297            .family("control_concurrency")
1298            .expect("reified control should produce a numeric gauge family");
1299        let metric = family.metrics().next().unwrap();
1300        assert_eq!(metric.labels().get("phase"), Some("rampup"));
1301        assert_eq!(metric.labels().get("control"), Some("concurrency"));
1302
1303        // The non-reified control has no numeric family —
1304        // it's carried by the info family instead.
1305        assert!(snap.family("control_non_reified").is_none());
1306    }
1307
1308    #[tokio::test]
1309    async fn registry_snapshot_emits_info_family_for_every_control() {
1310        // SRD 23 §"Non-numeric controls: info family": every
1311        // control emits `control_info.<name>` with the current
1312        // symbolic value as a `value="..."` label, regardless of
1313        // whether a numeric projection exists. Numeric controls
1314        // get the info family in addition to `control.<name>`;
1315        // non-numeric controls only get the info family.
1316        let reg = ControlRegistry::new();
1317        reg.declare(
1318            ControlBuilder::new("concurrency", 4u32)
1319                .reify_as_gauge(|v| Some(*v as f64))
1320                .build(),
1321        );
1322        reg.declare(ControlBuilder::new("enabled", true).build());
1323        reg.declare(ControlBuilder::new("errors_policy", "retry".to_string()).build());
1324
1325        let base = crate::labels::Labels::of("phase", "bulk");
1326        let now = std::time::Instant::now();
1327        let snap = reg.snapshot_gauges(&base, now);
1328
1329        // Numeric control: both families.
1330        assert!(snap.family("control_concurrency").is_some());
1331        let info = snap
1332            .family("control_info_concurrency")
1333            .expect("info family should accompany the numeric gauge");
1334        let m = info.metrics().next().unwrap();
1335        assert_eq!(m.labels().get("value"), Some("4"));
1336        assert_eq!(m.labels().get("control"), Some("concurrency"));
1337
1338        // Bool control: info family only, carrying "true".
1339        assert!(snap.family("control_enabled").is_none());
1340        let info = snap
1341            .family("control_info_enabled")
1342            .expect("bool control should emit info family");
1343        let m = info.metrics().next().unwrap();
1344        assert_eq!(m.labels().get("value"), Some("true"));
1345
1346        // String control: Debug rendering wraps in quotes — the
1347        // operator-facing tooling is expected to strip those if
1348        // needed. What matters is the label survives and the
1349        // dimension exists.
1350        let info = snap
1351            .family("control_info_errors_policy")
1352            .expect("string control should emit info family");
1353        let m = info.metrics().next().unwrap();
1354        assert_eq!(m.labels().get("value"), Some("\"retry\""));
1355    }
1356
1357    #[tokio::test]
1358    async fn works_with_non_copy_value_type() {
1359        // Exercise the Clone bound with an owned Vec<String>
1360        // so the fanout/commit path genuinely clones the value.
1361        #[derive(Clone, Debug, PartialEq, Eq)]
1362        struct Targets {
1363            hosts: Vec<String>,
1364        }
1365
1366        let c: Control<Targets> = ControlBuilder::new(
1367            "targets",
1368            Targets {
1369                hosts: vec!["a".into()],
1370            },
1371        )
1372        .build();
1373
1374        let seen: Arc<Mutex<Option<Targets>>> = Arc::new(Mutex::new(None));
1375        let seen_c = seen.clone();
1376        c.register_applier(SyncApplier::new(move |t: Targets| {
1377            *seen_c.lock().unwrap() = Some(t);
1378            Ok(())
1379        }));
1380
1381        let new_val = Targets {
1382            hosts: vec!["a".into(), "b".into()],
1383        };
1384        c.set(new_val.clone(), ControlOrigin::Test).await.unwrap();
1385        assert_eq!(c.value(), new_val);
1386        assert_eq!(*seen.lock().unwrap(), Some(new_val));
1387    }
1388
1389    // ---- Final declarations -----------------------------------
1390
1391    #[tokio::test]
1392    async fn final_control_rejects_non_launch_writes() {
1393        let c: Control<u32> = ControlBuilder::new("concurrency", 1u32)
1394            .final_at_scope("ddl_phase")
1395            .build();
1396        assert!(c.is_final());
1397        assert_eq!(c.final_scope(), Some("ddl_phase"));
1398
1399        // Runtime origins (Test, Tui, Cli, ...) all reject.
1400        for origin in [ControlOrigin::Test, ControlOrigin::Tui, ControlOrigin::Cli] {
1401            match c.set(42, origin.clone()).await {
1402                Err(SetError::FinalViolation { scope }) => {
1403                    assert_eq!(scope, "ddl_phase");
1404                }
1405                other => panic!("expected FinalViolation for {origin:?}, got {other:?}"),
1406            }
1407        }
1408        // Value unchanged across all rejected writes.
1409        assert_eq!(c.value(), 1u32);
1410        assert_eq!(c.get().rev, 0);
1411    }
1412
1413    #[tokio::test]
1414    async fn final_control_accepts_launch_seed() {
1415        // Launch-origin is the one write that goes through so
1416        // the initial value can be installed. Tests the flow
1417        // where the control is declared final but the scenario
1418        // startup still needs to seed a non-default initial.
1419        let c: Control<u32> = ControlBuilder::new("concurrency", 4u32)
1420            .final_at_scope("ddl_phase")
1421            .build();
1422        let rev = c.set(1u32, ControlOrigin::Launch).await.unwrap();
1423        assert!(rev >= 1);
1424        assert_eq!(c.value(), 1u32);
1425    }
1426
1427    #[tokio::test]
1428    async fn final_rejection_runs_before_validation() {
1429        // A final control with a validator that would reject
1430        // the attempted value must still surface
1431        // FinalViolation — the pin takes precedence so operators
1432        // see the real reason the write was rejected.
1433        let c: Control<u32> = ControlBuilder::new("concurrency", 1u32)
1434            .final_at_scope("ddl_phase")
1435            .validator(|v| {
1436                if *v == 99 {
1437                    Err("no 99s".into())
1438                } else {
1439                    Ok(())
1440                }
1441            })
1442            .build();
1443        match c.set(99, ControlOrigin::Test).await {
1444            Err(SetError::FinalViolation { scope }) => {
1445                assert_eq!(scope, "ddl_phase");
1446            }
1447            other => panic!("expected FinalViolation, got {other:?}"),
1448        }
1449    }
1450
1451    #[tokio::test]
1452    async fn nonfinal_control_is_default() {
1453        let c = build_u32("concurrency", 1);
1454        assert!(!c.is_final());
1455        assert!(c.final_scope().is_none());
1456    }
1457
1458    // ---- Branch scope ----------------------------------------
1459
1460    #[test]
1461    fn default_branch_scope_is_local() {
1462        let c = build_u32("c", 1);
1463        assert_eq!(c.branch_scope(), BranchScope::Local);
1464    }
1465
1466    #[test]
1467    fn branch_scope_subtree_is_preserved_in_erased() {
1468        let c: Control<u32> = ControlBuilder::new("hdr_sigdigs", 3u32)
1469            .branch_scope(BranchScope::Subtree)
1470            .build();
1471        assert_eq!(c.branch_scope(), BranchScope::Subtree);
1472
1473        let reg = ControlRegistry::new();
1474        reg.declare(c);
1475        let erased = reg.get_erased("hdr_sigdigs").unwrap();
1476        assert_eq!(erased.branch_scope(), BranchScope::Subtree);
1477    }
1478
1479    // ---- from_f64 / set_f64 (type-erased write path) ----------
1480
1481    #[tokio::test]
1482    async fn set_f64_rejects_without_converter() {
1483        let c: Control<u32> = ControlBuilder::new("concurrency", 4u32).build();
1484        let reg = ControlRegistry::new();
1485        reg.declare(c);
1486        let erased = reg.get_erased("concurrency").unwrap();
1487        assert!(!erased.accepts_f64_writes());
1488        match erased.set_f64(8.0, ControlOrigin::Test).await {
1489            Err(SetError::ValidationFailed(msg)) => {
1490                assert!(msg.contains("no f64 setter"), "got: {msg}");
1491            }
1492            other => panic!("expected ValidationFailed, got {other:?}"),
1493        }
1494    }
1495
1496    #[tokio::test]
1497    async fn set_f64_writes_through_converter_to_typed_control() {
1498        let c: Control<u32> = ControlBuilder::new("concurrency", 4u32)
1499            .from_f64(|v| {
1500                if v < 0.0 || v > u32::MAX as f64 {
1501                    Err(format!("concurrency out of range: {v}"))
1502                } else {
1503                    Ok(v as u32)
1504                }
1505            })
1506            .build();
1507        let reg = ControlRegistry::new();
1508        reg.declare(c.clone());
1509        let erased = reg.get_erased("concurrency").unwrap();
1510        assert!(erased.accepts_f64_writes());
1511
1512        let rev = erased.set_f64(64.0, ControlOrigin::Test).await.unwrap();
1513        assert!(rev >= 1);
1514        assert_eq!(c.value(), 64u32);
1515    }
1516
1517    #[tokio::test]
1518    async fn set_f64_converter_error_is_surfaced() {
1519        let c: Control<u32> = ControlBuilder::new("concurrency", 4u32)
1520            .from_f64(|v| {
1521                if v < 0.0 {
1522                    Err(format!("negative: {v}"))
1523                } else {
1524                    Ok(v as u32)
1525                }
1526            })
1527            .build();
1528        let reg = ControlRegistry::new();
1529        reg.declare(c.clone());
1530        let erased = reg.get_erased("concurrency").unwrap();
1531        match erased.set_f64(-1.0, ControlOrigin::Test).await {
1532            Err(SetError::ValidationFailed(msg)) => {
1533                assert!(msg.contains("negative"), "got: {msg}");
1534            }
1535            other => panic!("expected ValidationFailed, got {other:?}"),
1536        }
1537        // Committed value unchanged on failed conversion.
1538        assert_eq!(c.value(), 4u32);
1539    }
1540
1541    #[tokio::test]
1542    async fn set_f64_respects_final_scope() {
1543        let c: Control<u32> = ControlBuilder::new("concurrency", 4u32)
1544            .from_f64(|v| Ok(v as u32))
1545            .final_at_scope("ddl_phase")
1546            .build();
1547        let reg = ControlRegistry::new();
1548        reg.declare(c.clone());
1549        let erased = reg.get_erased("concurrency").unwrap();
1550        match erased.set_f64(8.0, ControlOrigin::Test).await {
1551            Err(SetError::FinalViolation { scope }) => {
1552                assert_eq!(scope, "ddl_phase");
1553            }
1554            other => panic!("expected FinalViolation, got {other:?}"),
1555        }
1556    }
1557}