unifier-cli 0.5.0

Filesystem postbox for inter-process communication via a Unix tree
Documentation
//! Span-tree logging for Jan cron scripts and agent pipelines.
//!
//! Spans live under `logs/` in the store root.  Each span is a JSON file
//! `logs/<id>.json`; log events are embedded in the span they belong to.
//! Reading back a span tree walks the directory and reconstructs parent→child
//! relationships in memory.
//!
//! Layout:
//! ```text
//! logs/
//!   <uuid>.json    # one file per span
//! ```
//!
//! Span JSON schema:
//! ```json
//! {
//!   "id":         "<uuid>",
//!   "name":       "fetch-prices",
//!   "parent_id":  "<uuid> | null",
//!   "started_at": "<rfc3339>",
//!   "ended_at":   "<rfc3339> | null",
//!   "fields":     { "key": "value", ... },
//!   "events":     [
//!     { "message": "...", "fields": {...}, "timestamp": "<rfc3339>" }
//!   ]
//! }
//! ```

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;

// ── Data types ───────────────────────────────────────────────────────────────

/// A single log event embedded in a span.
#[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,
}

/// A tracing span.  Persisted as `logs/<id>.json`.
#[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(),
        }
    }
}

// ── File helpers ─────────────────────────────────────────────────────────────

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)?)
}

/// Load every span from `logs/`.
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);
        }
    }
    // Stable order: started_at ascending
    spans.sort_by(|a, b| a.started_at.cmp(&b.started_at));
    Ok(spans)
}

// ── Public API ───────────────────────────────────────────────────────────────

/// Start a new span, optionally nested under `parent_id`.
/// Returns the new span's UUID.
pub fn span_start(
    home: &UnifierHome,
    name: &str,
    parent_id: Option<Uuid>,
    fields: HashMap<String, String>,
) -> Result<Uuid> {
    if let Some(pid) = parent_id {
        // Verify parent exists.
        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)
}

/// End a span (record `ended_at`).
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)
}

/// Add a log event to a 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)
}

/// Set key/value fields on an existing 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)
}

// ── Tree rendering ────────────────────────────────────────────────────────────

/// Render a span tree as indented text, starting from roots (or the given span).
///
/// Output example:
/// ```text
/// ▶ fetch-prices  [started 2026-09-03T12:00:00Z, open]
///   ▶ parse-csv   [started ..., ended ...]
///     · row-count=42
///     · 12:00:01Z  parsed headers
///   ▶ insert-db   [started ..., ended ...]
/// ```
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
        ));
        // Fields
        for (k, v) in &span.fields {
            out.push_str(&format!("{}  · {}={}\n", indent, k, v));
        }
        // Embedded log events
        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);
    }
}

/// Shorten an RFC3339 timestamp to the time portion for compact display.
fn short_ts(ts: &str) -> &str {
    // "2026-09-03T12:00:00.123456789+00:00" → keep up to 'Z' or '+'/'-' after T
    if let Some(t_pos) = ts.find('T') {
        let after_t = &ts[t_pos + 1..];
        // strip offset suffix
        let end = after_t
            .find(['+', '-', 'Z'])
            .map(|i| t_pos + 1 + i)
            .unwrap_or(ts.len());
        &ts[t_pos + 1..end]
    } else {
        ts
    }
}

// ── Key/value field parsing ───────────────────────────────────────────────────

/// Parse `key=value` pairs from a slice of strings.
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)
}

// ── Tests ─────────────────────────────────────────────────────────────────────

#[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");
    }
}