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