Skip to main content

guinea_trace/
sink.rs

1//! Where records go.
2
3use std::cell::RefCell;
4use std::collections::VecDeque;
5use std::rc::Rc;
6use std::sync::atomic::{AtomicUsize, Ordering};
7use std::sync::{Mutex, OnceLock, PoisonError};
8
9use tracing::Level;
10
11use crate::{Cause, Point, Record, Trace};
12
13pub type Observer = Rc<dyn Fn(&Trace)>;
14
15thread_local! {
16    static OBSERVER: RefCell<Option<Observer>> = const { RefCell::new(None) };
17}
18
19static OBSERVED_THREADS: AtomicUsize = AtomicUsize::new(0);
20
21/// Hands every record produced on this thread to `observer`, replacing any
22/// observer already set.
23pub fn observe(observer: impl Fn(&Trace) + 'static) {
24    let before = OBSERVER.with(|slot| slot.borrow_mut().replace(Rc::new(observer)));
25    if before.is_none() {
26        OBSERVED_THREADS.fetch_add(1, Ordering::Relaxed);
27    }
28}
29
30pub fn stop_observing() {
31    if OBSERVER.with(|slot| slot.borrow_mut().take()).is_some()
32        && OBSERVED_THREADS.fetch_sub(1, Ordering::Relaxed) == 1
33    {
34        take_elsewhere();
35    }
36}
37
38/// Whether devtools are watching this thread.
39pub fn is_observed() -> bool {
40    OBSERVER.with(|slot| slot.borrow().is_some())
41}
42
43/// Whether devtools are watching any thread.
44pub fn is_observed_anywhere() -> bool {
45    OBSERVED_THREADS.load(Ordering::Relaxed) > 0
46}
47
48/// Whether a point recorded anywhere goes somewhere: to devtools, or to
49/// `tracing` as `guinea::` events.
50pub fn is_recorded_anywhere() -> bool {
51    is_observed_anywhere() || tracing::enabled!(target: "guinea", Level::DEBUG)
52}
53
54/// How many records [`take_elsewhere`] keeps between two takes; past it the
55/// oldest go.
56pub const ELSEWHERE_LIMIT: usize = 16_384;
57
58/// Records made on threads nobody observes, while some thread observes.
59#[derive(Debug, Default)]
60pub struct Elsewhere {
61    /// Oldest first, each with the thread it was made on, as [`thread_id`]
62    /// names it.
63    pub records: Vec<(u32, Trace)>,
64    /// How many were let go since the last take, to keep the newest.
65    pub dropped: u64,
66}
67
68#[derive(Default)]
69struct Queue {
70    records: VecDeque<(u32, Trace)>,
71    dropped: u64,
72}
73
74fn elsewhere() -> &'static Mutex<Queue> {
75    static ELSEWHERE: OnceLock<Mutex<Queue>> = OnceLock::new();
76    ELSEWHERE.get_or_init(Mutex::default)
77}
78
79fn keep_elsewhere(trace: Trace) {
80    let mut queue = elsewhere().lock().unwrap_or_else(PoisonError::into_inner);
81
82    if queue.records.len() == ELSEWHERE_LIMIT {
83        queue.records.pop_front();
84        queue.dropped += 1;
85    }
86    queue.records.push_back((thread_id(), trace));
87}
88
89/// What other threads recorded since the last take.
90///
91/// Kept while any thread observes, from the threads that do not; a thread
92/// that observes hears its own records and leaves none here.
93pub fn take_elsewhere() -> Elsewhere {
94    let mut queue = elsewhere().lock().unwrap_or_else(PoisonError::into_inner);
95    let taken = std::mem::take(&mut *queue);
96
97    Elsewhere {
98        records: taken.records.into(),
99        dropped: taken.dropped,
100    }
101}
102
103/// This thread, as [`Elsewhere`] names it: the operating system's id on
104/// Windows, what a stack sampler names it by; elsewhere a number this
105/// process gives each thread once, from 1.
106pub fn thread_id() -> u32 {
107    thread_local! {
108        static ID: u32 = os_thread_id();
109    }
110    ID.with(|id| *id)
111}
112
113#[cfg(windows)]
114fn os_thread_id() -> u32 {
115    #[link(name = "kernel32")]
116    unsafe extern "system" {
117        fn GetCurrentThreadId() -> u32;
118    }
119
120    unsafe { GetCurrentThreadId() }
121}
122
123#[cfg(not(windows))]
124fn os_thread_id() -> u32 {
125    static NEXT: std::sync::atomic::AtomicU32 = std::sync::atomic::AtomicU32::new(1);
126    NEXT.fetch_add(1, Ordering::Relaxed)
127}
128
129fn observer() -> Option<Observer> {
130    OBSERVER.with(|slot| slot.borrow().clone())
131}
132
133pub(crate) fn wanted() -> bool {
134    is_recorded_anywhere()
135}
136
137/// Whether `target` is one [`emit`] writes points under, so a layer that
138/// turns `tracing` events into points can leave them alone.
139pub fn is_point_target(target: &str) -> bool {
140    target.starts_with("guinea::")
141}
142
143macro_rules! point {
144    ($target:literal, $record:expr $(, $($field:tt)+)?) => {
145        tracing::debug!(
146            target: $target,
147            id = $record.id.get(),
148            parent = $record.parent.map(Cause::get)
149            $(, $($field)+)?
150        )
151    };
152}
153
154fn write(record: &Record) {
155    match &record.point {
156        Point::Action { message } => point!("guinea::action", record, action = %message),
157        Point::Send { actor, message } => {
158            point!("guinea::send", record, actor = %actor, msg = %message)
159        }
160        Point::Handle { actor, message } => {
161            point!("guinea::handle", record, actor = %actor, msg = %message)
162        }
163        Point::Spawn {
164            actor,
165            actor_id,
166            output,
167        } => point!(
168            "guinea::spawn",
169            record,
170            actor = %actor,
171            actor_id,
172            output = %output
173        ),
174        Point::Settled {
175            actor,
176            actor_id,
177            output,
178            took_us,
179        } => point!(
180            "guinea::settled",
181            record,
182            actor = %actor,
183            actor_id,
184            output = %output,
185            took_us
186        ),
187        Point::Cancelled {
188            actor,
189            actor_id,
190            output,
191            took_us,
192        } => point!(
193            "guinea::cancelled",
194            record,
195            actor = %actor,
196            actor_id,
197            output = %output,
198            took_us
199        ),
200        Point::Source {
201            actor,
202            actor_id,
203            output,
204        } => point!(
205            "guinea::source",
206            record,
207            actor = %actor,
208            actor_id,
209            output = %output
210        ),
211        Point::Arrived {
212            actor,
213            actor_id,
214            output,
215            source,
216        } => point!(
217            "guinea::arrived",
218            record,
219            actor = %actor,
220            actor_id,
221            output = %output,
222            source
223        ),
224        Point::Pull {
225            actor,
226            actor_id,
227            output,
228            source,
229        } => point!(
230            "guinea::pull",
231            record,
232            actor = %actor,
233            actor_id,
234            output = %output,
235            source
236        ),
237        Point::Closed {
238            actor,
239            actor_id,
240            output,
241            took_us,
242            gone,
243        } => point!(
244            "guinea::closed",
245            record,
246            actor = %actor,
247            actor_id,
248            output = %output,
249            took_us,
250            gone
251        ),
252        Point::Publish {
253            event,
254            bus,
255            subscribers,
256        } => point!(
257            "guinea::publish",
258            record,
259            event = %event,
260            bus = %bus,
261            subscribers
262        ),
263        Point::Deliver { event, bus } => {
264            point!("guinea::deliver", record, event = %event, bus = %bus)
265        }
266        Point::Push { reducer } => point!("guinea::push", record, reducer = %reducer),
267        Point::Navigate { root, to } => point!("guinea::navigate", record, root = %root, to = %to),
268        Point::Tick {
269            timer,
270            name,
271            file,
272            line,
273        } => point!("guinea::tick", record, timer, name = *name, file = *file, line = *line),
274        Point::Store {
275            op,
276            path,
277            field,
278            outside,
279        } => point!(
280            "guinea::store",
281            record,
282            op = %op,
283            path = %path,
284            field = field.as_deref(),
285            outside
286        ),
287        Point::Render { segment, took_us } => {
288            point!("guinea::render", record, segment = %segment, took_us)
289        }
290        Point::Span {
291            name,
292            file,
293            line,
294            module,
295            fields,
296            ..
297        } => point!(
298            "guinea::span",
299            record,
300            span = %name,
301            module = *module,
302            file = *file,
303            line = *line,
304            fields = %fields
305        ),
306        Point::Note(text) => point!("guinea::note", record, "{text}"),
307        Point::Log { .. } => {}
308    }
309}
310
311pub(crate) fn emit(trace: Trace) {
312    match &trace {
313        Trace::Begin(record) | Trace::Mark(record) => write(record),
314        Trace::End { id, took } => tracing::trace!(
315            target: "guinea::end",
316            id = id.get(),
317            took_us = took.as_micros() as u64
318        ),
319    }
320
321    match observer() {
322        Some(observer) => observer(&trace),
323        None if is_observed_anywhere() => keep_elsewhere(trace),
324        None => {}
325    }
326}