hydro_lang/sim/compiled.rs
1//! Interfaces for compiled Hydro simulators and concrete simulation instances.
2//!
3//! # Quiescence and observation soundness
4//!
5//! The scheduler distinguishes two kinds of simulation work:
6//! - **Deterministic work**: running the top-level async dataflows, which simply propagate
7//! whatever data is already in flight. This makes no `nondet!` decisions, so running it can
8//! never change which executions are explored.
9//! - **Nondeterministic work**: running ticks and observations, whose behavior depends on
10//! decisions drawn from the bolero driver (batch boundaries, snapshot versions, message
11//! orderings). Each decision forks the space of possible executions.
12//!
13//! The simulation is **quiescent** when neither kind of work can make progress without new
14//! external input. Test-side observations (the methods on [`SimReceiver`] /
15//! [`SimClusterReceiver`]) interact with the scheduler while waiting, and the key soundness
16//! question is: *when is it okay for an observation to let nondeterministic work run?*
17//!
18//! **Waiting for a message is always sound.** If the message eventually arrives, the work
19//! that ran was necessary to produce it (schedules that run *extra* work are also valid
20//! executions and are explored separately). If the simulation instead quiesces without
21//! producing the message, the assertion fails and the instance ends, so nothing can observe
22//! the overrun. This is why [`SimReceiver::next`], [`SimReceiver::collect_n`], and the
23//! `assert_yields*` prefix checks are safe to use in the middle of a test.
24//!
25//! **Observing the *absence* of a message is dangerous.** Proving that "no more messages can
26//! arrive" requires driving the simulation all the way to quiescence, running *all* pending
27//! nondeterministic work. A later assertion may have needed to observe a state where that
28//! work had not yet run — e.g., `assert_yields_only([1, 2])` followed by reading a counter
29//! must be able to see the counter *before* the ticks that count `1` and `2` have fired.
30//! Forcing quiescence at the first assertion would make some executions unobservable, and
31//! extra messages produced by the forced work could surface at a *later* assertion,
32//! misattributing the failure. Absence-observing APIs therefore proceed in phases:
33//!
34//! 1. **Settle** (see `SettlePauseGuard::poll_settle`): the scheduler runs only deterministic work, pausing
35//! just before nondeterministic work. If the simulation reaches quiescence this way, the
36//! end-of-stream check is *free* — no decision was forced, no execution was cut off — and
37//! the test simply continues.
38//! 2. If nondeterministic work is pending, the check would overrun. What happens next depends
39//! on the API and engine:
40//! - The assertion APIs ([`SimReceiver::assert_no_more`], `assert_yields_only*`,
41//! `collect_n_only`) under [`CompiledSim::exhaustive`] **fork** the search on a bolero
42//! decision: one instance performs the check and then ends (via a discard panic, like
43//! `sim::continue_if!`), while sibling instances skip the check entirely and continue. The
44//! exhaustive driver enumerates the checking instance *first*, so a failing check is
45//! found before any instance runs past it — with a decision trace that leads exactly to
46//! the failing assertion. Since nothing after the check runs in the checking instance,
47//! the overrun it performs is unobservable, and the continuing instances never quiesce,
48//! so every downstream state remains reachable.
49//! - Otherwise (fuzz / RNG / replay engines, or the drain-everything APIs
50//! [`SimReceiver::try_next`], [`SimReceiver::collect`], and `collect_sorted` in every
51//! mode), the pending work runs and the instance is **tainted**
52//! (`QuiescenceState::tainted`). Reads of the now-quiescent state remain sound (they
53//! observe a fully-drained simulation that can no longer advance), so tests may drain
54//! multiple output ports at the end. But once new input is sent, the instance is
55//! **poisoned** (`QuiescenceState::poisoned`): any further receive panics (see
56//! `guard_not_poisoned`), because a failure observed after the forced overrun could
57//! have been caused by it and attributed to the wrong assertion.
58//!
59//! NOTE: This module runs inside bolero's `catch_unwind` scope, which silently
60//! swallows panics. Internal invariant checks should use `abort_assert!`
61//! rather than `panic!`/`assert!`.
62//!
63//! TODO(mingwei): Panics inside the tick DFIR (generated code in the dylib) are
64//! also caught by bolero's `catch_unwind`. Consider a mechanism to detect and
65//! propagate those as well.
66
67/// Like `assert!`, but calls `std::process::abort()` instead of `panic!()`.
68/// Use for internal invariants that must not be silently caught by bolero.
69macro_rules! abort_assert {
70 ($cond:expr, $($arg:tt)*) => {
71 if !$cond {
72 eprintln!("Simulator internal error: {}", format!($($arg)*));
73 std::process::abort();
74 }
75 };
76}
77
78use core::{fmt, panic};
79use std::cell::{Cell, RefCell};
80use std::collections::{HashMap, VecDeque};
81use std::fmt::Debug;
82use std::panic::RefUnwindSafe;
83use std::path::Path;
84use std::pin::{Pin, pin};
85use std::rc::Rc;
86use std::task::{Poll, ready};
87
88use bytes::Bytes;
89use colored::Colorize;
90use dfir_rs::scheduled::context::DfirErased;
91use dfir_rs::util::unsync::mpsc::{Receiver as UnsyncReceiver, Sender as UnsyncSender};
92use futures::StreamExt;
93use libloading::Library;
94use tokio::sync::{Mutex, Notify};
95
96use super::runtime::{
97 Hooks, InlineHooks, ObservationHooks, ScriptTarget, ScriptedHookControl, ScriptedHookRegistry,
98 ScriptedInlineHooks, ScriptedObservationHooks, ScriptedTickHooks, SimLocation,
99};
100use super::{SimClusterReceiver, SimClusterSender, SimReceiver, SimSender};
101use crate::compile::builder::ExternalPortId;
102use crate::compile::trybuild::generate::BuiltArtifact;
103use crate::live_collections::stream::{ExactlyOnce, NoOrder, Ordering, Retries, TotalOrder};
104use crate::location::dynamic::LocationId;
105use crate::sim::graph::{SimExternalPort, SimExternalPortRegistry};
106use crate::sim::runtime::{
107 InlineHook, ObservationHook, ScriptedObservationHook, ScriptedTickInputHook, TickInputHook,
108};
109
110struct QuiescenceState {
111 /// Set to true when the scheduler reaches quiescence; reset to false when new input is sent.
112 quiescent: Cell<bool>,
113 /// Notified when the scheduler reaches quiescence (wakes receivers waiting for data).
114 quiescence_notify: Notify,
115 /// Notified when new input is sent, signaling the scheduler to resume.
116 resume_notify: Notify,
117 /// When nonzero, the scheduler must not start nondeterministic work (ticks /
118 /// observations): once only such work remains, it sets `nondet_pending` and pauses until
119 /// resumed. Used by receivers to query whether the simulation can quiesce
120 /// deterministically. This is a count (not a bool) because multiple settling futures can
121 /// be in flight at once (e.g. `select!`/`join!` between two receiver awaits): the
122 /// scheduler must stay paused until *every* one of them has finished settling.
123 pause_nondet: Cell<usize>,
124 /// Set while the scheduler is paused because nondeterministic work is ready to run but
125 /// `pause_nondet` is set.
126 nondet_pending: Cell<bool>,
127 /// Wakers for test-side tasks waiting for the scheduler to settle (either quiesce or set
128 /// `nondet_pending`) while `pause_nondet` is set. Also used by scripting futures that
129 /// need to be woken when the scheduler parks.
130 settle_wakers: RefCell<Vec<std::task::Waker>>,
131 /// Set when an observation *forced* the simulation to quiesce (running pending
132 /// nondeterministic work) outside of exhaustive mode's forking. Further observations of
133 /// the quiescent state remain sound, but once new input is sent (see `poisoned`), later
134 /// observations could misattribute failures caused by the forced overrun.
135 tainted: Cell<bool>,
136 /// Set when new input is sent after `tainted`; all further receives panic.
137 poisoned: Cell<bool>,
138}
139
140impl QuiescenceState {
141 /// Signal that new input has been sent, waking the scheduler if it was quiescent.
142 fn resume(&self) {
143 if self.tainted.get() {
144 self.poisoned.set(true);
145 }
146 self.quiescent.set(false);
147 // `notify_one` (rather than `notify_waiters`) stores a permit if the scheduler driver
148 // is not currently parked on [`Self::resumed`], so a resume that fires before the
149 // driver parks (e.g. input sent while the driver is polling the thunk) is not lost.
150 self.resume_notify.notify_one();
151 }
152
153 /// Whether the scheduler is currently quiescent (no more progress possible without input).
154 fn is_quiescent(&self) -> bool {
155 self.quiescent.get()
156 }
157
158 /// Returns a future that completes when the scheduler next reaches quiescence.
159 fn notified(&self) -> tokio::sync::futures::Notified<'_> {
160 self.quiescence_notify.notified()
161 }
162
163 /// Wakes test-side tasks waiting for the scheduler to settle.
164 fn wake_settled(&self) {
165 for waker in self.settle_wakers.borrow_mut().drain(..) {
166 waker.wake();
167 }
168 }
169
170 /// Enter quiescence, waking receivers waiting for data (their streams end). The scheduler
171 /// driver is responsible for parking until [`Self::resume`] is called with new input.
172 fn enter_quiescence(&self) {
173 self.quiescent.set(true);
174 self.quiescence_notify.notify_waiters();
175 self.wake_settled();
176 }
177
178 /// Completes when new input arrives (via [`Self::resume`]).
179 async fn resumed(&self) {
180 self.resume_notify.notified().await;
181 }
182
183 /// Registers a waker to be woken the next time the scheduler parks (quiescence or
184 /// settle-pause). Used by scripting futures: while the scheduler is running, the test
185 /// body is re-polled after every step anyway, so a waker is only needed for the parked
186 /// cases. Duplicate registrations are harmless.
187 fn push_park_waker(&self, waker: &std::task::Waker) {
188 self.settle_wakers.borrow_mut().push(waker.clone());
189 }
190}
191
192/// The **current group** of scripted decisions: consecutive decision calls in the test body
193/// that target different hooks of the same tick form a group, describing one execution of
194/// that tick. At most one group's decisions are ever installed at a time; the first decision
195/// call of the *next* group suspends until the current group's tick execution has consumed
196/// every installed decision.
197pub(crate) struct CurrentGroup {
198 /// The scheduler action the group's decisions apply to.
199 target: ScriptTarget,
200 /// The registry keys (handle ID plus cluster member) with an installed decision in
201 /// this group.
202 members: Vec<(usize, Option<u32>)>,
203 /// Set when the scheduler starts a step. Until then, consecutive decisions for different
204 /// hooks of this tick may join the group in the same poll of the test body.
205 sealed: bool,
206}
207
208/// Coordinates the script protocol between test-side hook handles and the scheduler.
209#[derive(Default)]
210pub(crate) struct ScriptCoordinator {
211 /// `Some` means exactly one decision group is outstanding. The scheduler clears it only
212 /// after that group's tick executes, so the test cannot replace an unconsumed group.
213 current: Option<CurrentGroup>,
214 /// Set by the scheduler at each quiescence: `true` when the outstanding group is stuck
215 /// even though every queued decision is satisfiable, because none of them can trigger
216 /// the tick (and no unscripted input on the tick can trigger it either — otherwise the
217 /// tick would be runnable and the simulation would not be quiescent). Selects the
218 /// stuck-script error style rendered at the suspended test-side await; `false` means
219 /// some decision is waiting on input that can never arrive.
220 stuck_cannot_trigger: bool,
221}
222
223impl ScriptCoordinator {
224 /// Describes the not-yet-consumed decisions of the current group, one per line
225 /// (without a trailing newline), for error messages. `None` when no group is
226 /// outstanding or every decision is consumed.
227 fn describe_unconsumed(&self, hooks: &ScriptedHookRegistry) -> Option<String> {
228 let group = self.current.as_ref()?;
229 let mut out = String::new();
230 for key in &group.members {
231 let hook = hooks.get(key).unwrap().borrow();
232 if let Some(decision) = hook.describe_decision() {
233 use std::fmt::Write;
234 if !out.is_empty() {
235 out.push('\n');
236 }
237 let member = key
238 .1
239 .map(|m| format!(" (cluster member {m})"))
240 .unwrap_or_default();
241 write!(
242 out,
243 " {} is waiting on the hook at {}{}, which has {}",
244 decision,
245 hook.location_meta().location,
246 member,
247 hook.describe_pending()
248 .as_deref()
249 .unwrap_or("no pending input"),
250 )
251 .unwrap();
252 }
253 }
254 (!out.is_empty()).then_some(out)
255 }
256}
257
258/// The per-instance scripting context, resolved through the task-local sim connections.
259///
260/// The three `Rc`s are genuinely distinct (not one shared allocation) because they have
261/// different owners and lifetimes: the hook registry only materializes when the dylib is
262/// launched (it is part of the `DylibResult`), while the coordinator and quiescence
263/// state live in the pre-launch `SimConnections` and are independently shared with
264/// receivers and test-side handles (quiescence is also used by non-scripting paths).
265/// This struct is the bundle of all three, assembled by-clone at resolution time.
266pub(crate) struct ScriptCtx {
267 hooks: Rc<ScriptedHookRegistry>,
268 coordinator: Rc<RefCell<ScriptCoordinator>>,
269 quiescence: Rc<QuiescenceState>,
270}
271
272/// The result of attempting to schedule one decision; see
273/// [`ScriptCtx::try_schedule_decision`].
274pub(crate) enum ScheduleDecision {
275 /// The decision was installed into the current group.
276 Installed,
277 /// The previous group has not been consumed yet; the decision blob is handed back and
278 /// the caller should retry after the scheduler makes progress.
279 Wait(Vec<u8>),
280}
281
282const UNBOUND_HOOK_ERROR: &str = "this sim hook handle is not bound to any operator in the simulated flow; \
283 attach it with `nondet!(... hook = handle)` at the operator it should control";
284
285impl ScriptCtx {
286 /// Resolves a hook handle's scripted hook instance; `member` selects a cluster
287 /// member's instance (`None` for hooks on processes). Panics if the handle is not
288 /// bound to an operator or the member does not exist.
289 #[track_caller]
290 pub(crate) fn control(
291 &self,
292 hook_id: usize,
293 member: Option<u32>,
294 ) -> Rc<RefCell<dyn ScriptedHookControl>> {
295 if let Some(hook) = self.hooks.get(&(hook_id, member)) {
296 return hook.clone();
297 }
298
299 // The exact instance is missing; distinguish the misuse cases from a handle
300 // that was never bound at all.
301 let bound_members: Vec<u32> = self
302 .hooks
303 .range((hook_id, None)..=(hook_id, Some(u32::MAX)))
304 .filter_map(|((_, m), _)| *m)
305 .collect();
306 match member {
307 None if !bound_members.is_empty() => panic!(
308 "this sim hook handle is bound to an operator running on a cluster, where every member has its own independent instance to script; \
309 select one with `.on(member_id)` (members: {:?})",
310 bound_members
311 ),
312 Some(m) if self.hooks.contains_key(&(hook_id, None)) => panic!(
313 "`.on({m})` was used on a sim hook handle bound to an operator running on a process, which has no cluster members; \
314 script the handle without `.on(..)`"
315 ),
316 Some(m) if !bound_members.is_empty() => panic!(
317 "`.on({m})` does not name a member of the cluster this sim hook handle is bound to (members: {:?})",
318 bound_members
319 ),
320 _ => panic!("{}", UNBOUND_HOOK_ERROR),
321 }
322 }
323
324 /// Whether the simulation is currently quiescent (no more progress possible).
325 pub(crate) fn is_quiescent(&self) -> bool {
326 self.quiescence.is_quiescent()
327 }
328
329 /// See [`QuiescenceState::push_park_waker`].
330 pub(crate) fn push_park_waker(&self, waker: &std::task::Waker) {
331 self.quiescence.push_park_waker(waker);
332 }
333
334 /// Attempts to install a decision (bincode-serialized; the handle and hook statically
335 /// know the matching type) for the hook instance `(hook_id, member)` under the group
336 /// protocol: join the current group if this decision belongs to it, open a new group
337 /// if the previous one has been consumed, or hand the decision back to be retried
338 /// once the previous group's tick execution has happened. Each cluster member's tick
339 /// is its own scheduler action, so decisions for different members never share a
340 /// group.
341 #[track_caller]
342 pub(crate) fn try_schedule_decision(
343 &self,
344 hook_id: usize,
345 member: Option<u32>,
346 decision_blob: Vec<u8>,
347 ) -> Result<ScheduleDecision, String> {
348 let key = (hook_id, member);
349 let hook = self.control(hook_id, member);
350 let target = hook.borrow().target();
351
352 let mut coordinator = self.coordinator.borrow_mut();
353
354 enum Action {
355 Join,
356 NewGroup,
357 Wait,
358 }
359
360 let action = match &coordinator.current {
361 None => Action::NewGroup,
362 Some(group)
363 if !group.sealed
364 && matches!(target, ScriptTarget::Tick { .. })
365 && group.target == target
366 && !group.members.contains(&key) =>
367 {
368 Action::Join
369 }
370 Some(_) => Action::Wait,
371 };
372
373 match action {
374 Action::Join => {
375 coordinator.current.as_mut().unwrap().members.push(key);
376 }
377 Action::NewGroup => {
378 coordinator.current = Some(CurrentGroup {
379 target,
380 members: vec![key],
381 sealed: false,
382 });
383 }
384 Action::Wait => {
385 // The previous group's execution hasn't happened yet; hand the decision
386 // back to be retried. The waiting hook stays subject to the boundary scan:
387 // buffered input held across this wait must be declared with an explicit
388 // pause (the waiting decision names a *later* execution).
389 if self.quiescence.is_quiescent() {
390 let stuck = coordinator.describe_unconsumed(&self.hooks);
391 let stuck = stuck.as_deref().unwrap_or(" (unknown decision)");
392 let header = if coordinator.stuck_cannot_trigger {
393 "a previously scripted decision group can never run: none of its tick's hooks can trigger it (no scripted decision triggers, and no unscripted input has data)"
394 } else {
395 "a previously scripted decision can never be satisfied (the simulation has no more work it can do)"
396 };
397 return Err(format!("cannot script this decision: {header}:\n{stuck}"));
398 }
399 return Ok(ScheduleDecision::Wait(decision_blob));
400 }
401 }
402 drop(coordinator);
403
404 hook.borrow_mut().install_decision(&decision_blob);
405 // Installing a decision can make a tick runnable; wake the scheduler if parked.
406 self.quiescence.resume();
407 Ok(ScheduleDecision::Installed)
408 }
409}
410
411/// Resolves the per-instance scripting context. Panics if called outside a simulation.
412pub(crate) fn script_ctx() -> ScriptCtx {
413 CURRENT_SIM_CONNECTIONS.with(|connections| {
414 let connections = connections.borrow();
415 ScriptCtx {
416 hooks: connections.scripted_hooks.clone(),
417 coordinator: connections.script_coordinator.clone(),
418 quiescence: connections.quiescence.clone(),
419 }
420 })
421}
422
423/// Renders the stuck-script error for a quiescent simulation with an outstanding group.
424/// Two distinct failure styles: a decision that is *unsatisfiable* (waiting on input that
425/// can never arrive), vs decisions that are all satisfiable but *cannot trigger* their
426/// tick (none of them triggers, and no unscripted input on the tick has data).
427fn render_stuck_script_error(cannot_trigger: bool, stuck: &str) -> String {
428 if cannot_trigger {
429 format!(
430 "the simulation has stopped, but scripted decisions are still pending: none of the tick's hooks can trigger it (no scripted decision triggers, and no unscripted input has data):\n{stuck}\nhelp: script a decision that triggers the tick, or drive an unscripted input, so the tick can run"
431 )
432 } else {
433 format!("a scripted decision can never be satisfied:\n{stuck}")
434 }
435}
436
437/// Renders the stuck-script error for the current instance (see
438/// [`render_stuck_script_error`]); the scheduler classified the failure style when it
439/// reached quiescence.
440pub(crate) fn script_stuck_error(stuck: &str) -> String {
441 let cannot_trigger = CURRENT_SIM_CONNECTIONS.with(|connections| {
442 let connections = connections.borrow();
443 let coordinator = connections.script_coordinator.borrow();
444 coordinator.stuck_cannot_trigger
445 });
446 render_stuck_script_error(cannot_trigger, stuck)
447}
448
449/// If a scripted group is outstanding, returns a description of its decisions (used by
450/// output awaits and `pause_until` waits, which are script barriers: they must not
451/// resolve until every decision scripted so far has run).
452pub(crate) fn script_unconsumed_description() -> Option<String> {
453 CURRENT_SIM_CONNECTIONS.with(|connections| {
454 let connections = connections.borrow();
455 let coordinator = connections.script_coordinator.borrow();
456 coordinator.current.as_ref()?;
457 Some(
458 coordinator
459 .describe_unconsumed(&connections.scripted_hooks)
460 .unwrap_or_else(|| " (unknown decision)".to_owned()),
461 )
462 })
463}
464
465/// Tracks a pending "settle" pause request to the scheduler (see
466/// [`QuiescenceState::pause_nondet`]), releasing it if the requesting future is dropped
467/// mid-settle (e.g. by `select!`) so the scheduler is not left paused forever. Pause
468/// requests are counted, so concurrent settling futures each hold their own request.
469struct SettlePauseGuard {
470 quiescence: Rc<QuiescenceState>,
471 active: bool,
472}
473
474impl SettlePauseGuard {
475 fn new(quiescence: Rc<QuiescenceState>) -> Self {
476 SettlePauseGuard {
477 quiescence,
478 active: false,
479 }
480 }
481
482 fn acquire(&mut self) {
483 abort_assert!(!self.active, "settle pause acquired twice");
484 self.quiescence
485 .pause_nondet
486 .set(self.quiescence.pause_nondet.get() + 1);
487 self.active = true;
488 }
489
490 fn release(&mut self) {
491 abort_assert!(self.active, "settle pause released without being acquired");
492 self.active = false;
493 self.quiescence
494 .pause_nondet
495 .set(self.quiescence.pause_nondet.get() - 1);
496 }
497
498 /// Polls the "settle" handshake with the scheduler: deterministic (non-tick) work is
499 /// allowed to run, but the scheduler pauses instead of starting nondeterministic work
500 /// (ticks / observations). Resolves to `true` if the simulation reached quiescence
501 /// deterministically, or `false` if nondeterministic work is pending (in which case the
502 /// scheduler is resumed).
503 fn poll_settle(&mut self, cx: &mut std::task::Context<'_>) -> Poll<bool> {
504 let quiescence = self.quiescence.clone();
505 if !self.active {
506 if quiescence.is_quiescent() {
507 return Poll::Ready(true);
508 }
509 self.acquire();
510 }
511
512 if quiescence.is_quiescent() {
513 self.release();
514 Poll::Ready(true)
515 } else if quiescence.nondet_pending.get() {
516 self.release();
517 // `notify_one` (permit-based): the driver only parks *between* thunk polls, so it
518 // is not parked right now — the permit ensures this resume is not lost.
519 quiescence.resume_notify.notify_one();
520 Poll::Ready(false)
521 } else {
522 // This may push a duplicate waker if we are re-polled without an intervening
523 // `wake_settled` (e.g. a `join!` sibling waking the shared task), but duplicates
524 // are harmless (waking is idempotent) and are cleared at the next `wake_settled`,
525 // so deduplicating here isn't worth the scan on every poll.
526 quiescence
527 .settle_wakers
528 .borrow_mut()
529 .push(cx.waker().clone());
530 Poll::Pending
531 }
532 }
533}
534
535impl Drop for SettlePauseGuard {
536 fn drop(&mut self) {
537 if self.active {
538 self.release();
539 // Resume the scheduler in case this was the last pause request (otherwise it
540 // would stay parked forever with nobody left to resume it). `notify_one`
541 // (permit-based) so the resume is not lost if the driver has not parked yet. If
542 // other settlers still hold requests, this wakeup is spurious but harmless: the
543 // scheduler re-checks `pause_nondet > 0` before starting any nondeterministic
544 // work, so it immediately re-parks without running anything.
545 self.quiescence.resume_notify.notify_one();
546 }
547 }
548}
549
550/// Panics if the simulation has been poisoned: an earlier observation forced the simulation
551/// to quiesce (running pending nondeterministic work), and new input has been sent since, so
552/// further observations could misattribute failures caused by the forced overrun.
553fn guard_not_poisoned(quiescence: &QuiescenceState) {
554 assert!(
555 !quiescence.poisoned.get(),
556 "cannot receive more simulator output: an earlier observation (such as `try_next`, `collect`, or a quiescence assertion outside exhaustive mode) forced the simulation to quiesce by running pending nondeterministic work, and new input has been sent since. Failures observed now could be misattributed, so either restructure the test to make quiescence-forcing observations its last step, or insert an explicit `sim::quiesce().await` phase barrier before sending more input."
557 );
558}
559
560/// Runs the simulation to quiescence, as an explicit *phase barrier* between rounds of a
561/// multi-phase test.
562///
563/// All pending nondeterministic work (ticks / observations) is forced to run until no more
564/// progress is possible without new input. This deliberately narrows the explored executions:
565/// inputs sent after the barrier will never interleave with work from before it, modeling
566/// scenarios where new stimuli (such as timer ticks) arrive long after the system settles.
567/// Pair such tests with a separate barrier-free test if interleaved executions should also be
568/// explored.
569///
570/// Because the barrier is explicit, observations after it are *intended* to see the fully
571/// settled state, so — unlike [`SimReceiver::try_next`] / [`SimReceiver::collect`] forcing
572/// quiescence implicitly — it does not restrict what the test may do afterwards: receives
573/// after the barrier observe only buffered output (plus whatever later input produces), and
574/// failures cannot be misattributed across it.
575pub async fn quiesce() {
576 let quiescence =
577 CURRENT_SIM_CONNECTIONS.with(|connections| connections.borrow().quiescence.clone());
578 guard_not_poisoned(&quiescence);
579
580 let mut notified_fut = pin!(None);
581 std::future::poll_fn(|cx| {
582 if quiescence.is_quiescent() {
583 // A stuck scripted decision makes this a *dirty* quiescence: report it here
584 // rather than letting the barrier silently pass.
585 if let Some(stuck) = script_unconsumed_description() {
586 panic!("{}", script_stuck_error(&stuck));
587 }
588 return Poll::Ready(());
589 }
590 // Registered before the scheduler can run (single-threaded), so the quiescence
591 // notification cannot be missed.
592 if notified_fut.is_none() {
593 notified_fut.set(Some(quiescence.notified()));
594 }
595 let () = ready!(notified_fut.as_mut().as_pin_mut().unwrap().poll(cx));
596 Poll::Ready(())
597 })
598 .await;
599
600 // The barrier subsumes any quiescence forced by earlier observations in this phase:
601 // everything before it has fully settled, and the test has explicitly opted into
602 // observing only post-quiescence states from here on.
603 quiescence.tainted.set(false);
604}
605
606/// Receives the next message from `receiver` while trying not to overrun the simulation:
607/// first the simulation *settles* (deterministic work runs, but the scheduler pauses before
608/// nondeterministic work). If a message arrives, it is returned; if the simulation settles to
609/// quiescence, returns `None` without having run any nondeterministic work. Otherwise the
610/// scheduler is resumed and pending nondeterministic work runs until a message arrives or the
611/// simulation quiesces; quiescing this way *taints* the simulation (see
612/// [`QuiescenceState::tainted`]).
613async fn try_next_bytes(
614 receiver: &Mutex<UnsyncReceiver<Bytes>>,
615 quiescence: &Rc<QuiescenceState>,
616) -> Option<Bytes> {
617 guard_not_poisoned(quiescence);
618
619 let mut receiver_stream = receiver.lock().await;
620 let mut settle_guard = SettlePauseGuard::new(quiescence.clone());
621 // `Some` once the settle phase has concluded that nondeterministic work is pending and
622 // we have started forcing it to run.
623 let mut notified_fut = pin!(None);
624
625 std::future::poll_fn(|cx| {
626 // **Scripted-decision barrier**: an output await completes only after every
627 // decision scripted so far has been consumed, so every point where the test body
628 // resumes is a clean synchronization point (the script written so far has fully
629 // happened). If the simulation runs out of work while a scripted decision is still
630 // waiting, that decision can never be honored — panic instead of yielding output
631 // or end-of-stream, so a stuck script cannot masquerade as a completed one.
632 if let Some(stuck) = script_unconsumed_description() {
633 assert!(!quiescence.is_quiescent(), "{}", script_stuck_error(&stuck));
634 quiescence.push_park_waker(cx.waker());
635 return Poll::Pending;
636 }
637
638 // A message may become available at any point (including from deterministic work
639 // while settling), so always check the stream first.
640 match receiver_stream.poll_next_unpin(cx) {
641 Poll::Ready(Some(bytes)) => return Poll::Ready(Some(bytes)),
642 Poll::Ready(None) => return Poll::Ready(None),
643 Poll::Pending => {}
644 }
645
646 if notified_fut.is_none() {
647 match settle_guard.poll_settle(cx) {
648 // Deterministically quiescent: no more messages, and nothing was overrun.
649 Poll::Ready(true) => return Poll::Ready(None),
650 // Nondeterministic work is pending; start forcing it to run. The `Notified`
651 // is created here and polled (registered) below in this same synchronous
652 // poll — before the scheduler can run — and the simulation is not currently
653 // quiescent, so the quiescence notification cannot be missed.
654 Poll::Ready(false) => notified_fut.set(Some(quiescence.notified())),
655 Poll::Pending => return Poll::Pending,
656 }
657 }
658
659 // Let the scheduler run nondeterministic work until a message arrives or the
660 // simulation quiesces. Note that merely entering this phase does not taint: if a
661 // message arrives (the `Some` exit at the top), waiting was sound for the same
662 // reason as `SimReceiver::next` — the work that ran was needed to produce it. Only
663 // *observing quiescence* after forcing the pending work taints, since that is the
664 // overrun a later observation could misattribute.
665 let () = ready!(notified_fut.as_mut().as_pin_mut().unwrap().poll(cx));
666 quiescence.tainted.set(true);
667 Poll::Ready(None)
668 })
669 .await
670}
671
672struct SimConnections {
673 input_senders: HashMap<SimExternalPort, UnsyncSender<Bytes>>,
674 output_receivers: HashMap<SimExternalPort, Rc<Mutex<UnsyncReceiver<Bytes>>>>,
675 cluster_input_senders: HashMap<SimExternalPort, HashMap<u32, UnsyncSender<Bytes>>>,
676 cluster_output_receivers:
677 HashMap<SimExternalPort, HashMap<u32, Rc<Mutex<UnsyncReceiver<Bytes>>>>>,
678 external_registered: HashMap<ExternalPortId, SimExternalPort>,
679 quiescence: Rc<QuiescenceState>,
680 /// Every scripted hook (shared with the scheduler's tick lists), keyed by handle ID.
681 scripted_hooks: Rc<ScriptedHookRegistry>,
682 /// Coordinates the decision-group protocol between hook handles and the scheduler.
683 script_coordinator: Rc<RefCell<ScriptCoordinator>>,
684 log: bool,
685 /// Whether this instance is being executed by the exhaustive engine (see
686 /// [`CompiledSim::exhaustive`]), which affects how `assert_yields_only` explores
687 /// quiescence checks.
688 exhaustive: bool,
689}
690
691/// Implementation detail of [`crate::sim::continue_if!`](crate::continue_if); do not call directly.
692///
693/// If `condition` is false, aborts the current simulation instance by panicking with a special
694/// payload ([`bolero::generator::bolero_generator::any::Error`]) that bolero recognizes as an
695/// "invalid input" marker: the instance is discarded (not treated as a test failure, and never
696/// recorded as a reproducer) and exploration moves on to the next instance. If logging is
697/// enabled for the current instance, the failed assumption is logged first.
698#[doc(hidden)]
699#[track_caller]
700pub fn continue_if_impl(condition: bool, message: fmt::Arguments<'_>) {
701 if condition {
702 return;
703 }
704
705 let log = CURRENT_SIM_CONNECTIONS
706 .try_with(|connections| connections.borrow().log)
707 .unwrap_or(true);
708 if log {
709 eprintln!(
710 "{}",
711 render_continue_if_failure(std::panic::Location::caller(), message)
712 );
713 }
714
715 // Panics with `bolero_generator::any::Error`, which bolero's engines treat as an invalid
716 // input rather than a test failure. Both this function and bolero's `assume` are
717 // `#[track_caller]`, so the recorded location is the user's `continue_if!` call site.
718 bolero::generator::bolero_generator::any::assume(false, "simulation assumption failed");
719}
720
721/// Renders the log message for a failed assumption, echoing the source line with a caret
722/// pointing at the `continue_if!` call site, in the same style as the other simulator logs.
723fn render_continue_if_failure(
724 location: &std::panic::Location<'_>,
725 message: fmt::Arguments<'_>,
726) -> String {
727 use std::fmt::Write;
728
729 // `Location::file()` is relative to the directory the crate was compiled from (e.g. the
730 // workspace root), which may not match the current working directory (e.g. the crate
731 // root when running `cargo test`), so walk up from the current directory to find it.
732 let source_line = std::env::current_dir()
733 .ok()
734 .and_then(|cwd| {
735 cwd.ancestors()
736 .find_map(|base| std::fs::read_to_string(base.join(location.file())).ok())
737 })
738 .and_then(|content| {
739 content
740 .lines()
741 .nth((location.line() as usize).saturating_sub(1))
742 .map(|line| line.to_owned())
743 })
744 .unwrap_or_default();
745
746 let caret_indent = " ".repeat((location.column() as usize).saturating_sub(1));
747
748 let mut out = String::new();
749 let _ = writeln!(
750 out,
751 "\n{}",
752 "Condition failed (discarding simulation instance):"
753 .color(colored::Color::Yellow)
754 .bold()
755 );
756 let _ = writeln!(out, "{} {}", "-->".color(colored::Color::Blue), location);
757 let _ = writeln!(out, " {}{}", "|".color(colored::Color::Blue), source_line);
758 let _ = write!(
759 out,
760 " {}{}{}",
761 "|".color(colored::Color::Blue),
762 caret_indent,
763 format!("^ {}", message).color(colored::Color::Yellow)
764 );
765 out
766}
767
768tokio::task_local! {
769 static CURRENT_SIM_CONNECTIONS: RefCell<SimConnections>;
770}
771
772/// A handle to a compiled Hydro simulation, which can be instantiated and run.
773pub struct CompiledSim {
774 pub(super) _path: BuiltArtifact,
775 pub(super) lib: Library,
776 pub(super) externals_port_registry: SimExternalPortRegistry,
777 pub(super) unit_test_fuzz_iterations: usize,
778}
779
780#[sealed::sealed]
781/// A trait implemented by closures that can instantiate a compiled simulation.
782///
783/// This is needed to ensure [`RefUnwindSafe`] so instances can be created during fuzzing.
784pub trait Instantiator<'a>: RefUnwindSafe + Fn() -> CompiledSimInstance<'a> {}
785#[sealed::sealed]
786impl<'a, T: RefUnwindSafe + Fn() -> CompiledSimInstance<'a>> Instantiator<'a> for T {}
787
788fn null_handler(_args: fmt::Arguments<'_>) {}
789
790fn println_handler(args: fmt::Arguments<'_>) {
791 println!("{}", args);
792}
793
794fn eprintln_handler(args: fmt::Arguments<'_>) {
795 eprintln!("{}", args);
796}
797
798/// Creates a simulation instance, returning:
799/// - A list of async DFIRs to run (all process / cluster logic outside a tick)
800/// - A list of tick DFIRs to run (where the &'static str is for the tick location id)
801/// - A mapping of hooks for non-deterministic decisions at tick-input boundaries
802/// - A mapping of inline hooks for non-deterministic decisions inside ticks
803type SimLoaded<'a> = libloading::Symbol<
804 'a,
805 unsafe extern "Rust" fn(
806 should_color: bool,
807 external_out: &mut HashMap<usize, UnsyncReceiver<Bytes>>,
808 external_in: &mut HashMap<usize, UnsyncSender<Bytes>>,
809 cluster_external_out: &mut HashMap<usize, HashMap<u32, UnsyncReceiver<Bytes>>>,
810 cluster_external_in: &mut HashMap<usize, HashMap<u32, UnsyncSender<Bytes>>>,
811 println_handler: fn(fmt::Arguments<'_>),
812 eprintln_handler: fn(fmt::Arguments<'_>),
813 ) -> (
814 Vec<(LocationId, Option<u32>, DfirErased)>,
815 Vec<(LocationId, Option<u32>, DfirErased)>,
816 Hooks,
817 ObservationHooks,
818 InlineHooks,
819 ScriptedTickHooks,
820 ScriptedObservationHooks,
821 ScriptedInlineHooks,
822 ScriptedHookRegistry,
823 ),
824>;
825
826impl CompiledSim {
827 /// Executes the given closure with a single instance of the compiled simulation.
828 pub fn with_instance<T>(&self, thunk: impl FnOnce(CompiledSimInstance<'_>) -> T) -> T {
829 self.with_instantiator(|instantiator| thunk(instantiator()), true)
830 }
831
832 /// Executes the given closure with an [`Instantiator`], which can be called to create
833 /// independent instances of the simulation. This is useful for fuzzing, where we need to
834 /// re-execute the simulation several times with different decisions.
835 ///
836 /// The `always_log` parameter controls whether to log tick executions and stream releases. If
837 /// it is `true`, logging will always be enabled. If it is `false`, logging will only be
838 /// enabled if the `HYDRO_SIM_LOG` environment variable is set to `1`.
839 pub fn with_instantiator<T>(
840 &self,
841 thunk: impl FnOnce(&dyn Instantiator<'_>) -> T,
842 always_log: bool,
843 ) -> T {
844 let func: SimLoaded<'_> = unsafe { self.lib.get(b"__hydro_runtime").unwrap() };
845 let log = always_log || std::env::var("HYDRO_SIM_LOG").is_ok_and(|v| v == "1");
846 thunk(
847 &(|| CompiledSimInstance {
848 func: func.clone(),
849 externals_port_registry: self.externals_port_registry.clone(),
850 dylib_result: None,
851 log,
852 exhaustive: false,
853 deterministic: false,
854 }),
855 )
856 }
857
858 /// Uses a fuzzing strategy to explore possible executions of the simulation. The provided
859 /// closure will be repeatedly executed with instances of the Hydro program where the
860 /// batching boundaries, order of messages, and retries are varied.
861 ///
862 /// During development, you should run the test that invokes this function with the `cargo sim`
863 /// command, which will use `libfuzzer` to intelligently explore the execution space. If a
864 /// failure is found, a minimized test case will be produced in a `sim-failures` directory.
865 /// When running the test with `cargo test` (such as in CI), if a reproducer is found it will
866 /// be executed, and if no reproducer is found a small number of random executions will be
867 /// performed.
868 pub fn fuzz(&self, mut thunk: impl AsyncFnMut() + RefUnwindSafe) {
869 let caller_fn = crate::compile::ir::backtrace::Backtrace::get_backtrace(0)
870 .elements()
871 .into_iter()
872 .find(|e| {
873 !e.fn_name.starts_with("hydro_lang::sim::compiled")
874 && !e.fn_name.starts_with("hydro_lang::sim::flow")
875 && !e.fn_name.starts_with("fuzz<")
876 && !e.fn_name.starts_with("<hydro_lang::sim")
877 })
878 .unwrap();
879
880 let caller_path = Path::new(&caller_fn.filename.unwrap()).to_path_buf();
881 let repro_folder = caller_path.parent().unwrap().join("sim-failures");
882
883 let caller_fuzz_repro_path = repro_folder
884 .join(caller_fn.fn_name.replace("::", "__"))
885 .with_extension("bin");
886
887 if std::env::var("BOLERO_FUZZER").is_ok() {
888 let corpus_dir = std::env::current_dir().unwrap().join(".fuzz-corpus");
889 std::fs::create_dir_all(&corpus_dir).unwrap();
890 let libfuzzer_args = format!(
891 "{} {} -artifact_prefix={}/ -handle_abrt=0",
892 corpus_dir.to_str().unwrap(),
893 corpus_dir.to_str().unwrap(),
894 corpus_dir.to_str().unwrap(),
895 );
896
897 std::fs::create_dir_all(&repro_folder).unwrap();
898
899 if !std::env::var("HYDRO_NO_FAILURE_OUTPUT").is_ok_and(|v| v == "1") {
900 unsafe {
901 std::env::set_var(
902 "BOLERO_FAILURE_OUTPUT",
903 caller_fuzz_repro_path.to_str().unwrap(),
904 );
905 }
906 }
907
908 unsafe {
909 std::env::set_var("BOLERO_LIBFUZZER_ARGS", libfuzzer_args);
910 }
911
912 self.with_instantiator(
913 |instantiator| {
914 bolero::test(bolero::TargetLocation {
915 package_name: "",
916 manifest_dir: "",
917 module_path: "",
918 file: "",
919 line: 0,
920 item_path: "<unknown>::__bolero_item_path__",
921 test_name: None,
922 })
923 .run_with_replay(move |is_replay| {
924 let mut instance = instantiator();
925
926 if instance.log {
927 eprintln!(
928 "{}",
929 "\n==== New Simulation Instance ===="
930 .color(colored::Color::Cyan)
931 .bold()
932 );
933 }
934
935 if is_replay {
936 instance.log = true;
937 }
938
939 tokio::runtime::Builder::new_current_thread()
940 .build()
941 .unwrap()
942 .block_on(async { instance.run(&mut thunk).await })
943 })
944 },
945 false,
946 );
947 } else if let Ok(existing_bytes) = std::fs::read(&caller_fuzz_repro_path) {
948 self.fuzz_repro(existing_bytes, async |compiled| {
949 compiled.run_with_scheduler(thunk()).await
950 });
951 } else {
952 eprintln!(
953 "Running a fuzz test without `cargo sim` and no reproducer found at {}, using {} iterations with random inputs.",
954 caller_fuzz_repro_path.display(),
955 self.unit_test_fuzz_iterations,
956 );
957 self.with_instantiator(
958 |instantiator| {
959 bolero::test(bolero::TargetLocation {
960 package_name: "",
961 manifest_dir: "",
962 module_path: "",
963 file: ".",
964 line: 0,
965 item_path: "<unknown>::__bolero_item_path__",
966 test_name: None,
967 })
968 .with_iterations(self.unit_test_fuzz_iterations)
969 .run_with_replay(move |is_replay| {
970 let mut instance = instantiator();
971
972 if instance.log {
973 eprintln!(
974 "{}",
975 "\n==== New Simulation Instance ===="
976 .color(colored::Color::Cyan)
977 .bold()
978 );
979 }
980
981 if is_replay {
982 instance.log = true;
983 }
984
985 tokio::runtime::Builder::new_current_thread()
986 .build()
987 .unwrap()
988 .block_on(async { instance.run(&mut thunk).await })
989 })
990 },
991 false,
992 );
993 }
994 }
995
996 /// Executes the given closure with a single instance of the compiled simulation, using the
997 /// provided bytes as the source of fuzzing decisions. This can be used to manually reproduce a
998 /// failure found during fuzzing.
999 pub fn fuzz_repro<'a>(
1000 &'a self,
1001 bytes: Vec<u8>,
1002 thunk: impl AsyncFnOnce(CompiledSimInstance<'_>) + RefUnwindSafe,
1003 ) {
1004 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1005 self.with_instance(|instance| {
1006 bolero::bolero_engine::any::scope::with(
1007 Box::new(bolero::bolero_engine::driver::object::Object(
1008 bolero::bolero_engine::driver::bytes::Driver::new(
1009 bytes,
1010 &Default::default(),
1011 ),
1012 )),
1013 || {
1014 tokio::runtime::Builder::new_current_thread()
1015 .build()
1016 .unwrap()
1017 .block_on(async { instance.run_without_launching(thunk).await })
1018 },
1019 )
1020 })
1021 }));
1022
1023 if let Err(payload) = result {
1024 if payload
1025 .downcast_ref::<bolero::generator::bolero_generator::any::Error>()
1026 .is_some()
1027 {
1028 // A `continue_if!` failed (or the driver ran out of entropy) while replaying the
1029 // recorded bytes. Instances that fail an assumption are never recorded as
1030 // failures, so this means the reproducer is stale or does not correspond to
1031 // this program.
1032 panic!(
1033 "simulation assumption failed while replaying recorded fuzz decisions; the reproducer may be stale or may not correspond to this program"
1034 );
1035 }
1036 std::panic::resume_unwind(payload);
1037 }
1038 }
1039
1040 /// Exhaustively searches all possible executions of the simulation. The provided
1041 /// closure will be repeatedly executed with instances of the Hydro program where the
1042 /// batching boundaries, order of messages, and retries are varied.
1043 ///
1044 /// Exhaustive searching is feasible when the inputs to the Hydro program are finite and there
1045 /// are no dataflow loops that generate infinite messages. Exhaustive searching provides a
1046 /// stronger guarantee of correctness than fuzzing, but may take a long time to complete.
1047 /// Because no fuzzer is involved, you can run exhaustive tests with `cargo test`.
1048 ///
1049 /// Returns the number of distinct executions explored.
1050 pub fn exhaustive(&self, mut thunk: impl AsyncFnMut() + RefUnwindSafe) -> usize {
1051 if std::env::var("BOLERO_FUZZER").is_ok() {
1052 eprintln!(
1053 "Cannot run exhaustive tests with a fuzzer. Please use `cargo test` instead of `cargo sim`."
1054 );
1055 std::process::abort();
1056 }
1057
1058 let mut count = 0;
1059 let count_mut = &mut count;
1060
1061 let _span = tracing::debug_span!(target: "hydro_build", "sim_exhaustive").entered();
1062
1063 self.with_instantiator(
1064 |instantiator| {
1065 bolero::test(bolero::TargetLocation {
1066 package_name: "",
1067 manifest_dir: "",
1068 module_path: "",
1069 file: "",
1070 line: 0,
1071 item_path: "<unknown>::__bolero_item_path__",
1072 test_name: None,
1073 })
1074 .exhaustive()
1075 .run_with_replay(move |is_replay| {
1076 *count_mut += 1;
1077
1078 let mut instance = instantiator();
1079 instance.exhaustive = true;
1080 if instance.log {
1081 eprintln!(
1082 "{}",
1083 "\n==== New Simulation Instance ===="
1084 .color(colored::Color::Cyan)
1085 .bold()
1086 );
1087 }
1088
1089 if is_replay {
1090 instance.log = true;
1091 }
1092
1093 tokio::runtime::Builder::new_current_thread()
1094 .build()
1095 .unwrap()
1096 .block_on(async { instance.run(&mut thunk).await })
1097 })
1098 },
1099 false,
1100 );
1101
1102 count
1103 }
1104
1105 /// Runs the test body against exactly **one** execution of the program, with no fuzzer
1106 /// involved anywhere: if it passes once, it passes always, on every machine.
1107 ///
1108 /// Every source of variation must be pinned: inputs are already scripted (via
1109 /// `sim_input`), and every unsafe operator that receives data must be bound to a sim
1110 /// hook (see [`crate::sim_hooks`]) and scripted — encountering an unhooked operator
1111 /// with meaningful input panics, naming the operator. The scheduler needs no
1112 /// tie-breaking policy because at most one tick is ever runnable: scripted decisions
1113 /// activate one group at a time, so the *script* is the schedule.
1114 pub fn deterministic(&self, thunk: impl AsyncFnOnce() + RefUnwindSafe) {
1115 self.with_instance(|mut instance| {
1116 instance.deterministic = true;
1117
1118 // Deliberately do not install a Bolero entropy scope. Deterministic execution
1119 // must never draw entropy; Bolero's unset thread-local scope makes any accidental
1120 // draw fail immediately with `no scope set`.
1121 tokio::runtime::Builder::new_current_thread()
1122 .build()
1123 .unwrap()
1124 .block_on(instance.run(thunk));
1125 })
1126 }
1127}
1128
1129// This must be a tuple because it is referenced from generated code in `graph.rs`.
1130type DylibResult = (
1131 Vec<(LocationId, Option<u32>, DfirErased)>,
1132 Vec<(LocationId, Option<u32>, DfirErased)>,
1133 Hooks,
1134 ObservationHooks,
1135 InlineHooks,
1136 ScriptedTickHooks,
1137 ScriptedObservationHooks,
1138 ScriptedInlineHooks,
1139 ScriptedHookRegistry,
1140);
1141
1142/// A single instance of a compiled Hydro simulation, which provides methods to interactively
1143/// execute the simulation, feed inputs, and receive outputs.
1144pub struct CompiledSimInstance<'a> {
1145 func: SimLoaded<'a>,
1146 externals_port_registry: SimExternalPortRegistry,
1147 dylib_result: Option<DylibResult>,
1148 log: bool,
1149 exhaustive: bool,
1150 deterministic: bool,
1151}
1152
1153impl<'a> CompiledSimInstance<'a> {
1154 async fn run(self, thunk: impl AsyncFnOnce() + RefUnwindSafe) {
1155 self.run_without_launching(async |instance| {
1156 instance.run_with_scheduler(thunk()).await;
1157 })
1158 .await;
1159 }
1160
1161 async fn run_without_launching(
1162 mut self,
1163 thunk: impl AsyncFnOnce(CompiledSimInstance<'_>) + RefUnwindSafe,
1164 ) {
1165 let mut external_out: HashMap<usize, UnsyncReceiver<Bytes>> = HashMap::new();
1166 let mut external_in: HashMap<usize, UnsyncSender<Bytes>> = HashMap::new();
1167 let mut cluster_external_out: HashMap<usize, HashMap<u32, UnsyncReceiver<Bytes>>> =
1168 HashMap::new();
1169 let mut cluster_external_in: HashMap<usize, HashMap<u32, UnsyncSender<Bytes>>> =
1170 HashMap::new();
1171
1172 let mut dylib_result = unsafe {
1173 (self.func)(
1174 colored::control::SHOULD_COLORIZE.should_colorize(),
1175 &mut external_out,
1176 &mut external_in,
1177 &mut cluster_external_out,
1178 &mut cluster_external_in,
1179 if self.log {
1180 println_handler
1181 } else {
1182 null_handler
1183 },
1184 if self.log {
1185 eprintln_handler
1186 } else {
1187 null_handler
1188 },
1189 )
1190 };
1191
1192 let registered = &self.externals_port_registry.registered;
1193
1194 let quiescence = Rc::new(QuiescenceState {
1195 quiescent: Cell::new(false),
1196 quiescence_notify: Notify::new(),
1197 resume_notify: Notify::new(),
1198 pause_nondet: Cell::new(0),
1199 nondet_pending: Cell::new(false),
1200 settle_wakers: RefCell::new(vec![]),
1201 tainted: Cell::new(false),
1202 poisoned: Cell::new(false),
1203 });
1204
1205 let mut input_senders = HashMap::new();
1206 let mut output_receivers = HashMap::new();
1207 let mut cluster_input_senders = HashMap::new();
1208 let mut cluster_output_receivers = HashMap::new();
1209
1210 #[expect(
1211 clippy::disallowed_methods,
1212 reason = "inserts into maps also unordered"
1213 )]
1214 for sim_port in registered.values() {
1215 let usize_key = sim_port.into_inner();
1216 if let Some(sender) = external_in.remove(&usize_key) {
1217 input_senders.insert(*sim_port, sender);
1218 }
1219 if let Some(receiver) = external_out.remove(&usize_key) {
1220 output_receivers.insert(*sim_port, Rc::new(Mutex::new(receiver)));
1221 }
1222 if let Some(senders) = cluster_external_in.remove(&usize_key) {
1223 cluster_input_senders.insert(*sim_port, senders);
1224 }
1225 if let Some(receivers) = cluster_external_out.remove(&usize_key) {
1226 cluster_output_receivers.insert(
1227 *sim_port,
1228 receivers
1229 .into_iter()
1230 .map(|(member, r)| (member, Rc::new(Mutex::new(r))))
1231 .collect(),
1232 );
1233 }
1234 }
1235
1236 let scripted_hooks = Rc::new(std::mem::take(&mut dylib_result.8));
1237 self.dylib_result = Some(dylib_result);
1238
1239 CURRENT_SIM_CONNECTIONS
1240 .scope(
1241 RefCell::new(SimConnections {
1242 input_senders,
1243 output_receivers,
1244 cluster_input_senders,
1245 cluster_output_receivers,
1246 external_registered: self.externals_port_registry.registered.clone(),
1247 quiescence: quiescence.clone(),
1248 scripted_hooks,
1249 script_coordinator: Rc::new(RefCell::new(ScriptCoordinator::default())),
1250 log: self.log,
1251 exhaustive: self.exhaustive,
1252 }),
1253 async move {
1254 thunk(self).await;
1255 },
1256 )
1257 .await;
1258 }
1259
1260 /// Runs the simulation scheduler alongside the given future, until the future completes.
1261 ///
1262 /// The future always gets to run first; whenever it is blocked (e.g. waiting to receive
1263 /// simulation outputs), the scheduler runs a single step to completion. Steps are atomic
1264 /// with respect to the future: it is re-polled between every pair of scheduler steps, but
1265 /// never while a step is in flight. The [`LaunchedSim`] state struct lives across steps,
1266 /// in this function's frame.
1267 async fn run_with_scheduler(self, thunk: impl Future<Output = ()>) {
1268 self.run_with_scheduler_and_maybe_logger::<std::io::Empty>(None, thunk)
1269 .await;
1270 }
1271
1272 /// Runs the simulation scheduler alongside the given future, until the future completes,
1273 /// reporting the simulation trace to the given logger.
1274 ///
1275 /// The future always gets to run first; whenever it is blocked (e.g. waiting to receive
1276 /// simulation outputs), the scheduler runs a single step to completion. Steps are atomic
1277 /// with respect to the future: it is re-polled between every pair of scheduler steps, but
1278 /// never while a step is in flight.
1279 pub async fn run_with_scheduler_and_logger<W: std::io::Write>(
1280 self,
1281 log_writer: W,
1282 thunk: impl Future<Output = ()>,
1283 ) {
1284 self.run_with_scheduler_and_maybe_logger(Some(log_writer), thunk)
1285 .await;
1286 }
1287
1288 async fn run_with_scheduler_and_maybe_logger<W: std::io::Write>(
1289 self,
1290 log_override: Option<W>,
1291 thunk: impl Future<Output = ()>,
1292 ) {
1293 let mut sim = self.start(log_override);
1294 let mut thunk_fut = pin!(thunk);
1295 let mut thunk_complete = false;
1296 loop {
1297 // The thunk always gets to run first until it completes. Completion is itself a
1298 // script barrier: after the body returns, keep stepping until every decision it
1299 // installed has been consumed (or report a decision that can never be honored).
1300 if !thunk_complete && futures::poll!(thunk_fut.as_mut()).is_ready() {
1301 thunk_complete = true;
1302 }
1303
1304 if thunk_complete {
1305 let Some(stuck) = script_unconsumed_description() else {
1306 break;
1307 };
1308 assert!(
1309 !sim.quiescence.is_quiescent(),
1310 "{}",
1311 script_stuck_error(&stuck)
1312 );
1313 sim.step().await;
1314 continue;
1315 }
1316
1317 if sim.quiescence.is_quiescent() || sim.quiescence.nondet_pending.get() {
1318 // The scheduler is parked: either no step can make progress until the thunk
1319 // sends new input (quiescent), or nondeterministic work is ready but a
1320 // settling test-side observation has paused the scheduler (nondet_pending).
1321 // Park until either the thunk is woken independently or the scheduler is
1322 // resumed. (`resumed()` is permit-based, so a resume that fired while polling
1323 // the thunk above is not lost.)
1324 tokio::select! {
1325 biased;
1326 () = &mut thunk_fut => break,
1327 () = sim.quiescence.resumed() => {}
1328 }
1329 sim.quiescence.nondet_pending.set(false);
1330 } else {
1331 // Run a single scheduler step to completion. This is awaited directly (not
1332 // raced against the thunk), so a step is atomic: the thunk is never polled
1333 // while a step is in flight, and a step is never cancelled mid-execution.
1334 sim.step().await;
1335 }
1336 }
1337 }
1338
1339 /// Consumes this instance and constructs the [`LaunchedSim`] state struct, which is
1340 /// advanced incrementally via [`LaunchedSim::step`].
1341 fn start<W: std::io::Write>(mut self, log_override: Option<W>) -> LaunchedSim<W> {
1342 let (
1343 async_dfirs,
1344 tick_dfirs,
1345 mut hooks,
1346 mut observation_hooks,
1347 mut inline_hooks,
1348 mut scripted_hooks,
1349 mut scripted_observation_hooks,
1350 mut scripted_inline_hooks,
1351 _registry,
1352 ) = self.dylib_result.take().unwrap();
1353
1354 // The generated code keys hooks and tick DFIRs by the same locations, so we can
1355 // move each tick's / observation's hooks out of the maps and attach them
1356 // directly. This lets the scheduler's hot paths avoid keyed lookups entirely.
1357 let not_ready_ticks = tick_dfirs
1358 .into_iter()
1359 .map(|(location, cluster_id, dfir)| {
1360 let key = SimLocation {
1361 location,
1362 cluster_id,
1363 };
1364 let LocationId::Tick {
1365 tick: _,
1366 parent_location,
1367 } = &key.location
1368 else {
1369 unreachable!("tick DFIRs are always keyed by a tick location")
1370 };
1371 let parent_location = (**parent_location).clone();
1372 let tick = SimTick {
1373 parent_location,
1374 cluster_id,
1375 dfir,
1376 hooks: hooks.remove(&key).unwrap_or_default(),
1377 scripted_hooks: scripted_hooks.remove(&key).unwrap_or_default(),
1378 inline_hooks: inline_hooks.remove(&key).unwrap_or_default(),
1379 scripted_inline_hooks: scripted_inline_hooks.remove(&key).unwrap_or_default(),
1380 location: key.location,
1381 };
1382 abort_assert!(
1383 !(tick.hooks.is_empty() && tick.scripted_hooks.is_empty()),
1384 "every tick DFIR must have at least one hook"
1385 );
1386 tick
1387 })
1388 .collect();
1389
1390 let (quiescence, script_coordinator) = CURRENT_SIM_CONNECTIONS.with(|connections| {
1391 let connections = connections.borrow();
1392 (
1393 connections.quiescence.clone(),
1394 connections.script_coordinator.clone(),
1395 )
1396 });
1397
1398 let not_ready_observations = async_dfirs
1399 .iter()
1400 .flat_map(|(location, cluster_id, _)| {
1401 let key = SimLocation {
1402 location: location.clone(),
1403 cluster_id: *cluster_id,
1404 };
1405 let cluster_id = *cluster_id;
1406 let unscripted = observation_hooks
1407 .remove(&key)
1408 .unwrap_or_default()
1409 .into_iter()
1410 .map(|hook| ObservationSlot::Unscripted { hook });
1411 let scripted = scripted_observation_hooks
1412 .remove(&key)
1413 .unwrap_or_default()
1414 .into_iter()
1415 .map(|hook| {
1416 let ScriptTarget::Observation { hook_id, .. } = hook.borrow().target()
1417 else {
1418 unreachable!("observation-registered scripted hook had a tick target")
1419 };
1420 ObservationSlot::Scripted { hook_id, hook }
1421 });
1422 unscripted.chain(scripted).map(move |hook| SimObservation {
1423 location: key.location.clone(),
1424 cluster_id,
1425 hook,
1426 })
1427 })
1428 .collect();
1429
1430 debug_assert!(
1431 hooks.is_empty()
1432 && observation_hooks.is_empty()
1433 && inline_hooks.is_empty()
1434 && scripted_hooks.is_empty()
1435 && scripted_observation_hooks.is_empty()
1436 && scripted_inline_hooks.is_empty(),
1437 "all hooks should belong to either a tick DFIR or a top-level location"
1438 );
1439
1440 LaunchedSim {
1441 async_dfirs,
1442 possibly_ready_ticks: vec![],
1443 not_ready_ticks,
1444 current_scripted_tick: None,
1445 current_scripted_observation: None,
1446 script_coordinator,
1447 possibly_ready_observations: vec![],
1448 not_ready_observations,
1449 log: if self.log {
1450 if let Some(w) = log_override {
1451 LogKind::Custom(w)
1452 } else {
1453 LogKind::Stderr
1454 }
1455 } else {
1456 LogKind::Null
1457 },
1458 quiescence,
1459 deterministic: self.deterministic,
1460 }
1461 }
1462}
1463
1464impl<T, O: Ordering, R: Retries> Clone for SimReceiver<T, O, R> {
1465 fn clone(&self) -> Self {
1466 *self
1467 }
1468}
1469
1470impl<T, O: Ordering, R: Retries> Copy for SimReceiver<T, O, R> {}
1471
1472/// How a [`QuiescenceCheckFuture`] resolves the "did the stream end?" check of
1473/// `assert_no_more`. Decided once the simulation has settled (run out of deterministic
1474/// work).
1475#[derive(Clone, Copy)]
1476enum QuiescenceBranch {
1477 /// Skip the check and continue the test. Only taken in exhaustive mode, where a
1478 /// sibling instance performs the check instead.
1479 Continue,
1480 /// Perform the check, then end this simulation instance (exhaustive mode), letting
1481 /// sibling instances continue past this point without forcing quiescence.
1482 CheckThenEnd,
1483 /// Perform the check and keep running. Taken when the simulation is already quiescent
1484 /// (the check is free) and in non-exhaustive modes.
1485 CheckAndKeepRunning,
1486}
1487
1488/// Decides how to run the quiescence check when the simulation has pending nondeterministic
1489/// work (ticks / observations) that the check would force to run.
1490fn decide_quiescence_branch() -> QuiescenceBranch {
1491 let (exhaustive, log) = CURRENT_SIM_CONNECTIONS.with(|connections| {
1492 let connections = connections.borrow();
1493 (connections.exhaustive, connections.log)
1494 });
1495
1496 if !exhaustive {
1497 return QuiescenceBranch::CheckAndKeepRunning;
1498 }
1499
1500 // In exhaustive mode, fork the search on a bolero decision. The exhaustive driver
1501 // enumerates `false` first, so the instance that performs the quiescence check is
1502 // explored *before* any instance that continues past this assertion. This ensures that
1503 // if the stream has extra output, the failure is attributed to this assertion (with a
1504 // decision trace leading exactly to the check) rather than leaking the extra messages
1505 // into a later assertion.
1506 let continue_without_check: bool = bolero::any();
1507 if continue_without_check {
1508 if log {
1509 eprintln!(
1510 "\n{}",
1511 "Continuing past quiescence assertion without checking (checked by an earlier instance)"
1512 .color(colored::Color::Cyan)
1513 .bold()
1514 );
1515 }
1516 QuiescenceBranch::Continue
1517 } else {
1518 if log {
1519 eprintln!(
1520 "\n{}",
1521 "Checking that no more messages arrive (this instance will end after the check)"
1522 .color(colored::Color::Cyan)
1523 .bold()
1524 );
1525 }
1526 QuiescenceBranch::CheckThenEnd
1527 }
1528}
1529
1530/// Ends the current simulation instance after a passing quiescence check, by panicking with
1531/// [`bolero::generator::bolero_generator::any::Error`], which bolero's engines treat as an
1532/// invalid input rather than a test failure. The instance has verified everything up to and
1533/// including the quiescence check; sibling instances continue past the check instead.
1534fn end_instance_after_quiescence_check() -> ! {
1535 bolero::generator::bolero_generator::any::assume(
1536 false,
1537 "simulation instance ended after quiescence check",
1538 );
1539 unreachable!()
1540}
1541
1542pin_project_lite::pin_project! {
1543 // The "and then the stream ends" half of `assert_no_more` (and thus of
1544 // `assert_yields_only*` / `collect_n_only`). First lets the simulation *settle* (see
1545 // `poll_settle`): if it settles to quiescence, the check is free and the test simply
1546 // continues. Otherwise, in exhaustive mode the search forks into a checking instance and
1547 // continuing instances (see `SimReceiver::assert_no_more` and
1548 // `decide_quiescence_branch`); in non-exhaustive modes the check runs, forcing the
1549 // pending work (which taints the simulation, via `try_next_bytes`).
1550 //
1551 // See [`FutureTrackingCaller`] for why `poll` is `#[track_caller]`.
1552 struct QuiescenceCheckFuture<F: Future<Output = ()>> {
1553 #[pin]
1554 check: F,
1555 settle: SettlePauseGuard,
1556 branch: Option<QuiescenceBranch>,
1557 }
1558}
1559
1560impl<F: Future<Output = ()>> QuiescenceCheckFuture<F> {
1561 fn new(check: F) -> Self {
1562 QuiescenceCheckFuture {
1563 check,
1564 settle: SettlePauseGuard::new(
1565 CURRENT_SIM_CONNECTIONS.with(|connections| connections.borrow().quiescence.clone()),
1566 ),
1567 branch: None,
1568 }
1569 }
1570}
1571
1572impl<F: Future<Output = ()>> Future for QuiescenceCheckFuture<F> {
1573 type Output = ();
1574
1575 #[track_caller]
1576 fn poll(mut self: Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Self::Output> {
1577 let this = self.as_mut().project();
1578
1579 if this.branch.is_none() {
1580 *this.branch = Some(if ready!(this.settle.poll_settle(cx)) {
1581 // Settled to quiescence deterministically, so the check is free.
1582 QuiescenceBranch::CheckAndKeepRunning
1583 } else {
1584 // The check would force nondeterministic work to run.
1585 decide_quiescence_branch()
1586 });
1587 }
1588
1589 match this.branch.unwrap() {
1590 QuiescenceBranch::Continue => Poll::Ready(()),
1591 QuiescenceBranch::CheckAndKeepRunning => this.check.poll(cx),
1592 QuiescenceBranch::CheckThenEnd => {
1593 ready!(this.check.poll(cx));
1594 end_instance_after_quiescence_check()
1595 }
1596 }
1597 }
1598}
1599
1600impl<T, O: Ordering, R: Retries> SimReceiver<T, O, R> {
1601 fn connections(&self) -> (Rc<Mutex<UnsyncReceiver<Bytes>>>, Rc<QuiescenceState>) {
1602 CURRENT_SIM_CONNECTIONS.with(|connections| {
1603 let connections = connections.borrow();
1604 let port = connections.external_registered.get(&self.0).unwrap();
1605 (
1606 connections.output_receivers.get(port).unwrap().clone(),
1607 connections.quiescence.clone(),
1608 )
1609 })
1610 }
1611
1612 /// See [`try_next_bytes`].
1613 async fn try_next_impl(&self) -> Option<T> {
1614 let (receiver, quiescence) = self.connections();
1615 try_next_bytes(&receiver, &quiescence)
1616 .await
1617 .map(|bytes| (self.2)(&bytes))
1618 }
1619
1620 /// Asserts that the stream has ended and no more messages can possibly arrive.
1621 ///
1622 /// If the check cannot be answered without running pending nondeterministic work (such
1623 /// as ticks with buffered inputs):
1624 /// - Under [`CompiledSim::exhaustive`], the search forks: one instance performs the
1625 /// check and ends there, while sibling instances skip the check and continue.
1626 /// - In other modes, the pending work runs; afterwards, sending more input and then
1627 /// attempting to receive output will panic.
1628 pub fn assert_no_more(self) -> impl Future<Output = ()>
1629 where
1630 T: Debug,
1631 {
1632 QuiescenceCheckFuture::new(FutureTrackingCaller {
1633 future: async move {
1634 if let Some(next) = self.try_next_impl().await {
1635 return Err(format!(
1636 "Stream yielded unexpected message: {:?}, expected termination",
1637 next
1638 ));
1639 }
1640 Ok(())
1641 },
1642 })
1643 }
1644}
1645
1646impl<T> SimReceiver<T, TotalOrder, ExactlyOnce> {
1647 /// Receives the next message from the simulation output stream, waiting (and letting the
1648 /// scheduler run any pending simulation work) until one is available. If the simulation
1649 /// becomes quiescent without producing a message, the test fails.
1650 ///
1651 /// This is safe to use in the middle of a test; to observe the *absence* of a message,
1652 /// use [`Self::try_next`] or [`Self::assert_no_more`].
1653 pub fn next(&self) -> impl use<'_, T> + Future<Output = T> {
1654 // Waiting for a message never "overruns" the simulation, even though the scheduler
1655 // may run nondeterministic ticks while we wait: if a message arrives, some pending
1656 // work was necessary to produce it (schedules that run *extra* work are also valid
1657 // executions, explored separately), and if the simulation quiesces instead, the test
1658 // fails right here — so no later observation can be affected by the overrun (the
1659 // taint set by `try_next_impl` is unobservable). See the module docs for the full
1660 // soundness reasoning.
1661 FutureTrackingCaller {
1662 future: async move {
1663 self.try_next_impl().await.ok_or_else(|| {
1664 "Stream ended (simulation quiescent), but another message was expected"
1665 .to_owned()
1666 })
1667 },
1668 }
1669 }
1670
1671 /// Receives the next message from the simulation output stream, or returns `None` if no
1672 /// more messages can possibly arrive.
1673 ///
1674 /// If answering requires forcing pending nondeterministic work to run, then afterwards,
1675 /// sending more input and then attempting to receive output will panic. Prefer
1676 /// [`Self::next`] (or [`Self::assert_no_more`]) when possible.
1677 pub async fn try_next(&self) -> Option<T> {
1678 self.try_next_impl().await
1679 }
1680
1681 /// Receives the next `n` messages from the simulation output stream, waiting (and letting
1682 /// the scheduler run any pending simulation work) until they are available. If the
1683 /// simulation becomes quiescent before `n` messages arrive, the test fails.
1684 ///
1685 /// Like [`Self::next`], this is safe to use in the middle of a test. It does not check
1686 /// that the stream ends afterwards; use [`Self::collect_n_only`] for that.
1687 pub fn collect_n<C: Default + Extend<T>>(
1688 &self,
1689 n: usize,
1690 ) -> impl use<'_, T, C> + Future<Output = C> {
1691 FutureTrackingCaller {
1692 future: async move {
1693 let mut out = C::default();
1694 for i in 0..n {
1695 // Like `next`, waiting for each message is safe mid-test; the taint on a
1696 // forced `None` is unobservable because the test fails below.
1697 if let Some(v) = self.try_next_impl().await {
1698 out.extend([v]);
1699 } else {
1700 return Err(format!(
1701 "Stream ended (simulation quiescent) after {} messages, but {} were expected",
1702 i, n
1703 ));
1704 }
1705 }
1706 Ok(out)
1707 },
1708 }
1709 }
1710
1711 /// Receives the next `n` messages (like [`Self::collect_n`]) and then asserts that the
1712 /// stream ends (like [`Self::assert_no_more`], forking the search in exhaustive mode).
1713 pub async fn collect_n_only<C: Default + Extend<T>>(self, n: usize) -> C
1714 where
1715 T: Debug,
1716 {
1717 let out = self.collect_n(n).await;
1718 self.assert_no_more().await;
1719 out
1720 }
1721
1722 /// Collects all remaining messages from the simulation output stream into a collection,
1723 /// waiting until no more messages can possibly arrive.
1724 ///
1725 /// If this has to force pending nondeterministic work to run, it should be the last
1726 /// observation of the test: afterwards, sending more input and then attempting to
1727 /// receive output will panic. When the number of expected messages is known, prefer
1728 /// [`Self::collect_n`] / [`Self::collect_n_only`].
1729 pub async fn collect<C: Default + Extend<T>>(self) -> C {
1730 let mut out = C::default();
1731 while let Some(v) = self.try_next_impl().await {
1732 out.extend([v]);
1733 }
1734 out
1735 }
1736
1737 /// Asserts that the stream yields exactly the expected sequence of messages, in order.
1738 /// This does not check that the stream ends, use [`Self::assert_yields_only`] for that.
1739 ///
1740 /// Like [`Self::next`], this is safe to use in the middle of a test.
1741 pub fn assert_yields<T2: Debug, I: IntoIterator<Item = T2>>(
1742 &self,
1743 expected: I,
1744 ) -> impl use<'_, T, T2, I> + Future<Output = ()>
1745 where
1746 T: Debug + PartialEq<T2>,
1747 {
1748 FutureTrackingCaller {
1749 future: async {
1750 let mut expected: VecDeque<T2> = expected.into_iter().collect();
1751
1752 while !expected.is_empty() {
1753 // Like `next`, waiting for each expected message is safe mid-test; the
1754 // taint on a forced `None` is unobservable because the test fails below.
1755 if let Some(next) = self.try_next_impl().await {
1756 let next_expected = expected.pop_front().unwrap();
1757 if next != next_expected {
1758 return Err(format!(
1759 "Stream yielded unexpected message: {:?}, expected: {:?}",
1760 next, next_expected
1761 ));
1762 }
1763 } else {
1764 return Err(format!(
1765 "Stream ended early, still expected: {:?}",
1766 expected
1767 ));
1768 }
1769 }
1770
1771 Ok(())
1772 },
1773 }
1774 }
1775
1776 /// Asserts that the stream yields only the expected sequence of messages, in order,
1777 /// and then ends (like [`Self::assert_no_more`], forking the search in exhaustive mode).
1778 pub fn assert_yields_only<T2: Debug, I: IntoIterator<Item = T2>>(
1779 &self,
1780 expected: I,
1781 ) -> impl use<'_, T, T2, I> + Future<Output = ()>
1782 where
1783 T: Debug + PartialEq<T2>,
1784 {
1785 ChainedFuture {
1786 first: self.assert_yields(expected),
1787 second: self.assert_no_more(),
1788 first_done: false,
1789 }
1790 }
1791}
1792
1793pin_project_lite::pin_project! {
1794 // A future that tracks the location of the `.await` call for better panic messages.
1795 //
1796 // `#[track_caller]` is important for us to create assertion methods because it makes
1797 // the panic backtrace show up at that method (instead of inside the call tree within
1798 // that method). This is e.g. what `Option::unwrap` uses. Unfortunately, `#[track_caller]`
1799 // does not work correctly for async methods (or `dyn Future` either), so we have to
1800 // create these concrete future types that (1) have `#[track_caller]` on their `poll()`
1801 // method and (2) have the `panic!` triggered in their `poll()` method (or in a directly
1802 // nested concrete future).
1803 struct FutureTrackingCaller<F> {
1804 #[pin]
1805 future: F,
1806 }
1807}
1808
1809impl<T, F: Future<Output = Result<T, String>>> Future for FutureTrackingCaller<F> {
1810 type Output = T;
1811
1812 #[track_caller]
1813 fn poll(mut self: Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Self::Output> {
1814 match ready!(self.as_mut().project().future.poll(cx)) {
1815 Ok(v) => Poll::Ready(v),
1816 Err(e) => panic!("{}", e),
1817 }
1818 }
1819}
1820
1821pin_project_lite::pin_project! {
1822 // A future that first awaits the first future, then the second, propagating caller info.
1823 //
1824 // See [`FutureTrackingCaller`] for context.
1825 struct ChainedFuture<F1: Future<Output = ()>, F2: Future<Output = ()>> {
1826 #[pin]
1827 first: F1,
1828 #[pin]
1829 second: F2,
1830 first_done: bool,
1831 }
1832}
1833
1834impl<F1: Future<Output = ()>, F2: Future<Output = ()>> Future for ChainedFuture<F1, F2> {
1835 type Output = ();
1836
1837 #[track_caller]
1838 fn poll(mut self: Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Self::Output> {
1839 if !self.first_done {
1840 ready!(self.as_mut().project().first.poll(cx));
1841 *self.as_mut().project().first_done = true;
1842 }
1843
1844 self.as_mut().project().second.poll(cx)
1845 }
1846}
1847
1848impl<T> SimReceiver<T, NoOrder, ExactlyOnce> {
1849 /// Receives the next `n` messages, sorted, and then asserts that the stream ends (like
1850 /// [`SimReceiver::assert_no_more`], forking the search in exhaustive mode). If the
1851 /// simulation becomes quiescent before `n` messages arrive, the test fails.
1852 ///
1853 /// Unlike [`collect_n`](SimReceiver::collect_n) on ordered streams, there is no variant
1854 /// of this API that skips the end-of-stream check. On an unordered stream, the set of
1855 /// messages that arrives *first* is not well-defined, so observing a strict prefix of
1856 /// the output would be sensitive to arrival orders that the simulator does not explore
1857 /// (delivery into the port is FIFO, with no ordering hook); sorting normalizes the
1858 /// permutation of the received messages, but not the choice of *subset*. The quiescence
1859 /// check makes the observation sound: it proves the `n` messages are *all* the messages
1860 /// the program can produce from the input so far, a set which does not depend on
1861 /// arrival order.
1862 pub async fn collect_n_sorted_only<C: Default + Extend<T> + AsMut<[T]>>(self, n: usize) -> C
1863 where
1864 T: Debug + Ord,
1865 {
1866 let out = FutureTrackingCaller {
1867 future: async move {
1868 let mut out = C::default();
1869 for i in 0..n {
1870 // Like `next`, waiting for each message is safe mid-test; the taint on a
1871 // forced `None` is unobservable because the test fails below.
1872 if let Some(v) = self.try_next_impl().await {
1873 out.extend([v]);
1874 } else {
1875 return Err(format!(
1876 "Stream ended (simulation quiescent) after {} messages, but {} were expected",
1877 i, n
1878 ));
1879 }
1880 }
1881 out.as_mut().sort();
1882 Ok(out)
1883 },
1884 }
1885 .await;
1886 self.assert_no_more().await;
1887 out
1888 }
1889
1890 /// Receives the next message, and then asserts that the stream ends (like
1891 /// [`SimReceiver::assert_no_more`], forking the search in exhaustive mode). If the
1892 /// simulation becomes quiescent without producing a message, the test fails.
1893 ///
1894 /// This is a shortcut for [`Self::collect_n_sorted_only`] with `n = 1`. Unlike
1895 /// [`next`](SimReceiver::next) on ordered streams, there is no variant that skips the
1896 /// end-of-stream check, because on an unordered stream *which* message arrives first is
1897 /// not well-defined; the check proves the message is the *only* one the program can
1898 /// produce from the input so far.
1899 pub async fn next_only(self) -> T
1900 where
1901 T: Debug + Ord,
1902 {
1903 let mut out: Vec<T> = self.collect_n_sorted_only(1).await;
1904 out.remove(0)
1905 }
1906
1907 /// Collects all remaining messages from the simulation output stream into a collection,
1908 /// sorting them. This will wait until no more messages can possibly arrive.
1909 ///
1910 /// If this has to force pending nondeterministic work to run, it should be the last
1911 /// observation of the test; see [`collect`](SimReceiver::collect).
1912 pub async fn collect_sorted<C: Default + Extend<T> + AsMut<[T]>>(self) -> C
1913 where
1914 T: Ord,
1915 {
1916 let mut collected = C::default();
1917 while let Some(v) = self.try_next_impl().await {
1918 collected.extend([v]);
1919 }
1920 collected.as_mut().sort();
1921 collected
1922 }
1923
1924 /// Asserts that the stream yields exactly the expected sequence of messages, in some order.
1925 /// This does not check that the stream ends, use [`Self::assert_yields_only_unordered`] for that.
1926 ///
1927 /// Like [`SimReceiver::next`], this is safe to use in the middle of a test.
1928 pub fn assert_yields_unordered<T2: Debug, I: IntoIterator<Item = T2>>(
1929 &self,
1930 expected: I,
1931 ) -> impl use<'_, T, T2, I> + Future<Output = ()>
1932 where
1933 T: Debug + PartialEq<T2>,
1934 {
1935 FutureTrackingCaller {
1936 future: async {
1937 let mut expected: Vec<T2> = expected.into_iter().collect();
1938
1939 while !expected.is_empty() {
1940 // Like `next`, waiting for each expected message is safe mid-test; the
1941 // taint on a forced `None` is unobservable because the test fails below.
1942 if let Some(next) = self.try_next_impl().await {
1943 let idx = expected.iter().enumerate().find(|(_, e)| &next == *e);
1944 if let Some((i, _)) = idx {
1945 expected.swap_remove(i);
1946 } else {
1947 return Err(format!("Stream yielded unexpected message: {:?}", next));
1948 }
1949 } else {
1950 return Err(format!(
1951 "Stream ended early, still expected: {:?}",
1952 expected
1953 ));
1954 }
1955 }
1956
1957 Ok(())
1958 },
1959 }
1960 }
1961
1962 /// Asserts that the stream yields only the expected sequence of messages, in some order,
1963 /// and then ends (like [`Self::assert_no_more`], forking the search in exhaustive mode).
1964 pub fn assert_yields_only_unordered<T2: Debug, I: IntoIterator<Item = T2>>(
1965 &self,
1966 expected: I,
1967 ) -> impl use<'_, T, T2, I> + Future<Output = ()>
1968 where
1969 T: Debug + PartialEq<T2>,
1970 {
1971 ChainedFuture {
1972 first: self.assert_yields_unordered(expected),
1973 second: self.assert_no_more(),
1974 first_done: false,
1975 }
1976 }
1977}
1978
1979impl<T, O: Ordering, R: Retries> SimSender<T, O, R> {
1980 fn with_sink<Out>(&self, thunk: impl FnOnce(&dyn Fn(T)) -> Out) -> Out {
1981 let (sender, quiescence) = CURRENT_SIM_CONNECTIONS.with(|connections| {
1982 let connections = connections.borrow();
1983 (
1984 connections
1985 .input_senders
1986 .get(connections.external_registered.get(&self.0).unwrap())
1987 .unwrap()
1988 .clone(),
1989 connections.quiescence.clone(),
1990 )
1991 });
1992
1993 let encode = self.2;
1994 thunk(&move |t| {
1995 sender.try_send(encode(&t).into()).unwrap();
1996 quiescence.resume();
1997 })
1998 }
1999}
2000
2001impl<T, O: Ordering> SimSender<T, O, ExactlyOnce> {
2002 /// Sends several messages to the simulation input. The messages will be asynchronously
2003 /// processed as part of the simulation, in non-deterministic order.
2004 pub fn send_many_unordered<I: IntoIterator<Item = T>>(&self, iter: I) {
2005 self.with_sink(|send| {
2006 for t in iter {
2007 send(t);
2008 }
2009 })
2010 }
2011}
2012
2013impl<T> SimSender<T, TotalOrder, ExactlyOnce> {
2014 /// Sends a message to the simulation input. The message will be asynchronously processed
2015 /// as part of the simulation.
2016 pub fn send(&self, t: T) {
2017 self.with_sink(|send| send(t));
2018 }
2019
2020 /// Sends several messages to the simulation input. The messages will be asynchronously
2021 /// processed as part of the simulation.
2022 pub fn send_many<I: IntoIterator<Item = T>>(&self, iter: I) {
2023 self.with_sink(|send| {
2024 for t in iter {
2025 send(t);
2026 }
2027 })
2028 }
2029}
2030
2031impl<T, O: Ordering, R: Retries> Clone for SimClusterReceiver<T, O, R> {
2032 fn clone(&self) -> Self {
2033 *self
2034 }
2035}
2036
2037impl<T, O: Ordering, R: Retries> Copy for SimClusterReceiver<T, O, R> {}
2038
2039impl<T, O: Ordering, R: Retries> SimClusterReceiver<T, O, R> {
2040 fn member_connections(
2041 &self,
2042 member_id: u32,
2043 ) -> (Rc<Mutex<UnsyncReceiver<Bytes>>>, Rc<QuiescenceState>) {
2044 CURRENT_SIM_CONNECTIONS.with(|connections| {
2045 let connections = connections.borrow();
2046 let port = connections.external_registered.get(&self.0).unwrap();
2047 let receivers = connections.cluster_output_receivers.get(port).unwrap();
2048 (
2049 receivers[&member_id].clone(),
2050 connections.quiescence.clone(),
2051 )
2052 })
2053 }
2054
2055 /// See [`try_next_bytes`].
2056 async fn try_next_impl(&self, member_id: u32) -> Option<T> {
2057 let (receiver, quiescence) = self.member_connections(member_id);
2058 try_next_bytes(&receiver, &quiescence)
2059 .await
2060 .map(|bytes| (self.2)(&bytes))
2061 }
2062
2063 /// Asserts that the stream from a specific cluster member has ended and no more messages
2064 /// can possibly arrive.
2065 ///
2066 /// If the check cannot be answered without running pending nondeterministic work (such
2067 /// as ticks with buffered inputs):
2068 /// - Under [`CompiledSim::exhaustive`], the search forks: one instance performs the
2069 /// check and ends there, while sibling instances skip the check and continue.
2070 /// - In other modes, the pending work runs; afterwards, sending more input and then
2071 /// attempting to receive output will panic.
2072 pub fn assert_no_more(self, member_id: u32) -> impl Future<Output = ()>
2073 where
2074 T: Debug,
2075 {
2076 QuiescenceCheckFuture::new(FutureTrackingCaller {
2077 future: async move {
2078 if let Some(next) = self.try_next_impl(member_id).await {
2079 return Err(format!(
2080 "Stream yielded unexpected message: {:?}, expected termination",
2081 next
2082 ));
2083 }
2084 Ok(())
2085 },
2086 })
2087 }
2088}
2089
2090impl<T> SimClusterReceiver<T, TotalOrder, ExactlyOnce> {
2091 /// Receives the next value from a specific cluster member, waiting (and letting the
2092 /// scheduler run any pending simulation work) until one is available. If the simulation
2093 /// becomes quiescent without producing a value, the test fails.
2094 ///
2095 /// This is safe to use in the middle of a test; to observe the *absence* of a value,
2096 /// use [`Self::try_next`].
2097 pub fn next(&self, member_id: u32) -> impl use<'_, T> + Future<Output = T> {
2098 // See `SimReceiver::next` for why waiting for a value never "overruns" the
2099 // simulation.
2100 FutureTrackingCaller {
2101 future: async move {
2102 self.try_next_impl(member_id).await.ok_or_else(|| {
2103 "Stream ended (simulation quiescent), but another message was expected"
2104 .to_owned()
2105 })
2106 },
2107 }
2108 }
2109
2110 /// Receives the next value from a specific cluster member, or returns `None` if no more
2111 /// values can possibly arrive.
2112 ///
2113 /// If answering requires forcing pending nondeterministic work to run, then afterwards,
2114 /// sending more input and then attempting to receive output will panic. Prefer
2115 /// [`Self::next`] when possible.
2116 pub async fn try_next(&self, member_id: u32) -> Option<T> {
2117 self.try_next_impl(member_id).await
2118 }
2119
2120 /// Collects all remaining values from a specific cluster member into a collection,
2121 /// waiting until no more values can possibly arrive.
2122 ///
2123 /// If this has to force pending nondeterministic work to run, it should be the last
2124 /// observation of the test; see [`SimReceiver::collect`].
2125 pub async fn collect<C: Default + Extend<T>>(self, member_id: u32) -> C {
2126 let mut out = C::default();
2127 while let Some(v) = self.try_next_impl(member_id).await {
2128 out.extend([v]);
2129 }
2130 out
2131 }
2132}
2133
2134impl<T> SimClusterReceiver<T, NoOrder, ExactlyOnce> {
2135 /// Receives the next `n` values from a specific cluster member, sorted, and then
2136 /// asserts that the stream ends (like [`Self::assert_no_more`], forking the search in
2137 /// exhaustive mode). If the simulation becomes quiescent before `n` values arrive, the
2138 /// test fails.
2139 ///
2140 /// There is no variant of this API that skips the end-of-stream check; see
2141 /// [`SimReceiver::collect_n_sorted_only`] for why observing a strict prefix of an
2142 /// unordered stream would be unsound.
2143 pub async fn collect_n_sorted_only<C: Default + Extend<T> + AsMut<[T]>>(
2144 self,
2145 member_id: u32,
2146 n: usize,
2147 ) -> C
2148 where
2149 T: Debug + Ord,
2150 {
2151 let out = FutureTrackingCaller {
2152 future: async move {
2153 let mut out = C::default();
2154 for i in 0..n {
2155 // Like `SimReceiver::next`, waiting for each message is safe mid-test;
2156 // the taint on a forced `None` is unobservable because the test fails
2157 // below.
2158 if let Some(v) = self.try_next_impl(member_id).await {
2159 out.extend([v]);
2160 } else {
2161 return Err(format!(
2162 "Stream ended (simulation quiescent) after {} messages, but {} were expected",
2163 i, n
2164 ));
2165 }
2166 }
2167 out.as_mut().sort();
2168 Ok(out)
2169 },
2170 }
2171 .await;
2172 self.assert_no_more(member_id).await;
2173 out
2174 }
2175
2176 /// Receives the next value from a specific cluster member, and then asserts that the
2177 /// stream ends (like [`Self::assert_no_more`], forking the search in exhaustive mode).
2178 /// If the simulation becomes quiescent without producing a value, the test fails.
2179 ///
2180 /// This is a shortcut for [`Self::collect_n_sorted_only`] with `n = 1`; see
2181 /// [`SimReceiver::next_only`] for why there is no variant that skips the end-of-stream
2182 /// check.
2183 pub async fn next_only(self, member_id: u32) -> T
2184 where
2185 T: Debug + Ord,
2186 {
2187 let mut out: Vec<T> = self.collect_n_sorted_only(member_id, 1).await;
2188 out.remove(0)
2189 }
2190
2191 /// Collects all remaining values from a specific cluster member, sorted, waiting until no
2192 /// more values can possibly arrive.
2193 ///
2194 /// If this has to force pending nondeterministic work to run, it should be the last
2195 /// observation of the test; see [`SimReceiver::collect`].
2196 pub async fn collect_sorted<C: Default + Extend<T> + AsMut<[T]>>(self, member_id: u32) -> C
2197 where
2198 T: Ord,
2199 {
2200 let mut collected = C::default();
2201 while let Some(v) = self.try_next_impl(member_id).await {
2202 collected.extend([v]);
2203 }
2204 collected.as_mut().sort();
2205 collected
2206 }
2207}
2208
2209impl<T, O: Ordering, R: Retries> SimClusterSender<T, O, R> {
2210 fn with_sink<Out>(&self, thunk: impl FnOnce(&dyn Fn(u32, T)) -> Out) -> Out {
2211 let (senders, quiescence) = CURRENT_SIM_CONNECTIONS.with(|connections| {
2212 let connections = connections.borrow();
2213 (
2214 connections
2215 .cluster_input_senders
2216 .get(connections.external_registered.get(&self.0).unwrap())
2217 .unwrap()
2218 .clone(),
2219 connections.quiescence.clone(),
2220 )
2221 });
2222
2223 let encode = self.2;
2224 thunk(&move |member_id: u32, t: T| {
2225 senders[&member_id].try_send(encode(&t).into()).unwrap();
2226 quiescence.resume();
2227 })
2228 }
2229}
2230
2231impl<T, O: Ordering> SimClusterSender<T, O, ExactlyOnce> {
2232 /// Sends multiple values to specific cluster members. The messages will be asynchronously
2233 /// processed as part of the simulation, in non-deterministic order.
2234 pub fn send_many_unordered<I: IntoIterator<Item = (u32, T)>>(&self, iter: I) {
2235 self.with_sink(|send| {
2236 for (member_id, t) in iter {
2237 send(member_id, t);
2238 }
2239 })
2240 }
2241}
2242
2243impl<T> SimClusterSender<T, TotalOrder, ExactlyOnce> {
2244 /// Sends a value to a specific cluster member.
2245 pub fn send(&self, member_id: u32, t: T) {
2246 self.with_sink(|send| send(member_id, t));
2247 }
2248
2249 /// Sends multiple values to specific cluster members.
2250 pub fn send_many<I: IntoIterator<Item = (u32, T)>>(&self, iter: I) {
2251 self.with_sink(|send| {
2252 for (member_id, t) in iter {
2253 send(member_id, t);
2254 }
2255 })
2256 }
2257}
2258
2259enum LogKind<W: std::io::Write> {
2260 Null,
2261 Stderr,
2262 Custom(W),
2263}
2264
2265// via https://www.reddit.com/r/rust/comments/t69sld/is_there_a_way_to_allow_either_stdfmtwrite_or/
2266impl<W: std::io::Write> std::fmt::Write for LogKind<W> {
2267 fn write_str(&mut self, s: &str) -> Result<(), std::fmt::Error> {
2268 match self {
2269 LogKind::Null => Ok(()),
2270 LogKind::Stderr => {
2271 eprint!("{}", s);
2272 Ok(())
2273 }
2274 LogKind::Custom(w) => w.write_all(s.as_bytes()).map_err(|_| std::fmt::Error),
2275 }
2276 }
2277}
2278
2279/// A tick-scoped DFIR together with the hooks that feed it data.
2280struct SimTick {
2281 /// The tick's location, used to match this tick to an outstanding script group.
2282 location: LocationId,
2283 /// The location of the process/cluster the tick lives on, used to match this tick
2284 /// against the async DFIR that produces its input data.
2285 parent_location: LocationId,
2286 /// The cluster member ID, if the tick lives on a cluster.
2287 cluster_id: Option<u32>,
2288 /// The tick DFIR, executed once per tick.
2289 dfir: DfirErased,
2290 /// Hooks (e.g. from `batch`) resolved *before* the tick runs, deciding what data to
2291 /// release into it.
2292 hooks: Vec<Box<dyn TickInputHook>>,
2293 /// Scripted hooks (bound to test-side handles), also resolved before the tick runs.
2294 /// Kept separate from `hooks` so the scheduler can apply the script-specific rules
2295 /// (the boundary scan and `blocks_tick`), and shared (`Rc`) with the per-instance
2296 /// registry that test-side handles resolve through (see [`ScriptedRuntimeHook`]).
2297 scripted_hooks: Vec<Rc<RefCell<dyn ScriptedTickInputHook>>>,
2298 /// Hooks (e.g. from `assume_ordering` inside the tick) resolved *while* the tick DFIR
2299 /// is running, via a `tokio::select!` loop, for operators that block on ordering
2300 /// decisions mid-tick.
2301 inline_hooks: Vec<Box<dyn InlineHook>>,
2302 scripted_inline_hooks: Vec<Rc<RefCell<dyn crate::sim::runtime::ScriptedInlineHook>>>,
2303}
2304
2305impl SimTick {
2306 /// Whether the scheduler can execute this tick right now.
2307 fn can_run(&self) -> bool {
2308 // No scripted hook may have a queued decision that is not yet honorable
2309 // (such a decision names this tick's *next* execution, so the tick must wait
2310 // until it can be honored in full)...
2311 !self
2312 .scripted_hooks
2313 .iter()
2314 .any(|hook| hook.borrow().blocks_tick())
2315 // ...and at least one hook must be able to trigger the tick.
2316 && (self.hooks.iter().any(|hook| hook.can_trigger_tick())
2317 || self
2318 .scripted_hooks
2319 .iter()
2320 .any(|hook| hook.borrow().can_trigger_tick()))
2321 }
2322}
2323
2324/// A single top-level hook (e.g. from `assume_ordering` on a non-tick stream) that needs
2325/// scheduling decisions, but has no tick DFIR to execute. The scheduler just resolves the
2326/// hook.
2327///
2328/// Each top-level hook is its own observation ("its own virtual tick"), even when several
2329/// hooks live at the same location: unlike a tick's hooks, which one atomic tick
2330/// execution consumes together, co-located top-level hooks are causally independent
2331/// operators, so resolving them jointly would only couple their decisions. Grouping them
2332/// would both add redundant schedules (releasing jointly is equivalent to releasing in
2333/// consecutive steps, which is explored anyway) and *lose* schedules for hook kinds whose
2334/// decisions always release when resolved (a fold could never stay silent while a
2335/// co-located sibling acts). With one hook per observation, "act" and "stay silent" are
2336/// expressed purely by the scheduler picking or not picking the observation, and a picked
2337/// observation always makes a nontrivial decision.
2338struct SimObservation {
2339 /// The top-level location, used to match this observation against the async DFIR that
2340 /// produces its input data (and, for a scripted hook, against an outstanding script
2341 /// group).
2342 location: LocationId,
2343 /// The cluster member ID, if the location is a cluster.
2344 cluster_id: Option<u32>,
2345 /// The hook resolved when the scheduler selects this observation.
2346 hook: ObservationSlot,
2347}
2348
2349/// The single hook of a [`SimObservation`]: either an ordinary autonomous hook, or a
2350/// scripted hook (bound to a test-side handle), tagged with its hook ID so a script group
2351/// can be matched to exactly this observation.
2352enum ObservationSlot {
2353 /// An ordinary autonomous hook, owned by the scheduler.
2354 Unscripted { hook: Box<dyn ObservationHook> },
2355 /// A hook bound to a test-side handle, shared (`Rc`) with the per-instance registry.
2356 Scripted {
2357 /// The bound handle's ID, used to match a script group to this observation.
2358 hook_id: usize,
2359 hook: Rc<RefCell<dyn ScriptedObservationHook>>,
2360 },
2361}
2362
2363impl SimObservation {
2364 /// Whether the scheduler can resolve this observation's hook right now.
2365 fn can_run(&self) -> bool {
2366 match &self.hook {
2367 // Running an observation *is* releasing, so any pending input makes an
2368 // unscripted observation runnable.
2369 ObservationSlot::Unscripted { hook } => hook.has_pending_input(),
2370 ObservationSlot::Scripted { hook, .. } => hook.borrow().can_fire(),
2371 }
2372 }
2373}
2374
2375/// A running simulation, which manages the async DFIRs, tick DFIRs, and hook-based
2376/// scheduling decisions for non-deterministic operators like `batch` and `assume_ordering`.
2377///
2378/// This struct holds all simulator state across scheduler steps. Each [`Self::step`] performs
2379/// one of three kinds of work:
2380/// - **Async DFIRs**: long-running top-level dataflows (one per process/cluster member) that
2381/// produce data consumed by ticks and observations.
2382/// - **Ticks**: tick-scoped DFIRs that execute a single tick. Before running, their associated
2383/// hooks (e.g. from `batch`) are resolved to decide what data to release into the tick.
2384/// - **Observations**: top-level locations that have hooks (e.g. from `assume_ordering` on a
2385/// non-tick stream) needing decisions, but no tick DFIR to execute. The scheduler just
2386/// resolves their hooks.
2387struct LaunchedSim<W: std::io::Write> {
2388 /// Top-level async DFIRs, one per process/cluster member. These run continuously and
2389 /// produce data that feeds into ticks and observations.
2390 async_dfirs: Vec<(LocationId, Option<u32>, DfirErased)>,
2391 /// Ticks whose parent async DFIR has made progress, so they may be ready to run.
2392 /// The scheduler further filters these by checking whether their hooks have pending decisions.
2393 possibly_ready_ticks: Vec<SimTick>,
2394 /// Ticks whose parent async DFIR has not yet made progress since they were last checked.
2395 not_ready_ticks: Vec<SimTick>,
2396 /// The tick owned by the one sealed, outstanding scripted decision group. It is kept
2397 /// outside the ordinary ready lists until it executes and consumes that group.
2398 current_scripted_tick: Option<SimTick>,
2399 current_scripted_observation: Option<SimObservation>,
2400 /// Coordinates the decision group shared with test-side hook handles.
2401 script_coordinator: Rc<RefCell<ScriptCoordinator>>,
2402 /// Observations whose async DFIR has made progress, so their hooks may have decisions
2403 /// to resolve.
2404 possibly_ready_observations: Vec<SimObservation>,
2405 /// Observations whose async DFIR has not yet made progress since they were last checked.
2406 not_ready_observations: Vec<SimObservation>,
2407 log: LogKind<W>,
2408 /// Represents quiescence state of the simulation.
2409 quiescence: Rc<QuiescenceState>,
2410 /// When true, this simulation runs in deterministic mode: no fuzzer entropy is ever
2411 /// drawn, every unsafe operator with meaningful input must be scripted, and at most
2412 /// one tick is ever runnable (see `SimFlow::deterministic`).
2413 deterministic: bool,
2414}
2415
2416impl<W: std::io::Write> LaunchedSim<W> {
2417 /// Runs a single step of the simulation scheduler.
2418 ///
2419 /// A step first advances all async DFIRs; if none of them made progress, it instead runs
2420 /// one ready tick or resolves one ready observation. If nothing at all can make progress,
2421 /// the simulation is quiescent: this signals waiting receivers and returns; the driver is
2422 /// responsible for parking until new external input arrives (see
2423 /// [`QuiescenceState::resumed`]).
2424 ///
2425 /// This future is always awaited to completion by the driver, so a step is atomic: user
2426 /// code never runs (and never observes intermediate state) while a step is in flight.
2427 async fn step(&mut self) {
2428 // A group remains joinable only while the test body is in the same synchronous poll
2429 // that created it. Starting any scheduler step seals it and moves its tick out of
2430 // the ordinary lists exactly once; `Some(current)` then means that tick exclusively
2431 // owns the one outstanding group until it executes.
2432 let outstanding_target = {
2433 let mut coordinator = self.script_coordinator.borrow_mut();
2434 coordinator.current.as_mut().map(|group| {
2435 group.sealed = true;
2436 group.target.clone()
2437 })
2438 };
2439 match outstanding_target {
2440 Some(ScriptTarget::Tick {
2441 location:
2442 SimLocation {
2443 location: group_location,
2444 cluster_id: group_cluster_id,
2445 },
2446 }) => {
2447 abort_assert!(
2448 self.current_scripted_observation.is_none(),
2449 "scripted observation remained active for a tick group"
2450 );
2451 if self.current_scripted_tick.is_none() {
2452 let matches_group = |tick: &SimTick| {
2453 tick.location == group_location && tick.cluster_id == group_cluster_id
2454 };
2455 self.current_scripted_tick = self
2456 .possibly_ready_ticks
2457 .iter()
2458 .position(matches_group)
2459 .map(|index| self.possibly_ready_ticks.swap_remove(index))
2460 .or_else(|| {
2461 self.not_ready_ticks
2462 .iter()
2463 .position(matches_group)
2464 .map(|index| self.not_ready_ticks.swap_remove(index))
2465 });
2466 }
2467 let tick = self.current_scripted_tick.as_ref().unwrap();
2468 abort_assert!(
2469 tick.location == group_location && tick.cluster_id == group_cluster_id,
2470 "outstanding scripted group changed before its tick executed"
2471 );
2472 }
2473 Some(ScriptTarget::Observation {
2474 location:
2475 SimLocation {
2476 location: group_location,
2477 cluster_id: group_cluster_id,
2478 },
2479 hook_id,
2480 }) => {
2481 abort_assert!(
2482 self.current_scripted_tick.is_none(),
2483 "scripted tick remained active for an observation group"
2484 );
2485 if self.current_scripted_observation.is_none() {
2486 let matches_group = |observation: &SimObservation| {
2487 observation.location == group_location
2488 && observation.cluster_id == group_cluster_id
2489 && matches!(observation.hook, ObservationSlot::Scripted { hook_id: id, .. } if id == hook_id)
2490 };
2491 self.current_scripted_observation = self
2492 .possibly_ready_observations
2493 .iter()
2494 .position(matches_group)
2495 .map(|index| self.possibly_ready_observations.swap_remove(index))
2496 .or_else(|| {
2497 self.not_ready_observations
2498 .iter()
2499 .position(matches_group)
2500 .map(|index| self.not_ready_observations.swap_remove(index))
2501 });
2502 }
2503 abort_assert!(
2504 self.current_scripted_observation.is_some(),
2505 "outstanding scripted group did not match an observation"
2506 );
2507 }
2508 None => abort_assert!(
2509 self.current_scripted_tick.is_none() && self.current_scripted_observation.is_none(),
2510 "scripted action remained active without an outstanding group"
2511 ),
2512 }
2513
2514 let mut any_made_progress = false;
2515 for (loc, c_id, dfir) in &mut self.async_dfirs {
2516 if dfir.run_tick().await {
2517 any_made_progress = true;
2518
2519 // This async DFIR may have produced new data, so the ticks and observations
2520 // it feeds may now be ready.
2521 self.possibly_ready_ticks
2522 .extend(self.not_ready_ticks.extract_if(.., |tick| {
2523 tick.parent_location == *loc && tick.cluster_id == *c_id
2524 }));
2525 self.possibly_ready_observations.extend(
2526 self.not_ready_observations
2527 .extract_if(.., |obs| obs.location == *loc && obs.cluster_id == *c_id),
2528 );
2529 }
2530 }
2531
2532 if any_made_progress {
2533 return;
2534 }
2535
2536 // The **boundary scan**: the async dataflows have stopped making progress and we
2537 // are about to consider running ticks — the first moment where a missing scripted
2538 // decision could influence what happens next. Check ticks exposed by async progress,
2539 // plus the active scripted tick (which lives outside the ordinary ready lists).
2540 for tick in self
2541 .possibly_ready_ticks
2542 .iter()
2543 .chain(self.current_scripted_tick.iter())
2544 {
2545 for hook in &tick.scripted_hooks {
2546 if let Err(message) = hook.borrow().boundary_check() {
2547 panic!("{}", message);
2548 }
2549 }
2550 }
2551
2552 for observation in self
2553 .possibly_ready_observations
2554 .iter()
2555 .chain(self.current_scripted_observation.iter())
2556 {
2557 if let ObservationSlot::Scripted { hook, .. } = &observation.hook
2558 && let Err(message) = hook.borrow().boundary_check()
2559 {
2560 panic!("{}", message);
2561 }
2562 }
2563
2564 // A fully scripted tick needs at least one decision that can eventually trigger
2565 // it. There is exactly one outstanding group, so only its owned tick can contain
2566 // a newly installed group in which no decision can trigger.
2567 if let Some(tick) = &self.current_scripted_tick
2568 && tick.hooks.is_empty()
2569 {
2570 let has_pending_decision = tick
2571 .scripted_hooks
2572 .iter()
2573 .any(|hook| hook.borrow().has_decision());
2574 let any_pending_decision_can_eventually_trigger =
2575 tick.scripted_hooks.iter().any(|hook| {
2576 let hook = hook.borrow();
2577 // A decision that is not yet honorable may become honorable and
2578 // trigger once more data arrives, so it does not fail this check.
2579 hook.has_decision() && (hook.blocks_tick() || hook.can_trigger_tick())
2580 });
2581
2582 if has_pending_decision && !any_pending_decision_can_eventually_trigger {
2583 let mut details = String::new();
2584 let member = tick
2585 .cluster_id
2586 .map(|m| format!(" (cluster member {m})"))
2587 .unwrap_or_default();
2588 for hook in &tick.scripted_hooks {
2589 let hook = hook.borrow();
2590 if let Some(decision) = hook.describe_decision() {
2591 let loc = ScriptedHookControl::location_meta(&*hook).location;
2592 use std::fmt::Write;
2593 write!(details, "\n {} on the hook at {}{}", decision, loc, member)
2594 .unwrap();
2595 }
2596 }
2597 panic!(
2598 "none of the scripted decisions in this group can trigger their tick, so the tick can never run; at least one decision in the group must trigger it:{}",
2599 details
2600 );
2601 }
2602 }
2603
2604 use bolero::generator::*;
2605
2606 // Send anything that can't make a scheduling decision back to the not-ready lists.
2607 self.not_ready_ticks.extend(
2608 self.possibly_ready_ticks
2609 .extract_if(.., |tick| !tick.can_run()),
2610 );
2611 self.not_ready_observations.extend(
2612 self.possibly_ready_observations
2613 .extract_if(.., |obs| !obs.can_run()),
2614 );
2615
2616 let scripted_tick_runnable = self
2617 .current_scripted_tick
2618 .as_ref()
2619 .is_some_and(SimTick::can_run);
2620 let scripted_observation_runnable = self
2621 .current_scripted_observation
2622 .as_ref()
2623 .is_some_and(SimObservation::can_run);
2624
2625 if self.possibly_ready_ticks.is_empty()
2626 && !scripted_tick_runnable
2627 && !scripted_observation_runnable
2628 && self.possibly_ready_observations.is_empty()
2629 {
2630 // Classify why the outstanding scripted group (if any) is stuck, so the
2631 // suspended test-side await renders the right error: `true` when every
2632 // queued decision is satisfiable but none can trigger the tick — given
2633 // quiescence, no unscripted input on the tick can trigger it either, or the
2634 // tick would be runnable.
2635 self.script_coordinator.borrow_mut().stuck_cannot_trigger =
2636 self.current_scripted_tick.as_ref().is_some_and(|tick| {
2637 let mut queued = tick
2638 .scripted_hooks
2639 .iter()
2640 .filter(|hook| hook.borrow().has_decision())
2641 .peekable();
2642 queued.peek().is_some() && queued.all(|hook| !hook.borrow().blocks_tick())
2643 });
2644
2645 // Signal quiescence, waking receivers waiting for data (their streams end). The
2646 // driver is responsible for parking until new input arrives.
2647 self.quiescence.enter_quiescence();
2648 } else if self.quiescence.pause_nondet.get() > 0 {
2649 // The test is querying whether the simulation can quiesce without
2650 // nondeterministic work (see `SettlePauseGuard::poll_settle`). Report that
2651 // ticks/observations are pending and pause; the driver parks until the test
2652 // decides how to proceed.
2653 self.quiescence.nondet_pending.set(true);
2654 self.quiescence.wake_settled();
2655 } else {
2656 let ordinary_tick_count = self.possibly_ready_ticks.len();
2657 let scripted_tick_index = ordinary_tick_count;
2658 let observation_start = scripted_tick_index + usize::from(scripted_tick_runnable);
2659 let scripted_observation_index =
2660 observation_start + self.possibly_ready_observations.len();
2661 let candidate_count =
2662 scripted_observation_index + usize::from(scripted_observation_runnable);
2663 let next_tick_or_obs = if self.deterministic {
2664 for tick in self.possibly_ready_ticks.iter().chain(
2665 self.current_scripted_tick
2666 .iter()
2667 .filter(|_| scripted_tick_runnable),
2668 ) {
2669 for hook in &tick.hooks {
2670 assert!(
2671 hook.only_one_possible_decision(),
2672 "{}",
2673 crate::sim::runtime::render_unhooked_nondet_error(hook.location_meta())
2674 );
2675 }
2676 }
2677 for obs in &self.possibly_ready_observations {
2678 if let ObservationSlot::Unscripted { hook } = &obs.hook
2679 && !hook.only_one_possible_decision()
2680 {
2681 panic!(
2682 "{}",
2683 crate::sim::runtime::render_unhooked_nondet_error(hook.location_meta())
2684 );
2685 }
2686 }
2687 // Each action on its own may be free of choices, but the order in
2688 // which they run is not determined, and it can be observable.
2689 assert!(
2690 candidate_count <= 1,
2691 "deterministic simulation reached a state with more than one runnable tick/observation; the order in which they run is not deterministic\nhelp: script the involved operators so the schedule is explicit, or run under `fuzz` / `exhaustive` instead"
2692 );
2693 0
2694 } else {
2695 (0..candidate_count).any()
2696 };
2697
2698 if next_tick_or_obs < observation_start {
2699 let is_scripted_tick = next_tick_or_obs == scripted_tick_index;
2700 let mut tick = if is_scripted_tick {
2701 self.current_scripted_tick.take().unwrap()
2702 } else {
2703 self.possibly_ready_ticks.remove(next_tick_or_obs)
2704 };
2705
2706 match &mut self.log {
2707 LogKind::Null => {}
2708 LogKind::Stderr => {
2709 if let Some(cid) = &tick.cluster_id {
2710 eprintln!(
2711 "\n{}",
2712 format!("Running Tick (Cluster Member {})", cid)
2713 .color(colored::Color::Magenta)
2714 .bold()
2715 )
2716 } else {
2717 eprintln!("\n{}", "Running Tick".color(colored::Color::Magenta).bold())
2718 }
2719 }
2720 LogKind::Custom(writer) => {
2721 writeln!(
2722 writer,
2723 "\n{}",
2724 "Running Tick".color(colored::Color::Magenta).bold()
2725 )
2726 .unwrap();
2727 }
2728 }
2729
2730 let mut asterisk_indenter = |_line_no, write: &mut dyn std::fmt::Write| {
2731 write.write_str(&"*".color(colored::Color::Magenta).bold())?;
2732 write.write_str(" ")
2733 };
2734
2735 let mut tick_decision_writer = (!matches!(self.log, LogKind::Null)).then(|| {
2736 indenter::indented(&mut self.log).with_format(indenter::Format::Custom {
2737 inserter: &mut asterisk_indenter,
2738 })
2739 });
2740
2741 run_hooks(
2742 tick_decision_writer.as_mut(),
2743 &mut tick.hooks,
2744 &tick.scripted_hooks,
2745 );
2746
2747 let run_tick_future = tick.dfir.run_tick();
2748 if !tick.inline_hooks.is_empty() || !tick.scripted_inline_hooks.is_empty() {
2749 let mut run_tick_future_pinned = pin!(run_tick_future);
2750 let deterministic = self.deterministic;
2751
2752 loop {
2753 tokio::select! {
2754 biased;
2755 r = &mut run_tick_future_pinned => {
2756 abort_assert!(r, "runnable tick's DFIR run_tick() returned false");
2757 break;
2758 }
2759 () = async {} => {
2760 for hook in &tick.scripted_inline_hooks {
2761 if hook.borrow().has_pending_input() {
2762 let run = hook.borrow_mut().run_decision(
2763 tick_decision_writer
2764 .as_mut()
2765 .map(|w| w as &mut dyn std::fmt::Write),
2766 );
2767 // The error is reported here, on the host side of
2768 // the dylib boundary (unwinding across it aborts).
2769 if let Err(message) = run {
2770 panic!("{}", message);
2771 }
2772 }
2773 }
2774 if !tick.inline_hooks.is_empty() {
2775 bolero_generator::any::scope::borrow_with(|driver| {
2776 for hook in tick.inline_hooks.iter_mut() {
2777 if hook.has_pending_input() {
2778 // In deterministic mode there is no fuzzer
2779 // to decide for this operator; it may only
2780 // proceed when exactly one outcome is
2781 // possible.
2782 if deterministic {
2783 assert!(
2784 hook.only_one_possible_decision(),
2785 "{}",
2786 crate::sim::runtime::render_unhooked_nondet_error(
2787 hook.location_meta()
2788 )
2789 );
2790 }
2791 hook.autonomous_decision(driver);
2792 hook.release_decision(
2793 tick_decision_writer
2794 .as_mut()
2795 .map(|w| w as &mut dyn std::fmt::Write),
2796 );
2797 }
2798 }
2799 });
2800 }
2801 }
2802 }
2803 }
2804 } else {
2805 let made_progress = run_tick_future.await;
2806 abort_assert!(
2807 made_progress,
2808 "runnable tick's DFIR run_tick() returned false"
2809 );
2810 }
2811
2812 if is_scripted_tick {
2813 for hook in &tick.scripted_inline_hooks {
2814 abort_assert!(
2815 !hook.borrow().has_decision(),
2816 "tick completed without consuming a scripted inline decision"
2817 );
2818 }
2819 let group = self.script_coordinator.borrow_mut().current.take();
2820 abort_assert!(
2821 group.is_some(),
2822 "scripted tick executed without an outstanding group"
2823 );
2824 }
2825 self.possibly_ready_ticks.push(tick);
2826 } else {
2827 let is_scripted_observation = next_tick_or_obs == scripted_observation_index;
2828 let observation = if is_scripted_observation {
2829 self.current_scripted_observation.as_mut().unwrap()
2830 } else {
2831 &mut self.possibly_ready_observations[next_tick_or_obs - observation_start]
2832 };
2833 let log_writer = (!matches!(self.log, LogKind::Null)).then_some(&mut self.log);
2834 match &mut observation.hook {
2835 ObservationSlot::Unscripted { hook } => {
2836 run_observation_hook(log_writer, &mut **hook);
2837 }
2838 ObservationSlot::Scripted { hook, .. } => {
2839 abort_assert!(
2840 hook.borrow().can_fire(),
2841 "scripted observation ran without a releasing decision"
2842 );
2843 hook.borrow_mut()
2844 .run_decision(log_writer.map(|w| w as &mut dyn std::fmt::Write));
2845 }
2846 }
2847 if is_scripted_observation {
2848 let group = self.script_coordinator.borrow_mut().current.take();
2849 abort_assert!(group.is_some(), "scripted observation ran without a group");
2850 let observation = self.current_scripted_observation.take().unwrap();
2851 self.possibly_ready_observations.push(observation);
2852 }
2853 }
2854 }
2855 }
2856}
2857
2858fn run_hooks<W: std::fmt::Write>(
2859 mut tick_decision_writer: Option<&mut W>,
2860 hooks: &mut [Box<dyn TickInputHook>],
2861 scripted_hooks: &[Rc<RefCell<dyn ScriptedTickInputHook>>],
2862) {
2863 // Scripted hooks own and release their decisions without entropy. Run them completely
2864 // before considering regular hooks; only regular hooks need a Bolero driver.
2865 let mut made_triggering_decision = false;
2866 for hook in scripted_hooks {
2867 let mut hook = hook.borrow_mut();
2868 // Whether a scripted decision triggers is known before running it.
2869 made_triggering_decision |= hook.can_trigger_tick();
2870 hook.run_decision(
2871 tick_decision_writer
2872 .as_deref_mut()
2873 .map(|w| w as &mut dyn std::fmt::Write),
2874 );
2875 }
2876
2877 if !hooks.is_empty() {
2878 let mut decided = vec![false; hooks.len()];
2879 let mut remaining_decision_count = hooks.len();
2880 bolero::generator::bolero_generator::any::scope::borrow_with(|driver| {
2881 // First, resolve every hook that faces no choice (its decision consumes no
2882 // entropy). Doing this before the second pass lets the final undecided hook
2883 // be forced to trigger when no earlier hook made a triggering decision.
2884 for (hook, decided) in hooks.iter_mut().zip(decided.iter_mut()) {
2885 if hook.only_one_possible_decision() {
2886 // The no-choice decision can still trigger the tick (the passthrough
2887 // singleton always releases the latest value), so its result counts.
2888 made_triggering_decision |= hook.autonomous_decision(driver, false);
2889 *decided = true;
2890 remaining_decision_count -= 1;
2891 }
2892 }
2893
2894 for (hook, decided) in hooks.iter_mut().zip(decided.iter()) {
2895 if !decided {
2896 made_triggering_decision |= hook.autonomous_decision(
2897 driver,
2898 !made_triggering_decision && remaining_decision_count == 1,
2899 );
2900 remaining_decision_count -= 1;
2901 }
2902
2903 hook.release_decision(
2904 tick_decision_writer
2905 .as_deref_mut()
2906 .map(|w| w as &mut dyn std::fmt::Write),
2907 );
2908 }
2909 });
2910 }
2911
2912 abort_assert!(
2913 made_triggering_decision,
2914 "runnable tick had no hook make a triggering decision"
2915 );
2916}
2917
2918/// Resolves a single unscripted observation hook. The observation was only scheduled
2919/// because it has pending input (running an observation *is* releasing), so its
2920/// autonomous decision must stage a release — running an observation without releasing
2921/// would be a wasted schedule step the exploration must not contain.
2922fn run_observation_hook<W: std::fmt::Write>(
2923 writer: Option<&mut W>,
2924 hook: &mut dyn ObservationHook,
2925) {
2926 bolero::generator::bolero_generator::any::scope::borrow_with(|driver| {
2927 hook.autonomous_decision(driver);
2928 });
2929 // `release_decision` panics if the autonomous decision staged nothing, so a
2930 // contract violation cannot pass silently.
2931 hook.release_decision(writer.map(|w| w as &mut dyn std::fmt::Write));
2932}