Skip to main content

ax_task/sched/system/task_system/
lifecycle.rs

1//! Thread exit callbacks, registry reaping, and resource release.
2
3use core::sync::atomic::Ordering;
4
5use super::*;
6
7impl TaskSystem {
8    /// Marks an unmanaged, non-queued thread exited and queues its exit hook.
9    /// Managed creation tokens exclusively own cancellation of their tasks.
10    pub fn mark_exited(&self, thread: ThreadId) -> Result<(), TaskError> {
11        crate::runtime::delivery::work::validate_task_work_context()?;
12        let core = {
13            let state = self.state.lock();
14            Arc::clone(&state.thread_record(thread)?.core)
15        };
16        if core.execution.is_some() {
17            return Err(TaskError::NotReady);
18        }
19        self.mark_unqueued_exited(&core)
20    }
21
22    /// Consumes exit authority held by a validated caller or cancellation worker.
23    pub(super) fn mark_unqueued_exited(&self, core: &Arc<ThreadCore>) -> Result<(), TaskError> {
24        let thread = core.id();
25        let mut scheduler_exit = core
26            .close_owned_scheduler_activity()
27            .ok_or(TaskError::ThreadBusy)?;
28        let exited_core = {
29            let mut state = self.state.lock();
30            let record = state.thread_record_mut(thread)?;
31            if !Arc::ptr_eq(&record.core, core) {
32                return Err(TaskError::StaleThreadId);
33            }
34            let mut sched = record.sched.lock();
35            if record.activation.is_some() {
36                return Err(TaskError::ThreadBusy);
37            }
38            if sched.placement.queued_cpu().is_some() {
39                return Err(TaskError::AlreadyQueued);
40            }
41            if sched.placement.on_cpu().is_some() {
42                return Err(TaskError::ThreadBusy);
43            }
44            if sched.pi.blocked_on.is_some() || !sched.pi.donors.is_empty() {
45                return Err(TaskError::InvalidPiState);
46            }
47            let lifecycle = sched.lifecycle.state();
48            if !crate::thread::transition_is_valid(lifecycle, ThreadState::Exited) {
49                return Err(TaskError::InvalidTransition {
50                    from: lifecycle,
51                    to: ThreadState::Exited,
52                });
53            }
54            record.callbacks.validate_prepare_exit()?;
55            if let Some(owner) = sched.deadline.bandwidth.reservation_owner() {
56                let remote = self
57                    .cpu_remotes
58                    .get(owner.as_usize())
59                    .ok_or(TaskError::InvalidCpu(owner.as_u32()))?;
60                if !remote.is_online() {
61                    task_runtime::fatal_invariant(0x444c_1203, core.id().as_u64() as usize);
62                }
63                let mut transaction = OwnerRqTxn::begin(self, remote);
64                Self::detach_owner_deadline_bandwidth_in_rq(
65                    &record.core,
66                    &mut sched,
67                    remote,
68                    &mut transaction,
69                );
70                transaction.commit();
71                // A stale physical clockevent is harmless, but the owning CPU
72                // must promptly publish the new earliest scheduler deadline
73                // instead of waiting for that stale edge to fire.
74                remote.request_scheduler_work();
75            }
76            sched.placement.cancel_remote_handoff_for_exit();
77            sched
78                .transition(&record.core, ThreadState::Exited)
79                .unwrap_or_else(|_| {
80                    task_runtime::fatal_invariant(0x4558_000a, core.id().as_u64() as usize)
81                });
82            scheduler_exit.seal();
83            record
84                .callbacks
85                .prepare_exit(record.extension.is_some() || record.core.execution.is_some())
86                .unwrap_or_else(|_| {
87                    task_runtime::fatal_invariant(0x4558_000b, core.id().as_u64() as usize)
88                });
89            let exited_core = Arc::clone(&record.core);
90            drop(sched);
91            state.queue_exited_thread(thread);
92            let mut root_domain = self.root_domain.lock();
93            let released = state
94                .release_deadline_reservation_on_exit(thread)
95                .unwrap_or_else(|_| {
96                    task_runtime::fatal_invariant(0x4558_000c, core.id().as_u64() as usize)
97                });
98            root_domain.release_deadline(released);
99            exited_core
100        };
101        exited_core.notify_affinity_waiters();
102        self.task_work.publish();
103        Ok(())
104    }
105
106    /// Runs pending exit callbacks from an ordinary task-context safe point.
107    ///
108    /// Context-switch tail only proves that the exited stack is inactive; its
109    /// inherited IRQ and scheduler guards are still live. Calling an OS exit
110    /// hook there can acquire a sleepable lock and recursively enter the
111    /// scheduler. This bounded pass claims each callback under the registry
112    /// lock, invokes it without scheduler locks, and only then makes the record
113    /// eligible for reaping.
114    pub fn dispatch_exit_callbacks(&self, limit: usize) -> Result<usize, TaskError> {
115        if task_runtime::in_hard_irq() {
116            return Err(TaskError::UnsafeContext);
117        }
118        let _consumer = self.task_work.try_claim_consumer()?;
119        self.dispatch_exit_callbacks_inner(limit)
120    }
121
122    pub(super) fn dispatch_exit_callbacks_inner(&self, limit: usize) -> Result<usize, TaskError> {
123        let mut dispatched = 0;
124        while dispatched < limit {
125            let callback = {
126                let mut state = self.state.lock();
127                state.claim_pending_exit_callback()?
128            };
129            let Some(super::registry::ExitCallbackClaim { extension, core }) = callback else {
130                break;
131            };
132            // SAFETY: the registry record keeps the claimed extension live,
133            // and ThreadExtension construction validated this callback table.
134            if let Some(extension) = extension {
135                unsafe { (extension.ops().on_exit)(extension.data(), core.id()) };
136            }
137            if let Some(execution) = core.execution.as_ref() {
138                execution.finish();
139            }
140            self.state.lock().finish_exit_callback(core.id())?;
141            dispatched += 1;
142        }
143        Ok(dispatched)
144    }
145
146    /// Removes an exited registry record and makes its slot reusable.
147    pub fn reap_thread(&self, thread: ThreadId) -> Result<(), TaskError> {
148        crate::runtime::delivery::work::validate_task_work_context()?;
149        let record = {
150            let mut state = self.state.lock();
151            let mut root_domain = self.root_domain.lock();
152            let (record, released) = state.remove_exited_thread(thread)?;
153            root_domain.release_deadline(released);
154            record
155        };
156        self.release_thread_record(record);
157        Ok(())
158    }
159
160    /// Atomically removes an exited thread while consuming its owning handle.
161    ///
162    /// Keeping `handle` alive until registry removal prevents the detached
163    /// reaper on another CPU from winning between a handle drop and an ID-based
164    /// reap. Retryable failures return the same handle to the caller.
165    pub fn reap_thread_handle(&self, handle: ThreadHandle) -> Result<(), OwnedThreadReapError> {
166        if let Err(error) = crate::runtime::delivery::work::validate_task_work_context() {
167            return Err(OwnedThreadReapError::new(error, handle));
168        }
169        let record = {
170            let mut state = self.state.lock();
171            let mut root_domain = self.root_domain.lock();
172            match state.remove_exited_thread_with_handle(&handle) {
173                Ok((record, released)) => {
174                    root_domain.release_deadline(released);
175                    record
176                }
177                Err(error) => return Err(OwnedThreadReapError::new(error, handle)),
178            }
179        };
180        drop(handle);
181        self.release_thread_record(record);
182        Ok(())
183    }
184
185    /// Reaps exited records for which no external strong handle remains.
186    ///
187    /// This bounded task-context pass is the detached-thread reaper. Joinable
188    /// threads remain registered because their [`ThreadHandle`] contributes a
189    /// strong reference. Late IRQ wake handles likewise delay resource release
190    /// until their final reference reaches the task-context reaper.
191    pub fn reap_unreferenced_exited(&self, limit: usize) -> Result<usize, TaskError> {
192        if task_runtime::in_hard_irq() {
193            return Err(TaskError::UnsafeContext);
194        }
195        let _consumer = self.task_work.try_claim_consumer()?;
196        self.reap_unreferenced_exited_inner(limit)
197    }
198
199    pub(super) fn reap_unreferenced_exited_inner(&self, limit: usize) -> Result<usize, TaskError> {
200        let mut reaped = 0;
201        while reaped < limit {
202            let removed = {
203                let mut state = self.state.lock();
204                let mut root_domain = self.root_domain.lock();
205                let removed = state.take_unreferenced_exited()?;
206                if let Some((_, released)) = &removed {
207                    root_domain.release_deadline(*released);
208                }
209                removed
210            };
211            let Some((record, _released)) = removed else {
212                break;
213            };
214            self.release_thread_record(record);
215            reaped += 1;
216        }
217        Ok(reaped)
218    }
219
220    pub(super) fn reclaim_exited_execution(&self) -> Result<bool, TaskError> {
221        let detached = self.state.lock().take_exited_execution()?;
222        let Some(detached) = detached else {
223            return Ok(false);
224        };
225        let address_space = detached.resources.release();
226        self.release_address_space_token(address_space);
227        detached
228            .handle
229            .core
230            .execution_reclaimed
231            .store(true, Ordering::Release);
232        drop(detached.handle);
233        Ok(true)
234    }
235
236    pub(super) fn release_thread_record(&self, mut record: ThreadRecord) {
237        let address_space = record.resources.release();
238        record
239            .core
240            .execution_reclaimed
241            .store(true, Ordering::Release);
242        drop(record.extension.take());
243        self.release_address_space_token(address_space);
244    }
245
246    pub(super) fn release_unpublished_thread(&self, record: DetachedThreadRecord) {
247        let address_space = record.release();
248        self.release_address_space_token(address_space);
249    }
250
251    /// Releases a construction transaction that failed before thread registry
252    /// publication.
253    ///
254    /// Thread-private destruction is a one-way ownership transfer. Only the
255    /// independent active-mm token may outlive this call through the runtime's
256    /// explicit last-CPU readiness edge.
257    pub fn release_unpublished_resources(&self, resources: ThreadResources) {
258        self.release_unpublished_thread(DetachedThreadRecord::new(resources, None))
259    }
260
261    pub(crate) fn release_address_space_token(
262        &self,
263        address_space: crate::runtime::resource::AddressSpaceToken,
264    ) {
265        if address_space.is_none() {
266            return;
267        }
268        let handle = address_space.handle();
269        match task_runtime::destroy_address_space(handle) {
270            AddressSpaceDestroyOutcome::Released => return,
271            AddressSpaceDestroyOutcome::Active => {}
272        }
273        self.state
274            .lock()
275            .pending_address_space_reclaims
276            .push(address_space);
277        match task_runtime::arm_address_space_reclaim(handle) {
278            AddressSpaceReclaimArmOutcome::Ready => self.task_work.publish(),
279            AddressSpaceReclaimArmOutcome::Armed => {}
280        }
281    }
282}