use std::cell::RefCell;
use std::cmp::Ordering as CmpOrdering;
use std::collections::BinaryHeap;
use std::future::Future;
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::mpsc::{Receiver, Sender, TryRecvError};
use std::sync::{Arc, LazyLock, Weak};
use std::task::{Context as TaskContext, Poll, Waker};
use std::thread::JoinHandle;
use std::time::{Duration, Instant};
use parking_lot::Mutex;
use crate::effect::Disposable;
use crate::fiber::{Fiber, UndoMeta};
thread_local! {
static CURRENT_FIBER: RefCell<Option<Weak<Fiber>>> = const { RefCell::new(None) };
}
pub fn with_current_fiber<R>(fiber: &Arc<Fiber>, f: impl FnOnce() -> R) -> R {
let prev = CURRENT_FIBER.with(|slot| slot.borrow_mut().replace(Arc::downgrade(fiber)));
let out = f();
CURRENT_FIBER.with(|slot| *slot.borrow_mut() = prev);
out
}
fn current_fiber_scope() -> Option<Arc<Fiber>> {
CURRENT_FIBER
.with(|slot| slot.borrow().clone())
.and_then(|weak| weak.upgrade())
}
type Job = Box<dyn FnOnce() + Send>;
struct Entry {
deadline: Instant,
job: Job,
seq: u64,
}
impl Ord for Entry {
fn cmp(&self, other: &Self) -> CmpOrdering {
other
.deadline
.cmp(&self.deadline)
.then_with(|| other.seq.cmp(&self.seq))
}
}
impl PartialOrd for Entry {
fn partial_cmp(&self, other: &Self) -> Option<CmpOrdering> {
Some(self.cmp(other))
}
}
impl Eq for Entry {}
impl PartialEq for Entry {
fn eq(&self, other: &Self) -> bool {
self.deadline == other.deadline && self.seq == other.seq
}
}
#[derive(Default)]
struct Wheel {
heap: BinaryHeap<Entry>,
}
static WHEEL: LazyLock<Arc<Mutex<Wheel>>> =
LazyLock::new(|| Arc::new(Mutex::new(Wheel::default())));
static TIMER_THREAD: LazyLock<JoinHandle<()>> = LazyLock::new(|| {
std::thread::Builder::new()
.name("cordis-timer".into())
.spawn(run_timer_thread)
.expect("spawn cordis-timer thread")
});
static NEXT_SEQ: AtomicU64 = AtomicU64::new(1);
fn schedule(deadline: Instant, job: Job) {
let entry =
Entry { deadline, job, seq: NEXT_SEQ.fetch_add(1, Ordering::Relaxed) };
WHEEL.lock().heap.push(entry);
timer_thread().thread().unpark();
}
fn timer_thread() -> &'static JoinHandle<()> {
&TIMER_THREAD
}
fn run_timer_thread() {
loop {
let sleep_for: Option<Duration> = {
let w = WHEEL.lock();
w.heap
.peek()
.map(|top| top.deadline.saturating_duration_since(Instant::now()))
};
if sleep_for != Some(Duration::ZERO) {
match sleep_for {
Some(d) => std::thread::park_timeout(d),
None => std::thread::park(),
}
}
let due: Vec<Job> = {
let mut w = WHEEL.lock();
let now = Instant::now();
let mut fired = Vec::new();
while let Some(top) = w.heap.peek() {
if top.deadline > now {
break;
}
fired.push(w.heap.pop().expect("peeked entry exists").job);
}
fired
};
for job in due {
catch_panic(job);
}
}
}
fn catch_panic(job: Job) {
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(job));
if let Err(err) = result {
tracing::warn!(error = ?err, "cordis timer callback panicked; timer thread survives");
}
}
struct HandleInner {
disposed: AtomicBool,
on_dispose: Mutex<Option<Box<dyn Fn() + Send + Sync>>>,
}
impl Default for HandleInner {
fn default() -> Self {
Self { disposed: AtomicBool::new(false), on_dispose: Mutex::new(None) }
}
}
impl HandleInner {
fn set_on_dispose(&self, hook: Box<dyn Fn() + Send + Sync>) {
*self.on_dispose.lock() = Some(hook);
}
fn trigger_dispose(&self) {
if self.disposed.swap(true, Ordering::AcqRel) {
return;
}
if let Some(hook) = self.on_dispose.lock().take() {
hook();
}
}
}
pub struct EffectHandle {
inner: Arc<HandleInner>,
}
impl Clone for EffectHandle {
fn clone(&self) -> Self {
Self { inner: self.inner.clone() }
}
}
impl EffectHandle {
pub fn is_cancelled(&self) -> bool {
self.inner.disposed.load(Ordering::Acquire)
}
fn cancelled(&self) -> bool {
self.inner.disposed.load(Ordering::Acquire)
}
}
impl Disposable for EffectHandle {
fn dispose(self: Box<Self>) {
self.inner.trigger_dispose();
}
}
const UNDO_LABEL_PREFIX: &str = "timer:";
fn scoped_registration(label: &str, register: impl FnOnce() -> EffectHandle) -> EffectHandle {
let handle = register();
match current_fiber_scope() {
Some(fiber) => {
let meta = UndoMeta::new(format!("{UNDO_LABEL_PREFIX}{label}"));
let dispose_handle = handle.clone();
fiber.push_undo_labeled(
meta,
Box::new(move || Disposable::dispose(Box::new(dispose_handle))),
);
}
None => tracing::warn!(
label = %label,
"cordis timer registered outside a fiber scope; nothing will auto-cancel it"
),
}
handle
}
pub fn timeout<F>(delay: Duration, callback: F) -> EffectHandle
where
F: FnOnce() + Send + 'static,
{
scoped_registration("timeout", || {
let flag = Arc::new(HandleInner::default());
let fire_flag = flag.clone();
schedule(
Instant::now() + delay,
Box::new(move || {
if !fire_flag.disposed.load(Ordering::Acquire) {
callback();
}
}),
);
EffectHandle { inner: flag }
})
}
struct SleepState {
done: AtomicBool,
cancelled: AtomicBool,
waker: Mutex<Option<Waker>>,
}
impl SleepState {
fn new() -> Self {
Self {
done: AtomicBool::new(false),
cancelled: AtomicBool::new(false),
waker: Mutex::new(None),
}
}
fn resolved(&self) -> bool {
self.done.load(Ordering::Acquire) || self.cancelled.load(Ordering::Acquire)
}
}
pub fn sleep(delay: Duration) -> (EffectHandle, impl Future<Output = ()> + Send) {
let state = Arc::new(SleepState::new());
let handle = scoped_registration("sleep", || {
let job_state = state.clone();
schedule(
Instant::now() + delay,
Box::new(move || {
job_state.done.store(true, Ordering::Release);
if let Some(wk) = job_state.waker.lock().take() {
wk.wake();
}
}),
);
let inner = Arc::new(HandleInner::default());
let hook_state = state.clone();
inner.set_on_dispose(Box::new(move || {
hook_state.cancelled.store(true, Ordering::Release);
if let Some(wk) = hook_state.waker.lock().take() {
wk.wake();
}
}));
EffectHandle { inner }
});
let fut_state = state;
(
handle,
async move {
core::future::poll_fn(move |cx| {
if fut_state.resolved() {
return Poll::Ready(());
}
*fut_state.waker.lock() = Some(cx.waker().clone());
if fut_state.resolved() {
return Poll::Ready(());
}
Poll::Pending
})
.await;
},
)
}
pub fn interval<F>(delay: Duration, callback: F) -> EffectHandle
where
F: FnMut() + Send + 'static,
{
scoped_registration("interval", || {
let flag = Arc::new(HandleInner::default());
type CallbackCell = Arc<Mutex<Option<Box<dyn FnMut() + Send>>>>;
let cell: CallbackCell = Arc::new(Mutex::new(Some(Box::new(callback))));
fn rearm(
flag: Arc<HandleInner>,
cell: CallbackCell,
delay: Duration,
) {
schedule(
Instant::now() + delay,
Box::new(move || {
if flag.disposed.load(Ordering::Acquire) {
return; }
{
let mut guard = cell.lock();
if let Some(cb) = guard.as_mut() {
cb();
}
}
rearm(flag, cell, delay);
}),
);
}
rearm(flag.clone(), cell, delay);
EffectHandle { inner: flag }
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct InactiveEffect;
impl std::fmt::Display for InactiveEffect {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "timer effect went inactive (disposed)")
}
}
impl std::error::Error for InactiveEffect {}
pub type TickResult = Result<(), InactiveEffect>;
pub trait Stream {
type Item;
fn poll_next(self: Pin<&mut Self>, cx: &mut TaskContext<'_>) -> Poll<Option<Self::Item>>;
}
enum Tick {
Fire,
FinalErr,
}
pub struct Interval {
rx: Receiver<Tick>,
state: Arc<HandleInner>,
waker: Arc<Mutex<Option<Waker>>>,
final_emitted: AtomicBool,
}
impl Interval {
pub fn is_cancelled(&self) -> bool {
self.state.disposed.load(Ordering::Acquire)
}
}
impl Drop for Interval {
fn drop(&mut self) {
self.state.trigger_dispose();
}
}
impl Stream for Interval {
type Item = TickResult;
fn poll_next(
self: Pin<&mut Self>,
cx: &mut TaskContext<'_>,
) -> Poll<Option<Self::Item>> {
let this = unsafe { self.get_unchecked_mut() };
loop {
match this.rx.try_recv() {
Ok(Tick::Fire) => {
if this.state.disposed.load(Ordering::Acquire) {
continue; }
return Poll::Ready(Some(Ok(())));
}
Ok(Tick::FinalErr) => {
if this.final_emitted.swap(true, Ordering::AcqRel) {
continue;
}
return Poll::Ready(Some(Err(InactiveEffect)));
}
Err(TryRecvError::Empty) | Err(TryRecvError::Disconnected) => {
if this.final_emitted.load(Ordering::Acquire) {
return Poll::Ready(None);
}
if this.state.disposed.load(Ordering::Acquire) {
while matches!(this.rx.try_recv(), Ok(Tick::Fire)) {}
if !this.final_emitted.swap(true, Ordering::AcqRel) {
return Poll::Ready(Some(Err(InactiveEffect)));
}
return Poll::Ready(None);
}
*this.waker.lock() = Some(cx.waker().clone());
if this.state.disposed.load(Ordering::Acquire) {
continue;
}
match this.rx.try_recv() {
Ok(Tick::FinalErr) => {
if !this.final_emitted.swap(true, Ordering::AcqRel) {
return Poll::Ready(Some(Err(InactiveEffect)));
}
}
Ok(Tick::Fire) => return Poll::Ready(Some(Ok(()))),
Err(_) => return Poll::Pending,
}
}
}
}
}
}
pub fn interval_stream(delay: Duration) -> Interval {
let (tx, rx) = std::sync::mpsc::channel::<Tick>();
let flag = Arc::new(HandleInner::default());
let waker_slot: Arc<Mutex<Option<Waker>>> = Arc::new(Mutex::new(None));
scoped_registration("interval_stream", || {
fn rearm(
flag: Arc<HandleInner>,
tx: Sender<Tick>,
waker_slot: Arc<Mutex<Option<Waker>>>,
delay: Duration,
) {
schedule(
Instant::now() + delay,
Box::new(move || {
if flag.disposed.load(Ordering::Acquire) {
let _ = tx.send(Tick::FinalErr);
if let Some(wk) = waker_slot.lock().take() {
wk.wake();
}
return;
}
let _ = tx.send(Tick::Fire);
if let Some(wk) = waker_slot.lock().take() {
wk.wake();
}
rearm(flag, tx, waker_slot, delay);
}),
);
}
let hook_flag = flag.clone();
let hook_tx = tx.clone();
let hook_waker = waker_slot.clone();
flag.set_on_dispose(Box::new(move || {
hook_flag.disposed.store(true, Ordering::Release);
let _ = hook_tx.send(Tick::FinalErr);
if let Some(wk) = hook_waker.lock().take() {
wk.wake();
}
}));
rearm(flag.clone(), tx, waker_slot.clone(), delay);
EffectHandle { inner: flag.clone() }
});
Interval { rx, state: flag, waker: waker_slot, final_emitted: AtomicBool::new(false) }
}
pub struct Scheduled<T> {
tx: Sender<T>,
rx: Receiver<T>,
handle: EffectHandle,
submit: Box<dyn Fn(T) + Send + Sync>,
}
impl<T> Scheduled<T> {
pub fn call(&self, value: T) {
if self.handle.cancelled() {
return;
}
(self.submit)(value);
}
pub fn receive(&mut self) -> Option<T> {
self.rx.try_recv().ok()
}
pub fn receive_timeout(&mut self, timeout: Duration) -> Option<T> {
self.rx.recv_timeout(timeout).ok()
}
pub fn is_cancelled(&self) -> bool {
self.handle.is_cancelled()
}
pub fn cancel(&self) {
self.handle.inner.trigger_dispose();
}
}
impl<T: Send + 'static> Disposable for Scheduled<T> {
fn dispose(self: Box<Self>) {
self.handle.inner.trigger_dispose();
}
}
pub fn debounce<T: Send + 'static>(delay: Duration) -> Scheduled<T> {
let (tx, rx) = std::sync::mpsc::channel::<T>();
let flag = Arc::new(HandleInner::default());
let latest: Arc<Mutex<Option<T>>> = Arc::new(Mutex::new(None));
let generation = Arc::new(AtomicU64::new(0));
let handle = scoped_registration("debounce", || {
let hook_latest = latest.clone();
let hook_gen = generation.clone();
let hook_flag = flag.clone();
flag.set_on_dispose(Box::new(move || {
hook_flag.disposed.store(true, Ordering::Release);
*hook_latest.lock() = None;
hook_gen.fetch_add(1, Ordering::SeqCst);
}));
EffectHandle { inner: flag.clone() }
});
let submit = {
let latest = latest.clone();
let generation = generation.clone();
let flag = flag.clone();
let tx = tx.clone();
Box::new(move |value: T| {
if flag.disposed.load(Ordering::Acquire) {
return;
}
*latest.lock() = Some(value);
let my_gen = generation.fetch_add(1, Ordering::SeqCst) + 1;
let emit_latest = latest.clone();
let emit_gen = generation.clone();
let emit_tx = tx.clone();
let emit_flag = flag.clone();
schedule(
Instant::now() + delay,
Box::new(move || {
if emit_gen.load(Ordering::SeqCst) != my_gen {
return;
}
if emit_flag.disposed.load(Ordering::Acquire) {
return;
}
if let Some(v) = emit_latest.lock().take() {
let _ = emit_tx.send(v);
}
}),
);
}) as Box<dyn Fn(T) + Send + Sync>
};
Scheduled { tx, rx, handle, submit }
}
pub fn throttle<T: Send + 'static>(delay: Duration, no_trailing: bool) -> Scheduled<T> {
let (tx, rx) = std::sync::mpsc::channel::<T>();
let flag = Arc::new(HandleInner::default());
let pending: Arc<Mutex<Option<T>>> = Arc::new(Mutex::new(None));
let window_open = Arc::new(AtomicBool::new(false));
let handle = scoped_registration("throttle", || {
let hook_pending = pending.clone();
let hook_window = window_open.clone();
let hook_flag = flag.clone();
flag.set_on_dispose(Box::new(move || {
hook_flag.disposed.store(true, Ordering::Release);
*hook_pending.lock() = None;
hook_window.store(false, Ordering::Release);
}));
EffectHandle { inner: flag.clone() }
});
let submit = {
let pending = pending.clone();
let window_open = window_open.clone();
let flag = flag.clone();
let tx = tx.clone();
Box::new(move |value: T| {
if flag.disposed.load(Ordering::Acquire) {
return;
}
if window_open.swap(true, Ordering::SeqCst) {
*pending.lock() = Some(value);
return;
}
let _ = tx.send(value);
let close_pending = pending.clone();
let close_window = window_open.clone();
let close_flag = flag.clone();
let close_tx = tx.clone();
schedule(
Instant::now() + delay,
Box::new(move || {
close_window.store(false, Ordering::Release);
if close_flag.disposed.load(Ordering::Acquire) || no_trailing {
return;
}
if let Some(v) = close_pending.lock().take() {
let _ = close_tx.send_for_throttle_trailing(v);
}
}),
);
}) as Box<dyn Fn(T) + Send + Sync>
};
Scheduled { tx, rx, handle, submit }
}
trait SendForThrottleTrailing<T> {
fn send_for_throttle_trailing(&self, value: T) -> Result<(), std::sync::mpsc::SendError<T>>;
}
impl<T> SendForThrottleTrailing<T> for Sender<T> {
fn send_for_throttle_trailing(&self, value: T) -> Result<(), std::sync::mpsc::SendError<T>> {
self.send(value)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn block_on_bounded<F: Future>(fut: F, budget: Duration) -> F::Output {
let started = Instant::now();
let waker = Waker::noop();
let mut cx = TaskContext::from_waker(waker);
let mut fut = Box::pin(fut);
loop {
match fut.as_mut().poll(&mut cx) {
Poll::Ready(v) => return v,
Poll::Pending => {
assert!(
started.elapsed() <= budget,
"future did not resolve within {budget:?}"
);
std::thread::sleep(Duration::from_millis(2));
}
}
}
}
fn poll_stream_once(stream: &mut Interval) -> Option<TickResult> {
let waker = Waker::noop();
let mut cx = TaskContext::from_waker(waker);
match Stream::poll_next(Pin::new(stream), &mut cx) {
Poll::Ready(item) => item,
Poll::Pending => None,
}
}
#[test]
fn timeout_fires_once_and_disposes_with_fiber() {
let fiber = Arc::new(Fiber::new());
let hits = Arc::new(AtomicU64::new(0));
let h = hits.clone();
let handle = with_current_fiber(&fiber, || {
timeout(Duration::from_millis(20), move || {
h.fetch_add(1, Ordering::SeqCst);
})
});
assert!(!handle.is_cancelled());
assert_eq!(hits.load(Ordering::SeqCst), 0, "nothing fired before the deadline");
std::thread::sleep(Duration::from_millis(60));
assert_eq!(hits.load(Ordering::SeqCst), 1, "callback must run exactly once");
block_on_bounded(fiber.dispose(), Duration::from_secs(2)).unwrap();
assert!(handle.is_cancelled(), "fiber disposal must cancel timers");
assert_eq!(hits.load(Ordering::SeqCst), 1, "no extra fire after disposal");
}
#[test]
fn timeout_dispose_before_deadline_prevents_fire() {
let fiber = Arc::new(Fiber::new());
let hits = Arc::new(AtomicU64::new(0));
let h = hits.clone();
let handle = with_current_fiber(&fiber, || {
timeout(Duration::from_millis(60), move || {
h.fetch_add(1, Ordering::SeqCst);
})
});
Disposable::dispose(Box::new(handle));
std::thread::sleep(Duration::from_millis(90));
assert_eq!(hits.load(Ordering::SeqCst), 0, "disposed timeout must never fire");
}
#[test]
fn sleep_resolves_and_dispose_resolves_early() {
let fiber = Arc::new(Fiber::new());
let (handle, fut) = with_current_fiber(&fiber, || sleep(Duration::from_millis(30)));
block_on_bounded(fut, Duration::from_secs(2));
assert!(!handle.is_cancelled());
let fiber2 = Arc::new(Fiber::new());
let (handle2, fut2) = with_current_fiber(&fiber2, || sleep(Duration::from_millis(500)));
let started = Instant::now();
let poller = std::thread::spawn(move || {
block_on_bounded(fut2, Duration::from_secs(2));
started.elapsed()
});
std::thread::sleep(Duration::from_millis(30));
Disposable::dispose(Box::new(handle2));
let elapsed = poller.join().expect("poller thread");
assert!(
elapsed < Duration::from_millis(400),
"disposal must resolve the pending sleep early (took {elapsed:?})"
);
}
#[test]
fn interval_ticks_repeatedly_and_stops_on_dispose() {
let fiber = Arc::new(Fiber::new());
let hits = Arc::new(AtomicU64::new(0));
let h = hits.clone();
let handle = with_current_fiber(&fiber, || {
interval(Duration::from_millis(10), move || {
h.fetch_add(1, Ordering::SeqCst);
})
});
std::thread::sleep(Duration::from_millis(55));
let count = hits.load(Ordering::SeqCst);
assert!(count >= 2, "interval must tick repeatedly (got {count})");
Disposable::dispose(Box::new(handle));
std::thread::sleep(Duration::from_millis(60));
assert_eq!(hits.load(Ordering::SeqCst), count, "ticks must stop after disposal");
}
#[test]
fn interval_stream_final_err_on_dispose() {
let fiber = Arc::new(Fiber::new());
let mut stream =
with_current_fiber(&fiber, || interval_stream(Duration::from_millis(10)));
let mut live_ticks = 0u32;
let deadline = Instant::now() + Duration::from_secs(2);
while live_ticks < 2 {
assert!(Instant::now() < deadline, "timed out collecting live ticks");
if let Some(item) = poll_stream_once(&mut stream) {
assert_eq!(item, Ok(()), "live ticks must be Ok");
live_ticks += 1;
} else {
std::thread::sleep(Duration::from_millis(2));
}
}
block_on_bounded(fiber.dispose(), Duration::from_secs(2)).unwrap();
let final_item =
poll_stream_once(&mut stream).expect("final err item must arrive after disposal");
assert_eq!(final_item, Err(InactiveEffect));
assert!(
poll_stream_once(&mut stream).is_none(),
"stream must terminate after the final error"
);
assert!(stream.is_cancelled());
}
#[test]
fn interval_stream_discards_stale_live_ticks_before_final_err() {
let fiber = Arc::new(Fiber::new());
let mut stream =
with_current_fiber(&fiber, || interval_stream(Duration::from_millis(5)));
std::thread::sleep(Duration::from_millis(18));
block_on_bounded(fiber.dispose(), Duration::from_secs(2)).unwrap();
let mut saw_err = false;
let deadline = Instant::now() + Duration::from_secs(2);
while Instant::now() < deadline {
match poll_stream_once(&mut stream) {
Some(Err(InactiveEffect)) => {
saw_err = true;
break;
}
Some(Ok(())) => {} None => std::thread::sleep(Duration::from_millis(2)),
}
}
assert!(saw_err, "teardown must be observable as the final error");
assert!(poll_stream_once(&mut stream).is_none());
}
#[test]
fn debounce_collapses_bursts() {
let fiber = Arc::new(Fiber::new());
let mut sched =
with_current_fiber(&fiber, || debounce::<u32>(Duration::from_millis(40)));
for i in 0..5 {
sched.call(i);
std::thread::sleep(Duration::from_millis(4));
}
let delivered = sched.receive_timeout(Duration::from_secs(2));
assert_eq!(delivered, Some(4), "debounce must deliver only the trailing value");
let extra = sched.receive_timeout(Duration::from_millis(120));
assert_eq!(extra, None, "one burst collapses into exactly one delivery");
block_on_bounded(fiber.dispose(), Duration::from_secs(2)).unwrap();
assert!(sched.is_cancelled());
sched.call(9);
assert_eq!(sched.receive_timeout(Duration::from_millis(50)), None);
}
#[test]
fn throttle_trailing_edge_respected() {
let fiber = Arc::new(Fiber::new());
let mut sched =
with_current_fiber(&fiber, || throttle::<u32>(Duration::from_millis(50), false));
for i in 0..5 {
sched.call(i);
std::thread::sleep(Duration::from_millis(4));
}
let leading = sched.receive_timeout(Duration::from_secs(1));
assert_eq!(leading, Some(0), "leading edge passes immediately");
let trailing = sched.receive_timeout(Duration::from_secs(2));
assert_eq!(trailing, Some(4), "trailing edge must respect the last value");
let extra = sched.receive_timeout(Duration::from_millis(120));
assert_eq!(extra, None, "exactly leading + trailing per burst");
block_on_bounded(fiber.dispose(), Duration::from_secs(2)).unwrap();
assert!(sched.is_cancelled());
}
#[test]
fn throttle_no_trailing_drops_rest_of_burst() {
let fiber = Arc::new(Fiber::new());
let mut sched =
with_current_fiber(&fiber, || throttle::<u32>(Duration::from_millis(50), true));
for i in 0..4 {
sched.call(i);
std::thread::sleep(Duration::from_millis(4));
}
assert_eq!(sched.receive_timeout(Duration::from_secs(1)), Some(0));
assert_eq!(
sched.receive_timeout(Duration::from_millis(150)),
None,
"no_trailing must drop every value after the leading one"
);
}
#[test]
fn fiber_death_cancels_all_timers() {
let fiber = Arc::new(Fiber::new());
let (t_handle, timeout_hits, interval_hits, i_handle, mut stream) =
with_current_fiber(&fiber, || {
let th = Arc::new(AtomicU64::new(0));
let th2 = th.clone();
let t = timeout(Duration::from_millis(70), move || {
th2.fetch_add(1, Ordering::SeqCst);
});
let ih = Arc::new(AtomicU64::new(0));
let ih2 = ih.clone();
let i = interval(Duration::from_millis(15), move || {
ih2.fetch_add(1, Ordering::SeqCst);
});
let s = interval_stream(Duration::from_millis(12));
(t, th, ih, i, s)
});
let _ = (&t_handle, &i_handle);
std::thread::sleep(Duration::from_millis(50));
let pre_interval_hits = interval_hits.load(Ordering::SeqCst);
assert!(pre_interval_hits >= 1, "interval should tick before fiber death");
assert_eq!(timeout_hits.load(Ordering::SeqCst), 0, "timeout still pending");
block_on_bounded(fiber.dispose(), Duration::from_secs(2)).unwrap();
std::thread::sleep(Duration::from_millis(150));
assert_eq!(
timeout_hits.load(Ordering::SeqCst),
0,
"pending timeout must never fire after fiber death"
);
assert_eq!(
interval_hits.load(Ordering::SeqCst),
pre_interval_hits,
"interval must stop ticking after fiber death"
);
let mut saw_final_err = false;
let deadline = Instant::now() + Duration::from_secs(2);
while Instant::now() < deadline {
match poll_stream_once(&mut stream) {
Some(Err(InactiveEffect)) => {
saw_final_err = true;
break;
}
Some(Ok(())) => {} None => std::thread::sleep(Duration::from_millis(2)),
}
}
assert!(saw_final_err, "disposed interval_stream must yield Err(InactiveEffect)");
assert!(poll_stream_once(&mut stream).is_none(), "then terminate");
}
#[test]
fn orphan_registration_warns_but_disposable() {
let hits = Arc::new(AtomicU64::new(0));
let h = hits.clone();
let handle = timeout(Duration::from_millis(30), move || {
h.fetch_add(1, Ordering::SeqCst);
});
Disposable::dispose(Box::new(handle));
std::thread::sleep(Duration::from_millis(60));
assert_eq!(hits.load(Ordering::SeqCst), 0);
}
#[test]
fn undo_labels_are_recorded_on_the_fiber() {
let fiber = Arc::new(Fiber::new());
with_current_fiber(&fiber, || {
timeout(Duration::from_millis(500), || {});
interval(Duration::from_millis(500), || {});
});
let labels = fiber.pending_undo_labels();
assert!(
labels.iter().any(|l| l.contains("timer:timeout")),
"labels: {labels:?}"
);
assert!(
labels.iter().any(|l| l.contains("timer:interval")),
"labels: {labels:?}"
);
}
}