Skip to main content

strat9_kernel/process/scheduler/
core_impl.rs

1use super::{runtime_ops::idle_task_main, *};
2
3/// Create a `SchedulerCpu` for the given CPU index (creates its idle task).
4pub(super) fn create_cpu_scheduler(cpu_idx: usize) -> SchedulerCpu {
5    crate::serial_println!(
6        "[trace][sched] create_cpu_scheduler cpu={} create idle begin",
7        cpu_idx
8    );
9    let idle_task = Task::new_kernel_task(idle_task_main, "idle", TaskPriority::Idle)
10        .expect("Failed to create idle task");
11    crate::serial_println!(
12        "[trace][sched] create_cpu_scheduler cpu={} create idle done id={}",
13        cpu_idx,
14        idle_task.id.as_u64()
15    );
16    idle_task.set_sched_policy(crate::process::sched::SchedPolicy::Idle);
17    let mut class_rqs = PerCpuClassRqSet::new();
18    class_rqs.enqueue(crate::process::sched::SchedClassId::Idle, idle_task.clone());
19    SchedulerCpu {
20        class_rqs,
21        current_task: None,
22        current_runtime: crate::process::sched::CurrentRuntime::new(),
23        idle_task,
24        task_to_requeue: None,
25        task_to_drop: None,
26        need_resched: false,
27        class_table: crate::process::sched::SchedClassTable::default(),
28    }
29}
30
31impl GlobalSchedState {
32    /// Create a new global scheduler state (no per-CPU runqueues : those live in LOCAL_SCHEDULERS).
33    pub fn new() -> Self {
34        crate::serial_println!("[trace][sched] GlobalSchedState::new enter");
35        GlobalSchedState {
36            all_tasks: BTreeMap::new(),
37            task_cpu: BTreeMap::new(),
38            wake_deadlines: BTreeMap::new(),
39            wake_deadline_of: BTreeMap::new(),
40            zombies: BTreeMap::new(),
41            class_table: crate::process::sched::SchedClassTable::default(),
42        }
43    }
44
45    /// Performs the member add operation.
46    pub(crate) fn member_add(
47        map: &mut BTreeMap<Pid, alloc::vec::Vec<TaskId>>,
48        key: Pid,
49        task_id: TaskId,
50    ) {
51        let members = map.entry(key).or_default();
52        if !members.iter().any(|id| *id == task_id) {
53            members.push(task_id);
54        }
55    }
56
57    /// Performs the member remove operation.
58    pub(crate) fn member_remove(
59        map: &mut BTreeMap<Pid, alloc::vec::Vec<TaskId>>,
60        key: Pid,
61        task_id: TaskId,
62    ) {
63        let mut clear = false;
64        if let Some(members) = map.get_mut(&key) {
65            members.retain(|id| *id != task_id);
66            clear = members.is_empty();
67        }
68        if clear {
69            map.remove(&key);
70        }
71    }
72
73    /// Performs the register identity locked operation.
74    pub(crate) fn register_identity_locked(identity: &mut SchedIdentity, task: &Arc<Task>) {
75        let task_id = task.id;
76        let pid = task.pid;
77        let pgid = task.pgid.load(Ordering::Relaxed);
78        let sid = task.sid.load(Ordering::Relaxed);
79        crate::serial_println!(
80            "[trace][sched] register_identity enter tid={} pid={} pgid={} sid={}",
81            task_id.as_u64(),
82            pid,
83            pgid,
84            sid
85        );
86        identity.pid_to_pgid.insert(pid, pgid);
87        crate::serial_println!(
88            "[trace][sched] register_identity pid_to_pgid inserted pid={}",
89            pid
90        );
91        identity.pid_to_sid.insert(pid, sid);
92        crate::serial_println!(
93            "[trace][sched] register_identity pid_to_sid inserted pid={}",
94            pid
95        );
96        Self::member_add(&mut identity.pgid_members, pgid, task_id);
97        Self::member_add(&mut identity.sid_members, sid, task_id);
98        crate::serial_println!(
99            "[trace][sched] register_identity done tid={}",
100            task_id.as_u64()
101        );
102    }
103
104    /// Performs the unregister identity locked operation.
105    pub(crate) fn unregister_identity_locked(
106        identity: &mut SchedIdentity,
107        task_id: TaskId,
108        pid: Pid,
109        tid: Tid,
110    ) {
111        identity.pid_to_task.remove(&pid);
112        identity.tid_to_task.remove(&tid);
113        if let Some(pgid) = identity.pid_to_pgid.remove(&pid) {
114            Self::member_remove(&mut identity.pgid_members, pgid, task_id);
115        }
116        if let Some(sid) = identity.pid_to_sid.remove(&pid) {
117            Self::member_remove(&mut identity.sid_members, sid, task_id);
118        }
119    }
120
121    /// Add a task to the scheduler
122    pub fn add_task(&mut self, task: Arc<Task>) -> Option<usize> {
123        let cpu_index = self.select_cpu_for_task(&task);
124        self.add_task_on_cpu(task, cpu_index)
125    }
126
127    /// Performs the add task with parent operation.
128    pub fn add_task_with_parent(&mut self, task: Arc<Task>, parent: TaskId) -> Option<usize> {
129        let child = task.id;
130        let cpu_index = self.select_cpu_for_task(&task);
131        let ipi = self.add_task_on_cpu(task, cpu_index);
132        {
133            let mut identity = SCHED_IDENTITY.write();
134            identity.parent_of.insert(child, parent);
135            identity.children_of.entry(parent).or_default().push(child);
136        }
137        ipi
138    }
139
140    /// Performs the add task on cpu operation.
141    fn add_task_on_cpu(&mut self, task: Arc<Task>, cpu_index: usize) -> Option<usize> {
142        let task_id = task.id;
143        crate::serial_println!(
144            "[trace][sched] add_task_on_cpu enter tid={} cpu={}",
145            task_id.as_u64(),
146            cpu_index
147        );
148        task.set_state(TaskState::Ready);
149
150        self.insert_all_task_locked(task_id, task.clone());
151        self.task_cpu.insert(task_id, cpu_index);
152        task.home_cpu
153            .store(cpu_index, core::sync::atomic::Ordering::Relaxed);
154
155        // Lock order: IDENTITY write (rank 2). Caller holds GLOBAL (rank 1).
156        lockdep_acquire(LockRank::IdentityW, None);
157        {
158            let mut identity = SCHED_IDENTITY.write();
159            identity.pid_to_task.insert(task.pid, task_id);
160            identity.tid_to_task.insert(task.tid, task_id);
161            Self::register_identity_locked(&mut identity, &task);
162        }
163        lockdep_release(LockRank::IdentityW);
164
165        // Lock order: LOCAL (rank 4). Caller holds GLOBAL (rank 1).
166        lockdep_acquire(LockRank::Local, Some(cpu_index));
167        {
168            let class = self.class_table.class_for_task(&task);
169            if let Some(ref mut local_cpu) = *LOCAL_SCHEDULERS[cpu_index].lock() {
170                local_cpu.class_rqs.enqueue(class, task);
171                local_cpu.need_resched = true;
172            }
173        }
174        lockdep_release(LockRank::Local);
175
176        sched_trace(format_args!(
177            "enqueue task={} cpu={}",
178            task_id.as_u64(),
179            cpu_index
180        ));
181        if cpu_index != current_cpu_index() {
182            Some(cpu_index)
183        } else {
184            None
185        }
186    }
187
188    pub(super) fn insert_all_task_locked(&mut self, task_id: TaskId, task: Arc<Task>) {
189        assert_eq!(
190            task.id,
191            task_id,
192            "scheduler corruption: insert_all_task_locked task.id={} != task_id={}",
193            task.id.as_u64(),
194            task_id.as_u64()
195        );
196        if self.all_tasks.contains_key(&task_id) {
197            unsafe {
198                crate::arch::serial::putc(b'D');
199            }
200            crate::serial_force_println!(
201                "[RACE] insert_all_task_locked: duplicate tid={} all_tasks={}",
202                task_id.as_u64(),
203                self.all_tasks.len(),
204            );
205            panic!(
206                "scheduler corruption: duplicate insert_all_task_locked tid={}",
207                task_id.as_u64()
208            );
209        }
210        self.all_tasks.insert(task_id, task);
211    }
212
213    pub(super) fn remove_all_task_locked(&mut self, task_id: TaskId) -> Option<Arc<Task>> {
214        self.all_tasks.remove(&task_id)
215    }
216
217    /// Performs the clear task wake deadline locked operation.
218    pub fn clear_task_wake_deadline_locked(&mut self, id: TaskId) -> bool {
219        if let Some(task) = self.all_tasks.get(&id) {
220            task.wake_deadline_ns.store(0, Ordering::Relaxed);
221            true
222        } else {
223            false
224        }
225    }
226
227    /// Sets task wake deadline locked.
228    pub fn set_task_wake_deadline_locked(&mut self, id: TaskId, deadline: u64) -> bool {
229        if deadline == 0 {
230            return self.clear_task_wake_deadline_locked(id);
231        }
232        if let Some(task) = self.all_tasks.get(&id) {
233            task.wake_deadline_ns.store(deadline, Ordering::Relaxed);
234            true
235        } else {
236            false
237        }
238    }
239
240    /// Performs the wake task locked operation.
241    ///
242    /// Returns `(was_woken, ipi_cpu)`. The caller must send a resched IPI to
243    /// `ipi_cpu` after releasing the scheduler lock.
244    ///
245    /// NOTE: `blocked_tasks` lives in the separate `BLOCKED_TASKS` lock.
246    /// This method handles waking the task directly if it is already in `BLOCKED_TASKS`,
247    /// or setting `wake_pending` as a fallback if the task is still transitioning to Blocked.
248    pub fn wake_task_locked(&mut self, id: TaskId) -> (bool, Option<usize>) {
249        self.clear_task_wake_deadline_locked(id);
250
251        // Lock order: BLOCKED (rank 3) → LOCAL (rank 4). Caller holds GLOBAL (rank 1).
252        let mut woken = false;
253        let mut ipi_cpu = None;
254        {
255            lockdep_acquire(LockRank::Blocked, None);
256            let mut blocked = super::BLOCKED_TASKS.lock();
257            if let Some(task) = blocked.remove(&id) {
258                task.set_state(TaskState::Ready);
259
260                // --- CPU placement: prefer last_cpu (cache warmth), then home_cpu ---
261                let last = task.last_cpu.load(Ordering::Relaxed);
262                let home = task.home_cpu.load(Ordering::Relaxed);
263                let n = active_cpu_count();
264
265                // Try last_cpu first: valid, online, and reasonably loaded.
266                let cpu_index = if last < n {
267                    let last_ok = {
268                        lockdep_acquire(LockRank::Local, Some(last));
269                        let ok = LOCAL_SCHEDULERS[last]
270                            .lock()
271                            .as_ref()
272                            .map(|c| c.class_rqs.runnable_len() <= 2)
273                            .unwrap_or(false);
274                        lockdep_release(LockRank::Local);
275                        ok
276                    };
277                    if last_ok {
278                        last
279                    } else if home < n {
280                        home
281                    } else {
282                        0
283                    }
284                } else if home < n {
285                    home
286                } else {
287                    0
288                };
289
290                let class = {
291                    use crate::process::sched::SchedClassId;
292                    match task.sched_policy() {
293                        crate::process::sched::SchedPolicy::RealTimeRR { .. }
294                        | crate::process::sched::SchedPolicy::RealTimeFifo { .. } => {
295                            SchedClassId::RealTime
296                        }
297                        crate::process::sched::SchedPolicy::Fair(_) => SchedClassId::Fair,
298                        crate::process::sched::SchedPolicy::Idle => SchedClassId::Idle,
299                    }
300                };
301
302                lockdep_acquire(LockRank::Local, Some(cpu_index));
303                if let Some(ref mut local_cpu) = *super::LOCAL_SCHEDULERS[cpu_index].lock() {
304                    local_cpu.class_rqs.enqueue(class, task.clone());
305                    local_cpu.need_resched = true;
306                }
307                lockdep_release(LockRank::Local);
308
309                ipi_cpu = if cpu_index != current_cpu_index() {
310                    Some(cpu_index)
311                } else {
312                    None
313                };
314                woken = true;
315            }
316            lockdep_release(LockRank::Blocked);
317            drop(blocked);
318        }
319
320        if woken {
321            return (true, ipi_cpu);
322        }
323
324        // Fallback: task not yet in BLOCKED_TASKS (still transitioning to
325        // Blocked). Set wake_pending so block_current_task skips blocking.
326        if let Some(task) = self.all_tasks.get(&id) {
327            task.wake_pending
328                .store(true, core::sync::atomic::Ordering::Release);
329            (true, None)
330        } else {
331            (false, None)
332        }
333    }
334
335    /// Attempts to reap child locked.
336    /// If `target` is `Some(tid)`, only reaps that child; otherwise, reaps any child.
337    /// Returns `WaitChildResult` indicating the outcome.
338    /// Must be called with the scheduler lock held.
339    ///
340    pub fn try_reap_child_locked(
341        &mut self,
342        parent: TaskId,
343        target: Option<TaskId>,
344    ) -> WaitChildResult {
345        // First, check children under SCHED_IDENTITY lock.
346        let target_is_child = {
347            let identity = SCHED_IDENTITY.read();
348            let Some(children_view) = identity.children_of.get(&parent) else {
349                return WaitChildResult::NoChildren;
350            };
351
352            if children_view.is_empty() {
353                return WaitChildResult::NoChildren;
354            }
355
356            let target_is_child = if let Some(target_id) = target {
357                children_view.iter().any(|&id| id == target_id)
358            } else {
359                true
360            };
361            target_is_child
362        };
363
364        if !target_is_child {
365            return WaitChildResult::NoChildren;
366        }
367
368        // Find the zombie child : re-check children under SCHED_IDENTITY.
369        let zombie = {
370            let identity = SCHED_IDENTITY.read();
371            let children = match identity.children_of.get(&parent) {
372                Some(c) => c.clone(),
373                None => return WaitChildResult::NoChildren,
374            };
375            children
376                .iter()
377                .copied()
378                .find(|id| target.map_or(true, |t| t == *id) && self.zombies.contains_key(id))
379        };
380
381        if let Some(child) = zombie {
382            let (status, child_pid) = self.zombies.remove(&child).unwrap_or((0, 0));
383            // Remove from all_tasks now so that pick_next_task (if it races with
384            // reaping) will see was_registered=false and skip cleanup_task_resources.
385            let reaped_task = self.remove_all_task_locked(child);
386            if let Some(task) = reaped_task.as_ref() {
387                super::task_ops::cleanup_task_resources(task);
388            }
389            let child_tid = reaped_task.as_ref().map(|t| t.tid);
390            if child_pid != 0 {
391                if let Some(tid) = child_tid {
392                    let mut identity = SCHED_IDENTITY.write();
393                    Self::unregister_identity_locked(&mut identity, child, child_pid, tid);
394                }
395            }
396            {
397                let mut identity = SCHED_IDENTITY.write();
398                if let Some(children) = identity.children_of.get_mut(&parent) {
399                    children.retain(|&id| id != child);
400                    if children.is_empty() {
401                        identity.children_of.remove(&parent);
402                    }
403                }
404                identity.parent_of.remove(&child);
405            }
406            return WaitChildResult::Reaped {
407                child,
408                pid: child_pid,
409                status,
410            };
411        }
412
413        WaitChildResult::StillRunning
414    }
415
416    /// Select the least-loaded CPU for a newly created task.
417    ///
418    /// Called with `GLOBAL_SCHED_STATE` held.  Acquires each LOCAL sequentially.
419    /// Lock order: GLOBAL (held) → LOCAL[cpu] (rank 4).  Never two LOCALs simultaneously.
420    ///
421    /// Placement strategy:
422    /// 1. Early boot (all idle): CPU 0.
423    /// 2. Soft affinity: prefer CPUs in the task's `affinity_mask` (from silo config).
424    ///    Tie-break with `last_cpu` (cache warmth).  If no affinity-eligible CPU
425    ///    has strictly lower load, fall back to unrestricted search.
426    /// 3. Unrestricted: pick the globally least-loaded CPU.
427    fn select_cpu_for_task(&self, task: &Task) -> usize {
428        let n = active_cpu_count();
429        let all_idle = (0..n).all(|i| {
430            lockdep_acquire(LockRank::Local, Some(i));
431            let result = LOCAL_SCHEDULERS[i]
432                .lock()
433                .as_ref()
434                .map(|cpu| cpu.current_task.is_none())
435                .unwrap_or(true);
436            lockdep_release(LockRank::Local);
437            result
438        });
439        if all_idle {
440            crate::serial_println!("[trace][sched] select_cpu_for_task early-boot best=0");
441            return 0;
442        }
443
444        let affinity = task.affinity_mask.load(Ordering::Relaxed);
445        let last = task.last_cpu.load(Ordering::Relaxed);
446        let has_affinity = affinity != 0;
447
448        // Pass 1: find least-loaded CPU among affinity-eligible CPUs.
449        let mut best_aff = None;
450        let mut best_aff_load = usize::MAX;
451        let mut best_aff_last = None; // best among those matching last_cpu
452        let mut best_aff_last_load = usize::MAX;
453
454        for idx in 0..n {
455            let eligible = if has_affinity {
456                (affinity & (1u64 << idx)) != 0
457            } else {
458                true
459            };
460            if !eligible {
461                continue;
462            }
463            lockdep_acquire(LockRank::Local, Some(idx));
464            let load = {
465                let guard = LOCAL_SCHEDULERS[idx].lock();
466                if let Some(ref cpu) = *guard {
467                    let mut l = cpu.class_rqs.runnable_len();
468                    if let Some(current) = cpu.current_task.as_ref() {
469                        if self.class_table.class_for_task(current)
470                            != crate::process::sched::SchedClassId::Idle
471                        {
472                            l += 1;
473                        }
474                    }
475                    l
476                } else {
477                    0
478                }
479            };
480            lockdep_release(LockRank::Local);
481
482            if load < best_aff_load {
483                best_aff = Some(idx);
484                best_aff_load = load;
485            }
486            // Prefer last_cpu when loads are equal (cache warmth).
487            if idx == last && load <= best_aff_last_load {
488                best_aff_last = Some(idx);
489                best_aff_last_load = load;
490            }
491        }
492
493        // If we have an affinity-eligible candidate, prefer it.
494        // Use last_cpu tie-breaker when loads are equal.
495        if let Some(aff_cpu) = best_aff {
496            // If last_cpu is also eligible and has equal or lower load, prefer it.
497            if let Some(lc) = best_aff_last {
498                if best_aff_last_load <= best_aff_load {
499                    crate::serial_println!(
500                        "[trace][sched] select_cpu_for_task affinity+last_cpu cpu={} load={}",
501                        lc,
502                        best_aff_last_load
503                    );
504                    return lc;
505                }
506            }
507            crate::serial_println!(
508                "[trace][sched] select_cpu_for_task affinity cpu={} load={}",
509                aff_cpu,
510                best_aff_load
511            );
512            return aff_cpu;
513        }
514
515        // Pass 2 (fallback): unrestricted least-loaded.
516        let mut best = 0usize;
517        let mut best_load = usize::MAX;
518        for idx in 0..n {
519            lockdep_acquire(LockRank::Local, Some(idx));
520            let load = {
521                let guard = LOCAL_SCHEDULERS[idx].lock();
522                if let Some(ref cpu) = *guard {
523                    let mut l = cpu.class_rqs.runnable_len();
524                    if let Some(current) = cpu.current_task.as_ref() {
525                        if self.class_table.class_for_task(current)
526                            != crate::process::sched::SchedClassId::Idle
527                        {
528                            l += 1;
529                        }
530                    }
531                    l
532                } else {
533                    0
534                }
535            };
536            lockdep_release(LockRank::Local);
537            if load < best_load {
538                best = idx;
539                best_load = load;
540            }
541        }
542        crate::serial_println!(
543            "[trace][sched] select_cpu_for_task fallback best={} load={}",
544            best,
545            best_load
546        );
547        best
548    }
549
550    /// Migrate all ready tasks to match the current class table.
551    ///
552    /// Called with `GLOBAL_SCHED_STATE` held.  Acquires each `LOCAL_SCHEDULERS[cpu]`
553    /// sequentially (never two at once).  Lock order: GLOBAL (held) → LOCAL[cpu] (rank 4).
554    pub fn migrate_ready_tasks_for_new_class_table(&mut self) {
555        let mut ready: Vec<(TaskId, Arc<Task>, usize)> = Vec::new();
556        for (id, task) in self.all_tasks.iter() {
557            let state = task.get_state();
558            if state != TaskState::Ready {
559                continue;
560            }
561            let cpu = self.task_cpu.get(id).copied().unwrap_or(0);
562            ready.push((*id, task.clone(), cpu));
563        }
564
565        for (id, task, cpu_idx) in ready {
566            lockdep_acquire(LockRank::Local, Some(cpu_idx));
567            let mut guard = LOCAL_SCHEDULERS[cpu_idx].lock();
568            let Some(ref mut cpu) = *guard else {
569                lockdep_release(LockRank::Local);
570                continue;
571            };
572            if cpu.class_rqs.remove(id) {
573                let class = self.class_table.class_for_task(&task);
574                cpu.class_rqs.enqueue(class, task);
575                cpu.need_resched = true;
576            }
577            lockdep_release(LockRank::Local);
578        }
579    }
580}
581
582//  Per-CPU hot-path helpers
583//
584// These functions operate primarily on `SchedulerCpu` (acquired via
585// `LOCAL_SCHEDULERS[cpu_index]`).  Most never touch the global `SCHEDULER`
586// lock.  The one exception is `steal_task_local`, which does a **non-blocking**
587// `GLOBAL_SCHED_STATE.try_lock_no_irqsave()` to update `task_cpu` after a successful
588// steal.  This is an intentional lock-order inversion (LOCAL held, then GLOBAL
589// attempted) that is safe because the try-lock never blocks : if GLOBAL_SCHED_STATE is
590// contended, we simply skip stealing.
591//
592// Lock order for steal: own LOCAL held => try_lock GLOBAL => try_lock sibling
593// LOCALs.  Never blocking-wait, so no deadlock possible.
594
595/// Steal a task from the busiest sibling CPU using per-CPU LOCAL locks.
596///
597/// Called with `cpu` borrowed from `LOCAL_SCHEDULERS[cpu_index]` (our own
598/// LOCAL lock already held). Uses `try_lock_no_irqsave` on sibling entries :
599/// if a sibling or the global scheduler state is contended, we skip stealing
600/// rather than waiting.
601pub(super) fn steal_task_local(cpu: &mut SchedulerCpu, cpu_index: usize) -> Option<Arc<Task>> {
602    let now_tick = TICK_COUNT.load(Ordering::Relaxed);
603    if now_tick < LAST_STEAL_TICK[cpu_index].load(Ordering::Relaxed) + STEAL_COOLDOWN_TICKS {
604        return None;
605    }
606    lockdep_assert_held(LockRank::Local);
607
608    // Best-effort only: if a cold path is holding the global scheduler, skip
609    // stealing instead of blocking the hot path.
610    let mut scheduler = match GLOBAL_SCHED_STATE.try_lock_no_irqsave() {
611        Some(g) => g,
612        None => return None,
613    };
614    lockdep_acquire(LockRank::Global, None);
615    let result = steal_task_inner(&mut scheduler, cpu, cpu_index, now_tick);
616    lockdep_release(LockRank::Global);
617    result
618}
619
620/// Inner implementation of steal, separated for clean lockdep release.
621fn steal_task_inner(
622    scheduler: &mut Option<GlobalSchedState>,
623    cpu: &mut SchedulerCpu,
624    cpu_index: usize,
625    now_tick: u64,
626) -> Option<Arc<Task>> {
627    let sched = scheduler.as_mut()?;
628
629    let n = active_cpu_count();
630    let my_load = cpu.class_rqs.runnable_len();
631
632    let mut best_cpu = None;
633    let mut best_load = 0usize;
634
635    for i in 0..n {
636        if i == cpu_index {
637            continue;
638        }
639        // try_lock_no_irqsave: returns immediately if contended (no deadlock).
640        if let Some(guard) = LOCAL_SCHEDULERS[i].try_lock_no_irqsave() {
641            if let Some(ref sib) = *guard {
642                let load = sib.class_rqs.runnable_len();
643                if load > best_load {
644                    best_load = load;
645                    best_cpu = Some(i);
646                }
647            }
648        }
649    }
650
651    if best_load < my_load.saturating_add(STEAL_IMBALANCE_MIN) {
652        return None;
653    }
654    let steal_from = best_cpu?;
655
656    // Re-acquire the sibling lock to perform the steal.
657    if let Some(mut guard) = LOCAL_SCHEDULERS[steal_from].try_lock_no_irqsave() {
658        if let Some(ref mut sib) = *guard {
659            if sib.class_rqs.runnable_len() < 2 {
660                return None;
661            }
662            if let Some(task) = sib.class_rqs.steal_candidate(&sib.class_table) {
663                // Soft affinity check: skip if the stealing CPU is not in the
664                // task's affinity mask.  Re-enqueue on the source to avoid
665                // dropping the task.
666                let affinity = task.affinity_mask.load(Ordering::Relaxed);
667                if affinity != 0 && (affinity & (1u64 << cpu_index)) == 0 {
668                    // Re-enqueue on source CPU.
669                    let class = sib.class_table.class_for_task(&task);
670                    sib.class_rqs.enqueue(class, task);
671                    return None;
672                }
673                sched.task_cpu.insert(task.id, cpu_index);
674                task.home_cpu
675                    .store(cpu_index, core::sync::atomic::Ordering::Relaxed);
676                if cpu_is_valid(cpu_index) {
677                    CPU_STEAL_IN_COUNT[cpu_index].fetch_add(1, Ordering::Relaxed);
678                }
679                if cpu_is_valid(steal_from) {
680                    CPU_STEAL_OUT_COUNT[steal_from].fetch_add(1, Ordering::Relaxed);
681                }
682                LAST_STEAL_TICK[cpu_index].store(now_tick, Ordering::Relaxed);
683                return Some(task);
684            }
685        }
686    }
687    None
688}
689
690/// Pick the next task using only per-CPU LOCAL state.
691///
692/// Handles current task disposition (re-queue, drop-for-cleanup, or ignore if
693/// Blocked), then picks from the local class_rqs, falls back to work-stealing,
694/// and finally returns the idle task.
695///
696/// **Dead tasks**: if the current task is Dead, it goes into `task_to_drop`.
697/// Global map cleanup (`all_tasks`, `task_cpu`, etc.) must have been performed
698/// by the caller (e.g., `exit_current_task`) BEFORE reaching this point.
699pub(super) fn pick_next_task_local(cpu: &mut SchedulerCpu, cpu_index: usize) -> Arc<Task> {
700    // Step 1: dispose of the current task.
701    if let Some(task) = cpu.current_task.take() {
702        match task.get_state() {
703            TaskState::Running => {
704                task.set_state(TaskState::Ready);
705                if !Arc::ptr_eq(&task, &cpu.idle_task) {
706                    // Defer re-queue to finish_switch (not yet safe to enqueue :
707                    // another CPU could steal it before our context is saved).
708                    cpu.task_to_requeue = Some(task);
709                }
710            }
711            TaskState::Dead => {
712                // Global maps already cleaned by exit_current_task / kill_task.
713                // Defer the Arc drop so KernelStack::drop => buddy_alloc runs
714                // outside any lock.
715                cpu.task_to_drop = Some(task);
716            }
717            TaskState::Blocked | TaskState::Ready => {
718                // Blocked: moved to blocked_tasks by block_current_task : do nothing.
719                // Ready: shouldn't normally occur for current_task; safe to ignore.
720            }
721        }
722    }
723
724    // Step 2: pick from local class_rqs.
725    unsafe {
726        core::arch::asm!("out 0xe9, al", in("al") b'O', options(nomem, nostack));
727    }
728    let next = if let Some(next) = cpu.class_rqs.pick_next(&cpu.class_table) {
729        unsafe {
730            core::arch::asm!("out 0xe9, al", in("al") b'J', options(nomem, nostack));
731        }
732        next
733    } else {
734        // Step 3: local queues empty — try work-stealing before idle.
735        unsafe {
736            core::arch::asm!("out 0xe9, al", in("al") b'S', options(nomem, nostack));
737        }
738        if let Some(stolen) = steal_task_local(cpu, cpu_index) {
739            unsafe {
740                core::arch::asm!("out 0xe9, al", in("al") b's', options(nomem, nostack));
741            }
742            stolen
743        } else {
744            unsafe {
745                core::arch::asm!("out 0xe9, al", in("al") b'j', options(nomem, nostack));
746            }
747            // Step 4: idle fallback.
748            cpu.idle_task.clone()
749        }
750    };
751
752    next.set_state(TaskState::Running);
753    next.last_cpu.store(cpu_index, Ordering::Relaxed);
754    cpu.current_task = Some(next.clone());
755    cpu.current_runtime = crate::process::sched::CurrentRuntime::new();
756    next
757}
758
759/// Prepare a LOCAL-only context switch.
760///
761/// Updates the TSS, SYSCALL RSP, and CR3 for the next task. Returns the
762/// raw pointer pair needed by `do_switch_context`.
763///
764/// Returns `None` if there is no task to switch to (same task or invalid context).
765pub(super) fn yield_cpu_local(cpu: &mut SchedulerCpu, cpu_index: usize) -> Option<SwitchTarget> {
766    let current = cpu.current_task.as_ref()?.clone();
767
768    let next = pick_next_task_local(cpu, cpu_index);
769
770    if Arc::ptr_eq(&current, &next) {
771        return None;
772    }
773    if cpu_is_valid(cpu_index) {
774        CPU_SWITCH_COUNT[cpu_index].fetch_add(1, Ordering::Relaxed);
775    }
776
777    if let Err(e) = validate_task_context(&next) {
778        let bad_rsp = unsafe { (*next.context.get()).saved_rsp };
779        let stk_base = next.kernel_stack.virt_base.as_u64();
780        let stk_top = stk_base + next.kernel_stack.size as u64;
781        crate::serial_println!(
782            "[sched-local] WARN: invalid ctx task='{}' id={} cpu={}: {} \
783             rsp={:#x} stack=[{:#x}..{:#x}] : restoring current",
784            next.name,
785            next.id.as_u64(),
786            cpu_index,
787            e,
788            bad_rsp,
789            stk_base,
790            stk_top,
791        );
792
793        // Restore invariants: undo what pick_next_task_local mutated.
794        let is_idle = Arc::ptr_eq(&next, &cpu.idle_task);
795        drop(cpu.task_to_drop.take());
796        if let Some(prev) = cpu.task_to_requeue.take() {
797            prev.set_state(TaskState::Running);
798            cpu.current_task = Some(prev);
799        } else {
800            current.set_state(TaskState::Running);
801            cpu.current_task = Some(current.clone());
802        }
803        if !is_idle {
804            next.set_state(TaskState::Ready);
805            let class = cpu.class_table.class_for_task(&next);
806            cpu.class_rqs.enqueue(class, next);
807        }
808        return None;
809    }
810
811    // Update TSS.rsp0 and SYSCALL kernel RSP for the new task.
812    let stack_top = next.kernel_stack.virt_base.as_u64() + next.kernel_stack.size as u64;
813    crate::arch::tss::set_kernel_stack(crate::arch::xshim::VirtAddr::new(stack_top));
814    crate::arch::syscall::set_kernel_rsp(stack_top);
815
816    // Switch CR3 if the new task has a different address space.
817    // SAFETY: The new task's address space has a valid PML4 with the kernel half mapped.
818    unsafe {
819        next.process.address_space_arc().switch_to();
820    }
821
822    Some(SwitchTarget {
823        old_rsp_ptr: unsafe { &raw mut (*current.context.get()).saved_rsp },
824        new_rsp_ptr: unsafe { &raw const (*next.context.get()).saved_rsp },
825        old_fpu_ptr: current.fpu_state.get() as *mut u8,
826        new_fpu_ptr: next.fpu_state.get() as *const u8,
827        old_xcr0: current
828            .xcr0_mask
829            .load(core::sync::atomic::Ordering::Relaxed),
830        new_xcr0: next.xcr0_mask.load(core::sync::atomic::Ordering::Relaxed),
831    })
832}
833
834/// Post-switch cleanup using only LOCAL state: re-enqueue the previous task
835/// and optionally extract the task-to-drop for deferred deallocation.
836pub(super) fn drain_post_switch_local(
837    cpu: &mut SchedulerCpu,
838    take_drop: bool,
839) -> Option<Arc<Task>> {
840    let task_to_drop = if take_drop {
841        cpu.task_to_drop.take()
842    } else {
843        None
844    };
845    if let Some(task) = cpu.task_to_requeue.take() {
846        let class = cpu.class_table.class_for_task(&task);
847        cpu.class_rqs.enqueue(class, task);
848    }
849    task_to_drop
850}