Skip to main content

ax_task/sched/system/task_system/
lifecycle.rs

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