Skip to main content

henad_compute/cpu/
sim_thread.rs

1//! The CPU runner [`SimThread`]. It steps a [`SimState`] and publishes [`Snapshot`]s for the host to draw.
2//!
3//! A capped run steps a batch of ticks at a target rate in ticks per second (TPS), and an uncapped one steps as fast
4//! as it can. [`crate::runner`] decides whether the loop runs on its own thread or from the host's frame loop.
5
6use crate::fault::FaultSink;
7use henad_core::action::Schedule;
8use henad_core::model::SimState;
9use henad_core::params::ParamValue;
10
11use crate::runner::{Driver, Pace, RUN_TO_PUBLISH_INTERVAL, SharedSlot, SimLoop, SnapshotSlot};
12use crate::snapshot::{CpuLayers, GridSnapshot, PointSnapshot, Snapshot, SnapshotView};
13use std::time::Duration;
14use web_time::Instant;
15
16/// Returns the wall-clock seconds between capped batches.
17fn capped_batch_interval_secs(target_tps: f64, ticks_per_snapshot: u32) -> f64 {
18    let tps = if target_tps.is_finite() && target_tps > 0.0 {
19        target_tps
20    } else {
21        1.0
22    };
23    f64::from(ticks_per_snapshot.max(1)) / tps
24}
25
26/// Callback run after every publish, for an idle UI to come and collect the snapshot.
27///
28/// Without it a publish while the UI is idle, like a single step or the final one after pause,
29/// sits unread until some unrelated input event wakes the event loop. The callback must not block.
30#[cfg(not(all(target_arch = "wasm32", target_feature = "atomics")))]
31pub type WakeFn = std::sync::Arc<dyn Fn() + Send + Sync>;
32
33/// Callback run after every publish, with no `Send` or `Sync` bound.
34///
35/// An `egui::Context` is not `Send` under atomics, and no other thread calls this callback.
36#[cfg(all(target_arch = "wasm32", target_feature = "atomics"))]
37pub type WakeFn = std::sync::Arc<dyn Fn()>;
38
39/// Commands sent from the UI thread to the simulation thread.
40#[derive(Debug)]
41pub enum SimCommand {
42    /// Start stepping. It cancels a pending [`Self::RunTo`].
43    Play,
44    /// Stop stepping, and publish the state the run stopped at.
45    Pause,
46    /// Run one step and publish.
47    StepOnce,
48    /// Set the target rate in ticks per second.
49    SetTargetTps(f64),
50    /// Step as fast as possible when `true`, ignoring the target rate.
51    SetUncapped(bool),
52    /// Set the number of ticks one capped batch steps, at least 1.
53    SetTicksPerSnapshot(u32),
54    /// Set one parameter of the running model.
55    SetParam {
56        /// Position of the parameter in the full parameter list.
57        index: usize,
58        /// New value.
59        value: ParamValue,
60    },
61    /// Run the model's declared action at this index, once.
62    Act(usize),
63    /// Turn the layout on or off, with a time budget per publish in milliseconds.
64    ///
65    /// The layout relaxes on every publish that follows a tick. While paused, it relaxes only if `while_paused` is set.
66    SetLayout {
67        /// Whether the layout runs.
68        on: bool,
69        /// Time in milliseconds that one publish may spend relaxing the layout, capped at
70        /// [`MAX_VIEW_BUDGET_MS`](crate::runner::MAX_VIEW_BUDGET_MS).
71        budget_ms: f32,
72        /// Whether the layout keeps relaxing while paused.
73        while_paused: bool,
74    },
75    /// Replace the actions fired at their ticks, and fire those due at the current tick at once unless it has fired.
76    ///
77    /// Each later action fires once, after the step that reaches its tick. A tick has fired once a step has reached
78    /// it or an earlier schedule has fired it. No tick fires twice.
79    SetSchedule(Schedule),
80    /// Step as fast as possible to this tick, then pause and publish.
81    ///
82    /// `Play`, `Pause` and `StepOnce` cancel it. A tick at or behind the current one pauses at once.
83    RunTo(u64),
84    /// Stop the loop. A runner sends it when dropped.
85    Shutdown,
86}
87
88/// Shortest time between two publishes that no command forced, whatever the rate of the sim.
89const PUBLISH_INTERVAL: Duration = Duration::from_millis(16);
90
91/// Largest number of steps in one uncapped pump, so a bad estimate cannot buy a long stall.
92const MAX_UNCAPPED_STEPS: u32 = 4096;
93
94/// Wall-clock time in milliseconds that one uncapped pump aims to fill.
95///
96/// The frame driver returns the frame to the host between pumps, and one step per pump would pin a fast model to the
97/// refresh rate. Matching the driver's own budget keeps a frame to one pump.
98const UNCAPPED_PUMP_MS: f64 = crate::runner::PUMP_BUDGET_MS;
99
100/// Returns the number of steps that fit [`UNCAPPED_PUMP_MS`], from the measured cost of a step.
101///
102/// `engine_ms` is `None` until a step has been timed, and one step is enough to measure with.
103fn uncapped_steps_for(engine_ms: Option<f64>, ticks_per_snapshot: u32) -> u32 {
104    let Some(engine_ms) = engine_ms else {
105        return 1;
106    };
107    let fits = if engine_ms > 0.0 {
108        (UNCAPPED_PUMP_MS / engine_ms)
109            .floor()
110            .clamp(1.0, f64::from(MAX_UNCAPPED_STEPS)) as u32
111    } else {
112        MAX_UNCAPPED_STEPS
113    };
114    let stride = ticks_per_snapshot.max(1);
115    if fits < stride { fits } else { fits - fits % stride }
116}
117
118/// Steps a `SimState` and publishes snapshots. [`Driver`] decides what drives it.
119struct Loop {
120    state: Box<dyn SimState>,
121    slot: SharedSlot,
122    wake: Option<WakeFn>,
123    running: bool,
124    target_tps: f64,
125    uncapped: bool,
126    ticks_per_snapshot: u32,
127    step_count: u64,
128    tps_timer: Instant,
129    actual_tps: f64,
130    last_publish: Instant,
131    /// Number of snapshots published.
132    serial: u64,
133    /// Whether the state has a layout that is switched on.
134    layout_on: bool,
135    /// Whether the layout keeps relaxing while paused.
136    relax_paused: bool,
137    /// Whether a tick has run since the last publish.
138    ticked: bool,
139    /// Engine time per tick in milliseconds, as an exponential moving average. `None` until the first step has been
140    /// timed.
141    engine_ms: Option<f64>,
142    /// When the next capped batch falls due.
143    next_step_at: Instant,
144    /// Actions fired after the step that reaches their tick.
145    schedule: Schedule,
146    /// Highest tick whose actions have had their turn to fire. `None` until the first step or schedule.
147    fired_through: Option<u64>,
148    /// Tick a pending [`SimCommand::RunTo`] stops at.
149    run_to_target: Option<u64>,
150}
151
152impl SimLoop for Loop {
153    type Command = SimCommand;
154
155    fn handle_command(&mut self, cmd: SimCommand) -> bool {
156        match cmd {
157            SimCommand::Play => {
158                self.run_to_target = None;
159                self.running = true;
160                self.reset_tps_window();
161                self.next_step_at = Instant::now();
162            }
163            SimCommand::Pause => {
164                self.run_to_target = None;
165                self.running = false;
166                // A stopped sim runs at no rate. The GPU runner reports a pause the same way.
167                self.actual_tps = 0.0;
168                // Publish a final snapshot, so the UI shows the state it stopped at.
169                self.force_publish_snapshot();
170            }
171            SimCommand::StepOnce => {
172                self.run_to_target = None;
173                self.timed_step();
174                // One step is not a rate. `update_tps` would divide it by however long the pause
175                // before it lasted and report a fraction of a tick per second.
176                self.reset_tps_window();
177                self.force_publish_snapshot();
178            }
179            SimCommand::SetTargetTps(tps) => {
180                self.target_tps = tps;
181                self.reclamp_deadline();
182            }
183            SimCommand::SetUncapped(v) => {
184                self.uncapped = v;
185            }
186            SimCommand::SetTicksPerSnapshot(v) => {
187                self.ticks_per_snapshot = v.max(1);
188                self.reclamp_deadline();
189            }
190            SimCommand::SetParam { index, value } => {
191                if !self.state.set_param(index, &value) {
192                    log::warn!("Failed to set param index {index} to {value:?}");
193                }
194            }
195            SimCommand::SetLayout {
196                on,
197                budget_ms,
198                while_paused,
199            } => {
200                let budget_ms = budget_ms.min(crate::runner::MAX_VIEW_BUDGET_MS);
201                // Called before the `&&`. Inside it, a switch-off would short-circuit and never reach the state.
202                let accepted = self.state.set_layout(on, budget_ms);
203                self.layout_on = on && accepted;
204                self.relax_paused = while_paused;
205                self.force_publish_snapshot();
206            }
207            SimCommand::Act(index) => {
208                if self.state.act(index) {
209                    // The tick has not moved, so nothing else would publish what the action did.
210                    self.force_publish_snapshot();
211                } else {
212                    log::warn!("Model has no action at index {index}");
213                }
214            }
215            SimCommand::SetSchedule(schedule) => {
216                self.schedule = schedule;
217                if self.fired_through.is_none_or(|through| through < self.state.tick()) {
218                    self.fire_due();
219                }
220                self.force_publish_snapshot();
221            }
222            SimCommand::RunTo(target) => {
223                self.run_to_target = Some(target);
224                self.running = false;
225                self.reset_tps_window();
226            }
227            SimCommand::Shutdown => return true,
228        }
229        false
230    }
231
232    fn pump(&mut self) -> Pace {
233        if let Some(target) = self.run_to_target {
234            return self.advance_to_target(target);
235        }
236        if !self.running {
237            return self.relax_while_paused();
238        }
239        if self.uncapped {
240            for _ in 0..uncapped_steps_for(self.engine_ms, self.ticks_per_snapshot) {
241                self.timed_step();
242            }
243            self.update_tps();
244            self.maybe_publish_snapshot();
245            return Pace::Now;
246        }
247
248        let now = Instant::now();
249        if now < self.next_step_at {
250            return Pace::After(self.next_step_at - now);
251        }
252        // Advance from the previous deadline, so the batch's own execution time doesn't stretch
253        // every period. Resync if the sim is running behind.
254        let interval = self.batch_interval();
255        self.next_step_at += interval;
256        let now = Instant::now();
257        if self.next_step_at + interval < now {
258            self.next_step_at = now + interval;
259        }
260        for _ in 0..self.ticks_per_snapshot {
261            self.timed_step();
262        }
263        self.update_tps();
264        self.maybe_publish_snapshot();
265
266        let now = Instant::now();
267        if now >= self.next_step_at {
268            Pace::Now
269        } else {
270            Pace::After(self.next_step_at - now)
271        }
272    }
273}
274
275impl Loop {
276    /// Keeps publishing while paused, so the layout can keep relaxing if requested. Otherwise returns [`Pace::Idle`].
277    fn relax_while_paused(&mut self) -> Pace {
278        if !(self.layout_on && self.relax_paused) {
279            return Pace::Idle;
280        }
281        let since = Instant::now().duration_since(self.last_publish);
282        if since < PUBLISH_INTERVAL {
283            return Pace::After(PUBLISH_INTERVAL.saturating_sub(since));
284        }
285        self.force_publish_snapshot();
286        Pace::After(PUBLISH_INTERVAL)
287    }
288
289    /// Steps uncapped toward `target`, never past it, and pauses there with a final publish.
290    fn advance_to_target(&mut self, target: u64) -> Pace {
291        let remaining = target.saturating_sub(self.state.tick());
292        let steps = u64::from(uncapped_steps_for(self.engine_ms, 1)).min(remaining);
293        for _ in 0..steps {
294            self.timed_step();
295        }
296        if steps == remaining {
297            self.run_to_target = None;
298            self.actual_tps = 0.0;
299            self.force_publish_snapshot();
300            return self.relax_while_paused();
301        }
302        self.update_tps();
303        if Instant::now().duration_since(self.last_publish) >= RUN_TO_PUBLISH_INTERVAL {
304            self.force_publish_snapshot();
305        }
306        Pace::Now
307    }
308
309    fn batch_interval(&self) -> std::time::Duration {
310        std::time::Duration::from_secs_f64(capped_batch_interval_secs(self.target_tps, self.ticks_per_snapshot))
311    }
312
313    /// Pulls the next deadline in to at most one batch interval from now. It never moves the deadline later.
314    ///
315    /// Re-anchoring it to now would let a slider drag fire a batch per event and outrun the cap.
316    fn reclamp_deadline(&mut self) {
317        let limit = Instant::now() + self.batch_interval();
318        if self.next_step_at > limit {
319            self.next_step_at = limit;
320        }
321    }
322
323    /// Steps once, folds the step's cost into the smoothed engine time, then fires the actions due.
324    ///
325    /// The first sample is taken whole. Easing it in from zero would leave `uncapped_steps_for`
326    /// reading far too fast, and a frame would be spent paying for that.
327    fn timed_step(&mut self) {
328        let t0 = Instant::now();
329        self.state.step();
330        self.step_count += 1;
331        self.ticked = true;
332        let sample = t0.elapsed().as_secs_f64() * 1000.0;
333        // Exponential moving average with a weight of 0.1.
334        self.engine_ms = Some(match self.engine_ms {
335            Some(prev) => prev + 0.1 * (sample - prev),
336            None => sample,
337        });
338        self.fire_due();
339    }
340
341    /// Fires the scheduled actions due at the state's current tick, and records that tick as fired.
342    fn fire_due(&mut self) {
343        self.fired_through = Some(self.state.tick());
344        for refused in self.schedule.run_due(&mut *self.state) {
345            log::warn!("Model refused action '{}' at tick {}", refused.id, refused.tick);
346        }
347    }
348
349    /// Starts the window `update_tps` divides by, and drops the steps it had counted.
350    fn reset_tps_window(&mut self) {
351        self.tps_timer = Instant::now();
352        self.step_count = 0;
353    }
354
355    fn update_tps(&mut self) {
356        let elapsed = self.tps_timer.elapsed().as_secs_f64();
357        if elapsed >= 1.0 {
358            self.actual_tps = self.step_count as f64 / elapsed;
359            self.step_count = 0;
360            self.tps_timer = Instant::now();
361        }
362    }
363
364    fn maybe_publish_snapshot(&mut self) {
365        let now = Instant::now();
366        if now.duration_since(self.last_publish) < PUBLISH_INTERVAL {
367            return;
368        }
369        self.last_publish = now;
370        self.publish_snapshot();
371    }
372
373    fn force_publish_snapshot(&mut self) {
374        self.last_publish = Instant::now();
375        self.publish_snapshot();
376    }
377
378    /// Builds a snapshot and publishes it.
379    ///
380    /// The snapshot is built outside the lock. Otherwise the UI thread would block on `take_snapshot` for the whole
381    /// grid copy.
382    fn publish_snapshot(&mut self) {
383        let spare = crate::runner::claim_spare(&self.slot);
384        let engine_ms = self.engine_ms.unwrap_or(0.0);
385        self.serial += 1;
386        // Actions and setting changes also publish, and must not move the nodes of a paused network.
387        let relax = self.layout_on && (self.ticked || self.relax_paused);
388        self.ticked = false;
389        let snap = build_snapshot(spare, &mut *self.state, self.actual_tps, engine_ms, self.serial, relax);
390        crate::runner::publish(&self.slot, snap);
391        // After the lock, so waking the UI can never make it block on us.
392        if let Some(wake) = &self.wake {
393            wake();
394        }
395    }
396}
397
398/// Handle on a running CPU simulation.
399///
400/// The sim runs on its own thread on native targets. In a browser [`Self::update`] steps it from the host's
401/// frame loop.
402pub struct SimThread {
403    driver: Driver<Loop>,
404    slot: SharedSlot,
405}
406
407impl std::fmt::Debug for SimThread {
408    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
409        f.debug_struct("SimThread").finish_non_exhaustive()
410    }
411}
412
413impl SimThread {
414    /// Starts a paused runner over `state`, capped at `target_tps` ticks per second.
415    ///
416    /// `wake` is `None` only for a headless caller that polls on its own schedule.
417    ///
418    /// A panic out of the loop is stored in `faults`. The GPU runner reads the same sink off its `GpuContext`.
419    pub fn new(mut state: Box<dyn SimState>, target_tps: f64, wake: Option<WakeFn>, faults: FaultSink) -> Self {
420        // The first snapshot gives the UI something to draw before play is pressed.
421        let slot = SnapshotSlot::with_initial(build_snapshot(None, &mut *state, 0.0, 0.0, 0, false));
422        let now = Instant::now();
423        let sim = Loop {
424            state,
425            slot: SharedSlot::clone(&slot),
426            wake: wake.clone(),
427            running: false,
428            target_tps,
429            uncapped: false,
430            ticks_per_snapshot: 1,
431            step_count: 0,
432            tps_timer: now,
433            actual_tps: 0.0,
434            last_publish: now,
435            serial: 0,
436            layout_on: false,
437            relax_paused: false,
438            ticked: false,
439            engine_ms: None,
440            next_step_at: now,
441            schedule: Schedule::default(),
442            fired_through: None,
443            run_to_target: None,
444        };
445
446        let driver = Driver::spawn(sim, move |fault| {
447            faults.set_once(fault);
448            if let Some(wake) = &wake {
449                wake();
450            }
451        });
452
453        Self { driver, slot }
454    }
455
456    /// Sends a command to the loop.
457    pub fn send(&mut self, cmd: SimCommand) {
458        self.driver.send(cmd);
459    }
460
461    /// Takes the latest snapshot, or `None` when nothing new has been published since the last take.
462    pub fn take_snapshot(&mut self) -> Option<Snapshot> {
463        crate::runner::take_snapshot(&self.slot)
464    }
465
466    /// Stores a taken snapshot as the spare for the next publish to refill.
467    ///
468    /// This is only an optimisation. Dropping the snapshot instead only means the next publish allocates.
469    pub fn recycle(&mut self, snap: Snapshot) {
470        crate::runner::recycle(&self.slot, snap);
471    }
472
473    /// Sends [`SimCommand::Play`].
474    pub fn play(&mut self) {
475        self.send(SimCommand::Play);
476    }
477
478    /// Sends [`SimCommand::Pause`].
479    pub fn pause(&mut self) {
480        self.send(SimCommand::Pause);
481    }
482
483    /// Sends [`SimCommand::StepOnce`].
484    pub fn step_once(&mut self) {
485        self.send(SimCommand::StepOnce);
486    }
487
488    /// Sends [`SimCommand::SetSchedule`].
489    pub fn set_schedule(&mut self, schedule: Schedule) {
490        self.send(SimCommand::SetSchedule(schedule));
491    }
492
493    /// Sends [`SimCommand::RunTo`].
494    pub fn run_to(&mut self, tick: u64) {
495        self.send(SimCommand::RunTo(tick));
496    }
497
498    /// Advances the sim when the driver runs inside the frame loop, and does nothing when the driver has its own
499    /// thread.
500    pub fn update(&mut self, dt: f64) {
501        self.driver.update(dt);
502    }
503}
504
505impl Drop for SimThread {
506    fn drop(&mut self) {
507        self.driver.shutdown(SimCommand::Shutdown);
508    }
509}
510
511fn refill<T: Copy>(dst: &mut Vec<T>, src: &[T]) {
512    dst.clear();
513    dst.extend_from_slice(src);
514}
515
516/// Builds a snapshot of `state`, refilling the buffers of `reuse`.
517///
518/// `reuse` comes back from the UI thread through `recycle`, and a publish is then a copy without a fresh
519/// multi-megabyte allocation. Every view is consulted, so a composite model publishes its field and its agents.
520fn build_snapshot(
521    reuse: Option<Snapshot>,
522    state: &mut dyn SimState,
523    actual_tps: f64,
524    engine_ms: f64,
525    serial: u64,
526    relax: bool,
527) -> Snapshot {
528    // The model turns its state into something drawable here rather than every tick.
529    let view_started = Instant::now();
530    state.prepare_view();
531    if relax {
532        state.relax_layout();
533    }
534    let view_ms = view_started.elapsed().as_secs_f64() * 1000.0;
535    // The recycled layers are destructured up front, so both layers can claim buffers without moving `recycled` twice.
536    let recycled = match reuse.map(|s| s.view) {
537        Some(SnapshotView::Cpu(layers)) => layers,
538        _ => CpuLayers::default(),
539    };
540    let mut cells = recycled.grid.map(|g| g.cells).unwrap_or_default();
541    let mut spare_edges = recycled.edges;
542    let (mut pos_x, mut pos_y, mut color) = match recycled.points {
543        Some(p) => (p.pos_x, p.pos_y, p.color),
544        None => (Vec::new(), Vec::new(), Vec::new()),
545    };
546
547    let grid = state.grid_view().map(|gv| {
548        refill(&mut cells, gv.cells);
549        GridSnapshot {
550            width: gv.width,
551            height: gv.height,
552            cells: std::mem::take(&mut cells),
553            palette: gv.palette,
554        }
555    });
556
557    let points = state.point_view().map(|pv| {
558        refill(&mut pos_x, pv.pos_x);
559        refill(&mut pos_y, pv.pos_y);
560        refill(&mut color, pv.color.unwrap_or(&[]));
561        PointSnapshot {
562            pos_x: std::mem::take(&mut pos_x),
563            pos_y: std::mem::take(&mut pos_y),
564            world_w: pv.world_w,
565            world_h: pv.world_h,
566            color: std::mem::take(&mut color),
567            palette: pv.palette,
568        }
569    });
570
571    let edges = state.edge_view().map(|ev| {
572        let mut snap = spare_edges.take().unwrap_or_default();
573        // Copied only if the version or the length changed.
574        if snap.version != ev.version || snap.src.len() != ev.src.len() {
575            refill(&mut snap.src, ev.src);
576            refill(&mut snap.dst, ev.dst);
577            refill(&mut snap.color, ev.color.unwrap_or(&[]));
578            snap.version = ev.version;
579        }
580        snap.palette = ev.palette;
581        snap.directed = ev.directed;
582        snap
583    });
584
585    let view = SnapshotView::Cpu(CpuLayers { grid, points, edges });
586
587    Snapshot {
588        tick: state.tick(),
589        serial,
590        population: state.population(),
591        heap_bytes: state.heap_bytes(),
592        actual_tps,
593        engine_ms,
594        view_ms,
595        view,
596        stats: state.stats(),
597    }
598}
599
600#[cfg(all(test, not(target_arch = "wasm32")))]
601mod pacing_timing_tests {
602    use super::{SimCommand, SimThread};
603    use crate::fault::{FaultSink, STEPPING};
604    use crate::snapshot::Snapshot;
605    use henad_core::action::{Schedule, Scheduled};
606    use henad_core::model::SimState;
607    use henad_core::params::ParamValue;
608    use henad_core::view::StatEntry;
609    use std::sync::atomic::{AtomicU64, Ordering};
610    use std::sync::{Arc, Mutex};
611    use std::time::Duration;
612    use web_time::Instant;
613
614    /// Longest time a test waits for the loop to reach the state it expects.
615    const DEADLINE: Duration = Duration::from_secs(10);
616
617    /// Polls `condition` until it holds, and returns whether it held before [`DEADLINE`].
618    fn wait_until(condition: impl Fn() -> bool) -> bool {
619        let deadline = Instant::now() + DEADLINE;
620        loop {
621            if condition() {
622                return true;
623            }
624            if Instant::now() >= deadline {
625                return false;
626            }
627            std::thread::sleep(Duration::from_millis(10));
628        }
629    }
630
631    struct Counter(Arc<AtomicU64>);
632
633    impl SimState for Counter {
634        fn step(&mut self) {
635            self.0.fetch_add(1, Ordering::Relaxed);
636        }
637        fn tick(&self) -> u64 {
638            self.0.load(Ordering::Relaxed)
639        }
640        fn stats(&self) -> Vec<StatEntry> {
641            Vec::new()
642        }
643        fn set_param(&mut self, _index: usize, _value: &ParamValue) -> bool {
644            false
645        }
646        fn population(&self) -> u64 {
647            0
648        }
649        fn heap_bytes(&self) -> usize {
650            0
651        }
652    }
653
654    /// Counts steps as [`Counter`] does, and records the tick each action runs at.
655    struct ActionRecorder {
656        ticks: Arc<AtomicU64>,
657        fired: Arc<Mutex<Vec<u64>>>,
658    }
659
660    impl SimState for ActionRecorder {
661        fn step(&mut self) {
662            self.ticks.fetch_add(1, Ordering::Relaxed);
663        }
664        fn tick(&self) -> u64 {
665            self.ticks.load(Ordering::Relaxed)
666        }
667        fn stats(&self) -> Vec<StatEntry> {
668            Vec::new()
669        }
670        fn set_param(&mut self, _index: usize, _value: &ParamValue) -> bool {
671            false
672        }
673        fn act(&mut self, _index: usize) -> bool {
674            self.fired.lock().expect("action log").push(self.tick());
675            true
676        }
677        fn population(&self) -> u64 {
678            0
679        }
680        fn heap_bytes(&self) -> usize {
681            0
682        }
683    }
684
685    /// Records every layout switch that the loop passes on to the state.
686    struct LayoutSwitches(Arc<Mutex<Vec<bool>>>);
687
688    impl SimState for LayoutSwitches {
689        fn step(&mut self) {}
690        fn tick(&self) -> u64 {
691            0
692        }
693        fn stats(&self) -> Vec<StatEntry> {
694            Vec::new()
695        }
696        fn set_param(&mut self, _index: usize, _value: &ParamValue) -> bool {
697            false
698        }
699        fn set_layout(&mut self, on: bool, _budget_ms: f32) -> bool {
700            self.0.lock().expect("switch log").push(on);
701            true
702        }
703        fn population(&self) -> u64 {
704            0
705        }
706        fn heap_bytes(&self) -> usize {
707            0
708        }
709    }
710
711    /// Counts its steps as its tick, and how many times the layout is switched and relaxes.
712    struct Relaxes {
713        ticks: u64,
714        switches: Arc<AtomicU64>,
715        relaxes: Arc<AtomicU64>,
716    }
717
718    impl Relaxes {
719        /// Returns a thread over a new state, with its count of layout switches and its count of relaxes.
720        fn spawn() -> (SimThread, Arc<AtomicU64>, Arc<AtomicU64>) {
721            let switches = Arc::new(AtomicU64::new(0));
722            let relaxes = Arc::new(AtomicU64::new(0));
723            let state = Self {
724                ticks: 0,
725                switches: Arc::clone(&switches),
726                relaxes: Arc::clone(&relaxes),
727            };
728            let thread = SimThread::new(Box::new(state), 50.0, None, FaultSink::new());
729            (thread, switches, relaxes)
730        }
731    }
732
733    impl SimState for Relaxes {
734        fn step(&mut self) {
735            self.ticks += 1;
736        }
737        fn tick(&self) -> u64 {
738            self.ticks
739        }
740        fn stats(&self) -> Vec<StatEntry> {
741            Vec::new()
742        }
743        fn set_param(&mut self, _index: usize, _value: &ParamValue) -> bool {
744            false
745        }
746        fn set_layout(&mut self, _on: bool, _budget_ms: f32) -> bool {
747            // Released, so a test that sees the switch also sees every relax before it.
748            self.switches.fetch_add(1, Ordering::Release);
749            true
750        }
751        fn relax_layout(&mut self) {
752            self.relaxes.fetch_add(1, Ordering::Relaxed);
753        }
754        fn population(&self) -> u64 {
755            0
756        }
757        fn heap_bytes(&self) -> usize {
758            0
759        }
760    }
761
762    /// A model author's bug, from the engine's point of view.
763    struct DivideByZero(u64);
764
765    impl SimState for DivideByZero {
766        fn step(&mut self) {
767            let zero: u64 = std::hint::black_box(0);
768            self.0 = 1 / zero;
769        }
770        fn tick(&self) -> u64 {
771            self.0
772        }
773        fn stats(&self) -> Vec<StatEntry> {
774            Vec::new()
775        }
776        fn set_param(&mut self, _index: usize, _value: &ParamValue) -> bool {
777            false
778        }
779        fn population(&self) -> u64 {
780            0
781        }
782        fn heap_bytes(&self) -> usize {
783            0
784        }
785    }
786
787    /// A panicking kernel must not take the thread with it and leave the UI polling a viewport
788    /// that never updates again. The panic still prints. This test is noisy by design.
789    #[test]
790    fn a_panicking_step_lands_in_the_sink_instead_of_killing_the_thread() {
791        let faults = FaultSink::new();
792        let wakes = Arc::new(AtomicU64::new(0));
793        let counter = Arc::clone(&wakes);
794
795        let mut thread = SimThread::new(
796            Box::new(DivideByZero(0)),
797            1000.0,
798            Some(Arc::new(move || {
799                counter.fetch_add(1, Ordering::Relaxed);
800            })),
801            faults.clone(),
802        );
803        thread.play();
804
805        for _ in 0..200 {
806            if faults.is_set() {
807                break;
808            }
809            std::thread::sleep(std::time::Duration::from_millis(10));
810        }
811        let fault = faults.take().expect("the panic should have reached the sink");
812        assert_eq!(fault.during, STEPPING);
813        assert!(fault.to_string().contains("divide by zero"), "{fault}");
814        // Without the wake the UI would sit idle and never come and look.
815        assert!(wakes.load(Ordering::Relaxed) > 0, "the UI was never woken");
816    }
817
818    #[test]
819    fn capped_batching_holds_the_target_rate() {
820        let ticks = Arc::new(AtomicU64::new(0));
821        let mut thread = SimThread::new(Box::new(Counter(Arc::clone(&ticks))), 50.0, None, FaultSink::new());
822        thread.send(SimCommand::SetTicksPerSnapshot(10));
823        thread.play();
824        std::thread::sleep(std::time::Duration::from_millis(1000));
825        thread.pause();
826
827        let n = ticks.load(Ordering::Relaxed);
828        assert!((20..=150).contains(&n), "ran {n} ticks in 1s at 50 TPS");
829    }
830
831    /// Blocks until `wakes` reaches `want`, or gives up.
832    fn wait_for_wakes(wakes: &Arc<AtomicU64>, want: u64) -> u64 {
833        for _ in 0..200 {
834            let seen = wakes.load(Ordering::Relaxed);
835            if seen >= want {
836                return seen;
837            }
838            std::thread::sleep(std::time::Duration::from_millis(10));
839        }
840        wakes.load(Ordering::Relaxed)
841    }
842
843    /// Waits out a window in which a loop that has stopped has to stay still.
844    ///
845    /// Only a check that nothing happens needs it, and each caller has already seen the loop stop.
846    fn settle() {
847        std::thread::sleep(std::time::Duration::from_millis(100));
848    }
849
850    /// A pause reports a rate of zero, as the GPU runner does, and a step after a long pause is not
851    /// divided by the whole pause.
852    #[test]
853    fn a_pause_and_a_step_after_it_report_no_rate() {
854        let ticks = Arc::new(AtomicU64::new(0));
855        let mut thread = SimThread::new(Box::new(Counter(Arc::clone(&ticks))), 1000.0, None, FaultSink::new());
856
857        thread.play();
858        // A TPS window is a second wide, and nothing is reported until one closes.
859        assert!(
860            snapshot_where(&mut thread, |snap| snap.actual_tps > 0.0).is_some(),
861            "the window never closed, so the test proves nothing"
862        );
863
864        thread.pause();
865        // Every snapshot of the running sim since the window closed reports a rate. The first snapshot without a rate
866        // is the one the pause published.
867        let paused = snapshot_where(&mut thread, |snap| snap.actual_tps == 0.0).expect("a paused sim reported a rate");
868
869        // A single step must not be divided by this stale window.
870        std::thread::sleep(std::time::Duration::from_millis(1100));
871        thread.step_once();
872        let stepped = snapshot_at(&mut thread, paused.tick + 1).expect("a step publishes");
873        assert_eq!(stepped.actual_tps, 0.0, "one step over a long pause reported a rate");
874    }
875
876    /// Switching the layout off must reach the state.
877    /// Otherwise the next publish would still run the layout, and a paused network would keep moving whenever an
878    /// action or a step publishes.
879    #[test]
880    fn switching_the_layout_off_reaches_the_state() {
881        let switches = Arc::new(Mutex::new(Vec::new()));
882        let mut thread = SimThread::new(
883            Box::new(LayoutSwitches(Arc::clone(&switches))),
884            50.0,
885            None,
886            FaultSink::new(),
887        );
888        thread.send(SimCommand::SetLayout {
889            on: true,
890            budget_ms: 1.0,
891            while_paused: false,
892        });
893        thread.send(SimCommand::SetLayout {
894            on: false,
895            budget_ms: 1.0,
896            while_paused: false,
897        });
898        let log = || switches.lock().expect("switch log").clone();
899        assert!(wait_until(|| log().len() >= 2), "only {:?} reached the state", log());
900        assert_eq!(log(), [true, false]);
901    }
902
903    /// While paused, the layout moves on a step, and otherwise only when asked to relax while paused.
904    #[test]
905    fn a_paused_layout_relaxes_on_a_step_or_when_asked() {
906        let (mut thread, switches, relaxes) = Relaxes::spawn();
907        let count = || relaxes.load(Ordering::Relaxed);
908        let layout = |while_paused| SimCommand::SetLayout {
909            on: true,
910            budget_ms: 1.0,
911            while_paused,
912        };
913        assert!(
914            thread.take_snapshot().is_some(),
915            "a new thread publishes its first state"
916        );
917
918        thread.send(layout(false));
919        assert!(
920            snapshot_where(&mut thread, |_| true).is_some(),
921            "a layout switch publishes"
922        );
923        assert_eq!(count(), 0, "a paused publish relaxed");
924
925        thread.step_once();
926        assert!(snapshot_at(&mut thread, 1).is_some(), "a step publishes");
927        assert_eq!(count(), 1, "a step relaxes once and no more");
928
929        thread.send(layout(true));
930        assert!(
931            wait_until(|| count() >= 3),
932            "relaxing while paused stopped at {}",
933            count()
934        );
935
936        thread.send(layout(false));
937        // The switch reaches the state while the loop handles the command, and no relax runs after it.
938        assert!(
939            wait_until(|| switches.load(Ordering::Acquire) == 3),
940            "the switch never reached the state"
941        );
942        let stopped = count();
943        settle();
944        assert_eq!(count(), stopped, "the layout kept relaxing once told to stop");
945    }
946
947    /// A run to a tick ends paused, and a paused layout relaxes after it as it does after Pause.
948    #[test]
949    fn a_paused_layout_keeps_relaxing_after_a_run_to() {
950        let (mut thread, _, relaxes) = Relaxes::spawn();
951        let count = || relaxes.load(Ordering::Relaxed);
952
953        thread.send(SimCommand::SetLayout {
954            on: true,
955            budget_ms: 1.0,
956            while_paused: true,
957        });
958        thread.run_to(5);
959        // The target's snapshot is published after the run has stopped, so every relax from here on is a paused one.
960        assert!(
961            snapshot_at(&mut thread, 5).is_some(),
962            "the run never published its target"
963        );
964        let reached = count();
965        assert!(
966            wait_until(|| count() > reached),
967            "the layout stopped relaxing at the run's target"
968        );
969    }
970
971    /// A snapshot nobody is told about is a snapshot nobody draws. A step has to wake the UI, or the
972    /// viewport refreshes only once the mouse moves.
973    #[test]
974    fn a_publish_while_paused_wakes_the_ui() {
975        let ticks = Arc::new(AtomicU64::new(0));
976        let wakes = Arc::new(AtomicU64::new(0));
977        let counter = Arc::clone(&wakes);
978
979        let mut thread = SimThread::new(
980            Box::new(Counter(Arc::clone(&ticks))),
981            50.0,
982            Some(Arc::new(move || {
983                counter.fetch_add(1, Ordering::Relaxed);
984            })),
985            FaultSink::new(),
986        );
987
988        thread.step_once();
989        assert!(
990            wait_for_wakes(&wakes, 1) >= 1,
991            "a single step published without waking the UI"
992        );
993
994        // Pause force-publishes a final snapshot too.
995        let before = wakes.load(Ordering::Relaxed);
996        thread.pause();
997        assert!(
998            wait_for_wakes(&wakes, before + 1) > before,
999            "pausing published a final snapshot without waking the UI"
1000        );
1001    }
1002
1003    /// Takes snapshots until one satisfies `wanted`, or gives up after [`DEADLINE`].
1004    fn snapshot_where(thread: &mut SimThread, wanted: impl Fn(&Snapshot) -> bool) -> Option<Snapshot> {
1005        let deadline = Instant::now() + DEADLINE;
1006        while Instant::now() < deadline {
1007            if let Some(snap) = thread.take_snapshot()
1008                && wanted(&snap)
1009            {
1010                return Some(snap);
1011            }
1012            std::thread::sleep(Duration::from_millis(10));
1013        }
1014        None
1015    }
1016
1017    /// Takes snapshots until one reports `tick`, or gives up after [`DEADLINE`].
1018    fn snapshot_at(thread: &mut SimThread, tick: u64) -> Option<Snapshot> {
1019        snapshot_where(thread, |snap| snap.tick == tick)
1020    }
1021
1022    /// Returns a schedule of one action at each of `ticks`, in the order given.
1023    fn schedule_at(ticks: &[u64]) -> Schedule {
1024        let entries = ticks
1025            .iter()
1026            .map(|&tick| Scheduled {
1027                index: 0,
1028                id: "mark".to_owned(),
1029                tick,
1030            })
1031            .collect();
1032        Schedule::from_entries(entries)
1033    }
1034
1035    /// Returns a thread over an [`ActionRecorder`] state capped at 1000 TPS, its tick and its action log.
1036    fn action_recorder() -> (SimThread, Arc<AtomicU64>, Arc<Mutex<Vec<u64>>>) {
1037        let ticks = Arc::new(AtomicU64::new(0));
1038        let fired = Arc::new(Mutex::new(Vec::new()));
1039        let state = ActionRecorder {
1040            ticks: Arc::clone(&ticks),
1041            fired: Arc::clone(&fired),
1042        };
1043        let thread = SimThread::new(Box::new(state), 1000.0, None, FaultSink::new());
1044        (thread, ticks, fired)
1045    }
1046
1047    fn fired_ticks(fired: &Arc<Mutex<Vec<u64>>>) -> Vec<u64> {
1048        fired.lock().expect("action log").clone()
1049    }
1050
1051    #[test]
1052    fn run_to_stops_at_the_target_and_pauses() {
1053        let ticks = Arc::new(AtomicU64::new(0));
1054        // Capped at 1 TPS, so only an uncapped run reaches the target in time.
1055        let mut thread = SimThread::new(Box::new(Counter(Arc::clone(&ticks))), 1.0, None, FaultSink::new());
1056        thread.run_to(20_000);
1057        let reached = snapshot_at(&mut thread, 20_000).expect("the run never published its target");
1058        assert_eq!(reached.actual_tps, 0.0, "a paused run reported a rate");
1059        settle();
1060        assert_eq!(ticks.load(Ordering::Relaxed), 20_000, "the run went past its target");
1061        assert!(thread.take_snapshot().is_none(), "a paused run kept publishing");
1062
1063        thread.run_to(10);
1064        assert!(
1065            snapshot_at(&mut thread, 20_000).is_some(),
1066            "a run to a tick behind the current one pauses and publishes"
1067        );
1068        settle();
1069        assert_eq!(ticks.load(Ordering::Relaxed), 20_000);
1070    }
1071
1072    #[test]
1073    fn a_schedule_fires_once_per_tick_after_the_step_reaching_it() {
1074        let (mut thread, ticks, fired) = action_recorder();
1075        thread.set_schedule(schedule_at(&[3, 1, 3, 9, 10, 25, 60]));
1076        thread.run_to(10);
1077        assert!(snapshot_at(&mut thread, 10).is_some());
1078        assert_eq!(
1079            fired_ticks(&fired),
1080            [1, 3, 3, 9, 10],
1081            "the step reaching the target fires its actions"
1082        );
1083
1084        thread.step_once();
1085        thread.step_once();
1086        assert!(snapshot_at(&mut thread, 12).is_some());
1087        assert_eq!(
1088            fired_ticks(&fired),
1089            [1, 3, 3, 9, 10],
1090            "the target's actions fired twice"
1091        );
1092
1093        thread.play();
1094        assert!(
1095            wait_until(|| ticks.load(Ordering::Relaxed) >= 30),
1096            "the sim never played"
1097        );
1098        // Ten seconds ahead at 1000 TPS. The run to it stops Play long before the sim gets there.
1099        let end = ticks.load(Ordering::Relaxed) + 10_000;
1100        thread.run_to(end);
1101        assert!(snapshot_at(&mut thread, end).is_some());
1102        assert_eq!(fired_ticks(&fired), [1, 3, 3, 9, 10, 25, 60]);
1103    }
1104
1105    #[test]
1106    fn set_schedule_fires_what_is_due_now() {
1107        let (mut thread, ticks, fired) = action_recorder();
1108        assert!(
1109            thread.take_snapshot().is_some(),
1110            "a new thread publishes its first state"
1111        );
1112        thread.set_schedule(schedule_at(&[0, 1, 0]));
1113        assert!(snapshot_at(&mut thread, 0).is_some(), "setting a schedule publishes");
1114        assert_eq!(
1115            fired_ticks(&fired),
1116            [0, 0],
1117            "both actions due at tick 0 fire before any step"
1118        );
1119        assert_eq!(ticks.load(Ordering::Relaxed), 0);
1120
1121        thread.step_once();
1122        assert!(snapshot_at(&mut thread, 1).is_some());
1123        assert_eq!(fired_ticks(&fired), [0, 0, 1]);
1124    }
1125
1126    /// A second schedule must not fire the actions of the tick the loop sits at. That tick has fired already.
1127    #[test]
1128    fn a_replaced_schedule_never_fires_a_tick_again() {
1129        let (mut thread, _, fired) = action_recorder();
1130        thread.set_schedule(schedule_at(&[0, 4]));
1131        thread.set_schedule(schedule_at(&[0, 4, 9]));
1132        thread.run_to(4);
1133        assert!(snapshot_at(&mut thread, 4).is_some());
1134        assert_eq!(
1135            fired_ticks(&fired),
1136            [0, 4],
1137            "a second schedule at tick 0 fired it again"
1138        );
1139
1140        thread.set_schedule(schedule_at(&[4, 9]));
1141        thread.run_to(9);
1142        assert!(snapshot_at(&mut thread, 9).is_some());
1143        assert_eq!(
1144            fired_ticks(&fired),
1145            [0, 4, 9],
1146            "a schedule at a tick a step reached fired it again"
1147        );
1148
1149        // Tick 12 fired nothing, and has had its turn all the same.
1150        thread.run_to(12);
1151        assert!(snapshot_at(&mut thread, 12).is_some());
1152        thread.set_schedule(schedule_at(&[12, 13]));
1153        thread.step_once();
1154        assert!(snapshot_at(&mut thread, 13).is_some());
1155        assert_eq!(fired_ticks(&fired), [0, 4, 9, 13]);
1156    }
1157
1158    #[test]
1159    fn play_cancels_a_run_to() {
1160        let ticks = Arc::new(AtomicU64::new(0));
1161        let mut thread = SimThread::new(Box::new(Counter(Arc::clone(&ticks))), 20.0, None, FaultSink::new());
1162        thread.run_to(u64::MAX);
1163        // A thousand ticks take 50 s at the cap, so reaching them within the deadline means the run is uncapped.
1164        assert!(
1165            wait_until(|| ticks.load(Ordering::Relaxed) > 1000),
1166            "the run to a tick never ran uncapped"
1167        );
1168
1169        thread.play();
1170        // Once Play resets the TPS window, the next rate measured is the capped one. An uncapped run would report
1171        // thousands of ticks per second.
1172        let playing = snapshot_where(&mut thread, |snap| snap.actual_tps > 0.0 && snap.actual_tps <= 60.0);
1173        thread.pause();
1174        assert!(playing.is_some(), "Play never brought the sim back to its 20 TPS cap");
1175    }
1176}
1177
1178#[cfg(test)]
1179mod tests {
1180    use super::{MAX_UNCAPPED_STEPS, UNCAPPED_PUMP_MS, capped_batch_interval_secs, uncapped_steps_for};
1181
1182    /// Batching must not multiply the tick rate by the batch size.
1183    #[test]
1184    fn batching_does_not_change_effective_tick_rate() {
1185        for &tps in &[1.0, 30.0, 250.0, 1000.0] {
1186            for &batch in &[1, 2, 10, 137, 1000] {
1187                let interval = capped_batch_interval_secs(tps, batch);
1188                let effective = f64::from(batch) / interval;
1189                assert!(
1190                    (effective - tps).abs() < 1e-9,
1191                    "tps {tps}, batch {batch}: effective {effective}"
1192                );
1193            }
1194        }
1195    }
1196
1197    #[test]
1198    fn interval_is_batch_size_over_tps() {
1199        assert!((capped_batch_interval_secs(30.0, 10) - 1.0 / 3.0).abs() < 1e-12);
1200        assert!((capped_batch_interval_secs(60.0, 1) - 1.0 / 60.0).abs() < 1e-12);
1201    }
1202
1203    /// Guards against a `Duration::from_secs_f64` panic on a degenerate target rate.
1204    #[test]
1205    fn non_positive_tps_yields_a_finite_interval() {
1206        for &tps in &[0.0, -5.0, f64::NAN, f64::INFINITY] {
1207            let secs = capped_batch_interval_secs(tps, 4);
1208            assert!(secs.is_finite() && secs > 0.0, "tps {tps} gave {secs}");
1209            assert!(std::time::Duration::from_secs_f64(secs) > std::time::Duration::ZERO);
1210        }
1211    }
1212
1213    #[test]
1214    fn zero_ticks_per_snapshot_is_treated_as_one() {
1215        assert!((capped_batch_interval_secs(50.0, 0) - capped_batch_interval_secs(50.0, 1)).abs() < 1e-12);
1216    }
1217
1218    /// The frame is returned to the host between pumps, so a pump has to be worth a frame's work.
1219    #[test]
1220    fn an_uncapped_pump_fills_the_budget() {
1221        // A step costing a tenth of the budget earns ten steps.
1222        assert_eq!(uncapped_steps_for(Some(UNCAPPED_PUMP_MS / 10.0), 1), 10);
1223        // A step costing more than the budget still earns one step.
1224        assert_eq!(uncapped_steps_for(Some(UNCAPPED_PUMP_MS * 5.0), 1), 1);
1225    }
1226
1227    /// A publish lands on a stride boundary, so what fits is rounded down to whole strides.
1228    #[test]
1229    fn an_uncapped_pump_runs_whole_snapshot_strides() {
1230        assert_eq!(uncapped_steps_for(Some(UNCAPPED_PUMP_MS / 10.0), 5), 10);
1231        assert_eq!(uncapped_steps_for(Some(UNCAPPED_PUMP_MS / 12.0), 5), 10);
1232    }
1233
1234    /// The floor is one step. A whole stride would make a slow model with a large stride run a whole stride
1235    /// anyway, however long that took.
1236    #[test]
1237    fn a_stride_too_slow_for_the_budget_is_not_run_whole() {
1238        // Three steps fit, where a stride is a hundred.
1239        assert_eq!(uncapped_steps_for(Some(UNCAPPED_PUMP_MS / 3.0), 100), 3);
1240        // Not even one fits.
1241        assert_eq!(uncapped_steps_for(Some(UNCAPPED_PUMP_MS * 2.0), 100), 1);
1242    }
1243
1244    /// Checks the pump size before any step is timed, and when a step measures as free.
1245    #[test]
1246    fn an_unmeasured_step_is_bounded() {
1247        assert_eq!(uncapped_steps_for(None, 1), 1);
1248        assert_eq!(uncapped_steps_for(None, 100), 1);
1249        assert_eq!(uncapped_steps_for(Some(0.0), 1), MAX_UNCAPPED_STEPS);
1250        assert!(uncapped_steps_for(Some(0.0), 100) <= MAX_UNCAPPED_STEPS);
1251    }
1252}
1253
1254#[cfg(test)]
1255mod snapshot_tests {
1256    use super::build_snapshot;
1257    use crate::snapshot::SnapshotView;
1258    use henad_core::model::SimState;
1259    use henad_core::params::ParamValue;
1260    use henad_core::view::{EdgeView, GridView, PointView, StatEntry};
1261
1262    const PALETTE: &[[u8; 4]] = &[[1, 2, 3, 4], [5, 6, 7, 8]];
1263
1264    /// A model with nodes and edges.
1265    struct Graph {
1266        pos: Vec<f32>,
1267        src: Vec<u32>,
1268        dst: Vec<u32>,
1269        color: Vec<u8>,
1270        version: u64,
1271    }
1272
1273    impl Graph {
1274        fn new(edges: usize, version: u64) -> Self {
1275            Self {
1276                pos: vec![0.0; 8],
1277                src: (0..edges as u32).collect(),
1278                dst: (0..edges as u32).map(|i| i + 1).collect(),
1279                color: vec![0; edges],
1280                version,
1281            }
1282        }
1283    }
1284
1285    impl SimState for Graph {
1286        fn step(&mut self) {}
1287        fn tick(&self) -> u64 {
1288            0
1289        }
1290        fn point_view(&self) -> Option<PointView<'_>> {
1291            Some(PointView {
1292                pos_x: &self.pos,
1293                pos_y: &self.pos,
1294                world_w: 1.0,
1295                world_h: 1.0,
1296                color: None,
1297                palette: PALETTE,
1298            })
1299        }
1300        fn edge_view(&self) -> Option<EdgeView<'_>> {
1301            Some(EdgeView {
1302                src: &self.src,
1303                dst: &self.dst,
1304                color: Some(&self.color),
1305                palette: PALETTE,
1306                directed: false,
1307                version: self.version,
1308            })
1309        }
1310        fn stats(&self) -> Vec<StatEntry> {
1311            Vec::new()
1312        }
1313        fn set_param(&mut self, _index: usize, _value: &ParamValue) -> bool {
1314            false
1315        }
1316        fn population(&self) -> u64 {
1317            self.pos.len() as u64
1318        }
1319        fn heap_bytes(&self) -> usize {
1320            0
1321        }
1322    }
1323
1324    /// A model with a field and agents. `build_snapshot` has to publish both.
1325    struct Composite {
1326        cells: Vec<u8>,
1327        pos_x: Vec<f32>,
1328        pos_y: Vec<f32>,
1329        color: Vec<u8>,
1330        with_color: bool,
1331    }
1332
1333    impl Composite {
1334        fn new(agents: usize, with_color: bool) -> Self {
1335            Self {
1336                cells: vec![1; 12],
1337                pos_x: (0..agents).map(|i| i as f32).collect(),
1338                pos_y: (0..agents).map(|i| i as f32 * 2.0).collect(),
1339                color: (0..agents).map(|i| (i % 2) as u8).collect(),
1340                with_color,
1341            }
1342        }
1343    }
1344
1345    impl SimState for Composite {
1346        fn step(&mut self) {}
1347        fn tick(&self) -> u64 {
1348            0
1349        }
1350        fn grid_view(&self) -> Option<GridView<'_>> {
1351            Some(GridView {
1352                width: 4,
1353                height: 3,
1354                cells: &self.cells,
1355                palette: PALETTE,
1356            })
1357        }
1358        fn point_view(&self) -> Option<PointView<'_>> {
1359            Some(PointView {
1360                pos_x: &self.pos_x,
1361                pos_y: &self.pos_y,
1362                world_w: 4.0,
1363                world_h: 3.0,
1364                color: self.with_color.then_some(&self.color),
1365                palette: PALETTE,
1366            })
1367        }
1368        fn stats(&self) -> Vec<StatEntry> {
1369            Vec::new()
1370        }
1371        fn set_param(&mut self, _index: usize, _value: &ParamValue) -> bool {
1372            false
1373        }
1374        fn population(&self) -> u64 {
1375            self.pos_x.len() as u64
1376        }
1377        fn heap_bytes(&self) -> usize {
1378            0
1379        }
1380    }
1381
1382    fn layers(view: &SnapshotView) -> &crate::snapshot::CpuLayers {
1383        match view {
1384            SnapshotView::Cpu(l) => l,
1385            SnapshotView::Gpu(_) => panic!("expected a CPU snapshot"),
1386        }
1387    }
1388
1389    /// Publishing reaches `point_view` when there is a grid too. Otherwise a composite model silently
1390    /// drops every agent.
1391    #[test]
1392    fn a_composite_model_publishes_both_layers() {
1393        let mut state = Composite::new(3, true);
1394        let snap = build_snapshot(None, &mut state, 0.0, 0.0, 0, false);
1395        let layers = layers(&snap.view);
1396
1397        let grid = layers.grid.as_ref().expect("field layer was dropped");
1398        assert_eq!((grid.width, grid.height), (4, 3));
1399        assert_eq!(grid.cells.len(), 12);
1400
1401        let points = layers.points.as_ref().expect("agent layer was dropped");
1402        assert_eq!(points.pos_x, vec![0.0, 1.0, 2.0]);
1403        assert_eq!(points.pos_y, vec![0.0, 2.0, 4.0]);
1404        assert_eq!(points.color, vec![0, 1, 0]);
1405    }
1406
1407    /// An absent lane must arrive empty. The renderer treats an empty lane as uniform.
1408    #[test]
1409    fn a_model_without_a_color_lane_publishes_an_empty_one() {
1410        let mut state = Composite::new(2, false);
1411        let snap = build_snapshot(None, &mut state, 0.0, 0.0, 0, false);
1412        let points = layers(&snap.view).points.as_ref().expect("agent layer was dropped");
1413        assert!(points.color.is_empty());
1414        assert_eq!(points.pos_x.len(), 2);
1415    }
1416
1417    /// The colour lane has to recycle alongside the position lanes.
1418    #[test]
1419    fn recycling_reuses_the_color_lane_across_a_length_change() {
1420        let mut big = Composite::new(64, true);
1421        let first = build_snapshot(None, &mut big, 0.0, 0.0, 0, false);
1422        let capacity = layers(&first.view)
1423            .points
1424            .as_ref()
1425            .map(|p| p.color.capacity())
1426            .unwrap_or_default();
1427        assert!(capacity >= 64);
1428
1429        let mut small = Composite::new(5, true);
1430        let second = build_snapshot(Some(first), &mut small, 0.0, 0.0, 0, false);
1431        let points = layers(&second.view).points.as_ref().expect("agent layer was dropped");
1432        assert_eq!(points.color, vec![0, 1, 0, 1, 0]);
1433        assert_eq!(points.color.capacity(), capacity, "the color lane reallocated");
1434        assert_eq!(points.pos_x.len(), 5);
1435    }
1436
1437    #[test]
1438    fn an_unchanged_edge_list_is_handed_back_untouched() {
1439        let mut model = Graph::new(500, 7);
1440        let first = build_snapshot(None, &mut model, 0.0, 0.0, 1, false);
1441        let ptr = layers(&first.view).edges.as_ref().expect("edges").src.as_ptr();
1442
1443        // The version is unchanged, so nothing is copied.
1444        model.src[0] = 999;
1445        let second = build_snapshot(Some(first), &mut model, 0.0, 0.0, 2, false);
1446        let edges = layers(&second.view).edges.as_ref().expect("edges");
1447        assert_eq!(edges.src.as_ptr(), ptr, "the edge list reallocated");
1448        assert_eq!(edges.src[0], 0, "an unchanged version was copied anyway");
1449    }
1450
1451    #[test]
1452    fn a_changed_edge_list_is_refilled_into_the_same_room() {
1453        let mut model = Graph::new(500, 7);
1454        let first = build_snapshot(None, &mut model, 0.0, 0.0, 1, false);
1455        let capacity = layers(&first.view).edges.as_ref().expect("edges").src.capacity();
1456
1457        model.version = 8;
1458        model.src[0] = 999;
1459        let second = build_snapshot(Some(first), &mut model, 0.0, 0.0, 2, false);
1460        let edges = layers(&second.view).edges.as_ref().expect("edges");
1461        assert_eq!(edges.src[0], 999, "the change did not reach the snapshot");
1462        assert_eq!(edges.version, 8);
1463        assert_eq!(edges.src.capacity(), capacity, "the edge list reallocated");
1464    }
1465
1466    #[test]
1467    fn an_edge_list_that_changed_length_is_refilled() {
1468        let mut model = Graph::new(500, 7);
1469        let first = build_snapshot(None, &mut model, 0.0, 0.0, 1, false);
1470
1471        let mut shorter = Graph::new(3, 7);
1472        let second = build_snapshot(Some(first), &mut shorter, 0.0, 0.0, 2, false);
1473        let edges = layers(&second.view).edges.as_ref().expect("edges");
1474        assert_eq!(edges.src.len(), 3, "a shorter list was passed through whole");
1475    }
1476
1477    #[test]
1478    fn a_model_without_edges_publishes_none() {
1479        let mut graph = Graph::new(4, 1);
1480        let first = build_snapshot(None, &mut graph, 0.0, 0.0, 1, false);
1481        assert!(layers(&first.view).edges.is_some());
1482
1483        let mut plain = Composite::new(4, true);
1484        let second = build_snapshot(Some(first), &mut plain, 0.0, 0.0, 2, false);
1485        assert!(
1486            layers(&second.view).edges.is_none(),
1487            "the edge layer outlived its model"
1488        );
1489    }
1490
1491    #[test]
1492    fn the_serial_is_whatever_the_publish_was_given() {
1493        let mut model = Graph::new(2, 1);
1494        assert_eq!(build_snapshot(None, &mut model, 0.0, 0.0, 41, false).serial, 41);
1495        assert_eq!(build_snapshot(None, &mut model, 0.0, 0.0, 42, false).serial, 42);
1496    }
1497}