Skip to main content

ax_task/runtime/service/
ktimer.rs

1//! PREEMPT_RT-style per-CPU `ktimers/%u` service threads.
2
3use crate::{
4    runtime::{
5        context::{
6            RuntimeIrqGuard, current_cpu_remote, runtime_current_cpu_mut, runtime_task_system,
7            validate_task_context,
8        },
9        cpu::CpuRemote,
10        switch::dispatch::yield_current_cpu,
11        task_runtime,
12    },
13    sched::{CpuId, CpuSet, RtPriority, SchedulePolicy},
14    thread::{TaskError, ThreadBuilder, ThreadId, current::current_thread_handle},
15};
16
17/// Creates the current CPU's shutdown-lifetime soft-timer worker.
18///
19/// The runtime calls this exactly once after the rq becomes online and before
20/// enabling local timer IRQs. The thread is bound to that CPU, runs at Linux's
21/// low FIFO kernel-worker priority, and sleeps on the CPU's sticky IRQ event.
22pub fn start_current_ktimer_service() -> Result<(), TaskError> {
23    validate_task_context()?;
24    let system = runtime_task_system()?;
25    let remote = {
26        let _irq = RuntimeIrqGuard::enter();
27        let remote = current_cpu_remote().ok_or(TaskError::NotInitialized)?;
28        remote.begin_ktimer_worker_install()?;
29        remote
30    };
31    let owner = remote.owner();
32    let mut affinity = CpuSet::empty(system.cpu_topology_len());
33    if !affinity.insert(owner) {
34        remote.cancel_ktimer_worker_install();
35        return Err(TaskError::InvalidConfiguration);
36    }
37    let policy = SchedulePolicy::fifo(
38        RtPriority::new(1).expect("Linux low FIFO priority must remain representable"),
39    );
40    let worker = match ThreadBuilder::new(alloc::format!("ktimers/{}", owner.as_u32()))
41        .policy(policy)
42        .affinity(affinity)
43        .spawn(move || ktimer_service_entry(owner))
44    {
45        Ok(worker) => worker,
46        Err(error) => {
47            remote.cancel_ktimer_worker_install();
48            return Err(error);
49        }
50    };
51    worker.detach();
52    Ok(())
53}
54
55fn current_ktimer_remote(expected: CpuId) -> Result<&'static CpuRemote, TaskError> {
56    let _irq = RuntimeIrqGuard::enter();
57    let remote = current_cpu_remote().ok_or(TaskError::NotInitialized)?;
58    if remote.owner() != expected {
59        return Err(TaskError::CpuOwnerMismatch {
60            expected: expected.as_u32(),
61            actual: remote.owner().as_u32(),
62        });
63    }
64    Ok(remote)
65}
66
67fn ktimer_service_entry(owner: CpuId) {
68    if ktimer_service_loop(owner).is_err() {
69        task_runtime::fatal_invariant(0x4b54_0030, owner.as_u32() as usize);
70    }
71}
72
73fn ktimer_service_loop(owner: CpuId) -> Result<(), TaskError> {
74    let remote = current_ktimer_remote(owner)?;
75    let current = current_thread_handle()?;
76    let waiter = crate::sync::irq::worker::IrqWorkerWaiter::new(current.wake_handle());
77    remote.finish_ktimer_worker_install(current.id());
78
79    loop {
80        let Some(claim) = remote.claim_ktimer_work() else {
81            waiter.wait(remote.ktimer_event())?;
82            continue;
83        };
84        let mut processed = 0;
85        let pending = loop {
86            let pass = service_current_ktimer_pass(owner)?;
87            processed += pass.processed;
88            if !pass.pending || pass.processed == 0 || processed >= pass.limit {
89                break pass.pending;
90            }
91        };
92        if pending {
93            // The worker owns forward progress while the claim is active. One
94            // publication at the batch boundary is enough to retain a bounded
95            // remainder across the voluntary scheduling point.
96            remote.publish_ktimer_work();
97        }
98        remote.complete_ktimer_work(claim);
99        if pending {
100            yield_current_cpu()?;
101        }
102    }
103}
104
105struct KtimerServicePass {
106    processed: usize,
107    pending: bool,
108    limit: usize,
109}
110
111fn service_current_ktimer_pass(owner: CpuId) -> Result<KtimerServicePass, TaskError> {
112    let system = runtime_task_system()?;
113    let (processed, pending, limit, mut kernel_timer, mut task_timer, completed_timer) = {
114        let mut irq = RuntimeIrqGuard::enter();
115        let mut cpu = runtime_current_cpu_mut(&mut irq)?;
116        if cpu.owner() != owner {
117            return Err(TaskError::CpuOwnerMismatch {
118                expected: owner.as_u32(),
119                actual: cpu.owner().as_u32(),
120            });
121        }
122        let mut batch = system.service_ktimer_work(cpu.as_mut())?;
123        if let Some(update) = batch.update() {
124            task_runtime::publish_scheduler_deadline(update);
125        }
126        (
127            batch.processed(),
128            batch.pending(),
129            cpu.batch_limit(),
130            batch.take_kernel_timer(),
131            batch.take_task_timer(),
132            batch.take_completed_timer(),
133        )
134    };
135    drop(completed_timer);
136    if let Some(event) = task_timer.take() {
137        let handle = match event.thread().map(|thread| system.thread_handle(thread)) {
138            Some(Ok(handle)) => Some(handle),
139            Some(Err(TaskError::StaleThreadId)) | None => None,
140            Some(Err(_)) => task_runtime::fatal_invariant(
141                0x4b54_0031,
142                event.thread().map_or(0, ThreadId::as_u64) as usize,
143            ),
144        };
145        let (wake, update) = {
146            let mut irq = RuntimeIrqGuard::enter();
147            // Once an expiration leaves the deadline base, this worker owns
148            // its only completion token. Losing the pinned CPU or returning a
149            // recoverable error would orphan that token.
150            let mut cpu = runtime_current_cpu_mut(&mut irq).unwrap_or_else(|_| {
151                task_runtime::fatal_invariant(0x4b54_0032, owner.as_u32() as usize)
152            });
153            if cpu.owner() != owner {
154                task_runtime::fatal_invariant(0x4b54_0033, cpu.owner().as_u32() as usize);
155            }
156            system
157                .complete_task_timer_execution(cpu.as_mut(), event, handle.as_ref())
158                .unwrap_or_else(|_| {
159                    task_runtime::fatal_invariant(0x4b54_0034, owner.as_u32() as usize)
160                })
161        };
162        if let Some(update) = update {
163            task_runtime::publish_scheduler_deadline(update);
164        }
165        if let Some(wake) = wake {
166            let _wake_result = wake.wake();
167        }
168    }
169    if let Some(mut timer) = kernel_timer.take() {
170        let action = timer.invoke_soft();
171        let mut completion = {
172            let mut irq = RuntimeIrqGuard::enter();
173            // The callback entry and its queue tombstone must complete as one
174            // ownership transaction after arbitrary callback code returns.
175            let mut cpu = runtime_current_cpu_mut(&mut irq).unwrap_or_else(|_| {
176                task_runtime::fatal_invariant(0x4b54_0035, owner.as_u32() as usize)
177            });
178            if cpu.owner() != owner {
179                task_runtime::fatal_invariant(0x4b54_0036, cpu.owner().as_u32() as usize);
180            }
181            system
182                .complete_kernel_timer_execution(cpu.as_mut(), timer, action)
183                .unwrap_or_else(|_| {
184                    task_runtime::fatal_invariant(0x4b54_0037, owner.as_u32() as usize)
185                })
186        };
187        if let Some(update) = completion.update() {
188            task_runtime::publish_scheduler_deadline(update);
189        }
190        drop(completion.take_completed());
191    }
192    Ok(KtimerServicePass {
193        processed,
194        pending,
195        limit,
196    })
197}