Skip to main content

turnframe_telemetry/
composite.rs

1//! Combining observers, and the observer tests assert against.
2//!
3//! An application usually wants more than one destination for the same signal:
4//! a counter for the dashboard, a log line for the investigation, and — in a
5//! test — a list it can make assertions about. [`CompositeObserver`] fans one
6//! signal out to several observers; [`RecordingObserver`] keeps every signal it
7//! was given so another crate's tests can check what a run reported.
8
9use std::fmt;
10use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
11use std::time::Duration;
12
13use turnframe_core::observe::{Observer, Signal, SignalLabels};
14
15/// Fans every signal out to several observers, in registration order.
16///
17/// ```rust
18/// use std::sync::Arc;
19///
20/// use turnframe_core::observe::{Observer, Signal};
21/// use turnframe_telemetry::{CompositeObserver, MetricsObserver, RecordingObserver};
22///
23/// let recorder = Arc::new(RecordingObserver::new());
24/// let observer = CompositeObserver::new()
25///     .with(MetricsObserver::new())
26///     .with_shared(Arc::clone(&recorder));
27///
28/// observer.observe(&Signal::TurnReceived);
29/// assert_eq!(recorder.count(Signal::TurnReceived), 1);
30/// ```
31#[derive(Clone, Default)]
32pub struct CompositeObserver {
33    observers: Vec<Arc<dyn Observer>>,
34}
35
36impl CompositeObserver {
37    /// An empty composite. Signals given to it go nowhere.
38    #[must_use]
39    pub fn new() -> Self {
40        Self {
41            observers: Vec::new(),
42        }
43    }
44
45    /// Adds an observer, taking ownership of it.
46    #[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    /// Adds an observer the caller keeps a handle to, which is what a test
53    /// needs in order to inspect a [`RecordingObserver`] afterwards.
54    ///
55    /// Generic over the concrete observer so `Arc::clone(&handle)` can be
56    /// passed straight in; use [`CompositeObserver::push`] for a handle that is
57    /// already erased to `Arc<dyn Observer>`.
58    #[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    /// Adds an already-erased observer to an existing composite, for a
65    /// registry assembled at run time.
66    pub fn push(&mut self, observer: Arc<dyn Observer>) {
67        self.observers.push(observer);
68    }
69
70    /// How many observers are registered.
71    #[must_use]
72    pub fn len(&self) -> usize {
73        self.observers.len()
74    }
75
76    /// Returns `true` when no observer is registered.
77    #[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/// One signal as a [`RecordingObserver`] kept it.
112#[derive(Debug, Clone, PartialEq, Eq)]
113#[non_exhaustive]
114pub struct RecordedSignal {
115    /// The signal.
116    pub signal: Signal,
117    /// The labels it arrived with; empty when it arrived without any.
118    pub labels: SignalLabels,
119    /// The measured duration, for a latency signal.
120    pub duration: Option<Duration>,
121}
122
123/// An observer that keeps everything it was given, for assertions in tests.
124///
125/// It is cheap, thread-safe and never panics: a poisoned lock is recovered
126/// rather than propagated, because a telemetry sink must not turn one test
127/// failure into a cascade.
128///
129/// ```rust
130/// use turnframe_core::ids::WorkflowKey;
131/// use turnframe_core::observe::{Observer, Signal, SignalLabels};
132/// use turnframe_telemetry::RecordingObserver;
133///
134/// let observer = RecordingObserver::new();
135/// observer.observe_labeled(
136///     &Signal::CommandExecuted,
137///     &SignalLabels::workflow(WorkflowKey::from("trip")),
138/// );
139///
140/// assert_eq!(observer.count(Signal::CommandExecuted), 1);
141/// assert!(observer.contains(Signal::CommandExecuted));
142/// assert_eq!(observer.labels_of(Signal::CommandExecuted).len(), 1);
143/// ```
144#[derive(Debug, Default)]
145pub struct RecordingObserver {
146    records: Mutex<Vec<RecordedSignal>>,
147}
148
149impl RecordingObserver {
150    /// An observer with nothing recorded yet.
151    #[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    /// Everything recorded so far, in arrival order.
161    #[must_use]
162    pub fn records(&self) -> Vec<RecordedSignal> {
163        self.guard().clone()
164    }
165
166    /// The signals recorded so far, in arrival order.
167    #[must_use]
168    pub fn signals(&self) -> Vec<Signal> {
169        self.guard().iter().map(|record| record.signal).collect()
170    }
171
172    /// How many times a signal was recorded.
173    #[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    /// Whether a signal was recorded at least once.
182    #[must_use]
183    pub fn contains(&self, signal: Signal) -> bool {
184        self.guard().iter().any(|record| record.signal == signal)
185    }
186
187    /// The label sets a signal was recorded with, in arrival order.
188    #[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    /// The durations a latency signal was recorded with, in arrival order.
198    #[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    /// How many signals were recorded in total.
208    #[must_use]
209    pub fn len(&self) -> usize {
210        self.guard().len()
211    }
212
213    /// Whether nothing was recorded.
214    #[must_use]
215    pub fn is_empty(&self) -> bool {
216        self.guard().is_empty()
217    }
218
219    /// Drops everything recorded so far.
220    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}