ax_task/sched/system/task_system/delivery/
admission.rs1use super::*;
4
5impl TaskSystem {
6 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 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 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}