use std::sync::{Arc, RwLock};
use std::time::{Duration, Instant};
use crate::cadence_reporter::CadenceReporter;
use crate::component::Component;
use crate::labels::Labels;
use crate::snapshot::{CounterValue, Metric, MetricFamily, MetricPoint, MetricSet, MetricValue};
#[derive(Clone, Debug, Default)]
pub struct Selection {
family: Option<String>,
label_eq: Vec<(String, String)>,
label_contains: Vec<(String, String)>,
}
impl Selection {
pub fn all() -> Self {
Self::default()
}
pub fn family(name: impl Into<String>) -> Self {
Self {
family: Some(name.into()),
..Default::default()
}
}
pub fn with_family(mut self, name: impl Into<String>) -> Self {
self.family = Some(name.into());
self
}
pub fn with_label(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
self.label_eq.push((key.into(), value.into()));
self
}
pub fn with_label_containing(
mut self,
key: impl Into<String>,
substring: impl Into<String>,
) -> Self {
self.label_contains.push((key.into(), substring.into()));
self
}
pub fn matches_family(&self, family_name: &str) -> bool {
self.family
.as_deref()
.map(|f| f == family_name)
.unwrap_or(true)
}
pub fn matches_labels(&self, labels: &Labels) -> bool {
for (k, v) in &self.label_eq {
if labels.get(k) != Some(v.as_str()) {
return false;
}
}
for (k, sub) in &self.label_contains {
match labels.get(k) {
Some(value) if value.contains(sub.as_str()) => {}
_ => return false,
}
}
true
}
}
#[derive(Debug, PartialEq, Eq)]
pub enum SelectError {
NoMatch,
MultipleMatches(usize),
}
impl std::fmt::Display for SelectError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::NoMatch => write!(f, "selection matched no metric instance"),
Self::MultipleMatches(n) => {
write!(f, "selection matched {n} instances, expected exactly one")
}
}
}
}
impl std::error::Error for SelectError {}
pub struct MetricsQuery {
reporter: Arc<CadenceReporter>,
component_root: Arc<RwLock<Component>>,
}
impl MetricsQuery {
pub fn new(reporter: Arc<CadenceReporter>, component_root: Arc<RwLock<Component>>) -> Self {
Self {
reporter,
component_root,
}
}
pub fn reporter(&self) -> &Arc<CadenceReporter> {
&self.reporter
}
pub fn running_phase_count(&self) -> usize {
self.component_root
.read()
.map(|c| c.running_descendant_count())
.unwrap_or(0)
}
pub fn now(&self, selection: &Selection) -> MetricSet {
let smallest = self.reporter.declared_cadences().smallest();
if smallest.is_zero() {
return MetricSet::at(Instant::now(), Duration::ZERO);
}
self.cadence_window(smallest, selection)
}
pub fn resolve(&self, selection: Selection) -> MetricHandle {
let cadence = self.reporter.declared_cadences().smallest();
MetricHandle {
reporter: self.reporter.clone(),
selection,
cadence,
}
}
pub fn resolve_at(&self, selection: Selection, cadence: Duration) -> MetricHandle {
MetricHandle {
reporter: self.reporter.clone(),
selection,
cadence,
}
}
pub fn cadence_window(&self, cadence: Duration, selection: &Selection) -> MetricSet {
let mut out = MetricSet::at(Instant::now(), cadence);
for component in self.reporter.component_labels() {
let Some(snap) = self.reporter.latest(&component, cadence) else {
continue;
};
for family in snap.families() {
if !selection.matches_family(family.name()) {
continue;
}
for metric in family.metrics() {
if !selection.matches_labels(metric.labels()) {
continue;
}
insert_metric_into(&mut out, family, metric);
}
}
}
out
}
fn finest_cadence_covering(&self, span: Duration) -> Option<Duration> {
let cap = crate::cadence_reporter::HISTORY_RING_CAP as u128;
let mut chosen: Option<Duration> = None;
for layer in self.reporter.layers() {
if layer.hidden {
continue;
}
chosen = Some(layer.interval);
if layer.interval.as_nanos().saturating_mul(cap) >= span.as_nanos() {
break;
}
}
chosen
}
pub fn increase_over(&self, span: Duration, selection: &Selection) -> MetricSet {
let mut out = MetricSet::at(Instant::now(), span);
let Some(cadence) = self.finest_cadence_covering(span) else {
return out;
};
let windows = ((span.as_nanos().max(1)) / (cadence.as_nanos().max(1))).max(1) as usize;
let now = Instant::now();
for component in self.reporter.component_labels() {
let ring = self.reporter.ring(&component, cadence);
if ring.is_empty() {
continue;
}
let end = ring.len();
let start = end.saturating_sub(windows);
let per = coalesce_component_windows(
ring[start..end].iter().map(|a| a.as_ref()),
selection,
now,
span,
);
let baseline = start.checked_sub(1).map(|i| ring[i].clone());
for family in per.families() {
for metric in family.metrics() {
insert_counter_increase_into(&mut out, family, metric, baseline.as_deref());
}
}
}
out
}
pub fn distribution_over(&self, span: Duration, selection: &Selection) -> MetricSet {
let mut out = MetricSet::at(Instant::now(), span);
let Some(cadence) = self.finest_cadence_covering(span) else {
return out;
};
let windows = ((span.as_nanos().max(1)) / (cadence.as_nanos().max(1))).max(1) as usize;
let now = Instant::now();
for component in self.reporter.component_labels() {
let ring = self.reporter.ring(&component, cadence);
if ring.is_empty() {
continue;
}
let end = ring.len();
let start = end.saturating_sub(windows);
let per = coalesce_component_windows(
ring[start..end].iter().map(|a| a.as_ref()),
selection,
now,
span,
);
for family in per.families() {
for metric in family.metrics() {
if matches!(
metric.point().map(|p| p.value()),
Some(MetricValue::Histogram(_)) | Some(MetricValue::BucketedHistogram(_))
) {
insert_metric_into(&mut out, family, metric);
}
}
}
}
out
}
pub fn session_lifetime(&self, selection: &Selection) -> MetricSet {
let session_age = self.reporter.started_at().elapsed();
let now = Instant::now();
let mut out = MetricSet::at(now, session_age);
let largest = self.reporter.layers().last().map(|l| l.interval);
for component in self.reporter.component_labels() {
let mut sources: Vec<MetricSet> = Vec::new();
for layer in self.reporter.layers() {
if let Some(pre) = self.reporter.prebuffer(&component, layer.interval) {
sources.push(pre);
}
if Some(layer.interval) == largest
&& let Some(latest) = self.reporter.latest(&component, layer.interval)
{
sources.push((*latest).clone());
}
}
if sources.is_empty() {
continue;
}
let per = coalesce_component_windows(sources.iter(), selection, now, session_age);
for family in per.families() {
for metric in family.metrics() {
insert_metric_into(&mut out, family, metric);
}
}
}
out
}
pub fn select_one<F>(&self, mode: F) -> Result<MetricSet, SelectError>
where
F: FnOnce(&Self) -> MetricSet,
{
let snap = mode(self);
let total: usize = snap.families().map(|f| f.len()).sum();
match total {
0 => Err(SelectError::NoMatch),
1 => Ok(snap),
n => Err(SelectError::MultipleMatches(n)),
}
}
}
pub struct MetricHandle {
reporter: Arc<CadenceReporter>,
selection: Selection,
cadence: Duration,
}
impl MetricHandle {
pub fn read_now(&self) -> MetricSet {
let mut out = MetricSet::at(Instant::now(), self.cadence);
for component in self.reporter.component_labels() {
let Some(snap) = self.reporter.latest(&component, self.cadence) else {
continue;
};
for family in snap.families() {
if !self.selection.matches_family(family.name()) {
continue;
}
for metric in family.metrics() {
if !self.selection.matches_labels(metric.labels()) {
continue;
}
insert_metric_into(&mut out, family, metric);
}
}
}
out
}
pub fn refresh(&mut self) {}
pub fn selection(&self) -> &Selection {
&self.selection
}
pub fn cadence(&self) -> Duration {
self.cadence
}
pub fn source_count(&self) -> usize {
self.reporter.component_labels().len()
}
}
fn insert_metric_into(out: &mut MetricSet, family: &MetricFamily, metric: &Metric) {
insert_metric_with_mode(out, family, metric, crate::snapshot::CombineMode::Aggregate);
}
fn insert_metric_with_mode(
out: &mut MetricSet,
family: &MetricFamily,
metric: &Metric,
mode: crate::snapshot::CombineMode,
) {
let Some(point) = metric.point() else { return };
let existing = out
.family(family.name())
.and_then(|f| f.metric_with_labels(metric.labels()))
.is_some();
if existing {
let mut tmp = MetricSet::at(out.captured_at(), out.interval());
tmp.insert_metric(
family.name().to_string(),
family.r#type(),
metric.labels().clone(),
point.value().clone(),
point.timestamp().unwrap_or(out.captured_at()),
);
let merged = MetricSet::coalesce_with_mode(
std::slice::from_ref(out)
.iter()
.chain(std::slice::from_ref(&tmp).iter())
.cloned()
.collect::<Vec<_>>()
.as_slice(),
mode,
);
*out = merged;
} else {
out.insert_metric(
family.name().to_string(),
family.r#type(),
metric.labels().clone(),
point.value().clone(),
point.timestamp().unwrap_or(out.captured_at()),
);
}
}
fn coalesce_component_windows<'a>(
sources: impl IntoIterator<Item = &'a MetricSet>,
selection: &Selection,
captured_at: Instant,
interval: Duration,
) -> MetricSet {
let mut per = MetricSet::at(captured_at, interval);
for src in sources {
for family in src.families() {
if !selection.matches_family(family.name()) {
continue;
}
for metric in family.metrics() {
if !selection.matches_labels(metric.labels()) {
continue;
}
insert_metric_with_mode(
&mut per,
family,
metric,
crate::snapshot::CombineMode::Coalesce,
);
}
}
}
per
}
fn insert_counter_increase_into(
out: &mut MetricSet,
family: &MetricFamily,
metric: &Metric,
baseline: Option<&MetricSet>,
) {
let Some(point) = metric.point() else { return };
let MetricValue::Counter(c) = point.value() else {
return;
};
let base = baseline
.and_then(|b| b.family(family.name()))
.and_then(|f| f.metric_with_labels(metric.labels()))
.and_then(|m| m.point())
.and_then(|p| match p.value() {
MetricValue::Counter(bc) => Some(bc.cumulative),
_ => None,
})
.unwrap_or(0);
let increase = c.cumulative.saturating_sub(base);
let inc = Metric::single(
metric.labels().clone(),
MetricPoint::new(
MetricValue::Counter(CounterValue::new(increase)),
point.timestamp().unwrap_or(out.captured_at()),
),
);
insert_metric_into(out, family, &inc);
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cadence::{CadenceTree, Cadences};
use crate::component::{Component, ComponentState, InstrumentRef, attach};
use crate::instruments::counter::Counter;
use crate::snapshot::MetricValue;
use std::collections::HashMap;
fn build_one_component_query() -> (Arc<RwLock<Component>>, Arc<CadenceReporter>, MetricsQuery) {
let root = Component::root(Labels::of("session", "s1"), HashMap::new());
let phase = Arc::new(RwLock::new(Component::new(
Labels::of("phase", "load"),
HashMap::new(),
)));
attach(&root, &phase);
{
let mut p = phase.write().unwrap();
p.set_state(ComponentState::Running);
let counter = Arc::new(Counter::new(Labels::of("name", "ops")));
counter.inc_by(7);
p.register_instrument("ops", InstrumentRef::Counter(counter))
.unwrap();
}
let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
let tree = CadenceTree::plan_default(cadences);
let reporter = Arc::new(CadenceReporter::new(tree));
let query = MetricsQuery::new(reporter.clone(), root.clone());
(root, reporter, query)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn now_reads_smallest_cadence_window() {
let (_root, reporter, query) = build_one_component_query();
assert!(
query.now(&Selection::family("ops")).is_empty(),
"pre-close now should be empty"
);
let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
let mut s = MetricSet::new(Duration::from_millis(100));
s.insert_counter("ops", Labels::default(), 42, Instant::now());
reporter.ingest(&labels, s);
reporter.flush_for_tests();
let snap = query.now(&Selection::family("ops"));
let total = match snap
.family("ops")
.unwrap()
.metrics()
.next()
.unwrap()
.point()
.unwrap()
.value()
{
MetricValue::Counter(c) => c.cumulative,
_ => panic!("not a counter"),
};
assert_eq!(total, 42);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn cadence_window_returns_latest_closed_snapshot() {
let (_root, reporter, query) = build_one_component_query();
let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
let mut s = MetricSet::new(Duration::from_millis(100));
s.insert_counter("ops", Labels::default(), 99, Instant::now());
reporter.ingest(&labels, s);
reporter.flush_for_tests();
let snap = query.cadence_window(Duration::from_millis(100), &Selection::family("ops"));
let f = snap
.family("ops")
.expect("ops family in cadence_window result");
match f.metrics().next().unwrap().point().unwrap().value() {
MetricValue::Counter(c) => assert_eq!(c.cumulative, 99),
_ => panic!("not a counter"),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn session_lifetime_does_not_overcount_a_single_counter() {
let (_root, reporter, query) = build_one_component_query();
let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
let mut s = MetricSet::new(Duration::from_millis(100));
s.insert_counter("ops", Labels::default(), 42, Instant::now());
reporter.ingest(&labels, s);
reporter.flush_for_tests();
let snap = query.session_lifetime(&Selection::family("ops"));
let cumulative = match snap
.family("ops")
.unwrap()
.metrics()
.next()
.unwrap()
.point()
.unwrap()
.value()
{
MetricValue::Counter(c) => c.cumulative,
_ => panic!("not a counter"),
};
assert_eq!(
cumulative, 42,
"session_lifetime cumulative overcounted (got {cumulative})"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn increase_over_gives_the_span_increment() {
let (_root, reporter, query) = build_one_component_query();
let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
let mut s1 = MetricSet::new(Duration::from_millis(100));
s1.insert_counter("ops", Labels::default(), 10, Instant::now());
reporter.ingest(&labels, s1);
reporter.flush_for_tests();
std::thread::sleep(Duration::from_millis(2)); let mut s2 = MetricSet::new(Duration::from_millis(100));
s2.insert_counter("ops", Labels::default(), 20, Instant::now());
reporter.ingest(&labels, s2);
reporter.flush_for_tests();
let snap = query.increase_over(Duration::from_millis(250), &Selection::family("ops"));
let value = match snap
.family("ops")
.unwrap()
.metrics()
.next()
.unwrap()
.point()
.unwrap()
.value()
{
MetricValue::Counter(c) => c.cumulative,
_ => panic!("not a counter"),
};
assert_eq!(
value, 20,
"span increment over the recent window (got {value})"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn increase_over_subtracts_prior_cumulative() {
let (_root, reporter, query) = build_one_component_query();
let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
for v in [100u64, 110, 120] {
let mut s = MetricSet::new(Duration::from_millis(100));
s.insert_counter("ops", Labels::default(), v, Instant::now());
reporter.ingest(&labels, s);
reporter.flush_for_tests();
std::thread::sleep(Duration::from_millis(2));
}
let read = |span| match query
.increase_over(span, &Selection::family("ops"))
.family("ops")
.unwrap()
.metrics()
.next()
.unwrap()
.point()
.unwrap()
.value()
{
MetricValue::Counter(c) => c.cumulative,
_ => panic!("not a counter"),
};
assert_eq!(
read(Duration::from_millis(100)),
10,
"last window increase = 120−110"
);
assert_eq!(
read(Duration::from_millis(200)),
20,
"last two windows increase = 120−100"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn distribution_over_merges_histogram_windows() {
use hdrhistogram::Histogram as HdrHistogram;
let (_root, reporter, query) = build_one_component_query();
let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
for _ in 0..2 {
let mut s = MetricSet::new(Duration::from_millis(100));
let mut h = HdrHistogram::<u64>::new_with_bounds(1, 3_600_000_000_000, 3).unwrap();
h.record(1_000_000).unwrap();
h.record(2_000_000).unwrap();
s.insert_histogram("latency", Labels::default(), h, Instant::now());
reporter.ingest(&labels, s);
reporter.flush_for_tests();
std::thread::sleep(Duration::from_millis(2));
}
let snap =
query.distribution_over(Duration::from_millis(250), &Selection::family("latency"));
let count = match snap
.family("latency")
.unwrap()
.metrics()
.next()
.unwrap()
.point()
.unwrap()
.value()
{
MetricValue::Histogram(h) => h.count,
_ => panic!("not a histogram"),
};
assert_eq!(count, 4, "two windows of 2 samples merge to 4");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn selection_filter_excludes_non_matching_labels() {
let (_root, reporter, query) = build_one_component_query();
let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
let mut s = MetricSet::new(Duration::from_millis(100));
s.insert_counter("ops", Labels::of("kind", "a"), 5, Instant::now());
s.insert_counter("ops", Labels::of("kind", "b"), 9, Instant::now());
reporter.ingest(&labels, s);
reporter.flush_for_tests();
let snap = query.cadence_window(
Duration::from_millis(100),
&Selection::family("ops").with_label("kind", "b"),
);
let f = snap.family("ops").expect("ops family");
assert_eq!(f.len(), 1);
match f.metrics().next().unwrap().point().unwrap().value() {
MetricValue::Counter(c) => assert_eq!(c.cumulative, 9),
_ => panic!("not a counter"),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn select_one_errors_on_zero_matches() {
let (_root, _reporter, query) = build_one_component_query();
let result = query.select_one(|q| {
q.cadence_window(
Duration::from_millis(100),
&Selection::family("nonexistent"),
)
});
assert_eq!(result.unwrap_err(), SelectError::NoMatch);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn select_one_succeeds_on_exact_match() {
let (_root, reporter, query) = build_one_component_query();
let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
let mut s = MetricSet::new(Duration::from_millis(100));
s.insert_counter("ops", Labels::default(), 1, Instant::now());
reporter.ingest(&labels, s);
reporter.flush_for_tests();
let result = query.select_one(|q| {
q.cadence_window(Duration::from_millis(100), &Selection::family("ops"))
});
assert!(result.is_ok());
}
}