use crate::http::{self, Request, Response};
use crate::paths::Paths;
use crate::report::Report;
use anyhow::{Context, Result};
use std::io::{BufReader, Write};
use std::net::{TcpListener, TcpStream};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
const INDEX_HTML: &str = include_str!("../../assets/ui/index.html");
const APP_CSS: &str = include_str!("../../assets/ui/app.css");
const APP_JS: &str = include_str!("../../assets/ui/app.js");
const CSP_APP: &str = "default-src 'none'; script-src 'self'; style-src 'self'; \
img-src 'self' data:; connect-src 'self'; base-uri 'none'; \
form-action 'none'; frame-ancestors 'none'";
const CSP_DATA: &str = "default-src 'none'; base-uri 'none'; form-action 'none'; \
frame-ancestors 'none'";
const MAX_EVIDENCE_BYTES: u64 = 2 * 1024 * 1024;
const MAX_THREADS: usize = 64;
const READ_TIMEOUT: Duration = Duration::from_secs(5);
const WRITE_TIMEOUT: Duration = Duration::from_secs(30);
const VERSION_TTL: Duration = Duration::from_millis(250);
struct Server {
paths: Paths,
port: u16,
version: Mutex<Option<(Instant, String)>>,
}
pub fn ignore_sigpipe() {
#[cfg(unix)]
unsafe {
libc::signal(libc::SIGPIPE, libc::SIG_IGN);
}
}
pub fn serve(paths: Paths, listener: TcpListener) -> Result<()> {
let port = listener.local_addr()?.port();
let server = Arc::new(Server { paths, port, version: Mutex::new(None) });
let live = Arc::new(AtomicUsize::new(0));
for stream in listener.incoming() {
let Ok(stream) = stream else { continue };
if live.load(Ordering::Relaxed) >= MAX_THREADS {
shed(&stream);
continue;
}
let server = Arc::clone(&server);
let live = Arc::clone(&live);
live.fetch_add(1, Ordering::Relaxed);
std::thread::spawn(move || {
handle(&server, stream);
live.fetch_sub(1, Ordering::Relaxed);
});
}
Ok(())
}
fn shed(stream: &TcpStream) {
let _ = stream.set_read_timeout(Some(Duration::from_millis(250)));
let _ = stream.set_write_timeout(Some(Duration::from_millis(250)));
let mut reader = BufReader::new(stream);
let _ = http::parse_request(&mut reader);
let mut out = stream;
let _ = respond(&mut out, &error_response(503, "too many connections"), false);
let _ = stream.shutdown(std::net::Shutdown::Both);
}
fn handle(server: &Server, stream: TcpStream) {
let _ = stream.set_read_timeout(Some(READ_TIMEOUT));
let _ = stream.set_write_timeout(Some(WRITE_TIMEOUT));
let started = Instant::now();
let mut reader = BufReader::new(&stream);
let (response, method, path, head_only) = match http::parse_request(&mut reader) {
Ok(req) => {
let method = req.method.clone();
let path = req.path.clone();
let head_only = req.is_head();
(route(server, &req), method, path, head_only)
}
Err(e) => {
let status = e.status();
if matches!(e, http::Error::Io(_)) {
return;
}
(error_response(status, "bad request"), "-".into(), "-".into(), false)
}
};
let mut out = &stream;
let bytes = respond(&mut out, &response, head_only).unwrap_or(0);
eprintln!(
"{method} {path} {} {bytes}b {}ms",
response.status,
started.elapsed().as_millis()
);
}
fn respond(w: &mut impl Write, r: &Response, head_only: bool) -> std::io::Result<usize> {
http::write_response(w, r, head_only)
}
fn error_response(status: u16, message: &str) -> Response {
Response::new(
status,
"application/json; charset=utf-8",
serde_json::json!({ "error": message }).to_string().into_bytes(),
)
.with_header("Content-Security-Policy", CSP_DATA)
}
fn route(server: &Server, req: &Request) -> Response {
if let Some(rejection) = check_origin(server, req) {
return rejection;
}
match req.path.as_str() {
"/" | "/index.html" => {
Response::html(INDEX_HTML).with_header("Content-Security-Policy", CSP_APP)
}
"/app.css" => Response::new(200, "text/css; charset=utf-8", APP_CSS.as_bytes().to_vec())
.with_header("Content-Security-Policy", CSP_APP),
"/app.js" => Response::new(
200,
"text/javascript; charset=utf-8",
APP_JS.as_bytes().to_vec(),
)
.with_header("Content-Security-Policy", CSP_APP),
"/api/version" => json_route(|| Ok(serde_json::json!({ "version": version(server)? }))),
"/api/overview" => json_route(|| {
let r = Report::build(&server.paths, None)?;
Ok(serde_json::to_value(r)?)
}),
"/api/insights" => json_route(|| {
let i = crate::report::insights::Insights::build(&server.paths)?;
Ok(serde_json::to_value(i)?)
}),
p if p.starts_with("/api/spec/") => {
let slug = &p["/api/spec/".len()..];
if !is_safe_segment(slug) {
return error_response(404, "no such spec");
}
json_route(|| {
let r = Report::build(&server.paths, Some(slug))?;
Ok(serde_json::to_value(r)?)
})
}
p if p.starts_with("/api/run/") => run_route(server, &p["/api/run/".len()..]),
_ => error_response(404, "no such route"),
}
}
fn json_route(build: impl FnOnce() -> Result<serde_json::Value>) -> Response {
match build().and_then(|v| Ok(serde_json::to_string(&v)?)) {
Ok(body) => Response::json(body).with_header("Content-Security-Policy", CSP_DATA),
Err(e) => error_response(503, &format!("{e:#}")),
}
}
fn run_route(server: &Server, rest: &str) -> Response {
let (id, tail) = match rest.split_once('/') {
Some((id, tail)) => (id, Some(tail)),
None => (rest, None),
};
if !is_run_id(id) {
return error_response(404, "no such run");
}
match tail {
None => json_route(|| {
let run = crate::run::Run::load(&server.paths, id)?;
let scan = crate::trajectory::scan(&run.trajectory_path())
.unwrap_or(crate::trajectory::Scan { events: vec![], anomalies: vec![] });
Ok(serde_json::json!({
"run": run.meta,
"gates": run.gate_results().unwrap_or_default(),
"events": scan.events,
"anomalies": scan.anomalies,
"evidence": evidence_listing(&run),
}))
}),
Some(t) => match t.strip_prefix("evidence/") {
Some(name) => evidence_route(server, id, name),
None => error_response(404, "no such route"),
},
}
}
fn evidence_listing(run: &crate::run::Run) -> Vec<serde_json::Value> {
let dir = run.dir.join("evidence");
let Ok(entries) = std::fs::read_dir(&dir) else { return vec![] };
let mut out: Vec<_> = entries
.filter_map(|e| e.ok())
.filter(|e| e.file_type().map(|t| t.is_file()).unwrap_or(false))
.filter_map(|e| {
let name = e.file_name().to_str()?.to_string();
let bytes = e.metadata().ok()?.len();
Some(serde_json::json!({ "name": name, "bytes": bytes }))
})
.collect();
out.sort_by_key(|v| v["name"].as_str().unwrap_or_default().to_string());
out
}
fn evidence_route(server: &Server, id: &str, name: &str) -> Response {
if !is_safe_segment(name) {
return error_response(404, "no such evidence file");
}
let Ok(run) = crate::run::Run::load(&server.paths, id) else {
return error_response(404, "no such run");
};
let dir = run.dir.join("evidence");
let listed = std::fs::read_dir(&dir)
.ok()
.into_iter()
.flatten()
.filter_map(|e| e.ok())
.any(|e| e.file_name().to_str() == Some(name));
if !listed {
return error_response(404, "no such evidence file");
}
let (Ok(root), Ok(target)) = (dir.canonicalize(), dir.join(name).canonicalize()) else {
return error_response(404, "no such evidence file");
};
if !target.starts_with(&root) || !target.is_file() {
return error_response(404, "no such evidence file");
}
match read_tail(&target) {
Ok((body, total)) => {
let truncated = total > body.len() as u64;
let mut r = Response::plain(200, body)
.with_header("Content-Security-Policy", CSP_DATA)
.with_header("X-Keel-Total-Bytes", &total.to_string());
if truncated {
r = r.with_header("X-Keel-Truncated", "true");
}
r
}
Err(_) => error_response(404, "no such evidence file"),
}
}
fn read_tail(path: &std::path::Path) -> std::io::Result<(Vec<u8>, u64)> {
use std::io::{Read, Seek, SeekFrom};
let mut f = std::fs::File::open(path)?;
let total = f.metadata()?.len();
if total > MAX_EVIDENCE_BYTES {
f.seek(SeekFrom::Start(total - MAX_EVIDENCE_BYTES))?;
}
let mut buf = Vec::new();
f.take(MAX_EVIDENCE_BYTES).read_to_end(&mut buf)?;
Ok((buf, total))
}
fn check_origin(server: &Server, req: &Request) -> Option<Response> {
let expected = [
format!("127.0.0.1:{}", server.port),
format!("[::1]:{}", server.port),
];
match req.header("host") {
Some(h) if expected.iter().any(|e| e == h) => {}
_ => return Some(error_response(403, "unexpected Host")),
}
if let Some(origin) = req.header("origin") {
let ok = expected.iter().any(|e| origin == format!("http://{e}"));
if !ok {
return Some(error_response(403, "cross-origin request"));
}
}
None
}
fn is_safe_segment(s: &str) -> bool {
!s.is_empty()
&& s.len() <= 128
&& s != "."
&& s != ".."
&& !s.contains('/')
&& !s.contains('\\')
&& !s.chars().any(|c| c.is_control())
}
fn is_run_id(s: &str) -> bool {
let b = s.as_bytes();
b.len() == 14
&& b[..4].iter().all(u8::is_ascii_digit)
&& b[4] == b'-'
&& b[5..7].iter().all(u8::is_ascii_digit)
&& b[7] == b'-'
&& b[8..10].iter().all(u8::is_ascii_digit)
&& b[10] == b'-'
&& b[11..].iter().all(u8::is_ascii_hexdigit)
}
fn version(server: &Server) -> Result<String> {
let mut cache = server.version.lock().unwrap_or_else(|e| e.into_inner());
if let Some((at, ref token)) = *cache
&& at.elapsed() < VERSION_TTL
{
return Ok(token.clone());
}
let mut hasher = crate::hashing::SetHasher::new();
for root in [server.paths.specs(), server.paths.runs()] {
stamp(&root, 3, &mut hasher);
}
let token = crate::hashing::short(&hasher.finish()).to_string();
*cache = Some((Instant::now(), token.clone()));
Ok(token)
}
fn stamp(dir: &std::path::Path, depth: usize, hasher: &mut crate::hashing::SetHasher) {
if depth == 0 {
return;
}
let Ok(entries) = std::fs::read_dir(dir) else { return };
let mut rows: Vec<(String, String)> = Vec::new();
for e in entries.filter_map(|e| e.ok()) {
let path = e.path();
let Ok(meta) = e.metadata() else { continue };
if meta.is_dir() {
stamp(&path, depth - 1, hasher);
continue;
}
let modified = meta
.modified()
.ok()
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map(|d| d.as_nanos())
.unwrap_or(0);
rows.push((path.display().to_string(), format!("{modified}:{}", meta.len())));
}
rows.sort();
for (path, stamp) in rows {
hasher.add(&path, stamp.as_bytes());
}
}
pub fn bind(port: u16) -> Result<TcpListener> {
TcpListener::bind(("127.0.0.1", port)).with_context(|| {
format!("binding 127.0.0.1:{port} — is another `keel serve` already running?")
})
}