Skip to main content

ai_agents_observability/
manager.rs

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/// Opaque marker for scoping a later observation snapshot.
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub struct ObservabilityCursor {
29    ingested_events: u64,
30    dropped_events: u64,
31}
32
33/// Retained post-cursor events and the report derived from that exact event set.
34#[derive(Debug, Clone)]
35pub struct ScopedObservationSnapshot {
36    /// Redacted events still retained in the rolling metrics window.
37    pub events: Vec<ObservationEvent>,
38    /// Report derived only from the returned events.
39    pub report: ObservabilityReport,
40}
41
42/// Central collector that receives events, applies privacy rules, aggregates metrics, and exports reports.
43pub 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    /// Creates a shared manager with bounded event buffering.
57    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    /// Returns the immutable configuration used by this manager.
74    pub fn config(&self) -> &ObservabilityConfig {
75        &self.config
76    }
77
78    /// Starts a measured span for an LLM or tool wrapper.
79    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    /// Records hook-style lifecycle events that are not LLM or tool wrapper calls.
92    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    /// Records an event that should be finalized when its runtime branch resolves.
137    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    /// Finalizes all pending events for a runtime branch and ingests them normally.
151    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    /// Queues a completed event without blocking the observed call path.
195    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    /// Drains pending queued events into aggregation and raw buffers.
215    pub async fn flush(&self) -> Result<()> {
216        self.drain_pending();
217        Ok(())
218    }
219
220    /// Returns configured aggregate metrics after draining pending events.
221    pub fn get_metrics(&self) -> Vec<AggregatedMetrics> {
222        self.drain_pending();
223        self.aggregator.aggregate_configured()
224    }
225
226    /// Captures a cursor after draining all events already queued for ingestion.
227    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    /// Returns retained raw events after redaction and queue draining.
236    pub fn raw_events(&self) -> Vec<ObservationEvent> {
237        self.drain_pending();
238        self.raw_events.read().iter().cloned().collect()
239    }
240
241    /// Drains once and snapshots retained events after the cursor with a report from that exact event set.
242    /// Rolling-window eviction is preserved, while dropped events are reported as the counter delta since the cursor.
243    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    /// Builds the user-facing report from the current rolling event window.
256    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    /// Writes configured report, aggregate, raw event, and Prometheus files.
267    pub async fn export(&self) -> Result<ExportResult> {
268        export_observability(self).map_err(ObservabilityError::Io)
269    }
270
271    /// Returns the total number of events dropped by bounded buffers.
272    pub fn dropped_events(&self) -> u64 {
273        self.dropped_events.load(Ordering::Relaxed)
274    }
275
276    /// Returns the redactor used by wrappers for safe payload summaries.
277    pub fn redactor(&self) -> &Redactor {
278        &self.redactor
279    }
280
281    /// Converts a completed SpanGuard into an ObservationEvent.
282    //
283    // Keep this public signature stable for custom span integrations.
284    //
285    #[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    /// Drains queued events into the synchronous aggregation path.
321    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    /// Enriches, costs, redacts, aggregates, and optionally stores one event.
329    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    /// Applies token count switches before reports and cost estimates read usage.
363    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    /// Retains a redacted raw event when raw event export is enabled.
375    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    /// Renders current aggregate metrics in Prometheus text exposition format.
395    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    /// Returns true when a format is enabled in export.formats.
482    pub fn wants_format(&self, format: ExportFormat) -> bool {
483        self.config.export.formats.contains(&format)
484    }
485}
486
487/// Escapes label values for Prometheus text output.
488fn 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
500/// Builds the base event dimensions from the current span context.
501fn 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
518/// Resolves the language dimension by checking configured context paths in order.
519pub 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
534/// Looks up a top-level or dotted path in a JSON context map.
535fn 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
548/// Generates a session ID for observed runtime sessions that do not have one yet.
549pub 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}