use std::path::PathBuf;
use std::sync::{Arc, LazyLock};
use axum::Router;
use axum::extract::State;
use axum::http::{StatusCode, header};
use axum::response::{IntoResponse, Response};
use axum::routing::post;
use tokio::sync::Semaphore;
use yaml_rust2::parser::{Event, EventReceiver, Parser};
use yaml_rust2::{Yaml, YamlLoader};
use mira_core::frame::{self, Expand};
use mira_core::query::{self, Op, Search, Signal, Target, Term, Value};
use mira_core::series::{self, SeriesQuery};
use crate::pipeline;
#[derive(Clone, Default)]
pub struct Api {
pub data_dir: Arc<PathBuf>,
pub open: pipeline::OpenSlots,
pub alerts: Arc<crate::alert::Engine>,
}
impl Api {
pub(crate) async fn open(&self, signal: &str) -> Vec<Arc<mira_core::signal::Open>> {
let Some(i) = pipeline::SIGNALS.iter().position(|s| *s == signal) else {
return Vec::new();
};
self.open[i].fresh().await
}
pub(crate) async fn open_all(&self) -> Vec<Vec<Arc<mira_core::signal::Open>>> {
let mut out = Vec::with_capacity(pipeline::SIGNALS.len());
for s in pipeline::SIGNALS {
out.push(self.open(s).await);
}
out
}
}
pub fn router(api: Api) -> Router {
Router::new()
.route("/api/v1/query", post(query_handler))
.route("/api/v1/metrics/query", post(series_handler))
.route("/api/v1/metrics/names", post(names_handler))
.route("/api/v1/correlate", post(correlate_handler))
.route("/api/v1/map", post(map_handler))
.route("/api/v1/entities", post(entities_handler))
.with_state(api)
}
const MAX_LIMIT: usize = 10_000;
async fn query_handler(State(api): State<Api>, body: String) -> Response {
let q = match parse_search(&body, now_nanos()) {
Ok(q) => q,
Err(e) => return bad_request(&e),
};
let dir = api.data_dir.clone();
let open = api.open(q.signal.dir()).await;
run("rows", move || query::search_open(&dir, &q, &open)).await
}
async fn series_handler(State(api): State<Api>, body: String) -> Response {
let q = match parse_series(&body, now_nanos()) {
Ok(q) => q,
Err(e) => return bad_request(&e),
};
let dir = api.data_dir.clone();
let open = api.open("metrics").await;
run("series", move || series::series_open(&dir, &q, &open)).await
}
async fn names_handler(State(api): State<Api>, body: String) -> Response {
let now = now_nanos();
let doc = match window(if body.trim().is_empty() { "{}" } else { &body }, now) {
Ok(w) => w,
Err(e) => return bad_request(&e),
};
let dir = api.data_dir.clone();
let open = api.open("metrics").await;
run("names", move || {
series::names_open(&dir, doc.0, doc.1, &open)
})
.await
}
async fn correlate_handler(State(api): State<Api>, body: String) -> Response {
let now = now_nanos();
let (q, ops) = match parse_correlate(&body, now) {
Ok(v) => v,
Err(e) => return bad_request(&e),
};
let dir = api.data_dir.clone();
let anchored = api.open(q.signal.dir()).await;
let all = api.open_all().await;
run("frame", move || correlate(&dir, &q, &ops, &anchored, &all)).await
}
pub(crate) fn correlate(
dir: &std::path::Path,
q: &Search,
ops: &[Expand],
anchored: &[Arc<mira_core::signal::Open>],
all: &[Vec<Arc<mira_core::signal::Open>>],
) -> mira_core::error::Result<query::Results> {
let (f, a) = frame::anchor(dir, q, anchored)?;
let traces = all.get(1).map_or(&[][..], Vec::as_slice);
let (f, w) = frame::expand(dir, &f, ops, traces)?;
let names = frame::names_of(dir, &f, all)?;
let mut j = mira_core::json::Json::new();
f.write_json(&mut j, &names);
Ok(query::Results {
json: j.into_string(),
stats: query::Stats {
blocks_total: a.blocks_total + w.blocks_total,
blocks_scanned: a.blocks_scanned + w.blocks_scanned,
rows_scanned: a.rows_scanned + w.rows_scanned,
rows_matched: a.rows_matched + w.rows_matched,
..Default::default()
},
next: None,
cursors: Vec::new(),
})
}
async fn map_handler(State(api): State<Api>, body: String) -> Response {
let now = now_nanos();
let doc = match parse(if body.trim().is_empty() { "{}" } else { &body }) {
Ok(d) => d,
Err(e) => return bad_request(&e),
};
let (from, to, max_spans) = match map_doc(&doc, now) {
Ok(v) => v,
Err(e) => return bad_request(&e),
};
let dir = api.data_dir.clone();
let open = api.open("traces").await;
run("map", move || frame::map(&dir, from, to, max_spans, &open)).await
}
async fn entities_handler(State(api): State<Api>, body: String) -> Response {
let now = now_nanos();
let (from, to) = match window(if body.trim().is_empty() { "{}" } else { &body }, now) {
Ok(w) => w,
Err(e) => return bad_request(&e),
};
let dir = api.data_dir.clone();
let open = api.open_all().await;
run("entities", move || frame::entities(&dir, from, to, &open)).await
}
static SCANS: LazyLock<Semaphore> =
LazyLock::new(|| Semaphore::new(std::thread::available_parallelism().map_or(4, |n| n.get())));
pub(crate) async fn scan<T: Send + 'static>(
f: impl FnOnce() -> T + Send + 'static,
) -> Result<T, tokio::task::JoinError> {
let _permit = SCANS.acquire().await.ok();
tokio::task::spawn_blocking(f).await
}
async fn run(
field: &'static str,
f: impl FnOnce() -> mira_core::error::Result<query::Results> + Send + 'static,
) -> Response {
let t = std::time::Instant::now();
match scan(f).await {
Ok(Ok(r)) => json_ok(envelope(field, &r, t.elapsed())),
Ok(Err(e)) => (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()).into_response(),
Err(e) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("query task panicked: {e}"),
)
.into_response(),
}
}
pub fn envelope(field: &str, r: &query::Results, elapsed: std::time::Duration) -> String {
let next = r
.next
.map(|c| format!(",\"next\":\"{c}\""))
.unwrap_or_default();
let tail = next + &cursors(&r.cursors);
format!(
"{{\"{field}\":{},\"stats\":{{\"blocks_total\":{},\"blocks_scanned\":{},\
\"rows_scanned\":{},\"rows_matched\":{},\"elapsed_us\":{}}}{}}}",
r.json,
r.stats.blocks_total,
r.stats.blocks_scanned,
r.stats.rows_scanned,
r.stats.rows_matched,
elapsed.as_micros(),
tail
)
}
fn cursors(cs: &[mira_core::query::Cursor]) -> String {
if cs.is_empty() {
return String::new();
}
let mut s = String::from(",\"cursors\":[");
for (i, c) in cs.iter().enumerate() {
if i > 0 {
s.push(',');
}
s.push('"');
s.push_str(&c.to_string());
s.push('"');
}
s.push(']');
s
}
pub(crate) fn json_ok(body: String) -> Response {
(
StatusCode::OK,
[(header::CONTENT_TYPE, "application/json")],
body,
)
.into_response()
}
fn bad_request(msg: &str) -> Response {
error(StatusCode::BAD_REQUEST, msg)
}
pub(crate) fn error(code: StatusCode, msg: &str) -> Response {
let mut j = mira_core::json::Json::new();
j.obj(|j| {
j.key("error");
j.str(msg);
});
(
code,
[(header::CONTENT_TYPE, "application/json")],
j.into_string(),
)
.into_response()
}
pub fn now_nanos() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos() as i64
}
pub fn parse_search(text: &str, now: i64) -> Result<Search, String> {
search_doc(&parse(text)?, now)
}
pub fn search_doc(doc: &Yaml, now: i64) -> Result<Search, String> {
known(
doc,
&["signal", "from", "to", "where", "limit", "after", "cursors"],
)?;
let signal = match doc["signal"].as_str() {
Some(s) => Signal::parse(s).ok_or(format!("unknown signal {s:?}"))?,
None => Signal::Logs,
};
let (from, to) = bounds(doc, now)?;
let limit = positive(doc, "limit", 100, MAX_LIMIT)?;
let after = match &doc["after"] {
Yaml::BadValue | Yaml::Null => None,
y => Some(
y.as_str()
.ok_or("after: quote the cursor; it is a string")?
.parse()?,
),
};
Ok(Search {
signal,
from,
to,
terms: terms(doc)?,
limit,
after,
cursors: flag(doc, "cursors")?,
})
}
fn flag(doc: &Yaml, key: &str) -> Result<bool, String> {
match &doc[key] {
Yaml::BadValue | Yaml::Null => Ok(false),
Yaml::Boolean(b) => Ok(*b),
y => y
.as_str()
.and_then(|s| s.parse().ok())
.ok_or(format!("{key}: expected \"true\" or \"false\"")),
}
}
pub fn parse_series(text: &str, now: i64) -> Result<SeriesQuery, String> {
series_doc(&parse(text)?, now)
}
pub fn series_doc(doc: &Yaml, now: i64) -> Result<SeriesQuery, String> {
known(
doc,
&["name", "from", "to", "where", "max_series", "max_points"],
)?;
let (from, to) = bounds(doc, now)?;
Ok(SeriesQuery {
name: doc["name"].as_str().map(str::to_owned),
from,
to,
terms: terms(doc)?,
max_series: positive(doc, "max_series", 200, 2_000)?,
max_points: positive(doc, "max_points", 5_000, 100_000)?,
})
}
pub fn parse_correlate(text: &str, now: i64) -> Result<(Search, Vec<Expand>), String> {
correlate_doc(&parse(text)?, now)
}
pub fn correlate_doc(doc: &Yaml, now: i64) -> Result<(Search, Vec<Expand>), String> {
known(
doc,
&["signal", "from", "to", "where", "limit", "after", "expand"],
)?;
Ok((correlate_search(doc, now)?, expands(&doc["expand"])?))
}
pub fn map_doc(doc: &Yaml, now: i64) -> Result<(i64, i64, usize), String> {
known(doc, &["from", "to", "max_spans"])?;
let (from, to) = bounds(doc, now)?;
Ok((from, to, positive(doc, "max_spans", 1_000_000, 20_000_000)?))
}
fn correlate_search(doc: &Yaml, now: i64) -> Result<Search, String> {
let mut map = doc.as_hash().cloned().unwrap_or_default();
map.remove(&Yaml::String("expand".into()));
search_doc(&Yaml::Hash(map), now)
}
fn expands(y: &Yaml) -> Result<Vec<Expand>, String> {
let items = match y {
Yaml::BadValue | Yaml::Null => return Ok(Vec::new()),
Yaml::Array(a) => a,
_ => return Err("`expand` must be a list of steps".into()),
};
items
.iter()
.map(|s| {
let s = s.as_str().ok_or("each `expand` step must be a string")?;
match s.split_once(':') {
Some(("around", d)) => {
Ok(Expand::Around(crate::config::duration(d)?.as_nanos() as i64))
}
None if s == "traces" => Ok(Expand::Traces),
None if s == "peers" => Ok(Expand::Peers),
_ => Err(format!(
"unknown expand step {s:?}; expected traces, peers, or around:<duration>"
)),
}
})
.collect()
}
pub fn window(text: &str, now: i64) -> Result<(i64, i64), String> {
window_doc(&parse(text)?, now)
}
pub fn window_doc(doc: &Yaml, now: i64) -> Result<(i64, i64), String> {
known(doc, &["from", "to"])?;
bounds(doc, now)
}
pub fn known(doc: &Yaml, keys: &[&str]) -> Result<(), String> {
if doc.is_badvalue() || doc.is_null() {
return Ok(());
}
let map = doc.as_hash().ok_or("a query must be a mapping")?;
for k in map.keys() {
let k = k.as_str().ok_or("query keys must be strings")?;
if !keys.contains(&k) {
return Err(format!(
"unknown query key {k:?}; expected one of {}",
keys.join(" ")
));
}
}
Ok(())
}
pub fn parse(text: &str) -> Result<Yaml, String> {
if has_alias(text) {
return Err(
"KYAML has no anchors or aliases: write the value out instead of \
referring to an anchor with `*`"
.into(),
);
}
let docs = match YamlLoader::load_from_str(text) {
Ok(docs) => docs,
Err(e) => YamlLoader::load_from_str(&fold_surrogates(text))
.map_err(|_| format!("not valid KYAML: {e}"))?,
};
docs.into_iter().next().ok_or("empty query".into())
}
fn has_alias(text: &str) -> bool {
if !text.contains('*') {
return false;
}
#[derive(Default)]
struct Spy(bool);
impl EventReceiver for Spy {
fn on_event(&mut self, ev: Event) {
self.0 |= matches!(ev, Event::Alias(_));
}
}
let mut spy = Spy::default();
let _ = Parser::new_from_str(text).load(&mut spy, true);
spy.0
}
fn fold_surrogates(text: &str) -> String {
fn hex4(s: &str) -> Option<u32> {
u32::from_str_radix(s.strip_prefix("\\u")?.get(..4)?, 16).ok()
}
let mut out = String::with_capacity(text.len());
let mut rest = text;
while let Some(i) = rest.find("\\u") {
out.push_str(&rest[..i]);
let folded = match (hex4(&rest[i..]), rest.get(i + 6..).and_then(hex4)) {
(Some(h @ 0xD800..=0xDBFF), Some(l @ 0xDC00..=0xDFFF)) => {
char::from_u32(0x10000 + ((h - 0xD800) << 10) + (l - 0xDC00))
}
_ => None,
};
match folded {
Some(c) => {
out.push(c);
rest = &rest[i + 12..];
}
None => {
out.push_str("\\u");
rest = &rest[i + 2..];
}
}
}
out.push_str(rest);
out
}
pub fn bounds(doc: &Yaml, now: i64) -> Result<(i64, i64), String> {
let from = time_field(&doc["from"], now, now - 3_600_000_000_000)?;
let to = time_field(&doc["to"], now, now)?;
if from > to {
return Err(format!("from ({from}) is after to ({to})"));
}
Ok((from, to))
}
fn terms(doc: &Yaml) -> Result<Vec<Term>, String> {
match &doc["where"] {
Yaml::BadValue | Yaml::Null => Ok(Vec::new()),
Yaml::Array(a) => a.iter().map(parse_term).collect(),
_ => Err("`where` must be a list of terms".into()),
}
}
fn positive(doc: &Yaml, key: &str, default: usize, max: usize) -> Result<usize, String> {
match &doc[key] {
Yaml::BadValue | Yaml::Null => Ok(default),
Yaml::Integer(n) if *n > 0 => Ok((*n as usize).min(max)),
other => Err(format!("{key} must be a positive integer, got {other:?}")),
}
}
fn parse_term(y: &Yaml) -> Result<Term, String> {
let map = y.as_hash().ok_or("each `where` term must be a mapping")?;
let mut target = None;
let mut opval = None;
for (k, v) in map {
let k = k.as_str().ok_or("term keys must be strings")?;
match k {
"attr" => {
target = Some(Target::Attr(
v.as_str().ok_or("`attr` must be a string")?.to_owned(),
));
}
"field" => {
target = Some(Target::Field(
v.as_str().ok_or("`field` must be a string")?.to_owned(),
));
}
other => {
let op = Op::parse(other).ok_or(format!(
"unknown term key {other:?}; expected attr, field, or one of \
eq ne lt lte gt gte contains"
))?;
opval = Some((op, scalar(v)?));
}
}
}
let target = target.ok_or("a term needs `attr` or `field`")?;
let (op, value) = opval.ok_or("a term needs an operator, such as `eq`")?;
Ok(Term { target, op, value })
}
fn scalar(y: &Yaml) -> Result<Value, String> {
Ok(match y {
Yaml::String(s) => Value::Str(s.clone()),
Yaml::Integer(i) => Value::Int(*i),
Yaml::Boolean(b) => Value::Bool(*b),
Yaml::Real(r) => Value::Double(r.parse().map_err(|_| format!("{r:?} is not a number"))?),
other => return Err(format!("{other:?} is not a comparable value")),
})
}
fn time_field(y: &Yaml, now: i64, default: i64) -> Result<i64, String> {
match y {
Yaml::BadValue | Yaml::Null => Ok(default),
Yaml::Integer(n) => Ok(*n),
Yaml::String(s) if s == "now" => Ok(now),
Yaml::String(s) => {
let (sign, rest) = match s.strip_prefix('-') {
Some(r) => (-1i64, r),
None => (1, s.strip_prefix('+').unwrap_or(s)),
};
let d = crate::config::duration(rest)?;
i64::try_from(d.as_nanos())
.ok()
.and_then(|ns| now.checked_add(sign * ns))
.ok_or_else(|| format!("{s:?} is further from now than a timestamp reaches"))
}
other => Err(format!("{other:?} is not a time")),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_time_further_from_now_than_a_timestamp_reaches_is_rejected() {
let now = 1_757_000_000_000_000_000;
let at = |s: &str| time_field(&Yaml::String(s.to_owned()), now, 0);
for s in ["-10000000000s", "+9000000000s", "1000000000000000000d"] {
assert!(at(s).is_err(), "{s:?} must not wrap");
}
assert_eq!(at("-1h").unwrap(), now - 3_600_000_000_000);
assert_eq!(at("now").unwrap(), now);
assert_eq!(time_field(&Yaml::BadValue, now, 42).unwrap(), 42);
}
#[test]
fn json_from_a_browser_parses_as_kyaml() {
let now = 1_000_000_000_000_000_000;
let q = parse_search(
r#"{"signal":"traces","from":"-15m","to":"now","limit":50,
"where":[{"attr":"service.name","eq":"checkout"},
{"field":"duration_nano","gte":500000000},
{"field":"name","contains":"GET"}]}"#,
now,
)
.unwrap();
assert_eq!(q.signal, Signal::Traces);
assert_eq!(q.to, now);
assert_eq!(q.from, now - 900_000_000_000);
assert_eq!(q.limit, 50);
assert_eq!(q.terms.len(), 3);
assert!(matches!(&q.terms[0].target, Target::Attr(k) if k == "service.name"));
assert_eq!(q.terms[1].op, Op::Gte);
assert_eq!(q.terms[1].value, Value::Int(500_000_000));
assert_eq!(q.terms[2].op, Op::Contains);
}
#[test]
fn json_surrogate_escapes_are_folded_into_the_character() {
let q = parse_search(
r#"{"where":[{"field":"body","contains":"caf\u00e9 \ud83d\ude00"}]}"#,
0,
)
.unwrap();
assert_eq!(q.terms[0].value, Value::Str("café 😀".into()));
let err =
parse_search(r#"{"where":[{"field":"body","contains":"\ud83d"}]}"#, 0).unwrap_err();
assert!(err.contains("not valid KYAML"), "{err}");
}
#[test]
fn kyaml_with_comments_and_trailing_commas_parses_the_same() {
let q = parse_search(
r#"{
"signal": "logs",
# only the errors
"where": [
{ "field": "severity_number", "gte": 17 },
],
}"#,
0,
)
.unwrap();
assert_eq!(q.signal, Signal::Logs);
assert_eq!(q.terms.len(), 1);
assert_eq!(q.limit, 100);
assert_eq!(q.from, -3_600_000_000_000);
}
#[test]
fn malformed_queries_say_what_is_wrong() {
let err = |s: &str| parse_search(s, 0).unwrap_err();
assert!(err(r#"{"signal":"jaeger"}"#).contains("unknown signal"));
assert!(err(r#"{"where":[{"attr":"a","like":"b"}]}"#).contains("unknown term key"));
assert!(err(r#"{"where":[{"eq":"b"}]}"#).contains("needs `attr` or `field`"));
assert!(err(r#"{"where":[{"attr":"a"}]}"#).contains("needs an operator"));
assert!(err(r#"{"from":"now","to":"-1h"}"#).contains("is after"));
assert!(err(r#"{"limit":0}"#).contains("positive integer"));
assert!(err(r#"{"signal":"logs","filters":[]}"#).contains("unknown query key"));
assert!(err(r#"{"query":{"signal":"logs"}}"#).contains("unknown query key"));
assert!(err("[1,2,3]").contains("must be a mapping"));
let step = parse_series(r#"{"name":"m","step":"1m"}"#, 0).unwrap_err();
assert!(step.contains("step"), "{step}");
let w = window(r#"{"from":0,"limit":5}"#, 0).unwrap_err();
assert!(w.contains("limit"), "{w}");
assert!(err(r#"{"where":[{"attr":"a","eq":[1,2]}]}"#).contains("not a comparable value"));
assert!(err(r#"{"where":[{"attr":"a","eq":{}}]}"#).contains("not a comparable value"));
assert!(err(r#"{"from":[1]}"#).contains("is not a time"));
assert!(err(r#"{"to":{"at":1}}"#).contains("is not a time"));
}
#[test]
fn every_document_the_clients_send_is_accepted() {
for d in [
r#"{"signal":"logs","from":"-1h","to":"now","where":[],"limit":200}"#,
r#"{"signal":"traces","from":0,"to":"now","limit":2000,
"where":[{"field":"trace_id","eq":"ab"}]}"#,
r#"{"signal":"logs","limit":100,"after":"1757241600000000000.2718281828.7.41"}"#,
] {
assert!(parse_search(d, 0).is_ok(), "{d}");
}
for d in [
r#"{"name":"m","from":"-1h","to":"now","where":[]}"#,
r#"{"name":"m","from":"-1h","to":"now","max_series":64,"max_points":400,"where":[]}"#,
] {
assert!(parse_series(d, 0).is_ok(), "{d}");
}
for d in [
r#"{"signal":"logs","from":"-1h","to":"now","where":[],"limit":200,
"expand":["traces","peers"]}"#,
r#"{"expand":[]}"#,
"{}",
] {
assert!(parse_correlate(d, 0).is_ok(), "{d}");
}
assert!(window(r#"{"from":"-1h","to":"now"}"#, 0).is_ok());
assert!(window("{}", 0).is_ok());
assert!(search_doc(&Yaml::BadValue, 0).is_ok());
assert!(correlate_doc(&Yaml::BadValue, 0).is_ok());
assert!(map_doc(&Yaml::BadValue, 0).is_ok());
}
#[test]
fn an_expansion_walk_parses_in_order_and_says_what_it_does_not_know() {
let (q, ops) = parse_correlate(
r#"{"signal":"traces","expand":["around:2s","traces","peers"]}"#,
0,
)
.unwrap();
assert_eq!(q.signal.dir(), "traces");
assert_eq!(
ops,
[Expand::Around(2_000_000_000), Expand::Traces, Expand::Peers]
);
for (d, want) in [
(r#"{"expand":"traces"}"#, "list of steps"),
(r#"{"expand":[7]}"#, "must be a string"),
(r#"{"expand":["sideways"]}"#, "sideways"),
(r#"{"expand":["around:soon"]}"#, "soon"),
(r#"{"expand":["peers:1"]}"#, "peers:1"),
(r#"{"expanded":[]}"#, "unknown query key"),
] {
assert!(parse_correlate(d, 0).unwrap_err().contains(want), "{d}");
}
}
#[test]
fn an_alias_bomb_is_refused_before_it_is_expanded() {
let bomb = r#"{"a":&a "lol",
"b":&b [*a,*a,*a,*a,*a,*a,*a,*a,*a],
"c":&c [*b,*b,*b,*b,*b,*b,*b,*b,*b],
"d":&d [*c,*c,*c,*c,*c,*c,*c,*c,*c],
"e":&e [*d,*d,*d,*d,*d,*d,*d,*d,*d],
"f":&f [*e,*e,*e,*e,*e,*e,*e,*e,*e],
"g":[*f,*f,*f,*f,*f,*f,*f,*f,*f]}"#;
let err = parse(bomb).unwrap_err();
assert!(err.contains("alias"), "{err}");
assert!(err.contains("KYAML"), "{err}");
assert!(parse_search(bomb, 0).is_err());
assert!(parse_series(bomb, 0).is_err());
assert!(window(bomb, 0).is_err());
assert!(parse(r#"{"signal":&s "logs"}"#).is_ok());
let three = r#"{"a":&a "lol","b":&b [*a,*a,*a,*a,*a,*a,*a,*a,*a],
"c":&c [*b,*b,*b,*b,*b,*b,*b,*b,*b],"d":[*c,*c,*c,*c,*c,*c,*c,*c,*c]}"#;
let expanded = YamlLoader::load_from_str(three).unwrap().remove(0);
assert_eq!(expanded["d"].as_vec().unwrap().len(), 9);
assert_eq!(expanded["d"][8][8].as_vec().unwrap().len(), 9);
assert_eq!(expanded["d"][8][8][8].as_str(), Some("lol"));
}
#[test]
fn a_star_inside_a_string_is_not_an_alias() {
let q = parse_search(r#"{"where":[{"field":"body","contains":"rate *"}]}"#, 0).unwrap();
assert_eq!(q.terms[0].value, Value::Str("rate *".into()));
}
#[test]
fn limit_is_capped_rather_than_refused() {
let q = parse_search(r#"{"limit":9999999}"#, 0).unwrap();
assert_eq!(q.limit, MAX_LIMIT);
}
async fn answer(res: Response) -> (StatusCode, String, String) {
let status = res.status();
let ct = res
.headers()
.get(header::CONTENT_TYPE)
.and_then(|v| v.to_str().ok())
.unwrap_or_default()
.to_owned();
let body = axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap();
(status, ct, String::from_utf8(body.to_vec()).unwrap())
}
#[tokio::test]
async fn every_endpoint_refuses_a_document_it_cannot_honour_with_a_json_400() {
let api = Api::default();
let st = || State(api.clone());
for (res, want) in [
(
query_handler(st(), r#"{"signal":"jaeger"}"#.into()).await,
"unknown signal",
),
(
series_handler(st(), r#"{"max_series":0}"#.into()).await,
"max_series must be a positive integer",
),
(
names_handler(st(), r#"{"from":"yesterday"}"#.into()).await,
"yesterday",
),
(
correlate_handler(st(), r#"{"expand":"traces"}"#.into()).await,
"`expand` must be a list of steps",
),
(
map_handler(st(), "{not: [kyaml".into()).await,
"not valid KYAML",
),
(
map_handler(st(), r#"{"max_spans":0}"#.into()).await,
"max_spans must be a positive integer",
),
(
entities_handler(st(), r#"{"limit":5}"#.into()).await,
"unknown query key \\\"limit\\\"",
),
] {
let (status, ct, body) = answer(res).await;
assert_eq!(status, StatusCode::BAD_REQUEST, "{want}: {body}");
assert_eq!(ct, "application/json", "{want}");
assert!(body.starts_with(r#"{"error":""#), "{want}: {body}");
assert!(body.contains(want), "{body}");
}
let dir = std::env::temp_dir().join(format!("mira-api-empty-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let api = Api {
data_dir: Arc::new(dir.clone()),
..Default::default()
};
let st = || State(api.clone());
for res in [
names_handler(st(), String::new()).await,
map_handler(st(), " ".into()).await,
entities_handler(st(), String::new()).await,
] {
let (status, ct, body) = answer(res).await;
assert_eq!(status, StatusCode::OK, "{body}");
assert_eq!(ct, "application/json");
assert!(body.contains("\"stats\":"), "{body}");
}
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_failed_or_panicking_scan_answers_500_with_the_reason() {
let dir = std::env::temp_dir().join(format!("mira-api-broken-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("logs"), b"not a directory").unwrap();
let api = Api {
data_dir: Arc::new(dir.clone()),
..Default::default()
};
let (status, _, body) = answer(query_handler(State(api), "{}".into()).await).await;
assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR, "{body}");
assert!(!body.is_empty(), "a 500 with no reason is not a report");
assert!(!body.contains("panicked"), "{body}");
let (status, _, body) = answer(run("rows", || panic!("a page fault, say")).await).await;
assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR);
assert!(body.contains("query task panicked"), "{body}");
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn an_unknown_signal_has_no_open_block_rather_than_an_index_out_of_range() {
let dir = std::env::temp_dir().join(format!("mira-api-slots-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let pcfg = Arc::new(pipeline::Config {
data_dir: dir.clone(),
node: 0x51,
wal: Some(Arc::new(mira_core::wal::Wal::open(&dir, 0x51).unwrap())),
max_block_age: std::time::Duration::from_secs(3_600),
..Default::default()
});
let (ingest, logs_slot, flusher) = pipeline::spawn::<mira_core::logs::LogsBuilder>(&pcfg);
assert!(
ingest
.submit(crate::e2e::logs_export("checkout", 1_000, 3))
.await
.is_ok()
);
let api = Api {
data_dir: Arc::new(dir.clone()),
open: [logs_slot, Default::default(), Default::default()],
..Default::default()
};
let logs = api.open("logs").await;
assert_eq!(logs.len(), 1, "the served slot answers with its open block");
assert_eq!(logs[0].sealed.num_rows, 3);
for absent in ["profiles", "", "traces", "metrics"] {
assert!(api.open(absent).await.is_empty(), "{absent}");
}
let all = api.open_all().await;
assert_eq!(all.len(), pipeline::SIGNALS.len());
let i = pipeline::SIGNALS.iter().position(|s| *s == "logs").unwrap();
for (j, blocks) in all.iter().enumerate() {
assert_eq!(blocks.len(), usize::from(j == i), "slot {j}");
}
drop(ingest);
flusher.await.unwrap();
let _ = std::fs::remove_dir_all(&dir);
}
}