timer: send_after / cancel_timer message-delivery substrate
Add a cancellable message-delivery timer on the existing timer.rs min-heap,
the substrate the gen_server time idioms (idle/receive timeout, periodic tick,
debounce/backoff) need.
- Reason::Send { fire }: a type-erased delivery thunk. send_after (Pid<A>, via
send_to) and send_after_named (Name<M>, via send) capture dest+msg and
resolve the address on fire, not at arm time, so a dead target / restarted
name is observed when it fires; a failed resolve or closed inbox is dropped
(Erlang erlang:send_after semantics).
- Cancellation via an "armed" set keyed on the entry seq, exposed as an opaque
TimerId. Only Send timers use it; Sleep/WaitTimeout keep their inert-stale
behaviour untouched. pop_due fires a Send only while still armed and removes
it, so cancel returns true iff it landed before the fire (the race signal).
cancel is unscoped: any holder of the id can cancel (e.g. racing two servers
and cancelling the slow path).
- peek_deadline contract documented as "<= true next deadline" so a future
hierarchical timing wheel can back Timers without touching send_after or the
scheduler idle path. call_timeout left on recv_timeout (blast radius).
Tests: Timers-level (fire/cancel/race/clear/ordering) plus scheduler-level
delivery for both Name<M> and Pid<A>, cancel-prevents-delivery, and silent
drop on unresolved name / dead pid. Green on rq-mutex/rq-mpmc/rq-striped.
This commit is contained in:
+4
-3
@@ -65,11 +65,12 @@ pub use registry::{
|
||||
};
|
||||
pub use runtime::{init, Config, Runtime};
|
||||
pub use scheduler::{
|
||||
block_on_io, request_stop, run, self_pid, sleep, spawn, spawn_addr, spawn_under, wait_readable,
|
||||
wait_readable_timeout, wait_writable, wait_writable_timeout, yield_now, FdArm, JoinError,
|
||||
JoinHandle,
|
||||
block_on_io, cancel_timer, request_stop, run, self_pid, send_after, send_after_named, sleep,
|
||||
spawn, spawn_addr, spawn_under, wait_readable, wait_readable_timeout, wait_writable,
|
||||
wait_writable_timeout, yield_now, FdArm, JoinError, JoinHandle,
|
||||
};
|
||||
pub use supervisor::{ChildSpec, OneForOne, Restart, Signal, Strategy};
|
||||
pub use timer::TimerId;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// check!()
|
||||
|
||||
@@ -1232,6 +1232,14 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
// The callback may call unpark_at itself.
|
||||
target.on_timeout(entry.pid, epoch);
|
||||
}
|
||||
// A `send_after` deadline: run the captured delivery thunk.
|
||||
// It resolves the destination through the registry and
|
||||
// sends now (a send can unpark a receiver) — same as any
|
||||
// other in-loop unpark. The timers lock is already
|
||||
// released; lock order Leaf -> Channel is preserved by the
|
||||
// send itself. `pop_due` only returns still-armed Sends, so
|
||||
// a cancelled one never reaches here.
|
||||
crate::timer::Reason::Send { fire } => fire(),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+52
-1
@@ -9,7 +9,7 @@
|
||||
|
||||
use crate::actor::current_pid;
|
||||
use crate::channel::Sender;
|
||||
use crate::pid::Pid;
|
||||
use crate::pid::{Name, Pid};
|
||||
use crate::runtime::{
|
||||
self, RuntimeInner, YieldIntent, RUNTIME,
|
||||
};
|
||||
@@ -388,6 +388,57 @@ pub fn insert_wait_timer(
|
||||
});
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// send_after / cancel_timer — message-delivery timers (Erlang send_after).
|
||||
//
|
||||
// Arm a timer that delivers `msg` to an address after `after`, returning a
|
||||
// `TimerId`. The destination is resolved *on fire*, not at arm time: a
|
||||
// `Pid<A>` that has since died yields `SendError::Dead`, a `Name<M>` resolves
|
||||
// to whoever currently holds it (so a restarted server is reached). Either way
|
||||
// a failed resolve / closed inbox is dropped, matching `erlang:send_after`.
|
||||
// `cancel_timer` prevents an as-yet-unfired delivery; it returns whether the
|
||||
// timer was still armed.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Deliver `msg` to the exact actor named by `dest` (identity-bound, no
|
||||
/// redirect — see [`send_to`](crate::registry::send_to)) after `after`.
|
||||
/// Returns a [`TimerId`](crate::timer::TimerId) for [`cancel_timer`].
|
||||
pub fn send_after<A: crate::pid::Addressable>(
|
||||
after: std::time::Duration,
|
||||
dest: Pid<A>,
|
||||
msg: A::Msg,
|
||||
) -> crate::timer::TimerId {
|
||||
let deadline = crate::timer::deadline_from_now(after);
|
||||
let fire = Box::new(move || {
|
||||
let _ = crate::registry::send_to(dest, msg);
|
||||
});
|
||||
with_runtime(|inner| inner.timers.lock().unwrap().insert_send(deadline, dest.erase(), fire))
|
||||
}
|
||||
|
||||
/// Deliver `msg` to whichever actor holds the name `dest` at fire time
|
||||
/// (re-resolving [`send`](crate::registry::send) semantics) after `after`.
|
||||
/// Returns a [`TimerId`](crate::timer::TimerId) for [`cancel_timer`].
|
||||
pub fn send_after_named<M: Send + 'static>(
|
||||
after: std::time::Duration,
|
||||
dest: Name<M>,
|
||||
msg: M,
|
||||
) -> crate::timer::TimerId {
|
||||
let deadline = crate::timer::deadline_from_now(after);
|
||||
// Informational only (who armed it); not used for delivery.
|
||||
let armed_by = current_pid().unwrap_or(Pid::new(0, 0));
|
||||
let fire = Box::new(move || {
|
||||
let _ = crate::registry::send(dest, msg);
|
||||
});
|
||||
with_runtime(|inner| inner.timers.lock().unwrap().insert_send(deadline, armed_by, fire))
|
||||
}
|
||||
|
||||
/// Cancel a timer armed by [`send_after`] / [`send_after_named`]. Returns
|
||||
/// `true` if it was still pending (delivery now prevented), `false` if it had
|
||||
/// already fired or been cancelled.
|
||||
pub fn cancel_timer(id: crate::timer::TimerId) -> bool {
|
||||
with_runtime(|inner| inner.timers.lock().unwrap().cancel(id))
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// block_on_io / wait_readable / wait_writable / read / write
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
+77
-9
@@ -15,12 +15,14 @@
|
||||
//! `BinaryHeap` is a max-heap; entries are wrapped in `Reverse` to get
|
||||
//! min-heap behaviour.
|
||||
//!
|
||||
//! No cancellation. When a non-timer wakeup happens (e.g. lock granted
|
||||
//! before timeout), the timer entry is left in the heap. It will be popped
|
||||
//! eventually and the dispatch will observe "actor is no longer parked /
|
||||
//! the wait's epoch was consumed" and no-op. Cost is ~32 bytes per stale
|
||||
//! entry plus a few cycles on pop; acceptable given the upper bound is "one
|
||||
//! entry per parked actor".
|
||||
//! Cancellation is selective. A `Sleep` / `WaitTimeout` entry is left in the
|
||||
//! heap on a non-timer wakeup (lock granted before timeout): it is popped
|
||||
//! eventually and no-ops because a stale unpark fails its epoch CAS — cheap
|
||||
//! (~32 bytes per stale entry plus a few cycles on pop), bounded by one entry
|
||||
//! per parked actor. A `Send` entry is different: running its thunk delivers a
|
||||
//! real message, so a stale one is *not* inert. `send_after` therefore carries
|
||||
//! true cancellation via the `armed` set keyed on the entry's `seq`; `pop_due`
|
||||
//! fires a `Send` only while it is still armed, and `cancel` removes the arm.
|
||||
//!
|
||||
//! Stale pids (slot reused since the timer was inserted) are filtered on
|
||||
//! pop by the scheduler — same convention as the run queue.
|
||||
@@ -51,8 +53,28 @@ pub enum Reason {
|
||||
target: Arc<dyn TimerTarget>,
|
||||
epoch: u32,
|
||||
},
|
||||
/// `send_after`: deliver a message to an address at the deadline,
|
||||
/// cancellable. The destination (a `Pid<A>` / `Name<M>`) and the message
|
||||
/// are captured inside `fire`, which resolves the address through the
|
||||
/// registry and sends *when run* — so a target that died or, for a name,
|
||||
/// was restarted is observed at fire time, not arm time. A failed resolve
|
||||
/// or send is dropped (Erlang `erlang:send_after` semantics).
|
||||
///
|
||||
/// Unlike `Sleep` / `WaitTimeout`, a stale `Send` is **not** inert — running
|
||||
/// the thunk delivers a real message — so these are the only timers that
|
||||
/// carry true cancellation (the `armed` set on [`Timers`], keyed by the
|
||||
/// entry's `seq`). `pop_due` fires the thunk only for an entry still armed.
|
||||
Send { fire: Box<dyn FnOnce() + Send> },
|
||||
}
|
||||
|
||||
/// Opaque handle to an armed `send_after` timer, returned by
|
||||
/// [`Timers::insert_send`] and consumed by [`Timers::cancel`]. The inner value
|
||||
/// is the entry's insertion `seq`; callers must treat it as opaque so the
|
||||
/// backing structure can change (e.g. a future hierarchical timing wheel) with
|
||||
/// no API churn.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
|
||||
pub struct TimerId(u64);
|
||||
|
||||
/// Callback the scheduler invokes when a `WaitTimeout` entry pops.
|
||||
///
|
||||
/// Implementors: do not touch `SchedulerState` other than via the public
|
||||
@@ -97,13 +119,21 @@ impl PartialOrd for Entry {
|
||||
pub struct Timers {
|
||||
/// Reverse-wrapped so the smallest deadline is at the top.
|
||||
heap: BinaryHeap<Reverse<Entry>>,
|
||||
/// Monotonic counter for the tiebreaker `seq` field.
|
||||
/// Monotonic counter for the tiebreaker `seq` field (and the `TimerId` of a
|
||||
/// `Send` timer — the two are the same value).
|
||||
next_seq: u64,
|
||||
/// Presence set of *live* `Send` timers, keyed by `seq`. Populated on
|
||||
/// `insert_send`, removed on fire (in `pop_due`) and on `cancel`. A `Send`
|
||||
/// entry fires only while present, so a `cancel` that lands before the
|
||||
/// entry pops prevents delivery; a `cancel` after it has fired finds
|
||||
/// nothing (the race signal). Bounded by armed-but-not-yet-resolved timers
|
||||
/// and self-collecting — no sweep. `Sleep` / `WaitTimeout` never touch it.
|
||||
armed: std::collections::HashSet<u64>,
|
||||
}
|
||||
|
||||
impl Timers {
|
||||
pub fn new() -> Self {
|
||||
Self { heap: BinaryHeap::new(), next_seq: 0 }
|
||||
Self { heap: BinaryHeap::new(), next_seq: 0, armed: std::collections::HashSet::new() }
|
||||
}
|
||||
|
||||
/// Insert a `Sleep` timer. Convenience for the common case.
|
||||
@@ -111,6 +141,33 @@ impl Timers {
|
||||
self.insert(deadline, pid, Reason::Sleep { epoch });
|
||||
}
|
||||
|
||||
/// Arm a cancellable `send_after` timer: run `fire` at `deadline` unless
|
||||
/// [`cancel`](Self::cancel)led first. `pid` is informational only (the
|
||||
/// destination, or who armed it — useful for introspection); it is *not*
|
||||
/// used to wake anyone, the delivery lives entirely inside `fire`. Returns
|
||||
/// a [`TimerId`] for cancellation.
|
||||
pub fn insert_send(
|
||||
&mut self,
|
||||
deadline: Instant,
|
||||
pid: Pid,
|
||||
fire: Box<dyn FnOnce() + Send>,
|
||||
) -> TimerId {
|
||||
let seq = self.next_seq;
|
||||
self.next_seq = self.next_seq.wrapping_add(1);
|
||||
self.armed.insert(seq);
|
||||
self.heap.push(Reverse(Entry { deadline, seq, pid, reason: Reason::Send { fire } }));
|
||||
TimerId(seq)
|
||||
}
|
||||
|
||||
/// Cancel an armed `send_after` timer. Returns `true` if the timer was
|
||||
/// still armed (delivery is now prevented), `false` if it had already
|
||||
/// fired or been cancelled. The heap entry, if still pending, is left to be
|
||||
/// discarded when its deadline passes — `pop_due` drops any `Send` entry
|
||||
/// whose `seq` is no longer armed.
|
||||
pub fn cancel(&mut self, id: TimerId) -> bool {
|
||||
self.armed.remove(&id.0)
|
||||
}
|
||||
|
||||
/// Insert an arbitrary timer entry.
|
||||
pub fn insert(&mut self, deadline: Instant, pid: Pid, reason: Reason) {
|
||||
let seq = self.next_seq;
|
||||
@@ -127,6 +184,7 @@ impl Timers {
|
||||
/// discarded so it can't keep the runtime alive.
|
||||
pub fn clear(&mut self) {
|
||||
self.heap.clear();
|
||||
self.armed.clear();
|
||||
}
|
||||
|
||||
/// Soonest pending deadline, or `None` if the heap is empty.
|
||||
@@ -136,11 +194,21 @@ impl Timers {
|
||||
|
||||
/// Pop every entry whose deadline is ≤ `now`, in deadline order.
|
||||
/// The scheduler dispatches each entry by inspecting `entry.reason`.
|
||||
///
|
||||
/// A due `Send` entry is returned only if it is still armed; a cancelled
|
||||
/// one is silently dropped here (its `seq` was already removed from
|
||||
/// `armed` by [`cancel`](Self::cancel)). Returning it removes it from
|
||||
/// `armed`, so a later `cancel` of a fired timer reports `false`.
|
||||
pub fn pop_due(&mut self, now: Instant) -> Vec<Entry> {
|
||||
let mut out = Vec::new();
|
||||
while let Some(r) = self.heap.peek() {
|
||||
if r.0.deadline <= now {
|
||||
out.push(self.heap.pop().unwrap().0);
|
||||
let entry = self.heap.pop().unwrap().0;
|
||||
if matches!(entry.reason, Reason::Send { .. }) && !self.armed.remove(&entry.seq) {
|
||||
// Cancelled before it came due: discard, do not deliver.
|
||||
continue;
|
||||
}
|
||||
out.push(entry);
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user