1use 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
16fn 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#[cfg(not(all(target_arch = "wasm32", target_feature = "atomics")))]
31pub type WakeFn = std::sync::Arc<dyn Fn() + Send + Sync>;
32
33#[cfg(all(target_arch = "wasm32", target_feature = "atomics"))]
37pub type WakeFn = std::sync::Arc<dyn Fn()>;
38
39#[derive(Debug)]
41pub enum SimCommand {
42 Play,
44 Pause,
46 StepOnce,
48 SetTargetTps(f64),
50 SetUncapped(bool),
52 SetTicksPerSnapshot(u32),
54 SetParam {
56 index: usize,
58 value: ParamValue,
60 },
61 Act(usize),
63 SetLayout {
67 on: bool,
69 budget_ms: f32,
72 while_paused: bool,
74 },
75 SetSchedule(Schedule),
80 RunTo(u64),
84 Shutdown,
86}
87
88const PUBLISH_INTERVAL: Duration = Duration::from_millis(16);
90
91const MAX_UNCAPPED_STEPS: u32 = 4096;
93
94const UNCAPPED_PUMP_MS: f64 = crate::runner::PUMP_BUDGET_MS;
99
100fn 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
118struct 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 serial: u64,
133 layout_on: bool,
135 relax_paused: bool,
137 ticked: bool,
139 engine_ms: Option<f64>,
142 next_step_at: Instant,
144 schedule: Schedule,
146 fired_through: Option<u64>,
148 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 self.actual_tps = 0.0;
168 self.force_publish_snapshot();
170 }
171 SimCommand::StepOnce => {
172 self.run_to_target = None;
173 self.timed_step();
174 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 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 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 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 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 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 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 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 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 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 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 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 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 if let Some(wake) = &self.wake {
393 wake();
394 }
395 }
396}
397
398pub 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 pub fn new(mut state: Box<dyn SimState>, target_tps: f64, wake: Option<WakeFn>, faults: FaultSink) -> Self {
420 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 pub fn send(&mut self, cmd: SimCommand) {
458 self.driver.send(cmd);
459 }
460
461 pub fn take_snapshot(&mut self) -> Option<Snapshot> {
463 crate::runner::take_snapshot(&self.slot)
464 }
465
466 pub fn recycle(&mut self, snap: Snapshot) {
470 crate::runner::recycle(&self.slot, snap);
471 }
472
473 pub fn play(&mut self) {
475 self.send(SimCommand::Play);
476 }
477
478 pub fn pause(&mut self) {
480 self.send(SimCommand::Pause);
481 }
482
483 pub fn step_once(&mut self) {
485 self.send(SimCommand::StepOnce);
486 }
487
488 pub fn set_schedule(&mut self, schedule: Schedule) {
490 self.send(SimCommand::SetSchedule(schedule));
491 }
492
493 pub fn run_to(&mut self, tick: u64) {
495 self.send(SimCommand::RunTo(tick));
496 }
497
498 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
516fn 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 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 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 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 const DEADLINE: Duration = Duration::from_secs(10);
616
617 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 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 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 struct Relaxes {
713 ticks: u64,
714 switches: Arc<AtomicU64>,
715 relaxes: Arc<AtomicU64>,
716 }
717
718 impl Relaxes {
719 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 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 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 #[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 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 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 fn settle() {
847 std::thread::sleep(std::time::Duration::from_millis(100));
848 }
849
850 #[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 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 let paused = snapshot_where(&mut thread, |snap| snap.actual_tps == 0.0).expect("a paused sim reported a rate");
868
869 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 #[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 #[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 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 #[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 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 #[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 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 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 fn snapshot_at(thread: &mut SimThread, tick: u64) -> Option<Snapshot> {
1019 snapshot_where(thread, |snap| snap.tick == tick)
1020 }
1021
1022 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 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 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 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 #[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 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 assert!(
1165 wait_until(|| ticks.load(Ordering::Relaxed) > 1000),
1166 "the run to a tick never ran uncapped"
1167 );
1168
1169 thread.play();
1170 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 #[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 #[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 #[test]
1220 fn an_uncapped_pump_fills_the_budget() {
1221 assert_eq!(uncapped_steps_for(Some(UNCAPPED_PUMP_MS / 10.0), 1), 10);
1223 assert_eq!(uncapped_steps_for(Some(UNCAPPED_PUMP_MS * 5.0), 1), 1);
1225 }
1226
1227 #[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 #[test]
1237 fn a_stride_too_slow_for_the_budget_is_not_run_whole() {
1238 assert_eq!(uncapped_steps_for(Some(UNCAPPED_PUMP_MS / 3.0), 100), 3);
1240 assert_eq!(uncapped_steps_for(Some(UNCAPPED_PUMP_MS * 2.0), 100), 1);
1242 }
1243
1244 #[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 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 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 #[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 #[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 #[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 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}