use crate::{
runtime::{
context::{
RuntimeIrqGuard, current_cpu_remote, runtime_current_cpu_mut, runtime_task_system,
validate_task_context,
},
cpu::CpuRemote,
switch::dispatch::yield_current_cpu,
task_runtime,
},
sched::{CpuId, CpuSet, RtPriority, SchedulePolicy},
thread::{TaskError, ThreadBuilder, ThreadId, current::current_thread_handle},
};
pub fn start_current_ktimer_service() -> Result<(), TaskError> {
validate_task_context()?;
let system = runtime_task_system()?;
let remote = {
let _irq = RuntimeIrqGuard::enter();
let remote = current_cpu_remote().ok_or(TaskError::NotInitialized)?;
remote.begin_ktimer_worker_install()?;
remote
};
let owner = remote.owner();
let mut affinity = CpuSet::empty(system.cpu_topology_len());
if !affinity.insert(owner) {
remote.cancel_ktimer_worker_install();
return Err(TaskError::InvalidConfiguration);
}
let policy = SchedulePolicy::fifo(
RtPriority::new(1).expect("Linux low FIFO priority must remain representable"),
);
let worker = match ThreadBuilder::new(alloc::format!("ktimers/{}", owner.as_u32()))
.policy(policy)
.affinity(affinity)
.spawn(move || ktimer_service_entry(owner))
{
Ok(worker) => worker,
Err(error) => {
remote.cancel_ktimer_worker_install();
return Err(error);
}
};
worker.detach_permanent();
Ok(())
}
fn current_ktimer_remote(expected: CpuId) -> Result<&'static CpuRemote, TaskError> {
let _irq = RuntimeIrqGuard::enter();
let remote = current_cpu_remote().ok_or(TaskError::NotInitialized)?;
if remote.owner() != expected {
return Err(TaskError::CpuOwnerMismatch {
expected: expected.as_u32(),
actual: remote.owner().as_u32(),
});
}
Ok(remote)
}
fn ktimer_service_entry(owner: CpuId) {
if ktimer_service_loop(owner).is_err() {
task_runtime::fatal_invariant(0x4b54_0030, owner.as_u32() as usize);
}
}
fn ktimer_service_loop(owner: CpuId) -> Result<(), TaskError> {
let remote = current_ktimer_remote(owner)?;
let current = current_thread_handle()?;
let waiter = crate::sync::irq::worker::IrqWorkerWaiter::new(current.wake_handle());
remote.finish_ktimer_worker_install(current.id());
loop {
let Some(claim) = remote.claim_ktimer_work() else {
waiter.wait(remote.ktimer_event())?;
continue;
};
let mut processed = 0;
let pending = loop {
let pass = service_current_ktimer_pass(owner)?;
processed += pass.processed;
if !pass.pending || pass.processed == 0 || processed >= pass.limit {
break pass.pending;
}
};
if pending {
remote.publish_ktimer_work();
}
remote.complete_ktimer_work(claim);
if pending {
yield_current_cpu()?;
}
}
}
struct KtimerServicePass {
processed: usize,
pending: bool,
limit: usize,
}
fn service_current_ktimer_pass(owner: CpuId) -> Result<KtimerServicePass, TaskError> {
let system = runtime_task_system()?;
let (processed, pending, limit, mut kernel_timer, mut task_timer, completed_timer) = {
let mut irq = RuntimeIrqGuard::enter();
let mut cpu = runtime_current_cpu_mut(&mut irq)?;
if cpu.owner() != owner {
return Err(TaskError::CpuOwnerMismatch {
expected: owner.as_u32(),
actual: cpu.owner().as_u32(),
});
}
let mut batch = system.service_ktimer_work(cpu.as_mut())?;
if let Some(update) = batch.update() {
task_runtime::publish_scheduler_deadline(update);
}
(
batch.processed(),
batch.pending(),
cpu.batch_limit(),
batch.take_kernel_timer(),
batch.take_task_timer(),
batch.take_completed_timer(),
)
};
drop(completed_timer);
if let Some(event) = task_timer.take() {
let handle = match event.thread().map(|thread| system.thread_handle(thread)) {
Some(Ok(handle)) => Some(handle),
Some(Err(TaskError::StaleThreadId)) | None => None,
Some(Err(_)) => task_runtime::fatal_invariant(
0x4b54_0031,
event.thread().map_or(0, ThreadId::as_u64) as usize,
),
};
let (wake, update) = {
let mut irq = RuntimeIrqGuard::enter();
let mut cpu = runtime_current_cpu_mut(&mut irq).unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x4b54_0032, owner.as_u32() as usize)
});
if cpu.owner() != owner {
task_runtime::fatal_invariant(0x4b54_0033, cpu.owner().as_u32() as usize);
}
system
.complete_task_timer_execution(cpu.as_mut(), event, handle.as_ref())
.unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x4b54_0034, owner.as_u32() as usize)
})
};
if let Some(update) = update {
task_runtime::publish_scheduler_deadline(update);
}
if let Some(wake) = wake {
let _wake_result = wake.wake();
}
}
if let Some(mut timer) = kernel_timer.take() {
let action = timer.invoke_soft();
let mut completion = {
let mut irq = RuntimeIrqGuard::enter();
let mut cpu = runtime_current_cpu_mut(&mut irq).unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x4b54_0035, owner.as_u32() as usize)
});
if cpu.owner() != owner {
task_runtime::fatal_invariant(0x4b54_0036, cpu.owner().as_u32() as usize);
}
system
.complete_kernel_timer_execution(cpu.as_mut(), timer, action)
.unwrap_or_else(|_| {
task_runtime::fatal_invariant(0x4b54_0037, owner.as_u32() as usize)
})
};
if let Some(update) = completion.update() {
task_runtime::publish_scheduler_deadline(update);
}
drop(completion.take_completed());
}
Ok(KtimerServicePass {
processed,
pending,
limit,
})
}