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::common::v1::{AnyValue, InstrumentationScope, KeyValue, any_value};
use mira_proto::logs::v1::{LogRecord, ResourceLogs, ScopeLogs};
use mira_proto::resource::v1::Resource;
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:?}");
}
#[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 = Command::new(MIRA)
.args([
"--grpc",
"127.0.0.1:0",
"--http",
"127.0.0.1:0",
"--data-dir",
])
.arg(&dir)
.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]));
}
});
}
let logged = |marker: &str| {
(0..600).any(|_| {
if log.lock().unwrap().contains(marker) {
return true;
}
std::thread::sleep(Duration::from_millis(50));
false
})
};
let up = logged("mira listening");
let port = up.then(|| log.lock().unwrap().clone()).and_then(|l| {
let tail = l.split("http=127.0.0.1:").nth(1)?.to_owned();
tail.chars()
.take_while(char::is_ascii_digit)
.collect::<String>()
.parse::<u16>()
.ok()
});
let posted = port.map(|p| post(p, "/v1/logs", &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);
}
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"))
}
fn post(port: u16, path: &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: application/x-protobuf\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 = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos() as u64;
ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
resource: Some(Resource {
attributes: vec![KeyValue {
key: "service.name".into(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("mira.cli".into())),
}),
}],
..Default::default()
}),
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()
}],
}
}