use std::future::IntoFuture;
use std::sync::Arc;
use std::time::Duration;
use axum::Router;
use axum::body::Body;
use axum::http::{Request, StatusCode};
use prost::Message;
use tower::ServiceExt;
use mira_proto::collector::logs::v1::ExportLogsServiceRequest;
use mira_proto::collector::trace::v1::ExportTraceServiceRequest;
use mira_proto::common::v1::{AnyValue, InstrumentationScope, KeyValue, any_value};
use mira_proto::logs::v1::{LogRecord, ResourceLogs, ScopeLogs};
use mira_proto::resource::v1::Resource;
use mira_proto::trace::v1::{ResourceSpans, ScopeSpans, Span};
use crate::{api, pipeline, receiver};
pub struct Captured(Option<Arc<std::sync::Mutex<Vec<u8>>>>);
impl Captured {
pub fn text(&self) -> String {
let buf = self.0.as_ref().expect("a capture always has a buffer");
String::from_utf8_lossy(&buf.lock().expect("log buffer")).into_owned()
}
}
impl std::io::Write for Captured {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
if let Some(sink) = &self.0 {
sink.lock().expect("log buffer").extend_from_slice(buf);
}
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
thread_local! {
static SINK: std::cell::RefCell<Option<Arc<std::sync::Mutex<Vec<u8>>>>> =
const { std::cell::RefCell::new(None) };
}
struct ToCallingThread;
impl tracing_subscriber::fmt::MakeWriter<'_> for ToCallingThread {
type Writer = Captured;
fn make_writer(&self) -> Captured {
Captured(SINK.with(|s| s.borrow().clone()))
}
}
pub struct CaptureGuard(());
impl Drop for CaptureGuard {
fn drop(&mut self) {
SINK.with(|s| *s.borrow_mut() = None);
}
}
pub fn capture() -> (CaptureGuard, Captured) {
static INSTALLED: std::sync::Once = std::sync::Once::new();
INSTALLED.call_once(|| {
tracing::subscriber::set_global_default(
tracing_subscriber::fmt()
.with_writer(ToCallingThread)
.with_ansi(false)
.without_time()
.with_max_level(tracing::Level::TRACE)
.finish(),
)
.expect("nothing else in this test binary sets a global subscriber");
});
let buf = Arc::new(std::sync::Mutex::new(Vec::new()));
SINK.with(|s| *s.borrow_mut() = Some(Arc::clone(&buf)));
(CaptureGuard(()), Captured(Some(buf)))
}
#[test]
fn the_log_capture_hands_the_writer_and_the_reader_one_buffer() {
use std::io::Write;
use tracing_subscriber::fmt::MakeWriter;
let (guard, log) = capture();
let mut w = ToCallingThread.make_writer();
assert_eq!(w.write(b"one ").unwrap(), 4);
w.write_all(b"two").unwrap();
w.flush().unwrap();
assert_eq!(log.text(), "one two");
tracing::info!(marker = "captured", "hello");
drop(guard);
let text = log.text();
assert!(text.starts_with("one two"), "{text}");
assert!(text.contains(r#"marker="captured""#), "{text}");
let mut off = ToCallingThread.make_writer();
assert_eq!(off.write(b"nobody is listening").unwrap(), 19);
off.flush().unwrap();
assert_eq!(
log.text(),
text,
"a discarded line reached a dropped buffer"
);
}
fn kv(k: &str, v: &str) -> KeyValue {
KeyValue {
key: k.into(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue(v.into())),
}),
}
}
fn fresh_dir(name: &str) -> std::path::PathBuf {
let root = std::env::temp_dir().join(format!("mira-e2e-{}-{name}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
std::fs::create_dir_all(&root).unwrap();
root
}
fn wire(name: &str) -> (receiver::Receivers, api::Api, std::path::PathBuf) {
wire_with(name, true)
}
fn wire_with(name: &str, wal: bool) -> (receiver::Receivers, api::Api, std::path::PathBuf) {
let root = fresh_dir(name);
let cfg = Arc::new(pipeline::Config {
data_dir: root.clone(),
max_block_age: Duration::from_millis(100),
wal: wal.then(|| {
Arc::new(mira_core::wal::Wal::open(&root, mira_core::block::node_id("mira")).unwrap())
}),
..Default::default()
});
let (logs, o_logs, _) = pipeline::spawn::<mira_core::logs::LogsBuilder>(&cfg);
let (traces, o_traces, _) = pipeline::spawn::<mira_core::traces::TracesBuilder>(&cfg);
let (metrics, o_metrics, _) = pipeline::spawn::<mira_core::metrics::MetricsBuilder>(&cfg);
let recv = receiver::Receivers {
logs,
traces,
metrics,
max_request_bytes: crate::config::Config::default().max_request_bytes,
};
let api = api::Api {
data_dir: Arc::new(root.clone()),
open: [o_logs, o_traces, o_metrics],
alerts: Arc::default(),
};
(recv, api, root)
}
fn boot(name: &str) -> (Router, std::path::PathBuf) {
boot_with(name, true)
}
fn boot_sealing(name: &str) -> (Router, std::path::PathBuf) {
boot_with(name, false)
}
fn boot_alerting(name: &str, rules: &str) -> (Router, api::Api, std::path::PathBuf) {
let (recv, mut api, root) = wire_with(name, true);
api.alerts = Arc::new(crate::alert::Engine::new(
crate::alert::Rules::parse(rules).expect("rules"),
));
let app = receiver::http_router(recv)
.merge(api::router(api.clone()).layer(axum::middleware::from_fn(crate::timed)))
.merge(crate::alert::router(api.clone()));
(app, api, root)
}
fn boot_with(name: &str, wal: bool) -> (Router, std::path::PathBuf) {
let (recv, api, root) = wire_with(name, wal);
(router_for(recv, api), root)
}
fn router_for(recv: receiver::Receivers, api: api::Api) -> Router {
receiver::http_router(recv)
.merge(api::router(api.clone()).layer(axum::middleware::from_fn(crate::timed)))
.merge(crate::mcp::router(api.clone()))
.merge(crate::alert::router(api))
.merge(crate::ui::router())
}
fn boot_grpc(name: &str) -> (Router, Router, std::path::PathBuf) {
let (recv, api, root) = wire(name);
let grpc = tonic::service::Routes::default()
.add_service(recv.logs_server())
.add_service(recv.traces_server())
.add_service(recv.metrics_server())
.into_axum_router();
(grpc, api::router(api), root)
}
async fn get(
app: &Router,
path: &str,
if_none_match: Option<&str>,
) -> (StatusCode, Vec<u8>, String) {
let mut req = Request::builder().method("GET").uri(path);
if let Some(tag) = if_none_match {
req = req.header("if-none-match", tag);
}
let res = app
.clone()
.oneshot(req.body(Body::empty()).unwrap())
.await
.unwrap();
let status = res.status();
let etag = res
.headers()
.get("etag")
.map(|v| v.to_str().unwrap().to_owned())
.unwrap_or_default();
let bytes = axum::body::to_bytes(res.into_body(), 16 << 20)
.await
.unwrap();
(status, bytes.to_vec(), etag)
}
async fn post(app: &Router, path: &str, content_type: &str, body: Vec<u8>) -> (StatusCode, String) {
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri(path)
.header("content-type", content_type)
.body(Body::from(body))
.unwrap(),
)
.await
.unwrap();
let status = res.status();
let bytes = axum::body::to_bytes(res.into_body(), 16 << 20)
.await
.unwrap();
(status, String::from_utf8_lossy(&bytes).into_owned())
}
async fn post_enc(
app: &Router,
path: &str,
content_type: &str,
content_encoding: &str,
body: Vec<u8>,
) -> (StatusCode, String) {
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri(path)
.header("content-type", content_type)
.header("content-encoding", content_encoding)
.body(Body::from(body))
.unwrap(),
)
.await
.unwrap();
let status = res.status();
let bytes = axum::body::to_bytes(res.into_body(), 16 << 20)
.await
.unwrap();
(status, String::from_utf8_lossy(&bytes).into_owned())
}
async fn grpc(app: &Router, service: &str, encoding: Option<&str>, payload: Vec<u8>) -> String {
use http_body_util::BodyExt;
let mut frame = vec![encoding.is_some() as u8];
frame.extend_from_slice(&(payload.len() as u32).to_be_bytes());
frame.extend_from_slice(&payload);
let mut req = Request::builder()
.method("POST")
.uri(format!("/opentelemetry.proto.collector.{service}/Export"))
.header("content-type", "application/grpc")
.header("te", "trailers");
if let Some(e) = encoding {
req = req.header("grpc-encoding", e);
}
let res = app
.clone()
.oneshot(req.body(Body::from(frame)).unwrap())
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let status =
|h: &axum::http::HeaderMap| h.get("grpc-status").map(|v| v.to_str().unwrap().to_owned());
if let Some(s) = status(res.headers()) {
return s;
}
let body = res.into_body().collect().await.unwrap();
body.trailers().and_then(status).unwrap_or_default()
}
fn gzip(body: &[u8]) -> Vec<u8> {
use std::io::Write;
let mut e = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::fast());
e.write_all(body).unwrap();
e.finish().unwrap()
}
async fn otlp(app: &Router, path: &str, msg: impl Message) {
let (status, body) = post(app, path, "application/x-protobuf", msg.encode_to_vec()).await;
assert_eq!(status, StatusCode::OK, "{path} rejected the export: {body}");
}
async fn query(app: &Router, doc: &str) -> String {
let (status, body) = post(app, "/api/v1/query", "application/json", doc.into()).await;
assert_eq!(status, StatusCode::OK, "query failed: {body}");
body
}
pub(crate) fn logs_export(service: &str, base_ts: u64, n: usize) -> ExportLogsServiceRequest {
ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
resource: Some(Resource {
attributes: vec![kv("service.name", service), kv("deployment.env", "prod")],
..Default::default()
}),
scope_logs: vec![ScopeLogs {
scope: Some(InstrumentationScope {
name: "mira.e2e".into(),
..Default::default()
}),
log_records: (0..n)
.map(|i| LogRecord {
time_unix_nano: base_ts + i as u64,
severity_number: if i % 2 == 0 { 17 } else { 9 },
severity_text: if i % 2 == 0 { "ERROR" } else { "INFO" }.into(),
body: Some(AnyValue {
value: Some(any_value::Value::StringValue(format!(
"{service} handled request {i}"
))),
}),
attributes: vec![kv(
"http.method",
if i % 2 == 0 { "POST" } else { "GET" },
)],
trace_id: vec![0xab; 16].into(),
span_id: vec![0xcd; 8].into(),
..Default::default()
})
.collect(),
..Default::default()
}],
..Default::default()
}],
}
}
#[tokio::test]
async fn a_volume_that_cannot_take_a_block_is_retryable_and_a_too_wide_export_is_not() {
let (app, root) = boot_sealing("unwritable");
std::fs::write(root.join(".tmp"), b"not a directory").unwrap();
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/v1/logs")
.header("content-type", "application/x-protobuf")
.body(Body::from(
logs_export("checkout", 1_000, 1).encode_to_vec(),
))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::SERVICE_UNAVAILABLE);
assert_eq!(
res.headers()
.get("retry-after")
.map(|v| v.to_str().unwrap()),
Some("1"),
"an exporter with no backoff of its own needs the pushback"
);
let mut wide = logs_export("checkout", 1_000, 1);
wide.resource_logs[0].scope_logs[0].log_records[0].attributes = (0..70_000)
.map(|i| kv(&format!("k{i}"), "v"))
.collect::<Vec<_>>();
let (status, body) = post(
&app,
"/v1/logs",
"application/x-protobuf",
wide.encode_to_vec(),
)
.await;
assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR, "{body}");
let _ = std::fs::remove_dir_all(&root);
}
#[tokio::test]
async fn otlp_logs_are_queryable_the_moment_the_export_is_acknowledged() {
let (app, root) = boot("logs");
otlp(&app, "/v1/logs", logs_export("checkout", 1_000, 6)).await;
otlp(&app, "/v1/logs", logs_export("payments", 5_000, 4)).await;
let none = query(&app, r#"{"signal":"logs","from":0,"to":500}"#).await;
assert!(none.contains(r#""rows":[]"#), "{none}");
let all = query(&app, r#"{"signal":"logs","from":0,"to":100000}"#).await;
assert!(all.contains("checkout handled request 0"), "{all}");
assert!(all.contains("payments handled request 3"), "{all}");
let first_body = all.find("payments handled request 3").unwrap();
assert!(first_body < all.find("checkout handled request 0").unwrap());
let checkout = query(
&app,
r#"{"signal":"logs","from":0,"to":100000,
"where":[{"attr":"service.name","eq":"checkout"}]}"#,
)
.await;
assert_eq!(checkout.matches("handled request").count(), 6, "{checkout}");
assert!(!checkout.contains("payments"), "{checkout}");
let errs = query(
&app,
r#"{"signal":"logs","from":0,"to":100000,
"where":[{"attr":"http.method","eq":"POST"},
{"field":"severity_number","gte":17}]}"#,
)
.await;
assert_eq!(errs.matches("handled request").count(), 5, "{errs}");
assert!(
checkout.contains(r#""otel.scope.name":"mira.e2e""#),
"{checkout}"
);
assert!(
checkout.contains(r#""deployment.env":"prod""#),
"{checkout}"
);
assert!(!checkout.contains(r#""resource_id""#), "{checkout}");
assert!(checkout.contains(r#""blocks_total":1"#), "{checkout}");
assert!(checkout.contains(r#""rows_matched":6"#), "{checkout}");
assert_eq!(
mira_core::block::scan(&root, "logs").unwrap().len(),
0,
"nothing has been published yet, so every row above came from the open block"
);
}
#[tokio::test]
async fn a_paged_read_reassembles_the_one_shot_answer() {
let (app, _root) = boot("paging");
for base in [1_000, 2_000, 3_000] {
otlp(&app, "/v1/logs", logs_export("checkout", base, 5)).await;
}
let rows = |body: &str| {
let s = body.find(r#""rows":["#).unwrap() + 8;
body[s..body.find(r#"],"stats":"#).unwrap()].to_owned()
};
let all = query(&app, r#"{"signal":"logs","from":0,"to":100000}"#).await;
assert_eq!(all.matches("handled request").count(), 15, "{all}");
assert!(
!all.contains(r#""next""#),
"a short page is the last page: {all}"
);
let mut pages = Vec::new();
let mut after = String::new();
for _ in 0..10 {
let body = query(
&app,
&format!(r#"{{"signal":"logs","from":0,"to":100000,"limit":4{after}}}"#),
)
.await;
pages.push(rows(&body));
let Some(i) = body.find(r#""next":""#) else {
break;
};
let c = &body[i + 8..];
after = format!(r#","after":"{}""#, &c[..c.find('"').unwrap()]);
}
assert_eq!(pages.len(), 4, "15 rows, 4 a page");
assert_eq!(
pages.join(","),
rows(&all),
"pages must reassemble the whole"
);
for (doc, want) in [
(r#"{"signal":"logs","after":1.2}"#, "it is a string"),
(
r#"{"signal":"logs","after":"1.2.3"}"#,
"pass back the `next` field verbatim",
),
] {
let (status, body) = post(&app, "/api/v1/query", "application/json", doc.into()).await;
assert_eq!(status, StatusCode::BAD_REQUEST, "{body}");
assert!(body.contains(want), "{body}");
}
}
#[tokio::test]
async fn a_cursor_taken_from_the_open_block_survives_the_seal() {
let (app, root) = boot("seal-cursor");
for base in [1_000, 2_000, 3_000] {
otlp(&app, "/v1/logs", logs_export("checkout", base, 5)).await;
}
let rows = |body: &str| {
let s = body.find(r#""rows":["#).unwrap() + 8;
body[s..body.find(r#"],"stats":"#).unwrap()].to_owned()
};
let next = |body: &str| {
body.find(r#""next":""#).map(|i| {
let c = &body[i + 8..];
c[..c.find('"').unwrap()].to_owned()
})
};
let first = query(&app, r#"{"signal":"logs","from":0,"to":100000,"limit":4}"#).await;
assert_eq!(
mira_core::block::scan(&root, "logs").unwrap().len(),
0,
"the point of this test is that page one predates the block"
);
let cursor = next(&first).expect("15 rows, 4 a page");
tokio::time::sleep(Duration::from_millis(200)).await;
otlp(&app, "/v1/logs", logs_export("checkout", 9_000, 1)).await;
seal(&root, "logs").await;
let mut pages = vec![rows(&first)];
let mut after = format!(r#","after":"{cursor}""#);
for _ in 0..10 {
let body = query(
&app,
&format!(r#"{{"signal":"logs","from":0,"to":8000,"limit":4{after}}}"#),
)
.await;
pages.push(rows(&body));
let Some(c) = next(&body) else { break };
after = format!(r#","after":"{c}""#);
}
let joined = pages.join(",");
assert_eq!(
joined.matches("handled request").count(),
15,
"a row was duplicated or dropped across the seal: {joined}"
);
let one_shot = query(&app, r#"{"signal":"logs","from":0,"to":8000}"#).await;
assert_eq!(joined, rows(&one_shot), "pages must reassemble the whole");
}
#[tokio::test]
async fn otlp_spans_are_queryable_and_ids_come_back_as_hex() {
let (app, _root) = boot("traces");
let trace_id = vec![
0x4b, 0xf9, 0x2f, 0x35, 0x77, 0xb3, 0x4d, 0xa6, 0xa3, 0xce, 0x92, 0x9d, 0x0e, 0x0e, 0x47,
0x36,
];
otlp(
&app,
"/v1/traces",
ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: Some(Resource {
attributes: vec![kv("service.name", "frontend")],
..Default::default()
}),
scope_spans: vec![ScopeSpans {
scope: Some(InstrumentationScope {
name: "mira.e2e".into(),
..Default::default()
}),
spans: (0..3)
.map(|i| Span {
trace_id: trace_id.clone().into(),
span_id: vec![i as u8 + 1; 8].into(),
name: format!("GET /checkout/{i}"),
start_time_unix_nano: 1_000 + i * 10,
end_time_unix_nano: 1_500 + i * 10,
attributes: vec![kv("http.route", "/checkout/:id")],
..Default::default()
})
.collect(),
..Default::default()
}],
..Default::default()
}],
},
)
.await;
let hex = "4bf92f3577b34da6a3ce929d0e0e4736";
let found = query(
&app,
&format!(
r#"{{"signal":"traces","from":0,"to":100000,
"where":[{{"field":"trace_id","eq":"{hex}"}}]}}"#
),
)
.await;
assert_eq!(found.matches("GET /checkout/").count(), 3, "{found}");
assert!(found.contains(&format!(r#""trace_id":"{hex}""#)), "{found}");
assert!(found.contains(r#""span_id":"0101010101010101""#), "{found}");
assert!(found.contains(r#""http.route":"/checkout/:id""#), "{found}");
}
#[tokio::test]
async fn a_frame_walk_over_http_turns_one_log_line_into_a_call_chain() {
let (app, _root) = boot("frame");
let hex = "4bf92f3577b34da6a3ce929d0e0e4736";
let trace_id: Vec<u8> = (0..16)
.map(|i| u8::from_str_radix(&hex[i * 2..i * 2 + 2], 16).unwrap())
.collect();
let res = |svc: &str| Resource {
attributes: vec![kv("service.name", svc), kv("service.instance.id", "7f3a")],
..Default::default()
};
let scope = || {
Some(InstrumentationScope {
name: "mira.e2e".into(),
..Default::default()
})
};
let span = |svc: &str, id: u8, parent: u8, at: u64| ResourceSpans {
resource: Some(res(svc)),
scope_spans: vec![ScopeSpans {
scope: scope(),
spans: vec![Span {
trace_id: trace_id.clone().into(),
span_id: vec![id; 8].into(),
parent_span_id: if parent == 0 {
Vec::new()
} else {
vec![parent; 8]
}
.into(),
name: format!("{svc} work"),
start_time_unix_nano: at,
end_time_unix_nano: at + 500,
..Default::default()
}],
..Default::default()
}],
..Default::default()
};
otlp(
&app,
"/v1/traces",
ExportTraceServiceRequest {
resource_spans: vec![
span("gateway", 1, 0, 1_000),
span("api", 2, 1, 1_100),
span("db", 3, 2, 1_200),
],
},
)
.await;
otlp(
&app,
"/v1/logs",
ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
resource: Some(res("api")),
scope_logs: vec![ScopeLogs {
scope: scope(),
log_records: vec![LogRecord {
time_unix_nano: 50_000,
severity_number: 17,
severity_text: "ERROR".into(),
body: Some(AnyValue {
value: Some(any_value::Value::StringValue("checkout failed".into())),
}),
trace_id: trace_id.clone().into(),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
},
)
.await;
let call = |path: &'static str, doc: String| {
let app = app.clone();
async move {
let (status, body) = post(&app, path, "application/json", doc.into()).await;
assert_eq!(status, StatusCode::OK, "{path}: {body}");
body
}
};
let f = call(
"/api/v1/correlate",
r#"{"signal":"logs","from":40000,"to":60000,
"where":[{"field":"severity_text","eq":"ERROR"}],
"expand":["traces","peers"]}"#
.into(),
)
.await;
assert!(f.contains(&format!(r#""{hex}""#)), "{f}");
for svc in ["gateway", "api", "db"] {
assert!(f.contains(&format!(r#""name":"{svc}""#)), "{svc}: {f}");
}
assert!(f.contains(r#""from":"1000""#), "{f}");
assert!(f.contains(r#""truncated":false"#), "{f}");
let m = call("/api/v1/map", r#"{"from":0,"to":100000}"#.into()).await;
let key = |svc: &str| {
let at = m.find(&format!(r#","name":"{svc}""#)).expect(&m);
let k = m[..at].rsplit_once(r#""key":""#).expect(&m).1;
k.trim_end_matches('"').to_owned()
};
for (from, to) in [
("entry".to_owned(), key("gateway")),
(key("gateway"), key("api")),
(key("api"), key("db")),
] {
let edge = format!(r#""from":"{from}","to":"{to}""#);
assert!(m.contains(&edge), "{edge} missing: {m}");
}
assert!(m.contains(r#""unresolved":0"#), "{m}");
let e = call("/api/v1/entities", r#"{"from":0,"to":100000}"#.into()).await;
assert_eq!(e.matches(r#""name":"#).count(), 3, "{e}");
}
#[tokio::test]
async fn a_metric_series_survives_being_split_across_two_blocks() {
use mira_proto::collector::metrics::v1::ExportMetricsServiceRequest;
use mira_proto::metrics::v1::metric::Data;
use mira_proto::metrics::v1::number_data_point::Value as NumValue;
use mira_proto::metrics::v1::{
AggregationTemporality, HistogramDataPoint, Metric, NumberDataPoint, ResourceMetrics,
ScopeMetrics, Sum,
};
let (app, _root) = boot_sealing("metrics");
let export = |base: u64| ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
resource: Some(Resource {
attributes: vec![kv("service.name", "checkout")],
..Default::default()
}),
scope_metrics: vec![ScopeMetrics {
scope: Some(InstrumentationScope {
name: "mira.e2e".into(),
..Default::default()
}),
metrics: vec![Metric {
name: "http.server.requests".into(),
unit: "1".into(),
data: Some(Data::Sum(Sum {
aggregation_temporality: AggregationTemporality::Cumulative as i32,
is_monotonic: true,
data_points: ["GET", "POST"]
.iter()
.enumerate()
.flat_map(|(m, method)| {
(0..2).map(move |i| NumberDataPoint {
time_unix_nano: base + i as u64,
attributes: vec![kv("http.method", method)],
value: Some(NumValue::AsInt(
(base as i64) + i + m as i64 * 100,
)),
..Default::default()
})
})
.collect(),
})),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
};
otlp(&app, "/v1/metrics", export(1_000)).await;
otlp(&app, "/v1/metrics", export(5_000)).await;
let (status, body) = post(
&app,
"/api/v1/metrics/query",
"application/json",
br#"{"name":"http.server.requests","from":0,"to":100000}"#.to_vec(),
)
.await;
assert_eq!(status, StatusCode::OK, "{body}");
assert_eq!(
body.matches(r#""name":"http.server.requests""#).count(),
2,
"{body}"
);
assert!(
body.contains(r#"["1000","1000"],["1001","1001"],["5000","5000"],["5001","5001"]"#),
"{body}"
);
assert!(
body.contains(r#"["1000","1100"],["1001","1101"],["5000","5100"],["5001","5101"]"#),
"{body}"
);
assert!(
body.contains(r#""kind":"sum","temporality":2,"monotonic":true"#),
"{body}"
);
assert!(
body.contains(
r#""attributes":{"http.method":"GET","otel.scope.name":"mira.e2e","service.name":"checkout"}"#
),
"{body}"
);
assert!(body.contains(r#""blocks_scanned":2"#), "{body}");
let (status, capped) = post(
&app,
"/api/v1/metrics/query",
"application/json",
br#"{"name":"http.server.requests","from":0,"to":100000,"max_points":2}"#.to_vec(),
)
.await;
assert_eq!(status, StatusCode::OK, "{capped}");
assert!(
capped.contains(r#"["5000","5000"],["5001","5001"]"#),
"{capped}"
);
assert!(
capped.contains(r#"["5000","5100"],["5001","5101"]"#),
"{capped}"
);
assert!(
!capped.contains(r#"["1000","#),
"the older half must be gone\n{capped}"
);
assert_eq!(
capped.matches(r#""dropped_points":2"#).count(),
2,
"{capped}"
);
let (_, only_post) = post(
&app,
"/api/v1/metrics/query",
"application/json",
br#"{"from":0,"to":100000,"where":[{"attr":"http.method","eq":"POST"},
{"attr":"service.name","eq":"checkout"}]}"#
.to_vec(),
)
.await;
assert_eq!(only_post.matches(r#""points""#).count(), 1, "{only_post}");
assert!(only_post.contains(r#""http.method":"POST""#), "{only_post}");
let (status, recent) = post(&app, "/api/v1/metrics/names", "application/json", vec![]).await;
assert_eq!(status, StatusCode::OK, "{recent}");
assert!(
recent.contains(r#""names":[],"stats":{"blocks_total":2,"blocks_scanned":0"#),
"{recent}"
);
let (_, names) = post(
&app,
"/api/v1/metrics/names",
"application/json",
br#"{"from":0,"to":100000}"#.to_vec(),
)
.await;
assert!(
names.contains(r#"{"name":"http.server.requests","unit":"1","kind":"sum"}"#),
"{names}"
);
otlp(
&app,
"/v1/metrics",
ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
scope_metrics: vec![ScopeMetrics {
metrics: vec![Metric {
name: "http.server.duration".into(),
unit: "ms".into(),
data: Some(Data::Histogram(mira_proto::metrics::v1::Histogram {
aggregation_temporality: AggregationTemporality::Delta as i32,
data_points: vec![HistogramDataPoint {
time_unix_nano: 9_000,
count: 42,
sum: Some(1234.5),
explicit_bounds: vec![1.0, 5.0, 10.0],
bucket_counts: vec![10, 20, 10, 2],
..Default::default()
}],
})),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
},
)
.await;
let (_, hist) = post(
&app,
"/api/v1/metrics/query",
"application/json",
br#"{"name":"http.server.duration","from":0,"to":100000}"#.to_vec(),
)
.await;
assert!(
hist.contains(r#""name":"http.server.duration.count""#),
"{hist}"
);
assert!(
hist.contains(r#""name":"http.server.duration.sum""#),
"{hist}"
);
assert!(hist.contains(r#"["9000","42"]"#), "{hist}");
assert!(hist.contains(r#"["9000",1234.5]"#), "{hist}");
}
#[tokio::test]
async fn a_histogram_carries_its_exemplars_into_both_derived_series() {
use mira_proto::collector::metrics::v1::ExportMetricsServiceRequest;
use mira_proto::metrics::v1::metric::Data;
use mira_proto::metrics::v1::{
AggregationTemporality, HistogramDataPoint, Metric, ResourceMetrics, ScopeMetrics,
};
let (app, _root) = boot("exemplars");
otlp(
&app,
"/v1/metrics",
ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
scope_metrics: vec![ScopeMetrics {
metrics: vec![Metric {
name: "rpc.duration".into(),
unit: "ms".into(),
data: Some(Data::Histogram(mira_proto::metrics::v1::Histogram {
aggregation_temporality: AggregationTemporality::Delta as i32,
data_points: vec![HistogramDataPoint {
time_unix_nano: 9_000,
count: 3,
sum: Some(30.0),
explicit_bounds: vec![10.0],
bucket_counts: vec![2, 1],
exemplars: vec![mira_proto::metrics::v1::Exemplar {
time_unix_nano: 8_900,
trace_id: vec![0xab; 16].into(),
span_id: vec![0xcd; 8].into(),
value: Some(
mira_proto::metrics::v1::exemplar::Value::AsDouble(29.5),
),
..Default::default()
}],
..Default::default()
}],
})),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
},
)
.await;
let (_, hist) = post(
&app,
"/api/v1/metrics/query",
"application/json",
br#"{"name":"rpc.duration","from":0,"to":100000}"#.to_vec(),
)
.await;
assert_eq!(
hist.matches(r#""trace_id":"abababababababababababababababab""#)
.count(),
2,
"{hist}"
);
assert!(hist.contains(r#""span_id":"cdcdcdcdcdcdcdcd""#), "{hist}");
assert!(hist.contains(r#""double":29.5"#), "{hist}");
assert!(hist.contains(r#""time_unix_nano":"8900""#), "{hist}");
}
#[tokio::test]
async fn a_malformed_query_returns_a_readable_json_error() {
let (app, _root) = boot("badq");
let (status, body) = post(
&app,
"/api/v1/query",
"application/json",
br#"{"signal":"logs","where":[{"attr":"a","matches":"b"}]}"#.to_vec(),
)
.await;
assert_eq!(status, StatusCode::BAD_REQUEST);
assert!(body.starts_with(r#"{"error":"#), "{body}");
assert!(body.contains("unknown term key"), "{body}");
}
#[tokio::test]
async fn the_ui_is_served_from_the_binary_and_revalidates() {
let (app, _root) = boot("ui");
let (status, body, _) = get(&app, "/", None).await;
assert_eq!(status, StatusCode::OK);
assert!(
String::from_utf8_lossy(&body).contains(r#"<div id="app">"#),
"index.html is not the built bundle"
);
let (status, body, etag) = get(&app, "/app.js", None).await;
assert_eq!(status, StatusCode::OK);
assert!(!body.is_empty() && !etag.is_empty());
let (status, body, _) = get(&app, "/app.js", Some(&etag)).await;
assert_eq!(status, StatusCode::NOT_MODIFIED);
assert!(body.is_empty(), "a 304 must not carry the body");
let (status, _, _) = get(&app, "/logs", None).await;
assert_eq!(status, StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn an_mcp_client_can_list_the_tools_and_call_them() {
let (app, _root) = boot("mcp");
let rpc = async |body: &str| {
let (status, out) = post(&app, "/mcp", "application/json", body.into()).await;
assert_eq!(status, StatusCode::OK, "{out}");
out
};
let init = rpc(r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{}}"#).await;
assert!(init.contains(r#""id":1"#), "{init}");
assert!(init.contains(r#""protocolVersion""#), "{init}");
assert!(init.contains(r#""name":"mira""#), "{init}");
let (status, body) = post(
&app,
"/mcp",
"application/json",
br#"{"jsonrpc":"2.0","method":"notifications/initialized"}"#.to_vec(),
)
.await;
assert_eq!(status, StatusCode::ACCEPTED);
assert!(body.is_empty(), "a notification takes no reply: {body}");
let tools = rpc(r#"{"jsonrpc":"2.0","id":2,"method":"tools/list"}"#).await;
for t in ["query_records", "get_trace", "query_metric", "list_metrics"] {
assert!(tools.contains(&format!(r#""name":"{t}""#)), "{tools}");
}
otlp(&app, "/v1/logs", logs_export("checkout", 1_000, 4)).await;
let rows = rpc(r#"{"jsonrpc":"2.0","id":3,"method":"tools/call","params":{
"name":"query_records","arguments":{
"signal":"logs","from":0,"to":100000,
"where":[{"field":"severity_text","eq":"ERROR"}]}}}"#)
.await;
assert!(rows.contains(r#""isError":false"#), "{rows}");
assert!(rows.contains(r#"blocks_scanned"#), "{rows}");
assert_eq!(
rows.matches("checkout handled request").count(),
2,
"{rows}"
);
let bad = rpc(r#"{"jsonrpc":"2.0","id":4,"method":"tools/call",
"params":{"name":"get_trace","arguments":{"trace_id":"nope"}}}"#)
.await;
assert!(bad.contains(r#""isError":true"#), "{bad}");
assert!(bad.contains("hex trace id"), "{bad}");
let missing = rpc(r#"{"jsonrpc":"2.0","id":5,"method":"frobnicate"}"#).await;
assert!(missing.contains(r#""code":-32601"#), "{missing}");
otlp(
&app,
"/v1/traces",
ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: Some(Resource::default()),
scope_spans: vec![ScopeSpans {
spans: (0..3)
.map(|i| Span {
trace_id: vec![0x5a; 16].into(),
span_id: vec![i as u8 + 1; 8].into(),
name: format!("GET /pay/{i}"),
start_time_unix_nano: 1_000 + i * 10,
end_time_unix_nano: 1_500 + i * 10,
..Default::default()
})
.collect(),
..Default::default()
}],
..Default::default()
}],
},
)
.await;
let trace = rpc(r#"{"jsonrpc":"2.0","id":6,"method":"tools/call","params":{
"name":"get_trace","arguments":{
"trace_id":"5a5a5a5a5a5a5a5a5a5a5a5a5a5a5a5a"}}}"#)
.await;
assert!(trace.contains(r#""isError":false"#), "{trace}");
assert_eq!(trace.matches("GET /pay/").count(), 3, "{trace}");
}
#[tokio::test]
async fn otlp_json_bodies_decode_with_the_deviations_the_spec_requires() {
let (app, _root) = boot("json");
let body = r#"{
"resourceLogs": [{
"resource": {"attributes": [{"key":"service.name","value":{"stringValue":"json-svc"}}]},
"scopeLogs": [{
"scope": {"name":"mira.json","version":"1.2.3"},
"logRecords": [{
"timeUnixNano": "1000000",
"observedTimeUnixNano": "1000001",
"severityNumber": "SEVERITY_NUMBER_ERROR",
"severityText": "ERROR",
"body": {"stringValue": "json log one"},
"traceId": "5a5a5a5a5a5a5a5a5a5a5a5a5a5a5a5a",
"spanId": "0102030405060708",
"attributes": [
{"key":"http.status_code","value":{"intValue":"200"}},
{"key":"retry","value":{"boolValue":true}},
{"key":"ratio","value":{"doubleValue":1.5}},
{"key":"blob","value":{"bytesValue":"aGVsbG8="}}
]
}]
}]
}]
}"#;
let (status, answer) = post(&app, "/v1/logs", "application/json", body.into()).await;
assert_eq!(status, StatusCode::OK, "{answer}");
assert_eq!(answer, "{}");
let rows = query(&app, r#"{"signal":"logs","from":0,"to":9000000}"#).await;
assert!(rows.contains("json log one"), "{rows}");
assert!(
rows.contains(r#""trace_id":"5a5a5a5a5a5a5a5a5a5a5a5a5a5a5a5a""#),
"{rows}"
);
assert!(rows.contains(r#""span_id":"0102030405060708""#), "{rows}");
assert!(rows.contains(r#""severity_number":17"#), "{rows}");
assert!(rows.contains(r#""service.name":"json-svc""#), "{rows}");
assert!(rows.contains(r#""otel.scope.version":"1.2.3""#), "{rows}");
assert!(rows.contains(r#""http.status_code":"200""#), "{rows}");
let by_int = query(
&app,
r#"{"signal":"logs","from":0,"to":9000000,
"where":[{"attr":"http.status_code","eq":200}]}"#,
)
.await;
assert!(by_int.contains("json log one"), "{by_int}");
assert!(rows.contains(r#""retry":true"#), "{rows}");
assert!(rows.contains(r#""ratio":1.5"#), "{rows}");
let body = r#"{
"resource_spans": [{
"resource": {"attributes": [{"key":"service.name","value":{"string_value":"json-svc"}}]},
"scope_spans": [{
"scope": {"name":"mira.json"},
"spans": [{
"trace_id": "aabbccddeeff00112233445566778899",
"span_id": "1111111111111111",
"parent_span_id": "",
"name": "GET /json",
"kind": 2,
"start_time_unix_nano": 2000000,
"end_time_unix_nano": 2000500,
"status": {"code": "STATUS_CODE_ERROR", "message": "boom"},
"events": [{"time_unix_nano": "2000100", "name": "cache.miss"}],
"links": [{"trace_id": "99887766554433221100ffeeddccbbaa",
"span_id": "2222222222222222"}]
}]
}]
}]
}"#;
let (status, answer) = post(&app, "/v1/traces", "application/json", body.into()).await;
assert_eq!(status, StatusCode::OK, "{answer}");
let spans = query(
&app,
r#"{"signal":"traces","from":0,"to":9000000,
"where":[{"field":"trace_id","eq":"aabbccddeeff00112233445566778899"}]}"#,
)
.await;
assert!(spans.contains("GET /json"), "{spans}");
assert!(spans.contains(r#""span_id":"1111111111111111""#), "{spans}");
assert!(!spans.contains("parent_span_id"), "{spans}");
assert!(spans.contains(r#""events":[{"#), "{spans}");
assert!(spans.contains(r#""name":"cache.miss""#), "{spans}");
assert!(
spans.contains(r#""links":[{"trace_id":"99887766554433221100ffeeddccbbaa""#),
"{spans}"
);
let body = r#"{
"resourceMetrics": [{
"scopeMetrics": [{
"metrics": [{
"name": "json.requests",
"unit": "1",
"sum": {
"isMonotonic": true,
"aggregationTemporality": "AGGREGATION_TEMPORALITY_CUMULATIVE",
"dataPoints": [{"timeUnixNano": "3000000", "asInt": "42",
"attributes": [{"key":"route","value":{"stringValue":"/json"}}]}]
}
}]
}]
}]
}"#;
let (status, answer) = post(&app, "/v1/metrics", "application/json", body.into()).await;
assert_eq!(status, StatusCode::OK, "{answer}");
let (status, names) = post(
&app,
"/api/v1/metrics/names",
"application/json",
br#"{"from":0,"to":9000000}"#.to_vec(),
)
.await;
assert_eq!(status, StatusCode::OK, "{names}");
assert!(names.contains("json.requests"), "{names}");
}
#[tokio::test]
async fn otlp_json_refusals_are_json() {
let (app, _root) = boot("json-errors");
let (status, _) = post(&app, "/v1/logs", "text/plain", b"hello".to_vec()).await;
assert_eq!(status, StatusCode::UNSUPPORTED_MEDIA_TYPE);
let short = r#"{"resourceLogs":[{"scopeLogs":[{"logRecords":[
{"timeUnixNano":"1","traceId":"abcd"}]}]}]}"#;
let (status, body) = post(&app, "/v1/logs", "application/json", short.into()).await;
assert_eq!(status, StatusCode::BAD_REQUEST, "{body}");
assert!(body.starts_with(r#"{"code":"#), "{body}");
assert!(body.contains("traceId"), "{body}");
let (status, body) = post(&app, "/v1/logs", "application/json", b"{[".to_vec()).await;
assert_eq!(status, StatusCode::BAD_REQUEST, "{body}");
assert!(body.starts_with(r#"{"code":"#), "{body}");
let (status, body) = post(&app, "/v1/logs", "application/x-protobuf", vec![]).await;
assert_eq!(status, StatusCode::OK, "{body}");
}
#[tokio::test]
async fn gzipped_exports_are_accepted_on_both_body_encodings() {
let (app, _root) = boot("gzip");
let msg = logs_export("checkout", 1_000, 4).encode_to_vec();
let (status, body) = post_enc(
&app,
"/v1/logs",
"application/x-protobuf",
"gzip",
gzip(&msg),
)
.await;
assert_eq!(status, StatusCode::OK, "{body}");
let doc = br#"{"resourceLogs":[{"scopeLogs":[{"logRecords":[
{"timeUnixNano":"2000","body":{"stringValue":"gzipped json arrived"}}]}]}]}"#;
let (status, body) = post_enc(&app, "/v1/logs", "application/json", "gzip", gzip(doc)).await;
assert_eq!(status, StatusCode::OK, "{body}");
assert_eq!(body, "{}", "a JSON client gets a JSON answer");
let all = query(&app, r#"{"signal":"logs","from":0,"to":100000}"#).await;
assert!(all.contains("checkout handled request 3"), "{all}");
assert!(all.contains("gzipped json arrived"), "{all}");
let (status, body) = post_enc(
&app,
"/v1/logs",
"application/x-protobuf",
"identity",
msg.clone(),
)
.await;
assert_eq!(status, StatusCode::OK, "{body}");
let (status, body) = post_enc(
&app,
"/v1/logs",
"application/x-protobuf",
"br",
msg.clone(),
)
.await;
assert_eq!(status, StatusCode::UNSUPPORTED_MEDIA_TYPE, "{body}");
let (status, body) = post_enc(&app, "/v1/logs", "application/x-protobuf", "gzip", msg).await;
assert_eq!(status, StatusCode::BAD_REQUEST, "{body}");
assert!(body.contains("decompress"), "{body}");
let (status, body) = post_enc(
&app,
"/v1/logs",
"application/json",
"gzip",
gzip(&vec![b'0'; 80 << 20]),
)
.await;
assert_eq!(status, StatusCode::PAYLOAD_TOO_LARGE, "{body}");
assert!(body.starts_with(r#"{"code":"#), "{body}");
}
#[tokio::test]
async fn one_size_limit_governs_both_transports_and_gzip() {
const LIMIT: usize = 4 << 10;
let (mut recv, _api, _root) = wire("max-request-bytes");
recv.max_request_bytes = LIMIT;
let grpc_app = tonic::service::Routes::default()
.add_service(recv.logs_server())
.into_axum_router();
let http = receiver::http_router(recv.clone());
let over = vec![b'0'; LIMIT + 1];
let (status, body) = post(&http, "/v1/logs", "application/json", over.clone()).await;
assert_eq!(status, StatusCode::PAYLOAD_TOO_LARGE, "{body}");
let (status, body) = post_enc(&http, "/v1/logs", "application/json", "gzip", gzip(&over)).await;
assert_eq!(status, StatusCode::PAYLOAD_TOO_LARGE, "{body}");
assert!(body.contains("max_request_bytes"), "{body}");
assert_eq!(
grpc(&grpc_app, "logs.v1.LogsService", None, over).await,
"11"
);
let msg = logs_export("checkout", 1_000, 3).encode_to_vec();
assert!(msg.len() < LIMIT, "{} bytes", msg.len());
assert_eq!(
grpc(&grpc_app, "logs.v1.LogsService", None, msg.clone()).await,
"0"
);
let (status, body) = post(&http, "/v1/logs", "application/x-protobuf", msg).await;
assert_eq!(status, StatusCode::OK, "{body}");
}
#[tokio::test]
async fn a_gzipped_grpc_export_is_accepted_on_every_signal() {
let (grpc_app, api, _root) = boot_grpc("grpc-gzip");
let logs = logs_export("checkout", 1_000, 3).encode_to_vec();
let svc = "logs.v1.LogsService";
assert_eq!(grpc(&grpc_app, svc, Some("gzip"), gzip(&logs)).await, "0");
assert_eq!(grpc(&grpc_app, svc, None, logs).await, "0");
for svc in [
"trace.v1.TraceService",
"metrics.v1.MetricsService",
"logs.v1.LogsService",
] {
assert_eq!(
grpc(&grpc_app, svc, Some("gzip"), gzip(&[])).await,
"0",
"{svc} refused a gzipped export"
);
}
assert_eq!(grpc(&grpc_app, svc, Some("deflate"), gzip(&[])).await, "12");
let all = query(&api, r#"{"signal":"logs","from":0,"to":100000}"#).await;
assert_eq!(all.matches("handled request").count(), 6, "{all}");
}
#[tokio::test]
async fn a_rule_counts_what_the_query_api_would_have_returned() {
let now = crate::api::now_nanos() as u64;
let (app, api, _) = boot_alerting(
"alerting",
r#"{
"rules": [
{ "name": "checkout-error-rate", "over": "1m", "when": "ratio > 5%",
"severity": "critical",
"query": { "signal": "logs", "where": [
{ "attr": "service.name", "eq": "checkout" },
{ "field": "severity_number", "gte": 17 } ] },
"of": { "signal": "logs", "where": [
{ "attr": "service.name", "eq": "checkout" } ] } },
{ "name": "ghost-traffic", "over": "1m", "when": "count > 0",
"query": { "signal": "logs", "where": [
{ "attr": "service.name", "eq": "not-deployed" } ] } }
]
}"#,
);
api.alerts.tick(&api).await;
let quiet = String::from_utf8(get(&app, "/api/v1/alerts", None).await.1).unwrap();
assert_eq!(quiet.matches(r#""state":"ok""#).count(), 2, "{quiet}");
otlp(&app, "/v1/logs", logs_export("checkout", now, 10)).await;
otlp(&app, "/v1/logs", logs_export("payments", now, 4)).await;
api.alerts.tick(&api).await;
let (status, bytes, _) = get(&app, "/api/v1/alerts", None).await;
assert_eq!(status, StatusCode::OK);
let body = String::from_utf8(bytes).unwrap();
assert!(body.contains(r#""matched":5"#), "{body}");
assert!(body.contains(r#""total":10"#), "{body}");
assert!(body.contains(r#""value":0.5"#), "{body}");
assert!(body.contains(r#""state":"firing""#), "{body}");
assert!(body.contains(r#""severity":"critical""#), "{body}");
assert!(body.contains(r#""name":"ghost-traffic""#), "{body}");
assert_eq!(body.matches(r#""state":"ok""#).count(), 1, "{body}");
assert!(body.contains(r#""error":null"#), "{body}");
let window = format!(
r#""signal":"logs","from":{},"to":{},
"where":[{{"attr":"service.name","eq":"checkout"}},
{{"field":"severity_number","gte":17}}]"#,
now.saturating_sub(60_000_000_000),
now + 60_000_000_000
);
let counted = query(&app, &format!("{{{window},\"limit\":100}}")).await;
assert!(counted.contains(r#""rows_matched":5"#), "{counted}");
}
fn webhook(
calls: usize,
) -> (
String,
std::net::SocketAddr,
std::sync::mpsc::Receiver<String>,
) {
use std::io::{Read, Write};
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("bind");
let addr = listener.local_addr().expect("addr");
let url = format!("http://{addr}/hook");
let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
for sock in listener.incoming().take(calls) {
let mut sock = sock.expect("accept");
let mut head = Vec::new();
let mut byte = [0u8; 1];
while !head.ends_with(b"\r\n\r\n") {
if sock.read(&mut byte).unwrap_or(0) == 0 {
break;
}
head.push(byte[0]);
}
let text = String::from_utf8_lossy(&head).to_lowercase();
let len: usize = text
.split("content-length:")
.nth(1)
.and_then(|t| t.split("\r\n").next())
.and_then(|t| t.trim().parse().ok())
.unwrap_or(0);
let mut body = vec![0u8; len];
sock.read_exact(&mut body).expect("body");
let _ = sock.write_all(b"HTTP/1.1 200 OK\r\ncontent-length: 0\r\n\r\n");
let _ = tx.send(String::from_utf8_lossy(&body).into_owned());
}
});
(url, addr, rx)
}
#[tokio::test]
async fn a_transition_pages_once_and_a_dead_target_does_not_swallow_a_live_one() {
let now = crate::api::now_nanos() as u64;
let (hook, hook_addr, got) = webhook(3);
drop(std::net::TcpStream::connect(hook_addr).expect("probe"));
assert_eq!(
got.recv_timeout(Duration::from_secs(5)).expect("probe"),
"",
"a client that sends nothing must not wedge the receiver"
);
let rules = format!(
r#"{{ "link_base": "http://mira.test",
"notify": [ {{ "name": "dead", "url": "http://127.0.0.1:1/x", "format": "json" }},
{{ "name": "oncall", "url": "{hook}", "format": "slack" }} ],
"rules": [
{{ "name": "checkout-errors", "over": "1m", "when": "ratio > 40%",
"severity": "critical", "notify": ["dead", "oncall"],
"query": {{ "signal": "logs", "where": [
{{ "attr": "service.name", "eq": "checkout" }},
{{ "field": "severity_number", "gte": 17 }} ] }},
"of": {{ "signal": "logs" }} }} ] }}"#
);
let (app, api, _) = boot_alerting("webhook", &rules);
api.alerts.tick(&api).await;
otlp(&app, "/v1/logs", logs_export("checkout", now, 10)).await;
api.alerts.tick(&api).await;
let fired = got.recv_timeout(Duration::from_secs(5)).expect("fire");
assert!(
fired.starts_with(r#"{"text":"[FIRING] checkout-errors"#),
"{fired}"
);
assert!(
fired.contains(r"50.00% > 40.00% over 1m (5 of 10 records)"),
"{fired}"
);
let expect = "<http://mira.test/#/logs?q=attr%3Aservice.name%3Dcheckout".to_owned()
+ "%20field%3Aseverity_number%3E%3D17&range=-60s|open in Mira>";
assert!(fired.contains(&expect), "{fired}");
api.alerts.tick(&api).await;
otlp(&app, "/v1/logs", logs_export("payments", now, 20)).await;
api.alerts.tick(&api).await;
let resolved = got.recv_timeout(Duration::from_secs(5)).expect("resolve");
assert!(
resolved.starts_with(r#"{"text":"[RESOLVED] checkout-errors"#),
"{resolved}"
);
assert!(
resolved.contains(r"16.67% > 40.00% over 1m (5 of 30 records)"),
"{resolved}"
);
let body = String::from_utf8(get(&app, "/api/v1/alerts", None).await.1).unwrap();
assert!(body.contains(r#""state":"ok""#), "{body}");
assert!(body.contains(r#""error":null"#), "{body}");
}
#[tokio::test]
async fn a_node_with_rules_evaluates_them_without_being_asked() {
let (_, quiet, _) = boot_alerting("timer-off", r#"{ "rules": [] }"#);
let (guard, quiet_log) = capture();
crate::alert::spawn(quiet.clone());
drop(guard);
assert!(quiet.alerts.json().contains(r#""alerts":[]"#));
assert_eq!(quiet_log.text(), "", "no rules, no banner and no evaluator");
let (app, api, _) = boot_alerting(
"timer-on",
r#"{ "every": "50ms", "rules": [
{ "name": "any-log", "over": "1m", "when": "count >= 1",
"query": { "signal": "logs" } } ] }"#,
);
otlp(
&app,
"/v1/logs",
logs_export("checkout", crate::api::now_nanos() as u64, 4),
)
.await;
let (guard, log) = capture();
crate::alert::spawn(api.clone());
drop(guard);
let banner = log.text();
assert!(banner.contains("rules=1"), "{banner}");
assert!(banner.contains("targets=0"), "{banner}");
assert!(banner.contains("alerting"), "{banner}");
let deadline = std::time::Instant::now() + Duration::from_secs(5);
loop {
let body = api.alerts.json();
if body.contains(r#""state":"firing""#) {
assert!(body.contains(r#""matched":4"#), "{body}");
break;
}
assert!(
std::time::Instant::now() < deadline,
"timer never fired: {body}"
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
#[tokio::test]
async fn a_rule_whose_query_fails_records_the_failure_instead_of_a_count() {
let now = crate::api::now_nanos();
let (app, api, root) = boot_alerting(
"corrupt",
r#"{ "rules": [
{ "name": "broken", "over": "1m", "when": "count >= 0",
"query": { "signal": "logs" } },
{ "name": "live", "over": "1m", "when": "count >= 0",
"query": { "signal": "traces" } } ] }"#,
);
const NANOS_PER_HOUR: i64 = 3_600 * 1_000_000_000;
let dir = root
.join("logs")
.join(format!("p={}", now.div_euclid(NANOS_PER_HOUR)))
.join(format!(
"{:020}-{:020}-{:08x}-{:012}-{:020}",
now - 1_000,
now,
0u32,
1u64,
0u64
));
std::fs::create_dir_all(&dir).expect("block dir");
std::fs::write(dir.join("logs.arrow"), b"not an arrow file").expect("corrupt table");
let (guard, log) = capture();
api.alerts.tick(&api).await;
drop(guard);
let text = log.text();
assert!(text.contains("alert rule failed"), "{text}");
assert!(text.contains("rule=broken"), "{text}");
assert!(text.contains("rule=live"), "{text}");
assert!(text.contains(r#"state="firing""#), "{text}");
let body = String::from_utf8(get(&app, "/api/v1/alerts", None).await.1).unwrap();
assert_eq!(body.matches(r#""state":"firing""#).count(), 1, "{body}");
assert_eq!(body.matches(r#""state":"ok""#).count(), 1, "{body}");
assert!(body.contains("logs.arrow"), "{body}");
assert_eq!(body.matches(r#""error":null"#).count(), 1, "{body}");
let _ = std::fs::remove_dir_all(&root);
}
fn rows_of(body: &str) -> String {
let s = body.find(r#""rows":["#).expect("no rows array") + 8;
body[s..body.find(r#"],"stats":"#).expect("no stats object")].to_owned()
}
fn bodies_of(rows: &str) -> Vec<String> {
let mut out = Vec::new();
let mut rest = rows;
while let Some(i) = rest.find(r#""body":""#) {
let s = &rest[i + 8..];
let e = s.find('"').expect("unterminated body");
out.push(s[..e].to_owned());
rest = &s[e..];
}
out
}
fn node_cfg(root: &std::path::Path, wal: bool) -> pipeline::Config {
pipeline::Config {
data_dir: root.to_path_buf(),
max_block_age: Duration::from_secs(600),
wal: wal.then(|| {
Arc::new(mira_core::wal::Wal::open(root, mira_core::block::node_id("mira")).unwrap())
}),
..Default::default()
}
}
struct Node {
app: Router,
api: api::Api,
flushers: [pipeline::Flushers; 3],
}
impl Node {
async fn stop(self) {
let Node { app, api, flushers } = self;
drop((app, api));
for h in flushers {
tokio::time::timeout(Duration::from_secs(20), h)
.await
.expect("a flusher did not stop; some Ingest clone outlived the router")
.expect("flusher panicked");
}
}
async fn kill(self) {
let Node { app, api, flushers } = self;
for mut h in flushers {
h.abort();
let _ = h.await;
}
drop((app, api));
forget_open_blocks();
}
}
fn forget_open_blocks() {
for r in &pipeline::REJECTS {
r.forget_open();
}
}
async fn restart_with(cfg: pipeline::Config) -> Node {
let root = cfg.data_dir.clone();
let id = cfg.node;
let cfg = Arc::new(cfg);
let (logs, o_logs, h_logs) = pipeline::spawn::<mira_core::logs::LogsBuilder>(&cfg);
let (traces, o_traces, h_traces) = pipeline::spawn::<mira_core::traces::TracesBuilder>(&cfg);
let (metrics, o_metrics, h_metrics) =
pipeline::spawn::<mira_core::metrics::MetricsBuilder>(&cfg);
crate::replay(&root, id, logs.clone(), traces.clone(), metrics.clone())
.await
.expect("replay");
let recv = receiver::Receivers {
logs,
traces,
metrics,
max_request_bytes: crate::config::Config::default().max_request_bytes,
};
let api = api::Api {
data_dir: Arc::new(root),
open: [o_logs, o_traces, o_metrics],
alerts: Arc::default(),
};
Node {
app: router_for(recv, api.clone()),
api,
flushers: [h_logs, h_traces, h_metrics],
}
}
async fn restart(root: &std::path::Path, wal: bool) -> Node {
restart_with(node_cfg(root, wal)).await
}
async fn seal(root: &std::path::Path, signal: &str) -> Vec<mira_core::block::BlockRef> {
let before = mira_core::block::scan(root, signal).unwrap().len();
let deadline = std::time::Instant::now() + Duration::from_secs(20);
loop {
let blocks = mira_core::block::scan(root, signal).unwrap();
if blocks.len() > before {
return blocks;
}
assert!(
std::time::Instant::now() < deadline,
"{signal} never sealed; everything below this would have proved nothing"
);
tokio::time::sleep(Duration::from_millis(5)).await;
}
}
async fn serve(app: Router) -> String {
let l = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = l.local_addr().unwrap().to_string();
tokio::spawn(axum::serve(l, app).into_future());
addr
}
#[tokio::test]
async fn a_record_reads_back_identically_at_every_stage_of_its_life() {
let root = fresh_dir("lifecycle");
let n = restart_with(pipeline::Config {
max_block_age: Duration::from_millis(300),
..node_cfg(&root, true)
})
.await;
let doc = r#"{"signal":"logs","from":0,"to":8000,"limit":1000}"#;
let id = mira_core::block::node_id("mira");
otlp(&n.app, "/v1/logs", logs_export("checkout", 1_000, 200)).await;
assert!(
mira_core::block::scan(&root, "logs").unwrap().is_empty(),
"publishing before the acknowledgement is the other contract, and \
boot_sealing is the test for it"
);
let framed: u64 = std::fs::read_dir(root.join(".wal"))
.unwrap()
.map(|e| e.unwrap().metadata().unwrap().len())
.sum();
assert!(framed > 0, "acknowledged without being logged");
let open = rows_of(&query(&n.app, doc).await);
assert_eq!(open.matches("handled request").count(), 200);
let published = seal(&root, "logs").await;
assert_eq!(published.len(), 1);
assert_eq!(
published[0].wal_hi, 1,
"a block published under a log claims the frames it holds, or the next \
boot replays them into a second copy"
);
assert_eq!(
rows_of(&query(&n.app, doc).await),
open,
"the seal changed the answer"
);
let dir = published[0].dir.clone();
let table = dir.join("logs.arrow");
let plain_bytes = std::fs::metadata(&table).unwrap().len();
let mapped = mira_core::block::open_table(&table).unwrap();
let rows_of_table = |t: &mira_core::block::MappedTable| -> usize {
t.batches.iter().map(|b| b.num_rows()).sum()
};
let mapped_rows = rows_of_table(&mapped);
assert_eq!(mapped_rows, 200);
let hot = mapped.zero_copy_ratio();
assert_eq!(
hot.0, hot.1,
"a hot block copied buffers out of its mapping"
);
assert_eq!(
mira_core::block::compact(&root, "logs", id, i64::MAX).unwrap(),
1
);
assert!(
dir.join("cold").exists(),
"the marker is written last, so its absence means an unfinished rewrite"
);
assert!(
std::fs::metadata(&table).unwrap().len() < plain_bytes,
"a table that did not shrink was not ZSTD-rewritten"
);
let cold = mira_core::block::open_table(&table)
.unwrap()
.zero_copy_ratio();
assert!(
cold.0 < cold.1 && cold.0 < hot.0,
"a compressed block cannot be zero-copy — decompression has to allocate — \
so a cold table reporting {cold:?} against the hot {hot:?} means the \
rewrite left plain batches behind"
);
assert_eq!(
rows_of_table(&mapped),
mapped_rows,
"the rename pulled the mapping out from under a live reader"
);
assert_eq!(
rows_of(&query(&n.app, doc).await),
open,
"compaction changed the answer"
);
assert_eq!(
mira_core::block::compact(&root, "logs", id, i64::MAX).unwrap(),
0,
"the cold block was compacted a second time"
);
assert_eq!(
mira_core::block::expire(&root, "logs", i64::MAX).unwrap(),
1
);
assert!(mira_core::block::scan(&root, "logs").unwrap().is_empty());
assert!(
rows_of(&query(&n.app, doc).await).is_empty(),
"an expired block is still being answered from"
);
assert_eq!(
rows_of_table(&mapped),
mapped_rows,
"the unlink took the mapping with it"
);
n.stop().await;
}
#[tokio::test]
async fn a_block_seals_on_size_and_on_shutdown_and_neither_is_replayed_afterwards() {
let doc = r#"{"signal":"logs","from":0,"to":8000,"limit":100}"#;
let root = fresh_dir("seal-size");
let n = restart_with(pipeline::Config {
target_block_bytes: 1,
..node_cfg(&root, true)
})
.await;
otlp(&n.app, "/v1/logs", logs_export("checkout", 1_000, 4)).await;
assert_eq!(seal(&root, "logs").await.len(), 1);
n.stop().await;
let root = fresh_dir("seal-shutdown");
let n = restart(&root, true).await;
otlp(&n.app, "/v1/logs", logs_export("checkout", 1_000, 5)).await;
let before = rows_of(&query(&n.app, doc).await);
assert!(
mira_core::block::scan(&root, "logs").unwrap().is_empty(),
"something sealed early and the assertion below is now free"
);
n.stop().await;
assert_eq!(
mira_core::block::scan(&root, "logs").unwrap().len(),
1,
"the shutdown cost the block that was filling"
);
let node = mira_core::block::node_id("mira");
let watermarks = mira_core::block::wal_watermarks(&root, node).unwrap();
assert_eq!(watermarks[0], 1, "the block did not claim its frame");
let mut seen: Vec<u64> = Vec::new();
let mut sink = |_: mira_core::wal::Signal, seq: u64, _: &[u8]| {
seen.push(seq);
Ok(())
};
let again = mira_core::wal::Wal::replay(&root, node, watermarks, &mut sink).unwrap();
assert_eq!(
(again.replayed, again.skipped),
(0, 1),
"a frame inside a published block was replayed anyway"
);
let all = mira_core::wal::Wal::replay(&root, node, [0; 3], &mut sink).unwrap();
assert_eq!((all.replayed, all.skipped), (1, 0));
assert_eq!(
seen,
[0],
"the watermark is not one past the frame the block published"
);
let n = restart(&root, true).await;
let after = query(&n.app, doc).await;
assert_eq!(after.matches("handled request").count(), 5, "{after}");
assert_eq!(rows_of(&after), before, "the restart moved a row");
n.stop().await;
assert_eq!(
mira_core::block::scan(&root, "logs").unwrap().len(),
1,
"the second shutdown published a duplicate or an empty block"
);
}
#[tokio::test]
async fn a_kill_with_the_block_open_loses_no_acknowledged_export_and_duplicates_none() {
let root = fresh_dir("crash-open");
let doc = r#"{"signal":"logs","from":0,"to":8000,"limit":100}"#;
let n = restart(&root, true).await;
for base in [1_000u64, 2_000, 3_000] {
otlp(&n.app, "/v1/logs", logs_export("checkout", base, 5)).await;
}
let before = rows_of(&query(&n.app, doc).await);
assert_eq!(before.matches("handled request").count(), 15);
assert!(
mira_core::block::scan(&root, "logs").unwrap().is_empty(),
"the data reached disk before the kill, so this proves nothing"
);
n.kill().await;
let n = restart(&root, true).await;
let after = query(&n.app, doc).await;
assert_eq!(
after.matches("handled request").count(),
15,
"recovery lost or duplicated a record: {after}"
);
assert_eq!(rows_of(&after), before);
n.stop().await;
assert_eq!(
mira_core::block::wal_watermarks(&root, mira_core::block::node_id("mira")).unwrap()[0],
3
);
let n = restart(&root, true).await;
let third = query(&n.app, doc).await;
assert_eq!(
third.matches("handled request").count(),
15,
"the second boot replayed frames the first one had published: {third}"
);
n.stop().await;
}
#[tokio::test]
async fn a_kill_between_the_seal_and_the_rename_drops_the_staging_dir_and_not_the_data() {
let root = fresh_dir("crash-staging");
let doc = r#"{"signal":"logs","from":0,"to":8000,"limit":100}"#;
let id = mira_core::block::node_id("mira");
let n = restart(&root, true).await;
otlp(&n.app, "/v1/logs", logs_export("checkout", 1_000, 5)).await;
n.kill().await;
let staged = root.join(".tmp").join(format!(
"logs-{id:08x}-000000000000-{:020}-{:020}",
1_000, 1_004
));
std::fs::create_dir_all(&staged).unwrap();
std::fs::write(staged.join("logs.arrow"), b"half a table").unwrap();
let n = restart(&root, true).await;
let body = query(&n.app, doc).await;
assert_eq!(
body.matches("handled request").count(),
5,
"the export was acknowledged on the log and the log did not bring it back: {body}"
);
assert!(
!staged.exists(),
"the staging directory outlived the process that made it; nothing will \
ever look in .tmp again, so that is leaked disk for the volume's life"
);
n.stop().await;
assert_eq!(mira_core::block::scan(&root, "logs").unwrap().len(), 1);
}
#[tokio::test]
async fn a_torn_wal_tail_costs_only_the_frame_that_was_in_flight() {
let root = fresh_dir("crash-torn");
let doc = r#"{"signal":"logs","from":0,"to":8000,"limit":100}"#;
let n = restart(&root, true).await;
otlp(&n.app, "/v1/logs", logs_export("checkout", 1_000, 5)).await;
otlp(&n.app, "/v1/logs", logs_export("payments", 2_000, 5)).await;
n.kill().await;
let seg = std::fs::read_dir(root.join(".wal"))
.unwrap()
.map(|e| e.unwrap().path())
.find(|p| p.extension().is_some_and(|e| e == "wal"))
.expect("a segment");
let len = std::fs::metadata(&seg).unwrap().len();
std::fs::OpenOptions::new()
.write(true)
.open(&seg)
.unwrap()
.set_len(len - 4)
.unwrap();
let n = restart(&root, true).await;
let body = query(&n.app, doc).await;
assert_eq!(
body.matches("checkout handled request").count(),
5,
"the whole frame in front of the tear did not survive: {body}"
);
assert_eq!(
body.matches("payments handled request").count(),
0,
"half a frame was replayed; the checksum exists to stop exactly that"
);
n.stop().await;
assert_eq!(
mira_core::block::wal_watermarks(&root, mira_core::block::node_id("mira")).unwrap()[0],
1,
"the recovered block claimed a frame the tear had swallowed"
);
}
#[tokio::test]
async fn a_torn_wal_tail_does_not_stop_the_node_from_booting() {
let root = fresh_dir("crash-torn-boot");
let id = mira_core::block::node_id("mira");
let n = restart(&root, true).await;
otlp(&n.app, "/v1/logs", logs_export("checkout", 1_000, 5)).await;
n.kill().await;
let seg = std::fs::read_dir(root.join(".wal"))
.unwrap()
.map(|e| e.unwrap().path())
.find(|p| p.extension().is_some_and(|e| e == "wal"))
.expect("a segment");
let len = std::fs::metadata(&seg).unwrap().len();
std::fs::OpenOptions::new()
.write(true)
.open(&seg)
.unwrap()
.set_len(len - 4)
.unwrap();
mira_core::wal::Wal::open(&root, id).expect(
"a hard kill leaves a torn tail, so refusing to open the log refuses \
every restart after a crash",
);
}
#[tokio::test]
async fn a_kill_inside_the_cold_rewrite_reads_correctly_and_is_finished_next_sweep() {
let (app, root) = boot_sealing("crash-compact");
let doc = r#"{"signal":"logs","from":0,"to":8000,"limit":100}"#;
let id = mira_core::block::node_id("mira");
otlp(&app, "/v1/logs", logs_export("checkout", 1_000, 8)).await;
let dir = mira_core::block::scan(&root, "logs").unwrap()[0]
.dir
.clone();
let plain = rows_of(&query(&app, doc).await);
assert_eq!(plain.matches("handled request").count(), 8);
let keep = root.join("uncompressed");
std::fs::create_dir_all(&keep).unwrap();
for e in std::fs::read_dir(&dir).unwrap() {
let p = e.unwrap().path();
if p.extension().is_some_and(|x| x == "arrow") {
std::fs::copy(&p, keep.join(p.file_name().unwrap())).unwrap();
}
}
assert_eq!(
mira_core::block::compact(&root, "logs", id, i64::MAX).unwrap(),
1
);
std::fs::copy(keep.join("logs.arrow"), dir.join("logs.arrow")).unwrap();
std::fs::remove_file(dir.join("cold")).unwrap();
assert_eq!(
rows_of(&query(&app, doc).await),
plain,
"a half-compacted block did not read back"
);
assert_eq!(
mira_core::block::compact(&root, "logs", id, i64::MAX).unwrap(),
1,
"the unfinished rewrite was not retried"
);
assert!(dir.join("cold").exists());
assert_eq!(rows_of(&query(&app, doc).await), plain);
forget_open_blocks();
}
#[tokio::test]
async fn the_engine_the_api_and_both_ui_transports_answer_row_for_row() {
let (app, root) = boot_sealing("parity");
let addr = serve(app.clone()).await;
let doc = r#"{"signal":"logs","from":0,"to":8000,"limit":100}"#;
otlp(&app, "/v1/logs", logs_export("checkout", 1_000, 6)).await;
otlp(&app, "/v1/logs", logs_export("payments", 2_000, 6)).await;
let q = crate::api::parse_search(doc, crate::api::now_nanos()).expect("parse");
let direct = mira_core::query::search(&root, &q).expect("search");
let engine = rows_of(&crate::api::envelope("rows", &direct, Duration::ZERO));
assert_eq!(engine.matches("handled request").count(), 12, "{engine}");
let over_http = query(&app, doc).await;
assert_eq!(rows_of(&over_http), engine, "the API moved a row");
let (local, remote) = {
let (d, b, a) = (root.clone(), doc.to_owned(), addr.clone());
tokio::task::spawn_blocking(move || {
let route = "/api/v1/query";
(
crate::tui::Source::Local(d).post(route, &b).expect("local"),
crate::tui::Source::Remote(a)
.post(route, &b)
.expect("remote"),
)
})
.await
.unwrap()
};
for (what, parsed) in [("local", &local), ("remote", &remote)] {
let rows = parsed["rows"].as_vec().expect(what);
assert_eq!(rows.len(), 12, "{what}");
for (i, row) in rows.iter().enumerate() {
let body = row["body"].as_str().expect("body");
assert!(engine.contains(body), "{what} row {i} is not the engine's");
}
let first = &rows[0];
assert_eq!(
first["body"].as_str().unwrap(),
"payments handled request 5",
"{what} did not order newest-first"
);
assert_eq!(first["time_unix_nano"].as_i64(), None, "{what}");
assert_eq!(
first["time_unix_nano"]
.as_str()
.and_then(|s| s.parse().ok()),
Some(2_005i64),
"{what}"
);
assert_eq!(first["severity_number"].as_i64(), Some(9), "{what}");
assert_eq!(parsed["stats"]["rows_matched"].as_i64(), Some(12), "{what}");
}
assert!(
over_http.contains(r#""time_unix_nano":"2005""#),
"{over_http}"
);
assert!(over_http.contains(r#""severity_number":9,"#), "{over_http}");
assert!(
!over_http.contains(r#""time_unix_nano":2005"#),
"a nanosecond timestamp went out as a JSON number: {over_http}"
);
forget_open_blocks();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_reader_under_a_live_writer_never_sees_a_row_vanish_or_arrive_out_of_order() {
let (app, _root) = boot("parity-live");
let addr = serve(app.clone()).await;
let doc = r#"{"signal":"logs","from":0,"to":900000,"limit":10000}"#;
let (batches, per) = (12usize, 5usize);
let probe = |addr: String| {
let app = app.clone();
let doc = doc.to_owned();
async move {
let http = bodies_of(&rows_of(&query(&app, &doc).await));
let tui = tokio::task::spawn_blocking(move || {
crate::tui::Source::Remote(addr).post("/api/v1/query", &doc)
})
.await
.unwrap()
.expect("remote");
let tui = tui["rows"]
.as_vec()
.expect("rows")
.iter()
.map(|r| r["body"].as_str().expect("body").to_owned())
.collect::<Vec<_>>();
[http, tui]
}
};
let mut seen: Vec<String> = Vec::new();
let check = |answers: [Vec<String>; 2], seen: &mut Vec<String>| {
for bodies in answers {
assert!(
bodies.ends_with(seen),
"a row disappeared or moved under a live writer.\n had: {seen:?}\n got: {bodies:?}"
);
*seen = bodies;
}
};
check(probe(addr.clone()).await, &mut seen);
assert!(seen.is_empty(), "the writer has not started yet");
let writer = {
let app = app.clone();
tokio::spawn(async move {
for b in 0..batches {
let batch = logs_export(&format!("svc{b}"), 1_000 + b as u64 * 100, per);
otlp(&app, "/v1/logs", batch).await;
tokio::time::sleep(Duration::from_millis(5)).await;
}
})
};
let mut rounds = 0;
loop {
let done = writer.is_finished();
check(probe(addr.clone()).await, &mut seen);
rounds += 1;
if done {
break;
}
tokio::time::sleep(Duration::from_millis(3)).await;
}
writer.await.unwrap();
assert!(rounds >= 2, "the reader never overlapped the writer");
let all = query(&app, doc).await;
assert_eq!(
all.matches("handled request").count(),
batches * per,
"the writer's last batch is missing: {all}"
);
let rows = rows_of(&all);
assert_eq!(
rows.matches(r#""time_unix_nano":""#).count(),
bodies_of(&rows).len()
);
forget_open_blocks();
}
#[tokio::test]
async fn a_count_at_the_threshold_fires_gte_and_not_gt() {
let now = crate::api::now_nanos() as u64;
let (app, api, _root) = boot_alerting(
"alert-threshold",
r#"{ "rules": [
{ "name": "at", "over": "1m", "when": "count >= 5",
"query": { "signal": "logs", "where": [
{ "attr": "service.name", "eq": "checkout" },
{ "field": "severity_number", "gte": 17 } ] } },
{ "name": "over", "over": "1m", "when": "count > 5",
"query": { "signal": "logs", "where": [
{ "attr": "service.name", "eq": "checkout" },
{ "field": "severity_number", "gte": 17 } ] } } ] }"#,
);
otlp(&app, "/v1/logs", logs_export("checkout", now, 10)).await;
api.alerts.tick(&api).await;
let at_five = String::from_utf8(get(&app, "/api/v1/alerts", None).await.1).unwrap();
assert!(at_five.contains(r#""matched":5"#), "{at_five}");
assert!(
at_five.contains(r#""name":"at","state":"firing""#),
"a count exactly at a `>=` threshold did not fire: {at_five}"
);
assert!(
at_five.contains(r#""name":"over","state":"ok""#),
"a count exactly at a `>` threshold fired: {at_five}"
);
otlp(&app, "/v1/logs", logs_export("checkout", now + 100, 2)).await;
api.alerts.tick(&api).await;
let at_six = String::from_utf8(get(&app, "/api/v1/alerts", None).await.1).unwrap();
assert!(at_six.contains(r#""matched":6"#), "{at_six}");
assert!(
at_six.contains(r#""name":"over","state":"firing""#),
"{at_six}"
);
assert!(
at_six.contains(r#""name":"at","state":"firing""#),
"{at_six}"
);
forget_open_blocks();
}
#[tokio::test]
async fn an_empty_window_never_breaches_in_either_metric() {
let now = crate::api::now_nanos() as u64;
let (app, api, _root) = boot_alerting(
"alert-empty",
r#"{ "rules": [
{ "name": "ratio", "over": "1m", "when": "ratio > 0%",
"query": { "signal": "logs", "where": [
{ "field": "severity_number", "gte": 17 } ] },
"of": { "signal": "logs" } },
{ "name": "count", "over": "1m", "when": "count > 0",
"query": { "signal": "logs" } },
{ "name": "no-match", "over": "1m", "when": "count > 0",
"query": { "signal": "logs", "where": [
{ "attr": "service.name", "eq": "not-deployed" } ] } } ] }"#,
);
api.alerts.tick(&api).await;
let nothing = api.alerts.json();
assert_eq!(nothing.matches(r#""state":"ok""#).count(), 3, "{nothing}");
assert_eq!(nothing.matches(r#""value":0,"#).count(), 3, "{nothing}");
assert_eq!(nothing.matches(r#""error":null"#).count(), 3, "{nothing}");
otlp(
&app,
"/v1/logs",
logs_export("checkout", now - 2 * 3_600_000_000_000, 20),
)
.await;
api.alerts.tick(&api).await;
let stale = api.alerts.json();
assert_eq!(
stale.matches(r#""state":"ok""#).count(),
3,
"a rule counted data older than its own window: {stale}"
);
assert_eq!(stale.matches(r#""matched":0"#).count(), 3, "{stale}");
otlp(&app, "/v1/logs", logs_export("checkout", now, 4)).await;
api.alerts.tick(&api).await;
let live = api.alerts.json();
assert!(
live.contains(r#""name":"count","state":"firing""#),
"{live}"
);
assert!(live.contains(r#""name":"no-match","state":"ok""#), "{live}");
assert_eq!(live.matches(r#""error":null"#).count(), 3, "{live}");
forget_open_blocks();
}
#[tokio::test]
async fn an_alert_window_that_straddles_a_block_boundary_counts_each_record_once() {
let now = crate::api::now_nanos() as u64;
let root = fresh_dir("alert-straddle");
let engine = crate::alert::Engine::new(
crate::alert::Rules::parse(
r#"{ "rules": [
{ "name": "errors", "over": "5m", "when": "count >= 1",
"query": { "signal": "logs", "where": [
{ "attr": "service.name", "eq": "checkout" },
{ "field": "severity_number", "gte": 17 } ] } } ] }"#,
)
.expect("rules"),
);
let n = restart(&root, true).await;
otlp(&n.app, "/v1/logs", logs_export("checkout", now, 6)).await;
n.stop().await;
assert_eq!(mira_core::block::scan(&root, "logs").unwrap().len(), 1);
let n = restart(&root, true).await;
otlp(&n.app, "/v1/logs", logs_export("checkout", now + 1_000, 4)).await;
engine.tick(&n.api).await;
let body = engine.json();
assert_eq!(
mira_core::block::scan(&root, "logs").unwrap().len(),
1,
"the second half sealed, so nothing straddled and this proved nothing"
);
assert!(
body.contains(r#""matched":5"#),
"3 on disk plus 2 open is 5; a rule that saw one side saw 3 or 2: {body}"
);
n.stop().await;
}
#[tokio::test]
async fn two_evaluations_of_one_window_agree_and_the_transitions_run_both_ways() {
let now = crate::api::now_nanos() as u64;
let (app, api, _root) = boot_alerting(
"alert-determinism",
r#"{ "rules": [
{ "name": "checkout-errors", "over": "5m", "when": "ratio > 40%",
"severity": "critical",
"query": { "signal": "logs", "where": [
{ "attr": "service.name", "eq": "checkout" },
{ "field": "severity_number", "gte": 17 } ] },
"of": { "signal": "logs" } } ] }"#,
);
let stable = |body: &str| {
let i = body.find(r#""evaluated_at":""#).expect("stamp");
let rest = &body[i + 16..];
format!("{}{}", &body[..i], &rest[rest.find('"').unwrap()..])
};
let firing_since = |body: &str| {
let i = body.find(r#""firing_since":""#).expect("firing_since") + 16;
body[i..i + body[i..].find('"').unwrap()]
.parse::<i64>()
.expect("a nanosecond stamp")
};
otlp(&app, "/v1/logs", logs_export("checkout", now, 10)).await;
api.alerts.tick(&api).await;
let first = api.alerts.json();
assert!(first.contains(r#""state":"firing""#), "{first}");
assert!(first.contains(r#""value":0.5"#), "{first}");
api.alerts.tick(&api).await;
let second = api.alerts.json();
assert_eq!(
stable(&first),
stable(&second),
"the same window evaluated twice disagreed with itself"
);
assert_ne!(first, second, "the evaluation stamp did not move");
assert!(second.contains(r#""over_nano":"300000000000""#), "{second}");
assert!(second.contains(r#""for_nano":"0""#), "{second}");
assert!(second.contains(r#""matched":5,"total":10"#), "{second}");
let fired_at = firing_since(&second);
assert!(fired_at > 0, "{second}");
otlp(&app, "/v1/logs", logs_export("payments", now, 20)).await;
api.alerts.tick(&api).await;
let resolved = api.alerts.json();
assert!(resolved.contains(r#""state":"ok""#), "{resolved}");
assert!(
resolved.contains(r#""since":null,"firing_since":null"#),
"the resolve left the firing clock running: {resolved}"
);
let mut spike = logs_export("checkout", now + 1_000, 40);
for r in &mut spike.resource_logs[0].scope_logs[0].log_records {
r.severity_number = 17;
}
otlp(&app, "/v1/logs", spike).await;
api.alerts.tick(&api).await;
let again = api.alerts.json();
assert!(again.contains(r#""state":"firing""#), "{again}");
assert!(again.contains(r#""matched":45,"total":70"#), "{again}");
assert!(
firing_since(&again) > fired_at,
"the second incident reported the first one's start time"
);
forget_open_blocks();
}
async fn replica(name: &str) -> (Router, String) {
let root = fresh_dir(name);
let id = mira_core::block::node_id(name);
let n = restart_with(pipeline::Config {
data_dir: root.clone(),
node: id,
max_block_age: Duration::from_secs(600),
wal: Some(Arc::new(mira_core::wal::Wal::open(&root, id).unwrap())),
..Default::default()
})
.await;
let app = n.app.clone();
let addr = serve(n.app.clone()).await;
drop(n);
(app, addr)
}
fn proxied(replicas: Vec<String>) -> Router {
crate::proxy::router(
crate::proxy::Proxy::new(replicas, crate::config::Config::default().max_request_bytes)
.expect("at least one replica"),
)
}
#[tokio::test]
async fn a_proxy_answers_the_probes_kubernetes_gates_its_endpoints_on() {
let p = proxied(vec!["http://127.0.0.1:1".into()]);
for path in ["/health", "/readyz"] {
let (status, body, _) = get(&p, path, None).await;
assert_eq!(status, StatusCode::OK, "{path}");
assert_eq!(String::from_utf8(body).unwrap(), r#"{"status":"ok"}"#);
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_proxy_splits_an_export_across_replicas_and_pages_them_back_as_one_store() {
let (a, addr_a) = replica("proxy-a").await;
let (b, addr_b) = replica("proxy-b").await;
let p = proxied(vec![format!("http://{addr_a}"), format!("http://{addr_b}")]);
const SERVICES: [&str; 8] = [
"checkout",
"payments",
"search",
"cart",
"auth",
"billing",
"shipping",
"inventory",
];
let mut batch = ExportLogsServiceRequest::default();
for s in SERVICES {
batch
.resource_logs
.extend(logs_export(s, 1_000, 5).resource_logs);
}
otlp(&p, "/v1/logs", batch).await;
let window = r#"{"signal":"logs","from":0,"to":100000"#;
let (mut on_a, mut on_b) = (0, 0);
for s in SERVICES {
let doc =
format!(r#"{window},"limit":500,"where":[{{"attr":"service.name","eq":"{s}"}}]}}"#);
let here = bodies_of(&rows_of(&query(&a, &doc).await)).len();
let there = bodies_of(&rows_of(&query(&b, &doc).await)).len();
assert_eq!(
(here.min(there), here.max(there)),
(0, 5),
"{s} was split across both replicas: {here} and {there}"
);
on_a += here;
on_b += there;
}
assert_eq!(on_a + on_b, 40);
assert!(on_a > 0 && on_b > 0, "everything landed on one replica");
let all = query(&p, &format!(r#"{window},"limit":500}}"#)).await;
let merged = bodies_of(&rows_of(&all));
assert_eq!(merged.len(), 40, "{all}");
assert!(!all.contains(r#""next""#), "a short page is the last page");
assert!(all.contains(r#""rows_matched":40"#), "{all}");
assert!(!all.contains(r#""cursors""#), "{all}");
let with = query(&p, &format!(r#"{window},"limit":500,"cursors":"true"}}"#)).await;
let i = with
.find(r#""cursors":["#)
.unwrap_or_else(|| panic!("{with}"));
let list = &with[i + 11..];
let cursors: Vec<&str> = list[..list.find(']').expect("unclosed")]
.split(',')
.collect();
assert_eq!(cursors.len(), 40, "{with}");
assert_eq!(
cursors
.iter()
.collect::<std::collections::HashSet<_>>()
.len(),
40,
"the merge placed one row twice: {with}"
);
assert_eq!(bodies_of(&rows_of(&with)), merged, "{with}");
let mut pages: Vec<String> = Vec::new();
let mut after = String::new();
for _ in 0..12 {
let body = query(&p, &format!(r#"{window},"limit":7{after}}}"#)).await;
pages.extend(bodies_of(&rows_of(&body)));
let Some(i) = body.find(r#""next":""#) else {
break;
};
let c = &body[i + 8..];
after = format!(r#","after":"{}""#, &c[..c.find('"').unwrap()]);
}
assert_eq!(pages, merged, "the pages did not reassemble the whole");
let (status, body) = post(&p, "/api/v1/map", "application/json", "{}".into()).await;
assert_eq!(status, StatusCode::NOT_IMPLEMENTED, "{body}");
assert!(body.contains("/api/v1/map"), "{body}");
let (status, body) = post(
&p,
"/v1/logs",
"application/json",
br#"{"resourceLogs":[]}"#.to_vec(),
)
.await;
assert_eq!(status, StatusCode::OK, "{body}");
assert_eq!(body, "{}");
let (status, body) = post(&p, "/v1/logs", "application/x-protobuf", vec![0xff, 0xff]).await;
assert_eq!(status, StatusCode::BAD_REQUEST, "{body}");
use mira_proto::collector::metrics::v1::ExportMetricsServiceRequest;
use mira_proto::metrics::v1::metric::Data;
use mira_proto::metrics::v1::number_data_point::Value as NumValue;
use mira_proto::metrics::v1::{Gauge, Metric, NumberDataPoint, ResourceMetrics, ScopeMetrics};
otlp(
&p,
"/v1/traces",
ExportTraceServiceRequest {
resource_spans: SERVICES
.iter()
.map(|s| ResourceSpans {
resource: Some(Resource {
attributes: vec![kv("service.name", s)],
..Default::default()
}),
scope_spans: vec![ScopeSpans {
spans: vec![Span {
trace_id: vec![0xab; 16].into(),
span_id: vec![0xcd; 8].into(),
name: format!("GET /{s}"),
start_time_unix_nano: 1_000,
end_time_unix_nano: 1_500,
..Default::default()
}],
..Default::default()
}],
..Default::default()
})
.collect(),
},
)
.await;
otlp(
&p,
"/v1/metrics",
ExportMetricsServiceRequest {
resource_metrics: SERVICES
.iter()
.map(|s| ResourceMetrics {
resource: Some(Resource {
attributes: vec![kv("service.name", s)],
..Default::default()
}),
scope_metrics: vec![ScopeMetrics {
metrics: vec![Metric {
name: "http.server.requests".into(),
data: Some(Data::Gauge(Gauge {
data_points: vec![NumberDataPoint {
time_unix_nano: 1_000,
value: Some(NumValue::AsInt(1)),
..Default::default()
}],
})),
..Default::default()
}],
..Default::default()
}],
..Default::default()
})
.collect(),
},
)
.await;
let doc = r#"{"signal":"traces","from":0,"to":100000,"limit":500}"#;
let both = query(&p, doc).await;
assert!(both.contains(r#""rows_matched":8"#), "{both}");
for r in [&a, &b] {
let one = query(r, doc).await;
assert!(
(1..8).any(|k| one.contains(&format!(r#""rows_matched":{k}"#))),
"one replica answered the whole export: {one}"
);
}
let point = br#"{"name":"http.server.requests","from":0,"to":100000}"#.to_vec();
let (status, body) = post(
&p,
"/api/v1/metrics/query",
"application/json",
point.clone(),
)
.await;
assert_eq!(status, StatusCode::NOT_IMPLEMENTED, "{body}");
let mut seen = 0;
for r in [&a, &b] {
let (status, one) = post(
r,
"/api/v1/metrics/query",
"application/json",
point.clone(),
)
.await;
assert_eq!(status, StatusCode::OK, "{one}");
let here = SERVICES
.iter()
.filter(|s| one.contains(&format!(r#""service.name":"{s}""#)))
.count();
assert!(
here > 0 && here < 8,
"one replica holds {here} series: {one}"
);
seen += here;
}
assert_eq!(seen, 8, "the metric export did not arrive whole");
forget_open_blocks();
}
async fn stub(status: StatusCode, body: &'static str) -> String {
let addr = serve(Router::new().fallback(move || async move { (status, body) })).await;
format!("http://{addr}")
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_replica_that_will_not_answer_fails_the_request_instead_of_shortening_it() {
const DEAD: &str = "http://127.0.0.1:1";
let doc = r#"{"signal":"logs","limit":5}"#;
let p = proxied(vec![DEAD.into()]);
let (status, body) = post(&p, "/api/v1/query", "application/json", doc.into()).await;
assert_eq!(status, StatusCode::BAD_GATEWAY, "{body}");
assert!(body.contains("127.0.0.1:1"), "{body}");
let p = proxied(vec![stub(StatusCode::BAD_REQUEST, "signal: nope").await]);
let (status, body) = post(&p, "/api/v1/query", "application/json", doc.into()).await;
assert_eq!(status, StatusCode::BAD_GATEWAY, "{body}");
assert!(
body.contains("HTTP 400") && body.contains("signal: nope"),
"{body}"
);
let p = proxied(vec![stub(StatusCode::OK, r#"{"rows":[]}"#).await]);
let (status, body) = post(&p, "/api/v1/query", "application/json", doc.into()).await;
assert_eq!(status, StatusCode::BAD_GATEWAY, "{body}");
let addr = serve(Router::new().fallback(|| async { vec![0x80u8, 0xff] })).await;
let p = proxied(vec![format!("http://{addr}")]);
let (status, body) = post(&p, "/api/v1/query", "application/json", doc.into()).await;
assert_eq!(status, StatusCode::BAD_GATEWAY, "{body}");
let p = proxied(vec![DEAD.into()]);
for bad in [
"{",
r#"{"signal":"nope"}"#,
r#"{"cursors":"false","limit":5}"#,
] {
let (status, body) = post(&p, "/api/v1/query", "application/json", bad.into()).await;
assert_eq!(status, StatusCode::BAD_REQUEST, "{bad} gave {body}");
}
let mut batch = ExportLogsServiceRequest::default();
for s in [
"checkout",
"payments",
"search",
"cart",
"auth",
"billing",
"shipping",
"inventory",
] {
batch
.resource_logs
.extend(logs_export(s, 1_000, 1).resource_logs);
}
let body = batch.encode_to_vec();
let p = proxied(vec![DEAD.into(), DEAD.into()]);
let (status, why) = post(&p, "/v1/logs", "application/x-protobuf", body.clone()).await;
assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{why}");
let p = proxied(vec![DEAD.into(), stub(StatusCode::BAD_REQUEST, "no").await]);
let (status, why) = post(&p, "/v1/logs", "application/x-protobuf", body.clone()).await;
assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{why}");
let p = proxied(vec![stub(StatusCode::BAD_REQUEST, "no").await]);
let (status, why) = post(&p, "/v1/logs", "application/x-protobuf", body).await;
assert_eq!(status, StatusCode::BAD_REQUEST, "{why}");
}