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