use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering::Relaxed};
use arrow_array::cast::AsArray;
use arrow_array::types::{
Float64Type, Int32Type, Int64Type, TimestampNanosecondType, UInt8Type, UInt16Type, UInt32Type,
UInt64Type,
};
use arrow_array::{Array, FixedSizeBinaryArray, RecordBatch, StringArray};
use arrow_schema::DataType;
use mira_proto::common::v1::AnyValue;
use prost::Message;
use crate::block::{self, Src};
use crate::error::Result;
use crate::json::Json;
use crate::schema::AttrType;
use crate::signal::Open;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Op {
Eq,
Ne,
Lt,
Lte,
Gt,
Gte,
Contains,
}
impl Op {
pub fn parse(s: &str) -> Option<Op> {
Some(match s {
"eq" | "=" | "==" => Op::Eq,
"ne" | "!=" => Op::Ne,
"lt" | "<" => Op::Lt,
"lte" | "<=" => Op::Lte,
"gt" | ">" => Op::Gt,
"gte" | ">=" => Op::Gte,
"contains" | "~" => Op::Contains,
_ => return None,
})
}
fn test_ord(self, ord: std::cmp::Ordering) -> bool {
use std::cmp::Ordering::*;
match self {
Op::Eq => ord == Equal,
Op::Ne => ord != Equal,
Op::Lt => ord == Less,
Op::Lte => ord != Greater,
Op::Gt => ord == Greater,
Op::Gte => ord != Less,
Op::Contains => false,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum Value {
Str(String),
Int(i64),
Double(f64),
Bool(bool),
}
impl Value {
fn as_i64(&self) -> Option<i64> {
match self {
Value::Int(i) => Some(*i),
Value::Double(d) if d.fract() == 0.0 => Some(*d as i64),
Value::Str(s) => s.parse().ok(),
Value::Double(_) | Value::Bool(_) => None,
}
}
fn as_f64(&self) -> Option<f64> {
match self {
Value::Int(i) => Some(*i as f64),
Value::Double(d) => Some(*d),
Value::Str(s) => s.parse().ok(),
Value::Bool(_) => None,
}
}
fn as_bool(&self) -> Option<bool> {
match self {
Value::Bool(b) => Some(*b),
Value::Str(s) if s == "true" => Some(true),
Value::Str(s) if s == "false" => Some(false),
Value::Str(_) | Value::Int(_) | Value::Double(_) => None,
}
}
fn as_str(&self) -> Option<&str> {
match self {
Value::Str(s) => Some(s),
_ => None,
}
}
}
#[derive(Debug, Clone)]
pub enum Target {
Field(String),
Attr(String),
}
#[derive(Debug, Clone)]
pub struct Term {
pub target: Target,
pub op: Op,
pub value: Value,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Signal {
Logs,
Traces,
}
impl Signal {
pub fn parse(s: &str) -> Option<Signal> {
match s {
"logs" => Some(Signal::Logs),
"traces" | "spans" => Some(Signal::Traces),
_ => None,
}
}
pub fn dir(self) -> &'static str {
match self {
Signal::Logs => "logs",
Signal::Traces => "traces",
}
}
fn root(self) -> &'static str {
match self {
Signal::Logs => "logs",
Signal::Traces => "spans",
}
}
fn attrs(self) -> &'static str {
match self {
Signal::Logs => "log_attrs",
Signal::Traces => "span_attrs",
}
}
fn time_col(self) -> &'static str {
match self {
Signal::Logs => "time_unix_nano",
Signal::Traces => "start_time_unix_nano",
}
}
}
#[derive(Debug, Clone)]
pub struct Search {
pub signal: Signal,
pub from: i64,
pub to: i64,
pub terms: Vec<Term>,
pub limit: usize,
pub after: Option<Cursor>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Cursor {
pub ts: i64,
pub node: u32,
pub seq: u64,
pub row: u32,
}
impl Cursor {
fn key(&self) -> std::cmp::Reverse<(i64, u32, u64, u32)> {
std::cmp::Reverse((self.ts, self.node, self.seq, self.row))
}
}
impl std::fmt::Display for Cursor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}.{}.{}.{}", self.ts, self.node, self.seq, self.row)
}
}
impl std::str::FromStr for Cursor {
type Err = String;
fn from_str(s: &str) -> std::result::Result<Cursor, String> {
let bad = || format!("{s:?} is not a cursor; pass back the `next` field verbatim");
let mut p = s.split('.');
let mut next = |f: &dyn Fn(&str) -> bool| p.next().filter(|v| f(v)).ok_or_else(bad);
let ts = next(&|v: &str| !v.is_empty())?.parse().map_err(|_| bad())?;
let digits = |v: &str| !v.is_empty() && v.bytes().all(|b| b.is_ascii_digit());
let c = Cursor {
ts,
node: next(&digits)?.parse().map_err(|_| bad())?,
seq: next(&digits)?.parse().map_err(|_| bad())?,
row: next(&digits)?.parse().map_err(|_| bad())?,
};
match p.next() {
Some(_) => Err(bad()),
None => Ok(c),
}
}
}
pub(crate) struct Hit {
ts: i64,
pub(crate) block: usize,
pub(crate) row: u32,
}
#[derive(Debug, Default, Clone)]
pub struct Stats {
pub blocks_total: usize,
pub blocks_scanned: usize,
pub rows_scanned: usize,
pub rows_matched: usize,
pub dropped_series: usize,
}
pub struct Results {
pub json: String,
pub stats: Stats,
pub next: Option<Cursor>,
}
fn trace_needle(q: &Search) -> Option<[u8; 16]> {
q.terms.iter().find_map(|t| match (&t.target, t.op) {
(Target::Field(f), Op::Eq) if f == "trace_id" => unhex(t.value.as_str()?)?.try_into().ok(),
_ => None,
})
}
struct AttrProbe {
hash: (u64, u64),
numeric: bool,
}
impl AttrProbe {
fn maybe(&self, f: &crate::bloom::Filter) -> bool {
f.may_contain(self.hash) || (self.numeric && f.flags & crate::bloom::HAS_DOUBLE != 0)
}
}
fn attr_probes(q: &Search) -> Vec<AttrProbe> {
q.terms
.iter()
.filter(|t| t.op == Op::Eq)
.filter_map(|t| match &t.target {
Target::Attr(key) => Some(AttrProbe {
hash: crate::bloom::attr_hash(key, canon(&t.value).as_bytes()),
numeric: t.value.as_f64().is_some(),
}),
Target::Field(_) => None,
})
.collect()
}
fn range_probes(q: &Search) -> Vec<crate::zone::Probe> {
q.terms
.iter()
.filter(|t| !matches!(t.op, Op::Ne | Op::Contains))
.filter(|t| matches!(t.value, Value::Int(_) | Value::Double(_)))
.map(|t| crate::zone::Probe {
key: match &t.target {
Target::Attr(k) => crate::zone::attr_key(k),
Target::Field(f) => crate::zone::field_key(f),
},
op: t.op,
int: t.value.as_i64(),
float: t.value.as_f64(),
})
.collect()
}
fn canon(v: &Value) -> String {
match v {
Value::Str(s) => s.clone(),
Value::Int(i) => i.to_string(),
Value::Bool(b) => b.to_string(),
Value::Double(d) if d.fract() == 0.0 => (*d as i64).to_string(),
Value::Double(d) => d.to_string(),
}
}
pub fn search(root: &Path, q: &Search) -> Result<Results> {
search_open(root, q, &[])
}
pub fn search_open(root: &Path, q: &Search, open_blocks: &[Arc<Open>]) -> Result<Results> {
let disk = block::scan(root, q.signal.dir())?;
let mut refs = block::sources(&disk, open_blocks);
let mut stats = Stats {
blocks_total: refs.len(),
..Default::default()
};
refs.retain(|b| b.overlaps(q.from, q.to));
if let Some(c) = &q.after {
refs.retain(|b| b.min_ts <= c.ts);
}
refs.sort_by_key(|b| std::cmp::Reverse((b.max_ts, b.seq)));
let refs = &refs[..];
let scan = Scan::new(q);
let mut hits: Vec<Hit> = Vec::new();
let mut open: Vec<Option<Block>> = Vec::with_capacity(refs.len());
let mut width = 1;
let mut i = 0;
while i < refs.len() {
if hits.len() >= q.limit && hits.last().is_some_and(|w| refs[i].max_ts < w.ts) {
break;
}
let asked = (i + width).min(refs.len());
let answers = scan.wave(refs, i, asked);
let end = i + answers.len();
for done in answers {
let (block, found) = done?;
if let Some(b) = &block {
stats.blocks_scanned += 1;
stats.rows_scanned += b.root.num_rows();
}
stats.rows_matched += found.len();
hits.extend(found);
open.push(block);
}
if hits.len() > q.limit {
hits.select_nth_unstable_by_key(q.limit, |h| cursor(&refs[h.block], h).key());
hits.truncate(q.limit);
}
hits.sort_unstable_by_key(|h| cursor(&refs[h.block], h).key());
i = end;
width = (width * 2).min(MAX_FANOUT);
}
let mut j = Json::new();
j.arr(|j| {
for h in &hits {
let b = open[h.block].as_ref().expect("a hit implies an open block");
b.emit_row(j, h.row);
}
});
Ok(Results {
json: j.into_string(),
stats,
next: (hits.len() == q.limit)
.then(|| hits.last().map(|h| cursor(&refs[h.block], h)))
.flatten(),
})
}
const MAX_FANOUT: usize = 16;
static SPARE: std::sync::LazyLock<AtomicUsize> = std::sync::LazyLock::new(|| {
AtomicUsize::new(
std::thread::available_parallelism()
.map_or(1, |n| n.get())
.saturating_sub(1)
.min(MAX_FANOUT),
)
});
pub(crate) struct Helpers(usize);
impl Helpers {
pub(crate) fn claim(want: usize) -> Helpers {
let mut got = 0;
let _ = SPARE.fetch_update(Relaxed, Relaxed, |n| {
got = n.min(want);
(got > 0).then(|| n - got)
});
Helpers(got)
}
}
impl Drop for Helpers {
fn drop(&mut self) {
SPARE.fetch_add(self.0, Relaxed);
}
}
pub(crate) struct Scan<'a> {
q: &'a Search,
needle: Option<[u8; 16]>,
probes: Vec<AttrProbe>,
ranges: Vec<crate::zone::Probe>,
after: std::cmp::Reverse<(i64, u32, u64, u32)>,
}
impl<'a> Scan<'a> {
pub(crate) fn new(q: &'a Search) -> Scan<'a> {
Scan {
q,
needle: trace_needle(q),
probes: attr_probes(q),
ranges: range_probes(q),
after: q.after.map_or(
std::cmp::Reverse((i64::MAX, u32::MAX, u64::MAX, u32::MAX)),
|c| c.key(),
),
}
}
pub(crate) fn wave(&self, refs: &[Src<'_>], from: usize, to: usize) -> Vec<Result<Scanned>> {
let helpers = Helpers::claim(to - from - 1);
if helpers.0 == 0 {
return vec![self.block(from, &refs[from])];
}
let to = from + 1 + helpers.0;
std::thread::scope(|s| {
let rest: Vec<_> = (from + 1..to)
.map(|k| s.spawn(move || self.block(k, &refs[k])))
.collect();
let mut out = Vec::with_capacity(to - from);
out.push(self.block(from, &refs[from]));
out.extend(
rest.into_iter()
.map(|h| h.join().unwrap_or_else(|e| std::panic::resume_unwind(e))),
);
out
})
}
pub(crate) fn block(&self, i: usize, bref: &Src<'_>) -> Result<Scanned> {
let no_trace = bref.dir.zip(self.needle.as_ref()).is_some_and(|(dir, id)| {
std::fs::read(dir.join(crate::bloom::TRACE_IDX))
.is_ok_and(|f| !crate::bloom::may_contain(&f, id))
});
if no_trace {
return Ok((None, Vec::new()));
}
let no_attr = !self.probes.is_empty()
&& bref.dir.is_some_and(|dir| {
std::fs::read(dir.join(crate::bloom::ATTR_IDX)).is_ok_and(|bytes| {
crate::bloom::Filter::open(&bytes)
.is_some_and(|f| !self.probes.iter().all(|p| p.maybe(&f)))
})
});
if no_attr {
return Ok((None, Vec::new()));
}
let out_of_range = !self.ranges.is_empty()
&& bref.dir.is_some_and(|dir| {
std::fs::read(dir.join(crate::zone::ZONE_IDX)).is_ok_and(|bytes| {
crate::zone::Map::open(&bytes)
.is_some_and(|m| !self.ranges.iter().all(|p| p.maybe(&m)))
})
});
if out_of_range {
return Ok((None, Vec::new()));
}
let Some(b) = Block::open(bref, self.q.signal)? else {
return Ok((None, Vec::new()));
};
let sel = b.select(self.q, bref);
let hits = b.time().map_or_else(Vec::new, |time| {
sel.iter()
.filter_map(|&row| {
let h = Hit {
ts: time[row as usize],
block: i,
row,
};
(cursor(bref, &h).key() > self.after).then_some(h)
})
.collect()
});
Ok((Some(b), hits))
}
}
pub(crate) type Scanned = (Option<Block>, Vec<Hit>);
fn cursor(bref: &Src, h: &Hit) -> Cursor {
Cursor {
ts: h.ts,
node: bref.node,
seq: bref.seq,
row: h.row,
}
}
pub(crate) struct Block {
signal: Signal,
pub(crate) root: RecordBatch,
attrs: Option<RecordBatch>,
resource_attrs: Option<RecordBatch>,
scope_attrs: Option<RecordBatch>,
children: Vec<Child>,
}
struct Child {
label: &'static str,
rows: RecordBatch,
attrs: Option<RecordBatch>,
by_parent: Vec<Vec<u32>>,
parent_of_id: Vec<u32>,
}
fn child_tables(signal: Signal) -> &'static [(&'static str, &'static str, &'static str)] {
match signal {
Signal::Logs => &[],
Signal::Traces => &[
("events", "span_events", "span_event_attrs"),
("links", "span_links", "span_link_attrs"),
],
}
}
impl Block {
fn open(bref: &Src, signal: Signal) -> Result<Option<Block>> {
let load = |name: &str| bref.load(name);
let Some(root) = load(signal.root())? else {
return Ok(None);
};
let mut children = Vec::new();
for &(label, table, attrs) in child_tables(signal) {
if let Some(rows) = load(table)? {
children.push(Child {
label,
by_parent: crate::series::index_by_parent(&rows),
parent_of_id: index_parent_of_id(&rows),
rows,
attrs: load(attrs)?,
});
}
}
Ok(Some(Block {
signal,
attrs: load(signal.attrs())?,
resource_attrs: load("resource_attrs")?,
scope_attrs: load("scope_attrs")?,
children,
root,
}))
}
fn time(&self) -> Option<&[i64]> {
self.root
.column_by_name(self.signal.time_col())
.map(|c| &**c.as_primitive::<TimestampNanosecondType>().values())
}
fn select(&self, q: &Search, bref: &Src) -> Vec<u32> {
let n = self.root.num_rows();
let Some(time) = self.time() else {
return Vec::new();
};
let mut sel: Vec<u32> = if q.from <= bref.min_ts && q.to >= bref.max_ts {
(0..n as u32).collect()
} else {
(0..n as u32)
.filter(|&i| (q.from..=q.to).contains(&time[i as usize]))
.collect()
};
for term in &q.terms {
if sel.is_empty() {
break;
}
match &term.target {
Target::Field(name) => {
let pred = self
.root
.column_by_name(name)
.and_then(|c| field_pred(c.as_ref(), term.op, &term.value));
match pred {
Some(p) => sel.retain(|&i| p(i)),
None => sel.clear(),
}
}
Target::Attr(key) => {
let matched = self.attr_rows(key, term.op, &term.value, n);
sel.retain(|&i| matched[i as usize]);
}
}
}
sel
}
fn attr_rows(&self, key: &str, op: Op, value: &Value, n: usize) -> Vec<bool> {
let mut out = vec![false; n];
if let Some(a) = &self.attrs {
for pid in attr_parents(a, key, op, value) {
if let Some(slot) = out.get_mut(pid as usize) {
*slot = true;
}
}
}
for (table, fk) in [
(&self.resource_attrs, "resource_id"),
(&self.scope_attrs, "scope_id"),
] {
let (Some(a), Some(col)) = (table, self.root.column_by_name(fk)) else {
continue;
};
let ids = attr_parents(a, key, op, value);
let Some(&top) = ids.iter().max() else {
continue;
};
let mut wanted = vec![false; top as usize + 1];
for id in ids {
wanted[id as usize] = true;
}
for (i, &id) in col.as_primitive::<UInt16Type>().values().iter().enumerate() {
if wanted.get(id as usize).copied().unwrap_or(false) {
out[i] = true;
}
}
}
for c in &self.children {
let Some(a) = &c.attrs else { continue };
for id in attr_parents(a, key, op, value) {
if let Some(slot) = c
.parent_of_id
.get(id as usize)
.and_then(|&root| out.get_mut(root as usize))
{
*slot = true;
}
}
}
out
}
fn emit_row(&self, j: &mut Json, row: u32) {
j.obj(|j| {
emit_fields(j, &self.root, row);
j.key("attributes");
emit_attrs(
j,
&[
(&self.resource_attrs, self.fk(row, "resource_id")),
(&self.scope_attrs, self.fk(row, "scope_id")),
(&self.attrs, Some(row)),
],
);
for c in &self.children {
self.emit_children(j, c, row);
}
});
}
fn emit_children(&self, j: &mut Json, c: &Child, row: u32) {
let hits = match c.by_parent.get(row as usize) {
Some(h) if !h.is_empty() => h,
_ => return,
};
j.key(c.label);
j.arr(|j| {
for &r in hits {
let r = r as usize;
j.obj(|j| {
emit_fields(j, &c.rows, r as u32);
let id = c
.rows
.column_by_name("id")
.map(|col| col.as_primitive::<UInt32Type>().value(r));
j.key("attributes");
emit_attrs(j, &[(&c.attrs, id)]);
});
}
});
}
fn fk(&self, row: u32, name: &str) -> Option<u32> {
self.root
.column_by_name(name)
.map(|c| c.as_primitive::<UInt16Type>().value(row as usize) as u32)
}
}
fn index_parent_of_id(b: &RecordBatch) -> Vec<u32> {
let (Some(ids), Some(parents)) = (b.column_by_name("id"), b.column_by_name("parent_id")) else {
return Vec::new();
};
let ids = ids.as_primitive::<UInt32Type>().values();
let parents = parents.as_primitive::<UInt32Type>().values();
let mut out = vec![u32::MAX; ids.iter().copied().max().unwrap_or(0) as usize + 1];
for (&id, &parent) in ids.iter().zip(parents) {
out[id as usize] = parent;
}
out
}
pub(crate) fn emit_fields(j: &mut Json, b: &RecordBatch, row: u32) {
for (i, f) in b.schema().fields().iter().enumerate() {
if matches!(
f.name().as_str(),
"id" | "parent_id" | "resource_id" | "scope_id"
) {
continue;
}
let col = b.column(i);
if col.is_null(row as usize) {
continue;
}
j.key(f.name());
if f.name() == "body_ser" {
emit_any(j, col.as_binary::<i32>().value(row as usize));
} else {
emit_value(j, col.as_ref(), row as usize);
}
}
}
fn emit_attrs(j: &mut Json, levels: &[(&Option<RecordBatch>, Option<u32>)]) {
j.obj(|j| {
let mut merged: Vec<(&str, &RecordBatch, usize)> = Vec::new();
for &(table, parent) in levels {
let (Some(a), Some(parent)) = (table, parent) else {
continue;
};
let parents = a.column(0).as_primitive::<UInt32Type>().values();
for r in (0..a.num_rows()).filter(|&r| parents[r] == parent) {
merged.push((attr_key(a, r), a, r));
}
}
merged.sort_by_key(|(k, _, _)| *k);
for (i, &(k, a, r)) in merged.iter().enumerate() {
if merged.get(i + 1).is_some_and(|nxt| nxt.0 == k) {
continue;
}
j.key(k);
emit_attr(j, a, r);
}
});
}
pub(crate) fn attr_parents(a: &RecordBatch, key: &str, op: Op, value: &Value) -> Vec<u32> {
let keys = a.column(1).as_dictionary::<UInt16Type>();
let Some(code) = dict_index(keys.values().as_string::<i32>(), key) else {
return Vec::new();
};
let codes = keys.keys().values();
let parents = a.column(0).as_primitive::<UInt32Type>();
let types = a.column(2).as_primitive::<UInt8Type>();
(0..a.num_rows())
.filter(|&i| codes[i] == code && attr_matches(a, types.value(i), i, op, value))
.map(|i| parents.value(i))
.collect()
}
pub(crate) fn attr_key(a: &RecordBatch, row: usize) -> &str {
let d = a.column(1).as_dictionary::<UInt16Type>();
d.values()
.as_string::<i32>()
.value(d.keys().value(row) as usize)
}
fn attr_matches(a: &RecordBatch, ty: u8, row: usize, op: Op, v: &Value) -> bool {
const STR: u8 = AttrType::Str as u8;
const INT: u8 = AttrType::Int as u8;
const DOUBLE: u8 = AttrType::Double as u8;
const BOOL: u8 = AttrType::Bool as u8;
match ty {
STR => {
let strs = crate::attrs::str_column(a);
let s = strs.value(row);
match op {
Op::Eq | Op::Ne => op.test_ord(s.cmp(canon(v).as_str())),
Op::Contains => s.contains(canon(v).as_str()),
_ => match (s.parse::<f64>(), v.as_f64()) {
(Ok(x), Some(y)) if !matches!(v, Value::Str(_)) => {
x.partial_cmp(&y).is_some_and(|o| op.test_ord(o))
}
_ => op.test_ord(s.cmp(canon(v).as_str())),
},
}
}
INT => {
let x = a.column(4).as_primitive::<Int64Type>().value(row);
v.as_i64().is_some_and(|y| op.test_ord(x.cmp(&y)))
}
DOUBLE => {
let x = a.column(5).as_primitive::<Float64Type>().value(row);
v.as_f64()
.and_then(|y| x.partial_cmp(&y))
.is_some_and(|o| op.test_ord(o))
}
BOOL => {
let x = a.column(6).as_boolean().value(row);
v.as_bool().is_some_and(|y| op.test_ord(x.cmp(&y)))
}
_ => false,
}
}
pub(crate) fn dict_index(values: &StringArray, needle: &str) -> Option<u16> {
(0..values.len())
.find(|&i| values.value(i) == needle)
.map(|i| i as u16)
}
fn field_pred<'a>(col: &'a dyn Array, op: Op, v: &Value) -> Option<Box<dyn Fn(u32) -> bool + 'a>> {
macro_rules! int_col {
($t:ty) => {{
let a = col.as_primitive::<$t>();
let target = v.as_i64()?;
Some(Box::new(move |i: u32| {
!a.is_null(i as usize) && op.test_ord((a.value(i as usize) as i64).cmp(&target))
}) as Box<dyn Fn(u32) -> bool + 'a>)
}};
}
match col.data_type() {
DataType::Timestamp(_, _) => int_col!(TimestampNanosecondType),
DataType::Int64 => int_col!(Int64Type),
DataType::Int32 => int_col!(Int32Type),
DataType::UInt64 => int_col!(UInt64Type),
DataType::UInt32 => int_col!(UInt32Type),
DataType::UInt16 => int_col!(UInt16Type),
DataType::UInt8 => int_col!(UInt8Type),
DataType::Float64 => {
let a = col.as_primitive::<Float64Type>();
let target = v.as_f64()?;
Some(Box::new(move |i: u32| {
!a.is_null(i as usize)
&& a.value(i as usize)
.partial_cmp(&target)
.is_some_and(|o| op.test_ord(o))
}))
}
DataType::Boolean => {
let a = col.as_boolean();
let target = v.as_bool()?;
Some(Box::new(move |i: u32| {
!a.is_null(i as usize) && op.test_ord(a.value(i as usize).cmp(&target))
}))
}
DataType::Utf8 => {
let a = col.as_string::<i32>();
let target = v.as_str()?.to_owned();
Some(Box::new(move |i: u32| {
if a.is_null(i as usize) {
return false;
}
let s = a.value(i as usize);
match op {
Op::Contains => s.contains(&target),
_ => op.test_ord(s.cmp(target.as_str())),
}
}))
}
DataType::Dictionary(_, _) => {
let d = col.as_dictionary::<UInt16Type>();
let values = d.values().as_string::<i32>();
let needle = v.as_str()?;
let ok: Vec<bool> = (0..values.len())
.map(|i| {
let s = values.value(i);
match op {
Op::Contains => s.contains(needle),
_ => op.test_ord(s.cmp(needle)),
}
})
.collect();
let codes = d.keys();
Some(Box::new(move |i: u32| {
!codes.is_null(i as usize)
&& ok
.get(codes.value(i as usize) as usize)
.copied()
.unwrap_or(false)
}))
}
DataType::FixedSizeBinary(_) => {
let a = col.as_any().downcast_ref::<FixedSizeBinaryArray>()?;
let target = unhex(v.as_str()?)?;
Some(Box::new(move |i: u32| {
!a.is_null(i as usize) && op.test_ord(a.value(i as usize).cmp(target.as_slice()))
}))
}
_ => None,
}
}
pub fn unhex(s: &str) -> Option<Vec<u8>> {
if s.len() % 2 != 0 {
return None;
}
let b = s.as_bytes();
(0..b.len() / 2)
.map(|i| {
let hi = (b[i * 2] as char).to_digit(16)?;
let lo = (b[i * 2 + 1] as char).to_digit(16)?;
Some((hi * 16 + lo) as u8)
})
.collect()
}
fn emit_value(j: &mut Json, col: &dyn Array, row: usize) {
match col.data_type() {
DataType::Timestamp(_, _) => {
j.i64_str(col.as_primitive::<TimestampNanosecondType>().value(row));
}
DataType::Int64 => j.i64_str(col.as_primitive::<Int64Type>().value(row)),
DataType::Int32 => j.i64(col.as_primitive::<Int32Type>().value(row) as i64),
DataType::UInt64 => j.u64_str(col.as_primitive::<UInt64Type>().value(row)),
DataType::UInt32 => j.u64(col.as_primitive::<UInt32Type>().value(row) as u64),
DataType::UInt16 => j.u64(col.as_primitive::<UInt16Type>().value(row) as u64),
DataType::UInt8 => j.u64(col.as_primitive::<UInt8Type>().value(row) as u64),
DataType::Float64 => j.f64(col.as_primitive::<Float64Type>().value(row)),
DataType::Boolean => j.bool(col.as_boolean().value(row)),
DataType::Utf8 => j.str(col.as_string::<i32>().value(row)),
DataType::Binary => j.hex(col.as_binary::<i32>().value(row)),
DataType::FixedSizeBinary(_) => j.hex(col.as_fixed_size_binary().value(row)),
DataType::Dictionary(_, _) => {
let d = col.as_dictionary::<UInt16Type>();
j.str(
d.values()
.as_string::<i32>()
.value(d.keys().value(row) as usize),
);
}
DataType::List(_) => {
let inner = col.as_list::<i32>().value(row);
j.arr(|j| {
for i in 0..inner.len() {
if inner.is_null(i) {
j.null();
} else {
emit_value(j, inner.as_ref(), i);
}
}
});
}
_ => j.null(),
}
}
pub(crate) fn emit_attr(j: &mut Json, a: &RecordBatch, row: usize) {
const STR: u8 = AttrType::Str as u8;
const INT: u8 = AttrType::Int as u8;
const DOUBLE: u8 = AttrType::Double as u8;
const BOOL: u8 = AttrType::Bool as u8;
const BYTES: u8 = AttrType::Bytes as u8;
const SLICE: u8 = AttrType::Slice as u8;
const MAP: u8 = AttrType::Map as u8;
match a.column(2).as_primitive::<UInt8Type>().value(row) {
STR => j.str(crate::attrs::str_column(a).value(row)),
INT => j.i64_str(a.column(4).as_primitive::<Int64Type>().value(row)),
DOUBLE => j.f64(a.column(5).as_primitive::<Float64Type>().value(row)),
BOOL => j.bool(a.column(6).as_boolean().value(row)),
BYTES => j.hex(a.column(7).as_binary::<i32>().value(row)),
SLICE | MAP => emit_any(j, a.column(8).as_binary::<i32>().value(row)),
_ => j.null(),
}
}
pub(crate) fn emit_any(j: &mut Json, bytes: &[u8]) {
match AnyValue::decode(bytes) {
Ok(v) => emit_any_value(j, v.value.as_ref()),
Err(_) => j.null(),
}
}
fn emit_any_value(j: &mut Json, v: Option<&mira_proto::common::v1::any_value::Value>) {
use mira_proto::common::v1::any_value::Value as Av;
match v {
None => j.null(),
Some(Av::StringValue(s)) => j.str(s),
Some(Av::IntValue(i)) => j.i64_str(*i),
Some(Av::DoubleValue(d)) => j.f64(*d),
Some(Av::BoolValue(b)) => j.bool(*b),
Some(Av::BytesValue(b)) => j.hex(b),
Some(Av::ArrayValue(a)) => j.arr(|j| {
for e in &a.values {
emit_any_value(j, e.value.as_ref());
}
}),
Some(Av::KvlistValue(m)) => j.obj(|j| {
for e in &m.values {
j.key(&e.key);
emit_any_value(j, e.value.as_ref().and_then(|v| v.value.as_ref()));
}
}),
}
}
#[cfg(test)]
mod tests {
use super::*;
use arrow_array::builder::StringDictionaryBuilder;
use arrow_array::{
BinaryArray, BooleanArray, Float64Array, Int32Array, Int64Array, ListArray,
TimestampNanosecondArray, UInt8Array, UInt16Array, UInt32Array, UInt64Array,
};
use std::sync::Arc;
fn hits(col: &dyn Array, op: Op, v: &Value) -> Vec<u32> {
match field_pred(col, op, v) {
Some(p) => (0..col.len() as u32).filter(|&i| p(i)).collect(),
None => vec![u32::MAX],
}
}
#[test]
fn every_column_type_compares_the_same_way() {
let cols: Vec<(&str, Arc<dyn Array>)> = vec![
(
"timestamp",
Arc::new(TimestampNanosecondArray::from(vec![
Some(1),
Some(2),
Some(3),
None,
])),
),
(
"i64",
Arc::new(Int64Array::from(vec![Some(1), Some(2), Some(3), None])),
),
(
"i32",
Arc::new(Int32Array::from(vec![Some(1), Some(2), Some(3), None])),
),
(
"u64",
Arc::new(UInt64Array::from(vec![Some(1), Some(2), Some(3), None])),
),
(
"u32",
Arc::new(UInt32Array::from(vec![Some(1), Some(2), Some(3), None])),
),
(
"u16",
Arc::new(UInt16Array::from(vec![Some(1), Some(2), Some(3), None])),
),
(
"u8",
Arc::new(UInt8Array::from(vec![Some(1), Some(2), Some(3), None])),
),
(
"f64",
Arc::new(Float64Array::from(vec![
Some(1.0),
Some(2.0),
Some(3.0),
None,
])),
),
];
for (name, col) in &cols {
let c = col.as_ref();
assert_eq!(hits(c, Op::Eq, &Value::Int(2)), [1], "{name} eq");
assert_eq!(hits(c, Op::Ne, &Value::Int(2)), [0, 2], "{name} ne");
assert_eq!(hits(c, Op::Lt, &Value::Int(2)), [0], "{name} lt");
assert_eq!(hits(c, Op::Lte, &Value::Int(2)), [0, 1], "{name} lte");
assert_eq!(hits(c, Op::Gt, &Value::Int(2)), [2], "{name} gt");
assert_eq!(hits(c, Op::Gte, &Value::Int(2)), [1, 2], "{name} gte");
assert!(!hits(c, Op::Ne, &Value::Int(9)).contains(&3), "{name} null");
assert_eq!(hits(c, Op::Eq, &Value::Str("2".into())), [1], "{name} str");
assert_eq!(
hits(c, Op::Eq, &Value::Bool(true)),
[u32::MAX],
"{name} bool"
);
}
assert_eq!(
hits(cols[1].1.as_ref(), Op::Gt, &Value::Double(1.5)),
[u32::MAX]
);
assert_eq!(hits(cols[1].1.as_ref(), Op::Gt, &Value::Double(2.0)), [2]);
let f = Float64Array::from(vec![Some(1.5), Some(f64::NAN)]);
assert_eq!(hits(&f, Op::Gt, &Value::Double(1.0)), [0]);
assert_eq!(
hits(&f, Op::Eq, &Value::Double(f64::NAN)),
Vec::<u32>::new()
);
let b = BooleanArray::from(vec![Some(true), Some(false), None]);
assert_eq!(hits(&b, Op::Eq, &Value::Bool(true)), [0]);
assert_eq!(hits(&b, Op::Eq, &Value::Str("false".into())), [1]);
assert_eq!(hits(&b, Op::Eq, &Value::Int(1)), [u32::MAX]);
let s = StringArray::from(vec![Some("alpha"), Some("beta"), None]);
assert_eq!(hits(&s, Op::Eq, &Value::Str("beta".into())), [1]);
assert_eq!(hits(&s, Op::Contains, &Value::Str("et".into())), [1]);
assert_eq!(hits(&s, Op::Lt, &Value::Str("b".into())), [0]);
assert_eq!(hits(&s, Op::Eq, &Value::Int(1)), [u32::MAX]);
let mut d = StringDictionaryBuilder::<UInt16Type>::new();
for v in ["ERROR", "INFO", "ERROR"] {
d.append_value(v);
}
let d = d.finish();
assert_eq!(hits(&d, Op::Eq, &Value::Str("ERROR".into())), [0, 2]);
assert_eq!(
hits(&d, Op::Eq, &Value::Str("TRACE".into())),
Vec::<u32>::new()
);
let ids = FixedSizeBinaryArray::try_from_iter([[1u8, 2], [3, 4]].into_iter()).unwrap();
assert_eq!(hits(&ids, Op::Eq, &Value::Str("0102".into())), [0]);
assert_eq!(hits(&ids, Op::Gt, &Value::Str("0102".into())), [1]);
assert_eq!(hits(&ids, Op::Eq, &Value::Str("zz".into())), [u32::MAX]);
assert_eq!(hits(&ids, Op::Eq, &Value::Str("010".into())), [u32::MAX]);
let l = ListArray::from_iter_primitive::<Int64Type, _, _>(vec![Some(vec![Some(1)])]);
assert_eq!(hits(&l, Op::Eq, &Value::Int(1)), [u32::MAX]);
}
#[test]
fn every_column_type_materializes_as_the_json_type_it_is() {
let cell = |col: &dyn Array| {
let mut j = Json::new();
j.arr(|j| emit_value(j, col, 0));
let s = j.into_string();
s[1..s.len() - 1].to_owned()
};
let l = ListArray::from_iter_primitive::<Int64Type, _, _>(vec![Some(vec![
Some(1),
None,
Some(3),
])]);
let mut d = StringDictionaryBuilder::<UInt16Type>::new();
d.append_value("ERROR");
let cases: Vec<(Arc<dyn Array>, &str)> = vec![
(
Arc::new(TimestampNanosecondArray::from(vec![
1_700_000_000_000_000_001i64,
])),
r#""1700000000000000001""#,
),
(Arc::new(Int64Array::from(vec![-7i64])), r#""-7""#),
(Arc::new(Int32Array::from(vec![-7i32])), "-7"),
(
Arc::new(UInt64Array::from(vec![u64::MAX])),
r#""18446744073709551615""#,
),
(Arc::new(UInt32Array::from(vec![7u32])), "7"),
(Arc::new(UInt16Array::from(vec![7u16])), "7"),
(Arc::new(UInt8Array::from(vec![7u8])), "7"),
(Arc::new(Float64Array::from(vec![0.5f64])), "0.5"),
(Arc::new(BooleanArray::from(vec![true])), "true"),
(Arc::new(StringArray::from(vec!["a\"b"])), r#""a\"b""#),
(
Arc::new(BinaryArray::from(vec![&b"\xab\xcd"[..]])),
r#""abcd""#,
),
(
Arc::new(
FixedSizeBinaryArray::try_from_iter([[0xabu8, 0xcd]].into_iter()).unwrap(),
),
r#""abcd""#,
),
(Arc::new(d.finish()), r#""ERROR""#),
(Arc::new(l), r#"["1",null,"3"]"#),
];
for (col, want) in cases {
assert_eq!(cell(col.as_ref()), want, "{:?}", col.data_type());
}
let m = arrow_array::Int8Array::from(vec![1i8]);
assert_eq!(cell(&m), "null");
}
#[test]
fn a_scalar_coerces_to_what_the_column_needs_or_to_nothing() {
for (s, want) in [
("eq", Op::Eq),
("=", Op::Eq),
("==", Op::Eq),
("ne", Op::Ne),
("!=", Op::Ne),
("lt", Op::Lt),
("<", Op::Lt),
("lte", Op::Lte),
("<=", Op::Lte),
("gt", Op::Gt),
(">", Op::Gt),
("gte", Op::Gte),
(">=", Op::Gte),
("contains", Op::Contains),
("~", Op::Contains),
] {
assert_eq!(Op::parse(s), Some(want), "{s}");
}
assert_eq!(Op::parse("=~"), None);
use std::cmp::Ordering::*;
for ord in [Less, Equal, Greater] {
assert!(!Op::Contains.test_ord(ord));
}
assert_eq!(Value::Int(3).as_i64(), Some(3));
assert_eq!(Value::Double(3.0).as_i64(), Some(3));
assert_eq!(Value::Double(3.5).as_i64(), None);
assert_eq!(Value::Str("3".into()).as_i64(), Some(3));
assert_eq!(Value::Str("3.5".into()).as_i64(), None);
assert_eq!(Value::Bool(true).as_i64(), None);
assert_eq!(Value::Int(3).as_f64(), Some(3.0));
assert_eq!(Value::Double(3.5).as_f64(), Some(3.5));
assert_eq!(Value::Str("3.5".into()).as_f64(), Some(3.5));
assert_eq!(Value::Str("x".into()).as_f64(), None);
assert_eq!(Value::Bool(true).as_f64(), None);
assert_eq!(Value::Bool(false).as_bool(), Some(false));
assert_eq!(Value::Str("true".into()).as_bool(), Some(true));
assert_eq!(Value::Str("false".into()).as_bool(), Some(false));
assert_eq!(Value::Str("TRUE".into()).as_bool(), None);
assert_eq!(Value::Int(1).as_bool(), None);
assert_eq!(Value::Double(1.0).as_bool(), None);
assert_eq!(Value::Str("x".into()).as_str(), Some("x"));
assert_eq!(Value::Int(1).as_str(), None);
assert_eq!(unhex("0aFf"), Some(vec![0x0a, 0xff]));
assert_eq!(unhex(""), Some(vec![]));
assert_eq!(unhex("abc"), None);
assert_eq!(unhex("0g"), None);
assert_eq!(unhex("0 1"), None);
let vals = StringArray::from(vec!["a", "b"]);
assert_eq!(dict_index(&vals, "b"), Some(1));
assert_eq!(dict_index(&vals, "c"), None);
}
#[test]
fn every_any_value_arm_renders_and_damage_renders_as_null() {
use mira_proto::common::v1::any_value::Value as Av;
use mira_proto::common::v1::{ArrayValue, KeyValue, KeyValueList};
let render = |v: Option<Av>| {
let mut j = Json::new();
emit_any(&mut j, &AnyValue { value: v }.encode_to_vec());
j.into_string()
};
assert_eq!(render(None), "null");
assert_eq!(render(Some(Av::StringValue("s".into()))), r#""s""#);
assert_eq!(render(Some(Av::IntValue(-1))), r#""-1""#);
assert_eq!(render(Some(Av::DoubleValue(0.5))), "0.5");
assert_eq!(render(Some(Av::BoolValue(true))), "true");
assert_eq!(
render(Some(Av::BytesValue(vec![0xbe, 0xef].into()))),
r#""beef""#
);
assert_eq!(
render(Some(Av::ArrayValue(ArrayValue {
values: vec![
AnyValue { value: None },
AnyValue {
value: Some(Av::KvlistValue(KeyValueList {
values: vec![KeyValue {
key: "k".into(),
value: Some(AnyValue {
value: Some(Av::IntValue(2)),
}),
}],
})),
},
],
}))),
r#"[null,{"k":"2"}]"#
);
let mut j = Json::new();
emit_any(&mut j, &[0xff, 0xff, 0xff]);
assert_eq!(j.into_string(), "null");
}
#[test]
fn a_child_row_finds_its_parent_by_its_own_id() {
let ids: Arc<dyn Array> = Arc::new(UInt32Array::from(vec![2u32, 0]));
let parents: Arc<dyn Array> = Arc::new(UInt32Array::from(vec![7u32, 5]));
let b = RecordBatch::try_from_iter([("id", ids.clone()), ("parent_id", parents)]).unwrap();
let idx = index_parent_of_id(&b);
assert_eq!(idx.len(), 3, "dense from zero to the largest id");
assert_eq!(idx[0], 5);
assert_eq!(idx[1], u32::MAX, "an id no row claims has no parent");
assert_eq!(idx[2], 7);
let orphan = RecordBatch::try_from_iter([("id", ids)]).unwrap();
assert!(index_parent_of_id(&orphan).is_empty());
}
#[test]
fn a_scan_narrows_by_time_then_by_terms_and_the_most_specific_level_wins() {
use mira_proto::collector::logs::v1::ExportLogsServiceRequest;
use mira_proto::common::v1::any_value::Value as Av;
use mira_proto::common::v1::{
AnyValue, ArrayValue, InstrumentationScope, KeyValue, KeyValueList,
};
use mira_proto::logs::v1::{LogRecord, ResourceLogs, ScopeLogs};
use mira_proto::resource::v1::Resource;
let kv = |k: &str, v: Av| KeyValue {
key: k.into(),
value: Some(AnyValue { value: Some(v) }),
};
let base = 1_700_000_000_000_000_000u64;
let req = ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
resource: Some(Resource {
attributes: vec![
kv("service.name", Av::StringValue("checkout".into())),
kv("deploy.env", Av::StringValue("prod".into())),
],
..Default::default()
}),
scope_logs: vec![ScopeLogs {
scope: Some(InstrumentationScope {
name: "mira.test".into(),
attributes: vec![kv("scope.kind", Av::StringValue("lib".into()))],
..Default::default()
}),
log_records: (0..3)
.map(|i| LogRecord {
time_unix_nano: base + i * 1_000_000_000,
severity_number: 9 + i as i32,
severity_text: "INFO".into(),
event_name: if i == 0 {
"user.login".into()
} else {
String::new()
},
body: Some(AnyValue {
value: Some(if i == 2 {
Av::KvlistValue(KeyValueList {
values: vec![kv(
"msg",
Av::StringValue("structured body".into()),
)],
})
} else {
Av::StringValue(format!("line {i}"))
}),
}),
attributes: vec![
kv("deploy.env", Av::StringValue("canary".into())),
kv("attempt", Av::IntValue(i as i64)),
kv("payload", Av::BytesValue(vec![0xde, 0xad].into())),
kv(
"tags",
Av::ArrayValue(ArrayValue {
values: vec![
AnyValue {
value: Some(Av::StringValue("a".into())),
},
AnyValue {
value: Some(Av::IntValue(7)),
},
],
}),
),
kv(
"gen_ai.input.messages",
Av::KvlistValue(KeyValueList {
values: vec![kv("role", Av::StringValue("user".into()))],
}),
),
KeyValue {
key: "trace.hint".into(),
value: None,
},
],
..Default::default()
})
.collect(),
..Default::default()
}],
..Default::default()
}],
};
let dir = std::env::temp_dir().join(format!("mira-scan-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut b = crate::logs::LogsBuilder::new();
b.append_request(&req).unwrap();
let sealed = b.finish().unwrap();
let bref =
crate::block::publish(&dir, "logs", crate::block::node_id("a"), 0, 0, &sealed).unwrap();
let base = base as i64;
let scan = |from: i64, to: i64, terms: Vec<Term>| {
search(
&dir,
&Search {
signal: Signal::Logs,
from,
to,
terms,
limit: 100,
after: None,
},
)
.unwrap()
};
let field = |n: &str, op: Op, v: Value| Term {
target: Target::Field(n.into()),
op,
value: v,
};
let attr = |n: &str, op: Op, v: Value| Term {
target: Target::Attr(n.into()),
op,
value: v,
};
let r = scan(base + 500_000_000, base + 1_500_000_000, vec![]);
assert_eq!(r.stats.rows_matched, 1, "{}", r.json);
assert!(r.json.contains("line 1"), "{}", r.json);
let all = base + 10_000_000_000;
assert_eq!(scan(base, all, vec![]).stats.rows_matched, 3);
assert_eq!(
scan(
base,
all,
vec![field("duration_nano", Op::Gt, Value::Int(0))]
)
.stats
.rows_matched,
0
);
assert_eq!(
scan(
base,
all,
vec![field("severity_number", Op::Eq, Value::Bool(true))]
)
.stats
.rows_matched,
0
);
assert_eq!(
scan(
base,
all,
vec![
field("severity_number", Op::Gt, Value::Int(99)),
attr("service.name", Op::Eq, Value::Str("checkout".into())),
]
)
.stats
.rows_matched,
0
);
for (k, v) in [("service.name", "checkout"), ("scope.kind", "lib")] {
assert_eq!(
scan(base, all, vec![attr(k, Op::Eq, Value::Str(v.into()))])
.stats
.rows_matched,
3,
"{k}"
);
}
assert_eq!(
scan(base, all, vec![attr("attempt", Op::Gte, Value::Int(1))])
.stats
.rows_matched,
2
);
assert_eq!(
scan(
base,
all,
vec![attr("payload", Op::Eq, Value::Str("dead".into()))]
)
.stats
.rows_matched,
0
);
let r = scan(
base,
all,
vec![attr("payload", Op::Contains, Value::Str("dead".into()))],
);
assert_eq!(r.stats.blocks_scanned, 1, "no probe, so the block is read");
assert_eq!(r.stats.rows_matched, 0);
let row = scan(base, all, vec![]).json;
assert!(row.contains(r#""deploy.env":"canary""#), "{row}");
assert!(!row.contains("prod"), "{row}");
assert!(row.contains(r#""service.name":"checkout""#), "{row}");
assert!(row.contains(r#""payload":"dead""#), "{row}");
assert!(row.contains(r#""tags":["a","7"]"#), "{row}");
assert!(
row.contains(r#""gen_ai.input.messages":{"role":"user"}"#),
"{row}"
);
assert!(
row.contains(r#""body_ser":{"msg":"structured body"}"#),
"{row}"
);
assert!(row.contains(r#""event_name":"user.login""#), "{row}");
assert_eq!(row.matches("event_name").count(), 1, "{row}");
assert!(row.contains(r#""trace.hint":null"#), "{row}");
let table = bref.dir.join("logs.arrow");
let mut old = block::open_table_opt(&table).unwrap().unwrap().batches[0].clone();
old.remove_column(old.schema().index_of("event_name").unwrap());
let staged = bref.dir.join("logs.arrow.new");
crate::block::write_table(&staged, &old).unwrap();
std::fs::rename(&staged, &table).unwrap();
let r = scan(base, all, vec![]);
assert_eq!(r.stats.rows_matched, 3, "{}", r.json);
assert!(!r.json.contains("event_name"), "{}", r.json);
assert!(
r.json.contains(r#""body_ser":{"msg":"structured body"}"#),
"{}",
r.json
);
let mut old = block::open_table_opt(&table).unwrap().unwrap().batches[0].clone();
old.remove_column(old.schema().index_of("time_unix_nano").unwrap());
crate::block::write_table(&staged, &old).unwrap();
std::fs::rename(&staged, &table).unwrap();
let r = scan(base, all, vec![]);
assert_eq!(r.json, "[]");
assert_eq!(r.stats.rows_matched, 0);
std::fs::remove_file(bref.dir.join("logs.arrow")).unwrap();
let r = scan(base, all, vec![]);
assert_eq!((r.stats.blocks_total, r.stats.blocks_scanned), (1, 0));
assert_eq!(r.json, "[]");
}
#[test]
fn an_attribute_on_a_span_event_or_link_selects_the_span_it_hangs_off() {
use mira_proto::collector::trace::v1::ExportTraceServiceRequest;
use mira_proto::common::v1::any_value::Value as Av;
use mira_proto::common::v1::{AnyValue, InstrumentationScope, KeyValue};
use mira_proto::resource::v1::Resource;
use mira_proto::trace::v1::span::{Event, Link};
use mira_proto::trace::v1::{ResourceSpans, ScopeSpans, Span};
let kv = |k: &str, v: &str| KeyValue {
key: k.into(),
value: Some(AnyValue {
value: Some(Av::StringValue(v.into())),
}),
};
let event = |name: &str, attrs: Vec<KeyValue>| Event {
time_unix_nano: 10_050,
name: name.into(),
attributes: attrs,
..Default::default()
};
let span = |i: u64, events: Vec<Event>, links: Vec<Link>| Span {
trace_id: vec![1u8; 16].into(),
span_id: vec![i as u8 + 1; 8].into(),
name: format!("span {i}"),
start_time_unix_nano: 10_000 + i,
end_time_unix_nano: 10_100 + i,
attributes: vec![kv("http.method", "GET")],
events,
links,
..Default::default()
};
let req = ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: Some(Resource {
attributes: vec![kv("service.name", "checkout")],
..Default::default()
}),
scope_spans: vec![ScopeSpans {
scope: Some(InstrumentationScope {
name: "mira.test".into(),
..Default::default()
}),
spans: vec![
span(
0,
vec![
event("cache.miss", vec![kv("cache.key", "cart:7")]),
event("retrying", vec![]),
],
vec![],
),
span(
1,
vec![event(
"exception",
vec![kv("exception.type", "NullPointerException")],
)],
vec![],
),
span(
2,
vec![],
vec![Link {
trace_id: vec![9u8; 16].into(),
span_id: vec![8u8; 8].into(),
attributes: vec![kv("link.kind", "follows_from")],
..Default::default()
}],
),
],
..Default::default()
}],
..Default::default()
}],
};
let dir = std::env::temp_dir().join(format!("mira-child-attr-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut b = crate::traces::TracesBuilder::new();
b.append_request(&req).unwrap();
let sealed = b.finish().unwrap();
let bref = crate::block::publish(&dir, "traces", crate::block::node_id("a"), 0, 0, &sealed)
.unwrap();
let find = |key: &str, value: &str| {
search(
&dir,
&Search {
signal: Signal::Traces,
from: 0,
to: i64::MAX,
terms: vec![Term {
target: Target::Attr(key.into()),
op: Op::Eq,
value: Value::Str(value.into()),
}],
limit: 10,
after: None,
},
)
.unwrap()
};
for (key, value, want) in [
("exception.type", "NullPointerException", "span 1"),
("cache.key", "cart:7", "span 0"),
("link.kind", "follows_from", "span 2"),
] {
let r = find(key, value);
assert_eq!(r.stats.rows_matched, 1, "{key}: {}", r.json);
assert!(
r.json.contains(&format!(r#""name":"{want}""#)),
"{}",
r.json
);
}
assert_eq!(find("http.method", "GET").stats.rows_matched, 3);
assert_eq!(find("service.name", "checkout").stats.rows_matched, 3);
assert_eq!(find("exception.type", "IOError").stats.rows_matched, 0);
let r = find("exception.type", "NullPointerException");
assert!(
r.json
.contains(r#""attributes":{"exception.type":"NullPointerException"}"#),
"{}",
r.json
);
let strip = |table: &str, col: &str| {
let path = bref.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 = bref.dir.join(format!("{table}.staged"));
crate::block::write_table(&staged, &b).unwrap();
std::fs::rename(&staged, &path).unwrap();
};
strip("span_events", "parent_id");
let r = find("exception.type", "NullPointerException");
assert_eq!(r.stats.rows_matched, 0, "{}", r.json);
assert!(
!r.json.contains("span"),
"an unjoinable event picked a span"
);
assert_eq!(find("link.kind", "follows_from").stats.rows_matched, 1);
assert_eq!(find("http.method", "GET").stats.rows_matched, 3);
strip("spans", "start_time_unix_nano");
let r = find("http.method", "GET");
assert_eq!(r.stats.blocks_scanned, 1, "the block was still opened");
assert_eq!(r.json, "[]");
assert_eq!(r.stats.rows_matched, 0);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn an_attribute_v0_cannot_filter_still_renders_and_still_matches_nothing() {
use mira_proto::collector::logs::v1::ExportLogsServiceRequest;
use mira_proto::common::v1::any_value::Value as Av;
use mira_proto::common::v1::{AnyValue, ArrayValue, KeyValue, KeyValueList};
use mira_proto::logs::v1::{LogRecord, ResourceLogs, ScopeLogs};
let kv = |k: &str, v: Option<Av>| KeyValue {
key: k.into(),
value: v.map(|v| AnyValue { value: Some(v) }),
};
let req = ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
scope_logs: vec![ScopeLogs {
log_records: vec![LogRecord {
time_unix_nano: 1_000,
attributes: vec![
kv("note", None),
kv("payload", Some(Av::BytesValue(vec![0xde, 0xad].into()))),
kv(
"tags",
Some(Av::ArrayValue(ArrayValue {
values: vec![AnyValue {
value: Some(Av::StringValue("a".into())),
}],
})),
),
kv(
"gen_ai.input.messages",
Some(Av::KvlistValue(KeyValueList {
values: vec![kv("role", Some(Av::StringValue("user".into())))],
})),
),
kv("level", Some(Av::StringValue("a dead role".into()))),
],
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
};
let dir = std::env::temp_dir().join(format!("mira-unfilterable-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut b = crate::logs::LogsBuilder::new();
b.append_request(&req).unwrap();
let sealed = b.finish().unwrap();
crate::block::publish(&dir, "logs", crate::block::node_id("a"), 0, 0, &sealed).unwrap();
let scan = |terms: Vec<Term>| {
search(
&dir,
&Search {
signal: Signal::Logs,
from: 0,
to: i64::MAX,
terms,
limit: 10,
after: None,
},
)
.unwrap()
};
let contains = |k: &str, v: &str| {
vec![Term {
target: Target::Attr(k.into()),
op: Op::Contains,
value: Value::Str(v.into()),
}]
};
assert_eq!(scan(contains("level", "dead")).stats.rows_matched, 1);
for (key, needle) in [
("note", ""),
("payload", "dead"),
("tags", "a"),
("gen_ai.input.messages", "role"),
] {
let r = scan(contains(key, needle));
assert_eq!(r.stats.rows_matched, 0, "{key}: {}", r.json);
assert_eq!(r.stats.blocks_scanned, 1, "{key}");
}
let row = scan(Vec::new()).json;
assert!(row.contains(r#""note":null"#), "{row}");
assert!(row.contains(r#""payload":"dead""#), "{row}");
assert!(row.contains(r#""tags":["a"]"#), "{row}");
assert!(
row.contains(r#""gen_ai.input.messages":{"role":"user"}"#),
"{row}"
);
let ids = Arc::new(UInt32Array::from(vec![0u32, 1])) as Arc<dyn Array>;
let only_ids = RecordBatch::try_from_iter(vec![("id", ids.clone())]).unwrap();
let only_parents = RecordBatch::try_from_iter(vec![("parent_id", ids)]).unwrap();
assert!(index_parent_of_id(&only_ids).is_empty());
assert!(index_parent_of_id(&only_parents).is_empty());
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn rows_matched_counts_what_is_left_behind_the_cursor() {
use mira_proto::collector::logs::v1::ExportLogsServiceRequest;
use mira_proto::common::v1::any_value::Value as Av;
use mira_proto::common::v1::{AnyValue, InstrumentationScope};
use mira_proto::logs::v1::{LogRecord, ResourceLogs, ScopeLogs};
let base = 1_700_000_000_000_000_000u64;
let req = ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
scope_logs: vec![ScopeLogs {
scope: Some(InstrumentationScope {
name: "mira.test".into(),
..Default::default()
}),
log_records: (0..5)
.map(|i| LogRecord {
time_unix_nano: base + i * 1_000_000_000,
body: Some(AnyValue {
value: Some(Av::StringValue(format!("line {i}"))),
}),
..Default::default()
})
.collect(),
..Default::default()
}],
..Default::default()
}],
};
let dir = std::env::temp_dir().join(format!("mira-paged-count-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut b = crate::logs::LogsBuilder::new();
b.append_request(&req).unwrap();
let sealed = b.finish().unwrap();
crate::block::publish(&dir, "logs", crate::block::node_id("a"), 0, 0, &sealed).unwrap();
let page = |after: Option<Cursor>| {
search(
&dir,
&Search {
signal: Signal::Logs,
from: 0,
to: i64::MAX,
terms: Vec::new(),
limit: 2,
after,
},
)
.unwrap()
};
let mut cursor = None;
for want in [5, 3, 1] {
let r = page(cursor);
assert_eq!(r.stats.rows_matched, want, "{}", r.json);
cursor = r.next;
}
assert_eq!(cursor, None);
let _ = std::fs::remove_dir_all(&dir);
}
}