Something went wrong. Try again.
A video game where you play as a misaligned AI, deceiving and building power. An experiment in spec-driven development.
Something went wrong. Try again.
7.3 kB · 211 lines
Rust
at commit e957ce7b
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212//! 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<E> { 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<E> { pending: Vec<Scheduled<E>>, /// Monotonic counter giving every insertion a stable tie-break. next_seq: u64,}
impl<E> Default for Schedule<E> { fn default() -> Self { Self { pending: Vec::new(), next_seq: 0, } }}
impl<E> Schedule<E> { 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<E> { // Partition into due / not-due, preserving the rest. let mut due: Vec<Scheduled<E>> = Vec::new(); let mut keep: Vec<Scheduled<E>> = 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<u64> { 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<u32> = 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<u32> = 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<u32> = 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<u32> = 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<u64> = 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<String> = 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<String> = 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<String> = Schedule::new(); fresh.at(30, "b".into()); assert_eq!(fresh.next_tick(), Some(30)); }}