use std::collections::HashSet;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Mutex, OnceLock};
use std::time::Duration;
use web_time::Instant;
const NO_TIMEOUT: Duration = Duration::from_secs(60 * 60 * 24 * 365);
static PENDING: AtomicUsize = AtomicUsize::new(0);
static SEQ: AtomicU64 = AtomicU64::new(0);
fn cancelled() -> &'static Mutex<HashSet<Instant>> {
static SET: OnceLock<Mutex<HashSet<Instant>>> = OnceLock::new();
SET.get_or_init(|| Mutex::new(HashSet::new()))
}
#[derive(Debug)]
pub struct CancellableDeadline {
deadline: Instant,
registered: bool,
}
impl CancellableDeadline {
pub fn new(timeout: Option<Duration>) -> Self {
let base = Instant::now()
.checked_add(timeout.unwrap_or(NO_TIMEOUT))
.unwrap_or_else(Instant::now);
let nudge = Duration::from_nanos(SEQ.fetch_add(1, Ordering::Relaxed) % 1_000_000);
Self {
deadline: base.checked_add(nudge).unwrap_or(base),
registered: false,
}
}
pub fn deadline(&self) -> Instant {
self.deadline
}
pub fn cancel(&mut self) {
if self.registered {
return;
}
let mut set = cancelled().lock().unwrap_or_else(|p| p.into_inner());
if set.insert(self.deadline) {
PENDING.fetch_add(1, Ordering::Release);
}
self.registered = true;
}
pub fn is_cancelled(&self) -> bool {
self.registered
}
}
impl Drop for CancellableDeadline {
fn drop(&mut self) {
if !self.registered {
return;
}
let mut set = cancelled().lock().unwrap_or_else(|p| p.into_inner());
if set.remove(&self.deadline) {
PENDING.fetch_sub(1, Ordering::Release);
}
}
}
#[inline]
pub(crate) fn is_cancelled(deadline: Instant) -> bool {
if PENDING.load(Ordering::Acquire) == 0 {
return false;
}
cancelled()
.lock()
.unwrap_or_else(|p| p.into_inner())
.contains(&deadline)
}
thread_local! {
static ACTIVE: std::cell::Cell<Option<Instant>> = const { std::cell::Cell::new(None) };
}
pub(crate) struct DeadlineScope {
prev: Option<Instant>,
}
impl DeadlineScope {
pub(crate) fn enter(deadline: Option<Instant>) -> Self {
let prev = ACTIVE.with(|a| a.get());
ACTIVE.with(|a| a.set(deadline.or(prev)));
Self { prev }
}
}
impl Drop for DeadlineScope {
fn drop(&mut self) {
let prev = self.prev;
ACTIVE.with(|a| a.set(prev));
}
}
pub(crate) fn active_deadline() -> Option<Instant> {
ACTIVE.with(|a| a.get())
}
#[inline]
pub fn deadline_reached(deadline: Instant) -> bool {
Instant::now() >= deadline || is_cancelled(deadline)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn cancel_is_visible_and_forgotten_on_drop() {
let mut a = CancellableDeadline::new(None);
let b = CancellableDeadline::new(None);
assert_ne!(a.deadline(), b.deadline());
assert!(!deadline_reached(a.deadline()));
a.cancel();
assert!(deadline_reached(a.deadline()));
assert!(!deadline_reached(b.deadline()));
let d = a.deadline();
drop(a);
assert!(!is_cancelled(d));
}
#[test]
fn timeout_still_fires() {
let a = CancellableDeadline::new(Some(Duration::ZERO));
std::thread::sleep(Duration::from_millis(1));
assert!(deadline_reached(a.deadline()));
}
}