//! Start actors and control them: `run`, `spawn`, `sleep`, timers, and IO waits. //! //! ## What is an actor? //! //! An actor in smarm is a *green thread*: a lightweight, cooperatively //! scheduled unit of execution with its own stack, running alongside //! thousands of others on a handful of OS threads. You give it a closure; //! smarm gives it a [`Pid`] and a [`JoinHandle`]. Actors talk to each other by //! sending messages over [`channel`](mod@crate::channel)s, not by sharing memory, //! so most of the concurrency bugs that come from shared mutable state simply //! don't arise. //! //! Everything in smarm (channels, [`Mutex`](crate::Mutex), timers, IO waits, //! `gen_server`) needs to run on top of a scheduler. [`run`] starts one and //! blocks the calling OS thread until the actor you give it finishes; from //! inside that actor (or any actor it spawns), call [`spawn`] to start more. //! //! ``` //! use smarm::{run, spawn}; //! //! run(|| { //! let handle = spawn(|| { //! println!("hello from another actor"); //! }); //! handle.join().unwrap(); //! }); //! ``` //! //! ## Waiting for an actor: `JoinHandle` //! //! [`spawn`] returns a [`JoinHandle`], your only handle on the actor's //! outcome. Call [`JoinHandle::join`] to block the calling actor until the //! spawned one finishes: //! //! - if it returned normally, `join` returns `Ok(())`; //! - if it panicked, `join` returns `Err(`[`JoinError`]`)` carrying the panic //! payload, so a crash in one actor never silently vanishes and never //! crashes the process; the caller decides what to do with it (log it, //! propagate it, ignore it). //! //! If you don't need the result, drop the `JoinHandle` (or never bind it): //! the actor keeps running independently. To watch an actor without //! blocking, or to react to failures across a whole tree of actors, see //! [`monitor`](mod@crate::monitor), [`link`](mod@crate::link), and //! [`supervisor`](crate::supervisor) instead. //! //! ## Sleeping, timeouts, and IO //! //! [`sleep`] parks only the calling actor, not the OS thread underneath it, //! so thousands of sleeping actors cost nothing but the timer entry. //! [`send_after`] and [`send_after_named`] schedule a message to be delivered //! later (Erlang-style `send_after`), and [`cancel_timer`] can call one off //! before it fires. //! //! For blocking or file-descriptor-based IO, see [`block_on_io`], //! [`wait_readable`], and [`wait_writable`]: they park the calling actor and //! resume it when the work completes or the fd is ready, again without //! blocking an OS thread. //! //! ## Stopping an actor from the outside //! //! [`request_stop`] asks an actor to cooperatively unwind: it's the //! mechanism `gen_server` shutdown, timeouts, and supervisor restarts are //! built on. It's best-effort: an actor that never yields, allocates, or //! blocks (a tight loop with nothing else in it) has no opportunity to //! notice the request. use crate::actor::current_pid; use crate::channel::Sender; use crate::pid::{Name, Pid}; use crate::runtime::{self, RuntimeInner, YieldIntent, RUNTIME}; use crate::supervisor::Signal; use std::sync::atomic::Ordering; use std::sync::{Arc, Weak}; // --------------------------------------------------------------------------- // with_runtime / try_with_runtime // --------------------------------------------------------------------------- // Borrow the current runtime. Panics if called outside `Runtime::run()`. // // Preemption is disabled for the whole span. `f` holds a thread-local borrow // of `RUNTIME`; if a preemption-driven context switch moved the actor to a // different OS thread in the middle of `f`, the borrow guard would be // released on the wrong thread's copy of the thread-local, corrupting its // borrow count. `f` is also always runtime bookkeeping that should run to // completion without the actor being suspended or unwound partway through. #[inline(never)] pub(crate) fn with_runtime(f: impl FnOnce(&Arc) -> R) -> R { crate::context::tls_fence(); let prev = crate::preempt::preemption_swap(false); let result = RUNTIME.with(|r| { let b = r.borrow(); let inner = match b.as_ref() { Some(inner) => inner, None => panic!("smarm: not inside Runtime::run()"), }; f(inner) }); crate::preempt::preemption_swap(prev); result } // Borrow the runtime if present, otherwise `None`. Used on cleanup paths // (e.g. a channel's Drop impl during teardown) that may run after the // runtime has already gone away. Same preemption gate as `with_runtime`. #[inline(never)] pub(crate) fn try_with_runtime(f: impl FnOnce(&Arc) -> R) -> Option { crate::context::tls_fence(); let prev = crate::preempt::preemption_swap(false); let result = RUNTIME.with(|r| r.borrow().as_ref().map(f)); crate::preempt::preemption_swap(prev); result } // --------------------------------------------------------------------------- // JoinHandle / JoinError // --------------------------------------------------------------------------- /// The spawned actor panicked. Returned by [`JoinHandle::join`]; `payload` /// is exactly what the panic carried (the value passed to `panic!`, or /// whatever a library panicked with), the same payload you'd get from /// [`std::thread::JoinHandle::join`]. Downcast it if you need to inspect it: /// /// ``` /// use smarm::{run, spawn}; /// /// run(|| { /// let h = spawn(|| panic!("boom")); /// let err = h.join().unwrap_err(); /// let msg = err.payload.downcast_ref::<&str>().copied().unwrap_or("?"); /// assert_eq!(msg, "boom"); /// }); /// ``` #[derive(Debug)] pub struct JoinError { pub payload: Box, } /// A handle to a spawned actor, returned by [`spawn`], [`spawn_under`], and /// friends. Use [`join`](Self::join) to wait for the actor to finish and /// collect its outcome, or [`pid`](Self::pid) to get its identity for use /// with [`request_stop`], [`monitor`](crate::monitor::monitor), or /// [`link`](crate::link::link). /// /// If you never call `join` (or drop the handle instead), the actor is not /// affected: it keeps running and its resources are still reclaimed when it /// finishes. `join` is how you find out *what happened*, not a requirement /// for the actor to make progress or clean up. pub struct JoinHandle { pid: Pid, consumed: bool, } impl JoinHandle { /// The identity of the actor this handle refers to. pub fn pid(&self) -> Pid { self.pid } /// Block the calling actor until the spawned actor finishes, then /// report how it finished: `Ok(())` if it returned normally or stopped /// cooperatively via [`request_stop`], `Err(`[`JoinError`]`)` if it /// panicked. pub fn join(mut self) -> Result<(), JoinError> { use crate::actor::Outcome; let me = match current_pid() { Some(pid) => pid, None => panic!("join() called outside an actor"), }; loop { // Check-Done-or-register-waiter is atomic under the target's cold // lock; finalize publishes Done and takes the waiter list under // the same lock, so we either see the outcome or are woken. let outcome = with_runtime(|inner| { let slot = match inner.slot_at(self.pid) { Some(slot) => slot, None => panic!("join: pid index out of range: {:?}", self.pid), }; let mut cold = slot.cold.lock(); match slot.status_for(self.pid) { // Our outstanding handle pins the slot: it cannot be // reclaimed (generation cannot change) while we hold it. crate::slot_state::Status::Stale => { panic!("join: target slot has been reused") } crate::slot_state::Status::Done => Some(match cold.outcome.take() { Some(outcome) => outcome, None => panic!("Done slot must have outcome"), }), crate::slot_state::Status::Live => { // begin_wait is lock-free, legal under the cold lock; // registering under it makes the epoch atomic with // the check-Done-or-register linearization point. cold.waiters.push((me, begin_wait())); None } } }); match outcome { Some(o) => { self.consumed = true; self.decrement_handle_count(); return match o { Outcome::Exit => Ok(()), Outcome::Panic(p) => Err(JoinError { payload: p }), // A cooperative stop carries no panic payload to // propagate; the *reason* is observable via monitors // (DownReason::Stopped). join() therefore reports Ok. Outcome::Stopped => Ok(()), }; } None => { let _np = NoPreempt::enter(); park_current(); } } } } fn decrement_handle_count(&mut self) { with_runtime(|inner| { let should_reclaim = match inner.slot_at(self.pid) { Some(slot) => { let mut cold = slot.cold.lock(); match slot.status_for(self.pid) { crate::slot_state::Status::Stale => false, status => { cold.outstanding_handles = cold.outstanding_handles.saturating_sub(1); cold.outstanding_handles == 0 && status == crate::slot_state::Status::Done } } } None => false, }; if should_reclaim { // Re-verified inside; benign if finalize's reclaim won a race. crate::runtime::reclaim_slot(inner, self.pid); } }); } } impl Drop for JoinHandle { fn drop(&mut self) { if !self.consumed { // May be called outside run() if handle is dropped after teardown. if try_with_runtime(|_| ()).is_some() { self.decrement_handle_count(); } } } } // --------------------------------------------------------------------------- // spawn / spawn_under / self_pid // --------------------------------------------------------------------------- /// Per-spawn stack shape overrides (RFC 019). `None` fields resolve to the /// runtime's [`Config`](crate::runtime::Config) defaults at spawn time, so /// struct-update syntax works anywhere without a runtime handle: /// /// ``` /// use smarm::SpawnOpts; /// let opts = SpawnOpts { stack_reserve: Some(8 * 1024 * 1024), ..SpawnOpts::default() }; /// assert_eq!(opts.stack_reserve, Some(8 * 1024 * 1024)); /// ``` /// /// Both sizes are page-rounded. The reserve is *virtual* (demand-paged): /// an 8 MiB reserve costs address space, not memory — RSS follows touched /// pages. The guard is PROT_NONE below the stack; raise it for FFI code /// with unusually large C frames. Custom-shaped stacks bypass the recycle /// pool: they are mmapped fresh at spawn and munmapped at death. #[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] pub struct SpawnOpts { /// Usable stack reservation. `None` ⇒ [`Config::stack_reserve`](crate::runtime::Config::stack_reserve). pub stack_reserve: Option, /// PROT_NONE guard below the stack. `None` ⇒ [`Config::stack_guard`](crate::runtime::Config::stack_guard). pub guard_size: Option, } /// Why [`try_spawn`] could not start an actor. /// /// Marked `non_exhaustive`: today the only refusal is a full slab, but a /// future variant (say, a shutdown-in-progress refusal) must not be a /// breaking change for shed-path `match`es. #[non_exhaustive] #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub enum SpawnError { /// The fixed actor slab ([`Config::max_actors`](crate::runtime::Config::max_actors)) /// is full: every slot is claimed /// by a live actor. This is a routine overload condition, not an /// invariant violation — shed the unit of work (close the socket, /// return a 503) and try again once actors have died. AtCapacity, } impl core::fmt::Display for SpawnError { fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { match self { SpawnError::AtCapacity => { write!( f, "actor slab at capacity (`Config::max_actors` live actors)" ) } } } } impl std::error::Error for SpawnError {} /// Start a new actor running `f`, and return a [`JoinHandle`] for it. /// /// The new actor runs concurrently with its caller and with every other /// actor in the runtime; smarm schedules it cooperatively across the /// available OS threads. `spawn` can be called from inside `run`'s closure, /// or from inside any actor (spawning a child from a child works the same /// way); it cannot be called before `run` has started or after it returns. /// /// The returned [`JoinHandle`] is how you learn how the actor finished. If /// you don't need that, it's fine to drop it: the actor still runs to /// completion either way. /// /// Panics if called outside `Runtime::run()`. pub fn spawn(f: impl FnOnce() + Send + 'static) -> JoinHandle { let parent = current_pid().unwrap_or_else(|| { // Outside an actor but inside run(): the initial spawn. with_runtime // panics with "not inside Runtime::run()" if there's no runtime at all. with_runtime(|_| crate::runtime::ROOT_PID) }); spawn_under(parent, f) } /// [`spawn`] with per-actor stack shape overrides (RFC 019). pub fn spawn_with(opts: SpawnOpts, f: impl FnOnce() + Send + 'static) -> JoinHandle { let parent = current_pid().unwrap_or_else(|| with_runtime(|_| crate::runtime::ROOT_PID)); spawn_under_with(parent, opts, f) } /// Like [`spawn`], but explicitly attaches the new actor to `supervisor` /// instead of the calling actor. Ordinary code should reach for [`spawn`]; /// this exists for supervision trees (see [`supervisor`](crate::supervisor)) /// and other cases that need to place a child under a specific ancestor /// rather than its true caller. pub fn spawn_under(supervisor: Pid, f: impl FnOnce() + Send + 'static) -> JoinHandle { spawn_under_with(supervisor, SpawnOpts::default(), f) } /// [`spawn_under`] with per-actor stack shape overrides (RFC 019). pub fn spawn_under_with( supervisor: Pid, opts: SpawnOpts, f: impl FnOnce() + Send + 'static, ) -> JoinHandle { let supervisor = supervisor.erase(); // Stack + closure boxing happen before the slot locks are taken; the // pool lock inside acquire_stack is dropped before any mmap, so no // syscall ever stalls another scheduler thread. let stack = with_runtime(|inner| crate::runtime::acquire_stack(inner, opts)); let sp = init_actor_stack(stack.top(), crate::actor::trampoline); let closure: crate::runtime::Closure = Box::new(f); let pid = with_runtime(|inner| { let idx = inner.allocate_slot(); // panics loudly on slab exhaustion crate::runtime::install_actor(inner, idx, sp, stack, supervisor, closure, None) }); JoinHandle { pid, consumed: false, } } /// [`spawn`] and [`monitor`](crate::monitor()) the child in one step, with no /// window in which the child can die unobserved. /// /// `spawn` followed by `monitor(h.pid())` races: on a multi-scheduler /// runtime the child can run to completion before the monitor registers, and /// a monitor on a dead pid delivers [`DownReason::NoProc`](crate::DownReason::NoProc) /// — the real reason (Exit vs Panic) is lost. Here the monitor is registered /// on the child's slot *before* the child is published to any run queue, so /// the `Down` always carries the child's actual termination reason. This is /// Erlang's `spawn_monitor/1`. pub fn spawn_monitor(f: impl FnOnce() + Send + 'static) -> (JoinHandle, crate::monitor::Monitor) { spawn_monitor_with(SpawnOpts::default(), f) } /// [`spawn_monitor`] with per-actor stack shape overrides (RFC 019). pub fn spawn_monitor_with( opts: SpawnOpts, f: impl FnOnce() + Send + 'static, ) -> (JoinHandle, crate::monitor::Monitor) { let parent = current_pid().unwrap_or_else(|| with_runtime(|_| crate::runtime::ROOT_PID)); let (tx, rx) = crate::channel::channel::(); let stack = with_runtime(|inner| crate::runtime::acquire_stack(inner, opts)); let sp = init_actor_stack(stack.top(), crate::actor::trampoline); let closure: crate::runtime::Closure = Box::new(f); let (pid, id) = with_runtime(|inner| { let idx = inner.allocate_slot(); // panics loudly on slab exhaustion let id = inner.alloc_monitor_id(); let pid = crate::runtime::install_actor(inner, idx, sp, stack, parent, closure, Some((id, tx))); (pid, id) }); ( JoinHandle { pid, consumed: false, }, crate::monitor::Monitor { id, target: pid, rx, }, ) } /// [`spawn`] that reports a full actor slab instead of panicking. /// /// Behaviour parity with [`spawn`] in every case except one: when the fixed /// slab ([`Config::max_actors`](crate::runtime::Config::max_actors)) is /// full, this returns [`Err(SpawnError::AtCapacity)`](SpawnError::AtCapacity) /// where `spawn` panics the calling actor. Use it at load-shedding call /// sites — an accept loop spawning one actor per connection, a request /// admission point — where "at capacity" is a routine overload condition to /// handle (reject the unit of work), not an invariant violation. Internal /// and bounded spawn sites should keep [`spawn`]: there, the panic is a /// correct loud invariant check. /// /// The claim is atomic (claim-or-report): there is no /// check-then-spawn race against other spawners for the last slot, so no /// headroom margin is needed. pub fn try_spawn(f: impl FnOnce() + Send + 'static) -> Result { let parent = current_pid().unwrap_or_else(|| with_runtime(|_| crate::runtime::ROOT_PID)); try_spawn_under_with(parent, SpawnOpts::default(), f) } /// [`try_spawn`] with an explicit supervisor and per-actor stack shape /// overrides — the full-control core the other `try_` surface is built on /// (mirrors [`spawn_under_with`]). pub fn try_spawn_under_with( supervisor: Pid, opts: SpawnOpts, f: impl FnOnce() + Send + 'static, ) -> Result { let supervisor = supervisor.erase(); // Slot FIRST — deliberately the reverse of `spawn`'s stack-first order: // under overload the Err arm is the HOT path, and a rejection must cost // one mutex pop, not an mmap/pool-pop + init + recycle per shed unit of // work. The claim is a single atomic pop (no TOCTOU; see // `try_allocate_slot`). let idx = match with_runtime(|inner| inner.try_allocate_slot()) { Some(idx) => idx, None => return Err(SpawnError::AtCapacity), }; // Between claim and install the slot is owned by this frame alone; if // stack allocation panics in that window the slot must go back or it // leaks for the life of the runtime (and would trip the run()-teardown // slot-leak debug_assert). struct ReturnOnUnwind(Option); impl Drop for ReturnOnUnwind { fn drop(&mut self) { if let Some(idx) = self.0 { with_runtime(|inner| inner.return_vacant_slot(idx)); } } } let mut claimed = ReturnOnUnwind(Some(idx)); let stack = with_runtime(|inner| crate::runtime::acquire_stack(inner, opts)); let sp = init_actor_stack(stack.top(), crate::actor::trampoline); let closure: crate::runtime::Closure = Box::new(f); claimed.0 = None; // install_actor takes ownership of the slot from here let pid = with_runtime(|inner| { crate::runtime::install_actor(inner, idx, sp, stack, supervisor, closure, None) }); Ok(JoinHandle { pid, consumed: false, }) } /// Spawn an actor that other actors can message directly by its [`Pid`], /// rather than only by holding on to a channel `Sender` you passed it /// yourself. /// /// `body` receives the [`Receiver`](crate::channel::Receiver) smarm /// creates for it; `spawn_addr` publishes the matching `Sender` and hands /// back the actor's typed address. That address is usable the instant you /// hold it: an immediate [`send_to`](crate::send_to) on the returned pid /// always finds the inbox, even if `body` hasn't started running yet. /// /// The spawned actor is detached (there is no [`JoinHandle`] to join): its /// lifetime is up to its own logic, for example running until it receives a /// stop message, or until the actor decides to return. This mirrors how /// [`GenServerBuilder::start`](crate::GenServerBuilder::start) works. /// /// Panics if called outside `Runtime::run()`. pub fn spawn_addr( body: impl FnOnce(crate::channel::Receiver) + Send + 'static, ) -> Pid { let (tx, rx) = crate::channel::channel::(); let handle = spawn(move || body(rx)); let pid = handle.pid(); // Publish the sender for `pid` before returning the typed address. `handle` // drops at end of scope (detached). crate::registry::install_for::(pid, tx); crate::pid::assert_type::(pid) } /// [`spawn_addr`] with per-actor stack shape overrides (RFC 019). pub fn spawn_addr_with( opts: SpawnOpts, body: impl FnOnce(crate::channel::Receiver) + Send + 'static, ) -> Pid { let (tx, rx) = crate::channel::channel::(); let handle = spawn_with(opts, move || body(rx)); let pid = handle.pid(); crate::registry::install_for::(pid, tx); crate::pid::assert_type::(pid) } use crate::context::init_actor_stack; /// The identity of the actor currently running. Use it to hand your own /// address to another actor (for a reply, a monitor, or a link). /// /// Panics if called outside an actor (for example, from the closure passed /// to [`run`] itself, before any [`spawn`]). pub fn self_pid() -> Pid { match current_pid() { Some(pid) => pid, None => panic!("self_pid() called outside an actor"), } } // --------------------------------------------------------------------------- // yield_now / park_current / unpark // --------------------------------------------------------------------------- /// Voluntarily give up the CPU so another runnable actor gets a turn, then /// resume as soon as the scheduler gets back around to you. Use this in a /// long-running, allocation-free loop that you want to stay cooperative with /// the rest of the runtime (see also the [`check!`](crate::check) macro, /// which does the same thing conditionally, only when your timeslice has /// actually run out). pub fn yield_now() { runtime::set_yield_intent(YieldIntent::Yield); unsafe { crate::context::switch_to_scheduler() }; // Observation point: we may have been resumed only to be cancelled. crate::preempt::check_cancelled(); } // Suspend the current actor until something wakes it (a message arrives, a // timer fires, a lock is granted, and so on). This is the low-level parking // primitive that every blocking smarm operation (channel recv, sleep, // Mutex::lock, IO waits, JoinHandle::join) is built on; application code // should reach for one of those rather than calling this directly. // // Checks for a pending cooperative-stop request both before parking (a stop // requested while merely queued to run would otherwise have no future wake // to catch it) and after resuming (so a stop that arrived while parked is // noticed as soon as we wake, and unwinds from here exactly like any other // blocking call would). pub fn park_current() { crate::preempt::check_cancelled(); runtime::set_yield_intent(YieldIntent::Park); unsafe { crate::context::switch_to_scheduler() }; crate::preempt::check_cancelled(); } // Wake `pid` unconditionally, regardless of what it's currently waiting for. // Reserved for terminal wakes (`request_stop`); anything waking an actor from // a specific registered wait (a channel send, a mutex grant, a timer, an IO // completion) must use `unpark_at` instead, so a stale wakeup can never be // mistaken for the one the actor is actually waiting on. pub fn unpark(pid: Pid) { let _ = try_with_runtime(|inner| inner.unpark(pid)); } // Wake `pid` only if its current wait is still the one this waker // registered for (an "epoch-matched" unpark). Every registration-based // waker (channel senders, mutex grants, wait-timers, IO completions, joiner // wakes) must use this rather than the unconditional `unpark`. pub(crate) fn unpark_at(pid: Pid, epoch: u32) { let _ = try_with_runtime(|inner| inner.unpark_at(pid, epoch)); } // The current actor's runtime as a `Weak`, for a waker that must reach the // runtime from a foreign thread later. A channel captures this ONCE, the first // time its receiver parks, so a cross-thread `send` can wake without the // `RUNTIME` thread-local (unset off a scheduler thread). `None` off a // scheduler thread, where there is nothing to capture. // // Deliberately not called per park: `Arc::downgrade` plus the matching drop is // a locked RMW pair on one globally shared counter, and the park/unpark // round-trip is the hot path of every channel workload. pub(crate) fn runtime_weak() -> Option> { try_with_runtime(Arc::downgrade) } // Epoch-matched wake of `pid` from a waker that may or may not be on a // scheduler thread. On a scheduler thread we take the thread-local path // (preemption-gated, slot-eligible); off one that path is a silent no-op, so // we reach the runtime through `rt` — the `Weak` the waker captured while it // was in-runtime. Mirrors the IO backend's cross-context wake (io.rs, RFC 018). pub(crate) fn unpark_at_via( pid: Pid, epoch: u32, rt: impl FnOnce() -> Option>, ) { if try_with_runtime(|inner| inner.unpark_at(pid, epoch)).is_some() { return; } // Off a scheduler thread only: `rt` re-takes the waker's lock to read the // captured `Weak`, which is why it is a closure and not a value — the // in-runtime path above must not pay for it. if let Some(inner) = rt().as_ref().and_then(Weak::upgrade) { inner.unpark_at(pid, epoch); } } // Open a new wait for the current actor and return its wait identity // ("epoch"). Call once per wait, before registering with any waker. Lock-free, // so it's legal to call while already holding another internal lock. pub(crate) fn begin_wait() -> u32 { let me = match current_pid() { Some(pid) => pid, None => panic!("begin_wait() called outside an actor"), }; with_runtime(|inner| inner.begin_wait(me)) } // Close the current actor's wait without parking on it: the no-park exit // used when a multi-arm `select` finds an arm already ready at registration // time. Leftover registrations from the other arms are left to self-clean // when their wakers try to use them. pub(crate) fn retire_wait() { let me = match current_pid() { Some(pid) => pid, None => panic!("retire_wait() called outside an actor"), }; with_runtime(|inner| inner.retire_wait(me)); crate::preempt::check_cancelled(); } /// Ask `pid` to stop cooperatively. /// /// This sets a flag on the target actor and wakes it so it notices promptly; /// the actor itself decides when it's safe to actually unwind, at its next /// natural checkpoint (a blocking call returning, a `check!()`, or an /// allocation). Once it does, it terminates as if it panicked, except that /// [`JoinHandle::join`] reports it as a normal, non-error exit: cooperative /// stop is a controlled shutdown, not a failure. /// /// This is the *hard* stop — OTP's `exit(Pid, kill)`. It is what a supervisor /// falls back to when a child overstays its [`Shutdown`](crate::supervisor::Shutdown) /// grace period. For a stop the target gets to prepare for, use /// [`request_shutdown`]; for structured teardown, reach for /// [`GenServerRef::shutdown`](crate::GenServerRef::shutdown) or a /// [`supervisor`](crate::supervisor) instead of calling this directly. /// /// Because it's cooperative, an actor stuck in a tight loop with no /// blocking call, no [`check!`](crate::check), and no allocation cannot be /// stopped, for the same reason it cannot be preempted. Calling this on an /// actor that has already finished is a harmless no-op. pub fn request_stop(pid: Pid) { let pid = pid.erase(); let _ = try_with_runtime(|inner| request_stop_inner(inner, pid)); } // The core of `request_stop`, taking the runtime directly so it can also be // driven from inside the runtime itself (the RuntimeHandle path, supervisor // sweeps) without // re-borrowing the thread-local. Sets the stop flag under the target's lock // (a generation mismatch, or no live actor there, makes it a no-op) and // wakes the target. pub(crate) fn request_stop_inner(inner: &RuntimeInner, pid: Pid) { if let Some(slot) = inner.slot_at(pid) { { let cold = slot.cold.lock(); // Verify under the cold lock: generation can't change while // we hold it (reclaim takes the same lock). if slot.generation() == pid.generation() { if let Some(actor) = cold.actor.as_ref() { actor.stop.store(true, std::sync::atomic::Ordering::Relaxed); } } } inner.unpark(pid); } } /// Ask an actor to shut down gracefully — OTP's `exit(Pid, shutdown)`, where /// [`request_stop`] is `exit(Pid, kill)`. /// /// If the target has called [`trap_exit`](crate::trap_exit), it receives an /// [`ExitSignal`](crate::ExitSignal) with reason /// [`DownReason::Shutdown`](crate::DownReason::Shutdown) on its trap inbox and /// keeps running: the request is advisory, and the target is expected to wind /// down and exit normally in its own time (a supervisor bounds that time with /// its child's [`Shutdown`](crate::supervisor::Shutdown) policy and falls back /// to `request_stop`). A target that is not trapping is stopped exactly as by /// `request_stop`. A dead pid is a no-op. /// /// The signal's `from` is the calling actor, or `ROOT_PID` when driven from /// outside the runtime (see [`RuntimeHandle::request_shutdown`](crate::RuntimeHandle::request_shutdown)). pub fn request_shutdown(pid: Pid) { let pid = pid.erase(); let from = current_pid().unwrap_or(crate::runtime::ROOT_PID); let _ = try_with_runtime(|inner| request_shutdown_inner(inner, pid, from)); } // The core of `request_shutdown`. Reads the target's trap sender under its // cold lock (generation-verified), then acts outside the lock: a trap send // may unpark the receiver, and `request_stop_inner` re-takes the lock. pub(crate) fn request_shutdown_inner(inner: &RuntimeInner, pid: Pid, from: Pid) { request_shutdown_inner_probe(inner, pid, from); } /// [`request_shutdown_inner`], reporting what it found: `Some(true)` if the /// target was trapping (got the signal), `Some(false)` if it was stopped /// outright, `None` if there was nothing live at `pid`. pub(crate) fn request_shutdown_inner_probe( inner: &RuntimeInner, pid: Pid, from: Pid, ) -> Option { let trap = match inner.slot_at(pid) { Some(slot) => { let cold = slot.cold.lock(); if slot.generation() == pid.generation() { cold.actor.as_ref().map(|a| a.trap.clone()) } else { None // stale pid: nothing there to shut down } } None => None, }; match trap { Some(Some(tx)) => { let _ = tx.send(crate::link::ExitSignal { from, reason: crate::monitor::DownReason::Shutdown, }); Some(true) } Some(None) => { request_stop_inner(inner, pid); Some(false) } None => None, } } // --------------------------------------------------------------------------- // NoPreempt // --------------------------------------------------------------------------- /// A guard that disables preemption for its lifetime, restoring the /// previous setting on drop. Internal-use: application code has no need to /// disable preemption directly. See [`check!`](crate::check) for the /// user-facing side of preemption. pub struct NoPreempt(bool); impl NoPreempt { pub fn enter() -> Self { NoPreempt(crate::preempt::preemption_swap(false)) } } impl Drop for NoPreempt { fn drop(&mut self) { crate::preempt::preemption_swap(self.0); } } // --------------------------------------------------------------------------- // sleep / insert_wait_timer // --------------------------------------------------------------------------- /// Suspend the calling actor for `duration`. Unlike /// [`std::thread::sleep`], this parks only the actor, not the underlying OS /// thread, so every other actor (including others sharing the same OS /// thread) keeps running normally while this one waits. /// /// ``` /// use smarm::{run, sleep}; /// use std::time::Duration; /// /// run(|| { /// sleep(Duration::from_millis(1)); /// }); /// ``` /// /// Panics if called outside an actor. pub fn sleep(duration: std::time::Duration) { let me = match current_pid() { Some(pid) => pid, None => panic!("sleep() called outside an actor"), }; let _np = NoPreempt::enter(); let epoch = begin_wait(); let deadline = crate::timer::deadline_from_now(duration); with_runtime(|inner| match inner.timers.lock() { Ok(mut timers) => timers.insert_sleep(deadline, me, epoch), Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"), }); park_current(); } /// Like [`sleep`], but the deadline is always measured in real (wall-clock) /// time. Ordinary code should use [`sleep`]; this variant exists for smarm's /// own profiling and measurement tooling, which can otherwise stretch or /// compress simulated time. Without that tooling active, `sleep_wall` and /// `sleep` behave identically. pub fn sleep_wall(duration: std::time::Duration) { let me = match current_pid() { Some(pid) => pid, None => panic!("sleep_wall() called outside an actor"), }; let _np = NoPreempt::enter(); let epoch = begin_wait(); let deadline = crate::timer::deadline_from_now(duration); with_runtime(|inner| match inner.timers.lock() { Ok(mut timers) => timers.insert_sleep_wall(deadline, me, epoch), Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"), }); park_current(); } // Building block for bounded waits elsewhere in the crate (Mutex::lock_timeout, // Receiver::recv_timeout, select_timeout): arm a timer that, on expiry, asks // `target` whether this particular wait is still pending and should be woken // with a timeout. Not part of the public API; application code wants // `sleep`, `send_after`, or one of the `*_timeout` methods instead. pub fn insert_wait_timer( deadline: std::time::Instant, pid: Pid, target: std::sync::Arc, epoch: u32, ) { with_runtime(|inner| match inner.timers.lock() { Ok(mut timers) => timers.insert( deadline, pid, crate::timer::Reason::WaitTimeout { target, epoch }, ), Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"), }); } // --------------------------------------------------------------------------- // send_after / cancel_timer: message-delivery timers, in the same spirit as // Erlang's `erlang:send_after/3`. Each schedules `msg` for delivery after a // delay and returns a TimerId; `cancel_timer` can call one off before it // fires. The destination is resolved when the timer actually fires, not when // it's armed, so a `Name`-addressed timer always reaches whoever currently // holds that name, even if the original holder has since restarted. If the // destination is gone by fire time, the message is silently dropped, exactly // as a live send to a dead address would be. // --------------------------------------------------------------------------- /// Deliver `msg` to the exact actor identified by `dest` after `after` has /// elapsed. Unlike [`send_after_named`], this targets one specific actor: if /// that actor is gone by the time the timer fires, the message is dropped /// (it is never redirected to a different actor, even one that inherited the /// same name). Returns a [`TimerId`](crate::timer::TimerId) you can pass to /// [`cancel_timer`] to call it off early. pub fn send_after( 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(deadline, dest.erase(), fire), Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"), }) } /// Deliver `msg` to whichever actor holds the name `dest` when the timer /// fires, not necessarily whoever holds it now: if the named actor restarts /// (for example under a supervisor) and re-registers before the deadline, /// the message reaches the new instance. Returns a /// [`TimerId`](crate::timer::TimerId) you can pass to [`cancel_timer`]. pub fn send_after_named( 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(deadline, armed_by, fire), Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"), }) } /// Like [`send_after`], but the deadline is always measured in real /// (wall-clock) time rather than being subject to smarm's own profiling and /// measurement tooling. Use this for deadlines that need to reflect the /// outside world (a protocol timeout, a wall-clock schedule) rather than /// simulated workload pacing. Without that tooling active, it behaves /// identically to [`send_after`]. 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-clock-anchored [`send_after_named`]: re-resolving name delivery, /// with the same real-time deadline guarantee as [`send_after_wall`]. 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 already holds, after a delay, // rather than resolving a registry address at fire time. Used internally by // gen_server's timer support, which needs the fire to land on the server // loop's own dedicated channel instead of its public inbox. Otherwise // identical to `send_after`: same cancellation via `cancel_timer`, and a // send to a channel whose receiver is gone is silently dropped. pub(crate) fn send_after_to( after: std::time::Duration, tx: Sender, msg: T, ) -> crate::timer::TimerId { let deadline = crate::timer::deadline_from_now(after); let armed_by = current_pid().unwrap_or(Pid::new(0, 0)); let fire = Box::new(move || { let _ = tx.send(msg); }); with_runtime(|inner| match inner.timers.lock() { Ok(mut timers) => timers.insert_send(deadline, armed_by, fire), Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"), }) } /// Call off a timer armed by [`send_after`] or [`send_after_named`] before /// it fires. Returns `true` if the timer was still pending and delivery is /// now prevented, `false` if it had already fired or was already cancelled. pub fn cancel_timer(id: crate::timer::TimerId) -> bool { with_runtime(|inner| match inner.timers.lock() { Ok(mut timers) => timers.cancel(id), Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"), }) } // --------------------------------------------------------------------------- // block_on_io / wait_readable / wait_writable / read / write // --------------------------------------------------------------------------- /// Run a blocking closure (blocking file IO, a synchronous library call, /// anything that isn't itself actor-aware) on a dedicated worker thread, /// while the calling actor parks and every other actor keeps making /// progress. When `f` completes, the calling actor resumes with its result. /// /// Reach for this whenever you need to call into something that would /// otherwise block the underlying OS thread outright: a blocking C library, /// a synchronous filesystem call, DNS resolution via the system resolver, /// and so on. For plain readiness-based network IO on a file descriptor, /// [`wait_readable`] / [`wait_writable`] are cheaper since they don't need a /// dedicated thread. /// /// If `f` panics, that panic is carried over and re-raised in the calling /// actor, exactly as if the call had been made inline. /// /// Panics if called outside an actor. pub fn block_on_io(f: F) -> T where F: FnOnce() -> T + Send + 'static, T: Send + 'static, { let me = match current_pid() { Some(pid) => pid, None => panic!("block_on_io() called outside an actor"), }; let work: Box crate::io::IoResult + Send> = Box::new(move || { let v: T = f(); Ok(Box::new(v) as Box) }); { let _np = NoPreempt::enter(); let epoch = begin_wait(); with_runtime(|inner| { let mut io = match inner.io.lock() { Ok(io) => io, Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"), }; match io.as_mut() { Some(io) => { // RFC 018: count the op in flight BEFORE submit — the // pool decrements on completion, and an increment that // trailed the completion would underflow. Under the io // lock, so ordered against the same-lock submit. inner.io_outstanding.fetch_add(1, Ordering::AcqRel); io.submit(me, epoch, work); } None => panic!("io thread not started"), } }); park_current(); } let result = with_runtime(|inner| { let slot = match inner.slot_at(me) { Some(slot) => slot, None => panic!("block_on_io: own slot vanished"), }; let mut cold = slot.cold.lock(); debug_assert_eq!( slot.generation(), me.generation(), "block_on_io: own slot reused mid-park" ); match cold.pending_io_result.take() { Some(result) => result, None => panic!("block_on_io: resumed without a result"), } }); match result { Ok(any) => match any.downcast::() { Ok(typed) => *typed, Err(_) => panic!("block_on_io: type mismatch"), }, Err(payload) => std::panic::resume_unwind(payload), } } /// Park the calling actor until `fd` becomes readable. Other actors keep /// running while you wait; when the kernel reports the fd ready, the actor /// resumes and you perform the actual `read(2)` yourself (see [`read`] for /// a convenience wrapper that does both steps). /// /// Only one actor may wait on a given fd for a given direction at a time. /// Panics if called outside an actor. pub fn wait_readable(fd: std::os::fd::RawFd) -> std::io::Result<()> { wait_fd(fd, true, false) } /// Park the calling actor until `fd` becomes writable. See [`wait_readable`] /// for the read-side counterpart; the same notes apply. pub fn wait_writable(fd: std::os::fd::RawFd) -> std::io::Result<()> { wait_fd(fd, false, true) } fn wait_fd(fd: std::os::fd::RawFd, readable: bool, writable: bool) -> std::io::Result<()> { let me = match current_pid() { Some(pid) => pid, None => panic!("wait_*() called outside an actor"), }; let _np = NoPreempt::enter(); let epoch = begin_wait(); with_runtime(|inner| { let mut io = match inner.io.lock() { Ok(io) => io, Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"), }; match io.as_mut() { Some(io) => { // RFC 018: count the waiter BEFORE the ADD (mirror of // submit); roll back if the registration fails so a // rejected wait leaves the verdict counters clean. inner.io_fd_waiters.fetch_add(1, Ordering::AcqRel); let r = io.epoll_register(fd, me, epoch, readable, writable); if r.is_err() { inner.io_fd_waiters.fetch_sub(1, Ordering::AcqRel); } r } None => panic!("io thread not started"), } })?; // If a terminal stop unwinds us out of the park below, the registration // must not outlive us: a stale registration would fail every future // `wait_*` on this fd, and the kernel-side registration would leak until // the fd happens to be reused. Clean up only if the entry is still // ours; a wakeup racing the stop may have already consumed it, possibly // leaving a different actor's fresh registration in its place, which // must not be disturbed. struct Dereg { fd: std::os::fd::RawFd, me: Pid, epoch: u32, } impl Drop for Dereg { fn drop(&mut self) { with_runtime(|inner| { let mut io = match inner.io.lock() { Ok(io) => io, Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"), }; if let Some(io) = io.as_mut() { // `cancel_waiter` removes + DELs iff still ours, all // under the waiters lock (the ADD/DEL serialization); // decrement only when we actually removed it — a // FdReady that consumed it already did the decrement. if io.cancel_waiter(self.fd, self.me, self.epoch) { inner.io_fd_waiters.fetch_sub(1, Ordering::AcqRel); } } }); } } let guard = Dereg { fd, me, epoch }; park_current(); // Normal wake: the ready-fd path already removed the entry and // deregistered from epoll before waking us, so the guard's check on drop // is a guaranteed no-op. Skip it on this hot path (Dereg owns no // resource itself, so forgetting it leaks nothing). std::mem::forget(guard); Ok(()) } // --------------------------------------------------------------------------- // FdArm: fd readiness as a select arm // --------------------------------------------------------------------------- /// A file-descriptor readiness condition usable as an arm of /// [`select`](crate::select) / [`select_timeout`](crate::select_timeout), /// so you can wait on "this fd is readable" alongside ordinary channel /// receivers in the same call. Build one with [`FdArm::readable`] or /// [`FdArm::writable`]. /// /// Only one actor may wait on a given fd for a given direction at a time. A /// single fd open in both directions needs two `FdArm`s (or `dup` the fd) if /// you want to wait on both; you cannot wait on read and write readiness /// with one arm. pub struct FdArm { fd: std::os::fd::RawFd, readable: bool, writable: bool, } impl FdArm { /// An arm that becomes ready when `fd` is readable. pub fn readable(fd: std::os::fd::RawFd) -> Self { FdArm { fd, readable: true, writable: false, } } /// An arm that becomes ready when `fd` is writable. pub fn writable(fd: std::os::fd::RawFd) -> Self { FdArm { fd, readable: false, writable: true, } } } impl crate::channel::sealed::Sealed for FdArm {} impl crate::channel::Selectable for FdArm { // Ready-now check is a zero-timeout poll(2); if the requested events are // already pending, the wait is retired without registering. Otherwise // register with the IO thread. Any registration failure (a closed fd, // too many fds registered, a second waiter already on this fd) is // surfaced as an error rather than silently treated as "always ready", // so a caller never spins on an fd that genuinely can't be registered. fn sel_register(&self, pid: Pid, epoch: u32) -> std::io::Result { if poll_events(self.fd, self.readable, self.writable)? { return Ok(false); } with_runtime(|inner| { let mut io = match inner.io.lock() { Ok(io) => io, Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"), }; match io.as_mut() { Some(io) => { inner.io_fd_waiters.fetch_add(1, Ordering::AcqRel); let r = io.epoll_register(self.fd, pid, epoch, self.readable, self.writable); if r.is_err() { inner.io_fd_waiters.fetch_sub(1, Ordering::AcqRel); } r } None => panic!("io thread not started"), } })?; Ok(true) } // Same zero-timeout poll used for classification: a pure function of // fd state. An error here (fd closed mid-wait) is reported as ready // rather than pending forever: the caller's own read/write then // surfaces the real error, so a dead fd is an observable event, not a // silent hang. fn sel_ready(&self) -> bool { poll_events(self.fd, self.readable, self.writable).unwrap_or(true) } // Mirrors wait_fd's cleanup guard: remove the registration only if it's // still ours. A wakeup racing this cleanup may have already consumed // it, possibly leaving a different actor's fresh registration on the // same fd, which must not be touched. fn sel_unregister(&self, pid: Pid, epoch: u32) { with_runtime(|inner| { let mut io = match inner.io.lock() { Ok(io) => io, Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"), }; if let Some(io) = io.as_mut() { if io.cancel_waiter(self.fd, pid, epoch) { inner.io_fd_waiters.fetch_sub(1, Ordering::AcqRel); } } }); } fn sel_eager_cleanup(&self) -> bool { true } } // Zero-timeout poll(2): are any of the requested events (or an error/hangup // condition, so the caller's read/write fails loudly instead of parking // forever) pending on `fd` right now? fn poll_events(fd: std::os::fd::RawFd, readable: bool, writable: bool) -> std::io::Result { let mut events: libc::c_short = 0; if readable { events |= libc::POLLIN; } if writable { events |= libc::POLLOUT; } let mut pfd = libc::pollfd { fd, events, revents: 0, }; loop { let r = unsafe { libc::poll(&mut pfd, 1, 0) }; if r < 0 { let e = std::io::Error::last_os_error(); if e.kind() == std::io::ErrorKind::Interrupted { continue; } return Err(e); } if r == 0 { return Ok(false); } if pfd.revents & libc::POLLNVAL != 0 { return Err(std::io::Error::from_raw_os_error(libc::EBADF)); } return Ok(pfd.revents & (events | libc::POLLERR | libc::POLLHUP) != 0); } } /// Like [`wait_readable`], but gives up after `timeout` instead of waiting /// indefinitely: `Ok(true)` means the fd became ready, `Ok(false)` means the /// timeout elapsed first. pub fn wait_readable_timeout( fd: std::os::fd::RawFd, timeout: std::time::Duration, ) -> std::io::Result { let arm = FdArm::readable(fd); Ok(crate::channel::try_select_timeout(&[&arm], timeout)?.is_some()) } /// Like [`wait_writable`], but gives up after `timeout` instead of waiting /// indefinitely: `Ok(true)` means the fd became ready, `Ok(false)` means the /// timeout elapsed first. pub fn wait_writable_timeout( fd: std::os::fd::RawFd, timeout: std::time::Duration, ) -> std::io::Result { let arm = FdArm::writable(fd); Ok(crate::channel::try_select_timeout(&[&arm], timeout)?.is_some()) } /// Convenience wrapper: park until `fd` is readable, then perform the /// `read(2)` into `buf`. Equivalent to calling [`wait_readable`] yourself /// followed by a raw read, provided as a shorthand for the common case. pub fn read(fd: std::os::fd::RawFd, buf: &mut [u8]) -> std::io::Result { wait_readable(fd)?; let n = unsafe { libc::read(fd, buf.as_mut_ptr() as *mut _, buf.len()) }; if n < 0 { Err(std::io::Error::last_os_error()) } else { Ok(n as usize) } } /// Convenience wrapper: park until `fd` is writable, then perform the /// `write(2)` of `buf`. Equivalent to calling [`wait_writable`] yourself /// followed by a raw write, provided as a shorthand for the common case. pub fn write(fd: std::os::fd::RawFd, buf: &[u8]) -> std::io::Result { wait_writable(fd)?; let n = unsafe { libc::write(fd, buf.as_ptr() as *const _, buf.len()) }; if n < 0 { Err(std::io::Error::last_os_error()) } else { Ok(n as usize) } } // --------------------------------------------------------------------------- // register_supervisor_channel: internal wiring used by supervisor.rs // --------------------------------------------------------------------------- pub fn register_supervisor_channel(pid: Pid, sender: Sender) { with_runtime(|inner| { let slot = inner .slot_at(pid) .unwrap_or_else(|| panic!("register_supervisor_channel: pid {:?} not found", pid)); let mut cold = slot.cold.lock(); assert_eq!( slot.generation(), pid.generation(), "register_supervisor_channel: pid {:?} not found", pid ); cold.supervisor_channel = Some(sender); }); } // --------------------------------------------------------------------------- // run(): the single-threaded convenience entry point // --------------------------------------------------------------------------- /// Start the smarm runtime on a single OS thread, run `f` as the first /// (root) actor, and block until every actor it transitively spawned has /// finished. /// /// This is the simplest way to get started, and is all you need for most /// programs: `f` typically calls [`spawn`] to create more actors and waits /// on their [`JoinHandle`]s. If you want smarm to schedule actors across /// multiple OS threads instead, use /// [`crate::runtime::init`] with a /// [`Config`](crate::runtime::Config) that requests more than one scheduler /// thread, then call [`Runtime::run`](crate::runtime::Runtime::run) on it; /// `run` here is exactly that, pinned to one thread. pub fn run(f: F) { crate::runtime::init(crate::runtime::Config::exact(1)).run(f); } #[cfg(all(test, not(loom)))] mod send_after_to_tests { use super::*; use crate::channel::channel; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; use std::time::Duration; // A fired send_after_to lands its message on the caller's own channel. #[test] fn delivers_onto_the_channel() { run(|| { let (tx, rx) = channel::(); send_after_to(Duration::from_millis(10), tx, 42); // recv parks until the scheduler fires the timer thunk. assert_eq!(rx.recv().unwrap(), 42); }); } // cancel_timer before the deadline prevents delivery and reports the race // win; the channel then closes with no message once the lone sender (moved // into the now-discarded thunk) is gone. #[test] fn cancel_prevents_delivery() { run(|| { let (tx, rx) = channel::(); let id = send_after_to(Duration::from_millis(50), tx, 7); assert!(cancel_timer(id), "cancel before fire should win the race"); // No delivery: the discarded thunk drops the only sender, so recv // sees a closed channel rather than the value. assert!(rx.recv().is_err()); }); } // The thunk runs on the scheduler thread; a closed receiver makes the send // a harmless no-op (Erlang send_after semantics) rather than a panic. #[test] fn send_to_closed_channel_is_harmless() { let reached = Arc::new(AtomicBool::new(false)); let r2 = reached.clone(); run(move || { let (tx, rx) = channel::(); send_after_to(Duration::from_millis(10), tx, 1); drop(rx); // receiver gone before the timer fires crate::sleep(Duration::from_millis(30)); r2.store(true, Ordering::SeqCst); }); assert!( reached.load(Ordering::SeqCst), "runtime survived the dead-channel fire" ); } }