1use std::collections::HashMap;
6use std::io;
7use std::time::Duration;
8
9use super::{Event, Interest, RawFd, Reactor, Token};
10
11pub struct KqueueReactor {
16 max_events: usize,
18 registrations: HashMap<Token, Registration>,
20 fd_to_token: HashMap<RawFd, Token>,
22 next_token: usize,
24 #[cfg(any(
26 target_os = "macos",
27 target_os = "freebsd",
28 target_os = "openbsd",
29 target_os = "netbsd"
30 ))]
31 kq: RawFd,
32}
33
34struct Registration {
36 fd: RawFd,
37 interest: Interest,
38}
39
40impl KqueueReactor {
41 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 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 Ok(Self {
81 max_events,
82 registrations: HashMap::new(),
83 fd_to_token: HashMap::new(),
84 next_token: 1,
85 })
86 }
87 }
88
89 #[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 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 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 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(®.fd);
303
304 #[cfg(any(
305 target_os = "macos",
306 target_os = "freebsd",
307 target_os = "openbsd",
308 target_os = "netbsd"
309 ))]
310 {
311 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}