1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
39pub enum ComponentState {
40 Starting,
42 Running,
44 Stopping,
46 Stopped,
49}
50
51#[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 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
78pub struct RegisteredInstrument {
88 pub family: String,
89 pub unit: Option<String>,
90 pub instrument: InstrumentRef,
91}
92
93pub trait DynamicCapture: Send + Sync {
103 fn capture_into(&self, out: &mut MetricSet, now: Instant, drain: bool);
108}
109
110pub struct Component {
116 labels: Labels,
118 effective_labels: Labels,
121 props: HashMap<String, String>,
124 parent: Option<Weak<RwLock<Component>>>,
126 children: Vec<Arc<RwLock<Component>>>,
128 state: ComponentState,
130 instruments: Vec<RegisteredInstrument>,
138 dynamic_capture: Option<Arc<dyn DynamicCapture>>,
143 last_capture_instant: Mutex<Option<Instant>>,
151 controls: crate::controls::ControlRegistry,
156 cells: std::sync::Arc<crate::cells::CellMap>,
167 live: Option<std::sync::Arc<()>>,
177 live_children: std::collections::HashMap<String, Vec<std::sync::Weak<()>>>,
182}
183
184impl Component {
185 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 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 pub fn labels(&self) -> &Labels {
216 &self.labels
217 }
218
219 pub fn effective_labels(&self) -> &Labels {
221 &self.effective_labels
222 }
223
224 pub fn state(&self) -> ComponentState {
226 self.state
227 }
228
229 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 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 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 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 pub fn instruments(&self) -> &[RegisteredInstrument] {
304 &self.instruments
305 }
306
307 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 pub fn set_dynamic_capture(&mut self, hook: Arc<dyn DynamicCapture>) {
329 self.dynamic_capture = Some(hook);
330 }
331
332 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 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 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 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 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 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 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 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 pub fn set_prop(&mut self, name: &str, value: &str) {
499 self.props.insert(name.to_string(), value.to_string());
500 }
501
502 pub fn child_count(&self) -> usize {
504 self.children.len()
505 }
506
507 pub fn children(&self) -> impl Iterator<Item = &Arc<RwLock<Component>>> {
509 self.children.iter()
510 }
511
512 pub fn controls(&self) -> &crate::controls::ControlRegistry {
516 &self.controls
517 }
518
519 pub fn cells(&self) -> std::sync::Arc<crate::cells::CellMap> {
532 self.cells.clone()
533 }
534
535 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 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 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 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 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 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
683fn 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
701pub fn scope_close(
717 component: &Arc<RwLock<Component>>,
718 cadence_reporter: &crate::cadence_reporter::CadenceReporter,
719 interval: Duration,
720) {
721 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 let mut g = component.write().unwrap_or_else(|e| e.into_inner());
739 g.state = ComponentState::Stopped;
740}
741
742pub 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
787pub 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 }
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
829pub 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
845pub 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
864pub 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 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
958pub 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
969pub 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 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 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
1014pub 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 #[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 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 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 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 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 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 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 fn sample_tree() -> Arc<RwLock<Component>> {
1316 let root = Component::root(
1317 Labels::empty().with("session", "test-session"),
1318 HashMap::new(),
1319 );
1320 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 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 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 #[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 ¤t {
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 #[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 #[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 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 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 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1775 async fn capture_delta_auto_uses_fallback_on_first_call() {
1776 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 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 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 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}