use std::sync::Arc;
use std::time::Duration;
use tracing::Instrument as _;
use crate::backend::DurableBackendEnum;
use crate::backend::local::now_unix_millis;
#[derive(Debug)]
pub struct DurableTimerService {
backend: Arc<DurableBackendEnum>,
poll_interval: Duration,
}
impl DurableTimerService {
#[must_use]
pub fn new(backend: Arc<DurableBackendEnum>, poll_interval: Duration) -> Self {
Self {
backend,
poll_interval: poll_interval.max(Duration::from_millis(1)),
}
}
#[tracing::instrument(name = "durable.timer.run", skip_all)]
pub async fn run(self) {
let mut tick = tokio::time::interval(self.poll_interval);
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
tick.tick().await;
self.fire_due()
.instrument(tracing::info_span!("durable.timer.run.iter"))
.await;
}
}
#[tracing::instrument(name = "durable.timer.fire_due", skip_all)]
pub(crate) async fn fire_due(&self) {
let now = now_unix_millis();
let due = match self.backend.due_timers(now).await {
Ok(due) => due,
Err(error) => {
tracing::warn!(%error, "durable timer poll failed; will retry next tick");
return;
}
};
for timer in due {
match self.backend.mark_timer_fired(timer).await {
Ok(true) => tracing::debug!(timer_id = %timer.as_uuid(), "durable timer fired"),
Ok(false) => {}
Err(error) => {
tracing::warn!(%error, timer_id = %timer.as_uuid(), "failed to mark timer fired");
}
}
}
}
}