Skip to main content

starnix_core/mm/
futex_table.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::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
19/// A table of futexes.
20///
21/// Each 32-bit aligned address in an address space can potentially have an associated futex that
22/// userspace can wait upon. This table is a sparse representation that has an actual WaitQueue
23/// only for those addresses that have ever actually had a futex operation performed on them.
24pub struct FutexTable<Key: FutexKey> {
25    /// The futexes associated with each address in each VMO.
26    ///
27    /// This HashMap is populated on-demand when futexes are used.
28    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    /// Wait on the futex at the given address given a boot deadline.
39    ///
40    /// See FUTEX_WAIT when passed a deadline in CLOCK_REALTIME.
41    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        // As the state is locked, no wake can happen before the waiter is registered.
53        // If the addr is remapped, we will read stale data, but we will not miss a futex wake.
54        // Acquire ordering to synchronize with userspace modifications to the value on other
55        // threads.
56        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            // If wait returned an error (e.g., ETIMEDOUT, EINTR), we must explicitly
82            // remove our waiter from the queue to prevent a memory leak.
83            // If it succeeded, the waker has already removed us from the queue.
84            self.state.lock().remove_boot_waiter_from_queue(key, &waiter);
85        })
86    }
87
88    /// Wait on the futex at the given address.
89    ///
90    /// See FUTEX_WAIT.
91    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        // As the state is locked, no wake can happen before the waiter is registered.
102        // If the addr is remapped, we will read stale data, but we will not miss a futex wake.
103        // Acquire ordering to synchronize with userspace modifications to the value on other
104        // threads.
105        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            // If block_until returned an error (e.g., ETIMEDOUT, EINTR), we must explicitly
121            // remove our waiter from the queue to prevent a memory leak.
122            // If it succeeded, the waker has already removed us from the queue.
123            self.state.lock().remove_waiter_from_queue(key, &event);
124        })
125    }
126
127    /// Wake the given number of waiters on futex at the given address. Returns the number of
128    /// waiters actually woken.
129    ///
130    /// See FUTEX_WAKE.
131    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    /// Requeue the waiters to another address.
144    ///
145    /// See FUTEX_CMP_REQUEUE
146    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            // Use acquire ordering here to synchronize with mutex impls that store w/ release
162            // ordering.
163            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    /// Lock the futex at the given address.
173    ///
174    /// See FUTEX_LOCK_PI.
175    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        // As the state is locked, no unlock can happen before the waiter is registered.
184        // If the addr is remapped, we will read stale data, but we will not miss a futex unlock.
185        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        // Use a relaxed ordering because the compare/exchange below creates a synchronization
191        // point with userspace threads in the success case. No synchronization is required in
192        // failure cases.
193        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                // From <https://man7.org/linux/man-pages/man2/futex.2.html>:
198                //
199                //   EDEADLK
200                //          (FUTEX_LOCK_PI, FUTEX_LOCK_PI2, FUTEX_TRYLOCK_PI,
201                //          FUTEX_CMP_REQUEUE_PI) The futex word at uaddr is already
202                //          locked by the caller.
203                return error!(EDEADLOCK);
204            }
205
206            if current_value == 0 {
207                // Use acq/rel ordering to synchronize with acquire ordering on userspace lock ops
208                // and with the release ordering on userspace unlock ops.
209                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            // Use acq/rel ordering to synchronize with acquire ordering on userspace lock ops and
220            // with the release ordering on userspace unlock ops.
221            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        // ESRCH  (FUTEX_LOCK_PI, FUTEX_LOCK_PI2, FUTEX_TRYLOCK_PI,
242        //        FUTEX_CMP_REQUEUE_PI) The thread ID in the futex word at
243        //        uaddr does not exist.
244        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                // If block_with_owner_until returned an error (e.g., ETIMEDOUT), or if we
254                // failed to find the new owner (ESRCH), we must explicitly remove our waiter
255                // from the PI-mutex queue to prevent a memory leak.
256                self.state.lock().remove_rt_mutex_waiter_from_queue(key, &event);
257            })
258    }
259
260    /// Unlock the futex at the given address.
261    ///
262    /// See FUTEX_UNLOCK_PI.
263    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        // Use a relaxed ordering because the compare/exchange below creates a synchronization
272        // point with userspace threads in the success case. No synchronization is required in
273        // failure cases.
274        let current_value = mm.atomic_load_u32_relaxed(addr)?;
275        if current_value & FUTEX_TID_MASK != tid {
276            // From <https://man7.org/linux/man-pages/man2/futex.2.html>:
277            //
278            //   EPERM  (FUTEX_UNLOCK_PI) The caller does not own the lock
279            //          represented by the futex word.
280            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            // Use acq/rel ordering to synchronize with acquire ordering on userspace lock ops and
288            // with the release ordering on userspace unlock ops.
289            match mm.atomic_compare_exchange_u32_acq_rel(addr, current_value, target_value) {
290                CompareExchangeResult::Success => (),
291                // From <https://man7.org/linux/man-pages/man2/futex.2.html>:
292                //
293                //   EINVAL (FUTEX_LOCK_PI, FUTEX_LOCK_PI2, FUTEX_TRYLOCK_PI,
294                //       FUTEX_UNLOCK_PI) The kernel detected an inconsistency
295                //       between the user-space state at uaddr and the kernel
296                //       state.  This indicates either state corruption or that the
297                //       kernel found a waiter on uaddr which is waiting via
298                //       FUTEX_WAIT or FUTEX_WAIT_BITSET.
299                CompareExchangeResult::Stale { .. } => return error!(EINVAL),
300                // From <https://man7.org/linux/man-pages/man2/futex.2.html>:
301                //
302                //   EACCES No read access to the memory of a futex word.
303                CompareExchangeResult::Error(_) => return error!(EACCES),
304            }
305
306            let Some(mut waiter) = maybe_waiter else {
307                // We can stop trying to notify a thread if there are no more waiters.
308                break;
309            };
310
311            if waiter.notifiable.notify() {
312                break;
313            }
314
315            // If we couldn't notify the waiter, then we need to pull the next thread off the
316            // waiter list.
317        }
318
319        Ok(())
320    }
321}
322
323impl FutexTable<SharedFutexKey> {
324    /// Wait on the futex at the given offset in the memory.
325    ///
326    /// Returns a receiver that will be signaled when the futex is woken, and an
327    /// `Arc<()>` token that must be kept alive by the caller for the duration of the
328    /// wait. If the caller drops the token (e.g., if the external client
329    /// disconnects), the waiter is marked as stale and will be garbage-collected by the
330    /// next futex operation on this table.
331    ///
332    /// See FUTEX_WAIT.
333    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        // As the state is locked, no wake can happen before the waiter is registered.
343        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    /// Wake the given number of waiters on futex at the given offset in the memory. Returns the
355    /// number of waiters actually woken.
356    ///
357    /// See FUTEX_WAKE.
358    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        // If/when we move from a single table mutex to a mutex per futex, we'll likely want to
384        // define a consistent SharedFutexKey sort order independent of which is "first" and which
385        // is "second" in this call. Then we can acquire each of the two mutexes corresponding to
386        // each of the two futexes per that sort order. This way, we can be holding both mutexes to
387        // make the requeue atomic despite each futex having its own mutex, while avoiding
388        // deadlocks. But for now we lock the whole FutexTable.
389        let mut state = self.state.lock();
390        if let Some(expected) = expected_value {
391            // The state being locked is how this is included in the set of atomic changes.
392            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            // TODO: This read should be atomic.
404            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    // No chance of collisions since koids are never reused:
438    // https://fuchsia.dev/fuchsia-src/concepts/kernel/concepts#kernel_object_ids
439    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    /// Returns the FutexWaiters for a given address, creating an empty one if none is registered.
473    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                // Wake up at most `wake_count` waiters.
504                woken = entry.get_mut().notify(FUTEX_BITSET_MATCH_ANY, wake_count);
505
506                // Dequeue up to `requeue_count` waiters to requeue below.
507                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    /// Returns the RT-Mutex waiters queue for a given address, creating an empty queue if none is
524    /// registered.
525    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    /// Pop the next RT-Mutex for the given address.
530    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                // Clean up the hash map entry if the queue is empty. We do this
537                // regardless of whether `pop_front` returned a waiter or `None`,
538                // effectively garbage collecting erroneously empty map entries.
539                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    /// Removes a standard `FUTEX_WAIT` waiter from the queue.
550    ///
551    /// This uses a two-step approach:
552    /// 1. O(1) Fast path: Check the `key` where the waiter originally went to sleep.
553    /// 2. O(N) Fallback: If not found (e.g. moved via `FUTEX_REQUEUE`), scan all futexes.
554    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    /// Removes a `FUTEX_WAIT_BITSET` waiter (with `FUTEX_CLOCK_REALTIME`).
579    ///
580    /// Like `remove_waiter_from_queue`, it tries the fast O(1) lookup on the original `key` first,
581    /// and falls back to an O(N) scan across all queues in case of a requeue.
582    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    /// Removes a PI-mutex (`FUTEX_LOCK_PI`) waiter.
607    ///
608    /// Operates on the separate `rt_mutex_waiters` map using the same two-step
609    /// O(1)/O(N) algorithm as the other removal methods to handle edge cases where
610    /// PI-mutexes might be requeued (e.g. if `FUTEX_CMP_REQUEUE_PI` is used).
611    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
643/// Abstraction over a process waiting on a Futex that can be notified.
644enum FutexNotifiable {
645    /// An internal process waiting on a Futex.
646    Internal(Weak<InterruptibleEvent>),
647    // An internal process waiting on a Futex with a boot deadline.
648    InternalBoot(Weak<Waiter>),
649    /// An external process waiting on a Futex.
650    // The sender needs to be an option so that one can send the notification while only holding a
651    // mut reference on the ExternalWaiter.
652    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    /// Tries to notify the process. Returns `true` is the process have been notified. Returns
669    /// `false` otherwise. This means the process is stale and will never be available again.
670    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            // The send will fail if the receiver is gone, which means nothing was actualling
753            // waiting on the futex.
754            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    /// The tid, possibly with the FUTEX_WAITERS bit set.
792    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        } // token is dropped here, so it becomes stale
921
922        assert_eq!(state.waiters.len(), 1);
923
924        // Trigger a cleanup with a placeholder event
925        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}