Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions components/axpoll/src/axtest.rs
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,27 @@ fn axpoll_wakes_only_matching_interests() {
ax_assert_eq!(unsafe { poll_set.wake(IoEvents::IN | IoEvents::OUT) }, 0);
}

#[axtest]
fn axpoll_exclusive_wake_keeps_other_matching_waiters() {
let poll_set = PollSet::new();
let first_counter = WakeCounter::new();
let second_counter = WakeCounter::new();
let first_waker = counter_waker(&first_counter);
let second_waker = counter_waker(&second_counter);

unsafe {
poll_set.register(&first_waker, IoEvents::IN);
poll_set.register(&second_waker, IoEvents::IN);
}

ax_assert_eq!(unsafe { poll_set.wake_one(IoEvents::IN) }, 1);
ax_assert_eq!(first_counter.count() + second_counter.count(), 1);
ax_assert_eq!(unsafe { poll_set.wake_one(IoEvents::IN) }, 1);
ax_assert_eq!(first_counter.count(), 1);
ax_assert_eq!(second_counter.count(), 1);
ax_assert_eq!(unsafe { poll_set.wake_one(IoEvents::IN) }, 0);
}

#[axtest]
fn axpoll_capacity_overwrite_and_drop_rules_hold() {
let poll_set = PollSet::new();
Expand Down
43 changes: 43 additions & 0 deletions components/axpoll/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,24 @@ impl Inner {
}
old.cursor = 0;
}

fn take_one_ready(&mut self, ready: IoEvents) -> Option<Entry> {
let len = self.len();
let mut selected = None;
let mut keep_len = 0;

for index in 0..len {
let entry = unsafe { self.entries[index].assume_init_read() };
if selected.is_none() && entry.interests.intersects(ready) {
selected = Some(entry);
} else {
self.entries[keep_len].write(entry);
keep_len += 1;
}
}
self.cursor = keep_len;
selected
}
}

impl Drop for Inner {
Expand Down Expand Up @@ -213,6 +231,31 @@ impl PollSet {
woke
}

/// Wakes one registered waker whose interests intersect `ready`.
///
/// Matching wakers that are not selected remain registered. This is used
/// by wait queues with Linux-style exclusive wakeup semantics, where one
/// readiness transition should give one waiter a chance to consume work.
///
/// # Safety
///
/// This method is task/deferred-context only. Callers must not invoke it
/// from hard IRQ, NMI, or trap callbacks. The readiness state represented
/// by `ready` must be published before this method is called, and callers
/// must not hold locks that may be re-entered by waker execution or poll
/// wakeup paths.
pub unsafe fn wake_one(&self, ready: IoEvents) -> usize {
let Some(inner) = self.0.get() else {
return 0;
};
let ready_entry = inner.lock().take_one_ready(ready);
let Some(entry) = ready_entry else {
return 0;
};
entry.wake();
1
}

/// Wakes up registered wakers whose interests intersect `ready` from IRQ context.
///
/// Unlike [`wake`](Self::wake), this does not allocate a replacement
Expand Down
21 changes: 21 additions & 0 deletions components/axpoll/tests/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,27 @@ fn wake_only_matching_interests() {
assert_eq!(write_counter.count(), 1);
}

#[test]
fn wake_one_keeps_remaining_matching_waiters_registered() {
let ps = PollSet::new();
let first_counter = Counter::new();
let second_counter = Counter::new();
let first_waker = Waker::from(first_counter.clone());
let second_waker = Waker::from(second_counter.clone());

unsafe {
ps.register(&first_waker, IoEvents::IN);
ps.register(&second_waker, IoEvents::IN);
}

assert_eq!(unsafe { ps.wake_one(IoEvents::IN) }, 1);
assert_eq!(first_counter.count() + second_counter.count(), 1);
assert_eq!(unsafe { ps.wake_one(IoEvents::IN) }, 1);
assert_eq!(first_counter.count(), 1);
assert_eq!(second_counter.count(), 1);
assert_eq!(unsafe { ps.wake_one(IoEvents::IN) }, 0);
}

#[test]
fn concurrent_registers_preserve_interests() {
const NUM_WAITERS: usize = 64;
Expand Down
24 changes: 24 additions & 0 deletions os/StarryOS/kernel/src/axtest_exports.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,14 @@ pub fn pipe_resize_rejects_oversized_pipe() -> bool {
super::file::resize_rejects_oversized_pipe_for_test()
}

pub fn pipe_linux_io_semantics_hold() -> bool {
super::file::pipe_linux_io_semantics_hold_for_test()
}

pub fn interrupted_pipe_write_preserves_partial_progress() -> bool {
super::file::interrupted_pipe_write_preserves_partial_progress_for_test()
}

pub fn fcntl_setpipe_size_returns_capacity() -> bool {
super::syscall::fcntl_setpipe_size_returns_capacity_for_test()
}
Expand All @@ -58,6 +66,22 @@ pub fn concurrent_epoll_reverse_add_is_serialized() -> bool {
super::file::concurrent_reverse_add_is_serialized_for_test()
}

pub fn epoll_level_aliases_rotate_in_linux_callback_order() -> bool {
super::file::level_aliases_rotate_in_linux_callback_order_for_test()
}

pub fn epoll_edge_readiness_requires_a_new_notification() -> bool {
super::file::edge_readiness_requires_a_new_notification_for_test()
}

pub fn epoll_edge_callback_does_not_reenter_target() -> bool {
super::file::edge_callback_does_not_reenter_target_for_test()
}

pub fn epoll_hup_does_not_synthesize_readable() -> bool {
super::file::epoll_hup_does_not_synthesize_readable_for_test()
}

pub fn process_mem_stats_formats_linux_fields() -> bool {
use super::mm::ProcessMemStats;

Expand Down
Loading