use std::collections::HashMap;
use std::fs;
use chrono::Utc;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use crate::constants::LOGS;
use crate::error::{Error, Result};
use crate::fs_text::write_text_atomic;
use crate::home::UnifierHome;
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct LogEvent {
pub message: String,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub fields: HashMap<String, String>,
pub timestamp: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct Span {
pub id: Uuid,
pub name: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub parent_id: Option<Uuid>,
pub started_at: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub ended_at: Option<String>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub fields: HashMap<String, String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub events: Vec<LogEvent>,
}
impl Span {
fn new(name: &str, parent_id: Option<Uuid>) -> Self {
Self {
id: Uuid::new_v4(),
name: name.to_string(),
parent_id,
started_at: Utc::now().to_rfc3339(),
ended_at: None,
fields: HashMap::new(),
events: Vec::new(),
}
}
}
fn logs_dir(home: &UnifierHome) -> std::path::PathBuf {
home.path().join(LOGS)
}
fn span_path(home: &UnifierHome, id: Uuid) -> std::path::PathBuf {
logs_dir(home).join(format!("{}.json", id.hyphenated()))
}
fn save_span(home: &UnifierHome, span: &Span) -> Result<()> {
let dir = logs_dir(home);
fs::create_dir_all(&dir)?;
let json = serde_json::to_string_pretty(span)?;
write_text_atomic(&span_path(home, span.id), &json)?;
Ok(())
}
fn load_span(home: &UnifierHome, id: Uuid) -> Result<Span> {
let path = span_path(home, id);
if !path.is_file() {
return Err(Error::msg(format!("span not found: {id}")));
}
let text = fs::read_to_string(&path)?;
Ok(serde_json::from_str(&text)?)
}
pub fn load_all_spans(home: &UnifierHome) -> Result<Vec<Span>> {
let dir = logs_dir(home);
if !dir.is_dir() {
return Ok(vec![]);
}
let mut spans = Vec::new();
for entry in fs::read_dir(&dir)? {
let entry = entry?;
if !entry.file_type()?.is_file() {
continue;
}
let name = entry.file_name().to_string_lossy().into_owned();
let Some(stem) = name.strip_suffix(".json") else {
continue;
};
let Ok(_id) = Uuid::parse_str(stem) else {
continue;
};
let text = fs::read_to_string(entry.path())?;
if let Ok(span) = serde_json::from_str::<Span>(&text) {
spans.push(span);
}
}
spans.sort_by(|a, b| a.started_at.cmp(&b.started_at));
Ok(spans)
}
pub fn span_start(
home: &UnifierHome,
name: &str,
parent_id: Option<Uuid>,
fields: HashMap<String, String>,
) -> Result<Uuid> {
if let Some(pid) = parent_id {
load_span(home, pid)?;
}
let mut span = Span::new(name, parent_id);
span.fields = fields;
let id = span.id;
save_span(home, &span)?;
Ok(id)
}
pub fn span_end(home: &UnifierHome, id: Uuid) -> Result<()> {
let mut span = load_span(home, id)?;
if span.ended_at.is_some() {
return Err(Error::msg(format!("span {id} already ended")));
}
span.ended_at = Some(Utc::now().to_rfc3339());
save_span(home, &span)
}
pub fn log_event(
home: &UnifierHome,
span_id: Uuid,
message: &str,
fields: HashMap<String, String>,
) -> Result<()> {
let mut span = load_span(home, span_id)?;
span.events.push(LogEvent {
message: message.to_string(),
fields,
timestamp: Utc::now().to_rfc3339(),
});
save_span(home, &span)
}
pub fn span_set_fields(
home: &UnifierHome,
id: Uuid,
fields: HashMap<String, String>,
) -> Result<()> {
let mut span = load_span(home, id)?;
span.fields.extend(fields);
save_span(home, &span)
}
pub fn render_tree(spans: &[Span], root_id: Option<Uuid>) -> String {
let mut out = String::new();
render_children(spans, root_id, 0, &mut out);
out
}
fn render_children(spans: &[Span], parent: Option<Uuid>, depth: usize, out: &mut String) {
let indent = " ".repeat(depth);
let children: Vec<&Span> = spans.iter().filter(|s| s.parent_id == parent).collect();
for span in children {
let status = match &span.ended_at {
Some(t) => format!("ended {}", short_ts(t)),
None => "open".to_string(),
};
out.push_str(&format!(
"{}▶ {} [started {}, {}]\n",
indent,
span.name,
short_ts(&span.started_at),
status
));
for (k, v) in &span.fields {
out.push_str(&format!("{} · {}={}\n", indent, k, v));
}
for ev in &span.events {
let field_str: Vec<String> =
ev.fields.iter().map(|(k, v)| format!("{k}={v}")).collect();
if field_str.is_empty() {
out.push_str(&format!(
"{} · {} {}\n",
indent,
short_ts(&ev.timestamp),
ev.message
));
} else {
out.push_str(&format!(
"{} · {} {} ({})\n",
indent,
short_ts(&ev.timestamp),
ev.message,
field_str.join(", ")
));
}
}
render_children(spans, Some(span.id), depth + 1, out);
}
}
fn short_ts(ts: &str) -> &str {
if let Some(t_pos) = ts.find('T') {
let after_t = &ts[t_pos + 1..];
let end = after_t
.find(['+', '-', 'Z'])
.map(|i| t_pos + 1 + i)
.unwrap_or(ts.len());
&ts[t_pos + 1..end]
} else {
ts
}
}
pub fn parse_fields(pairs: &[String]) -> Result<HashMap<String, String>> {
let mut map = HashMap::new();
for pair in pairs {
let (k, v) = pair
.split_once('=')
.ok_or_else(|| Error::msg(format!("field must be key=value, got: {pair}")))?;
map.insert(k.to_string(), v.to_string());
}
Ok(map)
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::tempdir;
fn home(dir: &std::path::Path) -> UnifierHome {
UnifierHome::resolve(Some(dir.to_path_buf()), None).unwrap()
}
#[test]
fn span_start_end_roundtrip() {
let tmp = tempdir().unwrap();
let h = home(tmp.path());
let id = span_start(&h, "test", None, HashMap::new()).unwrap();
let span = load_span(&h, id).unwrap();
assert_eq!(span.name, "test");
assert!(span.ended_at.is_none());
span_end(&h, id).unwrap();
let span = load_span(&h, id).unwrap();
assert!(span.ended_at.is_some());
}
#[test]
fn nested_spans() {
let tmp = tempdir().unwrap();
let h = home(tmp.path());
let root = span_start(&h, "root", None, HashMap::new()).unwrap();
let child = span_start(&h, "child", Some(root), HashMap::new()).unwrap();
let spans = load_all_spans(&h).unwrap();
assert_eq!(spans.len(), 2);
let c = spans.iter().find(|s| s.id == child).unwrap();
assert_eq!(c.parent_id, Some(root));
}
#[test]
fn log_event_appended() {
let tmp = tempdir().unwrap();
let h = home(tmp.path());
let id = span_start(&h, "work", None, HashMap::new()).unwrap();
log_event(&h, id, "fetched 10 rows", HashMap::new()).unwrap();
let span = load_span(&h, id).unwrap();
assert_eq!(span.events.len(), 1);
assert_eq!(span.events[0].message, "fetched 10 rows");
}
#[test]
fn render_tree_smoke() {
let tmp = tempdir().unwrap();
let h = home(tmp.path());
let root = span_start(&h, "pipeline", None, HashMap::new()).unwrap();
let _child = span_start(&h, "step-1", Some(root), HashMap::new()).unwrap();
span_end(&h, root).unwrap();
let spans = load_all_spans(&h).unwrap();
let tree = render_tree(&spans, None);
assert!(tree.contains("pipeline"));
assert!(tree.contains("step-1"));
}
#[test]
fn double_end_is_error() {
let tmp = tempdir().unwrap();
let h = home(tmp.path());
let id = span_start(&h, "x", None, HashMap::new()).unwrap();
span_end(&h, id).unwrap();
assert!(span_end(&h, id).is_err());
}
#[test]
fn parse_fields_ok() {
let pairs = vec!["k=v".to_string(), "foo=bar=baz".to_string()];
let m = parse_fields(&pairs).unwrap();
assert_eq!(m["k"], "v");
assert_eq!(m["foo"], "bar=baz");
}
}