//! Deterministic scheduled-event queue on the sim tick clock (DESIGN.md //! "The flow law"). //! //! The flow systems need to decouple *when* from *what*: a message is read at //! the recipient's next desk block (messages.md), an account flow lands on //! its cadence (economy.md), a captured recording is processed some ticks //! after it lands (intel.md). Rather than every subsystem polling //! `tick.is_multiple_of(..)` by hand, they push a payload to fire at a future //! tick and drain what is due. //! //! Ordering is deterministic: events fire in `(tick, insertion sequence)` //! order, so two events due on the same tick resolve in the order they were //! scheduled — no RNG, no wall clock, stable across replays and save/load //! (constitution: architectural guardrails). Recurrence is a domain concern: //! a recurring flow reschedules itself when it fires (see `reschedule` in the //! tests), keeping the substrate minimal and obviously correct. //! //! Generic over the payload `E` so each domain schedules its own event type. /// One scheduled item: fire `event` at `tick`; `seq` breaks ties in /// insertion order. #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] struct Scheduled { tick: u64, seq: u64, event: E, } /// A deterministic queue of events keyed on the sim tick. #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Schedule { pending: Vec>, /// Monotonic counter giving every insertion a stable tie-break. next_seq: u64, } impl Default for Schedule { fn default() -> Self { Self { pending: Vec::new(), next_seq: 0, } } } impl Schedule { pub fn new() -> Self { Self::default() } /// Schedule `event` to fire at `tick`. If `tick` is already in the past /// relative to the next `due` call, it fires at that next call (never /// dropped). pub fn at(&mut self, tick: u64, event: E) { let seq = self.next_seq; self.next_seq += 1; self.pending.push(Scheduled { tick, seq, event }); } /// Schedule `event` to fire `delay` ticks after `now` (`delay` of 0 fires /// on the next `due(now)` for a later `now`, or immediately if `now` /// hasn't advanced — callers pass the current tick). pub fn after(&mut self, now: u64, delay: u64, event: E) { self.at(now.saturating_add(delay), event); } /// Remove and return every event whose tick is `<= now`, in /// `(tick, seq)` order. Draining is the only way events leave the queue; /// nothing fires by the passage of time alone. pub fn due(&mut self, now: u64) -> Vec { // Partition into due / not-due, preserving the rest. let mut due: Vec> = Vec::new(); let mut keep: Vec> = Vec::with_capacity(self.pending.len()); for s in self.pending.drain(..) { if s.tick <= now { due.push(s); } else { keep.push(s); } } self.pending = keep; due.sort_by(|a, b| a.tick.cmp(&b.tick).then(a.seq.cmp(&b.seq))); due.into_iter().map(|s| s.event).collect() } /// The earliest tick with a pending event, if any — lets a caller sleep /// until the next thing happens instead of polling every tick. pub fn next_tick(&self) -> Option { self.pending.iter().map(|s| s.tick).min() } /// Drop pending events whose payload fails `keep` (cancellation — a /// message recalled, a flow closed). Returns how many were removed. pub fn retain(&mut self, keep: impl Fn(&E) -> bool) -> usize { let before = self.pending.len(); self.pending.retain(|s| keep(&s.event)); before - self.pending.len() } pub fn len(&self) -> usize { self.pending.len() } pub fn is_empty(&self) -> bool { self.pending.is_empty() } } #[cfg(test)] mod tests { use super::*; #[test] fn fires_only_what_is_due() { let mut s: Schedule<&str> = Schedule::new(); s.at(10, "a"); s.at(20, "b"); assert_eq!(s.due(5), Vec::<&str>::new(), "nothing due yet"); assert_eq!(s.due(10), vec!["a"], "tick 10 fires a, not b"); assert_eq!(s.len(), 1); assert_eq!(s.due(100), vec!["b"], "b fires once its tick passes"); assert!(s.is_empty()); } #[test] fn same_tick_fires_in_insertion_order() { let mut s: Schedule = Schedule::new(); s.at(10, 1); s.at(10, 2); s.at(10, 3); assert_eq!(s.due(10), vec![1, 2, 3], "stable (tick, seq) ordering"); } #[test] fn cross_tick_ordering_is_by_tick_then_seq() { let mut s: Schedule = Schedule::new(); s.at(30, 30); s.at(10, 10); s.at(20, 20); s.at(10, 11); // second at tick 10, later seq assert_eq!(s.due(100), vec![10, 11, 20, 30]); } #[test] fn past_ticks_fire_at_next_due_never_dropped() { let mut s: Schedule<&str> = Schedule::new(); s.at(5, "late"); // scheduled in the past relative to now=9 assert_eq!(s.due(9), vec!["late"], "past-due events still fire"); } #[test] fn after_offsets_from_now() { let mut s: Schedule<&str> = Schedule::new(); s.after(100, 50, "reply"); // fires at 150 assert!(s.due(140).is_empty()); assert_eq!(s.due(150), vec!["reply"]); } #[test] fn next_tick_reports_the_soonest() { let mut s: Schedule = Schedule::new(); assert_eq!(s.next_tick(), None); s.at(40, 1); s.at(15, 2); assert_eq!(s.next_tick(), Some(15)); } #[test] fn retain_cancels_matching_events() { let mut s: Schedule = Schedule::new(); s.at(10, 1); s.at(10, 2); s.at(10, 3); let removed = s.retain(|e| *e % 2 == 1); // keep odds assert_eq!(removed, 1); assert_eq!(s.due(10), vec![1, 3]); } #[test] fn recurrence_by_reschedule() { // A recurring flow: fire every 20 ticks by rescheduling on fire. let mut s: Schedule = Schedule::new(); let period = 20; s.at(period, 0); // first fire let mut fires = Vec::new(); for now in [20u64, 40, 60] { for count in s.due(now) { fires.push(now); s.at(now + period, count + 1); // reschedule } } assert_eq!(fires, vec![20, 40, 60]); } #[test] fn serde_roundtrips_pending_and_seq() { let mut s: Schedule = Schedule::new(); s.at(10, "a".into()); s.at(30, "b".into()); let _ = s.due(10); // consume "a", advance internal state let json = serde_json::to_string(&s).unwrap(); let mut back: Schedule = serde_json::from_str(&json).unwrap(); assert_eq!(back.len(), 1); assert_eq!(back.due(30), vec!["b".to_string()]); // The seq counter survives, so post-load insertions still tie-break // deterministically after the restored ones. let mut fresh: Schedule = Schedule::new(); fresh.at(30, "b".into()); assert_eq!(fresh.next_tick(), Some(30)); } }