ax_task/sched/system/task_system/
placement.rs1use super::*;
4
5impl TaskSystem {
6 pub(super) fn affinity_schedule_state_locked(
13 &self,
14 core: &Arc<ThreadCore>,
15 sched: &ThreadSchedState,
16 ) -> Result<(SchedulePolicy, SchedulingEntity), TaskError> {
17 if let Some(active) = core.sched().active_option(sched) {
18 return Ok((active.policy(), active.entity().clone()));
19 }
20 let owner = sched
21 .placement
22 .control_owner()
23 .ok_or(TaskError::InvalidConfiguration)?;
24 let remote = self
25 .cpu_remote(owner)
26 .ok_or(TaskError::CpuOffline(owner.as_u32()))?;
27 let transaction = OwnerRqTxn::begin(self, remote);
28 let state = transaction
29 .scheduling_state(core.id())
30 .ok_or(TaskError::InvalidConfiguration);
31 transaction.commit();
32 state
33 }
34
35 pub(super) fn complete_affinity_if_satisfied_locked(
36 core: &Arc<ThreadCore>,
37 sched: &ThreadSchedState,
38 ) -> bool {
39 if sched.lifecycle.state() == ThreadState::Exited || sched.placement.has_pending_migration()
40 {
41 return false;
42 }
43 let placement_is_allowed = [sched.placement.queued_cpu(), sched.placement.on_cpu()]
44 .into_iter()
45 .flatten()
46 .all(|cpu| sched.affinity.requested_affinity.contains(cpu));
47 if !placement_is_allowed {
48 return false;
49 }
50 core.publish_affinity_completion(sched.affinity.affinity_generation)
51 }
52
53 pub(super) fn prepare_owner_migration(
54 &self,
55 core: &Arc<ThreadCore>,
56 source: CpuId,
57 target: CpuId,
58 ) -> Result<PreparedMigrationDelivery, TaskError> {
59 let target_remote = self
60 .cpu_remotes
61 .get(target.as_usize())
62 .ok_or(TaskError::InvalidCpu(target.as_u32()))?;
63 PreparedMigrationDelivery::prepare(target_remote, core, source, target)
64 }
65
66 pub(super) fn publish_owner_deadline_refresh(&self, core: &Arc<ThreadCore>, owner: CpuId) {
67 let remote = &self.cpu_remotes[owner.as_usize()];
68 let publication = remote
69 .begin_owner_delivery()
70 .unwrap_or_else(|| task_runtime::fatal_invariant(0x444c_0010, owner.as_u32() as usize));
71 if !core.reserve_scheduler_inbox_delivery() {
72 return;
73 }
74 let pointer = Arc::as_ptr(core);
75 unsafe { Arc::increment_strong_count(pointer) };
78 let node = unsafe { Pin::new_unchecked((*pointer).deadline_refresh_node()) };
80 let message = InboxMessage::deadline_refresh_with_payload(
81 core.id(),
82 owner,
83 0,
84 pointer.expose_provenance(),
85 );
86 if publication.publish_owner_control(node, message) != PublishResult::Published {
87 unsafe { Arc::decrement_strong_count(pointer) };
89 core.cancel_scheduler_inbox_delivery();
90 }
91 }
92
93 pub(super) fn publish_owner_affinity_retry(
94 &self,
95 core: &Arc<ThreadCore>,
96 owner: CpuId,
97 target: CpuId,
98 ) -> Result<(), TaskError> {
99 let remote = self
100 .cpu_remote(owner)
101 .ok_or(TaskError::CpuOffline(owner.as_u32()))?;
102 if !core.reserve_scheduler_inbox_delivery() {
103 return Ok(());
104 }
105 let pointer = Arc::as_ptr(core);
106 unsafe { Arc::increment_strong_count(pointer) };
109 let node = unsafe { Pin::new_unchecked((*pointer).affinity_update_node()) };
111 let message = InboxMessage::affinity_update_with_payload(
112 core.id(),
113 owner,
114 target,
115 pointer.expose_provenance(),
116 );
117 if remote.publish_owner_control(node, message) != PublishResult::Published {
118 unsafe { Arc::decrement_strong_count(pointer) };
121 core.cancel_scheduler_inbox_delivery();
122 }
123 Ok(())
124 }
125
126 pub fn request_thread_affinity(
128 &self,
129 thread: ThreadId,
130 affinity: CpuSet,
131 ) -> Result<ThreadAffinityChange, TaskError> {
132 if task_runtime::in_hard_irq() {
133 return Err(TaskError::UnsafeContext);
134 }
135 validate_affinity(&affinity, self.config.cpu_count())?;
136 let mut state = self.state.lock();
137 let root_domain = self.root_domain.lock();
138 let record = state.thread_record(thread)?;
139 let core = Arc::clone(&record.core);
140 let mut sched = record.sched.lock();
141 if sched.lifecycle.state() == ThreadState::Exited {
142 return Err(TaskError::NotReady);
143 }
144 let is_deadline = matches!(sched.policy.base, SchedulePolicy::Deadline(_))
145 || matches!(sched.policy.requested_policy(), SchedulePolicy::Deadline(_));
146 if is_deadline && !affinity.covers(&root_domain.online) {
147 return Err(TaskError::DeadlineAffinity);
148 }
149 let timer_cpu = core.sleep_timer_cpu();
150 if timer_cpu.is_some_and(|cpu| !affinity.contains(cpu)) {
151 return Err(TaskError::ActiveTimerAffinity);
152 }
153 let (policy, entity) = self.affinity_schedule_state_locked(&core, &sched)?;
154 let preferred = sched
155 .placement
156 .control_owner()
157 .or_else(|| core.wake_cpu_hint());
158 drop(root_domain);
159 let target = timer_cpu
160 .or_else(|| match policy {
161 SchedulePolicy::Fair { .. } => state.select_initial_fair_cpu(&affinity, preferred),
162 SchedulePolicy::Fifo { .. }
163 | SchedulePolicy::RoundRobin { .. }
164 | SchedulePolicy::Deadline(_) => {
165 self.select_priority_cpu(policy, Some(&entity), &affinity, preferred, None)
166 }
167 SchedulePolicy::KernelStop => self.select_fallback_active_cpu(&affinity, None),
168 })
169 .ok_or(TaskError::InvalidConfiguration)?;
170 let replacement = if record.activation.is_some() {
173 Some(self.prepare_owner_migration(
174 &core,
175 sched.placement.assigned_cpu().expect("new task CPU"),
176 target,
177 )?)
178 } else {
179 None
180 };
181 let generation = sched
182 .affinity
183 .affinity_generation
184 .checked_add(1)
185 .ok_or(TaskError::InvalidConfiguration)?;
186 sched.affinity.affinity_generation = generation;
187 sched.affinity.requested_affinity = Arc::new(affinity);
188 if sched.affinity.migration_depth == 0 {
189 sched.affinity.affinity = Arc::clone(&sched.affinity.requested_affinity);
190 }
191 let owner = sched.placement.control_owner();
196 let target = owner
197 .filter(|owner| sched.affinity.affinity.contains(*owner))
198 .unwrap_or(target);
199 core.set_wake_cpu_hint(target);
200 let completed = Self::complete_affinity_if_satisfied_locked(&core, &sched);
201 drop(sched);
202 let publication = owner.map_or(Ok(()), |owner| {
203 state.publish_affinity_update(&core, owner, target)
204 });
205 let previous = replacement.and_then(|delivery| {
206 state
207 .thread_record_mut(thread)
208 .expect("locked task identity")
209 .activation
210 .replace(delivery)
211 });
212 drop(state);
213 drop(previous);
214 if completed {
215 core.notify_affinity_waiters();
216 }
217 publication?;
218 Ok(ThreadAffinityChange::new(core, generation))
219 }
220
221 pub fn set_current_affinity(
228 &self,
229 cpu: Pin<&mut CpuLocal>,
230 affinity: CpuSet,
231 ) -> Result<bool, TaskError> {
232 self.ensure_owner_cpu_context(&cpu)?;
233 validate_affinity(&affinity, self.config.cpu_count())?;
234 self.ensure_owner_cpu_online(&cpu)?;
235 let root_domain = self.root_domain.lock();
236 let core = cpu.current_core().ok_or(TaskError::NoRunnableThread)?;
237 let current = core.id();
238 let mut sched = core.sched().lock();
239 if sched.placement.queued_cpu() != Some(cpu.owner())
240 || sched.placement.on_cpu() != Some(cpu.owner())
241 {
242 return Err(TaskError::InvalidConfiguration);
243 }
244 let is_deadline = matches!(sched.policy.base, SchedulePolicy::Deadline(_))
245 || matches!(sched.policy.requested_policy(), SchedulePolicy::Deadline(_));
246 if is_deadline && !affinity.covers(&root_domain.online) {
247 return Err(TaskError::DeadlineAffinity);
248 }
249 let timer_cpu = core.sleep_timer_cpu();
250 if timer_cpu.is_some_and(|timer_cpu| !affinity.contains(timer_cpu)) {
251 return Err(TaskError::ActiveTimerAffinity);
252 }
253 let owner = cpu.owner();
254 if sched.affinity.migration_depth != 0 && !affinity.contains(owner) {
255 return Err(TaskError::UnsafeContext);
256 }
257 let must_migrate = !affinity.contains(owner);
258 let remote = Arc::clone(cpu.remote());
259 let mut transaction = OwnerRqTxn::begin(self, &remote);
260 if transaction.current_thread() != Some(current)
261 || transaction
262 .current_core()
263 .is_none_or(|current_core| !Arc::ptr_eq(¤t_core, &core))
264 {
265 transaction.commit();
266 return Err(TaskError::InvalidConfiguration);
267 }
268 let selection = transaction
269 .scheduling_state(current)
270 .ok_or(TaskError::InvalidConfiguration)
271 .and_then(|(policy, entity)| {
272 timer_cpu
273 .or_else(|| {
274 self.select_priority_cpu(
275 policy,
276 Some(&entity),
277 &affinity,
278 Some(owner),
279 must_migrate.then_some(owner),
280 )
281 })
282 .ok_or(TaskError::InvalidConfiguration)
283 })
284 .and_then(|target| {
285 sched
286 .affinity
287 .affinity_generation
288 .checked_add(1)
289 .map(|generation| (target, generation))
290 .ok_or(TaskError::InvalidConfiguration)
291 });
292 let (target, generation) = match selection {
293 Ok(selection) => selection,
294 Err(error) => {
295 transaction.commit();
296 return Err(error);
297 }
298 };
299 sched.affinity.affinity_generation = generation;
300 sched.affinity.requested_affinity = Arc::new(affinity);
301 if sched.affinity.migration_depth == 0 {
302 sched.affinity.affinity = Arc::clone(&sched.affinity.requested_affinity);
303 }
304 transaction.update_thread_affinity(current, Arc::clone(&sched.affinity.affinity));
305 sched
306 .placement
307 .request_migration(must_migrate.then_some(target));
308 core.set_wake_cpu_hint(if must_migrate { target } else { owner });
309 let completed = Self::complete_affinity_if_satisfied_locked(&core, &sched);
310 transaction.commit();
311 drop(sched);
312 drop(root_domain);
313 if completed {
314 core.notify_affinity_waiters();
315 }
316 if must_migrate {
317 cpu.request_reschedule(RescheduleKind::Immediate);
318 }
319 Ok(must_migrate)
320 }
321}