1mod 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#[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 pub at: Duration,
65 pub point: Point,
66}
67
68#[derive(Clone, Debug, PartialEq)]
70pub enum Trace {
71 Begin(Record),
74 End { id: Cause, took: Duration },
75 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
88pub fn started_at() -> SystemTime {
90 start().1
91}
92
93thread_local! {
94 static CURRENT: Cell<Option<Cause>> = const { Cell::new(None) };
95}
96
97pub fn current() -> Option<Cause> {
99 CURRENT.with(Cell::get)
100}
101
102thread_local! {
103 static WATCHING: Cell<Duration> = const { Cell::new(Duration::ZERO) };
110}
111
112fn 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
123pub fn mark(point: impl FnOnce() -> Point) -> Cause {
125 mark_under(current(), point)
126}
127
128pub 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
145pub fn enter(point: impl FnOnce() -> Point) -> Entered {
148 enter_under(current(), point)
149}
150
151pub 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 watched: WATCHING.with(Cell::get),
175 }
176}
177
178pub fn reserve() -> Cause {
182 Cause::next()
183}
184
185pub 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
201pub fn end(id: Cause, took: Duration) {
203 if sink::wanted() {
204 observed(|| sink::emit(Trace::End { id, took }));
205 }
206}
207
208pub fn resume(cause: Option<Cause>) -> Resumed {
211 Resumed {
212 previous: CURRENT.with(|current| current.replace(cause)),
213 }
214}
215
216pub 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 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 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 #[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 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}