Skip to main content

lean_ctx/core/ocla/
tracing.rs

1use std::collections::VecDeque;
2use std::sync::{Arc, Mutex, MutexGuard, OnceLock, PoisonError};
3use std::time::{SystemTime, UNIX_EPOCH};
4
5use serde::Serialize;
6
7const MAX_SPANS: usize = 2048;
8
9#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
10pub enum SpanStatus {
11    Ok,
12    Error(String),
13}
14
15#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
16pub struct OclaSpan {
17    pub span_id: String,
18    pub trace_id: String,
19    pub parent_span_id: Option<String>,
20    pub operation: String,
21    pub start_ns: u64,
22    pub end_ns: Option<u64>,
23    pub status: SpanStatus,
24    pub attributes: Vec<(String, String)>,
25}
26
27#[derive(Clone, Debug)]
28pub struct SpanCollector {
29    spans: Arc<Mutex<VecDeque<OclaSpan>>>,
30}
31
32impl SpanCollector {
33    pub fn new() -> Self {
34        Self {
35            spans: Arc::new(Mutex::new(VecDeque::with_capacity(MAX_SPANS))),
36        }
37    }
38
39    fn lock(&self) -> MutexGuard<'_, VecDeque<OclaSpan>> {
40        self.spans.lock().unwrap_or_else(PoisonError::into_inner)
41    }
42
43    fn register(&self, span: OclaSpan) {
44        let mut spans = self.lock();
45        if spans.len() == MAX_SPANS {
46            spans.pop_front();
47        }
48        spans.push_back(span);
49    }
50
51    fn finish(&self, span_id: &str, end_ns: u64) {
52        let mut spans = self.lock();
53        if let Some(span) = spans.iter_mut().find(|span| span.span_id == span_id) {
54            span.end_ns = Some(end_ns);
55        }
56    }
57
58    fn set_status(&self, span_id: &str, status: SpanStatus) {
59        let mut spans = self.lock();
60        if let Some(span) = spans.iter_mut().find(|span| span.span_id == span_id) {
61            span.status = status;
62        }
63    }
64
65    fn add_attribute(&self, span_id: &str, key: String, value: String) {
66        let mut spans = self.lock();
67        if let Some(span) = spans.iter_mut().find(|span| span.span_id == span_id) {
68            span.attributes.push((key, value));
69        }
70    }
71
72    fn spans_for_trace(&self, trace_id: &str) -> Vec<OclaSpan> {
73        let spans = self.lock();
74        spans
75            .iter()
76            .filter(|span| span.trace_id == trace_id)
77            .cloned()
78            .collect()
79    }
80
81    pub(crate) fn span_count(&self) -> usize {
82        self.lock().len()
83    }
84}
85
86struct ActiveSpan {
87    trace_id: String,
88    span_id: String,
89}
90
91thread_local! {
92    static ACTIVE_SPANS: std::cell::RefCell<Vec<ActiveSpan>> = const {
93        std::cell::RefCell::new(Vec::new())
94    };
95}
96
97static COLLECTOR: OnceLock<SpanCollector> = OnceLock::new();
98
99fn collector() -> &'static SpanCollector {
100    COLLECTOR.get_or_init(SpanCollector::new)
101}
102
103pub(crate) fn initialized_collector() -> Option<&'static SpanCollector> {
104    COLLECTOR.get()
105}
106
107fn next_span_id() -> String {
108    let mut bytes = [0_u8; 8];
109    getrandom::fill(&mut bytes).expect("CSPRNG unavailable");
110    format!("{:016x}", u64::from_be_bytes(bytes))
111}
112
113fn now_ns() -> u64 {
114    SystemTime::now()
115        .duration_since(UNIX_EPOCH)
116        .unwrap_or_default()
117        .as_nanos() as u64
118}
119
120/// RAII handle that closes its span when dropped.
121pub struct SpanGuard {
122    collector: &'static SpanCollector,
123    span_id: String,
124}
125
126impl SpanGuard {
127    pub fn set_status(&self, status: SpanStatus) {
128        self.collector.set_status(&self.span_id, status);
129    }
130
131    pub fn add_attribute(&self, key: impl Into<String>, value: impl Into<String>) {
132        self.collector
133            .add_attribute(&self.span_id, key.into(), value.into());
134    }
135}
136
137impl Drop for SpanGuard {
138    fn drop(&mut self) {
139        self.collector.finish(&self.span_id, now_ns());
140        ACTIVE_SPANS.with(|active| {
141            let mut active = active.borrow_mut();
142            if let Some(index) = active.iter().rposition(|span| span.span_id == self.span_id) {
143                active.remove(index);
144            }
145        });
146    }
147}
148
149pub fn start_span(trace_id: &str, operation: &str) -> SpanGuard {
150    let span_id = next_span_id();
151    let parent_span_id = ACTIVE_SPANS.with(|active| {
152        active
153            .borrow()
154            .iter()
155            .rev()
156            .find(|span| span.trace_id == trace_id)
157            .map(|span| span.span_id.clone())
158    });
159    collector().register(OclaSpan {
160        span_id: span_id.clone(),
161        trace_id: trace_id.to_string(),
162        parent_span_id,
163        operation: operation.to_string(),
164        start_ns: now_ns(),
165        end_ns: None,
166        status: SpanStatus::Ok,
167        attributes: Vec::new(),
168    });
169    ACTIVE_SPANS.with(|active| {
170        active.borrow_mut().push(ActiveSpan {
171            trace_id: trace_id.to_string(),
172            span_id: span_id.clone(),
173        });
174    });
175    SpanGuard {
176        collector: collector(),
177        span_id,
178    }
179}
180
181pub fn spans_for_trace(trace_id: &str) -> Vec<OclaSpan> {
182    collector().spans_for_trace(trace_id)
183}
184
185/// Summary of trace spans that contributed tool savings.
186#[derive(Debug, Clone, Serialize)]
187pub struct TraceSavingsSummary {
188    /// Trace identifier used for the join.
189    pub trace_id: String,
190    /// Number of spans associated with the trace.
191    pub span_count: usize,
192    /// Tool names recorded on associated spans.
193    pub tool_names: Vec<String>,
194}
195
196/// Join trace spans with savings events by trace identifier.
197pub fn trace_savings_summary(trace_id: &str) -> TraceSavingsSummary {
198    let spans = spans_for_trace(trace_id);
199    let tool_names = spans
200        .iter()
201        .flat_map(|span| span.attributes.iter())
202        .filter_map(|(key, value)| (key == "tool").then_some(value.clone()))
203        .collect();
204    TraceSavingsSummary {
205        trace_id: trace_id.to_owned(),
206        span_count: spans.len(),
207        tool_names,
208    }
209}
210
211fn status_value(status: &SpanStatus) -> serde_json::Value {
212    match status {
213        SpanStatus::Ok => serde_json::json!({"code": "STATUS_OK"}),
214        SpanStatus::Error(message) => {
215            serde_json::json!({"code": "STATUS_ERROR", "message": message})
216        }
217    }
218}
219
220pub fn export_trace(trace_id: &str) -> serde_json::Value {
221    let spans: Vec<serde_json::Value> = spans_for_trace(trace_id)
222        .into_iter()
223        .map(|span| {
224            let attributes: Vec<serde_json::Value> = span
225                .attributes
226                .into_iter()
227                .map(|(key, value)| {
228                    serde_json::json!({
229                        "key": key,
230                        "value": {"stringValue": value}
231                    })
232                })
233                .collect();
234            let mut exported = serde_json::json!({
235                "traceId": span.trace_id,
236                "spanId": span.span_id,
237                "name": span.operation,
238                "startTimeUnixNano": span.start_ns.to_string(),
239                "status": status_value(&span.status),
240                "attributes": attributes,
241            });
242            if let Some(parent_span_id) = span.parent_span_id {
243                exported["parentSpanId"] = serde_json::Value::String(parent_span_id);
244            }
245            if let Some(end_ns) = span.end_ns {
246                exported["endTimeUnixNano"] = serde_json::Value::String(end_ns.to_string());
247            }
248            exported
249        })
250        .collect();
251    serde_json::json!({
252        "resourceSpans": [{
253            "scopeSpans": [{"spans": spans}]
254        }]
255    })
256}
257
258impl Default for SpanCollector {
259    fn default() -> Self {
260        Self::new()
261    }
262}
263
264#[cfg(test)]
265mod tests {
266    use super::*;
267
268    #[test]
269    fn test_trace_savings_summary() {
270        let trace_id = "trace-savings-summary";
271        let span = start_span(trace_id, "read");
272        span.add_attribute("tool", "ctx_read");
273        drop(span);
274
275        let summary = trace_savings_summary(trace_id);
276
277        assert_eq!(summary.trace_id, trace_id);
278        assert_eq!(summary.span_count, 1);
279        assert_eq!(summary.tool_names, vec!["ctx_read"]);
280    }
281
282    fn unique_trace(label: &str) -> String {
283        format!("{label}-{}", next_span_id())
284    }
285
286    #[test]
287    fn span_closes_with_monotonic_wall_clock_timing() {
288        let trace_id = unique_trace("timing");
289        drop(start_span(&trace_id, "test.operation"));
290        let span = spans_for_trace(&trace_id).pop().expect("span retained");
291        assert!(span.end_ns.expect("end time") >= span.start_ns);
292    }
293
294    #[test]
295    fn nested_spans_link_to_the_active_parent() {
296        let trace_id = unique_trace("parent");
297        let parent = start_span(&trace_id, "parent");
298        let parent_id = parent.span_id.clone();
299        let child = start_span(&trace_id, "child");
300        let child_id = child.span_id.clone();
301        drop(child);
302        drop(parent);
303        let child = spans_for_trace(&trace_id)
304            .into_iter()
305            .find(|span| span.span_id == child_id)
306            .expect("child retained");
307        assert_eq!(child.parent_span_id.as_deref(), Some(parent_id.as_str()));
308    }
309
310    #[test]
311    fn traces_are_grouped_and_exported_as_otel_json() {
312        let trace_id = unique_trace("export");
313        let guard = start_span(&trace_id, "export.operation");
314        guard.set_status(SpanStatus::Error("failed".into()));
315        guard.add_attribute("component", "ocla");
316        drop(guard);
317        assert_eq!(spans_for_trace(&trace_id).len(), 1);
318        let exported = export_trace(&trace_id);
319        let span = &exported["resourceSpans"][0]["scopeSpans"][0]["spans"][0];
320        assert_eq!(span["name"], "export.operation");
321        assert_eq!(span["status"]["code"], "STATUS_ERROR");
322        assert_eq!(span["attributes"][0]["value"]["stringValue"], "ocla");
323    }
324
325    #[test]
326    fn collector_evicts_oldest_spans_at_capacity() {
327        let collector = SpanCollector::new();
328        for index in 0..=MAX_SPANS {
329            collector.register(OclaSpan {
330                span_id: index.to_string(),
331                trace_id: "overflow".into(),
332                parent_span_id: None,
333                operation: "test".into(),
334                start_ns: index as u64,
335                end_ns: Some(index as u64),
336                status: SpanStatus::Ok,
337                attributes: Vec::new(),
338            });
339        }
340        let spans = collector.spans_for_trace("overflow");
341        assert_eq!(spans[0].span_id, "1");
342        assert_eq!(spans[MAX_SPANS - 1].span_id, MAX_SPANS.to_string());
343    }
344}