1use std::collections::HashMap;
16use std::collections::VecDeque;
17use std::sync::Arc;
18use std::sync::Mutex;
19use std::sync::atomic::{AtomicU8, Ordering};
20
21use crate::PluginSeam;
22
23#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug)]
31#[repr(u8)]
32pub enum TraceLevel {
33 Off = 0,
35 Error = 1,
36 Warn = 2,
37 Info = 3,
38 Debug = 4,
39 Trace = 5,
40}
41
42impl TraceLevel {
43 pub fn as_str(self) -> &'static str {
45 match self {
46 TraceLevel::Off => "off",
47 TraceLevel::Error => "error",
48 TraceLevel::Warn => "warn",
49 TraceLevel::Info => "info",
50 TraceLevel::Debug => "debug",
51 TraceLevel::Trace => "trace",
52 }
53 }
54
55 pub const ALL: [TraceLevel; 6] = [
58 TraceLevel::Off,
59 TraceLevel::Error,
60 TraceLevel::Warn,
61 TraceLevel::Info,
62 TraceLevel::Debug,
63 TraceLevel::Trace,
64 ];
65
66 pub fn cycle_next(self) -> TraceLevel {
69 TraceLevel::from_u8((self.as_u8() + 1) % Self::ALL.len() as u8)
70 }
71
72 pub fn parse(s: &str) -> Option<TraceLevel> {
76 TraceLevel::ALL.into_iter().find(|l| l.as_str() == s)
77 }
78
79 #[inline]
81 pub fn as_u8(self) -> u8 {
82 self as u8
83 }
84
85 #[inline]
88 pub fn from_u8(v: u8) -> Self {
89 match v {
90 1 => TraceLevel::Error,
91 2 => TraceLevel::Warn,
92 3 => TraceLevel::Info,
93 4 => TraceLevel::Debug,
94 5 => TraceLevel::Trace,
95 _ => TraceLevel::Off,
96 }
97 }
98}
99
100#[derive(Clone)]
108pub struct HotGate {
109 level: Arc<AtomicU8>,
110}
111
112impl HotGate {
113 #[inline]
115 pub fn level(&self) -> TraceLevel {
116 TraceLevel::from_u8(self.level.load(Ordering::Relaxed))
117 }
118
119 #[inline]
123 pub fn records_calls(&self) -> bool {
124 self.level() >= TraceLevel::Debug
125 }
126
127 pub fn disabled() -> Self {
130 Self {
131 level: Arc::new(AtomicU8::new(TraceLevel::Off.as_u8())),
132 }
133 }
134
135 #[inline]
137 fn store(&self, level: TraceLevel) {
138 self.level.store(level.as_u8(), Ordering::Relaxed);
139 }
140}
141
142#[derive(Clone, Copy, PartialEq, Eq, Debug)]
144pub enum Direction {
145 HostImport,
147 GuestExport,
149}
150
151#[derive(Clone, PartialEq, Eq, Debug)]
153pub enum TraceOutcome {
154 Ok { micros: u64, fuel_delta: u64 },
156 Trap { kind: String, func: String },
159 Denied { capability: String },
161}
162
163#[derive(Clone, Debug)]
168pub struct PluginTraceRecord {
169 pub plugin: u32,
170 pub seam: PluginSeam,
171 pub direction: Direction,
172 pub call: std::borrow::Cow<'static, str>,
174 pub level: TraceLevel,
175 pub outcome: TraceOutcome,
176 pub detail: Option<String>,
179}
180
181#[derive(Clone, Debug)]
185pub struct PluginTracePushed {
186 pub record: PluginTraceRecord,
187}
188
189lattice_protocol::register_event!(
190 PluginTracePushed,
191 "plugin.trace-pushed",
192 "Fired when the PluginTracer appends a boundary-trace record.",
193 "lattice-plugin-host",
194);
195
196type TracePublisher = Box<dyn Fn(PluginTraceRecord) + Send + Sync>;
200
201pub type PluginTracerHandle = std::sync::Arc<PluginTracer>;
204
205struct TracerState {
206 global: Mutex<VecDeque<PluginTraceRecord>>,
208 per_plugin: Mutex<HashMap<u32, VecDeque<PluginTraceRecord>>>,
210 default_level: Mutex<TraceLevel>,
212 plugin_levels: Mutex<HashMap<u32, TraceLevel>>,
214 gates: Mutex<HashMap<u32, HotGate>>,
219 capacity: usize,
221 publisher: Mutex<Option<TracePublisher>>,
222}
223
224pub struct PluginTracer {
228 state: std::sync::Arc<TracerState>,
229}
230
231impl PluginTracer {
232 pub fn new(default_level: TraceLevel, capacity: usize) -> Self {
234 Self {
235 state: std::sync::Arc::new(TracerState {
236 global: Mutex::new(VecDeque::with_capacity(capacity.min(1024))),
237 per_plugin: Mutex::new(HashMap::new()),
238 default_level: Mutex::new(default_level),
239 plugin_levels: Mutex::new(HashMap::new()),
240 gates: Mutex::new(HashMap::new()),
241 capacity,
242 publisher: Mutex::new(None),
243 }),
244 }
245 }
246
247 pub fn with_defaults() -> Self {
250 Self::new(TraceLevel::Info, 10_000)
251 }
252
253 pub fn set_event_publisher(&self, publisher: TracePublisher) {
255 *lock(&self.state.publisher) = Some(publisher);
256 }
257
258 pub fn set_default_level(&self, level: TraceLevel) {
262 *lock(&self.state.default_level) = level;
263 let overrides = lock(&self.state.plugin_levels);
264 for (plugin, gate) in lock(&self.state.gates).iter() {
265 if !overrides.contains_key(plugin) {
266 gate.store(level);
267 }
268 }
269 }
270
271 pub fn set_plugin_level(&self, plugin: u32, level: TraceLevel) {
275 lock(&self.state.plugin_levels).insert(plugin, level);
276 if let Some(gate) = lock(&self.state.gates).get(&plugin) {
277 gate.store(level);
278 }
279 }
280
281 pub fn hot_gate(&self, plugin: u32) -> HotGate {
286 let effective = self.plugin_level(plugin);
287 lock(&self.state.gates)
288 .entry(plugin)
289 .or_insert_with(|| HotGate {
290 level: Arc::new(AtomicU8::new(effective.as_u8())),
291 })
292 .clone()
293 }
294
295 pub fn plugin_level(&self, plugin: u32) -> TraceLevel {
297 lock(&self.state.plugin_levels)
298 .get(&plugin)
299 .copied()
300 .unwrap_or_else(|| *lock(&self.state.default_level))
301 }
302
303 pub fn trace(&self, record: PluginTraceRecord) {
307 if record.level > self.plugin_level(record.plugin) {
308 return;
309 }
310 push_bounded(
311 lock(&self.state.per_plugin)
312 .entry(record.plugin)
313 .or_insert_with(|| VecDeque::with_capacity(self.state.capacity.min(1024))),
314 record.clone(),
315 self.state.capacity,
316 );
317 push_bounded(
318 &mut lock(&self.state.global),
319 record.clone(),
320 self.state.capacity,
321 );
322 if let Some(publisher) = lock(&self.state.publisher).as_ref() {
323 publisher(record);
324 }
325 }
326
327 pub fn snapshot_plugin(&self, plugin: u32) -> Vec<PluginTraceRecord> {
330 lock(&self.state.per_plugin)
331 .get(&plugin)
332 .map(|ring| ring.iter().cloned().collect())
333 .unwrap_or_default()
334 }
335
336 pub fn snapshot_global(&self) -> Vec<PluginTraceRecord> {
338 lock(&self.state.global).iter().cloned().collect()
339 }
340
341 pub fn forget_plugin(&self, plugin: u32) {
344 lock(&self.state.per_plugin).remove(&plugin);
345 lock(&self.state.plugin_levels).remove(&plugin);
346 lock(&self.state.gates).remove(&plugin);
347 }
348}
349
350fn push_bounded(
352 ring: &mut VecDeque<PluginTraceRecord>,
353 record: PluginTraceRecord,
354 capacity: usize,
355) {
356 if ring.len() >= capacity {
357 ring.pop_front();
358 }
359 ring.push_back(record);
360}
361
362fn lock<T>(m: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
365 m.lock().unwrap_or_else(|e| e.into_inner())
366}
367
368#[cfg(test)]
369mod tests {
370 #![allow(clippy::unwrap_used)]
371 use std::sync::Arc;
372
373 use super::*;
374
375 fn rec(plugin: u32, level: TraceLevel) -> PluginTraceRecord {
376 PluginTraceRecord {
377 plugin,
378 seam: PluginSeam::Grammar,
379 direction: Direction::GuestExport,
380 call: "apply-motion".into(),
381 level,
382 outcome: TraceOutcome::Ok {
383 micros: 12,
384 fuel_delta: 100,
385 },
386 detail: None,
387 }
388 }
389
390 #[test]
391 fn a_record_at_or_below_the_gate_is_kept() {
392 let t = PluginTracer::new(TraceLevel::Info, 8);
393 t.trace(rec(1, TraceLevel::Error)); t.trace(rec(1, TraceLevel::Info)); assert_eq!(t.snapshot_plugin(1).len(), 2);
396 assert_eq!(t.snapshot_global().len(), 2);
397 }
398
399 #[test]
400 fn a_record_above_the_gate_is_dropped() {
401 let t = PluginTracer::new(TraceLevel::Info, 8);
402 t.trace(rec(1, TraceLevel::Debug)); t.trace(rec(1, TraceLevel::Trace)); assert!(t.snapshot_plugin(1).is_empty());
405 assert!(t.snapshot_global().is_empty());
406 }
407
408 #[test]
409 fn a_per_plugin_override_raises_verbosity_for_only_that_plugin() {
410 let t = PluginTracer::new(TraceLevel::Info, 8);
411 t.set_plugin_level(1, TraceLevel::Trace);
412 t.trace(rec(1, TraceLevel::Debug)); t.trace(rec(2, TraceLevel::Debug)); assert_eq!(t.snapshot_plugin(1).len(), 1);
415 assert!(t.snapshot_plugin(2).is_empty());
416 }
417
418 #[test]
419 fn off_gate_silences_everything() {
420 let t = PluginTracer::new(TraceLevel::Info, 8);
421 t.set_plugin_level(1, TraceLevel::Off);
422 t.trace(rec(1, TraceLevel::Error)); assert!(t.snapshot_plugin(1).is_empty());
424 }
425
426 #[test]
427 fn rings_are_bounded_evicting_oldest() {
428 let t = PluginTracer::new(TraceLevel::Trace, 3);
429 for _ in 0..5 {
430 t.trace(rec(1, TraceLevel::Info));
431 }
432 assert_eq!(t.snapshot_plugin(1).len(), 3, "per-plugin ring capped at 3");
433 assert_eq!(t.snapshot_global().len(), 3, "global ring capped at 3");
434 }
435
436 #[test]
437 fn publish_fires_on_every_kept_append() {
438 let t = PluginTracer::new(TraceLevel::Info, 8);
439 let count = Arc::new(std::sync::atomic::AtomicUsize::new(0));
440 let c = count.clone();
441 t.set_event_publisher(Box::new(move |_rec| {
442 c.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
443 }));
444 t.trace(rec(1, TraceLevel::Info)); t.trace(rec(1, TraceLevel::Debug)); assert_eq!(count.load(std::sync::atomic::Ordering::Relaxed), 1);
447 }
448
449 #[test]
450 fn trace_level_round_trips_through_the_gate_byte() {
451 for lvl in [
452 TraceLevel::Off,
453 TraceLevel::Error,
454 TraceLevel::Warn,
455 TraceLevel::Info,
456 TraceLevel::Debug,
457 TraceLevel::Trace,
458 ] {
459 assert_eq!(TraceLevel::from_u8(lvl.as_u8()), lvl);
460 }
461 assert_eq!(TraceLevel::from_u8(200), TraceLevel::Off);
463 }
464
465 #[test]
466 fn cycle_next_walks_all_and_wraps() {
467 let mut seen = Vec::new();
468 let mut lvl = TraceLevel::Off;
469 for _ in 0..TraceLevel::ALL.len() {
470 seen.push(lvl);
471 lvl = lvl.cycle_next();
472 }
473 assert_eq!(
474 seen,
475 TraceLevel::ALL.to_vec(),
476 "cycle visits every level in order"
477 );
478 assert_eq!(lvl, TraceLevel::Off, "Trace wraps back to Off");
479 assert_eq!(TraceLevel::Info.as_str(), "info");
480 }
481
482 #[test]
483 fn parse_is_the_inverse_of_as_str() {
484 for lvl in TraceLevel::ALL {
485 assert_eq!(TraceLevel::parse(lvl.as_str()), Some(lvl));
486 }
487 assert_eq!(TraceLevel::parse("loud"), None);
488 }
489
490 #[test]
491 fn a_hot_gate_seeds_to_the_effective_level() {
492 let t = PluginTracer::new(TraceLevel::Info, 8);
493 assert_eq!(t.hot_gate(1).level(), TraceLevel::Info);
495 assert!(!t.hot_gate(1).records_calls());
497 t.set_plugin_level(2, TraceLevel::Trace);
499 let g2 = t.hot_gate(2);
500 assert_eq!(g2.level(), TraceLevel::Trace);
501 assert!(g2.records_calls());
502 }
503
504 #[test]
505 fn setting_a_plugin_level_republishes_to_its_live_gate() {
506 let t = PluginTracer::new(TraceLevel::Info, 8);
507 let gate = t.hot_gate(1);
508 assert!(!gate.records_calls(), "starts off at the Info default");
509 t.set_plugin_level(1, TraceLevel::Debug);
510 assert!(
511 gate.records_calls(),
512 "the already-handed-out gate sees the raise"
513 );
514 t.set_plugin_level(1, TraceLevel::Off);
515 assert_eq!(gate.level(), TraceLevel::Off, "and the lowering");
516 }
517
518 #[test]
519 fn raising_the_default_republishes_only_to_unoverridden_gates() {
520 let t = PluginTracer::new(TraceLevel::Info, 8);
521 let g1 = t.hot_gate(1); let g2 = t.hot_gate(2);
523 t.set_plugin_level(2, TraceLevel::Off); t.set_default_level(TraceLevel::Trace);
525 assert_eq!(
526 g1.level(),
527 TraceLevel::Trace,
528 "the unoverridden gate follows the default"
529 );
530 assert_eq!(
531 g2.level(),
532 TraceLevel::Off,
533 "the overridden gate is untouched"
534 );
535 }
536
537 #[test]
538 fn hot_gate_is_cached_so_clones_share_the_atomic() {
539 let t = PluginTracer::new(TraceLevel::Info, 8);
540 let a = t.hot_gate(1);
541 let b = t.hot_gate(1);
542 t.set_plugin_level(1, TraceLevel::Trace);
543 assert!(a.records_calls());
544 assert!(b.records_calls(), "the second handle sees the same store");
545 }
546
547 #[test]
548 fn a_disabled_gate_never_records() {
549 let g = HotGate::disabled();
550 assert_eq!(g.level(), TraceLevel::Off);
551 assert!(!g.records_calls());
552 }
553
554 #[test]
555 fn a_poisoned_publisher_mutex_still_records_and_never_panics() {
556 let t = PluginTracer::new(TraceLevel::Info, 8);
561 t.set_event_publisher(Box::new(|_| panic!("boom in the publisher closure")));
562 let poisoned = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
565 t.trace(rec(1, TraceLevel::Info));
566 }));
567 assert!(poisoned.is_err(), "the panicking publisher unwound");
568 assert_eq!(
569 t.snapshot_plugin(1).len(),
570 1,
571 "record #1 landed pre-publish"
572 );
573 t.set_event_publisher(Box::new(|_| {}));
576 t.trace(rec(1, TraceLevel::Info));
577 assert_eq!(
578 t.snapshot_plugin(1).len(),
579 2,
580 "the tracer recovered the poisoned mutex and kept recording"
581 );
582 }
583
584 #[test]
585 fn forget_plugin_drops_its_ring_and_override() {
586 let t = PluginTracer::new(TraceLevel::Trace, 8);
587 t.set_plugin_level(1, TraceLevel::Trace);
588 t.trace(rec(1, TraceLevel::Trace));
589 assert_eq!(t.snapshot_plugin(1).len(), 1);
590 let gate = t.hot_gate(1);
592 assert_eq!(gate.level(), TraceLevel::Trace);
593 t.forget_plugin(1);
594 assert!(t.snapshot_plugin(1).is_empty());
595 assert_eq!(t.snapshot_global().len(), 1);
597 assert_eq!(
599 t.hot_gate(1).level(),
600 TraceLevel::Trace,
601 "default is Trace here"
602 );
603 }
604}