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.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 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 generation = sched
171 .affinity
172 .affinity_generation
173 .checked_add(1)
174 .ok_or(TaskError::InvalidConfiguration)?;
175 sched.affinity.affinity_generation = generation;
176 sched.affinity.affinity = Arc::new(affinity);
177 let owner = sched.placement.control_owner();
182 let target = owner
183 .filter(|owner| sched.affinity.affinity.contains(*owner))
184 .unwrap_or(target);
185 core.set_wake_cpu_hint(target);
186 let completed = Self::complete_affinity_if_satisfied_locked(&core, &sched);
187 drop(sched);
188 let publication = owner.map_or(Ok(()), |owner| {
189 state.publish_affinity_update(&core, owner, target)
190 });
191 drop(state);
192 if completed {
193 core.notify_affinity_waiters();
194 }
195 publication?;
196 Ok(ThreadAffinityChange::new(core, generation))
197 }
198
199 pub fn set_current_affinity(
206 &self,
207 cpu: Pin<&mut CpuLocal>,
208 affinity: CpuSet,
209 ) -> Result<bool, TaskError> {
210 self.ensure_owner_cpu_context(&cpu)?;
211 validate_affinity(&affinity, self.config.cpu_count())?;
212 self.ensure_owner_cpu_online(&cpu)?;
213 let root_domain = self.root_domain.lock();
214 let core = cpu.current_core().ok_or(TaskError::NoRunnableThread)?;
215 let current = core.id();
216 let mut sched = core.sched().lock();
217 if sched.placement.queued_cpu() != Some(cpu.owner())
218 || sched.placement.on_cpu() != Some(cpu.owner())
219 {
220 return Err(TaskError::InvalidConfiguration);
221 }
222 let is_deadline = matches!(sched.policy.base, SchedulePolicy::Deadline(_))
223 || matches!(sched.policy.requested_policy(), SchedulePolicy::Deadline(_));
224 if is_deadline && !affinity.covers(&root_domain.online) {
225 return Err(TaskError::DeadlineAffinity);
226 }
227 let timer_cpu = core.sleep_timer_cpu();
228 if timer_cpu.is_some_and(|timer_cpu| !affinity.contains(timer_cpu)) {
229 return Err(TaskError::ActiveTimerAffinity);
230 }
231 let owner = cpu.owner();
232 let must_migrate = !affinity.contains(owner);
233 let remote = Arc::clone(cpu.remote());
234 let mut transaction = OwnerRqTxn::begin(self, &remote);
235 if transaction.current_thread() != Some(current)
236 || transaction
237 .current_core()
238 .is_none_or(|current_core| !Arc::ptr_eq(¤t_core, &core))
239 {
240 transaction.commit();
241 return Err(TaskError::InvalidConfiguration);
242 }
243 let selection = transaction
244 .scheduling_state(current)
245 .ok_or(TaskError::InvalidConfiguration)
246 .and_then(|(policy, entity)| {
247 timer_cpu
248 .or_else(|| {
249 self.select_priority_cpu(
250 policy,
251 Some(&entity),
252 &affinity,
253 Some(owner),
254 must_migrate.then_some(owner),
255 )
256 })
257 .ok_or(TaskError::InvalidConfiguration)
258 })
259 .and_then(|target| {
260 sched
261 .affinity
262 .affinity_generation
263 .checked_add(1)
264 .map(|generation| (target, generation))
265 .ok_or(TaskError::InvalidConfiguration)
266 });
267 let (target, generation) = match selection {
268 Ok(selection) => selection,
269 Err(error) => {
270 transaction.commit();
271 return Err(error);
272 }
273 };
274 sched.affinity.affinity_generation = generation;
275 sched.affinity.affinity = Arc::new(affinity);
276 transaction.update_thread_affinity(current, Arc::clone(&sched.affinity.affinity));
277 sched
278 .placement
279 .request_migration(must_migrate.then_some(target));
280 core.set_wake_cpu_hint(if must_migrate { target } else { owner });
281 let completed = Self::complete_affinity_if_satisfied_locked(&core, &sched);
282 transaction.commit();
283 drop(sched);
284 drop(root_domain);
285 if completed {
286 core.notify_affinity_waiters();
287 }
288 if must_migrate {
289 cpu.request_reschedule(RescheduleKind::Immediate);
290 }
291 Ok(must_migrate)
292 }
293}