valence_telemetry/
recording.rs1use std::collections::HashMap;
4use std::sync::{Arc, Mutex};
5
6use crate::TelemetrySink;
7
8#[derive(Debug, Clone, PartialEq, Eq)]
10pub struct RecordedCounter {
11 pub name: String,
13 pub labels: Vec<(String, String)>,
15 pub delta: u64,
17}
18
19#[derive(Debug, Clone, PartialEq)]
21pub struct RecordedGauge {
22 pub name: String,
24 pub labels: Vec<(String, String)>,
26 pub value: f64,
28}
29
30#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct RecordedEvent {
33 pub schema: String,
35 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#[derive(Debug, Clone)]
48pub struct RecordingSink {
49 inner: Arc<Mutex<Inner>>,
50}
51
52impl RecordingSink {
53 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 pub fn counters(&self) -> Vec<RecordedCounter> {
68 self.lock_inner().counters.clone()
69 }
70
71 pub fn events(&self) -> Vec<RecordedEvent> {
73 self.lock_inner().events.clone()
74 }
75
76 pub fn gauges(&self) -> Vec<RecordedGauge> {
78 self.lock_inner().gauges.clone()
79 }
80
81 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 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}