Skip to main content

dcp/reactor/
kqueue.rs

1//! kqueue reactor implementation for macOS/BSD.
2//!
3//! Uses BSD's kqueue interface for event-driven I/O.
4
5use std::collections::HashMap;
6use std::io;
7use std::time::Duration;
8
9use super::{Event, Interest, RawFd, Reactor, Token};
10
11/// kqueue reactor for macOS/BSD
12///
13/// Uses BSD's kqueue interface for event-driven I/O.
14/// This is a poll-based reactor (not true async I/O).
15pub struct KqueueReactor {
16    /// Maximum events per poll
17    max_events: usize,
18    /// Registered file descriptors
19    registrations: HashMap<Token, Registration>,
20    /// Reverse mapping from fd to token
21    fd_to_token: HashMap<RawFd, Token>,
22    /// Next token to assign
23    next_token: usize,
24    /// kqueue file descriptor (on macOS/BSD)
25    #[cfg(any(
26        target_os = "macos",
27        target_os = "freebsd",
28        target_os = "openbsd",
29        target_os = "netbsd"
30    ))]
31    kq: RawFd,
32}
33
34/// Registration state for a file descriptor
35struct Registration {
36    fd: RawFd,
37    interest: Interest,
38}
39
40impl KqueueReactor {
41    /// Create a new kqueue reactor
42    pub fn new(max_events: usize) -> io::Result<Self> {
43        #[cfg(any(
44            target_os = "macos",
45            target_os = "freebsd",
46            target_os = "openbsd",
47            target_os = "netbsd"
48        ))]
49        {
50            let kq = unsafe { libc::kqueue() };
51            if kq < 0 {
52                return Err(io::Error::last_os_error());
53            }
54
55            // Set close-on-exec
56            unsafe {
57                let flags = libc::fcntl(kq, libc::F_GETFD);
58                if flags >= 0 {
59                    libc::fcntl(kq, libc::F_SETFD, flags | libc::FD_CLOEXEC);
60                }
61            }
62
63            Ok(Self {
64                max_events,
65                registrations: HashMap::new(),
66                fd_to_token: HashMap::new(),
67                next_token: 1,
68                kq,
69            })
70        }
71
72        #[cfg(not(any(
73            target_os = "macos",
74            target_os = "freebsd",
75            target_os = "openbsd",
76            target_os = "netbsd"
77        )))]
78        {
79            // Stub implementation for non-BSD platforms
80            Ok(Self {
81                max_events,
82                registrations: HashMap::new(),
83                fd_to_token: HashMap::new(),
84                next_token: 1,
85            })
86        }
87    }
88
89    /// Register kevent changes
90    #[cfg(any(
91        target_os = "macos",
92        target_os = "freebsd",
93        target_os = "openbsd",
94        target_os = "netbsd"
95    ))]
96    fn kevent_register(&self, fd: RawFd, interest: Interest, add: bool) -> io::Result<()> {
97        let mut changes = Vec::new();
98        let flags = if add {
99            libc::EV_ADD | libc::EV_CLEAR
100        } else {
101            libc::EV_DELETE
102        };
103
104        if interest.readable || !add {
105            changes.push(libc::kevent {
106                ident: fd as usize,
107                filter: libc::EVFILT_READ,
108                flags: if add && interest.readable {
109                    flags
110                } else {
111                    libc::EV_DELETE
112                },
113                fflags: 0,
114                data: 0,
115                udata: std::ptr::null_mut(),
116            });
117        }
118
119        if interest.writable || !add {
120            changes.push(libc::kevent {
121                ident: fd as usize,
122                filter: libc::EVFILT_WRITE,
123                flags: if add && interest.writable {
124                    flags
125                } else {
126                    libc::EV_DELETE
127                },
128                fflags: 0,
129                data: 0,
130                udata: std::ptr::null_mut(),
131            });
132        }
133
134        if changes.is_empty() {
135            return Ok(());
136        }
137
138        let result = unsafe {
139            libc::kevent(
140                self.kq,
141                changes.as_ptr(),
142                changes.len() as i32,
143                std::ptr::null_mut(),
144                0,
145                std::ptr::null(),
146            )
147        };
148
149        if result < 0 {
150            Err(io::Error::last_os_error())
151        } else {
152            Ok(())
153        }
154    }
155}
156
157#[cfg(any(
158    target_os = "macos",
159    target_os = "freebsd",
160    target_os = "openbsd",
161    target_os = "netbsd"
162))]
163impl Drop for KqueueReactor {
164    fn drop(&mut self) {
165        unsafe {
166            libc::close(self.kq);
167        }
168    }
169}
170
171impl Reactor for KqueueReactor {
172    fn poll(&mut self, timeout: Option<Duration>) -> io::Result<Vec<Event>> {
173        #[cfg(any(
174            target_os = "macos",
175            target_os = "freebsd",
176            target_os = "openbsd",
177            target_os = "netbsd"
178        ))]
179        {
180            let timespec = timeout.map(|d| libc::timespec {
181                tv_sec: d.as_secs() as libc::time_t,
182                tv_nsec: d.subsec_nanos() as libc::c_long,
183            });
184
185            let mut kevents: Vec<libc::kevent> =
186                vec![unsafe { std::mem::zeroed() }; self.max_events];
187
188            let n = unsafe {
189                libc::kevent(
190                    self.kq,
191                    std::ptr::null(),
192                    0,
193                    kevents.as_mut_ptr(),
194                    self.max_events as i32,
195                    timespec
196                        .as_ref()
197                        .map(|t| t as *const _)
198                        .unwrap_or(std::ptr::null()),
199                )
200            };
201
202            if n < 0 {
203                let err = io::Error::last_os_error();
204                if err.kind() == io::ErrorKind::Interrupted {
205                    return Ok(Vec::new());
206                }
207                return Err(err);
208            }
209
210            // Group events by fd
211            let mut fd_events: HashMap<RawFd, Event> = HashMap::new();
212
213            for i in 0..n as usize {
214                let kevent = &kevents[i];
215                let fd = kevent.ident as RawFd;
216
217                if let Some(&token) = self.fd_to_token.get(&fd) {
218                    let event = fd_events.entry(fd).or_insert_with(|| Event::new(token));
219
220                    if kevent.filter == libc::EVFILT_READ {
221                        event.readable = true;
222                    }
223                    if kevent.filter == libc::EVFILT_WRITE {
224                        event.writable = true;
225                    }
226                    if (kevent.flags & libc::EV_ERROR) != 0 {
227                        event.error = true;
228                    }
229                    if (kevent.flags & libc::EV_EOF) != 0 {
230                        event.closed = true;
231                    }
232                }
233            }
234
235            Ok(fd_events.into_values().collect())
236        }
237
238        #[cfg(not(any(
239            target_os = "macos",
240            target_os = "freebsd",
241            target_os = "openbsd",
242            target_os = "netbsd"
243        )))]
244        {
245            // Stub for non-BSD: just sleep and return empty
246            if let Some(duration) = timeout {
247                std::thread::sleep(duration.min(Duration::from_millis(100)));
248            }
249            Ok(Vec::new())
250        }
251    }
252
253    fn register(&mut self, fd: RawFd, interest: Interest) -> io::Result<Token> {
254        let token = Token(self.next_token);
255        self.next_token += 1;
256
257        #[cfg(any(
258            target_os = "macos",
259            target_os = "freebsd",
260            target_os = "openbsd",
261            target_os = "netbsd"
262        ))]
263        {
264            self.kevent_register(fd, interest, true)?;
265        }
266
267        self.registrations
268            .insert(token, Registration { fd, interest });
269        self.fd_to_token.insert(fd, token);
270
271        Ok(token)
272    }
273
274    fn modify(&mut self, token: Token, interest: Interest) -> io::Result<()> {
275        let reg = self
276            .registrations
277            .get_mut(&token)
278            .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "Token not registered"))?;
279
280        #[cfg(any(
281            target_os = "macos",
282            target_os = "freebsd",
283            target_os = "openbsd",
284            target_os = "netbsd"
285        ))]
286        {
287            // Remove old interest, add new
288            self.kevent_register(reg.fd, reg.interest, false)?;
289            self.kevent_register(reg.fd, interest, true)?;
290        }
291
292        reg.interest = interest;
293        Ok(())
294    }
295
296    fn deregister(&mut self, token: Token) -> io::Result<()> {
297        let reg = self
298            .registrations
299            .remove(&token)
300            .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "Token not registered"))?;
301
302        self.fd_to_token.remove(&reg.fd);
303
304        #[cfg(any(
305            target_os = "macos",
306            target_os = "freebsd",
307            target_os = "openbsd",
308            target_os = "netbsd"
309        ))]
310        {
311            // Ignore errors on deregister - fd might already be closed
312            let _ = self.kevent_register(reg.fd, reg.interest, false);
313        }
314
315        Ok(())
316    }
317
318    fn supports_async_io(&self) -> bool {
319        false
320    }
321
322    fn name(&self) -> &'static str {
323        "kqueue"
324    }
325}
326
327#[cfg(test)]
328mod tests {
329    use super::*;
330
331    #[test]
332    fn test_kqueue_creation() {
333        let result = KqueueReactor::new(256);
334        assert!(result.is_ok());
335    }
336
337    #[test]
338    fn test_kqueue_name() {
339        let reactor = KqueueReactor::new(256).unwrap();
340        assert_eq!(reactor.name(), "kqueue");
341    }
342
343    #[test]
344    fn test_kqueue_no_async_io() {
345        let reactor = KqueueReactor::new(256).unwrap();
346        assert!(!reactor.supports_async_io());
347    }
348}