Skip to main content

guinea_trace/
lib.rs

1//! What caused what.
2//!
3//! Every observable point - an action from the UI, a message sent to an actor,
4//! its handling, background work and its result, a publication and each of its
5//! deliveries, a reducer update, a navigation - is a [`Record`] with an id and
6//! the id of the point that caused it. Following `parent` from any record
7//! gives its provenance; following it the other way gives everything it set
8//! off.
9//!
10//! The cause is carried implicitly on the thread ([`current`]) and explicitly
11//! wherever work crosses a queue or a thread: an actor's mailbox, a background
12//! task, a hop back onto the UI thread. [`Cause`] is `Copy + Send` for that.
13//!
14//! Records go to the observers on the thread that produced them (devtools),
15//! and to `tracing` as one event each: the target names the kind of point,
16//! `guinea::send`, and the fields carry `id`, `parent` and what the point
17//! holds. `guinea::tick=off` silences one kind. [`json`] writes them, and the
18//! application's own events with the point they happened under, as JSON
19//! lines.
20
21mod json;
22mod point;
23mod sink;
24
25pub use json::{Json, json};
26pub use point::{Bus, Point, StoreOp};
27pub use sink::{
28    Observer, is_observed, is_observed_anywhere, is_point_target, observe, stop_observing,
29};
30
31use std::cell::Cell;
32use std::num::NonZeroU64;
33use std::sync::OnceLock;
34use std::sync::atomic::{AtomicU64, Ordering};
35use std::time::{Duration, Instant, SystemTime};
36
37/// One observed point, named for being the cause of what follows it.
38#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)]
39pub struct Cause(NonZeroU64);
40
41impl Cause {
42    fn next() -> Self {
43        static NEXT: AtomicU64 = AtomicU64::new(1);
44        Cause(NonZeroU64::new(NEXT.fetch_add(1, Ordering::Relaxed)).expect("ids start at 1"))
45    }
46
47    pub fn get(self) -> u64 {
48        self.0.get()
49    }
50}
51
52impl std::fmt::Display for Cause {
53    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
54        write!(f, "#{}", self.0)
55    }
56}
57
58#[derive(Clone, Debug, PartialEq)]
59pub struct Record {
60    pub id: Cause,
61    pub parent: Option<Cause>,
62    /// Since the process first traced anything.
63    pub at: Duration,
64    pub point: Point,
65}
66
67/// What an observer is told.
68#[derive(Clone, Debug, PartialEq)]
69pub enum Trace {
70    /// A point that stays current until [`Trace::End`]: what happens inside it
71    /// is caused by it.
72    Begin(Record),
73    End { id: Cause, took: Duration },
74    /// A point with no extent: a send, a push, a spawn.
75    Mark(Record),
76}
77
78fn start() -> &'static (Instant, SystemTime) {
79    static START: OnceLock<(Instant, SystemTime)> = OnceLock::new();
80    START.get_or_init(|| (Instant::now(), SystemTime::now()))
81}
82
83fn now() -> Duration {
84    start().0.elapsed()
85}
86
87/// The wall clock at the moment [`Record::at`] counts from.
88pub fn started_at() -> SystemTime {
89    start().1
90}
91
92thread_local! {
93    static CURRENT: Cell<Option<Cause>> = const { Cell::new(None) };
94}
95
96/// What is happening on this thread right now, if anything is.
97pub fn current() -> Option<Cause> {
98    CURRENT.with(Cell::get)
99}
100
101thread_local! {
102    /// What recording has cost since the point that is current now opened.
103    ///
104    /// Observing is not free - a record is built, given to devtools, written
105    /// through `tracing` - and it all happens inside whatever is being
106    /// measured. Charged to the observer instead, so a duration means the
107    /// same whether or not anyone is watching.
108    static WATCHING: Cell<Duration> = const { Cell::new(Duration::ZERO) };
109}
110
111/// Runs `recording` and charges what it took to the observer rather than to
112/// the point that is open.
113fn observed<R>(recording: impl FnOnce() -> R) -> R {
114    let started = Instant::now();
115    let done = recording();
116
117    WATCHING.with(|watching| watching.set(watching.get() + started.elapsed()));
118
119    done
120}
121
122/// Records a point caused by whatever is current, without making it current.
123pub fn mark(point: impl FnOnce() -> Point) -> Cause {
124    mark_under(current(), point)
125}
126
127/// Records a point caused by `parent`, for work that crossed a queue or a
128/// thread and brought its cause along.
129pub fn mark_under(parent: Option<Cause>, point: impl FnOnce() -> Point) -> Cause {
130    let id = Cause::next();
131    if sink::wanted() {
132        observed(|| {
133            sink::emit(Trace::Mark(Record {
134                id,
135                parent,
136                at: now(),
137                point: point(),
138            }));
139        });
140    }
141    id
142}
143
144/// Records a point caused by whatever is current, and makes it current until
145/// the guard drops.
146pub fn enter(point: impl FnOnce() -> Point) -> Entered {
147    enter_under(current(), point)
148}
149
150/// [`enter`], with the cause given rather than taken from the thread.
151pub fn enter_under(parent: Option<Cause>, point: impl FnOnce() -> Point) -> Entered {
152    let id = Cause::next();
153    let started = now();
154    if sink::wanted() {
155        observed(|| {
156            sink::emit(Trace::Begin(Record {
157                id,
158                parent,
159                at: started,
160                point: point(),
161            }));
162        });
163    }
164    let previous = CURRENT.with(|current| current.replace(Some(id)));
165
166    Entered {
167        id,
168        previous,
169        started,
170        // What observing had cost before this point opened. Whatever is
171        // added to it while the point is open was spent watching it, not
172        // doing it.
173        watched: WATCHING.with(Cell::get),
174    }
175}
176
177/// Makes `cause` current until the guard drops, without recording anything:
178/// for code that continues a point recorded elsewhere.
179pub fn resume(cause: Option<Cause>) -> Resumed {
180    Resumed {
181        previous: CURRENT.with(|current| current.replace(cause)),
182    }
183}
184
185/// Runs `future` with `cause` current on whichever thread polls it.
186///
187/// A guard held across `.await` would stay on the thread the task started on;
188/// this sets the cause for each poll instead.
189pub fn within<F: Future>(cause: Option<Cause>, future: F) -> Within<F> {
190    Within {
191        cause,
192        future: Box::pin(future),
193    }
194}
195
196pub struct Within<F> {
197    cause: Option<Cause>,
198    future: std::pin::Pin<Box<F>>,
199}
200
201impl<F: Future> Future for Within<F> {
202    type Output = F::Output;
203
204    fn poll(
205        mut self: std::pin::Pin<&mut Self>,
206        cx: &mut std::task::Context<'_>,
207    ) -> std::task::Poll<F::Output> {
208        let _resumed = resume(self.cause);
209        self.future.as_mut().poll(cx)
210    }
211}
212
213#[must_use = "the point stops being current when this is dropped"]
214pub struct Entered {
215    id: Cause,
216    previous: Option<Cause>,
217    started: Duration,
218    /// What observing had cost by the time this point opened.
219    watched: Duration,
220}
221
222impl Entered {
223    pub fn id(&self) -> Cause {
224        self.id
225    }
226}
227
228impl Drop for Entered {
229    fn drop(&mut self) {
230        CURRENT.with(|current| current.set(self.previous));
231
232        if sink::wanted() {
233            // Everything observing cost while this point was open comes off
234            // what the point is said to have taken. It stays on the running
235            // total, so the point above this one discounts it too - it was
236            // open for all of it as well.
237            let watching = WATCHING.with(Cell::get).saturating_sub(self.watched);
238            let took = now().saturating_sub(self.started).saturating_sub(watching);
239
240            observed(|| sink::emit(Trace::End { id: self.id, took }));
241        }
242    }
243}
244
245#[must_use = "the cause stops being current when this is dropped"]
246pub struct Resumed {
247    previous: Option<Cause>,
248}
249
250impl Drop for Resumed {
251    fn drop(&mut self) {
252        CURRENT.with(|current| current.set(self.previous));
253    }
254}
255
256#[cfg(test)]
257mod tests {
258    use super::*;
259    use std::cell::RefCell;
260    use std::rc::Rc;
261
262    fn collect() -> Rc<RefCell<Vec<Trace>>> {
263        let seen = Rc::new(RefCell::new(Vec::new()));
264        let sink = seen.clone();
265        observe(move |trace| sink.borrow_mut().push(trace.clone()));
266        seen
267    }
268
269    fn parent_of(seen: &[Trace], id: Cause) -> Option<Cause> {
270        seen.iter().find_map(|trace| match trace {
271            Trace::Begin(record) | Trace::Mark(record) if record.id == id => Some(record.parent),
272            _ => None,
273        })?
274    }
275
276    /// Observing costs time, and it is spent inside whatever is open. Left
277    /// in, a duration would say how long the work took *while watched*,
278    /// which is not a number anyone wants.
279    #[test]
280    fn what_a_point_took_leaves_out_what_watching_it_cost() {
281        let slow = Rc::new(RefCell::new(Vec::new()));
282        let sink = slow.clone();
283        observe(move |trace| {
284            // An observer that takes its time, so the cost is unmistakable.
285            std::thread::sleep(std::time::Duration::from_millis(2));
286            if let Trace::End { took, .. } = trace {
287                sink.borrow_mut().push(*took);
288            }
289        });
290
291        {
292            let _action = enter(|| Point::Action { message: "Save" });
293            for _ in 0..5 {
294                mark(|| Point::Push { reducer: "Metrics" });
295            }
296        }
297        stop_observing();
298
299        let took = *slow.borrow().first().expect("the action ended");
300        assert!(
301            took < std::time::Duration::from_millis(5),
302            "five marks at two milliseconds of observer each were charged to the action: {took:?}"
303        );
304    }
305
306    #[test]
307    fn what_happens_inside_a_point_is_caused_by_it() {
308        let seen = collect();
309
310        let action = enter(|| Point::Action { message: "Kill" });
311        let send = mark(|| Point::Send {
312            actor: "ProcessActor",
313            message: "Kill",
314        });
315        let handled = enter_under(Some(send), || Point::Handle {
316            actor: "ProcessActor",
317            message: "Kill",
318        });
319        let publish = mark(|| Point::Publish {
320            event: "ProcessKilled",
321            bus: Bus::Global,
322            subscribers: 2,
323        });
324        let handle_id = handled.id();
325        let action_id = action.id();
326        drop(handled);
327        drop(action);
328        stop_observing();
329
330        let seen = seen.borrow();
331        assert_eq!(parent_of(&seen, send), Some(action_id));
332        assert_eq!(parent_of(&seen, handle_id), Some(send));
333        assert_eq!(parent_of(&seen, publish), Some(handle_id));
334        assert_eq!(parent_of(&seen, action_id), None);
335        assert!(matches!(seen.last(), Some(Trace::End { id, .. }) if *id == action_id));
336        assert_eq!(current(), None, "everything entered has been left");
337    }
338
339    #[test]
340    fn a_resumed_cause_is_the_parent_of_what_follows_and_goes_away_after() {
341        let seen = collect();
342        let spawn = mark(|| Point::Spawn {
343            actor: "Poller",
344            actor_id: 1,
345            output: "Tick",
346        });
347        {
348            let _resumed = resume(Some(spawn));
349            mark(|| Point::Push { reducer: "Metrics" });
350        }
351        let after = mark(|| Point::Push { reducer: "Metrics" });
352        stop_observing();
353
354        let seen = seen.borrow();
355        let pushes: Vec<Option<Cause>> = seen
356            .iter()
357            .filter_map(|trace| match trace {
358                Trace::Mark(record) if matches!(record.point, Point::Push { .. }) => {
359                    Some(record.parent)
360                }
361                _ => None,
362            })
363            .collect();
364        assert_eq!(pushes, [Some(spawn), None]);
365        assert_eq!(parent_of(&seen, after), None);
366    }
367
368    type Event = (String, Vec<(String, String)>);
369
370    #[derive(Clone, Default)]
371    struct Written(std::sync::Arc<std::sync::Mutex<Vec<Event>>>);
372
373    struct Fields(Vec<(String, String)>);
374
375    impl tracing::field::Visit for Fields {
376        fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
377            self.0.push((field.name().to_string(), format!("{value:?}")));
378        }
379    }
380
381    impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for Written {
382        fn on_event(&self, event: &tracing::Event<'_>, _: tracing_subscriber::layer::Context<'_, S>) {
383            let mut fields = Fields(Vec::new());
384            event.record(&mut fields);
385
386            let target = event.metadata().target().to_string();
387            self.0.lock().unwrap().push((target, fields.0));
388        }
389    }
390
391    #[test]
392    fn a_point_goes_to_tracing_as_its_kind_with_what_it_holds_as_fields() {
393        use tracing_subscriber::layer::SubscriberExt;
394
395        let written = Written::default();
396        let subscriber = tracing_subscriber::registry().with(written.clone());
397
398        tracing::subscriber::with_default(subscriber, || {
399            let send = mark(|| Point::Send {
400                actor: "ProcessActor",
401                message: "Kill",
402            });
403            mark_under(Some(send), || Point::Log {
404                level: tracing::Level::INFO,
405                target: "app",
406                file: None,
407                line: None,
408                module: None,
409                text: "already written by whoever logged it".into(),
410            });
411        });
412
413        let written = written.0.lock().unwrap();
414        let fields: Vec<(&str, &str)> = written[0]
415            .1
416            .iter()
417            .map(|(name, value)| (name.as_str(), value.as_str()))
418            .collect();
419
420        assert_eq!(written.len(), 1, "a log point is not written back: {written:?}");
421        assert_eq!(written[0].0, "guinea::send");
422        assert!(fields.contains(&("actor", "ProcessActor")), "{fields:?}");
423        assert!(fields.contains(&("msg", "Kill")), "{fields:?}");
424        assert!(!fields.iter().any(|(name, _)| *name == "parent"), "no cause, no parent: {fields:?}");
425        assert!(!fields.iter().any(|(name, _)| *name == "message"), "no prose: {fields:?}");
426    }
427
428    #[test]
429    fn nothing_is_built_while_nobody_listens() {
430        let built = Cell::new(false);
431        mark(|| {
432            built.set(true);
433            Point::Note("unused".into())
434        });
435        assert!(!built.get());
436    }
437}