use super::*;
impl TaskSystem {
pub fn thread_state(&self, thread: ThreadId) -> Result<ThreadState, TaskError> {
Ok(self
.state
.lock()
.thread_record(thread)?
.sched
.lock()
.lifecycle
.state())
}
pub fn thread_runtime(&self, thread: ThreadId) -> Result<ThreadRuntimeSnapshot, TaskError> {
let (core, sched_cell) = {
let state = self.state.lock();
let record = state.thread_record(thread)?;
(Arc::clone(&record.core), Arc::clone(&record.sched))
};
let sched = sched_cell.lock();
let snapshot = if let Some(cpu) = sched.placement.assigned_cpu() {
let remote = self
.cpu_remotes
.get(cpu.as_usize())
.ok_or(TaskError::InvalidCpu(cpu.as_u32()))?;
let transaction = OwnerRqTxn::begin(self, remote);
let rq_state = transaction.task_state(thread, &sched.placement);
let running_interval_ns = if rq_state.is_current() {
let dispatch = transaction
.current()
.filter(|dispatch| dispatch.thread() == thread);
Some(
dispatch
.unwrap_or_else(|| {
task_runtime::fatal_invariant(0x5251_1210, thread.as_u64() as usize)
})
.runtime_interval_ns(transaction.clock().task().as_nanos()),
)
} else {
None
};
let snapshot = core.runtime_snapshot(running_interval_ns);
transaction.commit();
snapshot
} else {
core.runtime_snapshot(None)
};
Ok(snapshot)
}
pub fn replace_current_address_space(
&self,
cpu: Pin<&mut CpuLocal>,
address_space: &mut crate::runtime::resource::AddressSpaceToken,
) -> Result<crate::runtime::resource::AddressSpaceToken, TaskError> {
self.ensure_owner_cpu_context(&cpu)?;
if address_space.is_none() {
return Err(TaskError::InvalidConfiguration);
}
let mut state = self.state.lock();
state.ensure_cpu_online(&cpu)?;
let owner = cpu.owner();
let current = cpu.current().ok_or(TaskError::NoRunnableThread)?;
let record = state.thread_record_mut(current)?;
let mut sched = record.sched.lock();
if sched.lifecycle.state() != ThreadState::Running
|| sched.placement.queued_cpu() != Some(owner)
|| sched.placement.on_cpu() != Some(owner)
{
return Err(TaskError::InvalidConfiguration);
}
let next_handle = address_space.handle();
let next_membarrier_state = task_runtime::address_space_membarrier_state(next_handle);
let binding =
crate::runtime::switch::ThreadRuntimeBinding::new(sched.runtime.context, next_handle);
let remote = Arc::clone(cpu.remote());
let mut transaction = OwnerRqTxn::begin(self, &remote);
transaction.update_current_runtime_binding(current, binding, next_membarrier_state);
let next = core::mem::replace(
address_space,
crate::runtime::resource::AddressSpaceToken::NONE,
);
transaction.commit();
let previous = record.resources.replace_address_space(next);
sched.runtime.address_space = next_handle;
Ok(previous)
}
pub fn detach_current_address_space(
&self,
cpu: Pin<&mut CpuLocal>,
) -> Result<crate::runtime::resource::AddressSpaceToken, TaskError> {
self.ensure_owner_cpu_context(&cpu)?;
let mut state = self.state.lock();
state.ensure_cpu_online(&cpu)?;
let owner = cpu.owner();
let current = cpu.current().ok_or(TaskError::NoRunnableThread)?;
let record = state.thread_record_mut(current)?;
let mut sched = record.sched.lock();
if sched.lifecycle.state() != ThreadState::Running
|| sched.placement.queued_cpu() != Some(owner)
|| sched.placement.on_cpu() != Some(owner)
|| record.resources.address_space().is_none()
{
return Err(TaskError::InvalidConfiguration);
}
let binding = crate::runtime::switch::ThreadRuntimeBinding::new(
sched.runtime.context,
crate::runtime::resource::AddressSpaceHandle::NONE,
);
let remote = Arc::clone(cpu.remote());
let mut transaction = OwnerRqTxn::begin(self, &remote);
transaction.update_current_runtime_binding(
current,
binding,
crate::runtime::resource::AddressSpaceMembarrierState::NONE,
);
let previous = record.resources.take_address_space();
transaction.commit();
sched.runtime.address_space = crate::runtime::resource::AddressSpaceHandle::NONE;
Ok(previous)
}
pub fn thread_handle(&self, thread: ThreadId) -> Result<ThreadHandle, TaskError> {
let state = self.state.lock();
let record = state.thread_record(thread)?;
Ok(ThreadHandle::from_core(Arc::clone(&record.core)))
}
pub fn thread_extension<'thread>(
&self,
handle: &'thread ThreadHandle,
) -> Result<Option<ThreadExtensionBorrow<'thread>>, TaskError> {
let view = self.thread_extension_view(handle)?;
Ok(view.map(|view| ThreadExtensionBorrow::new(view, handle)))
}
pub fn thread_extension_lease(
&self,
handle: ThreadHandle,
) -> Result<Option<ThreadExtensionLease>, TaskError> {
let view = self.thread_extension_view(&handle)?;
Ok(view.map(|view| ThreadExtensionLease::new(view, handle)))
}
fn thread_extension_view(
&self,
handle: &ThreadHandle,
) -> Result<Option<ThreadExtensionView>, TaskError> {
let state = self.state.lock();
let record = state.thread_record(handle.id())?;
if !Arc::ptr_eq(&record.core, &handle.core) {
return Err(TaskError::StaleThreadId);
}
Ok(handle.extension_view())
}
pub fn set_thread_policy(
&self,
thread: ThreadId,
policy: SchedulePolicy,
) -> Result<(), TaskError> {
policy.validate()?;
let mut affinity = CpuSet::empty(self.config.cpu_count());
let core = {
let state = self.state.lock();
Arc::clone(&state.thread_record(thread)?.core)
};
let _activity = core.try_scheduler_activity().ok_or(TaskError::NotReady)?;
let state = self.state.lock();
let mut root_domain = self.root_domain.lock();
let record = state.thread_record(thread)?;
if !Arc::ptr_eq(&record.core, &core) {
return Err(TaskError::StaleThreadId);
}
let sched_cell = Arc::clone(&record.sched);
let mut sched = sched_cell.lock();
if sched.lifecycle.state() == ThreadState::Exited {
return Err(TaskError::NotReady);
}
affinity.copy_from_set(&sched.affinity.requested_affinity)?;
let applied_reservation = sched.deadline.bandwidth.reservation_scaled();
let pending_reservation = sched
.policy
.pending_update()
.map_or(0, |pending| pending.reservation_scaled);
let owner = sched
.placement
.assigned_cpu()
.ok_or(TaskError::InvalidPiState)?;
let reservation_owner = sched.deadline.bandwidth.reservation_owner();
if let Some(reservation_owner) = reservation_owner
&& owner != reservation_owner
{
task_runtime::fatal_invariant(0x444c_1201, core.id().as_u64() as usize);
}
let remote = self
.cpu_remotes
.get(owner.as_usize())
.ok_or(TaskError::InvalidCpu(owner.as_u32()))?;
let reservation = root_domain.deadline_reservation_for(policy, &affinity)?;
let pending = sched.policy.prepare_update(policy, reservation)?;
let old_held = applied_reservation.max(pending_reservation);
let new_held = applied_reservation.max(reservation);
root_domain.replace_deadline_utilization(old_held, new_held)?;
sched.policy.publish_update(pending);
drop(state);
let applied = self
.apply_owner_policy_update_locked(remote, &core, &mut sched, pending.generation)
.unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x5251_1208, core.id().as_u64() as usize)
});
Self::finish_policy_admission_locked(&mut root_domain, &core, applied.commit);
drop(root_domain);
drop(sched);
Self::notify_policy_generation(&core, applied.commit);
self.recompute_pi_after_policy_update(core.id())
.unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x5049_1216, core.id().as_u64() as usize)
});
let owner_work_required =
applied.scheduler_deadline_refresh_required || applied.rt_period_started;
match (applied.reschedule, owner_work_required) {
(Some(kind), true) => {
remote.request_remote_reschedule_with_scheduler_work(kind);
}
(Some(kind), false) => {
remote.request_remote_reschedule(kind);
}
(None, true) => {
remote.kick_scheduler_work();
}
(None, false) => {}
}
Ok(())
}
pub fn thread_affinity(&self, thread: ThreadId) -> Result<CpuSet, TaskError> {
Ok(self
.state
.lock()
.thread_record(thread)?
.sched
.lock()
.affinity
.requested_affinity
.as_ref()
.clone())
}
}