Skip to main content

rustpython_host_env/
select.rs

1use core::mem::MaybeUninit;
2use core::time::Duration;
3use std::io;
4use std::time::Instant;
5
6#[cfg(unix)]
7pub use libc::{
8    EINTR, FD_SETSIZE, PIPE_BUF, POLLERR, POLLHUP, POLLIN, POLLNVAL, POLLOUT, POLLPRI, POLLRDBAND,
9    POLLRDNORM, POLLWRBAND, POLLWRNORM,
10};
11
12#[cfg(any(target_os = "linux", target_os = "android", target_os = "redox"))]
13pub use libc::{
14    EPOLL_CLOEXEC, EPOLLERR, EPOLLET, EPOLLEXCLUSIVE, EPOLLHUP, EPOLLIN, EPOLLMSG, EPOLLONESHOT,
15    EPOLLOUT, EPOLLPRI, EPOLLRDBAND, EPOLLRDHUP, EPOLLRDNORM, EPOLLWAKEUP, EPOLLWRBAND,
16    EPOLLWRNORM,
17};
18
19#[cfg(unix)]
20pub mod platform {
21    pub use libc::pollfd;
22    pub use libc::{FD_ISSET, FD_SET, FD_SETSIZE, FD_ZERO, fd_set, select, timeval};
23    use std::io;
24    pub use std::os::unix::io::RawFd;
25
26    #[must_use]
27    pub const fn check_err(x: i32) -> bool {
28        x < 0
29    }
30
31    pub fn last_select_error() -> io::Error {
32        io::Error::last_os_error()
33    }
34}
35
36#[allow(non_snake_case)]
37#[cfg(windows)]
38pub mod platform {
39    pub use WinSock::{FD_SET as fd_set, FD_SETSIZE, SOCKET as RawFd, TIMEVAL as timeval, select};
40    use std::io;
41    use windows_sys::Win32::Networking::WinSock;
42
43    /// # Safety
44    ///
45    /// `set` must be a valid mutable pointer to an initialized WinSock fd_set.
46    pub unsafe fn FD_SET(fd: RawFd, set: *mut fd_set) {
47        let mut slot = unsafe { (&raw mut (*set).fd_array).cast::<RawFd>() };
48        let fd_count = unsafe { (*set).fd_count };
49        for _ in 0..fd_count {
50            if unsafe { *slot } == fd {
51                return;
52            }
53            slot = unsafe { slot.add(1) };
54        }
55        if fd_count < FD_SETSIZE {
56            unsafe {
57                *slot = fd as RawFd;
58                (*set).fd_count += 1;
59            }
60        }
61    }
62
63    /// # Safety
64    ///
65    /// `set` must be a valid mutable pointer to a WinSock fd_set.
66    pub unsafe fn FD_ZERO(set: *mut fd_set) {
67        unsafe { (*set).fd_count = 0 };
68    }
69
70    /// # Safety
71    ///
72    /// `set` must be a valid mutable pointer to an initialized WinSock fd_set.
73    pub unsafe fn FD_ISSET(fd: RawFd, set: *mut fd_set) -> bool {
74        use WinSock::__WSAFDIsSet;
75        unsafe { __WSAFDIsSet(fd as _, set) != 0 }
76    }
77
78    #[must_use]
79    pub fn check_err(x: i32) -> bool {
80        x == WinSock::SOCKET_ERROR
81    }
82
83    pub fn last_select_error() -> io::Error {
84        io::Error::from_raw_os_error(unsafe { WinSock::WSAGetLastError() })
85    }
86}
87
88#[cfg(target_os = "wasi")]
89pub mod platform {
90    pub use libc::{FD_SETSIZE, timeval};
91    use std::io;
92    pub use std::os::fd::RawFd;
93
94    pub const fn check_err(x: i32) -> bool {
95        x < 0
96    }
97
98    #[repr(C)]
99    pub struct fd_set {
100        __nfds: usize,
101        __fds: [libc::c_int; FD_SETSIZE],
102    }
103
104    #[allow(non_snake_case)]
105    /// # Safety
106    ///
107    /// `set` must be a valid pointer to an initialized fd_set.
108    pub unsafe fn FD_ISSET(fd: RawFd, set: *const fd_set) -> bool {
109        let set = unsafe { &*set };
110        for p in &set.__fds[..set.__nfds] {
111            if *p == fd {
112                return true;
113            }
114        }
115        false
116    }
117
118    #[allow(non_snake_case)]
119    /// # Safety
120    ///
121    /// `set` must be a valid mutable pointer to an initialized fd_set.
122    pub unsafe fn FD_SET(fd: RawFd, set: *mut fd_set) {
123        let set = unsafe { &mut *set };
124        for p in &set.__fds[..set.__nfds] {
125            if *p == fd {
126                return;
127            }
128        }
129        let n = set.__nfds;
130        if n < FD_SETSIZE {
131            set.__fds[n] = fd;
132            set.__nfds = n + 1;
133        }
134    }
135
136    #[allow(non_snake_case)]
137    /// # Safety
138    ///
139    /// `set` must be a valid mutable pointer to an fd_set.
140    pub unsafe fn FD_ZERO(set: *mut fd_set) {
141        unsafe { (*set).__nfds = 0 };
142    }
143
144    unsafe extern "C" {
145        pub fn select(
146            nfds: libc::c_int,
147            readfds: *mut fd_set,
148            writefds: *mut fd_set,
149            errorfds: *mut fd_set,
150            timeout: *const timeval,
151        ) -> libc::c_int;
152    }
153
154    pub fn last_select_error() -> io::Error {
155        io::Error::last_os_error()
156    }
157}
158
159pub use platform::{RawFd, timeval};
160
161#[cfg(unix)]
162pub type PollFd = platform::pollfd;
163
164#[repr(transparent)]
165pub struct FdSet(MaybeUninit<platform::fd_set>);
166
167impl FdSet {
168    pub fn new() -> Self {
169        let mut fdset = MaybeUninit::zeroed();
170        unsafe { platform::FD_ZERO(fdset.as_mut_ptr()) };
171        Self(fdset)
172    }
173
174    pub fn insert(&mut self, fd: RawFd) {
175        unsafe { platform::FD_SET(fd, self.0.as_mut_ptr()) };
176    }
177
178    pub fn contains(&mut self, fd: RawFd) -> bool {
179        unsafe { platform::FD_ISSET(fd, self.0.as_mut_ptr()) }
180    }
181
182    pub fn clear(&mut self) {
183        unsafe { platform::FD_ZERO(self.0.as_mut_ptr()) };
184    }
185
186    pub fn highest(&mut self) -> Option<RawFd> {
187        (0..platform::FD_SETSIZE as RawFd)
188            .rev()
189            .find(|&fd| self.contains(fd))
190    }
191}
192
193impl Default for FdSet {
194    fn default() -> Self {
195        Self::new()
196    }
197}
198
199pub fn select(
200    nfds: libc::c_int,
201    readfds: &mut FdSet,
202    writefds: &mut FdSet,
203    errfds: &mut FdSet,
204    timeout: Option<&mut timeval>,
205) -> io::Result<i32> {
206    let timeout = match timeout {
207        Some(tv) => tv as *mut timeval,
208        None => core::ptr::null_mut(),
209    };
210    let ret = unsafe {
211        platform::select(
212            nfds,
213            readfds.0.as_mut_ptr(),
214            writefds.0.as_mut_ptr(),
215            errfds.0.as_mut_ptr(),
216            timeout,
217        )
218    };
219    if platform::check_err(ret) {
220        Err(platform::last_select_error())
221    } else {
222        Ok(ret)
223    }
224}
225
226pub fn sec_to_timeval(sec: f64) -> timeval {
227    timeval {
228        tv_sec: sec.trunc() as _,
229        tv_usec: (sec.fract() * 1e6) as _,
230    }
231}
232
233pub fn duration_to_timeval(d: Duration) -> timeval {
234    let mut tv = timeval {
235        tv_sec: 0 as _,
236        tv_usec: d.subsec_micros() as _,
237    };
238    tv.tv_sec = saturate_secs(d.as_secs(), tv.tv_sec);
239    tv
240}
241
242fn saturate_secs<T>(secs: u64, _sample: T) -> T
243where
244    T: TryFrom<u64> + TryFrom<i32>,
245{
246    T::try_from(secs).unwrap_or_else(|_| {
247        T::try_from(i32::MAX).unwrap_or_else(|_| {
248            T::try_from(0u64).unwrap_or_else(|_| unreachable!("tv_sec holds 0"))
249        })
250    })
251}
252
253#[derive(Copy, Clone, Debug, PartialEq, Eq)]
254pub enum WaitKind {
255    Read,
256    Write,
257}
258
259#[derive(Copy, Clone, Debug, PartialEq, Eq)]
260pub enum WaitFd {
261    Ready,
262    Timeout,
263}
264
265#[derive(Debug)]
266pub enum WaitFdError {
267    Interrupted,
268    Io(io::Error),
269}
270
271#[cfg(windows)]
272type WaitFdArg = std::os::windows::io::RawSocket;
273#[cfg(not(windows))]
274type WaitFdArg = RawFd;
275
276pub fn wait_fd(
277    fd: WaitFdArg,
278    kind: WaitKind,
279    deadline: Option<Instant>,
280) -> Result<WaitFd, WaitFdError> {
281    #[cfg(unix)]
282    {
283        wait_fd_poll(fd, kind, deadline)
284    }
285    #[cfg(windows)]
286    {
287        wait_fd_select(fd as RawFd, kind, deadline)
288    }
289    #[cfg(not(any(unix, windows)))]
290    {
291        wait_fd_select(fd, kind, deadline)
292    }
293}
294
295#[cfg(unix)]
296fn wait_fd_poll(
297    fd: RawFd,
298    kind: WaitKind,
299    deadline: Option<Instant>,
300) -> Result<WaitFd, WaitFdError> {
301    let events = match kind {
302        WaitKind::Read => POLLIN | POLLPRI,
303        WaitKind::Write => POLLOUT,
304    };
305    let mut fds = [PollFd {
306        fd,
307        events,
308        revents: 0,
309    }];
310    loop {
311        let (timeout, is_capped) = match deadline {
312            None => (-1, false),
313            Some(deadline) => {
314                match deadline.checked_duration_since(Instant::now()) {
315                    // Deadline already passed: still poll once with 0 so a
316                    // readable/writable socket is not reported as timed out
317                    // after scheduling delay between starting the timer and
318                    // entering poll.
319                    None => (0, false),
320                    Some(remaining) => match duration_as_millis_ceiling(remaining) {
321                        Some(ms) => (ms, false),
322                        None => (i32::MAX, true),
323                    },
324                }
325            }
326        };
327        match poll_fds(&mut fds, timeout) {
328            Ok(0) if is_capped => {}
329            Ok(0) => return Ok(WaitFd::Timeout),
330            Ok(_) if fds[0].revents & POLLNVAL != 0 => {
331                return Err(WaitFdError::Io(io::Error::from_raw_os_error(libc::EBADF)));
332            }
333            Ok(_) => return Ok(WaitFd::Ready),
334            Err(err) if err.kind() == io::ErrorKind::Interrupted => {
335                return Err(WaitFdError::Interrupted);
336            }
337            Err(err) => return Err(WaitFdError::Io(err)),
338        }
339    }
340}
341
342#[cfg(not(unix))]
343fn wait_fd_select(
344    fd: RawFd,
345    kind: WaitKind,
346    deadline: Option<Instant>,
347) -> Result<WaitFd, WaitFdError> {
348    let mut reads = FdSet::new();
349    let mut writes = FdSet::new();
350    let mut errs = FdSet::new();
351    match kind {
352        WaitKind::Read => {
353            reads.insert(fd);
354            errs.insert(fd);
355        }
356        WaitKind::Write => {
357            writes.insert(fd);
358            errs.insert(fd);
359        }
360    }
361    let mut timeout = match deadline {
362        None => None,
363        Some(deadline) => {
364            let remaining = deadline
365                .checked_duration_since(Instant::now())
366                .unwrap_or(Duration::ZERO);
367            Some(duration_to_timeval(remaining))
368        }
369    };
370    let nfds = cfg_select! {
371        windows => 0,
372        _ => fd.saturating_add(1) as libc::c_int,
373    };
374    match select(nfds, &mut reads, &mut writes, &mut errs, timeout.as_mut()) {
375        Ok(0) => Ok(WaitFd::Timeout),
376        Ok(_) => Ok(WaitFd::Ready),
377        Err(err) if err.kind() == io::ErrorKind::Interrupted => Err(WaitFdError::Interrupted),
378        Err(err) => Err(WaitFdError::Io(err)),
379    }
380}
381
382/// Convert a duration to a `poll(2)` millisecond timeout, rounding toward +∞.
383#[cfg(unix)]
384pub fn duration_as_millis_ceiling(d: Duration) -> Option<i32> {
385    let mut ms = d.as_millis();
386    if Duration::from_millis(ms.min(u128::from(u64::MAX)) as u64) < d {
387        ms = ms.saturating_add(1);
388    }
389    i32::try_from(ms).ok()
390}
391
392#[cfg(unix)]
393#[inline]
394pub fn search_poll_fd(fds: &[PollFd], fd: i32) -> Result<usize, usize> {
395    fds.binary_search_by_key(&fd, |pfd| pfd.fd)
396}
397
398#[cfg(unix)]
399pub fn insert_poll_fd(fds: &mut Vec<PollFd>, fd: i32, events: i16) {
400    match search_poll_fd(fds, fd) {
401        Ok(i) => fds[i].events = events,
402        Err(i) => fds.insert(
403            i,
404            PollFd {
405                fd,
406                events,
407                revents: 0,
408            },
409        ),
410    }
411}
412
413#[cfg(unix)]
414pub fn get_poll_fd_mut(fds: &mut [PollFd], fd: i32) -> Option<&mut PollFd> {
415    search_poll_fd(fds, fd).ok().map(move |i| &mut fds[i])
416}
417
418#[cfg(unix)]
419pub fn remove_poll_fd(fds: &mut Vec<PollFd>, fd: i32) -> Option<PollFd> {
420    search_poll_fd(fds, fd).ok().map(|i| fds.remove(i))
421}
422
423#[cfg(unix)]
424pub fn poll_fds(fds: &mut [PollFd], timeout: i32) -> std::io::Result<i32> {
425    let res = unsafe { libc::poll(fds.as_mut_ptr(), fds.len() as _, timeout) };
426    if res < 0 {
427        Err(std::io::Error::last_os_error())
428    } else {
429        Ok(res)
430    }
431}
432
433#[cfg(any(target_os = "linux", target_os = "android", target_os = "redox"))]
434pub mod epoll {
435    use std::os::fd::{AsFd, IntoRawFd, OwnedFd};
436
437    pub use rustix::event::Timespec;
438    pub use rustix::event::epoll::{Event, EventData, EventFlags};
439
440    #[derive(Debug)]
441    pub enum WaitError {
442        Interrupted,
443        Io(std::io::Error),
444    }
445
446    pub fn create() -> std::io::Result<OwnedFd> {
447        rustix::event::epoll::create(rustix::event::epoll::CreateFlags::CLOEXEC).map_err(Into::into)
448    }
449
450    pub fn close(fd: OwnedFd) -> nix::Result<()> {
451        nix::unistd::close(fd.into_raw_fd())
452    }
453
454    pub fn add<F: AsFd>(epoll: &OwnedFd, fd: F, data: u64, events: u32) -> std::io::Result<()> {
455        rustix::event::epoll::add(
456            epoll,
457            fd,
458            EventData::new_u64(data),
459            EventFlags::from_bits_retain(events),
460        )
461        .map_err(Into::into)
462    }
463
464    pub fn modify<F: AsFd>(epoll: &OwnedFd, fd: F, data: u64, events: u32) -> std::io::Result<()> {
465        rustix::event::epoll::modify(
466            epoll,
467            fd,
468            EventData::new_u64(data),
469            EventFlags::from_bits_retain(events),
470        )
471        .map_err(Into::into)
472    }
473
474    pub fn delete<F: AsFd>(epoll: &OwnedFd, fd: F) -> std::io::Result<()> {
475        rustix::event::epoll::delete(epoll, fd).map_err(Into::into)
476    }
477
478    pub fn wait(
479        epoll: &OwnedFd,
480        events: &mut Vec<Event>,
481        timeout: Option<&Timespec>,
482    ) -> Result<usize, WaitError> {
483        events.clear();
484        match rustix::event::epoll::wait(epoll, rustix::buffer::spare_capacity(events), timeout) {
485            Ok(n) => {
486                unsafe { events.set_len(n) };
487                Ok(n)
488            }
489            Err(rustix::io::Errno::INTR) => Err(WaitError::Interrupted),
490            Err(err) => Err(WaitError::Io(err.into())),
491        }
492    }
493}
494
495#[cfg(any(
496    target_os = "macos",
497    target_os = "ios",
498    target_os = "freebsd",
499    target_os = "netbsd",
500    target_os = "openbsd",
501    target_os = "dragonfly",
502))]
503pub mod kqueue {
504    use alloc::sync::{Arc, Weak};
505    use core::sync::atomic::{AtomicI32, Ordering};
506    use core::time::Duration;
507    use parking_lot::Mutex;
508    use std::io;
509    use std::os::fd::BorrowedFd;
510
511    pub use libc::{
512        EV_ADD, EV_CLEAR, EV_DELETE, EV_DISABLE, EV_ENABLE, EV_EOF, EV_ERROR, EV_FLAG1, EV_ONESHOT,
513        EV_SYSFLAGS, EVFILT_AIO, EVFILT_PROC, EVFILT_READ, EVFILT_SIGNAL, EVFILT_TIMER,
514        EVFILT_VNODE, EVFILT_WRITE, NOTE_ATTRIB, NOTE_CHILD, NOTE_DELETE, NOTE_EXEC, NOTE_EXIT,
515        NOTE_EXTEND, NOTE_FORK, NOTE_LINK, NOTE_LOWAT, NOTE_PCTRLMASK, NOTE_PDATAMASK, NOTE_RENAME,
516        NOTE_REVOKE, NOTE_TRACK, NOTE_TRACKERR, NOTE_WRITE,
517    };
518
519    // NetBSD widths differ from i16/u16; the cast is a no-op elsewhere.
520    #[allow(clippy::unnecessary_cast)]
521    pub const DEFAULT_FILTER: i16 = EVFILT_READ as i16;
522    #[allow(clippy::unnecessary_cast)]
523    pub const DEFAULT_FLAGS: u16 = EV_ADD as u16;
524
525    #[derive(Copy, Clone, Debug)]
526    pub struct Timespec {
527        pub sec: i64,
528        pub nsec: i64,
529    }
530
531    impl Timespec {
532        pub fn from_secs(secs: f64) -> Option<Self> {
533            if !secs.is_finite() || secs > i64::MAX as f64 {
534                return None;
535            }
536            let mut sec = secs.trunc() as i64;
537            let mut nsec = ((secs - sec as f64) * 1e9).round() as i64;
538            if nsec >= 1_000_000_000 {
539                sec = sec.saturating_add(1);
540                nsec -= 1_000_000_000;
541            }
542            Some(Self { sec, nsec })
543        }
544
545        pub fn from_duration(d: Duration) -> Self {
546            Self {
547                sec: d.as_secs() as i64,
548                nsec: i64::from(d.subsec_nanos()),
549            }
550        }
551
552        pub fn to_duration(self) -> Option<Duration> {
553            if self.sec < 0 || !(0..1_000_000_000).contains(&self.nsec) {
554                return None;
555            }
556            Some(Duration::new(self.sec as u64, self.nsec as u32))
557        }
558
559        fn to_libc(self) -> libc::timespec {
560            libc::timespec {
561                tv_sec: self.sec,
562                tv_nsec: self.nsec,
563            }
564        }
565    }
566
567    #[derive(Copy, Clone, Debug, Default)]
568    pub struct Event {
569        pub ident: usize,
570        pub filter: i16,
571        pub flags: u16,
572        pub fflags: u32,
573        pub data: isize,
574        pub udata: usize,
575    }
576
577    impl Event {
578        pub fn to_libc(self) -> libc::kevent {
579            // Field widths and optional `ext` differ across BSDs; assign
580            // rather than using a struct literal.
581            let mut ev: libc::kevent = unsafe { core::mem::zeroed() };
582            ev.ident = self.ident as _;
583            ev.filter = self.filter as _;
584            ev.flags = self.flags as _;
585            ev.fflags = self.fflags as _;
586            ev.data = self.data as _;
587            ev.udata = self.udata as *mut libc::c_void;
588            ev
589        }
590
591        #[allow(clippy::unnecessary_cast)]
592        pub fn from_libc(e: libc::kevent) -> Self {
593            Self {
594                ident: e.ident as usize,
595                filter: e.filter as i16,
596                flags: e.flags as u16,
597                fflags: e.fflags as u32,
598                data: e.data as isize,
599                udata: e.udata as usize,
600            }
601        }
602    }
603
604    static OPEN: Mutex<Vec<Weak<AtomicI32>>> = Mutex::new(Vec::new());
605
606    fn register_open(cell: &Arc<AtomicI32>) {
607        let mut open = OPEN.lock();
608        open.retain(|w| w.strong_count() > 0);
609        open.push(Arc::downgrade(cell));
610    }
611
612    pub fn create() -> io::Result<Arc<AtomicI32>> {
613        let fd = unsafe { libc::kqueue() };
614        if fd < 0 {
615            return Err(io::Error::last_os_error());
616        }
617        let borrowed = unsafe { BorrowedFd::borrow_raw(fd) };
618        if let Err(err) = crate::posix::set_inheritable(borrowed, false) {
619            let _ = unsafe { libc::close(fd) };
620            return Err(err);
621        }
622        let cell = Arc::new(AtomicI32::new(fd));
623        register_open(&cell);
624        Ok(cell)
625    }
626
627    pub fn from_fd(fd: i32) -> Arc<AtomicI32> {
628        let cell = Arc::new(AtomicI32::new(fd));
629        register_open(&cell);
630        cell
631    }
632
633    pub fn close(cell: &AtomicI32) -> io::Result<()> {
634        let fd = cell.swap(-1, Ordering::SeqCst);
635        if fd < 0 {
636            return Ok(());
637        }
638        let ret = unsafe { libc::close(fd) };
639        if ret < 0 {
640            Err(io::Error::last_os_error())
641        } else {
642            Ok(())
643        }
644    }
645
646    pub fn fd(cell: &AtomicI32) -> i32 {
647        cell.load(Ordering::SeqCst)
648    }
649
650    pub fn kevent(
651        kq: i32,
652        changelist: &[Event],
653        eventlist: &mut [Event],
654        timeout: Option<&Timespec>,
655    ) -> io::Result<usize> {
656        let chl: Vec<libc::kevent> = changelist.iter().copied().map(Event::to_libc).collect();
657        let mut evl = vec![unsafe { core::mem::zeroed() }; eventlist.len()];
658        let timeout = timeout.map(|t| t.to_libc());
659        let timeout = timeout.as_ref().map_or(core::ptr::null(), |t| t);
660        let ret = unsafe {
661            libc::kevent(
662                kq,
663                chl.as_ptr(),
664                chl.len() as _,
665                evl.as_mut_ptr(),
666                evl.len() as _,
667                timeout,
668            )
669        };
670        if ret < 0 {
671            return Err(io::Error::last_os_error());
672        }
673        let n = ret as usize;
674        for (dst, src) in eventlist.iter_mut().zip(evl.into_iter().take(n)) {
675            *dst = Event::from_libc(src);
676        }
677        Ok(n)
678    }
679
680    pub fn mark_closed_after_fork() {
681        // After fork only this thread exists. If the parent held OPEN,
682        // the child's copy stays locked until it is released.
683        if OPEN.try_lock().is_none() {
684            unsafe { OPEN.force_unlock() };
685        }
686        let mut open = OPEN.lock();
687        for weak in open.drain(..) {
688            if let Some(cell) = weak.upgrade() {
689                cell.store(-1, Ordering::SeqCst);
690            }
691        }
692    }
693}
694
695#[cfg(test)]
696mod tests {
697    use super::*;
698
699    #[test]
700    fn duration_to_timeval_uses_micros() {
701        let tv = duration_to_timeval(Duration::from_micros(1_500_250));
702        assert_eq!(tv.tv_sec as u64, 1);
703        assert_eq!(tv.tv_usec as u32, 500_250);
704    }
705
706    #[cfg(unix)]
707    #[test]
708    fn wait_fd_sees_ready_socket_after_deadline() {
709        let mut fds = [0; 2];
710        assert_eq!(unsafe { libc::pipe(fds.as_mut_ptr()) }, 0);
711        let (r, w): (i32, i32) = fds.into();
712        let n = unsafe { libc::write(w, b"x".as_ptr().cast(), 1) };
713        assert_eq!(n, 1);
714        let past = Instant::now().checked_sub(Duration::from_secs(1)).unwrap();
715        let result = wait_fd(r, WaitKind::Read, Some(past));
716        unsafe {
717            libc::close(r);
718            libc::close(w);
719        }
720        assert!(matches!(result, Ok(WaitFd::Ready)));
721    }
722
723    #[cfg(unix)]
724    #[test]
725    fn wait_fd_times_out_when_deadline_passed_and_not_ready() {
726        let mut fds = [0; 2];
727        assert_eq!(unsafe { libc::pipe(fds.as_mut_ptr()) }, 0);
728        let (r, w): (i32, i32) = fds.into();
729        let past = Instant::now().checked_sub(Duration::from_secs(1)).unwrap();
730        let result = wait_fd(r, WaitKind::Read, Some(past));
731        unsafe {
732            libc::close(r);
733            libc::close(w);
734        }
735        assert!(matches!(result, Ok(WaitFd::Timeout)));
736    }
737
738    #[cfg(unix)]
739    #[test]
740    fn poll_timeout_rounds_fractional_millis_up() {
741        assert_eq!(duration_as_millis_ceiling(Duration::ZERO), Some(0));
742        assert_eq!(
743            duration_as_millis_ceiling(Duration::from_micros(100)),
744            Some(1)
745        );
746        assert_eq!(
747            duration_as_millis_ceiling(Duration::from_millis(1)),
748            Some(1)
749        );
750        assert_eq!(
751            duration_as_millis_ceiling(Duration::from_micros(1100)),
752            Some(2)
753        );
754    }
755
756    #[cfg(any(
757        target_os = "macos",
758        target_os = "ios",
759        target_os = "freebsd",
760        target_os = "netbsd",
761        target_os = "openbsd",
762        target_os = "dragonfly",
763    ))]
764    #[test]
765    fn timespec_carries_nanosecond_overflow() {
766        let ts = kqueue::Timespec::from_secs(0.999_999_999_9).unwrap();
767        assert_eq!(ts.sec, 1);
768        assert!(ts.nsec < 1_000_000_000);
769        assert!(kqueue::Timespec::from_secs(f64::INFINITY).is_none());
770        assert!(kqueue::Timespec::from_secs(f64::NAN).is_none());
771    }
772}