mod json;
mod point;
mod sink;
pub use json::{Json, json};
pub use point::{Bus, Point, StoreOp};
pub use sink::{
Observer, is_observed, is_observed_anywhere, is_point_target, observe, stop_observing,
};
use std::cell::Cell;
use std::num::NonZeroU64;
use std::sync::OnceLock;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant, SystemTime};
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct Cause(NonZeroU64);
impl Cause {
fn next() -> Self {
static NEXT: AtomicU64 = AtomicU64::new(1);
Cause(NonZeroU64::new(NEXT.fetch_add(1, Ordering::Relaxed)).expect("ids start at 1"))
}
pub fn get(self) -> u64 {
self.0.get()
}
}
impl std::fmt::Display for Cause {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "#{}", self.0)
}
}
#[derive(Clone, Debug, PartialEq)]
pub struct Record {
pub id: Cause,
pub parent: Option<Cause>,
pub at: Duration,
pub point: Point,
}
#[derive(Clone, Debug, PartialEq)]
pub enum Trace {
Begin(Record),
End { id: Cause, took: Duration },
Mark(Record),
}
fn start() -> &'static (Instant, SystemTime) {
static START: OnceLock<(Instant, SystemTime)> = OnceLock::new();
START.get_or_init(|| (Instant::now(), SystemTime::now()))
}
fn now() -> Duration {
start().0.elapsed()
}
pub fn started_at() -> SystemTime {
start().1
}
thread_local! {
static CURRENT: Cell<Option<Cause>> = const { Cell::new(None) };
}
pub fn current() -> Option<Cause> {
CURRENT.with(Cell::get)
}
thread_local! {
static WATCHING: Cell<Duration> = const { Cell::new(Duration::ZERO) };
}
fn observed<R>(recording: impl FnOnce() -> R) -> R {
let started = Instant::now();
let done = recording();
WATCHING.with(|watching| watching.set(watching.get() + started.elapsed()));
done
}
pub fn mark(point: impl FnOnce() -> Point) -> Cause {
mark_under(current(), point)
}
pub fn mark_under(parent: Option<Cause>, point: impl FnOnce() -> Point) -> Cause {
let id = Cause::next();
if sink::wanted() {
observed(|| {
sink::emit(Trace::Mark(Record {
id,
parent,
at: now(),
point: point(),
}));
});
}
id
}
pub fn enter(point: impl FnOnce() -> Point) -> Entered {
enter_under(current(), point)
}
pub fn enter_under(parent: Option<Cause>, point: impl FnOnce() -> Point) -> Entered {
let id = Cause::next();
let started = now();
if sink::wanted() {
observed(|| {
sink::emit(Trace::Begin(Record {
id,
parent,
at: started,
point: point(),
}));
});
}
let previous = CURRENT.with(|current| current.replace(Some(id)));
Entered {
id,
previous,
started,
watched: WATCHING.with(Cell::get),
}
}
pub fn resume(cause: Option<Cause>) -> Resumed {
Resumed {
previous: CURRENT.with(|current| current.replace(cause)),
}
}
pub fn within<F: Future>(cause: Option<Cause>, future: F) -> Within<F> {
Within {
cause,
future: Box::pin(future),
}
}
pub struct Within<F> {
cause: Option<Cause>,
future: std::pin::Pin<Box<F>>,
}
impl<F: Future> Future for Within<F> {
type Output = F::Output;
fn poll(
mut self: std::pin::Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<F::Output> {
let _resumed = resume(self.cause);
self.future.as_mut().poll(cx)
}
}
#[must_use = "the point stops being current when this is dropped"]
pub struct Entered {
id: Cause,
previous: Option<Cause>,
started: Duration,
watched: Duration,
}
impl Entered {
pub fn id(&self) -> Cause {
self.id
}
}
impl Drop for Entered {
fn drop(&mut self) {
CURRENT.with(|current| current.set(self.previous));
if sink::wanted() {
let watching = WATCHING.with(Cell::get).saturating_sub(self.watched);
let took = now().saturating_sub(self.started).saturating_sub(watching);
observed(|| sink::emit(Trace::End { id: self.id, took }));
}
}
}
#[must_use = "the cause stops being current when this is dropped"]
pub struct Resumed {
previous: Option<Cause>,
}
impl Drop for Resumed {
fn drop(&mut self) {
CURRENT.with(|current| current.set(self.previous));
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::cell::RefCell;
use std::rc::Rc;
fn collect() -> Rc<RefCell<Vec<Trace>>> {
let seen = Rc::new(RefCell::new(Vec::new()));
let sink = seen.clone();
observe(move |trace| sink.borrow_mut().push(trace.clone()));
seen
}
fn parent_of(seen: &[Trace], id: Cause) -> Option<Cause> {
seen.iter().find_map(|trace| match trace {
Trace::Begin(record) | Trace::Mark(record) if record.id == id => Some(record.parent),
_ => None,
})?
}
#[test]
fn what_a_point_took_leaves_out_what_watching_it_cost() {
let slow = Rc::new(RefCell::new(Vec::new()));
let sink = slow.clone();
observe(move |trace| {
std::thread::sleep(std::time::Duration::from_millis(2));
if let Trace::End { took, .. } = trace {
sink.borrow_mut().push(*took);
}
});
{
let _action = enter(|| Point::Action { message: "Save" });
for _ in 0..5 {
mark(|| Point::Push { reducer: "Metrics" });
}
}
stop_observing();
let took = *slow.borrow().first().expect("the action ended");
assert!(
took < std::time::Duration::from_millis(5),
"five marks at two milliseconds of observer each were charged to the action: {took:?}"
);
}
#[test]
fn what_happens_inside_a_point_is_caused_by_it() {
let seen = collect();
let action = enter(|| Point::Action { message: "Kill" });
let send = mark(|| Point::Send {
actor: "ProcessActor",
message: "Kill",
});
let handled = enter_under(Some(send), || Point::Handle {
actor: "ProcessActor",
message: "Kill",
});
let publish = mark(|| Point::Publish {
event: "ProcessKilled",
bus: Bus::Global,
subscribers: 2,
});
let handle_id = handled.id();
let action_id = action.id();
drop(handled);
drop(action);
stop_observing();
let seen = seen.borrow();
assert_eq!(parent_of(&seen, send), Some(action_id));
assert_eq!(parent_of(&seen, handle_id), Some(send));
assert_eq!(parent_of(&seen, publish), Some(handle_id));
assert_eq!(parent_of(&seen, action_id), None);
assert!(matches!(seen.last(), Some(Trace::End { id, .. }) if *id == action_id));
assert_eq!(current(), None, "everything entered has been left");
}
#[test]
fn a_resumed_cause_is_the_parent_of_what_follows_and_goes_away_after() {
let seen = collect();
let spawn = mark(|| Point::Spawn {
actor: "Poller",
actor_id: 1,
output: "Tick",
});
{
let _resumed = resume(Some(spawn));
mark(|| Point::Push { reducer: "Metrics" });
}
let after = mark(|| Point::Push { reducer: "Metrics" });
stop_observing();
let seen = seen.borrow();
let pushes: Vec<Option<Cause>> = seen
.iter()
.filter_map(|trace| match trace {
Trace::Mark(record) if matches!(record.point, Point::Push { .. }) => {
Some(record.parent)
}
_ => None,
})
.collect();
assert_eq!(pushes, [Some(spawn), None]);
assert_eq!(parent_of(&seen, after), None);
}
type Event = (String, Vec<(String, String)>);
#[derive(Clone, Default)]
struct Written(std::sync::Arc<std::sync::Mutex<Vec<Event>>>);
struct Fields(Vec<(String, String)>);
impl tracing::field::Visit for Fields {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
self.0.push((field.name().to_string(), format!("{value:?}")));
}
}
impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for Written {
fn on_event(&self, event: &tracing::Event<'_>, _: tracing_subscriber::layer::Context<'_, S>) {
let mut fields = Fields(Vec::new());
event.record(&mut fields);
let target = event.metadata().target().to_string();
self.0.lock().unwrap().push((target, fields.0));
}
}
#[test]
fn a_point_goes_to_tracing_as_its_kind_with_what_it_holds_as_fields() {
use tracing_subscriber::layer::SubscriberExt;
let written = Written::default();
let subscriber = tracing_subscriber::registry().with(written.clone());
tracing::subscriber::with_default(subscriber, || {
let send = mark(|| Point::Send {
actor: "ProcessActor",
message: "Kill",
});
mark_under(Some(send), || Point::Log {
level: tracing::Level::INFO,
target: "app",
file: None,
line: None,
module: None,
text: "already written by whoever logged it".into(),
});
});
let written = written.0.lock().unwrap();
let fields: Vec<(&str, &str)> = written[0]
.1
.iter()
.map(|(name, value)| (name.as_str(), value.as_str()))
.collect();
assert_eq!(written.len(), 1, "a log point is not written back: {written:?}");
assert_eq!(written[0].0, "guinea::send");
assert!(fields.contains(&("actor", "ProcessActor")), "{fields:?}");
assert!(fields.contains(&("msg", "Kill")), "{fields:?}");
assert!(!fields.iter().any(|(name, _)| *name == "parent"), "no cause, no parent: {fields:?}");
assert!(!fields.iter().any(|(name, _)| *name == "message"), "no prose: {fields:?}");
}
#[test]
fn nothing_is_built_while_nobody_listens() {
let built = Cell::new(false);
mark(|| {
built.set(true);
Point::Note("unused".into())
});
assert!(!built.get());
}
}