turnframe_telemetry/
composite.rs1use std::fmt;
10use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
11use std::time::Duration;
12
13use turnframe_core::observe::{Observer, Signal, SignalLabels};
14
15#[derive(Clone, Default)]
32pub struct CompositeObserver {
33 observers: Vec<Arc<dyn Observer>>,
34}
35
36impl CompositeObserver {
37 #[must_use]
39 pub fn new() -> Self {
40 Self {
41 observers: Vec::new(),
42 }
43 }
44
45 #[must_use]
47 pub fn with(mut self, observer: impl Observer + 'static) -> Self {
48 self.observers.push(Arc::new(observer));
49 self
50 }
51
52 #[must_use]
59 pub fn with_shared<O: Observer + 'static>(mut self, observer: Arc<O>) -> Self {
60 self.observers.push(observer);
61 self
62 }
63
64 pub fn push(&mut self, observer: Arc<dyn Observer>) {
67 self.observers.push(observer);
68 }
69
70 #[must_use]
72 pub fn len(&self) -> usize {
73 self.observers.len()
74 }
75
76 #[must_use]
78 pub fn is_empty(&self) -> bool {
79 self.observers.is_empty()
80 }
81}
82
83impl fmt::Debug for CompositeObserver {
84 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
85 f.debug_struct("CompositeObserver")
86 .field("observers", &self.observers.len())
87 .finish()
88 }
89}
90
91impl Observer for CompositeObserver {
92 fn observe(&self, signal: &Signal) {
93 for observer in &self.observers {
94 observer.observe(signal);
95 }
96 }
97
98 fn observe_labeled(&self, signal: &Signal, labels: &SignalLabels) {
99 for observer in &self.observers {
100 observer.observe_labeled(signal, labels);
101 }
102 }
103
104 fn observe_duration(&self, signal: &Signal, duration: Duration, labels: &SignalLabels) {
105 for observer in &self.observers {
106 observer.observe_duration(signal, duration, labels);
107 }
108 }
109}
110
111#[derive(Debug, Clone, PartialEq, Eq)]
113#[non_exhaustive]
114pub struct RecordedSignal {
115 pub signal: Signal,
117 pub labels: SignalLabels,
119 pub duration: Option<Duration>,
121}
122
123#[derive(Debug, Default)]
145pub struct RecordingObserver {
146 records: Mutex<Vec<RecordedSignal>>,
147}
148
149impl RecordingObserver {
150 #[must_use]
152 pub fn new() -> Self {
153 Self::default()
154 }
155
156 fn guard(&self) -> MutexGuard<'_, Vec<RecordedSignal>> {
157 self.records.lock().unwrap_or_else(PoisonError::into_inner)
158 }
159
160 #[must_use]
162 pub fn records(&self) -> Vec<RecordedSignal> {
163 self.guard().clone()
164 }
165
166 #[must_use]
168 pub fn signals(&self) -> Vec<Signal> {
169 self.guard().iter().map(|record| record.signal).collect()
170 }
171
172 #[must_use]
174 pub fn count(&self, signal: Signal) -> usize {
175 self.guard()
176 .iter()
177 .filter(|record| record.signal == signal)
178 .count()
179 }
180
181 #[must_use]
183 pub fn contains(&self, signal: Signal) -> bool {
184 self.guard().iter().any(|record| record.signal == signal)
185 }
186
187 #[must_use]
189 pub fn labels_of(&self, signal: Signal) -> Vec<SignalLabels> {
190 self.guard()
191 .iter()
192 .filter(|record| record.signal == signal)
193 .map(|record| record.labels.clone())
194 .collect()
195 }
196
197 #[must_use]
199 pub fn durations_of(&self, signal: Signal) -> Vec<Duration> {
200 self.guard()
201 .iter()
202 .filter(|record| record.signal == signal)
203 .filter_map(|record| record.duration)
204 .collect()
205 }
206
207 #[must_use]
209 pub fn len(&self) -> usize {
210 self.guard().len()
211 }
212
213 #[must_use]
215 pub fn is_empty(&self) -> bool {
216 self.guard().is_empty()
217 }
218
219 pub fn clear(&self) {
221 self.guard().clear();
222 }
223
224 fn push(&self, signal: Signal, labels: &SignalLabels, duration: Option<Duration>) {
225 self.guard().push(RecordedSignal {
226 signal,
227 labels: labels.clone(),
228 duration,
229 });
230 }
231}
232
233impl Observer for RecordingObserver {
234 fn observe(&self, signal: &Signal) {
235 self.push(*signal, &SignalLabels::none(), None);
236 }
237
238 fn observe_labeled(&self, signal: &Signal, labels: &SignalLabels) {
239 self.push(*signal, labels, None);
240 }
241
242 fn observe_duration(&self, signal: &Signal, duration: Duration, labels: &SignalLabels) {
243 self.push(*signal, labels, Some(duration));
244 }
245}
246
247#[cfg(test)]
248mod tests {
249 use turnframe_core::ids::WorkflowKey;
250
251 use super::*;
252
253 #[test]
254 fn a_composite_fans_out_to_every_observer() {
255 let first = Arc::new(RecordingObserver::new());
256 let second = Arc::new(RecordingObserver::new());
257 let composite = CompositeObserver::new()
258 .with_shared(Arc::clone(&first))
259 .with_shared(Arc::clone(&second));
260
261 assert_eq!(composite.len(), 2);
262 assert!(!composite.is_empty());
263
264 composite.observe(&Signal::TurnReceived);
265 composite.observe_labeled(
266 &Signal::CommandExecuted,
267 &SignalLabels::workflow(WorkflowKey::from("trip")),
268 );
269 composite.observe_duration(
270 &Signal::TurnDuration,
271 Duration::from_millis(12),
272 &SignalLabels::none(),
273 );
274
275 for recorder in [&first, &second] {
276 assert_eq!(recorder.len(), 3);
277 assert_eq!(recorder.count(Signal::TurnReceived), 1);
278 assert_eq!(
279 recorder.labels_of(Signal::CommandExecuted)[0]
280 .workflow
281 .as_ref()
282 .map(WorkflowKey::as_str),
283 Some("trip")
284 );
285 assert_eq!(
286 recorder.durations_of(Signal::TurnDuration),
287 vec![Duration::from_millis(12)]
288 );
289 }
290 }
291
292 #[test]
293 fn an_empty_composite_is_a_sink() {
294 let composite = CompositeObserver::new();
295 assert!(composite.is_empty());
296 assert_eq!(composite.len(), 0);
297 composite.observe(&Signal::TurnReceived);
298 assert_eq!(
299 format!("{composite:?}"),
300 "CompositeObserver { observers: 0 }"
301 );
302 }
303
304 #[test]
305 fn with_takes_ownership_of_an_observer() {
306 let composite = CompositeObserver::new().with(RecordingObserver::new());
307 assert_eq!(composite.len(), 1);
308 composite.observe(&Signal::TurnCompleted);
309 }
310
311 #[test]
312 fn push_adds_to_an_existing_composite() {
313 let recorder = Arc::new(RecordingObserver::new());
314 let erased: Arc<dyn Observer> = recorder.clone();
315 let mut composite = CompositeObserver::new();
316 composite.push(erased);
317 composite.observe(&Signal::QuestionAnswered);
318 assert!(recorder.contains(Signal::QuestionAnswered));
319 }
320
321 #[test]
322 fn a_recording_observer_keeps_order_and_clears() {
323 let observer = RecordingObserver::new();
324 assert!(observer.is_empty());
325 observer.observe(&Signal::TurnReceived);
326 observer.observe(&Signal::TurnCompleted);
327 assert_eq!(
328 observer.signals(),
329 vec![Signal::TurnReceived, Signal::TurnCompleted]
330 );
331 assert_eq!(observer.records().len(), 2);
332 assert!(!observer.contains(Signal::TurnFailed));
333 observer.clear();
334 assert!(observer.is_empty());
335 assert_eq!(observer.count(Signal::TurnReceived), 0);
336 }
337
338 #[test]
339 fn durations_are_only_kept_for_measured_signals() {
340 let observer = RecordingObserver::new();
341 observer.observe_labeled(&Signal::TurnDuration, &SignalLabels::none());
342 assert!(observer.durations_of(Signal::TurnDuration).is_empty());
343 observer.observe_duration(
344 &Signal::TurnDuration,
345 Duration::from_micros(900),
346 &SignalLabels::none(),
347 );
348 assert_eq!(
349 observer.durations_of(Signal::TurnDuration),
350 vec![Duration::from_micros(900)]
351 );
352 }
353}