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
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 for (k, v) in &span.fields {
234 out.push_str(&format!("{} · {}={}\n", indent, k, v));
235 }
236 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
260fn short_ts(ts: &str) -> &str {
262 if let Some(t_pos) = ts.find('T') {
264 let after_t = &ts[t_pos + 1..];
265 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
276pub 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#[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}