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 key in recheck_list {
403                let wait_object = state.wait_objects.get_mut(&key).unwrap();
404                self.do_recheck(current_task, wait_object, key)?;
405            }
406        }
407
408        let pending_list = self.wait_until_pending_event(current_task, max_events, deadline)?;
409
410        // Process the pending list and add processed ReadyItem
411        // entries to the rearm_list for the next wait.
412        let mut result = vec![];
413        let mut state = self.state.lock();
414        let state = &mut *state;
415        for pending_event in pending_list.iter().unique_by(|e| e.key) {
416            // The wait could have been deleted by here,
417            // so ignore the None case.
418            let Some(wait) = state.wait_objects.get_mut(&pending_event.key) else { continue };
419
420            let reported_events = pending_event.events & wait.events;
421            result.push(EpollEvent::new(reported_events, wait.data));
422
423            // Files marked with `EPOLLONESHOT` should only notify
424            // once and need to be rearmed manually with epoll_ctl_mod().
425            if wait.events.contains(FdEvents::EPOLLONESHOT) {
426                continue;
427            }
428
429            self.wait_on_file_edge_triggered(current_task, pending_event.key, wait)?;
430
431            if !wait.events.contains(FdEvents::EPOLLET) {
432                state.recheck_list.push(pending_event.key);
433            }
434
435            // TODO: is this really only supposed to happen for level-triggered events?
436            if !wait.events.contains(FdEvents::EPOLLET) {
437                // When this is the first time epoll_wait on this epoll fd, create and
438                // hold a wake lease until the next epoll_wait.
439                if wait.events.contains(FdEvents::EPOLLWAKEUP) {
440                    if let ReadyItemKey::Usize(key) = pending_event.key {
441                        let origin = WakeupSourceOrigin::Epoll(key);
442                        current_task
443                            .kernel()
444                            .suspend_resume_manager
445                            .activate_wakeup_source_with_actor(
446                                origin.clone(),
447                                Some(current_task.command()),
448                            );
449                        wait.active_wakeup_source = Some(origin);
450                    }
451                }
452            }
453        }
454
455        Ok(result)
456    }
457}
458
459impl FileOps for EpollFileObject {
460    fileops_impl_nonseekable!();
461    fileops_impl_noop_sync!();
462    fileops_impl_dataless!();
463
464    fn wait_async(
465        &self,
466        _file: &FileObject,
467        _current_task: &CurrentTask,
468        waiter: &Waiter,
469        events: FdEvents,
470        handler: EventHandler,
471    ) -> Option<WaitCanceler> {
472        Some(self.waiters.wait_async_fd_events(waiter, events, handler))
473    }
474
475    fn query_events(
476        &self,
477        _file: &FileObject,
478        current_task: &CurrentTask,
479    ) -> Result<FdEvents, Errno> {
480        let mut events = FdEvents::empty();
481        let state = self.state.lock();
482        if !state.processing_list.is_empty() || !self.waitable_state.lock().trigger_list.is_empty()
483        {
484            events |= FdEvents::POLLIN;
485        } else {
486            for key in &state.recheck_list {
487                let wait_object = state.wait_objects.get(key).unwrap();
488                let Some(target) = wait_object.target() else { continue };
489                // Target might be itself an epoll object. Because there is no loop,
490                // this allow_subclass is safe.
491                let _token = allow_subclass();
492                if !(target.query_events(current_task)? & wait_object.events).is_empty() {
493                    events |= FdEvents::POLLIN;
494                    break;
495                }
496            }
497        }
498        Ok(events)
499    }
500
501    fn close(self: Box<Self>, _file: &FileObjectState, current_task: &CurrentTask) {
502        let mut guard = self.state.lock();
503        for wait_object in guard.wait_objects.values_mut() {
504            wait_object.deactivate_wakeup_source(current_task);
505        }
506    }
507}
508
509#[derive(Clone)]
510pub struct EpollEventHandler {
511    key: ReadyItemKey,
512    waitable_state: Arc<LockDepMutex<EpollWaitableState, EpollWaitableStateLock>>,
513    waiters: Arc<WaitQueue>,
514}
515
516impl EpollEventHandler {
517    pub fn handle(self, events: FdEvents) {
518        {
519            let mut waitable_state = self.waitable_state.lock();
520            waitable_state.trigger_list.push_back(ReadyItem { key: self.key, events });
521        }
522        self.waiters.notify_fd_events(FdEvents::POLLIN);
523    }
524}
525
526impl EpollFileObject {
527    fn new_wait_handler(&self, key: ReadyItemKey) -> EventHandler {
528        EventHandler::Epoll(EpollEventHandler {
529            key,
530            waitable_state: Arc::clone(&self.waitable_state),
531            waiters: Arc::clone(&self.waiters),
532        })
533    }
534}
535
536#[cfg(test)]
537mod tests {
538    use super::*;
539    use crate::fs::fuchsia::create_fuchsia_pipe;
540    use crate::task::Waiter;
541    use crate::task::dynamic_thread_spawner::SpawnRequestBuilder;
542    use crate::testing::{anon_test_file, spawn_kernel_and_run};
543    use crate::vfs::buffers::{VecInputBuffer, VecOutputBuffer};
544    use crate::vfs::eventfd::{EventFdType, new_eventfd};
545    use crate::vfs::fs_registry::FsRegistry;
546    use crate::vfs::pipe::{new_pipe, register_pipe_fs};
547    use crate::vfs::socket::{SocketDomain, SocketType, UnixSocket};
548    use starnix_lifecycle::AtomicCounter;
549    use starnix_sync::Mutex;
550    use starnix_uapi::vfs::{EpollEvent, FdEvents};
551    use syncio::Zxio;
552
553    #[::fuchsia::test]
554    async fn test_epoll_read_ready() {
555        static WRITE_COUNT: AtomicCounter<usize> = AtomicCounter::<usize>::new_const(0);
556        const EVENT_DATA: u64 = 42;
557
558        spawn_kernel_and_run(async |current_task| {
559            let kernel = current_task.kernel();
560            register_pipe_fs(kernel.expando.get::<FsRegistry>().as_ref());
561
562            let (pipe_out, pipe_in) = new_pipe(&current_task).unwrap();
563
564            let test_string = "hello starnix".to_string();
565            let test_len = test_string.len();
566
567            let epoll_file_handle = EpollFileObject::new_file(&current_task);
568            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
569            epoll_file
570                .add(
571                    &current_task,
572                    &pipe_out,
573                    &epoll_file_handle,
574                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
575                )
576                .unwrap();
577
578            let (sender, receiver) = std::sync::mpsc::channel();
579            let value = test_string.clone();
580            let closure = move |task: &CurrentTask| {
581                let bytes_written =
582                    pipe_in.write(&task, &mut VecInputBuffer::new(value.as_bytes())).unwrap();
583                assert_eq!(bytes_written, test_len);
584                WRITE_COUNT.add(bytes_written);
585                sender.send(()).unwrap();
586            };
587            let req = SpawnRequestBuilder::new().with_sync_closure(closure).build();
588            kernel.kthreads.spawner().spawn_from_request(req);
589            let events =
590                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::INFINITE).unwrap();
591            receiver.recv().unwrap();
592            assert_eq!(1, events.len());
593            let event = &events[0];
594            assert!(event.events().contains(FdEvents::POLLIN));
595            assert_eq!(event.data(), EVENT_DATA);
596
597            let mut buffer = VecOutputBuffer::new(test_len);
598            let bytes_read = pipe_out.read(&current_task, &mut buffer).unwrap();
599            assert_eq!(bytes_read, WRITE_COUNT.get());
600            assert_eq!(bytes_read, test_len);
601            assert_eq!(buffer.data(), test_string.as_bytes());
602        })
603        .await;
604    }
605
606    #[::fuchsia::test]
607    async fn test_epoll_ready_then_wait() {
608        const EVENT_DATA: u64 = 42;
609
610        spawn_kernel_and_run(async |current_task| {
611            let kernel = current_task.kernel();
612            register_pipe_fs(kernel.expando.get::<FsRegistry>().as_ref());
613
614            let (pipe_out, pipe_in) = new_pipe(&current_task).unwrap();
615
616            let test_string = "hello starnix".to_string();
617            let test_bytes = test_string.as_bytes();
618            let test_len = test_bytes.len();
619
620            assert_eq!(
621                pipe_in.write(&current_task, &mut VecInputBuffer::new(test_bytes)).unwrap(),
622                test_bytes.len()
623            );
624
625            let epoll_file_handle = EpollFileObject::new_file(&current_task);
626            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
627            epoll_file
628                .add(
629                    &current_task,
630                    &pipe_out,
631                    &epoll_file_handle,
632                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
633                )
634                .unwrap();
635
636            let events =
637                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::INFINITE).unwrap();
638            assert_eq!(1, events.len());
639            let event = &events[0];
640            assert!(event.events().contains(FdEvents::POLLIN));
641            assert_eq!(event.data(), EVENT_DATA);
642
643            let mut buffer = VecOutputBuffer::new(test_len);
644            let bytes_read = pipe_out.read(&current_task, &mut buffer).unwrap();
645            assert_eq!(bytes_read, test_len);
646            assert_eq!(buffer.data(), test_bytes);
647        })
648        .await;
649    }
650
651    #[::fuchsia::test]
652    async fn test_epoll_ctl_cancel() {
653        spawn_kernel_and_run(async |current_task| {
654            for do_cancel in [true, false] {
655                let event = new_eventfd(&current_task, 0, EventFdType::Counter, true);
656                let waiter = Waiter::new();
657
658                let epoll_file_handle = EpollFileObject::new_file(&current_task);
659                let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
660                const EVENT_DATA: u64 = 42;
661                epoll_file
662                    .add(
663                        &current_task,
664                        &event,
665                        &epoll_file_handle,
666                        EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
667                    )
668                    .unwrap();
669
670                if do_cancel {
671                    epoll_file.delete(&current_task, &event).unwrap();
672                }
673
674                let wait_canceler = event
675                    .wait_async(&current_task, &waiter, FdEvents::POLLIN, EventHandler::None)
676                    .expect("wait_async");
677                if do_cancel {
678                    wait_canceler.cancel();
679                }
680
681                let add_val = 1u64;
682
683                assert_eq!(
684                    event
685                        .write(&current_task, &mut VecInputBuffer::new(&add_val.to_ne_bytes()))
686                        .unwrap(),
687                    std::mem::size_of::<u64>()
688                );
689
690                let events =
691                    epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
692
693                if do_cancel {
694                    assert_eq!(0, events.len());
695                } else {
696                    assert_eq!(1, events.len());
697                    let event = &events[0];
698                    assert!(event.events().contains(FdEvents::POLLIN));
699                    assert_eq!(event.data(), EVENT_DATA);
700                }
701            }
702        })
703        .await;
704    }
705
706    #[::fuchsia::test]
707    async fn test_multiple_events() {
708        spawn_kernel_and_run(async |current_task| {
709            let (client1, server1) = zx::Socket::create_stream();
710            let (client2, server2) = zx::Socket::create_stream();
711            let pipe1 = create_fuchsia_pipe(&current_task, client1, OpenFlags::RDWR)
712                .expect("create_fuchsia_pipe");
713            let pipe2 = create_fuchsia_pipe(&current_task, client2, OpenFlags::RDWR)
714                .expect("create_fuchsia_pipe");
715            let server1_zxio = Zxio::create(server1.into_handle()).expect("Zxio::create");
716            let server2_zxio = Zxio::create(server2.into_handle()).expect("Zxio::create");
717
718            let poll = || {
719                let epoll_object = EpollFileObject::new_file(&current_task);
720                let epoll_file = epoll_object.downcast_file::<EpollFileObject>().unwrap();
721                epoll_file
722                    .add(&current_task, &pipe1, &epoll_object, EpollEvent::new(FdEvents::POLLIN, 1))
723                    .expect("epoll_file.add");
724                epoll_file
725                    .add(&current_task, &pipe2, &epoll_object, EpollEvent::new(FdEvents::POLLIN, 2))
726                    .expect("epoll_file.add");
727                epoll_file.wait(&current_task, 2, zx::MonotonicInstant::ZERO).expect("wait")
728            };
729
730            let fds = poll();
731            assert!(fds.is_empty());
732
733            assert_eq!(server1_zxio.write(&[0]).expect("write"), 1);
734
735            let fds = poll();
736            assert_eq!(fds.len(), 1);
737            assert_eq!(fds[0].events(), FdEvents::POLLIN);
738            assert_eq!(fds[0].data(), 1);
739            assert_eq!(pipe1.read(&current_task, &mut VecOutputBuffer::new(64)).expect("read"), 1);
740
741            let fds = poll();
742            assert!(fds.is_empty());
743
744            assert_eq!(server2_zxio.write(&[0]).expect("write"), 1);
745
746            let fds = poll();
747            assert_eq!(fds.len(), 1);
748            assert_eq!(fds[0].events(), FdEvents::POLLIN);
749            assert_eq!(fds[0].data(), 2);
750            assert_eq!(pipe2.read(&current_task, &mut VecOutputBuffer::new(64)).expect("read"), 1);
751
752            let fds = poll();
753            assert!(fds.is_empty());
754        })
755        .await;
756    }
757
758    #[::fuchsia::test]
759    async fn test_cancel_after_notify() {
760        spawn_kernel_and_run(async |current_task| {
761            let event = new_eventfd(&current_task, 0, EventFdType::Counter, true);
762            let epoll_file_handle = EpollFileObject::new_file(&current_task);
763            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
764
765            // Add a thing
766            const EVENT_DATA: u64 = 42;
767            epoll_file
768                .add(
769                    &current_task,
770                    &event,
771                    &epoll_file_handle,
772                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
773                )
774                .unwrap();
775
776            // Make the thing send a notification, wait for it
777            let add_val = 1u64;
778            assert_eq!(
779                event
780                    .write(&current_task, &mut VecInputBuffer::new(&add_val.to_ne_bytes()))
781                    .unwrap(),
782                std::mem::size_of::<u64>()
783            );
784
785            assert_eq!(
786                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap().len(),
787                1
788            );
789
790            // Remove the thing
791            epoll_file.delete(&current_task, &event).unwrap();
792
793            // Wait for new notifications
794            assert_eq!(
795                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap().len(),
796                0
797            );
798            // That shouldn't crash
799        })
800        .await;
801    }
802
803    #[::fuchsia::test]
804    async fn test_add_then_modify() {
805        spawn_kernel_and_run(async |current_task| {
806            let (socket1, _socket2) = UnixSocket::new_pair(
807                &current_task,
808                SocketDomain::Unix,
809                SocketType::Stream,
810                OpenFlags::RDWR,
811            )
812            .expect("Failed to create socket pair.");
813
814            let epoll_file_handle = EpollFileObject::new_file(&current_task);
815            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
816
817            const EVENT_DATA: u64 = 42;
818            epoll_file
819                .add(
820                    &current_task,
821                    &socket1,
822                    &epoll_file_handle,
823                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
824                )
825                .unwrap();
826            assert_eq!(
827                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap().len(),
828                0
829            );
830
831            let read_write_event = FdEvents::POLLIN | FdEvents::POLLOUT;
832            epoll_file
833                .modify(&current_task, &socket1, EpollEvent::new(read_write_event, EVENT_DATA))
834                .unwrap();
835            let triggered_events =
836                epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
837            assert_eq!(1, triggered_events.len());
838            let event = &triggered_events[0];
839            assert_eq!(event.events(), FdEvents::POLLOUT);
840            assert_eq!(event.data(), EVENT_DATA);
841        })
842        .await;
843    }
844
845    #[::fuchsia::test]
846    async fn test_waiter_removal() {
847        spawn_kernel_and_run(async |current_task| {
848            let event = new_eventfd(&current_task, 0, EventFdType::Counter, true);
849            let epoll_file_handle = EpollFileObject::new_file(&current_task);
850            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
851
852            const EVENT_DATA: u64 = 42;
853            epoll_file
854                .add(
855                    &current_task,
856                    &event,
857                    &epoll_file_handle,
858                    EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
859                )
860                .unwrap();
861
862            std::mem::drop(event);
863
864            assert!(epoll_file.waiters.is_empty());
865        })
866        .await;
867    }
868
869    /// A test file whose event query results can be dynamically controlled.
870    ///
871    /// This allows tests to simulate scenarios where `query_events` succeeds initially
872    /// (e.g., during epoll registration and the first `wait`), but returns an error on a
873    /// subsequent call (e.g., when `do_recheck` re-evaluates `recheck_list`).
874    #[derive(Clone)]
875    struct ControlledEventsFile {
876        wait_queue: Arc<WaitQueue>,
877        events: Arc<Mutex<Result<FdEvents, Errno>>>,
878    }
879
880    impl ControlledEventsFile {
881        fn new(events: FdEvents) -> Self {
882            Self {
883                wait_queue: Arc::new(WaitQueue::default()),
884                events: Arc::new(Mutex::new(Ok(events))),
885            }
886        }
887
888        fn set_events(&self, events: Result<FdEvents, Errno>) {
889            *self.events.lock() = events;
890        }
891    }
892
893    impl FileOps for ControlledEventsFile {
894        fileops_impl_nonseekable!();
895        fileops_impl_noop_sync!();
896        fileops_impl_dataless!();
897
898        fn wait_async(
899            &self,
900            _file: &FileObject,
901            _current_task: &CurrentTask,
902            waiter: &Waiter,
903            events: FdEvents,
904            handler: EventHandler,
905        ) -> Option<WaitCanceler> {
906            Some(self.wait_queue.wait_async_fd_events(waiter, events, handler))
907        }
908
909        fn query_events(
910            &self,
911            _file: &FileObject,
912            _current_task: &CurrentTask,
913        ) -> Result<FdEvents, Errno> {
914            self.events.lock().clone()
915        }
916    }
917
918    #[::fuchsia::test]
919    async fn test_epoll_recheck_failure_deactivates_wakeup_sources() {
920        spawn_kernel_and_run(async |current_task| {
921            let epoll_file_handle = EpollFileObject::new_file(&current_task);
922            let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
923
924            let ops1 = ControlledEventsFile::new(FdEvents::POLLIN);
925            let file1 = anon_test_file(&current_task, Box::new(ops1.clone()), OpenFlags::RDWR);
926            let ops2 = ControlledEventsFile::new(FdEvents::POLLIN);
927            let file2 = anon_test_file(&current_task, Box::new(ops2.clone()), OpenFlags::RDWR);
928
929            epoll_file
930                .add(
931                    &current_task,
932                    &file1,
933                    &epoll_file_handle,
934                    EpollEvent::new(FdEvents::POLLIN | FdEvents::EPOLLWAKEUP, 1),
935                )
936                .unwrap();
937            epoll_file
938                .add(
939                    &current_task,
940                    &file2,
941                    &epoll_file_handle,
942                    EpollEvent::new(FdEvents::POLLIN | FdEvents::EPOLLWAKEUP, 2),
943                )
944                .unwrap();
945
946            let events = epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).unwrap();
947            assert_eq!(events.len(), 2);
948            assert!(!current_task.kernel().suspend_resume_manager.lock().can_suspend());
949
950            ops1.set_events(error!(EPERM));
951            ops2.set_events(error!(EPERM));
952
953            assert!(epoll_file.wait(&current_task, 10, zx::MonotonicInstant::ZERO).is_err());
954            assert!(current_task.kernel().suspend_resume_manager.lock().can_suspend());
955        })
956        .await;
957    }
958}