1use 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
24const MAX_NESTED_DEPTH: u32 = 5;
27
28struct 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 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
57pub type EpollKey = usize;
60
61#[derive(Default)]
64pub struct EpollFileObject {
65 waiter: Waiter,
66 state: LockDepMutex<EpollState, EpollStateLock>,
68 waitable_state: Arc<LockDepMutex<EpollWaitableState, EpollWaitableStateLock>>,
69 waiters: Arc<WaitQueue>,
72}
73
74#[derive(Default)]
75struct EpollState {
76 wait_objects: HashMap<ReadyItemKey, WaitObject>,
79 processing_list: VecDeque<ReadyItem>,
87 recheck_list: Vec<ReadyItemKey>,
93}
94
95#[derive(Default)]
96struct EpollWaitableState {
97 trigger_list: VecDeque<ReadyItem>,
100}
101
102impl EpollFileObject {
103 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 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 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 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 let _token = allow_subclass();
197 child_file.check_eloop(parent, depth_left - 1)?;
198 }
199
200 Ok(())
201 }
202
203 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 if let Some(epoll_to_add) = file.downcast_file::<EpollFileObject>() {
213 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 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 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 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 wait_object.deactivate_wakeup_source(current_task);
276 Ok(())
277 } else {
278 error!(ENOENT)
279 }
280 }
281
282 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 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 if let Some(target) = wait.target.upgrade() {
308 let events = {
309 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 self.wait_on_file(current_task, pending.key, wait)?;
320 }
321 }
322 }
323 }
324 }
325 Ok(())
326 }
327
328 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; }
344
345 if !pending_list.is_empty() {
346 wait_deadline = zx::MonotonicInstant::ZERO;
350 }
351
352 match self.waiter.wait_until(current_task, wait_deadline) {
366 Err(err) if err == ETIMEDOUT => break,
367 Err(err) if err == EINTR => {
368 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 result => result?,
382 }
383 }
384
385 Ok(pending_list)
386 }
387
388 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 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 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 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 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 if !wait.events.contains(FdEvents::EPOLLET) {
442 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 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(¤t_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(¤t_task);
573 let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
574 epoll_file
575 .add(
576 ¤t_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(¤t_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(¤t_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(¤t_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(¤t_task, &mut VecInputBuffer::new(test_bytes)).unwrap(),
627 test_bytes.len()
628 );
629
630 let epoll_file_handle = EpollFileObject::new_file(¤t_task);
631 let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
632 epoll_file
633 .add(
634 ¤t_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(¤t_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(¤t_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(¤t_task, 0, EventFdType::Counter, true);
661 let waiter = Waiter::new();
662
663 let epoll_file_handle = EpollFileObject::new_file(¤t_task);
664 let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
665 const EVENT_DATA: u64 = 42;
666 epoll_file
667 .add(
668 ¤t_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(¤t_task, &event).unwrap();
677 }
678
679 let wait_canceler = event
680 .wait_async(¤t_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(¤t_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(¤t_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(¤t_task, client1, OpenFlags::RDWR)
717 .expect("create_fuchsia_pipe");
718 let pipe2 = create_fuchsia_pipe(¤t_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(¤t_task);
725 let epoll_file = epoll_object.downcast_file::<EpollFileObject>().unwrap();
726 epoll_file
727 .add(¤t_task, &pipe1, &epoll_object, EpollEvent::new(FdEvents::POLLIN, 1))
728 .expect("epoll_file.add");
729 epoll_file
730 .add(¤t_task, &pipe2, &epoll_object, EpollEvent::new(FdEvents::POLLIN, 2))
731 .expect("epoll_file.add");
732 epoll_file.wait(¤t_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(¤t_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(¤t_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(¤t_task, 0, EventFdType::Counter, true);
767 let epoll_file_handle = EpollFileObject::new_file(¤t_task);
768 let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
769
770 const EVENT_DATA: u64 = 42;
772 epoll_file
773 .add(
774 ¤t_task,
775 &event,
776 &epoll_file_handle,
777 EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
778 )
779 .unwrap();
780
781 let add_val = 1u64;
783 assert_eq!(
784 event
785 .write(¤t_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(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap().len(),
792 1
793 );
794
795 epoll_file.delete(¤t_task, &event).unwrap();
797
798 assert_eq!(
800 epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap().len(),
801 0
802 );
803 })
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 ¤t_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(¤t_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 ¤t_task,
826 &socket1,
827 &epoll_file_handle,
828 EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
829 )
830 .unwrap();
831 assert_eq!(
832 epoll_file.wait(¤t_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(¤t_task, &socket1, EpollEvent::new(read_write_event, EVENT_DATA))
839 .unwrap();
840 let triggered_events =
841 epoll_file.wait(¤t_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(¤t_task, 0, EventFdType::Counter, true);
854 let epoll_file_handle = EpollFileObject::new_file(¤t_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 ¤t_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 #[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(¤t_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(¤t_task, Box::new(ops1.clone()), OpenFlags::RDWR);
931 let ops2 = ControlledEventsFile::new(FdEvents::POLLIN);
932 let file2 = anon_test_file(¤t_task, Box::new(ops2.clone()), OpenFlags::RDWR);
933
934 epoll_file
935 .add(
936 ¤t_task,
937 &file1,
938 &epoll_file_handle,
939 EpollEvent::new(FdEvents::POLLIN | FdEvents::EPOLLWAKEUP, 1),
940 )
941 .unwrap();
942 epoll_file
943 .add(
944 ¤t_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(¤t_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(¤t_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 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(¤t_task).unwrap();
975 let epoll_file_handle = EpollFileObject::new_file(¤t_task);
976 let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
977 const EVENT_DATA: u64 = 100;
978 epoll_file
979 .add(
980 ¤t_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 pipe_in.write(¤t_task, &mut VecInputBuffer::new(&[1, 2])).unwrap();
991
992 let events = epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
994 assert_eq!(events.len(), 1);
995 assert_eq!(events[0].data(), EVENT_DATA);
996
997 assert!(epoll_file.state.lock().recheck_list.contains(&key));
999
1000 let events = epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1002 assert_eq!(events.len(), 1);
1003 assert_eq!(events[0].data(), EVENT_DATA);
1004
1005 let mut out_buf = VecOutputBuffer::new(10);
1007 assert_eq!(pipe_out.read(¤t_task, &mut out_buf).unwrap(), 2);
1008
1009 let events = epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1011 assert_eq!(events.len(), 0);
1012
1013 pipe_in.write(¤t_task, &mut VecInputBuffer::new(&[3])).unwrap();
1015 let events = epoll_file.wait(¤t_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 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(¤t_task).unwrap();
1036 let epoll_inner_handle = EpollFileObject::new_file(¤t_task);
1037 let epoll_inner = epoll_inner_handle.downcast_file::<EpollFileObject>().unwrap();
1038 let epoll_outer_handle = EpollFileObject::new_file(¤t_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 ¤t_task,
1047 &pipe_out,
1048 &epoll_inner_handle,
1049 EpollEvent::new(FdEvents::POLLIN, INNER_DATA),
1050 )
1051 .unwrap();
1052 epoll_outer
1053 .add(
1054 ¤t_task,
1055 &epoll_inner_handle,
1056 &epoll_outer_handle,
1057 EpollEvent::new(FdEvents::POLLIN, OUTER_DATA),
1058 )
1059 .unwrap();
1060
1061 pipe_in.write(¤t_task, &mut VecInputBuffer::new(&[1, 2])).unwrap();
1063
1064 let inner_events =
1066 epoll_inner.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1067 assert_eq!(inner_events.len(), 1);
1068
1069 let outer_events =
1072 epoll_outer.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1073 assert_eq!(outer_events.len(), 1);
1074 assert_eq!(outer_events[0].data(), OUTER_DATA);
1075
1076 let mut out_buf = VecOutputBuffer::new(10);
1078 assert_eq!(pipe_out.read(¤t_task, &mut out_buf).unwrap(), 2);
1079
1080 pipe_in.write(¤t_task, &mut VecInputBuffer::new(&[3, 4])).unwrap();
1082
1083 let outer_events2 =
1085 epoll_outer.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1086 assert_eq!(outer_events2.len(), 1);
1087 assert_eq!(outer_events2[0].data(), OUTER_DATA);
1088
1089 let inner_events2 =
1091 epoll_inner.wait(¤t_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 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(¤t_task).unwrap();
1107 let epoll_file_handle = EpollFileObject::new_file(¤t_task);
1108 let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
1109 const EVENT_DATA: u64 = 42;
1110 epoll_file
1111 .add(
1112 ¤t_task,
1113 &pipe_out,
1114 &epoll_file_handle,
1115 EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
1116 )
1117 .unwrap();
1118
1119 pipe_in.write(¤t_task, &mut VecInputBuffer::new(&[1])).unwrap();
1120
1121 let waiter = Waiter::new();
1123 let canceler = epoll_file_handle
1124 .wait_async(¤t_task, &waiter, FdEvents::POLLIN, EventHandler::None)
1125 .expect("wait_async on epoll_file");
1126
1127 assert!(!epoll_file.waiters.is_empty());
1128
1129 let events = epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1131 assert_eq!(events.len(), 1);
1132
1133 canceler.cancel();
1135 assert!(epoll_file.waiters.is_empty());
1136
1137 let events = epoll_file.wait(¤t_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 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(¤t_task).unwrap();
1152 let epoll_file_handle = EpollFileObject::new_file(¤t_task);
1153 let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
1154 const EVENT_DATA: u64 = 42;
1155 epoll_file
1156 .add(
1157 ¤t_task,
1158 &pipe_out,
1159 &epoll_file_handle,
1160 EpollEvent::new(FdEvents::POLLIN, EVENT_DATA),
1161 )
1162 .unwrap();
1163
1164 pipe_in.write(¤t_task, &mut VecInputBuffer::new(&[1, 2])).unwrap();
1166
1167 let events1 = epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1169 assert_eq!(events1.len(), 1);
1170
1171 let events2 = epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1174 assert_eq!(events2.len(), 1);
1175
1176 let mut buf = VecOutputBuffer::new(10);
1178 assert_eq!(pipe_out.read(¤t_task, &mut buf).unwrap(), 2);
1179
1180 let events3 = epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1183 assert_eq!(events3.len(), 0);
1184
1185 pipe_in.write(¤t_task, &mut VecInputBuffer::new(&[3])).unwrap();
1187 let events4 = epoll_file.wait(¤t_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 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(¤t_task).unwrap();
1202 let epoll_file_handle = EpollFileObject::new_file(¤t_task);
1203 let epoll_file = epoll_file_handle.downcast_file::<EpollFileObject>().unwrap();
1204 epoll_file
1205 .add(
1206 ¤t_task,
1207 &pipe_out,
1208 &epoll_file_handle,
1209 EpollEvent::new(FdEvents::POLLIN, 42),
1210 )
1211 .unwrap();
1212
1213 pipe_in.write(¤t_task, &mut VecInputBuffer::new(&[1, 2])).unwrap();
1214
1215 let events = epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1217 assert_eq!(events.len(), 1);
1218
1219 let queried = epoll_file_handle.query_events(¤t_task).unwrap();
1221 assert!(queried.contains(FdEvents::POLLIN));
1222
1223 let mut buf = VecOutputBuffer::new(10);
1225 assert_eq!(pipe_out.read(¤t_task, &mut buf).unwrap(), 2);
1226
1227 let queried = epoll_file_handle.query_events(¤t_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 spawn_kernel_and_run(async |current_task| {
1239 let epoll_file_handle = EpollFileObject::new_file(¤t_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(¤t_task, Box::new(ops1.clone()), OpenFlags::RDWR);
1244 let ops2 = ControlledEventsFile::new(FdEvents::POLLIN);
1245 let file2 = anon_test_file(¤t_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 ¤t_task,
1253 &file1,
1254 &epoll_file_handle,
1255 EpollEvent::new(FdEvents::POLLIN, 1),
1256 )
1257 .unwrap();
1258 epoll_file
1259 .add(
1260 ¤t_task,
1261 &file2,
1262 &epoll_file_handle,
1263 EpollEvent::new(FdEvents::POLLIN, 2),
1264 )
1265 .unwrap();
1266
1267 let events = epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1269 assert_eq!(events.len(), 2);
1270
1271 assert_eq!(epoll_file.state.lock().recheck_list.len(), 2);
1273
1274 ops1.set_events(error!(EPERM));
1276
1277 assert!(epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).is_err());
1279
1280 assert!(
1282 epoll_file.state.lock().recheck_list.contains(&key2),
1283 "file2 should remain on recheck_list despite file1's error"
1284 );
1285
1286 epoll_file.delete(¤t_task, &file1).unwrap();
1288
1289 let events = epoll_file.wait(¤t_task, 10, zx::MonotonicInstant::ZERO).unwrap();
1291 assert_eq!(events.len(), 1);
1292 assert_eq!(events[0].data(), 2);
1293 })
1294 .await;
1295 }
1296}