use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::path::Path;
use std::sync::Arc;
use arrow_array::cast::AsArray;
use arrow_array::types::{
Float64Type, Int64Type, TimestampNanosecondType, UInt8Type, UInt16Type, UInt32Type, UInt64Type,
};
use arrow_array::{Array, RecordBatch};
use crate::block::{self, Src};
use crate::error::Result;
use crate::json::Json;
use crate::query::{Results, Stats, Term, attr_key, attr_parents, dict_index, emit_attr};
use crate::schema::MetricKind;
use crate::signal::Open;
#[derive(Debug, Clone)]
pub struct SeriesQuery {
pub name: Option<String>,
pub from: i64,
pub to: i64,
pub terms: Vec<Term>,
pub max_series: usize,
pub max_points: usize,
}
#[derive(Debug, Clone, Copy)]
enum Pt {
Int(i64),
Double(f64),
}
struct Series {
desc: String,
attrs: String,
points: Vec<(i64, Pt)>,
dropped: usize,
exemplars: Vec<(i64, String)>,
}
impl Series {
fn compact(&mut self, max: usize) {
self.points.sort_unstable_by_key(|p| std::cmp::Reverse(p.0));
self.dropped += self.points.len() - max;
self.points.truncate(max);
}
}
const MAX_EXEMPLARS: usize = 64;
type Attr = (String, String);
type Prefix = (String, Vec<Attr>);
const DP_TABLES: [&str; 4] = ["number_dp", "hist_dp", "exp_hist_dp", "summary_dp"];
const COUNT: &str = ".count";
const SUM: &str = ".sum";
const DERIVED: [&str; 2] = [COUNT, SUM];
pub fn series(root: &Path, q: &SeriesQuery) -> Result<Results> {
series_open(root, q, &[])
}
pub fn series_open(root: &Path, q: &SeriesQuery, open: &[Arc<Open>]) -> Result<Results> {
let disk = block::scan(root, "metrics")?;
let refs = block::sources(&disk, open);
let mut stats = Stats {
blocks_total: refs.len(),
..Default::default()
};
let mut out: BTreeMap<String, Series> = BTreeMap::new();
let mut dropped: BTreeSet<String> = BTreeSet::new();
for bref in &refs {
if bref.max_ts < q.from || bref.min_ts > q.to {
continue;
}
stats.blocks_scanned += 1;
collect_block(bref, q, &mut out, &mut dropped, &mut stats)?;
}
stats.dropped_series = dropped.len();
for s in out.values_mut() {
if s.points.len() > q.max_points {
s.compact(q.max_points);
}
s.points.sort_unstable_by_key(|p| p.0);
}
let mut j = Json::new();
j.arr(|j| {
for s in out.values() {
let pts = &s.points;
j.obj(|j| {
j.raw(&s.desc);
j.key("attributes");
j.raw(&s.attrs);
if s.dropped > 0 {
j.key("dropped_points");
j.u64(s.dropped as u64);
}
j.key("points");
j.arr(|j| {
for &(ts, v) in pts {
j.arr(|j| {
j.i64_str(ts);
match v {
Pt::Int(i) => j.i64_str(i),
Pt::Double(d) => j.f64(d),
}
});
}
});
if !s.exemplars.is_empty() {
let mut ex = s.exemplars.clone();
ex.sort_unstable_by_key(|e| e.0);
j.key("exemplars");
j.arr(|j| {
for (_, rendered) in &ex {
j.raw(rendered);
}
});
}
});
}
});
Ok(Results {
json: j.into_string(),
stats,
next: None,
})
}
pub fn names(root: &Path, from: i64, to: i64) -> Result<Results> {
names_open(root, from, to, &[])
}
pub fn names_open(root: &Path, from: i64, to: i64, open: &[Arc<Open>]) -> Result<Results> {
let disk = block::scan(root, "metrics")?;
let refs = block::sources(&disk, open);
let mut stats = Stats {
blocks_total: refs.len(),
..Default::default()
};
let mut seen: HashMap<String, (String, u8)> = HashMap::new();
for bref in &refs {
if bref.max_ts < from || bref.min_ts > to {
continue;
}
stats.blocks_scanned += 1;
let Some(m) = load(bref, "metrics")? else {
continue;
};
stats.rows_scanned += m.num_rows();
for r in 0..m.num_rows() {
let name = dict_str(&m, "name", r).unwrap_or("").to_owned();
let unit = dict_str(&m, "unit", r).unwrap_or("").to_owned();
let kind = u8_col(&m, "kind", r);
seen.entry(name).or_insert((unit, kind));
}
}
stats.rows_matched = seen.len();
let mut names: Vec<&String> = seen.keys().collect();
names.sort_unstable();
let mut j = Json::new();
j.arr(|j| {
for n in names {
let (unit, kind) = &seen[n];
j.obj(|j| {
j.key("name");
j.str(n);
j.key("unit");
j.str(unit);
j.key("kind");
j.str(kind_name(*kind));
});
}
});
Ok(Results {
json: j.into_string(),
stats,
next: None,
})
}
fn kind_name(k: u8) -> &'static str {
match k {
x if x == MetricKind::Gauge as u8 => "gauge",
x if x == MetricKind::Sum as u8 => "sum",
x if x == MetricKind::Histogram as u8 => "histogram",
x if x == MetricKind::ExponentialHistogram as u8 => "exponential_histogram",
x if x == MetricKind::Summary as u8 => "summary",
_ => "unset",
}
}
fn load(bref: &Src, name: &str) -> Result<Option<RecordBatch>> {
bref.load(name)
}
pub(crate) fn index_by_parent(b: &RecordBatch) -> Vec<Vec<u32>> {
let Some(col) = b.column_by_name("parent_id") else {
return Vec::new();
};
let parents = col.as_primitive::<UInt32Type>().values();
let mut out: Vec<Vec<u32>> =
vec![Vec::new(); parents.iter().copied().max().unwrap_or(0) as usize + 1];
for (r, &p) in parents.iter().enumerate() {
out[p as usize].push(r as u32);
}
out
}
fn exemplar_time(b: &RecordBatch, row: u32) -> i64 {
b.column_by_name("time_unix_nano").map_or(0, |c| {
c.as_primitive::<TimestampNanosecondType>()
.value(row as usize)
})
}
fn dict_str<'a>(b: &'a RecordBatch, col: &str, row: usize) -> Option<&'a str> {
let c = b.column_by_name(col)?;
let d = c.as_dictionary::<UInt16Type>();
if d.is_null(row) {
return None;
}
Some(
d.values()
.as_string::<i32>()
.value(d.keys().value(row) as usize),
)
}
fn u8_col(b: &RecordBatch, col: &str, row: usize) -> u8 {
b.column_by_name(col)
.map_or(0, |c| c.as_primitive::<UInt8Type>().value(row))
}
fn bound(out: &mut BTreeMap<String, Series>, dropped: &mut BTreeSet<String>, max: usize) {
while out.len() > max {
let Some((k, _)) = out.pop_last() else { break };
if dropped.len() < max {
dropped.insert(k);
}
}
}
fn collect_block(
bref: &Src,
q: &SeriesQuery,
out: &mut BTreeMap<String, Series>,
dropped: &mut BTreeSet<String>,
stats: &mut Stats,
) -> Result<()> {
let Some(metrics) = load(bref, "metrics")? else {
return Ok(());
};
let n_metrics = metrics.num_rows();
let metric_attrs = load(bref, "metric_attrs")?;
let dp_attrs = load(bref, "dp_attrs")?;
let exemplars = load(bref, "exemplars")?;
let by_point = exemplars.as_ref().map(index_by_parent).unwrap_or_default();
let resource_attrs = load(bref, "resource_attrs")?;
let scope_attrs = load(bref, "scope_attrs")?;
let mut wanted = vec![true; n_metrics];
let mut only_suffix = "";
if let Some(want) = &q.name {
let d = metrics
.column_by_name("name")
.map(|c| c.as_dictionary::<UInt16Type>());
let Some(d) = d else { return Ok(()) };
let names = d.values().as_string::<i32>();
let mut code = dict_index(names, want);
if code.is_none() {
for s in DERIVED {
if let Some(base) = want.strip_suffix(s) {
code = dict_index(names, base);
only_suffix = s;
break;
}
}
}
let Some(code) = code else { return Ok(()) };
let codes = d.keys().values();
for (r, w) in wanted.iter_mut().enumerate() {
*w = codes[r] == code;
}
}
let above: Vec<Vec<bool>> = q
.terms
.iter()
.map(|t| {
let key = match &t.target {
crate::query::Target::Attr(k) | crate::query::Target::Field(k) => k,
};
let mut hit = vec![false; n_metrics];
if let Some(a) = &metric_attrs {
for pid in attr_parents(a, key, t.op, &t.value) {
if let Some(s) = hit.get_mut(pid as usize) {
*s = true;
}
}
}
for (table, fk) in [(&resource_attrs, "resource_id"), (&scope_attrs, "scope_id")] {
let (Some(a), Some(col)) = (table, metrics.column_by_name(fk)) else {
continue;
};
let ids = attr_parents(a, key, t.op, &t.value);
let Some(&top) = ids.iter().max() else {
continue;
};
let mut want = vec![false; top as usize + 1];
for id in ids {
want[id as usize] = true;
}
for (r, &id) in col.as_primitive::<UInt16Type>().values().iter().enumerate() {
if want.get(id as usize).copied().unwrap_or(false) {
hit[r] = true;
}
}
}
hit
})
.collect();
let dp_hit: Vec<Vec<bool>> = q
.terms
.iter()
.map(|t| {
let key = match &t.target {
crate::query::Target::Attr(k) | crate::query::Target::Field(k) => k,
};
let mut hit = Vec::new();
if let Some(a) = &dp_attrs {
for pid in attr_parents(a, key, t.op, &t.value) {
if hit.len() <= pid as usize {
hit.resize(pid as usize + 1, false);
}
hit[pid as usize] = true;
}
}
hit
})
.collect();
let mut prefix: Vec<Option<Prefix>> = vec![None; n_metrics];
for table in DP_TABLES {
let Some(dp) = load(bref, table)? else {
continue;
};
let n = dp.num_rows();
stats.rows_scanned += n;
let (Some(time), Some(mid), Some(did)) = (
dp.column_by_name("time_unix_nano")
.map(|c| &**c.as_primitive::<TimestampNanosecondType>().values()),
dp.column_by_name("metric_id")
.map(|c| &**c.as_primitive::<UInt32Type>().values()),
dp.column_by_name("id")
.map(|c| &**c.as_primitive::<UInt32Type>().values()),
) else {
continue;
};
let vals = Values::for_table(table);
for r in 0..n {
let m = mid[r] as usize;
if !(q.from..=q.to).contains(&time[r]) || !wanted.get(m).copied().unwrap_or(false) {
continue;
}
let d = did[r] as usize;
if !(0..q.terms.len()).all(|t| {
above[t].get(m).copied().unwrap_or(false)
|| dp_hit[t].get(d).copied().unwrap_or(false)
}) {
continue;
}
stats.rows_matched += 1;
let (desc_prefix, upper) = prefix[m].get_or_insert_with(|| {
(
describe(&metrics, m),
upper_attrs(&metrics, m, &metric_attrs, &resource_attrs, &scope_attrs),
)
});
let own = own_attrs(&dp_attrs, d as u32);
let attrs = merge(upper, &own);
for (suffix, v) in vals.at(&dp, r) {
if !only_suffix.is_empty() && suffix != only_suffix {
continue;
}
let key = format!("{desc_prefix}\u{1}{suffix}\u{1}{attrs}");
let s = out.entry(key).or_insert_with(|| Series {
desc: with_suffix(desc_prefix, suffix),
attrs: attrs.clone(),
points: Vec::new(),
dropped: 0,
exemplars: Vec::new(),
});
s.points.push((time[r], v));
if s.points.len() >= 2 * q.max_points {
s.compact(q.max_points);
}
if let (Some(ex), Some(rows)) = (&exemplars, by_point.get(d)) {
for &er in rows {
if s.exemplars.len() >= MAX_EXEMPLARS {
break;
}
let mut j = Json::new();
j.obj(|j| crate::query::emit_fields(j, ex, er));
s.exemplars.push((exemplar_time(ex, er), j.into_string()));
}
}
bound(out, dropped, q.max_series);
}
}
}
Ok(())
}
fn describe(metrics: &RecordBatch, row: usize) -> String {
let mut j = Json::new();
j.key("name");
j.str(dict_str(metrics, "name", row).unwrap_or(""));
j.key("unit");
j.str(dict_str(metrics, "unit", row).unwrap_or(""));
j.key("kind");
j.str(kind_name(u8_col(metrics, "kind", row)));
j.key("temporality");
j.u64(u8_col(metrics, "temporality", row) as u64);
j.key("monotonic");
j.bool(
metrics
.column_by_name("is_monotonic")
.is_some_and(|c| c.as_boolean().value(row)),
);
j.into_string()
}
fn with_suffix(desc: &str, suffix: &str) -> String {
if suffix.is_empty() {
return desc.to_owned();
}
match desc.find("\",\"unit\"") {
Some(i) => format!("{}{suffix}{}", &desc[..i], &desc[i..]),
None => desc.to_owned(),
}
}
fn upper_attrs(
metrics: &RecordBatch,
row: usize,
metric_attrs: &Option<RecordBatch>,
resource_attrs: &Option<RecordBatch>,
scope_attrs: &Option<RecordBatch>,
) -> Vec<Attr> {
let fk = |name: &str| {
metrics
.column_by_name(name)
.map(|c| c.as_primitive::<UInt16Type>().value(row) as u32)
};
let mut v = Vec::new();
for (table, parent) in [
(resource_attrs, fk("resource_id")),
(scope_attrs, fk("scope_id")),
(metric_attrs, Some(row as u32)),
] {
let (Some(a), Some(p)) = (table, parent) else {
continue;
};
collect_attrs(a, p, &mut v);
}
dedup_last(&mut v);
v
}
fn own_attrs(dp_attrs: &Option<RecordBatch>, dp_id: u32) -> Vec<Attr> {
let mut v = Vec::new();
if let Some(a) = dp_attrs {
collect_attrs(a, dp_id, &mut v);
}
dedup_last(&mut v);
v
}
fn collect_attrs(a: &RecordBatch, parent: u32, out: &mut Vec<Attr>) {
let parents = a.column(0).as_primitive::<UInt32Type>().values();
for r in 0..a.num_rows() {
if parents[r] != parent {
continue;
}
let mut j = Json::new();
emit_attr(&mut j, a, r);
out.push((attr_key(a, r).to_owned(), j.into_string()));
}
}
fn dedup_last(v: &mut Vec<Attr>) {
v.sort_by(|a, b| a.0.cmp(&b.0));
v.dedup_by(|a, b| {
if a.0 == b.0 {
std::mem::swap(a, b);
true
} else {
false
}
});
}
fn merge(upper: &[Attr], own: &[Attr]) -> String {
let mut j = Json::new();
j.obj(|j| {
let (mut i, mut k) = (0, 0);
while i < upper.len() || k < own.len() {
let take_own = match (upper.get(i), own.get(k)) {
(Some(u), Some(o)) => {
if u.0 == o.0 {
i += 1;
}
o.0 <= u.0
}
(None, Some(_)) => true,
_ => false,
};
let (key, val) = if take_own {
k += 1;
&own[k - 1]
} else {
i += 1;
&upper[i - 1]
};
j.key(key);
j.raw(val);
}
});
j.into_string()
}
enum Values {
Number,
CountSum,
None,
}
impl Values {
fn for_table(table: &str) -> Values {
match table {
"number_dp" => Values::Number,
"hist_dp" | "exp_hist_dp" | "summary_dp" => Values::CountSum,
_ => Values::None,
}
}
fn at(&self, dp: &RecordBatch, r: usize) -> Vec<(&'static str, Pt)> {
match self {
Values::Number => {
if let Some(c) = dp.column_by_name("int") {
let a = c.as_primitive::<Int64Type>();
if !a.is_null(r) {
return vec![("", Pt::Int(a.value(r)))];
}
}
if let Some(c) = dp.column_by_name("double") {
let a = c.as_primitive::<Float64Type>();
if !a.is_null(r) {
return vec![("", Pt::Double(a.value(r)))];
}
}
Vec::new()
}
Values::CountSum => {
let mut v = Vec::with_capacity(2);
if let Some(c) = dp.column_by_name("count") {
let a = c.as_primitive::<UInt64Type>();
if !a.is_null(r) {
v.push((COUNT, Pt::Int(a.value(r) as i64)));
}
}
if let Some(c) = dp.column_by_name("sum") {
let a = c.as_primitive::<Float64Type>();
if !a.is_null(r) {
v.push((SUM, Pt::Double(a.value(r))));
}
}
v
}
Values::None => Vec::new(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use arrow_array::{ArrayRef, Float64Array, Int64Array, UInt32Array, UInt64Array};
use mira_proto::collector::metrics::v1::ExportMetricsServiceRequest;
use mira_proto::common::v1::any_value::Value as AnyVal;
use mira_proto::common::v1::{AnyValue, InstrumentationScope, KeyValue};
use mira_proto::metrics::v1::metric::Data;
use mira_proto::metrics::v1::number_data_point::Value as NumValue;
use mira_proto::metrics::v1::{
Exemplar, ExponentialHistogram, ExponentialHistogramDataPoint, Gauge, Histogram,
HistogramDataPoint, Metric, NumberDataPoint, ResourceMetrics, ScopeMetrics, Sum, Summary,
SummaryDataPoint, exemplar,
};
use mira_proto::resource::v1::Resource;
fn tmp(name: &str) -> std::path::PathBuf {
let d = std::env::temp_dir().join(format!("mira-{name}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&d);
d
}
fn kv(k: &str, v: &str) -> KeyValue {
KeyValue {
key: k.into(),
value: Some(AnyValue {
value: Some(AnyVal::StringValue(v.into())),
}),
}
}
fn publish(
dir: &std::path::Path,
seq: u64,
resource: Option<Resource>,
scope: Option<InstrumentationScope>,
metrics: Vec<Metric>,
) -> std::path::PathBuf {
let req = ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
resource,
scope_metrics: vec![ScopeMetrics {
scope,
metrics,
..Default::default()
}],
..Default::default()
}],
};
let mut b = crate::metrics::MetricsBuilder::new();
b.append_request(&req).unwrap();
let sealed = b.finish().unwrap();
crate::block::publish(dir, "metrics", crate::block::node_id("a"), seq, 0, &sealed)
.unwrap()
.dir
}
fn gauge(name: &str, points: &[(u64, i64)]) -> Metric {
Metric {
name: name.into(),
data: Some(Data::Gauge(Gauge {
data_points: points
.iter()
.map(|&(t, v)| NumberDataPoint {
time_unix_nano: t,
value: Some(NumValue::AsInt(v)),
..Default::default()
})
.collect(),
})),
..Default::default()
}
}
fn query(dir: &std::path::Path, q: &SeriesQuery) -> Results {
series(dir, q).unwrap()
}
fn wide() -> SeriesQuery {
SeriesQuery {
name: None,
from: 0,
to: i64::MAX,
terms: Vec::new(),
max_series: 100,
max_points: 100,
}
}
#[test]
fn a_derived_series_name_queries_back_to_the_series_it_names() {
let dir = std::env::temp_dir().join(format!("mira-derived-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let req = ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
scope_metrics: vec![ScopeMetrics {
metrics: vec![
Metric {
name: "http.server.duration".into(),
data: Some(Data::Histogram(Histogram {
data_points: vec![HistogramDataPoint {
time_unix_nano: 1_000,
count: 3,
sum: Some(1.5),
..Default::default()
}],
..Default::default()
})),
..Default::default()
},
Metric {
name: "queue.depth.count".into(),
data: Some(Data::Gauge(Gauge {
data_points: vec![NumberDataPoint {
time_unix_nano: 1_000,
value: Some(NumValue::AsInt(42)),
..Default::default()
}],
})),
..Default::default()
},
],
..Default::default()
}],
..Default::default()
}],
};
let mut b = crate::metrics::MetricsBuilder::new();
b.append_request(&req).unwrap();
let sealed = b.finish().unwrap();
crate::block::publish(&dir, "metrics", crate::block::node_id("a"), 0, 0, &sealed).unwrap();
let ask = |name: &str| {
series(
&dir,
&SeriesQuery {
name: Some(name.into()),
from: 0,
to: 10_000,
terms: Vec::new(),
max_series: 100,
max_points: 100,
},
)
.unwrap()
.json
};
let both = ask("http.server.duration");
for half in [COUNT, SUM] {
assert!(
both.contains(&format!(r#""name":"http.server.duration{half}""#)),
"{both}"
);
}
let c = ask("http.server.duration.count");
assert!(c.contains(r#""name":"http.server.duration.count""#), "{c}");
assert!(!c.contains(SUM), "{c}");
assert!(c.contains(r#"["1000","3"]"#), "{c}");
let s = ask("http.server.duration.sum");
assert!(s.contains(r#""name":"http.server.duration.sum""#), "{s}");
assert!(!s.contains(COUNT), "{s}");
assert!(s.contains(r#"["1000",1.5]"#), "{s}");
let q = ask("queue.depth.count");
assert!(q.contains(r#""name":"queue.depth.count""#), "{q}");
assert!(q.contains(r#"["1000","42"]"#), "{q}");
assert_eq!(ask("queue.depth.count.sum"), "[]");
assert_eq!(ask("nope.count"), "[]");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn the_series_cap_bounds_the_map_and_says_what_it_refused() {
let dir = std::env::temp_dir().join(format!("mira-cap-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let req = ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
scope_metrics: vec![ScopeMetrics {
metrics: vec![Metric {
name: "rpc.duration".into(),
data: Some(Data::Gauge(Gauge {
data_points: (0..10)
.map(|i| NumberDataPoint {
time_unix_nano: 1_000,
value: Some(NumValue::AsInt(i)),
attributes: vec![mira_proto::common::v1::KeyValue {
key: "peer".into(),
value: Some(mira_proto::common::v1::AnyValue {
value: Some(
mira_proto::common::v1::any_value::Value::StringValue(
format!("p{i}"),
),
),
}),
}],
..Default::default()
})
.collect(),
})),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
};
let mut b = crate::metrics::MetricsBuilder::new();
b.append_request(&req).unwrap();
let sealed = b.finish().unwrap();
crate::block::publish(&dir, "metrics", crate::block::node_id("a"), 0, 0, &sealed).unwrap();
let ask = |max_series: usize| {
series(
&dir,
&SeriesQuery {
name: None,
from: 0,
to: 10_000,
terms: Vec::new(),
max_series,
max_points: 100,
},
)
.unwrap()
};
let r = ask(100);
assert_eq!(r.json.matches(r#""peer""#).count(), 10, "{}", r.json);
assert_eq!(r.stats.dropped_series, 0);
let r = ask(8);
for i in 0..10 {
assert_eq!(
r.json.contains(&format!(r#""peer":"p{i}""#)),
i < 8,
"p{i}: {}",
r.json
);
}
assert_eq!(r.stats.dropped_series, 2);
let r = ask(4);
assert_eq!(r.json.matches(r#""peer""#).count(), 4, "{}", r.json);
assert_eq!(r.stats.dropped_series, 4);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_point_yields_the_value_of_whichever_column_its_table_actually_set() {
for (table, want) in [
("number_dp", "number"),
("hist_dp", "countsum"),
("exp_hist_dp", "countsum"),
("summary_dp", "countsum"),
("logs", "none"),
] {
let got = match Values::for_table(table) {
Values::Number => "number",
Values::CountSum => "countsum",
Values::None => "none",
};
assert_eq!(got, want, "{table}");
}
fn shape(v: Vec<(&'static str, Pt)>) -> Vec<String> {
v.into_iter()
.map(|(s, p)| match p {
Pt::Int(i) => format!("{s}=i{i}"),
Pt::Double(d) => format!("{s}=d{d}"),
})
.collect()
}
let num = RecordBatch::try_from_iter(vec![
(
"int",
Arc::new(Int64Array::from(vec![Some(7), None, None])) as ArrayRef,
),
(
"double",
Arc::new(Float64Array::from(vec![None, Some(0.5), None])) as ArrayRef,
),
])
.unwrap();
assert_eq!(shape(Values::Number.at(&num, 0)), ["=i7"]);
assert_eq!(shape(Values::Number.at(&num, 1)), ["=d0.5"]);
assert!(shape(Values::Number.at(&num, 2)).is_empty());
let only_double = RecordBatch::try_from_iter(vec![(
"double",
Arc::new(Float64Array::from(vec![1.25])) as ArrayRef,
)])
.unwrap();
assert_eq!(shape(Values::Number.at(&only_double, 0)), ["=d1.25"]);
let cs = RecordBatch::try_from_iter(vec![
(
"count",
Arc::new(UInt64Array::from(vec![None, Some(4), Some(4), None])) as ArrayRef,
),
(
"sum",
Arc::new(Float64Array::from(vec![Some(1.5), None, Some(2.5), None])) as ArrayRef,
),
])
.unwrap();
assert_eq!(shape(Values::CountSum.at(&cs, 0)), [".sum=d1.5"]);
assert_eq!(shape(Values::CountSum.at(&cs, 1)), [".count=i4"]);
assert_eq!(
shape(Values::CountSum.at(&cs, 2)),
[".count=i4", ".sum=d2.5"]
);
assert!(shape(Values::CountSum.at(&cs, 3)).is_empty());
assert!(
!shape(Values::CountSum.at(&cs, 1))
.iter()
.any(|s| s.starts_with(SUM))
);
assert!(shape(Values::None.at(&cs, 2)).is_empty());
}
#[test]
fn every_metric_kind_reports_the_name_a_dropdown_shows() {
let dir = tmp("series-kinds");
let point = NumberDataPoint {
time_unix_nano: 1_000,
value: Some(NumValue::AsInt(1)),
..Default::default()
};
let named = |name: &str, unit: &str, data: Option<Data>| Metric {
name: name.into(),
unit: unit.into(),
data,
..Default::default()
};
publish(
&dir,
0,
None,
None,
vec![
named("a.declared", "1", None),
named(
"b.gauge",
"By",
Some(Data::Gauge(Gauge {
data_points: vec![point.clone()],
})),
),
named(
"c.sum",
"1",
Some(Data::Sum(Sum {
data_points: vec![point],
..Default::default()
})),
),
named(
"d.hist",
"s",
Some(Data::Histogram(Histogram {
data_points: vec![HistogramDataPoint {
time_unix_nano: 1_000,
count: 2,
sum: Some(3.0),
..Default::default()
}],
..Default::default()
})),
),
named(
"e.exp",
"s",
Some(Data::ExponentialHistogram(ExponentialHistogram {
data_points: vec![ExponentialHistogramDataPoint {
time_unix_nano: 1_000,
count: 2,
sum: Some(3.0),
..Default::default()
}],
..Default::default()
})),
),
named(
"f.summary",
"s",
Some(Data::Summary(Summary {
data_points: vec![SummaryDataPoint {
time_unix_nano: 1_000,
count: 2,
sum: 3.0,
..Default::default()
}],
})),
),
],
);
let r = names(&dir, 0, 10_000).unwrap();
assert_eq!(
r.json,
concat!(
r#"[{"name":"a.declared","unit":"1","kind":"unset"},"#,
r#"{"name":"b.gauge","unit":"By","kind":"gauge"},"#,
r#"{"name":"c.sum","unit":"1","kind":"sum"},"#,
r#"{"name":"d.hist","unit":"s","kind":"histogram"},"#,
r#"{"name":"e.exp","unit":"s","kind":"exponential_histogram"},"#,
r#"{"name":"f.summary","unit":"s","kind":"summary"}]"#,
)
);
let json = query(&dir, &wide()).json;
for (name, kind) in [
("b.gauge", "gauge"),
("c.sum", "sum"),
("d.hist.count", "histogram"),
("d.hist.sum", "histogram"),
("e.exp.count", "exponential_histogram"),
("e.exp.sum", "exponential_histogram"),
("f.summary.count", "summary"),
("f.summary.sum", "summary"),
] {
assert!(
json.contains(&format!(r#""name":"{name}","unit":"#))
&& json.contains(&format!(r#""kind":"{kind}""#)),
"{name}/{kind}: {json}"
);
}
assert!(!json.contains("a.declared"), "{json}");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_series_is_bounded_by_the_window_asked_for_and_the_points_it_can_chart() {
let dir = tmp("series-bounds");
let pts: Vec<(u64, i64)> = (1..=6).map(|i| (i as u64 * 1_000, i)).collect();
publish(&dir, 0, None, None, vec![gauge("m", &pts)]);
publish(
&dir,
1,
None,
None,
vec![gauge("m", &[(9_000_000_000_000, 99)])],
);
let mut q = wide();
q.to = 10_000;
q.max_points = 4;
let r = query(&dir, &q);
assert_eq!(r.stats.blocks_total, 2);
assert_eq!(r.stats.blocks_scanned, 1, "the far block was opened");
assert!(r.json.contains(r#""dropped_points":2"#), "{}", r.json);
assert!(
r.json
.contains(r#""points":[["3000","3"],["4000","4"],["5000","5"],["6000","6"]]"#),
"{}",
r.json
);
let mut q = wide();
q.to = 10_000;
let r = query(&dir, &q);
assert!(!r.json.contains("dropped_points"), "{}", r.json);
assert!(r.json.contains(r#"["1000","1"]"#), "{}", r.json);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_term_is_satisfied_at_whichever_level_carries_the_attribute() {
let dir = tmp("series-levels");
publish(
&dir,
0,
Some(Resource {
attributes: vec![kv("service.name", "api")],
..Default::default()
}),
Some(InstrumentationScope {
attributes: vec![kv("otel.lib", "sdk")],
..Default::default()
}),
vec![
Metric {
metadata: vec![kv("tier", "gold")],
..gauge("m1", &[(1_000, 1)])
},
Metric {
..gauge("m2", &[(1_000, 2)])
},
],
);
publish(
&dir,
1,
None,
None,
vec![Metric {
metadata: vec![kv("tier", "gold")],
..gauge("m3", &[(1_000, 3)])
}],
);
let with_pod = |name: &str, v: i64, pod: &str| Metric {
data: Some(Data::Gauge(Gauge {
data_points: vec![NumberDataPoint {
time_unix_nano: 1_000,
value: Some(NumValue::AsInt(v)),
attributes: vec![kv("pod", pod)],
..Default::default()
}],
})),
..Metric {
name: name.into(),
..Default::default()
}
};
publish(
&dir,
2,
Some(Resource {
attributes: vec![kv("service.name", "db")],
..Default::default()
}),
None,
vec![with_pod("m4", 4, "p4")],
);
let ask = |key: &str, want: &str| {
let mut q = wide();
q.to = 10_000;
q.terms = vec![Term {
target: crate::query::Target::Attr(key.into()),
op: crate::query::Op::Eq,
value: crate::query::Value::Str(want.into()),
}];
let r = query(&dir, &q);
let mut names: Vec<String> = Vec::new();
for m in ["m1", "m2", "m3", "m4"] {
if r.json.contains(&format!(r#""name":"{m}""#)) {
names.push(m.to_owned());
}
}
names
};
assert_eq!(ask("service.name", "api"), ["m1", "m2"]);
assert_eq!(ask("service.name", "db"), ["m4"]);
assert_eq!(ask("otel.lib", "sdk"), ["m1", "m2"]);
assert_eq!(ask("tier", "gold"), ["m1", "m3"]);
assert_eq!(ask("pod", "p4"), ["m4"]);
assert!(ask("tier", "bronze").is_empty());
assert!(ask("nope", "x").is_empty());
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn attributes_from_every_level_merge_with_the_most_specific_winning() {
let dir = tmp("series-merge");
publish(
&dir,
0,
Some(Resource {
attributes: vec![kv("env", "prod"), kv("region", "eu")],
..Default::default()
}),
None,
vec![Metric {
metadata: vec![kv("env", "staging")],
data: Some(Data::Gauge(Gauge {
data_points: vec![NumberDataPoint {
time_unix_nano: 1_000,
value: Some(NumValue::AsInt(1)),
attributes: vec![kv("env", "canary"), kv("pod", "x")],
..Default::default()
}],
})),
..Metric {
name: "m".into(),
..Default::default()
}
}],
);
let json = query(&dir, &wide()).json;
assert!(
json.contains(r#""attributes":{"env":"canary","pod":"x","region":"eu"}"#),
"{json}"
);
assert!(
!json.contains("prod") && !json.contains("staging"),
"{json}"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn exemplars_are_capped_so_one_exporter_cannot_inflate_every_chart_response() {
let dir = tmp("series-exemplars");
publish(
&dir,
0,
None,
None,
vec![Metric {
data: Some(Data::Gauge(Gauge {
data_points: vec![NumberDataPoint {
time_unix_nano: 1_000,
value: Some(NumValue::AsInt(1)),
exemplars: (0..MAX_EXEMPLARS as i64 + 6)
.map(|i| Exemplar {
time_unix_nano: 1_000 + i as u64,
value: Some(exemplar::Value::AsInt(i)),
trace_id: vec![7u8; 16].into(),
..Default::default()
})
.collect(),
..Default::default()
}],
})),
..Metric {
name: "m".into(),
..Default::default()
}
}],
);
let json = query(&dir, &wide()).json;
assert_eq!(
json.matches(r#""time_unix_nano""#).count(),
MAX_EXEMPLARS,
"{json}"
);
assert!(json.contains(r#""trace_id":"07070707"#), "{json}");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_table_the_reader_did_not_write_contributes_nothing_instead_of_panicking() {
let dir = tmp("series-ragged");
publish(&dir, 0, None, None, vec![gauge("ok", &[(1_000, 1)])]);
let unlinked = publish(&dir, 1, None, None, vec![gauge("gone", &[(1_000, 2)])]);
let ragged = publish(&dir, 2, None, None, vec![gauge("ragged", &[(1_000, 3)])]);
std::fs::remove_file(unlinked.join("metrics.arrow")).unwrap();
let table = ragged.join("number_dp.arrow");
let mut b = crate::block::open_table_opt(&table)
.unwrap()
.unwrap()
.batches[0]
.clone();
b.remove_column(b.schema().index_of("id").unwrap());
let staged = ragged.join("number_dp.arrow.staged");
crate::block::write_table(&staged, &b).unwrap();
std::fs::rename(&staged, &table).unwrap();
let r = query(&dir, &wide());
assert_eq!(r.stats.blocks_scanned, 3, "all three were in the window");
assert!(r.json.contains(r#""name":"ok""#), "{}", r.json);
assert!(!r.json.contains("gone"), "{}", r.json);
assert!(!r.json.contains("ragged"), "{}", r.json);
let n = names(&dir, 0, 10_000).unwrap();
assert_eq!(
n.json,
concat!(
r#"[{"name":"ok","unit":"","kind":"gauge"},"#,
r#"{"name":"ragged","unit":"","kind":"gauge"}]"#,
)
);
let no_parent = RecordBatch::try_from_iter(vec![(
"id",
Arc::new(UInt32Array::from(vec![0u32, 1])) as ArrayRef,
)])
.unwrap();
assert!(index_by_parent(&no_parent).is_empty());
assert_eq!(
with_suffix(r#""kind":"histogram""#, COUNT),
r#""kind":"histogram""#
);
assert_eq!(
with_suffix(r#""name":"x","unit":"s""#, ""),
r#""name":"x","unit":"s""#
);
assert_eq!(
with_suffix(r#""name":"x","unit":"s""#, COUNT),
r#""name":"x.count","unit":"s""#
);
let _ = std::fs::remove_dir_all(&dir);
}
}