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.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        // 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 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        // An unpublished task can have its affinity changed through a published
171        // OS identity. Reserve the replacement target before committing the mask.
172        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        // The affinity mask is task metadata, but physical placement belongs
192        // to one runqueue owner. A remote writer only publishes a reconciliation
193        // request; it never rewrites Queued/Running or the independent
194        // switch-tail `on_cpu` publication in place.
195        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    /// Updates the owner CPU's running thread without publishing a self inbox.
222    ///
223    /// The caller owns `cpu` in an IRQ-off scheduler-safe window. A `true`
224    /// result means the current thread must schedule out before the operation
225    /// can return to its caller; switch tail will publish the detached context
226    /// to the selected destination CPU.
227    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(&current_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}