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
120pub 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#[derive(Debug, Clone, Serialize)]
187pub struct TraceSavingsSummary {
188 pub trace_id: String,
190 pub span_count: usize,
192 pub tool_names: Vec<String>,
194}
195
196pub 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}