use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use super::{MatchOp as MatcherOp, Matcher, QueryError as DataSourceError};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MetricFamilyMeta {
pub name: String,
pub ty: MetricType,
pub unit: Option<String>,
pub help: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum MetricType {
Counter,
Gauge,
Histogram,
GaugeHistogram,
Summary,
Info,
StateSet,
Unknown,
}
impl MetricType {
pub fn as_str(&self) -> &'static str {
match self {
MetricType::Counter => "counter",
MetricType::Gauge => "gauge",
MetricType::Histogram => "histogram",
MetricType::GaugeHistogram => "gaugehistogram",
MetricType::Summary => "summary",
MetricType::Info => "info",
MetricType::StateSet => "stateset",
MetricType::Unknown => "unknown",
}
}
pub fn parse(s: &str) -> MetricType {
match s.trim().to_ascii_lowercase().as_str() {
"counter" => MetricType::Counter,
"gauge" => MetricType::Gauge,
"histogram" => MetricType::Histogram,
"gaugehistogram" => MetricType::GaugeHistogram,
"summary" => MetricType::Summary,
"info" => MetricType::Info,
"stateset" => MetricType::StateSet,
_ => MetricType::Unknown,
}
}
pub fn has_derived_series(&self) -> bool {
matches!(
self,
MetricType::Histogram | MetricType::GaugeHistogram | MetricType::Summary,
)
}
pub fn implied_label(&self) -> Option<&'static str> {
match self {
MetricType::Histogram | MetricType::GaugeHistogram => Some("le"),
MetricType::Summary => Some("quantile"),
_ => None,
}
}
}
pub type LabelSet = Vec<(String, String)>;
#[derive(Debug, Clone, PartialEq)]
pub struct ExemplarPoint {
pub series: LabelSet,
pub sample_timestamp_ms: i64,
pub value: f64,
pub timestamp_ms: Option<i64>,
pub labels: Vec<(String, String)>,
}
pub trait MetricCatalog: Send + Sync {
fn metric_families(&self) -> Result<Vec<MetricFamilyMeta>, DataSourceError>;
fn label_keys(&self, family_filter: Option<&str>) -> Result<Vec<String>, DataSourceError>;
fn label_values(
&self,
key: &str,
family_filter: Option<&str>,
) -> Result<Vec<String>, DataSourceError>;
fn series(&self, matchers: &[Matcher]) -> Result<Vec<LabelSet>, DataSourceError>;
fn exemplars(
&self,
matchers: &[Matcher],
time_range: Option<(i64, i64)>,
) -> Result<Vec<ExemplarPoint>, DataSourceError> {
let _ = (matchers, time_range);
Ok(Vec::new())
}
}
pub struct CachedCatalog<C: MetricCatalog + ?Sized> {
inner: Arc<C>,
ttl: Duration,
mtime_fn: Option<Box<dyn Fn() -> Option<Instant> + Send + Sync>>,
state: Mutex<CacheState>,
}
#[derive(Default)]
struct CacheState {
generation: u64,
families: Option<CacheEntry<Vec<MetricFamilyMeta>>>,
label_keys: std::collections::HashMap<Option<String>, CacheEntry<Vec<String>>>,
label_values: std::collections::HashMap<(String, Option<String>), CacheEntry<Vec<String>>>,
series: std::collections::HashMap<String, CacheEntry<Vec<LabelSet>>>,
exemplars: std::collections::HashMap<String, CacheEntry<Vec<ExemplarPoint>>>,
}
struct CacheEntry<T> {
generation: u64,
filled_at: Instant,
backend_mtime: Option<Instant>,
value: T,
}
impl<C: MetricCatalog + ?Sized + 'static> CachedCatalog<C> {
pub fn new(inner: Arc<C>) -> Self {
Self {
inner,
ttl: Duration::from_secs(1),
mtime_fn: None,
state: Mutex::new(CacheState::default()),
}
}
pub fn with_ttl(mut self, ttl: Duration) -> Self {
self.ttl = ttl;
self
}
pub fn with_mtime_fn<F>(mut self, mtime_fn: F) -> Self
where
F: Fn() -> Option<Instant> + Send + Sync + 'static,
{
self.mtime_fn = Some(Box::new(mtime_fn));
self
}
pub fn invalidate(&self) {
if let Ok(mut s) = self.state.lock() {
s.generation = s.generation.wrapping_add(1);
}
}
fn is_fresh<T>(&self, entry: &CacheEntry<T>, gen_now: u64) -> bool {
if entry.generation != gen_now {
return false;
}
if !self.ttl.is_zero() && entry.filled_at.elapsed() >= self.ttl {
return false;
}
if let Some(f) = &self.mtime_fn
&& let Some(now_mtime) = f()
{
match entry.backend_mtime {
Some(prev) if now_mtime > prev => return false,
None => return false,
_ => {}
}
}
true
}
fn current_mtime(&self) -> Option<Instant> {
self.mtime_fn.as_ref().and_then(|f| f())
}
}
impl<C: MetricCatalog + ?Sized + 'static> MetricCatalog for CachedCatalog<C> {
fn metric_families(&self) -> Result<Vec<MetricFamilyMeta>, DataSourceError> {
let gen_now;
{
let state = self
.state
.lock()
.map_err(|_| DataSourceError::new("cache poisoned"))?;
gen_now = state.generation;
if let Some(entry) = &state.families
&& self.is_fresh(entry, gen_now)
{
return Ok(entry.value.clone());
}
}
let value = self.inner.metric_families()?;
let entry = CacheEntry {
generation: gen_now,
filled_at: Instant::now(),
backend_mtime: self.current_mtime(),
value: value.clone(),
};
if let Ok(mut state) = self.state.lock() {
state.families = Some(entry);
}
Ok(value)
}
fn label_keys(&self, family_filter: Option<&str>) -> Result<Vec<String>, DataSourceError> {
let key = family_filter.map(|s| s.to_string());
let gen_now;
{
let state = self
.state
.lock()
.map_err(|_| DataSourceError::new("cache poisoned"))?;
gen_now = state.generation;
if let Some(entry) = state.label_keys.get(&key)
&& self.is_fresh(entry, gen_now)
{
return Ok(entry.value.clone());
}
}
let value = self.inner.label_keys(family_filter)?;
let entry = CacheEntry {
generation: gen_now,
filled_at: Instant::now(),
backend_mtime: self.current_mtime(),
value: value.clone(),
};
if let Ok(mut state) = self.state.lock() {
state.label_keys.insert(key, entry);
}
Ok(value)
}
fn label_values(
&self,
key: &str,
family_filter: Option<&str>,
) -> Result<Vec<String>, DataSourceError> {
let cache_key = (key.to_string(), family_filter.map(|s| s.to_string()));
let gen_now;
{
let state = self
.state
.lock()
.map_err(|_| DataSourceError::new("cache poisoned"))?;
gen_now = state.generation;
if let Some(entry) = state.label_values.get(&cache_key)
&& self.is_fresh(entry, gen_now)
{
return Ok(entry.value.clone());
}
}
let value = self.inner.label_values(key, family_filter)?;
let entry = CacheEntry {
generation: gen_now,
filled_at: Instant::now(),
backend_mtime: self.current_mtime(),
value: value.clone(),
};
if let Ok(mut state) = self.state.lock() {
state.label_values.insert(cache_key, entry);
}
Ok(value)
}
fn series(&self, matchers: &[Matcher]) -> Result<Vec<LabelSet>, DataSourceError> {
let cache_key = encode_matchers(matchers);
let gen_now;
{
let state = self
.state
.lock()
.map_err(|_| DataSourceError::new("cache poisoned"))?;
gen_now = state.generation;
if let Some(entry) = state.series.get(&cache_key)
&& self.is_fresh(entry, gen_now)
{
return Ok(entry.value.clone());
}
}
let value = self.inner.series(matchers)?;
let entry = CacheEntry {
generation: gen_now,
filled_at: Instant::now(),
backend_mtime: self.current_mtime(),
value: value.clone(),
};
if let Ok(mut state) = self.state.lock() {
state.series.insert(cache_key, entry);
}
Ok(value)
}
fn exemplars(
&self,
matchers: &[Matcher],
time_range: Option<(i64, i64)>,
) -> Result<Vec<ExemplarPoint>, DataSourceError> {
let mut cache_key = encode_matchers(matchers);
if let Some((s, e)) = time_range {
cache_key.push_str(&format!("@{s}..{e}"));
}
let gen_now;
{
let state = self
.state
.lock()
.map_err(|_| DataSourceError::new("cache poisoned"))?;
gen_now = state.generation;
if let Some(entry) = state.exemplars.get(&cache_key)
&& self.is_fresh(entry, gen_now)
{
return Ok(entry.value.clone());
}
}
let value = self.inner.exemplars(matchers, time_range)?;
let entry = CacheEntry {
generation: gen_now,
filled_at: Instant::now(),
backend_mtime: self.current_mtime(),
value: value.clone(),
};
if let Ok(mut state) = self.state.lock() {
state.exemplars.insert(cache_key, entry);
}
Ok(value)
}
}
fn encode_matchers(matchers: &[Matcher]) -> String {
let mut sorted: Vec<&Matcher> = matchers.iter().collect();
sorted.sort_by(|a, b| a.label.cmp(&b.label));
let mut out = String::new();
for m in sorted {
let op = match m.op {
MatcherOp::Eq => "=",
MatcherOp::Ne => "!=",
MatcherOp::EqRegex => "=~",
MatcherOp::NeRegex => "!~",
};
out.push_str(&m.label);
out.push_str(op);
out.push_str(&m.value);
out.push('\x1f'); }
out
}
#[cfg(test)]
mod tests {
use super::*;
use MatcherOp;
use std::sync::atomic::{AtomicUsize, Ordering};
struct MockCatalog {
families: Vec<MetricFamilyMeta>,
keys: Vec<(Option<String>, Vec<String>)>,
values: Vec<(String, Option<String>, Vec<String>)>,
series: Vec<LabelSet>,
family_calls: AtomicUsize,
key_calls: AtomicUsize,
value_calls: AtomicUsize,
series_calls: AtomicUsize,
}
impl MockCatalog {
fn new() -> Self {
Self {
families: vec![
MetricFamilyMeta {
name: "ops_total".into(),
ty: MetricType::Counter,
unit: None,
help: Some("ops".into()),
},
MetricFamilyMeta {
name: "latency".into(),
ty: MetricType::Histogram,
unit: Some("seconds".into()),
help: None,
},
],
keys: vec![
(None, vec!["phase".into(), "scenario".into()]),
(Some("ops_total".into()), vec!["phase".into()]),
],
values: vec![("phase".into(), None, vec!["setup".into(), "run".into()])],
series: vec![vec![
("__name__".into(), "ops_total".into()),
("phase".into(), "setup".into()),
]],
family_calls: AtomicUsize::new(0),
key_calls: AtomicUsize::new(0),
value_calls: AtomicUsize::new(0),
series_calls: AtomicUsize::new(0),
}
}
}
impl MetricCatalog for MockCatalog {
fn metric_families(&self) -> Result<Vec<MetricFamilyMeta>, DataSourceError> {
self.family_calls.fetch_add(1, Ordering::SeqCst);
Ok(self.families.clone())
}
fn label_keys(&self, filter: Option<&str>) -> Result<Vec<String>, DataSourceError> {
self.key_calls.fetch_add(1, Ordering::SeqCst);
for (f, v) in &self.keys {
if f.as_deref() == filter {
return Ok(v.clone());
}
}
Ok(Vec::new())
}
fn label_values(
&self,
key: &str,
filter: Option<&str>,
) -> Result<Vec<String>, DataSourceError> {
self.value_calls.fetch_add(1, Ordering::SeqCst);
for (k, f, v) in &self.values {
if k == key && f.as_deref() == filter {
return Ok(v.clone());
}
}
Ok(Vec::new())
}
fn series(&self, _matchers: &[Matcher]) -> Result<Vec<LabelSet>, DataSourceError> {
self.series_calls.fetch_add(1, Ordering::SeqCst);
Ok(self.series.clone())
}
}
#[test]
fn metric_type_round_trips_via_str() {
for ty in [
MetricType::Counter,
MetricType::Gauge,
MetricType::Histogram,
MetricType::GaugeHistogram,
MetricType::Summary,
MetricType::Info,
MetricType::StateSet,
MetricType::Unknown,
] {
assert_eq!(MetricType::parse(ty.as_str()), ty);
}
}
#[test]
fn metric_type_parse_unknown_keyword_yields_unknown() {
assert_eq!(MetricType::parse("garbage"), MetricType::Unknown);
assert_eq!(MetricType::parse(""), MetricType::Unknown);
}
#[test]
fn metric_type_implied_labels_match_openmetrics() {
assert_eq!(MetricType::Histogram.implied_label(), Some("le"));
assert_eq!(MetricType::GaugeHistogram.implied_label(), Some("le"));
assert_eq!(MetricType::Summary.implied_label(), Some("quantile"));
assert_eq!(MetricType::Counter.implied_label(), None);
assert_eq!(MetricType::Gauge.implied_label(), None);
}
#[test]
fn metric_type_has_derived_series_only_for_histograms_and_summary() {
assert!(MetricType::Histogram.has_derived_series());
assert!(MetricType::GaugeHistogram.has_derived_series());
assert!(MetricType::Summary.has_derived_series());
assert!(!MetricType::Counter.has_derived_series());
assert!(!MetricType::Gauge.has_derived_series());
assert!(!MetricType::Info.has_derived_series());
}
#[test]
fn cache_serves_repeat_calls_from_cache() {
let inner = Arc::new(MockCatalog::new());
let cache = CachedCatalog::new(inner.clone());
let _ = cache.metric_families().unwrap();
let _ = cache.metric_families().unwrap();
let _ = cache.metric_families().unwrap();
assert_eq!(
inner.family_calls.load(Ordering::SeqCst),
1,
"cache should have served repeated calls"
);
}
#[test]
fn cache_invalidate_forces_refetch() {
let inner = Arc::new(MockCatalog::new());
let cache = CachedCatalog::new(inner.clone());
let _ = cache.metric_families().unwrap();
cache.invalidate();
let _ = cache.metric_families().unwrap();
assert_eq!(
inner.family_calls.load(Ordering::SeqCst),
2,
"invalidate should have forced a refetch"
);
}
#[test]
fn cache_zero_ttl_still_caches_via_generation() {
let inner = Arc::new(MockCatalog::new());
let cache = CachedCatalog::new(inner.clone()).with_ttl(Duration::ZERO);
let _ = cache.metric_families().unwrap();
let _ = cache.metric_families().unwrap();
assert_eq!(
inner.family_calls.load(Ordering::SeqCst),
1,
"TTL=0 alone still caches via generation"
);
}
#[test]
fn cache_distinguishes_label_keys_by_family_filter() {
let inner = Arc::new(MockCatalog::new());
let cache = CachedCatalog::new(inner.clone());
let _ = cache.label_keys(None).unwrap();
let _ = cache.label_keys(Some("ops_total")).unwrap();
let _ = cache.label_keys(None).unwrap();
let _ = cache.label_keys(Some("ops_total")).unwrap();
assert_eq!(
inner.key_calls.load(Ordering::SeqCst),
2,
"different family filters should be cached separately"
);
}
#[test]
fn cache_distinguishes_label_values_by_key_and_filter() {
let inner = Arc::new(MockCatalog::new());
let cache = CachedCatalog::new(inner.clone());
let _ = cache.label_values("phase", None).unwrap();
let _ = cache.label_values("phase", Some("ops_total")).unwrap();
let _ = cache.label_values("phase", None).unwrap();
assert_eq!(inner.value_calls.load(Ordering::SeqCst), 2);
}
#[test]
fn cache_series_keyed_by_matchers_independent_of_order() {
let inner = Arc::new(MockCatalog::new());
let cache = CachedCatalog::new(inner.clone());
let m1 = vec![
Matcher {
label: "a".into(),
op: MatcherOp::Eq,
value: "1".into(),
},
Matcher {
label: "b".into(),
op: MatcherOp::Eq,
value: "2".into(),
},
];
let m2 = vec![
Matcher {
label: "b".into(),
op: MatcherOp::Eq,
value: "2".into(),
},
Matcher {
label: "a".into(),
op: MatcherOp::Eq,
value: "1".into(),
},
];
let _ = cache.series(&m1).unwrap();
let _ = cache.series(&m2).unwrap();
assert_eq!(
inner.series_calls.load(Ordering::SeqCst),
1,
"matcher-order shouldn't fragment the cache"
);
}
#[test]
fn cache_mtime_hook_invalidates_on_advance() {
let inner = Arc::new(MockCatalog::new());
let mtime = Arc::new(Mutex::new(Instant::now()));
let mtime_clone = mtime.clone();
let cache = CachedCatalog::new(inner.clone())
.with_mtime_fn(move || Some(*mtime_clone.lock().unwrap()));
let _ = cache.metric_families().unwrap();
let _ = cache.metric_families().unwrap();
assert_eq!(inner.family_calls.load(Ordering::SeqCst), 1);
*mtime.lock().unwrap() += Duration::from_millis(1);
let _ = cache.metric_families().unwrap();
assert_eq!(
inner.family_calls.load(Ordering::SeqCst),
2,
"advanced mtime should have invalidated the cache"
);
}
#[test]
fn encode_matchers_sorts_by_label() {
let m1 = vec![
Matcher {
label: "a".into(),
op: MatcherOp::Eq,
value: "1".into(),
},
Matcher {
label: "b".into(),
op: MatcherOp::EqRegex,
value: ".*".into(),
},
];
let m2 = vec![
Matcher {
label: "b".into(),
op: MatcherOp::EqRegex,
value: ".*".into(),
},
Matcher {
label: "a".into(),
op: MatcherOp::Eq,
value: "1".into(),
},
];
assert_eq!(encode_matchers(&m1), encode_matchers(&m2));
}
}