1use 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
21pub 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
38pub fn is_observed() -> bool {
40 OBSERVER.with(|slot| slot.borrow().is_some())
41}
42
43pub fn is_observed_anywhere() -> bool {
45 OBSERVED_THREADS.load(Ordering::Relaxed) > 0
46}
47
48pub fn is_recorded_anywhere() -> bool {
51 is_observed_anywhere() || tracing::enabled!(target: "guinea", Level::DEBUG)
52}
53
54pub const ELSEWHERE_LIMIT: usize = 16_384;
57
58#[derive(Debug, Default)]
60pub struct Elsewhere {
61 pub records: Vec<(u32, Trace)>,
64 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
89pub 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
103pub 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
137pub 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}