Skip to main content

hydro_lang/sim/
compiled.rs

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