1use 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#[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#[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
82fn 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
109pub 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 spans.sort_by(|a, b| a.started_at.cmp(&b.started_at));
135 Ok(spans)
136}
137
138pub 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 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
159pub 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
169pub 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
185pub 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
196pub 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 for (k, v) in &span.fields {
231 out.push_str(&format!("{} · {}={}\n", indent, k, v));
232 }
233 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
258fn short_ts(ts: &str) -> &str {
260 if let Some(t_pos) = ts.find('T') {
262 let after_t = &ts[t_pos + 1..];
263 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
274pub 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#[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}