1use crate::mm::memory::MemoryObject;
6use crate::mm::{CompareExchangeResult, ProtectionFlags};
7use crate::task::{CurrentTask, EventHandler, SignalHandler, SignalHandlerInner, Task, Waiter};
8use futures::channel::oneshot;
9use starnix_sync::{FutexTableStateLock, InterruptibleEvent, LockDepMutex};
10use starnix_types::futex_address::FutexAddress;
11use starnix_uapi::errors::Errno;
12use starnix_uapi::user_address::UserAddress;
13use starnix_uapi::{FUTEX_BITSET_MATCH_ANY, FUTEX_TID_MASK, FUTEX_WAITERS, errno, error};
14use std::collections::hash_map::Entry;
15use std::collections::{HashMap, VecDeque};
16use std::hash::Hash;
17use std::sync::{Arc, Weak};
18
19pub struct FutexTable<Key: FutexKey> {
25 state: LockDepMutex<FutexTableState<Key>, FutexTableStateLock>,
29}
30
31impl<Key: FutexKey> Default for FutexTable<Key> {
32 fn default() -> Self {
33 Self { state: LockDepMutex::new(FutexTableState::default()) }
34 }
35}
36
37impl<Key: FutexKey> FutexTable<Key> {
38 pub fn wait_boot(
42 &self,
43 current_task: &CurrentTask,
44 addr: UserAddress,
45 value: u32,
46 mask: u32,
47 deadline: zx::BootInstant,
48 timer_slack: zx::BootDuration,
49 ) -> Result<(), Errno> {
50 let addr = FutexAddress::try_from(addr)?;
51 let mut state = self.state.lock();
52 let loaded_value = current_task.mm()?.atomic_load_u32_acquire(addr)?;
57 if value != loaded_value {
58 return error!(EAGAIN);
59 }
60
61 let key = Key::get(current_task, addr)?;
62 let waiter = Arc::new(Waiter::new());
63 let timer = zx::BootTimer::create();
64 let signal_handler = SignalHandler {
65 inner: SignalHandlerInner::None,
66 event_handler: EventHandler::None,
67 err_code: Some(errno!(ETIMEDOUT)),
68 };
69 waiter
70 .wake_on_zircon_signals(&timer, zx::Signals::TIMER_SIGNALED, signal_handler)
71 .expect("wait can only fail in OOM conditions");
72 timer
73 .set(deadline, timer_slack)
74 .expect("timer set cannot fail with valid handles and slack");
75 state.get_waiters_or_default(key.clone()).add(FutexWaiter {
76 mask,
77 notifiable: FutexNotifiable::new_internal_boot(Arc::downgrade(&waiter)),
78 });
79 std::mem::drop(state);
80 waiter.wait(current_task).inspect_err(|_| {
81 self.state.lock().remove_boot_waiter_from_queue(key, &waiter);
85 })
86 }
87
88 pub fn wait(
92 &self,
93 current_task: &CurrentTask,
94 addr: UserAddress,
95 value: u32,
96 mask: u32,
97 deadline: zx::MonotonicInstant,
98 ) -> Result<(), Errno> {
99 let addr = FutexAddress::try_from(addr)?;
100 let mut state = self.state.lock();
101 let loaded_value = current_task.mm()?.atomic_load_u32_acquire(addr)?;
106 if value != loaded_value {
107 return error!(EAGAIN);
108 }
109
110 let key = Key::get(current_task, addr)?;
111 let event = InterruptibleEvent::new();
112 let guard = event.begin_wait();
113 state.get_waiters_or_default(key.clone()).add(FutexWaiter {
114 mask,
115 notifiable: FutexNotifiable::new_internal(Arc::downgrade(&event)),
116 });
117 std::mem::drop(state);
118
119 current_task.block_until(guard, deadline).inspect_err(|_| {
120 self.state.lock().remove_waiter_from_queue(key, &event);
124 })
125 }
126
127 pub fn wake(
132 &self,
133 task: &Task,
134 addr: UserAddress,
135 count: usize,
136 mask: u32,
137 ) -> Result<usize, Errno> {
138 let addr = FutexAddress::try_from(addr)?;
139 let key = Key::get(task, addr)?;
140 Ok(self.state.lock().wake(key, count, mask))
141 }
142
143 pub fn requeue(
147 &self,
148 current_task: &CurrentTask,
149 addr: UserAddress,
150 wake_count: usize,
151 requeue_count: usize,
152 new_addr: UserAddress,
153 expected_value: Option<u32>,
154 ) -> Result<usize, Errno> {
155 let addr = FutexAddress::try_from(addr)?;
156 let new_addr = FutexAddress::try_from(new_addr)?;
157 let key = Key::get(current_task, addr)?;
158 let new_key = Key::get(current_task, new_addr)?;
159 let mut state = self.state.lock();
160 if let Some(expected) = expected_value {
161 let value = current_task.mm()?.atomic_load_u32_acquire(addr)?;
164 if value != expected {
165 return error!(EAGAIN);
166 }
167 }
168
169 Ok(state.requeue(key, new_key, wake_count, requeue_count))
170 }
171
172 pub fn lock_pi(
176 &self,
177 current_task: &CurrentTask,
178 addr: UserAddress,
179 deadline: zx::MonotonicInstant,
180 ) -> Result<(), Errno> {
181 let addr = FutexAddress::try_from(addr)?;
182 let mut state = self.state.lock();
183 let key = Key::get(current_task, addr)?;
186
187 let tid = current_task.get_tid() as u32;
188 let mm = current_task.mm()?;
189
190 let mut current_value = mm.atomic_load_u32_relaxed(addr)?;
194 let new_owner_tid = loop {
195 let new_owner_tid = current_value & FUTEX_TID_MASK;
196 if new_owner_tid == tid {
197 return error!(EDEADLOCK);
204 }
205
206 if current_value == 0 {
207 match mm.atomic_compare_exchange_weak_u32_acq_rel(addr, current_value, tid) {
210 CompareExchangeResult::Success => return Ok(()),
211 CompareExchangeResult::Stale { observed } => {
212 current_value = observed;
213 continue;
214 }
215 CompareExchangeResult::Error(e) => return Err(e),
216 }
217 }
218
219 let target_value = current_value | FUTEX_WAITERS;
222 match mm.atomic_compare_exchange_u32_acq_rel(addr, current_value, target_value) {
223 CompareExchangeResult::Success => (),
224 CompareExchangeResult::Stale { observed } => {
225 current_value = observed;
226 continue;
227 }
228 CompareExchangeResult::Error(e) => return Err(e),
229 }
230 break new_owner_tid;
231 };
232
233 let event = InterruptibleEvent::new();
234 let guard = event.begin_wait();
235 let notifiable = FutexNotifiable::new_internal(Arc::downgrade(&event));
236 state
237 .get_rt_mutex_waiters_or_default(key.clone())
238 .push_back(RtMutexWaiter { tid, notifiable });
239 std::mem::drop(state);
240
241 current_task
245 .get_task(new_owner_tid as i32)
246 .ok()
247 .and_then(|o| o.running_state().unwrap().thread.get().map(|t| Arc::clone(&t.thread)))
248 .map_or_else(
249 || error!(ESRCH),
250 |owner| current_task.block_with_owner_until(guard, &owner, deadline),
251 )
252 .inspect_err(|_| {
253 self.state.lock().remove_rt_mutex_waiter_from_queue(key, &event);
257 })
258 }
259
260 pub fn unlock_pi(&self, current_task: &CurrentTask, addr: UserAddress) -> Result<(), Errno> {
264 let addr = FutexAddress::try_from(addr)?;
265 let mut state = self.state.lock();
266 let tid = current_task.get_tid() as u32;
267 let mm = current_task.mm()?;
268
269 let key = Key::get(current_task, addr)?;
270
271 let current_value = mm.atomic_load_u32_relaxed(addr)?;
275 if current_value & FUTEX_TID_MASK != tid {
276 return error!(EPERM);
281 }
282
283 loop {
284 let maybe_waiter = state.pop_rt_mutex_waiter(key.clone());
285 let target_value = if let Some(waiter) = &maybe_waiter { waiter.tid } else { 0 };
286
287 match mm.atomic_compare_exchange_u32_acq_rel(addr, current_value, target_value) {
290 CompareExchangeResult::Success => (),
291 CompareExchangeResult::Stale { .. } => return error!(EINVAL),
300 CompareExchangeResult::Error(_) => return error!(EACCES),
304 }
305
306 let Some(mut waiter) = maybe_waiter else {
307 break;
309 };
310
311 if waiter.notifiable.notify() {
312 break;
313 }
314
315 }
318
319 Ok(())
320 }
321}
322
323impl FutexTable<SharedFutexKey> {
324 pub fn external_wait(
334 &self,
335 memory: MemoryObject,
336 offset: u64,
337 value: u32,
338 mask: u32,
339 ) -> Result<(Arc<()>, oneshot::Receiver<()>), Errno> {
340 let key = SharedFutexKey::new(&memory, offset);
341 let mut state = self.state.lock();
342 Self::external_check_futex_value(&memory, offset, value)?;
344
345 let token = Arc::new(());
346 let (sender, receiver) = oneshot::channel::<()>();
347 state.get_waiters_or_default(key).add(FutexWaiter {
348 mask,
349 notifiable: FutexNotifiable::new_external(Arc::downgrade(&token), sender),
350 });
351 Ok((token, receiver))
352 }
353
354 pub fn external_wake(
359 &self,
360 memory: MemoryObject,
361 offset: u64,
362 count: usize,
363 mask: u32,
364 ) -> Result<usize, Errno> {
365 Ok(self.state.lock().wake(SharedFutexKey::new(&memory, offset), count, mask))
366 }
367
368 pub fn external_requeue(
369 &self,
370 first_memory: MemoryObject,
371 first_offset: u64,
372 second_memory: Option<MemoryObject>,
373 second_offset: u64,
374 wake_count: usize,
375 requeue_count: usize,
376 expected_value: Option<u32>,
377 ) -> Result<usize, Errno> {
378 let first_key = SharedFutexKey::new(&first_memory, first_offset);
379 let second_key = match second_memory.as_ref() {
380 Some(second_memory) => SharedFutexKey::new(second_memory, second_offset),
381 None => SharedFutexKey::new(&first_memory, second_offset),
382 };
383 let mut state = self.state.lock();
390 if let Some(expected) = expected_value {
391 Self::external_check_futex_value(&first_memory, first_offset, expected)?;
393 }
394 Ok(state.requeue(first_key, second_key, wake_count, requeue_count))
395 }
396
397 fn external_check_futex_value(
398 memory: &MemoryObject,
399 offset: u64,
400 value: u32,
401 ) -> Result<(), Errno> {
402 let loaded_value = {
403 let mut buf = [0u8; 4];
405 memory.read(&mut buf, offset).map_err(|_| errno!(EINVAL))?;
406 u32::from_ne_bytes(buf)
407 };
408 if loaded_value != value {
409 return error!(EAGAIN);
410 }
411 Ok(())
412 }
413}
414
415pub trait FutexKey: Sized + Ord + Hash + Clone {
416 fn get(task: &Task, addr: FutexAddress) -> Result<Self, Errno>;
417 fn get_table_from_task(task: &Task) -> Result<Arc<FutexTable<Self>>, Errno>;
418}
419
420#[derive(Debug, Clone, Eq, Hash, PartialEq, Ord, PartialOrd)]
421pub struct PrivateFutexKey {
422 addr: FutexAddress,
423}
424
425impl FutexKey for PrivateFutexKey {
426 fn get(_task: &Task, addr: FutexAddress) -> Result<Self, Errno> {
427 Ok(PrivateFutexKey { addr })
428 }
429
430 fn get_table_from_task(task: &Task) -> Result<Arc<FutexTable<Self>>, Errno> {
431 Ok(task.mm()?.futex.clone())
432 }
433}
434
435#[derive(Debug, Clone, Eq, Hash, PartialEq, Ord, PartialOrd)]
436pub struct SharedFutexKey {
437 koid: zx::Koid,
440 offset: u64,
441}
442
443impl FutexKey for SharedFutexKey {
444 fn get(task: &Task, addr: FutexAddress) -> Result<Self, Errno> {
445 let (memory, offset) = task.mm()?.get_mapping_memory(addr.into(), ProtectionFlags::READ)?;
446 Ok(SharedFutexKey::new(&memory, offset))
447 }
448
449 fn get_table_from_task(task: &Task) -> Result<Arc<FutexTable<Self>>, Errno> {
450 Ok(task.kernel().shared_futexes.clone())
451 }
452}
453
454impl SharedFutexKey {
455 fn new(memory: &MemoryObject, offset: u64) -> Self {
456 Self { koid: memory.get_koid(), offset }
457 }
458}
459
460struct FutexTableState<Key: FutexKey> {
461 waiters: HashMap<Key, FutexWaiters>,
462 rt_mutex_waiters: HashMap<Key, VecDeque<RtMutexWaiter>>,
463}
464
465impl<Key: FutexKey> Default for FutexTableState<Key> {
466 fn default() -> Self {
467 Self { waiters: Default::default(), rt_mutex_waiters: Default::default() }
468 }
469}
470
471impl<Key: FutexKey> FutexTableState<Key> {
472 fn get_waiters_or_default(&mut self, key: Key) -> &mut FutexWaiters {
474 self.waiters.entry(key).or_default()
475 }
476
477 fn wake(&mut self, key: Key, count: usize, mask: u32) -> usize {
478 let entry = self.waiters.entry(key);
479 match entry {
480 Entry::Vacant(_) => 0,
481 Entry::Occupied(mut entry) => {
482 let count = entry.get_mut().notify(mask, count);
483 if entry.get().is_empty() {
484 entry.remove();
485 }
486 count
487 }
488 }
489 }
490
491 fn requeue(
492 &mut self,
493 key: Key,
494 new_key: Key,
495 wake_count: usize,
496 requeue_count: usize,
497 ) -> usize {
498 let woken;
499 let to_requeue;
500 match self.waiters.entry(key) {
501 Entry::Vacant(_) => return 0,
502 Entry::Occupied(mut entry) => {
503 woken = entry.get_mut().notify(FUTEX_BITSET_MATCH_ANY, wake_count);
505
506 to_requeue = entry.get_mut().split_for_requeue(requeue_count);
508
509 if entry.get().is_empty() {
510 entry.remove();
511 }
512 }
513 }
514
515 let requeued = to_requeue.0.len();
516 if !to_requeue.is_empty() {
517 self.get_waiters_or_default(new_key).transfer(to_requeue);
518 }
519
520 woken + requeued
521 }
522
523 fn get_rt_mutex_waiters_or_default(&mut self, key: Key) -> &mut VecDeque<RtMutexWaiter> {
526 self.rt_mutex_waiters.entry(key).or_default()
527 }
528
529 fn pop_rt_mutex_waiter(&mut self, key: Key) -> Option<RtMutexWaiter> {
531 let entry = self.rt_mutex_waiters.entry(key);
532 match entry {
533 Entry::Vacant(_) => None,
534 Entry::Occupied(mut entry) => {
535 let mut waiter = entry.get_mut().pop_front();
536 if entry.get().is_empty() {
540 entry.remove();
541 } else if let Some(waiter) = &mut waiter {
542 waiter.tid |= FUTEX_WAITERS;
543 }
544 waiter
545 }
546 }
547 }
548
549 fn remove_waiter_from_queue(&mut self, key: Key, event: &Arc<InterruptibleEvent>) {
555 if let Entry::Occupied(mut entry) = self.waiters.entry(key) {
556 if entry.get_mut().remove_waiter(event) {
557 if entry.get().is_empty() {
558 entry.remove();
559 }
560 return;
561 }
562 }
563
564 let mut key_to_remove = None;
565 for (key, waiters) in self.waiters.iter_mut() {
566 if waiters.remove_waiter(event) {
567 if waiters.is_empty() {
568 key_to_remove = Some(key.clone());
569 }
570 break;
571 }
572 }
573 if let Some(key) = key_to_remove {
574 self.waiters.remove(&key);
575 }
576 }
577
578 fn remove_boot_waiter_from_queue(&mut self, key: Key, waiter: &Arc<Waiter>) {
583 if let Entry::Occupied(mut entry) = self.waiters.entry(key) {
584 if entry.get_mut().remove_boot_waiter(waiter) {
585 if entry.get().is_empty() {
586 entry.remove();
587 }
588 return;
589 }
590 }
591
592 let mut key_to_remove = None;
593 for (key, waiters) in self.waiters.iter_mut() {
594 if waiters.remove_boot_waiter(waiter) {
595 if waiters.is_empty() {
596 key_to_remove = Some(key.clone());
597 }
598 break;
599 }
600 }
601 if let Some(key) = key_to_remove {
602 self.waiters.remove(&key);
603 }
604 }
605
606 fn remove_rt_mutex_waiter_from_queue(&mut self, key: Key, event: &Arc<InterruptibleEvent>) {
612 let predicate =
613 |w: &RtMutexWaiter| !w.notifiable.matches_event(event) && !w.notifiable.is_stale();
614
615 if let Entry::Occupied(mut entry) = self.rt_mutex_waiters.entry(key) {
616 let len_before = entry.get().len();
617 entry.get_mut().retain(predicate);
618 if entry.get().len() < len_before {
619 if entry.get().is_empty() {
620 entry.remove();
621 }
622 return;
623 }
624 }
625
626 let mut key_to_remove = None;
627 for (key, waiters) in self.rt_mutex_waiters.iter_mut() {
628 let len_before = waiters.len();
629 waiters.retain(predicate);
630 if waiters.len() < len_before {
631 if waiters.is_empty() {
632 key_to_remove = Some(key.clone());
633 }
634 break;
635 }
636 }
637 if let Some(key) = key_to_remove {
638 self.rt_mutex_waiters.remove(&key);
639 }
640 }
641}
642
643enum FutexNotifiable {
645 Internal(Weak<InterruptibleEvent>),
647 InternalBoot(Weak<Waiter>),
649 External(Weak<()>, Option<oneshot::Sender<()>>),
653}
654
655impl FutexNotifiable {
656 fn new_internal(event: Weak<InterruptibleEvent>) -> Self {
657 Self::Internal(event)
658 }
659
660 fn new_internal_boot(waiter: Weak<Waiter>) -> Self {
661 Self::InternalBoot(waiter)
662 }
663
664 fn new_external(token: Weak<()>, sender: oneshot::Sender<()>) -> Self {
665 Self::External(token, Some(sender))
666 }
667
668 fn notify(&mut self) -> bool {
671 match self {
672 Self::Internal(event) => {
673 if let Some(event) = event.upgrade() {
674 event.notify();
675 true
676 } else {
677 false
678 }
679 }
680 Self::InternalBoot(waiter) => {
681 if let Some(waiter) = waiter.upgrade() {
682 waiter.notify();
683 true
684 } else {
685 false
686 }
687 }
688 Self::External(_, sender) => {
689 if let Some(sender) = sender.take() {
690 sender.send(()).is_ok()
691 } else {
692 false
693 }
694 }
695 }
696 }
697
698 fn matches_event(&self, event: &Arc<InterruptibleEvent>) -> bool {
699 match self {
700 Self::Internal(weak) => {
701 if let Some(strong) = weak.upgrade() {
702 Arc::ptr_eq(&strong, event)
703 } else {
704 false
705 }
706 }
707 _ => false,
708 }
709 }
710
711 fn matches_waiter(&self, waiter: &Arc<Waiter>) -> bool {
712 match self {
713 Self::InternalBoot(weak) => {
714 if let Some(strong) = weak.upgrade() {
715 Arc::ptr_eq(&strong, waiter)
716 } else {
717 false
718 }
719 }
720 _ => false,
721 }
722 }
723
724 fn is_stale(&self) -> bool {
725 match self {
726 Self::Internal(weak) => weak.strong_count() == 0,
727 Self::External(weak, _) => weak.strong_count() == 0,
728 Self::InternalBoot(weak) => weak.strong_count() == 0,
729 }
730 }
731}
732
733struct FutexWaiter {
734 mask: u32,
735 notifiable: FutexNotifiable,
736}
737
738#[derive(Default)]
739struct FutexWaiters(VecDeque<FutexWaiter>);
740
741impl FutexWaiters {
742 fn add(&mut self, waiter: FutexWaiter) {
743 self.0.push_back(waiter);
744 }
745
746 fn notify(&mut self, mask: u32, count: usize) -> usize {
747 let mut woken = 0;
748 self.0.retain_mut(|waiter| {
749 if woken == count || waiter.mask & mask == 0 {
750 return true;
751 }
752 if waiter.notifiable.notify() {
755 woken += 1;
756 }
757 false
758 });
759 woken
760 }
761
762 fn transfer(&mut self, mut other: Self) {
763 self.0.append(&mut other.0);
764 }
765
766 fn is_empty(&self) -> bool {
767 self.0.is_empty()
768 }
769
770 fn remove_waiter(&mut self, event: &Arc<InterruptibleEvent>) -> bool {
771 let initial_len = self.0.len();
772 self.0.retain(|w| !w.notifiable.matches_event(event) && !w.notifiable.is_stale());
773 self.0.len() < initial_len
774 }
775
776 fn remove_boot_waiter(&mut self, waiter: &Arc<Waiter>) -> bool {
777 let initial_len = self.0.len();
778 self.0.retain(|w| !w.notifiable.matches_waiter(waiter) && !w.notifiable.is_stale());
779 self.0.len() < initial_len
780 }
781
782 fn split_for_requeue(&mut self, count: usize) -> Self {
783 let count = std::cmp::min(count, self.0.len());
784 let tail = self.0.split_off(count);
785 let head = std::mem::replace(&mut self.0, tail);
786 FutexWaiters(head)
787 }
788}
789
790struct RtMutexWaiter {
791 tid: u32,
793
794 notifiable: FutexNotifiable,
795}
796
797#[cfg(test)]
798mod tests {
799 use super::*;
800 use starnix_sync::InterruptibleEvent;
801 use starnix_uapi::restricted_aspace::RESTRICTED_ASPACE_BASE;
802 use starnix_uapi::user_address::UserAddress;
803
804 #[fuchsia::test]
805 fn test_remove_waiter_simple() {
806 let mut state = FutexTableState::<PrivateFutexKey>::default();
807 let key = PrivateFutexKey {
808 addr: FutexAddress::try_from(UserAddress::from(
809 (RESTRICTED_ASPACE_BASE + 0x1000) as u64,
810 ))
811 .unwrap(),
812 };
813 let event = Arc::new(InterruptibleEvent::new());
814
815 state.get_waiters_or_default(key.clone()).add(FutexWaiter {
816 mask: u32::MAX,
817 notifiable: FutexNotifiable::new_internal(Arc::downgrade(&event)),
818 });
819
820 assert_eq!(state.waiters.len(), 1);
821 state.remove_waiter_from_queue(key, &event);
822 assert_eq!(state.waiters.len(), 0);
823 }
824
825 #[fuchsia::test]
826 fn test_remove_waiter_requeued() {
827 let mut state = FutexTableState::<PrivateFutexKey>::default();
828 let key1 = PrivateFutexKey {
829 addr: FutexAddress::try_from(UserAddress::from(
830 (RESTRICTED_ASPACE_BASE + 0x1000) as u64,
831 ))
832 .unwrap(),
833 };
834 let key2 = PrivateFutexKey {
835 addr: FutexAddress::try_from(UserAddress::from(
836 (RESTRICTED_ASPACE_BASE + 0x2000) as u64,
837 ))
838 .unwrap(),
839 };
840 let event = Arc::new(InterruptibleEvent::new());
841
842 state.get_waiters_or_default(key2.clone()).add(FutexWaiter {
843 mask: u32::MAX,
844 notifiable: FutexNotifiable::new_internal(Arc::downgrade(&event)),
845 });
846
847 assert_eq!(state.waiters.len(), 1);
848 state.remove_waiter_from_queue(key1, &event);
849 assert_eq!(state.waiters.len(), 0);
850 }
851
852 #[fuchsia::test]
853 fn test_remove_rt_mutex_waiter() {
854 let mut state = FutexTableState::<PrivateFutexKey>::default();
855 let key = PrivateFutexKey {
856 addr: FutexAddress::try_from(UserAddress::from(
857 (RESTRICTED_ASPACE_BASE + 0x1000) as u64,
858 ))
859 .unwrap(),
860 };
861 let event = Arc::new(InterruptibleEvent::new());
862
863 state.get_rt_mutex_waiters_or_default(key.clone()).push_back(RtMutexWaiter {
864 tid: 1,
865 notifiable: FutexNotifiable::new_internal(Arc::downgrade(&event)),
866 });
867
868 assert_eq!(state.rt_mutex_waiters.len(), 1);
869 state.remove_rt_mutex_waiter_from_queue(key, &event);
870 assert_eq!(state.rt_mutex_waiters.len(), 0);
871 }
872
873 #[fuchsia::test]
874 fn test_split_for_requeue_fairness() {
875 let mut waiters = FutexWaiters::default();
876 let e1 = Arc::new(InterruptibleEvent::new());
877 let e2 = Arc::new(InterruptibleEvent::new());
878 let e3 = Arc::new(InterruptibleEvent::new());
879
880 waiters.add(FutexWaiter {
881 mask: 1,
882 notifiable: FutexNotifiable::new_internal(Arc::downgrade(&e1)),
883 });
884 waiters.add(FutexWaiter {
885 mask: 2,
886 notifiable: FutexNotifiable::new_internal(Arc::downgrade(&e2)),
887 });
888 waiters.add(FutexWaiter {
889 mask: 3,
890 notifiable: FutexNotifiable::new_internal(Arc::downgrade(&e3)),
891 });
892
893 let split = waiters.split_for_requeue(2);
894
895 assert_eq!(split.0.len(), 2);
896 assert_eq!(split.0[0].mask, 1);
897 assert_eq!(split.0[1].mask, 2);
898
899 assert_eq!(waiters.0.len(), 1);
900 assert_eq!(waiters.0[0].mask, 3);
901 }
902
903 #[fuchsia::test]
904 fn test_stale_external_waiter_cleanup() {
905 let mut state = FutexTableState::<PrivateFutexKey>::default();
906 let key = PrivateFutexKey {
907 addr: FutexAddress::try_from(UserAddress::from(
908 (RESTRICTED_ASPACE_BASE + 0x1000) as u64,
909 ))
910 .unwrap(),
911 };
912
913 {
914 let token = Arc::new(());
915 let (sender, _receiver) = oneshot::channel::<()>();
916 state.get_waiters_or_default(key.clone()).add(FutexWaiter {
917 mask: u32::MAX,
918 notifiable: FutexNotifiable::new_external(Arc::downgrade(&token), sender),
919 });
920 } assert_eq!(state.waiters.len(), 1);
923
924 let dummy_event = Arc::new(InterruptibleEvent::new());
926 state.remove_waiter_from_queue(key, &dummy_event);
927
928 assert_eq!(state.waiters.len(), 0, "Stale external waiter should be removed");
929 }
930}