use reverie::Error;
use reverie::Guest;
use reverie::Tool;
use reverie::syscalls::Addr;
use reverie::syscalls::AddrMut;
use reverie::syscalls::Errno;
use reverie::syscalls::MemoryAccess;
use reverie::syscalls::Syscall;
use reverie::syscalls::SyscallInfo;
use crate::digest::Digest;
use crate::types::DetTid;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct BufferExtent {
pub addr: u64,
pub len: u64,
}
pub(crate) fn rng_readv_extents(
iovecs: &[crate::iovecs::ImportedIovec],
written: usize,
) -> Result<Vec<BufferExtent>, Error> {
let mut output = rng_observation_vec(iovecs.len(), "imported geometry")?;
let mut remaining = written;
for iov in iovecs {
let take = remaining.min(iov.len);
remaining -= take;
if take > 0 {
output.push(BufferExtent {
addr: iov.base as u64,
len: take as u64,
});
}
}
Ok(output)
}
fn rng_observation_vec<T>(capacity: usize, purpose: &str) -> Result<Vec<T>, Error> {
let mut output = Vec::new();
output.try_reserve_exact(capacity).map_err(|error| {
Error::Tool(anyhow::anyhow!(
"RNG vector observation after output commit: cannot reserve {purpose}: {error}"
))
})?;
Ok(output)
}
fn observed_extents(
memory: &impl MemoryAccess,
call: &Syscall,
ret: i64,
rng_output: Option<&[BufferExtent]>,
) -> Result<Vec<BufferExtent>, Error> {
if let Some(output) = rng_output {
if ret <= 0 {
return Ok(Vec::new());
}
let mut observed = rng_observation_vec(output.len(), "observed geometry")?;
let mut remaining = ret as u64;
for extent in output {
let take = remaining.min(extent.len);
remaining -= take;
if take > 0 {
observed.push(BufferExtent {
addr: extent.addr,
len: take,
});
}
}
return Ok(observed);
}
extents(memory, call, ret)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Direction {
In,
Out,
}
impl Direction {
fn as_str(self) -> &'static str {
match self {
Direction::In => "in",
Direction::Out => "out",
}
}
}
fn clamp(addr: Option<u64>, capacity: usize, moved: i64) -> Vec<BufferExtent> {
let Some(addr) = addr else { return Vec::new() };
let moved = u64::try_from(moved).unwrap_or(0);
let len = moved.min(capacity as u64);
if len == 0 {
return Vec::new();
}
vec![BufferExtent { addr, len }]
}
fn whole(addr: Option<u64>, len: u64) -> Vec<BufferExtent> {
match addr {
Some(addr) if len > 0 => vec![BufferExtent { addr, len }],
_ => Vec::new(),
}
}
fn iovec_extents<M: MemoryAccess>(
memory: &M,
iov_addr: usize,
iov_count: usize,
moved: i64,
) -> Result<Vec<BufferExtent>, Error> {
let remaining = u64::try_from(moved).unwrap_or(0);
if iov_addr == 0 || iov_count == 0 || remaining == 0 {
return Ok(Vec::new());
}
let iov_count = iov_count.min(libc::UIO_MAXIOV as usize);
let iov_address: AddrMut<'_, libc::iovec> = AddrMut::from_raw(iov_addr).ok_or(Errno::EFAULT)?;
let mut iovecs: Vec<libc::iovec> = (0..iov_count)
.map(|_| unsafe { std::mem::zeroed() })
.collect();
memory.read_values(iov_address.into(), &mut iovecs)?;
Ok(iovec_extents_from_slice(&iovecs, moved))
}
fn iovec_extents_from_slice(iovecs: &[libc::iovec], moved: i64) -> Vec<BufferExtent> {
let mut remaining = u64::try_from(moved).unwrap_or(0);
let mut out = Vec::new();
for iov in iovecs {
if remaining == 0 {
break;
}
if iov.iov_base.is_null() || iov.iov_len == 0 {
continue;
}
let take = (iov.iov_len as u64).min(remaining);
out.push(BufferExtent {
addr: iov.iov_base as u64,
len: take,
});
remaining -= take;
}
out
}
fn iovec_extent_arguments(call: &Syscall) -> Option<(usize, usize)> {
match call {
Syscall::Readv(call) => Some((call.iov().map_or(0, |p| p.as_raw()), call.len())),
Syscall::Preadv(call) => Some((call.iov().map_or(0, |p| p.as_raw()), call.iov_len())),
Syscall::Preadv2(call) => Some((
call.iov().map_or(0, |p| p.as_raw()),
usize::try_from(call.iov_len()).unwrap_or(usize::MAX),
)),
Syscall::Writev(call) => Some((call.iov().map_or(0, |p| p.as_raw()), call.len())),
Syscall::Pwritev(call) => Some((call.iov().map_or(0, |p| p.as_raw()), call.iov_len())),
Syscall::Pwritev2(call) => Some((
call.iov().map_or(0, |p| p.as_raw()),
usize::try_from(call.iov_len()).unwrap_or(usize::MAX),
)),
_ => None,
}
}
fn msghdr_extents<M: MemoryAccess>(
memory: &M,
msg_addr: usize,
moved: i64,
) -> Result<Vec<BufferExtent>, Error> {
if msg_addr == 0 {
return Ok(Vec::new());
}
let address: AddrMut<'_, libc::msghdr> = AddrMut::from_raw(msg_addr).ok_or(Errno::EFAULT)?;
let message: libc::msghdr = memory.read_value(address)?;
iovec_extents(memory, message.msg_iov as usize, message.msg_iovlen, moved)
}
fn completed_mmsghdr_count(vlen: u32, completed: i64) -> usize {
usize::try_from(completed)
.unwrap_or(0)
.min(vlen as usize)
.min(libc::UIO_MAXIOV as usize)
}
fn mmsghdr_extents<M: MemoryAccess>(
memory: &M,
mmsg_addr: usize,
vlen: u32,
delivered: i64,
) -> Result<Vec<BufferExtent>, Error> {
let count = completed_mmsghdr_count(vlen, delivered);
if mmsg_addr == 0 || count == 0 {
return Ok(Vec::new());
}
let address: AddrMut<'_, libc::mmsghdr> = AddrMut::from_raw(mmsg_addr).ok_or(Errno::EFAULT)?;
let mut headers: Vec<libc::mmsghdr> =
(0..count).map(|_| unsafe { std::mem::zeroed() }).collect();
memory.read_values(address.into(), &mut headers)?;
let mut out = Vec::new();
for header in &headers {
out.extend(iovec_extents(
memory,
header.msg_hdr.msg_iov as usize,
header.msg_hdr.msg_iovlen,
i64::from(header.msg_len),
)?);
}
Ok(out)
}
#[cfg(test)]
pub(crate) const HASHED_SYSCALLS: &[reverie::syscalls::Sysno] = &[
reverie::syscalls::Sysno::read,
reverie::syscalls::Sysno::pread64,
reverie::syscalls::Sysno::recvfrom,
reverie::syscalls::Sysno::getrandom,
reverie::syscalls::Sysno::getcwd,
reverie::syscalls::Sysno::getdents64,
reverie::syscalls::Sysno::readlink,
reverie::syscalls::Sysno::readlinkat,
reverie::syscalls::Sysno::recvmsg,
reverie::syscalls::Sysno::recvmmsg,
reverie::syscalls::Sysno::readv,
reverie::syscalls::Sysno::preadv,
reverie::syscalls::Sysno::preadv2,
reverie::syscalls::Sysno::write,
reverie::syscalls::Sysno::pwrite64,
reverie::syscalls::Sysno::sendto,
reverie::syscalls::Sysno::sendmsg,
reverie::syscalls::Sysno::sendmmsg,
reverie::syscalls::Sysno::writev,
reverie::syscalls::Sysno::pwritev,
reverie::syscalls::Sysno::pwritev2,
reverie::syscalls::Sysno::poll,
reverie::syscalls::Sysno::ppoll,
];
fn ret_gates_output(call: &Syscall) -> bool {
!matches!(call, Syscall::Poll(_) | Syscall::Ppoll(_))
}
fn extents<M: MemoryAccess>(
memory: &M,
call: &Syscall,
ret: i64,
) -> Result<Vec<BufferExtent>, Error> {
if ret <= 0 && ret_gates_output(call) {
return Ok(Vec::new());
}
if let Some((iov_addr, iov_count)) = iovec_extent_arguments(call) {
return iovec_extents(memory, iov_addr, iov_count, ret);
}
let raw = |a: Option<AddrMut<'_, u8>>| a.map(|p| p.as_raw() as u64);
Ok(match call {
Syscall::Read(c) => clamp(raw(c.buf()), c.len(), ret),
Syscall::Pread64(c) => clamp(raw(c.buf()), c.len(), ret),
Syscall::Recvfrom(c) => clamp(c.buf().map(|p| p.as_raw() as u64), c.len(), ret),
Syscall::Getrandom(c) => clamp(raw(c.buf()), c.buflen(), ret),
Syscall::Getcwd(c) => clamp(c.buf().map(|p| p.as_raw() as u64), c.size(), ret),
Syscall::Getdents64(c) => clamp(
c.dirent().map(|p| p.as_raw() as u64),
c.count() as usize,
ret,
),
Syscall::Readlink(c) => clamp(c.buf().map(|p| p.as_raw() as u64), c.bufsize(), ret),
Syscall::Readlinkat(c) => clamp(c.buf().map(|p| p.as_raw() as u64), c.buf_len(), ret),
Syscall::Recvmsg(c) => msghdr_extents(memory, c.msg().map_or(0, |p| p.as_raw()), ret)?,
Syscall::Recvmmsg(c) => {
mmsghdr_extents(memory, c.mmsg().map_or(0, |p| p.as_raw()), c.vlen(), ret)?
}
Syscall::Write(c) => clamp(c.buf().map(|p| p.as_raw() as u64), c.len(), ret),
Syscall::Pwrite64(c) => clamp(c.buf().map(|p| p.as_raw() as u64), c.len(), ret),
Syscall::Sendto(c) => clamp(c.buf().map(|p| p.as_raw() as u64), c.size(), ret),
Syscall::Sendmsg(c) => msghdr_extents(memory, c.msg().map_or(0, |p| p.as_raw()), ret)?,
Syscall::Sendmmsg(c) => {
mmsghdr_extents(memory, c.msgvec().map_or(0, |p| p.as_raw()), c.vlen(), ret)?
}
Syscall::Poll(c) => whole(
c.fds().map(|p| p.as_raw() as u64),
c.nfds() * std::mem::size_of::<libc::pollfd>() as u64,
),
Syscall::Ppoll(c) => whole(
c.fds().map(|p| p.as_raw() as u64),
c.nfds() * std::mem::size_of::<libc::pollfd>() as u64,
),
_ => Vec::new(),
})
}
fn direction(call: &Syscall) -> Direction {
match call {
Syscall::Write(_)
| Syscall::Pwrite64(_)
| Syscall::Sendto(_)
| Syscall::Sendmsg(_)
| Syscall::Sendmmsg(_)
| Syscall::Writev(_)
| Syscall::Pwritev(_)
| Syscall::Pwritev2(_) => Direction::Out,
_ => Direction::In,
}
}
const CHUNK_CAP: usize = 8;
const CHUNK_MIN: usize = 256;
fn extent_digests<G, T>(
guest: &mut G,
addr: u64,
len: u64,
rng_output: bool,
) -> Result<(Digest, usize, Vec<String>), Error>
where
G: Guest<T>,
T: Tool,
{
let size = len as usize;
let mut buf = if rng_output {
let mut buf = rng_observation_vec(size, "digest payload")?;
buf.resize(size, 0u8);
buf
} else {
vec![0u8; size]
};
if size > 0 {
let start = Addr::<u8>::from_raw(addr as usize).ok_or(Errno::EFAULT)?;
guest.memory().read_values(start, buf.as_mut_slice())?;
}
let whole = Digest::new(buf.as_slice());
let (chunk, chunks) = chunk_digests(buf.as_slice());
Ok((whole, chunk, chunks))
}
fn chunk_digests(buf: &[u8]) -> (usize, Vec<String>) {
let size = buf.len();
let chunk = std::cmp::max(CHUNK_MIN, size.div_ceil(CHUNK_CAP));
if size <= chunk {
return (chunk, Vec::new());
}
let chunks = buf
.chunks(chunk)
.map(|c| {
Digest::new(c)
.to_string()
.chars()
.take(8)
.collect::<String>()
})
.collect();
(chunk, chunks)
}
pub(crate) fn detlog_io_buffers<G, T>(
guest: &mut G,
call: &Syscall,
ret: i64,
dettid: DetTid,
rng_output: Option<&[BufferExtent]>,
) -> Result<(), Error>
where
G: Guest<T>,
T: Tool,
{
if !crate::detlog_observed!() {
return Ok(());
}
let dir = direction(call).as_str();
let name = call.name();
let fd = match crate::syscalls::helpers::get_fd(*call) {
Some(fd) => fd.to_string(),
None => "-".to_string(),
};
let moved_extents = {
let memory = guest.memory();
observed_extents(&memory, call, ret, rng_output)?
};
for extent in moved_extents {
let (whole, chunk, chunks) =
extent_digests(guest, extent.addr, extent.len, rng_output.is_some()).map_err(
|error| {
if rng_output.is_some() {
Error::Tool(anyhow::Error::new(error).context(format!(
"RNG vector observation after output commit: digest read at {:#x}+{}",
extent.addr, extent.len
)))
} else {
error
}
},
)?;
let located = if chunks.is_empty() {
String::new()
} else {
format!(" chunks={}:{}", chunk, chunks.join(","))
};
crate::detlog!(
"[iobuf][dtid {}] {} {} fd={} {:#x}+{}->{}{}",
dettid,
name,
dir,
fd,
extent.addr,
extent.len,
whole,
located
);
}
Ok(())
}
#[cfg(test)]
mod event_tests {
include!("io_buffers/user_access_event_tests.rs");
use std::io::IoSlice;
use std::io::IoSliceMut;
use std::sync::Arc;
use std::sync::Mutex;
use nix::fcntl::OFlag;
use reverie::GlobalRPC;
use reverie::GlobalTool;
use reverie::Pid;
use tokio::sync::oneshot;
use tracing::Event;
use tracing::Id;
use tracing::Level;
use tracing::Metadata;
use tracing::Subscriber;
use tracing::field::Field;
use tracing::field::Visit;
use tracing::span::Attributes;
use tracing::span::Record;
use super::*;
use crate::Config;
use crate::Detcore;
use crate::GlobalState;
use crate::ThreadState;
use crate::fd::DetFd;
use crate::fd::FdType;
use crate::resources::Permission;
use crate::resources::ResourceID;
use crate::tool_global::GlobalRequest;
use crate::tool_global::GlobalResponse;
use crate::tool_global::ResumeStatus;
use crate::types::DetPid;
use crate::types::LogicalTime;
use crate::types::OpenFileId;
const FD: i32 = 3;
const IOV: usize = 0x1000;
const FIRST_DEST: usize = 0x2000;
const RETRY_DEST: usize = 0x3000;
const PIPE_BYTES: &[u8; 4] = b"pipe";
const CANARY: u8 = 0xa5;
#[derive(Clone)]
struct EventMemory(Arc<Mutex<Vec<u8>>>, Arc<Mutex<MemoryReads>>);
#[derive(Default)]
struct MemoryReads {
imported_entries: usize,
import_error: Option<(usize, Errno)>,
import_failed: bool,
after_import_error: Vec<&'static str>,
observer_reads: Vec<(usize, usize)>,
digest_error: Option<Errno>,
user_copy_audit: bool,
copy_done: bool,
copy_actions: std::collections::VecDeque<(usize, Result<usize, Errno>)>,
copy_lengths: Vec<usize>,
after_copy: Vec<&'static str>,
}
impl EventMemory {
fn new() -> Self {
Self(
Arc::new(Mutex::new(vec![CANARY; 0x4000])),
Arc::new(Mutex::new(MemoryReads::default())),
)
}
fn put_iovec(&self, index: usize, base: usize, len: usize) {
let start = IOV + index * std::mem::size_of::<libc::iovec>();
let mut bytes = self.0.lock().unwrap();
bytes[start..start + 8].copy_from_slice(&base.to_ne_bytes());
bytes[start + 8..start + 16].copy_from_slice(&len.to_ne_bytes());
}
fn iovec(&self, index: usize) -> (usize, usize) {
let address = IOV + index * std::mem::size_of::<libc::iovec>();
let iov: libc::iovec = self
.read_value(Addr::<libc::iovec>::from_raw(address).unwrap())
.unwrap();
(iov.iov_base as usize, iov.iov_len)
}
fn bytes(&self, address: usize, len: usize) -> Vec<u8> {
self.0.lock().unwrap()[address..address + len].to_vec()
}
fn copy_read(&self, start: usize, buf: &mut [u8]) -> Result<(), Errno> {
let end = start.checked_add(buf.len()).ok_or(Errno::EFAULT)?;
let bytes = self.0.lock().unwrap();
buf.copy_from_slice(bytes.get(start..end).ok_or(Errno::EFAULT)?);
Ok(())
}
}
impl MemoryAccess for EventMemory {
fn write_with_user_access(
&mut self,
addr: AddrMut<u8>,
buf: &[u8],
) -> Result<usize, Errno> {
let (written, outcome) = {
let mut audit = self.1.lock().unwrap();
audit.copy_lengths.push(buf.len());
if audit.import_failed {
audit.after_import_error.push("user-write");
}
audit
.copy_actions
.pop_front()
.unwrap_or((buf.len(), Ok(buf.len())))
};
assert!(written <= buf.len());
let start = addr.as_raw();
let end = start.checked_add(written).ok_or(Errno::EFAULT)?;
self.0
.lock()
.unwrap()
.get_mut(start..end)
.ok_or(Errno::EFAULT)?
.copy_from_slice(&buf[..written]);
self.1.lock().unwrap().copy_done = true;
outcome
}
fn read_vectored(
&self,
_remote: &[IoSlice],
_local: &mut [IoSliceMut],
) -> Result<usize, Errno> {
panic!("event fixture must use its scalar memory override")
}
fn write_vectored(
&mut self,
_local: &[IoSlice],
_remote: &mut [IoSliceMut],
) -> Result<usize, Errno> {
panic!("event fixture must use its scalar memory override")
}
fn read<'a, A>(&self, addr: A, buf: &mut [u8]) -> Result<usize, Errno>
where
A: Into<Addr<'a, u8>>,
{
let start = addr.into().as_raw();
let mut reads = self.1.lock().unwrap();
reads.observer_reads.push((start, buf.len()));
if reads.import_failed {
reads.after_import_error.push("memory-read");
}
if reads.user_copy_audit && reads.copy_done {
reads.after_copy.push("memory-read");
}
if start == FIRST_DEST
&& let Some(error) = reads.digest_error
{
return Err(error);
}
drop(reads);
self.copy_read(start, buf)?;
Ok(buf.len())
}
fn read_exact_with_user_access<'a, A>(&self, addr: A, buf: &mut [u8]) -> Result<(), Errno>
where
A: Into<Addr<'a, u8>>,
{
{
let mut reads = self.1.lock().unwrap();
reads.imported_entries += 1;
if reads.user_copy_audit && reads.copy_done {
reads.after_copy.push("user-read");
}
if reads.import_failed {
reads.after_import_error.push("user-read");
}
if let Some((attempt, error)) = reads.import_error
&& reads.imported_entries == attempt
{
reads.import_failed = true;
return Err(error);
}
}
self.copy_read(addr.into().as_raw(), buf)
}
fn write(&mut self, addr: AddrMut<u8>, buf: &[u8]) -> Result<usize, Errno> {
assert!(
!self.1.lock().unwrap().user_copy_audit,
"random output used debugger write"
);
let start = addr.as_raw();
let end = start.checked_add(buf.len()).ok_or(Errno::EFAULT)?;
let mut bytes = self.0.lock().unwrap();
bytes
.get_mut(start..end)
.ok_or(Errno::EFAULT)?
.copy_from_slice(buf);
Ok(buf.len())
}
}
type RetryGate = (oneshot::Sender<()>, oneshot::Receiver<()>);
struct EventGuest {
config: Config,
thread: ThreadState<()>,
memory: EventMemory,
injected_iovecs: Vec<(usize, usize)>,
injected_zero_reads: usize,
polls: Mutex<Vec<u32>>,
request_marks: Mutex<Vec<bool>>,
releases: Mutex<usize>,
retry_gate: Mutex<Option<RetryGate>>,
backend_runtime_bootstrap: bool,
}
impl EventGuest {
fn audit_call(&self, operation: &'static str) -> bool {
let mut audit = self.memory.1.lock().unwrap();
if audit.user_copy_audit && audit.copy_done {
audit.after_copy.push(operation);
}
if audit.import_failed {
audit.after_import_error.push(operation);
}
audit.user_copy_audit
}
}
struct UnusedStack;
struct UnusedStackGuard;
impl Drop for UnusedStackGuard {
fn drop(&mut self) {}
}
impl reverie::Stack for UnusedStack {
type StackGuard = UnusedStackGuard;
fn size(&self) -> usize {
panic!("readv event must not use a guest stack")
}
fn capacity(&self) -> usize {
panic!("readv event must not use a guest stack")
}
fn push<'stack, T>(&mut self, _value: T) -> Addr<'stack, T> {
panic!("readv event must not use a guest stack")
}
fn reserve<'stack, T>(&mut self) -> AddrMut<'stack, T> {
panic!("readv event must not use a guest stack")
}
fn commit(self) -> Result<Self::StackGuard, Errno> {
panic!("readv event must not use a guest stack")
}
}
#[reverie::tool]
impl GlobalRPC<GlobalState> for EventGuest {
async fn send_rpc(
&self,
message: <GlobalState as GlobalTool>::Request,
) -> <GlobalState as GlobalTool>::Response {
let response = match message.2 {
GlobalRequest::RequestResources(request, _) => {
assert_eq!(request.resources.len(), 1);
assert_eq!(
request.resources.get(&ResourceID::InternalIOPolling),
Some(&Permission::W)
);
self.polls.lock().unwrap().push(request.poll_attempt);
self.request_marks
.lock()
.unwrap()
.push(request.backend_runtime_bootstrap);
match request.poll_attempt {
0 => {}
1 => {
assert_eq!(self.injected_iovecs, [(FIRST_DEST, 8)]);
let (parked, resume) = self.retry_gate.lock().unwrap().take().unwrap();
parked.send(()).unwrap();
resume.await.unwrap();
}
attempt => panic!("unexpected extra readv poll {attempt}"),
}
GlobalResponse::RequestResources(ResumeStatus::Normal)
}
GlobalRequest::ReleaseAllResources => {
self.audit_call("release");
*self.releases.lock().unwrap() += 1;
GlobalResponse::ReleaseAllResources(())
}
request => panic!("unexpected readv event RPC: {request:?}"),
};
(None, response)
}
fn config(&self) -> &Config {
&self.config
}
}
#[reverie::tool]
impl Guest<Detcore> for EventGuest {
type Memory = EventMemory;
type Stack = UnusedStack;
fn tid(&self) -> Pid {
Pid::from_raw(self.thread.dettid.as_raw())
}
fn pid(&self) -> Pid {
Pid::from_raw(self.thread.detpid.unwrap().as_raw())
}
fn ppid(&self) -> Option<Pid> {
None
}
fn is_backend_runtime_bootstrap(&self) -> bool {
self.backend_runtime_bootstrap
}
fn memory(&self) -> Self::Memory {
self.memory.clone()
}
fn thread_state_mut(&mut self) -> &mut ThreadState<()> {
&mut self.thread
}
fn thread_state(&self) -> &ThreadState<()> {
&self.thread
}
async fn regs(&mut self) -> libc::user_regs_struct {
if self.audit_call("regs") {
return libc::user_regs_struct {
rip: 0x4000,
eflags: 2,
..unsafe { std::mem::zeroed() }
};
}
panic!("buffer observation must not request register evidence")
}
async fn set_regs(&mut self, _regs: libc::user_regs_struct) -> Result<(), Error> {
assert!(self.audit_call("set-regs"));
Ok(())
}
fn detlog_memory_regions(&self) -> Option<Vec<reverie::DetlogMemoryRegion>> {
if !self.audit_call("memory-regions") {
return None;
}
Some(vec![
reverie::DetlogMemoryRegion {
kind: reverie::DetlogRegionKind::Stack,
start: FIRST_DEST as u64,
end: (FIRST_DEST + 8) as u64,
},
reverie::DetlogMemoryRegion {
kind: reverie::DetlogRegionKind::Heap,
start: RETRY_DEST as u64,
end: (RETRY_DEST + 8) as u64,
},
])
}
async fn stack(&mut self) -> Self::Stack {
panic!("readv event must not use a guest stack")
}
async fn daemonize(&mut self) {
panic!("readv event must not daemonize")
}
async fn inject<S: SyscallInfo>(&mut self, syscall: S) -> Result<i64, Errno> {
let (number, args) = syscall.into_parts();
if let Syscall::Read(call) = Syscall::from_raw(number, args) {
assert_eq!(call.len(), 0);
self.injected_zero_reads += 1;
if call.fd() == -1 {
return Err(Errno::EBADF);
}
assert_eq!(call.fd(), FD);
assert_ne!(
self.thread.with_detfd(FD, |fd| fd.ty()).unwrap(),
FdType::Rng
);
return Err(Errno::EINVAL);
}
let Syscall::Readv(call) = Syscall::from_raw(number, args) else {
panic!("unexpected injected syscall {number}");
};
assert_eq!(call.fd(), FD);
assert_eq!(call.iov().unwrap().as_raw(), IOV);
assert_eq!(call.len(), 1);
assert_eq!(
self.thread.with_detfd(FD, |fd| fd.ty()).unwrap(),
FdType::Pipe,
"an RNG readv must be emulated without injection"
);
let (base, len) = self.memory.iovec(0);
self.injected_iovecs.push((base, len));
match self.injected_iovecs.len() {
1 => Err(Errno::EAGAIN),
2 => {
assert!(len >= PIPE_BYTES.len());
self.memory
.write_exact(AddrMut::from_raw(base).unwrap(), PIPE_BYTES)?;
Ok(PIPE_BYTES.len() as i64)
}
attempt => panic!("unexpected extra readv injection {attempt}"),
}
}
async fn tail_inject<S: SyscallInfo>(&mut self, _syscall: S) -> reverie::Never {
panic!("readv must return through the event observer")
}
fn set_timer(&mut self, _schedule: reverie::TimerSchedule) -> Result<(), Error> {
if self.audit_call("timer") {
return Ok(());
}
panic!("event fixture has no PMU timeslice")
}
fn set_timer_precise(&mut self, _schedule: reverie::TimerSchedule) -> Result<(), Error> {
if self.audit_call("timer") {
return Ok(());
}
panic!("event fixture has no PMU timeslice")
}
fn read_clock(&mut self) -> Result<u64, Error> {
if self.audit_call("clock") {
return Ok(0);
}
panic!("event fixture must not read a host clock")
}
}
fn event_guest(ty: FdType, retry_gate: Option<RetryGate>) -> (Detcore, EventGuest) {
let config = Config {
seed: 0,
rng_seed: Some(0),
sequentialize_threads: true,
recordreplay_modes: false,
record_preemptions: false,
max_timeslice: None,
detlog_io_buffers: true,
detlog_heap: false,
detlog_stack: false,
detlog_regs: false,
backend_is_kvm: true,
syscall_clobbers_virtualized_by_backend: true,
..Config::default()
};
let pid = DetPid::from_raw(1);
let mut thread = ThreadState::new(pid, &config, ());
thread.detpid = Some(pid);
let fd = DetFd::new(FD, OFlag::O_RDONLY, ty, OpenFileId::new(pid, 99));
if ty == FdType::Pipe {
fd.set_physically_nonblocking();
}
thread
.file_metadata
.lock()
.unwrap()
.file_handles
.insert(FD, fd);
let tool = <Detcore as Tool>::new(Pid::from_raw(1), &config);
let guest = EventGuest {
config,
thread,
memory: EventMemory::new(),
injected_iovecs: Vec::new(),
injected_zero_reads: 0,
polls: Mutex::new(Vec::new()),
request_marks: Mutex::new(Vec::new()),
releases: Mutex::new(0),
retry_gate: Mutex::new(retry_gate),
backend_runtime_bootstrap: false,
};
(tool, guest)
}
#[derive(Clone)]
struct BufferLog(Arc<Mutex<Vec<String>>>, &'static str);
impl Default for BufferLog {
fn default() -> Self {
Self(Arc::default(), "[iobuf]")
}
}
struct MessageVisitor(Option<String>);
impl Visit for MessageVisitor {
fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
if field.name() == "message" {
self.0 = Some(format!("{value:?}"));
}
}
}
impl Subscriber for BufferLog {
fn enabled(&self, metadata: &Metadata<'_>) -> bool {
*metadata.level() == Level::INFO
}
fn new_span(&self, _span: &Attributes<'_>) -> Id {
Id::from_u64(1)
}
fn record(&self, _span: &Id, _values: &Record<'_>) {}
fn record_follows_from(&self, _span: &Id, _follows: &Id) {}
fn event(&self, event: &Event<'_>) {
let mut visitor = MessageVisitor(None);
event.record(&mut visitor);
if let Some(message) = visitor.0
&& message.contains(self.1)
{
self.0.lock().unwrap().push(message);
}
}
fn enter(&self, _span: &Id) {}
fn exit(&self, _span: &Id) {}
}
fn readv(count: usize) -> Syscall {
reverie::syscalls::Readv::new()
.with_fd(FD)
.with_iov(Addr::from_raw(IOV))
.with_len(count)
.into()
}
fn assert_extent(message: &str, address: usize, bytes: &[u8]) {
assert_named_extent(message, "readv", address, bytes);
}
fn assert_named_extent(message: &str, name: &str, address: usize, bytes: &[u8]) {
let expected = format!(
"[iobuf][dtid 1] {name} in fd={FD} {address:#x}+{}->{}",
bytes.len(),
Digest::new(bytes)
);
assert!(
message.contains(&expected),
"expected {expected}, got {message}"
);
}
fn rng_vector_calls(count: usize) -> [(&'static str, Syscall, bool); 4] {
let preadv = reverie::syscalls::Preadv::new()
.with_fd(FD)
.with_iov(Addr::from_raw(IOV))
.with_iov_len(count)
.with_pos_l(19)
.with_pos_h(u64::MAX);
let preadv2 = reverie::syscalls::Preadv2::new()
.with_fd(FD)
.with_iov(Addr::from_raw(IOV))
.with_iov_len(count as u64)
.with_pos_l(19)
.with_pos_h(u64::MAX)
.with_flags(0);
[
("readv", readv(count), true),
("preadv", preadv.into(), false),
("preadv2", preadv2.into(), false),
("preadv2", preadv2.with_pos_l(u64::MAX).into(), true),
]
}
#[tokio::test(flavor = "current_thread")]
async fn backend_runtime_bootstrap_syscall_is_handled_but_not_charged_to_guest_time() {
struct Observed {
result: i64,
bytes: Vec<u8>,
syscall_count: u64,
random_offset: u64,
clock_before: LogicalTime,
clock_after: LogicalTime,
system_before: LogicalTime,
system_after: LogicalTime,
process_system_before: LogicalTime,
process_system_after: LogicalTime,
}
async fn one_readv(tool: &Detcore, guest: &mut EventGuest) -> Observed {
let memory = guest.memory.clone();
memory.put_iovec(0, FIRST_DEST, 8);
let clock_before = guest.thread.thread_logical_time.as_nanos();
let system_before = guest.thread.thread_logical_time.system_cpu_time();
let process_system_before = guest.thread.process_cpu_time().system;
let result = tool.handle_syscall_event(guest, readv(1)).await.unwrap();
Observed {
result,
bytes: memory.bytes(FIRST_DEST, 8),
syscall_count: guest.thread.stats.syscall_count,
random_offset: guest
.thread
.with_detfd(FD, |fd| fd.random_device_offset())
.unwrap(),
clock_before,
clock_after: guest.thread.thread_logical_time.as_nanos(),
system_before,
system_after: guest.thread.thread_logical_time.system_cpu_time(),
process_system_before,
process_system_after: guest.thread.process_cpu_time().system,
}
}
let (tool, mut plain) = event_guest(FdType::Rng, None);
let guest_call = one_readv(&tool, &mut plain).await;
assert!(
guest_call.clock_after > guest_call.clock_before,
"a guest syscall must advance logical time"
);
assert!(guest_call.system_after > guest_call.system_before);
assert!(guest_call.process_system_after > guest_call.process_system_before);
let (tool, mut booting) = event_guest(FdType::Rng, None);
booting.backend_runtime_bootstrap = true;
let bootstrap_call = one_readv(&tool, &mut booting).await;
assert_eq!(bootstrap_call.result, guest_call.result);
assert_eq!(bootstrap_call.bytes, guest_call.bytes);
assert_eq!(bootstrap_call.random_offset, guest_call.random_offset);
assert_eq!(bootstrap_call.syscall_count, 1);
assert_eq!(*booting.releases.lock().unwrap(), 1);
assert_eq!(
bootstrap_call.clock_after, bootstrap_call.clock_before,
"a backend-bootstrap syscall must not advance guest logical time"
);
assert_eq!(bootstrap_call.system_after, bootstrap_call.system_before);
assert_eq!(
bootstrap_call.process_system_after,
bootstrap_call.process_system_before
);
assert_eq!(bootstrap_call.clock_before, guest_call.clock_before);
booting.backend_runtime_bootstrap = false;
let after_window = one_readv(&tool, &mut booting).await;
assert_eq!(after_window.syscall_count, 2);
assert_eq!(after_window.clock_before, bootstrap_call.clock_after);
assert_eq!(after_window.clock_after, guest_call.clock_after);
assert_eq!(after_window.system_after, guest_call.system_after);
assert_eq!(
after_window.process_system_after,
guest_call.process_system_after
);
assert!(after_window.clock_after > after_window.clock_before);
}
#[tokio::test(flavor = "current_thread")]
async fn backend_runtime_bootstrap_window_charges_time_reads_and_caps_uncharged_syscalls() {
use reverie::syscalls::Sysno;
use crate::syscall_time::MAX_UNCHARGED_BOOTSTRAP_SYSCALLS;
use crate::syscall_time::observes_virtual_time;
const USAGE: usize = RETRY_DEST;
fn thread_rusage() -> Syscall {
reverie::syscalls::Getrusage::new()
.with_who(libc::RUSAGE_THREAD)
.with_usage(AddrMut::from_raw(USAGE))
.into()
}
async fn advance(tool: &Detcore, guest: &mut EventGuest, call: Syscall) -> u64 {
guest.memory.put_iovec(0, FIRST_DEST, 8);
let before = guest.thread.thread_logical_time.as_nanos().as_nanos();
tool.handle_syscall_event(guest, call).await.unwrap();
guest.thread.thread_logical_time.as_nanos().as_nanos() - before
}
fn reported_system_ns(guest: &EventGuest) -> u64 {
let usage: libc::rusage = guest
.memory
.read_value(Addr::<libc::rusage>::from_raw(USAGE).unwrap())
.unwrap();
usage.ru_stime.tv_sec as u64 * 1_000_000_000 + usage.ru_stime.tv_usec as u64 * 1_000
}
const TIME_READS: [Sysno; 11] = [
Sysno::gettimeofday,
Sysno::time,
Sysno::clock_gettime,
Sysno::sysinfo,
Sysno::times,
Sysno::getrusage,
Sysno::timerfd_gettime,
Sysno::timer_gettime,
Sysno::getitimer,
Sysno::adjtimex,
Sysno::clock_adjtime,
];
for sysno in TIME_READS {
assert!(observes_virtual_time(sysno), "{sysno} reads virtual time");
}
assert!(!observes_virtual_time(Sysno::readv));
assert!(!observes_virtual_time(Sysno::clock_getres));
let (tool, mut plain) = event_guest(FdType::Rng, None);
let time_read_cost = advance(&tool, &mut plain, thread_rusage()).await;
let first_read_outside = reported_system_ns(&plain);
let readv_cost = advance(&tool, &mut plain, readv(1)).await;
assert!(time_read_cost > 0 && readv_cost > 0);
let (tool, mut booting) = event_guest(FdType::Rng, None);
booting.backend_runtime_bootstrap = true;
let first_read_cost = advance(&tool, &mut booting, thread_rusage()).await;
let first_read = reported_system_ns(&booting);
let second_read_cost = advance(&tool, &mut booting, thread_rusage()).await;
let second_read = reported_system_ns(&booting);
assert_eq!(
first_read_cost, time_read_cost,
"a time read inside the window must advance logical time as it does outside"
);
assert_eq!(second_read_cost, time_read_cost);
assert_eq!(first_read, first_read_outside);
assert_ne!(
second_read, first_read,
"two time reads inside the window must observe different times"
);
assert_eq!(second_read - first_read, time_read_cost);
for sysno in TIME_READS {
assert!(
booting.thread.charge_syscall_time(true, sysno),
"{sysno} inside the window must be charged"
);
}
assert_eq!(booting.thread.uncharged_bootstrap_syscalls, 0);
let cap_logs = BufferLog(Arc::default(), "reached its cap");
let _subscriber = tracing::subscriber::set_default(cap_logs.clone());
for index in 0..MAX_UNCHARGED_BOOTSTRAP_SYSCALLS {
if index == MAX_UNCHARGED_BOOTSTRAP_SYSCALLS - 1 {
assert!(
cap_logs.0.lock().unwrap().is_empty(),
"the cap line was logged before the window reached the cap"
);
}
assert_eq!(
advance(&tool, &mut booting, readv(1)).await,
0,
"uncharged syscall {index} of the window advanced logical time"
);
}
assert_eq!(
booting.thread.uncharged_bootstrap_syscalls,
MAX_UNCHARGED_BOOTSTRAP_SYSCALLS
);
let cap_line = vec![format!(
"[dtid {}] backend runtime bootstrap window reached its cap of {} uncharged syscalls; every further syscall in this window is charged",
booting.thread.dettid, MAX_UNCHARGED_BOOTSTRAP_SYSCALLS
)];
assert_eq!(
*cap_logs.0.lock().unwrap(),
cap_line,
"reaching the cap must log exactly one info line"
);
assert_eq!(
advance(&tool, &mut booting, readv(1)).await,
readv_cost,
"the first syscall past the cap must be charged one syscall cost"
);
assert_eq!(advance(&tool, &mut booting, readv(1)).await, readv_cost);
assert_eq!(
booting.thread.uncharged_bootstrap_syscalls,
MAX_UNCHARGED_BOOTSTRAP_SYSCALLS
);
assert_eq!(
*cap_logs.0.lock().unwrap(),
cap_line,
"syscalls charged past the cap must not log the cap line again"
);
booting.backend_runtime_bootstrap = false;
assert_eq!(advance(&tool, &mut booting, readv(1)).await, readv_cost);
assert_eq!(booting.thread.uncharged_bootstrap_syscalls, 0);
booting.backend_runtime_bootstrap = true;
assert_eq!(advance(&tool, &mut booting, readv(1)).await, 0);
assert_eq!(booting.thread.uncharged_bootstrap_syscalls, 1);
assert_eq!(*cap_logs.0.lock().unwrap(), cap_line);
for index in 1..MAX_UNCHARGED_BOOTSTRAP_SYSCALLS {
if index == MAX_UNCHARGED_BOOTSTRAP_SYSCALLS - 1 {
assert_eq!(
*cap_logs.0.lock().unwrap(),
cap_line,
"the reopened window logged the cap line before reaching the cap"
);
}
assert_eq!(
advance(&tool, &mut booting, readv(1)).await,
0,
"uncharged syscall {index} of the reopened window advanced logical time"
);
}
assert_eq!(
booting.thread.uncharged_bootstrap_syscalls,
MAX_UNCHARGED_BOOTSTRAP_SYSCALLS
);
let two_cap_lines = vec![cap_line[0].clone(), cap_line[0].clone()];
assert_eq!(
*cap_logs.0.lock().unwrap(),
two_cap_lines,
"a reopened window that reaches the cap must log its own info line"
);
assert_eq!(
advance(&tool, &mut booting, readv(1)).await,
readv_cost,
"the first syscall past the reopened window's cap must be charged"
);
assert_eq!(*cap_logs.0.lock().unwrap(), two_cap_lines);
}
#[tokio::test(flavor = "current_thread")]
async fn bootstrap_syscall_marks_its_resource_requests_only_while_it_is_uncharged() {
use crate::syscall_time::MAX_UNCHARGED_BOOTSTRAP_SYSCALLS;
async fn pipe_readv(
in_window: bool,
uncharged_so_far: u32,
) -> (Vec<bool>, u64, EventGuest) {
let (parked_tx, parked_rx) = oneshot::channel();
let (resume_tx, resume_rx) = oneshot::channel();
let (tool, mut guest) = event_guest(FdType::Pipe, Some((parked_tx, resume_rx)));
guest.backend_runtime_bootstrap = in_window;
guest.thread.uncharged_bootstrap_syscalls = uncharged_so_far;
guest.memory.put_iovec(0, FIRST_DEST, 8);
let before = guest.thread.thread_logical_time.as_nanos().as_nanos();
let event = tool.handle_syscall_event(&mut guest, readv(1));
let resume = async {
parked_rx.await.unwrap();
resume_tx.send(()).unwrap();
};
let (result, ()) = tokio::time::timeout(std::time::Duration::from_secs(2), async {
tokio::join!(event, resume)
})
.await
.expect("readv did not cross its retry resource wait");
assert_eq!(result.unwrap(), 4);
assert_eq!(*guest.polls.lock().unwrap(), [0, 1]);
assert!(
!guest.thread.in_uncharged_bootstrap_syscall,
"the mark must be cleared when the syscall finishes"
);
let advance = guest.thread.thread_logical_time.as_nanos().as_nanos() - before;
let marks = guest.request_marks.lock().unwrap().clone();
(marks, advance, guest)
}
let (guest_marks, guest_cost, _) = pipe_readv(false, 0).await;
assert!(guest_cost > 0, "a guest syscall must advance logical time");
assert_eq!(
guest_marks,
[false, false],
"a guest syscall's requests must not be marked"
);
let (window_marks, window_cost, booting) = pipe_readv(true, 0).await;
assert_eq!(window_cost, 0, "the window's syscall must be uncharged");
assert_eq!(booting.thread.uncharged_bootstrap_syscalls, 1);
assert_eq!(
window_marks,
[true, true],
"every request of an uncharged window syscall must be marked"
);
let (capped_marks, capped_cost, _) =
pipe_readv(true, MAX_UNCHARGED_BOOTSTRAP_SYSCALLS).await;
assert_eq!(
capped_cost, guest_cost,
"a window syscall past the cap must be charged"
);
assert_eq!(
capped_marks,
[false, false],
"a charged window syscall's requests must not be marked"
);
}
#[tokio::test(flavor = "current_thread")]
async fn pipe_readv_event_observes_destination_imported_after_retry_wait() {
let logs = BufferLog::default();
let _subscriber = tracing::subscriber::set_default(logs.clone());
let (parked_tx, parked_rx) = oneshot::channel();
let (resume_tx, resume_rx) = oneshot::channel();
let (tool, mut guest) = event_guest(FdType::Pipe, Some((parked_tx, resume_rx)));
let memory = guest.memory.clone();
memory.put_iovec(0, FIRST_DEST, 8);
let event = tool.handle_syscall_event(&mut guest, readv(1));
let mutate_during_retry = async {
parked_rx.await.unwrap();
assert_eq!(memory.iovec(0), (FIRST_DEST, 8));
assert_eq!(memory.bytes(FIRST_DEST, 8), [CANARY; 8]);
memory.put_iovec(0, RETRY_DEST, 8);
resume_tx.send(()).unwrap();
};
let (result, ()) = tokio::time::timeout(std::time::Duration::from_secs(2), async {
tokio::join!(event, mutate_during_retry)
})
.await
.expect("readv did not cross its retry resource wait");
assert_eq!(result.unwrap(), 4);
assert_eq!(guest.injected_iovecs, [(FIRST_DEST, 8), (RETRY_DEST, 8)]);
assert_eq!(*guest.polls.lock().unwrap(), [0, 1]);
assert_eq!(*guest.releases.lock().unwrap(), 1);
assert_eq!(guest.thread.stats.syscall_count, 1);
assert_eq!(memory.bytes(FIRST_DEST, 8), [CANARY; 8]);
assert_eq!(
memory.bytes(RETRY_DEST, 4).as_slice(),
PIPE_BYTES.as_slice()
);
assert_eq!(memory.bytes(RETRY_DEST + 4, 4), [CANARY; 4]);
let messages = logs.0.lock().unwrap();
assert_eq!(messages.len(), 1, "{messages:?}");
assert_extent(&messages[0], RETRY_DEST, PIPE_BYTES);
}
#[tokio::test(flavor = "current_thread")]
async fn rng_readv_event_observes_imported_array_after_output_overwrites_it() {
for (name, call, advances) in rng_vector_calls(2) {
let logs = BufferLog::default();
let _subscriber = tracing::subscriber::set_default(logs.clone());
let (tool, mut guest) = event_guest(FdType::Rng, None);
let memory = guest.memory.clone();
let second_iovec = IOV + std::mem::size_of::<libc::iovec>();
memory.put_iovec(0, second_iovec, 16);
memory.put_iovec(1, FIRST_DEST, 4);
let original_second_iovec = memory.bytes(second_iovec, 16);
let result = tool.handle_syscall_event(&mut guest, call).await;
assert_eq!(result.unwrap(), 20);
let expected = if advances {
[
41, 114, 187, 4, 77, 150, 223, 40, 113, 186, 3, 76, 149, 222, 39, 112, 185, 2,
75, 148,
]
} else {
[
148, 221, 38, 111, 184, 1, 74, 147, 220, 37, 110, 183, 0, 73, 146, 219, 36,
109, 182, 255,
]
};
assert_eq!(memory.bytes(second_iovec, 16), expected[..16]);
assert_ne!(memory.bytes(second_iovec, 16), original_second_iovec);
assert_eq!(memory.bytes(FIRST_DEST, 4), expected[16..]);
assert_eq!(memory.bytes(FIRST_DEST + 4, 4), [CANARY; 4]);
assert!(guest.injected_iovecs.is_empty());
assert!(guest.polls.lock().unwrap().is_empty());
assert_eq!(*guest.releases.lock().unwrap(), 1);
assert_eq!(guest.thread.stats.syscall_count, 1);
assert_eq!(
guest
.thread
.with_detfd(FD, |fd| fd.random_device_offset())
.unwrap(),
if advances { 20 } else { 0 }
);
let messages = logs.0.lock().unwrap();
assert_eq!(messages.len(), 2, "{messages:?}");
assert_named_extent(&messages[0], name, second_iovec, &expected[..16]);
assert_named_extent(&messages[1], name, FIRST_DEST, &expected[16..]);
}
}
#[tokio::test(flavor = "current_thread")]
async fn rng_readv_event_digest_failure_is_terminal_after_bytes_and_cursor_commit() {
for (name, call, advances) in rng_vector_calls(2) {
for fault in [Errno::EFAULT, Errno::EIO] {
let logs = BufferLog::default();
let _subscriber = tracing::subscriber::set_default(logs.clone());
let (tool, mut guest) = event_guest(FdType::Rng, None);
let memory = guest.memory.clone();
memory.put_iovec(0, RETRY_DEST, 3);
memory.put_iovec(1, FIRST_DEST, 5);
memory.1.lock().unwrap().digest_error = Some(fault);
let result = tool.handle_syscall_event(&mut guest, call).await;
let Err(Error::Tool(error)) = result else {
panic!("completed RNG readv digest failure became guest result: {result:?}");
};
assert!(
error
.to_string()
.contains("RNG vector observation after output commit")
);
assert!(error.chain().any(|cause| {
matches!(cause.downcast_ref::<Error>(), Some(Error::Errno(error)) if *error == fault)
}));
let expected = if advances {
[41, 114, 187, 4, 77, 150, 223, 40]
} else {
[148, 221, 38, 111, 184, 1, 74, 147]
};
assert_eq!(memory.bytes(RETRY_DEST, 3), expected[..3]);
assert_eq!(memory.bytes(RETRY_DEST + 3, 1), [CANARY]);
assert_eq!(memory.bytes(FIRST_DEST, 5), expected[3..]);
assert_eq!(memory.bytes(FIRST_DEST + 5, 1), [CANARY]);
assert_eq!(
guest
.thread
.with_detfd(FD, |fd| fd.random_device_offset())
.unwrap(),
if advances { 8 } else { 0 }
);
assert_eq!(*guest.releases.lock().unwrap(), 1);
assert!(guest.injected_iovecs.is_empty());
let reads = memory.1.lock().unwrap();
assert_eq!(reads.imported_entries, 2);
assert_eq!(reads.observer_reads, [(RETRY_DEST, 3), (FIRST_DEST, 5)]);
let messages = logs.0.lock().unwrap();
assert_eq!(messages.len(), 1);
assert_named_extent(&messages[0], name, RETRY_DEST, &expected[..3]);
}
}
}
#[tokio::test(flavor = "current_thread")]
async fn rng_readv_event_observer_configuration_and_subscription_are_inert() {
for (name, call, advances) in rng_vector_calls((1_usize << 32) | 2) {
for (configured, subscribed) in [(true, true), (false, true), (true, false)] {
let logs = BufferLog::default();
let dispatch = if subscribed {
tracing::Dispatch::new(logs.clone())
} else {
tracing::Dispatch::new(tracing::subscriber::NoSubscriber::default())
};
let _subscriber = tracing::dispatcher::set_default(&dispatch);
let (mut tool, mut guest) = event_guest(FdType::Rng, None);
tool.cfg.detlog_io_buffers = configured;
guest.config.detlog_io_buffers = configured;
let memory = guest.memory.clone();
memory.put_iovec(0, FIRST_DEST, 3);
memory.put_iovec(1, RETRY_DEST, 5);
let result = tool.handle_syscall_event(&mut guest, call).await;
assert_eq!(result.unwrap(), 8);
let expected = if advances {
[41, 114, 187, 4, 77, 150, 223, 40]
} else {
[148, 221, 38, 111, 184, 1, 74, 147]
};
assert_eq!(memory.bytes(FIRST_DEST, 3), expected[..3]);
assert_eq!(memory.bytes(FIRST_DEST + 3, 1), [CANARY]);
assert_eq!(memory.bytes(RETRY_DEST, 5), expected[3..]);
assert_eq!(memory.bytes(RETRY_DEST + 5, 1), [CANARY]);
assert_eq!(
guest
.thread
.with_detfd(FD, |fd| fd.random_device_offset())
.unwrap(),
if advances { 8 } else { 0 }
);
assert_eq!(*guest.releases.lock().unwrap(), 1);
assert!(guest.injected_iovecs.is_empty());
let reads = memory.1.lock().unwrap();
assert_eq!(reads.imported_entries, 2);
let messages = logs.0.lock().unwrap();
if configured && subscribed {
assert_eq!(reads.observer_reads, [(FIRST_DEST, 3), (RETRY_DEST, 5)]);
assert_eq!(messages.len(), 2);
assert_named_extent(&messages[0], name, FIRST_DEST, &expected[..3]);
assert_named_extent(&messages[1], name, RETRY_DEST, &expected[3..]);
} else {
assert!(
reads.observer_reads.is_empty(),
"disabled observer read guest bytes"
);
assert!(messages.is_empty(), "disabled observer emitted hashes");
}
}
}
}
#[tokio::test(flavor = "current_thread")]
async fn rng_readv_event_zero_and_malformed_requests_have_no_observer_effects() {
for (count, length, expected, imports) in [
(0, 3, Ok(0), 0),
(1025, 3, Err(Errno::EINVAL), 0),
(1, usize::MAX, Err(Errno::EINVAL), 1),
] {
for (_, call, _) in rng_vector_calls(count) {
let logs = BufferLog::default();
let _subscriber = tracing::subscriber::set_default(logs.clone());
let (tool, mut guest) = event_guest(FdType::Rng, None);
let memory = guest.memory.clone();
memory.put_iovec(0, FIRST_DEST, length);
let result = tool.handle_syscall_event(&mut guest, call).await;
assert_eq!(
result.map_err(|error| error.into_errno().unwrap()),
expected
);
assert_eq!(memory.bytes(FIRST_DEST, 8), [CANARY; 8]);
assert_eq!(
guest
.thread
.with_detfd(FD, |fd| fd.random_device_offset())
.unwrap(),
0
);
assert_eq!(*guest.releases.lock().unwrap(), 1);
let reads = memory.1.lock().unwrap();
assert_eq!(reads.imported_entries, imports);
assert!(reads.observer_reads.is_empty());
assert!(logs.0.lock().unwrap().is_empty());
}
}
}
#[tokio::test(flavor = "current_thread")]
async fn rng_scalar_zero_read_checks_access_without_touching_replay_placeholders() {
for mode in [
OFlag::O_RDONLY,
OFlag::O_RDWR,
OFlag::O_WRONLY,
OFlag::O_ACCMODE,
OFlag::O_PATH,
] {
for address in [0, usize::MAX] {
let logs = BufferLog::default();
let _subscriber = tracing::subscriber::set_default(logs.clone());
let (tool, mut guest) = event_guest(FdType::Rng, None);
let fd = DetFd::new(
FD,
mode,
FdType::Rng,
OpenFileId::new(DetTid::from_raw(1), 99),
);
fd.advance_random_device_offset(7);
guest
.thread
.file_metadata
.lock()
.unwrap()
.file_handles
.insert(FD, fd);
let call = reverie::syscalls::Read::new()
.with_fd(FD)
.with_buf(AddrMut::from_raw(address))
.with_len(0);
let result = tool.handle_syscall_event(&mut guest, call.into()).await;
let expected = if mode == OFlag::O_RDONLY || mode == OFlag::O_RDWR {
if address == 0 {
Ok(0)
} else {
Err(Errno::EFAULT)
}
} else {
Err(Errno::EBADF)
};
assert_eq!(
result.map_err(|error| error.into_errno().unwrap()),
expected
);
assert_eq!(
guest
.thread
.with_detfd(FD, |fd| fd.random_device_offset())
.unwrap(),
7
);
assert!(guest.injected_iovecs.is_empty());
assert_eq!(guest.injected_zero_reads, 0);
assert!(guest.polls.lock().unwrap().is_empty());
let reads = guest.memory.1.lock().unwrap();
assert_eq!(reads.imported_entries, 0);
assert!(reads.observer_reads.is_empty());
assert!(logs.0.lock().unwrap().is_empty());
}
}
}
#[tokio::test(flavor = "current_thread")]
async fn non_rng_zero_reads_keep_the_backend_result() {
for (fd, expected) in [(FD, Errno::EINVAL), (-1, Errno::EBADF)] {
let (tool, mut guest) = event_guest(FdType::Regular, None);
let call = reverie::syscalls::Read::new().with_fd(fd).with_len(0);
let result = tool.handle_syscall_event(&mut guest, call.into()).await;
assert!(matches!(result, Err(Error::Errno(error)) if error == expected));
assert_eq!(guest.injected_zero_reads, usize::from(fd == FD));
assert!(guest.memory.1.lock().unwrap().observer_reads.is_empty());
}
}
}
#[cfg(test)]
mod tests {
use reverie::syscalls;
use reverie::syscalls::LocalMemory;
#[test]
fn rng_observation_capacity_failure_is_a_tool_error() {
let error = rng_observation_vec::<u8>(usize::MAX, "digest payload").unwrap_err();
let Error::Tool(error) = error else {
panic!("post-commit allocation failure must not become guest errno");
};
assert!(error.to_string().contains(
"RNG vector observation after output commit: cannot reserve digest payload:"
));
assert!(error.to_string().contains("digest payload"));
}
#[test]
fn chunk_digests_name_the_first_differing_window() {
let len = 1468; let left = vec![0xABu8; len];
let mut right = left.clone();
let offset = 900;
right[offset] ^= 0xFF; let (chunk, l) = chunk_digests(&left);
let (_, r) = chunk_digests(&right);
assert_eq!(chunk, 256);
assert_eq!(l.len(), r.len());
let first_diff = l.iter().zip(&r).position(|(a, b)| a != b);
assert_eq!(
first_diff,
Some(offset / chunk),
"must name the window holding the flip"
);
assert_eq!(l.iter().zip(&r).filter(|(a, b)| a != b).count(), 1);
}
#[test]
fn identical_buffers_produce_identical_chunk_lists() {
let buf = vec![7u8; 1468];
assert_eq!(chunk_digests(&buf).1, chunk_digests(&buf.clone()).1);
}
#[test]
fn a_single_chunk_extent_emits_no_list() {
assert!(chunk_digests(&vec![0u8; CHUNK_MIN]).1.is_empty());
assert!(chunk_digests(&[]).1.is_empty());
}
#[test]
fn the_chunk_count_is_capped_and_still_spans_the_extent() {
for len in [2048usize, 65536, 1 << 20] {
let (chunk, chunks) = chunk_digests(&vec![0u8; len]);
assert!(
chunks.len() <= CHUNK_CAP,
"len {len} produced {} chunks",
chunks.len()
);
assert!(
chunk * chunks.len() >= len,
"chunks must span the whole extent"
);
}
}
use super::*;
#[test]
fn clamp_bounds_the_extent_by_what_was_moved_not_by_capacity() {
assert_eq!(
clamp(Some(0x1000), 4096, 10),
vec![BufferExtent {
addr: 0x1000,
len: 10
}]
);
}
#[test]
fn clamp_never_exceeds_capacity() {
assert_eq!(
clamp(Some(0x1000), 64, 4096),
vec![BufferExtent {
addr: 0x1000,
len: 64
}]
);
}
#[test]
fn nothing_is_hashed_for_a_null_buffer_or_an_empty_move() {
assert!(clamp(None, 4096, 10).is_empty());
assert!(clamp(Some(0x1000), 4096, 0).is_empty());
assert!(clamp(Some(0x1000), 4096, -1).is_empty());
}
#[test]
fn poll_hashes_the_whole_array_not_a_prefix() {
assert_eq!(
whole(Some(0x2000), 3 * std::mem::size_of::<libc::pollfd>() as u64),
vec![BufferExtent {
addr: 0x2000,
len: 24
}]
);
assert!(whole(None, 24).is_empty());
assert!(whole(Some(0x2000), 0).is_empty());
}
#[test]
fn the_poll_family_is_exempt_from_the_return_value_guard() {
let nfds = 3;
for call in [
Syscall::Poll(
syscalls::Poll::new()
.with_fds(AddrMut::from_raw(0x2000))
.with_nfds(nfds),
),
Syscall::Ppoll(
syscalls::Ppoll::new()
.with_fds(AddrMut::from_raw(0x2000))
.with_nfds(nfds),
),
] {
assert!(
!ret_gates_output(&call),
"{} returns a count of READY DESCRIPTORS, not a byte count, so a \
zero return must not suppress its extent",
call.name()
);
}
for call in [
Syscall::Read(syscalls::Read::new()),
Syscall::Recvmsg(syscalls::Recvmsg::new()),
Syscall::Recvmmsg(syscalls::Recvmmsg::new()),
Syscall::Sendmmsg(syscalls::Sendmmsg::new()),
] {
assert!(
ret_gates_output(&call),
"{} returns a byte or message count, so a zero return means \
nothing was written",
call.name()
);
}
}
#[test]
fn direction_separates_kernel_produced_from_guest_produced() {
assert_eq!(Direction::In.as_str(), "in");
assert_eq!(Direction::Out.as_str(), "out");
assert_eq!(
direction(&Syscall::Preadv2(syscalls::Preadv2::new())),
Direction::In
);
assert_eq!(
direction(&Syscall::Pwritev2(syscalls::Pwritev2::new())),
Direction::Out
);
assert_eq!(
direction(&Syscall::Sendmmsg(syscalls::Sendmmsg::new())),
Direction::Out
);
}
#[test]
fn vectored_v2_calls_are_hashed_over_their_complete_iovec_prefix() {
for sysno in [syscalls::Sysno::preadv2, syscalls::Sysno::pwritev2] {
assert!(
HASHED_SYSCALLS.contains(&sysno),
"{sysno:?} must remain in strict io-buffer observation"
);
}
let first = [0_u8; 4];
let second = [0_u8; 8];
let third = [0_u8; 16];
let iovecs = [
libc::iovec {
iov_base: first.as_ptr() as *mut libc::c_void,
iov_len: first.len(),
},
libc::iovec {
iov_base: second.as_ptr() as *mut libc::c_void,
iov_len: second.len(),
},
libc::iovec {
iov_base: third.as_ptr() as *mut libc::c_void,
iov_len: third.len(),
},
];
let iov = Addr::from_ptr(iovecs.as_ptr());
let calls = [
Syscall::Preadv2(syscalls::Preadv2::new().with_iov(iov).with_iov_len(3)),
Syscall::Pwritev2(syscalls::Pwritev2::new().with_iov(iov).with_iov_len(3)),
];
let expected = vec![
BufferExtent {
addr: first.as_ptr() as u64,
len: 4,
},
BufferExtent {
addr: second.as_ptr() as u64,
len: 2,
},
];
let memory = LocalMemory::new();
for call in calls {
assert_eq!(
iovec_extent_arguments(&call),
Some((iovecs.as_ptr() as usize, 3))
);
assert_eq!(extents(&memory, &call, 6).unwrap(), expected);
assert!(extents(&memory, &call, 0).unwrap().is_empty());
assert!(extents(&memory, &call, -1).unwrap().is_empty());
}
}
#[test]
fn vectored_io_extents_stop_at_the_completed_byte_count() {
let iovecs = [
libc::iovec {
iov_base: 0x1000usize as *mut libc::c_void,
iov_len: 4,
},
libc::iovec {
iov_base: 0x2000usize as *mut libc::c_void,
iov_len: 8,
},
libc::iovec {
iov_base: 0x3000usize as *mut libc::c_void,
iov_len: 16,
},
];
assert_eq!(
iovec_extents_from_slice(&iovecs, 6),
vec![
BufferExtent {
addr: 0x1000,
len: 4,
},
BufferExtent {
addr: 0x2000,
len: 2,
},
]
);
assert!(iovec_extents_from_slice(&iovecs, 0).is_empty());
assert!(iovec_extents_from_slice(&iovecs, -1).is_empty());
}
#[test]
fn sendmmsg_hashes_only_completed_messages_and_each_reported_byte_prefix() {
assert!(HASHED_SYSCALLS.contains(&syscalls::Sysno::sendmmsg));
let first = [0_u8; 4];
let second = [0_u8; 8];
let third = [0_u8; 16];
let first_iovecs = [
libc::iovec {
iov_base: first.as_ptr() as *mut libc::c_void,
iov_len: first.len(),
},
libc::iovec {
iov_base: second.as_ptr() as *mut libc::c_void,
iov_len: second.len(),
},
libc::iovec {
iov_base: third.as_ptr() as *mut libc::c_void,
iov_len: third.len(),
},
];
let fourth = [0_u8; 5];
let fifth = [0_u8; 7];
let second_iovecs = [
libc::iovec {
iov_base: fourth.as_ptr() as *mut libc::c_void,
iov_len: fourth.len(),
},
libc::iovec {
iov_base: fifth.as_ptr() as *mut libc::c_void,
iov_len: fifth.len(),
},
];
let mut headers: [libc::mmsghdr; 2] = unsafe { std::mem::zeroed() };
headers[0].msg_hdr.msg_iov = first_iovecs.as_ptr() as *mut libc::iovec;
headers[0].msg_hdr.msg_iovlen = first_iovecs.len();
headers[0].msg_len = 6;
headers[1].msg_hdr.msg_iov = second_iovecs.as_ptr() as *mut libc::iovec;
headers[1].msg_hdr.msg_iovlen = second_iovecs.len();
headers[1].msg_len = 9;
let call = Syscall::Sendmmsg(
syscalls::Sendmmsg::new()
.with_msgvec(Addr::from_ptr(headers.as_ptr().cast::<libc::msghdr>()))
.with_vlen(2),
);
let memory = LocalMemory::new();
let first_message = vec![
BufferExtent {
addr: first.as_ptr() as u64,
len: 4,
},
BufferExtent {
addr: second.as_ptr() as u64,
len: 2,
},
];
assert_eq!(extents(&memory, &call, 1).unwrap(), first_message);
let both_messages = vec![
BufferExtent {
addr: first.as_ptr() as u64,
len: 4,
},
BufferExtent {
addr: second.as_ptr() as u64,
len: 2,
},
BufferExtent {
addr: fourth.as_ptr() as u64,
len: 5,
},
BufferExtent {
addr: fifth.as_ptr() as u64,
len: 4,
},
];
assert_eq!(extents(&memory, &call, 2).unwrap(), both_messages);
assert!(extents(&memory, &call, 0).unwrap().is_empty());
assert!(extents(&memory, &call, -1).unwrap().is_empty());
assert_eq!(completed_mmsghdr_count(1, 2), 1);
}
}