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