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.
This commit is contained in:
@@ -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<AtomicBool>, Arc<AtomicBool>) {
|
||||
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::<i64>();
|
||||
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::<i64>();
|
||||
|
||||
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::<i64>();
|
||||
|
||||
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::<i64>();
|
||||
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::<i64>();
|
||||
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::<i64>();
|
||||
|
||||
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));
|
||||
}
|
||||
Reference in New Issue
Block a user