use std::io::{Read, Write};
use std::path::PathBuf;
use std::process::{Command, Stdio};
use std::sync::{Arc, Mutex};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use prost::Message;
use mira_proto::collector::logs::v1::ExportLogsServiceRequest;
use mira_proto::collector::metrics::v1::ExportMetricsServiceRequest;
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::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};
use mira_proto::resource::v1::Resource;
use mira_proto::trace::v1::{ResourceSpans, ScopeSpans, Span};
const MIRA: &str = env!("CARGO_BIN_EXE_mira");
fn mira(args: &[&str]) -> (Option<i32>, String, String) {
let out = Command::new(MIRA).args(args).output().unwrap();
(
out.status.code(),
String::from_utf8_lossy(&out.stdout).into_owned(),
String::from_utf8_lossy(&out.stderr).into_owned(),
)
}
#[test]
fn the_command_line_answers_before_it_starts_a_server() {
let (code, out, err) = mira(&["--version"]);
assert_eq!((code, err.as_str()), (Some(0), ""));
assert_eq!(out.trim(), format!("mira {}", env!("CARGO_PKG_VERSION")));
for flag in ["-h", "--help"] {
let (code, out, _) = mira(&[flag]);
assert_eq!(code, Some(0), "{flag}");
assert!(out.contains("mira mira"), "{flag}: {out:?}");
}
let (code, out, err) = mira(&["--nope"]);
assert_eq!(code, Some(1));
assert!(err.starts_with("mira: unknown flag --nope"), "{err:?}");
assert_eq!(out, "");
let (code, out, _) = mira(&["mira", "--help"]);
assert_eq!(code, Some(0));
assert!(out.contains("mira tui"), "{out:?}");
let (code, _, err) = mira(&["tui", "--data-dir", "/nonexistent"]);
assert_eq!(code, Some(1));
assert!(err.contains("needs stdin and stdout on a tty"), "{err:?}");
let (code, _, err) = mira(&["mira", "--nope"]);
assert_eq!(code, Some(1));
assert!(err.contains("unknown flag --nope"), "{err:?}");
for cmd in ["proxy", "offload"] {
for flag in ["-h", "--help"] {
let (code, out, err) = mira(&[cmd, flag]);
assert_eq!((code, err.as_str()), (Some(0), ""), "mira {cmd} {flag}");
assert!(out.contains("mira mira"), "mira {cmd} {flag}: {out:?}");
}
let (code, out, _) = mira(&[cmd, "-V"]);
assert_eq!(code, Some(0), "mira {cmd} -V");
assert_eq!(out.trim(), format!("mira {}", env!("CARGO_PKG_VERSION")));
}
let (code, out, _) = mira(&["update", "--help"]);
assert_eq!(code, Some(0));
assert!(out.contains("mira update ["), "{out:?}");
let (code, out, _) = mira(&["update", "--version", "v0.1.0", "--dry-run"]);
assert_eq!(code, Some(0));
assert!(out.contains("v0.1.0"), "{out:?}");
}
#[test]
fn a_sigterm_stops_the_server_with_the_data_on_disk() {
let dir = std::env::temp_dir().join(format!("mira-cli-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let (mut child, log) = spawn_logged(&[
"--grpc",
"127.0.0.1:0",
"--http",
"127.0.0.1:0",
"--data-dir",
dir.to_str().unwrap(),
]);
let logged = |marker: &str| logged(&log, marker);
let up = logged("mira listening");
let port = up.then(|| port_of(&log)).flatten();
let posted = port.map(|p| post(p, "/v1/logs", PROTOBUF, &one_log().encode_to_vec()));
let signalled = unsafe { libc::kill(child.id() as i32, libc::SIGTERM) } == 0;
let stopped = signalled && logged("stopped");
if !stopped {
let _ = child.kill();
}
let status = child.wait().unwrap();
let tail = log.lock().unwrap().clone();
assert!(up, "never came up:\n{tail}");
assert!(
posted
.as_deref()
.is_some_and(|r| r.starts_with("HTTP/1.1 200")),
"export rejected: {posted:?}\n{tail}"
);
assert!(stopped, "did not drain after SIGTERM:\n{tail}");
assert!(status.success(), "exited {status}:\n{tail}");
let block = first(&first(&dir.join("logs")));
assert!(block.join("logs.arrow").is_file(), "{block:?}");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn the_proxy_subcommand_serves_the_node_behind_it_and_the_node_refuses_to_be_one() {
let (code, out, err) = mira(&["--replica", "http://127.0.0.1:4318"]);
assert_eq!((code, out.as_str()), (Some(1), ""));
assert!(err.contains("is read by `mira proxy`"), "{err:?}");
let (code, _, err) = mira(&["proxy", "--http", "127.0.0.1:0"]);
assert_eq!(code, Some(1));
assert!(err.contains("--replica"), "{err:?}");
let (code, _, err) = mira(&["proxy", "--http", "not-an-address"]);
assert_eq!(code, Some(1));
assert!(err.contains("--http: invalid socket address"), "{err:?}");
let dir = std::env::temp_dir().join(format!("mira-cli-proxy-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let (mut node, nlog) = spawn_logged(&[
"--grpc",
"127.0.0.1:0",
"--http",
"127.0.0.1:0",
"--data-dir",
dir.to_str().unwrap(),
]);
let replica = logged(&nlog, "mira listening")
.then(|| port_of(&nlog))
.flatten()
.map(|p| format!("http://127.0.0.1:{p}"));
let (mut px, plog) = spawn_logged(&[
"proxy",
"--http",
"127.0.0.1:0",
"--replica",
replica.as_deref().unwrap_or("http://127.0.0.1:1"),
]);
let up = logged(&plog, "mira proxy listening");
let port = up.then(|| port_of(&plog)).flatten();
let exported = port.map(|p| post(p, "/v1/logs", PROTOBUF, &one_log().encode_to_vec()));
let read = port.map(|p| {
post(
p,
"/api/v1/query",
"application/json",
br#"{"signal":"logs","limit":5}"#,
)
});
let others: Vec<String> = port
.map(|p| {
vec![
post(p, "/v1/traces", PROTOBUF, &one_span().encode_to_vec()),
post(p, "/v1/metrics", PROTOBUF, &one_point().encode_to_vec()),
]
})
.unwrap_or_default();
for c in [&node, &px] {
unsafe { libc::kill(c.id() as i32, libc::SIGTERM) };
}
let (pstatus, nstatus) = (px.wait().unwrap(), node.wait().unwrap());
let (ntail, ptail) = (nlog.lock().unwrap().clone(), plog.lock().unwrap().clone());
assert!(replica.is_some(), "the replica never came up:\n{ntail}");
assert!(up, "the proxy never came up:\n{ptail}");
assert!(ptail.contains("replicas=http://127.0.0.1:"), "{ptail}");
assert!(
exported
.as_deref()
.is_some_and(|r| r.starts_with("HTTP/1.1 200")),
"export rejected: {exported:?}\n{ptail}"
);
assert!(
read.as_deref().is_some_and(|r| r.contains("\"hello\"")),
"the row did not come back through the proxy: {read:?}\n{ptail}"
);
assert_eq!(others.len(), 2, "{ptail}");
for r in &others {
assert!(
r.starts_with("HTTP/1.1 200"),
"export rejected: {r}\n{ptail}"
);
}
assert!(pstatus.success(), "the proxy exited {pstatus}:\n{ptail}");
assert!(nstatus.success(), "the node exited {nstatus}:\n{ntail}");
let _ = std::fs::remove_dir_all(&dir);
}
fn spawn_logged(args: &[&str]) -> (std::process::Child, Arc<Mutex<String>>) {
let mut child = Command::new(MIRA)
.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap();
let log = Arc::new(Mutex::new(String::new()));
for stream in [
Box::new(child.stdout.take().unwrap()) as Box<dyn Read + Send>,
Box::new(child.stderr.take().unwrap()),
] {
let sink = log.clone();
std::thread::spawn(move || {
let mut stream = stream;
let mut buf = [0u8; 4096];
while let Ok(n) = stream.read(&mut buf) {
if n == 0 {
return;
}
sink.lock()
.unwrap()
.push_str(&String::from_utf8_lossy(&buf[..n]));
}
});
}
(child, log)
}
fn logged(log: &Arc<Mutex<String>>, marker: &str) -> bool {
(0..600).any(|_| {
if log.lock().unwrap().contains(marker) {
return true;
}
std::thread::sleep(Duration::from_millis(50));
false
})
}
fn port_of(log: &Arc<Mutex<String>>) -> Option<u16> {
let tail = log
.lock()
.unwrap()
.split("http=127.0.0.1:")
.nth(1)?
.to_owned();
tail.chars()
.take_while(char::is_ascii_digit)
.collect::<String>()
.parse()
.ok()
}
fn first(dir: &PathBuf) -> PathBuf {
let mut entries: Vec<_> = std::fs::read_dir(dir)
.unwrap_or_else(|e| panic!("{dir:?}: {e}"))
.map(|e| e.unwrap().path())
.collect();
entries.sort();
entries
.into_iter()
.next()
.unwrap_or_else(|| panic!("{dir:?} is empty"))
}
const PROTOBUF: &str = "application/x-protobuf";
fn post(port: u16, path: &str, content_type: &str, body: &[u8]) -> String {
let mut s = std::net::TcpStream::connect(("127.0.0.1", port)).unwrap();
write!(
s,
"POST {path} HTTP/1.1\r\nHost: localhost\r\nContent-Type: {content_type}\r\n\
Content-Length: {}\r\nConnection: close\r\n\r\n",
body.len()
)
.unwrap();
s.write_all(body).unwrap();
let mut res = String::new();
s.set_read_timeout(Some(Duration::from_secs(30))).unwrap();
s.read_to_string(&mut res).unwrap();
res
}
fn one_log() -> ExportLogsServiceRequest {
let now = nanos();
ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
resource: Some(named("mira.cli")),
scope_logs: vec![ScopeLogs {
scope: Some(InstrumentationScope {
name: "mira.cli".into(),
..Default::default()
}),
log_records: vec![LogRecord {
time_unix_nano: now,
severity_number: 9,
severity_text: "INFO".into(),
body: Some(AnyValue {
value: Some(any_value::Value::StringValue("hello".into())),
}),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
}
}
fn one_span() -> ExportTraceServiceRequest {
let now = nanos();
ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: Some(named("mira.cli")),
scope_spans: vec![ScopeSpans {
spans: vec![Span {
trace_id: vec![0xab; 16].into(),
span_id: vec![0xcd; 8].into(),
name: "cli".into(),
start_time_unix_nano: now,
end_time_unix_nano: now + 1,
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
}
}
fn one_point() -> ExportMetricsServiceRequest {
ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
resource: Some(named("mira.cli")),
scope_metrics: vec![ScopeMetrics {
metrics: vec![Metric {
name: "cli.requests".into(),
data: Some(Data::Gauge(Gauge {
data_points: vec![NumberDataPoint {
time_unix_nano: nanos(),
value: Some(NumValue::AsInt(1)),
..Default::default()
}],
})),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
}
}
fn nanos() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos() as u64
}
fn named(service: &str) -> Resource {
Resource {
attributes: vec![KeyValue {
key: "service.name".into(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue(service.into())),
}),
}],
..Default::default()
}
}