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 handle = {
43 let state = self.state.lock();
44 state.ensure_cpu_online(&cpu)?;
45 let record = state.thread_record(thread)?;
46 if record.core.execution.is_some() {
49 return Err(TaskError::NotReady);
50 }
51 ThreadHandle::from_core(Arc::clone(&record.core))
52 };
53 self.stage_new_thread(&handle)?;
54 self.activate_staged_thread(cpu.as_mut(), &handle);
55 Ok(())
56 }
57
58 pub(crate) fn stage_new_thread(&self, handle: &ThreadHandle) -> Result<(), TaskError> {
60 let source = CpuId::new(unsafe { task_runtime::current_cpu_id() }.as_u32());
62 let mut state = self.state.lock();
63 let record = state.thread_record(handle.id())?;
64 let sched = record.sched.lock();
65 if sched.lifecycle.state() != ThreadState::New || record.activation.is_some() {
66 return Err(TaskError::NotReady);
67 }
68 let active = record.core.sched().active(&sched);
69 let target = if matches!(active.policy(), SchedulePolicy::Fair { .. }) {
70 state.select_initial_fair_cpu(&sched.affinity.affinity, Some(source))
71 } else {
72 self.select_priority_cpu(
73 active.policy(),
74 Some(active.entity()),
75 &sched.affinity.affinity,
76 Some(source),
77 None,
78 )
79 }
80 .ok_or(TaskError::InvalidConfiguration)?;
81 let delivery = self.prepare_owner_migration(&record.core, source, target)?;
82 drop(active);
83 drop(sched);
84 state.thread_record_mut(handle.id())?.activation = Some(delivery);
85 Ok(())
86 }
87
88 pub fn dequeue(&self, cpu: Pin<&mut CpuLocal>, thread: ThreadId) -> Result<(), TaskError> {
90 self.ensure_owner_cpu_context(&cpu)?;
91 let state = self.state.lock();
92 state.ensure_cpu_online(&cpu)?;
93 let record = state.thread_record(thread)?;
94 let mut sched = record.sched.lock();
95 let remote = Arc::clone(cpu.remote());
96 let mut transaction = OwnerRqTxn::begin(self, &remote);
97 if transaction.current_thread() == Some(thread)
98 || transaction.scheduling_entity(thread).is_none()
99 {
100 transaction.commit();
101 return Err(TaskError::NotReady);
102 }
103 let queued = transaction.deactivate_task(thread);
104 record
105 .core
106 .sched()
107 .install_active(&mut sched, queued.into_active());
108 sched.placement.deactivate(cpu.owner());
109 transaction.commit();
110 drop(sched);
111 drop(state);
112 Ok(())
113 }
114}