use std::collections::HashMap;
use std::path::Path;
use std::sync::Arc;
use arrow_array::cast::AsArray;
use arrow_array::types::{TimestampNanosecondType, UInt8Type, UInt16Type, UInt32Type, UInt64Type};
use arrow_array::{Array, ArrayRef, FixedSizeBinaryArray};
use crate::block::{self, Src};
use crate::error::Result;
use crate::json::Json;
use crate::query::{self, Search};
use crate::signal::Open;
#[derive(Debug, Clone, Default, PartialEq)]
pub struct Frame {
pub from: i64,
pub to: i64,
pub entities: Vec<u64>,
pub traces: Vec<[u8; 16]>,
pub truncated: bool,
}
pub const MAX_TRACES: usize = 1000;
pub const MAX_ENTITIES: usize = 256;
impl Frame {
fn add_trace(&mut self, id: &[u8]) {
let Ok(id) = <[u8; 16]>::try_from(id) else {
return;
};
if self.traces.contains(&id) {
return;
}
if self.traces.len() >= MAX_TRACES {
self.truncated = true;
return;
}
self.traces.push(id);
}
fn add_entity(&mut self, key: u64) {
if key == crate::identity::NO_IDENTITY || self.entities.contains(&key) {
return;
}
if self.entities.len() >= MAX_ENTITIES {
self.truncated = true;
return;
}
self.entities.push(key);
}
pub fn write_json(&self, j: &mut Json, names: &HashMap<u64, String>) {
j.obj(|j| {
j.key("from");
j.i64_str(self.from);
j.key("to");
j.i64_str(self.to);
j.key("entities");
j.arr(|j| {
for e in &self.entities {
j.obj(|j| {
j.key("key");
j.u64_str(*e);
j.key("name");
j.str(names.get(e).map_or(UNKNOWN, String::as_str));
});
}
});
j.key("traces");
j.arr(|j| {
for t in &self.traces {
j.hex(t);
}
});
j.key("truncated");
j.bool(self.truncated);
});
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Expand {
Traces,
Peers,
Around(i64),
}
pub fn anchor(root: &Path, q: &Search, open: &[Arc<Open>]) -> Result<(Frame, Stats)> {
let disk = block::scan(root, q.signal.dir())?;
let mut refs = block::sources(&disk, open);
let mut st = Stats {
blocks_total: refs.len(),
..Default::default()
};
refs.retain(|b| b.overlaps(q.from, q.to));
refs.sort_by_key(|b| std::cmp::Reverse((b.max_ts, b.seq)));
let mut f = Frame {
from: q.from,
to: q.to,
..Default::default()
};
let scan = query::Scan::new(q);
let mut i = 0;
while i < refs.len() && !f.truncated {
let answers = scan.wave(&refs, i, refs.len());
i += answers.len();
for done in answers {
let (Some(b), hits) = done? else { continue };
st.blocks_scanned += 1;
st.rows_scanned += b.root.num_rows();
st.rows_matched += hits.len();
let Some(h0) = hits.first() else { continue };
let keys = entity_keys(&refs[h0.block])?;
let traces = b.root.column_by_name("trace_id").and_then(binary);
let rids = b.root.column_by_name("resource_id");
for h in &hits {
let row = h.row as usize;
if let Some(t) = traces.filter(|t| t.is_valid(row)) {
f.add_trace(t.value(row));
}
if let Some(&k) =
rids.and_then(|c| keys.get(c.as_primitive::<UInt16Type>().value(row) as usize))
{
f.add_entity(k);
}
}
}
}
Ok((f, st))
}
pub fn expand(
root: &Path,
f: &Frame,
ops: &[Expand],
open: &[Arc<Open>],
) -> Result<(Frame, Stats)> {
let mut f = f.clone();
let mut st = Stats::default();
for op in ops {
match *op {
Expand::Around(d) => {
f.from = f.from.saturating_sub(d);
f.to = f.to.saturating_add(d);
}
op => walk_spans(root, &mut f, op, open, &mut st)?,
}
}
Ok((f, st))
}
fn walk_spans(
root: &Path,
f: &mut Frame,
op: Expand,
open: &[Arc<Open>],
st: &mut Stats,
) -> Result<()> {
if f.traces.is_empty() {
return Ok(());
}
let disk = block::scan(root, "traces")?;
let refs = block::sources(&disk, open);
st.blocks_total += refs.len();
let (mut lo, mut hi) = (i64::MAX, i64::MIN);
for bref in &refs {
if op == Expand::Peers && !bref.overlaps(f.from, f.to) {
continue;
}
if !may_hold(bref, &f.traces) {
continue;
}
let Some(spans) = bref.load("spans")? else {
continue;
};
let Some(ids) = spans.column_by_name("trace_id").and_then(binary) else {
continue;
};
st.blocks_scanned += 1;
st.rows_scanned += spans.num_rows();
let keys = entity_keys(bref)?;
let rids = spans.column_by_name("resource_id");
let start = spans
.column_by_name("start_time_unix_nano")
.map(|c| c.as_primitive::<TimestampNanosecondType>());
let dur = spans
.column_by_name("duration_nano")
.map(|c| c.as_primitive::<UInt64Type>());
for row in 0..spans.num_rows() {
if !ids.is_valid(row) || !f.traces.iter().any(|t| t[..] == *ids.value(row)) {
continue;
}
st.rows_matched += 1;
if op == Expand::Peers {
if let Some(&k) =
rids.and_then(|c| keys.get(c.as_primitive::<UInt16Type>().value(row) as usize))
{
f.add_entity(k);
}
} else {
let s = start.map_or(0, |c| c.value(row));
lo = lo.min(s);
hi = hi.max(s.saturating_add(dur.map_or(0, |c| c.value(row)) as i64));
}
}
}
if op == Expand::Traces && lo <= hi {
f.from = f.from.min(lo);
f.to = f.to.max(hi);
}
Ok(())
}
pub fn map(
root: &Path,
from: i64,
to: i64,
max_spans: usize,
open: &[Arc<Open>],
) -> Result<query::Results> {
let disk = block::scan(root, "traces")?;
let mut refs = block::sources(&disk, open);
let mut stats = query::Stats {
blocks_total: refs.len(),
..Default::default()
};
refs.retain(|b| b.overlaps(from, to));
refs.sort_by_key(|b| std::cmp::Reverse((b.max_ts, b.seq)));
let mut nodes: HashMap<u64, Node> = HashMap::new();
let mut edges: HashMap<(u64, u64), Edge> = HashMap::new();
let mut names: HashMap<u64, String> = HashMap::new();
let mut unresolved = 0u64;
for bref in &refs {
if stats.rows_scanned >= max_spans {
break;
}
let Some(spans) = bref.load("spans")? else {
continue;
};
let (Some(ids), Some(parents), Some(rids)) = (
spans.column_by_name("span_id").and_then(binary),
spans.column_by_name("parent_span_id").and_then(binary),
spans.column_by_name("resource_id"),
) else {
continue;
};
stats.blocks_scanned += 1;
stats.rows_scanned += spans.num_rows();
let keys = entity_keys(bref)?;
resource_names(bref, &keys, &mut names)?;
let rids = rids.as_primitive::<UInt16Type>();
let time = spans
.column_by_name("start_time_unix_nano")
.map(|c| c.as_primitive::<TimestampNanosecondType>());
let dur = spans
.column_by_name("duration_nano")
.map(|c| c.as_primitive::<UInt64Type>());
let status = spans
.column_by_name("status_code")
.map(|c| c.as_primitive::<UInt8Type>());
let key_of = |row: usize| keys.get(rids.value(row) as usize).copied().unwrap_or(0);
let mut owner: HashMap<&[u8], u64> = HashMap::with_capacity(spans.num_rows());
let mut live: Vec<usize> = Vec::with_capacity(spans.num_rows());
for row in 0..spans.num_rows() {
if ids.is_valid(row) {
owner.insert(ids.value(row), key_of(row));
}
if time.is_none_or(|t| t.value(row) >= from && t.value(row) <= to) {
live.push(row);
}
}
for row in live {
let key = key_of(row);
let bad = status.is_some_and(|s| s.value(row) == STATUS_ERROR);
let d = dur.map_or(0, |c| c.value(row));
let n = nodes.entry(key).or_default();
n.spans += 1;
n.errors += u64::from(bad);
n.nanos += d;
let caller = if parents.is_valid(row) {
owner.get(parents.value(row)).copied()
} else {
Some(ENTRY)
};
let Some(caller) = caller else {
unresolved += 1;
continue;
};
if caller == key {
continue;
}
let e = edges.entry((caller, key)).or_default();
e.calls += 1;
e.errors += u64::from(bad);
e.nanos += d;
e.max_nanos = e.max_nanos.max(d);
}
}
stats.rows_matched = edges.len();
let mut ns: Vec<_> = nodes.into_iter().collect();
ns.sort_unstable_by_key(|(k, _)| *k);
let mut es: Vec<_> = edges.into_iter().collect();
es.sort_unstable_by_key(|(k, _)| *k);
let mut j = Json::new();
j.obj(|j| {
j.key("nodes");
j.arr(|j| {
for (key, n) in &ns {
j.obj(|j| {
j.key("key");
j.u64_str(*key);
j.key("name");
j.str(names.get(key).map_or(UNKNOWN, String::as_str));
j.key("spans");
j.u64(n.spans);
j.key("errors");
j.u64(n.errors);
j.key("avg_nano");
j.u64_str(n.nanos / n.spans.max(1));
});
}
});
j.key("edges");
j.arr(|j| {
for ((a, b), e) in &es {
j.obj(|j| {
j.key("from");
if *a == ENTRY {
j.str("entry");
} else {
j.u64_str(*a);
}
j.key("to");
j.u64_str(*b);
j.key("calls");
j.u64(e.calls);
j.key("errors");
j.u64(e.errors);
j.key("avg_nano");
j.u64_str(e.nanos / e.calls.max(1));
j.key("max_nano");
j.u64_str(e.max_nanos);
});
}
});
j.key("unresolved");
j.u64(unresolved);
});
Ok(query::Results {
json: j.into_string(),
stats,
next: None,
})
}
pub fn entities(
root: &Path,
from: i64,
to: i64,
open: &[Vec<Arc<Open>>],
) -> Result<query::Results> {
let mut names: HashMap<u64, String> = HashMap::new();
let mut seen: HashMap<u64, u64> = HashMap::new();
let mut stats = query::Stats::default();
for (i, signal) in SIGNALS.iter().enumerate() {
let disk = block::scan(root, signal)?;
let refs = block::sources(&disk, open_for(open, i));
stats.blocks_total += refs.len();
for bref in &refs {
if !bref.overlaps(from, to) {
continue;
}
stats.blocks_scanned += 1;
let keys = entity_keys(bref)?;
resource_names(bref, &keys, &mut names)?;
for k in keys.iter().filter(|&&k| k != 0) {
*seen.entry(*k).or_default() += 1;
}
}
}
let mut out: Vec<_> = seen.into_iter().collect();
out.sort_unstable_by(|a, b| {
let name = |k: &u64| names.get(k).map_or(UNKNOWN, String::as_str);
name(&a.0).cmp(name(&b.0)).then(a.0.cmp(&b.0))
});
stats.rows_matched = out.len();
let mut j = Json::new();
j.arr(|j| {
for (key, blocks) in &out {
j.obj(|j| {
j.key("key");
j.u64_str(*key);
j.key("name");
j.str(names.get(key).map_or(UNKNOWN, String::as_str));
j.key("blocks");
j.u64(*blocks);
});
}
});
Ok(query::Results {
json: j.into_string(),
stats,
next: None,
})
}
pub fn names_of(root: &Path, f: &Frame, open: &[Vec<Arc<Open>>]) -> Result<HashMap<u64, String>> {
let mut names = HashMap::new();
if f.entities.is_empty() {
return Ok(names);
}
for (i, signal) in SIGNALS.iter().enumerate() {
let disk = block::scan(root, signal)?;
for bref in block::sources(&disk, open_for(open, i)) {
if !bref.overlaps(f.from, f.to) {
continue;
}
let keys = entity_keys(&bref)?;
if keys.iter().any(|k| f.entities.contains(k)) {
resource_names(&bref, &keys, &mut names)?;
}
}
}
names.retain(|k, _| f.entities.contains(k));
Ok(names)
}
#[derive(Debug, Default, Clone)]
pub struct Stats {
pub blocks_total: usize,
pub blocks_scanned: usize,
pub rows_scanned: usize,
pub rows_matched: usize,
}
const SIGNALS: [&str; 3] = ["logs", "traces", "metrics"];
fn open_for(open: &[Vec<Arc<Open>>], i: usize) -> &[Arc<Open>] {
open.get(i).map_or(&[], Vec::as_slice)
}
const STATUS_ERROR: u8 = 2;
const UNKNOWN: &str = "unknown";
const ENTRY: u64 = u64::MAX;
#[derive(Default)]
struct Node {
spans: u64,
errors: u64,
nanos: u64,
}
#[derive(Default)]
struct Edge {
calls: u64,
errors: u64,
nanos: u64,
max_nanos: u64,
}
fn entity_keys(bref: &Src<'_>) -> Result<Vec<u64>> {
let Some(r) = bref.load("resources")? else {
return Ok(Vec::new());
};
let (Some(ids), Some(keys)) = (r.column_by_name("id"), r.column_by_name("key")) else {
return Ok(Vec::new());
};
let ids = ids.as_primitive::<UInt16Type>();
let keys = keys.as_primitive::<UInt64Type>();
let mut out = vec![0u64; r.num_rows()];
for row in 0..r.num_rows() {
let i = ids.value(row) as usize;
if i < out.len() {
out[i] = keys.value(row);
}
}
Ok(out)
}
fn resource_names(bref: &Src<'_>, keys: &[u64], out: &mut HashMap<u64, String>) -> Result<()> {
let Some(a) = bref.load("resource_attrs")? else {
return Ok(());
};
let (Some(parents), Some(vals)) = (a.column_by_name("parent_id"), a.column_by_name("str"))
else {
return Ok(());
};
let parents = parents.as_primitive::<UInt32Type>();
let vals = crate::attrs::str_values(vals);
for row in 0..a.num_rows() {
if query::attr_key(&a, row) != "service.name" || !vals.is_valid(row) {
continue;
}
if let Some(&k) = keys.get(parents.value(row) as usize).filter(|&&k| k != 0) {
out.entry(k).or_insert_with(|| vals.value(row).to_owned());
}
}
Ok(())
}
fn may_hold(bref: &Src<'_>, traces: &[[u8; 16]]) -> bool {
let Some(dir) = bref.dir else { return true };
match std::fs::read(dir.join(crate::bloom::TRACE_IDX)) {
Ok(f) => traces.iter().any(|t| crate::bloom::may_contain(&f, t)),
Err(_) => true,
}
}
fn binary(c: &ArrayRef) -> Option<&FixedSizeBinaryArray> {
c.as_any().downcast_ref::<FixedSizeBinaryArray>()
}
#[cfg(test)]
mod tests {
use super::*;
use mira_proto::collector::trace::v1::ExportTraceServiceRequest;
use mira_proto::common::v1::any_value::Value as AnyVal;
use mira_proto::common::v1::{AnyValue, KeyValue};
use mira_proto::resource::v1::Resource;
use mira_proto::trace::v1::{ResourceSpans, ScopeSpans, Span};
const T0: u64 = 1_000_000_000;
const HOUR: u64 = 3_600_000_000_000;
fn kv(k: &str, v: &str) -> KeyValue {
KeyValue {
key: k.into(),
value: Some(AnyValue {
value: Some(AnyVal::StringValue(v.into())),
}),
}
}
fn span(trace: u8, id: [u8; 8], parent: Option<[u8; 8]>, start: u64, dur: u64) -> Span {
Span {
trace_id: [trace; 16].to_vec().into(),
span_id: id.to_vec().into(),
parent_span_id: parent.map(|p| p.to_vec()).unwrap_or_default().into(),
name: "GET /checkout".into(),
start_time_unix_nano: start,
end_time_unix_nano: start + dur,
..Default::default()
}
}
fn traces_block(root: &Path, seq: u64, services: Vec<(&str, Vec<Span>)>) -> std::path::PathBuf {
let mut b = crate::traces::TracesBuilder::new();
b.append_request(&ExportTraceServiceRequest {
resource_spans: services
.into_iter()
.map(|(name, spans)| ResourceSpans {
resource: Some(Resource {
attributes: vec![kv("service.name", name)],
..Default::default()
}),
scope_spans: vec![ScopeSpans {
spans,
..Default::default()
}],
..Default::default()
})
.collect(),
})
.unwrap();
let sealed = b.finish().unwrap();
block::publish(root, "traces", block::node_id("a"), seq, 0, &sealed)
.unwrap()
.dir
}
fn drop_column(dir: &Path, table: &str, col: &str) {
let path = dir.join(format!("{table}.arrow"));
let mut b = block::open_table_opt(&path).unwrap().unwrap().batches[0].clone();
b.remove_column(b.schema().index_of(col).unwrap());
let staged = dir.join(format!("{table}.staged"));
block::write_table(&staged, &b).unwrap();
std::fs::rename(&staged, &path).unwrap();
}
#[test]
fn a_block_the_reader_cannot_use_contributes_nothing_and_widens_nothing() {
use crate::query::{Search, Signal};
let root = std::env::temp_dir().join(format!("mira-frame-torn-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
let sid = |svc: u8| [svc, 1, 0, 0, 0, 0, 0, 0];
let chain = |start: u64| {
vec![
("gateway", vec![span(1, sid(1), None, start, 900)]),
("api", vec![span(1, sid(2), Some(sid(1)), start, 500)]),
("db", vec![span(1, sid(3), Some(sid(2)), start, 100)]),
]
};
traces_block(&root, 0, chain(T0 + 1));
let ghost = traces_block(
&root,
1,
vec![(
"ghost",
vec![span(
9,
[9, 1, 0, 0, 0, 0, 0, 0],
Some([8, 8, 0, 0, 0, 0, 0, 0]),
T0 + HOUR,
10,
)],
)],
);
drop_column(&ghost, "resources", "key");
drop_column(&ghost, "resource_attrs", "str");
let unlinked = traces_block(&root, 2, chain(T0 + 2 * HOUR));
std::fs::remove_file(unlinked.join("spans.arrow")).unwrap();
std::fs::remove_file(unlinked.join("resources.arrow")).unwrap();
let foreign = traces_block(&root, 3, chain(T0 + 3 * HOUR));
std::fs::remove_file(foreign.join(crate::bloom::TRACE_IDX)).unwrap();
drop_column(&foreign, "spans", "trace_id");
drop_column(&foreign, "spans", "parent_span_id");
let q = Search {
signal: Signal::Traces,
from: T0 as i64,
to: (T0 + 2) as i64,
terms: Vec::new(),
limit: 10,
after: None,
};
let (f, _) = anchor(&root, &q, &[]).unwrap();
assert_eq!(f.traces, vec![[1u8; 16]], "one trace in the window");
assert_eq!(f.entities.len(), 3, "gateway, api and db");
let (g, st) = expand(&root, &f, &[Expand::Traces], &[]).unwrap();
assert_eq!(st.blocks_total, 4);
assert_eq!(st.blocks_scanned, 1, "three unusable blocks were read");
assert_eq!((g.from, g.to), (T0 as i64, (T0 + 901) as i64));
assert_eq!(g.entities, f.entities, "`Traces` does not touch entities");
let all = entities(&root, 0, i64::MAX, &[]).unwrap();
assert_eq!(all.stats.blocks_scanned, 4);
assert_eq!(all.json.matches(r#""name""#).count(), 3, "{}", all.json);
assert_eq!(all.json.matches(r#""blocks":2"#).count(), 3, "{}", all.json);
assert!(!all.json.contains("ghost"), "{}", all.json);
for name in ["api", "db", "gateway"] {
assert!(all.json.contains(&format!(r#""name":"{name}""#)), "{name}");
}
let near = entities(&root, T0 as i64, (T0 + 2_000) as i64, &[]).unwrap();
assert_eq!(near.stats.blocks_scanned, 1);
assert_eq!(
near.json.matches(r#""blocks":1"#).count(),
3,
"{}",
near.json
);
let sorted = |f: &Frame| {
let mut v: Vec<String> = names_of(&root, f, &[]).unwrap().into_values().collect();
v.sort();
v
};
assert_eq!(sorted(&f), ["api", "db", "gateway"]);
let wide = Frame {
from: 0,
to: i64::MAX,
..f
};
assert_eq!(sorted(&wide), ["api", "db", "gateway"]);
let m = map(&root, 0, i64::MAX, 1_000, &[]).unwrap();
assert_eq!(m.stats.blocks_scanned, 2, "{}", m.json);
assert_eq!(m.stats.rows_matched, 3, "entry->gateway->api->db");
assert!(m.json.contains(r#""from":"entry""#), "{}", m.json);
assert!(m.json.contains(r#""unresolved":1"#), "{}", m.json);
assert_eq!(
m.json.matches(r#""name":"unknown""#).count(),
1,
"{}",
m.json
);
let one = map(&root, 0, i64::MAX, 1, &[]).unwrap();
assert_eq!(one.stats.blocks_scanned, 1, "{}", one.json);
assert!(!one.json.contains(r#""from":"entry""#), "{}", one.json);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn the_caps_are_what_makes_a_frame_a_frame() {
let mut f = Frame::default();
for i in 0..MAX_TRACES as u32 + 5 {
f.add_trace(&i.to_be_bytes().repeat(4));
}
assert_eq!(f.traces.len(), MAX_TRACES);
assert!(f.truncated, "a dropped trace has to be visible");
f.add_trace(&[1, 2, 3]);
assert_eq!(f.traces.len(), MAX_TRACES);
let mut f = Frame::default();
for i in 0..MAX_ENTITIES as u64 + 5 {
f.add_entity(i + 1);
}
assert_eq!(f.entities.len(), MAX_ENTITIES);
f.add_entity(0);
assert!(!f.entities.contains(&0), "the sentinel is not an entity");
}
#[test]
fn expansions_are_ordered_because_they_do_not_commute() {
let root = std::env::temp_dir().join("mira-frame-empty");
let _ = std::fs::remove_dir_all(&root);
let f = Frame {
from: 1_000,
to: 2_000,
..Default::default()
};
let (g, _) = expand(&root, &f, &[Expand::Around(500), Expand::Around(500)], &[]).unwrap();
assert_eq!((g.from, g.to), (0, 3_000));
let (h, st) = expand(&root, &f, &[Expand::Traces, Expand::Peers], &[]).unwrap();
assert_eq!(h, f);
assert_eq!(st.blocks_scanned, 0);
}
}