mod exec_identity;
mod parked;
use std::cmp::Ordering;
use std::collections::BTreeMap;
use std::collections::BTreeSet;
use std::collections::HashMap;
use std::collections::HashSet;
use std::collections::btree_map::Entry;
use std::fmt::Debug;
use std::fs;
use std::fs::File;
use std::io::Write;
use std::num::NonZeroUsize;
use std::os::fd::FromRawFd;
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::AtomicBool;
use std::sync::atomic::AtomicU16;
use std::sync::atomic::Ordering::SeqCst;
use std::task::Poll;
use std::time::SystemTime;
use anyhow::bail;
use chrono::DateTime;
use chrono::Utc;
use detcore_model::procfs::mount_ids_are_ordered_subset;
use detcore_model::summary::RunSummary;
use detcore_model::summary::TimesliceStats;
pub(crate) use exec_identity::reconnect_exec;
pub(crate) use exec_identity::retire_exec;
use nix::sys::signal;
use nix::sys::signal::Signal;
use nix::unistd::Pid;
pub(crate) use parked::parked_wait_request;
pub(crate) use parked::polled_read_request;
pub(crate) use parked::signal_dequeued;
use reverie::GlobalRPC;
use reverie::GlobalTool;
use reverie::Guest;
use reverie::Tid;
use reverie::syscalls::AddrMut;
use reverie::syscalls::CloneFlags;
use reverie::syscalls::MemoryAccess;
use reverie::syscalls::Sysno;
use serde::Deserialize;
use serde::Serialize;
use tracing::debug;
use tracing::error;
use tracing::info;
use tracing::trace;
use tracing::warn;
use crate::config::Config;
use crate::consts::ROOT_DETPID;
use crate::ivar::Ivar;
use crate::preemptions::PreemptionReader;
use crate::preemptions::ThreadHistory;
use crate::record_or_replay::RecordOrReplay;
use crate::resources::ChaosEpochTransition;
use crate::resources::Permission;
use crate::resources::ResourceID;
use crate::resources::Resources;
use crate::scheduler::AdmitIntent;
use crate::scheduler::AdmitSide;
use crate::scheduler::ConsumeResult;
use crate::scheduler::DEFAULT_PRIORITY;
use crate::scheduler::ExecReconnect;
use crate::scheduler::MaybePrintStack;
use crate::scheduler::Priority;
use crate::scheduler::SchedResponse;
use crate::scheduler::SchedValue;
use crate::scheduler::Scheduler;
use crate::scheduler::ThreadNextTurn;
use crate::scheduler::entropy_to_priority;
use crate::scheduler::parked::*;
use crate::scheduler::real_timer::DequeueAck;
use crate::scheduler::real_timer::ItimerSnapshot;
use crate::scheduler::real_timer::TimerFailure;
use crate::scheduler::runqueue::FIRST_PRIORITY;
use crate::scheduler::runqueue::LAST_PRIORITY;
use crate::scheduler::runqueue::REPLAY_DEFERRED_PRIORITY;
use crate::scheduler::runqueue::REPLAY_FOREGROUND_PRIORITY;
use crate::scheduler::runqueue::is_ordinary_priority;
use crate::scheduler::sched_loop;
use crate::scheduler::sched_loop_external;
use crate::tool_local::Detcore;
use crate::tool_local::ExecFdBlockingOverrides;
use crate::tool_local::RobustListWake;
use crate::types::*;
pub(crate) async fn yield_once() {
let mut yielded = false;
std::future::poll_fn(|context| {
if yielded {
Poll::Ready(())
} else {
yielded = true;
context.waker().wake_by_ref();
Poll::Pending
}
})
.await;
}
#[derive(Debug)]
struct InodePool {
inodes: HashMap<RawInode, DetInode>,
detinodes_info: HashMap<DetInode, DetInodeInfo>,
next_inode: u64,
}
#[derive(Debug)]
struct DetInodeInfo {
raw: RawInode,
mtime: Option<LogicalTime>,
}
pub const CANONICAL_FILE_MTIME_SECONDS: [i64; 2] = [0, 1];
#[derive(PartialEq, Debug, Eq, Clone, Copy, Serialize, Deserialize)]
pub enum ObservedMtime {
Unobserved,
HostSpecific,
Canonical(LogicalTime),
}
impl ObservedMtime {
pub fn from_host_mtime(tv_sec: i64, tv_nsec: i64) -> Self {
match u64::try_from(tv_sec) {
Ok(secs) if tv_nsec == 0 && CANONICAL_FILE_MTIME_SECONDS.contains(&tv_sec) => {
ObservedMtime::Canonical(LogicalTime::from_secs(secs))
}
_ => ObservedMtime::HostSpecific,
}
}
fn first_seen_mtime(self, epoch: LogicalTime) -> Option<LogicalTime> {
match self {
ObservedMtime::Unobserved => None,
ObservedMtime::HostSpecific => Some(epoch),
ObservedMtime::Canonical(mtime) => Some(mtime),
}
}
}
struct ChildRegistration {
parent_dettid: DetTid,
parent_detpid: DetPid,
child_dettid: DetTid,
child_tid_addr: usize,
flags: Option<CloneFlags>,
exit_signal: libc::c_int,
physical_ids: Option<(i32, i32)>,
maybe_priority: Option<Priority>,
parent_is_kernel_blocked: bool,
}
#[derive(Clone, Debug, PartialEq, Eq)]
struct PendingExecState {
caller: DetTid,
process: DetPid,
mm: MmId,
fd_blocking: ExecFdBlockingOverrides,
}
pub struct BackendFailureCleanup {
pub scheduler: Result<(), tokio::task::JoinError>,
pub preemption_recording: Result<(), String>,
}
#[derive(Clone, Copy)]
struct RpcIncarnation {
dettid: DetTid,
mm: MmId,
}
impl Default for InodePool {
fn default() -> Self {
InodePool::new()
}
}
impl InodePool {
fn new() -> Self {
InodePool {
inodes: HashMap::new(),
detinodes_info: HashMap::new(),
next_inode: 1,
}
}
fn add_inode(
&mut self,
raw_inode: RawInode,
observed: ObservedMtime,
epoch: LogicalTime,
) -> (DetInode, LogicalTime) {
let dino = match self.inodes.get(&raw_inode) {
Some(dino) => *dino,
None => {
let new = DetInode::mint(self.next_inode);
self.next_inode += 1;
assert!(self.inodes.insert(raw_inode, new).is_none());
let prev = self.detinodes_info.insert(
new,
DetInodeInfo {
raw: raw_inode,
mtime: None,
},
);
assert!(prev.is_none()); new
}
};
let info = self
.detinodes_info
.get_mut(&dino)
.expect("Internal invariant broken, det_ino missing entry");
if info.mtime.is_none() {
info.mtime = observed.first_seen_mtime(epoch);
}
(dino, info.mtime.unwrap_or(epoch))
}
fn remove_inode(&mut self, det_inode: DetInode) {
if let Some(info) = self.detinodes_info.remove(&det_inode) {
self.inodes.remove(&info.raw);
}
}
}
#[derive(Debug)]
struct DevicePool {
devices: HashMap<u64, u64>,
next_device: u64,
}
#[derive(Debug)]
enum MountIdPool {
Uninitialized,
Invalid,
Ready {
mount_ids: BTreeMap<u64, u64>,
mountinfo_order: Vec<u64>,
allow_visible_subsets: bool,
unlisted_order: Vec<u64>,
next_mount_id: u64,
},
}
pub struct MountIdentityProvenance {
pub mountinfo_order: Vec<u64>,
pub unlisted_order: Vec<u64>,
}
impl MountIdPool {
fn from_config(mount_ids: &[u64], captured: bool, unlisted_ids: &[u64]) -> Self {
if !captured {
if !mount_ids.is_empty() || !unlisted_ids.is_empty() {
return Self::Invalid;
}
return Self::Uninitialized;
}
Self::from_orders(mount_ids, unlisted_ids, true).unwrap_or(Self::Invalid)
}
fn from_orders(
mount_ids: &[u64],
unlisted_ids: &[u64],
allow_visible_subsets: bool,
) -> Option<Self> {
let mut seen = BTreeSet::new();
if !mount_ids.iter().all(|raw| seen.insert(*raw))
|| !unlisted_ids
.iter()
.all(|raw| *raw != 0 && seen.insert(*raw))
{
return None;
}
let mut mappings = BTreeMap::new();
for (index, raw) in mount_ids.iter().chain(unlisted_ids).enumerate() {
mappings.insert(*raw, u64::try_from(index).ok()?.checked_add(1)?);
}
let next_mount_id = u64::try_from(mappings.len()).ok()?.checked_add(1)?;
Some(Self::Ready {
mount_ids: mappings,
mountinfo_order: mount_ids.to_vec(),
allow_visible_subsets,
unlisted_order: unlisted_ids.to_vec(),
next_mount_id,
})
}
fn validate_mountinfo_order(&mut self, mountinfo_order: &[u64]) -> bool {
if matches!(self, Self::Uninitialized) {
*self = Self::from_orders(mountinfo_order, &[], false).unwrap_or(Self::Invalid);
}
let Self::Ready {
mountinfo_order: expected,
allow_visible_subsets,
..
} = self
else {
return false;
};
if *allow_visible_subsets {
mount_ids_are_ordered_subset(mountinfo_order, expected)
} else {
mountinfo_order == expected
}
}
fn determinize(&mut self, raw_mount_id: u64, mountinfo_order: Option<&[u64]>) -> Option<u64> {
if raw_mount_id == 0 {
return Some(0);
}
if let Some(order) = mountinfo_order {
if !self.validate_mountinfo_order(order) {
return None;
}
} else if matches!(self, Self::Uninitialized) {
return None;
}
let Self::Ready {
mount_ids,
unlisted_order,
next_mount_id,
..
} = self
else {
return None;
};
if let Some(virtual_mount_id) = mount_ids.get(&raw_mount_id) {
return Some(*virtual_mount_id);
}
let virtual_mount_id = *next_mount_id;
*next_mount_id = next_mount_id.checked_add(1)?;
mount_ids.insert(raw_mount_id, virtual_mount_id);
unlisted_order.push(raw_mount_id);
Some(virtual_mount_id)
}
fn provenance(&self) -> Result<Option<MountIdentityProvenance>, &'static str> {
match self {
Self::Uninitialized => Ok(None),
Self::Invalid => Err("mount identity provenance is invalid"),
Self::Ready {
mountinfo_order,
unlisted_order,
..
} => Ok(Some(MountIdentityProvenance {
mountinfo_order: mountinfo_order.clone(),
unlisted_order: unlisted_order.clone(),
})),
}
}
}
impl Default for DevicePool {
fn default() -> Self {
DevicePool::new()
}
}
impl DevicePool {
fn new() -> Self {
DevicePool {
devices: HashMap::new(),
next_device: 1,
}
}
fn determinize(&mut self, raw_device: u64) -> u64 {
match self.devices.get(&raw_device) {
Some(dev) => *dev,
None => {
let new = self.next_device;
self.next_device += 1;
self.devices.insert(raw_device, new);
new
}
}
}
}
#[derive(Debug)]
pub struct GlobalState {
sched: Arc<Mutex<Scheduler>>,
inodes: Arc<Mutex<InodePool>>,
devices: Arc<Mutex<DevicePool>>,
mount_ids: Mutex<MountIdPool>,
next_port: AtomicU16,
used_ports: Mutex<HashSet<u16>>,
unsupported_syscalls: Mutex<BTreeSet<String>>,
unsupported_syscall_report_fd: Option<Mutex<File>>,
open_file_to_port: Mutex<HashMap<OpenFileId, u16>>,
port_start_range: AtomicU16,
port_end_range: AtomicU16,
past_first_execve: AtomicBool,
pending_exec_states: Mutex<BTreeMap<DetPid, PendingExecState>>,
completed_exec_transfers: Mutex<BTreeMap<DetPid, exec_identity::ExecTransferReceipt>>,
post_exec_fd_blocking: Mutex<BTreeMap<DetTid, ExecFdBlockingOverrides>>,
sched_handle: Option<tokio::task::JoinHandle<()>>,
global_time: Arc<Mutex<GlobalTime>>,
cfg: Config,
preemptions_to_replay: Option<PreemptionReader>,
realtime_start: SystemTime,
}
impl Default for GlobalState {
fn default() -> Self {
panic!("Detcore GlobalState Default impl should not be called");
}
}
impl Drop for GlobalState {
fn drop(&mut self) {
if let Some(message) =
format_unsupported_syscall_warning(&self.unsupported_syscalls.lock().unwrap())
{
warn!("{}", message);
}
info!("detcore shut down, destroying global state");
}
}
impl GlobalState {
async fn lock_rpc_scheduler(
&self,
consuming_cleanup: bool,
) -> std::sync::MutexGuard<'_, Scheduler> {
std::future::poll_fn(|_| {
let sched = self.sched.lock().unwrap();
if !consuming_cleanup && sched.backend_failed() {
Poll::Pending
} else {
Poll::Ready(sched)
}
})
.await
}
pub fn mount_identity_provenance(
&self,
) -> Result<Option<MountIdentityProvenance>, &'static str> {
self.mount_ids.lock().unwrap().provenance()
}
fn initialize(cfg: &Config, spawn_scheduler: bool) -> Self {
let sched = Arc::new(Mutex::new(Scheduler::new(cfg)));
let global_time = Arc::new(Mutex::new(GlobalTime::new(cfg)));
let handle = if cfg.sequentialize_threads && spawn_scheduler {
info!("[scheduler] daemon task starting up, waiting for guest thread start..");
Some(tokio::spawn(sched_loop(sched.clone(), global_time.clone())))
} else {
None
};
let preemptions_to_replay: Option<PreemptionReader> = cfg
.replay_preemptions_from
.as_ref()
.map(|path| PreemptionReader::new(path));
let range = Self::read_port_range();
let unsupported_syscall_report_fd = cfg.unsupported_syscall_report_fd.and_then(|fd| {
let duplicate = unsafe { libc::fcntl(fd, libc::F_DUPFD_CLOEXEC, fd) };
if duplicate == -1 {
warn!(
"failed to duplicate unsupported-syscall report fd {fd}: {}",
std::io::Error::last_os_error()
);
None
} else {
Some(Mutex::new(unsafe { File::from_raw_fd(duplicate) }))
}
});
Self {
sched,
next_port: AtomicU16::new(range[0]),
used_ports: Mutex::new(HashSet::new()),
unsupported_syscalls: Mutex::new(BTreeSet::new()),
unsupported_syscall_report_fd,
port_start_range: AtomicU16::new(range[0]),
port_end_range: AtomicU16::new(range[1]),
open_file_to_port: Mutex::new(HashMap::new()),
past_first_execve: AtomicBool::new(false),
pending_exec_states: Mutex::new(BTreeMap::new()),
completed_exec_transfers: Mutex::new(BTreeMap::new()),
post_exec_fd_blocking: Mutex::new(BTreeMap::new()),
inodes: Arc::new(Mutex::new(InodePool::new())),
devices: Arc::new(Mutex::new(DevicePool::new())),
mount_ids: Mutex::new(MountIdPool::from_config(
&cfg.mountinfo_mount_ids,
cfg.mountinfo_mount_ids_captured,
&cfg.fdinfo_unlisted_mount_ids,
)),
sched_handle: handle,
cfg: cfg.clone(),
realtime_start: SystemTime::now(),
global_time,
preemptions_to_replay,
}
}
pub fn init_for_external_scheduler(cfg: &Config) -> Self {
assert!(
cfg.sequentialize_threads,
"an external scheduler is only meaningful when threads are sequentialized"
);
Self::initialize(cfg, false)
}
pub fn run_external_scheduler(
&self,
observer: Arc<dyn Fn(&'static str) + Send + Sync>,
) -> impl Future<Output = ()> + use<> {
info!("[scheduler] daemon task starting up, waiting for guest thread start..");
sched_loop_external(self.sched.clone(), self.global_time.clone(), observer)
}
pub fn complete_physical_process_exit(&self, raw_pid: i32) {
let detpid = DetPid::from_raw(raw_pid);
self.pending_exec_states.lock().unwrap().remove(&detpid);
self.completed_exec_transfers
.lock()
.unwrap()
.remove(&detpid);
self.post_exec_fd_blocking.lock().unwrap().remove(&detpid);
if self
.sched
.lock()
.unwrap()
.complete_physical_process_exit(detpid)
{
trace!(
"[detcore, dpid {}] backend completed final physical process exit",
detpid
);
}
}
pub fn release_all_physical_process_exits(&self) {
self.pending_exec_states.lock().unwrap().clear();
self.completed_exec_transfers.lock().unwrap().clear();
self.post_exec_fd_blocking.lock().unwrap().clear();
let released = self
.sched
.lock()
.unwrap()
.release_all_physical_process_exits();
if released != 0 {
trace!("released {released} final physical process-exit barrier(s)");
}
}
pub fn force_shutdown_with_error(&self) {
let start = std::time::Instant::now();
let sched = loop {
if start.elapsed().as_millis() > 1000 {
eprintln!(
"Could not acquire scheduler lock during forced shutdown (timeout)... proceeding anyway."
);
return;
}
match self.sched.try_lock() {
Ok(guard) => {
break guard;
}
Err(std::sync::TryLockError::WouldBlock) => {
std::thread::yield_now();
continue;
}
Err(e) => {
eprintln!(
"Could not acquire scheduler lock during forced shutdown ({})... proceeding anyway.",
e
);
return;
}
}
};
info!("Scheduler state at exit:\n{}", sched.full_summary());
}
pub async fn cancel_internal_scheduler(&mut self) {
if let Some(handle) = self.sched_handle.take() {
handle.abort();
match handle.await {
Ok(()) => {}
Err(error) if error.is_cancelled() => {}
Err(error) => panic!("cancelled scheduler task panicked: {error}"),
}
}
}
pub async fn clean_up_after_backend_failure(mut self) -> BackendFailureCleanup {
let scheduler = if let Some(handle) = self.sched_handle.take() {
handle.await
} else {
Ok(())
};
let writer = self
.sched
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.preemption_writer
.take();
let preemption_recording = if self.cfg.record_preemptions_to.is_some() {
writer.map_or(Ok(()), |writer| writer.flush())
} else {
drop(writer);
Ok(())
};
BackendFailureCleanup {
scheduler,
preemption_recording,
}
}
pub async fn clean_up(self, to_stderr: bool, print_summary_to_json_file: &Option<PathBuf>) {
self.clean_up_with_dispatch_stats(to_stderr, print_summary_to_json_file, None)
.await
}
pub async fn clean_up_with_dispatch_stats(
mut self,
to_stderr: bool,
print_summary_to_json_file: &Option<PathBuf>,
dispatch_stats: Option<reverie::DispatchStats>,
) {
if let Some(handle) = self.sched_handle.take() {
debug!("Global state cleanup, confirming scheduler has shut down...");
handle.await.expect("Global scheduler clean shutdown");
debug!("Global state cleanup, continuing...");
}
let banner =
" ------------------------------ hermit run report ------------------------------";
let recording_destination = self.cfg.record_preemptions_to.clone();
let (mut summary, info_reprio_descrip) = self.into_run_summary_for_log().unwrap();
summary.dispatch_stats = dispatch_stats;
if let Some(path) = print_summary_to_json_file {
let json = serde_json::to_string_pretty(&summary).unwrap();
fs::write(path, json + "\n").unwrap();
}
if to_stderr {
{
use std::io::Write;
let _ = write!(crate::util::RetryingStderr, "{}\n{}", banner, summary);
}
} else {
let rt = summary.realtime_elapsed.take();
log_run_summary(
banner,
&summary,
info_reprio_descrip.as_deref(),
recording_destination.as_deref(),
);
if let Some(x) = rt {
debug!("Nondeterministic realtime elapsed: {:?}", x);
}
}
}
#[cfg(test)]
fn into_run_summary(self) -> anyhow::Result<RunSummary> {
self.into_run_summary_for_log().map(|(summary, _)| summary)
}
fn into_run_summary_for_log(self) -> anyhow::Result<(RunSummary, Option<String>)> {
let (mut summary, info_reprio_descrip) = {
let mut sched = self.sched.lock().unwrap();
sched.generate_partial_run_summary_for_log(self.cfg.record_preemptions_to.as_ref())?
};
summary.realtime_elapsed = Some(self.realtime_start.elapsed()?);
if self.cfg.virtualize_time {
let final_time = self.global_time.lock().unwrap();
summary.virttime_final = final_time.as_nanos().as_nanos();
summary.virttime_elapsed = match final_time.elapsed_nanos() {
Some(elapsed) => elapsed.as_nanos(),
None => bail!(
"Internal invariant violated! Global time {} is before its epoch baseline",
final_time.as_nanos()
),
};
}
Ok((summary, info_reprio_descrip))
}
}
fn log_run_summary(
banner: &str,
summary: &RunSummary,
info_reprio_descrip: Option<&str>,
recording_destination: Option<&std::path::Path>,
) {
info!("\n{}\n{}", banner, summary.info(info_reprio_descrip));
debug!(
replayed_events = summary.schedevent_replayed,
?recording_destination,
"Run recording/replay bookkeeping"
);
}
#[reverie::global_tool]
impl GlobalTool for GlobalState {
type Config = Config;
type Request = (DetTime, MmId, GlobalRequest);
type Response = (Option<LogicalTime>, GlobalResponse);
async fn init_global_state(cfg: &Config) -> GlobalState {
GlobalState::initialize(cfg, true)
}
fn install_backend_signal_control(
&self,
control: Option<reverie::BackendSignalControl>,
) -> Result<reverie::BackendSignalControlMode, reverie::Error> {
self.sched.lock().unwrap().install_signal_control(control)
}
fn authorize_backend_signal_boundary(
&self,
task: reverie::SignalTaskIdentity,
) -> Result<Option<reverie::SignalDeliveryPermit>, reverie::Error> {
self.sched.lock().unwrap().authorize_signal_boundary(task)
}
async fn on_backend_signal_boundary(
&self,
receipt: reverie::SignalBoundaryReceipt,
) -> Result<(), reverie::Error> {
self.sched.lock().unwrap().consume_signal_boundary(receipt)
}
fn on_backend_process_retired(
&self,
event: reverie::BackendProcessRetirement,
) -> Result<(), reverie::Error> {
let (result, wakes) = {
let mut sched = self.sched.lock().unwrap();
let result = sched.consume_process_retirement(event);
(result, sched.take_signal_failure_wakes())
};
for wake in wakes {
let _ = wake.send(());
}
result
}
fn report_backend_failure(&self, event: reverie::BackendFailure) {
let (wake, deferred) = {
let mut sched = self.sched.lock().unwrap();
(
sched.report_backend_failure(event),
sched.take_signal_failure_wakes(),
)
};
for wake in deferred {
let _ = wake.send(());
}
if let Some(wake) = wake {
let _ = wake.send(());
}
}
async fn wait_for_backend_failure(&self) {
let wake = self.sched.lock().unwrap().backend_failure_waiter();
wake.await
.expect("GlobalState owns the failure sender until publication");
}
async fn on_backend_child_wait_event(
&self,
event: reverie::BackendChildWaitEvent,
) -> Result<(), reverie::Error> {
crate::scheduler::signal_control::ChildExitPublicationFuture::new(self.sched.clone(), event)
.await
}
async fn receive_rpc(&self, from: Tid, gr: Self::Request) -> Self::Response {
type R = GlobalResponse;
let dtid = DetTid::from_raw(from.into()); let (guest_time, request_mm, request) = gr;
let time_from_guest = guest_time.as_nanos();
match &request {
GlobalRequest::ReconnectExec { former, process } => {
return self
.recv_exec_transfer(from, guest_time, request_mm, *former, *process)
.await;
}
GlobalRequest::RetireExec { thread, signaled } => {
return self
.recv_retire_exec_transfer(
from,
guest_time,
request_mm,
thread.clone(),
*signaled,
)
.await;
}
_ => {}
}
if let GlobalRequest::SignalDequeued {
detpid,
identity,
dequeue,
} = &request
{
return self
.recv_signal_dequeued(dtid, request_mm, guest_time, *detpid, *identity, *dequeue)
.await;
}
if let GlobalRequest::CreateVforkChildThread(parent, process, child, _, flags, ..) =
&request
{
let mut sched = self.lock_rpc_scheduler(false).await;
if sched.transferred_exec_tid_requires_registration(dtid) {
if *child != dtid
|| !(flags.contains(CloneFlags::CLONE_VFORK)
|| (self.cfg.backend_serializes_fork_children
&& !flags.contains(CloneFlags::CLONE_THREAD)))
|| !sched.pending_vfork_registration_matches(
*parent,
*process,
*child,
request_mm,
flags.contains(CloneFlags::CLONE_VM),
)
{
return (None, R::ThreadExited);
}
sched.register_reused_transferred_exec_tid(dtid, request_mm);
}
}
if matches!(&request, GlobalRequest::StartNewThread(child, ..) if *child == dtid) {
loop {
{
let sched = self.lock_rpc_scheduler(false).await;
if !sched.transferred_exec_tid_requires_registration(dtid)
|| sched.thread_is_logically_killed(dtid)
{
break;
}
}
yield_once().await;
}
}
let is_deregister = matches!(&request, GlobalRequest::DeregisterThread(_));
let consuming_cleanup =
is_deregister || matches!(&request, GlobalRequest::RobustListWakes(_));
let (exec_reconnect, is_exec_caller_after_local_mm_swap) = {
let pending = self.pending_exec_states.lock().unwrap();
let reconnect = match &request {
GlobalRequest::CreateChildThread(child, process, _, None, _, _, _)
if *child == dtid && *child == *process =>
{
pending.get(process).cloned()
}
_ => None,
};
let is_exec_caller_after_local_mm_swap = pending.values().any(|state| {
state.caller == dtid && state.mm.for_exec(state.process) == request_mm
});
(reconnect, is_exec_caller_after_local_mm_swap)
};
let mut tombstoned_deregistration = None;
{
let sched = self.lock_rpc_scheduler(consuming_cleanup).await;
if exec_reconnect.is_none()
&& !is_exec_caller_after_local_mm_swap
&& !sched.rpc_incarnation_matches(dtid, request_mm)
{
debug!(
"[detcore, dtid {}] rejecting {:?} RPC from retired exec incarnation {:?}",
dtid, request, request_mm,
);
return if is_deregister {
(None, R::DeregisterThread(()))
} else {
(None, R::ThreadExited)
};
}
if let GlobalRequest::DeregisterThread(owner) = &request {
assert_eq!(
owner.dettid, dtid,
"deregistration must belong to its sender"
);
assert_eq!(owner.mm, request_mm, "deregistration must retain its MmId");
if !sched.thread_was_registered(dtid) && !sched.thread_is_logically_killed(dtid) {
assert!(
sched.backend_failed() || !owner.thread_start_entered,
"a started thread must have a scheduler registration before deregistration"
);
return (None, R::DeregisterThread(()));
}
}
let child = match &request {
GlobalRequest::CreateChildThread(child, ..)
| GlobalRequest::CreateVforkChildThread(_, _, child, ..) => Some(*child),
_ => None,
};
if sched.thread_is_logically_killed(dtid) && exec_reconnect.is_none() {
trace!(
"[detcore, dtid {}] rejecting RPC after permanent logical-thread removal",
dtid
);
if let GlobalRequest::DeregisterThread(deregistration) = &request {
tombstoned_deregistration = Some(deregistration.clone());
} else {
return (None, R::ThreadExited);
}
}
if child.is_some_and(|child| sched.thread_is_logically_killed(child))
&& exec_reconnect.is_none()
{
trace!(
"[detcore, dtid {}] rejecting registration that reuses a tombstoned child TID",
dtid
);
return (None, R::ThreadExited);
}
if let GlobalRequest::ResumeExec(process) = &request
&& (dtid != *process
|| sched.registered_process(dtid) != Some(*process)
|| !sched.next_turns.contains_key(&dtid))
{
return (None, R::ThreadExited);
}
self.acknowledge_exec_transfer(dtid, request_mm);
let is_thread_reconnect = matches!(
&request,
GlobalRequest::StartNewThread(child_dettid, ..) if *child_dettid == dtid
) && self.global_time.lock().unwrap().contains_thread(dtid);
if tombstoned_deregistration.is_none()
&& exec_reconnect.is_none()
&& !is_thread_reconnect
{
self.global_time.lock().unwrap().update_global_time(
dtid,
time_from_guest,
guest_time.inherited_nanos(),
);
}
}
if let Some(deregistration) = tombstoned_deregistration {
self.recv_deregister_thread(from, deregistration).await;
return (None, R::DeregisterThread(()));
}
#[allow(clippy::unit_arg)]
let resp = match request {
GlobalRequest::ReconnectExec { .. } | GlobalRequest::RetireExec { .. } => {
unreachable!("exec transfer handled before ordinary admission")
}
GlobalRequest::ResumeExec(process) => {
let mut resources = Resources::new(dtid);
resources.insert(ResourceID::MemAddrSpace(process), Permission::RW);
let (response, duration) = self
.recv_request_resources(from, process, resources, Some(request_mm))
.await;
match response {
SchedulerRpcResult::Continue(status) => R::ResumeExec(status, duration),
SchedulerRpcResult::ThreadExited => R::ThreadExited,
}
}
GlobalRequest::SignalDequeued { .. } => {
unreachable!("consuming path handled before ordinary cancellation")
}
GlobalRequest::ParkedRequest(rs, pid, capability) => {
let (response, _) = self
.recv_resources_with_origin(
from,
pid,
rs,
Some(request_mm),
RpcOrigin::DirectRequestResources,
capability,
)
.await;
match response {
SchedulerRpcResult::Continue(r) => R::ParkedRequest(r),
SchedulerRpcResult::ThreadExited => R::ThreadExited,
}
}
GlobalRequest::ResumeParkedRequest {
ticket,
current_site,
} => {
self.recv_resume_parked(from, request_mm, ticket, current_site)
.await
}
GlobalRequest::FinishParkedObservation {
wait,
lease,
site,
finish,
} => {
let ack = Ivar::new();
let intent = ControlIntent::Finish {
wait,
lease,
site,
finish,
ack: ack.clone(),
};
let posted = self
.sched
.lock()
.unwrap()
.post_control(dtid, request_mm, intent);
if let Err(error) = posted {
self.sched.lock().unwrap().fail_parked(dtid, error);
return (None, R::ThreadExited);
}
R::FinishParkedObservation(ack.await)
}
GlobalRequest::ParkedProtocolFailure(error) => {
self.sched.lock().unwrap().fail_parked(dtid, error);
return (None, R::ThreadExited);
}
GlobalRequest::RequestResources(rs, pid) => {
let (response, _endtime) = self
.recv_request_resources(from, pid, rs, Some(request_mm))
.await;
match response {
SchedulerRpcResult::Continue(response) => R::RequestResources(response),
SchedulerRpcResult::ThreadExited => R::ThreadExited,
}
}
GlobalRequest::ReleaseResources(rs) => {
R::ReleaseResources(self.recv_release_resources(from, rs).await)
}
GlobalRequest::ReleaseAllResources => {
R::ReleaseAllResources(self.recv_release_all_resources(from).await)
}
GlobalRequest::ReportUnsupportedSyscall(name) => {
let _sched = self.lock_rpc_scheduler(false).await;
let inserted = self
.unsupported_syscalls
.lock()
.unwrap()
.insert(name.clone());
if inserted
&& let Some(report) = &self.unsupported_syscall_report_fd
&& let Err(error) = writeln!(report.lock().unwrap(), "{name}")
{
warn!("failed to append unsupported-syscall report: {error}");
}
R::ReportUnsupportedSyscall(())
}
GlobalRequest::PrepareExec(process, mm, fd_blocking) => {
let mut sched = self.lock_rpc_scheduler(false).await;
if mm != request_mm {
return (None, R::ThreadExited);
}
if self.cfg.sequentialize_threads && !sched.prepare_exec_teardown(dtid, process, mm)
{
return (None, R::ThreadExited);
}
trace!(
"[detcore, dtid {}] preparing exec for process {} with mm {:?} and logically blocking descriptors {:?}",
dtid, process, mm, fd_blocking,
);
self.pending_exec_states.lock().unwrap().insert(
process,
PendingExecState {
caller: dtid,
process,
mm,
fd_blocking,
},
);
R::PrepareExec(())
}
GlobalRequest::CancelExec(process) => {
let mut sched = self.lock_rpc_scheduler(false).await;
let mut pending = self.pending_exec_states.lock().unwrap();
if pending
.get(&process)
.is_some_and(|state| state.caller == dtid)
{
let prepared = pending.remove(&process).unwrap();
sched.finish_exec_teardown(prepared.caller, process, prepared.mm, false);
}
R::CancelExec(())
}
GlobalRequest::MarkPastFirstExecve(detpid, signal_identity) => {
let mut sched = self.lock_rpc_scheduler(false).await;
let prepared = self.pending_exec_states.lock().unwrap().get(&dtid).cloned();
if self.cfg.kvm_shared_dequeue_timers {
let result = (|| {
let identity = signal_identity.ok_or(ProtocolFailure::Identity)?;
let pid = sched
.registered_process(dtid)
.ok_or(ProtocolFailure::Identity)?;
let mut pending = self.pending_exec_states.lock().unwrap();
if let Some(prepared) = pending.get(&pid) {
if prepared.caller != dtid || prepared.process != pid {
return Err(ProtocolFailure::Identity);
}
sched.complete_signal_exec(
pid,
dtid,
prepared.mm,
request_mm,
identity,
)?;
pending.remove(&pid);
} else {
sched
.real_timers
.validate_task(pid, dtid, request_mm, identity)?;
}
Ok(())
})();
if let Err(error) = result {
sched.fail_parked(dtid, error);
return (None, R::ThreadExited);
}
}
sched.blocked.timed_waiters.remove_posix_timers(detpid);
if let Some(prepared) = &prepared
&& self.cfg.sequentialize_threads
&& prepared.caller == dtid
&& prepared.process == dtid
&& prepared.mm.for_exec(dtid) == request_mm
{
sched.reconnect_after_exec(ExecReconnect {
caller: dtid,
new_leader: dtid,
detpid: dtid,
pre_exec_mm: prepared.mm,
post_exec_mm: request_mm,
child_tid_addr: 0,
reconnect_priority: None,
});
self.pending_exec_states.lock().unwrap().remove(&dtid);
}
self.past_first_execve.store(true, SeqCst);
let overrides = self
.post_exec_fd_blocking
.lock()
.unwrap()
.remove(&dtid)
.unwrap_or_default();
trace!(
"[detcore, dtid {}] restoring logically blocking descriptors after exec: {:?}",
dtid, overrides,
);
R::MarkPastFirstExecve(overrides)
}
GlobalRequest::CreateChildThread(
dettid,
parent_detpid,
ctid,
flags,
exit_signal,
physical_ids,
priority,
) => {
if let Some(prepared) = &exec_reconnect {
let mut sched = self.lock_rpc_scheduler(false).await;
let (pending, post_exec_mm) = {
let mut states = self.pending_exec_states.lock().unwrap();
let Some(pending) = states.remove(&parent_detpid) else {
return (None, R::ThreadExited);
};
assert_eq!(&pending, prepared);
let post_exec_mm = pending.mm.for_exec(pending.process);
(pending, post_exec_mm)
};
assert_eq!(pending.process, parent_detpid);
if let Some((physical_pid, physical_tid)) = physical_ids
&& let Err(open_error) = sched.register_physical_thread(
dettid,
post_exec_mm,
physical_pid,
physical_tid,
)
{
error!(
"[detcore, dtid {}] failed to register post-exec host process {} thread {}: {}",
dettid, physical_pid, physical_tid, open_error,
);
return (None, R::ThreadExited);
}
let retired = sched.reconnect_after_exec(ExecReconnect {
caller: pending.caller,
new_leader: dettid,
detpid: parent_detpid,
pre_exec_mm: pending.mm,
post_exec_mm,
child_tid_addr: ctid,
reconnect_priority: priority,
});
if pending.caller != dettid {
self.global_time
.lock()
.unwrap()
.reassign_thread(pending.caller, dettid);
}
if !pending.fd_blocking.is_empty() {
self.post_exec_fd_blocking
.lock()
.unwrap()
.insert(dettid, pending.fd_blocking);
}
debug!(
"[detcore, dtid {}] reconciled successful exec from caller {}; retired prior identities {:?}",
dtid, pending.caller, retired
);
R::CreateChildThread(Some(post_exec_mm))
} else {
match self
.recv_create_child_thread(
from,
request_mm,
ChildRegistration {
parent_dettid: DetTid::from_raw(from.into()),
parent_detpid,
child_dettid: dettid,
child_tid_addr: ctid,
flags,
exit_signal,
physical_ids,
maybe_priority: priority,
parent_is_kernel_blocked: false,
},
)
.await
{
SchedulerRpcResult::Continue(()) => R::CreateChildThread(None),
SchedulerRpcResult::ThreadExited => R::ThreadExited,
}
}
}
GlobalRequest::CreateVforkChildThread(
parent_dettid,
parent_detpid,
child_dettid,
ctid,
flags,
exit_signal,
priority,
) => match self
.recv_create_child_thread(
from,
request_mm,
ChildRegistration {
parent_dettid,
parent_detpid,
child_dettid,
child_tid_addr: ctid,
flags: Some(flags),
exit_signal,
physical_ids: None,
maybe_priority: priority,
parent_is_kernel_blocked: true,
},
)
.await
{
SchedulerRpcResult::Continue(()) => R::CreateChildThread(None),
SchedulerRpcResult::ThreadExited => R::ThreadExited,
},
GlobalRequest::StartNewThread(dettid, detpid, physical_ids, signal_identity) => {
match self
.recv_start_new_thread(
from,
dettid,
detpid,
request_mm,
physical_ids,
signal_identity,
)
.await
{
SchedulerRpcResult::Continue(history) => R::StartNewThread(history),
SchedulerRpcResult::ThreadExited => R::ThreadExited,
}
}
GlobalRequest::DeregisterThread(deregistration) => {
R::DeregisterThread(self.recv_deregister_thread(from, deregistration).await)
}
GlobalRequest::SetChildTidAddress(address) => {
let updated = self
.lock_rpc_scheduler(false)
.await
.set_child_tid_address(dtid, address);
if updated {
R::SetChildTidAddress(())
} else {
R::ThreadExited
}
}
GlobalRequest::FutexAction(dettid, action, futexid, init_read, mask) => R::FutexAction(
self.recv_futex_action(
RpcIncarnation {
dettid,
mm: request_mm,
},
action,
futexid,
init_read,
mask,
)
.await,
),
GlobalRequest::RobustListWakes(wakes) => {
R::RobustListWakes(self.recv_robust_list_wakes(wakes))
}
GlobalRequest::DeterminizeInode(ino, observed) => {
R::DeterminizeInode(self.recv_determinize_inode(from, ino, observed).await)
}
GlobalRequest::DeterminizeDevice(dev) => {
R::DeterminizeDevice(self.recv_determinize_device(from, dev).await)
}
GlobalRequest::DeterminizeMountId(raw_mount_id, fallback_order) => {
R::DeterminizeMountId(
self.recv_determinize_mount_id(from, raw_mount_id, fallback_order.as_deref())
.await,
)
}
GlobalRequest::ValidateMountIdOrder(mountinfo_order) => R::ValidateMountIdOrder(
self.recv_validate_mount_id_order(from, &mountinfo_order)
.await,
),
GlobalRequest::UnlinkInode(d_ino) => {
R::UnlinkInode(self.recv_unlink_inode(from, d_ino).await)
}
GlobalRequest::TouchFile(ino) => R::TouchFile(self.recv_touch_file(from, ino).await),
GlobalRequest::SetFileMtime(ino, mtime) => {
R::SetFileMtime(self.recv_set_file_mtime(from, ino, mtime).await)
}
GlobalRequest::GlobalTimeLowerBound => {
let ns = self.global_time.lock().unwrap().as_nanos();
R::GlobalTimeLowerBound(ns)
}
GlobalRequest::TraceSchedEvent(ev, detpid, command_bootstrap) => {
match self
.recv_trace_schedevent(ev, detpid, request_mm, command_bootstrap)
.await
{
SchedulerRpcResult::Continue(response) => R::TraceSchedEvent(response),
SchedulerRpcResult::ThreadExited => R::ThreadExited,
}
}
GlobalRequest::RegisterAlarm(dpid, dtid, duration, interval, sig) => {
let now = self.global_time.lock().unwrap().as_nanos();
match self
.recv_register_alarm(
dpid,
RpcIncarnation {
dettid: dtid,
mm: request_mm,
},
now,
duration,
interval,
sig,
)
.await
{
SchedulerRpcResult::Continue(remaining) => R::RegisterAlarm(remaining),
SchedulerRpcResult::ThreadExited => R::ThreadExited,
}
}
GlobalRequest::AlarmRemaining(dpid) => {
let now = self.global_time.lock().unwrap().as_nanos();
let mut sched = self.lock_rpc_scheduler(false).await;
match sched.itimer_snapshot(dpid, now) {
Ok(snapshot) => R::AlarmRemaining(snapshot),
Err(error) => {
sched.fail_parked(dtid, ProtocolFailure::Timer(error));
R::ThreadExited
}
}
}
GlobalRequest::RegisterPosixTimer(dpid, dtid, timer_id, deadline, interval, sig) => {
match self
.recv_register_posix_timer(
dpid,
RpcIncarnation {
dettid: dtid,
mm: request_mm,
},
timer_id,
deadline,
interval,
sig,
)
.await
{
SchedulerRpcResult::Continue(()) => R::RegisterPosixTimer(()),
SchedulerRpcResult::ThreadExited => R::ThreadExited,
}
}
GlobalRequest::ResolveKillTargets(dpid) => R::ResolveKillTargets(
self.lock_rpc_scheduler(false)
.await
.process_signal_targets(dpid),
),
GlobalRequest::NotifySignalPending(dettid, SigWrapper(signal), target_process) => {
let mut scheduler = self.lock_rpc_scheduler(false).await;
scheduler.notify_signal_pending(dettid, SigWrapper(signal));
if signal == libc::SIGKILL
&& let Some(detpid) = target_process
{
scheduler.note_process_sigkill(dettid, detpid);
}
R::NotifySignalPending(())
}
GlobalRequest::ThreadIsLive(dtid) => {
R::ThreadIsLive(self.lock_rpc_scheduler(false).await.thread_is_live(dtid))
}
GlobalRequest::ExactChildWaitState(parent, child) => R::ExactChildWaitState(
self.lock_rpc_scheduler(false)
.await
.exact_child_wait_state(parent, child),
),
GlobalRequest::ReadyChildWait(parent, selector) => {
let sched = self.lock_rpc_scheduler(false).await;
R::ReadyChildWait((
sched.ready_child_wait(parent, selector),
sched.has_child_wait_target(parent, selector),
))
}
GlobalRequest::ConsumeChildWait(parent, child) => R::ConsumeChildWait(
self.lock_rpc_scheduler(false)
.await
.consume_child_wait(parent, child),
),
GlobalRequest::ProcessGroup(process) => R::ProcessGroup(
self.lock_rpc_scheduler(false)
.await
.thread_tree
.process_group(process),
),
GlobalRequest::SetProcessGroup(process, group) => R::SetProcessGroup(
self.lock_rpc_scheduler(false)
.await
.thread_tree
.set_process_group(process, group),
),
GlobalRequest::CreateSession(process) => R::CreateSession(
self.lock_rpc_scheduler(false)
.await
.thread_tree
.create_session(process),
),
GlobalRequest::UnrecoverableShutdown => {
self.force_shutdown_with_error();
R::UnrecoverableShutdown(())
}
GlobalRequest::RequestPort(open_file_id) => {
let _sched = self.lock_rpc_scheduler(false).await;
let mut mut_used_ports = self.used_ports.lock().unwrap();
self.update_port_range();
let total_available =
self.port_end_range.load(SeqCst) - self.port_start_range.load(SeqCst);
let mut index = 0;
while (*mut_used_ports).contains(&self.next_port.load(SeqCst))
&& index < total_available
{
self.next_port.fetch_add(1, SeqCst);
if self.next_port.load(SeqCst) > self.port_end_range.load(SeqCst) {
self.next_port
.store(self.port_start_range.load(SeqCst), SeqCst);
}
index += 1;
}
if index == total_available {
R::PortFull
} else {
(*mut_used_ports).insert(self.next_port.load(SeqCst));
let mut open_file_to_port = self.open_file_to_port.lock().unwrap();
open_file_to_port.insert(open_file_id, self.next_port.load(SeqCst));
R::RequestPort(self.next_port.load(SeqCst))
}
}
GlobalRequest::AddUsedPort(port, open_file_id) => {
let _sched = self.lock_rpc_scheduler(false).await;
let mut used_ports = self.used_ports.lock().unwrap();
used_ports.insert(port);
let mut open_file_to_port = self.open_file_to_port.lock().unwrap();
open_file_to_port.insert(open_file_id, port);
R::AddUsedPort
}
GlobalRequest::ReleasePort(open_file_id) => {
let _sched = self.lock_rpc_scheduler(false).await;
let mut used_ports = self.used_ports.lock().unwrap();
let mut open_file_to_port = self.open_file_to_port.lock().unwrap();
let port = open_file_to_port.remove(&open_file_id);
if let Some(port) = port {
used_ports.remove(&port);
}
R::ReleasePort(port)
}
};
let sender_became_terminal =
if is_deregister || exec_reconnect.is_some() || is_exec_caller_after_local_mm_swap {
false
} else {
let sched = self.lock_rpc_scheduler(consuming_cleanup).await;
sched.thread_is_logically_killed(dtid)
|| !sched.rpc_incarnation_matches(dtid, request_mm)
};
if resp == R::ThreadExited || sender_became_terminal {
return (None, R::ThreadExited);
}
let time_from_sched = self.global_time.lock().unwrap().threads_time(dtid);
let time_update = match time_from_sched.cmp(&time_from_guest) {
Ordering::Equal => None,
Ordering::Less => {
panic!(
"internal error: thread time should never go down, only monotonically up: time in sched {}, thread local time was {}",
time_from_sched, time_from_guest
)
}
Ordering::Greater => Some(time_from_sched),
};
(time_update, resp)
}
}
impl GlobalState {
async fn recv_resources_with_origin(
&self,
from: Tid,
detpid: DetPid,
rs: Resources,
request_mm: Option<MmId>,
rpc: RpcOrigin,
capability: ControlCapability,
) -> (SchedulerRpcResult<ResourceReply>, Option<LogicalTime>) {
let dettid = DetTid::from_raw(from.into());
let resp2 = {
let mut sched = self.lock_rpc_scheduler(false).await;
if sched.thread_is_logically_killed(dettid)
|| request_mm.is_some_and(|mm| !sched.rpc_incarnation_matches(dettid, mm))
{
return (SchedulerRpcResult::ThreadExited, None);
}
let Some(nextturn) = sched.next_turns.get(&dettid).cloned() else {
panic!(
"Detcore internal error: no entry for dettid {} in next_turns during resource request.",
dettid
);
};
trace!(
"[detcore, dtid {}] ResourceRequest, filling request into {}",
&dettid, &nextturn.req
);
if let Some(mm) = request_mm
&& let Err(error) = sched.install_resource_origin(
dettid,
ResourceOrigin {
rpc,
mm,
control: capability,
},
)
{
sched.fail_parked(dettid, error);
return (SchedulerRpcResult::ThreadExited, None);
}
sched.request_put(&nextturn.req, rs.clone(), &self.global_time);
nextturn.resp
};
trace!(
"[detcore, dtid {}] waiting on {} for resources: {:?}",
dettid, &resp2, rs
);
let answer = resp2.get().await; self.finish_resource_response(from, detpid, rs, request_mm, answer)
.await
}
async fn finish_resource_response(
&self,
from: Tid,
detpid: DetPid,
rs: Resources,
request_mm: Option<MmId>,
answer: SchedResponse,
) -> (SchedulerRpcResult<ResourceReply>, Option<LogicalTime>) {
let dettid = DetTid::from_raw(from.as_raw());
let request_became_stale = {
let sched = self.lock_rpc_scheduler(false).await;
sched.thread_is_logically_killed(dettid)
|| request_mm.is_some_and(|mm| !sched.rpc_incarnation_matches(dettid, mm))
};
if request_became_stale {
trace!(
"[detcore, dtid {}] terminating pending request after logical removal",
dettid
);
return (SchedulerRpcResult::ThreadExited, None);
}
if let SchedResponse::ObserveSignal(control) = answer {
return (
SchedulerRpcResult::Continue(ResourceReply::ObserveSignal(control)),
None,
);
}
if let Some((true, process, mm)) = rs.exit_identity() {
info!(
"Scheduler authorized an exit-group scenario, from dettid {} / detpid {}",
dettid, detpid
);
{
let mut sched = self.lock_rpc_scheduler(false).await;
if sched.thread_is_logically_killed(dettid)
|| request_mm.is_some_and(|mm| !sched.rpc_incarnation_matches(dettid, mm))
{
return (SchedulerRpcResult::ThreadExited, None);
}
for tid in sched.thread_tree.my_thread_group(&dettid) {
if tid != dettid {
sched.logically_kill_thread(&tid, &process, mm);
}
}
}
}
match answer {
SchedResponse::ObserveSignal(_) => {
self.sched
.lock()
.unwrap()
.fail_parked(dettid, ProtocolFailure::UnexpectedControl);
(SchedulerRpcResult::ThreadExited, None)
}
SchedResponse::Go(Some(schedval)) => {
trace!(
"[dtid {}] resources granted, resuming normally: {:?}",
dettid, rs
);
let endtime_update = match schedval {
SchedValue::TimeOut => None,
SchedValue::Value(timeslice) => Some(LogicalTime::from_nanos(timeslice)),
};
(
SchedulerRpcResult::Continue(ResourceReply::Grant(ResumeStatus::Normal)),
endtime_update,
)
}
SchedResponse::Go(None) => {
trace!(
"[dtid {}] resources granted but no timeslice specified",
dettid,
);
(
SchedulerRpcResult::Continue(ResourceReply::Grant(ResumeStatus::Normal)),
None,
)
}
SchedResponse::Signaled(signal) => {
trace!(
"[dtid {}] resources granted but interrupted by signal",
dettid,
);
(
SchedulerRpcResult::Continue(ResourceReply::Grant(ResumeStatus::Signaled(
signal,
))),
None,
)
}
}
}
async fn recv_release_resources(&self, from: Tid, rs: Resources) {
trace!("[detcore] Resources released to pid {}: {:?}", from, rs);
}
async fn recv_release_all_resources(&self, from: Tid) {
trace!("[detcore] All resources held by pid {} released", from);
}
async fn recv_create_child_thread(
&self,
rpc_sender: Tid,
request_mm: MmId,
registration: ChildRegistration,
) -> SchedulerRpcResult<()> {
let ChildRegistration {
parent_dettid,
parent_detpid,
child_dettid,
child_tid_addr: ctid,
flags,
exit_signal,
physical_ids,
maybe_priority,
parent_is_kernel_blocked,
} = registration;
let initial_priority = if let Some(pr) = &self.preemptions_to_replay {
assert!(maybe_priority.is_none());
let prio = pr
.thread_initial_priority(&child_dettid)
.unwrap_or_else(|| {
warn!(
"Child thread {} not found in preemption history to replay",
child_dettid
);
DEFAULT_PRIORITY
});
if !is_ordinary_priority(prio) {
panic!(
"Read a bad initial_prority from file: {}\nFull file: {}",
prio,
pr.load_all(),
);
}
prio
} else {
let prio = maybe_priority.expect(
"create_child_thread must take an initial priority unless replaying preemptions",
);
if !is_ordinary_priority(prio) {
panic!(
"recv_create_child_thread received a bad prority argument : {}",
prio,
);
}
prio
};
{
let mut sched = self.lock_rpc_scheduler(false).await;
let sender = DetTid::from_raw(rpc_sender.into());
if sched.thread_is_logically_killed(sender)
|| !sched.rpc_incarnation_matches(sender, request_mm)
|| sched.thread_is_logically_killed(child_dettid)
{
return SchedulerRpcResult::ThreadExited;
}
if parent_is_kernel_blocked && self.cfg.sequentialize_threads {
sched.complete_vfork_registration(parent_dettid, child_dettid);
}
let child_mm = MmId::for_clone(
request_mm,
child_dettid,
flags.is_some_and(|flags| flags.contains(CloneFlags::CLONE_VM)),
);
sched.register_reused_transferred_exec_tid(child_dettid, child_mm);
let _entry = sched
.next_turns
.entry(child_dettid)
.or_insert_with(|| ThreadNextTurn {
dettid: child_dettid,
child_tid_addr: ctid,
req: Ivar::new(),
resp: Ivar::new(),
protocol: Default::default(),
});
{
let is_group_leader = if let Some(f) = flags {
!f.contains(CloneFlags::CLONE_THREAD)
} else {
true };
sched.thread_tree.add_child_with_wait_metadata(
parent_dettid,
child_dettid,
is_group_leader,
flags.is_some_and(|flags| flags.contains(CloneFlags::CLONE_PARENT)),
exit_signal,
);
}
if let Some((physical_pid, physical_tid)) = physical_ids {
let child_is_thread =
flags.is_some_and(|flags| flags.contains(CloneFlags::CLONE_THREAD));
let child_detpid = if child_is_thread {
parent_detpid
} else {
child_dettid
};
if let Err(open_error) = sched.register_physical_thread(
child_dettid,
child_mm,
physical_pid,
physical_tid,
) {
error!(
"[detcore, dtid {}] cannot bind host process {} thread {} during child registration: {}",
child_dettid, physical_pid, physical_tid, open_error,
);
sched.logically_kill_thread(&child_dettid, &child_detpid, child_mm);
return SchedulerRpcResult::ThreadExited;
}
}
sched.hb_note_spawn(child_dettid);
if self.cfg.replay_schedule_from.is_none() {
let old_prio = sched.priorities.insert(child_dettid, initial_priority);
assert!(old_prio.is_none());
} else {
if let std::collections::btree_map::Entry::Vacant(entry) =
sched.priorities.entry(child_dettid)
{
assert_eq!(parent_detpid, ROOT_DETPID);
entry.insert(initial_priority);
}
}
if let Some(pr) = &mut sched.preemption_writer {
pr.register_thread(child_dettid, initial_priority);
}
let intent = if self.cfg.sequentialize_threads && !parent_is_kernel_blocked {
AdmitIntent::PostFork(self.cfg.runs_post_fork)
} else {
AdmitIntent::Fixed(AdmitSide::Back)
};
sched.admit_to_run_queue(child_dettid, intent);
debug!(
"[detcore] CreateChildThread with dtid {}: admit child via {:?}.",
child_dettid, intent,
);
sched.started_up.try_put(());
}
if self.cfg.sequentialize_threads && !parent_is_kernel_blocked {
let mut rs = Resources::new(parent_detpid);
rs.insert(
ResourceID::ParentContinue {
parent: parent_dettid,
child: child_dettid,
},
Permission::W,
);
if matches!(
self.recv_grant_resources(
rpc_sender,
parent_detpid,
rs,
Some(request_mm),
RpcOrigin::ParentContinue
)
.await
.0,
SchedulerRpcResult::ThreadExited
) {
return SchedulerRpcResult::ThreadExited;
}
}
SchedulerRpcResult::Continue(())
}
async fn recv_start_new_thread(
&self,
from: Tid,
dettid: DetTid,
detpid: DetPid,
request_mm: MmId,
physical_ids: Option<(i32, i32)>,
signal_identity: Option<reverie::SignalTaskIdentity>,
) -> SchedulerRpcResult<Option<ThreadHistory>> {
let mut tries: u64 = 0;
let response_ivar = loop {
yield_once().await;
let mut sched = self.lock_rpc_scheduler(false).await;
if sched.thread_is_logically_killed(dettid)
|| !sched.rpc_incarnation_matches(dettid, request_mm)
{
return SchedulerRpcResult::ThreadExited;
}
if self.cfg.backend_requires_thread_directed_process_signals && physical_ids.is_none() {
error!(
"[detcore, dtid {}] backend requires a host thread ID at StartNewThread",
dettid,
);
sched.logically_kill_thread(&dettid, &detpid, request_mm);
return SchedulerRpcResult::ThreadExited;
}
let rsrcs = {
let mut s = HashMap::new();
s.insert(ResourceID::MemAddrSpace(detpid), Permission::RW); Resources {
tid: dettid,
resources: s,
poll_attempt: 0,
fyi: String::new(),
signal_interrupt_errno: None,
backend_runtime_bootstrap: false,
}
};
let nextturn = match sched.next_turns.entry(dettid) {
Entry::Vacant(_entry) => {
if tries == 0 {
trace!(
"[detcore, dtid {}] thread showed up early, no queue entry yet. Waiting...",
dettid
);
}
tries += 1;
continue;
}
Entry::Occupied(entry) => {
trace!(
"[detcore, dtid {}] handling StartNewThread rpc. Found next_turns entry (after {} tries)",
from, tries
);
entry.get().clone()
}
};
if let Some((physical_pid, physical_tid)) = physical_ids
&& let Err(open_error) =
sched.register_physical_thread(dettid, request_mm, physical_pid, physical_tid)
{
error!(
"[detcore, dtid {}] cannot bind host process {} thread {} for exact signal delivery: {}",
dettid, physical_pid, physical_tid, open_error,
);
sched.logically_kill_thread(&dettid, &detpid, request_mm);
return SchedulerRpcResult::ThreadExited;
}
if self.cfg.kvm_shared_dequeue_timers {
let binding = signal_identity
.ok_or(TimerFailure::Identity)
.and_then(|identity| {
if from.as_raw() != dettid.as_raw()
|| sched.registered_process(dettid) != Some(detpid)
{
return Err(TimerFailure::Identity);
}
sched.real_timers.bind(detpid, dettid, request_mm, identity)
});
if let Err(error) = binding {
sched.fail_parked(dettid, ProtocolFailure::Timer(error));
return SchedulerRpcResult::ThreadExited;
}
}
if let Err(error) = sched.install_resource_origin(
dettid,
ResourceOrigin {
rpc: RpcOrigin::ThreadStart,
mm: request_mm,
control: ControlCapability::None,
},
) {
sched.fail_parked(dettid, error);
return SchedulerRpcResult::ThreadExited;
}
sched.request_put(&nextturn.req, rsrcs, &self.global_time);
break nextturn.resp;
};
debug!(
"[detcore, dtid {}] New thread will now wait for response on {}...",
&dettid, &response_ivar
);
let answer = response_ivar.get().await;
if matches!(answer, SchedResponse::ObserveSignal(_)) {
self.sched
.lock()
.unwrap()
.fail_parked(dettid, ProtocolFailure::UnexpectedControl);
return SchedulerRpcResult::ThreadExited;
}
let request_became_stale = {
let sched = self.lock_rpc_scheduler(false).await;
sched.thread_is_logically_killed(dettid)
|| !sched.rpc_incarnation_matches(dettid, request_mm)
};
if request_became_stale {
return SchedulerRpcResult::ThreadExited;
}
info!(
"[detcore, dtid {}] New thread given go-ahead to proceed via {}",
&dettid, &response_ivar
);
if let Some(pr) = &self.preemptions_to_replay {
let (history, old_prio) = {
let mut sched = self.lock_rpc_scheduler(false).await;
if sched.thread_is_logically_killed(dettid)
|| !sched.rpc_incarnation_matches(dettid, request_mm)
{
return SchedulerRpcResult::ThreadExited;
}
let history = pr.extract_thread_record(&dettid).unwrap_or_else(|| {
warn!(
"Replaying preemptions, but no record found for thread {}",
dettid
);
ThreadHistory::new()
});
let old_prio = sched.priorities.insert(dettid, history.initial_priority());
(history, old_prio)
};
debug!(
"[replay-preemption] Enqueing new thread at priority {:?} (changed from {:?})",
history.initial_priority(),
old_prio,
);
SchedulerRpcResult::Continue(Some(history))
} else {
SchedulerRpcResult::Continue(None)
}
}
async fn recv_deregister_thread(&self, _from: Tid, mut deregistration: ThreadDeregistration) {
let mut sched = self.sched.lock().unwrap();
let ThreadDeregistration {
dettid, detpid, mm, ..
} = &deregistration;
let mut pending = self.pending_exec_states.lock().unwrap();
if pending.get(detpid).is_some_and(|state| {
state.caller == *dettid && (*mm == state.mm || *mm == state.mm.for_exec(state.process))
}) {
let prepared = pending.remove(detpid).unwrap();
sched.finish_exec_teardown(prepared.caller, *detpid, prepared.mm, false);
deregistration.mm = prepared.mm;
}
drop(pending);
assert!(self.cfg.sequentialize_threads);
self.account_deregistered_thread(&mut sched, deregistration);
}
fn account_deregistered_thread(
&self,
sched: &mut Scheduler,
deregistration: ThreadDeregistration,
) {
let ThreadDeregistration {
dettid,
detpid,
mm,
thread_start_entered: _,
timeslice_stats,
syscall_count,
chaos_epochs,
} = deregistration;
if !sched.rpc_incarnation_matches(dettid, mm) {
debug!(
"[detcore, dtid {}] ignoring deregistration from retired exec incarnation {:?}",
dettid, mm,
);
return;
}
self.post_exec_fd_blocking.lock().unwrap().remove(&dettid);
if !sched.note_deregistration_accounted(dettid) {
trace!(
"[detcore, dtid {}] acknowledging already-accounted deregistration",
dettid
);
return;
}
if let Some(writer) = &mut sched.preemption_writer {
for transition in chaos_epochs {
writer.insert_chaos_epoch(dettid, transition);
}
}
sched.record_timeslice_stats(dettid, timeslice_stats);
sched.record_syscall_count(dettid, syscall_count);
if !sched.defer_exec_sibling_retirement(dettid, detpid, mm)
&& !sched.thread_is_logically_killed(dettid)
{
sched.logically_kill_thread(&dettid, &detpid, mm);
}
trace!(
"[detcore, dtid {}] thread deregistered, removed from sched structures.",
dettid
);
}
async fn recv_futex_action(
&self,
caller: RpcIncarnation,
action: FutexAction,
futexid: FutexID,
init_read: i32,
mask: u32,
) -> Option<SchedValue> {
let RpcIncarnation { dettid, mm } = caller;
trace!("[detcore, dtid {}] Futex action: {:?}", &dettid, action);
let response_iv = {
let mut sched = self.lock_rpc_scheduler(false).await;
if sched.thread_is_logically_killed(dettid)
|| !sched.rpc_incarnation_matches(dettid, mm)
{
return Some(SchedValue::Value(nix::errno::Errno::EINTR as u64));
}
let Some(resp_iv) = sched
.next_turns
.get(&dettid)
.map(|nextturn| nextturn.resp.clone())
else {
trace!(
"[detcore, dtid {}] ignoring futex action after logical thread removal",
dettid
);
return Some(SchedValue::Value(nix::errno::Errno::EINTR as u64));
};
match action {
FutexAction::WaitRequest(maybe_timeout) => {
if sched.child_tid_was_cleared(futexid, init_read) {
trace!(
"[detcore, dtid {}] late wait on cleared child-TID futex {:?}",
dettid, futexid
);
return Some(SchedValue::Value(0));
}
if let Err(error) = sched.install_resource_origin(
dettid,
ResourceOrigin {
rpc: RpcOrigin::FutexAction,
mm,
control: ControlCapability::None,
},
) {
sched.fail_parked(dettid, error);
return Some(SchedValue::Value(nix::errno::Errno::EINTR as u64));
}
sched.sleep_futex_waiter(&dettid, futexid, maybe_timeout, mask);
}
FutexAction::WaitFinished => {
return None;
}
FutexAction::WakeRequest(num_threads) => {
let num = sched.wake_futex_waiters(dettid, futexid, num_threads, mask);
return Some(SchedValue::Value(num));
}
FutexAction::WakeFinished(_num_threads) => {
return None;
}
}
assert!(sched.run_queue.remove_tid(dettid));
resp_iv
};
match response_iv.get().await {
SchedResponse::Go(answer) => {
trace!(
"[detcore, dtid {}] Unblocked from futex_wait! ({})",
&dettid, &response_iv
);
answer
}
SchedResponse::Signaled(_) => Some(SchedValue::Value(nix::errno::Errno::EINTR as u64)),
SchedResponse::ObserveSignal(_) => {
self.sched
.lock()
.unwrap()
.fail_parked(dettid, ProtocolFailure::UnexpectedControl);
self.wait_for_backend_failure().await;
futures::future::pending().await
}
}
}
fn recv_robust_list_wakes(&self, wakes: Vec<(DetTid, FutexID)>) -> Vec<u64> {
let mut sched = self.sched.lock().unwrap();
sched.wake_futex_waiters_after_exit(&wakes)
}
fn epoch_logical_time(&self) -> LogicalTime {
let nanos = self
.cfg
.epoch
.timestamp_nanos_opt()
.expect("epoch cannot be represented in a timestamp with nanosecond precision")
as u64;
LogicalTime::from_nanos(nanos)
}
async fn recv_determinize_inode(
&self,
from: Tid,
ino: RawInode,
observed: ObservedMtime,
) -> (DetInode, LogicalTime) {
let _sched = self.lock_rpc_scheduler(false).await;
let epoch = self.epoch_logical_time();
let (dino, ns) = self.inodes.lock().unwrap().add_inode(ino, observed, epoch);
trace!(
"[detcore, dtid {}] resolved (raw) inode {:?} to {:?}, mtime {}",
from, ino, dino, ns
);
(dino, ns)
}
async fn recv_determinize_device(&self, from: Tid, raw_device: u64) -> u64 {
let _sched = self.lock_rpc_scheduler(false).await;
let det_device = self.devices.lock().unwrap().determinize(raw_device);
trace!(
"[detcore, dtid {}] resolved (raw) device {} to {}",
from, raw_device, det_device
);
det_device
}
async fn recv_determinize_mount_id(
&self,
from: Tid,
raw_mount_id: u64,
mountinfo_order: Option<&[u64]>,
) -> Option<u64> {
let _sched = self.lock_rpc_scheduler(false).await;
let virtual_mount_id = self
.mount_ids
.lock()
.unwrap()
.determinize(raw_mount_id, mountinfo_order);
trace!(
"[detcore, dtid {}] resolved fdinfo mount ID {} to {:?}",
from, raw_mount_id, virtual_mount_id
);
virtual_mount_id
}
async fn recv_validate_mount_id_order(&self, from: Tid, mountinfo_order: &[u64]) -> bool {
let _sched = self.lock_rpc_scheduler(false).await;
let valid = self
.mount_ids
.lock()
.unwrap()
.validate_mountinfo_order(mountinfo_order);
trace!(
"[detcore, dtid {}] validated mountinfo identity order: {}",
from, valid
);
valid
}
async fn recv_unlink_inode(&self, from: Tid, d_ino: DetInode) {
let _sched = self.lock_rpc_scheduler(false).await;
trace!("[detcore, dtid {}] unlink (det) inode {:?}", from, d_ino);
self.inodes.lock().unwrap().remove_inode(d_ino);
}
async fn recv_touch_file(&self, from: Tid, ino: RawInode) {
let _sched = self.lock_rpc_scheduler(false).await;
let mtime = if self.cfg.virtualize_time {
self.global_time.lock().unwrap().as_nanos()
} else {
let dt: DateTime<Utc> = Utc::now();
let nanos = dt.timestamp_nanos_opt().expect(
"current time cannot be represented in a timestamp with nanosecond precision",
) as u64;
LogicalTime::from_nanos(nanos)
};
trace!(
"[dtid {}] bumping mtime on file (rawinode {:?}) to {}",
from, ino, mtime,
);
self.set_inode_mtime(ino, mtime);
}
async fn recv_set_file_mtime(&self, from: Tid, ino: RawInode, mtime: LogicalTime) {
let _sched = self.lock_rpc_scheduler(false).await;
trace!(
"[dtid {}] setting mtime on file (rawinode {:?}) to {}",
from, ino, mtime,
);
self.set_inode_mtime(ino, mtime);
}
fn set_inode_mtime(&self, ino: RawInode, mtime: LogicalTime) {
let epoch = self.epoch_logical_time();
let mut mg = self.inodes.lock().unwrap();
let (dino, _) = mg.add_inode(ino, ObservedMtime::Unobserved, epoch);
let info = mg
.detinodes_info
.get_mut(&dino)
.expect("Invariant violation: det inode missing from map.");
info.mtime = Some(mtime);
}
async fn recv_trace_schedevent(
&self,
ev: SchedEvent,
detpid: DetPid,
request_mm: MmId,
command_bootstrap: bool,
) -> SchedulerRpcResult<TraceSchedEventResponse> {
let ev = {
let sched = self.lock_rpc_scheduler(false).await;
if !sched.rpc_incarnation_matches(ev.dettid, request_mm) {
return SchedulerRpcResult::ThreadExited;
}
let ev = {
if self.past_first_execve.load(SeqCst) {
ev
} else {
info!(
"Warning: erasing rip of pre-execve sched event! {:?}",
SchedEventForLog {
event: &ev,
command_bootstrap
}
);
SchedEvent {
end_rip: None,
start_rip: None,
..ev
}
}
};
if ev.op == Op::Syscall(Sysno::execve, SyscallPhase::Prehook) {
self.past_first_execve.store(true, SeqCst);
}
ev
};
let result = if self.cfg.replay_schedule_from.is_some() {
let (consumed, print_stack2) = {
let mut sched = self.lock_rpc_scheduler(false).await;
if sched.thread_is_logically_killed(ev.dettid)
|| !sched.rpc_incarnation_matches(ev.dettid, request_mm)
{
return SchedulerRpcResult::ThreadExited;
}
let consumed = sched.consume_schedevent(&ev);
let print_stack2 = if self.cfg.record_preemptions {
sched.record_event(&ev)
} else {
None
};
(consumed, print_stack2)
};
let ConsumeResult {
keep_running,
print_stack,
event_ix: _,
timeslice_remaining: mut end_of_timeslice,
} = consumed;
trace!(
"keep_running :{}, end_of_timeslice: {:?}",
keep_running, end_of_timeslice
);
if !keep_running {
trace!(
"[detcore, dtid {}] Thread yielding to follow replay schedule",
&ev.dettid,
);
let tid = reverie::Tid::from(ev.dettid.as_raw()); let mut rsrcs = Resources::new(ev.dettid);
rsrcs.insert(ResourceID::TraceReplay, Permission::RW);
let (response, timeslice) = self
.recv_grant_resources(
tid,
detpid,
rsrcs,
Some(request_mm),
RpcOrigin::TraceSchedEvent,
)
.await;
if response == SchedulerRpcResult::ThreadExited {
return SchedulerRpcResult::ThreadExited;
}
end_of_timeslice = timeslice;
trace!(
"[detcore, dtid {}] Thread reactivated after yielding for replay schedule",
&ev.dettid,
);
}
TraceSchedEventResponse {
print_stack_strace: print_stack.or(print_stack2),
timeslice: end_of_timeslice,
}
} else {
let print_stack_strace = {
let mut sched = self.lock_rpc_scheduler(false).await;
if sched.thread_is_logically_killed(ev.dettid)
|| !sched.rpc_incarnation_matches(ev.dettid, request_mm)
{
return SchedulerRpcResult::ThreadExited;
}
if self.cfg.record_preemptions {
sched.record_event(&ev)
} else {
None
}
};
TraceSchedEventResponse {
print_stack_strace,
timeslice: None,
}
};
if result.print_stack_strace.is_some()
&& let Some(sig) = &self.cfg.stacktrace_signal
{
let _sched = self.lock_rpc_scheduler(false).await;
trace!(
"[dtid {}] signaling thread with {} at the point of stack trace printing.",
ev.dettid, sig.0
);
let tid = Pid::from_raw(ev.dettid.as_raw());
match sig.signal() {
Some(named) => signal::kill(tid, named).unwrap(),
None => {
let rc = unsafe { libc::kill(tid.as_raw(), sig.raw()) };
assert_eq!(rc, 0, "raw kill of signal {} failed", sig.raw());
}
}
}
SchedulerRpcResult::Continue(result)
}
fn read_port_range() -> Vec<u16> {
let contents = fs::read_to_string("/proc/sys/net/ipv4/ip_local_port_range")
.expect("File should be present");
let range: Vec<u16> = contents
.split_whitespace()
.filter_map(|number| number.parse().ok())
.collect();
range
}
fn update_port_range(&self) {
let range = Self::read_port_range();
self.port_start_range.store(range[0], SeqCst);
self.port_end_range.store(range[1], SeqCst);
}
async fn recv_register_alarm(
&self,
detpid: DetPid,
caller: RpcIncarnation,
now: LogicalTime,
duration: LogicalTime,
interval: LogicalTime,
sig: SigWrapper,
) -> SchedulerRpcResult<(LogicalTime, LogicalTime)> {
let RpcIncarnation { dettid, mm } = caller;
let mut sched = self.lock_rpc_scheduler(false).await;
if sched.thread_is_logically_killed(dettid) || !sched.rpc_incarnation_matches(dettid, mm) {
return SchedulerRpcResult::ThreadExited;
}
match sched.replace_real_timer(detpid, dettid, now, duration, interval, alarm_signal(sig)) {
Ok(old) => SchedulerRpcResult::Continue(old),
Err(error) => {
sched.fail_parked(dettid, ProtocolFailure::Timer(error));
SchedulerRpcResult::ThreadExited
}
}
}
async fn recv_register_posix_timer(
&self,
detpid: DetPid,
caller: RpcIncarnation,
timer_id: i32,
deadline: Option<LogicalTime>,
interval: LogicalTime,
sig: SigWrapper,
) -> SchedulerRpcResult<()> {
let RpcIncarnation { dettid, mm } = caller;
let mut sched = self.lock_rpc_scheduler(false).await;
if sched.thread_is_logically_killed(dettid) || !sched.rpc_incarnation_matches(dettid, mm) {
return SchedulerRpcResult::ThreadExited;
}
sched.register_posix_timer(
detpid,
dettid,
timer_id,
deadline,
interval,
alarm_signal(sig),
);
SchedulerRpcResult::Continue(())
}
}
#[derive(PartialEq, Debug, Eq, Clone, Serialize, Deserialize)]
pub struct ThreadDeregistration {
pub(crate) dettid: DetTid,
pub(crate) detpid: DetPid,
pub(crate) mm: MmId,
pub(crate) thread_start_entered: bool,
pub(crate) timeslice_stats: TimesliceStats,
pub(crate) syscall_count: u64,
pub(crate) chaos_epochs: Vec<ChaosEpochTransition>,
}
#[derive(PartialEq, Debug, Eq, Clone, Serialize, Deserialize)]
#[allow(clippy::enum_variant_names)]
pub enum GlobalRequest {
RequestResources(Resources, DetPid),
ParkedRequest(Resources, DetPid, ControlCapability),
ResumeParkedRequest {
ticket: ResumeTicket,
current_site: reverie::CallbackSignalSite,
},
FinishParkedObservation {
wait: ContinuationId,
lease: reverie::ParkedObservationLease,
site: reverie::CallbackSignalSite,
finish: ObservationFinish,
},
ParkedProtocolFailure(ProtocolFailure),
SignalDequeued {
detpid: DetPid,
identity: reverie::SignalTaskIdentity,
dequeue: reverie::SignalDequeue,
},
ReleaseResources(Resources),
ReleaseAllResources,
ReportUnsupportedSyscall(String),
PrepareExec(DetPid, MmId, ExecFdBlockingOverrides),
CancelExec(DetPid),
ReconnectExec {
former: DetTid,
process: DetPid,
},
ResumeExec(DetPid),
RetireExec {
thread: ThreadDeregistration,
signaled: bool,
},
MarkPastFirstExecve(DetPid, Option<reverie::SignalTaskIdentity>),
CreateChildThread(
DetTid,
DetPid,
usize,
Option<CloneFlags>,
libc::c_int,
Option<(i32, i32)>,
Option<Priority>,
),
CreateVforkChildThread(
DetTid,
DetPid,
DetTid,
usize,
CloneFlags,
libc::c_int,
Option<Priority>,
),
StartNewThread(
DetTid,
DetPid,
Option<(i32, i32)>,
Option<reverie::SignalTaskIdentity>,
),
DeregisterThread(ThreadDeregistration),
SetChildTidAddress(usize),
FutexAction(DetTid, FutexAction, FutexID, i32, u32),
DeterminizeInode(RawInode, ObservedMtime),
DeterminizeDevice(u64),
DeterminizeMountId(u64, Option<Vec<u64>>),
ValidateMountIdOrder(Vec<u64>),
UnlinkInode(DetInode),
TouchFile(RawInode),
SetFileMtime(RawInode, LogicalTime),
GlobalTimeLowerBound,
TraceSchedEvent(SchedEvent, DetPid, bool),
RegisterAlarm(DetPid, DetTid, LogicalTime, LogicalTime, SigWrapper),
RegisterPosixTimer(
DetPid,
DetTid,
i32,
Option<LogicalTime>,
LogicalTime,
SigWrapper,
),
AlarmRemaining(DetPid),
ReadyChildWait(DetPid, ChildWaitSpec),
ConsumeChildWait(DetPid, DetPid),
ProcessGroup(DetPid),
SetProcessGroup(DetPid, DetPid),
CreateSession(DetPid),
ResolveKillTargets(DetPid),
NotifySignalPending(DetTid, SigWrapper, Option<DetPid>),
ThreadIsLive(DetTid),
ExactChildWaitState(DetPid, DetPid),
UnrecoverableShutdown,
RequestPort(OpenFileId),
AddUsedPort(u16, OpenFileId),
ReleasePort(OpenFileId),
RobustListWakes(Vec<(DetTid, FutexID)>),
}
#[allow(missing_docs, clippy::unit_arg)]
#[derive(PartialEq, Debug, Eq, Clone, Serialize, Deserialize)]
pub enum GlobalResponse {
ThreadExited,
RequestResources(ResumeStatus),
ParkedRequest(ResourceReply),
ResumeParkedRequest(ResourceReply),
FinishParkedObservation(Result<FinishAck, ProtocolFailure>),
SignalDequeued {
ack: Result<DequeueAck, TimerFailure>,
terminal: bool,
},
ReleaseResources(()),
ReleaseAllResources(()),
ReportUnsupportedSyscall(()),
PrepareExec(()),
CancelExec(()),
ReconnectExec(bool),
ResumeExec(ResumeStatus, Option<LogicalTime>),
RetireExec(bool),
MarkPastFirstExecve(ExecFdBlockingOverrides),
CreateChildThread(Option<MmId>),
StartNewThread(Option<ThreadHistory>),
DeregisterThread(()),
SetChildTidAddress(()),
FutexAction(Option<SchedValue>),
DeterminizeInode((DetInode, LogicalTime)),
DeterminizeDevice(u64),
DeterminizeMountId(Option<u64>),
ValidateMountIdOrder(bool),
UnlinkInode(()),
TouchFile(()),
SetFileMtime(()),
GlobalTimeLowerBound(LogicalTime),
TraceSchedEvent(TraceSchedEventResponse),
RegisterAlarm((LogicalTime, LogicalTime)),
RegisterPosixTimer(()),
ReadyChildWait((Option<DetPid>, bool)),
ConsumeChildWait(bool),
ProcessGroup(Option<DetPid>),
SetProcessGroup(bool),
CreateSession(bool),
AlarmRemaining(ItimerSnapshot),
ResolveKillTargets(Vec<DetTid>),
NotifySignalPending(()),
ThreadIsLive(bool),
ExactChildWaitState(ExactChildWaitState),
UnrecoverableShutdown(()),
RequestPort(u16),
AddUsedPort,
ReleasePort(Option<u16>),
PortFull,
RobustListWakes(Vec<u64>),
}
pub fn format_unsupported_syscall_warning(syscalls: &BTreeSet<String>) -> Option<String> {
if syscalls.is_empty() {
None
} else {
Some(format!(
"syscalls {} used but not yet supported",
syscalls.iter().cloned().collect::<Vec<_>>().join(",")
))
}
}
pub async fn prepare_exec<G, T>(guest: &mut G, mm: MmId, fd_blocking: ExecFdBlockingOverrides)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let detpid = guest.thread_state().detpid.expect("detpid unset");
let (_, response) =
send_and_update_time(guest, GlobalRequest::PrepareExec(detpid, mm, fd_blocking)).await;
assert_eq!(response, GlobalResponse::PrepareExec(()));
}
pub async fn cancel_exec<G, T>(guest: &mut G)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let detpid = guest.thread_state().detpid.expect("detpid unset");
let (_, response) = send_and_update_time(guest, GlobalRequest::CancelExec(detpid)).await;
assert_eq!(response, GlobalResponse::CancelExec(()));
}
pub async fn mark_past_first_execve<G, T>(guest: &mut G)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let signal_identity = guest
.config()
.kvm_shared_dequeue_timers
.then(|| guest.signal_task_identity())
.flatten();
let detpid = guest.thread_state().detpid.expect("detpid unset");
let (_, response) = send_and_update_time(
guest,
GlobalRequest::MarkPastFirstExecve(detpid, signal_identity),
)
.await;
let overrides = match response {
GlobalResponse::MarkPastFirstExecve(overrides) => overrides,
_ => unreachable!(),
};
if guest.config().kvm_shared_dequeue_timers {
guest.thread_state_mut().signal_task_identity = signal_identity;
}
if !overrides.is_empty() {
let dettid = guest.thread_state().dettid;
let metadata = Arc::clone(&guest.thread_state().file_metadata);
metadata
.lock()
.unwrap()
.apply_exec_blocking_overrides(dettid, overrides);
}
}
pub async fn report_unsupported_syscall<G, T>(guest: &mut G, sysno: Sysno)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let (_, response) = send_and_update_time(
guest,
GlobalRequest::ReportUnsupportedSyscall(sysno.to_string()),
)
.await;
assert_eq!(response, GlobalResponse::ReportUnsupportedSyscall(()));
}
pub(crate) async fn set_child_tid_address<G, T>(guest: &mut G, address: usize)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let (_, response) =
send_and_update_time(guest, GlobalRequest::SetChildTidAddress(address)).await;
assert_eq!(response, GlobalResponse::SetChildTidAddress(()));
}
pub async fn send_and_update_time<G, T>(
guest: &mut G,
request: GlobalRequest,
) -> (Option<LogicalTime>, GlobalResponse)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
assert_eq!(
DetTid::from_raw(guest.tid().as_raw()),
guest.thread_state().dettid,
"replacement image must reconnect before ordinary RPCs"
);
let mytime = guest.thread_state().thread_logical_time.clone();
let mm = guest.thread_state().mm_id;
let resp = guest.send_rpc((mytime, mm, request)).await;
if resp.1 == GlobalResponse::ThreadExited {
let dettid = guest.thread_state().dettid;
trace!(
"[detcore, dtid {}] exiting after terminal scheduler cancellation",
dettid
);
guest.retire_current_thread().await
}
if let Some(time) = resp.0 {
guest
.thread_state_mut()
.thread_logical_time
.advance_to(time);
}
resp
}
#[derive(PartialEq, Debug, Eq, Clone, Serialize, Deserialize)]
pub enum ResumeStatus {
Normal,
Signaled(Option<Vec<SigWrapper>>),
}
#[derive(PartialEq, Debug, Eq, Clone)]
enum SchedulerRpcResult<T> {
Continue(T),
ThreadExited,
}
pub async fn resource_request<G, T>(guest: &mut G, mut r: Resources) -> ResumeStatus
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
r.backend_runtime_bootstrap |=
guest.thread_state().in_uncharged_bootstrap_syscall && guest.is_backend_runtime_bootstrap();
if guest.config().sequentialize_threads {
if let Some(lease) = guest.signal_observation_lease() {
if let Some(site) = guest.parked_signal_site() {
return parked::capable_resource_request(
guest,
r,
ControlCapability::PublishOnly { lease, site },
)
.await;
}
parked::terminate_protocol(guest, ProtocolFailure::Identity).await;
}
let dettid = guest.thread_state().dettid;
let detpid = guest.thread_state().detpid.expect("detpid unset");
trace!(
"[detcore, dtid {}] BLOCKING on resource_request rpc... {:?}",
&dettid, r
);
let resp =
send_and_update_time(guest, GlobalRequest::RequestResources(r.clone(), detpid)).await;
match resp.1 {
GlobalResponse::RequestResources(x) => {
trace!(
"[detcore, dtid {}] UNBLOCKED, acquired resources: {:?}",
&dettid, r
);
x
}
_ => unreachable!(),
}
} else {
ResumeStatus::Normal
}
}
pub async fn resource_release_all<G, T>(guest: &mut G)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
if guest.config().sequentialize_threads {
let resp = send_and_update_time(guest, GlobalRequest::ReleaseAllResources).await;
match resp.1 {
GlobalResponse::ReleaseAllResources(x) => x,
_ => unreachable!(),
}
}
}
pub async fn thread_start_request<G, T>(
cfg: &Config,
guest: &mut G,
detpid: DetPid,
) -> Option<ThreadHistory>
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let dettid = guest.thread_state().dettid;
if cfg.sequentialize_threads {
trace!("[detcore, dtid {}] new thread BLOCKING on rpc...", &dettid);
let physical_ids = guest
.thread_state()
.physical_tid
.map(|tid| (guest.pid().as_raw(), tid));
let signal_identity = guest.signal_task_identity();
guest.thread_state_mut().signal_task_identity = signal_identity;
let resp = send_and_update_time(
guest,
GlobalRequest::StartNewThread(dettid, detpid, physical_ids, signal_identity),
)
.await;
match resp.1 {
GlobalResponse::StartNewThread(preempts) => {
trace!("[detcore, dtid {}] new thread UNBLOCKED (post-rpc)", dettid);
preempts
}
_ => unreachable!(),
}
} else {
None
}
}
pub(crate) fn child_tid_clear_address(flags: CloneFlags, address: usize) -> usize {
if flags.contains(CloneFlags::CLONE_CHILD_CLEARTID) {
address
} else {
0
}
}
pub async fn create_child_thread<G, T>(
guest: &mut G,
child_dettid: DetTid,
ctid: usize,
flags: Option<CloneFlags>,
exit_signal: libc::c_int,
physical_ids: Option<(i32, i32)>,
) -> Option<MmId>
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let starting_priority = if guest.config().replay_preemptions_from.is_some() {
None
} else if guest.config().replay_schedule_from.is_some() {
if child_dettid <= DetTid::from_raw(3) {
Some(REPLAY_FOREGROUND_PRIORITY)
} else {
Some(REPLAY_DEFERRED_PRIORITY)
}
} else if guest.config().chaos {
let entropy = guest
.thread_state_mut()
.chaos_prng_next_u64("child_priority");
if guest.config().chaos_target_races {
if entropy.is_multiple_of(2) {
Some(FIRST_PRIORITY)
} else {
Some(LAST_PRIORITY)
}
} else {
Some(entropy_to_priority(entropy))
}
} else {
Some(DEFAULT_PRIORITY)
};
let detpid = guest.thread_state().detpid.expect("detpid unset");
let resp = send_and_update_time(
guest,
GlobalRequest::CreateChildThread(
child_dettid,
detpid,
ctid,
flags,
exit_signal,
physical_ids,
starting_priority,
),
)
.await;
match resp.1 {
GlobalResponse::CreateChildThread(x) => x,
_ => unreachable!(),
}
}
pub async fn create_vfork_child_thread<G, T>(
guest: &mut G,
child_dettid: DetTid,
vfork: crate::tool_local::PendingVfork,
) where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let starting_priority = if guest.config().replay_preemptions_from.is_some() {
None
} else if guest.config().replay_schedule_from.is_some() {
Some(if child_dettid <= DetTid::from_raw(3) {
REPLAY_FOREGROUND_PRIORITY
} else {
REPLAY_DEFERRED_PRIORITY
})
} else if guest.config().chaos {
Some(entropy_to_priority(vfork.child_priority_entropy.expect(
"vfork child priority entropy missing in chaos mode",
)))
} else {
Some(DEFAULT_PRIORITY - 1)
};
let resp = send_and_update_time(
guest,
GlobalRequest::CreateVforkChildThread(
vfork.parent_dettid,
vfork.parent_detpid,
child_dettid,
vfork.child_tid_addr,
vfork.flags,
vfork.exit_signal,
starting_priority,
),
)
.await;
match resp.1 {
GlobalResponse::CreateChildThread(_) => (),
_ => unreachable!(),
}
}
pub(crate) async fn deregister_thread<R>(
threads_time: DetTime,
cfg: &Config,
reverie: &R,
thread: ThreadDeregistration,
) where
R: GlobalRPC<GlobalState>,
{
if cfg.sequentialize_threads {
let mm = thread.mm;
let resp = reverie
.send_rpc((threads_time, mm, GlobalRequest::DeregisterThread(thread)))
.await;
match resp.1 {
GlobalResponse::DeregisterThread(x) => x,
_ => unreachable!(),
}
}
}
pub(crate) async fn acknowledge_robust_list_exit_time<R>(
threads_time: DetTime,
reverie: &R,
mm: MmId,
) -> bool
where
R: GlobalRPC<GlobalState>,
{
let response = reverie
.send_rpc((threads_time, mm, GlobalRequest::RobustListWakes(Vec::new())))
.await;
match response.1 {
GlobalResponse::RobustListWakes(counts) => {
assert!(
counts.is_empty(),
"an empty exit-clock acknowledgement woke a waiter"
);
true
}
GlobalResponse::ThreadExited => false,
_ => unreachable!(),
}
}
pub(crate) async fn robust_list_wakes_after_exit<R>(
threads_time: DetTime,
reverie: &R,
mm: MmId,
wakes: Vec<(DetTid, RobustListWake)>,
) -> Vec<u64>
where
R: GlobalRPC<GlobalState>,
{
if wakes.is_empty() {
return Vec::new();
}
let response = reverie
.send_rpc((
threads_time,
mm,
GlobalRequest::RobustListWakes(
wakes
.into_iter()
.map(|(owner, wake)| (owner, wake.futex))
.collect(),
),
))
.await;
match response.1 {
GlobalResponse::RobustListWakes(counts) => counts,
_ => unreachable!(),
}
}
#[derive(PartialEq, Debug, Eq, Clone, Copy, Serialize, Deserialize)]
pub enum FutexAction {
WaitRequest(Option<LogicalTime>),
WaitFinished,
WakeRequest(i32),
WakeFinished(i32),
}
pub async fn futex_action<G, T>(
guest: &mut G,
futex_action: FutexAction,
futexid: &FutexID,
init_read: i32,
mask: u32,
) -> Option<SchedValue>
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
assert!(guest.config().sequentialize_threads);
let dettid = guest.thread_state().dettid;
let req = GlobalRequest::FutexAction(dettid, futex_action, *futexid, init_read, mask);
trace!(
"BLOCKING on futex_action: sending request to scheduler: {:?}",
req
);
let resp = send_and_update_time(guest, req.clone()).await;
match resp.1 {
GlobalResponse::FutexAction(answer) => {
trace!("UNBLOCKING after futex_action. Request was: {:?}", req);
answer
}
_ => unreachable!(),
}
}
pub async fn determinize_inode<G, T>(guest: &mut G, inode: RawInode) -> (DetInode, LogicalTime)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
determinize_inode_observing_mtime(guest, inode, ObservedMtime::Unobserved).await
}
pub async fn determinize_inode_observing_mtime<G, T>(
guest: &mut G,
inode: RawInode,
observed: ObservedMtime,
) -> (DetInode, LogicalTime)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let resp = send_and_update_time(guest, GlobalRequest::DeterminizeInode(inode, observed)).await;
match resp.1 {
GlobalResponse::DeterminizeInode(x) => x,
_ => unreachable!(),
}
}
pub async fn determinize_device<G, T>(guest: &mut G, raw_device: u64) -> u64
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let resp = send_and_update_time(guest, GlobalRequest::DeterminizeDevice(raw_device)).await;
match resp.1 {
GlobalResponse::DeterminizeDevice(x) => x,
_ => unreachable!(),
}
}
pub async fn determinize_mount_id<G, T>(
guest: &mut G,
raw_mount_id: u64,
mountinfo_order: Option<Vec<u64>>,
) -> Option<u64>
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let resp = send_and_update_time(
guest,
GlobalRequest::DeterminizeMountId(raw_mount_id, mountinfo_order),
)
.await;
match resp.1 {
GlobalResponse::DeterminizeMountId(value) => value,
_ => unreachable!(),
}
}
pub async fn validate_mountinfo_identity_order<G, T>(
guest: &mut G,
mountinfo_order: Vec<u64>,
) -> bool
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let resp =
send_and_update_time(guest, GlobalRequest::ValidateMountIdOrder(mountinfo_order)).await;
match resp.1 {
GlobalResponse::ValidateMountIdOrder(valid) => valid,
_ => unreachable!(),
}
}
#[allow(unused)]
pub async fn unlink_inode<G, T>(guest: &mut G, d_ino: DetInode)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let resp = send_and_update_time(guest, GlobalRequest::UnlinkInode(d_ino)).await;
match resp.1 {
GlobalResponse::UnlinkInode(x) => x,
_ => unreachable!(),
}
}
pub async fn touch_file<G, T>(guest: &mut G, inode: RawInode)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let resp = send_and_update_time(guest, GlobalRequest::TouchFile(inode)).await;
match resp.1 {
GlobalResponse::TouchFile(x) => x,
_ => unreachable!(),
}
}
pub async fn set_file_mtime<G, T>(guest: &mut G, inode: RawInode, mtime: LogicalTime)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let resp = send_and_update_time(guest, GlobalRequest::SetFileMtime(inode, mtime)).await;
match resp.1 {
GlobalResponse::SetFileMtime(x) => x,
_ => unreachable!(),
}
}
pub async fn global_time_lower_bound<G, T>(guest: &mut G) -> LogicalTime
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let resp = send_and_update_time(guest, GlobalRequest::GlobalTimeLowerBound).await;
match resp.1 {
GlobalResponse::GlobalTimeLowerBound(x) => x,
_ => unreachable!(),
}
}
pub async fn thread_observe_time<G, T>(guest: &mut G) -> LogicalTime
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
global_time_lower_bound(guest).await
}
fn write_backtrace<G, T>(guest: &mut G, m_path: Option<&PathBuf>)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
if let Some(backtrace) = guest.backtrace() {
if let Some(path) = m_path {
let file = File::create(path).expect("Failed to open preemption stacktrace log file");
serde_json::to_writer(file, &backtrace.force_pretty()).unwrap();
} else {
eprintln!("{}", backtrace.force_pretty());
}
} else {
warn!("Could not read backtrace!");
}
}
#[derive(PartialEq, Debug, Eq, Clone, Serialize, Deserialize)]
pub struct TraceSchedEventResponse {
print_stack_strace: MaybePrintStack,
timeslice: Option<LogicalTime>,
}
struct SchedEventForLog<'a> {
event: &'a SchedEvent,
command_bootstrap: bool,
}
struct CommandBootstrapInstructionPointer(NonZeroUsize);
impl std::fmt::Debug for CommandBootstrapInstructionPointer {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", crate::logdiff::host_addr(self.0.get()))
}
}
impl std::fmt::Debug for SchedEventForLog<'_> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
if !self.command_bootstrap {
return std::fmt::Debug::fmt(self.event, f);
}
f.debug_struct("SchedEvent")
.field("dettid", &self.event.dettid)
.field("op", &self.event.op)
.field("count", &self.event.count)
.field(
"start_rip",
&self.event.start_rip.map(CommandBootstrapInstructionPointer),
)
.field(
"end_rip",
&self.event.end_rip.map(CommandBootstrapInstructionPointer),
)
.field("end_time", &self.event.end_time)
.finish()
}
}
pub async fn trace_schedevent<G, T>(guest: &mut G, ev: SchedEvent, tag_end_rip: bool)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
assert!(guest.config().sequentialize_threads);
let ev = if tag_end_rip {
let end_rip = if let Some(r) = ev.end_rip {
r
} else {
let regs = guest.regs().await;
NonZeroUsize::new(regs.rip.try_into().unwrap()).unwrap()
};
SchedEvent {
end_rip: Some(end_rip),
..ev
}
} else {
ev
};
if tracing::enabled!(tracing::Level::TRACE)
&& let Some(rip) = ev.end_rip
{
let rip_addr = AddrMut::<u16>::from_raw(rip.into()).unwrap();
let rip_contents = guest.memory().read_value(rip_addr);
trace!(
"Tracing sched event, after which rip is {}, next two instruction bytes {:?}",
rip, rip_contents
);
}
let detpid = guest.thread_state().detpid.expect("detpid unset");
let command_bootstrap = guest.is_command_bootstrap();
let resp = send_and_update_time(
guest,
GlobalRequest::TraceSchedEvent(ev, detpid, command_bootstrap),
)
.await;
trace!("trace_schedevent result: {:?}", resp);
match resp {
(
_,
GlobalResponse::TraceSchedEvent(TraceSchedEventResponse {
print_stack_strace,
timeslice,
}),
) => {
if let Some(m_path) = print_stack_strace {
trace!("[trace_schedevent] writing stacktrace via Reverie...");
write_backtrace(guest, m_path.as_ref());
}
if let Some(timeslice) = timeslice
&& guest.thread_state().past_global_first_execve
{
let end_of_timeslice =
guest.thread_state().thread_logical_time.as_nanos() + timeslice;
trace!(
"[detcore][dettid {}] setting end_of_timeslice to {:?} as instructed by replayer",
guest.thread_state().dettid,
end_of_timeslice
);
guest.thread_state_mut().end_of_timeslice = Some(end_of_timeslice);
if guest.config().max_timeslice.is_some() {
guest.thread_state_mut().max_timeslice_end = Some(end_of_timeslice);
}
}
}
_ => {
unreachable!()
}
}
}
pub async fn register_alarm<G, T>(
guest: &mut G,
duration: LogicalTime,
interval: LogicalTime,
sig: Signal,
) -> (LogicalTime, LogicalTime)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let dettid = guest.thread_state().dettid;
let detpid = guest.thread_state().detpid.expect("detpid unset");
let resp = send_and_update_time(
guest,
GlobalRequest::RegisterAlarm(detpid, dettid, duration, interval, SigWrapper::from(sig)),
)
.await;
match resp.1 {
GlobalResponse::RegisterAlarm(x) => x,
_ => unreachable!(),
}
}
pub async fn alarm_remaining<G, T>(guest: &mut G) -> ItimerSnapshot
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let detpid = guest.thread_state().detpid.expect("detpid unset");
let resp = send_and_update_time(guest, GlobalRequest::AlarmRemaining(detpid)).await;
match resp.1 {
GlobalResponse::AlarmRemaining(remaining) => remaining,
_ => unreachable!(),
}
}
pub async fn register_posix_timer<G, T>(
guest: &mut G,
timer_id: i32,
deadline: Option<LogicalTime>,
interval: LogicalTime,
sig: Signal,
) where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let dettid = guest.thread_state().dettid;
let detpid = guest.thread_state().detpid.expect("detpid unset");
let resp = send_and_update_time(
guest,
GlobalRequest::RegisterPosixTimer(
detpid,
dettid,
timer_id,
deadline,
interval,
SigWrapper::from(sig),
),
)
.await;
match resp.1 {
GlobalResponse::RegisterPosixTimer(()) => {}
_ => unreachable!(),
}
}
pub async fn thread_is_live<G, T>(guest: &mut G, dettid: DetTid) -> bool
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let response = send_and_update_time(guest, GlobalRequest::ThreadIsLive(dettid)).await;
match response.1 {
GlobalResponse::ThreadIsLive(live) => live,
_ => unreachable!(),
}
}
pub async fn exact_child_wait_state<G, T>(guest: &mut G, child: DetPid) -> ExactChildWaitState
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let parent = guest.thread_state().detpid.expect("detpid unset");
let response =
send_and_update_time(guest, GlobalRequest::ExactChildWaitState(parent, child)).await;
match response.1 {
GlobalResponse::ExactChildWaitState(state) => state,
_ => unreachable!(),
}
}
pub async fn await_exact_child_physical_exit<G, T>(
guest: &mut G,
child: DetPid,
) -> ExactChildWaitState
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let mut state = exact_child_wait_state(guest, child).await;
if matches!(
state,
ExactChildWaitState::PhysicalExitPending | ExactChildWaitState::PhysicallyExited
) {
let dettid = guest.thread_state().dettid;
let mut resources = Resources::new(dettid);
resources.insert(ResourceID::WaitPhysicalChild(child), Permission::R);
resources.fyi("wait-child-physical-exit");
let _ = resource_request(guest, resources).await;
state = exact_child_wait_state(guest, child).await;
}
state
}
pub async fn wait_for_child_lifecycle<G, T>(guest: &mut G, spec: ChildWaitSpec) -> ResumeStatus
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let dettid = guest.thread_state().dettid;
let parent = guest.thread_state().detpid.expect("detpid unset");
let mut resources = Resources::new(dettid);
resources.insert(ResourceID::WaitChild { parent, spec }, Permission::R);
resources.fyi("wait-child-lifecycle");
resource_request(guest, resources).await
}
pub async fn ready_child_wait<G, T>(guest: &mut G, spec: ChildWaitSpec) -> (Option<DetPid>, bool)
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let parent = guest.thread_state().detpid.expect("detpid unset");
let response = send_and_update_time(guest, GlobalRequest::ReadyChildWait(parent, spec)).await;
match response.1 {
GlobalResponse::ReadyChildWait(snapshot) => snapshot,
_ => unreachable!(),
}
}
pub async fn process_group<G, T>(guest: &mut G, process: DetPid) -> Option<DetPid>
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let response = send_and_update_time(guest, GlobalRequest::ProcessGroup(process)).await;
match response.1 {
GlobalResponse::ProcessGroup(group) => group,
_ => unreachable!(),
}
}
pub async fn set_process_group<G, T>(guest: &mut G, process: DetPid, group: DetPid) -> bool
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let response =
send_and_update_time(guest, GlobalRequest::SetProcessGroup(process, group)).await;
match response.1 {
GlobalResponse::SetProcessGroup(updated) => updated,
_ => unreachable!(),
}
}
pub async fn create_session<G, T>(guest: &mut G, process: DetPid) -> bool
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let response = send_and_update_time(guest, GlobalRequest::CreateSession(process)).await;
match response.1 {
GlobalResponse::CreateSession(updated) => updated,
_ => unreachable!(),
}
}
pub async fn consume_child_wait<G, T>(guest: &mut G, child: DetPid) -> bool
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let parent = guest.thread_state().detpid.expect("detpid unset");
let response =
send_and_update_time(guest, GlobalRequest::ConsumeChildWait(parent, child)).await;
match response.1 {
GlobalResponse::ConsumeChildWait(consumed) => consumed,
_ => unreachable!(),
}
}
pub async fn resolve_kill_targets<G, T>(guest: &mut G, detpid: DetPid) -> Vec<DetTid>
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let response = send_and_update_time(guest, GlobalRequest::ResolveKillTargets(detpid)).await;
match response.1 {
GlobalResponse::ResolveKillTargets(targets) => targets,
_ => unreachable!(),
}
}
fn alarm_signal(sig: SigWrapper) -> Signal {
sig.signal().unwrap_or_else(|| {
panic!(
"timer registration received unnameable signal {}",
sig.raw()
)
})
}
pub async fn notify_signal_pending<G, T>(
guest: &mut G,
dettid: DetTid,
signal: SigWrapper,
target_process: Option<DetPid>,
) where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let response = send_and_update_time(
guest,
GlobalRequest::NotifySignalPending(dettid, signal, target_process),
)
.await;
match response.1 {
GlobalResponse::NotifySignalPending(()) => {}
_ => unreachable!(),
}
}
pub async fn unrecoverable_shutdown<G, T>(guest: &G, status: i32) -> !
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
if cfg!(debug_assertions) {
let mytime = guest.thread_state().thread_logical_time.clone();
let mm = guest.thread_state().mm_id;
let _ = guest
.send_rpc((mytime, mm, GlobalRequest::UnrecoverableShutdown))
.await;
}
std::process::exit(status);
}
#[cfg(test)]
mod tests {
#[test]
fn external_scheduler_announces_before_its_future_is_polled() {
#[derive(Clone, Default)]
struct Capture(std::sync::Arc<Mutex<Vec<String>>>);
struct Visitor(String);
impl tracing::field::Visit for Visitor {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
if field.name() == "message" {
use std::fmt::Write;
write!(self.0, "{value:?}").unwrap();
}
}
}
impl tracing::Subscriber for Capture {
fn enabled(&self, metadata: &tracing::Metadata<'_>) -> bool {
*metadata.level() == tracing::Level::INFO
}
fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
tracing::span::Id::from_u64(1)
}
fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
fn event(&self, event: &tracing::Event<'_>) {
let mut visitor = Visitor(String::new());
event.record(&mut visitor);
self.0.lock().unwrap().push(visitor.0);
}
fn enter(&self, _: &tracing::span::Id) {}
fn exit(&self, _: &tracing::span::Id) {}
}
let config = Config {
sequentialize_threads: true,
..Config::default()
};
let state = GlobalState::init_for_external_scheduler(&config);
let captured = Capture::default();
let scheduler = tracing::subscriber::with_default(captured.clone(), || {
state.run_external_scheduler(std::sync::Arc::new(|_| {}))
});
let announcement = "[scheduler] daemon task starting up, waiting for guest thread start..";
assert_eq!(*captured.0.lock().unwrap(), [announcement]);
drop(scheduler);
}
#[test]
fn schedule_event_host_markers_require_command_bootstrap_provenance() {
let event = SchedEvent::branches(DetTid::from_raw(3), 223)
.with_end_rip(std::num::NonZeroUsize::new(0x1234).unwrap())
.with_time(LogicalTime::from_nanos(2230));
let original = serde_json::to_string(&event).unwrap();
let plain = format!(
"{:?}",
super::SchedEventForLog {
event: &event,
command_bootstrap: false
}
);
assert_eq!(plain, format!("{event:?}"));
let marked = format!(
"{:?}",
super::SchedEventForLog {
event: &event,
command_bootstrap: true
}
);
assert_eq!(
marked,
"SchedEvent { dettid: DetPid(3), op: Branch, count: 223, start_rip: None, end_rip: Some(<hostaddr 0x1234>), end_time: Some(LogicalTime(2230)) }"
);
assert_eq!(serde_json::to_string(&event).unwrap(), original);
}
#[test]
fn summary_preemption_views_keep_counts_full_report_and_single_flush() {
let (_config, state, tid, _) = cancellation_test_state();
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("recording");
let mut writer = crate::preemptions::PreemptionWriter::new(Some(path.clone()));
writer.register_thread(tid, DEFAULT_PRIORITY);
writer.insert_reprioritization(tid, LogicalTime::from_nanos(100), 10, DEFAULT_PRIORITY, 20);
let mut scheduler = state.sched.lock().unwrap();
scheduler.preemption_writer = Some(writer);
let (summary, info_description) = scheduler
.generate_partial_run_summary_for_log(Some(&path))
.unwrap();
assert!(scheduler.preemption_writer.is_none());
let recorded = std::fs::read(&path).unwrap();
let parsed: serde_json::Value = serde_json::from_slice(&recorded).unwrap();
assert_eq!(
parsed["per_thread"][tid.to_string()]["prio_changes"],
serde_json::json!([[100, DEFAULT_PRIORITY]])
);
let count_line = "Record of 1 preemption and reprioritization events:\n";
assert_eq!(info_description.as_deref(), Some(count_line));
assert_eq!(
summary.reprio_descrip.as_deref(),
Some(format!("{count_line} (Writing to file {path:?})\n").as_str())
);
let full = summary.to_string();
let json = serde_json::to_vec(&summary).unwrap();
let info = summary.info(info_description.as_deref()).to_string();
assert!(info.contains(count_line));
assert!(!info.contains("Writing to file"));
assert!(full.contains(&format!("Writing to file {path:?}")));
let _ = scheduler.generate_partial_run_summary(None).unwrap();
assert_eq!(std::fs::read(&path).unwrap(), recorded);
assert_eq!(summary.to_string(), full);
assert_eq!(serde_json::to_vec(&summary).unwrap(), json);
}
#[test]
fn run_summary_info_keeps_semantics_and_debug_retains_bookkeeping() {
#[derive(Clone)]
struct Capture(std::sync::Arc<Mutex<Vec<(tracing::Level, String)>>>);
struct Visitor(String);
impl tracing::field::Visit for Visitor {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
use std::fmt::Write;
write!(self.0, "{}={value:?};", field.name()).unwrap();
}
}
impl tracing::Subscriber for Capture {
fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
true
}
fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
tracing::span::Id::from_u64(1)
}
fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
fn event(&self, event: &tracing::Event<'_>) {
let mut visitor = Visitor(String::new());
event.record(&mut visitor);
self.0
.lock()
.unwrap()
.push((*event.metadata().level(), visitor.0));
}
fn enter(&self, _: &tracing::span::Id) {}
fn exit(&self, _: &tracing::span::Id) {}
}
let summary = super::RunSummary {
sched_turns: 4,
schedevent_recorded: 11,
schedevent_replayed: 11,
schedevent_desynced: 2,
desync_descrip: Some("two real desyncs\n".into()),
num_processes: 1,
num_threads: 1,
threads_descrip: "[3]".into(),
syscalls: Some(3),
virttime_elapsed: 2_512_380,
virttime_final: 2_512_380,
timeslice_stats: TimesliceStats {
count: 1,
sum_ns: 512_380,
min_ns: 512_380,
max_ns: 512_380,
},
..Default::default()
};
let description = "Record of 7 preemption and reprioritization events:\n";
let captured = Capture(Default::default());
tracing::subscriber::with_default(captured.clone(), || {
super::log_run_summary(
"report",
&summary,
Some(description),
Some(std::path::Path::new("/host/recording")),
);
});
let events = captured.0.lock().unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[0].0, tracing::Level::INFO);
assert_eq!(events[1].0, tracing::Level::DEBUG);
let info = &events[0].1;
for text in [
"1 group leaders of 1 thread(s)",
"3 syscalls",
"4 turns, recorded 11 events (2 desynced)",
"two real desyncs",
"Record of 7 preemption",
"2_512_380ns",
"min=512380ns max=512380ns mean=512380ns count=1",
] {
assert!(info.contains(text), "missing semantic field {text}: {info}");
}
assert!(!info.contains("/host/recording"));
assert!(!info.contains("replayed"));
assert!(events[1].1.contains("replayed_events=11"));
assert!(events[1].1.contains("/host/recording"));
let baseline = summary.info(Some(description)).to_string();
let mutations: [fn(&mut super::RunSummary); 8] = [
|s| s.sched_turns += 1,
|s| s.schedevent_recorded += 1,
|s| s.schedevent_desynced += 1,
|s| s.num_threads += 1,
|s| s.syscalls = Some(4),
|s| s.virttime_elapsed += 1,
|s| s.timeslice_stats.count += 1,
|s| s.threads_descrip.push_str(",4"),
];
for mutate in mutations {
let mut changed = summary.clone();
mutate(&mut changed);
assert_ne!(changed.info(Some(description)).to_string(), baseline);
}
assert_ne!(
summary
.info(Some(
"Record of 8 preemption and reprioritization events:\n"
))
.to_string(),
baseline
);
}
mod backend_failure_tests;
mod exec_identity_tests;
mod exec_teardown_tests;
use std::collections::BTreeSet;
use std::os::fd::AsRawFd;
use std::os::fd::FromRawFd;
use std::os::fd::OwnedFd;
use std::sync::Mutex;
use std::time::Duration;
use nix::sys::signal::Signal;
use reverie::GlobalRPC;
use reverie::GlobalTool;
use reverie::Guest;
use reverie::Tid;
use reverie::syscalls::CloneFlags;
use super::FutexAction;
use super::GlobalRequest;
use super::GlobalResponse;
use super::GlobalState;
use super::MountIdPool;
use super::PendingExecState;
use super::ResumeStatus;
use super::RpcIncarnation;
use super::SchedulerRpcResult;
use super::SigWrapper;
use super::ThreadDeregistration;
use super::TimesliceStats;
use super::format_unsupported_syscall_warning;
use crate::Detcore;
use crate::config::Config;
use crate::config::RunsPostFork;
use crate::ivar::Ivar;
use crate::preemptions::PreemptionRecord;
use crate::resources::ExternalOpId;
use crate::resources::Permission;
use crate::resources::ResourceID;
use crate::resources::Resources;
use crate::scheduler::DEFAULT_PRIORITY;
use crate::scheduler::SchedRequest;
use crate::scheduler::SchedResponse;
use crate::scheduler::SchedValue;
use crate::scheduler::ThreadNextTurn;
use crate::tool_local::ExecFdBlockingOverrides;
use crate::types::DetPid;
use crate::types::DetTid;
use crate::types::DetTime;
use crate::types::FutexID;
use crate::types::LogicalTime;
use crate::types::MmId;
use crate::types::Op;
use crate::types::SchedEvent;
#[test]
fn fdinfo_mount_ids_preserve_raw_equivalence_and_distinctness() {
let mut pool = MountIdPool::from_config(&[10, 20, 30], true, &[]);
assert_eq!(pool.determinize(20, None), Some(2));
assert_eq!(pool.determinize(700, None), Some(4));
assert_eq!(pool.determinize(701, None), Some(5));
assert_eq!(pool.determinize(702, None), Some(6));
assert_eq!(pool.determinize(701, None), Some(5));
}
#[test]
fn fdinfo_mount_ids_seed_from_low_level_snapshot_and_refuse_drift() {
let mut pool = MountIdPool::from_config(&[], false, &[]);
assert_eq!(pool.determinize(20, Some(&[10, 20, 30])), Some(2));
assert_eq!(pool.determinize(700, Some(&[10, 20, 30])), Some(4));
assert_eq!(pool.determinize(20, Some(&[10, 20, 30])), Some(2));
assert_eq!(pool.determinize(20, Some(&[10, 99, 30])), None);
assert_eq!(pool.determinize(20, Some(&[10, 20])), None);
assert_eq!(pool.determinize(20, Some(&[10, 20, 30, 40])), None);
}
#[test]
fn mountinfo_reads_seed_once_and_refuse_later_table_changes() {
let mut pool = MountIdPool::from_config(&[], false, &[]);
assert!(pool.validate_mountinfo_order(&[10, 20, 30]));
assert!(pool.validate_mountinfo_order(&[10, 20, 30]));
assert!(!pool.validate_mountinfo_order(&[10, 99, 30]));
let mut empty = MountIdPool::from_config(&[], false, &[]);
assert!(empty.validate_mountinfo_order(&[]));
assert!(!empty.validate_mountinfo_order(&[10]));
}
#[test]
fn configured_mountinfo_order_accepts_only_ordered_known_subsets() {
let mut pool = MountIdPool::from_config(&[10, 20, 30, 40], true, &[]);
assert!(pool.validate_mountinfo_order(&[10, 30, 40]));
assert_eq!(pool.determinize(30, None), Some(3));
assert!(pool.validate_mountinfo_order(&[20, 40]));
assert!(!pool.validate_mountinfo_order(&[30, 20]));
assert!(!pool.validate_mountinfo_order(&[10, 10]));
}
#[test]
fn configured_mountinfo_order_refuses_new_namespace_ids() {
let mut pool = MountIdPool::from_config(&[10, 20, 30, 40], true, &[]);
assert!(pool.validate_mountinfo_order(&[10, 30, 40]));
assert!(
!pool.validate_mountinfo_order(&[10, 99, 40]),
"a namespace-local mount ID absent from captured provenance must fail closed"
);
}
#[test]
fn fdinfo_mount_ids_refuse_malformed_configured_provenance() {
let mut pool = MountIdPool::from_config(&[10, 10], true, &[]);
assert_eq!(pool.determinize(10, None), None);
}
#[test]
fn raw_zero_is_one_mount_identity_without_descriptor_type_partitioning() {
let mut pool = MountIdPool::from_config(&[10, 0], true, &[]);
assert_eq!(pool.determinize(0, None), Some(0));
assert_eq!(pool.determinize(20, None), Some(3));
assert_eq!(pool.determinize(0, None), Some(0));
}
#[test]
fn recorded_unlisted_mount_order_rebuilds_the_same_mapping() {
let mut recording = MountIdPool::from_config(&[], false, &[]);
assert!(recording.validate_mountinfo_order(&[10, 20]));
assert_eq!(recording.determinize(700, None), Some(3));
assert_eq!(recording.determinize(701, None), Some(4));
let provenance = recording.provenance().unwrap().unwrap();
let mut replay = MountIdPool::from_config(
&provenance.mountinfo_order,
true,
&provenance.unlisted_order,
);
assert_eq!(replay.determinize(701, None), Some(4));
assert_eq!(replay.determinize(700, None), Some(3));
}
fn cancellation_test_state() -> (Config, GlobalState, DetTid, DetPid) {
let config = Config {
sequentialize_threads: true,
cancel_killed_thread_rpcs: true,
..Config::default()
};
let state = GlobalState::initialize(&config, false);
let dettid = DetTid::from_raw(17);
let detpid = DetPid::from_raw(17);
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(dettid, dettid, true);
(config, state, dettid, detpid)
}
fn install_test_registration(state: &GlobalState, dettid: DetTid, request: Ivar<SchedRequest>) {
let mut scheduler = state.sched.lock().unwrap();
assert!(!scheduler.thread_is_logically_killed(dettid));
scheduler.next_turns.insert(
dettid,
ThreadNextTurn {
dettid,
child_tid_addr: 0,
req: request,
resp: Ivar::new(),
protocol: Default::default(),
},
);
scheduler.priorities.insert(dettid, DEFAULT_PRIORITY);
scheduler.runqueue_push_back(dettid);
}
struct ExternalRegistrationGuest<'a> {
global: &'a GlobalState,
config: &'a Config,
thread: crate::ThreadState<()>,
requests: Mutex<Vec<GlobalRequest>>,
post_exec: bool,
}
struct ExternalRegistrationStack;
struct ExternalRegistrationStackGuard;
impl Drop for ExternalRegistrationStackGuard {
fn drop(&mut self) {}
}
impl reverie::Stack for ExternalRegistrationStack {
type StackGuard = ExternalRegistrationStackGuard;
fn size(&self) -> usize {
panic!("external registration must not use a guest stack")
}
fn capacity(&self) -> usize {
panic!("external registration must not use a guest stack")
}
fn push<'stack, T>(&mut self, _value: T) -> reverie::syscalls::Addr<'stack, T> {
panic!("external registration must not use a guest stack")
}
fn reserve<'stack, T>(&mut self) -> reverie::syscalls::AddrMut<'stack, T> {
panic!("external registration must not use a guest stack")
}
fn commit(self) -> Result<Self::StackGuard, reverie::syscalls::Errno> {
panic!("external registration must not use a guest stack")
}
}
#[reverie::tool]
impl GlobalRPC<GlobalState> for ExternalRegistrationGuest<'_> {
async fn send_rpc(
&self,
message: <GlobalState as GlobalTool>::Request,
) -> <GlobalState as GlobalTool>::Response {
self.requests.lock().unwrap().push(message.2.clone());
self.global
.receive_rpc(Tid::from_raw(self.thread.dettid.as_raw()), message)
.await
}
fn config(&self) -> &Config {
self.config
}
}
#[reverie::tool]
impl Guest<Detcore> for ExternalRegistrationGuest<'_> {
type Memory = reverie::syscalls::LocalMemory;
type Stack = ExternalRegistrationStack;
fn tid(&self) -> reverie::Pid {
reverie::Pid::from_raw(self.thread.dettid.as_raw())
}
fn pid(&self) -> reverie::Pid {
reverie::Pid::from_raw(self.thread.detpid.unwrap().as_raw())
}
fn ppid(&self) -> Option<reverie::Pid> {
None
}
fn auxv(&self) -> reverie::Auxv {
assert!(self.post_exec, "external registration must not read auxv");
reverie::Auxv::from_entries([])
}
fn memory(&self) -> Self::Memory {
panic!("external registration must not access guest memory")
}
fn thread_state_mut(&mut self) -> &mut crate::ThreadState<()> {
&mut self.thread
}
fn thread_state(&self) -> &crate::ThreadState<()> {
&self.thread
}
async fn regs(&mut self) -> libc::user_regs_struct {
assert!(
self.post_exec,
"external registration must not read guest registers"
);
unsafe { std::mem::zeroed() }
}
async fn stack(&mut self) -> Self::Stack {
panic!("external registration must not use a guest stack")
}
async fn daemonize(&mut self) {
panic!("external registration must not daemonize")
}
async fn inject<S: reverie::syscalls::SyscallInfo>(
&mut self,
_syscall: S,
) -> Result<i64, reverie::syscalls::Errno> {
panic!("external registration must not inject a syscall")
}
async fn tail_inject<S: reverie::syscalls::SyscallInfo>(
&mut self,
_syscall: S,
) -> reverie::Never {
panic!("external registration unexpectedly retired its live parent")
}
fn set_timer(&mut self, _schedule: reverie::TimerSchedule) -> Result<(), reverie::Error> {
panic!("external registration must not set a timer")
}
fn set_timer_precise(
&mut self,
_schedule: reverie::TimerSchedule,
) -> Result<(), reverie::Error> {
panic!("external registration must not set a timer")
}
fn read_clock(&mut self) -> Result<u64, reverie::Error> {
panic!("external registration must not read a host clock")
}
}
struct RetirementGuest<'a> {
global: &'a GlobalState,
config: &'a Config,
thread: crate::ThreadState<()>,
requests: Mutex<Vec<GlobalRequest>>,
retired: std::sync::atomic::AtomicBool,
}
#[reverie::tool]
impl GlobalRPC<GlobalState> for RetirementGuest<'_> {
async fn send_rpc(
&self,
message: <GlobalState as GlobalTool>::Request,
) -> <GlobalState as GlobalTool>::Response {
self.requests.lock().unwrap().push(message.2.clone());
self.global
.receive_rpc(Tid::from_raw(self.thread.dettid.as_raw()), message)
.await
}
fn config(&self) -> &Config {
self.config
}
}
#[reverie::tool]
impl Guest<Detcore> for RetirementGuest<'_> {
type Memory = reverie::syscalls::LocalMemory;
type Stack = ExternalRegistrationStack;
fn tid(&self) -> reverie::Pid {
reverie::Pid::from_raw(self.thread.dettid.as_raw())
}
fn pid(&self) -> reverie::Pid {
reverie::Pid::from_raw(self.thread.detpid.unwrap().as_raw())
}
fn ppid(&self) -> Option<reverie::Pid> {
None
}
fn memory(&self) -> Self::Memory {
panic!("external registration must not access guest memory")
}
fn thread_state_mut(&mut self) -> &mut crate::ThreadState<()> {
&mut self.thread
}
fn thread_state(&self) -> &crate::ThreadState<()> {
&self.thread
}
async fn regs(&mut self) -> libc::user_regs_struct {
panic!("external registration must not read guest registers")
}
async fn stack(&mut self) -> Self::Stack {
panic!("external registration must not use a guest stack")
}
async fn daemonize(&mut self) {
panic!("external registration must not daemonize")
}
async fn inject<S: reverie::syscalls::SyscallInfo>(
&mut self,
_syscall: S,
) -> Result<i64, reverie::syscalls::Errno> {
panic!("external registration must not inject a syscall")
}
async fn retire_current_thread(&mut self) -> reverie::Never {
self.retired
.store(true, std::sync::atomic::Ordering::SeqCst);
futures::future::pending().await
}
async fn cancel_current_thread(&mut self) -> reverie::Never {
panic!("scheduler retirement must not request group cancellation")
}
async fn tail_inject<S: reverie::syscalls::SyscallInfo>(
&mut self,
_syscall: S,
) -> reverie::Never {
panic!("external registration unexpectedly retired its live parent")
}
fn set_timer(&mut self, _schedule: reverie::TimerSchedule) -> Result<(), reverie::Error> {
panic!("external registration must not set a timer")
}
fn set_timer_precise(
&mut self,
_schedule: reverie::TimerSchedule,
) -> Result<(), reverie::Error> {
panic!("external registration must not set a timer")
}
fn read_clock(&mut self) -> Result<u64, reverie::Error> {
panic!("external registration must not read a host clock")
}
}
#[tokio::test]
async fn thread_exited_uses_natural_retirement_for_killed_stale_and_failed_rpcs() {
use std::sync::atomic::Ordering;
use reverie::Tool;
for cause in ["killed", "stale image", "backend failure", "live"] {
let (config, state, tid, pid) = cancellation_test_state();
install_test_registration(&state, tid, Ivar::new());
let tool: Detcore = Detcore::new(Tid::from_raw(tid.as_raw()), &config);
let mut thread = tool.init_thread_state(Tid::from_raw(tid.as_raw()), None);
thread.detpid = Some(pid);
thread.thread_logical_time.add_syscall_with_cost(37);
match cause {
"killed" => {
state
.sched
.lock()
.unwrap()
.logically_kill_thread(&tid, &pid, thread.mm_id)
}
"stale image" => {
state
.sched
.lock()
.unwrap()
.install_test_exec_incarnation(tid, thread.mm_id);
thread.mm_id = thread.mm_id.for_exec(pid);
}
"backend failure" => state.report_backend_failure(reverie::BackendFailure {
pid: Tid::from_raw(pid.as_raw()),
tid: Tid::from_raw(tid.as_raw()),
phase: "natural retirement control",
}),
"live" => {}
_ => unreachable!(),
}
let mut guest = RetirementGuest {
global: &state,
config: &config,
thread,
requests: Mutex::new(Vec::new()),
retired: std::sync::atomic::AtomicBool::new(false),
};
let before = guest.thread.thread_logical_time.as_nanos();
{
let mut call = std::pin::pin!(super::send_and_update_time(
&mut guest,
GlobalRequest::GlobalTimeLowerBound
));
if cause == "live" {
assert!(matches!(
futures::poll!(call.as_mut()),
std::task::Poll::Ready((_, GlobalResponse::GlobalTimeLowerBound(_)))
));
} else {
assert!(futures::poll!(call.as_mut()).is_pending(), "{cause}");
}
}
assert_eq!(
guest.retired.load(Ordering::SeqCst),
matches!(cause, "killed" | "stale image"),
"{cause}"
);
assert_eq!(guest.requests.lock().unwrap().len(), 1);
assert_eq!(guest.thread.thread_logical_time.as_nanos(), before);
if cause == "backend failure" {
assert!(state.sched.lock().unwrap().backend_failed());
assert!(
futures::poll!(std::pin::pin!(state.wait_for_backend_failure())).is_ready()
);
}
}
}
async fn check_external_child_tid_registration(
flags: CloneFlags,
supplied_address: usize,
expected_address: usize,
) {
use reverie::Tool;
let config = Config {
sequentialize_threads: true,
cancel_killed_thread_rpcs: true,
runs_post_fork: crate::RunsPostFork::Parent,
..Config::default()
};
let state = GlobalState::initialize(&config, false);
let parent = DetTid::from_raw(17);
let parent_pid = DetPid::from_raw(17);
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(parent, parent, true);
install_test_registration(&state, parent, Ivar::new());
let tool = Detcore::new(reverie::Pid::from_raw(parent.as_raw()), &config);
let mut thread = tool.init_thread_state(Tid::from_raw(parent.as_raw()), None);
thread.detpid = Some(parent_pid);
let mut guest = ExternalRegistrationGuest {
global: &state,
config: &config,
thread,
requests: Mutex::new(Vec::new()),
post_exec: false,
};
let child = DetTid::from_raw(18);
let exit_signal = if flags.contains(CloneFlags::CLONE_THREAD) {
0
} else {
libc::SIGCHLD
};
tokio::time::timeout(Duration::from_secs(2), async {
let mut registration = std::pin::pin!(tool.register_external_child(
&mut guest,
Tid::from_raw(child.as_raw()),
supplied_address,
flags,
exit_signal,
None,
));
assert!(futures::poll!(registration.as_mut()).is_pending());
let committed = crate::scheduler::do_a_turn_blocking(
state.sched.clone(),
state.global_time.clone(),
&Err(crate::scheduler::SkipTurn),
)
.await
.expect("the parent continuation must commit");
assert_eq!(committed.tid, parent);
assert_eq!(
committed.resources,
std::collections::HashMap::from([(
crate::resources::ResourceID::ParentContinue { parent, child },
crate::resources::Permission::W,
)]),
);
registration.await;
})
.await
.expect("external child registration must return without running a guest");
assert_eq!(guest.thread.clone_flags, None);
assert_eq!(guest.thread.dettid, parent);
let mut scheduler = state.sched.lock().unwrap();
assert_eq!(scheduler.next_turns.len(), 2);
assert!(scheduler.next_turns.contains_key(&parent));
assert_eq!(
scheduler.next_turns[&child].child_tid_addr,
expected_address
);
assert_eq!(
*guest.requests.lock().unwrap(),
vec![GlobalRequest::CreateChildThread(
child,
parent_pid,
expected_address,
Some(flags),
exit_signal,
None,
Some(DEFAULT_PRIORITY),
)],
"external registration must send one correctly gated real RPC"
);
let child_pid = if flags.contains(CloneFlags::CLONE_THREAD) {
parent_pid
} else {
child
};
let child_mm = MmId::for_clone(
MmId::initial(parent_pid),
child,
flags.contains(CloneFlags::CLONE_VM),
);
scheduler.logically_kill_thread(&child, &child_pid, child_mm);
assert!(!scheduler.next_turns.contains_key(&child));
assert!(scheduler.next_turns.contains_key(&parent));
assert_eq!(
scheduler.child_tid_was_cleared(
FutexID::private(child_mm, supplied_address),
child.as_raw(),
),
expected_address != 0,
"exit must use only the registered child-TID address"
);
assert!(!scheduler.child_tid_was_cleared(FutexID::private(child_mm, 0), child.as_raw(),));
assert!(!scheduler.child_tid_was_cleared(
FutexID::private(child_mm, supplied_address),
parent.as_raw(),
));
}
#[tokio::test]
async fn external_registration_without_child_cleartid_disables_exit_wake() {
for kind in [
CloneFlags::empty(),
CloneFlags::CLONE_THREAD | CloneFlags::CLONE_VM | CloneFlags::CLONE_SIGHAND,
] {
for registration in [CloneFlags::empty(), CloneFlags::CLONE_CHILD_SETTID] {
check_external_child_tid_registration(kind | registration, 0x1234, 0).await;
}
}
}
#[tokio::test]
async fn external_registration_with_child_cleartid_preserves_exact_exit_wake() {
for kind in [
CloneFlags::empty(),
CloneFlags::CLONE_THREAD | CloneFlags::CLONE_VM | CloneFlags::CLONE_SIGHAND,
] {
for registration in [
CloneFlags::CLONE_CHILD_CLEARTID,
CloneFlags::CLONE_CHILD_CLEARTID | CloneFlags::CLONE_CHILD_SETTID,
] {
check_external_child_tid_registration(kind | registration, 0x1234, 0x1234).await;
check_external_child_tid_registration(kind | registration, 0, 0).await;
}
}
}
#[tokio::test]
async fn set_child_tid_address_rpc_updates_and_resets_the_registration() {
let (config, state, dettid, detpid) = cancellation_test_state();
install_test_registration(&state, dettid, Ivar::new());
state
.sched
.lock()
.unwrap()
.next_turns
.get_mut(&dettid)
.unwrap()
.child_tid_addr = 0x1234;
let response = state
.receive_rpc(
reverie::Tid::from_raw(dettid.as_raw()),
(
DetTime::new(&config),
MmId::initial(detpid),
GlobalRequest::SetChildTidAddress(0),
),
)
.await;
assert_eq!(response.1, GlobalResponse::SetChildTidAddress(()));
assert_eq!(
state
.sched
.lock()
.unwrap()
.next_turns
.get(&dettid)
.unwrap()
.child_tid_addr,
0
);
}
async fn child_start_clock_trajectory(
start_before_selection: bool,
child_first: bool,
) -> (Vec<LogicalTime>, Vec<DetTid>, bool) {
let config = Config {
sequentialize_threads: true,
runs_post_fork: if child_first {
RunsPostFork::Child
} else {
RunsPostFork::Parent
},
..Config::default()
};
let state = GlobalState::initialize(&config, false);
let parent = DetTid::from_raw(17);
let parent_pid = parent;
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(parent, parent, true);
let child = DetTid::from_raw(parent.as_raw() + 1);
let parent_mm = MmId::initial(parent_pid);
let child_mm = MmId::for_clone(parent_mm, child, false);
install_test_registration(&state, parent, Ivar::new());
let epoch = DetTime::new(&config).as_nanos();
let mut parent_clock = DetTime::new(&config);
parent_clock.advance_to(epoch + LogicalTime::from_nanos(1_001));
let child_clock = parent_clock.clone_for_child();
let mut registration = Box::pin(state.receive_rpc(
Tid::from_raw(parent.as_raw()),
(
parent_clock.clone(),
parent_mm,
GlobalRequest::CreateChildThread(
child,
parent_pid,
0,
Some(CloneFlags::empty()),
libc::SIGCHLD,
None,
Some(DEFAULT_PRIORITY),
),
),
));
assert!(futures::poll!(&mut registration).is_pending());
let child_request = state.sched.lock().unwrap().next_turns[&child].req.clone();
let mut startup = Box::pin(state.receive_rpc(
Tid::from_raw(child.as_raw()),
(
child_clock.clone(),
child_mm,
GlobalRequest::StartNewThread(child, child, None, None),
),
));
if start_before_selection {
assert!(futures::poll!(&mut startup).is_pending());
assert!(futures::poll!(&mut startup).is_pending());
assert!(child_request.try_read().is_some());
}
let skipped = Err(crate::scheduler::SkipTurn);
let mut turn = Box::pin(crate::scheduler::do_a_turn_blocking(
state.sched.clone(),
state.global_time.clone(),
&skipped,
));
if !start_before_selection && child_first {
assert!(futures::poll!(&mut turn).is_pending());
assert_eq!(child_request.to_string(), "<ivar HasWaiter>");
assert!(futures::poll!(&mut startup).is_pending());
assert!(futures::poll!(&mut startup).is_pending());
assert!(child_request.try_read().is_some());
}
let first = turn.await.expect("first post-fork turn must commit");
if !start_before_selection && !child_first {
assert!(child_request.try_read().is_none());
assert!(futures::poll!(&mut startup).is_pending());
assert!(futures::poll!(&mut startup).is_pending());
assert!(child_request.try_read().is_some());
}
let child_resource = ResourceID::MemAddrSpace(child);
let parent_resource = ResourceID::ParentContinue { parent, child };
let (first_tid, first_resource, first_permission, second_resource, second_permission) =
if child_first {
(
child,
&child_resource,
Permission::RW,
&parent_resource,
Permission::W,
)
} else {
(
parent,
&parent_resource,
Permission::W,
&child_resource,
Permission::RW,
)
};
assert_eq!(first.tid, first_tid);
assert_eq!(first.resources.len(), 1);
assert_eq!(first.resources.get(first_resource), Some(&first_permission));
if child_first {
assert_eq!(
startup.as_mut().await,
(None, GlobalResponse::StartNewThread(None))
);
} else {
assert_eq!(
registration.as_mut().await,
(None, GlobalResponse::CreateChildThread(None))
);
}
let first_time = state.sched.lock().unwrap().committed_time;
assert_eq!(first_time, epoch + LogicalTime::from_nanos(1_001));
assert_eq!(
state.global_time.lock().unwrap().threads_time(child),
child_clock.as_nanos()
);
let (mut first_clock, first_mm) = if child_first {
(child_clock, child_mm)
} else {
(parent_clock, parent_mm)
};
first_clock.advance_to(first_clock.as_nanos() + LogicalTime::from_nanos(1));
let mut resources = Resources::new(first_tid);
resources.insert(ResourceID::MemAddrSpace(first_tid), Permission::RW);
let mut work = Box::pin(state.receive_rpc(
Tid::from_raw(first_tid.as_raw()),
(
first_clock,
first_mm,
GlobalRequest::RequestResources(resources, first_tid),
),
));
assert!(futures::poll!(&mut work).is_pending());
let second = crate::scheduler::do_a_turn_blocking(
state.sched.clone(),
state.global_time.clone(),
&Ok(first),
)
.await
.expect("other thread's continuation must commit");
assert_eq!(second.resources.len(), 1);
assert_eq!(
second.resources.get(second_resource),
Some(&second_permission)
);
if child_first {
assert_eq!(
registration.await,
(None, GlobalResponse::CreateChildThread(None))
);
} else {
assert_eq!(startup.await, (None, GlobalResponse::StartNewThread(None)));
}
let mut scheduler = state.sched.lock().unwrap();
let next_time = scheduler.committed_time;
assert_eq!(next_time, epoch + LogicalTime::from_nanos(501_002));
let queue = scheduler.run_queue.tids().copied().collect();
let next_random = scheduler.child_runs_first_post_fork(RunsPostFork::Random);
(vec![first_time, next_time], queue, next_random)
}
#[tokio::test]
async fn child_start_clock_is_independent_of_first_rpc_arrival() {
for child_first in [true, false] {
assert_eq!(
child_start_clock_trajectory(true, child_first).await,
child_start_clock_trajectory(false, child_first).await,
);
}
}
#[tokio::test]
async fn vfork_registration_does_not_charge_inherited_work_before_startup() {
let (config, state, parent, parent_pid) = cancellation_test_state();
let child = DetTid::from_raw(parent.as_raw() + 1);
let mm = MmId::initial(parent_pid);
install_test_registration(&state, parent, Ivar::new());
let mut parent_clock = DetTime::new(&config);
let epoch = parent_clock.as_nanos();
parent_clock.advance_to(epoch + LogicalTime::from_nanos(1_001));
let child_clock = parent_clock.clone_for_child();
let mut resources = Resources::new(parent);
resources.insert(
ResourceID::BlockingVfork(ExternalOpId::new(parent, 1)),
Permission::RW,
);
let mut blocking = Box::pin(state.receive_rpc(
Tid::from_raw(parent.as_raw()),
(
parent_clock,
mm,
GlobalRequest::RequestResources(resources, parent_pid),
),
));
assert!(futures::poll!(&mut blocking).is_pending());
let background = crate::scheduler::do_a_turn_blocking(
state.sched.clone(),
state.global_time.clone(),
&Err(crate::scheduler::SkipTurn),
)
.await;
assert!(background.is_err());
assert_eq!(
blocking.await,
(None, GlobalResponse::RequestResources(ResumeStatus::Normal))
);
let before_child = state.global_time.lock().unwrap().as_nanos();
assert_eq!(before_child, epoch + LogicalTime::from_nanos(1_001));
let created = state
.receive_rpc(
Tid::from_raw(child.as_raw()),
(
child_clock.clone(),
mm,
GlobalRequest::CreateVforkChildThread(
parent,
parent_pid,
child,
0,
CloneFlags::CLONE_VFORK | CloneFlags::CLONE_VM,
libc::SIGCHLD,
Some(DEFAULT_PRIORITY - 1),
),
),
)
.await;
assert_eq!(created, (None, GlobalResponse::CreateChildThread(None)));
assert_eq!(state.global_time.lock().unwrap().as_nanos(), before_child);
assert_eq!(
state.global_time.lock().unwrap().threads_time(child),
child_clock.as_nanos()
);
let mut startup = Box::pin(state.receive_rpc(
Tid::from_raw(child.as_raw()),
(
child_clock,
mm,
GlobalRequest::StartNewThread(child, child, None, None),
),
));
assert!(futures::poll!(&mut startup).is_pending());
assert!(futures::poll!(&mut startup).is_pending());
let first = crate::scheduler::do_a_turn_blocking(
state.sched.clone(),
state.global_time.clone(),
&background,
)
.await
.expect("vfork child must receive its first turn");
assert_eq!(first.tid, child);
assert_eq!(startup.await, (None, GlobalResponse::StartNewThread(None)));
assert_eq!(state.global_time.lock().unwrap().as_nanos(), before_child);
}
#[tokio::test]
async fn first_nonstartup_rpc_counts_only_new_child_work() {
let config = Config {
sequentialize_threads: false,
..Config::default()
};
let state = GlobalState::initialize(&config, false);
let parent = DetTid::from_raw(17);
let child = DetTid::from_raw(18);
let mut parent_clock = DetTime::new(&config);
let epoch = parent_clock.as_nanos();
parent_clock.advance_to(epoch + LogicalTime::from_nanos(101));
let mut child_clock = parent_clock.clone_for_child();
let _ = state
.receive_rpc(
Tid::from_raw(parent.as_raw()),
(
parent_clock,
MmId::initial(parent),
GlobalRequest::GlobalTimeLowerBound,
),
)
.await;
child_clock.advance_to(child_clock.as_nanos() + LogicalTime::from_nanos(1));
let observed = state
.receive_rpc(
Tid::from_raw(child.as_raw()),
(
child_clock,
MmId::initial(child),
GlobalRequest::GlobalTimeLowerBound,
),
)
.await;
assert_eq!(
observed,
(
None,
GlobalResponse::GlobalTimeLowerBound(epoch + LogicalTime::from_nanos(102))
)
);
}
#[tokio::test]
async fn exec_reconnect_retains_inherited_work_accounting_across_local_reload() {
let (config, state, leader, detpid) = cancellation_test_state();
let ancestor = DetTid::from_raw(leader.as_raw() - 1);
let worker = DetTid::from_raw(leader.as_raw() + 1);
let old_mm = MmId::initial(detpid);
install_test_registration(&state, leader, Ivar::new());
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(leader, worker, false);
install_test_registration(&state, worker, Ivar::new());
let mut ancestor_clock = DetTime::new(&config);
let epoch = ancestor_clock.as_nanos();
ancestor_clock.advance_to(epoch + LogicalTime::from_nanos(1_000));
let mut leader_clock = ancestor_clock.clone_for_child();
leader_clock.advance_to(epoch + LogicalTime::from_nanos(1_100));
let mut worker_clock = leader_clock.clone_for_child();
worker_clock.advance_to(epoch + LogicalTime::from_nanos(1_350));
for (tid, clock) in [
(ancestor, ancestor_clock),
(leader, leader_clock),
(worker, worker_clock.clone()),
] {
let _ = state
.receive_rpc(
Tid::from_raw(tid.as_raw()),
(clock, old_mm, GlobalRequest::GlobalTimeLowerBound),
)
.await;
}
let total = state.global_time.lock().unwrap().as_nanos();
assert_eq!(total, epoch + LogicalTime::from_nanos(1_350));
let _ = state
.receive_rpc(
Tid::from_raw(worker.as_raw()),
(
worker_clock.clone(),
old_mm,
GlobalRequest::PrepareExec(detpid, old_mm, Default::default()),
),
)
.await;
let cancelled = state
.receive_rpc(
Tid::from_raw(worker.as_raw()),
(
worker_clock.clone(),
old_mm,
GlobalRequest::CancelExec(detpid),
),
)
.await;
assert_eq!(cancelled, (None, GlobalResponse::CancelExec(())));
assert!(state.pending_exec_states.lock().unwrap().is_empty());
assert_eq!(state.global_time.lock().unwrap().as_nanos(), total);
assert_eq!(
state.global_time.lock().unwrap().threads_time(worker),
worker_clock.as_nanos()
);
let prepared = state
.receive_rpc(
Tid::from_raw(worker.as_raw()),
(
worker_clock.clone(),
old_mm,
GlobalRequest::PrepareExec(detpid, old_mm, Default::default()),
),
)
.await;
assert_eq!(prepared, (None, GlobalResponse::PrepareExec(())));
let mut fresh = DetTime::new(&config);
let recreated = state
.receive_rpc(
Tid::from_raw(leader.as_raw()),
(
fresh.clone(),
MmId::initial(leader),
GlobalRequest::CreateChildThread(
leader,
detpid,
0,
None,
libc::SIGCHLD,
None,
Some(DEFAULT_PRIORITY),
),
),
)
.await;
assert_eq!(
recreated,
(
Some(worker_clock.as_nanos()),
GlobalResponse::CreateChildThread(Some(old_mm.for_exec(detpid)))
)
);
let stale = state
.receive_rpc(
Tid::from_raw(leader.as_raw()),
(
worker_clock.clone(),
old_mm,
GlobalRequest::RequestResources(Resources::new(leader), detpid),
),
)
.await;
assert_eq!(stale, (None, GlobalResponse::ThreadExited));
assert_eq!(state.global_time.lock().unwrap().as_nanos(), total);
fresh.advance_to(recreated.0.unwrap());
assert_eq!(fresh.inherited_nanos(), LogicalTime::ZERO);
let mut startup = Box::pin(state.receive_rpc(
Tid::from_raw(leader.as_raw()),
(
fresh.clone(),
old_mm.for_exec(detpid),
GlobalRequest::StartNewThread(leader, detpid, None, None),
),
));
assert!(futures::poll!(&mut startup).is_pending());
assert!(futures::poll!(&mut startup).is_pending());
let first = crate::scheduler::do_a_turn_blocking(
state.sched.clone(),
state.global_time.clone(),
&Err(crate::scheduler::SkipTurn),
)
.await
.expect("replacement leader must run");
assert_eq!(first.tid, leader);
assert_eq!(startup.await, (None, GlobalResponse::StartNewThread(None)));
assert_eq!(state.global_time.lock().unwrap().as_nanos(), total);
fresh.advance_to(fresh.as_nanos() + LogicalTime::from_nanos(1));
let observed = state
.receive_rpc(
Tid::from_raw(leader.as_raw()),
(
fresh.clone(),
old_mm.for_exec(detpid),
GlobalRequest::GlobalTimeLowerBound,
),
)
.await;
assert_eq!(
observed,
(
None,
GlobalResponse::GlobalTimeLowerBound(total + LogicalTime::from_nanos(1))
)
);
let global = state.global_time.lock().unwrap();
assert_eq!(global.threads_time(leader), fresh.as_nanos());
assert!(!global.contains_thread(worker));
}
#[test]
fn live_registration_without_next_turn_is_not_terminal() {
let (_, state, dettid, detpid) = cancellation_test_state();
install_test_registration(&state, dettid, Ivar::new());
state.sched.lock().unwrap().next_turns.remove(&dettid);
assert!(
!state
.sched
.lock()
.unwrap()
.thread_is_logically_killed(dettid),
"transient next-turn absence must not imply logical death"
);
state
.sched
.lock()
.unwrap()
.logically_kill_thread(&dettid, &detpid, MmId::initial(detpid));
assert!(
state
.sched
.lock()
.unwrap()
.thread_is_logically_killed(dettid),
"explicit logical death must install a permanent TID tombstone"
);
}
#[tokio::test]
async fn exec_reconnect_retires_siblings_and_reuses_live_scheduler_and_clock_state() {
let (config, state, dettid, detpid) = cancellation_test_state();
let old_mm = MmId::initial(detpid).for_exec(detpid);
install_test_registration(&state, dettid, Ivar::new());
let sibling = DetTid::from_raw(dettid.as_raw() + 1);
let sibling_request = Ivar::new();
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(dettid, sibling, false);
install_test_registration(&state, sibling, sibling_request.clone());
{
let mut scheduler = state.sched.lock().unwrap();
scheduler.next_turns.get_mut(&dettid).unwrap().resp =
Ivar::full(SchedResponse::Go(None));
scheduler
.next_turns
.get_mut(&sibling)
.unwrap()
.child_tid_addr = 0x1234;
}
let mut existing_time = DetTime::new(&config);
existing_time.add_syscall();
existing_time.add_syscall();
state.global_time.lock().unwrap().update_global_time(
dettid,
existing_time.as_nanos(),
LogicalTime::ZERO,
);
let (global_before, thread_before) = {
let global_time = state.global_time.lock().unwrap();
(global_time.as_nanos(), global_time.threads_time(dettid))
};
let fresh_local_time = DetTime::new(&config);
let physical_pid = std::process::id() as i32;
let physical_tid = unsafe { libc::syscall(libc::SYS_gettid) as i32 };
let physical_ids = Some((physical_pid, physical_tid));
state.pending_exec_states.lock().unwrap().insert(
detpid,
PendingExecState {
caller: dettid,
process: detpid,
mm: old_mm,
fd_blocking: Default::default(),
},
);
state
.sched
.lock()
.unwrap()
.install_test_exec_incarnation(dettid, old_mm);
let in_flight_exec_response = state
.receive_rpc(
reverie::Tid::from_raw(dettid.as_raw()),
(
existing_time.clone(),
old_mm.for_exec(detpid),
GlobalRequest::ReportUnsupportedSyscall("exec-in-flight".to_owned()),
),
)
.await;
assert_eq!(
in_flight_exec_response.1,
GlobalResponse::ReportUnsupportedSyscall(())
);
let create_response = state
.receive_rpc(
reverie::Tid::from_raw(dettid.as_raw()),
(
fresh_local_time.clone(),
MmId::initial(dettid),
GlobalRequest::CreateChildThread(
dettid,
detpid,
0,
None,
libc::SIGCHLD,
physical_ids,
Some(DEFAULT_PRIORITY),
),
),
)
.await;
assert_eq!(
create_response,
(
Some(thread_before),
GlobalResponse::CreateChildThread(Some(old_mm.for_exec(detpid)))
)
);
assert_eq!(
state.sched.lock().unwrap().physical_thread_identity(dettid),
Some((old_mm.for_exec(detpid), physical_pid, physical_tid)),
"post-exec host identity must be installed by CreateChildThread before admission"
);
let start_response = state
.receive_rpc(
reverie::Tid::from_raw(dettid.as_raw()),
(
fresh_local_time,
old_mm.for_exec(detpid),
GlobalRequest::StartNewThread(dettid, detpid, physical_ids, None),
),
)
.await;
assert_eq!(
start_response,
(Some(thread_before), GlobalResponse::StartNewThread(None))
);
let scheduler = state.sched.lock().unwrap();
assert!(!scheduler.thread_is_logically_killed(dettid));
assert!(scheduler.thread_is_logically_killed(sibling));
assert_eq!(scheduler.next_turns.len(), 1);
assert!(matches!(sibling_request.try_read(), Some(Err(_))));
assert!(
scheduler.child_tid_was_cleared(FutexID::private(old_mm, 0x1234), sibling.as_raw())
);
drop(scheduler);
let global_time = state.global_time.lock().unwrap();
assert_eq!(global_time.as_nanos(), global_before);
assert_eq!(global_time.threads_time(dettid), thread_before);
}
#[tokio::test]
async fn nonleader_exec_rebinds_caller_to_leader_and_preserves_its_clock() {
let (config, state, leader, detpid) = cancellation_test_state();
let worker = DetTid::from_raw(leader.as_raw() + 1);
let sibling = DetTid::from_raw(leader.as_raw() + 2);
let old_mm = MmId::initial(detpid).for_exec(detpid);
let leader_request = Ivar::new();
let worker_request = Ivar::new();
let sibling_request = Ivar::new();
install_test_registration(&state, leader, leader_request.clone());
{
let mut scheduler = state.sched.lock().unwrap();
scheduler.thread_tree.add_child(leader, worker, false);
scheduler.thread_tree.add_child(leader, sibling, false);
}
install_test_registration(&state, worker, worker_request.clone());
install_test_registration(&state, sibling, sibling_request.clone());
{
let mut scheduler = state.sched.lock().unwrap();
scheduler
.next_turns
.get_mut(&sibling)
.unwrap()
.child_tid_addr = 0x5678;
scheduler
.timeslices
.insert(leader, Some(LogicalTime::from_nanos(99)));
scheduler.install_test_vfork_barrier(leader, sibling);
}
let mut leader_clock = DetTime::new(&config);
leader_clock.add_syscall();
let mut worker_clock = DetTime::new(&config);
worker_clock.add_syscall();
worker_clock.add_syscall();
worker_clock.add_syscall();
{
let mut global_time = state.global_time.lock().unwrap();
global_time.update_global_time(leader, leader_clock.as_nanos(), LogicalTime::ZERO);
global_time.update_global_time(worker, worker_clock.as_nanos(), LogicalTime::ZERO);
}
let total_before = state.global_time.lock().unwrap().as_nanos();
let fd_blocking: ExecFdBlockingOverrides = [42].into_iter().collect();
state.pending_exec_states.lock().unwrap().insert(
detpid,
PendingExecState {
caller: worker,
process: detpid,
mm: old_mm,
fd_blocking: fd_blocking.clone(),
},
);
let fresh_local_time = DetTime::new(&config);
let create_response = state
.receive_rpc(
reverie::Tid::from_raw(leader.as_raw()),
(
fresh_local_time.clone(),
MmId::initial(leader),
GlobalRequest::CreateChildThread(
leader,
detpid,
0,
None,
libc::SIGCHLD,
None,
Some(DEFAULT_PRIORITY),
),
),
)
.await;
assert_eq!(
create_response,
(
Some(worker_clock.as_nanos()),
GlobalResponse::CreateChildThread(Some(old_mm.for_exec(detpid)))
)
);
assert!(state.pending_exec_states.lock().unwrap().is_empty());
assert_eq!(
state.post_exec_fd_blocking.lock().unwrap().get(&leader),
Some(&fd_blocking)
);
let late_old_request = state
.receive_rpc(
reverie::Tid::from_raw(leader.as_raw()),
(
leader_clock.clone(),
old_mm,
GlobalRequest::RequestResources(Resources::new(leader), detpid),
),
)
.await;
assert_eq!(late_old_request, (None, GlobalResponse::ThreadExited));
let admitted_before_fence = state
.recv_request_resources(
reverie::Tid::from_raw(leader.as_raw()),
detpid,
Resources::new(leader),
Some(old_mm),
)
.await;
assert_eq!(
admitted_before_fence,
(SchedulerRpcResult::ThreadExited, None)
);
let late_old_deregister = state
.receive_rpc(
reverie::Tid::from_raw(leader.as_raw()),
(
leader_clock.clone(),
old_mm,
GlobalRequest::DeregisterThread(ThreadDeregistration {
thread_start_entered: true,
dettid: leader,
detpid,
mm: old_mm,
timeslice_stats: TimesliceStats::default(),
syscall_count: 0,
chaos_epochs: Vec::new(),
}),
),
)
.await;
assert_eq!(
late_old_deregister,
(None, GlobalResponse::DeregisterThread(()))
);
let duplicate_create = state
.receive_rpc(
reverie::Tid::from_raw(leader.as_raw()),
(
fresh_local_time.clone(),
MmId::initial(leader),
GlobalRequest::CreateChildThread(
leader,
detpid,
0,
None,
libc::SIGCHLD,
None,
Some(DEFAULT_PRIORITY),
),
),
)
.await;
assert_eq!(duplicate_create, (None, GlobalResponse::ThreadExited));
{
let scheduler = state.sched.lock().unwrap();
assert!(!scheduler.thread_is_logically_killed(leader));
assert!(scheduler.thread_is_logically_killed(worker));
assert!(scheduler.thread_is_logically_killed(sibling));
assert_eq!(scheduler.next_turns.len(), 1);
assert!(scheduler.next_turns.contains_key(&leader));
assert!(!scheduler.timeslices.contains_key(&leader));
assert!(!scheduler.vfork_barrier_mentions(leader));
assert!(!scheduler.vfork_barrier_mentions(sibling));
assert!(matches!(leader_request.try_read(), Some(Err(_))));
assert!(matches!(worker_request.try_read(), Some(Err(_))));
assert!(matches!(sibling_request.try_read(), Some(Err(_))));
assert!(
scheduler
.child_tid_was_cleared(FutexID::private(old_mm, 0x5678), sibling.as_raw(),)
);
}
let turn_sched = state.sched.clone();
let turn_time = state.global_time.clone();
let turn = tokio::spawn(async move {
let last: Result<Resources, crate::scheduler::SkipTurn> =
Err(crate::scheduler::SkipTurn);
crate::scheduler::do_a_turn_blocking(turn_sched, turn_time, &last).await
});
let start_response = state
.receive_rpc(
reverie::Tid::from_raw(leader.as_raw()),
(
fresh_local_time,
old_mm.for_exec(detpid),
GlobalRequest::StartNewThread(leader, detpid, None, None),
),
)
.await;
assert!(
turn.await
.expect("replacement scheduler turn panicked")
.is_ok(),
"replacement leader did not survive the first step2 drain"
);
assert_eq!(
start_response,
(
Some(worker_clock.as_nanos()),
GlobalResponse::StartNewThread(None)
)
);
{
let scheduler = state.sched.lock().unwrap();
assert_eq!(
scheduler
.run_queue
.tids()
.filter(|dettid| **dettid == leader)
.count(),
1
);
assert!(!scheduler.run_queue.contains_tid(worker));
assert!(!scheduler.run_queue.contains_tid(sibling));
}
{
let global_time = state.global_time.lock().unwrap();
assert_eq!(global_time.as_nanos(), total_before);
assert_eq!(global_time.threads_time(leader), worker_clock.as_nanos());
assert!(!global_time.contains_thread(worker));
}
let mark_response = state
.receive_rpc(
reverie::Tid::from_raw(leader.as_raw()),
(
worker_clock,
old_mm.for_exec(detpid),
GlobalRequest::MarkPastFirstExecve(detpid, None),
),
)
.await;
assert_eq!(
mark_response.1,
GlobalResponse::MarkPastFirstExecve(fd_blocking)
);
assert!(state.post_exec_fd_blocking.lock().unwrap().is_empty());
}
#[tokio::test]
async fn post_exec_deletes_armed_disarmed_and_non_notifying_timers() {
use reverie::Tool;
let config = Config {
sequentialize_threads: true,
cancel_killed_thread_rpcs: true,
max_timeslice: None,
..Config::default()
};
let state = GlobalState::initialize(&config, false);
let leader = DetTid::from_raw(17);
let detpid = DetPid::from_raw(17);
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(leader, leader, true);
let next_request = Ivar::new();
install_test_registration(&state, leader, next_request.clone());
let tool = Detcore::new(reverie::Pid::from_raw(leader.as_raw()), &config);
let mut thread = tool.init_thread_state(Tid::from_raw(leader.as_raw()), None);
thread.detpid = Some(detpid);
thread.thread_logical_time.add_syscall_with_cost(137);
let mut guest = ExternalRegistrationGuest {
global: &state,
config: &config,
thread,
requests: Mutex::new(Vec::new()),
post_exec: true,
};
let t = LogicalTime::from_nanos;
let [periodic, disarmed, silent] = {
let mut timers = guest.thread.posix_timers.lock().unwrap();
let periodic = timers.create(Some(libc::SIGUSR2));
let disarmed = timers.create(Some(libc::SIGALRM));
let silent = timers.create(None);
timers.settime(periodic, 50, Some(t(100)), t(0));
timers.settime(silent, 0, Some(t(200)), t(0));
[periodic, disarmed, silent]
};
for _ in 0..2 {
let old_mm = guest.thread.mm_id;
super::prepare_exec(&mut guest, old_mm, Default::default()).await;
guest.thread.mm_id = old_mm.for_exec(detpid);
let before = guest.thread.thread_logical_time.as_nanos();
tokio::time::timeout(Duration::from_secs(2), tool.handle_post_exec(&mut guest))
.await
.expect("post-exec must finish within the caller's scheduler turn")
.unwrap();
assert_eq!(guest.thread.thread_logical_time.as_nanos(), before);
assert!(
next_request.try_read().is_none(),
"post-exec unexpectedly yielded"
);
let mut timers = guest.thread.posix_timers.lock().unwrap();
for id in [periodic, disarmed, silent] {
assert!(!timers.contains(id), "post-exec retained POSIX timer {id}");
assert_eq!(timers.gettime(id, t(10)), None);
assert_eq!(timers.settime(id, 0, Some(t(300)), t(10)), None);
assert_eq!(timers.signal(id), None);
assert!(!timers.remove(id));
}
}
let mut timers = guest.thread.posix_timers.lock().unwrap();
let new = timers.create(Some(libc::SIGUSR1));
assert_eq!(new, 3);
assert_eq!(timers.settime(new, 0, Some(t(300)), t(10)), Some((0, 0)));
assert_eq!(timers.gettime(new, t(20)), Some((280, 0)));
}
#[tokio::test]
async fn only_successful_exec_cancels_posix_deadlines() {
let (config, state, leader, detpid) = cancellation_test_state();
install_test_registration(&state, leader, Ivar::new());
let clock = DetTime::new(&config);
let mm = MmId::initial(detpid);
let deadline = LogicalTime::from_nanos(1_000_000);
{
let mut sched = state.sched.lock().unwrap();
sched.register_alarm(
detpid,
leader,
LogicalTime::ZERO,
deadline,
LogicalTime::ZERO,
Signal::SIGALRM,
);
sched.register_posix_timer(
detpid,
leader,
0,
Some(deadline),
deadline,
Signal::SIGUSR2,
);
}
let before: Vec<_> = state
.sched
.lock()
.unwrap()
.blocked
.timed_waiters
.iter()
.collect();
for request in [
GlobalRequest::PrepareExec(detpid, mm, Default::default()),
GlobalRequest::CancelExec(detpid),
GlobalRequest::PrepareExec(detpid, mm, Default::default()),
] {
state
.receive_rpc(
reverie::Tid::from_raw(leader.as_raw()),
(clock.clone(), mm, request),
)
.await;
assert_eq!(
state
.sched
.lock()
.unwrap()
.blocked
.timed_waiters
.iter()
.collect::<Vec<_>>(),
before
);
}
let response = state
.receive_rpc(
reverie::Tid::from_raw(leader.as_raw()),
(
clock,
mm.for_exec(detpid),
GlobalRequest::MarkPastFirstExecve(detpid, None),
),
)
.await;
assert_eq!(
response.1,
GlobalResponse::MarkPastFirstExecve(Default::default())
);
let sched = state.sched.lock().unwrap();
assert_eq!(
sched.blocked.timed_waiters.alarm_state(detpid),
Some((deadline, LogicalTime::ZERO))
);
assert_eq!(sched.blocked.timed_waiters.iter().count(), 1);
}
#[tokio::test]
async fn failed_exec_clears_prepared_state_without_retiring_siblings() {
let (config, state, leader, detpid) = cancellation_test_state();
let sibling = DetTid::from_raw(leader.as_raw() + 1);
install_test_registration(&state, leader, Ivar::new());
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(leader, sibling, false);
install_test_registration(&state, sibling, Ivar::new());
let clock = DetTime::new(&config);
let prepared = state
.receive_rpc(
reverie::Tid::from_raw(leader.as_raw()),
(
clock.clone(),
MmId::initial(leader),
GlobalRequest::PrepareExec(detpid, MmId::initial(detpid), Default::default()),
),
)
.await;
assert_eq!(prepared.1, GlobalResponse::PrepareExec(()));
assert!(
state
.pending_exec_states
.lock()
.unwrap()
.contains_key(&detpid)
);
let cancelled = state
.receive_rpc(
reverie::Tid::from_raw(leader.as_raw()),
(
clock,
MmId::initial(leader),
GlobalRequest::CancelExec(detpid),
),
)
.await;
assert_eq!(cancelled.1, GlobalResponse::CancelExec(()));
assert!(state.pending_exec_states.lock().unwrap().is_empty());
{
let scheduler = state.sched.lock().unwrap();
assert!(!scheduler.thread_is_logically_killed(leader));
assert!(!scheduler.thread_is_logically_killed(sibling));
assert_eq!(scheduler.next_turns.len(), 2);
}
state.pending_exec_states.lock().unwrap().insert(
detpid,
PendingExecState {
caller: leader,
process: detpid,
mm: MmId::initial(detpid),
fd_blocking: Default::default(),
},
);
state
.post_exec_fd_blocking
.lock()
.unwrap()
.insert(leader, [42].into_iter().collect());
state
.recv_deregister_thread(
reverie::Tid::from_raw(leader.as_raw()),
ThreadDeregistration {
thread_start_entered: true,
dettid: leader,
detpid,
mm: MmId::initial(detpid),
timeslice_stats: TimesliceStats::default(),
syscall_count: 0,
chaos_epochs: Vec::new(),
},
)
.await;
assert!(state.pending_exec_states.lock().unwrap().is_empty());
assert!(state.post_exec_fd_blocking.lock().unwrap().is_empty());
state.pending_exec_states.lock().unwrap().insert(
detpid,
PendingExecState {
caller: leader,
process: detpid,
mm: MmId::initial(detpid),
fd_blocking: Default::default(),
},
);
state
.recv_deregister_thread(
reverie::Tid::from_raw(leader.as_raw()),
ThreadDeregistration {
thread_start_entered: true,
dettid: leader,
detpid,
mm: MmId::initial(detpid).for_exec(detpid),
timeslice_stats: TimesliceStats::default(),
syscall_count: 0,
chaos_epochs: Vec::new(),
},
)
.await;
assert!(state.pending_exec_states.lock().unwrap().is_empty());
state.pending_exec_states.lock().unwrap().insert(
detpid,
PendingExecState {
caller: leader,
process: detpid,
mm: MmId::initial(detpid),
fd_blocking: Default::default(),
},
);
state
.post_exec_fd_blocking
.lock()
.unwrap()
.insert(leader, [42].into_iter().collect());
state.complete_physical_process_exit(detpid.as_raw());
assert!(state.pending_exec_states.lock().unwrap().is_empty());
assert!(state.post_exec_fd_blocking.lock().unwrap().is_empty());
}
#[test]
fn unsupported_syscall_report_duplicate_is_close_on_exec() {
let mut descriptors = [-1; 2];
assert_eq!(
unsafe { libc::pipe2(descriptors.as_mut_ptr(), libc::O_CLOEXEC) },
0
);
let _reader = unsafe { OwnedFd::from_raw_fd(descriptors[0]) };
let writer = unsafe { OwnedFd::from_raw_fd(descriptors[1]) };
let config = Config {
unsupported_syscall_report_fd: Some(writer.as_raw_fd()),
..Config::default()
};
let state = GlobalState::initialize(&config, false);
let duplicate = state
.unsupported_syscall_report_fd
.as_ref()
.expect("report writer should be duplicated")
.lock()
.unwrap();
let flags = unsafe { libc::fcntl(duplicate.as_raw_fd(), libc::F_GETFD) };
assert_ne!(flags, -1);
assert_ne!(flags & libc::FD_CLOEXEC, 0);
}
#[test]
fn device_pool_remaps_deterministically() {
use super::DevicePool;
let raw_root = 0x20; let raw_proc_run1 = 3_145_792; let raw_proc_run2 = 3_145_788;
let mut pool1 = DevicePool::new();
let root1 = pool1.determinize(raw_root);
let proc1 = pool1.determinize(raw_proc_run1);
let mut pool2 = DevicePool::new();
let root2 = pool2.determinize(raw_root);
let proc2 = pool2.determinize(raw_proc_run2);
assert_eq!(root1, root2);
assert_eq!(proc1, proc2);
assert_ne!(root1, proc1);
assert_eq!(root1, 1);
assert_eq!(proc1, 2);
assert_eq!(pool1.determinize(raw_root), root1);
assert_eq!(pool1.determinize(raw_proc_run1), proc1);
}
#[test]
fn mountinfo_prepopulation_is_reused_by_later_stat_observations() {
use super::DevicePool;
let mountinfo_devices = [libc::makedev(8, 1), libc::makedev(0, 44)];
let mut pool = DevicePool::new();
let rendered = mountinfo_devices
.into_iter()
.map(|raw| pool.determinize(raw))
.collect::<Vec<_>>();
assert_eq!(rendered, [1, 2]);
assert_eq!(pool.determinize(libc::makedev(0, 44)), rendered[1]);
assert_eq!(pool.determinize(libc::makedev(8, 1)), rendered[0]);
assert_eq!(pool.determinize(libc::makedev(259, 7)), 3);
}
#[tokio::test]
async fn late_futex_rpc_after_thread_removal_returns_eintr() {
let config = Config {
sequentialize_threads: true,
..Config::default()
};
let state = GlobalState::initialize(&config, false);
let dettid = DetTid::from_raw(17);
let detpid = DetPid::from_raw(17);
let response = state
.recv_futex_action(
RpcIncarnation {
dettid,
mm: MmId::initial(detpid),
},
FutexAction::WaitRequest(None),
FutexID::private(MmId::initial(detpid), 0x1000),
0,
u32::MAX,
)
.await;
assert!(matches!(
response,
Some(SchedValue::Value(value)) if value == nix::errno::Errno::EINTR as u64
));
}
#[tokio::test]
async fn late_child_tid_wait_after_exit_returns_spurious_wake() {
let config = Config {
sequentialize_threads: true,
..Config::default()
};
let state = GlobalState::initialize(&config, false);
let detpid = DetPid::from_raw(17);
let child = DetTid::from_raw(18);
let futex = FutexID::private(MmId::initial(detpid), 0x1000);
state.sched.lock().unwrap().next_turns.insert(
detpid,
ThreadNextTurn {
dettid: detpid,
child_tid_addr: 0,
req: Ivar::new(),
resp: Ivar::new(),
protocol: Default::default(),
},
);
state
.sched
.lock()
.unwrap()
.wake_futex_child_cleartid(futex, child);
assert!(
state
.sched
.lock()
.unwrap()
.child_tid_was_cleared(futex, child.as_raw())
);
assert!(
!state
.sched
.lock()
.unwrap()
.child_tid_was_cleared(futex, child.as_raw() + 1)
);
let response = state
.recv_futex_action(
RpcIncarnation {
dettid: detpid,
mm: MmId::initial(detpid),
},
FutexAction::WaitRequest(None),
futex,
child.as_raw(),
u32::MAX,
)
.await;
assert!(matches!(response, Some(SchedValue::Value(0))));
assert!(state.sched.lock().unwrap().blocked.futex_waiters.is_empty());
}
#[tokio::test]
async fn late_resource_request_after_logical_kill_is_cancelled() {
let (config, state, dettid, detpid) = cancellation_test_state();
install_test_registration(&state, dettid, Ivar::new());
let mut current_time = DetTime::new(&config);
current_time.add_syscall();
state.global_time.lock().unwrap().update_global_time(
dettid,
current_time.as_nanos(),
LogicalTime::ZERO,
);
state
.sched
.lock()
.unwrap()
.logically_kill_thread(&dettid, &detpid, MmId::initial(detpid));
let (global_before, thread_before) = {
let global_time = state.global_time.lock().unwrap();
(global_time.as_nanos(), global_time.threads_time(dettid))
};
let mut late_time = current_time;
late_time.add_syscall();
let response = state
.receive_rpc(
reverie::Tid::from_raw(dettid.as_raw()),
(
late_time,
MmId::initial(dettid),
GlobalRequest::RequestResources(Resources::new(dettid), detpid),
),
)
.await;
assert_eq!(response, (None, GlobalResponse::ThreadExited));
assert!(!state.sched.lock().unwrap().next_turns.contains_key(&dettid));
let global_time = state.global_time.lock().unwrap();
assert_eq!(global_time.as_nanos(), global_before);
assert_eq!(global_time.threads_time(dettid), thread_before);
}
#[tokio::test]
async fn duplicate_deregistration_is_acknowledged_without_clock_or_scheduler_mutation() {
let (config, state, dettid, detpid) = cancellation_test_state();
install_test_registration(&state, dettid, Ivar::new());
let mut current_time = DetTime::new(&config);
current_time.add_syscall();
state.global_time.lock().unwrap().update_global_time(
dettid,
current_time.as_nanos(),
LogicalTime::ZERO,
);
state
.sched
.lock()
.unwrap()
.logically_kill_thread(&dettid, &detpid, MmId::initial(detpid));
let (global_before, thread_before) = {
let global_time = state.global_time.lock().unwrap();
(global_time.as_nanos(), global_time.threads_time(dettid))
};
let mut late_time = current_time;
late_time.add_syscall();
let mut final_stats = TimesliceStats::default();
final_stats.record(7);
let first_response = state
.receive_rpc(
reverie::Tid::from_raw(dettid.as_raw()),
(
late_time.clone(),
MmId::initial(dettid),
GlobalRequest::DeregisterThread(ThreadDeregistration {
thread_start_entered: true,
dettid,
detpid,
mm: MmId::initial(detpid),
timeslice_stats: final_stats,
syscall_count: 17,
chaos_epochs: Vec::new(),
}),
),
)
.await;
assert_eq!(first_response, (None, GlobalResponse::DeregisterThread(())));
assert_eq!(
state
.sched
.lock()
.unwrap()
.per_thread_timeslice
.get(&dettid),
Some(&final_stats)
);
assert_eq!(
state.sched.lock().unwrap().per_thread_syscalls.get(&dettid),
Some(&17)
);
late_time.add_syscall();
let duplicate_response = state
.receive_rpc(
reverie::Tid::from_raw(dettid.as_raw()),
(
late_time,
MmId::initial(dettid),
GlobalRequest::DeregisterThread(ThreadDeregistration {
thread_start_entered: true,
dettid,
detpid,
mm: MmId::initial(detpid),
timeslice_stats: final_stats,
syscall_count: 99,
chaos_epochs: Vec::new(),
}),
),
)
.await;
assert_eq!(
duplicate_response,
(None, GlobalResponse::DeregisterThread(()))
);
assert_eq!(
state
.sched
.lock()
.unwrap()
.per_thread_timeslice
.get(&dettid),
Some(&final_stats)
);
assert_eq!(
state.sched.lock().unwrap().per_thread_syscalls.get(&dettid),
Some(&17),
"a duplicate deregistration must not double-count or replace final accounting"
);
let summary = state
.sched
.lock()
.unwrap()
.generate_partial_run_summary(None)
.unwrap();
assert_eq!(summary.syscalls, Some(17));
assert!(!state.sched.lock().unwrap().next_turns.contains_key(&dettid));
let global_time = state.global_time.lock().unwrap();
assert_eq!(global_time.as_nanos(), global_before);
assert_eq!(global_time.threads_time(dettid), thread_before);
}
#[tokio::test]
async fn child_registration_fails_closed_for_a_tombstoned_tid() {
let (config, state, parent, detpid) = cancellation_test_state();
install_test_registration(&state, parent, Ivar::new());
let child = DetTid::from_raw(18);
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(parent, child, false);
install_test_registration(&state, child, Ivar::new());
state
.sched
.lock()
.unwrap()
.logically_kill_thread(&child, &detpid, MmId::initial(detpid));
let response = state
.receive_rpc(
reverie::Tid::from_raw(parent.as_raw()),
(
DetTime::new(&config),
MmId::initial(parent),
GlobalRequest::CreateChildThread(
child,
detpid,
0,
None,
libc::SIGCHLD,
None,
Some(DEFAULT_PRIORITY),
),
),
)
.await;
assert_eq!(response, (None, GlobalResponse::ThreadExited));
let scheduler = state.sched.lock().unwrap();
assert!(scheduler.thread_is_logically_killed(child));
assert!(!scheduler.next_turns.contains_key(&child));
assert!(!scheduler.priorities.contains_key(&child));
assert!(!state.global_time.lock().unwrap().contains_thread(parent));
}
#[tokio::test]
async fn pending_resource_request_woken_by_logical_kill_is_terminal() {
let (_, state, dettid, detpid) = cancellation_test_state();
let request_seen = Ivar::new();
install_test_registration(&state, dettid, request_seen.clone());
let request = state.recv_request_resources(
reverie::Tid::from_raw(dettid.as_raw()),
detpid,
Resources::new(dettid),
None,
);
let kill_after_request = async {
while request_seen.try_read().is_none() {
tokio::task::yield_now().await;
}
state.sched.lock().unwrap().logically_kill_thread(
&dettid,
&detpid,
MmId::initial(detpid),
);
};
let (response, ()) = tokio::join!(request, kill_after_request);
assert_eq!(response, (SchedulerRpcResult::ThreadExited, None));
}
#[tokio::test]
async fn trace_replay_yield_propagates_terminal_scheduler_cancellation() {
let dettid = DetTid::from_raw(17);
let detpid = DetPid::from_raw(17);
let next_tid = DetTid::from_raw(18);
let event = SchedEvent {
dettid,
op: Op::OtherInstructions,
count: 1,
start_rip: None,
end_rip: None,
end_time: Some(LogicalTime::from_nanos(1)),
};
let next_event = SchedEvent {
dettid: next_tid,
end_time: Some(LogicalTime::from_nanos(2)),
..event.clone()
};
let trace_file = tempfile::NamedTempFile::new().unwrap();
std::fs::write(
trace_file.path(),
PreemptionRecord::from_sched_events(vec![event.clone(), next_event]).to_string(),
)
.unwrap();
let config = Config {
sequentialize_threads: true,
cancel_killed_thread_rpcs: true,
replay_schedule_from: Some(trace_file.path().to_path_buf()),
..Config::default()
};
let state = GlobalState::initialize(&config, false);
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(dettid, dettid, true);
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(dettid, next_tid, false);
let request_seen = Ivar::new();
install_test_registration(&state, dettid, request_seen.clone());
install_test_registration(&state, next_tid, Ivar::new());
let replay = state.recv_trace_schedevent(event, detpid, MmId::initial(detpid), false);
let kill_after_replay_yield = async {
while request_seen.try_read().is_none() {
tokio::task::yield_now().await;
}
state.sched.lock().unwrap().logically_kill_thread(
&dettid,
&detpid,
MmId::initial(detpid),
);
};
let (response, ()) = tokio::time::timeout(Duration::from_secs(1), async {
tokio::join!(replay, kill_after_replay_yield)
})
.await
.expect("trace replay cancellation did not terminate the pending scheduler RPC");
assert_eq!(response, SchedulerRpcResult::ThreadExited);
}
#[tokio::test]
async fn pending_start_request_woken_by_logical_kill_is_terminal() {
let (config, state, dettid, detpid) = cancellation_test_state();
let request_seen = Ivar::new();
install_test_registration(&state, dettid, request_seen.clone());
let request = state.receive_rpc(
reverie::Tid::from_raw(dettid.as_raw()),
(
DetTime::new(&config),
MmId::initial(dettid),
GlobalRequest::StartNewThread(dettid, detpid, None, None),
),
);
let kill_after_request = async {
while request_seen.try_read().is_none() {
tokio::task::yield_now().await;
}
state.sched.lock().unwrap().logically_kill_thread(
&dettid,
&detpid,
MmId::initial(detpid),
);
};
let (response, ()) = tokio::join!(request, kill_after_request);
assert_eq!(response, (None, GlobalResponse::ThreadExited));
assert!(!state.sched.lock().unwrap().priorities.contains_key(&dettid));
}
#[tokio::test]
async fn required_physical_thread_id_missing_is_terminal() {
let config = Config {
sequentialize_threads: true,
cancel_killed_thread_rpcs: true,
backend_requires_thread_directed_process_signals: true,
..Config::default()
};
let state = GlobalState::initialize(&config, false);
let dettid = DetTid::from_raw(17);
let detpid = DetPid::from_raw(17);
state
.sched
.lock()
.unwrap()
.thread_tree
.add_child(dettid, dettid, true);
install_test_registration(&state, dettid, Ivar::new());
let response = state
.receive_rpc(
reverie::Tid::from_raw(dettid.as_raw()),
(
DetTime::new(&config),
MmId::initial(detpid),
GlobalRequest::StartNewThread(dettid, detpid, None, None),
),
)
.await;
assert_eq!(response, (None, GlobalResponse::ThreadExited));
let scheduler = state.sched.lock().unwrap();
assert!(!scheduler.next_turns.contains_key(&dettid));
}
#[tokio::test]
async fn tombstoned_timer_registration_cannot_mutate_scheduler_state() {
let (_, state, dettid, detpid) = cancellation_test_state();
install_test_registration(&state, dettid, Ivar::new());
state
.sched
.lock()
.unwrap()
.logically_kill_thread(&dettid, &detpid, MmId::initial(detpid));
let now = LogicalTime::from_nanos(100);
let alarm = state
.recv_register_alarm(
detpid,
RpcIncarnation {
dettid,
mm: MmId::initial(detpid),
},
now,
LogicalTime::from_nanos(10),
LogicalTime::ZERO,
SigWrapper::from(Signal::SIGALRM),
)
.await;
assert_eq!(alarm, SchedulerRpcResult::ThreadExited);
let posix = state
.recv_register_posix_timer(
detpid,
RpcIncarnation {
dettid,
mm: MmId::initial(detpid),
},
1,
Some(now + LogicalTime::from_nanos(10)),
LogicalTime::ZERO,
SigWrapper::from(Signal::SIGALRM),
)
.await;
assert_eq!(posix, SchedulerRpcResult::ThreadExited);
assert!(state.sched.lock().unwrap().blocked.timed_waiters.is_empty());
}
#[tokio::test]
async fn pending_futex_request_woken_by_logical_kill_is_terminal() {
let (config, state, dettid, detpid) = cancellation_test_state();
install_test_registration(&state, dettid, Ivar::new());
let request = state.receive_rpc(
reverie::Tid::from_raw(dettid.as_raw()),
(
DetTime::new(&config),
MmId::initial(dettid),
GlobalRequest::FutexAction(
dettid,
FutexAction::WaitRequest(None),
FutexID::private(MmId::initial(detpid), 0x1000),
0,
u32::MAX,
),
),
);
let kill_after_wait = async {
while state.sched.lock().unwrap().blocked.futex_waiters.is_empty() {
tokio::task::yield_now().await;
}
state.sched.lock().unwrap().logically_kill_thread(
&dettid,
&detpid,
MmId::initial(detpid),
);
};
let (response, ()) = tokio::time::timeout(Duration::from_secs(1), async {
tokio::join!(request, kill_after_wait)
})
.await
.expect("futex teardown did not wake the blocked RPC");
assert_eq!(response, (None, GlobalResponse::ThreadExited));
}
#[tokio::test]
async fn parent_continue_propagates_terminal_scheduler_cancellation() {
let (config, state, parent, detpid) = cancellation_test_state();
let parent_request = Ivar::new();
install_test_registration(&state, parent, parent_request.clone());
let child = DetTid::from_raw(18);
let physical_pid = std::process::id() as i32;
let physical_tid = unsafe { libc::syscall(libc::SYS_gettid) as i32 };
let request = state.receive_rpc(
reverie::Tid::from_raw(parent.as_raw()),
(
DetTime::new(&config),
MmId::initial(parent),
GlobalRequest::CreateChildThread(
child,
detpid,
0,
None,
libc::SIGCHLD,
Some((physical_pid, physical_tid)),
Some(DEFAULT_PRIORITY),
),
),
);
let kill_after_parent_parks = async {
while parent_request.try_read().is_none() {
tokio::task::yield_now().await;
}
let mut scheduler = state.sched.lock().unwrap();
let (registered_mm, registered_pid, registered_tid) = scheduler
.physical_thread_identity(child)
.expect("parent registration must install the child pidfd before continuing");
assert_eq!(
registered_mm,
MmId::for_clone(MmId::initial(parent), child, false)
);
assert_eq!(
(registered_pid, registered_tid),
(physical_pid, physical_tid)
);
scheduler.logically_kill_thread(&parent, &detpid, MmId::initial(detpid));
};
let (response, ()) = tokio::join!(request, kill_after_parent_parks);
assert_eq!(response, (None, GlobalResponse::ThreadExited));
}
#[test]
fn unsupported_syscall_warning_is_sorted_and_aggregated() {
let syscalls = BTreeSet::from([
"vmsplice".to_owned(),
"getppid".to_owned(),
"getppid".to_owned(),
]);
assert_eq!(
format_unsupported_syscall_warning(&syscalls).as_deref(),
Some("syscalls getppid,vmsplice used but not yet supported")
);
assert_eq!(format_unsupported_syscall_warning(&BTreeSet::new()), None);
}
#[tokio::test]
async fn abnormal_cleanup_cancels_an_unstarted_scheduler() {
let config = Config {
sequentialize_threads: true,
..Config::default()
};
let mut state = GlobalState::initialize(&config, true);
state.cancel_internal_scheduler().await;
let summary_path = None;
let cleanup = state.clean_up(false, &summary_path);
assert!(
tokio::time::timeout(Duration::from_millis(100), cleanup)
.await
.is_ok(),
"cleanup waited for a scheduler whose guest never registered"
);
}
#[tokio::test]
async fn abnormal_cleanup_cancels_a_registered_scheduler() {
let config = Config {
sequentialize_threads: true,
..Config::default()
};
let mut state = GlobalState::initialize(&config, true);
let dettid = DetTid::from_raw(1);
{
let mut scheduler = state.sched.lock().unwrap();
scheduler.priorities.insert(dettid, DEFAULT_PRIORITY);
scheduler.next_turns.insert(
dettid,
ThreadNextTurn {
dettid,
child_tid_addr: 0,
req: Ivar::new(),
resp: Ivar::new(),
protocol: Default::default(),
},
);
scheduler.runqueue_push_back(dettid);
scheduler.started_up.put(());
}
tokio::task::yield_now().await;
state.cancel_internal_scheduler().await;
let summary_path = None;
assert!(
tokio::time::timeout(
Duration::from_millis(100),
state.clean_up(false, &summary_path),
)
.await
.is_ok(),
"cleanup waited after cancelling a registered scheduler"
);
}
#[test]
fn det_inodes_are_minted_not_passed_through() {
use crate::types::DetInode;
let mut pool = super::InodePool::new();
let t = LogicalTime::from_nanos(0);
let seen = super::ObservedMtime::Unobserved;
let host_a = 221_742_951; let host_b = 998_877_665;
let (a, _) = pool.add_inode(host_a, seen, t);
let (b, _) = pool.add_inode(host_b, seen, t);
assert_ne!(a.as_raw(), host_a, "det inode must not be the host inode");
assert_ne!(b.as_raw(), host_b, "det inode must not be the host inode");
assert_eq!(a, DetInode::mint(1), "minting starts at 1");
assert_eq!(b, DetInode::mint(2), "minting is monotonic");
let (a_again, _) = pool.add_inode(host_a, seen, t);
assert_eq!(a, a_again, "mapping must be stable per host inode");
}
#[test]
fn only_exact_canonical_host_mtimes_are_kept() {
use super::ObservedMtime;
assert_eq!(
ObservedMtime::from_host_mtime(1, 0),
ObservedMtime::Canonical(LogicalTime::from_secs(1))
);
assert_eq!(
ObservedMtime::from_host_mtime(0, 0),
ObservedMtime::Canonical(LogicalTime::from_secs(0))
);
for (secs, nanos) in [(1, 5), (0, 1), (2, 0), (-1, 0), (1_600_000_000, 0)] {
assert_eq!(
ObservedMtime::from_host_mtime(secs, nanos),
ObservedMtime::HostSpecific,
"{secs}.{nanos:09} is a real timestamp, not a canonical one"
);
}
}
#[test]
fn first_seen_mtime_is_resolved_by_the_first_stat() {
use super::ObservedMtime;
let epoch = LogicalTime::from_secs(1_798_761_600);
let canonical = ObservedMtime::Canonical(LogicalTime::from_secs(1));
let mut pool = super::InodePool::new();
let (_, mtime) = pool.add_inode(10, ObservedMtime::HostSpecific, epoch);
assert_eq!(mtime, epoch);
let (_, mtime) = pool.add_inode(11, canonical, epoch);
assert_eq!(mtime, LogicalTime::from_secs(1));
let (_, mtime) = pool.add_inode(12, ObservedMtime::Unobserved, epoch);
assert_eq!(mtime, epoch, "an unresolved mtime reads as the epoch");
let (_, mtime) = pool.add_inode(12, canonical, epoch);
assert_eq!(mtime, LogicalTime::from_secs(1));
let (_, mtime) = pool.add_inode(10, canonical, epoch);
assert_eq!(mtime, epoch);
let (_, mtime) = pool.add_inode(11, ObservedMtime::HostSpecific, epoch);
assert_eq!(mtime, LogicalTime::from_secs(1));
}
}
#[cfg(test)]
mod robust_exit_clock_tests {
use std::sync::Mutex;
use std::task::Poll;
use nix::sys::signal::Signal;
use reverie::ExitStatus;
use reverie::GlobalRPC;
use reverie::GlobalTool;
use reverie::Tid;
use reverie::Tool;
use super::GlobalRequest;
use super::GlobalResponse;
use super::GlobalState;
use crate::Detcore;
use crate::ThreadState;
use crate::config::Config;
use crate::ivar::Ivar;
use crate::resources::Resources;
use crate::scheduler::DEFAULT_PRIORITY;
use crate::scheduler::SkipTurn;
use crate::scheduler::ThreadNextTurn;
use crate::tool_local::RobustListExit;
use crate::tool_local::RobustListWake;
use crate::types::*;
#[derive(Debug, PartialEq)]
struct WakeObservation {
wakes: Vec<(DetTid, FutexID)>,
counts: Vec<u64>,
clocks: serde_json::Value,
turn: u64,
}
#[derive(Debug)]
struct RpcObservation {
sender: DetTid,
kind: &'static str,
accepted: bool,
clocks: serde_json::Value,
queued: Vec<DetTid>,
waiters: usize,
turn: u64,
}
struct ExitRpc<'a> {
state: &'a GlobalState,
sender: DetTid,
observations: &'a Mutex<Vec<WakeObservation>>,
rpc_observations: &'a Mutex<Vec<RpcObservation>>,
}
#[reverie::tool]
impl GlobalRPC<GlobalState> for ExitRpc<'_> {
async fn send_rpc(
&self,
request: <GlobalState as GlobalTool>::Request,
) -> <GlobalState as GlobalTool>::Response {
let kind = match &request.2 {
GlobalRequest::RobustListWakes(wakes) if wakes.is_empty() => "empty-wake",
GlobalRequest::RobustListWakes(_) => "wake",
GlobalRequest::DeregisterThread(_) => "deregister",
_ => panic!("unexpected exit RPC: {:?}", request.2),
};
let wakes = match &request.2 {
GlobalRequest::RobustListWakes(wakes) if !wakes.is_empty() => Some(wakes.clone()),
_ => None,
};
let response = self
.state
.receive_rpc(Tid::from_raw(self.sender.as_raw()), request)
.await;
let clocks = serde_json::to_value(&*self.state.global_time.lock().unwrap()).unwrap();
{
let sched = self.state.sched.lock().unwrap();
self.rpc_observations.lock().unwrap().push(RpcObservation {
sender: self.sender,
kind,
accepted: !matches!(response.1, GlobalResponse::ThreadExited),
clocks: clocks.clone(),
queued: sched.run_queue.tids().copied().collect(),
waiters: sched
.blocked
.futex_waiters
.values()
.map(Vec::len)
.sum::<usize>(),
turn: sched.turn,
});
}
if let Some(wakes) = wakes {
let GlobalResponse::RobustListWakes(counts) = &response.1 else {
panic!("a complete admitted exit batch was refused: {response:?}");
};
let clocks =
serde_json::to_value(&*self.state.global_time.lock().unwrap()).unwrap();
let turn = self.state.sched.lock().unwrap().turn;
self.observations.lock().unwrap().push(WakeObservation {
wakes,
counts: counts.clone(),
clocks,
turn,
});
}
response
}
fn config(&self) -> &Config {
&self.state.cfg
}
}
struct Fixture {
state: GlobalState,
tool: Detcore,
owners: [ThreadState<()>; 2],
initial_clocks: [DetTime; 2],
waiters: [DetTid; 2],
peer: DetTid,
futexes: [FutexID; 2],
observations: Mutex<Vec<WakeObservation>>,
rpc_observations: Mutex<Vec<RpcObservation>>,
}
impl Fixture {
fn new(
reason: RobustListExit,
equal_clocks: bool,
empty_owner: Option<usize>,
cancel_killed_thread_rpcs: bool,
) -> Self {
let config = Config {
sequentialize_threads: true,
cancel_killed_thread_rpcs,
..Config::default()
};
let state = GlobalState::initialize(&config, false);
let parent = DetTid::from_raw(1);
let leader = DetTid::from_raw(17);
let worker = DetTid::from_raw(18);
let waiters = [DetTid::from_raw(21), DetTid::from_raw(23)];
let peer = DetTid::from_raw(25);
let mm = MmId::initial(leader);
let mut first = ThreadState::new(leader, &config, ());
first.detpid = Some(leader);
let mut second = first.clone();
second.dettid = worker;
let mut inherited = DetTime::new(&config);
inherited.add_syscall_with_cost(1_000);
let first_initial = inherited.clone_for_child();
inherited.add_syscall_with_cost(370);
let second_initial = inherited.clone_for_child();
first.thread_logical_time = first_initial.clone();
second.thread_logical_time = second_initial.clone();
first
.thread_logical_time
.add_syscall_with_cost(if equal_clocks { 407 } else { 37 });
second.thread_logical_time.add_syscall_with_cost(37);
first.record_robust_list_head(Some(0x404100));
second.record_robust_list_head(Some(0x404200));
let object = SharedMemoryObjectId::Anonymous {
origin: MmId::initial(parent),
sequence: 1,
};
let futexes = [FutexID::shared(object, 0), FutexID::shared(object, 8)];
first.stage_robust_list_wakes(
reason,
vec![
(
worker,
if empty_owner == Some(1) {
Vec::new()
} else {
vec![RobustListWake { futex: futexes[1] }]
},
),
(
leader,
if empty_owner == Some(0) {
Vec::new()
} else {
vec![RobustListWake { futex: futexes[0] }]
},
),
],
);
{
let mut sched = state.sched.lock().unwrap();
sched.thread_tree.add_child(parent, parent, true);
sched.thread_tree.add_child(parent, leader, true);
sched.thread_tree.add_child(leader, worker, false);
for tid in [waiters[0], waiters[1], peer] {
sched.thread_tree.add_child(parent, tid, true);
}
for tid in [leader, worker, waiters[0], waiters[1], peer] {
sched.priorities.insert(tid, DEFAULT_PRIORITY);
sched.next_turns.insert(
tid,
ThreadNextTurn {
dettid: tid,
child_tid_addr: 0,
req: Ivar::new(),
resp: Ivar::new(),
protocol: Default::default(),
},
);
sched.install_test_exec_incarnation(
tid,
if tid == leader || tid == worker {
mm
} else {
MmId::initial(tid)
},
);
}
sched.runqueue_push_back(leader);
sched.runqueue_push_back(worker);
sched.runqueue_push_back(peer);
sched.next_turns[&peer].req.put(Ok(Resources::new(peer)));
for (waiter, futex) in waiters.into_iter().zip(futexes) {
sched.sleep_futex_waiter(&waiter, futex, None, u32::MAX);
}
}
{
let mut time = state.global_time.lock().unwrap();
for (tid, clock) in [(leader, &first_initial), (worker, &second_initial)] {
time.update_global_time(tid, clock.as_nanos(), clock.inherited_nanos());
}
}
let tool = Detcore::new(Tid::from_raw(leader.as_raw()), &config);
Self {
state,
tool,
owners: [first, second],
initial_clocks: [first_initial, second_initial],
waiters,
peer,
futexes,
observations: Mutex::new(Vec::new()),
rpc_observations: Mutex::new(Vec::new()),
}
}
fn assert_clocks(&self, completed: &[usize]) -> serde_json::Value {
let time = self.state.global_time.lock().unwrap();
let snapshot = serde_json::to_value(&*time).unwrap();
let epoch = DetTime::new(&self.state.cfg).as_nanos();
let mut expected = epoch;
for (index, owner) in self.owners.iter().enumerate() {
let clock = if completed.contains(&index) {
&owner.thread_logical_time
} else {
&self.initial_clocks[index]
};
assert_eq!(time.threads_time(owner.dettid), clock.as_nanos());
assert_eq!(
snapshot["inherited_time"][owner.dettid.as_raw().to_string()],
serde_json::to_value(clock.inherited_nanos()).unwrap()
);
expected = expected + (clock.as_nanos() - epoch - clock.inherited_nanos());
}
assert_eq!(
time.as_nanos(),
expected,
"only each owner's own uninherited work contributes"
);
snapshot
}
async fn exit(&self, index: usize, status: ExitStatus) {
self.exit_thread(self.owners[index].clone(), status).await;
}
async fn exit_thread(&self, thread: ThreadState<()>, status: ExitStatus) {
let rpc = ExitRpc {
state: &self.state,
sender: thread.dettid,
observations: &self.observations,
rpc_observations: &self.rpc_observations,
};
let exit =
self.tool
.on_exit_thread(Tid::from_raw(rpc.sender.as_raw()), &rpc, thread, status);
let mut exit = std::pin::pin!(exit);
assert!(
matches!(futures::poll!(exit.as_mut()), Poll::Ready(Ok(()))),
"exit callback yielded before its accounting and cleanup completed"
);
}
async fn nonmember_exit(&self, raw_tid: i32) {
let mut thread = self.owners[0].clone();
thread.dettid = DetTid::from_raw(raw_tid);
thread.thread_logical_time = DetTime::new(&self.state.cfg);
let tid = thread.dettid;
{
let mut sched = self.state.sched.lock().unwrap();
sched
.thread_tree
.add_child(self.owners[0].dettid, tid, false);
sched.priorities.insert(tid, DEFAULT_PRIORITY);
sched.next_turns.insert(
tid,
ThreadNextTurn {
dettid: tid,
child_tid_addr: 0,
req: Ivar::new(),
resp: Ivar::new(),
protocol: Default::default(),
},
);
sched.install_test_exec_incarnation(tid, thread.mm_id);
sched.runqueue_push_back(tid);
}
let before = self.rpc_observations.lock().unwrap().len();
self.exit_thread(thread, ExitStatus::Exited(0)).await;
let observations = self.rpc_observations.lock().unwrap();
assert_eq!(
observations.len(),
before + 1,
"a nonmember sent an acknowledgement or wake RPC"
);
assert_eq!(observations[before].sender, tid);
assert_eq!(observations[before].kind, "deregister");
assert!(observations[before].accepted);
}
}
#[tokio::test]
async fn robust_exit_acknowledgements_preserve_global_action_order_and_eligibility() {
for order in [[0, 1], [1, 0]] {
for empty_owner in [0, 1] {
let f = Fixture::new(RobustListExit::ExitGroup, false, Some(empty_owner), false);
f.nonmember_exit(30).await;
assert!(f.observations.lock().unwrap().is_empty());
let queue_before_first: Vec<_> = f
.state
.sched
.lock()
.unwrap()
.run_queue
.tids()
.copied()
.collect();
let first_start = f.rpc_observations.lock().unwrap().len();
f.exit(order[0], ExitStatus::Exited(0)).await;
let first_clocks = f.assert_clocks(&[order[0]]);
{
let observations = f.rpc_observations.lock().unwrap();
let actions = &observations[first_start..];
assert_eq!(
actions.iter().map(|a| a.kind).collect::<Vec<_>>(),
["empty-wake", "deregister"]
);
assert!(actions.iter().all(|a| a.accepted
&& a.sender == f.owners[order[0]].dettid
&& a.clocks == first_clocks
&& a.waiters == 2
&& a.turn == 0));
assert_eq!(
actions[0].queued, queue_before_first,
"empty acknowledgement changed scheduler eligibility"
);
}
f.nonmember_exit(31).await;
f.assert_clocks(&[order[0]]);
assert!(
f.observations.lock().unwrap().is_empty(),
"a nonmember completed the physical-exit barrier"
);
let queue_before_last: Vec<_> = f
.state
.sched
.lock()
.unwrap()
.run_queue
.tids()
.copied()
.collect();
let last_start = f.rpc_observations.lock().unwrap().len();
f.exit(order[1], ExitStatus::Exited(0)).await;
let final_clocks = f.assert_clocks(&[0, 1]);
{
let observations = f.rpc_observations.lock().unwrap();
let actions = &observations[last_start..];
assert_eq!(
actions.iter().map(|a| a.kind).collect::<Vec<_>>(),
["empty-wake", "wake", "deregister"]
);
assert!(actions.iter().all(|a| a.accepted
&& a.sender == f.owners[order[1]].dettid
&& a.clocks == final_clocks
&& a.turn == 0));
assert_eq!(actions[0].waiters, 2);
assert_eq!(actions[1].waiters, 1);
assert_eq!(actions[0].queued, queue_before_last);
assert_eq!(
actions[1].queued, queue_before_last,
"wake bypassed deferred admission"
);
}
f.nonmember_exit(32).await;
f.assert_clocks(&[0, 1]);
assert_eq!(f.observations.lock().unwrap().len(), 1);
let observations = f.rpc_observations.lock().unwrap();
for owner in &f.owners {
assert_eq!(
observations
.iter()
.filter(|a| a.sender == owner.dettid && a.kind == "empty-wake")
.count(),
1,
"each unique matching owner, including an empty-wake owner, must acknowledge before the batch clears"
);
}
}
}
}
#[tokio::test]
async fn robust_exit_callbacks_preserve_owner_clocks_across_arrival_orders() {
for reason in [
RobustListExit::ExitGroup,
RobustListExit::Signal(libc::SIGTERM),
] {
for equal_clocks in [false, true] {
for empty_owner in [None, Some(0), Some(1)] {
let mut results = Vec::new();
for order in [[0, 1], [1, 0]] {
let f = Fixture::new(reason, equal_clocks, empty_owner, false);
let status = match reason {
RobustListExit::ExitGroup => ExitStatus::Exited(0),
RobustListExit::Signal(_) => {
ExitStatus::Signaled(Signal::SIGTERM, false)
}
};
f.exit(order[0], status).await;
f.assert_clocks(&[order[0]]);
assert!(f.observations.lock().unwrap().is_empty());
{
let sched = f.state.sched.lock().unwrap();
assert_eq!(
sched
.blocked
.futex_waiters
.values()
.map(Vec::len)
.sum::<usize>(),
2
);
assert_eq!(sched.turn, 0);
}
{
let last = Err(SkipTurn);
let turn = crate::scheduler::do_a_turn_blocking(
f.state.sched.clone(),
f.state.global_time.clone(),
&last,
);
let mut turn = std::pin::pin!(turn);
assert!(matches!(futures::poll!(turn.as_mut()), Poll::Pending));
}
f.exit(order[0], status).await;
f.assert_clocks(&[order[0]]);
assert!(
f.observations.lock().unwrap().is_empty(),
"duplicate physical exit released an incomplete group"
);
f.exit(order[1], status).await;
let clocks = f.assert_clocks(&[0, 1]);
{
let observations = f.observations.lock().unwrap();
assert_eq!(observations.len(), 1);
let expected: Vec<_> = (0..2)
.filter(|i| Some(*i) != empty_owner)
.map(|i| (f.owners[i].dettid, f.futexes[i]))
.collect();
assert_eq!(observations[0].wakes, expected);
assert_eq!(observations[0].counts, vec![1; expected.len()]);
assert_eq!(
observations[0].clocks, clocks,
"all owner clocks must be accounted before wake admission"
);
assert_eq!(observations[0].turn, 0);
}
f.exit(order[1], status).await;
assert_eq!(f.assert_clocks(&[0, 1]), clocks);
assert_eq!(
f.observations.lock().unwrap().len(),
1,
"repeated cleanup emitted a second batch"
);
{
let sched = f.state.sched.lock().unwrap();
assert!(!sched.run_queue.contains_tid(f.waiters[0]));
assert!(!sched.run_queue.contains_tid(f.waiters[1]));
}
let result = crate::scheduler::do_a_turn_blocking(
f.state.sched.clone(),
f.state.global_time.clone(),
&Err(SkipTurn),
)
.await;
assert!(result.is_ok());
let sched = f.state.sched.lock().unwrap();
assert_eq!(sched.turn, 1);
let queued: Vec<_> = sched.run_queue.tids().copied().collect();
for (i, waiter) in f.waiters.iter().enumerate() {
assert_eq!(queued.contains(waiter), Some(i) != empty_owner);
}
assert!(queued.contains(&f.peer));
results.push((
clocks,
queued,
std::mem::take(&mut *f.observations.lock().unwrap()),
));
}
assert_eq!(
results[0], results[1],
"callback order changed final clocks, typed wake observations or the actual scheduler drain"
);
}
}
}
}
#[tokio::test]
async fn robust_exit_callbacks_keep_rejected_or_mismatched_batches_incomplete() {
for rejection in ["old-mm", "tombstone", "wrong-signal", "normal-exit"] {
let f = Fixture::new(
RobustListExit::Signal(libc::SIGTERM),
false,
None,
rejection == "tombstone",
);
let rejected = f.owners[0].dettid;
let mm = f.owners[0].mm_id;
f.exit(1, ExitStatus::Signaled(Signal::SIGTERM, false))
.await;
let replacement_request = Ivar::new();
if rejection == "old-mm" {
let mut sched = f.state.sched.lock().unwrap();
sched.install_test_exec_incarnation(rejected, mm.for_exec(rejected));
sched.next_turns.insert(
rejected,
ThreadNextTurn {
dettid: rejected,
child_tid_addr: 0,
req: replacement_request.clone(),
resp: Ivar::new(),
protocol: Default::default(),
},
);
let mut replacement_time = f.owners[0].thread_logical_time.clone();
replacement_time.add_syscall_with_cost(500);
f.state.global_time.lock().unwrap().update_global_time(
rejected,
replacement_time.as_nanos(),
replacement_time.inherited_nanos(),
);
} else if rejection == "tombstone" {
let mut sched = f.state.sched.lock().unwrap();
sched.logically_kill_thread(&rejected, &rejected, mm);
}
let before = serde_json::to_value(&*f.state.global_time.lock().unwrap()).unwrap();
let status = match rejection {
"wrong-signal" => ExitStatus::Signaled(Signal::SIGKILL, false),
"normal-exit" => ExitStatus::Exited(0),
_ => ExitStatus::Signaled(Signal::SIGTERM, false),
};
f.exit(0, status).await;
assert!(
f.observations.lock().unwrap().is_empty(),
"{rejection} released a group"
);
if rejection == "old-mm" || rejection == "tombstone" {
assert_eq!(
serde_json::to_value(&*f.state.global_time.lock().unwrap()).unwrap(),
before,
"rejected acknowledgement changed clock state"
);
}
let sched = f.state.sched.lock().unwrap();
assert_eq!(
sched
.blocked
.futex_waiters
.values()
.map(Vec::len)
.sum::<usize>(),
2,
"rejected group lost a real waiter"
);
assert!(
f.waiters
.iter()
.all(|tid| !sched.run_queue.contains_tid(*tid))
);
if rejection == "old-mm" {
assert!(sched.rpc_incarnation_matches(rejected, mm.for_exec(rejected)));
assert_eq!(
sched.next_turns[&rejected].req, replacement_request,
"old cleanup destroyed replacement registration"
);
}
}
}
#[tokio::test]
#[should_panic(expected = "Attempted to update tid 17 time")]
async fn robust_exit_clock_ack_still_refuses_a_backwards_owner_sample() {
let f = Fixture::new(RobustListExit::ExitGroup, false, None, false);
let owner = &f.owners[0];
let mut later = owner.thread_logical_time.clone();
later.add_syscall_with_cost(1);
f.state.global_time.lock().unwrap().update_global_time(
owner.dettid,
later.as_nanos(),
later.inherited_nanos(),
);
f.exit(0, ExitStatus::Exited(0)).await;
}
#[tokio::test]
async fn backend_failure_preserves_consuming_robust_exit_clock_accounting() {
for order in [[0, 1], [1, 0]] {
let f = Fixture::new(RobustListExit::ExitGroup, false, None, false);
let selected = {
let mut sched = f.state.sched.lock().unwrap();
sched.next_turns.get_mut(&f.peer).unwrap().req = Ivar::new();
sched.select_test_turn().unwrap()
};
let responses = {
let sched = f.state.sched.lock().unwrap();
f.waiters.map(|tid| sched.next_turns[&tid].resp.clone())
};
let mut turn = std::pin::pin!(crate::scheduler::finish_selected_turn(
f.state.sched.clone(),
f.state.global_time.clone(),
selected.0,
selected.1,
selected.2,
));
assert!(futures::poll!(turn.as_mut()).is_pending());
f.state.report_backend_failure(reverie::BackendFailure {
pid: Tid::from_raw(17),
tid: Tid::from_raw(18),
phase: "native robust cleanup control",
});
for index in order {
f.exit(index, ExitStatus::Exited(0)).await;
}
assert!(matches!(futures::poll!(turn.as_mut()), Poll::Ready(Err(_))));
f.assert_clocks(&[0, 1]);
let observations = f.observations.lock().unwrap();
assert_eq!(observations.len(), 1, "one complete batch");
assert_eq!(observations[0].counts, vec![1, 1]);
assert!(
responses
.iter()
.all(|response| response.try_read().is_none())
);
let mut sched = f.state.sched.lock().unwrap();
assert_eq!(sched.turn, 0);
for owner in &f.owners {
assert!(!sched.next_turns.contains_key(&owner.dettid));
assert!(!sched.note_deregistration_accounted(owner.dettid));
}
}
}
}