From 393cdd01f4b1381a97418b77772fd38d5331a29e Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 11 Jun 2026 11:28:58 +0000 Subject: [PATCH] feat(select): fd arms in select + timed fd waits (RFC 008 phase 1) FdArm composes fd readiness with channel arms on one wait epoch. Selectable grows fallible sel_register and an eager-cleanup hook; losing/stop-unwound/timed-out fd arms are unregistered (waiters entry + kernel ONESHOT) so the fd is never poisoned. try_select / try_select_timeout surface registration errors (EBADF, EMFILE, AlreadyExists) instead of RFC 008's permanently-ready lean, which would busy-loop a healthy-but-unregistrable fd; select/select_timeout stay infallible for channel-only arms. Adds wait_readable_timeout / wait_writable_timeout as one-arm selects. Known benign race (pre-existing, slightly widened): a queued FdReady racing the cleanup DEL can spuriously wake a fresh waiter on that fd; absorbed by select's defensive re-loop. Fixable by epoch-stamping completions. --- src/channel.rs | 179 +++++++++++++++++--- src/lib.rs | 8 +- src/scheduler.rs | 132 +++++++++++++++ tests/fd_select.rs | 403 +++++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 694 insertions(+), 28 deletions(-) create mode 100644 tests/fd_select.rs diff --git a/src/channel.rs b/src/channel.rs index 491a5b7..29ef06a 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -361,7 +361,7 @@ impl crate::timer::TimerTarget for RawMutex> { // select — ready-index wait over multiple receivers // --------------------------------------------------------------------------- -mod sealed { +pub(crate) mod sealed { pub trait Sealed {} } impl sealed::Sealed for Receiver {} @@ -371,28 +371,46 @@ impl sealed::Sealed for Receiver {} /// /// Contract (all under the arm's own lock): `sel_register` checks-or- /// registers atomically — if the arm is ready it does NOT register and -/// returns `false`; otherwise it publishes `(pid, epoch)` where its wakers -/// will find it. "Ready" means a receive would not park: a message is -/// queued, or the arm is closed. +/// returns `Ok(false)`; otherwise it publishes `(pid, epoch)` where its +/// wakers will find it and returns `Ok(true)`. "Ready" means a receive +/// would not park: a message is queued, or the arm is closed. `Err` means +/// the arm could not register at all (only fd arms can fail; channel +/// registration is infallible) — the wait must be retired and earlier +/// eager-cleanup arms unregistered. pub trait Selectable: sealed::Sealed { #[doc(hidden)] - fn sel_register(&self, pid: Pid, epoch: u32) -> bool; + fn sel_register(&self, pid: Pid, epoch: u32) -> std::io::Result; #[doc(hidden)] fn sel_ready(&self) -> bool; + /// Remove this arm's `(pid, epoch)` registration if — and only if — it + /// is still in place. Default no-op: a losing channel arm's stale + /// registration is inert (its wakers die at the epoch CAS; the next + /// wait overwrites the slot). Fd arms override this: their staleness + /// poisons the fd (waiters entry + kernel-side ONESHOT registration) + /// and needs an eager cleanup pass. + #[doc(hidden)] + fn sel_unregister(&self, _pid: Pid, _epoch: u32) {} + /// Whether this arm requires the eager cleanup pass at all. Gates the + /// post-wake `sel_unregister` sweep so channel-only selects keep + /// today's zero-cancellation hot path. + #[doc(hidden)] + fn sel_eager_cleanup(&self) -> bool { + false + } } impl Selectable for Receiver { - fn sel_register(&self, pid: Pid, epoch: u32) -> bool { + fn sel_register(&self, pid: Pid, epoch: u32) -> std::io::Result { let mut g = self.inner.lock(); if !g.queue.is_empty() || g.senders == 0 { - return false; + return Ok(false); } debug_assert!( g.parked_receiver.is_none_or(|(p, _)| p == pid), "channel has more than one receiver" ); g.parked_receiver = Some((pid, epoch)); - true + Ok(true) } fn sel_ready(&self) -> bool { @@ -429,30 +447,95 @@ impl Selectable for Receiver { /// self-clean at their wakers' failed CAS, or get overwritten by this /// receiver's next wait on that channel. /// -/// Panics if `arms` is empty, or when called outside an actor. +/// Panics if `arms` is empty, when called outside an actor, or if an fd +/// arm fails to register (EBADF, EMFILE, a second waiter on one fd — +/// see [`try_select`] for the fallible form; channel-only selects cannot +/// fail). pub fn select(arms: &[&dyn Selectable]) -> usize { + try_select(arms).expect("select(): fd arm failed to register (use try_select)") +} + +/// [`select`], fallible: `Err` when an arm fails to register (only fd +/// arms can — EBADF, EMFILE on the epoll set, or a second waiter on an +/// fd that already has one). On `Err` the wait is fully retired and no +/// registration is left behind: every arm registered before the failing +/// one has been unregistered. +pub fn try_select(arms: &[&dyn Selectable]) -> std::io::Result { assert!(!arms.is_empty(), "select() on an empty arm list"); let me = crate::actor::current_pid().expect("select() called outside an actor"); loop { let epoch = crate::scheduler::begin_wait(); - if let Some(i) = register_arms(me, epoch, arms) { - return i; + if let Some(i) = register_arms(me, epoch, arms)? { + return Ok(i); } + // Stale fd registrations are not harmless (a losing fd arm's + // waiters entry poisons the fd with AlreadyExists and its + // kernel-side ONESHOT registration can fire arbitrarily late), so + // selects containing fd arms run an eager cleanup pass after the + // park — including when a terminal stop unwinds out of it, via + // the guard. Channel-only selects skip all of it: `eager` is + // false, the guard is disarmed, and the loser-arm self-cleaning + // story is unchanged. + let eager = arms.iter().any(|a| a.sel_eager_cleanup()); + let mut guard = UnregisterGuard { arms, me, epoch, armed: eager }; + crate::scheduler::park_current(); + + if eager { + unregister_arms(arms, me, epoch); + } + guard.armed = false; + drop(guard); + // Woken precisely: an arm's send (message) or last-sender drop // (closure) consumed our epoch, and both leave their arm ready — // return the first one, in priority order (which may be a // different, higher-priority arm than the one that woke us; its // message stays queued and re-reports ready on the next call). + // Fd arms classify by a fresh zero-timeout poll, so they too are + // a pure function of state — independent of the registration the + // cleanup pass just removed. for (i, arm) in arms.iter().enumerate() { if arm.sel_ready() { - return i; + return Ok(i); } } // Unreachable by protocol (a stop wake unwinds out of // park_current). Defensive: re-open the wait and re-register — - // stale own-registrations are overwritten. + // stale own-registrations are overwritten (channels) or were + // removed by the cleanup pass above (fds). + } +} + +/// Eager-cleanup sweep: remove every fd arm's registration that is still +/// ours. No-op per channel arm (one virtual call); one io-lock visit per +/// fd arm. +fn unregister_arms(arms: &[&dyn Selectable], me: Pid, epoch: u32) { + for arm in arms { + if arm.sel_eager_cleanup() { + arm.sel_unregister(me, epoch); + } + } +} + +/// Stop-unwind twin of the explicit cleanup pass: a terminal stop unwinds +/// out of `park_current`, and a registered fd arm must not outlive its +/// actor (the generalization of `wait_fd`'s `Dereg`). Disarmed on the +/// normal path after the explicit pass runs; never armed when no fd arm +/// registered, keeping the channel-only path guard-free in effect. +struct UnregisterGuard<'a> { + arms: &'a [&'a dyn Selectable], + me: Pid, + epoch: u32, + armed: bool, +} + +impl Drop for UnregisterGuard<'_> { + fn drop(&mut self) { + if self.armed { + unregister_arms(self.arms, self.me, self.epoch); + } } } @@ -462,19 +545,35 @@ pub fn select(arms: &[&dyn Selectable]) -> usize { /// right after its registration wakes the caller through the protocol (the /// prep-to-park window is closed by RunningNotified). /// -/// `Some(i)` = arm `i` was ready, the pass stopped, and the wait has been -/// RETIRED (no park may follow): earlier arms hold live-epoch -/// registrations, so the epoch is bumped, a landed notification eaten, and -/// a pending stop re-observed — without which a stale arm wake could fault -/// a later one-shot park. `None` = every arm registered; the caller parks. -fn register_arms(me: Pid, epoch: u32, arms: &[&dyn Selectable]) -> Option { +/// `Ok(Some(i))` = arm `i` was ready, the pass stopped, and the wait has +/// been RETIRED (no park may follow): earlier arms hold live-epoch +/// registrations, so earlier *fd* arms are unregistered eagerly, then the +/// epoch is bumped, a landed notification eaten, and a pending stop +/// re-observed — without which a stale arm wake could fault a later +/// one-shot park. `Err` = an arm failed to register; identical unwind +/// (earlier fd arms unregistered, wait retired). `Ok(None)` = every arm +/// registered; the caller parks. +fn register_arms( + me: Pid, + epoch: u32, + arms: &[&dyn Selectable], +) -> std::io::Result> { for (i, arm) in arms.iter().enumerate() { - if !arm.sel_register(me, epoch) { + let registered = match arm.sel_register(me, epoch) { + Ok(r) => r, + Err(e) => { + unregister_arms(&arms[..i], me, epoch); + crate::scheduler::retire_wait(); + return Err(e); + } + }; + if !registered { + unregister_arms(&arms[..i], me, epoch); crate::scheduler::retire_wait(); - return Some(i); + return Ok(Some(i)); } } - None + Ok(None) } /// The [`select_timeout`] timer target: stateless, because precise wakes @@ -506,16 +605,29 @@ impl crate::timer::TimerTarget for SelectTimeout { /// `Duration::ZERO` is a valid timeout: it parks until the immediately-due /// timer is drained, then reports `None` unless an arm was already ready. /// -/// Panics if `arms` is empty, or when called outside an actor. +/// Panics if `arms` is empty, when called outside an actor, or if an fd +/// arm fails to register (see [`try_select_timeout`] for the fallible +/// form; channel-only selects cannot fail). pub fn select_timeout( arms: &[&dyn Selectable], timeout: std::time::Duration, ) -> Option { + try_select_timeout(arms, timeout) + .expect("select_timeout(): fd arm failed to register (use try_select_timeout)") +} + +/// [`select_timeout`], fallible: `Err` when an arm fails to register +/// (only fd arms can). On `Err` the wait is fully retired and no +/// registration — arm-side or kernel-side — is left behind. +pub fn try_select_timeout( + arms: &[&dyn Selectable], + timeout: std::time::Duration, +) -> std::io::Result> { assert!(!arms.is_empty(), "select_timeout() on an empty arm list"); let me = crate::actor::current_pid().expect("select_timeout() called outside an actor"); let epoch = crate::scheduler::begin_wait(); - if let Some(i) = register_arms(me, epoch, arms) { - return Some(i); // ready now: the timer was never armed + if let Some(i) = register_arms(me, epoch, arms)? { + return Ok(Some(i)); // ready now: the timer was never armed } // Arm the timer after the registration pass, outside every Channel @@ -524,7 +636,22 @@ pub fn select_timeout( let target: std::sync::Arc = std::sync::Arc::new(SelectTimeout); crate::scheduler::insert_wait_timer(deadline, me, target, epoch); + // Same eager-cleanup story as `try_select`: the timer arm needs none + // (stateless, stale entries die at the epoch CAS), channel arms need + // none, fd arms do — and a timer win in particular leaves every fd + // arm's registration behind, which without this pass would poison + // those fds until a kernel event happened to fire. + let eager = arms.iter().any(|a| a.sel_eager_cleanup()); + let mut guard = UnregisterGuard { arms, me, epoch, armed: eager }; + crate::scheduler::park_current(); + + if eager { + unregister_arms(arms, me, epoch); + } + guard.armed = false; + drop(guard); + // Woken precisely: an arm (ready below) or the timer (nothing ready). - arms.iter().position(|arm| arm.sel_ready()) + Ok(arms.iter().position(|arm| arm.sel_ready())) } diff --git a/src/lib.rs b/src/lib.rs index 56874f2..2ffb49c 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -45,7 +45,10 @@ static ALLOCATOR: preempt::PreemptingAllocator = preempt::PreemptingAllocator; // Public API re-exports // --------------------------------------------------------------------------- -pub use channel::{channel, select, select_timeout, Receiver, RecvError, RecvTimeoutError, Selectable, Sender}; +pub use channel::{ + channel, select, select_timeout, try_select, try_select_timeout, Receiver, RecvError, + RecvTimeoutError, Selectable, Sender, +}; pub use gen_server::{CallError, CallTimeoutError, CastError, GenServer, ServerBuilder, ServerCtx, ServerRef, Watcher}; pub use link::{link, trap_exit, unlink, ExitSignal}; pub use monitor::{demonitor, monitor, Down, DownReason, Monitor, MonitorId}; @@ -55,7 +58,8 @@ pub use registry::{name_of, register, unregister, whereis, RegisterError}; pub use runtime::{init, Config, Runtime}; pub use scheduler::{ block_on_io, request_stop, run, self_pid, sleep, spawn, spawn_under, wait_readable, - wait_writable, yield_now, JoinError, JoinHandle, + wait_readable_timeout, wait_writable, wait_writable_timeout, yield_now, FdArm, JoinError, + JoinHandle, }; pub use supervisor::{ChildSpec, OneForOne, Restart, Signal, Strategy}; diff --git a/src/scheduler.rs b/src/scheduler.rs index 1add4d3..583bdf9 100644 --- a/src/scheduler.rs +++ b/src/scheduler.rs @@ -432,6 +432,138 @@ fn wait_fd(fd: std::os::fd::RawFd, readable: bool, writable: bool) -> std::io::R Ok(()) } +// --------------------------------------------------------------------------- +// FdArm — fd readiness as a select arm (RFC 008) +// --------------------------------------------------------------------------- + +/// An fd-readiness arm for [`crate::select`] / [`crate::select_timeout`]: +/// ready when the fd is readable (resp. writable), composable with channel +/// receivers on one wait epoch. Phase-1 rules apply: one waiter per fd at a +/// time, one direction per arm (duplex on a single fd needs `dup`; epoll +/// registrations key on the open file description, so dup'd fds register +/// independently). +pub struct FdArm { + fd: std::os::fd::RawFd, + readable: bool, + writable: bool, +} + +impl FdArm { + pub fn readable(fd: std::os::fd::RawFd) -> Self { + FdArm { fd, readable: true, writable: false } + } + + 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 pending the wait is retired without registering (`Ok(false)`, + /// the channel-arm contract). Otherwise register with the io thread — + /// every failure surfaces as `Err` (EBADF including a closed-fd + /// POLLNVAL, EMFILE on the epoll set, AlreadyExists for a second + /// waiter on the fd): the fallible-out, nothing-left-behind rule, in + /// deviation from RFC 008's permanently-ready lean, which would spin a + /// consumer whose fd is healthy but unregistrable (EMFILE). + 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 = inner.io.lock().unwrap(); + io.as_mut() + .expect("io thread not started") + .epoll_register(self.fd, pid, epoch, self.readable, self.writable) + })?; + Ok(true) + } + + /// Classification is the same zero-timeout poll: a pure function of fd + /// state, independent of the registration the cleanup pass removed. An + /// error here (EBADF: fd closed mid-wait) reports READY — the + /// consumer's read/write surfaces the errno; a dead arm is an event, + /// not a hang. + fn sel_ready(&self) -> bool { + poll_events(self.fd, self.readable, self.writable).unwrap_or(true) + } + + /// `wait_fd`'s `Dereg` compare, verbatim: remove the waiters entry and + /// kernel-side registration iff the entry is still `(pid, epoch)`-ours. + /// A `FdReady` racing the wake may have consumed it (it removes + DELs + /// under the io lock), after which the fd may even carry ANOTHER + /// actor's fresh registration; in that case touch nothing. + fn sel_unregister(&self, pid: Pid, epoch: u32) { + with_runtime(|inner| { + let mut io = inner.io.lock().unwrap(); + if let Some(io) = io.as_mut() { + if io.waiters.get(&self.fd) == Some(&(pid, epoch)) { + io.waiters.remove(&self.fd); + io.epoll_deregister(self.fd); + } + } + }); + } + + fn sel_eager_cleanup(&self) -> bool { + true + } +} + +/// Zero-timeout `poll(2)`: are any of the requested events (or ERR/HUP, +/// which make the consumer's read/write fail loudly rather than park +/// forever) pending on `fd`? POLLNVAL maps to `Err(EBADF)`. +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); + } +} + +/// Wait until `fd` is readable or `timeout` elapses: `Ok(true)` = ready, +/// `Ok(false)` = timed out. A one-arm [`crate::try_select_timeout`]. +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()) +} + +/// Wait until `fd` is writable or `timeout` elapses: `Ok(true)` = ready, +/// `Ok(false)` = timed out. +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()) +} + 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()) }; diff --git a/tests/fd_select.rs b/tests/fd_select.rs new file mode 100644 index 0000000..38de6eb --- /dev/null +++ b/tests/fd_select.rs @@ -0,0 +1,403 @@ +//! RFC 008 — fd arms in select. Beyond the functional cases, the +//! *_stays_usable tests are the soundness probes for the one asymmetry the +//! RFC must close: a losing CHANNEL arm's stale registration is inert, but +//! a losing FD arm's registration (waiters entry + kernel ONESHOT) poisons +//! the fd with AlreadyExists until the eager cleanup pass removes it. Every +//! "loser" scenario therefore re-waits on the same fd afterwards and must +//! succeed — pre-cleanup, each of those re-waits errors or hangs. +//! +//! House pattern: actor panics are trampoline-caught and `run` returns +//! normally, so every test funnels its result into an outcome flag asserted +//! OUTSIDE `run` — an in-actor assertion alone passes vacuously. + +use smarm::{ + channel, run, select, select_timeout, spawn, try_select, wait_readable, + wait_readable_timeout, wait_writable_timeout, yield_now, FdArm, +}; +use std::os::fd::RawFd; +use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +// --------------------------------------------------------------------------- +// Pipe helper (as in io_epoll.rs) +// --------------------------------------------------------------------------- + +struct Pipe { + read: RawFd, + write: RawFd, +} + +impl Pipe { + fn new() -> Self { + let mut fds: [libc::c_int; 2] = [0; 2]; + let r = unsafe { libc::pipe2(fds.as_mut_ptr(), libc::O_CLOEXEC | libc::O_NONBLOCK) }; + assert_eq!(r, 0, "pipe2 failed"); + Pipe { read: fds[0], write: fds[1] } + } +} + +impl Drop for Pipe { + fn drop(&mut self) { + unsafe { + libc::close(self.read); + libc::close(self.write); + } + } +} + +fn raw_write(fd: RawFd, buf: &[u8]) -> isize { + unsafe { libc::write(fd, buf.as_ptr() as *const _, buf.len()) } +} + +fn raw_read(fd: RawFd, buf: &mut [u8]) -> isize { + unsafe { libc::read(fd, buf.as_mut_ptr() as *mut _, buf.len()) } +} + +fn flag() -> (Arc, Arc) { + let f = Arc::new(AtomicBool::new(false)); + (f.clone(), f) +} + +// --------------------------------------------------------------------------- +// Ready-now: data already pending retires the wait without parking. +// --------------------------------------------------------------------------- + +#[test] +fn fd_arm_ready_now_returns_without_parking() { + let (ok, ok2) = flag(); + run(move || { + let p = Pipe::new(); + assert_eq!(raw_write(p.write, b"x"), 1); + let (_tx, rx) = channel::(); + let fd_arm = FdArm::readable(p.read); + // fd arm at index 1: the ready-now path must also work for a + // non-first arm (and clean nothing — channel arms are inert). + let i = select(&[&rx, &fd_arm]); + assert_eq!(i, 1); + let mut buf = [0u8; 1]; + assert_eq!(raw_read(p.read, &mut buf), 1); + ok2.store(true, Ordering::SeqCst); + }); + assert!(ok.load(Ordering::SeqCst)); +} + +// --------------------------------------------------------------------------- +// Park-then-wake: fd arm wins against an idle channel arm. +// --------------------------------------------------------------------------- + +#[test] +fn fd_arm_parks_until_data_then_wins() { + let got = Arc::new(AtomicU32::new(0)); + let got2 = got.clone(); + run(move || { + let p = Pipe::new(); + let (rfd, wfd) = (p.read, p.write); + let (_tx_keepalive, rx) = channel::(); + + let h = spawn(move || { + let fd_arm = FdArm::readable(rfd); + let i = select(&[&fd_arm, &rx]); + assert_eq!(i, 0); + let mut buf = [0u8; 1]; + assert_eq!(raw_read(rfd, &mut buf), 1); + got2.store(buf[0] as u32, Ordering::SeqCst); + }); + yield_now(); // let it park + assert_eq!(raw_write(wfd, b"y"), 1); + let _ = h.join(); + }); + assert_eq!(got.load(Ordering::SeqCst), b'y' as u32); +} + +// --------------------------------------------------------------------------- +// THE asymmetry probe: channel arm wins, losing fd arm must be cleaned — +// the same actor (and the io thread) must be able to wait that fd again. +// --------------------------------------------------------------------------- + +#[test] +fn losing_fd_arm_is_cleaned_up_and_fd_stays_usable() { + let got = Arc::new(AtomicU32::new(0)); + let got2 = got.clone(); + run(move || { + let p = Pipe::new(); + let (rfd, wfd) = (p.read, p.write); + let (tx, rx) = channel::(); + + let h = spawn(move || { + let fd_arm = FdArm::readable(rfd); + // Channel wins; the fd arm registered and lost. + let i = select(&[&fd_arm, &rx]); + assert_eq!(i, 1); + assert_eq!(rx.try_recv().unwrap(), Some(7)); + + // Pre-cleanup this wait_readable fails AlreadyExists (the + // waiters entry is stale-ours) — the cleanup pass must have + // removed it. + wait_readable(rfd).unwrap(); + let mut buf = [0u8; 1]; + assert_eq!(raw_read(rfd, &mut buf), 1); + got2.store(buf[0] as u32, Ordering::SeqCst); + }); + yield_now(); // let it park in the select + tx.send(7).unwrap(); + yield_now(); // let it reach the second wait + assert_eq!(raw_write(wfd, b"z"), 1); + let _ = h.join(); + }); + assert_eq!(got.load(Ordering::SeqCst), b'z' as u32); +} + +// --------------------------------------------------------------------------- +// Ready-now on a LATER arm must unregister an earlier fd arm (the +// register_arms prefix-cleanup path: no park ever happens). +// --------------------------------------------------------------------------- + +#[test] +fn ready_now_later_arm_cleans_earlier_fd_arm() { + let (ok, ok2) = flag(); + run(move || { + let p = Pipe::new(); + let (rfd, wfd) = (p.read, p.write); + let (tx, rx) = channel::(); + tx.send(1).unwrap(); // arm 1 ready before the select + + let fd_arm = FdArm::readable(rfd); + let i = select(&[&fd_arm, &rx]); + assert_eq!(i, 1); + assert_eq!(rx.try_recv().unwrap(), Some(1)); + + // The fd arm registered (idle pipe), then arm 1 retired the wait. + // Its registration must have been removed in the same pass. + assert_eq!(raw_write(wfd, b"a"), 1); + wait_readable(rfd).unwrap(); + let mut buf = [0u8; 1]; + assert_eq!(raw_read(rfd, &mut buf), 1); + ok2.store(true, Ordering::SeqCst); + }); + assert!(ok.load(Ordering::SeqCst)); +} + +// --------------------------------------------------------------------------- +// Two fd arms in one select (phase-1: distinct fds, nothing special). +// --------------------------------------------------------------------------- + +#[test] +fn two_fd_arms_second_fires_first_stays_usable() { + let (ok, ok2) = flag(); + run(move || { + let pa = Pipe::new(); + let pb = Pipe::new(); + let (rfd_a, wfd_a) = (pa.read, pa.write); + let (rfd_b, wfd_b) = (pb.read, pb.write); + + let h = spawn(move || { + let a = FdArm::readable(rfd_a); + let b = FdArm::readable(rfd_b); + let i = select(&[&a, &b]); + assert_eq!(i, 1); + let mut buf = [0u8; 1]; + assert_eq!(raw_read(rfd_b, &mut buf), 1); + + // Arm a lost; its fd must be immediately re-waitable. + assert_eq!(raw_write(wfd_a, b"q"), 1); + wait_readable(rfd_a).unwrap(); + assert_eq!(raw_read(rfd_a, &mut buf), 1); + assert_eq!(buf[0], b'q'); + ok2.store(true, Ordering::SeqCst); + }); + yield_now(); + assert_eq!(raw_write(wfd_b, b"b"), 1); + let _ = h.join(); + }); + assert!(ok.load(Ordering::SeqCst)); +} + +// --------------------------------------------------------------------------- +// select_timeout: timer wins over an idle fd arm; the fd is left clean. +// --------------------------------------------------------------------------- + +#[test] +fn select_timeout_timer_beats_idle_fd_arm_and_fd_stays_usable() { + let (ok, ok2) = flag(); + run(move || { + let p = Pipe::new(); + let (rfd, wfd) = (p.read, p.write); + + let fd_arm = FdArm::readable(rfd); + let start = Instant::now(); + let r = select_timeout(&[&fd_arm], Duration::from_millis(30)); + assert!(r.is_none(), "idle fd must time out"); + assert!(start.elapsed() >= Duration::from_millis(30)); + + // Timer win is exactly the case where the fd arm's registration is + // left behind without an eager pass. + assert_eq!(raw_write(wfd, b"c"), 1); + wait_readable(rfd).unwrap(); + let mut buf = [0u8; 1]; + assert_eq!(raw_read(rfd, &mut buf), 1); + ok2.store(true, Ordering::SeqCst); + }); + assert!(ok.load(Ordering::SeqCst)); +} + +// --------------------------------------------------------------------------- +// Derived wrappers. +// --------------------------------------------------------------------------- + +#[test] +fn wait_readable_timeout_times_out_then_succeeds_with_data() { + let (ok, ok2) = flag(); + run(move || { + let p = Pipe::new(); + let (rfd, wfd) = (p.read, p.write); + + let start = Instant::now(); + assert_eq!(wait_readable_timeout(rfd, Duration::from_millis(30)).unwrap(), false); + assert!(start.elapsed() >= Duration::from_millis(30)); + + // Timed-out wait must leave the fd clean; ready path returns true. + assert_eq!(raw_write(wfd, b"d"), 1); + assert_eq!(wait_readable_timeout(rfd, Duration::from_secs(5)).unwrap(), true); + let mut buf = [0u8; 1]; + assert_eq!(raw_read(rfd, &mut buf), 1); + ok2.store(true, Ordering::SeqCst); + }); + assert!(ok.load(Ordering::SeqCst)); +} + +#[test] +fn wait_readable_timeout_wakes_on_late_data() { + let got = Arc::new(AtomicU32::new(0)); + let got2 = got.clone(); + run(move || { + let p = Pipe::new(); + let (rfd, wfd) = (p.read, p.write); + let h = spawn(move || { + assert_eq!(wait_readable_timeout(rfd, Duration::from_secs(5)).unwrap(), true); + let mut buf = [0u8; 1]; + assert_eq!(raw_read(rfd, &mut buf), 1); + got2.store(buf[0] as u32, Ordering::SeqCst); + }); + yield_now(); + assert_eq!(raw_write(wfd, b"e"), 1); + let _ = h.join(); + }); + assert_eq!(got.load(Ordering::SeqCst), b'e' as u32); +} + +#[test] +fn wait_writable_timeout_ready_now_on_empty_pipe() { + let (ok, ok2) = flag(); + run(move || { + let p = Pipe::new(); + // An empty pipe's write end is writable: ready-now path, no park. + assert_eq!(wait_writable_timeout(p.write, Duration::from_secs(5)).unwrap(), true); + ok2.store(true, Ordering::SeqCst); + }); + assert!(ok.load(Ordering::SeqCst)); +} + +// --------------------------------------------------------------------------- +// Error surface: registration failure is an Err from try_select, with the +// wait retired (the actor can immediately wait on something else). +// --------------------------------------------------------------------------- + +#[test] +fn try_select_surfaces_registration_error_and_retires_the_wait() { + let (ok, ok2) = flag(); + run(move || { + let bad: RawFd = { + let p = Pipe::new(); + p.read + }; // both ends closed by Drop: EBADF on registration + + let fd_arm = FdArm::readable(bad); + let err = try_select(&[&fd_arm]).unwrap_err(); + // EBADF, whether the pre-poll or epoll_ctl ADD reports it. + assert_eq!(err.raw_os_error(), Some(libc::EBADF)); + + // The wait was retired: a normal select right after works. + let (tx, rx) = channel::(); + tx.send(9).unwrap(); + let i = select(&[&rx]); + assert_eq!(i, 0); + assert_eq!(rx.try_recv().unwrap(), Some(9)); + ok2.store(true, Ordering::SeqCst); + }); + assert!(ok.load(Ordering::SeqCst)); +} + +// --------------------------------------------------------------------------- +// Stop-unwind: an actor stopped while parked in an fd-arm select must not +// poison the fd (the UnregisterGuard generalization of wait_fd's Dereg). +// Mirrors io_epoll.rs::stopped_waiter_does_not_poison_the_fd. +// --------------------------------------------------------------------------- + +#[test] +fn stopped_selector_does_not_poison_the_fd() { + let seen = Arc::new(AtomicU32::new(0)); + let seen_outer = seen.clone(); + run(move || { + let p = Pipe::new(); + let (rfd, wfd) = (p.read, p.write); + let (_tx_keepalive, rx) = channel::(); + + let h = spawn(move || { + let fd_arm = FdArm::readable(rfd); + select(&[&fd_arm, &rx]); + unreachable!("neither arm ever fires while this actor lives"); + }); + yield_now(); // let it reach the park + smarm::request_stop(h.pid()); + let _ = h.join(); // Ok(()): stopped, not panicked + + // Second waiter on the SAME fd must register and be woken. + let seen2 = seen.clone(); + let h2 = spawn(move || { + wait_readable(rfd).unwrap(); + let mut buf = [0u8; 1]; + assert_eq!(raw_read(rfd, &mut buf), 1); + seen2.store(buf[0] as u32, Ordering::SeqCst); + }); + yield_now(); + assert_eq!(raw_write(wfd, b"x"), 1); + let _ = h2.join(); + }); + assert_eq!(seen_outer.load(Ordering::SeqCst), b'x' as u32); +} + +// --------------------------------------------------------------------------- +// Phase-1 misuse: a second waiter on an fd that already has one is an Err +// (AlreadyExists), not a hang — and does not disturb the first waiter. +// --------------------------------------------------------------------------- + +#[test] +fn second_waiter_on_same_fd_errs_without_disturbing_the_first() { + let got = Arc::new(AtomicU32::new(0)); + let got2 = got.clone(); + let (ok, ok2) = flag(); + run(move || { + let p = Pipe::new(); + let (rfd, wfd) = (p.read, p.write); + + let h = spawn(move || { + wait_readable(rfd).unwrap(); + let mut buf = [0u8; 1]; + assert_eq!(raw_read(rfd, &mut buf), 1); + got2.store(buf[0] as u32, Ordering::SeqCst); + }); + yield_now(); // first waiter parked + + let fd_arm = FdArm::readable(rfd); + let err = try_select(&[&fd_arm]).unwrap_err(); + assert_eq!(err.kind(), std::io::ErrorKind::AlreadyExists); + + // First waiter still wakes normally. + assert_eq!(raw_write(wfd, b"w"), 1); + let _ = h.join(); + ok2.store(true, Ordering::SeqCst); + }); + assert_eq!(got.load(Ordering::SeqCst), b'w' as u32); + assert!(ok.load(Ordering::SeqCst)); +}