Skip to main content

ax_task/sched/system/task_system/
placement.rs

1//! Affinity updates and owner-to-owner placement delivery.
2
3use super::*;
4
5impl TaskSystem {
6    /// Reads the unique scheduler-class state while holding `p->pi_lock`.
7    ///
8    /// A detached or migrating task owns the state in its task-stable slot.
9    /// A queued or running task owns it in exactly one rq. This mirrors
10    /// Linux's `task_rq_lock()` rule instead of copying current state back into
11    /// the task merely to answer an affinity request.
12    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        // SAFETY: this count is transferred to the dedicated refresh node and
76        // consumed by exactly one owner-side inbox drain.
77        unsafe { Arc::increment_strong_count(pointer) };
78        // SAFETY: the transferred Arc count pins the embedded refresh node.
79        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            // SAFETY: rejected/coalesced publication retained no extra count.
88            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        // SAFETY: this count is transferred to the dedicated affinity node and
107        // consumed by exactly one later owner drain.
108        unsafe { Arc::increment_strong_count(pointer) };
109        // SAFETY: the transferred Arc count pins the embedded control node.
110        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            // SAFETY: rejected/coalesced publication did not consume this
119            // attempt's retained reference.
120            unsafe { Arc::decrement_strong_count(pointer) };
121            core.cancel_scheduler_inbox_delivery();
122        }
123        Ok(())
124    }
125
126    /// Publishes one affinity generation and returns its completion owner.
127    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        // The affinity mask is task metadata, but physical placement belongs
178        // to one runqueue owner. A remote writer only publishes a reconciliation
179        // request; it never rewrites Queued/Running or the independent
180        // switch-tail `on_cpu` publication in place.
181        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    /// Updates the owner CPU's running thread without publishing a self inbox.
200    ///
201    /// The caller owns `cpu` in an IRQ-off scheduler-safe window. A `true`
202    /// result means the current thread must schedule out before the operation
203    /// can return to its caller; switch tail will publish the detached context
204    /// to the selected destination CPU.
205    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(&current_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}