1use std::sync::{Arc, RwLock};
34use std::time::{Duration, Instant};
35
36use crate::cadence_reporter::CadenceReporter;
37use crate::component::Component;
38use crate::labels::Labels;
39use crate::snapshot::{CounterValue, Metric, MetricFamily, MetricPoint, MetricSet, MetricValue};
40
41#[derive(Clone, Debug, Default)]
44pub struct Selection {
45 family: Option<String>,
47 label_eq: Vec<(String, String)>,
50 label_contains: Vec<(String, String)>,
53}
54
55impl Selection {
56 pub fn all() -> Self {
57 Self::default()
58 }
59
60 pub fn family(name: impl Into<String>) -> Self {
62 Self {
63 family: Some(name.into()),
64 ..Default::default()
65 }
66 }
67
68 pub fn with_family(mut self, name: impl Into<String>) -> Self {
70 self.family = Some(name.into());
71 self
72 }
73
74 pub fn with_label(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
76 self.label_eq.push((key.into(), value.into()));
77 self
78 }
79
80 pub fn with_label_containing(
83 mut self,
84 key: impl Into<String>,
85 substring: impl Into<String>,
86 ) -> Self {
87 self.label_contains.push((key.into(), substring.into()));
88 self
89 }
90
91 pub fn matches_family(&self, family_name: &str) -> bool {
94 self.family
95 .as_deref()
96 .map(|f| f == family_name)
97 .unwrap_or(true)
98 }
99
100 pub fn matches_labels(&self, labels: &Labels) -> bool {
102 for (k, v) in &self.label_eq {
103 if labels.get(k) != Some(v.as_str()) {
104 return false;
105 }
106 }
107 for (k, sub) in &self.label_contains {
108 match labels.get(k) {
109 Some(value) if value.contains(sub.as_str()) => {}
110 _ => return false,
111 }
112 }
113 true
114 }
115}
116
117#[derive(Debug, PartialEq, Eq)]
119pub enum SelectError {
120 NoMatch,
122 MultipleMatches(usize),
125}
126
127impl std::fmt::Display for SelectError {
128 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
129 match self {
130 Self::NoMatch => write!(f, "selection matched no metric instance"),
131 Self::MultipleMatches(n) => {
132 write!(f, "selection matched {n} instances, expected exactly one")
133 }
134 }
135 }
136}
137
138impl std::error::Error for SelectError {}
139
140pub struct MetricsQuery {
145 reporter: Arc<CadenceReporter>,
146 component_root: Arc<RwLock<Component>>,
151}
152
153impl MetricsQuery {
154 pub fn new(reporter: Arc<CadenceReporter>, component_root: Arc<RwLock<Component>>) -> Self {
155 Self {
156 reporter,
157 component_root,
158 }
159 }
160
161 pub fn reporter(&self) -> &Arc<CadenceReporter> {
165 &self.reporter
166 }
167
168 pub fn running_phase_count(&self) -> usize {
174 self.component_root
175 .read()
176 .map(|c| c.running_descendant_count())
177 .unwrap_or(0)
178 }
179
180 pub fn now(&self, selection: &Selection) -> MetricSet {
207 let smallest = self.reporter.declared_cadences().smallest();
208 if smallest.is_zero() {
209 return MetricSet::at(Instant::now(), Duration::ZERO);
210 }
211 self.cadence_window(smallest, selection)
212 }
213
214 pub fn resolve(&self, selection: Selection) -> MetricHandle {
223 let cadence = self.reporter.declared_cadences().smallest();
224 MetricHandle {
225 reporter: self.reporter.clone(),
226 selection,
227 cadence,
228 }
229 }
230
231 pub fn resolve_at(&self, selection: Selection, cadence: Duration) -> MetricHandle {
235 MetricHandle {
236 reporter: self.reporter.clone(),
237 selection,
238 cadence,
239 }
240 }
241
242 pub fn cadence_window(&self, cadence: Duration, selection: &Selection) -> MetricSet {
252 let mut out = MetricSet::at(Instant::now(), cadence);
253 for component in self.reporter.component_labels() {
254 let Some(snap) = self.reporter.latest(&component, cadence) else {
255 continue;
256 };
257 for family in snap.families() {
258 if !selection.matches_family(family.name()) {
259 continue;
260 }
261 for metric in family.metrics() {
262 if !selection.matches_labels(metric.labels()) {
263 continue;
264 }
265 insert_metric_into(&mut out, family, metric);
266 }
267 }
268 }
269 out
270 }
271
272 fn finest_cadence_covering(&self, span: Duration) -> Option<Duration> {
278 let cap = crate::cadence_reporter::HISTORY_RING_CAP as u128;
279 let mut chosen: Option<Duration> = None;
280 for layer in self.reporter.layers() {
281 if layer.hidden {
282 continue;
283 }
284 chosen = Some(layer.interval);
285 if layer.interval.as_nanos().saturating_mul(cap) >= span.as_nanos() {
286 break;
287 }
288 }
289 chosen
290 }
291
292 pub fn increase_over(&self, span: Duration, selection: &Selection) -> MetricSet {
301 let mut out = MetricSet::at(Instant::now(), span);
302 let Some(cadence) = self.finest_cadence_covering(span) else {
303 return out;
304 };
305 let windows = ((span.as_nanos().max(1)) / (cadence.as_nanos().max(1))).max(1) as usize;
306 let now = Instant::now();
307 for component in self.reporter.component_labels() {
308 let ring = self.reporter.ring(&component, cadence);
309 if ring.is_empty() {
310 continue;
311 }
312 let end = ring.len();
313 let start = end.saturating_sub(windows);
314 let per = coalesce_component_windows(
315 ring[start..end].iter().map(|a| a.as_ref()),
316 selection,
317 now,
318 span,
319 );
320 let baseline = start.checked_sub(1).map(|i| ring[i].clone());
324 for family in per.families() {
325 for metric in family.metrics() {
326 insert_counter_increase_into(&mut out, family, metric, baseline.as_deref());
327 }
328 }
329 }
330 out
331 }
332
333 pub fn distribution_over(&self, span: Duration, selection: &Selection) -> MetricSet {
340 let mut out = MetricSet::at(Instant::now(), span);
341 let Some(cadence) = self.finest_cadence_covering(span) else {
342 return out;
343 };
344 let windows = ((span.as_nanos().max(1)) / (cadence.as_nanos().max(1))).max(1) as usize;
345 let now = Instant::now();
346 for component in self.reporter.component_labels() {
347 let ring = self.reporter.ring(&component, cadence);
348 if ring.is_empty() {
349 continue;
350 }
351 let end = ring.len();
352 let start = end.saturating_sub(windows);
353 let per = coalesce_component_windows(
354 ring[start..end].iter().map(|a| a.as_ref()),
355 selection,
356 now,
357 span,
358 );
359 for family in per.families() {
360 for metric in family.metrics() {
361 if matches!(
362 metric.point().map(|p| p.value()),
363 Some(MetricValue::Histogram(_)) | Some(MetricValue::BucketedHistogram(_))
364 ) {
365 insert_metric_into(&mut out, family, metric);
366 }
367 }
368 }
369 }
370 out
371 }
372
373 pub fn session_lifetime(&self, selection: &Selection) -> MetricSet {
389 let session_age = self.reporter.started_at().elapsed();
390 let now = Instant::now();
391 let mut out = MetricSet::at(now, session_age);
392 let largest = self.reporter.layers().last().map(|l| l.interval);
393
394 for component in self.reporter.component_labels() {
395 let mut sources: Vec<MetricSet> = Vec::new();
397 for layer in self.reporter.layers() {
398 if let Some(pre) = self.reporter.prebuffer(&component, layer.interval) {
399 sources.push(pre);
400 }
401 if Some(layer.interval) == largest
408 && let Some(latest) = self.reporter.latest(&component, layer.interval)
409 {
410 sources.push((*latest).clone());
411 }
412 }
413 if sources.is_empty() {
414 continue;
415 }
416 let per = coalesce_component_windows(sources.iter(), selection, now, session_age);
418 for family in per.families() {
419 for metric in family.metrics() {
420 insert_metric_into(&mut out, family, metric);
421 }
422 }
423 }
424
425 out
426 }
427
428 pub fn select_one<F>(&self, mode: F) -> Result<MetricSet, SelectError>
434 where
435 F: FnOnce(&Self) -> MetricSet,
436 {
437 let snap = mode(self);
438 let total: usize = snap.families().map(|f| f.len()).sum();
439 match total {
440 0 => Err(SelectError::NoMatch),
441 1 => Ok(snap),
442 n => Err(SelectError::MultipleMatches(n)),
443 }
444 }
445}
446
447pub struct MetricHandle {
462 reporter: Arc<CadenceReporter>,
463 selection: Selection,
464 cadence: Duration,
465}
466
467impl MetricHandle {
468 pub fn read_now(&self) -> MetricSet {
472 let mut out = MetricSet::at(Instant::now(), self.cadence);
473 for component in self.reporter.component_labels() {
474 let Some(snap) = self.reporter.latest(&component, self.cadence) else {
475 continue;
476 };
477 for family in snap.families() {
478 if !self.selection.matches_family(family.name()) {
479 continue;
480 }
481 for metric in family.metrics() {
482 if !self.selection.matches_labels(metric.labels()) {
483 continue;
484 }
485 insert_metric_into(&mut out, family, metric);
486 }
487 }
488 }
489 out
490 }
491
492 pub fn refresh(&mut self) {}
498
499 pub fn selection(&self) -> &Selection {
501 &self.selection
502 }
503
504 pub fn cadence(&self) -> Duration {
507 self.cadence
508 }
509
510 pub fn source_count(&self) -> usize {
513 self.reporter.component_labels().len()
514 }
515}
516
517fn insert_metric_into(out: &mut MetricSet, family: &MetricFamily, metric: &Metric) {
523 insert_metric_with_mode(out, family, metric, crate::snapshot::CombineMode::Aggregate);
524}
525
526fn insert_metric_with_mode(
534 out: &mut MetricSet,
535 family: &MetricFamily,
536 metric: &Metric,
537 mode: crate::snapshot::CombineMode,
538) {
539 let Some(point) = metric.point() else { return };
540 let existing = out
541 .family(family.name())
542 .and_then(|f| f.metric_with_labels(metric.labels()))
543 .is_some();
544 if existing {
545 let mut tmp = MetricSet::at(out.captured_at(), out.interval());
551 tmp.insert_metric(
552 family.name().to_string(),
553 family.r#type(),
554 metric.labels().clone(),
555 point.value().clone(),
556 point.timestamp().unwrap_or(out.captured_at()),
557 );
558 let merged = MetricSet::coalesce_with_mode(
559 std::slice::from_ref(out)
560 .iter()
561 .chain(std::slice::from_ref(&tmp).iter())
562 .cloned()
563 .collect::<Vec<_>>()
564 .as_slice(),
565 mode,
566 );
567 *out = merged;
568 } else {
569 out.insert_metric(
570 family.name().to_string(),
571 family.r#type(),
572 metric.labels().clone(),
573 point.value().clone(),
574 point.timestamp().unwrap_or(out.captured_at()),
575 );
576 }
577}
578
579fn coalesce_component_windows<'a>(
586 sources: impl IntoIterator<Item = &'a MetricSet>,
587 selection: &Selection,
588 captured_at: Instant,
589 interval: Duration,
590) -> MetricSet {
591 let mut per = MetricSet::at(captured_at, interval);
592 for src in sources {
593 for family in src.families() {
594 if !selection.matches_family(family.name()) {
595 continue;
596 }
597 for metric in family.metrics() {
598 if !selection.matches_labels(metric.labels()) {
599 continue;
600 }
601 insert_metric_with_mode(
602 &mut per,
603 family,
604 metric,
605 crate::snapshot::CombineMode::Coalesce,
606 );
607 }
608 }
609 }
610 per
611}
612
613fn insert_counter_increase_into(
620 out: &mut MetricSet,
621 family: &MetricFamily,
622 metric: &Metric,
623 baseline: Option<&MetricSet>,
624) {
625 let Some(point) = metric.point() else { return };
626 let MetricValue::Counter(c) = point.value() else {
627 return;
628 };
629 let base = baseline
630 .and_then(|b| b.family(family.name()))
631 .and_then(|f| f.metric_with_labels(metric.labels()))
632 .and_then(|m| m.point())
633 .and_then(|p| match p.value() {
634 MetricValue::Counter(bc) => Some(bc.cumulative),
635 _ => None,
636 })
637 .unwrap_or(0);
638 let increase = c.cumulative.saturating_sub(base);
639 let inc = Metric::single(
640 metric.labels().clone(),
641 MetricPoint::new(
642 MetricValue::Counter(CounterValue::new(increase)),
643 point.timestamp().unwrap_or(out.captured_at()),
644 ),
645 );
646 insert_metric_into(out, family, &inc);
647}
648
649#[cfg(test)]
654mod tests {
655 use super::*;
656 use crate::cadence::{CadenceTree, Cadences};
657 use crate::component::{Component, ComponentState, InstrumentRef, attach};
658 use crate::instruments::counter::Counter;
659 use crate::snapshot::MetricValue;
660 use std::collections::HashMap;
661
662 fn build_one_component_query() -> (Arc<RwLock<Component>>, Arc<CadenceReporter>, MetricsQuery) {
663 let root = Component::root(Labels::of("session", "s1"), HashMap::new());
664 let phase = Arc::new(RwLock::new(Component::new(
665 Labels::of("phase", "load"),
666 HashMap::new(),
667 )));
668 attach(&root, &phase);
669 {
670 let mut p = phase.write().unwrap();
671 p.set_state(ComponentState::Running);
672 let counter = Arc::new(Counter::new(Labels::of("name", "ops")));
673 counter.inc_by(7);
674 p.register_instrument("ops", InstrumentRef::Counter(counter))
675 .unwrap();
676 }
677
678 let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
679 let tree = CadenceTree::plan_default(cadences);
680 let reporter = Arc::new(CadenceReporter::new(tree));
681 let query = MetricsQuery::new(reporter.clone(), root.clone());
682 (root, reporter, query)
683 }
684
685 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
686 async fn now_reads_smallest_cadence_window() {
687 let (_root, reporter, query) = build_one_component_query();
691 assert!(
692 query.now(&Selection::family("ops")).is_empty(),
693 "pre-close now should be empty"
694 );
695
696 let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
697 let mut s = MetricSet::new(Duration::from_millis(100));
698 s.insert_counter("ops", Labels::default(), 42, Instant::now());
699 reporter.ingest(&labels, s);
700 reporter.flush_for_tests();
701
702 let snap = query.now(&Selection::family("ops"));
703 let total = match snap
704 .family("ops")
705 .unwrap()
706 .metrics()
707 .next()
708 .unwrap()
709 .point()
710 .unwrap()
711 .value()
712 {
713 MetricValue::Counter(c) => c.cumulative,
714 _ => panic!("not a counter"),
715 };
716 assert_eq!(total, 42);
717 }
718
719 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
720 async fn cadence_window_returns_latest_closed_snapshot() {
721 let (_root, reporter, query) = build_one_component_query();
722 let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
723 let mut s = MetricSet::new(Duration::from_millis(100));
725 s.insert_counter("ops", Labels::default(), 99, Instant::now());
726 reporter.ingest(&labels, s);
727 reporter.flush_for_tests();
728
729 let snap = query.cadence_window(Duration::from_millis(100), &Selection::family("ops"));
730 let f = snap
731 .family("ops")
732 .expect("ops family in cadence_window result");
733 match f.metrics().next().unwrap().point().unwrap().value() {
734 MetricValue::Counter(c) => assert_eq!(c.cumulative, 99),
735 _ => panic!("not a counter"),
736 }
737 }
738
739 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
740 async fn session_lifetime_does_not_overcount_a_single_counter() {
741 let (_root, reporter, query) = build_one_component_query();
746 let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
747 let mut s = MetricSet::new(Duration::from_millis(100));
748 s.insert_counter("ops", Labels::default(), 42, Instant::now());
749 reporter.ingest(&labels, s);
750 reporter.flush_for_tests();
751
752 let snap = query.session_lifetime(&Selection::family("ops"));
753 let cumulative = match snap
754 .family("ops")
755 .unwrap()
756 .metrics()
757 .next()
758 .unwrap()
759 .point()
760 .unwrap()
761 .value()
762 {
763 MetricValue::Counter(c) => c.cumulative,
764 _ => panic!("not a counter"),
765 };
766 assert_eq!(
767 cumulative, 42,
768 "session_lifetime cumulative overcounted (got {cumulative})"
769 );
770 }
771
772 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
773 async fn increase_over_gives_the_span_increment() {
774 let (_root, reporter, query) = build_one_component_query();
778 let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
779
780 let mut s1 = MetricSet::new(Duration::from_millis(100));
781 s1.insert_counter("ops", Labels::default(), 10, Instant::now());
782 reporter.ingest(&labels, s1);
783 reporter.flush_for_tests();
784 std::thread::sleep(Duration::from_millis(2)); let mut s2 = MetricSet::new(Duration::from_millis(100));
786 s2.insert_counter("ops", Labels::default(), 20, Instant::now());
787 reporter.ingest(&labels, s2);
788 reporter.flush_for_tests();
789
790 let snap = query.increase_over(Duration::from_millis(250), &Selection::family("ops"));
791 let value = match snap
792 .family("ops")
793 .unwrap()
794 .metrics()
795 .next()
796 .unwrap()
797 .point()
798 .unwrap()
799 .value()
800 {
801 MetricValue::Counter(c) => c.cumulative,
802 _ => panic!("not a counter"),
803 };
804 assert_eq!(
805 value, 20,
806 "span increment over the recent window (got {value})"
807 );
808 }
809
810 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
811 async fn increase_over_subtracts_prior_cumulative() {
812 let (_root, reporter, query) = build_one_component_query();
816 let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
817 for v in [100u64, 110, 120] {
818 let mut s = MetricSet::new(Duration::from_millis(100));
819 s.insert_counter("ops", Labels::default(), v, Instant::now());
820 reporter.ingest(&labels, s);
821 reporter.flush_for_tests();
822 std::thread::sleep(Duration::from_millis(2));
823 }
824 let read = |span| match query
825 .increase_over(span, &Selection::family("ops"))
826 .family("ops")
827 .unwrap()
828 .metrics()
829 .next()
830 .unwrap()
831 .point()
832 .unwrap()
833 .value()
834 {
835 MetricValue::Counter(c) => c.cumulative,
836 _ => panic!("not a counter"),
837 };
838 assert_eq!(
839 read(Duration::from_millis(100)),
840 10,
841 "last window increase = 120−110"
842 );
843 assert_eq!(
844 read(Duration::from_millis(200)),
845 20,
846 "last two windows increase = 120−100"
847 );
848 }
849
850 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
851 async fn distribution_over_merges_histogram_windows() {
852 use hdrhistogram::Histogram as HdrHistogram;
853 let (_root, reporter, query) = build_one_component_query();
856 let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
857 for _ in 0..2 {
858 let mut s = MetricSet::new(Duration::from_millis(100));
859 let mut h = HdrHistogram::<u64>::new_with_bounds(1, 3_600_000_000_000, 3).unwrap();
860 h.record(1_000_000).unwrap();
861 h.record(2_000_000).unwrap();
862 s.insert_histogram("latency", Labels::default(), h, Instant::now());
863 reporter.ingest(&labels, s);
864 reporter.flush_for_tests();
865 std::thread::sleep(Duration::from_millis(2));
866 }
867 let snap =
868 query.distribution_over(Duration::from_millis(250), &Selection::family("latency"));
869 let count = match snap
870 .family("latency")
871 .unwrap()
872 .metrics()
873 .next()
874 .unwrap()
875 .point()
876 .unwrap()
877 .value()
878 {
879 MetricValue::Histogram(h) => h.count,
880 _ => panic!("not a histogram"),
881 };
882 assert_eq!(count, 4, "two windows of 2 samples merge to 4");
883 }
884
885 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
886 async fn selection_filter_excludes_non_matching_labels() {
887 let (_root, reporter, query) = build_one_component_query();
888 let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
889 let mut s = MetricSet::new(Duration::from_millis(100));
890 s.insert_counter("ops", Labels::of("kind", "a"), 5, Instant::now());
891 s.insert_counter("ops", Labels::of("kind", "b"), 9, Instant::now());
892 reporter.ingest(&labels, s);
893 reporter.flush_for_tests();
894
895 let snap = query.cadence_window(
896 Duration::from_millis(100),
897 &Selection::family("ops").with_label("kind", "b"),
898 );
899 let f = snap.family("ops").expect("ops family");
900 assert_eq!(f.len(), 1);
901 match f.metrics().next().unwrap().point().unwrap().value() {
902 MetricValue::Counter(c) => assert_eq!(c.cumulative, 9),
903 _ => panic!("not a counter"),
904 }
905 }
906
907 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
908 async fn select_one_errors_on_zero_matches() {
909 let (_root, _reporter, query) = build_one_component_query();
910 let result = query.select_one(|q| {
911 q.cadence_window(
912 Duration::from_millis(100),
913 &Selection::family("nonexistent"),
914 )
915 });
916 assert_eq!(result.unwrap_err(), SelectError::NoMatch);
917 }
918
919 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
920 async fn select_one_succeeds_on_exact_match() {
921 let (_root, reporter, query) = build_one_component_query();
922 let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
923 let mut s = MetricSet::new(Duration::from_millis(100));
924 s.insert_counter("ops", Labels::default(), 1, Instant::now());
925 reporter.ingest(&labels, s);
926 reporter.flush_for_tests();
927
928 let result = query.select_one(|q| {
929 q.cadence_window(Duration::from_millis(100), &Selection::family("ops"))
930 });
931 assert!(result.is_ok());
932 }
933}