ax_task/sched/system/task_system/pi/
operations.rs1use 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 lock_raw_pi_mutex_waiters(lock)
21 };
22 let core = unsafe {
23 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 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 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 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 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 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 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 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 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 drop(lock_state);
320 let released = unsafe {
321 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 let _woken = wakes.wake_all();
377 Ok(())
378 }
379
380 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 lock_raw_pi_mutex_waiters(lock)
389 };
390 let mutex_core = unsafe {
391 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}