use std::num::NonZeroUsize;
use std::time::Duration;
use async_trait::async_trait;
use reverie::Errno;
use reverie::Error;
use reverie::Guest;
use reverie::Stack;
use reverie::syscalls;
use reverie::syscalls::Addr;
use reverie::syscalls::Displayable;
use reverie::syscalls::MapFlags;
use reverie::syscalls::MemoryAccess;
use reverie::syscalls::ProtFlags;
use reverie::syscalls::Syscall;
use reverie::syscalls::SyscallInfo;
use reverie::syscalls::Sysno;
use reverie::syscalls::Timespec;
use reverie::syscalls::WaitPidFlag;
use crate::fd::FdType;
use crate::record_or_replay::RecordOrReplay;
use crate::resources::ExternalOpId;
use crate::resources::Permission;
use crate::resources::ResourceID;
use crate::resources::Resources;
use crate::syscalls::threads::KernelSigaction;
use crate::syscalls::threads::KernelSigset;
use crate::syscalls::threads::WaitSignalDisposition;
use crate::syscalls::threads::block_signals_for_disposition;
use crate::syscalls::threads::blocked_signal_mask;
use crate::syscalls::threads::restore_signals_after_disposition;
use crate::syscalls::threads::wait_signal_disposition;
use crate::tool_global::ResumeStatus;
use crate::tool_global::resource_request;
use crate::tool_global::thread_observe_time;
use crate::tool_global::trace_schedevent;
use crate::tool_local::Detcore;
use crate::tool_local::finish_partial_record_or_replay_write;
use crate::types::DetTid;
use crate::types::LogicalTime;
use crate::types::OpenFileId;
use crate::types::SchedEvent;
use crate::types::SyscallPhase;
impl<T: RecordOrReplay> Detcore<T> {
pub(crate) async fn refuse_unserviceable_operation<G: Guest<Self>>(
&self,
guest: &mut G,
sysno: reverie::syscalls::Sysno,
fallback: Errno,
) -> Result<i64, Error> {
if !self.cfg.panic_on_unsupported_syscalls {
return Err(fallback.into());
}
if guest.config().shutdown_on_unsupported_syscall {
crate::tool_global::unrecoverable_shutdown(
guest,
detcore_model::HERMIT_POLICY_REFUSAL_EXIT,
)
.await;
}
if guest.config().exit_on_unsupported_syscall {
return Err(Error::Tool(anyhow::Error::new(
crate::UnsupportedSyscallError(sysno),
)));
}
panic!("unserviceable operation on syscall: {sysno:?}");
}
pub async fn record_or_replay_blocking<G: Guest<Self>>(
&self,
guest: &mut G,
call: Syscall,
) -> Result<i64, Error> {
let dettid = guest.thread_state().dettid;
let op_id = ExternalOpId::new(dettid, guest.thread_state().stats.syscall_count);
self.record_or_replay_blocking_resource(guest, call, ResourceID::BlockingExternalIO(op_id))
.await
}
pub async fn record_or_replay_rt_sigsuspend<G: Guest<Self>>(
&self,
guest: &mut G,
call: syscalls::RtSigsuspend,
) -> Result<i64, Error> {
let dettid = guest.thread_state().dettid;
let op_id = ExternalOpId::new(dettid, guest.thread_state().stats.syscall_count);
self.record_or_replay_blocking_resource(
guest,
call.into(),
ResourceID::BlockingRtSigsuspend(op_id),
)
.await
}
async fn record_or_replay_blocking_resource<G: Guest<Self>>(
&self,
guest: &mut G,
call: Syscall,
blocking_resource: ResourceID,
) -> Result<i64, Error> {
let dettid = guest.thread_state().dettid;
let op_id = match &blocking_resource {
ResourceID::BlockingExternalIO(op_id) | ResourceID::BlockingRtSigsuspend(op_id) => {
*op_id
}
_ => unreachable!("blocking syscall helper requires a blocking resource"),
};
debug_assert!(
!self.cfg.sequentialize_threads || !syscall_targets_internal_fd(guest, call),
"record_or_replay_blocking (BlockingExternalIO) reached for an internal pipe fd \
on syscall {}; internal fds must use the InternalIOPolling path",
call.name()
);
{
let mut rsrcs = Resources::new(dettid);
rsrcs.insert(blocking_resource, Permission::RW);
rsrcs.fyi(call.name());
resource_request(guest, rsrcs).await;
}
tracing::trace!(
"Guest proceeding to execute potentially blocking call {}...",
call.name()
);
let res = self
.record_or_replay_preserving_tool_errors(guest, call)
.await;
{
let mut rsrcs = Resources::new(dettid);
rsrcs.insert(ResourceID::BlockedExternalContinue(op_id), Permission::RW);
rsrcs.fyi(call.name());
resource_request(guest, rsrcs).await;
}
res
}
pub async fn execute_nonblockable_fd_syscall<
G: Guest<Self>,
C: SyscallInfo + NonblockableSyscall + Into<Syscall>,
>(
&self,
guest: &mut G,
call: C,
) -> Result<i64, Error> {
let wrapped: Syscall = call.into();
let action = match ioaction_based_on_fd_status(guest, call) {
Ok(action) => action,
Err(errno) => {
tracing::trace!(
"NonblockableSyscall: fd classification failed with {}; executing kernel-authoritatively: {}",
errno,
call.name()
);
return self.record_or_replay_blocking(guest, wrapped).await;
}
};
let internal_fd = syscall_targets_internal_fd(guest, wrapped);
if !self.cfg.sequentialize_threads
|| (self.cfg.recordreplay_modes && !internal_fd)
|| action == IOAction::Blocking
{
tracing::trace!(
"NonblockableSyscall: executing in blocking mode after all: {}",
call.name()
);
Ok(self.record_or_replay_blocking(guest, wrapped).await?)
} else if action == IOAction::NonblockizeRetry {
tracing::trace!(
"NonblockableSyscall: converting to nonblocking syscall (internal polling): {}",
call.name()
);
let mut rsrc = Resources::new(guest.thread_state().dettid);
rsrc.insert(ResourceID::InternalIOPolling, Permission::W);
rsrc.fyi(call.name());
let subtool = (self.cfg.recordreplay_modes && internal_fd).then_some(self);
Ok(retry_nonblocking_syscall(guest, call, rsrc, subtool).await?)
} else {
assert!(action == IOAction::PassThru);
tracing::trace!(
"NonblockableSyscall: just passing it through: {}",
call.name()
);
self.record_or_replay_preserving_tool_errors(guest, wrapped)
.await
}
}
pub async fn execute_blocking_pipe_writev<G: Guest<Self>>(
&self,
guest: &mut G,
call: syscalls::Writev,
expected_open_file: OpenFileId,
) -> Result<i64, Error> {
const MAX_IOVECS: usize = 1024;
const MAX_RW_COUNT: usize = 0x7fff_f000;
const PIPE_BUF: usize = 4096;
const STACK_IOVECS: usize = 32;
let Some(iov_addr) = call.iov() else {
return self.execute_nonblockable_fd_syscall(guest, call).await;
};
if call.len() == 0 || call.len() > MAX_IOVECS {
return self.execute_nonblockable_fd_syscall(guest, call).await;
}
let iovecs: Vec<(usize, usize)> = {
let mut raw_iovecs = vec![
libc::iovec {
iov_base: std::ptr::null_mut(),
iov_len: 0,
};
call.len()
];
guest.memory().read_values(iov_addr, &mut raw_iovecs)?;
raw_iovecs
.into_iter()
.map(|iovec| (iovec.iov_base as usize, iovec.iov_len))
.collect()
};
let requested = iovecs.iter().try_fold(0usize, |total, (_, length)| {
total.checked_add(*length).ok_or(Errno::EINVAL)
})?;
if requested > isize::MAX as usize {
return Err(Errno::EINVAL.into());
}
let target = requested.min(MAX_RW_COUNT);
if target == 0 {
return self.execute_nonblockable_fd_syscall(guest, call).await;
}
let atomic_pipe_write = target <= PIPE_BUF;
tracing::trace!(
"NonblockableSyscall: converting to nonblocking syscall (internal polling): writev"
);
let mut resources = pipe_writev_resources(guest.thread_state().dettid, call);
let subtool = self.cfg.recordreplay_modes.then_some(self);
let mut current = Syscall::Writev(call);
let mut written_total = 0usize;
let blocked_mask = blocked_signal_mask();
let mut stack = guest.stack().await;
let atomic_scratch_iov = if atomic_pipe_write && iovecs.len() <= STACK_IOVECS {
let mut raw_iovecs = [libc::iovec {
iov_base: std::ptr::null_mut(),
iov_len: 0,
}; STACK_IOVECS];
for (raw, (base, length)) in raw_iovecs.iter_mut().zip(&iovecs) {
raw.iov_base = *base as *mut libc::c_void;
raw.iov_len = *length;
}
let scratch_iov: Addr<libc::iovec> = stack.push(raw_iovecs).cast();
Some(scratch_iov.as_raw())
} else {
None
};
let blocked_mask_addr = stack.push(blocked_mask);
let old_mask_addr = stack.reserve::<KernelSigset>();
let action_addr = stack.reserve::<KernelSigaction>();
let _mask_guard = stack.commit()?;
let guest_signal_mask =
block_signals_for_disposition(guest, blocked_mask_addr, old_mask_addr).await?;
let result: Result<i64, Error> = loop {
let status = resource_request(guest, resources.clone()).await;
let disposition = match wait_signal_disposition(
guest,
status.clone(),
&guest_signal_mask,
action_addr,
true,
)
.await
{
Ok(disposition) => disposition,
Err(error) => break Err(error),
};
if matches!(status, ResumeStatus::Signaled(_))
&& let Some(result) = interrupted_write_result(&call, written_total, disposition)
{
break result;
}
if resources.poll_attempt > 0
&& !guest
.thread_state()
.with_detfd(call.fd(), |detfd| {
detfd.open_file_id() == expected_open_file
})
.unwrap_or(false)
{
break if written_total > 0 {
Ok(written_total as i64)
} else {
self.refuse_unserviceable_operation(guest, Sysno::writev, Errno::EOPNOTSUPP)
.await
};
}
let result = if atomic_pipe_write {
self.execute_atomic_pipe_writev_attempt(guest, call, &iovecs, atomic_scratch_iov)
.await
} else {
match subtool {
Some(detcore) => {
detcore
.record_or_replay_preserving_tool_errors(guest, current)
.await
}
None => guest.inject_with_retry(current).await.map_err(Error::from),
}
};
match result {
Ok(written) if written > 0 => {
let written = match usize::try_from(written) {
Ok(written) => written,
Err(_) => break Err(Errno::EIO.into()),
};
written_total = match written_total.checked_add(written) {
Some(written_total) => written_total,
None => break Err(Errno::EIO.into()),
};
if written_total >= target {
break Ok(written_total as i64);
}
if atomic_pipe_write {
break Ok(written_total as i64);
}
current = match remaining_writev_segment(
call.fd(),
&iovecs,
written_total,
target - written_total,
) {
Ok(Some(write)) => Syscall::Write(write),
Ok(None) => break Ok(written_total as i64),
Err(_) => break Ok(written_total as i64),
};
}
Ok(0) => break Ok(written_total as i64),
Err(Error::Errno(Errno::EAGAIN)) => {
if !atomic_pipe_write && matches!(current, Syscall::Writev(_)) {
current = match remaining_writev_segment(call.fd(), &iovecs, 0, target) {
Ok(Some(write)) => Syscall::Write(write),
Ok(None) => break Ok(0),
Err(error) => break Err(error.into()),
};
}
}
Err(error) => {
break finish_partial_record_or_replay_write(written_total as i64, error);
}
Ok(_) => break Err(Errno::EIO.into()),
}
resources.poll_attempt += 1;
tracing::trace!(
"Retry #{} for {}blocking pipe writev after EAGAIN: {}",
resources.poll_attempt,
if atomic_pipe_write { "atomic " } else { "" },
call.display(&guest.memory())
);
record_retry_event(guest, call).await;
};
restore_signals_after_disposition(guest, old_mask_addr).await?;
result
}
pub async fn execute_blocking_pipe_write<G: Guest<Self>>(
&self,
guest: &mut G,
call: syscalls::Write,
expected_open_file: OpenFileId,
) -> Result<i64, Error> {
const MAX_RW_COUNT: usize = 0x7fff_f000;
let target = call.len().min(MAX_RW_COUNT);
if target == 0 {
return self.execute_nonblockable_fd_syscall(guest, call).await;
}
tracing::trace!(
"NonblockableSyscall: converting to nonblocking syscall (internal polling): write"
);
let mut resources = Resources::new(guest.thread_state().dettid);
resources.insert(ResourceID::InternalIOPolling, Permission::W);
resources.fyi(call.name());
let subtool = self.cfg.recordreplay_modes.then_some(self);
let mut current = call;
let mut written_total = 0usize;
loop {
if resources.poll_attempt > 0
&& matches!(
resource_request(guest, resources.clone()).await,
ResumeStatus::Signaled(_)
)
{
break if written_total > 0 {
Ok(written_total as i64)
} else {
Err(call.signal_interrupt_errno().into())
};
}
if resources.poll_attempt > 0
&& !guest
.thread_state()
.with_detfd(call.fd(), |detfd| {
detfd.open_file_id() == expected_open_file
})
.unwrap_or(false)
{
break if written_total > 0 {
Ok(written_total as i64)
} else {
self.refuse_unserviceable_operation(guest, Sysno::write, Errno::EOPNOTSUPP)
.await
};
}
let result = match subtool {
Some(detcore) => {
detcore
.record_or_replay_preserving_tool_errors(guest, current)
.await
}
None => guest.inject_with_retry(current).await.map_err(Error::from),
};
match result {
Ok(written) if written > 0 => {
let written = usize::try_from(written).map_err(|_| Errno::EIO)?;
let remaining = target.checked_sub(written_total).ok_or(Errno::EIO)?;
if written > remaining {
break Err(Errno::EIO.into());
}
written_total = written_total.checked_add(written).ok_or(Errno::EIO)?;
if written_total == target {
break Ok(written_total as i64);
}
let Some(buffer) = call.buf() else {
break Err(Errno::EFAULT.into());
};
let Some(next_buffer) = buffer
.as_raw()
.checked_add(written_total)
.and_then(Addr::<u8>::from_raw)
else {
break finish_partial_record_or_replay_write(
written_total as i64,
Errno::EFAULT.into(),
);
};
current = call
.with_buf(Some(next_buffer))
.with_len(target - written_total);
}
Ok(0) => break Ok(written_total as i64),
Err(Error::Errno(Errno::EAGAIN)) => {}
Err(error) => {
break finish_partial_record_or_replay_write(written_total as i64, error);
}
Ok(_) => break Err(Errno::EIO.into()),
}
resources.poll_attempt += 1;
tracing::trace!(
"Retry #{} for blocking pipe write: {}",
resources.poll_attempt,
call.display(&guest.memory())
);
record_retry_event(guest, call).await;
}
}
async fn execute_atomic_pipe_writev_attempt<G: Guest<Self>>(
&self,
guest: &mut G,
call: syscalls::Writev,
iovecs: &[(usize, usize)],
stack_iov: Option<usize>,
) -> Result<i64, Error> {
if let Some(stack_iov) = stack_iov {
let scratch_call = call.with_iov(Addr::from_raw(stack_iov));
return if self.cfg.recordreplay_modes {
self.record_or_replay_preserving_tool_errors(guest, scratch_call)
.await
} else {
guest
.inject_with_retry(scratch_call)
.await
.map_err(Error::from)
};
}
let mapping_len = iovecs
.len()
.checked_mul(std::mem::size_of::<libc::iovec>())
.expect("validated iovec count cannot overflow scratch length");
let mapped = guest
.inject_with_retry(Syscall::Mmap(
syscalls::Mmap::new()
.with_addr(None)
.with_len(mapping_len)
.with_prot(ProtFlags::PROT_READ | ProtFlags::PROT_WRITE)
.with_flags(MapFlags::MAP_PRIVATE | MapFlags::MAP_ANONYMOUS)
.with_fd(-1)
.with_offset(0),
))
.await
.unwrap_or_else(|error| panic!("failed to map atomic writev scratch: {error}"));
let mapped = usize::try_from(mapped)
.unwrap_or_else(|_| panic!("atomic writev scratch mmap returned {mapped}"));
let scratch_iov = Addr::<libc::iovec>::from_raw(mapped)
.unwrap_or_else(|| panic!("atomic writev scratch mmap returned a null address"));
let mapping_addr: Addr<libc::c_void> = scratch_iov.cast();
let write_result = {
let raw_iovecs: Vec<libc::iovec> = iovecs
.iter()
.map(|(base, length)| libc::iovec {
iov_base: *base as *mut libc::c_void,
iov_len: *length,
})
.collect();
guest
.memory()
.write_values(unsafe { scratch_iov.into_mut() }, &raw_iovecs)
};
if let Err(write_error) = write_result {
guest
.inject_with_retry(Syscall::Munmap(
syscalls::Munmap::new()
.with_addr(Some(mapping_addr))
.with_len(mapping_len),
))
.await
.unwrap_or_else(|cleanup_error| {
panic!(
"failed to populate atomic writev scratch ({write_error}); cleanup failed ({cleanup_error})"
)
});
panic!("failed to populate atomic writev scratch: {write_error}");
}
let scratch_call = call.with_iov(Some(scratch_iov));
let result = if self.cfg.recordreplay_modes {
self.record_or_replay_preserving_tool_errors(guest, scratch_call)
.await
} else {
guest
.inject_with_retry(scratch_call)
.await
.map_err(Error::from)
};
guest
.inject_with_retry(Syscall::Munmap(
syscalls::Munmap::new()
.with_addr(Some(mapping_addr))
.with_len(mapping_len),
))
.await
.unwrap_or_else(|error| panic!("failed to unmap atomic writev scratch: {error}"));
result
}
pub fn maybe_set_nonblocking_fd<G: Guest<Self>>(&self, guest: &G, fd: i32) {
if self.cfg.sequentialize_threads && !self.cfg.debug_externalize_sockets {
guest
.thread_state()
.with_detfd(fd, |detfd| {
detfd.set_physically_nonblocking();
})
.unwrap();
}
}
}
fn remaining_writev_segment(
fd: i32,
iovecs: &[(usize, usize)],
mut consumed: usize,
remaining_limit: usize,
) -> Result<Option<syscalls::Write>, Errno> {
for (base, length) in iovecs {
if consumed >= *length {
consumed -= *length;
continue;
}
let base = base.checked_add(consumed).ok_or(Errno::EFAULT)?;
let buffer = Addr::<u8>::from_raw(base).ok_or(Errno::EFAULT)?;
return Ok(Some(
syscalls::Write::new()
.with_fd(fd)
.with_buf(Some(buffer))
.with_len((*length - consumed).min(remaining_limit)),
));
}
Ok(None)
}
#[derive(PartialEq, Eq, Debug)]
pub enum IOAction {
Blocking,
NonblockizeRetry,
PassThru,
}
pub fn ioaction_based_on_fd_status<
G: Guest<Detcore<T>>,
T: RecordOrReplay,
C: SyscallInfo + Into<Syscall>,
>(
guest: &mut G,
call: C,
) -> Result<IOAction, Errno> {
let wrapped: Syscall = call.into();
let fd = get_fd(wrapped).unwrap_or_else(|| panic!("Failed to get fd for {}", call.name()));
let (phys, virt) = guest.thread_state().with_detfd(fd, |detfd| {
(detfd.physically_nonblocking(), detfd.is_nonblocking())
})?;
tracing::trace!(
"Checking FD {} for nonblocking: physical {} / virtual {}",
fd,
phys,
virt
);
if virt && !phys {
panic!(
"Invariant violation, fd {}: we cannot simulate nonblocking behavior when set to blocking mode in the kernel.",
fd
);
} else if !virt && !phys {
Ok(IOAction::Blocking)
} else if virt && phys {
Ok(IOAction::PassThru)
} else {
Ok(IOAction::NonblockizeRetry)
}
}
pub fn syscall_targets_internal_fd<G: Guest<Detcore<T>>, T: RecordOrReplay>(
guest: &mut G,
call: Syscall,
) -> bool {
match get_fd(call) {
Some(fd) => guest
.thread_state()
.with_detfd(fd, |detfd| matches!(detfd.ty(), FdType::Pipe))
.unwrap_or(false),
None => false,
}
}
pub(crate) fn get_fd(s: Syscall) -> Option<i32> {
match s {
Syscall::Recvfrom(s) => Some(s.fd()),
Syscall::Recvmsg(s) => Some(s.sockfd()),
Syscall::Recvmmsg(s) => Some(s.fd()),
Syscall::Sendto(s) => Some(s.fd()),
Syscall::Sendmsg(s) => Some(s.fd()),
Syscall::Sendmmsg(s) => Some(s.sockfd()),
Syscall::Accept(s) => Some(s.sockfd()),
Syscall::Accept4(s) => Some(s.sockfd()),
Syscall::Connect(s) => Some(s.fd()),
Syscall::Bind(s) => Some(s.fd()),
Syscall::Listen(s) => Some(s.fd()),
Syscall::Getsockname(s) => Some(s.fd()),
Syscall::Getpeername(s) => Some(s.fd()),
Syscall::Setsockopt(s) => Some(s.fd()),
Syscall::Getsockopt(s) => Some(s.fd()),
Syscall::Read(s) => Some(s.fd()),
Syscall::Write(s) => Some(s.fd()),
Syscall::Close(s) => Some(s.fd()),
Syscall::Fstat(s) => Some(s.fd()),
Syscall::Lseek(s) => Some(s.fd()),
Syscall::Mmap(s) => Some(s.fd()),
Syscall::Ioctl(s) => Some(s.fd()),
Syscall::Pread64(s) => Some(s.fd()),
Syscall::Pwrite64(s) => Some(s.fd()),
Syscall::Readv(s) => Some(s.fd()),
Syscall::Writev(s) => Some(s.fd()),
Syscall::Shutdown(s) => Some(s.fd()),
Syscall::Fcntl(s) => Some(s.fd()),
Syscall::Flock(s) => Some(s.fd()),
Syscall::Fsync(s) => Some(s.fd()),
Syscall::Fdatasync(s) => Some(s.fd()),
Syscall::Ftruncate(s) => Some(s.fd()),
Syscall::Fchdir(s) => Some(s.fd()),
Syscall::Fchmod(s) => Some(s.fd()),
Syscall::Fchown(s) => Some(s.fd()),
Syscall::Fstatfs(s) => Some(s.fd()),
Syscall::Readahead(s) => Some(s.fd()),
Syscall::Fsetxattr(s) => Some(s.fd()),
Syscall::Fgetxattr(s) => Some(s.fd()),
Syscall::Flistxattr(s) => Some(s.fd()),
Syscall::Fremovexattr(s) => Some(s.fd()),
Syscall::Fadvise64(s) => Some(s.fd()),
Syscall::InotifyAddWatch(s) => Some(s.fd()),
Syscall::InotifyRmWatch(s) => Some(s.fd()),
Syscall::SyncFileRange(s) => Some(s.fd()),
Syscall::Vmsplice(s) => Some(s.fd()),
Syscall::Utimensat(s) => Some(s.dirfd()),
Syscall::Signalfd(s) => Some(s.fd()),
Syscall::Fallocate(s) => Some(s.fd()),
Syscall::TimerfdSettime(s) => Some(s.fd()),
Syscall::TimerfdGettime(s) => Some(s.fd()),
Syscall::Signalfd4(s) => Some(s.fd()),
Syscall::Preadv(s) => Some(s.fd()),
Syscall::Pwritev(s) => Some(s.fd()),
Syscall::Syncfs(s) => Some(s.fd()),
Syscall::Setns(s) => Some(s.fd()),
Syscall::FinitModule(s) => Some(s.fd()),
Syscall::Preadv2(s) => Some(s.fd()),
Syscall::Pwritev2(s) => Some(s.fd()),
Syscall::Openat(s) => Some(s.dirfd()),
Syscall::Mkdirat(s) => Some(s.dirfd()),
Syscall::Mknodat(s) => Some(s.dirfd()),
Syscall::Fchownat(s) => Some(s.dirfd()),
Syscall::Futimesat(s) => Some(s.dirfd()),
Syscall::Newfstatat(s) => Some(s.dirfd()),
Syscall::Unlinkat(s) => Some(s.dirfd()),
Syscall::Readlinkat(s) => Some(s.dirfd()),
Syscall::Fchmodat(s) => Some(s.dirfd()),
Syscall::Faccessat(s) => Some(s.dirfd()),
Syscall::NameToHandleAt(s) => Some(s.dirfd()),
Syscall::Execveat(s) => Some(s.dirfd()),
Syscall::Statx(s) => Some(s.dirfd()),
Syscall::Symlinkat(s) => Some(s.newdirfd()),
Syscall::PerfEventOpen(s) => Some(s.group_fd()),
Syscall::OpenByHandleAt(s) => Some(s.mount_fd()),
Syscall::EpollCtl(s) => Some(s.epfd()),
Syscall::EpollWait(s) => Some(s.epfd()),
Syscall::EpollPwait(s) => Some(s.epfd()),
Syscall::Dup2(_) => None,
Syscall::Sendfile(_) => None,
Syscall::Renameat(_) => None,
Syscall::Linkat(_) => None,
Syscall::FanotifyMark(_) => None,
Syscall::Renameat2(_) => None,
Syscall::Dup3(_) => None,
Syscall::KexecLoad(_) => None,
Syscall::Poll(_) => None,
Syscall::Ppoll(_) => None,
_ => None,
}
}
#[allow(clippy::double_must_use)]
#[async_trait]
pub trait NonblockableSyscall: SyscallInfo {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>);
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Ok(0)
}
fn signal_interrupt_errno(&self) -> Errno {
Errno::ERESTARTSYS
}
fn normalize_nonblocking_result(
&self,
res: Result<i64, Errno>,
_retried: bool,
) -> Result<i64, Errno> {
res
}
}
pub trait TimeoutableSyscall: SyscallInfo {
fn timeout_return_val(&self) -> Result<i64, Errno>;
}
fn interrupted_write_result<C: NonblockableSyscall>(
call: &C,
written_total: usize,
disposition: Option<WaitSignalDisposition>,
) -> Option<Result<i64, Error>> {
match disposition {
None => None,
Some(_) if written_total > 0 => Some(Ok(written_total as i64)),
Some(_) => Some(Err(call.signal_interrupt_errno().into())),
}
}
fn pipe_writev_resources(dettid: DetTid, call: reverie::syscalls::Writev) -> Resources {
let mut resources = Resources::new(dettid);
resources.insert(ResourceID::InternalIOPolling, Permission::W);
resources.fyi(call.name());
resources.set_signal_interrupt_errno(call.signal_interrupt_errno());
resources
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Poll {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
_guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
(self.with_timeout(0), None)
}
fn signal_interrupt_errno(&self) -> Errno {
Errno::EINTR
}
}
impl TimeoutableSyscall for reverie::syscalls::Poll {
fn timeout_return_val(&self) -> Result<i64, Errno> {
Ok(0)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Ppoll {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
let (tp, guard) = zero_timespec(guest).await;
let tp = unsafe { tp.into_mut() };
(self.with_timeout(Some(tp)), Some(guard))
}
fn signal_interrupt_errno(&self) -> Errno {
Errno::EINTR
}
}
impl TimeoutableSyscall for reverie::syscalls::Ppoll {
fn timeout_return_val(&self) -> Result<i64, Errno> {
Ok(0)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::EpollWait {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
_guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
(self.with_timeout(0), None)
}
fn signal_interrupt_errno(&self) -> Errno {
Errno::EINTR
}
}
impl TimeoutableSyscall for reverie::syscalls::EpollWait {
fn timeout_return_val(&self) -> Result<i64, Errno> {
Ok(0)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::EpollPwait {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
_guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
(self.with_timeout(0), None)
}
fn signal_interrupt_errno(&self) -> Errno {
Errno::EINTR
}
}
impl TimeoutableSyscall for reverie::syscalls::EpollPwait {
fn timeout_return_val(&self) -> Result<i64, Errno> {
Ok(0)
}
}
async fn zero_timespec<'stack, T: RecordOrReplay, G: Guest<Detcore<T>>>(
guest: &mut G,
) -> (Addr<'stack, Timespec>, <G::Stack as Stack>::StackGuard) {
let mut stack = guest.stack().await;
let tp_val = Timespec {
tv_sec: 0,
tv_nsec: 0,
};
let tp = stack.push(tp_val);
let guard = stack.commit().expect("stack.commit to succeed");
(tp, guard)
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Wait4 {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
_guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
let call2 = self.with_options(self.options() | WaitPidFlag::WNOHANG);
(call2, None)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Ok(0)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Futex {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
let (tp, guard) = zero_timespec(guest).await;
(self.with_timeout(Some(tp)), Some(guard))
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::ETIMEDOUT)
}
}
impl TimeoutableSyscall for reverie::syscalls::Futex {
fn timeout_return_val(&self) -> Result<i64, Errno> {
Err(Errno::ETIMEDOUT)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::RtSigtimedwait {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
let (tp, guard) = zero_timespec(guest).await;
(self.with_timeout(Some(tp)), Some(guard))
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN)
}
fn signal_interrupt_errno(&self) -> Errno {
Errno::EINTR
}
}
impl TimeoutableSyscall for reverie::syscalls::RtSigtimedwait {
fn timeout_return_val(&self) -> Result<i64, Errno> {
Err(Errno::EAGAIN)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Read {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Write {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Readv {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Writev {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
fn network_comm_syscall<T: RecordOrReplay, G: Guest<Detcore<T>>, C: SyscallInfo + Into<Syscall>>(
call: C,
guest: &mut G,
) -> (C, Option<<G::Stack as Stack>::StackGuard>) {
let fd = get_fd(call.into()).unwrap_or_else(|| {
panic!(
"network_comm_syscall called on invalid syscall / unknown fd: {}",
call.name()
);
});
guest
.thread_state()
.with_detfd(fd, |detfd| {
assert!(
detfd.physically_nonblocking(),
"expecting sockets/pipes to be physically nonblocking"
);
})
.unwrap();
(call, None)
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Accept4 {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
impl TimeoutableSyscall for reverie::syscalls::Accept4 {
fn timeout_return_val(&self) -> Result<i64, Errno> {
Ok(0)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Recvfrom {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Recvmsg {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Recvmmsg {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Sendto {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Sendmmsg {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Sendmsg {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN) || res == Err(Errno::EWOULDBLOCK)
}
}
#[async_trait]
impl NonblockableSyscall for reverie::syscalls::Connect {
async fn into_nonblocking<T: RecordOrReplay, G: Guest<Detcore<T>>>(
self,
guest: &mut G,
) -> (Self, Option<<G::Stack as Stack>::StackGuard>) {
network_comm_syscall(self, guest)
}
fn syscall_would_have_blocked(&self, res: Result<i64, Errno>) -> bool {
res == Err(Errno::EAGAIN)
|| res == Err(Errno::EWOULDBLOCK)
|| res == Err(Errno::EINPROGRESS)
|| res == Err(Errno::EALREADY)
}
fn normalize_nonblocking_result(
&self,
res: Result<i64, Errno>,
retried: bool,
) -> Result<i64, Errno> {
match (retried, res) {
(true, Err(Errno::EISCONN)) => Ok(0),
(_, res) => res,
}
}
}
pub async fn retry_nonblocking_syscall<T, G, C>(
guest: &mut G,
call: C,
rsrc: Resources,
subtool: Option<&Detcore<T>>,
) -> Result<i64, Error>
where
C: NonblockableSyscall + Into<Syscall>,
T: RecordOrReplay,
G: Guest<Detcore<T>>,
{
retry_nonblocking_syscall_helper(guest, call, rsrc, None, subtool).await
}
pub async fn retry_nonblocking_syscall_with_timeout<T, G, C>(
guest: &mut G,
call: C,
rsrc: Resources,
maybe_timeout: Option<LogicalTime>,
) -> Result<i64, Error>
where
C: NonblockableSyscall + TimeoutableSyscall + Into<Syscall>,
T: RecordOrReplay,
G: Guest<Detcore<T>>,
{
let maybe_tup = maybe_timeout.map(|t| (t, call.timeout_return_val()));
retry_nonblocking_syscall_helper(guest, call, rsrc, maybe_tup, None).await
}
async fn retry_nonblocking_syscall_helper<T, G, C>(
guest: &mut G,
call0: C,
rsrc: Resources,
maybe_timeout: Option<(LogicalTime, Result<i64, Errno>)>,
subtool: Option<&Detcore<T>>,
) -> Result<i64, Error>
where
C: NonblockableSyscall + Into<Syscall>,
T: RecordOrReplay,
G: Guest<Detcore<T>>,
{
let (call, _maybe_stackguard) = call0.into_nonblocking(guest).await;
let mut rsrc = rsrc.clone();
loop {
let resumed = match call.into() {
Syscall::Read(read) if maybe_timeout.is_none() && _maybe_stackguard.is_none() => {
crate::tool_global::polled_read_request(guest, read, rsrc.clone()).await
}
_ => resource_request(guest, rsrc.clone()).await,
};
if matches!(resumed, ResumeStatus::Signaled(_)) {
let errno = call.signal_interrupt_errno();
tracing::trace!(
"retry_nonblocking_syscall: interrupted by signal before retrying {}: {:?}",
call.display(&guest.memory()),
errno
);
return Err(errno.into());
}
let res = match subtool {
Some(detcore) => {
detcore
.record_or_replay_preserving_tool_errors(guest, call)
.await
}
None => guest.inject_with_retry(call).await.map_err(Error::from),
};
let syscall_result = match res {
Ok(value) => Ok(value),
Err(Error::Errno(error)) => Err(error),
Err(error) => return Err(error),
};
if call.syscall_would_have_blocked(syscall_result) {
rsrc.poll_attempt += 1;
if let Some((timeout, timeout_result)) = maybe_timeout {
let new_time = thread_observe_time(guest).await;
if new_time >= timeout {
tracing::trace!(
"Timing out syscall after #{} retries: {}",
rsrc.poll_attempt - 1,
call.display(&guest.memory())
);
return timeout_result.map_err(|e| e.into());
} else {
tracing::trace!(
"Retry #{} for syscall due to result {:?}, {} from timeout: {}",
rsrc.poll_attempt,
syscall_result,
timeout - new_time,
call.display(&guest.memory())
);
record_retry_event(guest, call).await;
}
} else {
tracing::trace!(
"Retry #{} for syscall due to result {:?}: {}",
rsrc.poll_attempt,
syscall_result,
call.display(&guest.memory())
);
record_retry_event(guest, call).await;
}
} else {
let res = call
.normalize_nonblocking_result(syscall_result, rsrc.poll_attempt > 0)
.map_err(|e| e.into());
tracing::trace!(
"retry_nonblocking_syscall: syscall completed after {} retries: {} = {:?}",
rsrc.poll_attempt,
call.display(&guest.memory()),
res
);
return res;
}
}
}
pub(crate) async fn record_retry_event<G, C, T>(guest: &mut G, call: C)
where
C: SyscallInfo,
T: RecordOrReplay,
G: Guest<Detcore<T>>,
{
let dettid = guest.thread_state().dettid;
let cfg = &guest.config();
if cfg.sequentialize_threads && cfg.should_trace_schedevent() {
trace_schedevent(
guest,
with_guest_time(
guest,
SchedEvent::syscall(dettid, call.number(), SyscallPhase::Polling),
),
true,
)
.await;
}
}
pub fn with_guest_time<G, T>(guest: &G, event: SchedEvent) -> SchedEvent
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
let dettime = &guest.thread_state().thread_logical_time;
event.with_dettime(dettime)
}
pub async fn with_guest_rip<G, T>(guest: &mut G, mut event: SchedEvent) -> SchedEvent
where
G: Guest<Detcore<T>>,
T: RecordOrReplay,
{
assert!(event.end_rip.is_none());
let regs = guest.regs().await;
let end_rip = NonZeroUsize::new(regs.rip.try_into().unwrap()).unwrap();
event.end_rip = Some(end_rip);
event
}
pub async fn millis_duration_to_absolute_timeout<G: Guest<Detcore<T>>, T: RecordOrReplay>(
guest: &mut G,
timeout_millis: i32,
) -> Option<LogicalTime> {
match positive_millis_as_nanos(timeout_millis) {
Some(timeout_nanos) => nanos_duration_to_absolute_timeout(guest, timeout_nanos).await,
None => None,
}
}
fn positive_millis_as_nanos(timeout_millis: i32) -> Option<u128> {
(timeout_millis > 0).then(|| (timeout_millis as u128) * 1_000_000)
}
pub async fn nanos_duration_to_absolute_timeout<G: Guest<Detcore<T>>, T: RecordOrReplay>(
guest: &mut G,
timeout_nanos: u128,
) -> Option<LogicalTime> {
if timeout_nanos > 0 {
let ns_delta = Duration::from_nanos(timeout_nanos as u64);
let base_time = thread_observe_time(guest).await;
let target_time = base_time + ns_delta;
Some(target_time)
} else {
None
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn millis_to_nanos_conversion() {
assert_eq!(positive_millis_as_nanos(-1), None);
assert_eq!(positive_millis_as_nanos(i32::MIN), None);
assert_eq!(positive_millis_as_nanos(0), None);
assert_eq!(positive_millis_as_nanos(1), Some(1_000_000));
assert_eq!(positive_millis_as_nanos(1_000), Some(1_000_000_000));
assert_eq!(
positive_millis_as_nanos(i32::MAX),
Some(i32::MAX as u128 * 1_000_000)
);
for millis in [1, 2, 7, 250, 1_000, 86_400_000] {
assert_eq!(
positive_millis_as_nanos(millis),
Some(millis as u128 * 1_000_000),
"1 ms must convert to 1_000_000 ns, not 1_000"
);
}
}
#[test]
fn connect_nonblocking_results() {
let call = reverie::syscalls::Connect::new();
assert!(call.syscall_would_have_blocked(Err(Errno::EINPROGRESS)));
assert!(call.syscall_would_have_blocked(Err(Errno::EALREADY)));
assert_eq!(
call.normalize_nonblocking_result(Err(Errno::EISCONN), true),
Ok(0)
);
assert_eq!(
call.normalize_nonblocking_result(Err(Errno::EISCONN), false),
Err(Errno::EISCONN)
);
}
#[test]
fn signal_interruption_errno_matches_linux_restart_policy() {
assert_eq!(
reverie::syscalls::Poll::new().signal_interrupt_errno(),
Errno::EINTR
);
assert_eq!(
reverie::syscalls::Ppoll::new().signal_interrupt_errno(),
Errno::EINTR
);
assert_eq!(
reverie::syscalls::EpollWait::new().signal_interrupt_errno(),
Errno::EINTR
);
let sigtimedwait = reverie::syscalls::RtSigtimedwait::new();
assert_eq!(sigtimedwait.signal_interrupt_errno(), Errno::EINTR);
assert!(sigtimedwait.syscall_would_have_blocked(Err(Errno::EAGAIN)));
assert_eq!(sigtimedwait.timeout_return_val(), Err(Errno::EAGAIN));
assert_eq!(
reverie::syscalls::Read::new().signal_interrupt_errno(),
Errno::ERESTARTSYS
);
assert_eq!(
reverie::syscalls::Writev::new().signal_interrupt_errno(),
Errno::ERESTARTSYS
);
assert_eq!(
reverie::syscalls::Futex::new().signal_interrupt_errno(),
Errno::ERESTARTSYS
);
}
#[test]
fn writev_signal_result_uses_disposition_and_progress() {
let call = reverie::syscalls::Writev::new();
for disposition in [
WaitSignalDisposition::Interrupt,
WaitSignalDisposition::Restart,
] {
assert!(matches!(
interrupted_write_result(&call, 0, Some(disposition)),
Some(Err(Error::Errno(Errno::ERESTARTSYS)))
));
assert!(matches!(
interrupted_write_result(&call, 17, Some(disposition)),
Some(Ok(17))
));
}
assert!(interrupted_write_result(&call, 0, None).is_none());
assert!(interrupted_write_result(&call, 17, None).is_none());
}
#[test]
fn pipe_writev_requests_signal_disposition() {
let dettid = DetTid::from_raw(42);
let call = reverie::syscalls::Writev::new();
let request = pipe_writev_resources(dettid, call);
assert_eq!(
request.signal_interrupt_errno(),
Some(Errno::ERESTARTSYS.into_raw())
);
assert_eq!(request.resources.len(), 1);
}
}