Skip to main content

vv_agent/
tracing.rs

1use std::collections::BTreeMap;
2use std::fs::{File, OpenOptions};
3use std::io::Write;
4use std::path::Path;
5use std::sync::Mutex;
6use std::time::{SystemTime, UNIX_EPOCH};
7
8use serde::{Deserialize, Serialize};
9use serde_json::{json, Value};
10
11pub trait TraceSink: Send + Sync {
12    fn on_span_start(&self, span: &Span);
13    fn on_span_end(&self, span: &Span);
14
15    fn flush(&self) -> Result<(), String> {
16        Ok(())
17    }
18}
19
20#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
21pub struct Span {
22    pub name: String,
23    pub trace_id: String,
24    pub span_id: String,
25    #[serde(default, skip_serializing_if = "Option::is_none")]
26    pub parent_id: Option<String>,
27    pub started_at: f64,
28    #[serde(default, skip_serializing_if = "Option::is_none")]
29    pub ended_at: Option<f64>,
30    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
31    pub metadata: BTreeMap<String, Value>,
32}
33
34impl Span {
35    pub fn new(trace_id: impl Into<String>, name: impl Into<String>) -> Self {
36        Self {
37            name: name.into(),
38            trace_id: trace_id.into(),
39            span_id: format!("span_{}", uuid::Uuid::new_v4().simple()),
40            parent_id: None,
41            started_at: timestamp_seconds(),
42            ended_at: None,
43            metadata: BTreeMap::new(),
44        }
45    }
46
47    pub fn with_parent_id(mut self, parent_id: impl Into<String>) -> Self {
48        self.parent_id = Some(parent_id.into());
49        self
50    }
51
52    pub fn with_metadata(mut self, key: impl Into<String>, value: Value) -> Self {
53        self.metadata.insert(key.into(), value);
54        self
55    }
56
57    pub fn finish(mut self) -> Self {
58        self.ended_at = Some(timestamp_seconds());
59        self
60    }
61}
62
63pub struct JsonlTraceExporter {
64    file: Mutex<File>,
65}
66
67impl JsonlTraceExporter {
68    pub fn new(path: impl AsRef<Path>) -> Result<Self, String> {
69        let file = OpenOptions::new()
70            .create(true)
71            .append(true)
72            .open(path)
73            .map_err(|error| error.to_string())?;
74        Ok(Self {
75            file: Mutex::new(file),
76        })
77    }
78
79    fn write_event(&self, event: &str, span: &Span) {
80        if let Ok(mut file) = self.file.lock() {
81            let _ = writeln!(
82                file,
83                "{}",
84                json!({
85                    "event": event,
86                    "timestamp": timestamp_seconds(),
87                    "span": span,
88                })
89            );
90        }
91    }
92}
93
94impl TraceSink for JsonlTraceExporter {
95    fn on_span_start(&self, span: &Span) {
96        self.write_event("span_start", span);
97    }
98
99    fn on_span_end(&self, span: &Span) {
100        self.write_event("span_end", span);
101    }
102
103    fn flush(&self) -> Result<(), String> {
104        self.file
105            .lock()
106            .map_err(|_| "trace exporter lock poisoned".to_string())?
107            .flush()
108            .map_err(|error| error.to_string())
109    }
110}
111
112fn timestamp_seconds() -> f64 {
113    SystemTime::now()
114        .duration_since(UNIX_EPOCH)
115        .map(|duration| duration.as_secs_f64())
116        .unwrap_or_default()
117}