use std::io::{Read, Write};
use std::net::{TcpStream, ToSocketAddrs};
use std::path::{Path, PathBuf};
use std::time::Duration;
use yaml_rust2::Yaml;
use crate::api;
pub enum Source {
Local(PathBuf),
Remote(String),
}
pub const QUERY: &str = "/api/v1/query";
pub const SERIES: &str = "/api/v1/metrics/query";
pub const NAMES: &str = "/api/v1/metrics/names";
pub const CORRELATE: &str = "/api/v1/correlate";
pub const MAP: &str = "/api/v1/map";
pub const ENTITIES: &str = "/api/v1/entities";
pub const STATS: &str = "/api/v1/stats";
pub const ALERTS: &str = "/api/v1/alerts";
impl Source {
pub fn label(&self) -> String {
match self {
Source::Local(p) => format!("local {}", p.display()),
Source::Remote(a) => format!("http {a}"),
}
}
pub fn post(&self, route: &str, body: &str) -> Result<Yaml, String> {
let text = match self {
Source::Local(dir) => local(dir, route, body)?,
Source::Remote(addr) => http(addr, "POST", route, Some(body))?,
};
Self::envelope(&text)
}
pub fn get(&self, route: &str) -> Result<Yaml, String> {
let text = match self {
Source::Local(_) => {
return Err(format!(
"{} reports a running node; this is a directory, so open it with --addr",
route.rsplit('/').next().unwrap_or(route)
));
}
Source::Remote(addr) => http(addr, "GET", route, None)?,
};
Self::envelope(&text)
}
fn envelope(text: &str) -> Result<Yaml, String> {
let doc = api::parse(text)?;
match doc["error"].as_str() {
Some(e) => Err(e.to_owned()),
None => Ok(doc),
}
}
}
fn local(dir: &Path, route: &str, body: &str) -> Result<String, String> {
let now = api::now_nanos();
let t = std::time::Instant::now();
let run = |field, r: mira_core::error::Result<mira_core::query::Results>| match r {
Ok(r) => Ok(api::envelope(field, &r, t.elapsed())),
Err(e) => Err(e.to_string()),
};
match route {
QUERY => run(
"rows",
mira_core::query::search(dir, &api::parse_search(body, now)?),
),
SERIES => run(
"series",
mira_core::series::series(dir, &api::parse_series(body, now)?),
),
NAMES => {
let (from, to) = api::window(body, now)?;
run("names", mira_core::series::names(dir, from, to))
}
CORRELATE => {
let (q, ops) = api::parse_correlate(body, now)?;
run("frame", api::correlate(dir, &q, &ops, &[], &[]))
}
MAP => {
let (from, to, max) = api::map_doc(&api::parse(body)?, now)?;
run("map", mira_core::frame::map(dir, from, to, max, &[]))
}
ENTITIES => {
let (from, to) = api::window(body, now)?;
run("entities", mira_core::frame::entities(dir, from, to, &[]))
}
other => Err(format!("no such route {other}")),
}
}
const CONNECT_TIMEOUT: Duration = Duration::from_secs(5);
fn http(addr: &str, method: &str, path: &str, body: Option<&str>) -> Result<String, String> {
let io = || -> Result<Vec<u8>, String> {
let mut last = format!("{addr}: resolved to no address");
let mut sock = None;
for sa in addr.to_socket_addrs().map_err(|e| format!("{addr}: {e}"))? {
match TcpStream::connect_timeout(&sa, CONNECT_TIMEOUT) {
Ok(c) => {
sock = Some(c);
break;
}
Err(e) => last = format!("{addr}: {e}"),
}
}
let mut s = sock.ok_or(last)?;
s.set_read_timeout(Some(Duration::from_secs(60)))
.map_err(|e| format!("{addr}: {e}"))?;
s.set_write_timeout(Some(Duration::from_secs(60)))
.map_err(|e| format!("{addr}: {e}"))?;
let body = body.unwrap_or("");
write!(
s,
"{method} {path} HTTP/1.1\r\nHost: {addr}\r\nContent-Type: application/json\r\n\
Content-Length: {}\r\nConnection: close\r\n\r\n{body}",
body.len()
)
.and_then(|()| s.flush())
.map_err(|e| format!("{addr}: sending the request: {e}"))?;
let mut buf = Vec::new();
s.read_to_end(&mut buf)
.map_err(|e| format!("{addr}: reading the response: {e}"))?;
Ok(buf)
};
let raw = io()?;
let split = raw
.windows(4)
.position(|w| w == b"\r\n\r\n")
.ok_or_else(|| format!("{addr}: response has no header terminator"))?;
let head = String::from_utf8_lossy(&raw[..split]);
let text = String::from_utf8_lossy(&raw[split + 4..]).into_owned();
let status = head
.lines()
.next()
.and_then(|l| l.split_whitespace().nth(1))
.unwrap_or("?");
match status {
"200" => Ok(text),
_ if text.trim_start().starts_with('{') => Ok(text),
_ => Err(format!("{addr} returned HTTP {status}: {}", text.trim())),
}
}
pub fn parse_addr(s: &str) -> Result<String, String> {
let s = s
.trim()
.trim_start_matches("http://")
.trim_end_matches('/')
.trim();
if s.starts_with("https://") {
return Err("--addr: Mira serves plain HTTP; there is no TLS client here".into());
}
if s.is_empty() {
return Err("--addr needs a host".into());
}
Ok(
match s
.rsplit(':')
.next()
.is_some_and(|p| p.parse::<u16>().is_ok())
{
true => s.to_owned(),
false => format!("{s}:4318"),
},
)
}
#[cfg(test)]
pub fn serve(replies: Vec<String>) -> String {
let l = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let addr = l.local_addr().unwrap().to_string();
std::thread::spawn(move || {
for reply in replies {
let Ok((mut s, _)) = l.accept() else { return };
let mut head = Vec::new();
let mut byte = [0u8; 1];
while std::io::Read::read(&mut s, &mut byte).unwrap_or(0) == 1 {
head.push(byte[0]);
if head.ends_with(b"\r\n\r\n") {
break;
}
}
let len: usize = String::from_utf8_lossy(&head)
.lines()
.find_map(|l| l.strip_prefix("Content-Length: ")?.trim().parse().ok())
.unwrap_or(0);
let mut body = vec![0u8; len];
let _ = std::io::Read::read_exact(&mut s, &mut body);
let _ = s.write_all(reply.as_bytes());
}
});
addr
}
#[cfg(test)]
pub fn ok(body: &str) -> String {
format!("HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n\r\n{body}")
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn addresses_normalise_to_host_port() {
assert_eq!(parse_addr("localhost").unwrap(), "localhost:4318");
assert_eq!(parse_addr("http://mira:4318/").unwrap(), "mira:4318");
assert_eq!(parse_addr("10.0.0.4:9000").unwrap(), "10.0.0.4:9000");
assert_eq!(parse_addr("[::1]:4318").unwrap(), "[::1]:4318");
assert_eq!(parse_addr("[::1]").unwrap(), "[::1]:4318");
assert!(parse_addr("https://mira").is_err());
assert!(parse_addr(" ").is_err());
}
#[test]
fn a_response_envelope_parses_with_the_kyaml_loader() {
let text = r#"{"rows":[{"body":"line1\nline2 \"q\" \u0001","severity_number":17,
"ratio":0.5,"ok":true,"gone":null,
"attributes":{"service.name":"checkout"}}],
"stats":{"blocks_total":69,"blocks_scanned":1,
"rows_scanned":100,"rows_matched":2}}"#;
let d = api::parse(text).unwrap();
let row = &d["rows"][0];
assert_eq!(row["body"].as_str().unwrap(), "line1\nline2 \"q\" \u{1}");
assert_eq!(row["severity_number"].as_i64().unwrap(), 17);
assert_eq!(row["ratio"].as_f64().unwrap(), 0.5);
assert!(row["ok"].as_bool().unwrap());
assert!(row["gone"].is_null());
assert_eq!(
row["attributes"]["service.name"].as_str().unwrap(),
"checkout"
);
assert_eq!(d["stats"]["blocks_scanned"].as_i64().unwrap(), 1);
}
#[test]
fn a_local_source_answers_out_of_a_directory_with_no_server() {
let dir = std::env::temp_dir().join(format!("mira-src-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let mut b = mira_core::logs::LogsBuilder::new();
b.append_request(&crate::e2e::logs_export("checkout", 1_000, 6))
.unwrap();
let sealed = b.finish().unwrap();
mira_core::block::publish(&dir, "logs", mira_core::block::node_id("a"), 0, 0, &sealed)
.unwrap();
let src = Source::Local(dir.clone());
assert!(src.label().starts_with("local "));
let d = src
.post(QUERY, r#"{"signal":"logs","from":0,"to":100000,"limit":2}"#)
.unwrap();
assert_eq!(d["rows"].as_vec().unwrap().len(), 2);
assert_eq!(d["stats"]["rows_matched"].as_i64().unwrap(), 6);
for (route, field, body) in [
(SERIES, "series", r#"{"name":"anything"}"#),
(NAMES, "names", "{}"),
] {
let d = src.post(route, body).unwrap();
assert!(d[field].as_vec().unwrap().is_empty(), "{route}");
assert_eq!(d["stats"]["blocks_total"].as_i64().unwrap(), 0);
}
let d = src.post(ENTITIES, r#"{"from":0,"to":100000}"#).unwrap();
assert_eq!(d["entities"][0]["name"].as_str().unwrap(), "checkout");
let d = src
.post(
CORRELATE,
r#"{"signal":"logs","from":0,"to":100000,"expand":["traces","peers"]}"#,
)
.unwrap();
assert_eq!(
d["frame"]["entities"][0]["name"].as_str().unwrap(),
"checkout"
);
assert!(!d["frame"]["truncated"].as_bool().unwrap());
let d = src.post(MAP, r#"{"from":0,"to":100000}"#).unwrap();
assert!(d["map"]["nodes"].as_vec().unwrap().is_empty());
assert_eq!(d["map"]["unresolved"].as_i64().unwrap(), 0);
for (route, body, want) in [
(CORRELATE, r#"{"expand":["sideways"]}"#, "sideways"),
(MAP, r#"{"limit":5}"#, "unknown query key"),
(ENTITIES, r#"{"to":"soon"}"#, "soon"),
] {
let e = src.post(route, body).unwrap_err();
assert!(e.contains(want), "{route}: {e}");
}
let e = src.post(QUERY, r#"{"signal":"nope"}"#).unwrap_err();
assert!(e.contains("nope"), "{e}");
assert!(
src.post("/api/v1/nope", "{}")
.unwrap_err()
.contains("route")
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_remote_source_reports_what_the_server_actually_said() {
let ok = "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n\r\n{\"rows\":[],\
\"stats\":{\"blocks_total\":0,\"blocks_scanned\":0,\"rows_scanned\":0,\
\"rows_matched\":0}}";
let bad = "HTTP/1.1 400 Bad Request\r\n\r\n{\"error\":\"unknown signal \\\"nope\\\"\"}";
let plain = "HTTP/1.1 502 Bad Gateway\r\n\r\nupstream is down";
let truncated = "HTTP/1.1 200 OK\r\nContent-Type: application/json";
let src = Source::Remote(serve(
[ok, bad, plain, truncated].map(str::to_owned).to_vec(),
));
assert!(src.label().starts_with("http "));
let d = src.post(QUERY, "{}").unwrap();
assert!(d["rows"].as_vec().unwrap().is_empty());
assert_eq!(
src.post(QUERY, "{}").unwrap_err(),
r#"unknown signal "nope""#
);
let e = src.post(QUERY, "{}").unwrap_err();
assert!(e.contains("502") && e.contains("upstream is down"), "{e}");
assert!(
src.post(QUERY, "{}").unwrap_err().contains("terminator"),
"a reply with no blank line is not an empty answer"
);
let dead = Source::Remote("127.0.0.1:1".into());
let e = dead.post(QUERY, "{}").unwrap_err();
assert!(e.starts_with("127.0.0.1:1: "), "{e}");
}
#[test]
fn a_request_that_cannot_be_written_reports_the_write_and_not_the_reply() {
let l = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let addr = l.local_addr().unwrap().to_string();
std::thread::spawn(move || {
for s in l.incoming() {
match s {
Ok(s) => drop(s.shutdown(std::net::Shutdown::Both)),
Err(_) => return,
}
}
});
let mut e = String::new();
for mib in [1usize, 8, 64] {
let body = format!(r#"{{"signal":"logs","q":"{}"}}"#, "x".repeat(mib << 20));
e = Source::Remote(addr.clone()).post(QUERY, &body).unwrap_err();
if e.contains("sending the request") {
break;
}
}
assert!(
e.starts_with(&format!("{addr}: sending the request: ")),
"the write failed and the message says which half died: {e}"
);
}
#[test]
fn a_name_resolving_to_several_addresses_is_tried_down_the_list() {
let ok = "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n\r\n{\"rows\":[],\
\"stats\":{\"blocks_total\":0,\"blocks_scanned\":0,\"rows_scanned\":0,\
\"rows_matched\":0}}";
let bound = serve(vec![ok.into()]);
let port = bound.rsplit(':').next().unwrap();
let d = Source::Remote(format!("localhost:{port}"))
.post(QUERY, "{}")
.unwrap();
assert!(d["rows"].as_vec().unwrap().is_empty());
}
}