Skip to main content

nmbrs_metrics/
component.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Runtime component tree for metrics ownership and dimensional labels.
5//!
6//! Every Polydat context layer (session, scenario, phase, dispenser) is a
7//! [`Component`] in a parent-child tree. Labels inherit downward —
8//! a phase component's effective labels include all ancestor labels.
9//! Properties walk upward — a child can query a prop set on any ancestor.
10//!
11//! ## Instrument ownership (consolidated 2026-05)
12//!
13//! Each component carries a single `Vec<RegisteredInstrument>` — the
14//! canonical store for every instrument hung on the node. Per-cycle
15//! callers (op-dispenser wrappers, the activity executor) hold typed
16//! `Arc<...>` references captured at registration time and never
17//! look up by family name on the hot path. The [`Component::find_instrument`]
18//! linear scan exists for diagnostics / introspection only.
19//!
20//! Dynamic instruments whose existence isn't known at init —
21//! per-error-type counters allocated on first sighting — register
22//! through the [`DynamicCapture`] hook installed via
23//! [`Component::set_dynamic_capture`]. Capture walks the registry
24//! first, then invokes the dynamic hook (if any).
25
26use std::collections::HashMap;
27use std::sync::{Arc, Mutex, RwLock, Weak};
28use std::time::{Duration, Instant};
29
30use crate::instruments::counter::Counter;
31use crate::instruments::gauge::ValueGauge;
32use crate::instruments::histogram::Histogram;
33use crate::instruments::timer::Timer;
34use crate::labels::Labels;
35use crate::snapshot::MetricSet;
36
37/// Lifecycle state of a component.
38#[derive(Debug, Clone, Copy, PartialEq, Eq)]
39pub enum ComponentState {
40    /// Component is being initialized.
41    Starting,
42    /// Component is actively running. Instruments are captured.
43    Running,
44    /// Component is shutting down. Final flush pending.
45    Stopping,
46    /// Component is done. Instruments no longer captured.
47    /// Cumulative view remains queryable in the store until detach.
48    Stopped,
49}
50
51/// A typed instrument reference owned by a [`Component`].
52///
53/// One variant per kind matches the [`crate::snapshot::MetricType`]
54/// axis (counter / gauge / histogram / timer). Capture dispatches
55/// on the variant to call the right kind-specific snapshot method.
56#[derive(Clone)]
57pub enum InstrumentRef {
58    Counter(Arc<Counter>),
59    Gauge(Arc<ValueGauge>),
60    Histogram(Arc<Histogram>),
61    Timer(Arc<Timer>),
62}
63
64impl InstrumentRef {
65    /// Labels recorded on the underlying instrument. The `name=...`
66    /// pair (used by [`split_name_label`]) is preserved here so
67    /// existing snapshots keep their shape.
68    pub fn labels(&self) -> &Labels {
69        match self {
70            Self::Counter(c) => c.labels(),
71            Self::Gauge(g) => g.labels(),
72            Self::Histogram(h) => h.labels(),
73            Self::Timer(t) => t.labels(),
74        }
75    }
76}
77
78/// One registry entry: the bare family name, optional OpenMetrics
79/// unit, and the typed instrument.
80///
81/// `family` is the bare name as given to
82/// [`Component::register_instrument`]. The `_<unit>` suffix per
83/// SRD-40a §4.3 is applied at capture time (in
84/// [`crate::snapshot::MetricSet::insert_metric_with_unit`]) so the
85/// `metric_family.name` ends up suffixed and `metric_family.unit`
86/// holds the unit. Unit `None` means the family is published as-is.
87pub struct RegisteredInstrument {
88    pub family: String,
89    pub unit: Option<String>,
90    pub instrument: InstrumentRef,
91}
92
93/// Hook for components that own a dynamically-extending set of
94/// instruments — e.g. per-error-type counters allocated lazily.
95///
96/// The registry-side `Vec<RegisteredInstrument>` is the canonical
97/// store for instruments known at init. Anything that needs to
98/// register more instruments after `register_on` has run installs
99/// a `DynamicCapture` via [`Component::set_dynamic_capture`]; the
100/// component's capture path invokes it after walking the registry
101/// so the dynamic samples ride the same cadence pipeline.
102pub trait DynamicCapture: Send + Sync {
103    /// Append the dynamic instruments' current samples into `out`.
104    /// `drain` mirrors the registry walk: `true` for the cadence
105    /// reporter's per-tick path (drain histograms, etc.); `false`
106    /// for the non-mutating "current" path.
107    fn capture_into(&self, out: &mut MetricSet, now: Instant, drain: bool);
108}
109
110/// A node in the runtime component tree.
111///
112/// Components form a hierarchy: Session → Scenario → Phase → Dispenser.
113/// Each component carries its own labels, inheritable properties, and
114/// its own instrument registry.
115pub struct Component {
116    /// This component's own labels (e.g., `phase="rampup"`).
117    labels: Labels,
118    /// Effective labels = all ancestor labels merged with own labels.
119    /// Computed on [`attach`] and cached.
120    effective_labels: Labels,
121    /// Inheritable properties. Queried via walk-up to first ancestor
122    /// that has the key set. Used for `hdr_digits`, `base_interval`, etc.
123    props: HashMap<String, String>,
124    /// Weak reference to parent for prop walk-up.
125    parent: Option<Weak<RwLock<Component>>>,
126    /// Child components. Populated at runtime as phases start.
127    children: Vec<Arc<RwLock<Component>>>,
128    /// Lifecycle state. Only RUNNING components are captured.
129    state: ComponentState,
130    /// Canonical instrument store for this component.
131    ///
132    /// Hot-path callers hold typed `Arc<...>` references obtained
133    /// at registration time and never look up by name per cycle.
134    /// Family-name lookup ([`find_instrument`]) is a linear scan and
135    /// is reserved for diagnostics / introspection — see the
136    /// [`find_instrument`] doc-comment.
137    instruments: Vec<RegisteredInstrument>,
138    /// Optional hook for instruments whose existence isn't known at
139    /// init (per-error-type counters, etc.). Capture walks
140    /// `instruments` first, then invokes this if present. See
141    /// [`DynamicCapture`].
142    dynamic_capture: Option<Arc<dyn DynamicCapture>>,
143    /// Wall-clock instant of the most recent `capture_delta` /
144    /// `capture_delta_auto` call. Used by `capture_delta_auto` to
145    /// compute the true elapsed-time interval for a phase-end
146    /// flush — eliminates the 1-second quantization that comes
147    /// from stamping the partial with the scheduler's nominal
148    /// `base_interval`. `None` until the first capture; the auto
149    /// path treats that as "use caller-supplied fallback".
150    last_capture_instant: Mutex<Option<Instant>>,
151    /// Dynamic-controls declared on this component (SRD 23).
152    /// Empty unless the code that instantiates the component
153    /// explicitly declares a control via
154    /// `component.controls().declare(...)`.
155    controls: crate::controls::ControlRegistry,
156    /// Data-materialised child cells, keyed by coordinate.
157    ///
158    /// Held behind an `Arc` and handed out BY VALUE, not as a borrow through
159    /// the component's guard: resolving a cell attaches a child, which takes
160    /// this component's WRITE lock. A caller that reached the map through a
161    /// read guard would still be holding it — a self-deadlock on the same
162    /// `RwLock`. Returning the `Arc` lets the guard drop before resolution.
163    ///
164    /// A cell's lifetime is this component's: a phase's cells go when the
165    /// phase subtree does.
166    cells: std::sync::Arc<crate::cells::CellMap>,
167    /// Liveness token. `Some` from construction until this component reaches
168    /// [`ComponentState::Stopped`], then dropped.
169    ///
170    /// A parent holds `Weak` clones of its children's tokens to detect
171    /// CONCURRENT siblings sharing a label set (see [`attach`]). Weak, so the
172    /// check needs no lock on any child and no cleanup hook anywhere: the token
173    /// dies when the component stops, and again if the component is simply
174    /// dropped without stopping. There is no counter to decrement and therefore
175    /// no way to leak a phantom sibling that blocks a legitimate re-attach.
176    live: Option<std::sync::Arc<()>>,
177    /// Live children indexed by own-label rendering, for the
178    /// concurrent-sibling check in [`attach`]. Holds `Weak` tokens and is
179    /// pruned lazily on the only path that reads a bucket, so it never keeps a
180    /// dead component alive and never needs an explicit teardown pass.
181    live_children: std::collections::HashMap<String, Vec<std::sync::Weak<()>>>,
182}
183
184impl Component {
185    /// Create a new detached component with the given labels and props.
186    ///
187    /// The component starts in [`ComponentState::Starting`]. Call
188    /// [`attach`] to wire it into the tree and compute effective labels.
189    pub fn new(labels: Labels, props: HashMap<String, String>) -> Self {
190        Self {
191            effective_labels: labels.clone(),
192            labels,
193            props,
194            parent: None,
195            children: Vec::new(),
196            state: ComponentState::Starting,
197            instruments: Vec::new(),
198            dynamic_capture: None,
199            last_capture_instant: Mutex::new(None),
200            controls: crate::controls::ControlRegistry::new(),
201            cells: std::sync::Arc::new(crate::cells::CellMap::new()),
202            live: Some(std::sync::Arc::new(())),
203            live_children: std::collections::HashMap::new(),
204        }
205    }
206
207    /// Create a root component (session level). No parent.
208    pub fn root(labels: Labels, props: HashMap<String, String>) -> Arc<RwLock<Self>> {
209        let mut component = Self::new(labels, props);
210        component.state = ComponentState::Running;
211        Arc::new(RwLock::new(component))
212    }
213
214    /// This component's own labels (not including ancestors).
215    pub fn labels(&self) -> &Labels {
216        &self.labels
217    }
218
219    /// Effective labels: all ancestor labels merged with own labels.
220    pub fn effective_labels(&self) -> &Labels {
221        &self.effective_labels
222    }
223
224    /// Current lifecycle state.
225    pub fn state(&self) -> ComponentState {
226        self.state
227    }
228
229    /// Transition to a new lifecycle state.
230    ///
231    /// Reaching [`ComponentState::Stopped`] drops the liveness token, which is
232    /// what releases this component's claim on its label set among its
233    /// siblings. Every stop path goes through here, so there is exactly one
234    /// place that can forget to do it.
235    pub fn set_state(&mut self, state: ComponentState) {
236        self.state = state;
237        if state == ComponentState::Stopped {
238            self.live = None;
239        }
240    }
241
242    /// A `Weak` handle to this component's liveness token, for a parent's
243    /// concurrent-sibling index. Dead once the component stops or is dropped.
244    fn live_handle(&self) -> std::sync::Weak<()> {
245        match &self.live {
246            Some(t) => std::sync::Arc::downgrade(t),
247            None => std::sync::Weak::new(),
248        }
249    }
250
251    /// Register an instrument under `family` on this component.
252    ///
253    /// Returns `Err` when `family` is already registered on this
254    /// component — duplicate-family declarations on the same
255    /// dimensional cell surface as a workload error here, before
256    /// any cycle runs (SRD-40b §7.2). The component's
257    /// `effective_labels` define the dimensional cell; the same
258    /// family on a different component is a different cell and
259    /// produces no collision.
260    ///
261    /// The collision check is a linear scan over the registry
262    /// Vec — see the storage-shape note on [`Self::instruments`].
263    pub fn register_instrument(
264        &mut self,
265        family: impl Into<String>,
266        instrument: InstrumentRef,
267    ) -> Result<(), String> {
268        self.register_instrument_with_unit(family, None, instrument)
269    }
270
271    /// Variant of [`Self::register_instrument`] that records an
272    /// OpenMetrics unit (`ms`, `bytes`, …).
273    ///
274    /// At capture time the unit drives the `_<unit>` suffix on
275    /// `metric_family.name` and populates the `unit` column per
276    /// SRD-40a §4.3 / SRD-40b §1. `None` is identical to the
277    /// no-unit `register_instrument` path.
278    pub fn register_instrument_with_unit(
279        &mut self,
280        family: impl Into<String>,
281        unit: Option<String>,
282        instrument: InstrumentRef,
283    ) -> Result<(), String> {
284        let family = family.into();
285        if self.instruments.iter().any(|ri| ri.family == family) {
286            return Err(format!(
287                "duplicate family name on dimensionally-same metric \
288                 context: {family}{}",
289                self.effective_labels.to_prometheus(),
290            ));
291        }
292        self.instruments.push(RegisteredInstrument {
293            family,
294            unit,
295            instrument,
296        });
297        Ok(())
298    }
299
300    /// Read-only view of every registered instrument on this
301    /// component, in insertion order. Walked by the cadence
302    /// reporter on every tick.
303    pub fn instruments(&self) -> &[RegisteredInstrument] {
304        &self.instruments
305    }
306
307    /// Linear scan by family name — diagnostic / rare-path only.
308    ///
309    /// Hot-path callers must use the typed `Arc<...>` they
310    /// captured at registration time. The Vec storage and linear
311    /// scan are deliberate: registration is once-at-init,
312    /// per-cycle access is pre-bound, and a HashMap probe would
313    /// add API + Hash bound for ~40 ns saved once per workload load.
314    ///
315    /// If you find yourself reaching for this on a per-cycle code
316    /// path, that's a design bug in the caller — pre-bind the
317    /// `Arc<...>` you got from [`register_instrument`] instead.
318    pub fn find_instrument(&self, family: &str) -> Option<&InstrumentRef> {
319        self.instruments
320            .iter()
321            .find(|ri| ri.family == family)
322            .map(|ri| &ri.instrument)
323    }
324
325    /// Install a [`DynamicCapture`] hook for instruments whose
326    /// existence isn't known at init time. Replaces any prior
327    /// installation. See the trait doc.
328    pub fn set_dynamic_capture(&mut self, hook: Arc<dyn DynamicCapture>) {
329        self.dynamic_capture = Some(hook);
330    }
331
332    /// Capture a delta snapshot covering `interval`.
333    ///
334    /// Drains histogram/timer reservoirs; counters report their
335    /// absolute running total (no draining). Called by the scheduler on
336    /// every tick — the result feeds the cadence reporter's
337    /// smallest-cadence accumulator. The caller-supplied
338    /// `interval` is recorded on the snapshot verbatim; the
339    /// scheduler passes its nominal `base_interval` so the
340    /// "canonical scheduler cadence" property is preserved
341    /// in storage even when wall-clock between ticks drifts
342    /// (drift surfaces via the scheduler's tick warning,
343    /// not by mutating the snapshot's interval).
344    ///
345    /// Also stamps `last_capture_instant` so the phase-end
346    /// flush path (`capture_delta_auto`) can compute true
347    /// elapsed time since the previous capture.
348    pub fn capture_delta(&self, interval: Duration) -> MetricSet {
349        let now = Instant::now();
350        let mut out = MetricSet::at(now, interval);
351        self.capture_registry_into(&mut out, now, true);
352        if let Some(hook) = &self.dynamic_capture {
353            hook.capture_into(&mut out, now, true);
354        }
355        *self
356            .last_capture_instant
357            .lock()
358            .unwrap_or_else(|e| e.into_inner()) = Some(now);
359        out
360    }
361
362    /// Phase-end variant of [`capture_delta`] that stamps the
363    /// snapshot with **real elapsed wall time** since the
364    /// previous capture, rather than a caller-supplied
365    /// nominal interval.
366    ///
367    /// Eliminates the 1-second quantization that surfaced as
368    /// the spurious `cycles_total_rate = 10000/8 = 1250.0`
369    /// cluster for short phases — a 7.84-second phase will now
370    /// carry a 7843-ms interval (subject to storage precision)
371    /// instead of being padded to the scheduler's nominal 1s
372    /// final-flush stamp.
373    ///
374    /// `fallback` covers the edge case where no prior capture
375    /// has happened (a phase that ended before the first
376    /// scheduler tick). The phase-end caller passes the
377    /// scheduler's nominal `base_interval` here — a phase
378    /// shorter than one tick still gets stamped with that
379    /// nominal duration rather than zero.
380    pub fn capture_delta_auto(&self, fallback: Duration) -> MetricSet {
381        let now = Instant::now();
382        let interval = {
383            let mut prev = self
384                .last_capture_instant
385                .lock()
386                .unwrap_or_else(|e| e.into_inner());
387            let elapsed = prev.map(|t| now.duration_since(t));
388            *prev = Some(now);
389            elapsed.unwrap_or(fallback)
390        };
391        let mut out = MetricSet::at(now, interval);
392        self.capture_registry_into(&mut out, now, true);
393        if let Some(hook) = &self.dynamic_capture {
394            hook.capture_into(&mut out, now, true);
395        }
396        out
397    }
398
399    /// Capture a non-mutating snapshot of current state.
400    ///
401    /// - Counters: absolute totals (atomic load).
402    /// - Gauges: current value.
403    /// - Histograms / Timers: non-draining clone (`peek_snapshot`).
404    ///
405    /// Never touches internal accumulators — callers may invoke
406    /// this arbitrarily often without perturbing the scheduler's
407    /// per-tick cascade.
408    pub fn capture_current(&self) -> MetricSet {
409        let now = Instant::now();
410        let mut out = MetricSet::at(now, Duration::ZERO);
411        self.capture_registry_into(&mut out, now, false);
412        if let Some(hook) = &self.dynamic_capture {
413            hook.capture_into(&mut out, now, false);
414        }
415        out
416    }
417
418    /// Walk the registered instruments and emit their samples into
419    /// `out`. `drain=true` drains histogram/timer reservoirs;
420    /// `drain=false` peeks without disturbing them. Counters always
421    /// report their absolute running total — `drain` does not apply.
422    fn capture_registry_into(&self, out: &mut MetricSet, now: Instant, drain: bool) {
423        for ri in &self.instruments {
424            let family = ri.family.clone();
425            let unit = ri.unit.as_deref();
426            match &ri.instrument {
427                InstrumentRef::Counter(c) => {
428                    let lbl = strip_name_label(c.labels());
429                    // A counter is its absolute running total — captured the
430                    // same on every path (the `drain` flag only governs
431                    // histogram/timer reservoirs). Per-interval deltas are
432                    // derived downstream by differencing samples. See the
433                    // cumulative-counter note.
434                    out.insert_counter_with_unit(family, unit, lbl, c.get(), now);
435                }
436                InstrumentRef::Gauge(g) => {
437                    let lbl = strip_name_label(g.labels());
438                    out.insert_gauge_with_unit(family, unit, lbl, g.get(), now);
439                }
440                InstrumentRef::Histogram(h) => {
441                    let lbl = strip_name_label(h.labels());
442                    let reservoir = if drain {
443                        h.snapshot()
444                    } else {
445                        h.peek_snapshot()
446                    };
447                    // `cumulative_count` is the instrument's lifetime total
448                    // (monotonic, never drained); the reservoir carries the
449                    // per-window distribution. See the cumulative-counter note.
450                    out.insert_histogram_with_unit_cumulative(
451                        family,
452                        unit,
453                        lbl,
454                        reservoir,
455                        h.total(),
456                        now,
457                    );
458                }
459                InstrumentRef::Timer(t) => {
460                    let lbl = strip_name_label(t.labels());
461                    let snap = if drain {
462                        t.snapshot()
463                    } else {
464                        t.peek_snapshot()
465                    };
466                    // `snap.count` is the Timer's absolute lifetime count.
467                    out.insert_histogram_with_unit_cumulative(
468                        family,
469                        unit,
470                        lbl,
471                        snap.histogram,
472                        snap.count,
473                        now,
474                    );
475                }
476            }
477        }
478    }
479
480    /// Get a property by name, walking up the tree.
481    ///
482    /// Checks this component's props first, then each ancestor in
483    /// order until found. Returns `None` if no ancestor has the key.
484    pub fn get_prop(&self, name: &str) -> Option<String> {
485        if let Some(value) = self.props.get(name) {
486            return Some(value.clone());
487        }
488        if let Some(ref parent_weak) = self.parent
489            && let Some(parent_arc) = parent_weak.upgrade()
490            && let Ok(parent) = parent_arc.read()
491        {
492            return parent.get_prop(name);
493        }
494        None
495    }
496
497    /// Set a property on this component.
498    pub fn set_prop(&mut self, name: &str, value: &str) {
499        self.props.insert(name.to_string(), value.to_string());
500    }
501
502    /// Number of child components.
503    pub fn child_count(&self) -> usize {
504        self.children.len()
505    }
506
507    /// Iterator over this component's direct children.
508    pub fn children(&self) -> impl Iterator<Item = &Arc<RwLock<Component>>> {
509        self.children.iter()
510    }
511
512    /// Borrow this component's dynamic-controls registry. Every
513    /// component carries one; empty until something declares a
514    /// control on it. See SRD 23.
515    pub fn controls(&self) -> &crate::controls::ControlRegistry {
516        &self.controls
517    }
518
519    /// This component's data-materialised cells. See [`crate::cells::CellMap`]:
520    /// one series per dimension instance is one CHILD per instance, not a label
521    /// bag on an instrument.
522    ///
523    /// Returns the `Arc` by value ON PURPOSE — see the field docs. Resolving
524    /// takes this component's write lock, so the caller must be able to drop
525    /// its read guard first:
526    ///
527    /// ```ignore
528    /// let cells = parent.read().unwrap().cells();  // guard dropped here
529    /// let cell  = cells.resolve(&parent, &coord);  // takes the write lock
530    /// ```
531    pub fn cells(&self) -> std::sync::Arc<crate::cells::CellMap> {
532        self.cells.clone()
533    }
534
535    /// Resolve a typed control by name, walking up the parent
536    /// chain. This component's registry is checked first; then
537    /// each ancestor in order. An ancestor declaration is only
538    /// honored if its [`BranchScope`] is `Subtree` —
539    /// [`BranchScope::Local`] declarations do not propagate to
540    /// descendants. Returns `None` if no in-scope declaration
541    /// matches the `<name, T>` pair.
542    ///
543    /// Mirrors [`Self::get_prop`] but for typed controls (SRD 23
544    /// §"Branch-scoped and final controls").
545    pub fn find_control_up<T>(&self, name: &str) -> Option<crate::controls::Control<T>>
546    where
547        T: Clone + Send + Sync + 'static,
548    {
549        if let Some(c) = self.controls.get::<T>(name) {
550            return Some(c);
551        }
552        if let Some(ref parent_weak) = self.parent
553            && let Some(parent_arc) = parent_weak.upgrade()
554            && let Ok(parent) = parent_arc.read()
555        {
556            return parent.find_control_up_subtree::<T>(name);
557        }
558        None
559    }
560
561    /// Ancestor-side recursion. Honors [`BranchScope::Subtree`]
562    /// on an ancestor's declaration; otherwise keeps walking.
563    fn find_control_up_subtree<T>(&self, name: &str) -> Option<crate::controls::Control<T>>
564    where
565        T: Clone + Send + Sync + 'static,
566    {
567        if let Some(erased) = self.controls.get_erased(name)
568            && erased.branch_scope() == crate::controls::BranchScope::Subtree
569            && let Some(c) = self.controls.get::<T>(name)
570        {
571            return Some(c);
572        }
573        if let Some(ref parent_weak) = self.parent
574            && let Some(parent_arc) = parent_weak.upgrade()
575            && let Ok(parent) = parent_arc.read()
576        {
577            return parent.find_control_up_subtree::<T>(name);
578        }
579        None
580    }
581
582    /// Erased variant of [`Self::find_control_up`] — returns
583    /// just the enumeration handle, useful for diagnostics
584    /// (`dryrun=controls`, TUI surfaces) that don't need the
585    /// typed value.
586    pub fn find_control_erased_up(
587        &self,
588        name: &str,
589    ) -> Option<std::sync::Arc<dyn crate::controls::ErasedControl>> {
590        if let Some(erased) = self.controls.get_erased(name) {
591            return Some(erased);
592        }
593        if let Some(ref parent_weak) = self.parent
594            && let Some(parent_arc) = parent_weak.upgrade()
595            && let Ok(parent) = parent_arc.read()
596        {
597            return parent.find_control_erased_up_subtree(name);
598        }
599        None
600    }
601
602    fn find_control_erased_up_subtree(
603        &self,
604        name: &str,
605    ) -> Option<std::sync::Arc<dyn crate::controls::ErasedControl>> {
606        if let Some(erased) = self.controls.get_erased(name)
607            && erased.branch_scope() == crate::controls::BranchScope::Subtree
608        {
609            return Some(erased);
610        }
611        if let Some(ref parent_weak) = self.parent
612            && let Some(parent_arc) = parent_weak.upgrade()
613            && let Ok(parent) = parent_arc.read()
614        {
615            return parent.find_control_erased_up_subtree(name);
616        }
617        None
618    }
619
620    /// SRD-89 — flatten the up-walk control resolution starting at `start`
621    /// into a name → erased-handle map, computed **without nested locks**.
622    ///
623    /// This has the same visibility as calling
624    /// [`Self::find_control_erased_up`] for every name — the start
625    /// component's own controls plus any `BranchScope::Subtree` control on an
626    /// ancestor, nearest-wins — but it acquires and releases **one** tier's
627    /// lock at a time (never holding a child's guard while reading a parent),
628    /// so it is immune to the writer-preferring `RwLock` starvation that a
629    /// nested up-walk hits under concurrent in-process execution (a hot-path
630    /// per-cycle `find_control_erased_up` deadlocks against the cadence
631    /// path's instrument-registration writes; this is built once per phase
632    /// and read lock-free thereafter — see SRD-89 §3c-i).
633    pub fn control_snapshot(
634        start: &std::sync::Arc<std::sync::RwLock<Component>>,
635    ) -> std::collections::HashMap<String, std::sync::Arc<dyn crate::controls::ErasedControl>> {
636        let mut map: std::collections::HashMap<
637            String,
638            std::sync::Arc<dyn crate::controls::ErasedControl>,
639        > = std::collections::HashMap::new();
640        let mut next = Some(start.clone());
641        let mut is_start = true;
642        while let Some(arc) = next {
643            let parent_next;
644            {
645                let g = arc.read().unwrap_or_else(|e| e.into_inner());
646                for handle in g.controls.list() {
647                    // The start component's own controls are always visible;
648                    // an ancestor's only if subtree-scoped. Nearest wins.
649                    if is_start || handle.branch_scope() == crate::controls::BranchScope::Subtree {
650                        map.entry(handle.name().to_string()).or_insert(handle);
651                    }
652                }
653                parent_next = g.parent.as_ref().and_then(|w| w.upgrade());
654            }
655            next = parent_next;
656            is_start = false;
657        }
658        map
659    }
660
661    /// Count of `Running`-state descendants (this component's
662    /// children, grandchildren, …). Used by callers that want a
663    /// structural "how many phases are in flight?" query against
664    /// the live component tree — e.g. the TUI's Focus-LOD
665    /// placeholder decision (SRD 62 §"Scenario done?").
666    ///
667    /// The component itself is NOT counted — the query is meant
668    /// to traverse from an activity root into its phases.
669    pub fn running_descendant_count(&self) -> usize {
670        let mut count = 0;
671        for child in &self.children {
672            if let Ok(c) = child.read() {
673                if c.state == ComponentState::Running {
674                    count += 1;
675                }
676                count += c.running_descendant_count();
677            }
678        }
679        count
680    }
681}
682
683/// Strip the legacy `name=...` label from an instrument's `Labels`,
684/// returning the dimensional residual that goes onto the captured
685/// `MetricFamily` row. The family name itself is provided
686/// separately by [`RegisteredInstrument::family`] — historical
687/// instruments embedded the family name as a `name=...` label, but
688/// that pair must NOT appear on the metric's `LabelSet` (label-set
689/// uniqueness within a family per OpenMetrics §4.5.1 would
690/// otherwise be polluted).
691fn strip_name_label(labels: &Labels) -> Labels {
692    let mut out = Labels::default();
693    for (k, v) in labels.iter() {
694        if k != "name" {
695            out = out.with(k, v);
696        }
697    }
698    out
699}
700
701/// SRD-40b §11 / SRD-42 §"Component lifecycle: scope_close flush" —
702/// fused teardown helper. Captures a final delta from this component's
703/// instruments, fires the cadence reporter's `scope_close` (which
704/// marks the partial, ingests, and closes the path), and transitions
705/// the component to [`ComponentState::Stopped`].
706///
707/// **Only acts on `Running` components.** Components that are
708/// `Starting`, `Stopping`, or already `Stopped` return without
709/// touching the reporter — calling `scope_close` twice on the same
710/// component is a no-op on the second call, matching SRD-40 §
711/// component lifecycle.
712///
713/// Components with no registered instruments still close the path so
714/// any in-flight prebuffer at that label set (e.g. ingests routed via
715/// a sibling layer) flushes through the cascade.
716pub fn scope_close(
717    component: &Arc<RwLock<Component>>,
718    cadence_reporter: &crate::cadence_reporter::CadenceReporter,
719    interval: Duration,
720) {
721    // Read-capture the delta first, then take a write guard to
722    // transition state. Read guard is released between the two so
723    // the write doesn't deadlock.
724    let (labels, delta) = {
725        let g = component.read().unwrap_or_else(|e| e.into_inner());
726        if g.state != ComponentState::Running {
727            return;
728        }
729        let delta = g.capture_delta(interval);
730        (g.effective_labels.clone(), delta)
731    };
732
733    cadence_reporter.scope_close(&labels, delta);
734
735    // Transition to Stopped so a subsequent capture pass skips this
736    // component (capture_tree only walks Running) and a second
737    // scope_close call is a no-op.
738    let mut g = component.write().unwrap_or_else(|e| e.into_inner());
739    g.state = ComponentState::Stopped;
740}
741
742// =========================================================================
743// Selector-based lookup (SRD 24)
744// =========================================================================
745
746/// Collect every component in a subtree whose `effective_labels`
747/// match the selector. Order is pre-order DFS: root first, then
748/// each child's subtree in insertion order.
749///
750/// Selector-based lookup takes an `Arc<RwLock<Component>>` root
751/// (rather than `&Component`) because the results must also be
752/// `Arc<RwLock<Component>>` — `find`'s Vec is a list of live
753/// handles callers can then mutate, not a snapshot. Scoping a
754/// query to a subtree is expressed by passing that subtree's
755/// `Arc` as the root.
756///
757/// The query is read-only against every visited component: a
758/// failed `read()` (e.g. poisoned lock) is treated as "this
759/// subtree is opaque for this query" and silently skipped.
760/// Collection order is stable across calls as long as no
761/// components are attached or detached mid-traversal.
762pub fn find(
763    root: &Arc<RwLock<Component>>,
764    sel: &crate::selector::Selector,
765) -> Vec<Arc<RwLock<Component>>> {
766    let mut out = Vec::new();
767    find_into(root, sel, &mut out);
768    out
769}
770
771fn find_into(
772    root: &Arc<RwLock<Component>>,
773    sel: &crate::selector::Selector,
774    out: &mut Vec<Arc<RwLock<Component>>>,
775) {
776    let Ok(guard) = root.read() else { return };
777    if sel.matches(&guard.effective_labels) {
778        out.push(root.clone());
779    }
780    let children = guard.children.clone();
781    drop(guard);
782    for child in &children {
783        find_into(child, sel, out);
784    }
785}
786
787/// Expect exactly one match. Returns
788/// [`crate::selector::LookupError::NotFound`] or
789/// [`crate::selector::LookupError::Ambiguous`] otherwise.
790///
791/// Short-circuits on the second hit — the Vec-returning [`find`]
792/// is preferable when all matches are wanted.
793pub fn find_one(
794    root: &Arc<RwLock<Component>>,
795    sel: &crate::selector::Selector,
796) -> Result<Arc<RwLock<Component>>, crate::selector::LookupError> {
797    let mut first: Option<Arc<RwLock<Component>>> = None;
798    let mut count = 0usize;
799    find_one_walk(root, sel, &mut first, &mut count);
800    match first {
801        None => Err(crate::selector::LookupError::NotFound),
802        Some(c) if count == 1 => Ok(c),
803        Some(_) => Err(crate::selector::LookupError::Ambiguous { count }),
804    }
805}
806
807fn find_one_walk(
808    root: &Arc<RwLock<Component>>,
809    sel: &crate::selector::Selector,
810    first: &mut Option<Arc<RwLock<Component>>>,
811    count: &mut usize,
812) {
813    let Ok(guard) = root.read() else { return };
814    if sel.matches(&guard.effective_labels) {
815        *count += 1;
816        if first.is_none() {
817            *first = Some(root.clone());
818        }
819        // Keep walking so `count` reflects the total — callers
820        // rely on the Ambiguous count.
821    }
822    let children = guard.children.clone();
823    drop(guard);
824    for child in &children {
825        find_one_walk(child, sel, first, count);
826    }
827}
828
829/// True if any component in the subtree matches. Short-circuits
830/// on the first hit.
831pub fn any(root: &Arc<RwLock<Component>>, sel: &crate::selector::Selector) -> bool {
832    any_walk(root, sel)
833}
834
835fn any_walk(root: &Arc<RwLock<Component>>, sel: &crate::selector::Selector) -> bool {
836    let Ok(guard) = root.read() else { return false };
837    if sel.matches(&guard.effective_labels) {
838        return true;
839    }
840    let children = guard.children.clone();
841    drop(guard);
842    children.iter().any(|c| any_walk(c, sel))
843}
844
845/// Count every matching component in the subtree.
846pub fn count(root: &Arc<RwLock<Component>>, sel: &crate::selector::Selector) -> usize {
847    let mut n = 0usize;
848    count_walk(root, sel, &mut n);
849    n
850}
851
852fn count_walk(root: &Arc<RwLock<Component>>, sel: &crate::selector::Selector, n: &mut usize) {
853    let Ok(guard) = root.read() else { return };
854    if sel.matches(&guard.effective_labels) {
855        *n += 1;
856    }
857    let children = guard.children.clone();
858    drop(guard);
859    for child in &children {
860        count_walk(child, sel, n);
861    }
862}
863
864/// Attach a child component to a parent.
865///
866/// Computes the child's effective labels by composing the parent's
867/// effective labels with the child's own labels. Adds the child to
868/// the parent's children list and sets the child's parent reference.
869///
870/// **Label-ownership invariant.** A dimensional label name is owned
871/// by exactly one component in any ancestor chain: once a name is
872/// set on a component at initialization, no descendant may redeclare
873/// it — neither with a differing value (which would silently corrupt
874/// the dimensional cell) nor with the same value (which makes
875/// ownership ambiguous). The session tier owns `session`, the
876/// execution tier owns `{exec_id, workload}`, the phase tier owns
877/// `{phase, …for_each}`, and so on down — each component declares
878/// ONLY the labels it introduces. This check enforces that at
879/// attach time (init, not per-cycle): a collision is a construction
880/// bug and panics with both label sets named, rather than letting
881/// the composition silently pick a winner.
882///
883/// **Concurrent-sibling invariant.** The rule above is *vertical* — it
884/// constrains a child against its ANCESTORS. Two SIBLINGS declaring the
885/// same own-labels compose byte-identical `effective_labels`, and the same
886/// family registered on each then yields two instruments sharing one metric
887/// identity, which the per-component duplicate-family check cannot see
888/// because they are different components.
889///
890/// That is rejected only when the siblings are alive AT THE SAME TIME.
891/// Sequential reuse is legitimate: an iteration whose values repeat (the fib
892/// comprehension yields `n=1` twice) re-materialises the same identity,
893/// which is one identity sampled again over time, not a second identity. An
894/// unconditional check was implemented and rejected for exactly that reason.
895///
896/// Liveness is tracked by a token each component holds until it reaches
897/// [`ComponentState::Stopped`]; the parent indexes `Weak` clones per
898/// own-label set. So the check needs no lock on any child, costs a lookup in
899/// one bucket, and cleans up with no teardown pass — a stopped or dropped
900/// component's claim simply expires.
901pub fn attach(parent: &Arc<RwLock<Component>>, child: &Arc<RwLock<Component>>) {
902    let parent_effective = {
903        let p = parent.read().unwrap_or_else(|e| e.into_inner());
904        p.effective_labels.clone()
905    };
906    let mut c = child.write().unwrap_or_else(|e| e.into_inner());
907    if let Some((k, _)) = c
908        .labels
909        .iter()
910        .find(|(k, _)| parent_effective.get(k).is_some())
911    {
912        panic!(
913            "component label-ownership violation: child re-declares label `{k}` \
914             already owned by an ancestor. Each label name must be set on exactly \
915             one tier and inherited downward (child {}, ancestors {}).",
916            c.labels.to_prometheus(),
917            parent_effective.to_prometheus(),
918        );
919    }
920    c.effective_labels = parent_effective.extend(&c.labels);
921    c.parent = Some(Arc::downgrade(parent));
922    let child_own = c.labels.to_prometheus();
923    let child_live = c.live_handle();
924    drop(c);
925
926    let mut p = parent.write().unwrap_or_else(|e| e.into_inner());
927    // Concurrent-sibling check — the horizontal half of the ownership rule,
928    // scoped to components that are alive AT THE SAME TIME.
929    //
930    // Sequential reuse of a label set is legitimate and must stay legal: an
931    // iteration whose values repeat (the fib comprehension yields `n=1` twice)
932    // re-materialises the same identity, which is one identity sampled again
933    // over time. Two components alive at once is a different thing entirely —
934    // each can register the same family, and the two instruments then share one
935    // metric identity with the per-component duplicate check unable to see it.
936    //
937    // Only this key's bucket is touched, and pruning happens on the same visit,
938    // so the check is O(live siblings sharing this exact label set) — normally
939    // zero — and the index self-cleans without a teardown pass.
940    let bucket = p.live_children.entry(child_own.clone()).or_default();
941    bucket.retain(|w| w.strong_count() > 0);
942    if !bucket.is_empty() {
943        let parent_labels = parent_effective.to_prometheus();
944        panic!(
945            "component sibling-identity violation: a LIVE sibling already \
946             declares the own-label set {child_own} under {parent_labels}. Both \
947             would compose byte-identical effective labels, so the same family \
948             registered on each yields two instruments sharing one metric \
949             identity — which the per-component duplicate-family check cannot \
950             see, because they are different components. (Re-using a label set \
951             AFTER the previous component stops is fine and is not this.)"
952        );
953    }
954    bucket.push(child_live);
955    p.children.push(child.clone());
956}
957
958/// Detach a child component from its parent.
959///
960/// Removes the child from the parent's children list and clears
961/// the child's parent reference.
962pub fn detach(parent: &Arc<RwLock<Component>>, child: &Arc<RwLock<Component>>) {
963    let mut p = parent.write().unwrap_or_else(|e| e.into_inner());
964    p.children.retain(|c| !Arc::ptr_eq(c, child));
965    let mut c = child.write().unwrap_or_else(|e| e.into_inner());
966    c.parent = None;
967}
968
969/// Walk the component tree and capture delta snapshots from all
970/// RUNNING components.
971///
972/// Returns one `(effective_labels, snapshot)` pair per captured
973/// component. Draining semantics — used by the scheduler tick.
974pub fn capture_tree(root: &Arc<RwLock<Component>>, interval: Duration) -> Vec<(Labels, MetricSet)> {
975    let mut results = Vec::new();
976    capture_recursive(root, interval, &mut results);
977    results
978}
979
980fn capture_recursive(
981    node: &Arc<RwLock<Component>>,
982    interval: Duration,
983    results: &mut Vec<(Labels, MetricSet)>,
984) {
985    // Take a read guard, snapshot the values we need, drop the
986    // guard before recursing so child locks don't nest on ours.
987    let Ok(guard) = node.read() else { return };
988    let state = guard.state;
989    let effective_labels = guard.effective_labels.clone();
990    let children = guard.children.clone();
991
992    if state == ComponentState::Running {
993        let snapshot = guard.capture_delta(interval);
994        if !snapshot.is_empty() {
995            results.push((effective_labels.clone(), snapshot));
996        }
997        // Reified control gauges — one per declared control that
998        // has a numeric projection. Published at every tick so
999        // they flow through the same sinks as regular metrics.
1000        let control_gauges = guard
1001            .controls
1002            .snapshot_gauges(&effective_labels, Instant::now());
1003        if !control_gauges.is_empty() {
1004            results.push((effective_labels, control_gauges));
1005        }
1006    }
1007
1008    drop(guard);
1009    for child in &children {
1010        capture_recursive(child, interval, results);
1011    }
1012}
1013
1014/// Non-mutating counterpart of [`capture_tree`]. Walks every RUNNING
1015/// component and returns absolute/peeked snapshots via
1016/// [`Component::capture_current`]. Safe to call arbitrarily often —
1017/// doesn't drain histogram/timer reservoirs.
1018pub fn capture_tree_current(root: &Arc<RwLock<Component>>) -> Vec<(Labels, MetricSet)> {
1019    let mut results = Vec::new();
1020    capture_current_recursive(root, &mut results);
1021    results
1022}
1023
1024fn capture_current_recursive(
1025    node: &Arc<RwLock<Component>>,
1026    results: &mut Vec<(Labels, MetricSet)>,
1027) {
1028    let Ok(guard) = node.read() else { return };
1029    let state = guard.state;
1030    let effective_labels = guard.effective_labels.clone();
1031    let children = guard.children.clone();
1032
1033    if state == ComponentState::Running {
1034        let snapshot = guard.capture_current();
1035        if !snapshot.is_empty() {
1036            results.push((effective_labels.clone(), snapshot));
1037        }
1038        let control_gauges = guard
1039            .controls
1040            .snapshot_gauges(&effective_labels, Instant::now());
1041        if !control_gauges.is_empty() {
1042            results.push((effective_labels, control_gauges));
1043        }
1044    }
1045
1046    drop(guard);
1047    for child in &children {
1048        capture_current_recursive(child, results);
1049    }
1050}
1051
1052#[cfg(test)]
1053mod tests {
1054    use super::*;
1055    use std::sync::atomic::{AtomicU64, Ordering};
1056
1057    fn new_counter(family: &str) -> Arc<Counter> {
1058        Arc::new(Counter::new(Labels::of("name", family)))
1059    }
1060
1061    // ── SRD-40b §7.2: register_instrument duplicate detection ──
1062
1063    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1064    async fn register_instrument_first_time_succeeds() {
1065        let mut c = Component::new(Labels::empty(), HashMap::new());
1066        assert!(
1067            c.register_instrument(
1068                "recall_at_10",
1069                InstrumentRef::Counter(new_counter("recall_at_10")),
1070            )
1071            .is_ok()
1072        );
1073        assert!(c.find_instrument("recall_at_10").is_some());
1074    }
1075
1076    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1077    async fn register_instrument_duplicate_errors() {
1078        let mut c = Component::new(Labels::empty(), HashMap::new());
1079        c.register_instrument(
1080            "recall_at_10",
1081            InstrumentRef::Counter(new_counter("recall_at_10")),
1082        )
1083        .unwrap();
1084        let err = c
1085            .register_instrument(
1086                "recall_at_10",
1087                InstrumentRef::Counter(new_counter("recall_at_10")),
1088            )
1089            .unwrap_err();
1090        assert!(err.contains("duplicate family"), "wrong message: {err}");
1091        assert!(
1092            err.contains("recall_at_10"),
1093            "family name not in error: {err}"
1094        );
1095    }
1096
1097    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1098    async fn register_instrument_distinct_names_succeed() {
1099        let mut c = Component::new(Labels::empty(), HashMap::new());
1100        c.register_instrument("a", InstrumentRef::Counter(new_counter("a")))
1101            .unwrap();
1102        c.register_instrument("b", InstrumentRef::Counter(new_counter("b")))
1103            .unwrap();
1104        c.register_instrument("c", InstrumentRef::Counter(new_counter("c")))
1105            .unwrap();
1106        assert_eq!(c.instruments().len(), 3);
1107    }
1108
1109    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1110    async fn register_instrument_error_carries_label_context() {
1111        // SRD-40b §7's contract: the error message names the
1112        // dimensional cell so the workload author can see WHICH
1113        // op-template's label set produced the collision.
1114        let labels = Labels::of("phase", "pvs_query").with("op", "select_ann");
1115        let mut c = Component::new(labels, HashMap::new());
1116        c.register_instrument("overscan", InstrumentRef::Counter(new_counter("overscan")))
1117            .unwrap();
1118        let err = c
1119            .register_instrument("overscan", InstrumentRef::Counter(new_counter("overscan")))
1120            .unwrap_err();
1121        assert!(err.contains("phase"), "missing phase label: {err}");
1122        assert!(err.contains("pvs_query"), "missing phase value: {err}");
1123        assert!(err.contains("op"), "missing op label: {err}");
1124    }
1125
1126    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1127    async fn register_instrument_isolated_per_component() {
1128        // Two components — registering the same family on each
1129        // is OK; dimensional uniqueness comes from the
1130        // component-tree structure, not a global registry.
1131        let mut a = Component::new(Labels::of("op", "foo"), HashMap::new());
1132        let mut b = Component::new(Labels::of("op", "bar"), HashMap::new());
1133        assert!(
1134            a.register_instrument("overscan", InstrumentRef::Counter(new_counter("overscan")),)
1135                .is_ok()
1136        );
1137        assert!(
1138            b.register_instrument("overscan", InstrumentRef::Counter(new_counter("overscan")),)
1139                .is_ok()
1140        );
1141    }
1142
1143    // Test helper: register a counter that records a fixed value.
1144    fn install_counter(c: &mut Component, family: &str, value: u64) -> Arc<Counter> {
1145        let counter = new_counter(family);
1146        counter.inc_by(value);
1147        c.register_instrument(family, InstrumentRef::Counter(counter.clone()))
1148            .unwrap();
1149        counter
1150    }
1151
1152    // ── DynamicCapture hook ──
1153
1154    struct DynamicCounter {
1155        inner: AtomicU64,
1156    }
1157    impl DynamicCapture for DynamicCounter {
1158        fn capture_into(&self, out: &mut MetricSet, now: Instant, _drain: bool) {
1159            let v = self.inner.load(Ordering::Relaxed);
1160            out.insert_counter("dynamic_counter", Labels::default(), v, now);
1161        }
1162    }
1163
1164    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1165    async fn dynamic_capture_runs_after_registry() {
1166        let mut c = Component::new(Labels::empty(), HashMap::new());
1167        install_counter(&mut c, "static_counter", 5);
1168        c.set_dynamic_capture(Arc::new(DynamicCounter {
1169            inner: AtomicU64::new(7),
1170        }));
1171        let snap = c.capture_current();
1172        assert!(snap.family("static_counter").is_some());
1173        assert!(snap.family("dynamic_counter").is_some());
1174    }
1175
1176    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1177    async fn component_attach_computes_effective_labels() {
1178        let root = Component::root(Labels::of("session", "s1"), HashMap::new());
1179        let child = Arc::new(RwLock::new(Component::new(
1180            Labels::of("phase", "rampup"),
1181            HashMap::new(),
1182        )));
1183        attach(&root, &child);
1184
1185        let c = child.read().unwrap();
1186        let eff = c.effective_labels();
1187        assert_eq!(eff.get("session"), Some("s1"));
1188        assert_eq!(eff.get("phase"), Some("rampup"));
1189    }
1190
1191    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1192    async fn prop_walk_up_inheritance() {
1193        let mut root_props = HashMap::new();
1194        root_props.insert("hdr_digits".to_string(), "4".to_string());
1195        let root = Component::root(Labels::of("session", "s1"), root_props);
1196
1197        let child = Arc::new(RwLock::new(Component::new(
1198            Labels::of("phase", "rampup"),
1199            HashMap::new(),
1200        )));
1201        attach(&root, &child);
1202
1203        let c = child.read().unwrap();
1204        assert_eq!(c.get_prop("hdr_digits").as_deref(), Some("4"));
1205        assert_eq!(c.get_prop("nonexistent").as_deref(), None);
1206    }
1207
1208    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1209    async fn prop_child_overrides_parent() {
1210        let mut root_props = HashMap::new();
1211        root_props.insert("hdr_digits".to_string(), "3".to_string());
1212        let root = Component::root(Labels::of("session", "s1"), root_props);
1213
1214        let mut child_props = HashMap::new();
1215        child_props.insert("hdr_digits".to_string(), "4".to_string());
1216        let child = Arc::new(RwLock::new(Component::new(
1217            Labels::of("phase", "rampup"),
1218            child_props,
1219        )));
1220        attach(&root, &child);
1221
1222        let c = child.read().unwrap();
1223        assert_eq!(c.get_prop("hdr_digits").as_deref(), Some("4"));
1224    }
1225
1226    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1227    async fn detach_removes_child() {
1228        let root = Component::root(Labels::of("session", "s1"), HashMap::new());
1229        let child = Arc::new(RwLock::new(Component::new(
1230            Labels::of("phase", "rampup"),
1231            HashMap::new(),
1232        )));
1233        attach(&root, &child);
1234        assert_eq!(root.read().unwrap().child_count(), 1);
1235
1236        detach(&root, &child);
1237        assert_eq!(root.read().unwrap().child_count(), 0);
1238    }
1239
1240    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1241    async fn capture_tree_collects_running_components() {
1242        let root = Component::root(Labels::of("session", "s1"), HashMap::new());
1243
1244        // Running child with a registered counter.
1245        let child1 = Arc::new(RwLock::new(Component::new(
1246            Labels::of("phase", "load"),
1247            HashMap::new(),
1248        )));
1249        attach(&root, &child1);
1250        {
1251            let mut c = child1.write().unwrap();
1252            c.set_state(ComponentState::Running);
1253            install_counter(&mut c, "test_counter", 42);
1254        }
1255
1256        // Stopped child with a registered counter — must NOT be captured.
1257        let child2 = Arc::new(RwLock::new(Component::new(
1258            Labels::of("phase", "done"),
1259            HashMap::new(),
1260        )));
1261        attach(&root, &child2);
1262        {
1263            let mut c = child2.write().unwrap();
1264            c.set_state(ComponentState::Stopped);
1265            install_counter(&mut c, "test_counter", 99);
1266        }
1267
1268        let captured = capture_tree(&root, Duration::from_secs(1));
1269        assert_eq!(captured.len(), 1);
1270        assert_eq!(captured[0].0.get("phase"), Some("load"));
1271    }
1272
1273    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1274    async fn capture_tree_walks_nested_children() {
1275        let root = Component::root(Labels::of("session", "s1"), HashMap::new());
1276
1277        let scenario = Arc::new(RwLock::new(Component::new(
1278            Labels::of("scenario", "default"),
1279            HashMap::new(),
1280        )));
1281        attach(&root, &scenario);
1282        scenario.write().unwrap().set_state(ComponentState::Running);
1283
1284        let phase = Arc::new(RwLock::new(Component::new(
1285            Labels::of("phase", "search"),
1286            HashMap::new(),
1287        )));
1288        attach(&scenario, &phase);
1289        {
1290            let mut p = phase.write().unwrap();
1291            p.set_state(ComponentState::Running);
1292            install_counter(&mut p, "test_counter", 10);
1293        }
1294
1295        let captured = capture_tree(&root, Duration::from_secs(1));
1296        assert_eq!(captured.len(), 1);
1297        let eff = &captured[0].0;
1298        assert_eq!(eff.get("session"), Some("s1"));
1299        assert_eq!(eff.get("scenario"), Some("default"));
1300        assert_eq!(eff.get("phase"), Some("search"));
1301    }
1302
1303    // =====================================================================
1304    // Selector-based lookup (SRD 24)
1305    // =====================================================================
1306
1307    /// Fixture: a small session tree with two activity subtrees,
1308    /// each holding a handful of phases with distinct label shapes.
1309    ///
1310    /// Models the production label convention: each label is a
1311    /// condensed `(semantic, instance)` pair, so the KEY is the
1312    /// tier's kind (`session` / `activity` / `phase`) and the VALUE
1313    /// is its instance. No tier redeclares a name an ancestor owns
1314    /// (the label-ownership invariant `attach` enforces).
1315    fn sample_tree() -> Arc<RwLock<Component>> {
1316        let root = Component::root(
1317            Labels::empty().with("session", "test-session"),
1318            HashMap::new(),
1319        );
1320        // Activity A: rampup + two ann_query phases at different k.
1321        let activity_a = Arc::new(RwLock::new(Component::new(
1322            Labels::empty().with("activity", "a"),
1323            HashMap::new(),
1324        )));
1325        attach(&root, &activity_a);
1326        let rampup = Arc::new(RwLock::new(Component::new(
1327            Labels::empty()
1328                .with("phase", "rampup")
1329                .with("profile", "label_00"),
1330            HashMap::new(),
1331        )));
1332        attach(&activity_a, &rampup);
1333        for k in ["10", "100"] {
1334            let aq = Arc::new(RwLock::new(Component::new(
1335                Labels::empty()
1336                    .with("phase", "ann_query")
1337                    .with("profile", "label_00")
1338                    .with("k", k),
1339                HashMap::new(),
1340            )));
1341            attach(&activity_a, &aq);
1342        }
1343        // Activity B: one phase with a different profile shape.
1344        let activity_b = Arc::new(RwLock::new(Component::new(
1345            Labels::empty().with("activity", "b"),
1346            HashMap::new(),
1347        )));
1348        attach(&root, &activity_b);
1349        let teardown = Arc::new(RwLock::new(Component::new(
1350            Labels::empty()
1351                .with("phase", "teardown")
1352                .with("profile", "label_99"),
1353            HashMap::new(),
1354        )));
1355        attach(&activity_b, &teardown);
1356        root
1357    }
1358
1359    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1360    async fn find_returns_every_match_in_preorder() {
1361        let root = sample_tree();
1362        let sel = crate::selector::Selector::new().present("phase");
1363        let hits = find(&root, &sel);
1364        assert_eq!(hits.len(), 4);
1365        let names: Vec<String> = hits
1366            .iter()
1367            .filter_map(|c| {
1368                c.read()
1369                    .ok()
1370                    .and_then(|g| g.effective_labels().get("phase").map(|s| s.to_string()))
1371            })
1372            .collect();
1373        assert_eq!(names, vec!["rampup", "ann_query", "ann_query", "teardown"],);
1374    }
1375
1376    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1377    async fn find_with_empty_selector_returns_everything() {
1378        let root = sample_tree();
1379        let all = find(&root, &crate::selector::Selector::new());
1380        // session root + 2 activities + (rampup + 2 ann_query + teardown) = 7
1381        assert_eq!(all.len(), 7);
1382    }
1383
1384    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1385    async fn find_with_glob_and_eq_conjunction() {
1386        let root = sample_tree();
1387        let sel = crate::selector::Selector::new().glob("phase", "ann_*");
1388        let hits = find(&root, &sel);
1389        assert_eq!(hits.len(), 2);
1390        for h in &hits {
1391            let g = h.read().unwrap();
1392            assert_eq!(g.effective_labels().get("phase"), Some("ann_query"));
1393        }
1394    }
1395
1396    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1397    async fn find_with_present_and_absent_clauses() {
1398        let root = sample_tree();
1399        let with_k = find(
1400            &root,
1401            &crate::selector::Selector::new()
1402                .present("phase")
1403                .present("k"),
1404        );
1405        assert_eq!(with_k.len(), 2);
1406
1407        let without_k = find(
1408            &root,
1409            &crate::selector::Selector::new()
1410                .present("phase")
1411                .absent("k"),
1412        );
1413        assert_eq!(without_k.len(), 2);
1414    }
1415
1416    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1417    async fn find_one_exact_match() {
1418        let root = sample_tree();
1419        let sel = crate::selector::Selector::new().eq("phase", "rampup");
1420        let c = find_one(&root, &sel).unwrap();
1421        assert_eq!(
1422            c.read().unwrap().effective_labels().get("phase"),
1423            Some("rampup"),
1424        );
1425    }
1426
1427    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1428    async fn find_one_not_found() {
1429        let root = sample_tree();
1430        let sel = crate::selector::Selector::new().eq("phase", "nowhere");
1431        match find_one(&root, &sel) {
1432            Err(crate::selector::LookupError::NotFound) => {}
1433            Err(other) => panic!("expected NotFound, got {other:?}"),
1434            Ok(_) => panic!("expected NotFound, got a match"),
1435        }
1436    }
1437
1438    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1439    async fn find_one_ambiguous_reports_count() {
1440        let root = sample_tree();
1441        let sel = crate::selector::Selector::new().eq("phase", "ann_query");
1442        match find_one(&root, &sel) {
1443            Err(crate::selector::LookupError::Ambiguous { count }) => {
1444                assert_eq!(count, 2);
1445            }
1446            Err(other) => panic!("expected Ambiguous, got {other:?}"),
1447            Ok(_) => panic!("expected Ambiguous, got a single match"),
1448        }
1449    }
1450
1451    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1452    async fn any_short_circuits_on_first_hit() {
1453        let root = sample_tree();
1454        assert!(any(
1455            &root,
1456            &crate::selector::Selector::new().eq("phase", "rampup")
1457        ));
1458        assert!(!any(
1459            &root,
1460            &crate::selector::Selector::new().eq("phase", "zzz")
1461        ));
1462    }
1463
1464    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1465    async fn count_matches_len_of_find() {
1466        let root = sample_tree();
1467        let sel = crate::selector::Selector::new().present("phase");
1468        assert_eq!(count(&root, &sel), find(&root, &sel).len());
1469    }
1470
1471    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1472    async fn query_from_subtree_is_scoped() {
1473        let root = sample_tree();
1474        let activity_a = root.read().unwrap().children.first().unwrap().clone();
1475        let hits = find(
1476            &activity_a,
1477            &crate::selector::Selector::new().present("phase"),
1478        );
1479        assert_eq!(hits.len(), 3);
1480    }
1481
1482    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1483    async fn effective_labels_include_inherited_session_label() {
1484        let root = sample_tree();
1485        let sel = crate::selector::Selector::new()
1486            .eq("session", "test-session")
1487            .present("phase");
1488        assert_eq!(count(&root, &sel), 4);
1489    }
1490
1491    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1492    async fn selector_macro_drives_find() {
1493        let root = sample_tree();
1494        let hits = find(&root, &crate::selector!(phase = "teardown"));
1495        assert_eq!(hits.len(), 1);
1496    }
1497
1498    // =====================================================================
1499    // Controls on components (SRD 23)
1500    // =====================================================================
1501
1502    #[tokio::test]
1503    async fn controls_declare_and_lookup_through_component() {
1504        let root = Component::root(Labels::of("session", "s"), HashMap::new());
1505        {
1506            let guard = root.read().unwrap();
1507            guard
1508                .controls()
1509                .declare(crate::controls::ControlBuilder::new("concurrency", 16u32).build());
1510        }
1511        let c: crate::controls::Control<u32> = {
1512            let guard = root.read().unwrap();
1513            guard.controls().get::<u32>("concurrency").unwrap()
1514        };
1515        c.set(32, crate::controls::ControlOrigin::Test)
1516            .await
1517            .unwrap();
1518        let reread: crate::controls::Control<u32> = {
1519            let guard = root.read().unwrap();
1520            guard.controls().get::<u32>("concurrency").unwrap()
1521        };
1522        assert_eq!(reread.value(), 32);
1523    }
1524
1525    #[tokio::test]
1526    async fn reified_control_gauges_flow_through_capture_tree() {
1527        let root = Component::root(Labels::empty().with("session", "s1"), HashMap::new());
1528        let phase = Arc::new(RwLock::new(Component::new(
1529            Labels::empty().with("phase", "rampup"),
1530            HashMap::new(),
1531        )));
1532        attach(&root, &phase);
1533        {
1534            let p = phase.write().unwrap();
1535            p.controls().declare(
1536                crate::controls::ControlBuilder::new("concurrency", 8u32)
1537                    .reify_as_gauge(|v| Some(*v as f64))
1538                    .build(),
1539            );
1540        }
1541        phase.write().unwrap().set_state(ComponentState::Running);
1542
1543        let c: crate::controls::Control<u32> = phase
1544            .read()
1545            .unwrap()
1546            .controls()
1547            .get::<u32>("concurrency")
1548            .unwrap();
1549        c.set(64, crate::controls::ControlOrigin::Test)
1550            .await
1551            .unwrap();
1552
1553        let captured = capture_tree(&root, Duration::from_secs(1));
1554        let mut found_value: Option<f64> = None;
1555        for (labels, set) in &captured {
1556            if labels.get("phase") != Some("rampup") {
1557                continue;
1558            }
1559            if let Some(fam) = set.family("control_concurrency")
1560                && let Some(m) = fam.metrics().next()
1561            {
1562                if let Some(p) = m.point()
1563                    && let crate::snapshot::MetricValue::Gauge(g) = p.value()
1564                {
1565                    found_value = Some(g.value);
1566                }
1567                assert_eq!(m.labels().get("phase"), Some("rampup"));
1568                assert_eq!(m.labels().get("control"), Some("concurrency"));
1569            }
1570        }
1571        assert_eq!(found_value, Some(64.0));
1572
1573        let current = capture_tree_current(&root);
1574        let mut saw_via_current = false;
1575        for (_, set) in &current {
1576            if set.family("control_concurrency").is_some() {
1577                saw_via_current = true;
1578            }
1579        }
1580        assert!(saw_via_current);
1581    }
1582
1583    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1584    async fn dryrun_controls_enumeration_over_tree() {
1585        let root = sample_tree();
1586        let phase_hits = find(&root, &crate::selector::Selector::new().present("phase"));
1587        for (idx, phase) in phase_hits.iter().enumerate() {
1588            let guard = phase.read().unwrap();
1589            guard.controls().declare(
1590                crate::controls::ControlBuilder::new("concurrency", (10 * (idx + 1)) as u32)
1591                    .build(),
1592            );
1593        }
1594
1595        let mut entries: Vec<(String, String)> = Vec::new();
1596        for c in find(&root, &crate::selector::Selector::new()) {
1597            let guard = c.read().unwrap();
1598            let labels = guard.effective_labels().clone();
1599            for ctl in guard.controls().list() {
1600                entries.push((
1601                    format!("{}/{}", labels.get("phase").unwrap_or("-"), ctl.name(),),
1602                    ctl.value_string(),
1603                ));
1604            }
1605        }
1606
1607        assert_eq!(entries.len(), phase_hits.len());
1608        for (key, value) in &entries {
1609            assert!(key.ends_with("/concurrency"), "key = {key}");
1610            assert!(value.parse::<u32>().is_ok(), "value = {value}");
1611        }
1612    }
1613
1614    // ---- Branch-scoped control walk-up (SRD 23) ------------------
1615
1616    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1617    async fn branch_scope_subtree_resolves_from_descendant() {
1618        use crate::controls::{BranchScope, ControlBuilder};
1619        let root = Component::root(Labels::empty().with("session", "s1"), HashMap::new());
1620        let phase = Arc::new(RwLock::new(Component::new(
1621            Labels::empty().with("phase", "rampup"),
1622            HashMap::new(),
1623        )));
1624        attach(&root, &phase);
1625
1626        root.read().unwrap().controls().declare(
1627            ControlBuilder::new("hdr_sigdigs", 3u32)
1628                .branch_scope(BranchScope::Subtree)
1629                .build(),
1630        );
1631
1632        let resolved = phase.read().unwrap().find_control_up::<u32>("hdr_sigdigs");
1633        assert!(
1634            resolved.is_some(),
1635            "Subtree-scoped control should be visible to descendant"
1636        );
1637        assert_eq!(resolved.unwrap().value(), 3u32);
1638    }
1639
1640    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1641    async fn branch_scope_local_does_not_leak_to_descendants() {
1642        use crate::controls::{BranchScope, ControlBuilder};
1643        let root = Component::root(Labels::empty().with("session", "s1"), HashMap::new());
1644        let phase = Arc::new(RwLock::new(Component::new(
1645            Labels::empty().with("phase", "rampup"),
1646            HashMap::new(),
1647        )));
1648        attach(&root, &phase);
1649
1650        root.read().unwrap().controls().declare(
1651            ControlBuilder::new("private", 99u32)
1652                .branch_scope(BranchScope::Local)
1653                .build(),
1654        );
1655
1656        let leaked = phase.read().unwrap().find_control_up::<u32>("private");
1657        assert!(
1658            leaked.is_none(),
1659            "Local-scoped control must not be visible to descendants"
1660        );
1661    }
1662
1663    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1664    async fn nearest_declaration_wins_during_walk_up() {
1665        use crate::controls::{BranchScope, ControlBuilder};
1666        let root = Component::root(Labels::empty().with("session", "s1"), HashMap::new());
1667        let phase = Arc::new(RwLock::new(Component::new(
1668            Labels::empty().with("phase", "rampup"),
1669            HashMap::new(),
1670        )));
1671        attach(&root, &phase);
1672
1673        root.read().unwrap().controls().declare(
1674            ControlBuilder::new("hdr_sigdigs", 3u32)
1675                .branch_scope(BranchScope::Subtree)
1676                .build(),
1677        );
1678        phase
1679            .read()
1680            .unwrap()
1681            .controls()
1682            .declare(ControlBuilder::new("hdr_sigdigs", 5u32).build());
1683
1684        let v = phase
1685            .read()
1686            .unwrap()
1687            .find_control_up::<u32>("hdr_sigdigs")
1688            .unwrap()
1689            .value();
1690        assert_eq!(v, 5u32, "phase override should win over session default");
1691    }
1692
1693    // =====================================================================
1694    // SRD-40b §11 / SRD-42 §"Component lifecycle: scope_close flush"
1695    // =====================================================================
1696
1697    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1698    async fn component_scope_close_flushes_running_component_marks_partial_and_stops() {
1699        use crate::cadence::{CadenceTree, Cadences};
1700        use crate::cadence_reporter::CadenceReporter;
1701
1702        let tree = CadenceTree::plan_default(Cadences::new(&[Duration::from_secs(1)]).unwrap());
1703        let reporter = CadenceReporter::new(tree);
1704
1705        // Build a phase component with a registered counter holding N=42.
1706        let root = Component::root(Labels::of("session", "s1"), HashMap::new());
1707        let phase = Arc::new(RwLock::new(Component::new(
1708            Labels::of("phase", "short"),
1709            HashMap::new(),
1710        )));
1711        attach(&root, &phase);
1712        {
1713            let mut p = phase.write().unwrap();
1714            p.set_state(ComponentState::Running);
1715            install_counter(&mut p, "test_counter", 42);
1716        }
1717
1718        scope_close(&phase, &reporter, Duration::from_millis(150));
1719        reporter.flush_for_tests();
1720
1721        // Component is now Stopped — second call must be a no-op.
1722        assert_eq!(phase.read().unwrap().state(), ComponentState::Stopped);
1723        scope_close(&phase, &reporter, Duration::from_millis(150));
1724        reporter.flush_for_tests();
1725
1726        let labels = phase.read().unwrap().effective_labels().clone();
1727        let latest = reporter
1728            .latest(&labels, Duration::from_secs(1))
1729            .expect("scope_close must publish the partial");
1730        assert!(latest.is_partial(), "snapshot must be marked partial");
1731        let f = latest
1732            .family("test_counter")
1733            .expect("test_counter family present");
1734        let m = f.metrics().next().unwrap();
1735        match m.point().unwrap().value() {
1736            crate::snapshot::MetricValue::Counter(c) => assert_eq!(c.cumulative, 42),
1737            v => panic!("expected counter, got {v:?}"),
1738        }
1739    }
1740
1741    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1742    async fn component_scope_close_skips_non_running_states() {
1743        use crate::cadence::{CadenceTree, Cadences};
1744        use crate::cadence_reporter::CadenceReporter;
1745
1746        let tree = CadenceTree::plan_default(Cadences::new(&[Duration::from_secs(1)]).unwrap());
1747        let reporter = CadenceReporter::new(tree);
1748
1749        let root = Component::root(Labels::of("session", "s1"), HashMap::new());
1750        let phase = Arc::new(RwLock::new(Component::new(
1751            Labels::of("phase", "starting"),
1752            HashMap::new(),
1753        )));
1754        attach(&root, &phase);
1755        assert_eq!(phase.read().unwrap().state(), ComponentState::Starting);
1756        {
1757            let mut p = phase.write().unwrap();
1758            install_counter(&mut p, "test_counter", 99);
1759        }
1760
1761        scope_close(&phase, &reporter, Duration::from_millis(150));
1762        reporter.flush_for_tests();
1763
1764        assert_eq!(phase.read().unwrap().state(), ComponentState::Starting);
1765        let labels = phase.read().unwrap().effective_labels().clone();
1766        assert!(
1767            reporter.latest(&labels, Duration::from_secs(1)).is_none(),
1768            "scope_close on a non-Running component must not publish"
1769        );
1770    }
1771
1772    // ── capture_delta_auto: real-elapsed interval ────────────
1773
1774    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1775    async fn capture_delta_auto_uses_fallback_on_first_call() {
1776        // No prior capture → fallback is the recorded interval.
1777        // Same shape the executor's phase-end flush sees on the
1778        // edge case where a phase ends before any scheduler tick.
1779        let c = Component::new(Labels::empty(), HashMap::new());
1780        let s = c.capture_delta_auto(Duration::from_millis(500));
1781        assert_eq!(s.interval(), Duration::from_millis(500));
1782    }
1783
1784    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1785    async fn capture_delta_auto_measures_real_elapsed_after_prior_capture() {
1786        // After a `capture_delta` records the watermark,
1787        // `capture_delta_auto` reports the actual wall-clock
1788        // delta between the two — NOT the prior `interval`
1789        // argument. This is what kills the 1s quantization in
1790        // the phase-end flush path: a phase-end auto-capture
1791        // ~80ms after the last scheduler tick stamps ~80ms,
1792        // not the nominal 1000ms fallback.
1793        let c = Component::new(Labels::empty(), HashMap::new());
1794        let _ = c.capture_delta(Duration::from_secs(1));
1795        std::thread::sleep(Duration::from_millis(80));
1796        let s = c.capture_delta_auto(Duration::from_secs(1));
1797        // Allow generous lower / upper bounds — CI machines
1798        // can be loaded — but anything within (60ms, 500ms) is
1799        // unambiguously distinct from the 1000ms fallback.
1800        assert!(
1801            s.interval() > Duration::from_millis(60),
1802            "interval should reflect real ~80ms elapsed, got {:?}",
1803            s.interval()
1804        );
1805        assert!(
1806            s.interval() < Duration::from_millis(500),
1807            "interval should be the real elapsed, not the 1s fallback: {:?}",
1808            s.interval()
1809        );
1810    }
1811
1812    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1813    async fn capture_delta_auto_chained_uses_inter_capture_elapsed() {
1814        // Two consecutive auto-captures: the second sees the
1815        // elapsed between auto calls, not the cumulative
1816        // since component creation.
1817        let c = Component::new(Labels::empty(), HashMap::new());
1818        let _first = c.capture_delta_auto(Duration::from_secs(1));
1819        std::thread::sleep(Duration::from_millis(50));
1820        let second = c.capture_delta_auto(Duration::from_secs(1));
1821        assert!(
1822            second.interval() < Duration::from_millis(500),
1823            "second auto-capture should measure inter-capture \
1824             elapsed, not cumulative: {:?}",
1825            second.interval()
1826        );
1827    }
1828}