diff --git a/src/lib.rs b/src/lib.rs index 2b5a367..2920bc2 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -80,7 +80,8 @@ pub use registry::{ }; pub use runtime::{init, Config, Runtime}; pub use scheduler::{ - block_on_io, cancel_timer, request_stop, run, self_pid, send_after, send_after_named, sleep, sleep_wall, + block_on_io, cancel_timer, request_stop, run, self_pid, send_after, send_after_named, + send_after_named_wall, send_after_wall, sleep, sleep_wall, spawn, spawn_addr, spawn_under, wait_readable, wait_readable_timeout, wait_writable, wait_writable_timeout, yield_now, FdArm, JoinError, JoinHandle, }; diff --git a/src/scheduler.rs b/src/scheduler.rs index 7fdc42e..40b8eff 100644 --- a/src/scheduler.rs +++ b/src/scheduler.rs @@ -497,6 +497,51 @@ pub fn send_after_named( }) } +/// Wall-anchored [`send_after`] (RFC 007): under causal profiling the timer +/// opts out of the virtual-time shift and fires at its raw deadline instead +/// of dilating with injected delay — the user-facing opt-out whose substrate +/// [`sleep_wall`] landed. Use it for deadlines that reflect the outside +/// world (protocol timeouts, wall-clock schedules) rather than workload +/// pacing. Without the `smarm-causal` feature it is identical to +/// [`send_after`]. Cancellation via [`cancel_timer`] is unchanged. +pub fn send_after_wall( + after: std::time::Duration, + dest: Pid, + 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| { + match inner.timers.lock() { + Ok(mut timers) => timers.insert_send_wall(deadline, dest.erase(), fire), + Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"), + } + }) +} + +/// Wall-anchored [`send_after_named`] (RFC 007): re-resolving name delivery, +/// raw-deadline anchor — see [`send_after_wall`] for the semantics. +pub fn send_after_named_wall( + after: std::time::Duration, + dest: Name, + 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| { + match inner.timers.lock() { + Ok(mut timers) => timers.insert_send_wall(deadline, armed_by, fire), + Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"), + } + }) +} + /// Deliver `msg` onto a channel the caller owns after `after`, rather than to a /// registry address. The sibling of [`send_after`] used by the gen_server timer /// layer (RFC 015 §5): arming a server timer must land the fire on the loop's diff --git a/src/timer.rs b/src/timer.rs index ea041bb..143cae0 100644 --- a/src/timer.rs +++ b/src/timer.rs @@ -193,6 +193,22 @@ impl Timers { TimerId(self.push(deadline, pid, Reason::Send { fire }, false)) } + /// Arm a *wall-anchored* cancellable `send_after` timer (RFC 007): the + /// same contract as [`insert_send`](Self::insert_send), but the entry + /// opts out of the virtual-time shift and fires at its raw deadline + /// regardless of injected delay — the `Send`-reason sibling of + /// [`insert_sleep_wall`](Self::insert_sleep_wall). Without the + /// `smarm-causal` feature this is identical to `insert_send`. + pub fn insert_send_wall( + &mut self, + deadline: Instant, + pid: Pid, + fire: Box, + ) -> TimerId { + self.armed.insert(self.next_seq); + TimerId(self.push(deadline, pid, Reason::Send { fire }, true)) + } + /// 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 diff --git a/tests/causal.rs b/tests/causal.rs index c9f1a1b..653700b 100644 --- a/tests/causal.rs +++ b/tests/causal.rs @@ -570,3 +570,81 @@ fn sleeping_actor_pays_injected_delay() { "sleeper paid nothing: slept {ms}ms, expected >= 150ms (raw 50 + injected 150)" ); } + +/// RFC 007 user-facing opt-out: a wall-anchored `send_after` entry fires at +/// its raw deadline while a virtual sibling armed at the same instant is +/// shifted by injected delay — the `Send`-reason mirror of +/// `wall_timer_ignores_injected_delay`. +#[test] +fn wall_send_after_ignores_injected_delay() { + let _s = serial(); + use smarm::pid::Pid; + use smarm::timer::Timers; + use std::time::Instant; + + let hz = smarm::causal::tsc_hz(); + let mut t = Timers::new(); + let now = Instant::now(); + let wall_fired = Arc::new(AtomicBool::new(false)); + let virt_fired = Arc::new(AtomicBool::new(false)); + let (w2, v2) = (wall_fired.clone(), virt_fired.clone()); + let _wid = t.insert_send_wall( + now + Duration::from_millis(50), + Pid::new(0, 0), + Box::new(move || w2.store(true, Ordering::Relaxed)), + ); + let _vid = t.insert_send( + now + Duration::from_millis(50), + Pid::new(1, 0), + Box::new(move || v2.store(true, Ordering::Relaxed)), + ); + + // 100ms of debt lands while both are pending. + smarm::causal::inject_delay_cycles_for_test((hz * 0.100) as u64); + + // Just past the raw deadline: only the wall send pops; run its thunk. + let due = t.pop_due(now + Duration::from_millis(60)); + assert_eq!(due.len(), 1, "exactly the wall send must fire at raw deadline"); + for e in due { + if let smarm::timer::Reason::Send { fire } = e.reason { + fire(); + } + } + assert!(wall_fired.load(Ordering::Relaxed), "wall send must deliver"); + assert!(!virt_fired.load(Ordering::Relaxed), "virtual send must not"); + assert!(!t.is_empty(), "virtual sibling must remain queued, shifted"); +} + +/// `cancel` is anchor-agnostic: a wall-anchored `send_after` cancels exactly +/// like a virtual one and never delivers, even with delay debt outstanding. +#[test] +fn wall_send_after_cancels() { + let _s = serial(); + use smarm::pid::Pid; + use smarm::timer::Timers; + use std::time::Instant; + + let hz = smarm::causal::tsc_hz(); + let mut t = Timers::new(); + let now = Instant::now(); + let delivered = Arc::new(AtomicBool::new(false)); + let d2 = delivered.clone(); + let id = t.insert_send_wall( + now + Duration::from_millis(20), + Pid::new(0, 0), + Box::new(move || d2.store(true, Ordering::Relaxed)), + ); + + smarm::causal::inject_delay_cycles_for_test((hz * 0.100) as u64); + + assert!(t.cancel(id), "cancel must find the wall-anchored send"); + for e in t.pop_due(now + Duration::from_secs(3600)) { + if let smarm::timer::Reason::Send { fire } = e.reason { + fire(); + } + } + assert!( + !delivered.load(Ordering::Relaxed), + "cancelled wall send_after delivered" + ); +} diff --git a/tests/timer.rs b/tests/timer.rs index 00d0c01..81fd0fd 100644 --- a/tests/timer.rs +++ b/tests/timer.rs @@ -405,3 +405,60 @@ fn send_after_to_dead_typed_pid_is_silent() { assert_eq!(report_rx.try_recv(), Ok(None)); }); } + +// --------------------------------------------------------------------------- +// Wall-anchored send_after (RFC 007 user-facing opt-out). The API exists in +// both feature configs; featureless it is behaviourally identical to +// `send_after` — these tests pin exactly that. +// --------------------------------------------------------------------------- + +#[test] +fn armed_wall_send_timer_is_returned_and_fires() { + let mut t = Timers::new(); + let now = Instant::now(); + let fired = Arc::new(AtomicBool::new(false)); + let f = fired.clone(); + let _id = t.insert_send_wall( + now + Duration::from_millis(10), + Pid::new(0, 0), + Box::new(move || f.store(true, Ordering::SeqCst)), + ); + + let mut due = t.pop_due(now + Duration::from_millis(20)); + assert_eq!(due.len(), 1, "an armed wall send timer should pop when due"); + run_fire(due.pop().unwrap()); + assert!(fired.load(Ordering::SeqCst), "running the thunk delivers"); + assert!(t.is_empty()); +} + +use smarm::send_after_named_wall; + +#[test] +fn send_after_named_wall_delivers_after_the_delay() { + const WPING: Name = Name::new("send_after_wall_ping"); + run(|| { + let (tx, rx) = channel::(); + register(WPING, tx).unwrap(); + let t0 = Instant::now(); + let _id = send_after_named_wall(Duration::from_millis(30), WPING, 99); + assert_eq!(rx.recv().unwrap(), 99); + assert!( + t0.elapsed() >= Duration::from_millis(25), + "delivered too early: {:?}", + t0.elapsed() + ); + }); +} + +#[test] +fn send_after_named_wall_cancels() { + const WC: Name = Name::new("send_after_wall_cancel"); + run(|| { + let (tx, rx) = channel::(); + register(WC, tx).unwrap(); + let id = send_after_named_wall(Duration::from_millis(50), WC, 7); + assert!(cancel_timer(id), "cancel before fire returns true"); + sleep(Duration::from_millis(90)); + assert_eq!(rx.try_recv(), Ok(None), "cancelled wall timer delivered"); + }); +}