Skip to main content

tokio_prompt_orchestrator/
observability.rs

1//! # Observability — OpenTelemetry-Compatible Tracing and Metrics
2//!
3//! Implements an OpenTelemetry-compatible data model from scratch (no external
4//! OTel crate) including spans, exporters, and a Prometheus-format metrics
5//! registry.
6//!
7//! ## Example
8//!
9//! ```rust
10//! use std::collections::HashMap;
11//! use std::sync::Arc;
12//! use tokio_prompt_orchestrator::observability::{
13//!     InMemoryExporter, ObservabilityRegistry, Tracer,
14//! };
15//!
16//! let exporter = Arc::new(InMemoryExporter::new());
17//! let tracer = Arc::new(Tracer::new(exporter.clone()));
18//! let registry = ObservabilityRegistry::new(tracer);
19//!
20//! let mut span = registry.tracer.start_span("my-operation");
21//! span.set_attribute("http.method", "GET".into());
22//! span.finish();
23//!
24//! registry.record_metric("requests_total", 1.0, HashMap::new());
25//! println!("{}", registry.to_prometheus());
26//! ```
27
28use std::collections::HashMap;
29use std::sync::atomic::{AtomicU64, Ordering};
30use std::sync::{Arc, Mutex};
31use std::time::{SystemTime, UNIX_EPOCH};
32
33// ── ID Generation ─────────────────────────────────────────────────────────────
34
35/// Global monotonic counter used to contribute uniqueness to generated IDs.
36static ID_COUNTER: AtomicU64 = AtomicU64::new(1);
37
38/// Returns the current time in nanoseconds since UNIX epoch.
39fn now_ns() -> u64 {
40    SystemTime::now()
41        .duration_since(UNIX_EPOCH)
42        .unwrap_or_default()
43        .as_nanos() as u64
44}
45
46/// Generates a pseudo-random u64 using an LCG seeded from the counter and
47/// current time.  Not cryptographically secure — suitable for trace IDs only.
48fn gen_id64() -> u64 {
49    let seq = ID_COUNTER.fetch_add(1, Ordering::Relaxed);
50    let time = now_ns();
51    // Mix with a simple LCG step.
52    
53    seq
54        .wrapping_mul(6_364_136_223_846_793_005)
55        .wrapping_add(time ^ 1_442_695_040_888_963_407)
56}
57
58/// Generates a pseudo-random u128 trace ID.
59fn 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// ── SpanContext ───────────────────────────────────────────────────────────────
66
67/// Immutable identity and sampling context for a span, equivalent to the
68/// OpenTelemetry `SpanContext`.
69#[derive(Debug, Clone, PartialEq, Eq)]
70pub struct SpanContext {
71    /// 128-bit trace identifier shared by all spans in a trace.
72    pub trace_id: u128,
73    /// 64-bit identifier unique within the trace.
74    pub span_id: u64,
75    /// `span_id` of the parent span, or `None` for a root span.
76    pub parent_span_id: Option<u64>,
77    /// Whether this span is sampled (i.e. should be exported).
78    pub sampled: bool,
79}
80
81impl SpanContext {
82    /// Creates a new root-level span context with freshly generated IDs.
83    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    /// Creates a child span context that inherits `trace_id` and records this
93    /// span as the parent.
94    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// ── AttributeValue ────────────────────────────────────────────────────────────
105
106/// A typed attribute value, matching the OpenTelemetry attribute type system.
107#[derive(Debug, Clone, PartialEq)]
108pub enum AttributeValue {
109    /// UTF-8 string attribute.
110    String(String),
111    /// 64-bit signed integer attribute.
112    Int(i64),
113    /// 64-bit floating-point attribute.
114    Float(f64),
115    /// Boolean attribute.
116    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// ── SpanEvent ─────────────────────────────────────────────────────────────────
150
151/// A timestamped event recorded within a span's lifetime.
152#[derive(Debug, Clone)]
153pub struct SpanEvent {
154    /// Human-readable event name.
155    pub name: String,
156    /// Nanoseconds since UNIX epoch at which the event occurred.
157    pub timestamp_ns: u64,
158    /// Arbitrary key-value attributes attached to this event.
159    pub attributes: HashMap<String, AttributeValue>,
160}
161
162// ── SpanStatus ────────────────────────────────────────────────────────────────
163
164/// The completion status of a span, mirroring the OpenTelemetry `StatusCode`.
165#[derive(Debug, Clone, PartialEq)]
166#[derive(Default)]
167pub enum SpanStatus {
168    /// No explicit status has been set.
169    #[default]
170    Unset,
171    /// The operation completed successfully.
172    Ok,
173    /// The operation failed with the given description.
174    Error(String),
175}
176
177
178// ── Span ──────────────────────────────────────────────────────────────────────
179
180/// A single unit of work in a distributed trace.
181#[derive(Debug)]
182pub struct Span {
183    /// Sampling and identity context.
184    pub context: SpanContext,
185    /// Human-readable operation name.
186    pub name: String,
187    /// Nanoseconds since UNIX epoch when the span started.
188    pub start_ns: u64,
189    /// Nanoseconds since UNIX epoch when the span ended, or `None` if still
190    /// in progress.
191    pub end_ns: Option<u64>,
192    /// Key-value attributes describing the operation.
193    pub attributes: HashMap<String, AttributeValue>,
194    /// Timestamped events that occurred during this span.
195    pub events: Vec<SpanEvent>,
196    /// Completion status of this span.
197    pub status: SpanStatus,
198}
199
200impl Span {
201    /// Creates a new in-progress span with the given context and name.
202    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    /// Sets a key-value attribute on this span.
215    pub fn set_attribute(&mut self, key: impl Into<String>, value: AttributeValue) {
216        self.attributes.insert(key.into(), value);
217    }
218
219    /// Appends a named event with optional attributes to this span.
220    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    /// Marks this span as finished by recording the current time as `end_ns`.
229    pub fn finish(&mut self) {
230        self.end_ns = Some(now_ns());
231    }
232
233    /// Returns the duration of this span in nanoseconds, or `None` if it has
234    /// not yet finished.
235    pub fn duration_ns(&self) -> Option<u64> {
236        self.end_ns.map(|end| end.saturating_sub(self.start_ns))
237    }
238}
239
240// ── SpanExporter ──────────────────────────────────────────────────────────────
241
242/// A sink that receives completed spans for storage or forwarding.
243pub trait SpanExporter: Send + Sync {
244    /// Called when a span is ready to be exported.
245    fn export(&self, span: Span);
246}
247
248// ── InMemoryExporter ─────────────────────────────────────────────────────────
249
250/// A [`SpanExporter`] that accumulates all exported spans in memory.
251///
252/// Useful for testing and introspection.
253#[derive(Debug, Default)]
254pub struct InMemoryExporter {
255    spans: Arc<Mutex<Vec<Span>>>,
256}
257
258impl InMemoryExporter {
259    /// Creates a new, empty in-memory exporter.
260    pub fn new() -> Self {
261        Self {
262            spans: Arc::new(Mutex::new(Vec::new())),
263        }
264    }
265
266    /// Returns a snapshot of all spans that have been exported so far.
267    pub fn get_spans(&self) -> Vec<String> {
268        // Return span names for easy inspection in tests.
269        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    /// Returns a clone of every exported [`Span`].
278    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    /// Returns the number of spans collected so far.
286    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// ── Tracer ────────────────────────────────────────────────────────────────────
302
303/// Creates and manages spans, routing finished spans to the configured exporter.
304#[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    /// Creates a tracer that will send completed spans to `exporter`.
317    pub fn new(exporter: Arc<dyn SpanExporter>) -> Self {
318        Self { exporter }
319    }
320
321    /// Starts a new root span with the given name.  The span is sampled by
322    /// default.
323    pub fn start_span(&self, name: &str) -> Span {
324        let ctx = SpanContext::new_root(true);
325        Span::new(ctx, name)
326    }
327
328    /// Starts a child span whose `parent_span_id` is set to the `parent`
329    /// span's `span_id`.
330    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    /// Finishes and exports `span`.
336    pub fn finish_span(&self, mut span: Span) {
337        span.finish();
338        self.exporter.export(span);
339    }
340}
341
342// ── MetricPoint ───────────────────────────────────────────────────────────────
343
344/// A single timestamped metric observation.
345#[derive(Debug, Clone)]
346pub struct MetricPoint {
347    /// Metric name (e.g. `"http_requests_total"`).
348    pub name: String,
349    /// Numeric value.
350    pub value: f64,
351    /// Prometheus-style label key-value pairs.
352    pub labels: HashMap<String, String>,
353    /// Nanoseconds since UNIX epoch.
354    pub timestamp_ns: u64,
355}
356
357// ── ObservabilityRegistry ─────────────────────────────────────────────────────
358
359/// Central registry that combines a [`Tracer`] with an in-process metrics
360/// store and Prometheus export.
361pub struct ObservabilityRegistry {
362    /// The tracer used to create spans.
363    pub tracer: Arc<Tracer>,
364    /// Accumulated metric observations.
365    pub metrics: Arc<Mutex<Vec<MetricPoint>>>,
366}
367
368impl ObservabilityRegistry {
369    /// Creates a new registry wrapping the given [`Tracer`].
370    pub fn new(tracer: Arc<Tracer>) -> Self {
371        Self {
372            tracer,
373            metrics: Arc::new(Mutex::new(Vec::new())),
374        }
375    }
376
377    /// Records a metric observation.
378    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    /// Serializes all recorded metrics in the
395    /// [Prometheus exposition format](https://prometheus.io/docs/instrumenting/exposition_formats/).
396    ///
397    /// Each metric is emitted as:
398    /// ```text
399    /// metric_name{label="value",...} <value> <timestamp_ms>
400    /// ```
401    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            // Build label set string.
407            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// ── Tests ─────────────────────────────────────────────────────────────────────
430
431#[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        // Duration should be non-negative.
445        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}