moonpool_sim/observability/
mod.rs1pub 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 #[test]
84 fn emit_without_layer_is_safe() {
85 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 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 #[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}