use alloc::{sync::Arc, vec::Vec};
use core::ffi::c_long;
use ax_runtime::hal::time::TimeValue;
use axpoll::IoEvents;
use bytemuck::AnyBitPattern;
use linux_raw_sys::general::ROBUST_LIST_LIMIT;
use starry_signal::{SignalInfo, Signo};
use super::{
AlarmTarget, AlarmToken, PendingTimerActions, ProcessData, Thread, UserTaskRef, ZombieSnapshot,
current_user_task, processes, publish_zombie, resolve_futex_for_process_teardown,
send_signal_to_process, send_signal_to_process_data, send_signal_to_thread, yield_now,
};
use crate::{
StarryError, StarryResult,
mm::{VmMutPtr, VmPtr},
task::{
PgidNumber, PidIdentity, PidNamespaceLifecycle, PidNamespaceRef, PidView, Process,
ProcessCpuTime, ProcessGroup, ROOT_PID_NS, Tgid, ThreadExit, Tid, TidNumber,
},
};
const FUTEX_OWNER_DIED: u32 = 0x40000000;
const FUTEX_TID_MASK: u32 = 0x3fffffff;
const FUTEX_WAITERS: u32 = 0x80000000;
pub fn decode_wait_status(raw: i32) -> (i32, i32) {
use linux_raw_sys::general::{CLD_DUMPED, CLD_EXITED, CLD_KILLED};
if raw & 0x7f == 0 {
(CLD_EXITED as i32, (raw >> 8) & 0xff)
} else {
let signum = raw & 0x7f;
if (raw & 0x80) != 0 {
(CLD_DUMPED as i32, signum)
} else {
(CLD_KILLED as i32, signum)
}
}
}
pub fn tasks() -> Vec<UserTaskRef> {
ROOT_PID_NS
.published_members()
.into_iter()
.filter(|identity| identity.has_role::<Tid>())
.filter_map(|identity| identity.live_task())
.collect()
}
pub(crate) fn get_task_by_number(tid: TidNumber) -> StarryResult<UserTaskRef> {
PidView::new(ROOT_PID_NS.clone())
.resolve_thread(tid)?
.live_task()
.ok_or(StarryError::NoSuchProcess)
}
pub(crate) fn get_user_task_by_number(tid: TidNumber) -> StarryResult<UserTaskRef> {
super::current_pid_view()
.resolve_thread(tid)?
.live_task()
.ok_or(StarryError::NoSuchProcess)
}
pub fn detach_live_tracees_of(tracer: &Arc<PidIdentity>) {
if !tracer.may_have_ptrace_tracees() {
return;
}
for tracee in processes() {
if !tracee
.ptrace_tracer_identity()
.is_some_and(|registered| Arc::ptr_eq(®istered, tracer))
{
continue;
}
tracee.clear_ptrace_stop();
tracee.clear_ptrace_traceme();
tracee.clear_ptrace_attached();
tracee.clear_ptrace_tracer();
tracee.set_ptrace_options(0);
}
}
pub(crate) fn get_process_group_by_number(pgid: PgidNumber) -> StarryResult<Arc<ProcessGroup>> {
PidView::new(ROOT_PID_NS.clone()).resolve_group(pgid)
}
pub fn task_cpu_time(task: &UserTaskRef) -> (TimeValue, TimeValue) {
task.as_thread().cpu_time_output()
}
fn apply_process_timer_actions(proc_data: &ProcessData, pending: PendingTimerActions) {
let pid = proc_data.proc.pid_number();
for signo in pending.signals() {
let _ = send_signal_to_process(pid, Some(SignalInfo::new_kernel(signo)));
}
pending.apply_alarms(AlarmTarget::Process(Arc::downgrade(&proc_data.identity())));
}
#[cfg(all(test, axtest))]
mod axtests {
use alloc::{string::ToString, sync::Arc};
use core::sync::atomic::{AtomicBool, Ordering};
use ax_runtime::{
task::{
diagnostics::{
begin_pi_schedule_test_probe, end_pi_schedule_test_probe,
pi_schedule_test_probe_snapshot,
},
sync::WaitQueue,
thread::current::current_thread_id,
},
thread::{join_thread, spawn_raw},
};
use crate::sync::Mutex;
fn wait_for(mut condition: impl FnMut() -> bool, message: &str) {
for _ in 0..1_000_000 {
if condition() {
return;
}
ax_std::thread::yield_now();
}
panic!("{message}");
}
#[axtest::axtest]
fn kernel_thread_retains_active_mm_membarrier_state() {
assert!(
ax_runtime::thread::kernel_thread_retains_active_mm_membarrier_state_for_test(),
"a kernel thread borrows the CPU's active mm and must retain its rq membarrier state",
);
}
#[axtest::axtest]
fn unchanged_pi_schedule_returns_before_the_owner_rq_transaction() {
let mutex = Arc::new(Mutex::new(()));
let owner_wait = Arc::new(WaitQueue::new());
let owner_locked = Arc::new(AtomicBool::new(false));
let release_owner = Arc::new(AtomicBool::new(false));
let waiter_done = Arc::new(AtomicBool::new(false));
let owner = {
let mutex = Arc::clone(&mutex);
let owner_wait = Arc::clone(&owner_wait);
let owner_locked = Arc::clone(&owner_locked);
let release_owner = Arc::clone(&release_owner);
spawn_raw(
move || {
begin_pi_schedule_test_probe(
current_thread_id().expect("PI owner must have a thread identity"),
);
let _guard = mutex.lock();
owner_locked.store(true, Ordering::Release);
owner_wait.wait_until(|| release_owner.load(Ordering::Acquire));
},
"pi-no-rq-owner".to_string(),
256 * 1024,
)
.expect("failed to spawn PI owner")
};
wait_for(
|| owner_locked.load(Ordering::Acquire),
"PI owner did not acquire the mutex",
);
let waiter = {
let mutex = Arc::clone(&mutex);
let waiter_done = Arc::clone(&waiter_done);
spawn_raw(
move || {
drop(mutex.lock());
waiter_done.store(true, Ordering::Release);
},
"pi-no-rq-waiter".to_string(),
256 * 1024,
)
.expect("failed to spawn PI waiter")
};
wait_for(
|| {
let snapshot = pi_schedule_test_probe_snapshot();
snapshot.recompute_attempts > 0
&& snapshot.no_rq_fast_returns + snapshot.owner_rq_transactions
>= snapshot.recompute_attempts
},
"equal-policy PI contention did not reach owner recompute",
);
let registered = pi_schedule_test_probe_snapshot();
assert_eq!(
registered.owner_rq_transactions, 0,
"unchanged PI state must return before the owner-rq transaction"
);
assert_eq!(
registered.no_rq_fast_returns, registered.recompute_attempts,
"every equal-policy PI recompute must resolve from task-owned state"
);
release_owner.store(true, Ordering::Release);
owner_wait.notify_all();
join_thread(owner).expect("PI owner must exit cleanly");
join_thread(waiter).expect("PI waiter must exit cleanly");
assert!(
waiter_done.load(Ordering::Acquire),
"PI waiter must acquire the mutex after owner release"
);
let completed = pi_schedule_test_probe_snapshot();
end_pi_schedule_test_probe();
assert!(
completed.recompute_attempts >= 2,
"PI registration and release must both recompute the owner schedule"
);
assert_eq!(
completed.no_rq_fast_returns, completed.recompute_attempts,
"unchanged registration and deboost must both avoid the owner rq"
);
assert_eq!(
completed.owner_rq_transactions, 0,
"unchanged registration and deboost must not enter the owner rq"
);
assert_eq!(
completed.waiter_registrations, 1,
"the probe must observe the PI waiter registration"
);
assert_eq!(
completed.parking_waiter_registrations, completed.waiter_registrations,
"Linux publishes the rtmutex wait state before linking the waiter"
);
}
}
fn poll_interval_timers(proc_data: &ProcessData, token: Option<&AlarmToken>) {
if !proc_data.has_active_interval_timers() {
return;
}
let snapshot = proc_data.cpu_time_snapshot();
if let Some(pending) = proc_data.poll_interval_timers(snapshot, token) {
apply_process_timer_actions(proc_data, pending);
}
}
pub(crate) fn poll_process_cpu_timers_from_scheduler_tick(proc_data: &ProcessData) {
if !proc_data.has_active_cpu_interval_timers() {
return;
}
let snapshot = proc_data.scheduler_tick_cpu_time_snapshot();
if let Some(pending) = proc_data.poll_cpu_interval_timers(snapshot) {
apply_process_timer_actions(proc_data, pending);
}
}
pub(crate) fn poll_process_timer_for_alarm(identity: &Arc<PidIdentity>, token: &AlarmToken) {
if let Some(proc_data) = identity.live_data() {
poll_interval_timers(&proc_data, Some(token));
proc_data.posix_timers().poll_expired_for(
AlarmTarget::Process(Arc::downgrade(identity)),
token,
|sig| {
let _ = send_signal_to_process(proc_data.proc.pid_number(), Some(sig));
},
);
}
}
#[repr(C)]
#[derive(Debug, Copy, Clone, AnyBitPattern)]
pub struct RobustList {
pub next: *mut RobustList,
}
#[repr(C)]
#[derive(Debug, Copy, Clone, AnyBitPattern)]
pub struct RobustListHead {
pub list: RobustList,
pub futex_offset: c_long,
pub list_op_pending: *mut RobustList,
}
fn robust_futex_address(entry: *mut RobustList, offset: i64) -> StarryResult<usize> {
let address = (entry as u64)
.checked_add_signed(offset)
.ok_or(StarryError::InvalidInput)?;
let address = usize::try_from(address).map_err(|_| StarryError::InvalidInput)?;
if address % size_of::<u32>() != 0 {
return Err(StarryError::InvalidInput);
}
Ok(address)
}
fn wake_robust_futex(proc_data: &ProcessData, address: usize) {
resolve_futex_for_process_teardown(proc_data, address).wake(1, u32::MAX);
}
fn handle_futex_death(
current: &UserTaskRef,
thr: &Thread,
entry: *mut RobustList,
offset: i64,
pending: bool,
) -> StarryResult<()> {
let address = robust_futex_address(entry, offset)?;
let futex_word = address as *mut u32;
let owner_tid = thr.user_tid().get() & FUTEX_TID_MASK;
let value = futex_word.vm_read(current)?;
let owner = value & FUTEX_TID_MASK;
if pending && owner == 0 {
wake_robust_futex(&thr.proc_data, address);
return Ok(());
}
if owner != owner_tid {
return Ok(());
}
futex_word.vm_write(current, (value & FUTEX_WAITERS) | FUTEX_OWNER_DIED)?;
if value & FUTEX_WAITERS != 0 {
wake_robust_futex(&thr.proc_data, address);
}
Ok(())
}
pub fn exit_robust_list(
current: &UserTaskRef,
thr: &Thread,
head: *const RobustListHead,
) -> crate::StarryResult<()> {
let mut limit = ROBUST_LIST_LIMIT;
let end_ptr = head.cast::<RobustList>() as *mut RobustList;
let head = head.vm_read(current)?;
let mut entry = head.list.next;
let offset = head.futex_offset;
let pending = (head.list_op_pending as usize & !1) as *mut RobustList;
while !core::ptr::eq(entry, end_ptr) {
if entry.is_null() {
break;
}
let Ok(node) = entry.vm_read(current) else {
debug!("robust list: failed to read entry {entry:?}");
break;
};
let next_entry = node.next;
if entry != pending {
handle_futex_death(current, thr, entry, offset, false).unwrap_or_else(|err| {
debug!("robust list: failed to clean entry {entry:?}: {err:?}");
});
}
entry = next_entry;
limit -= 1;
if limit == 0 {
debug!("robust list: entry limit reached");
break;
}
yield_now();
}
if !pending.is_null() && !core::ptr::eq(pending, end_ptr) {
handle_futex_death(current, thr, pending, offset, true).unwrap_or_else(|err| {
debug!("robust list: failed to clean pending entry {pending:?}: {err:?}");
});
}
Ok(())
}
ax_tracepoint::define_event_trace!(
sched_process_exit,
TP_kops(crate::tracepoint::KernelTraceAux),
TP_system(sched),
TP_PROTO(tid: u64, exit_code: i32),
TP_STRUCT__entry {
tid: u64,
exit_code: i32,
},
TP_fast_assign {
tid: tid,
exit_code: exit_code,
},
TP_ident(__entry),
TP_printk({ alloc::format!("tid={} exit_code={}", __entry.tid, __entry.exit_code,) })
);
fn emit_sched_process_exit(tid: TidNumber, exit_code: i32) {
trace_sched_process_exit(tid.get() as u64, exit_code);
}
fn close_process_relations_for_exit(
process: &Arc<Process>,
pid_namespace: &PidNamespaceRef,
) -> Vec<Arc<Process>> {
loop {
if pid_namespace.lifecycle() == PidNamespaceLifecycle::ShuttingDown {
return process
.begin_namespace_shutdown_relations()
.into_retained_children();
}
let orphan_reaper = super::orphan_reaper_for(process);
if let Some(relations) = process.try_begin_exit_relations(&orphan_reaper) {
return relations.into_reparented_children();
}
}
}
pub fn do_exit(exit_code: i32, group_exit: bool) {
let curr = current_user_task();
let thr = curr.as_thread();
if !thr.begin_exit() {
return;
}
info!("{} exit with code: {}", curr.id_name(), exit_code);
emit_sched_process_exit(thr.tid(), exit_code);
if group_exit && let Some(tids) = thr.proc_data.proc.start_group_exit(exit_code) {
let sig = SignalInfo::new_kernel(Signo::SIGKILL);
for tid in tids {
if tid == thr.tid_number() {
continue;
}
let _ = send_signal_to_thread(None, tid, Some(sig));
let _ = zap_thread(tid);
}
}
#[cfg(target_arch = "aarch64")]
crate::perf::task::on_task_exit(thr);
let head = thr.robust_list_head() as *const RobustListHead;
if !head.is_null()
&& let Err(err) = exit_robust_list(&curr, thr, head)
{
warn!("exit robust list failed: {err:?}");
}
let clear_child_tid = thr.clear_child_tid() as *mut u32;
if clear_child_tid.vm_write(&curr, 0).is_ok() {
resolve_futex_for_process_teardown(&thr.proc_data, clear_child_tid as usize)
.wake(1, u32::MAX);
yield_now();
}
let process = &thr.proc_data.proc;
crate::file::close_all_fds();
let retired_fs =
thr.with_current_scope_mut(|scope| ax_fs_ng::vfs::FS_CONTEXT.scope_mut(scope).take());
drop(retired_fs);
ax_runtime::thread::detach_current_address_space()
.unwrap_or_else(|error| panic!("failed to detach exiting task address space: {error}"));
let is_process_leader = thr.tid().pid_number() == process.pid().pid_number();
thr.commit_cpu_time_now();
let (utime, stime) = task_cpu_time(&curr);
let task_identity = thr.pid_identity();
let exit_path = if is_process_leader {
let (tid_lease, exit_path) = thr.retire_pid_retaining_tid();
thr.proc_data.retire_leader(thr.nice(), tid_lease);
exit_path
} else {
thr.retire_pid()
};
let task_generation = ax_cgroup::ProcessId::new(task_identity.id().get())
.expect("PID identity generation must be non-zero");
let (thread_exit, cgroup_exit) = thr.proc_data.finish_thread_exit(task_generation, || {
process.exit_thread(
thr.tid_number(),
exit_code,
ProcessCpuTime::new(utime, stime),
)
});
super::cgroup_exit_invariant::enforce(cgroup_exit);
if let ThreadExit::Last(exit_owner) = thread_exit {
debug_assert!(Arc::ptr_eq(exit_owner.process(), process));
thr.proc_data.release_cgroup_namespace();
thr.proc_data
.cancel_interval_timer_alarm()
.apply_cancellation();
thr.proc_data.posix_timers().clear();
crate::syscall::cleanup_aio_contexts_for_process(thr.proc_data.identity().id());
detach_live_tracees_of(&process.identity());
let process_identity_id = process.identity().id();
crate::syscall::release_pid_locks(process_identity_id);
crate::syscall::release_pid_flock_locks(process_identity_id);
let pid_ns = thr.active_pid_namespace();
let identity = thr.proc_data.identity();
let shutdown_executor = thr.pid_identity();
let namespace_shutdown = if pid_ns.init_identity() == Some(identity.id()) {
Some(
pid_ns
.begin_shutdown(identity.id(), shutdown_executor.id())
.expect("PID namespace init failed to enter shutdown"),
)
} else {
None
};
let children_snapshot = if namespace_shutdown.is_some() {
process
.begin_namespace_shutdown_relations()
.into_retained_children()
} else {
close_process_relations_for_exit(process, &pid_ns)
};
if let Some(shutdown) = namespace_shutdown.as_ref() {
let sig = SignalInfo::new_kernel(Signo::SIGKILL);
for victim in pid_ns.published_members() {
if victim.id() != identity.id()
&& victim.has_role::<Tgid>()
&& let Ok(victim_process) = victim.public_process()
{
let _ = send_signal_to_process(victim_process.pid_number(), Some(sig));
for tid in victim_process.threads() {
if let Ok(task) = get_task_by_number(tid) {
task.interrupt();
}
}
}
}
shutdown.wait_for_descendants_exit();
}
let zombie_cred = thr.cred();
let ptrace_tracer = thr.proc_data.ptrace_tracer_identity();
let is_clone_child = thr.proc_data.is_clone_child();
let wait_parent_tid = thr.proc_data.wait_parent_tid();
let (zombie_nice, leader_tid_lease) = thr.proc_data.take_retired_leader_for_zombie();
if let Ok(aspace) = thr.proc_data.pin_aspace() {
crate::syscall::clear_proc_shm(
process_identity_id,
process.identity().snapshot(),
&aspace,
);
} else {
warn!("shared-memory exit cleanup skipped for an unavailable MM");
}
thr.proc_data.retire_mm_owner();
publish_zombie(
&thr.proc_data,
ZombieSnapshot {
cred: zombie_cred,
nice: zombie_nice,
ptrace_tracer: ptrace_tracer.as_ref().map(|identity| identity.snapshot()),
is_clone_child,
wait_parent_tid,
cpu_time: exit_owner.cpu_time(),
tid_lease: leader_tid_lease,
tgid_lease: thr.proc_data.take_tgid_lease(),
},
)
.expect("last process thread must own one live PID identity");
if let Some(parent) = process.parent()
&& let Some(parent_data) = parent.identity().live_data()
{
if let Some(signo) = thr.proc_data.exit_signal() {
use starry_signal::Signo;
let child_uid = thr.cred().uid;
let (code, status) = decode_wait_status(exit_owner.exit_code());
let sig = if signo == Signo::SIGCHLD {
let child_pid = process
.identity()
.visible_number(&parent.identity().active_namespace())
.unwrap_or_else(|| {
panic!(
"child process must be visible to its parent: child id={:?} \
snapshot={:?}, parent id={:?} snapshot={:?} parent active \
ns={:?} lifecycle={:?}",
process.identity().id(),
process.identity().snapshot(),
parent.identity().id(),
parent.identity().snapshot(),
parent.identity().active_namespace().id(),
parent.identity().active_namespace().lifecycle(),
)
})
.get();
SignalInfo::new_sigchld(child_pid, child_uid, code, status)
} else {
SignalInfo::new_kernel(signo)
};
let _ = send_signal_to_process_data(&parent_data, Some(sig));
}
unsafe { parent_data.child_exit_event().wake(axpoll::IoEvents::IN) };
}
if let Some(tracer) = ptrace_tracer
&& process
.parent()
.is_none_or(|parent| !Arc::ptr_eq(&parent.identity(), &tracer))
&& let Some(data) = tracer.live_data()
{
unsafe { data.child_exit_event().wake(axpoll::IoEvents::IN) };
}
for child in children_snapshot {
let child_tid = TidNumber::from(child.pid_number().pid_number());
if let Ok(child_task) = get_task_by_number(child_tid) {
let child_thr = child_task.as_thread();
let sig = child_thr.pdeathsig();
if sig > 0
&& let Some(signo) = Signo::from_repr(sig as u8)
{
let _ = send_signal_to_process(
child.pid_number(),
Some(SignalInfo::new_kernel(signo)),
);
}
}
}
unsafe {
thr.proc_data
.exit_event()
.wake(IoEvents::IN | IoEvents::RDNORM);
};
thr.proc_data.notify_vfork_done();
}
thr.set_exit();
task_identity.notify_thread_pidfd_exit();
unsafe { thr.exit_event().wake(axpoll::IoEvents::IN) };
exit_path.complete();
unsafe { thr.proc_data.thread_exit_event().wake(axpoll::IoEvents::IN) };
}
pub fn zap_thread(tid: TidNumber) -> StarryResult<()> {
let task = get_task_by_number(tid)?;
let thr = task.as_thread();
thr.set_exit_request();
task.interrupt();
Ok(())
}
#[cfg(all(test, not(axtest)))]
fn decode_wait_status_rules_hold_for_test() -> bool {
use linux_raw_sys::general::{CLD_DUMPED, CLD_EXITED, CLD_KILLED};
let (code, status) = decode_wait_status(0);
assert!(code == CLD_EXITED as i32 && status == 0);
let (code, status) = decode_wait_status(0x0100); assert!(code == CLD_EXITED as i32 && status == 1);
let (code, status) = decode_wait_status(0xFF00); assert!(code == CLD_EXITED as i32 && status == 255);
let (code, status) = decode_wait_status(9); assert!(code == CLD_KILLED as i32 && status == 9);
let (code, status) = decode_wait_status(11); assert!(code == CLD_KILLED as i32 && status == 11);
let (code, status) = decode_wait_status(0x89); assert!(code == CLD_DUMPED as i32 && status == 9);
true
}
#[cfg(all(test, not(axtest)))]
mod tests {
#[test]
fn decode_wait_status_rules_hold() {
assert!(super::decode_wait_status_rules_hold_for_test());
}
}