tokio_prompt_orchestrator/
observability.rs1use std::collections::HashMap;
29use std::sync::atomic::{AtomicU64, Ordering};
30use std::sync::{Arc, Mutex};
31use std::time::{SystemTime, UNIX_EPOCH};
32
33static ID_COUNTER: AtomicU64 = AtomicU64::new(1);
37
38fn now_ns() -> u64 {
40 SystemTime::now()
41 .duration_since(UNIX_EPOCH)
42 .unwrap_or_default()
43 .as_nanos() as u64
44}
45
46fn gen_id64() -> u64 {
49 let seq = ID_COUNTER.fetch_add(1, Ordering::Relaxed);
50 let time = now_ns();
51 seq
54 .wrapping_mul(6_364_136_223_846_793_005)
55 .wrapping_add(time ^ 1_442_695_040_888_963_407)
56}
57
58fn gen_trace_id() -> u128 {
60 let hi = gen_id64() as u128;
61 let lo = gen_id64() as u128;
62 (hi << 64) | lo
63}
64
65#[derive(Debug, Clone, PartialEq, Eq)]
70pub struct SpanContext {
71 pub trace_id: u128,
73 pub span_id: u64,
75 pub parent_span_id: Option<u64>,
77 pub sampled: bool,
79}
80
81impl SpanContext {
82 pub fn new_root(sampled: bool) -> Self {
84 Self {
85 trace_id: gen_trace_id(),
86 span_id: gen_id64(),
87 parent_span_id: None,
88 sampled,
89 }
90 }
91
92 pub fn child(&self) -> Self {
95 Self {
96 trace_id: self.trace_id,
97 span_id: gen_id64(),
98 parent_span_id: Some(self.span_id),
99 sampled: self.sampled,
100 }
101 }
102}
103
104#[derive(Debug, Clone, PartialEq)]
108pub enum AttributeValue {
109 String(String),
111 Int(i64),
113 Float(f64),
115 Bool(bool),
117}
118
119impl From<&str> for AttributeValue {
120 fn from(s: &str) -> Self {
121 AttributeValue::String(s.to_owned())
122 }
123}
124
125impl From<String> for AttributeValue {
126 fn from(s: String) -> Self {
127 AttributeValue::String(s)
128 }
129}
130
131impl From<i64> for AttributeValue {
132 fn from(i: i64) -> Self {
133 AttributeValue::Int(i)
134 }
135}
136
137impl From<f64> for AttributeValue {
138 fn from(f: f64) -> Self {
139 AttributeValue::Float(f)
140 }
141}
142
143impl From<bool> for AttributeValue {
144 fn from(b: bool) -> Self {
145 AttributeValue::Bool(b)
146 }
147}
148
149#[derive(Debug, Clone)]
153pub struct SpanEvent {
154 pub name: String,
156 pub timestamp_ns: u64,
158 pub attributes: HashMap<String, AttributeValue>,
160}
161
162#[derive(Debug, Clone, PartialEq)]
166#[derive(Default)]
167pub enum SpanStatus {
168 #[default]
170 Unset,
171 Ok,
173 Error(String),
175}
176
177
178#[derive(Debug)]
182pub struct Span {
183 pub context: SpanContext,
185 pub name: String,
187 pub start_ns: u64,
189 pub end_ns: Option<u64>,
192 pub attributes: HashMap<String, AttributeValue>,
194 pub events: Vec<SpanEvent>,
196 pub status: SpanStatus,
198}
199
200impl Span {
201 fn new(context: SpanContext, name: impl Into<String>) -> Self {
203 Self {
204 context,
205 name: name.into(),
206 start_ns: now_ns(),
207 end_ns: None,
208 attributes: HashMap::new(),
209 events: Vec::new(),
210 status: SpanStatus::Unset,
211 }
212 }
213
214 pub fn set_attribute(&mut self, key: impl Into<String>, value: AttributeValue) {
216 self.attributes.insert(key.into(), value);
217 }
218
219 pub fn add_event(&mut self, name: impl Into<String>, attributes: HashMap<String, AttributeValue>) {
221 self.events.push(SpanEvent {
222 name: name.into(),
223 timestamp_ns: now_ns(),
224 attributes,
225 });
226 }
227
228 pub fn finish(&mut self) {
230 self.end_ns = Some(now_ns());
231 }
232
233 pub fn duration_ns(&self) -> Option<u64> {
236 self.end_ns.map(|end| end.saturating_sub(self.start_ns))
237 }
238}
239
240pub trait SpanExporter: Send + Sync {
244 fn export(&self, span: Span);
246}
247
248#[derive(Debug, Default)]
254pub struct InMemoryExporter {
255 spans: Arc<Mutex<Vec<Span>>>,
256}
257
258impl InMemoryExporter {
259 pub fn new() -> Self {
261 Self {
262 spans: Arc::new(Mutex::new(Vec::new())),
263 }
264 }
265
266 pub fn get_spans(&self) -> Vec<String> {
268 self.spans
270 .lock()
271 .unwrap_or_else(|e| e.into_inner())
272 .iter()
273 .map(|s| s.name.clone())
274 .collect()
275 }
276
277 pub fn drain_spans(&self) -> Vec<String> {
279 let mut guard = self.spans.lock().unwrap_or_else(|e| e.into_inner());
280 let names: Vec<String> = guard.iter().map(|s| s.name.clone()).collect();
281 guard.clear();
282 names
283 }
284
285 pub fn span_count(&self) -> usize {
287 self.spans
288 .lock()
289 .unwrap_or_else(|e| e.into_inner())
290 .len()
291 }
292}
293
294impl SpanExporter for InMemoryExporter {
295 fn export(&self, span: Span) {
296 let mut guard = self.spans.lock().unwrap_or_else(|e| e.into_inner());
297 guard.push(span);
298 }
299}
300
301#[derive(Clone)]
305pub struct Tracer {
306 exporter: Arc<dyn SpanExporter>,
307}
308
309impl std::fmt::Debug for Tracer {
310 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
311 f.debug_struct("Tracer").field("exporter", &"<dyn SpanExporter>").finish()
312 }
313}
314
315impl Tracer {
316 pub fn new(exporter: Arc<dyn SpanExporter>) -> Self {
318 Self { exporter }
319 }
320
321 pub fn start_span(&self, name: &str) -> Span {
324 let ctx = SpanContext::new_root(true);
325 Span::new(ctx, name)
326 }
327
328 pub fn start_child_span(&self, name: &str, parent: &Span) -> Span {
331 let ctx = parent.context.child();
332 Span::new(ctx, name)
333 }
334
335 pub fn finish_span(&self, mut span: Span) {
337 span.finish();
338 self.exporter.export(span);
339 }
340}
341
342#[derive(Debug, Clone)]
346pub struct MetricPoint {
347 pub name: String,
349 pub value: f64,
351 pub labels: HashMap<String, String>,
353 pub timestamp_ns: u64,
355}
356
357pub struct ObservabilityRegistry {
362 pub tracer: Arc<Tracer>,
364 pub metrics: Arc<Mutex<Vec<MetricPoint>>>,
366}
367
368impl ObservabilityRegistry {
369 pub fn new(tracer: Arc<Tracer>) -> Self {
371 Self {
372 tracer,
373 metrics: Arc::new(Mutex::new(Vec::new())),
374 }
375 }
376
377 pub fn record_metric(
379 &self,
380 name: &str,
381 value: f64,
382 labels: HashMap<String, String>,
383 ) {
384 let point = MetricPoint {
385 name: name.to_owned(),
386 value,
387 labels,
388 timestamp_ns: now_ns(),
389 };
390 let mut guard = self.metrics.lock().unwrap_or_else(|e| e.into_inner());
391 guard.push(point);
392 }
393
394 pub fn to_prometheus(&self) -> String {
402 let guard = self.metrics.lock().unwrap_or_else(|e| e.into_inner());
403 let mut out = String::new();
404
405 for point in guard.iter() {
406 let label_str = if point.labels.is_empty() {
408 String::new()
409 } else {
410 let pairs: Vec<String> = point
411 .labels
412 .iter()
413 .map(|(k, v)| format!("{}=\"{}\"", k, v))
414 .collect();
415 format!("{{{}}}", pairs.join(","))
416 };
417
418 let ts_ms = point.timestamp_ns / 1_000_000;
419 out.push_str(&format!(
420 "{}{} {} {}\n",
421 point.name, label_str, point.value, ts_ms
422 ));
423 }
424
425 out
426 }
427}
428
429#[cfg(test)]
432mod tests {
433 use super::*;
434 use std::collections::HashMap;
435
436 #[test]
437 fn span_duration_calculation() {
438 let ctx = SpanContext::new_root(true);
439 let mut span = Span::new(ctx, "test-span");
440 assert!(span.duration_ns().is_none(), "span should have no duration before finish");
441 span.finish();
442 let dur = span.duration_ns();
443 assert!(dur.is_some(), "finished span should have a duration");
444 assert!(dur.unwrap() < u64::MAX);
446 }
447
448 #[test]
449 fn child_span_inherits_trace_id() {
450 let root_ctx = SpanContext::new_root(true);
451 let child_ctx = root_ctx.child();
452
453 assert_eq!(
454 root_ctx.trace_id, child_ctx.trace_id,
455 "child should inherit trace_id"
456 );
457 assert_ne!(
458 root_ctx.span_id, child_ctx.span_id,
459 "child should have a distinct span_id"
460 );
461 assert_eq!(
462 child_ctx.parent_span_id,
463 Some(root_ctx.span_id),
464 "child parent_span_id should be root span_id"
465 );
466 }
467
468 #[test]
469 fn in_memory_exporter_collects_spans() {
470 let exporter = Arc::new(InMemoryExporter::new());
471 let tracer = Arc::new(Tracer::new(exporter.clone()));
472
473 let span1 = tracer.start_span("op-one");
474 tracer.finish_span(span1);
475
476 let span2 = tracer.start_span("op-two");
477 tracer.finish_span(span2);
478
479 assert_eq!(exporter.span_count(), 2);
480 let names = exporter.get_spans();
481 assert!(names.contains(&"op-one".to_string()));
482 assert!(names.contains(&"op-two".to_string()));
483 }
484
485 #[test]
486 fn child_span_via_tracer() {
487 let exporter = Arc::new(InMemoryExporter::new());
488 let tracer = Arc::new(Tracer::new(exporter.clone()));
489
490 let parent = tracer.start_span("parent");
491 let child = tracer.start_child_span("child", &parent);
492
493 assert_eq!(parent.context.trace_id, child.context.trace_id);
494 assert_eq!(child.context.parent_span_id, Some(parent.context.span_id));
495
496 tracer.finish_span(parent);
497 tracer.finish_span(child);
498 assert_eq!(exporter.span_count(), 2);
499 }
500
501 #[test]
502 fn prometheus_output_format() {
503 let exporter = Arc::new(InMemoryExporter::new());
504 let tracer = Arc::new(Tracer::new(exporter));
505 let registry = ObservabilityRegistry::new(tracer);
506
507 let mut labels = HashMap::new();
508 labels.insert("model".to_string(), "gpt-4".to_string());
509 registry.record_metric("tokens_total", 42.0, labels);
510 registry.record_metric("latency_ms", 123.5, HashMap::new());
511
512 let prom = registry.to_prometheus();
513 assert!(prom.contains("tokens_total{"), "should contain metric with labels");
514 assert!(prom.contains("model=\"gpt-4\""), "should contain label value");
515 assert!(prom.contains("42"), "should contain metric value");
516 assert!(prom.contains("latency_ms"), "should contain second metric");
517 assert!(prom.contains("123.5"), "should contain second metric value");
518 }
519
520 #[test]
521 fn root_span_has_no_parent() {
522 let ctx = SpanContext::new_root(true);
523 assert!(ctx.parent_span_id.is_none());
524 }
525
526 #[test]
527 fn span_attributes_and_events() {
528 let ctx = SpanContext::new_root(true);
529 let mut span = Span::new(ctx, "attr-test");
530 span.set_attribute("key", AttributeValue::String("value".to_string()));
531 span.set_attribute("count", AttributeValue::Int(7));
532 span.add_event("cache-hit", HashMap::new());
533 span.finish();
534
535 assert_eq!(span.attributes.len(), 2);
536 assert_eq!(span.events.len(), 1);
537 assert_eq!(span.events[0].name, "cache-hit");
538 }
539
540 #[test]
541 fn unique_trace_ids() {
542 let ctx1 = SpanContext::new_root(true);
543 let ctx2 = SpanContext::new_root(true);
544 assert_ne!(ctx1.trace_id, ctx2.trace_id);
545 assert_ne!(ctx1.span_id, ctx2.span_id);
546 }
547}