ax_task/sched/system/task_system/
lifecycle.rs1use core::sync::atomic::Ordering;
4
5use super::*;
6
7impl TaskSystem {
8 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 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 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 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 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 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 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 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 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}