Skip to main content

ax_task/sched/system/task_system/delivery/
admission.rs

1//! Admission under the owning scheduler transaction.
2
3use super::*;
4
5impl TaskSystem {
6    /// Enqueues a ready thread on an affinity-compatible owner CPU.
7    pub fn enqueue(&self, mut cpu: Pin<&mut CpuLocal>, thread: ThreadId) -> Result<(), TaskError> {
8        self.ensure_owner_cpu_context(&cpu)?;
9        let core = {
10            let state = self.state.lock();
11            state.ensure_cpu_online(&cpu)?;
12            Arc::clone(&state.thread_record(thread)?.core)
13        };
14        self.enqueue_owner_thread(cpu.as_mut(), core, EnqueueReason::Wake)?;
15        self.program_local_timer(cpu.as_mut(), SchedulerDeadlineDerivationSource::Enqueue)
16    }
17
18    /// Admits a new thread and commits its placement on an allowed active CPU.
19    ///
20    /// Rejected admission does not change lifecycle or placement. Success
21    /// guarantees either local
22    /// runqueue admission or an owned remote activation delivery. There is no
23    /// public state-only runnable transition to complete in a second call.
24    ///
25    /// Ordinary fair work is placed on the least-loaded allowed CPU, including
26    /// its current non-idle dispatch and migrations not yet consumed by the
27    /// destination owner. Other classes preserve owner-local placement unless
28    /// affinity requires a transfer. Remote placement uses the owner-only
29    /// owner-control inbox and never mutates another CPU's runqueue.
30    ///
31    /// # Errors
32    ///
33    /// Returns an error when the source CPU is offline, the thread is not a
34    /// new unqueued thread, no allowed CPU is online, or remote delivery
35    /// cannot be reserved. Failures after admission are runtime invariants.
36    pub fn start_thread(
37        &self,
38        mut cpu: Pin<&mut CpuLocal>,
39        thread: ThreadId,
40    ) -> Result<(), TaskError> {
41        self.ensure_owner_cpu_context(&cpu)?;
42        let owner = cpu.owner();
43        let migration = {
44            let state = self.state.lock();
45            state.ensure_cpu_online(&cpu)?;
46            let record = state.thread_record(thread)?;
47            let mut sched = record.sched.lock();
48            if sched.lifecycle.state() != ThreadState::New {
49                return Err(TaskError::NotReady);
50            }
51            if sched.placement.queued_cpu().is_some()
52                || sched.placement.on_cpu().is_some()
53                || sched.placement.has_pending_migration()
54            {
55                return Err(TaskError::AlreadyQueued);
56            }
57            let affinity = &sched.affinity.affinity;
58            let active = record.core.sched().active(&sched);
59            let policy = active.policy();
60            let load_aware = matches!(policy, SchedulePolicy::Fair { .. });
61            let target = if load_aware {
62                state.select_initial_fair_cpu(affinity, Some(owner))
63            } else if matches!(
64                policy,
65                SchedulePolicy::Fifo { .. }
66                    | SchedulePolicy::RoundRobin { .. }
67                    | SchedulePolicy::Deadline(_)
68            ) {
69                self.select_priority_cpu(policy, Some(active.entity()), affinity, Some(owner), None)
70            } else if affinity.contains(owner) {
71                Some(owner)
72            } else {
73                self.select_fallback_active_cpu(affinity, None)
74            }
75            .ok_or(TaskError::InvalidConfiguration)?;
76            drop(active);
77            let core = Arc::clone(&record.core);
78            if target == owner {
79                drop(sched);
80                drop(state);
81                self.start_owner_thread(cpu.as_mut(), core)?;
82                None
83            } else {
84                let carrier = self.prepare_owner_migration(&core, owner, target)?;
85                sched.transition(&core, ThreadState::Running)?;
86                sched.placement.begin_remote_wakeup(target);
87                record.core.set_wake_cpu_hint(target);
88                drop(sched);
89                Some((carrier, target))
90            }
91        };
92        if let Some((carrier, _target)) = migration {
93            carrier.commit();
94            return Ok(());
95        }
96        self.program_local_timer(cpu.as_mut(), SchedulerDeadlineDerivationSource::Placement)
97            .unwrap_or_else(|_| {
98                task_runtime::fatal_invariant(0x5354_0001, thread.as_u64() as usize)
99            });
100        Ok(())
101    }
102
103    /// Removes a ready thread from its owner run queue for migration or update.
104    pub fn dequeue(&self, cpu: Pin<&mut CpuLocal>, thread: ThreadId) -> Result<(), TaskError> {
105        self.ensure_owner_cpu_context(&cpu)?;
106        let state = self.state.lock();
107        state.ensure_cpu_online(&cpu)?;
108        let record = state.thread_record(thread)?;
109        let mut sched = record.sched.lock();
110        let remote = Arc::clone(cpu.remote());
111        let mut transaction = OwnerRqTxn::begin(self, &remote);
112        if transaction.current_thread() == Some(thread)
113            || transaction.scheduling_entity(thread).is_none()
114        {
115            transaction.commit();
116            return Err(TaskError::NotReady);
117        }
118        let queued = transaction.deactivate_task(thread);
119        record
120            .core
121            .sched()
122            .install_active(&mut sched, queued.into_active());
123        sched.placement.deactivate(cpu.owner());
124        transaction.commit();
125        drop(sched);
126        drop(state);
127        Ok(())
128    }
129}