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.iter().filter(|s| s.parent_id == parent).collect();
217    for span in children {
218        let status = match &span.ended_at {
219            Some(t) => format!("ended {}", short_ts(t)),
220            None => "open".to_string(),
221        };
222        out.push_str(&format!(
223            "{}▶ {}  [started {}, {}]\n",
224            indent,
225            span.name,
226            short_ts(&span.started_at),
227            status
228        ));
229        // Fields
230        for (k, v) in &span.fields {
231            out.push_str(&format!("{}  · {}={}\n", indent, k, v));
232        }
233        // Embedded log events
234        for ev in &span.events {
235            let field_str: Vec<String> =
236                ev.fields.iter().map(|(k, v)| format!("{k}={v}")).collect();
237            if field_str.is_empty() {
238                out.push_str(&format!(
239                    "{}  · {}  {}\n",
240                    indent,
241                    short_ts(&ev.timestamp),
242                    ev.message
243                ));
244            } else {
245                out.push_str(&format!(
246                    "{}  · {}  {}  ({})\n",
247                    indent,
248                    short_ts(&ev.timestamp),
249                    ev.message,
250                    field_str.join(", ")
251                ));
252            }
253        }
254        render_children(spans, Some(span.id), depth + 1, out);
255    }
256}
257
258/// Shorten an RFC3339 timestamp to the time portion for compact display.
259fn short_ts(ts: &str) -> &str {
260    // "2026-09-03T12:00:00.123456789+00:00" → keep up to 'Z' or '+'/'-' after T
261    if let Some(t_pos) = ts.find('T') {
262        let after_t = &ts[t_pos + 1..];
263        // strip offset suffix
264        let end = after_t
265            .find(['+', '-', 'Z'])
266            .map(|i| t_pos + 1 + i)
267            .unwrap_or(ts.len());
268        &ts[t_pos + 1..end]
269    } else {
270        ts
271    }
272}
273
274// ── Key/value field parsing ───────────────────────────────────────────────────
275
276/// Parse `key=value` pairs from a slice of strings.
277pub fn parse_fields(pairs: &[String]) -> Result<HashMap<String, String>> {
278    let mut map = HashMap::new();
279    for pair in pairs {
280        let (k, v) = pair
281            .split_once('=')
282            .ok_or_else(|| Error::msg(format!("field must be key=value, got: {pair}")))?;
283        map.insert(k.to_string(), v.to_string());
284    }
285    Ok(map)
286}
287
288// ── Tests ─────────────────────────────────────────────────────────────────────
289
290#[cfg(test)]
291mod tests {
292    use super::*;
293    use tempfile::tempdir;
294
295    fn home(dir: &std::path::Path) -> UnifierHome {
296        UnifierHome::resolve(Some(dir.to_path_buf()), None).unwrap()
297    }
298
299    #[test]
300    fn span_start_end_roundtrip() {
301        let tmp = tempdir().unwrap();
302        let h = home(tmp.path());
303        let id = span_start(&h, "test", None, HashMap::new()).unwrap();
304        let span = load_span(&h, id).unwrap();
305        assert_eq!(span.name, "test");
306        assert!(span.ended_at.is_none());
307        span_end(&h, id).unwrap();
308        let span = load_span(&h, id).unwrap();
309        assert!(span.ended_at.is_some());
310    }
311
312    #[test]
313    fn nested_spans() {
314        let tmp = tempdir().unwrap();
315        let h = home(tmp.path());
316        let root = span_start(&h, "root", None, HashMap::new()).unwrap();
317        let child = span_start(&h, "child", Some(root), HashMap::new()).unwrap();
318        let spans = load_all_spans(&h).unwrap();
319        assert_eq!(spans.len(), 2);
320        let c = spans.iter().find(|s| s.id == child).unwrap();
321        assert_eq!(c.parent_id, Some(root));
322    }
323
324    #[test]
325    fn log_event_appended() {
326        let tmp = tempdir().unwrap();
327        let h = home(tmp.path());
328        let id = span_start(&h, "work", None, HashMap::new()).unwrap();
329        log_event(&h, id, "fetched 10 rows", HashMap::new()).unwrap();
330        let span = load_span(&h, id).unwrap();
331        assert_eq!(span.events.len(), 1);
332        assert_eq!(span.events[0].message, "fetched 10 rows");
333    }
334
335    #[test]
336    fn render_tree_smoke() {
337        let tmp = tempdir().unwrap();
338        let h = home(tmp.path());
339        let root = span_start(&h, "pipeline", None, HashMap::new()).unwrap();
340        let _child = span_start(&h, "step-1", Some(root), HashMap::new()).unwrap();
341        span_end(&h, root).unwrap();
342        let spans = load_all_spans(&h).unwrap();
343        let tree = render_tree(&spans, None);
344        assert!(tree.contains("pipeline"));
345        assert!(tree.contains("step-1"));
346    }
347
348    #[test]
349    fn double_end_is_error() {
350        let tmp = tempdir().unwrap();
351        let h = home(tmp.path());
352        let id = span_start(&h, "x", None, HashMap::new()).unwrap();
353        span_end(&h, id).unwrap();
354        assert!(span_end(&h, id).is_err());
355    }
356
357    #[test]
358    fn parse_fields_ok() {
359        let pairs = vec!["k=v".to_string(), "foo=bar=baz".to_string()];
360        let m = parse_fields(&pairs).unwrap();
361        assert_eq!(m["k"], "v");
362        assert_eq!(m["foo"], "bar=baz");
363    }
364}