Skip to main content

moonpool_sim/observability/
mod.rs

1//! Trace-based observability for moonpool simulations.
2//!
3//! Processes and workloads emit ordinary `tracing` events — the exact same
4//! instrumentation production observability consumes (`fmt`, OpenTelemetry,
5//! Loki, ...). In simulation, [`SimulationLayer`] captures those events into
6//! a timeline and registered [`Invariant`]s cross-validate them after every
7//! simulation step. The point: anomalies like dual leaders or metastable
8//! failures are detected in simulation with the same signals you would use
9//! to find them in production — traces.
10//!
11//! # Emission convention
12//!
13//! Emit a plain `tracing` event with a constant message (the event name) and
14//! structured fields. No special markers, no derives:
15//!
16//! ```ignore
17//! tracing::info!(target: "raft", term, leader = %my_ip, "leader_elected");
18//! ```
19//!
20//! - Use `%` (`Display`) for strings and IPs — `?` (`Debug`) on a `String`
21//!   includes quotes.
22//! - The message must be a constant name (`"leader_elected"`), not an
23//!   interpolated sentence: events are grouped and queried by name.
24//!
25//! An event is captured iff it is `INFO` or more severe, carries a non-empty
26//! message, and fires inside a process/workload span (the orchestrator wraps
27//! every actor task in a span carrying its `ip`, which becomes the event's
28//! `source`). In production, where no [`SimulationLayer`] is installed, the
29//! same emission flows to whatever subscriber is configured.
30//!
31//! # Querying
32//!
33//! Invariants receive a [`TraceQuery`] and read events by name, extracting
34//! typed fields by key — the same way you would query a production trace
35//! store:
36//!
37//! ```ignore
38//! fn observe(&self, q: &dyn TraceQuery, _sim_time_ms: u64) {
39//!     for e in q.since("leader_elected", &self.cursor) {
40//!         let term = e.u64("term");
41//!         let leader = e.str("leader");
42//!         // assert one leader per term...
43//!     }
44//! }
45//! ```
46//!
47//! Sim-injected faults (partitions, kills, storage corruption) are recorded
48//! by the runner under the [`crate::SIM_FAULT_EVENT_NAME`] name with a
49//! `kind` field and `source = "sim"`, so invariants can correlate
50//! application anomalies with infrastructure faults in one timeline.
51
52pub mod event;
53pub mod fmt;
54pub mod init;
55pub mod invariant;
56pub mod layer;
57pub mod query;
58
59pub use event::{FieldValue, TraceEvent};
60pub use fmt::{Clock, SimTime};
61pub use init::init_sim_tracing;
62pub use invariant::{Invariant, invariant_fn};
63pub use layer::{InstallGuard, SimulationLayer, SimulationLayerHandle};
64pub use query::TraceQuery;
65
66#[cfg(test)]
67mod tests {
68    use super::*;
69    use std::cell::Cell;
70    use std::sync::Arc;
71    use std::sync::atomic::{AtomicUsize, Ordering};
72
73    fn in_actor_span(ip: &str, f: impl FnOnce()) {
74        let span = tracing::info_span!("process", ip = %ip);
75        let _enter = span.enter();
76        f();
77    }
78
79    // ------------------------------------------------------------------
80    // Capture-path tests
81    // ------------------------------------------------------------------
82
83    #[test]
84    fn emit_without_layer_is_safe() {
85        // No layer installed: tracing::info! does not panic and there's
86        // nothing to observe. We just check the macro compiles and runs.
87        tracing::info!(term = 1_u64, leader = %"10.0.1.1", "leader_elected");
88    }
89
90    #[test]
91    fn layer_captures_event_inside_actor_span() {
92        let (handle, _guard) = SimulationLayer::new().install();
93
94        handle.set_sim_time_ms(1234);
95        in_actor_span("10.0.1.1", || {
96            tracing::info!(term = 7_u64, leader = %"10.0.1.1", "leader_elected");
97        });
98
99        let entries = handle.snapshot("leader_elected");
100        assert_eq!(entries.len(), 1);
101        let e = &entries[0];
102        assert_eq!(e.name, "leader_elected");
103        assert_eq!(e.source, "10.0.1.1");
104        assert_eq!(e.seq, 0);
105        assert_eq!(e.time_ms, 1234);
106        assert_eq!(e.level, tracing::Level::INFO);
107        assert_eq!(e.u64("term"), Some(7));
108        assert_eq!(e.str("leader"), Some("10.0.1.1"));
109    }
110
111    #[test]
112    fn event_outside_actor_span_is_dropped() {
113        let (handle, _guard) = SimulationLayer::new().install();
114
115        tracing::info!(term = 7_u64, "leader_elected");
116
117        assert!(handle.snapshot("leader_elected").is_empty());
118    }
119
120    #[test]
121    fn debug_level_is_dropped() {
122        let (handle, _guard) = SimulationLayer::new().install();
123
124        in_actor_span("10.0.1.1", || {
125            tracing::debug!(term = 7_u64, "leader_elected");
126            tracing::trace!(term = 8_u64, "leader_elected");
127            tracing::warn!(term = 9_u64, "leader_elected");
128        });
129
130        let entries = handle.snapshot("leader_elected");
131        assert_eq!(entries.len(), 1, "only the WARN event is captured");
132        assert_eq!(entries[0].u64("term"), Some(9));
133        assert_eq!(entries[0].level, tracing::Level::WARN);
134    }
135
136    #[test]
137    fn event_without_message_is_dropped() {
138        let (handle, _guard) = SimulationLayer::new().install();
139
140        in_actor_span("10.0.1.1", || {
141            tracing::info!(term = 7_u64);
142        });
143
144        // No message → no name to group under; nothing captured anywhere.
145        assert_eq!(handle.len("leader_elected"), 0);
146    }
147
148    #[test]
149    fn nearest_enclosing_span_wins() {
150        let (handle, _guard) = SimulationLayer::new().install();
151
152        in_actor_span("10.0.1.1", || {
153            in_actor_span("10.0.1.2", || {
154                tracing::info!("ping_sent");
155            });
156        });
157
158        let entries = handle.snapshot("ping_sent");
159        assert_eq!(entries.len(), 1);
160        assert_eq!(entries[0].source, "10.0.1.2", "innermost ip attributes");
161    }
162
163    #[test]
164    fn field_accessors_extract_typed_values() {
165        let (handle, _guard) = SimulationLayer::new().install();
166
167        in_actor_span("10.0.1.1", || {
168            tracing::info!(
169                count = 5_u64,
170                delta = -3_i64,
171                ratio = 0.5_f64,
172                ok = true,
173                name = %"alpha",
174                detail = ?vec![1, 2],
175                "mixed_fields"
176            );
177        });
178
179        let entries = handle.snapshot("mixed_fields");
180        assert_eq!(entries.len(), 1);
181        let e = &entries[0];
182        assert_eq!(e.u64("count"), Some(5));
183        assert_eq!(e.i64("count"), Some(5), "u64 readable as i64 when it fits");
184        assert_eq!(e.i64("delta"), Some(-3));
185        assert_eq!(e.u64("delta"), None, "negative i64 not readable as u64");
186        assert!((e.f64("ratio").expect("ratio field") - 0.5).abs() < f64::EPSILON);
187        assert_eq!(e.bool("ok"), Some(true));
188        assert_eq!(e.str("name"), Some("alpha"), "% display value unquoted");
189        assert_eq!(e.str("detail"), Some("[1, 2]"), "? debug formatting");
190        assert_eq!(e.u64("missing"), None);
191    }
192
193    #[test]
194    fn seq_is_monotonic_across_names() {
195        let (handle, _guard) = SimulationLayer::new().install();
196
197        in_actor_span("10.0.1.1", || {
198            tracing::info!("event_a");
199            tracing::info!("event_b");
200            tracing::info!("event_a");
201        });
202
203        let a = handle.snapshot("event_a");
204        let b = handle.snapshot("event_b");
205        assert_eq!(a[0].seq, 0);
206        assert_eq!(b[0].seq, 1);
207        assert_eq!(a[1].seq, 2);
208    }
209
210    #[test]
211    fn cursor_since_returns_only_new_entries() {
212        let (handle, _guard) = SimulationLayer::new().install();
213
214        let cursor = Cell::new(0);
215        in_actor_span("10.0.1.1", || {
216            tracing::info!(n = 1_u64, "hb");
217        });
218        assert_eq!(handle.since("hb", &cursor).len(), 1);
219        assert!(handle.since("hb", &cursor).is_empty(), "cursor advanced");
220
221        in_actor_span("10.0.1.1", || {
222            tracing::info!(n = 2_u64, "hb");
223        });
224        let new = handle.since("hb", &cursor);
225        assert_eq!(new.len(), 1);
226        assert_eq!(new[0].u64("n"), Some(2));
227    }
228
229    #[test]
230    fn reset_for_seed_clears_state() {
231        let (handle, _guard) = SimulationLayer::new().install();
232
233        in_actor_span("10.0.1.1", || {
234            tracing::info!(n = 1_u64, "hb");
235        });
236        handle.set_sim_time_ms(500);
237        handle.reset_for_seed();
238        in_actor_span("10.0.1.1", || {
239            tracing::info!(n = 2_u64, "hb");
240        });
241
242        let entries = handle.snapshot("hb");
243        assert_eq!(entries.len(), 1, "reset cleared the first event");
244        assert_eq!(entries[0].seq, 0, "seq counter resets per seed");
245        assert_eq!(handle.current_sim_time_ms(), 0, "clock resets per seed");
246    }
247
248    #[test]
249    fn run_invariants_pumps_registered_invariants() {
250        let (handle, _guard) = SimulationLayer::new().install();
251
252        let observed = Arc::new(AtomicUsize::new(0));
253        let observed_clone = observed.clone();
254        let cursor = Cell::new(0);
255        handle.register(invariant_fn("counter", move |q, _t| {
256            let new = q.since("hb", &cursor);
257            observed_clone.fetch_add(new.len(), Ordering::Relaxed);
258        }));
259
260        in_actor_span("10.0.1.1", || {
261            tracing::info!(n = 1_u64, "hb");
262            tracing::info!(n = 2_u64, "hb");
263        });
264        assert_eq!(
265            observed.load(Ordering::Relaxed),
266            0,
267            "invariants do not run inside tracing dispatch"
268        );
269
270        handle.run_invariants();
271        assert_eq!(observed.load(Ordering::Relaxed), 2, "batched at pump time");
272
273        handle.run_invariants();
274        assert_eq!(observed.load(Ordering::Relaxed), 2, "cursor advanced");
275    }
276
277    #[test]
278    fn record_sim_fault_lands_in_timeline() {
279        use crate::chaos::{SIM_FAULT_EVENT_NAME, SimFaultEvent};
280
281        let (handle, _guard) = SimulationLayer::new().install();
282
283        handle.record_sim_fault(
284            42,
285            &SimFaultEvent::PartitionCreated {
286                from: "10.0.1.1".to_owned(),
287                to: "10.0.1.2".to_owned(),
288            },
289        );
290        handle.record_sim_fault(
291            43,
292            &SimFaultEvent::ProcessForceKill {
293                ip: "10.0.1.1".to_owned(),
294            },
295        );
296
297        let entries = handle.snapshot(SIM_FAULT_EVENT_NAME);
298        assert_eq!(entries.len(), 2);
299        assert_eq!(entries[0].source, "sim");
300        assert_eq!(entries[0].time_ms, 42);
301        assert_eq!(entries[0].str("kind"), Some("partition_created"));
302        assert_eq!(entries[0].str("from"), Some("10.0.1.1"));
303        assert_eq!(entries[0].str("to"), Some("10.0.1.2"));
304        assert_eq!(entries[1].str("kind"), Some("process_force_kill"));
305        assert_eq!(entries[1].str("ip"), Some("10.0.1.1"));
306    }
307
308    // ------------------------------------------------------------------
309    // SimTime fmt-integration tests (formatter, not the capture path)
310    // ------------------------------------------------------------------
311
312    #[test]
313    fn sim_time_format_writes_seconds_and_millis() {
314        use tracing_subscriber::fmt::format::Writer;
315        use tracing_subscriber::fmt::time::FormatTime;
316
317        let layer = SimulationLayer::new();
318        let handle = layer.handle();
319        handle.set_sim_time_ms(7042);
320
321        let st = SimTime::new(handle.clone());
322        let mut buf = String::new();
323        let mut writer = Writer::new(&mut buf);
324        st.format_time(&mut writer)
325            .expect("writing to a String never fails");
326        assert_eq!(buf, "sim+    7.042s");
327    }
328
329    #[test]
330    fn clock_trait_can_be_implemented_for_a_stub() {
331        use tracing_subscriber::fmt::format::Writer;
332        use tracing_subscriber::fmt::time::FormatTime;
333
334        struct FixedClock(u64);
335        impl crate::observability::Clock for FixedClock {
336            fn now_ms(&self) -> u64 {
337                self.0
338            }
339        }
340
341        let st = SimTime::new(FixedClock(123_456));
342        let mut buf = String::new();
343        let mut writer = Writer::new(&mut buf);
344        st.format_time(&mut writer)
345            .expect("writing to a String never fails");
346        assert_eq!(buf, "sim+  123.456s");
347    }
348
349    #[test]
350    fn fmt_layer_with_sim_time_prefixes_log_output() {
351        use std::io;
352        use std::sync::Mutex;
353
354        use tracing_subscriber::Layer as _;
355        use tracing_subscriber::layer::SubscriberExt;
356
357        #[derive(Clone, Default)]
358        struct VecWriter(Arc<Mutex<Vec<u8>>>);
359
360        impl io::Write for VecWriter {
361            fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
362                self.0
363                    .lock()
364                    .expect("test writer poisoned")
365                    .extend_from_slice(buf);
366                Ok(buf.len())
367            }
368            fn flush(&mut self) -> io::Result<()> {
369                Ok(())
370            }
371        }
372
373        impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for VecWriter {
374            type Writer = VecWriter;
375            fn make_writer(&'a self) -> Self::Writer {
376                self.clone()
377            }
378        }
379
380        let layer = SimulationLayer::new();
381        let handle = layer.handle();
382        handle.set_sim_time_ms(5_000);
383
384        let writer = VecWriter::default();
385        let buf = writer.0.clone();
386
387        let fmt_layer = tracing_subscriber::fmt::layer()
388            .with_writer(writer)
389            .with_ansi(false)
390            .with_timer(SimTime::new(handle.clone()))
391            .with_filter(tracing_subscriber::filter::LevelFilter::INFO);
392        let subscriber = tracing_subscriber::registry().with(layer).with(fmt_layer);
393
394        tracing::subscriber::with_default(subscriber, || {
395            tracing::info!("hello at 5s");
396            handle.set_sim_time_ms(12_345);
397            tracing::info!("hello at 12.345s");
398        });
399
400        let output = String::from_utf8(buf.lock().expect("buf").clone()).expect("utf-8 fmt output");
401        assert!(
402            output.contains("sim+    5.000s"),
403            "expected sim+5s prefix; got: {output}"
404        );
405        assert!(
406            output.contains("sim+   12.345s"),
407            "expected sim+12.345s prefix after clock advance; got: {output}"
408        );
409    }
410}