Skip to main content

unifier/
log.rs

1//! Span-tree logging for Jan cron scripts and agent pipelines.
2//!
3//! Spans live under `logs/` in the store root.  Each span is a JSON file
4//! `logs/<id>.json`; log events are embedded in the span they belong to.
5//! Reading back a span tree walks the directory and reconstructs parent→child
6//! relationships in memory.
7//!
8//! Layout:
9//! ```text
10//! logs/
11//!   <uuid>.json    # one file per span
12//! ```
13//!
14//! Span JSON schema:
15//! ```json
16//! {
17//!   "id":         "<uuid>",
18//!   "name":       "fetch-prices",
19//!   "parent_id":  "<uuid> | null",
20//!   "started_at": "<rfc3339>",
21//!   "ended_at":   "<rfc3339> | null",
22//!   "fields":     { "key": "value", ... },
23//!   "events":     [
24//!     { "message": "...", "fields": {...}, "timestamp": "<rfc3339>" }
25//!   ]
26//! }
27//! ```
28
29use std::collections::HashMap;
30use std::fs;
31
32use chrono::Utc;
33use serde::{Deserialize, Serialize};
34use uuid::Uuid;
35
36use crate::constants::LOGS;
37use crate::error::{Error, Result};
38use crate::fs_text::write_text_atomic;
39use crate::home::UnifierHome;
40
41// ── Data types ───────────────────────────────────────────────────────────────
42
43/// A single log event embedded in a span.
44#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
45pub struct LogEvent {
46    pub message: String,
47    #[serde(default, skip_serializing_if = "HashMap::is_empty")]
48    pub fields: HashMap<String, String>,
49    pub timestamp: String,
50}
51
52/// A tracing span.  Persisted as `logs/<id>.json`.
53#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
54pub struct Span {
55    pub id: Uuid,
56    pub name: String,
57    #[serde(skip_serializing_if = "Option::is_none")]
58    pub parent_id: Option<Uuid>,
59    pub started_at: String,
60    #[serde(skip_serializing_if = "Option::is_none")]
61    pub ended_at: Option<String>,
62    #[serde(default, skip_serializing_if = "HashMap::is_empty")]
63    pub fields: HashMap<String, String>,
64    #[serde(default, skip_serializing_if = "Vec::is_empty")]
65    pub events: Vec<LogEvent>,
66}
67
68impl Span {
69    fn new(name: &str, parent_id: Option<Uuid>) -> Self {
70        Self {
71            id: Uuid::new_v4(),
72            name: name.to_string(),
73            parent_id,
74            started_at: Utc::now().to_rfc3339(),
75            ended_at: None,
76            fields: HashMap::new(),
77            events: Vec::new(),
78        }
79    }
80}
81
82// ── File helpers ─────────────────────────────────────────────────────────────
83
84fn logs_dir(home: &UnifierHome) -> std::path::PathBuf {
85    home.path().join(LOGS)
86}
87
88fn span_path(home: &UnifierHome, id: Uuid) -> std::path::PathBuf {
89    logs_dir(home).join(format!("{}.json", id.hyphenated()))
90}
91
92fn save_span(home: &UnifierHome, span: &Span) -> Result<()> {
93    let dir = logs_dir(home);
94    fs::create_dir_all(&dir)?;
95    let json = serde_json::to_string_pretty(span)?;
96    write_text_atomic(&span_path(home, span.id), &json)?;
97    Ok(())
98}
99
100fn load_span(home: &UnifierHome, id: Uuid) -> Result<Span> {
101    let path = span_path(home, id);
102    if !path.is_file() {
103        return Err(Error::msg(format!("span not found: {id}")));
104    }
105    let text = fs::read_to_string(&path)?;
106    Ok(serde_json::from_str(&text)?)
107}
108
109/// Load every span from `logs/`.
110pub fn load_all_spans(home: &UnifierHome) -> Result<Vec<Span>> {
111    let dir = logs_dir(home);
112    if !dir.is_dir() {
113        return Ok(vec![]);
114    }
115    let mut spans = Vec::new();
116    for entry in fs::read_dir(&dir)? {
117        let entry = entry?;
118        if !entry.file_type()?.is_file() {
119            continue;
120        }
121        let name = entry.file_name().to_string_lossy().into_owned();
122        let Some(stem) = name.strip_suffix(".json") else {
123            continue;
124        };
125        let Ok(_id) = Uuid::parse_str(stem) else {
126            continue;
127        };
128        let text = fs::read_to_string(entry.path())?;
129        if let Ok(span) = serde_json::from_str::<Span>(&text) {
130            spans.push(span);
131        }
132    }
133    // Stable order: started_at ascending
134    spans.sort_by(|a, b| a.started_at.cmp(&b.started_at));
135    Ok(spans)
136}
137
138// ── Public API ───────────────────────────────────────────────────────────────
139
140/// Start a new span, optionally nested under `parent_id`.
141/// Returns the new span's UUID.
142pub fn span_start(
143    home: &UnifierHome,
144    name: &str,
145    parent_id: Option<Uuid>,
146    fields: HashMap<String, String>,
147) -> Result<Uuid> {
148    if let Some(pid) = parent_id {
149        // Verify parent exists.
150        load_span(home, pid)?;
151    }
152    let mut span = Span::new(name, parent_id);
153    span.fields = fields;
154    let id = span.id;
155    save_span(home, &span)?;
156    Ok(id)
157}
158
159/// End a span (record `ended_at`).
160pub fn span_end(home: &UnifierHome, id: Uuid) -> Result<()> {
161    let mut span = load_span(home, id)?;
162    if span.ended_at.is_some() {
163        return Err(Error::msg(format!("span {id} already ended")));
164    }
165    span.ended_at = Some(Utc::now().to_rfc3339());
166    save_span(home, &span)
167}
168
169/// Add a log event to a span.
170pub fn log_event(
171    home: &UnifierHome,
172    span_id: Uuid,
173    message: &str,
174    fields: HashMap<String, String>,
175) -> Result<()> {
176    let mut span = load_span(home, span_id)?;
177    span.events.push(LogEvent {
178        message: message.to_string(),
179        fields,
180        timestamp: Utc::now().to_rfc3339(),
181    });
182    save_span(home, &span)
183}
184
185/// Set key/value fields on an existing span.
186pub fn span_set_fields(
187    home: &UnifierHome,
188    id: Uuid,
189    fields: HashMap<String, String>,
190) -> Result<()> {
191    let mut span = load_span(home, id)?;
192    span.fields.extend(fields);
193    save_span(home, &span)
194}
195
196// ── Tree rendering ────────────────────────────────────────────────────────────
197
198/// Render a span tree as indented text, starting from roots (or the given span).
199///
200/// Output example:
201/// ```text
202/// ▶ fetch-prices  [started 2026-09-03T12:00:00Z, open]
203///   ▶ parse-csv   [started ..., ended ...]
204///     · row-count=42
205///     · 12:00:01Z  parsed headers
206///   ▶ insert-db   [started ..., ended ...]
207/// ```
208pub fn render_tree(spans: &[Span], root_id: Option<Uuid>) -> String {
209    let mut out = String::new();
210    render_children(spans, root_id, 0, &mut out);
211    out
212}
213
214fn render_children(spans: &[Span], parent: Option<Uuid>, depth: usize, out: &mut String) {
215    let indent = "  ".repeat(depth);
216    let children: Vec<&Span> = spans
217        .iter()
218        .filter(|s| s.parent_id == parent)
219        .collect();
220    for span in children {
221        let status = match &span.ended_at {
222            Some(t) => format!("ended {}", short_ts(t)),
223            None => "open".to_string(),
224        };
225        out.push_str(&format!(
226            "{}▶ {}  [started {}, {}]\n",
227            indent,
228            span.name,
229            short_ts(&span.started_at),
230            status
231        ));
232        // Fields
233        for (k, v) in &span.fields {
234            out.push_str(&format!("{}  · {}={}\n", indent, k, v));
235        }
236        // Embedded log events
237        for ev in &span.events {
238            let field_str: Vec<String> = ev.fields.iter().map(|(k, v)| format!("{k}={v}")).collect();
239            if field_str.is_empty() {
240                out.push_str(&format!(
241                    "{}  · {}  {}\n",
242                    indent,
243                    short_ts(&ev.timestamp),
244                    ev.message
245                ));
246            } else {
247                out.push_str(&format!(
248                    "{}  · {}  {}  ({})\n",
249                    indent,
250                    short_ts(&ev.timestamp),
251                    ev.message,
252                    field_str.join(", ")
253                ));
254            }
255        }
256        render_children(spans, Some(span.id), depth + 1, out);
257    }
258}
259
260/// Shorten an RFC3339 timestamp to the time portion for compact display.
261fn short_ts(ts: &str) -> &str {
262    // "2026-09-03T12:00:00.123456789+00:00" → keep up to 'Z' or '+'/'-' after T
263    if let Some(t_pos) = ts.find('T') {
264        let after_t = &ts[t_pos + 1..];
265        // strip offset suffix
266        let end = after_t
267            .find(['+', '-', 'Z'])
268            .map(|i| t_pos + 1 + i)
269            .unwrap_or(ts.len());
270        &ts[t_pos + 1..end]
271    } else {
272        ts
273    }
274}
275
276// ── Key/value field parsing ───────────────────────────────────────────────────
277
278/// Parse `key=value` pairs from a slice of strings.
279pub fn parse_fields(pairs: &[String]) -> Result<HashMap<String, String>> {
280    let mut map = HashMap::new();
281    for pair in pairs {
282        let (k, v) = pair
283            .split_once('=')
284            .ok_or_else(|| Error::msg(format!("field must be key=value, got: {pair}")))?;
285        map.insert(k.to_string(), v.to_string());
286    }
287    Ok(map)
288}
289
290// ── Tests ─────────────────────────────────────────────────────────────────────
291
292#[cfg(test)]
293mod tests {
294    use super::*;
295    use tempfile::tempdir;
296
297    fn home(dir: &std::path::Path) -> UnifierHome {
298        UnifierHome::resolve(Some(dir.to_path_buf()), None).unwrap()
299    }
300
301    #[test]
302    fn span_start_end_roundtrip() {
303        let tmp = tempdir().unwrap();
304        let h = home(tmp.path());
305        let id = span_start(&h, "test", None, HashMap::new()).unwrap();
306        let span = load_span(&h, id).unwrap();
307        assert_eq!(span.name, "test");
308        assert!(span.ended_at.is_none());
309        span_end(&h, id).unwrap();
310        let span = load_span(&h, id).unwrap();
311        assert!(span.ended_at.is_some());
312    }
313
314    #[test]
315    fn nested_spans() {
316        let tmp = tempdir().unwrap();
317        let h = home(tmp.path());
318        let root = span_start(&h, "root", None, HashMap::new()).unwrap();
319        let child = span_start(&h, "child", Some(root), HashMap::new()).unwrap();
320        let spans = load_all_spans(&h).unwrap();
321        assert_eq!(spans.len(), 2);
322        let c = spans.iter().find(|s| s.id == child).unwrap();
323        assert_eq!(c.parent_id, Some(root));
324    }
325
326    #[test]
327    fn log_event_appended() {
328        let tmp = tempdir().unwrap();
329        let h = home(tmp.path());
330        let id = span_start(&h, "work", None, HashMap::new()).unwrap();
331        log_event(&h, id, "fetched 10 rows", HashMap::new()).unwrap();
332        let span = load_span(&h, id).unwrap();
333        assert_eq!(span.events.len(), 1);
334        assert_eq!(span.events[0].message, "fetched 10 rows");
335    }
336
337    #[test]
338    fn render_tree_smoke() {
339        let tmp = tempdir().unwrap();
340        let h = home(tmp.path());
341        let root = span_start(&h, "pipeline", None, HashMap::new()).unwrap();
342        let _child = span_start(&h, "step-1", Some(root), HashMap::new()).unwrap();
343        span_end(&h, root).unwrap();
344        let spans = load_all_spans(&h).unwrap();
345        let tree = render_tree(&spans, None);
346        assert!(tree.contains("pipeline"));
347        assert!(tree.contains("step-1"));
348    }
349
350    #[test]
351    fn double_end_is_error() {
352        let tmp = tempdir().unwrap();
353        let h = home(tmp.path());
354        let id = span_start(&h, "x", None, HashMap::new()).unwrap();
355        span_end(&h, id).unwrap();
356        assert!(span_end(&h, id).is_err());
357    }
358
359    #[test]
360    fn parse_fields_ok() {
361        let pairs = vec!["k=v".to_string(), "foo=bar=baz".to_string()];
362        let m = parse_fields(&pairs).unwrap();
363        assert_eq!(m["k"], "v");
364        assert_eq!(m["foo"], "bar=baz");
365    }
366}