//! Telemetry and distributed tracing. //! //! This module provides distributed tracing capabilities for RustyTorch++, //! including trace context propagation, span instrumentation, and export //! to various tracing backends. use crate::MonitoringResult; use parking_lot::RwLock; use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; use tracing::info; /// Trace ID generator using atomic counter + timestamp for uniqueness static TRACE_COUNTER: AtomicU64 = AtomicU64::new(0); /// Generate a unique trace ID pub fn generate_trace_id() -> String { let counter = TRACE_COUNTER.fetch_add(1, Ordering::SeqCst); let timestamp = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map(|d| d.as_nanos()) .unwrap_or(0); format!("{:016x}{:016x}", timestamp as u64, counter) } /// Generate a unique span ID pub fn generate_span_id() -> String { let counter = TRACE_COUNTER.fetch_add(1, Ordering::SeqCst); format!("{counter:016x}") } /// Span context for distributed tracing. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct SpanContext { /// Unique trace identifier (propagated across services) pub trace_id: String, /// Unique span identifier pub span_id: String, /// Parent span ID (if any) pub parent_span_id: Option, /// Sampling decision pub sampled: bool, /// Baggage items (propagated context) #[serde(default)] pub baggage: HashMap, } impl SpanContext { /// Create a new root span context pub fn new_root() -> Self { Self { trace_id: generate_trace_id(), span_id: generate_span_id(), parent_span_id: None, sampled: true, baggage: HashMap::new(), } } /// Create a child span context pub fn child(&self) -> Self { Self { trace_id: self.trace_id.clone(), span_id: generate_span_id(), parent_span_id: Some(self.span_id.clone()), sampled: self.sampled, baggage: self.baggage.clone(), } } /// Add baggage item pub fn with_baggage(mut self, key: &str, value: &str) -> Self { self.baggage.insert(key.to_string(), value.to_string()); self } /// Serialize to W3C Trace Context format for HTTP headers pub fn to_traceparent(&self) -> String { let flags = if self.sampled { "01" } else { "00" }; format!("00-{}-{}-{}", self.trace_id, self.span_id, flags) } /// Parse from W3C Trace Context format pub fn from_traceparent(header: &str) -> Option { let parts: Vec<&str> = header.split('-').collect(); if parts.len() != 4 || parts[0] != "00" { return None; } Some(Self { trace_id: parts[1].to_string(), span_id: generate_span_id(), // New span ID for this service parent_span_id: Some(parts[2].to_string()), sampled: parts[3] == "01", baggage: HashMap::new(), }) } } /// Trace configuration. #[derive(Debug, Clone)] pub struct TraceConfig { /// Service name for tracing pub service_name: String, /// Sample rate (0.0 to 1.0) pub sample_rate: f64, /// Enable console output pub enable_console: bool, /// Enable JSON output pub enable_json: bool, /// Jaeger endpoint (if enabled) pub jaeger_endpoint: Option, /// OTLP endpoint (if enabled) pub otlp_endpoint: Option, /// Maximum spans to buffer pub max_spans_buffer: usize, /// Batch export interval pub export_interval: std::time::Duration, } impl Default for TraceConfig { fn default() -> Self { Self { service_name: "rtx-service".to_string(), sample_rate: 1.0, enable_console: true, enable_json: false, jaeger_endpoint: None, otlp_endpoint: None, max_spans_buffer: 10000, export_interval: std::time::Duration::from_secs(5), } } } impl TraceConfig { /// Create a production configuration pub fn production(service_name: &str) -> Self { Self { service_name: service_name.to_string(), sample_rate: 0.1, // Sample 10% in production enable_console: false, enable_json: true, jaeger_endpoint: None, otlp_endpoint: None, max_spans_buffer: 50000, export_interval: std::time::Duration::from_secs(10), } } /// Create a development configuration pub fn development(service_name: &str) -> Self { Self { service_name: service_name.to_string(), sample_rate: 1.0, // Sample everything in dev enable_console: true, enable_json: false, jaeger_endpoint: None, otlp_endpoint: None, max_spans_buffer: 1000, export_interval: std::time::Duration::from_secs(1), } } /// Enable Jaeger export pub fn with_jaeger(mut self, endpoint: &str) -> Self { self.jaeger_endpoint = Some(endpoint.to_string()); self } /// Enable OTLP export pub fn with_otlp(mut self, endpoint: &str) -> Self { self.otlp_endpoint = Some(endpoint.to_string()); self } } /// Recorded span data for export #[derive(Debug, Clone, Serialize, Deserialize)] pub struct SpanData { /// Span context pub context: SpanContext, /// Operation name pub operation_name: String, /// Start timestamp pub start_time: chrono::DateTime, /// End timestamp (if completed) pub end_time: Option>, /// Duration in microseconds pub duration_us: Option, /// Span status pub status: SpanStatus, /// Span attributes pub attributes: HashMap, /// Span events pub events: Vec, } /// Span status #[derive(Debug, Clone, Serialize, Deserialize, Default)] pub enum SpanStatus { /// Unset status #[default] Unset, /// Operation completed successfully Ok, /// Operation failed Error(String), } /// Span event (annotation) #[derive(Debug, Clone, Serialize, Deserialize)] pub struct SpanEvent { /// Event name pub name: String, /// Event timestamp pub timestamp: chrono::DateTime, /// Event attributes pub attributes: HashMap, } /// Telemetry manager for distributed tracing. #[derive(Debug)] pub struct TelemetryManager { config: TraceConfig, /// Buffer for completed spans spans_buffer: Arc>>, /// Active spans by trace ID active_spans: Arc>>, /// Whether the manager is enabled enabled: bool, } impl TelemetryManager { /// Create a new telemetry manager pub async fn new(config: TraceConfig) -> MonitoringResult { let manager = Self { config: config.clone(), spans_buffer: Arc::new(RwLock::new(Vec::new())), active_spans: Arc::new(RwLock::new(HashMap::new())), enabled: true, }; // Initialize tracing subscriber manager.init_tracing()?; info!( service = %config.service_name, sample_rate = %config.sample_rate, "Telemetry manager initialized" ); Ok(manager) } /// Initialize the tracing subscriber fn init_tracing(&self) -> MonitoringResult<()> { use tracing_subscriber::{EnvFilter, fmt, prelude::*}; let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")); let subscriber = tracing_subscriber::registry().with(filter); if self.config.enable_json { let json_layer = fmt::layer() .json() .with_target(true) .with_thread_ids(true) .with_file(true) .with_line_number(true); let _ = subscriber.with(json_layer).try_init(); } else if self.config.enable_console { let fmt_layer = fmt::layer() .with_target(true) .with_thread_ids(false) .with_file(false) .with_line_number(false); let _ = subscriber.with(fmt_layer).try_init(); } Ok(()) } /// Start a new span pub fn start_span(&self, operation_name: &str, parent: Option<&SpanContext>) -> SpanData { let context = match parent { Some(p) => p.child(), None => SpanContext::new_root(), }; let span = SpanData { context: context.clone(), operation_name: operation_name.to_string(), start_time: chrono::Utc::now(), end_time: None, duration_us: None, status: SpanStatus::Unset, attributes: HashMap::new(), events: Vec::new(), }; // Track active span self.active_spans .write() .insert(context.span_id.clone(), span.clone()); span } /// End a span and record it pub fn end_span(&self, mut span: SpanData, status: SpanStatus) { let end_time = chrono::Utc::now(); span.end_time = Some(end_time); span.duration_us = Some((end_time - span.start_time).num_microseconds().unwrap_or(0) as u64); span.status = status; // Remove from active spans self.active_spans.write().remove(&span.context.span_id); // Add to buffer (with size limit) let mut buffer = self.spans_buffer.write(); if buffer.len() >= self.config.max_spans_buffer { buffer.remove(0); // Remove oldest } buffer.push(span); } /// Add an event to a span pub fn add_span_event(&self, span_id: &str, name: &str, attributes: HashMap) { if let Some(span) = self.active_spans.write().get_mut(span_id) { span.events.push(SpanEvent { name: name.to_string(), timestamp: chrono::Utc::now(), attributes, }); } } /// Set span attribute pub fn set_span_attribute(&self, span_id: &str, key: &str, value: &str) { if let Some(span) = self.active_spans.write().get_mut(span_id) { span.attributes.insert(key.to_string(), value.to_string()); } } /// Get completed spans for export pub fn drain_spans(&self) -> Vec { let mut buffer = self.spans_buffer.write(); std::mem::take(&mut *buffer) } /// Get active span count pub fn active_span_count(&self) -> usize { self.active_spans.read().len() } /// Get buffered span count pub fn buffered_span_count(&self) -> usize { self.spans_buffer.read().len() } /// Export spans to Jaeger format (JSON) pub fn export_jaeger_json(&self) -> String { let spans = self.spans_buffer.read(); serde_json::to_string_pretty(&*spans).unwrap_or_else(|_| "[]".to_string()) } /// Check if telemetry is enabled pub fn is_enabled(&self) -> bool { self.enabled } /// Get the service name pub fn service_name(&self) -> &str { &self.config.service_name } } /// RAII guard for automatic span management pub struct SpanGuard<'a> { manager: &'a TelemetryManager, span: Option, } impl<'a> SpanGuard<'a> { /// Create a new span guard pub fn new( manager: &'a TelemetryManager, operation_name: &str, parent: Option<&SpanContext>, ) -> Self { let span = manager.start_span(operation_name, parent); Self { manager, span: Some(span), } } /// Get the span context pub fn context(&self) -> Option<&SpanContext> { self.span.as_ref().map(|s| &s.context) } /// Add an attribute to the span pub fn set_attribute(&self, key: &str, value: &str) { if let Some(span) = &self.span { self.manager .set_span_attribute(&span.context.span_id, key, value); } } /// Mark the span as successful pub fn set_ok(mut self) { if let Some(span) = self.span.take() { self.manager.end_span(span, SpanStatus::Ok); } } /// Mark the span as failed pub fn set_error(mut self, error: &str) { if let Some(span) = self.span.take() { self.manager .end_span(span, SpanStatus::Error(error.to_string())); } } } impl Drop for SpanGuard<'_> { fn drop(&mut self) { if let Some(span) = self.span.take() { // If not explicitly ended, mark as unset self.manager.end_span(span, SpanStatus::Unset); } } } /// Instrumentation macros for common operations #[macro_export] macro_rules! instrument_inference { ($telemetry:expr, $model_name:expr, $batch_size:expr) => {{ let guard = $crate::telemetry::SpanGuard::new($telemetry, "inference", None); guard.set_attribute("model", $model_name); guard.set_attribute("batch_size", &$batch_size.to_string()); guard }}; } #[cfg(test)] mod tests { use super::*; #[test] fn test_generate_trace_id() { let id1 = generate_trace_id(); let id2 = generate_trace_id(); assert_ne!(id1, id2); assert_eq!(id1.len(), 32); } #[test] fn test_span_context_new_root() { let ctx = SpanContext::new_root(); assert!(!ctx.trace_id.is_empty()); assert!(!ctx.span_id.is_empty()); assert!(ctx.parent_span_id.is_none()); assert!(ctx.sampled); } #[test] fn test_span_context_child() { let parent = SpanContext::new_root(); let child = parent.child(); assert_eq!(child.trace_id, parent.trace_id); assert_ne!(child.span_id, parent.span_id); assert_eq!(child.parent_span_id, Some(parent.span_id)); } #[test] fn test_traceparent_roundtrip() { let ctx = SpanContext::new_root(); let header = ctx.to_traceparent(); let parsed = SpanContext::from_traceparent(&header); assert!(parsed.is_some()); let parsed = parsed.unwrap(); assert_eq!(parsed.trace_id, ctx.trace_id); assert_eq!(parsed.parent_span_id, Some(ctx.span_id)); } #[test] fn test_trace_config_default() { let config = TraceConfig::default(); assert_eq!(config.sample_rate, 1.0); assert!(config.enable_console); } #[test] fn test_trace_config_production() { let config = TraceConfig::production("test-service"); assert_eq!(config.sample_rate, 0.1); assert!(!config.enable_console); assert!(config.enable_json); } #[tokio::test] async fn test_telemetry_manager_creation() { let config = TraceConfig::default(); let manager = TelemetryManager::new(config).await.unwrap(); assert!(manager.is_enabled()); } #[tokio::test] async fn test_span_lifecycle() { let config = TraceConfig::default(); let manager = TelemetryManager::new(config).await.unwrap(); let span = manager.start_span("test_operation", None); assert_eq!(manager.active_span_count(), 1); manager.end_span(span, SpanStatus::Ok); assert_eq!(manager.active_span_count(), 0); assert_eq!(manager.buffered_span_count(), 1); } #[tokio::test] async fn test_span_guard() { let config = TraceConfig::default(); let manager = TelemetryManager::new(config).await.unwrap(); { let guard = SpanGuard::new(&manager, "guarded_op", None); guard.set_attribute("key", "value"); assert_eq!(manager.active_span_count(), 1); guard.set_ok(); } assert_eq!(manager.active_span_count(), 0); assert_eq!(manager.buffered_span_count(), 1); } }