use super::event_queue::EventSender;
#[cfg(not(alloc_frugal))]
use super::types::Event;
#[cfg(not(alloc_frugal))]
use crate::compat::Condvar;
use crate::compat::{format, lock, Box, HashMap, Instant, MiniToString, Mutex, String, Vec};
use crate::core::ObjectId;
use alloc::sync::Arc;
use core::time::Duration;
#[cfg(not(alloc_frugal))]
use std::thread;
struct TimerEntry {
interval: Duration,
repeating: bool,
next_fire: Instant,
}
#[cfg(not(alloc_frugal))]
fn next_fire_deadline(interval: Duration) -> Result<Instant, String> {
Instant::now()
.checked_add(interval)
.ok_or_else(|| format!("timer interval {interval:?} overflows the platform clock"))
}
#[cfg(alloc_frugal)]
fn next_fire_deadline(interval: Duration) -> Result<Instant, String> {
const CEILING: Duration = Duration::from_secs(10 * 365 * 24 * 3600);
if interval > CEILING {
return Err(format!("timer interval {interval:?} overflows the platform clock"));
}
Ok(Instant::now() + interval)
}
#[derive(Default)]
struct TimerState {
timers: HashMap<(ObjectId, u32), TimerEntry>,
running: bool,
}
#[cfg(not(alloc_frugal))]
struct TimerShared {
state: Mutex<TimerState>,
signal: Condvar,
#[allow(dead_code)]
wake_count: core::sync::atomic::AtomicU64,
}
#[cfg(not(alloc_frugal))]
impl TimerShared {
fn notify(&self) {
self.signal.notify_one();
}
fn nearest_deadline(state: &TimerState) -> Option<Instant> {
state.timers.values().map(|entry| entry.next_fire).min()
}
}
pub struct TimerManager {
#[cfg(not(alloc_frugal))]
state: Arc<TimerShared>,
#[cfg(alloc_frugal)]
state: Arc<Mutex<TimerState>>,
#[cfg(not(alloc_frugal))]
thread_handle: Option<thread::JoinHandle<()>>,
#[cfg(alloc_frugal)]
#[cfg_attr(alloc_frugal, allow(dead_code))]
thread_handle: Option<()>,
#[cfg(alloc_frugal)]
sender: EventSender,
}
impl TimerManager {
#[cfg(not(alloc_frugal))]
pub fn new(sender: EventSender) -> Self {
let shared = Arc::new(TimerShared {
state: Mutex::new(TimerState { timers: HashMap::new(), running: true }),
signal: Condvar::new(),
wake_count: core::sync::atomic::AtomicU64::new(0),
});
let worker_shared = Arc::clone(&shared);
let worker_sender = sender;
let thread_handle = thread::spawn(move || loop {
worker_shared.wake_count.fetch_add(1, core::sync::atomic::Ordering::SeqCst);
let mut due_events: Vec<(ObjectId, u32)> = Vec::new();
{
let mut guard = lock(&worker_shared.state);
if !guard.running {
return;
}
let now = Instant::now();
let keys: Vec<(ObjectId, u32)> = guard.timers.keys().copied().collect();
for key in keys {
if let Some(entry) = guard.timers.get_mut(&key) {
if now >= entry.next_fire {
due_events.push(key);
if entry.repeating {
entry.next_fire = Instant::now() + entry.interval;
}
}
}
}
for &(target, id) in &due_events {
if let Some(entry) = guard.timers.get(&(target, id)) {
if !entry.repeating {
guard.timers.remove(&(target, id));
}
}
}
if due_events.is_empty() {
match TimerShared::nearest_deadline(&guard) {
Some(deadline) => {
let wait_for = deadline.saturating_duration_since(Instant::now());
let _ = worker_shared.signal.wait_timeout(guard, wait_for);
}
None => {
guard = worker_shared
.signal
.wait(guard)
.unwrap_or_else(|poisoned| poisoned.into_inner());
}
}
continue;
}
}
for (target, id) in due_events {
if worker_sender.post(target, Event::timer(id)).is_err() {
lock(&worker_shared.state).timers.remove(&(target, id));
}
}
});
Self { state: shared, thread_handle: Some(thread_handle) }
}
#[cfg(alloc_frugal)]
pub fn new(sender: EventSender) -> Self {
let state = Arc::new(Mutex::new(TimerState { timers: HashMap::new(), running: true }));
Self { state, sender, thread_handle: None }
}
#[cfg(not(alloc_frugal))]
fn lock_timers(&self) -> crate::compat::MutexGuard<'_, TimerState> {
lock(&self.state.state)
}
#[cfg(alloc_frugal)]
fn lock_timers(&self) -> crate::compat::MutexGuard<'_, TimerState> {
lock(&self.state)
}
#[cfg(not(alloc_frugal))]
fn notify_worker(&self) {
self.state.notify();
}
#[cfg(alloc_frugal)]
fn notify_worker(&self) {}
pub fn start_timer(
&self,
target: ObjectId,
id: u32,
interval: Duration,
repeating: bool,
) -> Result<(), String> {
if interval.is_zero() {
return Err("timer interval must be > 0".to_string());
}
let entry = TimerEntry { interval, repeating, next_fire: next_fire_deadline(interval)? };
self.lock_timers().timers.insert((target, id), entry);
self.notify_worker();
Ok(())
}
pub fn stop_timer(&self, target: ObjectId, id: u32) -> bool {
let removed = self.lock_timers().timers.remove(&(target, id)).is_some();
if removed {
self.notify_worker();
}
removed
}
pub fn stop_timers_for_target(&self, target: ObjectId) -> usize {
let mut guard = self.lock_timers();
let before = guard.timers.len();
guard.timers.retain(|(timer_target, _), _| *timer_target != target);
let removed = before.saturating_sub(guard.timers.len());
drop(guard);
if removed > 0 {
self.notify_worker();
}
removed
}
pub fn clear(&self) {
let had_timers = {
let mut guard = self.lock_timers();
let had = !guard.timers.is_empty();
guard.timers.clear();
had
};
if had_timers {
self.notify_worker();
}
}
#[cfg(alloc_frugal)]
pub fn pump(&self) {
let now = Instant::now();
let mut due_events = Vec::new();
{
let mut guard = lock(&self.state);
if !guard.running {
return;
}
let keys: Vec<(ObjectId, u32)> = guard.timers.keys().copied().collect();
for key in keys {
if let Some(entry) = guard.timers.get_mut(&key) {
if now >= entry.next_fire {
due_events.push(key);
if entry.repeating {
entry.next_fire = Instant::now() + entry.interval;
}
}
}
}
for &(target, id) in &due_events {
if let Some(entry) = guard.timers.get(&(target, id)) {
if !entry.repeating {
guard.timers.remove(&(target, id));
}
}
}
}
for (target, id) in due_events {
let _ = self.sender.post(target, crate::event::types::Event::timer(id));
}
}
}
#[cfg(not(alloc_frugal))]
impl Drop for TimerManager {
fn drop(&mut self) {
{
let mut guard = lock(&self.state.state);
guard.running = false;
guard.timers.clear();
}
self.state.signal.notify_all();
if let Some(handle) = self.thread_handle.take() {
if let Err(e) = handle.join() {
log::error!("[timer-manager] Thread join failed: {e:?}");
}
}
}
}
pub struct IdleTask {
pub id: u64,
pub callback: Box<dyn FnMut() + Send>,
pub threshold_frames: u32,
pub frames_since_run: u32,
}
impl IdleTask {
pub fn new<F>(id: u64, threshold_frames: u32, callback: F) -> Self
where
F: FnMut() + Send + 'static,
{
Self { id, callback: Box::new(callback), threshold_frames, frames_since_run: 0 }
}
pub fn tick(&mut self) -> bool {
self.frames_since_run += 1;
if self.frames_since_run >= self.threshold_frames {
self.frames_since_run = 0;
(self.callback)();
true
} else {
false
}
}
}
#[cfg(all(test, not(alloc_frugal), not(target_arch = "wasm32")))]
mod tests {
use super::*;
use crate::event::EventQueue;
#[cfg(not(alloc_frugal))]
use std::thread;
#[cfg(not(alloc_frugal))]
use std::time::Instant;
#[test]
fn one_shot_timer_emits_single_event() {
let queue = EventQueue::new();
let manager = TimerManager::new(queue.sender());
manager
.start_timer(7, 11, Duration::from_millis(20), false)
.expect("one-shot timer should start");
let deadline = Instant::now() + Duration::from_millis(400);
let mut hit_count = 0usize;
while Instant::now() < deadline {
if let Some((target, event, _priority)) = queue.dequeue() {
if target == 7 && matches!(event, Event::Timer { id: 11 }) {
hit_count += 1;
}
}
thread::sleep(Duration::from_millis(2));
}
assert_eq!(hit_count, 1);
}
#[test]
fn repeating_timer_can_be_stopped() {
let queue = EventQueue::new();
let manager = TimerManager::new(queue.sender());
manager
.start_timer(9, 3, Duration::from_millis(15), true)
.expect("repeating timer should start");
let deadline = Instant::now() + Duration::from_millis(300);
let mut hits = 0usize;
while Instant::now() < deadline && hits < 2 {
if let Some((target, event, _priority)) = queue.dequeue() {
if target == 9 && matches!(event, Event::Timer { id: 3 }) {
hits += 1;
}
}
thread::sleep(Duration::from_millis(2));
}
assert!(hits >= 2);
assert!(manager.stop_timer(9, 3));
let post_stop_deadline = Instant::now() + Duration::from_millis(120);
let mut post_stop_hits = 0usize;
while Instant::now() < post_stop_deadline {
if let Some((target, event, _priority)) = queue.dequeue() {
if target == 9 && matches!(event, Event::Timer { id: 3 }) {
post_stop_hits += 1;
}
}
thread::sleep(Duration::from_millis(2));
}
assert_eq!(post_stop_hits, 0);
}
#[test]
fn an_unrepresentable_interval_is_rejected_without_replacing() {
let queue = EventQueue::new();
let manager = TimerManager::new(queue.sender());
manager.start_timer(5, 1, Duration::from_millis(10), false).expect("a valid timer starts");
let result = manager.start_timer(5, 1, Duration::MAX, false);
assert!(result.is_err(), "an unrepresentable time must be an Err, not a panic");
assert!(manager.stop_timer(5, 1), "the previous timer must still be present");
}
#[test]
fn idle_manager_does_not_poll_periodically() {
let queue = EventQueue::new();
let manager = TimerManager::new(queue.sender());
thread::sleep(Duration::from_millis(150));
let iterations = manager.state.wake_count.load(core::sync::atomic::Ordering::SeqCst);
assert!(
iterations < 20,
"an idle TimerManager must not poll periodically; it ran {iterations} worker \
iterations in 150 ms, which is a fixed-interval scan rather than a park"
);
}
#[test]
fn a_short_timer_started_from_a_parked_worker_is_not_delayed() {
let queue = EventQueue::new();
let manager = TimerManager::new(queue.sender());
thread::sleep(Duration::from_millis(30));
let start = Instant::now();
manager
.start_timer(21, 5, Duration::from_millis(20), false)
.expect("a short timer should start");
let deadline = start + Duration::from_millis(500);
let mut fired_at = None;
while Instant::now() < deadline {
if let Some((target, event, _)) = queue.dequeue() {
if target == 21 && matches!(event, Event::Timer { id: 5 }) {
fired_at = Some(Instant::now());
break;
}
}
thread::sleep(Duration::from_millis(2));
}
let fired_at = fired_at.expect("the short timer must fire");
let latency = fired_at.saturating_duration_since(start);
assert!(
latency < Duration::from_millis(250),
"a 20 ms timer must not be delayed beyond tolerance; it took {latency:?} \
to be observed"
);
}
#[test]
fn concurrent_mutations_wake_a_parked_worker_without_deadlock() {
use std::sync::mpsc::channel;
let (done_tx, done_rx) = channel();
let worker = thread::spawn(move || {
let queue = EventQueue::new();
let manager = TimerManager::new(queue.sender());
thread::sleep(Duration::from_millis(20));
manager.start_timer(1, 1, Duration::from_millis(5), false).unwrap();
manager.stop_timers_for_target(1);
manager.start_timer(2, 2, Duration::from_millis(5), true).unwrap();
assert!(manager.stop_timer(2, 2));
manager.clear();
manager.start_timer(3, 3, Duration::from_millis(10), false).unwrap();
let deadline = Instant::now() + Duration::from_millis(500);
let mut fired = false;
while Instant::now() < deadline {
if let Some((target, event, _)) = queue.dequeue() {
if target == 3 && matches!(event, Event::Timer { id: 3 }) {
fired = true;
break;
}
}
thread::sleep(Duration::from_millis(2));
}
assert!(fired, "a timer started after the parked mutations must still fire");
drop(manager);
let _ = done_tx.send(());
});
match done_rx.recv_timeout(Duration::from_secs(5)) {
Ok(()) => {}
Err(_) => panic!("the timer worker deadlocked on a concurrent mutation or stop"),
}
worker.join().expect("the timer worker thread must join cleanly");
}
}
#[cfg(all(test, alloc_frugal))]
mod mini_tests {
use super::*;
use crate::event::EventQueue;
#[test]
fn mini_pump_fires_due_timers() {
let queue = EventQueue::new();
let manager = TimerManager::new(queue.sender());
manager
.start_timer(7, 11, Duration::from_millis(1), false)
.expect("one-shot timer should start");
let mut found = false;
for _ in 0..1000 {
manager.pump();
if let Some((target, event, _priority)) = queue.dequeue() {
if target == 7 && matches!(event, crate::event::types::Event::Timer { id: 11 }) {
found = true;
break;
}
}
std::thread::sleep(Duration::from_millis(1));
}
assert!(found, "mini pump should have posted the due timer event");
}
}