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 mut wakes = ThreadWakeBatch::new();
304        let selected_id = loop {
305            let mutex_core = lock.core();
306            let lock_state = lock_pi_mutex_waiters(lock);
307            let snapshot = mutex_core.owner_snapshot();
308            if snapshot.owner() != Some(old_owner.into()) {
309                return Err(TaskError::InvalidPiState);
310            }
311            if !snapshot.has_waiters() {
312                if !lock_state.waiters.is_empty() {
313                    return Err(TaskError::InvalidPiState);
314                }
315                // Match Linux rt_mutex_slowunlock(): once the final waiter
316                // cancels, drop wait_lock before the owner-word CAS. A waiter
317                // which marks the owner in that window makes the CAS fail, so
318                // slow unlock retakes wait_lock and selects the new top waiter.
319                drop(lock_state);
320                let released = unsafe {
321                    // SAFETY: `old_owner` came from the owner-authorized raw
322                    // release transition, and this loop retains that authority.
323                    mutex_core.try_release_for_thread(old_owner)
324                }?;
325                if released {
326                    return Ok(());
327                }
328                continue;
329            }
330            let selected_entry = lock_state
331                .waiters
332                .first_entry()
333                .ok_or(TaskError::InvalidPiState)?;
334            let selected_key = selected_entry.0;
335            let selected = selected_entry
336                .1
337                .waiter_core()
338                .ok_or(TaskError::InvalidPiState)?;
339            let selected_generation = selected_entry
340                .1
341                .wait_generation()
342                .ok_or(TaskError::InvalidPiState)?;
343            if selected.id() != selected_key.thread
344                || !selected.pi_wait_state().can_grant(selected_generation)
345            {
346                return Err(TaskError::InvalidPiState);
347            }
348            {
349                let mut old_owner_sched = old_owner_core.sched().lock();
350                if !old_owner_sched.pi.donors.contains(selected_key) {
351                    return Err(TaskError::InvalidPiState);
352                }
353                self.replace_owner_lock_top_locked(
354                    &old_owner_core,
355                    &mut old_owner_sched,
356                    Some(selected_entry),
357                    None,
358                )?;
359            }
360            mutex_core.publish_ownerless();
361            let selected_id = selected.id();
362            let _queued = wakes.push(ThreadWakeHandle::from_core(selected));
363            drop(lock_state);
364            break selected_id;
365        };
366        self.recompute_pi_cleanup_chain(old_owner, selected_id)
367            .unwrap_or_else(|_| {
368                task_runtime::fatal_invariant(0x5049_1211, old_owner.as_u64() as usize)
369            });
370
371        // The ownerless publication and waiter generation commit the mutex
372        // handoff. As with Linux wake_q, this delayed wake is only a scheduling
373        // hint: the selected waiter may claim, run, and exit before the batch
374        // drains. The batch retains its task reference and intentionally
375        // ignores a wake that no longer changes task state.
376        let _woken = wakes.wake_all();
377        Ok(())
378    }
379
380    /// Claims an ownerless handoff selected for this waiter.
381    pub fn pi_mutex_claim(&self, token: &PiWaitToken) -> Result<PiMutexClaimOutcome, TaskError> {
382        let _preempt = PreemptScope::enter();
383        let claimant = ThreadId::from(token.thread_id());
384        let claimant_core = self.pi_thread_core(claimant)?;
385        let lock = token.lock_raw();
386        let mut lock_state = unsafe {
387            // SAFETY: the borrowed token keeps the physical mutex core live.
388            lock_raw_pi_mutex_waiters(lock)
389        };
390        let mutex_core = unsafe {
391            // SAFETY: the token lifetime is borrowed from this mutex core.
392            lock.core()
393        };
394        if !mutex_core.owner_snapshot().is_ownerless() {
395            return Ok(PiMutexClaimOutcome::Retry);
396        }
397        let Some(registration) = self.claim_ownerless_lock_waiter(
398            &mut lock_state,
399            &claimant_core,
400            lock,
401            token.generation(),
402        )?
403        else {
404            return Ok(PiMutexClaimOutcome::Retry);
405        };
406        debug_assert_eq!(registration.generation, token.generation());
407        mutex_core.publish_owner(claimant.into(), !lock_state.waiters.is_empty());
408        claimant_core
409            .pi_wait_state()
410            .grant(token.generation())
411            .unwrap_or_else(|_| {
412                task_runtime::fatal_invariant(0x5049_1213, claimant.as_u64() as usize)
413            });
414        drop(lock_state);
415        self.recompute_pi_cleanup_chain(claimant, claimant)
416            .unwrap_or_else(|_| {
417                task_runtime::fatal_invariant(0x5049_1214, claimant.as_u64() as usize)
418            });
419        Ok(PiMutexClaimOutcome::Claimed)
420    }
421
422    pub(crate) fn pi_initial_owner_is_on_cpu(
423        &self,
424        token: &PiWaitToken,
425    ) -> Result<bool, TaskError> {
426        let Some(owner) = token.initial_owner() else {
427            return Ok(false);
428        };
429        let owner = token
430            .initial_owner_handle()
431            .filter(|handle| handle.id() == ThreadId::from(owner))
432            .ok_or(TaskError::InvalidPiState)?
433            .runtime_core_arc();
434        Ok(owner.sched().scheduler_fence_cpu().is_some())
435    }
436}