Skip to main content

starnix_core/vfs/
epoll.rs

1// Copyright 2021 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use crate::power::WakeupSourceOrigin;
6use crate::task::{
7    CurrentTask, EventHandler, ReadyItem, ReadyItemKey, WaitCanceler, WaitQueue, Waiter,
8};
9use crate::vfs::{
10    Anon, FileHandle, FileObject, FileObjectState, FileOps, WeakFileHandle, fileops_impl_dataless,
11    fileops_impl_nonseekable, fileops_impl_noop_sync,
12};
13use itertools::Itertools;
14use starnix_logging::log_warn;
15use starnix_sync::{EpollStateLock, EpollWaitableStateLock, LockDepMutex, allow_subclass};
16use starnix_uapi::error;
17use starnix_uapi::errors::{EINTR, ETIMEDOUT, Errno};
18use starnix_uapi::open_flags::OpenFlags;
19use starnix_uapi::vfs::{EpollEvent, FdEvents};
20use std::collections::hash_map::Entry;
21use std::collections::{HashMap, VecDeque};
22use std::sync::Arc;
23
24/// Maximum depth of epoll instances monitoring one another.
25/// From https://man7.org/linux/man-pages/man2/epoll_ctl.2.html
26const MAX_NESTED_DEPTH: u32 = 5;
27
28/// WaitObject represents a FileHandle that is being waited upon.
29/// The `data` field is a user defined quantity passed in
30/// via `sys_epoll_ctl`. Typically C programs could use this
31/// to store a pointer to the data that needs to be processed
32/// after an event.
33struct WaitObject {
34    target: WeakFileHandle,
35    events: FdEvents,
36    data: u64,
37    wait_canceler: Option<WaitCanceler>,
38    active_wakeup_source: Option<WakeupSourceOrigin>,
39}
40
41impl WaitObject {
42    /// Returns the target `FileHandle` of the `WaitObject`, or `None` if the file has been closed.
43    ///
44    /// It is fine for the `FileHandle` to be closed after being added to an epoll, and subsequent
45    /// epoll_waits end up timing out (importantly not returning EBADF).
46    fn target(&self) -> Option<FileHandle> {
47        self.target.upgrade()
48    }
49
50    fn deactivate_wakeup_source(&mut self, current_task: &CurrentTask) {
51        if let Some(origin) = self.active_wakeup_source.take() {
52            current_task.kernel().suspend_resume_manager.deactivate_wakeup_source(&origin);
53        }
54    }
55}
56
57/// EpollKey acts as an key to a map of WaitObject.
58/// In reality it is a pointer to a FileHandle object.
59pub type EpollKey = usize;
60
61/// EpollFileObject represents the FileObject used to
62/// implement epoll_create1/epoll_ctl/epoll_pwait.
63#[derive(Default)]
64pub struct EpollFileObject {
65    waiter: Waiter,
66    /// Mutable state of this epoll object.
67    state: LockDepMutex<EpollState, EpollStateLock>,
68    waitable_state: Arc<LockDepMutex<EpollWaitableState, EpollWaitableStateLock>>,
69    /// A list of waiters waiting for events from this
70    /// epoll instance.
71    waiters: Arc<WaitQueue>,
72}
73
74#[derive(Default)]
75struct EpollState {
76    /// Any file tracked by this epoll instance
77    /// will exist as a key in `wait_objects`.
78    wait_objects: HashMap<ReadyItemKey, WaitObject>,
79    /// processing_list is a FIFO of events that are being
80    /// processed.
81    ///
82    /// Objects from the `EpollFileObject`'s `trigger_list` are moved into this
83    /// list so that we can handle triggered events without holding its lock
84    /// longer than we need to. This reduces contention with waited-on objects
85    /// that tries to notify this epoll object on subscribed events.
86    processing_list: VecDeque<ReadyItem>,
87    /// recheck_list is the list of items that need to have query_events checked
88    /// at the start of the next EpollFileObject::wait call. This is only items
89    /// that were returned from the last wait call, because those are the only
90    /// ones that might need to be returned even if no events come in between
91    /// wait calls.
92    recheck_list: Vec<ReadyItemKey>,
93}
94
95#[derive(Default)]
96struct EpollWaitableState {
97    /// trigger_list is a FIFO of events that have
98    /// happened, but have not yet been processed.
99    trigger_list: VecDeque<ReadyItem>,
100}
101
102impl EpollFileObject {
103    /// Allocate a new, empty epoll object.
104    pub fn new_file(current_task: &CurrentTask) -> FileHandle {
105        let epoll = Box::new(EpollFileObject::default());
106
107        #[cfg(any(test, debug_assertions))]
108        {
109            let _l1 = epoll.state.lock();
110            let _l2 = epoll.waitable_state.lock();
111        }
112
113        Anon::new_private_file(current_task, epoll, OpenFlags::RDWR, "[eventpoll]")
114    }
115
116    fn wait_on_file(
117        &self,
118        current_task: &CurrentTask,
119        key: ReadyItemKey,
120        wait_object: &mut WaitObject,
121    ) -> Result<(), Errno> {
122        // First start the wait. If an event happens after this, we'll get it.
123        self.wait_on_file_edge_triggered(current_task, key, wait_object)?;
124
125        self.do_recheck(current_task, wait_object, key)?;
126
127        Ok(())
128    }
129
130    fn do_recheck(
131        &self,
132        current_task: &CurrentTask,
133        wait_object: &mut WaitObject,
134        key: ReadyItemKey,
135    ) -> Result<(), Errno> {
136        let Some(target) = wait_object.target() else { return Ok(()) };
137        let events = {
138            // Target might be itself an epoll object. Because there is no loop,
139            // this allow_subclass is safe.
140            let _token = allow_subclass();
141            target.query_events(current_task)?
142        };
143        if !(events & wait_object.events).is_empty() {
144            self.waiter.wake_immediately(events, self.new_wait_handler(key));
145            if let Some(wait_canceler) = wait_object.wait_canceler.take() {
146                wait_canceler.cancel();
147            } else {
148                log_warn!("wait canceler should have been set by `wait_on_file_edge_triggered`");
149            }
150        }
151        Ok(())
152    }
153
154    fn wait_on_file_edge_triggered(
155        &self,
156        current_task: &CurrentTask,
157        key: ReadyItemKey,
158        wait_object: &mut WaitObject,
159    ) -> Result<(), Errno> {
160        let Some(target) = wait_object.target() else {
161            return Ok(());
162        };
163
164        wait_object.wait_canceler = target.wait_async(
165            current_task,
166            &self.waiter,
167            wait_object.events,
168            self.new_wait_handler(key),
169        );
170        if wait_object.wait_canceler.is_none() {
171            return error!(EPERM);
172        }
173        Ok(())
174    }
175
176    /// Checks if adding self to the `epoll_file_object` at `epoll_file_handle` would cause a loop
177    /// or exceed max depth.
178    fn check_eloop(&self, parent: &FileHandle, depth_left: u32) -> Result<(), Errno> {
179        if depth_left == 0 {
180            return error!(ELOOP);
181        }
182
183        let state = self.state.lock();
184        for nested_object in state.wait_objects.values() {
185            let Some(child) = nested_object.target() else {
186                continue;
187            };
188            let Some(child_file) = child.downcast_file::<EpollFileObject>() else {
189                continue;
190            };
191
192            if Arc::ptr_eq(&child, parent) {
193                return error!(ELOOP);
194            }
195            // Child is not part of a loop, so subclassing is safe.
196            let _token = allow_subclass();
197            child_file.check_eloop(parent, depth_left - 1)?;
198        }
199
200        Ok(())
201    }
202
203    /// Asynchronously wait on certain events happening on a FileHandle.
204    pub fn add(
205        &self,
206        current_task: &CurrentTask,
207        file: &FileHandle,
208        epoll_file_handle: &FileHandle,
209        epoll_event: EpollEvent,
210    ) -> Result<(), Errno> {
211        // Check if adding this file would cause a cycle at a max depth of 5.
212        if let Some(epoll_to_add) = file.downcast_file::<EpollFileObject>() {
213            // We need to check for `MAX_NESTED_DEPTH - 1` because adding `epoll_to_add` to self
214            // would result in a total depth of one more.
215            epoll_to_add.check_eloop(epoll_file_handle, MAX_NESTED_DEPTH - 1)?;
216        }
217
218        let mut state = self.state.lock();
219        let key = file.id.as_epoll_key().into();
220        match state.wait_objects.entry(key) {
221            Entry::Occupied(_) => error!(EEXIST),
222            Entry::Vacant(entry) => {
223                let wait_object = entry.insert(WaitObject {
224                    target: Arc::downgrade(file),
225                    events: epoll_event.events() | FdEvents::POLLHUP | FdEvents::POLLERR,
226                    data: epoll_event.data(),
227                    wait_canceler: None,
228                    active_wakeup_source: None,
229                });
230                self.wait_on_file(current_task, key, wait_object)
231            }
232        }
233    }
234
235    /// Modify the events we are looking for on a Filehandle.
236    pub fn modify(
237        &self,
238        current_task: &CurrentTask,
239        file: &FileHandle,
240        epoll_event: EpollEvent,
241    ) -> Result<(), Errno> {
242        let mut state = self.state.lock();
243        let key = file.id.as_epoll_key();
244        state.recheck_list.retain(|x| *x != key.into());
245        let Some(wait_object) = state.wait_objects.get_mut(&key.into()) else {
246            return error!(ENOENT);
247        };
248        if let Some(wait_canceler) = wait_object.wait_canceler.take() {
249            wait_canceler.cancel();
250        }
251        wait_object.events = epoll_event.events() | FdEvents::POLLHUP | FdEvents::POLLERR;
252        wait_object.data = epoll_event.data();
253        // If the new epoll event doesn't include EPOLLWAKEUP, we need to take down the
254        // wake lease. This ensures that the system doesn't stay awake unnecessarily when
255        // the event no longer requires it to be awake.
256        if wait_object.events.contains(FdEvents::EPOLLWAKEUP)
257            && !epoll_event.events().contains(FdEvents::EPOLLWAKEUP)
258        {
259            wait_object.deactivate_wakeup_source(current_task);
260        }
261        self.wait_on_file(current_task, key.into(), wait_object)
262    }
263
264    /// Cancel an asynchronous wait on an object. Events triggered before
265    /// calling this will still be delivered.
266    pub fn delete(&self, current_task: &CurrentTask, file: &FileObject) -> Result<(), Errno> {
267        let mut state = self.state.lock();
268        let key = file.id.as_epoll_key().into();
269        if let Some(mut wait_object) = state.wait_objects.remove(&key) {
270            if let Some(wait_canceler) = wait_object.wait_canceler.take() {
271                wait_canceler.cancel();
272            }
273            state.recheck_list.retain(|x| *x != key);
274            // Deactivate the wake lock if it was active.
275            wait_object.deactivate_wakeup_source(current_task);
276            Ok(())
277        } else {
278            error!(ENOENT)
279        }
280    }
281
282    /// Stores events from the Epoll's trigger list to the parameter `pending_list`. This does not
283    /// actually invoke the waiter which is how items are added to the trigger list. The caller
284    /// will have to do that before calling if needed.
285    ///
286    /// If an event in the trigger list is stale, the event will be re-added to the waiter.
287    ///
288    /// Returns true if any events were added. False means there was nothing in the trigger list.
289    fn process_triggered_events(
290        &self,
291        current_task: &CurrentTask,
292        pending_list: &mut Vec<ReadyItem>,
293        max_events: usize,
294    ) -> Result<(), Errno> {
295        let mut state = self.state.lock();
296        // Move all the elements from `self.trigger_list` to this intermediary
297        // queue that we handle events from. This reduces the time spent holding
298        // `self.trigger_list`'s lock which reduces contention with objects that
299        // this epoll object has subscribed for notifications from.
300        state.processing_list.append(&mut self.waitable_state.lock().trigger_list);
301        while pending_list.len() < max_events && !state.processing_list.is_empty() {
302            if let Some(pending) = state.processing_list.pop_front() {
303                if let Some(wait) = state.wait_objects.get_mut(&pending.key) {
304                    // The weak pointer to the FileObject target can be gone if the file was closed
305                    // out from under us. If this happens it is not an error: ignore it and
306                    // continue.
307                    if let Some(target) = wait.target.upgrade() {
308                        let events = {
309                            // Target might be itself an epoll object. Because there is no loop,
310                            // this allow_subclass is safe.
311                            let _token = allow_subclass();
312                            target.query_events(current_task)?
313                        };
314                        let ready = ReadyItem { key: pending.key, events };
315                        if ready.events.intersects(wait.events) {
316                            pending_list.push(ready);
317                        } else {
318                            // Another thread already handled this event, wait for another one.
319                            self.wait_on_file(current_task, pending.key, wait)?;
320                        }
321                    }
322                }
323            }
324        }
325        Ok(())
326    }
327
328    /// Waits until an event exists in `pending_list` or until `timeout` has
329    /// been reached.
330    fn wait_until_pending_event(
331        &self,
332        current_task: &CurrentTask,
333        max_events: usize,
334        mut wait_deadline: zx::MonotonicInstant,
335    ) -> Result<Vec<ReadyItem>, Errno> {
336        let mut pending_list = Vec::new();
337
338        loop {
339            self.process_triggered_events(current_task, &mut pending_list, max_events)?;
340
341            if pending_list.len() == max_events {
342                break; // No input events or output list full, nothing more we can do.
343            }
344
345            if !pending_list.is_empty() {
346                // We now know we have at least one event to return. We shouldn't return
347                // immediately, in case there are more events available, but the next loop should
348                // wait with a 0 timeout to prevent further blocking.
349                wait_deadline = zx::MonotonicInstant::ZERO;
350            }
351
352            // Loop back to check if there are more items in the Waiter's queue. Every wait_until()
353            // call will process a single event. In order to drain as many events as we can that
354            // are synchronously available, keep trying until it reports empty.
355            //
356            // The handlers in the waits cause items to be appended to trigger_list. See the closure
357            // in `wait_on_file` to see how this happens.
358            //
359            // This wait may return EINTR for nonzero timeouts which is not an error. We must be
360            // careful not to lose events if this happens.
361            //
362            // The first time through this loop we'll use the timeout passed into this function so
363            // can get EINTR. But since we haven't done anything or accumulated any results yet it's
364            // OK to immediately return and no information will be lost.
365            match self.waiter.wait_until(current_task, wait_deadline) {
366                Err(err) if err == ETIMEDOUT => break,
367                Err(err) if err == EINTR => {
368                    // Terminating early will lose any events in the pending_list so that should
369                    // only be for unrecoverable errors (not EINTR). The only time there should be a
370                    // nonzero wait_deadline (and hence the ability to encounter EINTR) is when the
371                    // pending list is empty.
372                    debug_assert!(
373                        pending_list.is_empty(),
374                        "Got EINTR from wait of {}ns with {} items pending.",
375                        wait_deadline.into_nanos(),
376                        pending_list.len()
377                    );
378                    return Err(err);
379                }
380                // TODO check if this is supposed to actually fail!
381                result => result?,
382            }
383        }
384
385        Ok(pending_list)
386    }
387
388    /// Blocking wait on all waited upon events with a timeout.
389    pub fn wait(
390        &self,
391        current_task: &CurrentTask,
392        max_events: usize,
393        deadline: zx::MonotonicInstant,
394    ) -> Result<Vec<EpollEvent>, Errno> {
395        {
396            let mut state = self.state.lock();
397            let recheck_list = std::mem::take(&mut state.recheck_list);
398            for key in &recheck_list {
399                let wait_object = state.wait_objects.get_mut(key).unwrap();
400                wait_object.deactivate_wakeup_source(current_task);
401            }
402            for (idx, key) in recheck_list.iter().enumerate() {
403                let wait_object = state.wait_objects.get_mut(key).unwrap();
404                if let Err(err) = self.do_recheck(current_task, wait_object, *key) {
405                    // Put remaining unprocessed entries (including the failing one) back
406                    // onto `recheck_list` so they aren't lost and can be rechecked on the next wait.
407                    state.recheck_list.extend(recheck_list[idx..].iter().copied());
408                    return Err(err);
409                }
410            }
411        }
412
413        let pending_list = self.wait_until_pending_event(current_task, max_events, deadline)?;
414
415        // Process the pending list and add processed ReadyItem
416        // entries to the rearm_list for the next wait.
417        let mut result = vec![];
418        let mut state = self.state.lock();
419        let state = &mut *state;
420        for pending_event in pending_list.iter().unique_by(|e| e.key) {
421            // The wait could have been deleted by here,
422            // so ignore the None case.
423            let Some(wait) = state.wait_objects.get_mut(&pending_event.key) else { continue };
424
425            let reported_events = pending_event.events & wait.events;
426            result.push(EpollEvent::new(reported_events, wait.data));
427
428            // Files marked with `EPOLLONESHOT` should only notify
429            // once and need to be rearmed manually with epoll_ctl_mod().
430            if wait.events.contains(FdEvents::EPOLLONESHOT) {
431                continue;
432            }
433
434            self.wait_on_file_edge_triggered(current_task, pending_event.key, wait)?;
435
436            if !wait.events.contains(FdEvents::EPOLLET) {
437                state.recheck_list.push(pending_event.key);
438            }
439
440            // TODO: is this really only supposed to happen for level-triggered events?
441            if !wait.events.contains(FdEvents::EPOLLET) {
442                // When this is the first time epoll_wait on this epoll fd, create and
443                // hold a wake lease until the next epoll_wait.
444                if wait.events.contains(FdEvents::EPOLLWAKEUP) {
445                    if let ReadyItemKey::Usize(key) = pending_event.key {
446                        let origin = WakeupSourceOrigin::Epoll(key);
447                        current_task
448                            .kernel()
449                            .suspend_resume_manager
450                            .activate_wakeup_source_with_actor(
451                                origin.clone(),
452                                Some(current_task.command()),
453                            );
454                        wait.active_wakeup_source = Some(origin);
455                    }
456                }
457            }
458        }
459
460        Ok(result)
461    }
462}
463
464impl FileOps for EpollFileObject {
465    fileops_impl_nonseekable!();
466    fileops_impl_noop_sync!();
467    fileops_impl_dataless!();
468
469    fn wait_async(
470        &self,
471        _file: &FileObject,
472        _current_task: &CurrentTask,
473        waiter: &Waiter,
474        events: FdEvents,
475        handler: EventHandler,
476    ) -> Option<WaitCanceler> {
477        Some(self.waiters.wait_async_fd_events(waiter, events, handler))
478    }
479
480    fn query_events(
481        &self,
482        _file: &FileObject,
483        current_task: &CurrentTask,
484    ) -> Result<FdEvents, Errno> {
485        let mut events = FdEvents::empty();
486        let state = self.state.lock();
487        if !state.processing_list.is_empty() || !self.waitable_state.lock().trigger_list.is_empty()
488        {
489            events |= FdEvents::POLLIN;
490        } else {
491            for key in &state.recheck_list {
492                let wait_object = state.wait_objects.get(key).unwrap();
493                let Some(target) = wait_object.target() else { continue };
494                // Target might be itself an epoll object. Because there is no loop,
495                // this allow_subclass is safe.
496                let _token = allow_subclass();
497                if !(target.query_events(current_task)? & wait_object.events).is_empty() {
498                    events |= FdEvents::POLLIN;
499                    break;
500                }
501            }
502        }
503        Ok(events)
504    }
505
506    fn close(self: Box<Self>, _file: &FileObjectState, current_task: &CurrentTask) {
507        let mut guard = self.state.lock();
508        for wait_object in guard.wait_objects.values_mut() {
509            wait_object.deactivate_wakeup_source(current_task);
510        }
511    }
512}
513
514#[derive(Clone)]
515pub struct EpollEventHandler {
516    key: ReadyItemKey,
517    waitable_state: Arc<LockDepMutex<EpollWaitableState, EpollWaitableStateLock>>,
518    waiters: Arc<WaitQueue>,
519}
520
521impl EpollEventHandler {
522    pub fn handle(self, events: FdEvents) {
523        {
524            let mut waitable_state = self.waitable_state.lock();
525            waitable_state.trigger_list.push_back(ReadyItem { key: self.key, events });
526        }
527        self.waiters.notify_fd_events(FdEvents::POLLIN);
528    }
529}
530
531impl EpollFileObject {
532    fn new_wait_handler(&self, key: ReadyItemKey) -> EventHandler {
533        EventHandler::Epoll(EpollEventHandler {
534            key,
535            waitable_state: Arc::clone(&self.waitable_state),
536            waiters: Arc::clone(&self.waiters),
537        })
538    }
539}
540
541#[cfg(test)]
542mod tests {
543    use super::*;
544    use crate::fs::fuchsia::create_fuchsia_pipe;
545    use crate::task::Waiter;
546    use crate::task::dynamic_thread_spawner::SpawnRequestBuilder;
547    use crate::testing::{anon_test_file, spawn_kernel_and_run};
548    use crate::vfs::buffers::{VecInputBuffer, VecOutputBuffer};
549    use crate::vfs::eventfd::{EventFdType, new_eventfd};
550    use crate::vfs::fs_registry::FsRegistry;
551    use crate::vfs::pipe::{new_pipe, register_pipe_fs};
552    use crate::vfs::socket::{SocketDomain, SocketType, UnixSocket};
553    use starnix_lifecycle::AtomicCounter;
554    use starnix_sync::Mutex;
555    use starnix_uapi::vfs::{EpollEvent, FdEvents};
556    use syncio::Zxio;
557
558    #[::fuchsia::test]
559    async fn test_epoll_read_ready() {
560        static WRITE_COUNT: AtomicCounter<usize> = AtomicCounter::<usize>::new_const(0);
561        const EVENT_DATA: u64 = 42;
562
563        spawn_kernel_and_run(async |current_task| {
564            let kernel = current_task.kernel();
565            register_pipe_fs(kernel.expando.get::<FsRegistry>().as_ref());
566
567            let (pipe_out, pipe_in) = new_pipe(&current_task).unwrap();
568
569            let test_string = "hello starnix".to_string();
570            let test_len = test_string.len();
571
572            let epoll_file_handle = EpollFileObject::new_file(&current_task);
573            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
574            epoll_file
575                .add(
576                    &current_task,
577                    &pipe_out,
578                    &epoll_file_handle,
579                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
580                )
581                .unwrap();
582
583            let (sender, receiver) = std::sync::mpsc::channel();
584            let value = test_string.clone();
585            let closure = move |task: &CurrentTask| {
586                let bytes_written =
587                    pipe_in.write(&task, &mut VecInputBuffer::new(value.as_bytes())).unwrap();
588                assert_eq!(bytes_written, test_len);
589                WRITE_COUNT.add(bytes_written);
590                sender.send(()).unwrap();
591            };
592            let req = SpawnRequestBuilder::new().with_sync_closure(closure).build();
593            kernel.kthreads.spawner().spawn_from_request(req);
594            let events =
595                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::INFINITE).unwrap();
596            receiver.recv().unwrap();
597            assert_eq!(1, events.len());
598            let event = &events[0];
599            assert!(event.events().contains(FdEvents::POLLIN));
600            assert_eq!(event.data(), EVENT_DATA);
601
602            let mut buffer = VecOutputBuffer::new(test_len);
603            let bytes_read = pipe_out.read(&current_task, &mut buffer).unwrap();
604            assert_eq!(bytes_read, WRITE_COUNT.get());
605            assert_eq!(bytes_read, test_len);
606            assert_eq!(buffer.data(), test_string.as_bytes());
607        })
608        .await;
609    }
610
611    #[::fuchsia::test]
612    async fn test_epoll_ready_then_wait() {
613        const EVENT_DATA: u64 = 42;
614
615        spawn_kernel_and_run(async |current_task| {
616            let kernel = current_task.kernel();
617            register_pipe_fs(kernel.expando.get::<FsRegistry>().as_ref());
618
619            let (pipe_out, pipe_in) = new_pipe(&current_task).unwrap();
620
621            let test_string = "hello starnix".to_string();
622            let test_bytes = test_string.as_bytes();
623            let test_len = test_bytes.len();
624
625            assert_eq!(
626                pipe_in.write(&current_task, &mut VecInputBuffer::new(test_bytes)).unwrap(),
627                test_bytes.len()
628            );
629
630            let epoll_file_handle = EpollFileObject::new_file(&current_task);
631            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
632            epoll_file
633                .add(
634                    &current_task,
635                    &pipe_out,
636                    &epoll_file_handle,
637                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
638                )
639                .unwrap();
640
641            let events =
642                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::INFINITE).unwrap();
643            assert_eq!(1, events.len());
644            let event = &events[0];
645            assert!(event.events().contains(FdEvents::POLLIN));
646            assert_eq!(event.data(), EVENT_DATA);
647
648            let mut buffer = VecOutputBuffer::new(test_len);
649            let bytes_read = pipe_out.read(&current_task, &mut buffer).unwrap();
650            assert_eq!(bytes_read, test_len);
651            assert_eq!(buffer.data(), test_bytes);
652        })
653        .await;
654    }
655
656    #[::fuchsia::test]
657    async fn test_epoll_ctl_cancel() {
658        spawn_kernel_and_run(async |current_task| {
659            for do_cancel in [true, false] {
660                let event = new_eventfd(&current_task, 0, EventFdType::Counter, true);
661                let waiter = Waiter::new();
662
663                let epoll_file_handle = EpollFileObject::new_file(&current_task);
664                let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
665                const EVENT_DATA: u64 = 42;
666                epoll_file
667                    .add(
668                        &current_task,
669                        &event,
670                        &epoll_file_handle,
671                        EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
672                    )
673                    .unwrap();
674
675                if do_cancel {
676                    epoll_file.delete(&current_task, &event).unwrap();
677                }
678
679                let wait_canceler = event
680                    .wait_async(&current_task, &waiter, FdEvents::POLLIN, EventHandler::None)
681                    .expect("wait_async");
682                if do_cancel {
683                    wait_canceler.cancel();
684                }
685
686                let add_val = 1u64;
687
688                assert_eq!(
689                    event
690                        .write(&current_task, &mut VecInputBuffer::new(&add_val.to_ne_bytes()))
691                        .unwrap(),
692                    std::mem::size_of::<u64>()
693                );
694
695                let events =
696                    epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
697
698                if do_cancel {
699                    assert_eq!(0, events.len());
700                } else {
701                    assert_eq!(1, events.len());
702                    let event = &events[0];
703                    assert!(event.events().contains(FdEvents::POLLIN));
704                    assert_eq!(event.data(), EVENT_DATA);
705                }
706            }
707        })
708        .await;
709    }
710
711    #[::fuchsia::test]
712    async fn test_multiple_events() {
713        spawn_kernel_and_run(async |current_task| {
714            let (client1, server1) = zx::Socket::create_stream();
715            let (client2, server2) = zx::Socket::create_stream();
716            let pipe1 = create_fuchsia_pipe(&current_task, client1, OpenFlags::RDWR)
717                .expect("create_fuchsia_pipe");
718            let pipe2 = create_fuchsia_pipe(&current_task, client2, OpenFlags::RDWR)
719                .expect("create_fuchsia_pipe");
720            let server1_zxio = Zxio::create(server1.into_handle()).expect("Zxio::create");
721            let server2_zxio = Zxio::create(server2.into_handle()).expect("Zxio::create");
722
723            let poll = || {
724                let epoll_object = EpollFileObject::new_file(&current_task);
725                let epoll_file = epoll_object.downcast_file::<EpollFileObject>().unwrap();
726                epoll_file
727                    .add(&current_task, &pipe1, &epoll_object, EpollEvent::new(FdEvents::POLLIN, 1))
728                    .expect("epoll_file.add");
729                epoll_file
730                    .add(&current_task, &pipe2, &epoll_object, EpollEvent::new(FdEvents::POLLIN, 2))
731                    .expect("epoll_file.add");
732                epoll_file.wait(&current_task, 2, zx::MonotonicInstant::ZERO).expect("wait")
733            };
734
735            let fds = poll();
736            assert!(fds.is_empty());
737
738            assert_eq!(server1_zxio.write(&[0]).expect("write"), 1);
739
740            let fds = poll();
741            assert_eq!(fds.len(), 1);
742            assert_eq!(fds[0].events(), FdEvents::POLLIN);
743            assert_eq!(fds[0].data(), 1);
744            assert_eq!(pipe1.read(&current_task, &mut VecOutputBuffer::new(64)).expect("read"), 1);
745
746            let fds = poll();
747            assert!(fds.is_empty());
748
749            assert_eq!(server2_zxio.write(&[0]).expect("write"), 1);
750
751            let fds = poll();
752            assert_eq!(fds.len(), 1);
753            assert_eq!(fds[0].events(), FdEvents::POLLIN);
754            assert_eq!(fds[0].data(), 2);
755            assert_eq!(pipe2.read(&current_task, &mut VecOutputBuffer::new(64)).expect("read"), 1);
756
757            let fds = poll();
758            assert!(fds.is_empty());
759        })
760        .await;
761    }
762
763    #[::fuchsia::test]
764    async fn test_cancel_after_notify() {
765        spawn_kernel_and_run(async |current_task| {
766            let event = new_eventfd(&current_task, 0, EventFdType::Counter, true);
767            let epoll_file_handle = EpollFileObject::new_file(&current_task);
768            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
769
770            // Add a thing
771            const EVENT_DATA: u64 = 42;
772            epoll_file
773                .add(
774                    &current_task,
775                    &event,
776                    &epoll_file_handle,
777                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
778                )
779                .unwrap();
780
781            // Make the thing send a notification, wait for it
782            let add_val = 1u64;
783            assert_eq!(
784                event
785                    .write(&current_task, &mut VecInputBuffer::new(&add_val.to_ne_bytes()))
786                    .unwrap(),
787                std::mem::size_of::<u64>()
788            );
789
790            assert_eq!(
791                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap().len(),
792                1
793            );
794
795            // Remove the thing
796            epoll_file.delete(&current_task, &event).unwrap();
797
798            // Wait for new notifications
799            assert_eq!(
800                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap().len(),
801                0
802            );
803            // That shouldn't crash
804        })
805        .await;
806    }
807
808    #[::fuchsia::test]
809    async fn test_add_then_modify() {
810        spawn_kernel_and_run(async |current_task| {
811            let (socket1, _socket2) = UnixSocket::new_pair(
812                &current_task,
813                SocketDomain::Unix,
814                SocketType::Stream,
815                OpenFlags::RDWR,
816            )
817            .expect("Failed to create socket pair.");
818
819            let epoll_file_handle = EpollFileObject::new_file(&current_task);
820            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
821
822            const EVENT_DATA: u64 = 42;
823            epoll_file
824                .add(
825                    &current_task,
826                    &socket1,
827                    &epoll_file_handle,
828                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
829                )
830                .unwrap();
831            assert_eq!(
832                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap().len(),
833                0
834            );
835
836            let read_write_event = FdEvents::POLLIN | FdEvents::POLLOUT;
837            epoll_file
838                .modify(&current_task, &socket1, EpollEvent::new(read_write_event, EVENT_DATA))
839                .unwrap();
840            let triggered_events =
841                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
842            assert_eq!(1, triggered_events.len());
843            let event = &triggered_events[0];
844            assert_eq!(event.events(), FdEvents::POLLOUT);
845            assert_eq!(event.data(), EVENT_DATA);
846        })
847        .await;
848    }
849
850    #[::fuchsia::test]
851    async fn test_waiter_removal() {
852        spawn_kernel_and_run(async |current_task| {
853            let event = new_eventfd(&current_task, 0, EventFdType::Counter, true);
854            let epoll_file_handle = EpollFileObject::new_file(&current_task);
855            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
856
857            const EVENT_DATA: u64 = 42;
858            epoll_file
859                .add(
860                    &current_task,
861                    &event,
862                    &epoll_file_handle,
863                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
864                )
865                .unwrap();
866
867            std::mem::drop(event);
868
869            assert!(epoll_file.waiters.is_empty());
870        })
871        .await;
872    }
873
874    /// A test file whose event query results can be dynamically controlled.
875    ///
876    /// This allows tests to simulate scenarios where `query_events` succeeds initially
877    /// (e.g., during epoll registration and the first `wait`), but returns an error on a
878    /// subsequent call (e.g., when `do_recheck` re-evaluates `recheck_list`).
879    #[derive(Clone)]
880    struct ControlledEventsFile {
881        wait_queue: Arc<WaitQueue>,
882        events: Arc<Mutex<Result<FdEvents, Errno>>>,
883    }
884
885    impl ControlledEventsFile {
886        fn new(events: FdEvents) -> Self {
887            Self {
888                wait_queue: Arc::new(WaitQueue::default()),
889                events: Arc::new(Mutex::new(Ok(events))),
890            }
891        }
892
893        fn set_events(&self, events: Result<FdEvents, Errno>) {
894            *self.events.lock() = events;
895        }
896    }
897
898    impl FileOps for ControlledEventsFile {
899        fileops_impl_nonseekable!();
900        fileops_impl_noop_sync!();
901        fileops_impl_dataless!();
902
903        fn wait_async(
904            &self,
905            _file: &FileObject,
906            _current_task: &CurrentTask,
907            waiter: &Waiter,
908            events: FdEvents,
909            handler: EventHandler,
910        ) -> Option<WaitCanceler> {
911            Some(self.wait_queue.wait_async_fd_events(waiter, events, handler))
912        }
913
914        fn query_events(
915            &self,
916            _file: &FileObject,
917            _current_task: &CurrentTask,
918        ) -> Result<FdEvents, Errno> {
919            self.events.lock().clone()
920        }
921    }
922
923    #[::fuchsia::test]
924    async fn test_epoll_recheck_failure_deactivates_wakeup_sources() {
925        spawn_kernel_and_run(async |current_task| {
926            let epoll_file_handle = EpollFileObject::new_file(&current_task);
927            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
928
929            let ops1 = ControlledEventsFile::new(FdEvents::POLLIN);
930            let file1 = anon_test_file(&current_task, Box::new(ops1.clone()), OpenFlags::RDWR);
931            let ops2 = ControlledEventsFile::new(FdEvents::POLLIN);
932            let file2 = anon_test_file(&current_task, Box::new(ops2.clone()), OpenFlags::RDWR);
933
934            epoll_file
935                .add(
936                    &current_task,
937                    &file1,
938                    &epoll_file_handle,
939                    EpollEvent::new(FdEvents::POLLIN | FdEvents::EPOLLWAKEUP, 1),
940                )
941                .unwrap();
942            epoll_file
943                .add(
944                    &current_task,
945                    &file2,
946                    &epoll_file_handle,
947                    EpollEvent::new(FdEvents::POLLIN | FdEvents::EPOLLWAKEUP, 2),
948                )
949                .unwrap();
950
951            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
952            assert_eq!(events.len(), 2);
953            assert!(!current_task.kernel().suspend_resume_manager.lock().can_suspend());
954
955            ops1.set_events(error!(EPERM));
956            ops2.set_events(error!(EPERM));
957
958            assert!(epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).is_err());
959            assert!(current_task.kernel().suspend_resume_manager.lock().can_suspend());
960        })
961        .await;
962    }
963
964    #[::fuchsia::test]
965    async fn test_level_triggered_remains_ready_until_drained() {
966        // Validates standard level-triggered epoll semantics:
967        // A level-triggered file remains ready across consecutive wait() calls as long
968        // as data remains in its buffer. Once drained, wait() returns 0 events.
969        // Once new data arrives, wait() returns ready again.
970        spawn_kernel_and_run(async |current_task| {
971            let kernel = current_task.kernel();
972            register_pipe_fs(kernel.expando.get::<FsRegistry>().as_ref());
973
974            let (pipe_out, pipe_in) = new_pipe(&current_task).unwrap();
975            let epoll_file_handle = EpollFileObject::new_file(&current_task);
976            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
977            const EVENT_DATA: u64 = 100;
978            epoll_file
979                .add(
980                    &current_task,
981                    &pipe_out,
982                    &epoll_file_handle,
983                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
984                )
985                .unwrap();
986
987            let key: ReadyItemKey = pipe_out.id.as_epoll_key().into();
988
989            // Write 2 bytes.
990            pipe_in.write(&current_task, &mut VecInputBuffer::new(&[1, 2])).unwrap();
991
992            // First wait: harvests ready event.
993            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
994            assert_eq!(events.len(), 1);
995            assert_eq!(events[0].data(), EVENT_DATA);
996
997            // Level-triggered files are added to recheck_list so subsequent waits recheck readiness.
998            assert!(epoll_file.state.lock().recheck_list.contains(&key));
999
1000            // Second wait without reading data: the file is STILL READY and must be returned again.
1001            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1002            assert_eq!(events.len(), 1);
1003            assert_eq!(events[0].data(), EVENT_DATA);
1004
1005            // Drain the pipe completely.
1006            let mut out_buf = VecOutputBuffer::new(10);
1007            assert_eq!(pipe_out.read(&current_task, &mut out_buf).unwrap(), 2);
1008
1009            // Third wait: the file is no longer ready, wait returns 0 events.
1010            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1011            assert_eq!(events.len(), 0);
1012
1013            // Write new data: the wait fires and returns the event.
1014            pipe_in.write(&current_task, &mut VecInputBuffer::new(&[3])).unwrap();
1015            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1016            assert_eq!(events.len(), 1);
1017            assert_eq!(events[0].data(), EVENT_DATA);
1018        })
1019        .await;
1020    }
1021
1022    #[::fuchsia::test]
1023    async fn test_nested_epoll_delayed_drain() {
1024        // In nested epoll (ep_outer -> ep_inner -> target), when target has data:
1025        // 1. ep_inner.wait() returns target.
1026        // 2. Target is not drained immediately (delayed drain).
1027        // 3. ep_outer.wait() observes ep_inner is ready.
1028        // 4. Target is then drained by user space.
1029        // 5. Remote peer writes new data to target.
1030        // 6. ep_outer and ep_inner MUST both receive notifications for the new data and not hang.
1031        spawn_kernel_and_run(async |current_task| {
1032            let kernel = current_task.kernel();
1033            register_pipe_fs(kernel.expando.get::<FsRegistry>().as_ref());
1034
1035            let (pipe_out, pipe_in) = new_pipe(&current_task).unwrap();
1036            let epoll_inner_handle = EpollFileObject::new_file(&current_task);
1037            let epoll_inner = epoll_inner_handle.downcast_file::<EpollFileObject>().unwrap();
1038            let epoll_outer_handle = EpollFileObject::new_file(&current_task);
1039            let epoll_outer = epoll_outer_handle.downcast_file::<EpollFileObject>().unwrap();
1040
1041            const INNER_DATA: u64 = 1;
1042            const OUTER_DATA: u64 = 2;
1043
1044            epoll_inner
1045                .add(
1046                    &current_task,
1047                    &pipe_out,
1048                    &epoll_inner_handle,
1049                    EpollEvent::new(FdEvents::POLLIN, INNER_DATA),
1050                )
1051                .unwrap();
1052            epoll_outer
1053                .add(
1054                    &current_task,
1055                    &epoll_inner_handle,
1056                    &epoll_outer_handle,
1057                    EpollEvent::new(FdEvents::POLLIN, OUTER_DATA),
1058                )
1059                .unwrap();
1060
1061            // Step 1: Write initial data to pipe.
1062            pipe_in.write(&current_task, &mut VecInputBuffer::new(&[1, 2])).unwrap();
1063
1064            // Step 2: Harvest in ep_inner.wait().
1065            let inner_events =
1066                epoll_inner.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1067            assert_eq!(inner_events.len(), 1);
1068
1069            // Step 3: Do NOT drain pipe_out yet. It is still ready.
1070            // Step 4: ep_outer waits on ep_inner.
1071            let outer_events =
1072                epoll_outer.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1073            assert_eq!(outer_events.len(), 1);
1074            assert_eq!(outer_events[0].data(), OUTER_DATA);
1075
1076            // Step 5: User drains pipe_out now.
1077            let mut out_buf = VecOutputBuffer::new(10);
1078            assert_eq!(pipe_out.read(&current_task, &mut out_buf).unwrap(), 2);
1079
1080            // Step 6: Remote peer sends new data to pipe_out.
1081            pipe_in.write(&current_task, &mut VecInputBuffer::new(&[3, 4])).unwrap();
1082
1083            // Step 7: ep_outer must observe the new event without hanging.
1084            let outer_events2 =
1085                epoll_outer.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1086            assert_eq!(outer_events2.len(), 1);
1087            assert_eq!(outer_events2[0].data(), OUTER_DATA);
1088
1089            // ep_inner also sees the new event.
1090            let inner_events2 =
1091                epoll_inner.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1092            assert_eq!(inner_events2.len(), 1);
1093            assert_eq!(inner_events2[0].data(), INNER_DATA);
1094        })
1095        .await;
1096    }
1097
1098    #[::fuchsia::test]
1099    async fn test_epoll_wait_with_external_waiter_preserves_wait_canceler() {
1100        // Validates that external waiters registered via wait_async on an epoll object
1101        // correctly receive event notifications when target events occur.
1102        spawn_kernel_and_run(async |current_task| {
1103            let kernel = current_task.kernel();
1104            register_pipe_fs(kernel.expando.get::<FsRegistry>().as_ref());
1105
1106            let (pipe_out, pipe_in) = new_pipe(&current_task).unwrap();
1107            let epoll_file_handle = EpollFileObject::new_file(&current_task);
1108            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
1109            const EVENT_DATA: u64 = 42;
1110            epoll_file
1111                .add(
1112                    &current_task,
1113                    &pipe_out,
1114                    &epoll_file_handle,
1115                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
1116                )
1117                .unwrap();
1118
1119            pipe_in.write(&current_task, &mut VecInputBuffer::new(&[1])).unwrap();
1120
1121            // Register an external waiter on epoll_file.
1122            let waiter = Waiter::new();
1123            let canceler = epoll_file_handle
1124                .wait_async(&current_task, &waiter, FdEvents::POLLIN, EventHandler::None)
1125                .expect("wait_async on epoll_file");
1126
1127            assert!(!epoll_file.waiters.is_empty());
1128
1129            // First wait on epoll_file while external waiter is registered.
1130            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1131            assert_eq!(events.len(), 1);
1132
1133            // Cancel the external waiter.
1134            canceler.cancel();
1135            assert!(epoll_file.waiters.is_empty());
1136
1137            // Second wait continues to function properly.
1138            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1139            assert_eq!(events.len(), 1);
1140        })
1141        .await;
1142    }
1143
1144    #[::fuchsia::test]
1145    async fn test_multithreaded_epoll_concurrent_recheck() {
1146        // Validates multithreaded behavior when multiple workers use the same epoll instance.
1147        spawn_kernel_and_run(async |current_task| {
1148            let kernel = current_task.kernel();
1149            register_pipe_fs(kernel.expando.get::<FsRegistry>().as_ref());
1150
1151            let (pipe_out, pipe_in) = new_pipe(&current_task).unwrap();
1152            let epoll_file_handle = EpollFileObject::new_file(&current_task);
1153            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
1154            const EVENT_DATA: u64 = 42;
1155            epoll_file
1156                .add(
1157                    &current_task,
1158                    &pipe_out,
1159                    &epoll_file_handle,
1160                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
1161                )
1162                .unwrap();
1163
1164            // Write 2 bytes.
1165            pipe_in.write(&current_task, &mut VecInputBuffer::new(&[1, 2])).unwrap();
1166
1167            // Worker 1 harvests the event.
1168            let events1 = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1169            assert_eq!(events1.len(), 1);
1170
1171            // Worker 2 calls wait() without any read occurring yet:
1172            // Worker 2's do_recheck observes data is still present and returns it.
1173            let events2 = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1174            assert_eq!(events2.len(), 1);
1175
1176            // Now worker 1 drains the pipe.
1177            let mut buf = VecOutputBuffer::new(10);
1178            assert_eq!(pipe_out.read(&current_task, &mut buf).unwrap(), 2);
1179
1180            // Worker 2 calls wait():
1181            // Pre-sleep do_recheck finds pipe empty, returns 0 events.
1182            let events3 = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1183            assert_eq!(events3.len(), 0);
1184
1185            // Remote peer writes new data: wait fires for Worker 1.
1186            pipe_in.write(&current_task, &mut VecInputBuffer::new(&[3])).unwrap();
1187            let events4 = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1188            assert_eq!(events4.len(), 1);
1189        })
1190        .await;
1191    }
1192
1193    #[::fuchsia::test]
1194    async fn test_epoll_query_events_checks_recheck_list() {
1195        // When an outer epoll checks whether an inner epoll is ready, query_events()
1196        // inspects state.recheck_list directly if processing_list and trigger_list are empty.
1197        spawn_kernel_and_run(async |current_task| {
1198            let kernel = current_task.kernel();
1199            register_pipe_fs(kernel.expando.get::<FsRegistry>().as_ref());
1200
1201            let (pipe_out, pipe_in) = new_pipe(&current_task).unwrap();
1202            let epoll_file_handle = EpollFileObject::new_file(&current_task);
1203            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
1204            epoll_file
1205                .add(
1206                    &current_task,
1207                    &pipe_out,
1208                    &epoll_file_handle,
1209                    EpollEvent::new(FdEvents::POLLIN, 42),
1210                )
1211                .unwrap();
1212
1213            pipe_in.write(&current_task, &mut VecInputBuffer::new(&[1, 2])).unwrap();
1214
1215            // Harvesting the event leaves pipe_out on recheck_list.
1216            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1217            assert_eq!(events.len(), 1);
1218
1219            // query_events on the epoll file checks recheck_list and returns POLLIN.
1220            let queried = epoll_file_handle.query_events(&current_task).unwrap();
1221            assert!(queried.contains(FdEvents::POLLIN));
1222
1223            // Drain the pipe.
1224            let mut buf = VecOutputBuffer::new(10);
1225            assert_eq!(pipe_out.read(&current_task, &mut buf).unwrap(), 2);
1226
1227            // Now query_events sees the file on recheck_list is empty, so returns empty events.
1228            let queried = epoll_file_handle.query_events(&current_task).unwrap();
1229            assert!(!queried.contains(FdEvents::POLLIN));
1230        })
1231        .await;
1232    }
1233
1234    #[::fuchsia::test]
1235    async fn test_recheck_list_query_error_preserves_remaining_keys() {
1236        // In wait(), if do_recheck on one item fails with an error (e.g. query_events returns Err),
1237        // remaining unprocessed items must not be dropped from recheck_list.
1238        spawn_kernel_and_run(async |current_task| {
1239            let epoll_file_handle = EpollFileObject::new_file(&current_task);
1240            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
1241
1242            let ops1 = ControlledEventsFile::new(FdEvents::POLLIN);
1243            let file1 = anon_test_file(&current_task, Box::new(ops1.clone()), OpenFlags::RDWR);
1244            let ops2 = ControlledEventsFile::new(FdEvents::POLLIN);
1245            let file2 = anon_test_file(&current_task, Box::new(ops2.clone()), OpenFlags::RDWR);
1246
1247            let _key1: ReadyItemKey = file1.id.as_epoll_key().into();
1248            let key2: ReadyItemKey = file2.id.as_epoll_key().into();
1249
1250            epoll_file
1251                .add(
1252                    &current_task,
1253                    &file1,
1254                    &epoll_file_handle,
1255                    EpollEvent::new(FdEvents::POLLIN, 1),
1256                )
1257                .unwrap();
1258            epoll_file
1259                .add(
1260                    &current_task,
1261                    &file2,
1262                    &epoll_file_handle,
1263                    EpollEvent::new(FdEvents::POLLIN, 2),
1264                )
1265                .unwrap();
1266
1267            // First wait: both files are ready and returned.
1268            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1269            assert_eq!(events.len(), 2);
1270
1271            // Both files are on recheck_list.
1272            assert_eq!(epoll_file.state.lock().recheck_list.len(), 2);
1273
1274            // Make file1 return an error on query_events.
1275            ops1.set_events(error!(EPERM));
1276
1277            // Second wait: do_recheck fails on file1 and returns Err.
1278            assert!(epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).is_err());
1279
1280            // Unprocessed items must be preserved on recheck_list.
1281            assert!(
1282                epoll_file.state.lock().recheck_list.contains(&key2),
1283                "file2 should remain on recheck_list despite file1's error"
1284            );
1285
1286            // Delete file1 so it does not fail subsequent calls.
1287            epoll_file.delete(&current_task, &file1).unwrap();
1288
1289            // Third wait: file2 is still ready and should be returned now.
1290            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1291            assert_eq!(events.len(), 1);
1292            assert_eq!(events[0].data(), 2);
1293        })
1294        .await;
1295    }
1296}