Skip to main content

ax_task/thread/
park.rs

1//! Generation-checked thread park handshake.
2
3use core::sync::atomic::{AtomicU8, Ordering};
4
5use crate::{
6    runtime::switch::ScheduleDecision,
7    thread::ThreadId,
8    time::queue::{TaskDeadlineRegistration, TaskDeadlineToken},
9};
10
11const WAIT_WAKE_QUEUED: u8 = 0;
12const WAIT_WAKE_SELECTED: u8 = 1;
13const WAIT_WAKE_DELIVERED: u8 = 2;
14const WAIT_WAKE_CANCELLED: u8 = 3;
15const WAIT_WAKE_INACTIVE: u8 = 4;
16
17/// One queue notification claim state machine bound to an exact park attempt.
18///
19/// The containing wait queue owns entry order, while this atomic state owns
20/// selection against concurrent cleanup. The scheduler owns delivery after all
21/// fallible placement preparation, while timeout cleanup may close a selected
22/// claim before that delivery point. An unavailable scheduler owner may return
23/// a cancelled claim to the same wait entry for another selection attempt.
24#[derive(Debug)]
25pub(crate) struct WaitWakeClaim {
26    thread: ThreadId,
27    park_generation: u64,
28    state: AtomicU8,
29}
30
31impl WaitWakeClaim {
32    pub(crate) const fn new(thread: ThreadId, park_generation: u64) -> Self {
33        Self {
34            thread,
35            park_generation,
36            state: AtomicU8::new(WAIT_WAKE_QUEUED),
37        }
38    }
39
40    pub(crate) const fn thread(&self) -> ThreadId {
41        self.thread
42    }
43
44    pub(crate) const fn park_generation(&self) -> u64 {
45        self.park_generation
46    }
47
48    pub(crate) fn state(&self) -> WaitWakeClaimState {
49        match self.state.load(Ordering::Acquire) {
50            WAIT_WAKE_QUEUED => WaitWakeClaimState::Queued,
51            WAIT_WAKE_SELECTED => WaitWakeClaimState::Selected,
52            WAIT_WAKE_DELIVERED => WaitWakeClaimState::Delivered,
53            WAIT_WAKE_CANCELLED => WaitWakeClaimState::Cancelled,
54            WAIT_WAKE_INACTIVE => WaitWakeClaimState::Inactive,
55            _ => unreachable!("wait-wake claim state must be one of five closed states"),
56        }
57    }
58
59    pub(crate) fn is_active(&self) -> bool {
60        matches!(
61            self.state(),
62            WaitWakeClaimState::Queued
63                | WaitWakeClaimState::Selected
64                | WaitWakeClaimState::Cancelled
65        )
66    }
67
68    pub(crate) fn select(&self) -> bool {
69        self.state
70            .compare_exchange(
71                WAIT_WAKE_QUEUED,
72                WAIT_WAKE_SELECTED,
73                Ordering::AcqRel,
74                Ordering::Acquire,
75            )
76            .is_ok()
77    }
78
79    pub(crate) fn deliver_selected(&self) -> bool {
80        self.state
81            .compare_exchange(
82                WAIT_WAKE_SELECTED,
83                WAIT_WAKE_DELIVERED,
84                Ordering::AcqRel,
85                Ordering::Acquire,
86            )
87            .is_ok()
88    }
89
90    pub(crate) fn cancel_selected(&self) -> bool {
91        self.state
92            .compare_exchange(
93                WAIT_WAKE_SELECTED,
94                WAIT_WAKE_CANCELLED,
95                Ordering::AcqRel,
96                Ordering::Acquire,
97            )
98            .is_ok()
99    }
100
101    /// Returns a synchronously rejected delivery to its owning wait entry.
102    ///
103    /// The scheduler does not retain this claim after returning `Unavailable`.
104    /// This transition races cleanup through the same atomic state: either the
105    /// claim becomes queued again or cleanup closes it as inactive.
106    pub(crate) fn requeue_cancelled(&self) -> bool {
107        self.state
108            .compare_exchange(
109                WAIT_WAKE_CANCELLED,
110                WAIT_WAKE_QUEUED,
111                Ordering::AcqRel,
112                Ordering::Acquire,
113            )
114            .is_ok()
115    }
116
117    /// Closes this wait entry against later selection or unavailable requeue.
118    ///
119    /// Selection, scheduler delivery, and cleanup all transition this one
120    /// atomic state. The scheduler's task lock still owns runnable placement;
121    /// no separate wait-entry lock is needed to serialize these terminal CAS
122    /// operations.
123    pub(crate) fn deactivate(&self) -> bool {
124        loop {
125            let current = self.state.load(Ordering::Acquire);
126            match current {
127                WAIT_WAKE_DELIVERED => return true,
128                WAIT_WAKE_INACTIVE => return false,
129                WAIT_WAKE_QUEUED | WAIT_WAKE_SELECTED | WAIT_WAKE_CANCELLED => {
130                    if self
131                        .state
132                        .compare_exchange(
133                            current,
134                            WAIT_WAKE_INACTIVE,
135                            Ordering::AcqRel,
136                            Ordering::Acquire,
137                        )
138                        .is_ok()
139                    {
140                        return false;
141                    }
142                }
143                _ => unreachable!("wait-wake claim state must be one of five closed states"),
144            }
145        }
146    }
147}
148
149#[derive(Clone, Copy, Debug, Eq, PartialEq)]
150pub(crate) enum WaitWakeClaimState {
151    Queued,
152    Selected,
153    Delivered,
154    Cancelled,
155    Inactive,
156}
157
158#[derive(Clone, Copy, Debug, Eq, PartialEq)]
159pub(crate) enum WaitWakeDelivery {
160    Delivered,
161    Cancelled,
162    Exited,
163    Unavailable,
164}
165
166/// Move-only ownership of one park attempt and its optional timeout deadline.
167///
168/// ```compile_fail
169/// # use ax_task::thread::ParkTicket;
170/// fn duplicate(ticket: ParkTicket) {
171///     let first = ticket;
172///     let second = ticket;
173///     drop((first, second));
174/// }
175/// ```
176#[must_use = "a prepared park and its deadline must be committed or cancelled"]
177#[derive(Debug, Eq, PartialEq)]
178pub struct ParkTicket {
179    thread: ThreadId,
180    generation: u64,
181    deadline: Option<TaskDeadlineRegistration>,
182    resolved: bool,
183}
184
185impl ParkTicket {
186    pub(crate) const fn new(thread: ThreadId, generation: u64) -> Self {
187        Self {
188            thread,
189            generation,
190            deadline: None,
191            resolved: false,
192        }
193    }
194
195    /// Returns the thread that prepared this park attempt.
196    pub const fn thread(&self) -> ThreadId {
197        self.thread
198    }
199
200    /// Returns the monotonically increasing attempt generation.
201    pub const fn generation(&self) -> u64 {
202        self.generation
203    }
204
205    pub(crate) fn attach_deadline(
206        &mut self,
207        deadline: TaskDeadlineRegistration,
208    ) -> Result<(), TaskDeadlineRegistration> {
209        if self.deadline.is_some() {
210            Err(deadline)
211        } else {
212            self.deadline = Some(deadline);
213            Ok(())
214        }
215    }
216
217    pub(crate) const fn deadline(&self) -> Option<&TaskDeadlineRegistration> {
218        self.deadline.as_ref()
219    }
220
221    pub(crate) fn clear_deadline(&mut self, deadline: TaskDeadlineToken) -> bool {
222        if self
223            .deadline
224            .as_ref()
225            .is_some_and(|registration| registration.token() == deadline)
226        {
227            let _registration = self.deadline.take();
228            true
229        } else {
230            false
231        }
232    }
233
234    pub(crate) fn mark_resolved(&mut self) {
235        self.resolved = true;
236    }
237
238    pub(crate) const fn is_resolved(&self) -> bool {
239        self.resolved
240    }
241
242    pub(crate) const fn has_deadline(&self) -> bool {
243        self.deadline.is_some()
244    }
245}
246
247/// Result of publishing the `PARKING` phase.
248#[derive(Debug, Eq, PartialEq)]
249pub enum ParkPrepare {
250    /// A preceding notification was consumed, so the caller must not block.
251    Notified,
252    /// The caller published `PARKING` and may proceed to the commit phase.
253    Prepared(ParkTicket),
254}
255
256/// Result of rechecking a prepared park at the scheduler safe point.
257#[derive(Debug)]
258pub enum ParkCommit {
259    /// A concurrent notification cancelled the park before schedule-out.
260    Notified,
261    /// The thread committed `BLOCKED` and selected its replacement.
262    Blocked(ScheduleDecision),
263}
264
265#[cfg(test)]
266mod tests {
267    use super::*;
268
269    #[test]
270    fn cancelled_wait_wake_claim_can_retry_the_same_park() {
271        let claim = WaitWakeClaim::new(ThreadId::from_parts(7, 3), 11);
272
273        assert!(claim.select());
274        assert!(claim.cancel_selected());
275        assert!(claim.requeue_cancelled());
276        assert!(claim.select());
277        assert!(claim.deliver_selected());
278        assert_eq!(claim.state(), WaitWakeClaimState::Delivered);
279    }
280
281    #[test]
282    fn delivered_wait_wake_claim_cannot_be_requeued() {
283        let claim = WaitWakeClaim::new(ThreadId::from_parts(7, 3), 11);
284
285        assert!(claim.select());
286        assert!(claim.deliver_selected());
287        assert!(!claim.requeue_cancelled());
288        assert_eq!(claim.state(), WaitWakeClaimState::Delivered);
289    }
290
291    #[test]
292    fn deactivated_queued_claim_cannot_be_selected_or_requeued() {
293        let claim = WaitWakeClaim::new(ThreadId::from_parts(7, 3), 11);
294
295        assert!(!claim.deactivate());
296        assert!(!claim.select());
297        assert!(!claim.requeue_cancelled());
298        assert!(!claim.is_active());
299        assert_eq!(claim.state(), WaitWakeClaimState::Inactive);
300    }
301
302    #[test]
303    fn deactivation_cancels_a_selected_claim_before_delivery() {
304        let claim = WaitWakeClaim::new(ThreadId::from_parts(7, 3), 11);
305
306        assert!(claim.select());
307        assert!(!claim.deactivate());
308        assert!(!claim.deliver_selected());
309        assert_eq!(claim.state(), WaitWakeClaimState::Inactive);
310    }
311
312    #[test]
313    fn deactivation_observes_a_delivered_claim_idempotently() {
314        let claim = WaitWakeClaim::new(ThreadId::from_parts(7, 3), 11);
315
316        assert!(claim.select());
317        assert!(claim.deliver_selected());
318        assert!(claim.deactivate());
319        assert!(claim.deactivate());
320        assert_eq!(claim.state(), WaitWakeClaimState::Delivered);
321    }
322
323    #[test]
324    fn deactivated_cancelled_claim_cannot_be_requeued() {
325        let claim = WaitWakeClaim::new(ThreadId::from_parts(7, 3), 11);
326
327        assert!(claim.select());
328        assert!(claim.cancel_selected());
329        assert!(!claim.deactivate());
330        assert!(!claim.requeue_cancelled());
331        assert_eq!(claim.state(), WaitWakeClaimState::Inactive);
332    }
333}