1use crate::aggregator::{
2 AggregatedMetrics, MetricsAggregator, aggregate_events, enrich_dimensions,
3};
4use crate::config::{AggregationDimension, ExportFormat, ObservabilityConfig, UnknownPricePolicy};
5use crate::context::{SpanContext, current_observation_context};
6use crate::cost::CostEstimator;
7use crate::event::{
8 CostEstimate, EventStatus, EventType, ObservationEvent, ObservationPurpose,
9 ObservationTokenUsage,
10};
11use crate::export::{ExportResult, export_observability};
12use crate::redaction::Redactor;
13use crate::report::{ObservabilityReport, generate_report};
14use crate::span::SpanGuard;
15use crate::{ObservabilityError, Result};
16use chrono::Utc;
17use parking_lot::{Mutex, RwLock};
18use serde_json::Value;
19use std::collections::{HashMap, VecDeque};
20use std::sync::Arc;
21use std::sync::atomic::{AtomicU64, Ordering};
22use std::time::Duration;
23use tokio::sync::mpsc;
24use uuid::Uuid;
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub struct ObservabilityCursor {
29 ingested_events: u64,
30 dropped_events: u64,
31}
32
33#[derive(Debug, Clone)]
35pub struct ScopedObservationSnapshot {
36 pub events: Vec<ObservationEvent>,
38 pub report: ObservabilityReport,
40}
41
42pub struct ObservabilityManager {
44 config: ObservabilityConfig,
45 sender: mpsc::Sender<ObservationEvent>,
46 receiver: Mutex<mpsc::Receiver<ObservationEvent>>,
47 raw_events: RwLock<VecDeque<ObservationEvent>>,
48 pending_branch_events: RwLock<HashMap<String, Vec<ObservationEvent>>>,
49 aggregator: MetricsAggregator,
50 cost_estimator: CostEstimator,
51 redactor: Redactor,
52 dropped_events: AtomicU64,
53}
54
55impl ObservabilityManager {
56 pub fn new(config: ObservabilityConfig) -> Arc<Self> {
58 let _ = config.validate();
59 let (sender, receiver) = mpsc::channel(config.buffer.event_buffer.max(1));
60 Arc::new(Self {
61 cost_estimator: CostEstimator::new(config.cost.clone()),
62 redactor: Redactor::new(config.privacy.clone()),
63 aggregator: MetricsAggregator::new(config.aggregation.clone()),
64 sender,
65 receiver: Mutex::new(receiver),
66 raw_events: RwLock::new(VecDeque::new()),
67 pending_branch_events: RwLock::new(HashMap::new()),
68 dropped_events: AtomicU64::new(0),
69 config,
70 })
71 }
72
73 pub fn config(&self) -> &ObservabilityConfig {
75 &self.config
76 }
77
78 pub fn start_span(
80 self: &Arc<Self>,
81 event_type: EventType,
82 purpose: ObservationPurpose,
83 ) -> SpanGuard {
84 let mut context = current_observation_context()
85 .map(|ctx| ctx.child())
86 .unwrap_or_else(|| SpanContext::new_root("unknown"));
87 context.purpose = purpose;
88 SpanGuard::new(Arc::clone(self), context, event_type)
89 }
90
91 pub fn record_lifecycle_event(
93 &self,
94 event_type: EventType,
95 purpose: ObservationPurpose,
96 status: EventStatus,
97 duration_ms: u64,
98 tags: HashMap<String, String>,
99 payload: Option<Value>,
100 ) {
101 let context = current_observation_context()
102 .map(|ctx| ctx.child())
103 .unwrap_or_else(|| SpanContext::new_root("unknown"));
104 let mut dimensions = context_dimension_map(&context);
105 for (key, value) in &tags {
106 if key.starts_with("runtime.") {
107 dimensions.insert(key.clone(), value.clone());
108 if let Some(short_key) = key.strip_prefix("runtime.") {
109 dimensions.insert(short_key.to_string(), value.clone());
110 }
111 }
112 }
113 let event = ObservationEvent {
114 trace_id: context.trace_id,
115 span_id: context.span_id,
116 parent_span_id: context.parent_span_id,
117 turn_id: context.turn_id,
118 agent_id: context.agent_id,
119 actor_id: context.actor_id,
120 session_id: context.session_id,
121 event_type,
122 purpose,
123 status,
124 timestamp: Utc::now(),
125 duration_ms,
126 tokens: None,
127 cost: None,
128 error: None,
129 dimensions,
130 tags,
131 payload,
132 };
133 self.record_event(event);
134 }
135
136 pub fn record_pending_event(&self, branch_id: impl Into<String>, event: ObservationEvent) {
138 if !self.config.enabled {
139 return;
140 }
141 let mut pending = self.pending_branch_events.write();
142 let pending_count: usize = pending.values().map(Vec::len).sum();
143 if pending_count >= self.config.buffer.pending_branch_event_limit {
144 self.dropped_events.fetch_add(1, Ordering::Relaxed);
145 return;
146 }
147 pending.entry(branch_id.into()).or_default().push(event);
148 }
149
150 pub fn finalize_pending_branch(
152 &self,
153 branch_id: &str,
154 branch_status: impl Into<String>,
155 winner: bool,
156 extra_tags: HashMap<String, String>,
157 ) -> usize {
158 let mut events = self
159 .pending_branch_events
160 .write()
161 .remove(branch_id)
162 .unwrap_or_default();
163 let status = branch_status.into();
164 let count = events.len();
165 for event in &mut events {
166 event
167 .tags
168 .insert("runtime.branch_status".to_string(), status.clone());
169 event
170 .tags
171 .insert("runtime.winner".to_string(), winner.to_string());
172 event.tags.insert("winner".to_string(), winner.to_string());
173 event
174 .dimensions
175 .insert("branch_status".to_string(), status.clone());
176 event
177 .dimensions
178 .insert("runtime.branch_status".to_string(), status.clone());
179 event
180 .dimensions
181 .insert("runtime.winner".to_string(), winner.to_string());
182 event
183 .dimensions
184 .insert("winner".to_string(), winner.to_string());
185 for (key, value) in &extra_tags {
186 event.tags.insert(key.clone(), value.clone());
187 event.dimensions.insert(key.clone(), value.clone());
188 }
189 self.record_event(event.clone());
190 }
191 count
192 }
193
194 pub fn record_event(&self, event: ObservationEvent) {
196 if !self.config.enabled {
197 return;
198 }
199 match self.sender.try_send(event) {
200 Ok(()) => {}
201 Err(mpsc::error::TrySendError::Full(event)) => {
202 if self.config.buffer.drop_on_full {
203 self.dropped_events.fetch_add(1, Ordering::Relaxed);
204 } else {
205 self.ingest_event(event);
206 }
207 }
208 Err(mpsc::error::TrySendError::Closed(event)) => {
209 self.ingest_event(event);
210 }
211 }
212 }
213
214 pub async fn flush(&self) -> Result<()> {
216 self.drain_pending();
217 Ok(())
218 }
219
220 pub fn get_metrics(&self) -> Vec<AggregatedMetrics> {
222 self.drain_pending();
223 self.aggregator.aggregate_configured()
224 }
225
226 pub fn event_cursor(&self) -> ObservabilityCursor {
228 self.drain_pending();
229 ObservabilityCursor {
230 ingested_events: self.aggregator.cursor(),
231 dropped_events: self.dropped_events(),
232 }
233 }
234
235 pub fn raw_events(&self) -> Vec<ObservationEvent> {
237 self.drain_pending();
238 self.raw_events.read().iter().cloned().collect()
239 }
240
241 pub fn snapshot_since(&self, cursor: ObservabilityCursor) -> ScopedObservationSnapshot {
244 self.drain_pending();
245 let events = self.aggregator.events_since(cursor.ingested_events);
246 let dropped_events = self.dropped_events().saturating_sub(cursor.dropped_events);
247 let report = generate_report(
248 &events,
249 aggregate_events(&events, &self.config.aggregation.dimensions),
250 dropped_events,
251 );
252 ScopedObservationSnapshot { events, report }
253 }
254
255 pub fn generate_report(&self) -> ObservabilityReport {
257 self.drain_pending();
258 let events = self.aggregator.events();
259 generate_report(
260 &events,
261 aggregate_events(&events, &self.config.aggregation.dimensions),
262 self.dropped_events(),
263 )
264 }
265
266 pub async fn export(&self) -> Result<ExportResult> {
268 export_observability(self).map_err(ObservabilityError::Io)
269 }
270
271 pub fn dropped_events(&self) -> u64 {
273 self.dropped_events.load(Ordering::Relaxed)
274 }
275
276 pub fn redactor(&self) -> &Redactor {
278 &self.redactor
279 }
280
281 #[allow(clippy::too_many_arguments)]
286 pub fn build_event_from_span(
287 &self,
288 context: SpanContext,
289 event_type: EventType,
290 duration: Duration,
291 status: EventStatus,
292 tokens: Option<crate::event::ObservationTokenUsage>,
293 error: Option<crate::event::ObservationError>,
294 tags: HashMap<String, String>,
295 payload: Option<Value>,
296 ) -> ObservationEvent {
297 let dimensions = context_dimension_map(&context);
298 ObservationEvent {
299 trace_id: context.trace_id,
300 span_id: context.span_id,
301 parent_span_id: context.parent_span_id,
302 turn_id: context.turn_id,
303 agent_id: context.agent_id,
304 actor_id: context.actor_id,
305 session_id: context.session_id,
306 event_type,
307 purpose: context.purpose,
308 status,
309 timestamp: Utc::now(),
310 duration_ms: duration.as_millis() as u64,
311 tokens,
312 cost: None::<CostEstimate>,
313 error,
314 dimensions,
315 tags,
316 payload,
317 }
318 }
319
320 fn drain_pending(&self) {
322 let mut receiver = self.receiver.lock();
323 while let Ok(event) = receiver.try_recv() {
324 self.ingest_event(event);
325 }
326 }
327
328 fn ingest_event(&self, mut event: ObservationEvent) {
330 enrich_dimensions(&mut event);
331 event.tokens = event
332 .tokens
333 .take()
334 .map(|tokens| self.apply_token_config(tokens));
335 if event.cost.is_none() {
336 let (provider, model) = match &event.event_type {
337 EventType::LlmCall {
338 provider, model, ..
339 } => (Some(provider.as_str()), Some(model.as_str())),
340 _ => (None, None),
341 };
342 event.cost = self
343 .cost_estimator
344 .estimate(provider, model, event.tokens.as_ref());
345 if matches!(
346 self.config.cost.unknown_price_policy,
347 UnknownPricePolicy::Error
348 ) && event.tokens.is_some()
349 && event.cost.is_none()
350 && matches!(&event.event_type, EventType::LlmCall { .. })
351 {
352 event
353 .tags
354 .insert("cost_error".to_string(), "unknown_price".to_string());
355 }
356 }
357 let event = self.redactor.redact_event(event);
358 self.aggregator.record(event.clone());
359 self.store_raw_event(event);
360 }
361
362 fn apply_token_config(&self, mut tokens: ObservationTokenUsage) -> ObservationTokenUsage {
364 if !self.config.tokens.count_input {
365 tokens.input_tokens = 0;
366 }
367 if !self.config.tokens.count_output {
368 tokens.output_tokens = 0;
369 }
370 tokens.total_tokens = tokens.input_tokens + tokens.output_tokens;
371 tokens
372 }
373
374 fn store_raw_event(&self, event: ObservationEvent) {
376 if !self.config.export.write_raw_events {
377 return;
378 }
379 if self.config.buffer.raw_event_limit == 0 {
380 self.dropped_events.fetch_add(1, Ordering::Relaxed);
381 return;
382 }
383 let mut raw_events = self.raw_events.write();
384 if raw_events.len() >= self.config.buffer.raw_event_limit {
385 if self.config.buffer.drop_on_full {
386 self.dropped_events.fetch_add(1, Ordering::Relaxed);
387 return;
388 }
389 raw_events.pop_front();
390 }
391 raw_events.push_back(event);
392 }
393
394 pub fn render_prometheus(&self) -> String {
396 let report = self.generate_report();
397 let events = self.aggregator.events();
398 let llm_events: Vec<_> = events
399 .iter()
400 .filter(|event| matches!(&event.event_type, EventType::LlmCall { .. }))
401 .cloned()
402 .collect();
403 let tool_events: Vec<_> = events
404 .iter()
405 .filter(|event| matches!(&event.event_type, EventType::ToolCall { .. }))
406 .cloned()
407 .collect();
408 let by_model_purpose = aggregate_events(
409 &llm_events,
410 &[AggregationDimension::Model, AggregationDimension::Purpose],
411 );
412 let by_tool = aggregate_events(&tool_events, &[AggregationDimension::Tool]);
413 let mut output = String::new();
414 output.push_str(
415 "# HELP ai_agents_observation_events_total Total recorded observation events\n",
416 );
417 output.push_str("# TYPE ai_agents_observation_events_total counter\n");
418 output.push_str(&format!(
419 "ai_agents_observation_events_total {}\n",
420 report.summary.total_events
421 ));
422 output.push_str("# HELP ai_agents_observation_errors_total Total observation events with error status\n");
423 output.push_str("# TYPE ai_agents_observation_errors_total counter\n");
424 output.push_str(&format!(
425 "ai_agents_observation_errors_total {}\n",
426 report.summary.total_errors
427 ));
428 output.push_str(
429 "# HELP ai_agents_observation_cost_usd_total Estimated total LLM cost in USD\n",
430 );
431 output.push_str("# TYPE ai_agents_observation_cost_usd_total counter\n");
432 output.push_str(&format!(
433 "ai_agents_observation_cost_usd_total {:.8}\n",
434 report.summary.total_cost_usd
435 ));
436 output.push_str("# HELP ai_agents_observation_tokens_total Total observed LLM tokens\n");
437 output.push_str("# TYPE ai_agents_observation_tokens_total counter\n");
438 output.push_str(&format!(
439 "ai_agents_observation_tokens_total {}\n",
440 report.summary.total_tokens
441 ));
442 output.push_str("# HELP ai_agents_llm_calls_total LLM calls grouped by safe labels\n");
443 output.push_str("# TYPE ai_agents_llm_calls_total counter\n");
444 for metric in by_model_purpose {
445 let model = metric
446 .dimensions
447 .get("model")
448 .map(String::as_str)
449 .unwrap_or("unknown");
450 let purpose = metric
451 .dimensions
452 .get("purpose")
453 .map(String::as_str)
454 .unwrap_or("unknown");
455 output.push_str(&format!(
456 "ai_agents_llm_calls_total{{model=\"{}\",purpose=\"{}\"}} {}\n",
457 prometheus_label(model),
458 prometheus_label(purpose),
459 metric.count
460 ));
461 }
462 output.push_str("# HELP ai_agents_tool_calls_total Tool calls grouped by tool ID\n");
463 output.push_str("# TYPE ai_agents_tool_calls_total counter\n");
464 for metric in by_tool {
465 let tool = metric
466 .dimensions
467 .get("tool")
468 .map(String::as_str)
469 .unwrap_or("unknown");
470 if tool != "unknown" {
471 output.push_str(&format!(
472 "ai_agents_tool_calls_total{{tool=\"{}\"}} {}\n",
473 prometheus_label(tool),
474 metric.count
475 ));
476 }
477 }
478 output
479 }
480
481 pub fn wants_format(&self, format: ExportFormat) -> bool {
483 self.config.export.formats.contains(&format)
484 }
485}
486
487fn prometheus_label(value: &str) -> String {
489 value
490 .chars()
491 .flat_map(|ch| match ch {
492 '\\' => "\\\\".chars().collect::<Vec<_>>(),
493 '"' => "\\\"".chars().collect::<Vec<_>>(),
494 '\n' | '\r' | '\t' => "_".chars().collect::<Vec<_>>(),
495 _ => vec![ch],
496 })
497 .collect()
498}
499
500fn context_dimension_map(context: &SpanContext) -> HashMap<String, String> {
502 let mut dimensions = HashMap::new();
503 dimensions.insert("agent".to_string(), context.agent_id.clone());
504 dimensions.insert("purpose".to_string(), context.purpose.as_label());
505 if let Some(actor) = &context.actor_id {
506 dimensions.insert("actor".to_string(), actor.clone());
507 }
508 if let Some(state) = &context.state {
509 dimensions.insert("state".to_string(), state.clone());
510 }
511 if let Some(language) = &context.language {
512 dimensions.insert("language".to_string(), language.clone());
513 }
514 dimensions.extend(context.tags.clone());
515 dimensions
516}
517
518pub fn resolve_language_from_context(
520 config: &ObservabilityConfig,
521 context: &HashMap<String, Value>,
522) -> String {
523 for path in &config.language.paths {
524 if let Some(value) = get_dotted(context, path)
525 && let Some(language) = value.as_str()
526 && !language.trim().is_empty()
527 {
528 return language.to_string();
529 }
530 }
531 config.language.fallback.clone()
532}
533
534fn get_dotted<'a>(context: &'a HashMap<String, Value>, path: &str) -> Option<&'a Value> {
536 if let Some(value) = context.get(path) {
537 return Some(value);
538 }
539 let mut parts = path.split('.');
540 let first = parts.next()?;
541 let mut current = context.get(first)?;
542 for part in parts {
543 current = current.get(part)?;
544 }
545 Some(current)
546}
547
548pub fn new_session_id() -> String {
550 Uuid::new_v4().to_string()
551}
552
553#[cfg(test)]
554mod tests {
555 use super::*;
556 use crate::event::{ObservationTokenUsage, TokenUsageSource};
557
558 fn test_event() -> ObservationEvent {
559 ObservationEvent {
560 trace_id: "trace".to_string(),
561 span_id: Uuid::new_v4().to_string(),
562 parent_span_id: None,
563 turn_id: "turn".to_string(),
564 agent_id: "agent".to_string(),
565 actor_id: None,
566 session_id: None,
567 event_type: EventType::LlmCall {
568 provider: "openai".to_string(),
569 model: "test".to_string(),
570 alias: Some("default".to_string()),
571 streaming: false,
572 },
573 purpose: ObservationPurpose::MainResponse,
574 status: EventStatus::Success,
575 timestamp: Utc::now(),
576 duration_ms: 10,
577 tokens: Some(ObservationTokenUsage::new(
578 100,
579 25,
580 TokenUsageSource::Provider,
581 )),
582 cost: None,
583 error: None,
584 dimensions: HashMap::new(),
585 tags: HashMap::new(),
586 payload: None,
587 }
588 }
589
590 #[test]
591 fn token_count_flags_are_applied_before_report() {
592 let config = ObservabilityConfig {
593 enabled: true,
594 tokens: crate::config::TokenConfig {
595 count_input: false,
596 count_output: true,
597 ..Default::default()
598 },
599 cost: crate::config::CostConfig {
600 enabled: false,
601 ..Default::default()
602 },
603 ..Default::default()
604 };
605 let manager = ObservabilityManager::new(config);
606 manager.record_event(test_event());
607
608 let report = manager.generate_report();
609 assert_eq!(report.token_breakdown.total_input, 0);
610 assert_eq!(report.token_breakdown.total_output, 25);
611 assert_eq!(report.token_breakdown.total_tokens, 25);
612 }
613
614 #[test]
615 fn pending_branch_event_is_hidden_until_finalized() {
616 let config = ObservabilityConfig {
617 enabled: true,
618 export: crate::config::ExportConfig {
619 write_raw_events: true,
620 ..Default::default()
621 },
622 ..Default::default()
623 };
624 let manager = ObservabilityManager::new(config);
625 manager.record_pending_event("branch", test_event());
626
627 let mut tags = HashMap::new();
628 tags.insert("runtime.speculative".to_string(), "true".to_string());
629 tags.insert("speculative".to_string(), "true".to_string());
630
631 assert_eq!(manager.generate_report().summary.total_events, 0);
632 manager.finalize_pending_branch("branch", "discarded", false, tags);
633 let report = manager.generate_report();
634 assert_eq!(report.summary.total_events, 1);
635 assert_eq!(
636 manager.raw_events()[0].dimensions.get("branch_status"),
637 Some(&"discarded".to_string())
638 );
639 assert_eq!(
640 manager.raw_events()[0].dimensions.get("runtime.winner"),
641 Some(&"false".to_string())
642 );
643 assert_eq!(
644 manager.raw_events()[0].dimensions.get("speculative"),
645 Some(&"true".to_string())
646 );
647 assert_eq!(
648 manager.raw_events()[0]
649 .dimensions
650 .get("runtime.speculative"),
651 Some(&"true".to_string())
652 );
653 }
654
655 #[test]
656 fn pending_branch_events_are_bounded() {
657 let config = ObservabilityConfig {
658 enabled: true,
659 buffer: crate::config::BufferConfig {
660 pending_branch_event_limit: 1,
661 ..Default::default()
662 },
663 ..Default::default()
664 };
665 let manager = ObservabilityManager::new(config);
666 manager.record_pending_event("branch-a", test_event());
667 manager.record_pending_event("branch-b", test_event());
668
669 manager.finalize_pending_branch("branch-a", "committed", true, HashMap::new());
670 manager.finalize_pending_branch("branch-b", "committed", true, HashMap::new());
671 let report = manager.generate_report();
672 assert_eq!(report.summary.total_events, 1);
673 assert_eq!(report.dropped_events, 1);
674 }
675
676 #[test]
677 fn snapshot_since_scopes_report_and_events_to_cursor() {
678 let config = ObservabilityConfig {
679 enabled: true,
680 aggregation: crate::config::AggregationConfig {
681 dimensions: vec![AggregationDimension::Purpose],
682 window_size: 4,
683 ..Default::default()
684 },
685 ..Default::default()
686 };
687 let manager = ObservabilityManager::new(config);
688 let mut first = test_event();
689 first.turn_id = "turn-1".to_string();
690 manager.record_event(first);
691 let cursor = manager.event_cursor();
692
693 let mut second = test_event();
694 second.turn_id = "turn-2".to_string();
695 manager.record_event(second);
696
697 let snapshot = manager.snapshot_since(cursor);
698 assert_eq!(snapshot.events.len(), 1);
699 assert_eq!(snapshot.events[0].turn_id, "turn-2");
700 assert_eq!(
701 snapshot.report.summary.total_events,
702 snapshot.events.len() as u64
703 );
704 assert_eq!(
705 snapshot
706 .report
707 .configured
708 .iter()
709 .map(|metric| metric.count)
710 .sum::<u64>(),
711 snapshot.events.len() as u64
712 );
713 assert_eq!(manager.generate_report().summary.total_events, 2);
714 }
715
716 #[test]
717 fn snapshot_since_preserves_rolling_window_eviction() {
718 let config = ObservabilityConfig {
719 enabled: true,
720 aggregation: crate::config::AggregationConfig {
721 window_size: 2,
722 ..Default::default()
723 },
724 ..Default::default()
725 };
726 let manager = ObservabilityManager::new(config);
727 let cursor = manager.event_cursor();
728 for turn_id in ["turn-1", "turn-2", "turn-3"] {
729 let mut event = test_event();
730 event.turn_id = turn_id.to_string();
731 manager.record_event(event);
732 }
733
734 let snapshot = manager.snapshot_since(cursor);
735 assert_eq!(snapshot.events.len(), 2);
736 assert_eq!(snapshot.events[0].turn_id, "turn-2");
737 assert_eq!(snapshot.events[1].turn_id, "turn-3");
738 assert_eq!(snapshot.report.summary.total_events, 2);
739 }
740
741 #[test]
742 fn snapshot_since_exposes_only_the_dropped_event_delta() {
743 let config = ObservabilityConfig {
744 enabled: true,
745 buffer: crate::config::BufferConfig {
746 event_buffer: 1,
747 drop_on_full: true,
748 ..Default::default()
749 },
750 ..Default::default()
751 };
752 let manager = ObservabilityManager::new(config);
753 manager.record_event(test_event());
754 manager.record_event(test_event());
755 let cursor = manager.event_cursor();
756
757 manager.record_event(test_event());
758 manager.record_event(test_event());
759
760 let snapshot = manager.snapshot_since(cursor);
761 assert_eq!(snapshot.report.summary.total_events, 1);
762 assert_eq!(snapshot.events.len(), 1);
763 assert_eq!(snapshot.report.dropped_events, 1);
764 assert_eq!(manager.generate_report().dropped_events, 2);
765 }
766}