use std::collections::HashMap;
use std::sync::{Arc, Mutex, RwLock, Weak};
use std::time::{Duration, Instant};
use crate::instruments::counter::Counter;
use crate::instruments::gauge::ValueGauge;
use crate::instruments::histogram::Histogram;
use crate::instruments::timer::Timer;
use crate::labels::Labels;
use crate::snapshot::MetricSet;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ComponentState {
Starting,
Running,
Stopping,
Stopped,
}
#[derive(Clone)]
pub enum InstrumentRef {
Counter(Arc<Counter>),
Gauge(Arc<ValueGauge>),
Histogram(Arc<Histogram>),
Timer(Arc<Timer>),
}
impl InstrumentRef {
pub fn labels(&self) -> &Labels {
match self {
Self::Counter(c) => c.labels(),
Self::Gauge(g) => g.labels(),
Self::Histogram(h) => h.labels(),
Self::Timer(t) => t.labels(),
}
}
}
pub struct RegisteredInstrument {
pub family: String,
pub unit: Option<String>,
pub instrument: InstrumentRef,
}
pub trait DynamicCapture: Send + Sync {
fn capture_into(&self, out: &mut MetricSet, now: Instant, drain: bool);
}
pub struct Component {
labels: Labels,
effective_labels: Labels,
props: HashMap<String, String>,
parent: Option<Weak<RwLock<Component>>>,
children: Vec<Arc<RwLock<Component>>>,
state: ComponentState,
instruments: Vec<RegisteredInstrument>,
dynamic_capture: Option<Arc<dyn DynamicCapture>>,
last_capture_instant: Mutex<Option<Instant>>,
controls: crate::controls::ControlRegistry,
cells: std::sync::Arc<crate::cells::CellMap>,
live: Option<std::sync::Arc<()>>,
live_children: std::collections::HashMap<String, Vec<std::sync::Weak<()>>>,
}
impl Component {
pub fn new(labels: Labels, props: HashMap<String, String>) -> Self {
Self {
effective_labels: labels.clone(),
labels,
props,
parent: None,
children: Vec::new(),
state: ComponentState::Starting,
instruments: Vec::new(),
dynamic_capture: None,
last_capture_instant: Mutex::new(None),
controls: crate::controls::ControlRegistry::new(),
cells: std::sync::Arc::new(crate::cells::CellMap::new()),
live: Some(std::sync::Arc::new(())),
live_children: std::collections::HashMap::new(),
}
}
pub fn root(labels: Labels, props: HashMap<String, String>) -> Arc<RwLock<Self>> {
let mut component = Self::new(labels, props);
component.state = ComponentState::Running;
Arc::new(RwLock::new(component))
}
pub fn labels(&self) -> &Labels {
&self.labels
}
pub fn effective_labels(&self) -> &Labels {
&self.effective_labels
}
pub fn state(&self) -> ComponentState {
self.state
}
pub fn set_state(&mut self, state: ComponentState) {
self.state = state;
if state == ComponentState::Stopped {
self.live = None;
}
}
fn live_handle(&self) -> std::sync::Weak<()> {
match &self.live {
Some(t) => std::sync::Arc::downgrade(t),
None => std::sync::Weak::new(),
}
}
pub fn register_instrument(
&mut self,
family: impl Into<String>,
instrument: InstrumentRef,
) -> Result<(), String> {
self.register_instrument_with_unit(family, None, instrument)
}
pub fn register_instrument_with_unit(
&mut self,
family: impl Into<String>,
unit: Option<String>,
instrument: InstrumentRef,
) -> Result<(), String> {
let family = family.into();
if self.instruments.iter().any(|ri| ri.family == family) {
return Err(format!(
"duplicate family name on dimensionally-same metric \
context: {family}{}",
self.effective_labels.to_prometheus(),
));
}
self.instruments.push(RegisteredInstrument {
family,
unit,
instrument,
});
Ok(())
}
pub fn instruments(&self) -> &[RegisteredInstrument] {
&self.instruments
}
pub fn find_instrument(&self, family: &str) -> Option<&InstrumentRef> {
self.instruments
.iter()
.find(|ri| ri.family == family)
.map(|ri| &ri.instrument)
}
pub fn set_dynamic_capture(&mut self, hook: Arc<dyn DynamicCapture>) {
self.dynamic_capture = Some(hook);
}
pub fn capture_delta(&self, interval: Duration) -> MetricSet {
let now = Instant::now();
let mut out = MetricSet::at(now, interval);
self.capture_registry_into(&mut out, now, true);
if let Some(hook) = &self.dynamic_capture {
hook.capture_into(&mut out, now, true);
}
*self
.last_capture_instant
.lock()
.unwrap_or_else(|e| e.into_inner()) = Some(now);
out
}
pub fn capture_delta_auto(&self, fallback: Duration) -> MetricSet {
let now = Instant::now();
let interval = {
let mut prev = self
.last_capture_instant
.lock()
.unwrap_or_else(|e| e.into_inner());
let elapsed = prev.map(|t| now.duration_since(t));
*prev = Some(now);
elapsed.unwrap_or(fallback)
};
let mut out = MetricSet::at(now, interval);
self.capture_registry_into(&mut out, now, true);
if let Some(hook) = &self.dynamic_capture {
hook.capture_into(&mut out, now, true);
}
out
}
pub fn capture_current(&self) -> MetricSet {
let now = Instant::now();
let mut out = MetricSet::at(now, Duration::ZERO);
self.capture_registry_into(&mut out, now, false);
if let Some(hook) = &self.dynamic_capture {
hook.capture_into(&mut out, now, false);
}
out
}
fn capture_registry_into(&self, out: &mut MetricSet, now: Instant, drain: bool) {
for ri in &self.instruments {
let family = ri.family.clone();
let unit = ri.unit.as_deref();
match &ri.instrument {
InstrumentRef::Counter(c) => {
let lbl = strip_name_label(c.labels());
out.insert_counter_with_unit(family, unit, lbl, c.get(), now);
}
InstrumentRef::Gauge(g) => {
let lbl = strip_name_label(g.labels());
out.insert_gauge_with_unit(family, unit, lbl, g.get(), now);
}
InstrumentRef::Histogram(h) => {
let lbl = strip_name_label(h.labels());
let reservoir = if drain {
h.snapshot()
} else {
h.peek_snapshot()
};
out.insert_histogram_with_unit_cumulative(
family,
unit,
lbl,
reservoir,
h.total(),
now,
);
}
InstrumentRef::Timer(t) => {
let lbl = strip_name_label(t.labels());
let snap = if drain {
t.snapshot()
} else {
t.peek_snapshot()
};
out.insert_histogram_with_unit_cumulative(
family,
unit,
lbl,
snap.histogram,
snap.count,
now,
);
}
}
}
}
pub fn get_prop(&self, name: &str) -> Option<String> {
if let Some(value) = self.props.get(name) {
return Some(value.clone());
}
if let Some(ref parent_weak) = self.parent
&& let Some(parent_arc) = parent_weak.upgrade()
&& let Ok(parent) = parent_arc.read()
{
return parent.get_prop(name);
}
None
}
pub fn set_prop(&mut self, name: &str, value: &str) {
self.props.insert(name.to_string(), value.to_string());
}
pub fn child_count(&self) -> usize {
self.children.len()
}
pub fn children(&self) -> impl Iterator<Item = &Arc<RwLock<Component>>> {
self.children.iter()
}
pub fn controls(&self) -> &crate::controls::ControlRegistry {
&self.controls
}
pub fn cells(&self) -> std::sync::Arc<crate::cells::CellMap> {
self.cells.clone()
}
pub fn find_control_up<T>(&self, name: &str) -> Option<crate::controls::Control<T>>
where
T: Clone + Send + Sync + 'static,
{
if let Some(c) = self.controls.get::<T>(name) {
return Some(c);
}
if let Some(ref parent_weak) = self.parent
&& let Some(parent_arc) = parent_weak.upgrade()
&& let Ok(parent) = parent_arc.read()
{
return parent.find_control_up_subtree::<T>(name);
}
None
}
fn find_control_up_subtree<T>(&self, name: &str) -> Option<crate::controls::Control<T>>
where
T: Clone + Send + Sync + 'static,
{
if let Some(erased) = self.controls.get_erased(name)
&& erased.branch_scope() == crate::controls::BranchScope::Subtree
&& let Some(c) = self.controls.get::<T>(name)
{
return Some(c);
}
if let Some(ref parent_weak) = self.parent
&& let Some(parent_arc) = parent_weak.upgrade()
&& let Ok(parent) = parent_arc.read()
{
return parent.find_control_up_subtree::<T>(name);
}
None
}
pub fn find_control_erased_up(
&self,
name: &str,
) -> Option<std::sync::Arc<dyn crate::controls::ErasedControl>> {
if let Some(erased) = self.controls.get_erased(name) {
return Some(erased);
}
if let Some(ref parent_weak) = self.parent
&& let Some(parent_arc) = parent_weak.upgrade()
&& let Ok(parent) = parent_arc.read()
{
return parent.find_control_erased_up_subtree(name);
}
None
}
fn find_control_erased_up_subtree(
&self,
name: &str,
) -> Option<std::sync::Arc<dyn crate::controls::ErasedControl>> {
if let Some(erased) = self.controls.get_erased(name)
&& erased.branch_scope() == crate::controls::BranchScope::Subtree
{
return Some(erased);
}
if let Some(ref parent_weak) = self.parent
&& let Some(parent_arc) = parent_weak.upgrade()
&& let Ok(parent) = parent_arc.read()
{
return parent.find_control_erased_up_subtree(name);
}
None
}
pub fn control_snapshot(
start: &std::sync::Arc<std::sync::RwLock<Component>>,
) -> std::collections::HashMap<String, std::sync::Arc<dyn crate::controls::ErasedControl>> {
let mut map: std::collections::HashMap<
String,
std::sync::Arc<dyn crate::controls::ErasedControl>,
> = std::collections::HashMap::new();
let mut next = Some(start.clone());
let mut is_start = true;
while let Some(arc) = next {
let parent_next;
{
let g = arc.read().unwrap_or_else(|e| e.into_inner());
for handle in g.controls.list() {
if is_start || handle.branch_scope() == crate::controls::BranchScope::Subtree {
map.entry(handle.name().to_string()).or_insert(handle);
}
}
parent_next = g.parent.as_ref().and_then(|w| w.upgrade());
}
next = parent_next;
is_start = false;
}
map
}
pub fn running_descendant_count(&self) -> usize {
let mut count = 0;
for child in &self.children {
if let Ok(c) = child.read() {
if c.state == ComponentState::Running {
count += 1;
}
count += c.running_descendant_count();
}
}
count
}
}
fn strip_name_label(labels: &Labels) -> Labels {
let mut out = Labels::default();
for (k, v) in labels.iter() {
if k != "name" {
out = out.with(k, v);
}
}
out
}
pub fn scope_close(
component: &Arc<RwLock<Component>>,
cadence_reporter: &crate::cadence_reporter::CadenceReporter,
interval: Duration,
) {
let (labels, delta) = {
let g = component.read().unwrap_or_else(|e| e.into_inner());
if g.state != ComponentState::Running {
return;
}
let delta = g.capture_delta(interval);
(g.effective_labels.clone(), delta)
};
cadence_reporter.scope_close(&labels, delta);
let mut g = component.write().unwrap_or_else(|e| e.into_inner());
g.state = ComponentState::Stopped;
}
pub fn find(
root: &Arc<RwLock<Component>>,
sel: &crate::selector::Selector,
) -> Vec<Arc<RwLock<Component>>> {
let mut out = Vec::new();
find_into(root, sel, &mut out);
out
}
fn find_into(
root: &Arc<RwLock<Component>>,
sel: &crate::selector::Selector,
out: &mut Vec<Arc<RwLock<Component>>>,
) {
let Ok(guard) = root.read() else { return };
if sel.matches(&guard.effective_labels) {
out.push(root.clone());
}
let children = guard.children.clone();
drop(guard);
for child in &children {
find_into(child, sel, out);
}
}
pub fn find_one(
root: &Arc<RwLock<Component>>,
sel: &crate::selector::Selector,
) -> Result<Arc<RwLock<Component>>, crate::selector::LookupError> {
let mut first: Option<Arc<RwLock<Component>>> = None;
let mut count = 0usize;
find_one_walk(root, sel, &mut first, &mut count);
match first {
None => Err(crate::selector::LookupError::NotFound),
Some(c) if count == 1 => Ok(c),
Some(_) => Err(crate::selector::LookupError::Ambiguous { count }),
}
}
fn find_one_walk(
root: &Arc<RwLock<Component>>,
sel: &crate::selector::Selector,
first: &mut Option<Arc<RwLock<Component>>>,
count: &mut usize,
) {
let Ok(guard) = root.read() else { return };
if sel.matches(&guard.effective_labels) {
*count += 1;
if first.is_none() {
*first = Some(root.clone());
}
}
let children = guard.children.clone();
drop(guard);
for child in &children {
find_one_walk(child, sel, first, count);
}
}
pub fn any(root: &Arc<RwLock<Component>>, sel: &crate::selector::Selector) -> bool {
any_walk(root, sel)
}
fn any_walk(root: &Arc<RwLock<Component>>, sel: &crate::selector::Selector) -> bool {
let Ok(guard) = root.read() else { return false };
if sel.matches(&guard.effective_labels) {
return true;
}
let children = guard.children.clone();
drop(guard);
children.iter().any(|c| any_walk(c, sel))
}
pub fn count(root: &Arc<RwLock<Component>>, sel: &crate::selector::Selector) -> usize {
let mut n = 0usize;
count_walk(root, sel, &mut n);
n
}
fn count_walk(root: &Arc<RwLock<Component>>, sel: &crate::selector::Selector, n: &mut usize) {
let Ok(guard) = root.read() else { return };
if sel.matches(&guard.effective_labels) {
*n += 1;
}
let children = guard.children.clone();
drop(guard);
for child in &children {
count_walk(child, sel, n);
}
}
pub fn attach(parent: &Arc<RwLock<Component>>, child: &Arc<RwLock<Component>>) {
let parent_effective = {
let p = parent.read().unwrap_or_else(|e| e.into_inner());
p.effective_labels.clone()
};
let mut c = child.write().unwrap_or_else(|e| e.into_inner());
if let Some((k, _)) = c
.labels
.iter()
.find(|(k, _)| parent_effective.get(k).is_some())
{
panic!(
"component label-ownership violation: child re-declares label `{k}` \
already owned by an ancestor. Each label name must be set on exactly \
one tier and inherited downward (child {}, ancestors {}).",
c.labels.to_prometheus(),
parent_effective.to_prometheus(),
);
}
c.effective_labels = parent_effective.extend(&c.labels);
c.parent = Some(Arc::downgrade(parent));
let child_own = c.labels.to_prometheus();
let child_live = c.live_handle();
drop(c);
let mut p = parent.write().unwrap_or_else(|e| e.into_inner());
let bucket = p.live_children.entry(child_own.clone()).or_default();
bucket.retain(|w| w.strong_count() > 0);
if !bucket.is_empty() {
let parent_labels = parent_effective.to_prometheus();
panic!(
"component sibling-identity violation: a LIVE sibling already \
declares the own-label set {child_own} under {parent_labels}. Both \
would compose byte-identical effective labels, so the same family \
registered on each yields two instruments sharing one metric \
identity — which the per-component duplicate-family check cannot \
see, because they are different components. (Re-using a label set \
AFTER the previous component stops is fine and is not this.)"
);
}
bucket.push(child_live);
p.children.push(child.clone());
}
pub fn detach(parent: &Arc<RwLock<Component>>, child: &Arc<RwLock<Component>>) {
let mut p = parent.write().unwrap_or_else(|e| e.into_inner());
p.children.retain(|c| !Arc::ptr_eq(c, child));
let mut c = child.write().unwrap_or_else(|e| e.into_inner());
c.parent = None;
}
pub fn capture_tree(root: &Arc<RwLock<Component>>, interval: Duration) -> Vec<(Labels, MetricSet)> {
let mut results = Vec::new();
capture_recursive(root, interval, &mut results);
results
}
fn capture_recursive(
node: &Arc<RwLock<Component>>,
interval: Duration,
results: &mut Vec<(Labels, MetricSet)>,
) {
let Ok(guard) = node.read() else { return };
let state = guard.state;
let effective_labels = guard.effective_labels.clone();
let children = guard.children.clone();
if state == ComponentState::Running {
let snapshot = guard.capture_delta(interval);
if !snapshot.is_empty() {
results.push((effective_labels.clone(), snapshot));
}
let control_gauges = guard
.controls
.snapshot_gauges(&effective_labels, Instant::now());
if !control_gauges.is_empty() {
results.push((effective_labels, control_gauges));
}
}
drop(guard);
for child in &children {
capture_recursive(child, interval, results);
}
}
pub fn capture_tree_current(root: &Arc<RwLock<Component>>) -> Vec<(Labels, MetricSet)> {
let mut results = Vec::new();
capture_current_recursive(root, &mut results);
results
}
fn capture_current_recursive(
node: &Arc<RwLock<Component>>,
results: &mut Vec<(Labels, MetricSet)>,
) {
let Ok(guard) = node.read() else { return };
let state = guard.state;
let effective_labels = guard.effective_labels.clone();
let children = guard.children.clone();
if state == ComponentState::Running {
let snapshot = guard.capture_current();
if !snapshot.is_empty() {
results.push((effective_labels.clone(), snapshot));
}
let control_gauges = guard
.controls
.snapshot_gauges(&effective_labels, Instant::now());
if !control_gauges.is_empty() {
results.push((effective_labels, control_gauges));
}
}
drop(guard);
for child in &children {
capture_current_recursive(child, results);
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicU64, Ordering};
fn new_counter(family: &str) -> Arc<Counter> {
Arc::new(Counter::new(Labels::of("name", family)))
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn register_instrument_first_time_succeeds() {
let mut c = Component::new(Labels::empty(), HashMap::new());
assert!(
c.register_instrument(
"recall_at_10",
InstrumentRef::Counter(new_counter("recall_at_10")),
)
.is_ok()
);
assert!(c.find_instrument("recall_at_10").is_some());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn register_instrument_duplicate_errors() {
let mut c = Component::new(Labels::empty(), HashMap::new());
c.register_instrument(
"recall_at_10",
InstrumentRef::Counter(new_counter("recall_at_10")),
)
.unwrap();
let err = c
.register_instrument(
"recall_at_10",
InstrumentRef::Counter(new_counter("recall_at_10")),
)
.unwrap_err();
assert!(err.contains("duplicate family"), "wrong message: {err}");
assert!(
err.contains("recall_at_10"),
"family name not in error: {err}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn register_instrument_distinct_names_succeed() {
let mut c = Component::new(Labels::empty(), HashMap::new());
c.register_instrument("a", InstrumentRef::Counter(new_counter("a")))
.unwrap();
c.register_instrument("b", InstrumentRef::Counter(new_counter("b")))
.unwrap();
c.register_instrument("c", InstrumentRef::Counter(new_counter("c")))
.unwrap();
assert_eq!(c.instruments().len(), 3);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn register_instrument_error_carries_label_context() {
let labels = Labels::of("phase", "pvs_query").with("op", "select_ann");
let mut c = Component::new(labels, HashMap::new());
c.register_instrument("overscan", InstrumentRef::Counter(new_counter("overscan")))
.unwrap();
let err = c
.register_instrument("overscan", InstrumentRef::Counter(new_counter("overscan")))
.unwrap_err();
assert!(err.contains("phase"), "missing phase label: {err}");
assert!(err.contains("pvs_query"), "missing phase value: {err}");
assert!(err.contains("op"), "missing op label: {err}");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn register_instrument_isolated_per_component() {
let mut a = Component::new(Labels::of("op", "foo"), HashMap::new());
let mut b = Component::new(Labels::of("op", "bar"), HashMap::new());
assert!(
a.register_instrument("overscan", InstrumentRef::Counter(new_counter("overscan")),)
.is_ok()
);
assert!(
b.register_instrument("overscan", InstrumentRef::Counter(new_counter("overscan")),)
.is_ok()
);
}
fn install_counter(c: &mut Component, family: &str, value: u64) -> Arc<Counter> {
let counter = new_counter(family);
counter.inc_by(value);
c.register_instrument(family, InstrumentRef::Counter(counter.clone()))
.unwrap();
counter
}
struct DynamicCounter {
inner: AtomicU64,
}
impl DynamicCapture for DynamicCounter {
fn capture_into(&self, out: &mut MetricSet, now: Instant, _drain: bool) {
let v = self.inner.load(Ordering::Relaxed);
out.insert_counter("dynamic_counter", Labels::default(), v, now);
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn dynamic_capture_runs_after_registry() {
let mut c = Component::new(Labels::empty(), HashMap::new());
install_counter(&mut c, "static_counter", 5);
c.set_dynamic_capture(Arc::new(DynamicCounter {
inner: AtomicU64::new(7),
}));
let snap = c.capture_current();
assert!(snap.family("static_counter").is_some());
assert!(snap.family("dynamic_counter").is_some());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn component_attach_computes_effective_labels() {
let root = Component::root(Labels::of("session", "s1"), HashMap::new());
let child = Arc::new(RwLock::new(Component::new(
Labels::of("phase", "rampup"),
HashMap::new(),
)));
attach(&root, &child);
let c = child.read().unwrap();
let eff = c.effective_labels();
assert_eq!(eff.get("session"), Some("s1"));
assert_eq!(eff.get("phase"), Some("rampup"));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn prop_walk_up_inheritance() {
let mut root_props = HashMap::new();
root_props.insert("hdr_digits".to_string(), "4".to_string());
let root = Component::root(Labels::of("session", "s1"), root_props);
let child = Arc::new(RwLock::new(Component::new(
Labels::of("phase", "rampup"),
HashMap::new(),
)));
attach(&root, &child);
let c = child.read().unwrap();
assert_eq!(c.get_prop("hdr_digits").as_deref(), Some("4"));
assert_eq!(c.get_prop("nonexistent").as_deref(), None);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn prop_child_overrides_parent() {
let mut root_props = HashMap::new();
root_props.insert("hdr_digits".to_string(), "3".to_string());
let root = Component::root(Labels::of("session", "s1"), root_props);
let mut child_props = HashMap::new();
child_props.insert("hdr_digits".to_string(), "4".to_string());
let child = Arc::new(RwLock::new(Component::new(
Labels::of("phase", "rampup"),
child_props,
)));
attach(&root, &child);
let c = child.read().unwrap();
assert_eq!(c.get_prop("hdr_digits").as_deref(), Some("4"));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn detach_removes_child() {
let root = Component::root(Labels::of("session", "s1"), HashMap::new());
let child = Arc::new(RwLock::new(Component::new(
Labels::of("phase", "rampup"),
HashMap::new(),
)));
attach(&root, &child);
assert_eq!(root.read().unwrap().child_count(), 1);
detach(&root, &child);
assert_eq!(root.read().unwrap().child_count(), 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn capture_tree_collects_running_components() {
let root = Component::root(Labels::of("session", "s1"), HashMap::new());
let child1 = Arc::new(RwLock::new(Component::new(
Labels::of("phase", "load"),
HashMap::new(),
)));
attach(&root, &child1);
{
let mut c = child1.write().unwrap();
c.set_state(ComponentState::Running);
install_counter(&mut c, "test_counter", 42);
}
let child2 = Arc::new(RwLock::new(Component::new(
Labels::of("phase", "done"),
HashMap::new(),
)));
attach(&root, &child2);
{
let mut c = child2.write().unwrap();
c.set_state(ComponentState::Stopped);
install_counter(&mut c, "test_counter", 99);
}
let captured = capture_tree(&root, Duration::from_secs(1));
assert_eq!(captured.len(), 1);
assert_eq!(captured[0].0.get("phase"), Some("load"));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn capture_tree_walks_nested_children() {
let root = Component::root(Labels::of("session", "s1"), HashMap::new());
let scenario = Arc::new(RwLock::new(Component::new(
Labels::of("scenario", "default"),
HashMap::new(),
)));
attach(&root, &scenario);
scenario.write().unwrap().set_state(ComponentState::Running);
let phase = Arc::new(RwLock::new(Component::new(
Labels::of("phase", "search"),
HashMap::new(),
)));
attach(&scenario, &phase);
{
let mut p = phase.write().unwrap();
p.set_state(ComponentState::Running);
install_counter(&mut p, "test_counter", 10);
}
let captured = capture_tree(&root, Duration::from_secs(1));
assert_eq!(captured.len(), 1);
let eff = &captured[0].0;
assert_eq!(eff.get("session"), Some("s1"));
assert_eq!(eff.get("scenario"), Some("default"));
assert_eq!(eff.get("phase"), Some("search"));
}
fn sample_tree() -> Arc<RwLock<Component>> {
let root = Component::root(
Labels::empty().with("session", "test-session"),
HashMap::new(),
);
let activity_a = Arc::new(RwLock::new(Component::new(
Labels::empty().with("activity", "a"),
HashMap::new(),
)));
attach(&root, &activity_a);
let rampup = Arc::new(RwLock::new(Component::new(
Labels::empty()
.with("phase", "rampup")
.with("profile", "label_00"),
HashMap::new(),
)));
attach(&activity_a, &rampup);
for k in ["10", "100"] {
let aq = Arc::new(RwLock::new(Component::new(
Labels::empty()
.with("phase", "ann_query")
.with("profile", "label_00")
.with("k", k),
HashMap::new(),
)));
attach(&activity_a, &aq);
}
let activity_b = Arc::new(RwLock::new(Component::new(
Labels::empty().with("activity", "b"),
HashMap::new(),
)));
attach(&root, &activity_b);
let teardown = Arc::new(RwLock::new(Component::new(
Labels::empty()
.with("phase", "teardown")
.with("profile", "label_99"),
HashMap::new(),
)));
attach(&activity_b, &teardown);
root
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn find_returns_every_match_in_preorder() {
let root = sample_tree();
let sel = crate::selector::Selector::new().present("phase");
let hits = find(&root, &sel);
assert_eq!(hits.len(), 4);
let names: Vec<String> = hits
.iter()
.filter_map(|c| {
c.read()
.ok()
.and_then(|g| g.effective_labels().get("phase").map(|s| s.to_string()))
})
.collect();
assert_eq!(names, vec!["rampup", "ann_query", "ann_query", "teardown"],);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn find_with_empty_selector_returns_everything() {
let root = sample_tree();
let all = find(&root, &crate::selector::Selector::new());
assert_eq!(all.len(), 7);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn find_with_glob_and_eq_conjunction() {
let root = sample_tree();
let sel = crate::selector::Selector::new().glob("phase", "ann_*");
let hits = find(&root, &sel);
assert_eq!(hits.len(), 2);
for h in &hits {
let g = h.read().unwrap();
assert_eq!(g.effective_labels().get("phase"), Some("ann_query"));
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn find_with_present_and_absent_clauses() {
let root = sample_tree();
let with_k = find(
&root,
&crate::selector::Selector::new()
.present("phase")
.present("k"),
);
assert_eq!(with_k.len(), 2);
let without_k = find(
&root,
&crate::selector::Selector::new()
.present("phase")
.absent("k"),
);
assert_eq!(without_k.len(), 2);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn find_one_exact_match() {
let root = sample_tree();
let sel = crate::selector::Selector::new().eq("phase", "rampup");
let c = find_one(&root, &sel).unwrap();
assert_eq!(
c.read().unwrap().effective_labels().get("phase"),
Some("rampup"),
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn find_one_not_found() {
let root = sample_tree();
let sel = crate::selector::Selector::new().eq("phase", "nowhere");
match find_one(&root, &sel) {
Err(crate::selector::LookupError::NotFound) => {}
Err(other) => panic!("expected NotFound, got {other:?}"),
Ok(_) => panic!("expected NotFound, got a match"),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn find_one_ambiguous_reports_count() {
let root = sample_tree();
let sel = crate::selector::Selector::new().eq("phase", "ann_query");
match find_one(&root, &sel) {
Err(crate::selector::LookupError::Ambiguous { count }) => {
assert_eq!(count, 2);
}
Err(other) => panic!("expected Ambiguous, got {other:?}"),
Ok(_) => panic!("expected Ambiguous, got a single match"),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn any_short_circuits_on_first_hit() {
let root = sample_tree();
assert!(any(
&root,
&crate::selector::Selector::new().eq("phase", "rampup")
));
assert!(!any(
&root,
&crate::selector::Selector::new().eq("phase", "zzz")
));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn count_matches_len_of_find() {
let root = sample_tree();
let sel = crate::selector::Selector::new().present("phase");
assert_eq!(count(&root, &sel), find(&root, &sel).len());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn query_from_subtree_is_scoped() {
let root = sample_tree();
let activity_a = root.read().unwrap().children.first().unwrap().clone();
let hits = find(
&activity_a,
&crate::selector::Selector::new().present("phase"),
);
assert_eq!(hits.len(), 3);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn effective_labels_include_inherited_session_label() {
let root = sample_tree();
let sel = crate::selector::Selector::new()
.eq("session", "test-session")
.present("phase");
assert_eq!(count(&root, &sel), 4);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn selector_macro_drives_find() {
let root = sample_tree();
let hits = find(&root, &crate::selector!(phase = "teardown"));
assert_eq!(hits.len(), 1);
}
#[tokio::test]
async fn controls_declare_and_lookup_through_component() {
let root = Component::root(Labels::of("session", "s"), HashMap::new());
{
let guard = root.read().unwrap();
guard
.controls()
.declare(crate::controls::ControlBuilder::new("concurrency", 16u32).build());
}
let c: crate::controls::Control<u32> = {
let guard = root.read().unwrap();
guard.controls().get::<u32>("concurrency").unwrap()
};
c.set(32, crate::controls::ControlOrigin::Test)
.await
.unwrap();
let reread: crate::controls::Control<u32> = {
let guard = root.read().unwrap();
guard.controls().get::<u32>("concurrency").unwrap()
};
assert_eq!(reread.value(), 32);
}
#[tokio::test]
async fn reified_control_gauges_flow_through_capture_tree() {
let root = Component::root(Labels::empty().with("session", "s1"), HashMap::new());
let phase = Arc::new(RwLock::new(Component::new(
Labels::empty().with("phase", "rampup"),
HashMap::new(),
)));
attach(&root, &phase);
{
let p = phase.write().unwrap();
p.controls().declare(
crate::controls::ControlBuilder::new("concurrency", 8u32)
.reify_as_gauge(|v| Some(*v as f64))
.build(),
);
}
phase.write().unwrap().set_state(ComponentState::Running);
let c: crate::controls::Control<u32> = phase
.read()
.unwrap()
.controls()
.get::<u32>("concurrency")
.unwrap();
c.set(64, crate::controls::ControlOrigin::Test)
.await
.unwrap();
let captured = capture_tree(&root, Duration::from_secs(1));
let mut found_value: Option<f64> = None;
for (labels, set) in &captured {
if labels.get("phase") != Some("rampup") {
continue;
}
if let Some(fam) = set.family("control_concurrency")
&& let Some(m) = fam.metrics().next()
{
if let Some(p) = m.point()
&& let crate::snapshot::MetricValue::Gauge(g) = p.value()
{
found_value = Some(g.value);
}
assert_eq!(m.labels().get("phase"), Some("rampup"));
assert_eq!(m.labels().get("control"), Some("concurrency"));
}
}
assert_eq!(found_value, Some(64.0));
let current = capture_tree_current(&root);
let mut saw_via_current = false;
for (_, set) in ¤t {
if set.family("control_concurrency").is_some() {
saw_via_current = true;
}
}
assert!(saw_via_current);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn dryrun_controls_enumeration_over_tree() {
let root = sample_tree();
let phase_hits = find(&root, &crate::selector::Selector::new().present("phase"));
for (idx, phase) in phase_hits.iter().enumerate() {
let guard = phase.read().unwrap();
guard.controls().declare(
crate::controls::ControlBuilder::new("concurrency", (10 * (idx + 1)) as u32)
.build(),
);
}
let mut entries: Vec<(String, String)> = Vec::new();
for c in find(&root, &crate::selector::Selector::new()) {
let guard = c.read().unwrap();
let labels = guard.effective_labels().clone();
for ctl in guard.controls().list() {
entries.push((
format!("{}/{}", labels.get("phase").unwrap_or("-"), ctl.name(),),
ctl.value_string(),
));
}
}
assert_eq!(entries.len(), phase_hits.len());
for (key, value) in &entries {
assert!(key.ends_with("/concurrency"), "key = {key}");
assert!(value.parse::<u32>().is_ok(), "value = {value}");
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn branch_scope_subtree_resolves_from_descendant() {
use crate::controls::{BranchScope, ControlBuilder};
let root = Component::root(Labels::empty().with("session", "s1"), HashMap::new());
let phase = Arc::new(RwLock::new(Component::new(
Labels::empty().with("phase", "rampup"),
HashMap::new(),
)));
attach(&root, &phase);
root.read().unwrap().controls().declare(
ControlBuilder::new("hdr_sigdigs", 3u32)
.branch_scope(BranchScope::Subtree)
.build(),
);
let resolved = phase.read().unwrap().find_control_up::<u32>("hdr_sigdigs");
assert!(
resolved.is_some(),
"Subtree-scoped control should be visible to descendant"
);
assert_eq!(resolved.unwrap().value(), 3u32);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn branch_scope_local_does_not_leak_to_descendants() {
use crate::controls::{BranchScope, ControlBuilder};
let root = Component::root(Labels::empty().with("session", "s1"), HashMap::new());
let phase = Arc::new(RwLock::new(Component::new(
Labels::empty().with("phase", "rampup"),
HashMap::new(),
)));
attach(&root, &phase);
root.read().unwrap().controls().declare(
ControlBuilder::new("private", 99u32)
.branch_scope(BranchScope::Local)
.build(),
);
let leaked = phase.read().unwrap().find_control_up::<u32>("private");
assert!(
leaked.is_none(),
"Local-scoped control must not be visible to descendants"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn nearest_declaration_wins_during_walk_up() {
use crate::controls::{BranchScope, ControlBuilder};
let root = Component::root(Labels::empty().with("session", "s1"), HashMap::new());
let phase = Arc::new(RwLock::new(Component::new(
Labels::empty().with("phase", "rampup"),
HashMap::new(),
)));
attach(&root, &phase);
root.read().unwrap().controls().declare(
ControlBuilder::new("hdr_sigdigs", 3u32)
.branch_scope(BranchScope::Subtree)
.build(),
);
phase
.read()
.unwrap()
.controls()
.declare(ControlBuilder::new("hdr_sigdigs", 5u32).build());
let v = phase
.read()
.unwrap()
.find_control_up::<u32>("hdr_sigdigs")
.unwrap()
.value();
assert_eq!(v, 5u32, "phase override should win over session default");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn component_scope_close_flushes_running_component_marks_partial_and_stops() {
use crate::cadence::{CadenceTree, Cadences};
use crate::cadence_reporter::CadenceReporter;
let tree = CadenceTree::plan_default(Cadences::new(&[Duration::from_secs(1)]).unwrap());
let reporter = CadenceReporter::new(tree);
let root = Component::root(Labels::of("session", "s1"), HashMap::new());
let phase = Arc::new(RwLock::new(Component::new(
Labels::of("phase", "short"),
HashMap::new(),
)));
attach(&root, &phase);
{
let mut p = phase.write().unwrap();
p.set_state(ComponentState::Running);
install_counter(&mut p, "test_counter", 42);
}
scope_close(&phase, &reporter, Duration::from_millis(150));
reporter.flush_for_tests();
assert_eq!(phase.read().unwrap().state(), ComponentState::Stopped);
scope_close(&phase, &reporter, Duration::from_millis(150));
reporter.flush_for_tests();
let labels = phase.read().unwrap().effective_labels().clone();
let latest = reporter
.latest(&labels, Duration::from_secs(1))
.expect("scope_close must publish the partial");
assert!(latest.is_partial(), "snapshot must be marked partial");
let f = latest
.family("test_counter")
.expect("test_counter family present");
let m = f.metrics().next().unwrap();
match m.point().unwrap().value() {
crate::snapshot::MetricValue::Counter(c) => assert_eq!(c.cumulative, 42),
v => panic!("expected counter, got {v:?}"),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn component_scope_close_skips_non_running_states() {
use crate::cadence::{CadenceTree, Cadences};
use crate::cadence_reporter::CadenceReporter;
let tree = CadenceTree::plan_default(Cadences::new(&[Duration::from_secs(1)]).unwrap());
let reporter = CadenceReporter::new(tree);
let root = Component::root(Labels::of("session", "s1"), HashMap::new());
let phase = Arc::new(RwLock::new(Component::new(
Labels::of("phase", "starting"),
HashMap::new(),
)));
attach(&root, &phase);
assert_eq!(phase.read().unwrap().state(), ComponentState::Starting);
{
let mut p = phase.write().unwrap();
install_counter(&mut p, "test_counter", 99);
}
scope_close(&phase, &reporter, Duration::from_millis(150));
reporter.flush_for_tests();
assert_eq!(phase.read().unwrap().state(), ComponentState::Starting);
let labels = phase.read().unwrap().effective_labels().clone();
assert!(
reporter.latest(&labels, Duration::from_secs(1)).is_none(),
"scope_close on a non-Running component must not publish"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn capture_delta_auto_uses_fallback_on_first_call() {
let c = Component::new(Labels::empty(), HashMap::new());
let s = c.capture_delta_auto(Duration::from_millis(500));
assert_eq!(s.interval(), Duration::from_millis(500));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn capture_delta_auto_measures_real_elapsed_after_prior_capture() {
let c = Component::new(Labels::empty(), HashMap::new());
let _ = c.capture_delta(Duration::from_secs(1));
std::thread::sleep(Duration::from_millis(80));
let s = c.capture_delta_auto(Duration::from_secs(1));
assert!(
s.interval() > Duration::from_millis(60),
"interval should reflect real ~80ms elapsed, got {:?}",
s.interval()
);
assert!(
s.interval() < Duration::from_millis(500),
"interval should be the real elapsed, not the 1s fallback: {:?}",
s.interval()
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn capture_delta_auto_chained_uses_inter_capture_elapsed() {
let c = Component::new(Labels::empty(), HashMap::new());
let _first = c.capture_delta_auto(Duration::from_secs(1));
std::thread::sleep(Duration::from_millis(50));
let second = c.capture_delta_auto(Duration::from_secs(1));
assert!(
second.interval() < Duration::from_millis(500),
"second auto-capture should measure inter-capture \
elapsed, not cumulative: {:?}",
second.interval()
);
}
}