use super::*;
use crate::sched::{algorithm::SchedulerClass, system::OwnerRqTaskState};
impl TaskSystem {
pub(in crate::sched::system::task_system) fn resolved_pi_schedule_update(
&self,
base: SchedulePolicy,
base_entity: SchedulingEntity,
donor: Option<(PiWaitKey, PiDonation)>,
generation: u64,
) -> Result<PiScheduleUpdate, TaskError> {
let mut policy = base;
let mut effective_urgency = base_entity.scheduling_urgency(base);
let mut pi_donor = None;
let mut deadline_donor = None;
if let Some((_top, donor)) = donor.as_ref()
&& donor.boost_urgency < effective_urgency
&& let Some(inherited) = pi_inherited_policy(base, donor.policy)
{
policy = inherited;
effective_urgency = donor.boost_urgency;
pi_donor = Some(donor.root);
deadline_donor =
matches!(donor.policy, SchedulePolicy::Deadline(_)).then_some(donor.root);
}
let _ = effective_urgency;
let deadline_donor_core = deadline_donor.map(|donor_id| {
let (_, donor) = donor
.as_ref()
.filter(|(_, donor)| donor.root == donor_id)
.expect("resolved Deadline donor must retain its task reference");
donor.root_core.clone()
});
let deadline_donor_server = deadline_donor_core
.as_ref()
.map(|core| {
core.upgrade()
.ok_or(TaskError::InvalidPiState)
.map(|core| core.sched().deadline_server())
})
.transpose()?;
Ok(PiScheduleUpdate {
policy,
donor: pi_donor,
deadline_donor,
deadline_donor_core,
deadline_donor_server,
generation,
})
}
fn apply_pi_schedule_update_in_rq(
&self,
core: &Arc<ThreadCore>,
sched: &mut ThreadSchedState,
update: PiScheduleUpdate,
transaction: &mut OwnerRqTxn<'_>,
) -> PiRqFollowup {
let owner = sched
.placement
.assigned_cpu()
.expect("PI target must retain task_cpu()");
if transaction.owner() != owner {
task_runtime::fatal_invariant(0x5049_1206, core.id().as_u64() as usize);
}
let rq_state = transaction.task_state(core.id(), &sched.placement);
let owner_now_ns = transaction.clock().wall().as_nanos();
let source_fair = core
.sched()
.active_option(sched)
.and_then(|active| active.base_entity().fair())
.or_else(|| {
transaction
.base_scheduling_entity(core.id())
.and_then(|entity| entity.fair())
});
let fair_placement = match (source_fair, update.policy) {
(Some(_), SchedulePolicy::Fair { .. }) => Some(FairPolicyPlacement {
source_virtual_time: transaction.virtual_time(),
destination_virtual_time: transaction.virtual_time(),
}),
_ => None,
};
if rq_state.is_current() {
let active = transaction.detach_current_schedule(core.id());
let active =
apply_pi_schedule_update(sched, active, update, owner_now_ns, fair_placement)
.unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x5049_1207, core.id().as_u64() as usize)
});
let policy = active.policy();
let entity = active.entity().clone();
let rt_quota_exempt = sched.is_pi_boosted_rt_owner_for(policy);
let metadata = sched.rq_task_metadata().unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x5049_1208, core.id().as_u64() as usize)
});
transaction.install_current_schedule(
core.id(),
active,
Arc::clone(core),
rt_quota_exempt,
sched.affinity.affinity.is_migration_capable(),
metadata,
);
core.publish_effective_schedule(policy, &entity);
return PiRqFollowup::RemoteReschedule;
}
if rq_state.is_delayed_fair() {
let active = transaction
.take_delayed_fair_for_update(core.id())
.into_active();
let mut active =
apply_pi_schedule_update(sched, active, update, owner_now_ns, fair_placement)
.unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x5049_120c, core.id().as_u64() as usize)
});
let policy = active.policy();
let entity = active.entity().clone();
if entity.fair().is_some_and(|fair| fair.is_delayed()) {
let metadata = sched.rq_task_metadata().unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x5049_120d, core.id().as_u64() as usize)
});
let queued = QueuedThread::new(
core.id(),
active,
Arc::clone(core),
false,
sched.affinity.affinity.is_migration_capable(),
metadata,
);
let _entity = transaction.restore_delayed_fair_after_update(queued);
} else {
transaction
.finish_detached_delayed_fair(&mut active, self.config.timing_granularity_ns());
core.sched().install_active(sched, active);
sched.placement.finish_delayed_dequeue(owner);
}
core.publish_effective_schedule(policy, &entity);
return PiRqFollowup::SchedulerWork;
}
if rq_state.is_queued() {
let current_fair = transaction.current_fair_contender();
let active = transaction.reclassify_task(core.id()).into_active();
let active =
apply_pi_schedule_update(sched, active, update, owner_now_ns, fair_placement)
.unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x5049_1209, core.id().as_u64() as usize)
});
let policy = active.policy();
let entity = active.entity().clone();
let rt_quota_exempt = sched.is_pi_boosted_rt_owner_for(policy);
let metadata = sched.rq_task_metadata().unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x5049_120a, core.id().as_u64() as usize)
});
let _enqueue_consumed_by_remote_reschedule = transaction.enqueue_task(
QueuedThread::new(
core.id(),
active,
Arc::clone(core),
rt_quota_exempt,
sched.affinity.affinity.is_migration_capable(),
metadata,
),
EnqueueReason::PolicyChanged,
current_fair,
);
core.publish_effective_schedule(policy, &entity);
return PiRqFollowup::RemoteReschedule;
}
let active = core.sched().take_active(sched);
let active = apply_pi_schedule_update(sched, active, update, owner_now_ns, fair_placement)
.unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x5049_120b, core.id().as_u64() as usize)
});
core.publish_effective_schedule(active.policy(), active.entity());
core.sched().install_active(sched, active);
PiRqFollowup::SchedulerWork
}
pub(in crate::sched::system::task_system) fn recompute_pi_owner_locked(
&self,
core: &Arc<ThreadCore>,
sched: &mut ThreadSchedState,
donor: Option<(PiWaitKey, PiDonation)>,
) -> Result<bool, TaskError> {
record_pi_schedule_recompute_attempt(core.id());
if pi_schedule_update_unchanged_without_rq(
sched.policy.base,
core.effective_policy_snapshot(),
sched.pi.donor,
sched.pi.deadline_donor,
donor.as_ref(),
) {
record_pi_schedule_no_rq_fast_return(core.id());
return Ok(false);
}
let owner = sched
.placement
.assigned_cpu()
.ok_or(TaskError::InvalidPiState)?;
let remote = self
.cpu_remotes
.get(owner.as_usize())
.ok_or(TaskError::InvalidPiState)?;
if !remote.is_online() {
return Err(TaskError::CpuOffline(owner.as_u32()));
}
record_pi_schedule_owner_rq_transaction(core.id());
let mut transaction = OwnerRqTxn::begin(self, remote);
let owner_state = transaction.task_state(core.id(), &sched.placement);
let accounting_path = match owner_state {
OwnerRqTaskState::Current => PiOwnerRqAccountingPath::Running,
OwnerRqTaskState::Queued { .. } | OwnerRqTaskState::DelayedFair { .. } => {
PiOwnerRqAccountingPath::QueuedClassDequeue
}
OwnerRqTaskState::Inactive => PiOwnerRqAccountingPath::Inactive,
};
let owner_accounting_class =
match SchedulerClass::for_policy(core.effective_policy_snapshot()) {
SchedulerClass::Stop => None,
class => Some(class),
};
let current_accounting_class = transaction
.current()
.filter(|current| !current.is_dedicated_idle())
.map(|current| SchedulerClass::for_policy(current.schedule_policy()));
if owner_rq_needs_current_settlement(
accounting_path,
owner_accounting_class,
current_accounting_class,
) {
let _settled = transaction.settle_current(0);
}
let base_entity = core
.sched()
.active_option(sched)
.map(|active| active.base_entity().clone())
.or_else(|| transaction.base_scheduling_entity(core.id()));
let Some(base_entity) = base_entity else {
transaction.commit();
return Err(TaskError::InvalidPiState);
};
let Some(generation) = sched.policy.dispatch_generation.checked_add(1) else {
transaction.commit();
return Err(TaskError::InvalidConfiguration);
};
let update = match self.resolved_pi_schedule_update(
sched.policy.base,
base_entity,
donor,
generation,
) {
Ok(update) => update,
Err(error) => {
transaction.commit();
return Err(error);
}
};
let changed = core.effective_policy_snapshot() != update.policy
|| sched.pi.donor != update.donor
|| sched.pi.deadline_donor != update.deadline_donor;
#[cfg(feature = "qperf-metrics")]
if !changed {
crate::diagnostics::counters::record_pi_schedule_unchanged_after_rq();
}
let followup = if changed {
sched.policy.dispatch_generation = generation;
Some(self.apply_pi_schedule_update_in_rq(core, sched, update, &mut transaction))
} else {
None
};
transaction.commit();
match followup {
Some(PiRqFollowup::RemoteReschedule) => {
remote.request_remote_reschedule(RescheduleKind::Immediate)
}
Some(PiRqFollowup::SchedulerWork) => remote.request_scheduler_work(),
None => {}
}
Ok(changed)
}
}
fn record_pi_schedule_recompute_attempt(owner: ThreadId) {
#[cfg(feature = "qperf-metrics")]
crate::diagnostics::counters::record_pi_schedule_recompute_attempt();
#[cfg(axtest)]
super::axtest::record_recompute_attempt(owner);
#[cfg(not(axtest))]
let _ = owner;
}
fn record_pi_schedule_no_rq_fast_return(owner: ThreadId) {
#[cfg(feature = "qperf-metrics")]
crate::diagnostics::counters::record_pi_schedule_no_rq_fast_return();
#[cfg(axtest)]
super::axtest::record_no_rq_fast_return(owner);
#[cfg(not(axtest))]
let _ = owner;
}
fn record_pi_schedule_owner_rq_transaction(owner: ThreadId) {
#[cfg(feature = "qperf-metrics")]
crate::diagnostics::counters::record_pi_schedule_owner_rq_transaction();
#[cfg(axtest)]
super::axtest::record_owner_rq_transaction(owner);
#[cfg(not(axtest))]
let _ = owner;
}
fn pi_schedule_update_unchanged_without_rq(
base: SchedulePolicy,
effective: SchedulePolicy,
effective_donor: Option<ThreadId>,
effective_deadline_donor: Option<ThreadId>,
donor: Option<&(PiWaitKey, PiDonation)>,
) -> bool {
if matches!(base, SchedulePolicy::Deadline(_)) {
return false;
}
let base_urgency = base.scheduling_urgency();
let mut next_policy = base;
let mut next_donor = None;
if let Some((_top, donation)) = donor
&& donation.boost_urgency < base_urgency
&& let Some(inherited) = pi_inherited_policy(base, donation.policy)
{
if matches!(inherited, SchedulePolicy::Deadline(_)) {
return false;
}
next_policy = inherited;
next_donor = Some(donation.root);
}
effective == next_policy && effective_donor == next_donor && effective_deadline_donor.is_none()
}
fn pi_inherited_policy(base: SchedulePolicy, donor: SchedulePolicy) -> Option<SchedulePolicy> {
match donor {
SchedulePolicy::Deadline(policy) => Some(SchedulePolicy::Deadline(policy)),
SchedulePolicy::Fifo { priority } | SchedulePolicy::RoundRobin { priority, .. } => {
Some(match base {
SchedulePolicy::RoundRobin { quantum_ns, .. } => SchedulePolicy::RoundRobin {
priority,
quantum_ns,
},
SchedulePolicy::KernelStop | SchedulePolicy::Deadline(_) => return None,
SchedulePolicy::Fair { .. } | SchedulePolicy::Fifo { .. } => {
SchedulePolicy::Fifo { priority }
}
})
}
SchedulePolicy::Fair { .. } => matches!(base, SchedulePolicy::Fair { .. }).then_some(donor),
SchedulePolicy::KernelStop => None,
}
}