nmbrs_metrics/queryapi/
live.rs1use std::sync::Arc;
30use std::time::SystemTime;
31
32fn stamp_ms(now_ms: i64, captured_at: std::time::Instant) -> i64 {
39 now_ms - captured_at.elapsed().as_micros().div_ceil(1000) as i64
40}
41
42use super::{Matcher, MetricAccess, QueryError, Sample, Series, Vector};
43use crate::metrics_query::MetricsQuery;
44use crate::snapshot::MetricValue;
45
46pub struct MetricsQueryAccess {
49 query: Arc<MetricsQuery>,
50}
51
52impl MetricsQueryAccess {
53 pub fn new(query: Arc<MetricsQuery>) -> Self {
55 Self { query }
56 }
57
58 pub fn earliest_ms(&self) -> Option<i64> {
65 let reporter = self.query.reporter();
66 let cadence = reporter.declared_cadences().smallest();
67 if cadence.is_zero() {
68 return None;
69 }
70 let now_ms = SystemTime::now()
71 .duration_since(SystemTime::UNIX_EPOCH)
72 .map(|d| d.as_millis() as i64)
73 .unwrap_or(0);
74 let mut earliest: Option<i64> = None;
75 for component in reporter.component_labels() {
76 if let Some(oldest) = reporter
80 .window_view(&component, cadence)
81 .and_then(|v| v.ring.first().cloned())
82 {
83 let ms = stamp_ms(now_ms, oldest.captured_at());
84 earliest = Some(earliest.map_or(ms, |e: i64| e.min(ms)));
85 }
86 }
87 earliest
88 }
89}
90
91impl super::HorizonAware for MetricsQueryAccess {
92 fn earliest_ms(&self) -> Option<i64> {
93 MetricsQueryAccess::earliest_ms(self)
94 }
95}
96
97impl MetricAccess for MetricsQueryAccess {
98 fn select_range(
99 &self,
100 matchers: &[Matcher],
101 start_ms: i64,
102 end_ms: i64,
103 ) -> Result<Vector, QueryError> {
104 let reporter = self.query.reporter();
105 let cadence = reporter.declared_cadences().smallest();
106 if cadence.is_zero() {
107 return Ok(Vector::default());
108 }
109 let now_ms = SystemTime::now()
110 .duration_since(SystemTime::UNIX_EPOCH)
111 .map(|d| d.as_millis() as i64)
112 .unwrap_or(0);
113
114 let mut out: Vec<Series> = Vec::new();
123 for component in reporter.component_labels() {
124 let Some(view) = reporter.window_view(&component, cadence) else {
130 continue;
131 };
132 let mut windows: Vec<Arc<crate::snapshot::MetricSet>> =
133 Vec::with_capacity(view.ring.len() + 2);
134 windows.extend(view.ring.iter().cloned());
135 if let Some(l) = &view.latest {
136 windows.push(l.clone());
137 }
138 if let Some(p) = &view.prebuffer {
139 windows.push(p.clone());
140 }
141 windows.sort_by_key(|w| w.captured_at());
142 windows.dedup_by_key(|w| w.captured_at());
143
144 for window in &windows {
145 let window_ms = stamp_ms(now_ms, window.captured_at());
146 if window_ms < start_ms || window_ms > end_ms {
147 continue;
148 }
149 for family in window.families() {
150 let name = family.name();
151 for metric in family.metrics() {
152 let labels = series_labels(name, &component, metric.labels());
153 if !matchers.iter().all(|m| m.matches(&labels)) {
154 continue;
155 }
156 let Some(value) = metric.point().and_then(|p| value_to_f64(p.value()))
157 else {
158 continue;
159 };
160 push_sample(
161 &mut out,
162 labels,
163 Sample {
164 timestamp_ms: window_ms,
165 value,
166 },
167 );
168 }
169 }
170 }
171 }
172
173 for s in &mut out {
174 s.samples.sort_by_key(|x| x.timestamp_ms);
175 }
176 Ok(Vector::new(out))
177 }
178
179 }
181
182fn series_labels(
186 name: &str,
187 component: &crate::labels::Labels,
188 metric: &crate::labels::Labels,
189) -> Vec<(String, String)> {
190 let mut out = vec![("__name__".to_string(), name.to_string())];
191 let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
192 seen.insert("__name__".to_string());
193 for (k, v) in metric.iter().chain(component.iter()) {
194 if seen.insert(k.to_string()) {
195 out.push((k.to_string(), v.to_string()));
196 }
197 }
198 out
199}
200
201fn value_to_f64(v: &MetricValue) -> Option<f64> {
204 match v {
205 MetricValue::Counter(c) => Some(c.cumulative as f64),
209 MetricValue::Gauge(g) => Some(g.value),
210 MetricValue::Histogram(h) => Some(h.cumulative_count as f64),
213 MetricValue::BucketedHistogram(b) => Some(b.cumulative_count as f64),
214 _ => None,
215 }
216}
217
218fn push_sample(out: &mut Vec<Series>, labels: Vec<(String, String)>, sample: Sample) {
221 if let Some(existing) = out.iter_mut().find(|s| s.labels == labels) {
222 existing.samples.push(sample);
223 } else {
224 out.push(Series {
225 labels,
226 samples: vec![sample],
227 });
228 }
229}
230
231#[cfg(test)]
232mod tests {
233 use super::*;
234
235 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
236 async fn value_projection_per_variant() {
237 use crate::snapshot::{CounterValue, GaugeValue, HistogramValue};
238 use hdrhistogram::Histogram as HdrHistogram;
239 assert_eq!(
240 value_to_f64(&MetricValue::Counter(CounterValue::new(7))),
241 Some(7.0)
242 );
243 assert_eq!(
244 value_to_f64(&MetricValue::Gauge(GaugeValue::new(3.5))),
245 Some(3.5)
246 );
247 let mut hdr = HdrHistogram::<u64>::new(3).unwrap();
251 hdr.record(5).unwrap();
252 hdr.record(6).unwrap();
253 let hv = HistogramValue::from_hdr(hdr).with_cumulative_count(42);
254 assert_eq!(value_to_f64(&MetricValue::Histogram(hv)), Some(42.0));
255 }
256
257 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
258 async fn series_labels_promote_name_and_merge_scope() {
259 let component = crate::labels::Labels::empty().with("phase", "saturate");
260 let metric = crate::labels::Labels::empty().with("__name__", "stale");
261 let out = series_labels("errors_total", &component, &metric);
262 assert_eq!(out[0], ("__name__".to_string(), "errors_total".to_string()));
263 assert_eq!(out.iter().filter(|(k, _)| k == "__name__").count(), 1);
264 assert!(out.iter().any(|(k, v)| k == "phase" && v == "saturate"));
265 }
266
267 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
268 async fn select_instant_returns_a_live_counter() {
269 use crate::cadence::{CadenceTree, Cadences};
270 use crate::cadence_reporter::CadenceReporter;
271 use crate::labels::Labels;
272 use crate::snapshot::MetricSet;
273 use std::collections::HashMap;
274 use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
275
276 let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
278 let reporter = Arc::new(CadenceReporter::new(CadenceTree::plan_default(cadences)));
279 let comp = Labels::of("phase", "saturate");
280 let mut delta = MetricSet::new(Duration::from_millis(100));
281 delta.insert_counter("ops", Labels::default(), 7, Instant::now());
282 reporter.scope_close(&comp, delta);
283 reporter.flush_for_tests();
284
285 let root = crate::component::Component::root(Labels::of("session", "s1"), HashMap::new());
286 let access = MetricsQueryAccess::new(Arc::new(MetricsQuery::new(reporter.clone(), root)));
287
288 let now_ms = SystemTime::now()
289 .duration_since(UNIX_EPOCH)
290 .unwrap()
291 .as_millis() as i64;
292 let v = access
293 .select_instant(&[Matcher::eq("__name__", "ops")], now_ms, Some(60_000))
294 .expect("select");
295
296 assert_eq!(v.len(), 1, "one series, got {v:?}");
297 let s = &v.series()[0];
298 assert!(
299 s.labels
300 .iter()
301 .any(|(k, val)| k == "__name__" && val == "ops")
302 );
303 assert!(
304 s.labels
305 .iter()
306 .any(|(k, val)| k == "phase" && val == "saturate"),
307 "merged scope tag: {:?}",
308 s.labels
309 );
310 assert_eq!(s.samples.len(), 1, "instant vector: one sample per series");
311 assert_eq!(s.samples[0].value, 7.0);
312 }
313}