use std::cmp::Ordering;
use std::sync::atomic::{AtomicU64, Ordering as AtomicOrdering};
use std::time::{SystemTime, UNIX_EPOCH};
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct Timestamp(i64);
impl Timestamp {
pub fn from_micros(us: i64) -> Self {
Self(us)
}
pub fn from_secs(secs: i64) -> Self {
Self(secs * 1_000_000)
}
pub fn as_micros(&self) -> i64 {
self.0
}
pub fn as_secs(&self) -> i64 {
self.0 / 1_000_000
}
pub fn now() -> Self {
let d = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("system clock before UNIX epoch");
Self(d.as_micros() as i64)
}
}
pub trait Clock: 'static {
fn now(&self) -> Timestamp;
}
pub struct SystemClock;
impl Clock for SystemClock {
fn now(&self) -> Timestamp {
Timestamp::now()
}
}
pub struct TestClock {
now: Timestamp,
}
impl TestClock {
pub fn new(now: Timestamp) -> Self {
Self { now }
}
pub fn set(&mut self, now: Timestamp) {
self.now = now;
}
pub fn advance_secs(&mut self, secs: i64) {
self.now = Timestamp::from_micros(self.now.as_micros() + secs * 1_000_000);
}
pub fn advance_micros(&mut self, us: i64) {
self.now = Timestamp::from_micros(self.now.as_micros() + us);
}
}
impl Clock for TestClock {
fn now(&self) -> Timestamp {
self.now
}
}
static NEXT_DEADLINE_ID: AtomicU64 = AtomicU64::new(1);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct DeadlineId(u64);
impl DeadlineId {
fn next() -> Self {
Self(NEXT_DEADLINE_ID.fetch_add(1, AtomicOrdering::Relaxed))
}
}
struct DeadlineEntry<T> {
id: DeadlineId,
deadline: Timestamp,
event: T,
}
pub struct Deadlines<T: 'static> {
entries: Vec<DeadlineEntry<T>>,
}
impl<T: 'static> Deadlines<T> {
pub fn new() -> Self {
Self {
entries: Vec::new(),
}
}
pub fn schedule(&mut self, deadline: Timestamp, event: T) -> DeadlineId {
let id = DeadlineId::next();
let pos = self
.entries
.binary_search_by(|e| e.deadline.cmp(&deadline).then(Ordering::Less))
.unwrap_or_else(|i| i);
self.entries.insert(
pos,
DeadlineEntry {
id,
deadline,
event,
},
);
id
}
pub fn cancel(&mut self, id: DeadlineId) -> bool {
if let Some(pos) = self.entries.iter().position(|e| e.id == id) {
self.entries.remove(pos);
true
} else {
false
}
}
pub fn drain_overdue(&mut self, now: Timestamp) -> Vec<T> {
let split = self.entries.partition_point(|e| e.deadline <= now);
if split == 0 {
return Vec::new();
}
self.entries.drain(..split).map(|e| e.event).collect()
}
pub fn len(&self) -> usize {
self.entries.len()
}
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
pub fn next_deadline(&self) -> Option<Timestamp> {
self.entries.first().map(|e| e.deadline)
}
}
impl<T: 'static> Default for Deadlines<T> {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn timestamp_from_secs_roundtrip() {
let ts = Timestamp::from_secs(1000);
assert_eq!(ts.as_secs(), 1000);
assert_eq!(ts.as_micros(), 1_000_000_000);
}
#[test]
fn timestamp_from_micros() {
let ts = Timestamp::from_micros(123_456_789);
assert_eq!(ts.as_micros(), 123_456_789);
assert_eq!(ts.as_secs(), 123); }
#[test]
fn timestamp_ordering() {
let a = Timestamp::from_secs(100);
let b = Timestamp::from_secs(200);
assert!(a < b);
assert!(b > a);
assert_eq!(a, Timestamp::from_secs(100));
}
#[test]
fn timestamp_now_is_positive() {
let ts = Timestamp::now();
assert!(ts.as_micros() > 0);
}
#[test]
fn system_clock_returns_positive() {
let clock = SystemClock;
assert!(clock.now().as_micros() > 0);
}
#[test]
fn test_clock_manual_control() {
let mut clock = TestClock::new(Timestamp::from_secs(1000));
assert_eq!(clock.now().as_secs(), 1000);
clock.advance_secs(60);
assert_eq!(clock.now().as_secs(), 1060);
clock.set(Timestamp::from_secs(2000));
assert_eq!(clock.now().as_secs(), 2000);
}
#[test]
fn test_clock_advance_micros() {
let mut clock = TestClock::new(Timestamp::from_micros(0));
clock.advance_micros(500_000);
assert_eq!(clock.now().as_micros(), 500_000);
}
#[test]
fn deadline_ids_are_unique() {
let a = DeadlineId::next();
let b = DeadlineId::next();
assert_ne!(a, b);
}
#[derive(Debug, PartialEq)]
struct TestEvent(u32);
#[test]
fn schedule_and_drain() {
let mut deadlines = Deadlines::new();
deadlines.schedule(Timestamp::from_secs(100), TestEvent(1));
deadlines.schedule(Timestamp::from_secs(200), TestEvent(2));
deadlines.schedule(Timestamp::from_secs(300), TestEvent(3));
assert_eq!(deadlines.len(), 3);
let fired = deadlines.drain_overdue(Timestamp::from_secs(150));
assert_eq!(fired, vec![TestEvent(1)]);
assert_eq!(deadlines.len(), 2);
}
#[test]
fn drain_all_overdue_batch() {
let mut deadlines = Deadlines::new();
deadlines.schedule(Timestamp::from_secs(100), TestEvent(1));
deadlines.schedule(Timestamp::from_secs(200), TestEvent(2));
deadlines.schedule(Timestamp::from_secs(300), TestEvent(3));
let fired = deadlines.drain_overdue(Timestamp::from_secs(300));
assert_eq!(fired, vec![TestEvent(1), TestEvent(2), TestEvent(3)]);
assert!(deadlines.is_empty());
}
#[test]
fn drain_none_overdue() {
let mut deadlines = Deadlines::new();
deadlines.schedule(Timestamp::from_secs(200), TestEvent(1));
let fired = deadlines.drain_overdue(Timestamp::from_secs(100));
assert!(fired.is_empty());
assert_eq!(deadlines.len(), 1);
}
#[test]
fn drain_empty_queue() {
let mut deadlines: Deadlines<TestEvent> = Deadlines::new();
let fired = deadlines.drain_overdue(Timestamp::from_secs(100));
assert!(fired.is_empty());
}
#[test]
fn cancel_removes_entry() {
let mut deadlines = Deadlines::new();
let id = deadlines.schedule(Timestamp::from_secs(100), TestEvent(1));
deadlines.schedule(Timestamp::from_secs(200), TestEvent(2));
assert!(deadlines.cancel(id));
assert_eq!(deadlines.len(), 1);
let fired = deadlines.drain_overdue(Timestamp::from_secs(300));
assert_eq!(fired, vec![TestEvent(2)]);
}
#[test]
fn cancel_nonexistent_returns_false() {
let mut deadlines: Deadlines<TestEvent> = Deadlines::new();
let id = DeadlineId::next();
assert!(!deadlines.cancel(id));
}
#[test]
fn cancel_already_cancelled_returns_false() {
let mut deadlines = Deadlines::new();
let id = deadlines.schedule(Timestamp::from_secs(100), TestEvent(1));
assert!(deadlines.cancel(id));
assert!(!deadlines.cancel(id));
}
#[test]
fn sorted_insertion_order() {
let mut deadlines = Deadlines::new();
deadlines.schedule(Timestamp::from_secs(300), TestEvent(3));
deadlines.schedule(Timestamp::from_secs(100), TestEvent(1));
deadlines.schedule(Timestamp::from_secs(200), TestEvent(2));
let fired = deadlines.drain_overdue(Timestamp::from_secs(400));
assert_eq!(fired, vec![TestEvent(1), TestEvent(2), TestEvent(3)]);
}
#[test]
fn same_deadline_fires_all() {
let mut deadlines = Deadlines::new();
deadlines.schedule(Timestamp::from_secs(100), TestEvent(1));
deadlines.schedule(Timestamp::from_secs(100), TestEvent(2));
let fired = deadlines.drain_overdue(Timestamp::from_secs(100));
assert_eq!(fired.len(), 2);
}
#[test]
fn next_deadline_returns_earliest() {
let mut deadlines = Deadlines::new();
assert!(deadlines.next_deadline().is_none());
deadlines.schedule(Timestamp::from_secs(300), TestEvent(3));
deadlines.schedule(Timestamp::from_secs(100), TestEvent(1));
assert_eq!(deadlines.next_deadline(), Some(Timestamp::from_secs(100)));
}
#[test]
fn world_add_deadline_type_is_idempotent() {
let mut world = crate::world::World::new();
world.add_deadline_type::<TestEvent>();
world.add_deadline_type::<TestEvent>(); assert!(world.try_resource::<Deadlines<TestEvent>>().is_some());
}
#[test]
fn world_schedule_and_drain() {
let mut world = crate::world::World::new();
world.add_deadline_type::<TestEvent>();
world.schedule_deadline(Timestamp::from_secs(100), TestEvent(1));
world.schedule_deadline(Timestamp::from_secs(200), TestEvent(2));
world.drain_deadlines::<TestEvent>(Timestamp::from_secs(150));
world.update_events();
let events = world.resource::<crate::event::Events<TestEvent>>();
assert_eq!(events.len(), 1);
let fired: Vec<_> = events.read().collect();
assert_eq!(fired[0].0, 1);
assert_eq!(world.resource::<Deadlines<TestEvent>>().len(), 1);
}
#[test]
fn world_cancel_deadline() {
let mut world = crate::world::World::new();
world.add_deadline_type::<TestEvent>();
let id = world.schedule_deadline(Timestamp::from_secs(100), TestEvent(1));
assert!(world.cancel_deadline::<TestEvent>(id));
assert!(world.resource::<Deadlines<TestEvent>>().is_empty());
}
#[test]
fn world_batch_reconciliation() {
let mut world = crate::world::World::new();
world.add_deadline_type::<TestEvent>();
world.schedule_deadline(Timestamp::from_secs(100), TestEvent(1));
world.schedule_deadline(Timestamp::from_secs(200), TestEvent(2));
world.schedule_deadline(Timestamp::from_secs(300), TestEvent(3));
world.drain_deadlines::<TestEvent>(Timestamp::from_secs(1000));
world.update_events();
let events = world.resource::<crate::event::Events<TestEvent>>();
assert_eq!(events.len(), 3);
let values: Vec<u32> = events.read().map(|e| e.0).collect();
assert_eq!(values, vec![1, 2, 3]);
}
#[test]
fn world_drain_with_test_clock() {
let mut world = crate::world::World::new();
world.add_deadline_type::<TestEvent>();
let mut clock = TestClock::new(Timestamp::from_secs(0));
world.schedule_deadline(Timestamp::from_secs(60), TestEvent(1));
world.schedule_deadline(Timestamp::from_secs(120), TestEvent(2));
world.drain_deadlines::<TestEvent>(clock.now());
world.update_events();
assert!(
world
.resource::<crate::event::Events<TestEvent>>()
.is_empty()
);
clock.advance_secs(60);
world.drain_deadlines::<TestEvent>(clock.now());
world.update_events();
assert_eq!(world.resource::<crate::event::Events<TestEvent>>().len(), 1);
clock.advance_secs(60);
world.drain_deadlines::<TestEvent>(clock.now());
world.update_events();
assert_eq!(world.resource::<crate::event::Events<TestEvent>>().len(), 1);
assert!(world.resource::<Deadlines<TestEvent>>().is_empty());
}
#[test]
fn drain_all_deadlines_with_clock_resource() {
let mut world = crate::world::World::new();
world.add_deadline_type::<TestEvent>();
world
.insert_resource(Box::new(TestClock::new(Timestamp::from_secs(150))) as Box<dyn Clock>);
world.schedule_deadline(Timestamp::from_secs(100), TestEvent(1));
world.schedule_deadline(Timestamp::from_secs(200), TestEvent(2));
world.drain_all_deadlines();
world.update_events();
let events = world.resource::<crate::event::Events<TestEvent>>();
assert_eq!(events.len(), 1);
assert_eq!(events.read().next().unwrap().0, 1);
}
#[test]
fn drain_all_deadlines_no_clock_is_noop() {
let mut world = crate::world::World::new();
world.add_deadline_type::<TestEvent>();
world.schedule_deadline(Timestamp::from_secs(100), TestEvent(1));
world.drain_all_deadlines();
assert_eq!(world.resource::<Deadlines<TestEvent>>().len(), 1);
}
#[test]
fn schedule_run_fires_deadlines_readable_same_tick() {
use crate::schedule::Schedule;
let mut world = crate::world::World::new();
world.add_deadline_type::<TestEvent>();
world
.insert_resource(Box::new(TestClock::new(Timestamp::from_secs(200))) as Box<dyn Clock>);
world.schedule_deadline(Timestamp::from_secs(100), TestEvent(42));
world.insert_resource(0_u32); fn count_fired(
reader: crate::event::EventReader<'_, TestEvent>,
mut counter: crate::system_param::ResMut<'_, u32>,
) {
for _ in reader.read() {
*counter += 1;
}
}
let mut schedule = Schedule::new();
schedule.add_system::<(
crate::event::EventReader<'_, TestEvent>,
crate::system_param::ResMut<'_, u32>,
)>("update", "count_fired", count_fired);
schedule.run(&mut world);
assert_eq!(*world.resource::<u32>(), 1);
assert!(world.resource::<Deadlines<TestEvent>>().is_empty());
}
#[test]
fn schedule_run_multiple_deadline_types() {
use crate::schedule::Schedule;
#[derive(Debug, PartialEq)]
struct OtherEvent(u32);
let mut world = crate::world::World::new();
world.add_deadline_type::<TestEvent>();
world.add_deadline_type::<OtherEvent>();
world
.insert_resource(Box::new(TestClock::new(Timestamp::from_secs(500))) as Box<dyn Clock>);
world.schedule_deadline(Timestamp::from_secs(100), TestEvent(1));
world.schedule_deadline(Timestamp::from_secs(200), OtherEvent(2));
let mut schedule = Schedule::new();
schedule.run(&mut world);
assert!(world.resource::<Deadlines<TestEvent>>().is_empty());
assert!(world.resource::<Deadlines<OtherEvent>>().is_empty());
assert_eq!(world.resource::<crate::event::Events<TestEvent>>().len(), 1);
assert_eq!(
world.resource::<crate::event::Events<OtherEvent>>().len(),
1
);
}
}