Skip to main content

ax_task/sched/system/task_system/pi/
operations.rs

1//! Public PI mutex registration, cancellation, release, and claim transactions.
2
3use super::*;
4
5enum PiWaiterRemoval {
6    Removed(Option<ThreadId>),
7    HandoffPending,
8}
9
10impl TaskSystem {
11    fn remove_registered_waiter(
12        &self,
13        waiter_core: &Arc<ThreadCore>,
14        lock: PiMutexRaw,
15        generation: u64,
16        reject_handoff_top: bool,
17    ) -> Result<PiWaiterRemoval, TaskError> {
18        let mut lock_state = unsafe {
19            // SAFETY: the token or rollback caller retains the mutex identity.
20            lock_raw_pi_mutex_waiters(lock)
21        };
22        let core = unsafe {
23            // SAFETY: identical lifetime contract to the waiter-tree guard.
24            lock.core()
25        };
26        let snapshot = core.owner_snapshot();
27        if !snapshot.has_waiters() {
28            return Err(TaskError::InvalidPiState);
29        }
30        let registration = waiter_core
31            .sched()
32            .lock()
33            .pi
34            .blocked_on
35            .filter(|registration| {
36                registration.lock == lock && registration.generation == generation
37            })
38            .ok_or(TaskError::InvalidPiState)?;
39        if reject_handoff_top
40            && snapshot.is_ownerless()
41            && lock_state.waiters.first() == Some(registration.key)
42        {
43            return Ok(PiWaiterRemoval::HandoffPending);
44        }
45        if !lock_state.waiters.contains(registration.key) {
46            return Err(TaskError::InvalidPiState);
47        }
48        let owner = snapshot.owner().map(ThreadId::from);
49        let owner_core = owner.map(|owner| self.pi_thread_core(owner)).transpose()?;
50        self.remove_lock_waiter(
51            &mut lock_state,
52            owner_core.as_ref(),
53            waiter_core,
54            generation,
55        )?;
56        if lock_state.waiters.is_empty() {
57            if let Some(owner) = owner {
58                core.clear_waiters_bit(owner.into());
59            } else {
60                core.publish_unlocked();
61            }
62        }
63        Ok(PiWaiterRemoval::Removed(owner))
64    }
65
66    /// Registers one contender in the mutex-owned PI waiter tree.
67    pub fn pi_mutex_lock_slow(
68        &self,
69        lock: PiMutexRef<'_>,
70        waiter: ThreadId,
71        sequence: u64,
72    ) -> Result<PiMutexLockResult, TaskError> {
73        let _preempt = PreemptScope::enter();
74        let waiter_core = self.pi_thread_core(waiter)?;
75        let Some(_waiter_activity) = waiter_core.try_scheduler_activity() else {
76            return Err(TaskError::InvalidPiWaitState(
77                PiWaitStateError::ExitedParticipant,
78            ));
79        };
80        let (donation_policy, donation_root) = {
81            let sched = waiter_core.sched().lock();
82            if sched.lifecycle.state() == ThreadState::Exited {
83                return Err(TaskError::InvalidPiWaitState(
84                    PiWaitStateError::ExitedParticipant,
85                ));
86            }
87            if sched.pi.blocked_on.is_some() {
88                return Err(TaskError::InvalidPiWaitState(
89                    PiWaitStateError::WaiterAlreadyBlocked,
90                ));
91            }
92            (
93                waiter_core.effective_policy_snapshot(),
94                sched.pi.donor.unwrap_or(waiter_core.id()),
95            )
96        };
97        let lock_raw = lock.raw();
98        let mutex_core = lock.core();
99        let urgency = waiter_core.effective_pi_wait_urgency();
100        let donation =
101            self.pi_donation_from_snapshot(&waiter_core, donation_policy, donation_root)?;
102        let key = PiWaitKey::new(urgency, sequence, waiter);
103        let mut lock_state = lock_pi_mutex_waiters(lock);
104        loop {
105            let snapshot = mutex_core.owner_snapshot();
106            if snapshot.is_unlocked() {
107                if !lock_state.waiters.is_empty() {
108                    return Err(TaskError::InvalidPiWaitState(
109                        PiWaitStateError::StaleSchedulerOwnership,
110                    ));
111                }
112                if mutex_core.try_acquire_snapshot(snapshot, waiter.into()) {
113                    return Ok(PiMutexLockResult::Acquired);
114                }
115                continue;
116            }
117            let owner = snapshot.owner().map(ThreadId::from);
118            if owner == Some(waiter) {
119                return Err(TaskError::InvalidPiWaitState(
120                    PiWaitStateError::WaiterOwnsLock,
121                ));
122            }
123            if snapshot.has_waiters() != !lock_state.waiters.is_empty()
124                || snapshot.is_ownerless() && lock_state.waiters.is_empty()
125            {
126                return Err(TaskError::InvalidPiWaitState(
127                    PiWaitStateError::StaleSchedulerOwnership,
128                ));
129            }
130            if owner.is_none() && !snapshot.is_ownerless() {
131                return Err(TaskError::InvalidPiWaitState(
132                    PiWaitStateError::OwnerlessSelectionMissing,
133                ));
134            }
135            // Linux rereads the futex owner after an exit-race lookup fails:
136            // an owner may release this mutex and close its scheduler lifetime
137            // after our physical snapshot but before we lease its PI state.
138            // Only an unchanged snapshot still describes an exited owner.
139            let initial_owner = if let Some(owner) = owner {
140                match self.pi_thread_core(owner) {
141                    Ok(owner) => Some(owner),
142                    Err(_) if mutex_core.owner_snapshot() != snapshot => continue,
143                    Err(_) => {
144                        return Err(TaskError::InvalidPiWaitState(
145                            PiWaitStateError::ExitedParticipant,
146                        ));
147                    }
148                }
149            } else {
150                None
151            };
152
153            let _owner_activity = if let Some(owner) = initial_owner.as_ref() {
154                match owner.try_scheduler_activity() {
155                    Some(activity) => Some(activity),
156                    None if mutex_core.owner_snapshot() != snapshot => continue,
157                    None => {
158                        return Err(TaskError::InvalidPiWaitState(
159                            PiWaitStateError::ExitedParticipant,
160                        ));
161                    }
162                }
163            } else {
164                None
165            };
166            if !mutex_core.try_mark_waiters(snapshot) {
167                continue;
168            }
169            let generation = match waiter_core.pi_wait_state().begin() {
170                Ok(generation) => generation,
171                Err(error) => {
172                    if !snapshot.has_waiters() {
173                        mutex_core.clear_waiters_bit(
174                            owner.expect("an owned mutex must retain its owner").into(),
175                        );
176                    }
177                    return Err(error);
178                }
179            };
180            let owner_next_lock = match self.insert_lock_waiter(
181                &mut lock_state,
182                initial_owner.as_ref(),
183                &waiter_core,
184                PiWaitRegistration {
185                    lock: lock_raw,
186                    key,
187                    generation,
188                },
189                donation.with_wait_generation(generation),
190            ) {
191                Ok(owner_next_lock) => owner_next_lock,
192                Err(error) => {
193                    if lock_state.waiters.is_empty() {
194                        mutex_core.clear_waiters_bit(
195                            owner
196                                .expect("an empty waiter tree must retain owner")
197                                .into(),
198                        );
199                    }
200                    return Err(error);
201                }
202            };
203            // Linux snapshots `owner->pi_blocked_on->lock` while both the
204            // origin wait-lock and owner pi-lock protect the newly installed
205            // edge. A later `blocked_on` value can belong to an unrelated
206            // dependency created after the owner releases this mutex.
207            drop(lock_state);
208            let chain_result =
209                if let (Some(owner), Some(owner_next_lock)) = (owner, owner_next_lock) {
210                    self.recompute_pi_chain(owner, lock_raw, owner_next_lock, waiter)
211                } else {
212                    Ok(())
213                };
214            if let Err(error) = chain_result {
215                let rollback_owner = self
216                    .remove_registered_waiter(&waiter_core, lock_raw, generation, false)
217                    .unwrap_or_else(|_| {
218                        task_runtime::fatal_invariant(
219                            0x5049_120e,
220                            waiter_core.id().as_u64() as usize,
221                        )
222                    });
223                let PiWaiterRemoval::Removed(rollback_owner) = rollback_owner else {
224                    task_runtime::fatal_invariant(0x5049_1216, waiter_core.id().as_u64() as usize)
225                };
226                if let Some(rollback_owner) = rollback_owner {
227                    self.recompute_pi_cleanup_chain(rollback_owner, waiter)
228                        .unwrap_or_else(|_| {
229                            task_runtime::fatal_invariant(
230                                0x5049_120f,
231                                waiter_core.id().as_u64() as usize,
232                            )
233                        });
234                }
235                return Err(error);
236            }
237            // Linux takes a task_struct reference while the owner remains
238            // protected by wait_lock/pi_lock. Retain the equivalent typed
239            // scheduler handle before releasing the owner activity lease, so
240            // owner observation never has to resolve a stale numeric ID.
241            let initial_owner_handle = initial_owner
242                .as_ref()
243                .map(|owner| ThreadHandle::from_core(Arc::clone(owner)));
244            #[cfg(axtest)]
245            super::axtest::record_waiter_registration(owner, waiter_core.state());
246            drop(_owner_activity);
247            drop(_waiter_activity);
248
249            return Ok(PiMutexLockResult::Waiting(unsafe {
250                // SAFETY: both waiter-tree edges are committed and retain this
251                // physical lock identity until claim or cancellation.
252                PiWaitToken::from_registration(
253                    lock_raw,
254                    waiter.into(),
255                    initial_owner_handle,
256                    generation,
257                    core::ptr::NonNull::from(waiter_core.pi_wait_state()).cast(),
258                )
259            }));
260        }
261    }
262
263    /// Cancels a committed waiter which has not been selected for claim.
264    pub fn pi_wait_cancel(&self, token: PiWaitToken) -> Result<(), TaskError> {
265        match self.pi_wait_try_cancel(&token)? {
266            PiWaitCancelOutcome::Cancelled => Ok(()),
267            PiWaitCancelOutcome::HandoffPending => Err(TaskError::InvalidPiState),
268        }
269    }
270
271    /// Tries to cancel a committed waiter without consuming a published
272    /// ownerless handoff.
273    pub fn pi_wait_try_cancel(
274        &self,
275        token: &PiWaitToken,
276    ) -> Result<PiWaitCancelOutcome, TaskError> {
277        let _preempt = PreemptScope::enter();
278        let waiter = ThreadId::from(token.thread_id());
279        let waiter_core = self.pi_thread_core(waiter)?;
280        let removal = self.remove_registered_waiter(
281            &waiter_core,
282            token.lock_raw(),
283            token.generation(),
284            true,
285        )?;
286        let PiWaiterRemoval::Removed(owner) = removal else {
287            return Ok(PiWaitCancelOutcome::HandoffPending);
288        };
289        if let Some(owner) = owner {
290            self.recompute_pi_cleanup_chain(owner, waiter)?;
291        }
292        Ok(PiWaitCancelOutcome::Cancelled)
293    }
294
295    /// Publishes an ownerless handoff and wakes the current top waiter.
296    pub fn pi_mutex_release(
297        &self,
298        lock: PiMutexRef<'_>,
299        old_owner: ThreadId,
300    ) -> Result<(), TaskError> {
301        let _preempt = PreemptScope::enter();
302        let old_owner_core = self.pi_thread_core(old_owner)?;
303        let wake = loop {
304            let mutex_core = lock.core();
305            let lock_state = lock_pi_mutex_waiters(lock);
306            let snapshot = mutex_core.owner_snapshot();
307            if snapshot.owner() != Some(old_owner.into()) {
308                return Err(TaskError::InvalidPiState);
309            }
310            if !snapshot.has_waiters() {
311                if !lock_state.waiters.is_empty() {
312                    return Err(TaskError::InvalidPiState);
313                }
314                // Match Linux rt_mutex_slowunlock(): once the final waiter
315                // cancels, drop wait_lock before the owner-word CAS. A waiter
316                // which marks the owner in that window makes the CAS fail, so
317                // slow unlock retakes wait_lock and selects the new top waiter.
318                drop(lock_state);
319                let released = unsafe {
320                    // SAFETY: `old_owner` came from the owner-authorized raw
321                    // release transition, and this loop retains that authority.
322                    mutex_core.try_release_for_thread(old_owner)
323                }?;
324                if released {
325                    return Ok(());
326                }
327                continue;
328            }
329            let selected_entry = lock_state
330                .waiters
331                .first_entry()
332                .ok_or(TaskError::InvalidPiState)?;
333            let selected_key = selected_entry.0;
334            let selected = selected_entry
335                .1
336                .waiter_core()
337                .ok_or(TaskError::InvalidPiState)?;
338            let selected_generation = selected_entry
339                .1
340                .wait_generation()
341                .ok_or(TaskError::InvalidPiState)?;
342            if selected.id() != selected_key.thread
343                || !selected.pi_wait_state().can_grant(selected_generation)
344            {
345                return Err(TaskError::InvalidPiState);
346            }
347            {
348                let mut old_owner_sched = old_owner_core.sched().lock();
349                if !old_owner_sched.pi.donors.contains(selected_key) {
350                    return Err(TaskError::InvalidPiState);
351                }
352                self.replace_owner_lock_top_locked(
353                    &old_owner_core,
354                    &mut old_owner_sched,
355                    Some(selected_entry),
356                    None,
357                )?;
358            }
359            mutex_core.publish_ownerless();
360            // Snapshot the wake channel while wait_lock still protects this
361            // registration. A delayed lock wake must never notify a later wait.
362            let wake = PiHandoffWake::new(selected, selected_generation);
363            drop(lock_state);
364            break wake;
365        };
366        self.recompute_pi_cleanup_chain(old_owner, wake.core.id())
367            .unwrap_or_else(|_| {
368                task_runtime::fatal_invariant(0x5049_1211, old_owner.as_u64() as usize)
369            });
370
371        // A selected waiter may claim before delivery. The retained identity
372        // and generation make a delayed RT-lock wake harmless to its next wait.
373        wake.deliver(self);
374        Ok(())
375    }
376
377    /// Claims an ownerless handoff selected for this waiter.
378    pub fn pi_mutex_claim(&self, token: &PiWaitToken) -> Result<PiMutexClaimOutcome, TaskError> {
379        let _preempt = PreemptScope::enter();
380        let claimant = ThreadId::from(token.thread_id());
381        let claimant_core = self.pi_thread_core(claimant)?;
382        let lock = token.lock_raw();
383        let mut lock_state = unsafe {
384            // SAFETY: the borrowed token keeps the physical mutex core live.
385            lock_raw_pi_mutex_waiters(lock)
386        };
387        let mutex_core = unsafe {
388            // SAFETY: the token lifetime is borrowed from this mutex core.
389            lock.core()
390        };
391        if !mutex_core.owner_snapshot().is_ownerless() {
392            return Ok(PiMutexClaimOutcome::Retry);
393        }
394        let Some(registration) = self.claim_ownerless_lock_waiter(
395            &mut lock_state,
396            &claimant_core,
397            lock,
398            token.generation(),
399        )?
400        else {
401            return Ok(PiMutexClaimOutcome::Retry);
402        };
403        debug_assert_eq!(registration.generation, token.generation());
404        mutex_core.publish_owner(claimant.into(), !lock_state.waiters.is_empty());
405        claimant_core
406            .pi_wait_state()
407            .grant(token.generation())
408            .unwrap_or_else(|_| {
409                task_runtime::fatal_invariant(0x5049_1213, claimant.as_u64() as usize)
410            });
411        drop(lock_state);
412        self.recompute_pi_cleanup_chain(claimant, claimant)
413            .unwrap_or_else(|_| {
414                task_runtime::fatal_invariant(0x5049_1214, claimant.as_u64() as usize)
415            });
416        Ok(PiMutexClaimOutcome::Claimed)
417    }
418
419    pub(crate) fn pi_initial_owner_is_on_cpu(
420        &self,
421        token: &PiWaitToken,
422    ) -> Result<bool, TaskError> {
423        let Some(owner) = token.initial_owner() else {
424            return Ok(false);
425        };
426        let owner = token
427            .initial_owner_handle()
428            .filter(|handle| handle.id() == ThreadId::from(owner))
429            .ok_or(TaskError::InvalidPiState)?
430            .runtime_core_arc();
431        Ok(owner.sched().scheduler_fence_cpu().is_some())
432    }
433}