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, 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#[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 pub at: Duration,
64 pub point: Point,
65}
66
67#[derive(Clone, Debug, PartialEq)]
69pub enum Trace {
70 Begin(Record),
73 End { id: Cause, took: Duration },
74 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
87pub fn started_at() -> SystemTime {
89 start().1
90}
91
92thread_local! {
93 static CURRENT: Cell<Option<Cause>> = const { Cell::new(None) };
94}
95
96pub fn current() -> Option<Cause> {
98 CURRENT.with(Cell::get)
99}
100
101thread_local! {
102 static WATCHING: Cell<Duration> = const { Cell::new(Duration::ZERO) };
109}
110
111fn 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
122pub fn mark(point: impl FnOnce() -> Point) -> Cause {
124 mark_under(current(), point)
125}
126
127pub 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
144pub fn enter(point: impl FnOnce() -> Point) -> Entered {
147 enter_under(current(), point)
148}
149
150pub 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 watched: WATCHING.with(Cell::get),
174 }
175}
176
177pub fn resume(cause: Option<Cause>) -> Resumed {
180 Resumed {
181 previous: CURRENT.with(|current| current.replace(cause)),
182 }
183}
184
185pub 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 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 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 #[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 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}