use std::collections::HashMap;
pub mod error;
pub use error::TimerError;
const MAX_ROUNDS: u64 = i32::MAX as u64;
pub struct TimerWheel<V> {
slots: Vec<Slot<V>>,
mask: usize,
hand: usize,
next_id: u64,
id_to_slot: HashMap<u64, usize>,
}
struct Slot<V> {
entries: Vec<Entry<V>>,
}
struct Entry<V> {
id: u64,
rounds: u32,
value: V,
cancelled: bool,
}
impl<V> TimerWheel<V> {
pub fn new(num_slots: usize) -> Self {
let n = num_slots.max(2).next_power_of_two();
let mut slots = Vec::with_capacity(n);
for _ in 0..n {
slots.push(Slot {
entries: Vec::new(),
});
}
Self {
slots,
mask: n - 1,
hand: 0,
next_id: 1,
id_to_slot: HashMap::new(),
}
}
pub fn num_slots(&self) -> usize {
self.slots.len()
}
pub fn max_delay(&self) -> u64 {
self.slots.len() as u64 * MAX_ROUNDS
}
pub fn pending(&self) -> usize {
self.id_to_slot.len()
}
pub fn is_empty(&self) -> bool {
self.id_to_slot.is_empty()
}
pub fn slot_len(&self, slot: usize) -> usize {
self.slots.get(slot).map_or(0, |s| s.entries.len())
}
pub fn schedule(&mut self, delay_ticks: usize, value: V) -> u64 {
let d = self.clamp_delay(delay_ticks);
let id = self.next_id;
self.next_id += 1;
self.insert(id, d, value);
id
}
pub fn try_schedule(&mut self, delay_ticks: usize, value: V) -> Result<u64, TimerError> {
let max = self.max_delay();
if delay_ticks as u64 > max {
return Err(TimerError::DelayTooLong {
delay: delay_ticks as u64,
max,
});
}
Ok(self.schedule(delay_ticks, value))
}
pub fn cancel(&mut self, id: u64) -> bool {
let Some(slot) = self.id_to_slot.remove(&id) else {
return false;
};
for e in &mut self.slots[slot].entries {
if e.id == id && !e.cancelled {
e.cancelled = true;
return true;
}
}
false
}
pub fn reschedule(&mut self, id: u64, delay_ticks: usize) -> bool {
let Some(slot) = self.id_to_slot.remove(&id) else {
return false;
};
let Some(pos) = self.slots[slot]
.entries
.iter()
.position(|e| e.id == id && !e.cancelled)
else {
return false;
};
let entry = self.slots[slot].entries.swap_remove(pos);
let d = self.clamp_delay(delay_ticks);
self.insert(id, d, entry.value);
true
}
pub fn tick(&mut self) -> Vec<V> {
self.hand = (self.hand + 1) & self.mask;
let slot = self.hand;
let mut fired = Vec::new();
let entries = std::mem::take(&mut self.slots[slot].entries);
let mut survivors = Vec::new();
for mut e in entries {
if e.cancelled {
continue;
}
if e.rounds == 0 {
self.id_to_slot.remove(&e.id);
fired.push(e.value);
} else {
e.rounds -= 1;
survivors.push(e);
}
}
self.slots[slot].entries = survivors;
fired
}
pub fn advance(&mut self, ticks: usize) -> Vec<V> {
let mut fired = Vec::new();
for _ in 0..ticks {
fired.append(&mut self.tick());
}
fired
}
pub fn drain(&mut self) -> Vec<V> {
let mut out = Vec::with_capacity(self.id_to_slot.len());
for slot in &mut self.slots {
for e in std::mem::take(&mut slot.entries) {
if !e.cancelled {
out.push(e.value);
}
}
}
self.id_to_slot.clear();
out
}
pub fn clear(&mut self) {
for slot in &mut self.slots {
slot.entries.clear();
}
self.id_to_slot.clear();
self.hand = 0;
}
fn clamp_delay(&self, delay_ticks: usize) -> usize {
let max = self.max_delay().min(usize::MAX as u64) as usize;
delay_ticks.clamp(1, max)
}
fn insert(&mut self, id: u64, delay: usize, value: V) {
let n = self.slots.len();
let slot = self.hand.wrapping_add(delay) & self.mask;
let rounds = (delay.div_ceil(n) - 1) as u32;
self.slots[slot].entries.push(Entry {
id,
rounds,
value,
cancelled: false,
});
self.id_to_slot.insert(id, slot);
}
}
#[cfg(feature = "harness")]
pub mod recipe;
#[cfg(any(
feature = "hierarchical",
feature = "concurrent",
feature = "deadline-scheduler",
feature = "cron",
feature = "metrics",
))]
pub mod features;
#[cfg(feature = "concurrent")]
pub use features::concurrent::ConcurrentTimerWheel;
#[cfg(feature = "cron")]
pub use features::cron::{CronError, CronSchedule, CronScheduler};
#[cfg(feature = "deadline-scheduler")]
pub use features::deadline_scheduler::{Clock, DeadlineScheduler, MonotonicClock, TestClock};
#[cfg(feature = "hierarchical")]
pub use features::hierarchical::HierarchicalTimerWheel;
#[cfg(feature = "metrics")]
pub use features::metrics::{MeteredTimerWheel, TimerMetrics};
#[cfg(test)]
#[path = "wheel_tests.rs"]
mod wheel_tests;
#[cfg(test)]
#[path = "sample_app_tests.rs"]
mod sample_app_tests;