pub mod catalog;
mod hybrid;
mod live;
mod shapes;
#[cfg(feature = "sqlite")]
pub mod sqlite;
pub use catalog::{
CachedCatalog, ExemplarPoint, LabelSet, MetricCatalog, MetricFamilyMeta, MetricType,
};
pub use hybrid::{HorizonAware, HybridStore, Tier};
pub use live::MetricsQueryAccess;
pub use shapes::{MatchOp, Matcher, Sample, Series, Vector};
use std::sync::{Arc, LazyLock, OnceLock};
use arc_swap::ArcSwapOption;
#[derive(Debug, Clone)]
pub struct QueryError {
pub message: String,
}
impl QueryError {
pub fn new(msg: impl Into<String>) -> Self {
Self {
message: msg.into(),
}
}
}
impl std::fmt::Display for QueryError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "metrics query: {}", self.message)
}
}
impl std::error::Error for QueryError {}
pub trait MetricAccess: Send + Sync {
fn select_range(
&self,
matchers: &[Matcher],
start_ms: i64,
end_ms: i64,
) -> Result<Vector, QueryError>;
fn select_instant(
&self,
matchers: &[Matcher],
at_ms: i64,
lookback_ms: Option<i64>,
) -> Result<Vector, QueryError> {
let start_ms = at_ms - lookback_ms.unwrap_or(0);
let range = self.select_range(matchers, start_ms, at_ms)?;
let reduced = range
.into_series()
.into_iter()
.filter_map(|s| {
s.samples.last().copied().map(|last| Series {
labels: s.labels,
samples: vec![last],
})
})
.collect();
Ok(Vector::new(reduced))
}
}
struct LiveHolder(Arc<dyn MetricAccess>);
static LIVE: LazyLock<ArcSwapOption<LiveHolder>> = LazyLock::new(ArcSwapOption::empty);
pub fn install_live_access(service: Arc<dyn MetricAccess>) {
LIVE.store(Some(Arc::new(LiveHolder(service))));
}
pub fn uninstall_live_access() {
drop(LIVE.swap(None));
}
pub fn live_access() -> Option<Arc<dyn MetricAccess>> {
LIVE.load_full().map(|h| h.0.clone())
}
static READ_EXEC_ID_HOOK: OnceLock<fn() -> Option<u64>> = OnceLock::new();
pub fn install_read_exec_id_hook(hook: fn() -> Option<u64>) {
let _ = READ_EXEC_ID_HOOK.set(hook);
}
pub fn current_read_exec_id() -> Option<u64> {
READ_EXEC_ID_HOOK.get().and_then(|h| h())
}
pub struct ExecScopedAccess {
inner: Arc<dyn MetricAccess>,
}
impl ExecScopedAccess {
pub fn new(inner: Arc<dyn MetricAccess>) -> Self {
Self { inner }
}
fn scoped(&self, matchers: &[Matcher]) -> Option<Vec<Matcher>> {
let id = current_read_exec_id()?;
if matchers.iter().any(|m| m.label == "exec_id") {
return None; }
let mut v = matchers.to_vec();
v.push(Matcher::eq("exec_id", &id.to_string()));
Some(v)
}
}
impl MetricAccess for ExecScopedAccess {
fn select_range(
&self,
matchers: &[Matcher],
start_ms: i64,
end_ms: i64,
) -> Result<Vector, QueryError> {
match self.scoped(matchers) {
Some(scoped) => self.inner.select_range(&scoped, start_ms, end_ms),
None => self.inner.select_range(matchers, start_ms, end_ms),
}
}
fn select_instant(
&self,
matchers: &[Matcher],
at_ms: i64,
lookback_ms: Option<i64>,
) -> Result<Vector, QueryError> {
match self.scoped(matchers) {
Some(scoped) => self.inner.select_instant(&scoped, at_ms, lookback_ms),
None => self.inner.select_instant(matchers, at_ms, lookback_ms),
}
}
}
pub struct AccessProvider {
pub scheme: &'static str,
pub open: fn(target: &str) -> Result<Box<dyn MetricAccess>, QueryError>,
}
inventory::collect!(AccessProvider);
pub fn provider(scheme: &str) -> Option<&'static AccessProvider> {
inventory::iter::<AccessProvider>
.into_iter()
.find(|p| p.scheme == scheme)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn live_access_installs_and_reads_back() {
struct Stub;
impl MetricAccess for Stub {
fn select_instant(
&self,
_: &[Matcher],
_: i64,
_: Option<i64>,
) -> Result<Vector, QueryError> {
Ok(Vector::default())
}
fn select_range(&self, _: &[Matcher], _: i64, _: i64) -> Result<Vector, QueryError> {
Ok(Vector::default())
}
}
install_live_access(Arc::new(Stub));
assert!(live_access().is_some());
}
#[test]
fn unknown_provider_scheme_is_none() {
assert!(provider("no-such-scheme").is_none());
}
}