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
84pub fn now() -> Duration {
87 start().0.elapsed()
88}
89
90pub fn started_at() -> SystemTime {
92 start().1
93}
94
95thread_local! {
96 static CURRENT: Cell<Option<Cause>> = const { Cell::new(None) };
97}
98
99pub fn current() -> Option<Cause> {
101 CURRENT.with(Cell::get)
102}
103
104thread_local! {
105 static WATCHING: Cell<Duration> = const { Cell::new(Duration::ZERO) };
112}
113
114fn 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
125pub fn mark(point: impl FnOnce() -> Point) -> Cause {
127 mark_under(current(), point)
128}
129
130pub 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
147pub fn enter(point: impl FnOnce() -> Point) -> Entered {
150 enter_under(current(), point)
151}
152
153pub 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 watched: WATCHING.with(Cell::get),
177 }
178}
179
180pub fn reserve() -> Cause {
184 Cause::next()
185}
186
187pub 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
203pub fn end(id: Cause, took: Duration) {
205 if sink::wanted() {
206 observed(|| sink::emit(Trace::End { id, took }));
207 }
208}
209
210pub fn resume(cause: Option<Cause>) -> Resumed {
213 Resumed {
214 previous: CURRENT.with(|current| current.replace(cause)),
215 }
216}
217
218pub 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 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 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 #[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 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}