use super::catalog::{ExemplarPoint, LabelSet, MetricCatalog, MetricFamilyMeta, MetricType};
use super::{
MatchOp as MatcherOp, Matcher, MetricAccess, QueryError as DataSourceError, Sample, Series,
Vector,
};
use rusqlite::{Connection, OptionalExtension, params_from_iter, types::Value};
use std::path::PathBuf;
use std::sync::Mutex;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum ExecutionSelection {
#[default]
LatestPerInstance,
All,
Latest,
Specific(u64),
}
pub struct SqliteDataSource {
conn: Mutex<Connection>,
db_path: Option<PathBuf>,
selection: ExecutionSelection,
}
impl SqliteDataSource {
pub fn open(path: impl AsRef<std::path::Path>) -> Result<Self, DataSourceError> {
let path_ref = path.as_ref();
let conn = Connection::open_with_flags(
path_ref,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
)
.map_err(|e| DataSourceError::new(format!("open metrics.db: {e}")))?;
let mut src = Self::from_connection(conn)?;
src.db_path = Some(path_ref.to_path_buf());
Ok(src)
}
pub fn db_path(&self) -> Option<&std::path::Path> {
self.db_path.as_deref()
}
pub fn mtime_fn(
&self,
) -> Option<impl Fn() -> Option<std::time::Instant> + Send + Sync + 'static> {
let path = self.db_path.clone()?;
let anchor_instant = std::time::Instant::now();
let anchor_system = std::time::SystemTime::now();
Some(move || -> Option<std::time::Instant> {
let meta = std::fs::metadata(&path).ok()?;
let mtime = meta.modified().ok()?;
let delta = mtime.duration_since(anchor_system).ok()?;
anchor_instant.checked_add(delta)
})
}
pub fn from_connection(conn: Connection) -> Result<Self, DataSourceError> {
conn.execute_batch(
"PRAGMA cache_size = -65536;\
PRAGMA temp_store = MEMORY;\
PRAGMA mmap_size = 268435456;",
)
.map_err(|e| DataSourceError::new(format!("apply pragmas: {e}")))?;
register_regexp(&conn)?;
Ok(Self {
conn: Mutex::new(conn),
db_path: None,
selection: ExecutionSelection::default(),
})
}
pub fn with_execution_selection(mut self, selection: ExecutionSelection) -> Self {
self.selection = selection;
self
}
}
fn retain_latest_per_instance(series: Vec<Series>) -> Vec<Series> {
use std::collections::{HashMap, HashSet};
let mut best: HashMap<Vec<(String, String)>, (i64, usize)> = HashMap::new();
for (i, s) in series.iter().enumerate() {
let exec = s
.labels
.iter()
.find(|(k, _)| k == "exec_id")
.and_then(|(_, v)| v.parse::<i64>().ok())
.unwrap_or(0);
let mut logical: Vec<(String, String)> = s
.labels
.iter()
.filter(|(k, _)| k != "exec_id" && k != "session")
.cloned()
.collect();
logical.sort();
match best.get(&logical) {
Some(&(e, _)) if e >= exec => {}
_ => {
best.insert(logical, (exec, i));
}
}
}
let keep: HashSet<usize> = best.values().map(|(_, i)| *i).collect();
series
.into_iter()
.enumerate()
.filter_map(|(i, s)| keep.contains(&i).then_some(s))
.collect()
}
fn register_regexp(conn: &Connection) -> Result<(), DataSourceError> {
use rusqlite::functions::FunctionFlags;
use std::collections::HashMap;
use std::sync::Mutex as StdMutex;
let cache: StdMutex<HashMap<String, regex::Regex>> = StdMutex::new(HashMap::new());
conn.create_scalar_function(
"regexp",
2,
FunctionFlags::SQLITE_DETERMINISTIC | FunctionFlags::SQLITE_UTF8,
move |ctx| {
let pattern: String = ctx.get(0)?;
let value: String = ctx.get(1)?;
let mut guard = cache.lock().unwrap_or_else(|e| e.into_inner());
let re = match guard.get(&pattern) {
Some(r) => r.clone(),
None => {
let anchored = format!("^(?:{pattern})$");
let r = regex::Regex::new(&anchored).map_err(|e| {
rusqlite::Error::UserFunctionError(
format!("regexp pattern '{pattern}': {e}").into(),
)
})?;
guard.insert(pattern, r.clone());
r
}
};
Ok(re.is_match(&value))
},
)
.map_err(|e| DataSourceError::new(format!("register REGEXP function: {e}")))?;
Ok(())
}
impl MetricAccess for SqliteDataSource {
fn select_range(
&self,
matchers: &[Matcher],
start_ms: i64,
end_ms: i64,
) -> Result<Vector, DataSourceError> {
let conn = self
.conn
.lock()
.map_err(|_| DataSourceError::new("sqlite mutex poisoned"))?;
let Some(name_matcher) = matchers.iter().find(|m| m.label == "__name__") else {
return Ok(Vector::default());
};
let resolved = match name_matcher.op {
MatcherOp::Eq => resolve_family(&conn, &name_matcher.value)?,
_ => {
return Err(DataSourceError::new(
"non-Eq match on __name__ not supported by sqlite adapter yet",
));
}
};
let Some(resolved) = resolved else {
return Ok(Vector::default());
};
let other_matchers: Vec<&Matcher> = matchers
.iter()
.filter(|m| m.label != "__name__" && m.label != "exec_id")
.collect();
let label_filter = instance_label_filter_clause(&other_matchers)?;
let stat_expr = resolved.stat_expr;
let interval_proj = if resolved.is_rate {
", sv.interval_ms"
} else {
""
};
let exec_label_filter = matchers
.iter()
.find(|m| m.label == "exec_id" && m.op == MatcherOp::Eq)
.and_then(|m| m.value.parse::<i64>().ok())
.map(|n| format!("AND mi.exec_id = {n} "))
.unwrap_or_default();
let exec_filter = match self.selection {
ExecutionSelection::Specific(n) => format!("AND mi.exec_id = {n} "),
ExecutionSelection::Latest => "AND mi.exec_id = \
(SELECT MAX(exec_id) FROM metric_instance WHERE family_id = ?1) "
.to_string(),
ExecutionSelection::All | ExecutionSelection::LatestPerInstance => String::new(),
};
let exec_filter = format!("{exec_filter}{exec_label_filter}");
let sql = format!(
"SELECT mi.id, sv.timestamp_ms, {stat_expr}{interval_proj} \
FROM metric_instance mi \
JOIN sample_value sv ON sv.instance_id = mi.id \
WHERE mi.family_id = ?1 \
AND sv.timestamp_ms >= ?2 AND sv.timestamp_ms <= ?3 \
{label_filter} {exec_filter}\
ORDER BY mi.id, sv.timestamp_ms"
);
let mut params: Vec<Value> = vec![
Value::Integer(resolved.family_id),
Value::Integer(start_ms),
Value::Integer(end_ms),
];
for m in &other_matchers {
params.push(Value::Text(m.label.clone()));
params.push(Value::Text(m.value.clone()));
}
let mut stmt = conn
.prepare(&sql)
.map_err(|e| DataSourceError::new(format!("prepare fetch: {e}")))?;
let mut rows = stmt
.query(params_from_iter(params.iter()))
.map_err(|e| DataSourceError::new(format!("query fetch: {e}")))?;
let mut out: Vec<Series> = Vec::new();
let mut current_instance_id: Option<i64> = None;
let mut current_samples: Vec<Sample> = Vec::new();
let mut min_rate_interval_ms: Option<i64> = None;
while let Some(row) = rows
.next()
.map_err(|e| DataSourceError::new(format!("step fetch: {e}")))?
{
let instance_id: i64 = row
.get(0)
.map_err(|e| DataSourceError::new(format!("row.get(0): {e}")))?;
let timestamp_ms: i64 = row
.get(1)
.map_err(|e| DataSourceError::new(format!("row.get(1): {e}")))?;
let value: f64 = row
.get::<_, Option<f64>>(2)
.map_err(|e| DataSourceError::new(format!("row.get(2): {e}")))?
.unwrap_or(f64::NAN);
if resolved.is_rate {
let iv: i64 = row
.get(3)
.map_err(|e| DataSourceError::new(format!("row.get(3): {e}")))?;
if iv > 0 && iv < 1000 {
min_rate_interval_ms =
Some(min_rate_interval_ms.map(|m| m.min(iv)).unwrap_or(iv));
}
}
if Some(instance_id) != current_instance_id {
if let Some(prev) = current_instance_id.take() {
out.push(materialize_series(
&conn,
prev,
&resolved.virtual_name,
std::mem::take(&mut current_samples),
)?);
}
current_instance_id = Some(instance_id);
}
current_samples.push(Sample {
timestamp_ms,
value,
});
}
if let Some(iv) = min_rate_interval_ms {
eprintln!(
"warning: `{}` evaluated over samples with sub-1s windows \
(shortest seen: {iv}ms). Intervals are precisely measured \
(ms-resolution), so the numeric rate is honest — but a short \
window samples a brief slice of phase activity and is more \
susceptible to instantaneous noise (warmup, GC pause, single \
outlier op). For a steady-state view, prefer rate({}[30s]) or \
ensure the phase runs long enough to span ≥ 1 cadence window.",
resolved.virtual_name,
resolved.virtual_name.trim_end_matches("_rate"),
);
}
if let Some(last) = current_instance_id {
out.push(materialize_series(
&conn,
last,
&resolved.virtual_name,
current_samples,
)?);
}
if self.selection == ExecutionSelection::LatestPerInstance {
out = retain_latest_per_instance(out);
}
Ok(Vector::new(out))
}
}
impl MetricCatalog for SqliteDataSource {
fn metric_families(&self) -> Result<Vec<MetricFamilyMeta>, DataSourceError> {
let conn = self
.conn
.lock()
.map_err(|_| DataSourceError::new("sqlite mutex poisoned"))?;
let mut stmt = conn
.prepare("SELECT name, type, unit, help FROM metric_family ORDER BY name")
.map_err(|e| DataSourceError::new(format!("prepare families: {e}")))?;
let rows = stmt
.query_map([], |r| {
Ok((
r.get::<_, String>(0)?,
r.get::<_, String>(1)?,
r.get::<_, Option<String>>(2)?,
r.get::<_, Option<String>>(3)?,
))
})
.map_err(|e| DataSourceError::new(format!("query families: {e}")))?;
let mut out = Vec::new();
for row in rows {
let (name, ty_str, unit, help) =
row.map_err(|e| DataSourceError::new(format!("decode family row: {e}")))?;
out.push(MetricFamilyMeta {
name,
ty: MetricType::parse(&ty_str),
unit,
help,
});
}
Ok(out)
}
fn label_keys(&self, family_filter: Option<&str>) -> Result<Vec<String>, DataSourceError> {
let conn = self
.conn
.lock()
.map_err(|_| DataSourceError::new("sqlite mutex poisoned"))?;
let mut keys: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
let sql = match family_filter {
Some(_) => {
"SELECT DISTINCT il.key \
FROM instance_label il \
JOIN metric_instance mi ON mi.id = il.instance_id \
JOIN metric_family mf ON mf.id = mi.family_id \
WHERE mf.name = ?1 AND il.key != '__name__' \
ORDER BY il.key"
}
None => {
"SELECT DISTINCT key FROM instance_label \
WHERE key != '__name__' \
ORDER BY key"
}
};
let mut stmt = conn
.prepare(sql)
.map_err(|e| DataSourceError::new(format!("prepare label_keys: {e}")))?;
let mut rows: Box<dyn Iterator<Item = rusqlite::Result<String>>> = match family_filter {
Some(name) => Box::new(
stmt.query_map([name], |r| r.get::<_, String>(0))
.map_err(|e| DataSourceError::new(format!("query label_keys: {e}")))?
.collect::<Vec<_>>()
.into_iter(),
),
None => Box::new(
stmt.query_map([], |r| r.get::<_, String>(0))
.map_err(|e| DataSourceError::new(format!("query label_keys: {e}")))?
.collect::<Vec<_>>()
.into_iter(),
),
};
for row in &mut rows {
let k = row.map_err(|e| DataSourceError::new(format!("decode label_key: {e}")))?;
keys.insert(k);
}
Ok(keys.into_iter().collect())
}
fn label_values(
&self,
key: &str,
family_filter: Option<&str>,
) -> Result<Vec<String>, DataSourceError> {
let conn = self
.conn
.lock()
.map_err(|_| DataSourceError::new("sqlite mutex poisoned"))?;
let sql = match family_filter {
Some(_) => {
"SELECT DISTINCT il.value \
FROM instance_label il \
JOIN metric_instance mi ON mi.id = il.instance_id \
JOIN metric_family mf ON mf.id = mi.family_id \
WHERE il.key = ?1 AND mf.name = ?2 \
ORDER BY il.value"
}
None => {
"SELECT DISTINCT value FROM instance_label \
WHERE key = ?1 \
ORDER BY value"
}
};
let mut stmt = conn
.prepare(sql)
.map_err(|e| DataSourceError::new(format!("prepare label_values: {e}")))?;
let mut out = Vec::new();
let rows: Vec<rusqlite::Result<String>> = match family_filter {
Some(name) => stmt
.query_map([key, name], |r| r.get::<_, String>(0))
.map_err(|e| DataSourceError::new(format!("query label_values: {e}")))?
.collect(),
None => stmt
.query_map([key], |r| r.get::<_, String>(0))
.map_err(|e| DataSourceError::new(format!("query label_values: {e}")))?
.collect(),
};
for row in rows {
out.push(row.map_err(|e| DataSourceError::new(format!("decode value: {e}")))?);
}
Ok(out)
}
fn series(&self, matchers: &[Matcher]) -> Result<Vec<LabelSet>, DataSourceError> {
let conn = self
.conn
.lock()
.map_err(|_| DataSourceError::new("sqlite mutex poisoned"))?;
let name_matcher = matchers.iter().find(|m| m.label == "__name__");
let resolved = match name_matcher.map(|m| m.op) {
Some(MatcherOp::Eq) => resolve_family(&conn, &name_matcher.unwrap().value)?,
Some(_) => {
return Err(DataSourceError::new(
"non-Eq match on __name__ not supported by sqlite catalog yet",
));
}
None => None,
};
let other_matchers: Vec<&Matcher> =
matchers.iter().filter(|m| m.label != "__name__").collect();
let label_filter = instance_label_filter_clause(&other_matchers)?;
let sql_with_family = format!(
"SELECT mi.id FROM metric_instance mi \
WHERE mi.family_id = ?1 \
{label_filter} \
ORDER BY mi.id"
);
let sql_no_family = format!(
"SELECT mi.id FROM metric_instance mi \
WHERE 1=1 \
{label_filter} \
ORDER BY mi.id"
);
let sql = if resolved.is_some() {
sql_with_family
} else {
sql_no_family
};
let mut params: Vec<Value> = Vec::new();
if let Some(r) = &resolved {
params.push(Value::Integer(r.family_id));
}
for m in &other_matchers {
params.push(Value::Text(m.label.clone()));
params.push(Value::Text(m.value.clone()));
}
let mut stmt = conn
.prepare(&sql)
.map_err(|e| DataSourceError::new(format!("prepare series: {e}")))?;
let rows = stmt
.query_map(params_from_iter(params.iter()), |r| r.get::<_, i64>(0))
.map_err(|e| DataSourceError::new(format!("query series: {e}")))?;
let mut out = Vec::new();
for row in rows {
let instance_id =
row.map_err(|e| DataSourceError::new(format!("decode series row: {e}")))?;
let mut labels = materialize_instance_labels(&conn, instance_id)?;
if let Some(pos) = labels.iter().position(|(k, _)| k == "__name__")
&& pos != 0
{
let pair = labels.remove(pos);
labels.insert(0, pair);
}
out.push(labels);
}
Ok(out)
}
fn exemplars(
&self,
matchers: &[Matcher],
time_range: Option<(i64, i64)>,
) -> Result<Vec<ExemplarPoint>, DataSourceError> {
let conn = self
.conn
.lock()
.map_err(|_| DataSourceError::new("sqlite mutex poisoned"))?;
let name_matcher = matchers.iter().find(|m| m.label == "__name__");
let resolved = match name_matcher.map(|m| m.op) {
Some(MatcherOp::Eq) => resolve_family(&conn, &name_matcher.unwrap().value)?,
Some(_) => {
return Err(DataSourceError::new(
"non-Eq match on __name__ not supported by sqlite catalog yet",
));
}
None => None,
};
let other_matchers: Vec<&Matcher> =
matchers.iter().filter(|m| m.label != "__name__").collect();
let label_filter = instance_label_filter_clause(&other_matchers)?;
let (start_ms, end_ms) = time_range.unwrap_or((i64::MIN, i64::MAX));
let sql_with_family = format!(
"SELECT mi.id, \
e.sample_timestamp_ms, e.value, \
e.timestamp_ms, e.labels_spec \
FROM exemplar e \
JOIN metric_instance mi ON mi.id = e.instance_id \
WHERE mi.family_id = ?1 \
AND e.sample_timestamp_ms >= ?2 \
AND e.sample_timestamp_ms <= ?3 \
{label_filter} \
ORDER BY e.sample_timestamp_ms"
);
let sql_no_family = format!(
"SELECT mi.id, \
e.sample_timestamp_ms, e.value, \
e.timestamp_ms, e.labels_spec \
FROM exemplar e \
JOIN metric_instance mi ON mi.id = e.instance_id \
WHERE e.sample_timestamp_ms >= ?1 \
AND e.sample_timestamp_ms <= ?2 \
{label_filter} \
ORDER BY e.sample_timestamp_ms"
);
let mut params: Vec<Value> = Vec::new();
let sql = match &resolved {
Some(r) => {
params.push(Value::Integer(r.family_id));
params.push(Value::Integer(start_ms));
params.push(Value::Integer(end_ms));
sql_with_family
}
None => {
params.push(Value::Integer(start_ms));
params.push(Value::Integer(end_ms));
sql_no_family
}
};
for m in &other_matchers {
params.push(Value::Text(m.label.clone()));
params.push(Value::Text(m.value.clone()));
}
let mut stmt = conn
.prepare(&sql)
.map_err(|e| DataSourceError::new(format!("prepare exemplars: {e}")))?;
let rows = stmt
.query_map(params_from_iter(params.iter()), |r| {
Ok((
r.get::<_, i64>(0)?, r.get::<_, i64>(1)?, r.get::<_, f64>(2)?, r.get::<_, Option<i64>>(3)?, r.get::<_, String>(4)?, ))
})
.map_err(|e| DataSourceError::new(format!("query exemplars: {e}")))?;
let mut out = Vec::new();
for row in rows {
let (instance_id, sample_ts, value, ts, labels_spec) =
row.map_err(|e| DataSourceError::new(format!("decode exemplar: {e}")))?;
let mut series = materialize_instance_labels(&conn, instance_id)?;
if let Some(pos) = series.iter().position(|(k, _)| k == "__name__")
&& pos != 0
{
let pair = series.remove(pos);
series.insert(0, pair);
}
let labels = parse_labels_spec(&labels_spec);
out.push(ExemplarPoint {
series,
sample_timestamp_ms: sample_ts,
value,
timestamp_ms: ts,
labels,
});
}
Ok(out)
}
}
pub fn parse_labels_spec(spec: &str) -> Vec<(String, String)> {
let s = spec.trim();
if s.is_empty() {
return Vec::new();
}
let mut out = Vec::new();
let mut cur_key = String::new();
let mut cur_val = String::new();
let bytes = s.as_bytes();
let mut i = 0;
while i < bytes.len() {
while i < bytes.len() && bytes[i] != b'=' {
cur_key.push(bytes[i] as char);
i += 1;
}
if i >= bytes.len() {
break;
}
i += 1; let quoted = i < bytes.len() && bytes[i] == b'"';
if quoted {
i += 1;
}
while i < bytes.len() {
if quoted {
if bytes[i] == b'"' {
i += 1;
break;
}
} else if bytes[i] == b',' {
break;
}
cur_val.push(bytes[i] as char);
i += 1;
}
out.push((cur_key.trim().to_string(), cur_val.clone()));
cur_key.clear();
cur_val.clear();
while i < bytes.len() && (bytes[i] == b',' || bytes[i].is_ascii_whitespace()) {
i += 1;
}
}
out
}
fn materialize_instance_labels(
conn: &Connection,
instance_id: i64,
) -> Result<Vec<(String, String)>, DataSourceError> {
let mut stmt = conn
.prepare_cached(
"SELECT key, value FROM instance_label \
WHERE instance_id = ?1 \
ORDER BY key",
)
.map_err(|e| DataSourceError::new(format!("prepare label set: {e}")))?;
let rows = stmt
.query_map([instance_id], |r| {
Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))
})
.map_err(|e| DataSourceError::new(format!("query label set: {e}")))?;
let mut out = Vec::new();
for row in rows {
out.push(row.map_err(|e| DataSourceError::new(format!("decode label entry: {e}")))?);
}
Ok(out)
}
struct ResolvedName {
family_id: i64,
virtual_name: String,
stat_expr: &'static str,
is_rate: bool,
}
fn resolve_family(conn: &Connection, name: &str) -> Result<Option<ResolvedName>, DataSourceError> {
if let Some((family_id, family_type)) = lookup_family(conn, name)? {
let stat_expr = default_column_for_type(&family_type);
return Ok(Some(ResolvedName {
family_id,
virtual_name: name.to_string(),
stat_expr,
is_rate: false,
}));
}
for suffix in STAT_SUFFIXES {
if let Some(stripped) = name.strip_suffix(suffix.text)
&& let Some((family_id, family_type)) = lookup_family(conn, stripped)?
&& suffix.applies_to(&family_type)
{
return Ok(Some(ResolvedName {
family_id,
virtual_name: name.to_string(),
stat_expr: suffix.expr,
is_rate: suffix.text == "_rate",
}));
}
}
Ok(None)
}
fn lookup_family(conn: &Connection, name: &str) -> Result<Option<(i64, String)>, DataSourceError> {
conn.query_row(
"SELECT id, type FROM metric_family WHERE name = ?1",
rusqlite::params![name],
|row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)),
)
.optional()
.map_err(|e| DataSourceError::new(format!("family lookup: {e}")))
}
pub fn default_column_for_type(family_type: &str) -> &'static str {
match family_type {
"counter" => "sv.count",
"gauge" => "sv.mean",
"summary" => "sv.count",
"histogram" => "sv.count",
"gaugehistogram" => "sv.count",
"info" => "sv.count",
"stateset" => "sv.mean",
"unknown" => "sv.mean",
_ => "sv.mean", }
}
struct StatSuffix {
text: &'static str,
expr: &'static str,
applies_to_fn: fn(&str) -> bool,
}
impl StatSuffix {
fn applies_to(&self, family_type: &str) -> bool {
(self.applies_to_fn)(family_type)
}
}
fn applies_summary(t: &str) -> bool {
matches!(t, "summary" | "histogram" | "gaugehistogram")
}
fn applies_counted(t: &str) -> bool {
matches!(
t,
"counter" | "summary" | "histogram" | "gaugehistogram" | "info"
)
}
fn applies_any(_t: &str) -> bool {
true
}
const STAT_SUFFIXES: &[StatSuffix] = &[
StatSuffix {
text: "_p999",
expr: "sv.p999",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_p99",
expr: "sv.p99",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_p98",
expr: "sv.p98",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_p95",
expr: "sv.p95",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_p90",
expr: "sv.p90",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_p75",
expr: "sv.p75",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_p50",
expr: "sv.p50",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_count",
expr: "sv.count",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_sum",
expr: "sv.sum",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_min",
expr: "sv.min",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_max",
expr: "sv.max",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_mean",
expr: "sv.mean",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_stddev",
expr: "sv.stddev",
applies_to_fn: applies_summary,
},
StatSuffix {
text: "_rate",
expr: "((CAST(sv.count AS REAL) \
- COALESCE(LAG(CAST(sv.count AS REAL)) \
OVER (PARTITION BY sv.instance_id ORDER BY sv.timestamp_ms), 0)) \
* 1000.0 / NULLIF(sv.interval_ms, 0))",
applies_to_fn: applies_counted,
},
StatSuffix {
text: "_interval_ns",
expr: "(CAST(sv.interval_ms AS REAL) * 1000000.0)",
applies_to_fn: applies_any,
},
];
fn instance_label_filter_clause(matchers: &[&Matcher]) -> Result<String, DataSourceError> {
if matchers.is_empty() {
return Ok(String::new());
}
let mut parts: Vec<String> = Vec::with_capacity(matchers.len());
for (i, m) in matchers.iter().enumerate() {
let kparam = i * 2 + 4; let vparam = i * 2 + 5;
let cmp_clause = match m.op {
MatcherOp::Eq => format!("il.key = ?{kparam} AND il.value = ?{vparam}"),
MatcherOp::Ne => format!("il.key = ?{kparam} AND il.value != ?{vparam}"),
MatcherOp::EqRegex => format!("il.key = ?{kparam} AND il.value REGEXP ?{vparam}"),
MatcherOp::NeRegex => format!("il.key = ?{kparam} AND NOT (il.value REGEXP ?{vparam})"),
};
parts.push(format!(
"SELECT il.instance_id FROM instance_label il WHERE {cmp_clause}"
));
}
Ok(format!(" AND mi.id IN ({})", parts.join(" INTERSECT ")))
}
fn materialize_series(
conn: &Connection,
instance_id: i64,
virtual_name: &str,
samples: Vec<Sample>,
) -> Result<Series, DataSourceError> {
let mut labels = materialize_instance_labels(conn, instance_id)?;
labels.retain(|(k, _)| k != "__name__");
labels.insert(0, ("__name__".to_string(), virtual_name.to_string()));
Ok(Series { labels, samples })
}
inventory::submit! {
super::AccessProvider {
scheme: "sqlite",
open: |target| {
SqliteDataSource::open(target)
.map(|s| Box::new(s) as Box<dyn super::MetricAccess>)
},
}
}