Skip to main content

valence_telemetry/
recording.rs

1//! In-memory [`TelemetrySink`] for tests.
2
3use std::collections::HashMap;
4use std::sync::{Arc, Mutex};
5
6use crate::TelemetrySink;
7
8/// Captured counter increment.
9#[derive(Debug, Clone, PartialEq, Eq)]
10pub struct RecordedCounter {
11    /// Metric name.
12    pub name: String,
13    /// Label dimensions at emission time.
14    pub labels: Vec<(String, String)>,
15    /// Increment applied to the counter.
16    pub delta: u64,
17}
18
19/// Captured gauge sample.
20#[derive(Debug, Clone, PartialEq)]
21pub struct RecordedGauge {
22    /// Metric name.
23    pub name: String,
24    /// Label dimensions at emission time.
25    pub labels: Vec<(String, String)>,
26    /// Observed gauge value.
27    pub value: f64,
28}
29
30/// Captured structured log event.
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct RecordedEvent {
33    /// Event schema identifier.
34    pub schema: String,
35    /// Flattened field payload.
36    pub fields: HashMap<String, String>,
37}
38
39#[derive(Debug, Default)]
40struct Inner {
41    counters: Vec<RecordedCounter>,
42    gauges: Vec<RecordedGauge>,
43    events: Vec<RecordedEvent>,
44}
45
46/// Append-only in-memory sink for assertions in unit and integration tests.
47#[derive(Debug, Clone)]
48pub struct RecordingSink {
49    inner: Arc<Mutex<Inner>>,
50}
51
52impl RecordingSink {
53    /// Create an empty in-memory recording sink.
54    pub fn new() -> Self {
55        Self {
56            inner: Arc::new(Mutex::new(Inner::default())),
57        }
58    }
59
60    fn lock_inner(&self) -> std::sync::MutexGuard<'_, Inner> {
61        self.inner
62            .lock()
63            .unwrap_or_else(std::sync::PoisonError::into_inner)
64    }
65
66    /// Return all recorded counter increments.
67    pub fn counters(&self) -> Vec<RecordedCounter> {
68        self.lock_inner().counters.clone()
69    }
70
71    /// Return all recorded structured events.
72    pub fn events(&self) -> Vec<RecordedEvent> {
73        self.lock_inner().events.clone()
74    }
75
76    /// Return all recorded gauge samples.
77    pub fn gauges(&self) -> Vec<RecordedGauge> {
78        self.lock_inner().gauges.clone()
79    }
80
81    /// Counters matching name and label subset (test helper).
82    pub fn recorded_counters_matching(
83        &self,
84        name: &str,
85        labels: &[(&str, &str)],
86    ) -> Vec<RecordedCounter> {
87        self.counters()
88            .into_iter()
89            .filter(|c| {
90                c.name == name
91                    && labels
92                        .iter()
93                        .all(|(k, v)| c.labels.iter().any(|(lk, lv)| lk == k && lv == v))
94            })
95            .collect()
96    }
97
98    /// Events matching schema name (test helper).
99    pub fn recorded_events_for(&self, schema: &str) -> Vec<RecordedEvent> {
100        self.events()
101            .into_iter()
102            .filter(|e| e.schema == schema)
103            .collect()
104    }
105}
106
107impl TelemetrySink for RecordingSink {
108    fn record_counter(&self, name: &str, labels: &[(&str, &str)], delta: u64) {
109        let labels: Vec<(String, String)> = labels
110            .iter()
111            .map(|(k, v)| (k.to_string(), v.to_string()))
112            .collect();
113        let mut inner = self.lock_inner();
114        inner.counters.push(RecordedCounter {
115            name: name.to_string(),
116            labels,
117            delta,
118        });
119    }
120
121    fn record_gauge(&self, name: &str, labels: &[(&str, &str)], value: f64) {
122        let labels: Vec<(String, String)> = labels
123            .iter()
124            .map(|(k, v)| (k.to_string(), v.to_string()))
125            .collect();
126        let mut inner = self.lock_inner();
127        inner.gauges.push(RecordedGauge {
128            name: name.to_string(),
129            labels,
130            value,
131        });
132    }
133
134    fn log_event(&self, schema: &str, fields: &[(&str, &str)]) {
135        let fields: HashMap<String, String> = fields
136            .iter()
137            .map(|(k, v)| (k.to_string(), v.to_string()))
138            .collect();
139        let mut inner = self.lock_inner();
140        inner.events.push(RecordedEvent {
141            schema: schema.to_string(),
142            fields,
143        });
144    }
145}
146
147impl Default for RecordingSink {
148    fn default() -> Self {
149        Self::new()
150    }
151}
152
153#[cfg(test)]
154mod tests {
155    use super::*;
156
157    #[test]
158    fn captures_counters_and_events() {
159        let sink = RecordingSink::new();
160        sink.record_counter("valence_queries", &[("table", "user")], 1);
161        sink.log_event("valence.record.created", &[("table", "user")]);
162        assert_eq!(sink.counters().len(), 1);
163        assert_eq!(sink.events().len(), 1);
164    }
165}