Skip to main content

detcore/syscalls/
io.rs

1/*
2 * Copyright (c) Meta Platforms, Inc. and affiliates.
3 * All rights reserved.
4 *
5 * This source code is licensed under the BSD-style license found in the
6 * LICENSE file in the root directory of this source tree.
7 */
8
9//! System calls dealing with IO and networking.
10//!
11//! Of course this overlaps somewhat with "files.rs".
12
13use std::net::Ipv4Addr;
14use std::net::Ipv6Addr;
15use std::os::unix::io::RawFd;
16use std::time::Duration;
17
18use nix::fcntl::OFlag;
19use reverie::Errno;
20use reverie::Error;
21use reverie::Guest;
22use reverie::Stack;
23use reverie::syscalls;
24use reverie::syscalls::Addr;
25use reverie::syscalls::AddrMut;
26use reverie::syscalls::Displayable;
27use reverie::syscalls::MemoryAccess;
28use reverie::syscalls::Syscall;
29use reverie::syscalls::SyscallInfo;
30use reverie::syscalls::Timespec;
31use tracing::debug;
32use tracing::trace;
33
34use crate::config::SchedHeuristic;
35use crate::fd::FdType;
36use crate::record_or_replay::RecordOrReplay;
37use crate::resources::Permission;
38use crate::resources::ResourceID;
39use crate::resources::Resources;
40use crate::resources::SABRE_LOOPBACK_POLL_YIELD_FYI;
41use crate::scheduler::runqueue::FIRST_PRIORITY;
42use crate::syscalls::helpers::NonblockableSyscall;
43use crate::syscalls::helpers::millis_duration_to_absolute_timeout;
44use crate::syscalls::helpers::record_retry_event;
45use crate::syscalls::helpers::retry_nonblocking_syscall_with_timeout;
46use crate::syscalls::signal::read_kernel_sigset;
47use crate::tool_global::*;
48use crate::tool_local::Detcore;
49use crate::types::DetTid;
50use crate::types::LogicalTime;
51
52// Printing helper
53// TODO: this should be subsumed by better syscall printing.
54fn print_poll(call: &syscalls::Poll) {
55    let len = call.nfds();
56    debug!("POLL: on {} fds, timeout {}", len, call.timeout());
57    // TODO: nicer API for reading arrays from the guest:
58    unsafe {
59        for i in 0..len {
60            debug!(
61                "POLL: fd {} = {}",
62                i,
63                call.fds().unwrap().offset(i as isize)
64            );
65        }
66    }
67}
68
69/// Build the scheduler request for a zero-timeout poll.
70///
71/// An empty request only rotates the caller within its current priority band. That is not a
72/// yield when a busy poller has a higher priority than the producer whose readiness it is
73/// probing. `SchedYield` excludes the caller from the next selection, allowing exactly one
74/// other runnable guest to make progress before the nonblocking poll executes. Limit that strong
75/// yield to SaBRe tasks with a loopback peer: libcurl alternates its loopback socket with an
76/// internal wakeup fd, and either zero-timeout probe can otherwise starve the peer. General
77/// build-tool polling retains the existing empty-turn behavior.
78fn zero_timeout_poll_request(dettid: DetTid, yield_to_peer: bool) -> Resources {
79    let mut request = Resources::new(dettid);
80    if yield_to_peer {
81        request.insert(ResourceID::SchedYield, Permission::W);
82        request.fyi(SABRE_LOOPBACK_POLL_YIELD_FYI);
83    }
84    request
85}
86
87fn connect_result_allows_peer_classification(result: &Result<i64, Error>) -> bool {
88    match result {
89        Ok(_) => true,
90        Err(Error::Errno(errno)) => *errno == Errno::EINPROGRESS,
91        Err(_) => false,
92    }
93}
94
95const KERNEL_SIGSET_SIZE: usize = std::mem::size_of::<u64>();
96const PSELECT6_INTERNAL_MAX_NFDS: i32 = (std::mem::size_of::<libc::c_ulong>() * 8) as i32;
97
98#[derive(Clone, Copy)]
99#[repr(C)]
100struct Pselect6SigmaskArg {
101    sigmask: usize,
102    sigsetsize: usize,
103}
104
105fn pselect6_fd_set_len(nfds: i32) -> Result<usize, Errno> {
106    let nfds = usize::try_from(nfds).map_err(|_| Errno::EINVAL)?;
107    let bits_per_word = std::mem::size_of::<libc::c_ulong>() * 8;
108    Ok(nfds.div_ceil(bits_per_word) * std::mem::size_of::<libc::c_ulong>())
109}
110
111fn pselect6_probe_result(result: Result<i64, Errno>) -> Result<i64, Errno> {
112    match result {
113        // The injected syscall runs outside the guest's original restart frame.
114        // Do not expose this kernel-internal restart instruction at the
115        // rewritten pselect6 call site.
116        Err(Errno::ERESTARTSYS) => Err(Errno::EINTR),
117        result => result,
118    }
119}
120
121fn read_pselect6_fd_set<T, G>(
122    guest: &mut G,
123    address: Option<AddrMut<'_, libc::fd_set>>,
124    len: usize,
125) -> Result<Option<Vec<u8>>, Error>
126where
127    T: RecordOrReplay,
128    G: Guest<Detcore<T>>,
129{
130    let Some(address) = address else {
131        return Ok(None);
132    };
133    let mut bytes = vec![0; len];
134    if len != 0 {
135        guest
136            .memory()
137            .read_exact(address.cast(), &mut bytes)
138            .map_err(|_| Errno::EFAULT)?;
139    }
140    Ok(Some(bytes))
141}
142
143fn write_pselect6_fd_set<T, G>(
144    guest: &mut G,
145    address: Option<AddrMut<'_, libc::fd_set>>,
146    bytes: &Option<Vec<u8>>,
147) -> Result<(), Error>
148where
149    T: RecordOrReplay,
150    G: Guest<Detcore<T>>,
151{
152    if let (Some(address), Some(bytes)) = (address, bytes)
153        && !bytes.is_empty()
154    {
155        guest
156            .memory()
157            .write_exact(address.cast(), bytes)
158            .map_err(|_| Errno::EFAULT)?;
159    }
160    Ok(())
161}
162
163fn copy_pselect6_fd_set<T, G>(
164    guest: &mut G,
165    source: Option<AddrMut<'_, libc::fd_set>>,
166    destination: Option<AddrMut<'_, libc::fd_set>>,
167    len: usize,
168) -> Result<(), Error>
169where
170    T: RecordOrReplay,
171    G: Guest<Detcore<T>>,
172{
173    if let (Some(source), Some(destination)) = (source, destination)
174        && len != 0
175    {
176        let mut bytes = vec![0; len];
177        guest
178            .memory()
179            .read_exact(source.cast(), &mut bytes)
180            .map_err(|_| Errno::EFAULT)?;
181        guest
182            .memory()
183            .write_exact(destination.cast(), &bytes)
184            .map_err(|_| Errno::EFAULT)?;
185    }
186    Ok(())
187}
188
189fn ppoll_timeout_duration(timeout: Timespec) -> Result<Duration, Errno> {
190    let seconds = u64::try_from(timeout.tv_sec).map_err(|_| Errno::EINVAL)?;
191    let nanoseconds = u32::try_from(timeout.tv_nsec).map_err(|_| Errno::EINVAL)?;
192    if nanoseconds >= 1_000_000_000 {
193        return Err(Errno::EINVAL);
194    }
195    Ok(Duration::new(seconds, nanoseconds))
196}
197
198fn select_timeout_duration(timeout: libc::timeval) -> Result<Duration, Errno> {
199    let seconds = u64::try_from(timeout.tv_sec).map_err(|_| Errno::EINVAL)?;
200    let microseconds = u64::try_from(timeout.tv_usec).map_err(|_| Errno::EINVAL)?;
201    // Linux rejects select timeouts whose microsecond field is out of range.
202    if microseconds >= 1_000_000 {
203        return Err(Errno::EINVAL);
204    }
205    Ok(Duration::new(seconds, (microseconds * 1_000) as u32))
206}
207
208fn timespec_from_duration(duration: Duration) -> Timespec {
209    Timespec {
210        tv_sec: duration.as_secs() as libc::time_t,
211        tv_nsec: duration.subsec_nanos() as libc::c_long,
212    }
213}
214
215const SCM_TIMESTAMP_OLD: libc::c_int = 29;
216const SCM_TIMESTAMPNS_OLD: libc::c_int = 35;
217const SCM_TIMESTAMPING_OLD: libc::c_int = 37;
218const SCM_TIMESTAMP_NEW: libc::c_int = 63;
219const SCM_TIMESTAMPNS_NEW: libc::c_int = 64;
220const SCM_TIMESTAMPING_NEW: libc::c_int = 65;
221const MAX_CONTROL_BYTES: usize = 64 * 1024;
222
223#[derive(Clone, Copy)]
224enum SocketTimestampKind {
225    Timeval,
226    Timespec,
227    Timestamping,
228}
229
230#[derive(Clone, Copy)]
231struct SocketTimestampMessage {
232    data_offset: usize,
233    available_data_len: usize,
234    kind: SocketTimestampKind,
235}
236
237fn cmsg_align(length: usize) -> usize {
238    let alignment = std::mem::size_of::<usize>();
239    (length + alignment - 1) & !(alignment - 1)
240}
241
242fn read_control_value<T: Copy>(bytes: &[u8]) -> Option<T> {
243    if bytes.len() < std::mem::size_of::<T>() {
244        return None;
245    }
246    // SAFETY: the length check guarantees a complete T, and read_unaligned
247    // permits the control buffer's byte alignment.
248    Some(unsafe { bytes.as_ptr().cast::<T>().read_unaligned() })
249}
250
251fn write_control_value<T: Copy>(bytes: &mut [u8], value: T) -> bool {
252    if bytes.len() < std::mem::size_of::<T>() {
253        return false;
254    }
255    // SAFETY: the length check guarantees room for T, and write_unaligned
256    // permits the control buffer's byte alignment.
257    unsafe { bytes.as_mut_ptr().cast::<T>().write_unaligned(value) };
258    true
259}
260
261fn write_control_prefix<T: Copy>(bytes: &mut [u8], value: T) -> usize {
262    let value_len = std::mem::size_of::<T>();
263    let write_len = bytes.len().min(value_len);
264    // SAFETY: `value` is alive for this copy and the resulting byte view has
265    // exactly its initialized object representation.
266    let value_bytes =
267        unsafe { std::slice::from_raw_parts((&value as *const T).cast::<u8>(), value_len) };
268    bytes[..write_len].copy_from_slice(&value_bytes[..write_len]);
269    write_len
270}
271
272fn socket_timestamp_messages(control: &[u8]) -> Vec<SocketTimestampMessage> {
273    let header_len = cmsg_align(std::mem::size_of::<libc::cmsghdr>());
274    let mut messages = Vec::new();
275    let mut offset = 0usize;
276
277    while let Some(header_bytes) = control.get(offset..) {
278        let Some(header) = read_control_value::<libc::cmsghdr>(header_bytes) else {
279            break;
280        };
281        if header.cmsg_len < header_len {
282            break;
283        }
284        let Some(end) = offset.checked_add(header.cmsg_len) else {
285            break;
286        };
287
288        if header.cmsg_level == libc::SOL_SOCKET {
289            let data_offset = offset + header_len;
290            let declared_data_len = header.cmsg_len - header_len;
291            let available_data_len = control
292                .len()
293                .saturating_sub(data_offset)
294                .min(declared_data_len);
295            let kind = match header.cmsg_type {
296                SCM_TIMESTAMP_OLD | SCM_TIMESTAMP_NEW => Some(SocketTimestampKind::Timeval),
297                SCM_TIMESTAMPNS_OLD | SCM_TIMESTAMPNS_NEW => Some(SocketTimestampKind::Timespec),
298                SCM_TIMESTAMPING_OLD | SCM_TIMESTAMPING_NEW => {
299                    Some(SocketTimestampKind::Timestamping)
300                }
301                _ => None,
302            };
303            if let Some(kind) = kind {
304                messages.push(SocketTimestampMessage {
305                    data_offset,
306                    available_data_len,
307                    kind,
308                });
309            }
310        }
311
312        if end > control.len() {
313            break;
314        }
315
316        let step = cmsg_align(header.cmsg_len);
317        let Some(next) = offset.checked_add(step) else {
318            break;
319        };
320        if next <= offset {
321            break;
322        }
323        offset = next;
324    }
325    messages
326}
327
328// AUTONOMOUS-BOT-IMPLEMENTED
329// TODO-HUMAN-REVIEW(PR-901)
330fn canonicalize_socket_timestamps(control: &mut [u8], now: LogicalTime) -> usize {
331    let messages = socket_timestamp_messages(control);
332    let timespec = libc::timespec {
333        tv_sec: now.as_secs() as libc::time_t,
334        tv_nsec: now.subsec_nanos() as libc::c_long,
335    };
336    let timeval = libc::timeval {
337        tv_sec: timespec.tv_sec,
338        tv_usec: (timespec.tv_nsec / 1_000) as libc::suseconds_t,
339    };
340
341    for message in &messages {
342        let available_end = message
343            .data_offset
344            .saturating_add(message.available_data_len)
345            .min(control.len());
346        let data = &mut control[message.data_offset..available_end];
347        match message.kind {
348            SocketTimestampKind::Timeval => {
349                write_control_prefix(data, timeval);
350            }
351            SocketTimestampKind::Timespec => {
352                write_control_prefix(data, timespec);
353            }
354            SocketTimestampKind::Timestamping => {
355                let size = std::mem::size_of::<libc::timespec>();
356                let zero = libc::timespec {
357                    tv_sec: 0,
358                    tv_nsec: 0,
359                };
360                for slot in 0..3 {
361                    let start = slot * size;
362                    if start >= data.len() {
363                        break;
364                    }
365                    let end = (start + size).min(data.len());
366                    let slot_bytes = &mut data[start..end];
367                    if slot_bytes.len() != size {
368                        slot_bytes.fill(0);
369                        continue;
370                    }
371                    let original = read_control_value::<libc::timespec>(slot_bytes)
372                        .expect("complete timestamping slot");
373                    let replacement = if original.tv_sec == 0 && original.tv_nsec == 0 {
374                        zero
375                    } else {
376                        timespec
377                    };
378                    let _ = write_control_value(slot_bytes, replacement);
379                }
380            }
381        }
382    }
383    messages.len()
384}
385
386fn sanitize_ppoll_signal_mask(mask: u64) -> u64 {
387    let signal_bit = (reverie::PERF_EVENT_SIGNAL as usize) - 1;
388    mask & !(1_u64 << signal_bit)
389}
390
391fn ppoll_uses_kernel_wait(
392    sequentialize_threads: bool,
393    recordreplay_modes: bool,
394    has_signal_mask: bool,
395) -> bool {
396    !sequentialize_threads || (recordreplay_modes && !has_signal_mask)
397}
398
399impl<T: RecordOrReplay> Detcore<T> {
400    /// poll syscall (MAYHANG)
401    // TODO-HUMAN-REVIEW(PR-1023): Review zero-timeout poll scheduling across backends.
402    pub async fn handle_poll<G: Guest<Self>>(
403        &self,
404        guest: &mut G,
405
406        call: syscalls::Poll,
407    ) -> Result<i64, Error> {
408        if self.cfg.sequentialize_threads && call.timeout() == 0 {
409            // This cannot block, but still yield a scheduler turn so a polling thread cannot
410            // monopolize the guest between preemptions.
411            let yield_to_peer =
412                self.cfg.discover_live_file_metadata && guest.thread_state().has_loopback_peer();
413            resource_request(
414                guest,
415                zero_timeout_poll_request(guest.thread_state().dettid, yield_to_peer),
416            )
417            .await;
418            if self.cfg.recordreplay_modes {
419                Ok(self.record_or_replay(guest, call).await?)
420            } else {
421                Ok(guest.inject(call).await?)
422            }
423        } else if !self.cfg.sequentialize_threads || self.cfg.recordreplay_modes {
424            // In replay mode, we cannot assume the existence of FILES during replay.
425            // Thus we must record the poll and replay it from the trace.
426            Ok(self.handle_external_poll(guest, call).await?)
427        } else {
428            // TODO:
429            // if is-external-poll { self.handle_external_poll(guest, call) }
430            self.handle_internal_poll(guest, call).await
431        }
432    }
433
434    /// Record or replay a raw `select` or `pselect6` the way `poll` is: a zero
435    /// timeout cannot block and takes an ordinary scheduler turn; any other call
436    /// may wait in the kernel, so it runs as a blocking external operation. The
437    /// recorder captures the descriptor sets and remaining time, so replay never
438    /// asks the kernel about descriptors it did not recreate.
439    async fn record_or_replay_select_family<G: Guest<Self>>(
440        &self,
441        guest: &mut G,
442        call: Syscall,
443        zero_timeout: bool,
444    ) -> Result<i64, Error> {
445        if zero_timeout {
446            resource_request(guest, Resources::new(guest.thread_state().dettid)).await;
447            Ok(self.record_or_replay(guest, call).await?)
448        } else {
449            self.record_or_replay_blocking(guest, call).await
450        }
451    }
452
453    /// pselect6 syscall (MAYHANG).
454    // AUTONOMOUS-BOT-IMPLEMENTED
455    // TODO-HUMAN-REVIEW(#686): Review scratch fd sets and scheduler polling.
456    pub async fn handle_pselect6<G: Guest<Self>>(
457        &self,
458        guest: &mut G,
459        call: syscalls::Pselect6,
460    ) -> Result<i64, Error> {
461        if self.cfg.recordreplay_modes {
462            let zero_timeout = call.timeout().is_some_and(|timeout| {
463                guest
464                    .memory()
465                    .read_value(timeout)
466                    .is_ok_and(|timeout: Timespec| timeout.tv_sec == 0 && timeout.tv_nsec == 0)
467            });
468            return self
469                .record_or_replay_select_family(guest, Syscall::Pselect6(call), zero_timeout)
470                .await;
471        }
472        if !self.cfg.sequentialize_threads {
473            return Ok(guest.inject(call).await?);
474        }
475
476        if call.nfds() < 0 {
477            return Ok(guest.inject(call).await?);
478        }
479
480        // Linux copies pselect6's outer { sigmask, sigsetsize } wrapper before
481        // validating the timeout. Copy only the wrapper here; validation of the
482        // pointed-to signal mask remains below, after timeout validation.
483        let sigmask_argument = match call.sigmask() {
484            Some(argument) => {
485                // A split read can fall back to PTRACE_PEEKDATA for the final
486                // word, bypassing PROT_NONE or reporting EIO for an unmapped
487                // page. Have Linux validate both wrapper words first. It copies
488                // this wrapper before rejecting a malformed timeout, without
489                // reading the inner mask, changing it, waiting, or writing output.
490                let mut stack = guest.stack().await;
491                let validation_timeout = stack.reserve::<Timespec>();
492                let _guard = stack.commit()?;
493                guest.memory().write_value(
494                    validation_timeout,
495                    &Timespec {
496                        tv_sec: 0,
497                        tv_nsec: 1_000_000_000,
498                    },
499                )?;
500                let validation = syscalls::Pselect6::new()
501                    .with_nfds(0)
502                    .with_readfds(None)
503                    .with_writefds(None)
504                    .with_exceptfds(None)
505                    .with_timeout(Some(validation_timeout))
506                    .with_sigmask(Some(argument));
507                match guest.inject(validation).await {
508                    Err(Errno::EINVAL) => {}
509                    Err(errno) => return Err(errno.into()),
510                    // Success would mean the backend did not validate the probe.
511                    Ok(_) => return Err(Errno::EIO.into()),
512                }
513                let argument: Pselect6SigmaskArg = guest.memory().read_value(argument.cast())?;
514                Some(argument)
515            }
516            None => None,
517        };
518        let raw_timeout = match call.timeout() {
519            Some(timeout) => {
520                let timeout: Timespec = guest.memory().read_value(timeout)?;
521                Some(timeout)
522            }
523            None => None,
524        };
525        let timeout = raw_timeout.map(ppoll_timeout_duration).transpose()?;
526        if timeout == Some(Duration::ZERO) {
527            return Ok(guest.inject(call).await?);
528        }
529
530        // Linux clamps raw fd-set copies to the process fd table's current max_fds.
531        // Its initial table holds one machine word; larger nfds values can therefore
532        // require fewer bytes than a userspace calculation predicts. Keep those calls
533        // under kernel ownership rather than over-reading the guest bitmap.
534        if call.nfds() > PSELECT6_INTERNAL_MAX_NFDS {
535            return self
536                .record_or_replay_blocking(guest, Syscall::Pselect6(call))
537                .await;
538        }
539
540        // Linux wraps pselect6's temporary mask in { pointer, size }. Glibc supplies
541        // the wrapper even when the inner pointer is null. A real mask must stay in
542        // effect for the whole wait so an unblocked signal (make's jobserver unblocks
543        // SIGCHLD) can interrupt it. Previously that forced the external-blocking path,
544        // whose completion timing is host-decided and is a source of `make -jN`
545        // execution-log divergence. With SIGCHLD admission now deterministic (scheduler
546        // `sigchld_deferred`/`sigchld_ready`), honor the mask on each deterministic poll
547        // probe instead: a pending unblocked signal is observed at a scheduler-decided
548        // probe point rather than at host signal-arrival time.
549        let sigmask = if let Some(argument) = sigmask_argument {
550            if argument.sigmask != 0 {
551                if argument.sigsetsize != KERNEL_SIGSET_SIZE {
552                    return Err(Errno::EINVAL.into());
553                }
554                let mask_addr =
555                    Addr::<libc::sigset_t>::from_raw(argument.sigmask).ok_or(Errno::EFAULT)?;
556                let mask = read_kernel_sigset(guest, mask_addr).await?;
557                Some(sanitize_ppoll_signal_mask(mask))
558            } else {
559                None
560            }
561        } else {
562            None
563        };
564        // The inner mask was snapshotted above. Do not let later guest mutations of the
565        // outer wrapper change the meaning of a retry probe.
566        let call = call.with_sigmask(None);
567
568        self.handle_internal_pselect6(guest, call, timeout, sigmask)
569            .await
570    }
571
572    async fn handle_internal_pselect6<G: Guest<Self>>(
573        &self,
574        guest: &mut G,
575        call: syscalls::Pselect6,
576        timeout: Option<Duration>,
577        sigmask: Option<u64>,
578    ) -> Result<i64, Error> {
579        let len = pselect6_fd_set_len(call.nfds())?;
580        let deadline = match timeout {
581            Some(timeout) => Some(thread_observe_time(guest).await + timeout),
582            None => None,
583        };
584        let original_readfds = match read_pselect6_fd_set(guest, call.readfds(), len) {
585            Ok(value) => value,
586            Err(error) => {
587                self.write_pselect6_remaining(guest, call, deadline).await?;
588                return Err(error);
589            }
590        };
591        let original_writefds = match read_pselect6_fd_set(guest, call.writefds(), len) {
592            Ok(value) => value,
593            Err(error) => {
594                self.write_pselect6_remaining(guest, call, deadline).await?;
595                return Err(error);
596            }
597        };
598        let original_exceptfds = match read_pselect6_fd_set(guest, call.exceptfds(), len) {
599            Ok(value) => value,
600            Err(error) => {
601                self.write_pselect6_remaining(guest, call, deadline).await?;
602                return Err(error);
603            }
604        };
605
606        let mut stack = guest.stack().await;
607        let readfds = call.readfds().map(|_| stack.reserve::<libc::fd_set>());
608        let writefds = call.writefds().map(|_| stack.reserve::<libc::fd_set>());
609        let exceptfds = call.exceptfds().map(|_| stack.reserve::<libc::fd_set>());
610        // pselect6's timeout is a writable in-out kernel timespec, so the probe
611        // needs a mutable scratch cell (re-zeroed each iteration below to keep
612        // every probe a non-blocking poll).
613        let probe_timeout = stack.reserve::<Timespec>();
614        // Carry the temporary signal mask on every zero-timeout probe so the kernel
615        // applies it atomically: a pending, mask-unblocked signal makes the probe return
616        // EINTR at a deterministic scheduler point. The probe's wrapper points at scratch
617        // memory the guard keeps alive across each injection.
618        let probe_sigmask = sigmask.map(|mask| {
619            let sigset = stack.push(mask);
620            stack
621                .push(Pselect6SigmaskArg {
622                    sigmask: sigset.as_raw(),
623                    sigsetsize: KERNEL_SIGSET_SIZE,
624                })
625                .cast()
626        });
627        let _guard = stack.commit()?;
628        let probe = call
629            .with_readfds(readfds)
630            .with_writefds(writefds)
631            .with_exceptfds(exceptfds)
632            .with_timeout(Some(probe_timeout))
633            .with_sigmask(probe_sigmask);
634
635        let mut resources = Resources::new(guest.thread_state().dettid);
636        resources.insert(ResourceID::InternalIOPolling, Permission::W);
637        resources.fyi("pselect6");
638        // Keep the request metadata accurate, but do not make it eligible for
639        // the scheduler's ERESTARTSYS wakeup. A cross-task signal must first be
640        // checked against pselect6's snapshotted temporary mask and disposition.
641        resources.set_signal_interrupt_errno(Errno::EINTR);
642
643        loop {
644            if matches!(
645                resource_request(guest, resources.clone()).await,
646                ResumeStatus::Signaled(_)
647            ) {
648                self.write_pselect6_remaining(guest, call, deadline).await?;
649                return Err(Errno::EINTR.into());
650            }
651            guest.memory().write_value(
652                probe_timeout,
653                &Timespec {
654                    tv_sec: 0,
655                    tv_nsec: 0,
656                },
657            )?;
658            write_pselect6_fd_set(guest, probe.readfds(), &original_readfds)?;
659            write_pselect6_fd_set(guest, probe.writefds(), &original_writefds)?;
660            write_pselect6_fd_set(guest, probe.exceptfds(), &original_exceptfds)?;
661
662            let result = pselect6_probe_result(guest.inject(probe).await);
663            if result != Ok(0) {
664                let copy_result = if result.is_ok() {
665                    self.copy_pselect6_results(guest, probe, call, len)
666                } else {
667                    Ok(())
668                };
669                self.write_pselect6_remaining(guest, call, deadline).await?;
670                copy_result?;
671                return result.map_err(Into::into);
672            }
673
674            resources.poll_attempt += 1;
675            if let Some(deadline) = deadline
676                && thread_observe_time(guest).await >= deadline
677            {
678                let copy_result = self.copy_pselect6_results(guest, probe, call, len);
679                self.write_pselect6_remaining(guest, call, Some(deadline))
680                    .await?;
681                copy_result?;
682                return Ok(0);
683            }
684            trace!(
685                "Retry #{} for syscall due to result Ok(0): {}",
686                resources.poll_attempt,
687                probe.display(&guest.memory())
688            );
689            record_retry_event(guest, probe).await;
690        }
691    }
692
693    fn copy_pselect6_results<G: Guest<Self>>(
694        &self,
695        guest: &mut G,
696        probe: syscalls::Pselect6,
697        call: syscalls::Pselect6,
698        len: usize,
699    ) -> Result<(), Error> {
700        copy_pselect6_fd_set(guest, probe.readfds(), call.readfds(), len)?;
701        copy_pselect6_fd_set(guest, probe.writefds(), call.writefds(), len)?;
702        copy_pselect6_fd_set(guest, probe.exceptfds(), call.exceptfds(), len)
703    }
704
705    async fn write_pselect6_remaining<G: Guest<Self>>(
706        &self,
707        guest: &mut G,
708        call: syscalls::Pselect6,
709        deadline: Option<LogicalTime>,
710    ) -> Result<(), Error> {
711        if let (Some(timeout), Some(deadline)) = (call.timeout(), deadline) {
712            let now = thread_observe_time(guest).await;
713            let remaining = deadline.as_nanos().saturating_sub(now.as_nanos());
714            let remaining = Timespec {
715                tv_sec: (remaining / 1_000_000_000) as libc::time_t,
716                tv_nsec: (remaining % 1_000_000_000) as libc::c_long,
717            };
718            // pselect6's timeout is a writable in-out kernel timespec; reverie-syscalls
719            // now types it as `AddrMut<Timespec>`, so the remaining time can be written
720            // back directly without an unsafe pointer cast.
721            if let Err(error) = guest.memory().write_value(timeout, &remaining) {
722                // Linux preserves the pselect6 result when remaining-time copyout faults.
723                trace!(?error, "ignoring pselect6 timeout writeback failure");
724            }
725        }
726        Ok(())
727    }
728
729    /// select syscall (MAYHANG).
730    ///
731    /// `select` is the classic `timeval` sibling of `pselect6` (which is already
732    /// Determinized). It reuses the pselect6 fd-set scratch machinery, but takes
733    /// a `struct timeval` timeout (which Linux updates in place with the time not
734    /// slept) and carries no signal mask.
735    // AUTONOMOUS-BOT-IMPLEMENTED
736    // TODO-HUMAN-REVIEW(#800): Review select determinization mirroring pselect6.
737    pub async fn handle_select<G: Guest<Self>>(
738        &self,
739        guest: &mut G,
740        call: syscalls::Select,
741    ) -> Result<i64, Error> {
742        if self.cfg.recordreplay_modes {
743            let zero_timeout = call.timeout().is_some_and(|timeout| {
744                guest
745                    .memory()
746                    .read_value(timeout)
747                    .is_ok_and(|timeout: libc::timeval| timeout.tv_sec == 0 && timeout.tv_usec == 0)
748            });
749            return self
750                .record_or_replay_select_family(guest, Syscall::Select(call), zero_timeout)
751                .await;
752        }
753        if !self.cfg.sequentialize_threads {
754            return Ok(guest.inject(call).await?);
755        }
756
757        if call.nfds() < 0 {
758            return Ok(guest.inject(call).await?);
759        }
760
761        let raw_timeout = match call.timeout() {
762            Some(timeout) => {
763                let timeout: libc::timeval = guest.memory().read_value(timeout)?;
764                Some(timeout)
765            }
766            None => None,
767        };
768        if matches!(raw_timeout, Some(timeout) if timeout.tv_sec == 0 && timeout.tv_usec == 0) {
769            // A zero timeout is a pure non-blocking poll; the kernel can service it directly.
770            return Ok(guest.inject(call).await?);
771        }
772
773        // Mirror pselect6: keep large fd tables under kernel ownership rather than
774        // over-reading the guest bitmap (Linux clamps raw fd-set copies to max_fds).
775        if call.nfds() > PSELECT6_INTERNAL_MAX_NFDS {
776            return self
777                .record_or_replay_blocking(guest, Syscall::Select(call))
778                .await;
779        }
780
781        let timeout = raw_timeout.map(select_timeout_duration).transpose()?;
782        self.handle_internal_select(guest, call, timeout).await
783    }
784
785    async fn handle_internal_select<G: Guest<Self>>(
786        &self,
787        guest: &mut G,
788        call: syscalls::Select,
789        timeout: Option<Duration>,
790    ) -> Result<i64, Error> {
791        let len = pselect6_fd_set_len(call.nfds())?;
792        let deadline = match timeout {
793            Some(timeout) => Some(thread_observe_time(guest).await + timeout),
794            None => None,
795        };
796        let original_readfds = match read_pselect6_fd_set(guest, call.readfds(), len) {
797            Ok(value) => value,
798            Err(error) => {
799                self.write_select_remaining(guest, call, deadline).await?;
800                return Err(error);
801            }
802        };
803        let original_writefds = match read_pselect6_fd_set(guest, call.writefds(), len) {
804            Ok(value) => value,
805            Err(error) => {
806                self.write_select_remaining(guest, call, deadline).await?;
807                return Err(error);
808            }
809        };
810        let original_exceptfds = match read_pselect6_fd_set(guest, call.exceptfds(), len) {
811            Ok(value) => value,
812            Err(error) => {
813                self.write_select_remaining(guest, call, deadline).await?;
814                return Err(error);
815            }
816        };
817
818        let mut stack = guest.stack().await;
819        let readfds = call.readfds().map(|_| stack.reserve::<libc::fd_set>());
820        let writefds = call.writefds().map(|_| stack.reserve::<libc::fd_set>());
821        let exceptfds = call.exceptfds().map(|_| stack.reserve::<libc::fd_set>());
822        // select modifies its timeout in place, so the probe timeout must be a
823        // writable scratch cell. It is re-zeroed each iteration to keep every
824        // probe a non-blocking poll (a NULL timeout would block indefinitely).
825        let probe_timeout = stack.reserve::<libc::timeval>();
826        let _guard = stack.commit()?;
827        let probe = call
828            .with_readfds(readfds)
829            .with_writefds(writefds)
830            .with_exceptfds(exceptfds)
831            .with_timeout(Some(probe_timeout));
832
833        let mut resources = Resources::new(guest.thread_state().dettid);
834        resources.insert(ResourceID::InternalIOPolling, Permission::W);
835        resources.fyi("select");
836        // EINTR records select's interruption result. The scheduler only wakes
837        // ERESTARTSYS requests: blocked and ignored signals still need a target-
838        // side disposition check before this Signaled path can be used.
839        resources.set_signal_interrupt_errno(Errno::EINTR);
840
841        loop {
842            if matches!(
843                resource_request(guest, resources.clone()).await,
844                ResumeStatus::Signaled(_)
845            ) {
846                self.write_select_remaining(guest, call, deadline).await?;
847                return Err(Errno::EINTR.into());
848            }
849            guest.memory().write_value(
850                probe_timeout,
851                &libc::timeval {
852                    tv_sec: 0,
853                    tv_usec: 0,
854                },
855            )?;
856            write_pselect6_fd_set(guest, probe.readfds(), &original_readfds)?;
857            write_pselect6_fd_set(guest, probe.writefds(), &original_writefds)?;
858            write_pselect6_fd_set(guest, probe.exceptfds(), &original_exceptfds)?;
859
860            let result = guest.inject(probe).await;
861            if result != Ok(0) {
862                let copy_result = if result.is_ok() {
863                    self.copy_select_results(guest, probe, call, len)
864                } else {
865                    Ok(())
866                };
867                self.write_select_remaining(guest, call, deadline).await?;
868                copy_result?;
869                return result.map_err(Into::into);
870            }
871
872            resources.poll_attempt += 1;
873            if let Some(deadline) = deadline
874                && thread_observe_time(guest).await >= deadline
875            {
876                let copy_result = self.copy_select_results(guest, probe, call, len);
877                self.write_select_remaining(guest, call, Some(deadline))
878                    .await?;
879                copy_result?;
880                return Ok(0);
881            }
882            trace!(
883                "Retry #{} for syscall due to result Ok(0): {}",
884                resources.poll_attempt,
885                probe.display(&guest.memory())
886            );
887            record_retry_event(guest, probe).await;
888        }
889    }
890
891    fn copy_select_results<G: Guest<Self>>(
892        &self,
893        guest: &mut G,
894        probe: syscalls::Select,
895        call: syscalls::Select,
896        len: usize,
897    ) -> Result<(), Error> {
898        copy_pselect6_fd_set(guest, probe.readfds(), call.readfds(), len)?;
899        copy_pselect6_fd_set(guest, probe.writefds(), call.writefds(), len)?;
900        copy_pselect6_fd_set(guest, probe.exceptfds(), call.exceptfds(), len)
901    }
902
903    async fn write_select_remaining<G: Guest<Self>>(
904        &self,
905        guest: &mut G,
906        call: syscalls::Select,
907        deadline: Option<LogicalTime>,
908    ) -> Result<(), Error> {
909        if let (Some(timeout), Some(deadline)) = (call.timeout(), deadline) {
910            let now = thread_observe_time(guest).await;
911            let remaining = deadline.as_nanos().saturating_sub(now.as_nanos());
912            let remaining = libc::timeval {
913                tv_sec: (remaining / 1_000_000_000) as libc::time_t,
914                tv_usec: ((remaining % 1_000_000_000) / 1_000) as libc::suseconds_t,
915            };
916            // select's timeout is a writable in-out kernel timeval reporting the
917            // time not slept; derive it from deterministic virtual time.
918            if let Err(error) = guest.memory().write_value(timeout, &remaining) {
919                // Linux preserves the select result when remaining-time copyout faults.
920                trace!(?error, "ignoring select timeout writeback failure");
921            }
922        }
923        Ok(())
924    }
925
926    /// ppoll syscall (MAYHANG)
927    // TODO-HUMAN-REVIEW(PR-273)
928    pub async fn handle_ppoll<G: Guest<Self>>(
929        &self,
930        guest: &mut G,
931        call: syscalls::Ppoll,
932    ) -> Result<i64, Error> {
933        let timeout_address = call.timeout();
934        let timeout = match timeout_address {
935            Some(timeout) => Some(ppoll_timeout_duration(guest.memory().read_value(timeout)?)?),
936            None => None,
937        };
938
939        let result: Result<i64, Error> = if timeout == Some(Duration::ZERO) {
940            if self.cfg.recordreplay_modes {
941                resource_request(guest, Resources::new(guest.thread_state().dettid)).await;
942            }
943            let (probe, _probe_guard) = self.prepare_ppoll_probe(guest, call).await?;
944            let result = if self.cfg.recordreplay_modes {
945                Ok(self.record_or_replay(guest, probe).await?)
946            } else {
947                Ok(guest.inject_with_retry(probe).await?)
948            };
949            // Linux does not write back an initially zero timeout. Besides matching the
950            // kernel, omitting this write matters when the timeout aliases the pollfd array:
951            // the injected probe may have just stored revents in those same bytes.
952            result
953        } else if ppoll_uses_kernel_wait(
954            self.cfg.sequentialize_threads,
955            self.cfg.recordreplay_modes,
956            call.sigmask().is_some(),
957        ) {
958            // The kernel owns the blocking wait when threads are not sequentialized and for
959            // unmasked record/replay calls. A masked sequentialized wait must use the probe
960            // below so a call that would block keeps the strict fail-closed behavior.
961            // Use scratch memory only for the signal mask so raw ppoll can still update the
962            // guest timeout.
963            let mut signal_mask_guard = None;
964            let call = if let Some(signal_mask) = call.sigmask() {
965                if call.sigsetsize() != KERNEL_SIGSET_SIZE {
966                    return Err(Errno::EINVAL.into());
967                }
968                let signal_mask = read_kernel_sigset(guest, signal_mask).await?;
969                let mut stack = guest.stack().await;
970                let signal_mask = stack.push(sanitize_ppoll_signal_mask(signal_mask)).cast();
971                signal_mask_guard = Some(stack.commit()?);
972                call.with_sigmask(Some(signal_mask))
973            } else {
974                call
975            };
976            let result = Ok(self
977                .record_or_replay_blocking(guest, Syscall::Ppoll(call))
978                .await?);
979            drop(signal_mask_guard);
980            result
981        } else {
982            self.handle_internal_ppoll(guest, call, timeout).await
983        };
984
985        result
986    }
987
988    async fn prepare_ppoll_probe<G: Guest<Self>>(
989        &self,
990        guest: &mut G,
991        call: syscalls::Ppoll,
992    ) -> Result<(syscalls::Ppoll, <G::Stack as Stack>::StackGuard), Error> {
993        let signal_mask = match call.sigmask() {
994            Some(signal_mask) => {
995                if call.sigsetsize() != KERNEL_SIGSET_SIZE {
996                    return Err(Errno::EINVAL.into());
997                }
998                let signal_mask = read_kernel_sigset(guest, signal_mask).await?;
999                Some(sanitize_ppoll_signal_mask(signal_mask))
1000            }
1001            None => None,
1002        };
1003
1004        let mut stack = guest.stack().await;
1005        let timeout = stack.push(timespec_from_duration(Duration::ZERO));
1006        // The scratch stack guard outlives the injected syscall, so the pointee is writable.
1007        let timeout = unsafe { timeout.into_mut() };
1008        let mut probe = call.with_timeout(Some(timeout));
1009        if let Some(signal_mask) = signal_mask {
1010            probe = probe.with_sigmask(Some(stack.push(signal_mask).cast()));
1011        }
1012        let guard = stack.commit()?;
1013        Ok((probe, guard))
1014    }
1015
1016    /// Handle a guest-internal `ppoll` using zero-time kernel probes.
1017    async fn handle_internal_ppoll<G: Guest<Self>>(
1018        &self,
1019        guest: &mut G,
1020        call: syscalls::Ppoll,
1021        timeout: Option<Duration>,
1022    ) -> Result<i64, Error> {
1023        debug_assert_ne!(timeout, Some(Duration::ZERO));
1024        let timeout_address = call.timeout();
1025        let started_at = if timeout.is_some() {
1026            Some(thread_observe_time(guest).await)
1027        } else {
1028            None
1029        };
1030        let deadline = match (timeout, started_at) {
1031            (Some(duration), Some(started_at)) => Some(started_at + duration),
1032            (None, None) => None,
1033            _ => unreachable!(),
1034        };
1035
1036        // A zero probe can honor a temporary signal mask atomically. Keeping that mask
1037        // active while parked would require scheduler-level pending-signal state, so fail
1038        // closed rather than letting a masked signal interrupt a simulated wait.
1039        if call.sigmask().is_some() {
1040            let (probe, _probe_guard) = self.prepare_ppoll_probe(guest, call).await?;
1041            let result = if self.cfg.recordreplay_modes {
1042                self.record_or_replay(guest, probe).await
1043            } else {
1044                guest.inject_with_retry(probe).await
1045            };
1046            if probe.syscall_would_have_blocked(result) {
1047                return Err(Errno::ENOSYS.into());
1048            }
1049            let result = result.map_err(Into::into);
1050            if let (Some(timeout_address), Some(timeout), Some(started_at)) =
1051                (timeout_address, timeout, started_at)
1052            {
1053                self.write_ppoll_remaining(guest, timeout_address, timeout, started_at)
1054                    .await?;
1055            }
1056            return result;
1057        }
1058
1059        let mut rsrc = Resources::new(guest.thread_state().dettid);
1060        rsrc.insert(ResourceID::InternalIOPolling, Permission::W);
1061        rsrc.fyi("ppoll");
1062        let result = retry_nonblocking_syscall_with_timeout(guest, call, rsrc, deadline).await;
1063        if let (Some(timeout_address), Some(timeout), Some(started_at)) =
1064            (timeout_address, timeout, started_at)
1065        {
1066            self.write_ppoll_remaining(guest, timeout_address, timeout, started_at)
1067                .await?;
1068        }
1069        result
1070    }
1071
1072    async fn write_ppoll_remaining<G: Guest<Self>>(
1073        &self,
1074        guest: &mut G,
1075        timeout_address: AddrMut<'_, Timespec>,
1076        timeout: Duration,
1077        started_at: LogicalTime,
1078    ) -> Result<(), Error> {
1079        let now = thread_observe_time(guest).await;
1080        let elapsed = Duration::from_nanos(now.as_nanos().saturating_sub(started_at.as_nanos()));
1081        let remaining = timeout.saturating_sub(elapsed);
1082        if let Err(error) = guest
1083            .memory()
1084            .write_value(timeout_address, &timespec_from_duration(remaining))
1085        {
1086            // Linux preserves the ppoll result when remaining-time copyout faults.
1087            trace!(?error, "ignoring ppoll timeout writeback failure");
1088        }
1089        Ok(())
1090    }
1091
1092    /// Handle a guest-internal poll call that can be fully determinized.
1093    // TODO-HUMAN-REVIEW(PR-1052): Review scheduler fairness for zero-timeout poll.
1094    pub async fn handle_internal_poll<G: Guest<Self>>(
1095        &self,
1096        guest: &mut G,
1097        call: syscalls::Poll,
1098    ) -> Result<i64, Error> {
1099        let timeout_millis = call.timeout();
1100        if timeout_millis == 0 {
1101            // A nonblocking poll can still be the synchronization point in a
1102            // userspace polling loop. Yield once before probing so backends
1103            // without PMU preemption cannot let that loop starve its producer.
1104            let yield_to_peer =
1105                self.cfg.discover_live_file_metadata && guest.thread_state().has_loopback_peer();
1106            resource_request(
1107                guest,
1108                zero_timeout_poll_request(guest.thread_state().dettid, yield_to_peer),
1109            )
1110            .await;
1111            Ok(guest.inject(call).await?) // Already non-blocking.
1112        } else {
1113            let maybe_timeout_ns = millis_duration_to_absolute_timeout(guest, timeout_millis).await;
1114            let mut rsrc = Resources::new(guest.thread_state().dettid);
1115            rsrc.insert(ResourceID::InternalIOPolling, Permission::W);
1116            rsrc.fyi("poll");
1117            retry_nonblocking_syscall_with_timeout(guest, call, rsrc, maybe_timeout_ns).await
1118        }
1119    }
1120
1121    /// Handle a poll syscall that deponds on external, nondeterminstic IO.
1122    pub async fn handle_external_poll<G: Guest<Self>>(
1123        &self,
1124        guest: &mut G,
1125        call: syscalls::Poll,
1126    ) -> Result<i64, Error> {
1127        let len = call.nfds();
1128        let time_delta = Duration::from_millis(call.timeout() as u64);
1129
1130        if len == 0 && time_delta.is_zero() {
1131            let request = Self::sleep_request(guest, time_delta).await;
1132            resource_request(guest, request).await;
1133            Ok(0)
1134        } else {
1135            print_poll(&call);
1136            Ok(self
1137                .record_or_replay_blocking(guest, Syscall::Poll(call))
1138                .await?)
1139        }
1140    }
1141
1142    /// epoll_create1 syscall
1143    pub async fn handle_epoll_create1<G: Guest<Self>>(
1144        &self,
1145        guest: &mut G,
1146        call: syscalls::EpollCreate1,
1147    ) -> Result<i64, Error> {
1148        let dettid = guest.thread_state().dettid;
1149        resource_request(guest, Resources::new(dettid)).await; // empty request
1150        let fd = self.record_or_replay(guest, call).await? as RawFd;
1151        // Register the epoll fd in the DetFd table like every other
1152        // fd-creating syscall (openat, eventfd2, pipe2, socket, ...). Without
1153        // this, later operations that consult the table via `with_detfd` /
1154        // `dup_fd` (F_GETFL, F_SETFD, F_DUPFD[_CLOEXEC], dup, ...) would fail
1155        // with EBADF even though the underlying kernel fd is valid. This broke,
1156        // for example, running rustup proxies (cargo/rustc) under hermit, whose
1157        // tokio runtime dups its epoll fd at startup.
1158        //
1159        // EPOLL_CLOEXEC shares the same bit value as O_CLOEXEC, so we can carry
1160        // the cloexec flag straight across.
1161        self.add_fd(
1162            guest,
1163            fd,
1164            OFlag::from_bits_truncate(call.flags().bits()),
1165            FdType::Epoll,
1166        )
1167        .await?;
1168        Ok(fd as i64)
1169    }
1170
1171    /// Apply an ADD, MOD, or DEL mutation to an epoll interest list.
1172    ///
1173    /// Determinism: strict execution serializes this mutation, so its result depends only on the
1174    /// operation, event payload, and deterministically reconstructed epoll/file-descriptor state.
1175    /// Record/replay rebuilds that state by reinjecting the same control operations.
1176    pub async fn handle_epoll_ctl<G: Guest<Self>>(
1177        &self,
1178        guest: &mut G,
1179        call: syscalls::EpollCtl,
1180    ) -> Result<i64, Error> {
1181        let dettid = guest.thread_state().dettid;
1182        resource_request(guest, Resources::new(dettid)).await; // empty request
1183        Ok(self.record_or_replay(guest, call).await?)
1184    }
1185
1186    /// epoll_pwait syscall (MAYHANG)
1187    pub async fn handle_epoll_pwait<G: Guest<Self>>(
1188        &self,
1189        guest: &mut G,
1190        call: syscalls::EpollPwait,
1191    ) -> Result<i64, Error> {
1192        // This used to unconditionally inject the raw
1193        // call and wait for it to return. With an infinite timeout under
1194        // `--sequentialize-threads` that DEADLOCKS the whole guest: the calling
1195        // task holds the scheduler turn while blocked in the kernel, and the
1196        // only task that could ever satisfy the wait is sitting in the run queue
1197        // waiting for a turn that never comes. Observed as `cmake` configure
1198        // hanging forever at zero CPU with a grandchild frozen mid-`openat`;
1199        // hermit's own scheduler log ends at `COMMIT turn N, dettid <parent>`
1200        // injecting `epoll_pwait(..., -1, NULL, 8)` with `queue len 2`.
1201        //
1202        // WHO ACTUALLY REACHES THIS, measured rather than assumed. It is NOT
1203        // glibc's `epoll_wait(2)` on this architecture: glibc calls
1204        // `SYS_epoll_pwait` only where `__NR_epoll_wait` does not exist (arm64
1205        // and friends). x86_64 has it, and `strace` on glibc 2.34/x86_64 shows
1206        // a plain `epoll_wait` syscall, which `handle_epoll_wait` has always
1207        // handled correctly. The callers that land here are programs issuing
1208        // `epoll_pwait` DIRECTLY -- libuv does, which is how the original
1209        // `cmake` hang was found. With a NULL sigmask the two calls are
1210        // semantically identical, so route them together.
1211        //
1212        // A NON-NULL sigmask keeps the previous behavior: its whole purpose is
1213        // to swap the signal mask atomically for the duration of the wait, and
1214        // a timeout-0 polling loop cannot reproduce that atomicity. Such calls
1215        // remain able to block the scheduler; that is a known remaining gap
1216        // rather than something this change silently pretends to fix.
1217        if call.sigmask().is_some() {
1218            let dettid = guest.thread_state().dettid;
1219            resource_request(guest, Resources::new(dettid)).await; // empty request
1220            return Ok(self.record_or_replay(guest, call).await?);
1221        }
1222        if self.cfg.recordreplay_modes && call.timeout() == 0 {
1223            // Cannot block, but still yield a scheduler turn so a polling thread
1224            // cannot monopolize the guest between preemptions.
1225            resource_request(guest, Resources::new(guest.thread_state().dettid)).await;
1226            Ok(self.record_or_replay(guest, call).await?)
1227        } else if !self.cfg.sequentialize_threads || self.cfg.recordreplay_modes {
1228            Ok(self
1229                .record_or_replay_blocking(guest, Syscall::EpollPwait(call))
1230                .await?)
1231        } else {
1232            self.handle_internal_epoll_pwait(guest, call).await
1233        }
1234    }
1235
1236    /// Handle a guest-internal `epoll_pwait` (NULL sigmask) that can be fully
1237    /// determinized. Mirrors `handle_internal_epoll_wait`.
1238    pub async fn handle_internal_epoll_pwait<G: Guest<Self>>(
1239        &self,
1240        guest: &mut G,
1241        call: syscalls::EpollPwait,
1242    ) -> Result<i64, Error> {
1243        let timeout_millis = call.timeout();
1244        if timeout_millis == 0 {
1245            // Cannot block, but must still yield a scheduler turn: a
1246            // zero-timeout polling loop that never requests a resource can
1247            // monopolize the guest between preemptions and starve the producer
1248            // it is polling for. `handle_poll` takes a turn for every
1249            // sequential mode, and the record/replay arm of `handle_epoll_pwait`
1250            // does the same; before this PR routed NULL-sigmask `epoll_pwait`
1251            // here, the old handler always made an empty request. Omitting it
1252            // only on the plain-strict path would be a scheduling regression,
1253            // not a refactor.
1254            if self.cfg.sequentialize_threads {
1255                let yield_to_peer = self.cfg.discover_live_file_metadata
1256                    && guest.thread_state().has_loopback_peer();
1257                resource_request(
1258                    guest,
1259                    zero_timeout_poll_request(guest.thread_state().dettid, yield_to_peer),
1260                )
1261                .await;
1262            }
1263            Ok(guest.inject(call).await?) // Already non-blocking.
1264        } else {
1265            let maybe_timeout_ns = millis_duration_to_absolute_timeout(guest, timeout_millis).await;
1266            let mut rsrc = Resources::new(guest.thread_state().dettid);
1267            rsrc.insert(ResourceID::InternalIOPolling, Permission::W);
1268            rsrc.fyi("epoll_pwait");
1269            retry_nonblocking_syscall_with_timeout(guest, call, rsrc, maybe_timeout_ns).await
1270        }
1271    }
1272
1273    /// epoll_pwait2 syscall (MAYHANG).
1274    ///
1275    /// epoll_pwait2 is epoll_pwait with a `struct timespec *` timeout instead of
1276    /// an int-milliseconds timeout; recent glibc implements epoll_wait/
1277    /// epoll_pwait via epoll_pwait2 when the kernel supports it. The pinned
1278    /// Reverie revision has no typed variant, so it arrives as a raw
1279    /// `Syscall::Other` and is dispatched here by Sysno. Detcore treats it
1280    /// exactly like epoll_pwait: a scheduler yield point followed by
1281    /// record/replay-aware forwarding of the raw call.
1282    // AUTONOMOUS-BOT-IMPLEMENTED
1283    // TODO-HUMAN-REVIEW(#773)
1284    pub async fn handle_epoll_pwait2<G: Guest<Self>>(
1285        &self,
1286        guest: &mut G,
1287        call: Syscall,
1288    ) -> Result<i64, Error> {
1289        let dettid = guest.thread_state().dettid;
1290        resource_request(guest, Resources::new(dettid)).await; // empty request
1291        Ok(self.record_or_replay(guest, call).await?)
1292    }
1293
1294    /// epoll_wait syscall (MAYHANG)
1295    pub async fn handle_epoll_wait<G: Guest<Self>>(
1296        &self,
1297        guest: &mut G,
1298        call: syscalls::EpollWait,
1299    ) -> Result<i64, Error> {
1300        if self.cfg.recordreplay_modes && call.timeout() == 0 {
1301            // This cannot block, but still yield a scheduler turn so a polling thread cannot
1302            // monopolize the guest between preemptions.
1303            resource_request(guest, Resources::new(guest.thread_state().dettid)).await;
1304            Ok(self.record_or_replay(guest, call).await?)
1305        } else if !self.cfg.sequentialize_threads || self.cfg.recordreplay_modes {
1306            Ok(self
1307                .record_or_replay_blocking(guest, Syscall::EpollWait(call))
1308                .await?)
1309        } else {
1310            self.handle_internal_epoll_wait(guest, call).await
1311        }
1312    }
1313
1314    /// Handle a guest-internal `epoll_wait` call that can be fully determinized.
1315    pub async fn handle_internal_epoll_wait<G: Guest<Self>>(
1316        &self,
1317        guest: &mut G,
1318        call: syscalls::EpollWait,
1319    ) -> Result<i64, Error> {
1320        let timeout_millis = call.timeout();
1321        if timeout_millis == 0 {
1322            Ok(guest.inject(call).await?) // Already non-blocking.
1323        } else {
1324            let maybe_timeout_ns = millis_duration_to_absolute_timeout(guest, timeout_millis).await;
1325            let mut rsrc = Resources::new(guest.thread_state().dettid);
1326            rsrc.insert(ResourceID::InternalIOPolling, Permission::W);
1327            rsrc.fyi("epoll_wait");
1328            retry_nonblocking_syscall_with_timeout(guest, call, rsrc, maybe_timeout_ns).await
1329        }
1330    }
1331
1332    /// Connect system call (MAYHANG)
1333    /// Note that connect waits until a TCP handshake but does not wait for accept() on the other end.
1334    /// Nevertheless, it can block for a long time while waiting for connection, unless the socket
1335    /// is already nonblocking.
1336    pub async fn handle_connect<G: Guest<Self>>(
1337        &self,
1338        guest: &mut G,
1339        call: syscalls::Connect,
1340    ) -> Result<i64, Error> {
1341        let fd = call.fd();
1342        let uservaddr = call.uservaddr();
1343        let addrlen = call.addrlen();
1344
1345        if guest.config().sched_heuristic == SchedHeuristic::ConnectBind {
1346            trace!("Scheduling heuristic: reprioritizing connect");
1347            let resource = ResourceID::PriorityChangePoint(
1348                FIRST_PRIORITY,
1349                guest.thread_state().thread_logical_time.as_nanos(),
1350                guest.thread_state().committed_clock_value,
1351                Vec::new(),
1352            );
1353            let req = guest.thread_state().mk_request(resource, Permission::W);
1354            resource_request(guest, req).await;
1355        }
1356
1357        let result = self.execute_nonblockable_fd_syscall(guest, call).await;
1358        if self.cfg.discover_live_file_metadata
1359            && connect_result_allows_peer_classification(&result)
1360        {
1361            // This metadata is a SaBRe-only scheduling hint, not part of connect's semantics.
1362            // Let the kernel establish the authoritative result first, then classify the peer
1363            // best-effort so an invalid guest pointer or an untracked fd can never replace the
1364            // kernel's errno.
1365            let loopback_peer = (|| -> Result<Option<bool>, Error> {
1366                let Some(address) = uservaddr else {
1367                    return Ok(None);
1368                };
1369                let addrlen = usize::try_from(addrlen).unwrap_or(0);
1370                if addrlen < std::mem::size_of::<u16>() {
1371                    return Ok(None);
1372                }
1373                let family: u16 = guest.memory().read_value(address.cast())?;
1374                if family == libc::AF_INET as u16 {
1375                    if addrlen < std::mem::size_of::<libc::sockaddr_in>() {
1376                        return Ok(None);
1377                    }
1378                    let address: libc::sockaddr_in = guest.memory().read_value(address.cast())?;
1379                    Ok(Some(
1380                        Ipv4Addr::from(address.sin_addr.s_addr.to_ne_bytes()).is_loopback(),
1381                    ))
1382                } else if family == libc::AF_INET6 as u16 {
1383                    if addrlen < std::mem::size_of::<libc::sockaddr_in6>() {
1384                        return Ok(None);
1385                    }
1386                    let address: libc::sockaddr_in6 = guest.memory().read_value(address.cast())?;
1387                    Ok(Some(
1388                        Ipv6Addr::from(address.sin6_addr.s6_addr).is_loopback(),
1389                    ))
1390                } else {
1391                    // This includes a successful AF_UNSPEC disconnect and successful connects
1392                    // to non-IP families, neither of which has a loopback IP peer.
1393                    Ok(Some(false))
1394                }
1395            })()
1396            .ok()
1397            .flatten();
1398            if let Some(loopback_peer) = loopback_peer {
1399                let _ = guest
1400                    .thread_state()
1401                    .with_detfd(fd, |detfd| detfd.set_loopback_peer(loopback_peer));
1402            }
1403        }
1404
1405        result
1406    }
1407
1408    /// Handles sendto, sendmsg, and sendmmsg syscalls (MAYHANG).
1409    pub async fn handle_sendrecv<
1410        G: Guest<Self>,
1411        C: SyscallInfo + NonblockableSyscall + Into<Syscall>,
1412    >(
1413        &self,
1414        guest: &mut G,
1415        call: C,
1416    ) -> Result<i64, Error> {
1417        self.execute_nonblockable_fd_syscall(guest, call).await
1418    }
1419
1420    /// Sends one message and invalidates process-wide flock knowledge after success.
1421    ///
1422    /// The guest can mutate shared message and control memory while this helper
1423    /// deschedules. Parsing before the syscall would not prove which descriptors the
1424    /// kernel later transferred, so a successful unbound send conservatively makes
1425    /// every cached open-file-description lock mode unknown.
1426    pub async fn handle_sendmsg<G: Guest<Self>>(
1427        &self,
1428        guest: &mut G,
1429        call: syscalls::Sendmsg,
1430    ) -> Result<i64, Error> {
1431        let result = self.execute_nonblockable_fd_syscall(guest, call).await?;
1432        guest.thread_state().forget_flock_modes();
1433        Ok(result)
1434    }
1435
1436    /// Sends a message batch and invalidates process-wide flock knowledge when the
1437    /// kernel reports at least one message sent. This intentionally includes
1438    /// descriptors named only by an unsent tail message: the mutable guest array is
1439    /// not stable across a possible deschedule, so narrower attribution is unsafe.
1440    pub async fn handle_sendmmsg<G: Guest<Self>>(
1441        &self,
1442        guest: &mut G,
1443        call: syscalls::Sendmmsg,
1444    ) -> Result<i64, Error> {
1445        let result = self.execute_nonblockable_fd_syscall(guest, call).await?;
1446        if result > 0 {
1447            guest.thread_state().forget_flock_modes();
1448        }
1449        Ok(result)
1450    }
1451
1452    // TODO-HUMAN-REVIEW(PR-912): Review receive-time capture across socket aliases.
1453    async fn observe_socket_receive<G: Guest<Self>>(
1454        &self,
1455        guest: &mut G,
1456        fd: i32,
1457    ) -> Result<LogicalTime, Error> {
1458        let timestamp = thread_observe_time(guest).await;
1459        guest.thread_state().with_detfd(fd, |detfd| {
1460            detfd.set_socket_receive_timestamp(timestamp);
1461        })?;
1462        Ok(timestamp)
1463    }
1464
1465    // AUTONOMOUS-BOT-IMPLEMENTED
1466    // TODO-HUMAN-REVIEW(PR-901)
1467    /// Receive one message and replace host socket timestamps with logical time.
1468    pub async fn handle_recvmsg<G: Guest<Self>>(
1469        &self,
1470        guest: &mut G,
1471        call: syscalls::Recvmsg,
1472    ) -> Result<i64, Error> {
1473        // NETLINK_SOCK_DIAG replies carry host-assigned socket identities in
1474        // their msg_iov payload (not msg_control). Canonicalize the supported
1475        // fields before the binary reply reaches the guest.
1476        // The predicate is shared with the read/readv/recvfrom/recvmmsg paths so
1477        // the five receive syscalls cannot drift apart again.
1478        if self.sock_diag_reply_fd(guest, call.sockfd()) {
1479            return self.handle_sock_diag_recvmsg(guest, call).await;
1480        }
1481
1482        if !self.cfg.virtualize_time {
1483            return self.execute_nonblockable_fd_syscall(guest, call).await;
1484        }
1485
1486        let Some(message_address) = call.msg() else {
1487            return self
1488                .handle_socket_receive(guest, call, call.sockfd(), true)
1489                .await;
1490        };
1491        // Snapshot every input field before the receive. Linux permits the
1492        // control buffer to overlap this header, so rereading it afterward can
1493        // turn a successful consuming receive into an artificial EFAULT.
1494        let message: libc::msghdr = guest.memory().read_value(message_address)?;
1495        if message.msg_control.is_null() || message.msg_controllen == 0 {
1496            return self
1497                .handle_socket_receive(guest, call, call.sockfd(), true)
1498                .await;
1499        }
1500        let control_len = message.msg_controllen.min(MAX_CONTROL_BYTES);
1501        let control_address: AddrMut<'_, u8> =
1502            AddrMut::from_raw(message.msg_control as usize).ok_or(Errno::EFAULT)?;
1503        let mut control = vec![0; control_len];
1504        // Validate the output region before consuming a datagram.
1505        guest.memory().read_exact(control_address, &mut control)?;
1506
1507        let result = self.execute_nonblockable_fd_syscall(guest, call).await?;
1508        let now = self.observe_socket_receive(guest, call.sockfd()).await?;
1509        let mut control = vec![0; control_len];
1510        guest.memory().read_exact(control_address, &mut control)?;
1511        if socket_timestamp_messages(&control).is_empty() {
1512            return Ok(result);
1513        }
1514
1515        canonicalize_socket_timestamps(&mut control, now);
1516        guest.memory().write_exact(control_address, &control)?;
1517        Ok(result)
1518    }
1519
1520    // AUTONOMOUS-BOT-IMPLEMENTED
1521    // TODO-HUMAN-REVIEW(PR-1064)
1522    /// Whether replies received on `fd` must have their supported socket
1523    /// identities canonicalized.
1524    ///
1525    /// This is the single predicate every receive path consults. It exists as
1526    /// one function because the flag used to be tested inline in `recvmsg`
1527    /// alone, and four other receive syscalls reached the same dump without it
1528    /// (see [`Self::sanitize_sock_diag_segments`]).
1529    pub(crate) fn sock_diag_reply_fd<G: Guest<Self>>(&self, guest: &mut G, fd: RawFd) -> bool {
1530        self.cfg.virtualize_metadata
1531            && guest
1532                .thread_state()
1533                .with_detfd(fd, |detfd| detfd.is_sock_diag() || detfd.is_netlink_route())
1534                .unwrap_or(false)
1535    }
1536
1537    /// Whether this descriptor is specifically a `NETLINK_ROUTE` socket, which
1538    /// needs the link-counter sanitizer rather than the sock-diag one.
1539    fn netlink_route_reply_fd<G: Guest<Self>>(&self, guest: &mut G, fd: RawFd) -> bool {
1540        guest
1541            .thread_state()
1542            .with_detfd(fd, |detfd| detfd.is_netlink_route())
1543            .unwrap_or(false)
1544    }
1545
1546    /// Read an `iovec` array out of guest memory as plain `(address, capacity)`
1547    /// scalars, so no non-`Send` raw pointer is held across an await.
1548    fn read_iov_segments<G: Guest<Self>>(
1549        guest: &mut G,
1550        iov: usize,
1551        iovlen: usize,
1552    ) -> Result<Vec<(usize, usize)>, Error> {
1553        if iov == 0 || iovlen == 0 {
1554            return Ok(Vec::new());
1555        }
1556        let count = iovlen.min(libc::UIO_MAXIOV as usize);
1557        let address: AddrMut<'_, libc::iovec> = AddrMut::from_raw(iov).ok_or(Errno::EFAULT)?;
1558        // SAFETY: `iovec` is a plain C record; an all-zero value is a valid
1559        // staging value immediately overwritten by `read_values`.
1560        let mut iovecs: Vec<libc::iovec> =
1561            (0..count).map(|_| unsafe { std::mem::zeroed() }).collect();
1562        guest.memory().read_values(address.into(), &mut iovecs)?;
1563        Ok(iovecs
1564            .iter()
1565            .map(|iov| (iov.iov_base as usize, iov.iov_len))
1566            .collect())
1567    }
1568
1569    /// Canonicalize host-assigned identities in a `NETLINK_SOCK_DIAG` reply that
1570    /// the kernel has already written into guest memory.
1571    ///
1572    /// `segments` describes the destination buffers as `(address, capacity)` in
1573    /// the order the kernel filled them; `received` is the syscall's return
1574    /// value. The reply is gathered contiguously (it may be scattered across
1575    /// several buffers), sanitized via `crate::sock_diag` (fail-open,
1576    /// zero-only, never resizes), and written back preserving the original
1577    /// boundaries. Bounded by both each buffer's capacity and `received`, which
1578    /// under `MSG_TRUNC` can exceed the total capacity.
1579    ///
1580    /// Synchronous on purpose: every caller has already completed its receive,
1581    /// so nothing here awaits and no guest address outlives the borrow.
1582    fn sanitize_sock_diag_segments<G: Guest<Self>>(
1583        &self,
1584        guest: &mut G,
1585        fd: RawFd,
1586        segments: &[(usize, usize)],
1587        received: usize,
1588    ) -> Result<(), Error> {
1589        if received == 0 || segments.is_empty() {
1590            return Ok(());
1591        }
1592        let mut filled: Vec<(AddrMut<'_, u8>, usize)> = Vec::new();
1593        let mut buffer: Vec<u8> = Vec::with_capacity(received);
1594        let mut remaining = received;
1595        for &(base, capacity) in segments {
1596            if remaining == 0 {
1597                break;
1598            }
1599            if base == 0 || capacity == 0 {
1600                continue;
1601            }
1602            let take = capacity.min(remaining);
1603            let address: AddrMut<'_, u8> = AddrMut::from_raw(base).ok_or(Errno::EFAULT)?;
1604            let mut segment = vec![0u8; take];
1605            guest.memory().read_exact(address, &mut segment)?;
1606            buffer.extend_from_slice(&segment);
1607            filled.push((address, take));
1608            remaining -= take;
1609        }
1610
1611        // NETLINK_ROUTE and NETLINK_SOCK_DIAG replies need different
1612        // sanitizers: one zeroes live interface counters, the other determinizes
1613        // supported socket identities. The descriptor decides which, so a guest
1614        // holding both kinds of socket gets each handled correctly.
1615        let modified = if self.netlink_route_reply_fd(guest, fd) {
1616            crate::netlink_route::sanitize_route_link_stats(&mut buffer)
1617        } else {
1618            crate::sock_diag::sanitize_sock_diag_identities(&mut buffer)
1619        };
1620        if !modified {
1621            return Ok(());
1622        }
1623
1624        let mut offset = 0;
1625        for (address, len) in filled {
1626            guest
1627                .memory()
1628                .write_exact(address, &buffer[offset..offset + len])?;
1629            offset += len;
1630        }
1631        Ok(())
1632    }
1633
1634    /// `recvfrom` on a socket-diag descriptor: one destination buffer.
1635    ///
1636    /// `recv(2)` has no syscall of its own on x86_64 — glibc lowers it to
1637    /// `recvfrom` with a null address — so this covers `recv` as well, which is
1638    /// what Python's `socket.recv()` reaches.
1639    pub async fn handle_sock_diag_recvfrom<G: Guest<Self>>(
1640        &self,
1641        guest: &mut G,
1642        call: syscalls::Recvfrom,
1643    ) -> Result<i64, Error> {
1644        let fd = call.fd();
1645        let base = call.buf().map(|address| address.as_raw()).unwrap_or(0);
1646        let len = call.len();
1647        let result = self.handle_socket_receive(guest, call, fd, true).await?;
1648        let received = usize::try_from(result).unwrap_or(0);
1649        self.sanitize_sock_diag_segments(guest, fd, &[(base, len)], received)?;
1650        Ok(result)
1651    }
1652
1653    /// `read` on a socket-diag descriptor: one destination buffer.
1654    pub async fn handle_sock_diag_read<G: Guest<Self>>(
1655        &self,
1656        guest: &mut G,
1657        call: syscalls::Read,
1658    ) -> Result<i64, Error> {
1659        let fd = call.fd();
1660        let base = call.buf().map(|address| address.as_raw()).unwrap_or(0);
1661        let len = call.len();
1662        let result = self.handle_read(guest, call).await?;
1663        let received = usize::try_from(result).unwrap_or(0);
1664        self.sanitize_sock_diag_segments(guest, fd, &[(base, len)], received)?;
1665        Ok(result)
1666    }
1667
1668    /// Receive a `NETLINK_SOCK_DIAG` dump and zero the host-assigned socket
1669    /// inode numbers in the reply so `ss`-style enumeration is deterministic.
1670    ///
1671    /// The dump lands in `msg_iov` (netlink diag sockets carry no ancillary
1672    /// data), possibly scattered across several iovecs.
1673    async fn handle_sock_diag_recvmsg<G: Guest<Self>>(
1674        &self,
1675        guest: &mut G,
1676        call: syscalls::Recvmsg,
1677    ) -> Result<i64, Error> {
1678        let fd = call.sockfd();
1679        let Some(message_address) = call.msg() else {
1680            return self.handle_socket_receive(guest, call, fd, true).await;
1681        };
1682        // Snapshot the header's iovec pointer/count before the receive (they are
1683        // stable across it; the kernel fills the pointed-to buffers, not the
1684        // array). Scoped so the `msghdr`'s raw pointers are dropped before the
1685        // await below: holding one would make this future non-`Send`.
1686        let segments = {
1687            let message: libc::msghdr = guest.memory().read_value(message_address)?;
1688            Self::read_iov_segments(guest, message.msg_iov as usize, message.msg_iovlen)?
1689        };
1690
1691        let result = self.handle_socket_receive(guest, call, fd, true).await?;
1692        let received = usize::try_from(result).unwrap_or(0);
1693        self.sanitize_sock_diag_segments(guest, fd, &segments, received)?;
1694        Ok(result)
1695    }
1696
1697    /// `readv` on a socket-diag descriptor: an `iovec` array, no message header.
1698    pub async fn handle_sock_diag_readv<G: Guest<Self>>(
1699        &self,
1700        guest: &mut G,
1701        call: syscalls::Readv,
1702    ) -> Result<i64, Error> {
1703        let fd = call.fd();
1704        let iov = call.iov().map(|address| address.as_raw()).unwrap_or(0);
1705        let segments = Self::read_iov_segments(guest, iov, call.len())?;
1706        let result = self.handle_readv(guest, call).await?;
1707        let received = usize::try_from(result).unwrap_or(0);
1708        self.sanitize_sock_diag_segments(guest, fd, &segments, received)?;
1709        Ok(result)
1710    }
1711
1712    /// `recvmmsg` on a socket-diag descriptor.
1713    ///
1714    /// Each delivered `mmsghdr` is a separate datagram with its own byte count
1715    /// in `msg_len`, so each is gathered and sanitized independently; treating
1716    /// the batch as one buffer would let one message's length run into the
1717    /// next message's memory.
1718    pub async fn handle_sock_diag_recvmmsg<G: Guest<Self>>(
1719        &self,
1720        guest: &mut G,
1721        call: syscalls::Recvmmsg,
1722    ) -> Result<i64, Error> {
1723        let fd = call.fd();
1724        let Some(messages_address) = call.mmsg() else {
1725            return self.handle_recvmmsg(guest, call).await;
1726        };
1727        let vlen = call.vlen();
1728        if vlen == 0 || vlen > libc::UIO_MAXIOV as u32 {
1729            return self.handle_recvmmsg(guest, call).await;
1730        }
1731        let base = messages_address.as_raw();
1732
1733        // Snapshot each message's iovec geometry BEFORE the receive: the kernel
1734        // fills the pointed-to buffers and writes msg_len, but does not move the
1735        // iovec arrays themselves. Scoped, and reduced to plain scalars, so the
1736        // `mmsghdr` raw pointers are dropped before the await below: holding one
1737        // would make this future non-`Send`.
1738        let geometry: Vec<Vec<(usize, usize)>> = {
1739            // SAFETY: `mmsghdr` is a plain C record; an all-zero value is a
1740            // valid staging value immediately overwritten by `read_values`.
1741            let mut headers: Vec<libc::mmsghdr> = (0..vlen as usize)
1742                .map(|_| unsafe { std::mem::zeroed() })
1743                .collect();
1744            guest
1745                .memory()
1746                .read_values(messages_address.into(), &mut headers)?;
1747            let mut geometry = Vec::with_capacity(headers.len());
1748            for header in &headers {
1749                geometry.push(Self::read_iov_segments(
1750                    guest,
1751                    header.msg_hdr.msg_iov as usize,
1752                    header.msg_hdr.msg_iovlen,
1753                )?);
1754            }
1755            geometry
1756        };
1757
1758        let result = self.handle_recvmmsg(guest, call).await?;
1759        let delivered = usize::try_from(result).unwrap_or(0).min(geometry.len());
1760        if delivered == 0 {
1761            return Ok(result);
1762        }
1763
1764        // Re-read the array for the per-message byte counts the kernel just
1765        // wrote. Also scoped: nothing awaits past this point, but keeping the
1766        // raw pointers contained keeps the rule visible.
1767        let counts: Vec<usize> = {
1768            let address: AddrMut<'_, libc::mmsghdr> =
1769                AddrMut::from_raw(base).ok_or(Errno::EFAULT)?;
1770            // SAFETY: as above.
1771            let mut headers: Vec<libc::mmsghdr> = (0..vlen as usize)
1772                .map(|_| unsafe { std::mem::zeroed() })
1773                .collect();
1774            guest.memory().read_values(address.into(), &mut headers)?;
1775            headers.iter().map(|h| h.msg_len as usize).collect()
1776        };
1777        for (index, segments) in geometry.iter().enumerate().take(delivered) {
1778            self.sanitize_sock_diag_segments(guest, fd, segments, counts[index])?;
1779        }
1780        Ok(result)
1781    }
1782
1783    // TODO-HUMAN-REVIEW(PR-912): Review receive-time capture across socket aliases.
1784    /// Handle a socket receive and retain one timestamp for every alias of its open file.
1785    pub async fn handle_socket_receive<
1786        G: Guest<Self>,
1787        C: SyscallInfo + NonblockableSyscall + Into<Syscall>,
1788    >(
1789        &self,
1790        guest: &mut G,
1791        call: C,
1792        fd: i32,
1793        zero_delivers_packet: bool,
1794    ) -> Result<i64, Error> {
1795        let result = self.execute_nonblockable_fd_syscall(guest, call).await?;
1796        if self.cfg.virtualize_time && (result > 0 || (result == 0 && zero_delivers_packet)) {
1797            self.observe_socket_receive(guest, fd).await?;
1798        }
1799        Ok(result)
1800    }
1801
1802    // AUTONOMOUS-BOT-IMPLEMENTED
1803    // TODO-HUMAN-REVIEW(PR-901)
1804    /// Receive a message batch and replace every host socket timestamp with logical time.
1805    pub async fn handle_recvmmsg<G: Guest<Self>>(
1806        &self,
1807        guest: &mut G,
1808        call: syscalls::Recvmmsg,
1809    ) -> Result<i64, Error> {
1810        if !self.cfg.virtualize_time || call.vlen() > libc::UIO_MAXIOV as u32 {
1811            return self.execute_nonblockable_fd_syscall(guest, call).await;
1812        }
1813
1814        let Some(messages_address) = call.mmsg() else {
1815            return self.execute_nonblockable_fd_syscall(guest, call).await;
1816        };
1817        let controls = {
1818            // SAFETY: `mmsghdr` is a plain C record and an all-zero value is a valid
1819            // initialized staging value that is immediately overwritten by `read_values`.
1820            let mut messages: Vec<libc::mmsghdr> = (0..call.vlen())
1821                .map(|_| unsafe { std::mem::zeroed() })
1822                .collect();
1823            guest
1824                .memory()
1825                .read_values(messages_address.into(), &mut messages)?;
1826
1827            let mut controls = Vec::with_capacity(messages.len());
1828            for message in &messages {
1829                let header = &message.msg_hdr;
1830                if header.msg_control.is_null() || header.msg_controllen == 0 {
1831                    controls.push(None);
1832                    continue;
1833                }
1834                controls.push(Some((
1835                    header.msg_control as usize,
1836                    header.msg_controllen.min(MAX_CONTROL_BYTES),
1837                )));
1838            }
1839            controls
1840        };
1841
1842        let result = self.execute_nonblockable_fd_syscall(guest, call).await?;
1843        let delivered = usize::try_from(result).unwrap_or(0).min(controls.len());
1844        if delivered == 0 {
1845            return Ok(result);
1846        }
1847        let now = self.observe_socket_receive(guest, call.fd()).await?;
1848        let mut timestamped = Vec::new();
1849        for (address, length) in controls.into_iter().take(delivered).flatten() {
1850            let address: AddrMut<'_, u8> = AddrMut::from_raw(address).ok_or(Errno::EFAULT)?;
1851            let mut bytes = vec![0; length];
1852            guest.memory().read_exact(address, &mut bytes)?;
1853            if !socket_timestamp_messages(&bytes).is_empty() {
1854                timestamped.push((address, bytes));
1855            }
1856        }
1857        if timestamped.is_empty() {
1858            return Ok(result);
1859        }
1860
1861        for (address, mut bytes) in timestamped {
1862            canonicalize_socket_timestamps(&mut bytes, now);
1863            guest.memory().write_exact(address, &bytes)?;
1864        }
1865        Ok(result)
1866    }
1867}
1868
1869#[cfg(test)]
1870mod tests {
1871    use super::*;
1872
1873    #[test]
1874    fn zero_timeout_socket_poll_requests_a_strong_one_turn_yield() {
1875        let request = zero_timeout_poll_request(DetTid::from_raw(17), true);
1876
1877        assert_eq!(request.resources.len(), 1);
1878        assert_eq!(
1879            request.resources.get(&ResourceID::SchedYield),
1880            Some(&Permission::W)
1881        );
1882        assert_eq!(request.fyi, SABRE_LOOPBACK_POLL_YIELD_FYI);
1883    }
1884
1885    #[test]
1886    fn zero_timeout_non_socket_poll_keeps_the_existing_empty_turn() {
1887        let request = zero_timeout_poll_request(DetTid::from_raw(17), false);
1888
1889        assert!(request.resources.is_empty());
1890        assert!(request.fyi.is_empty());
1891    }
1892
1893    #[test]
1894    fn connect_peer_classification_never_overrides_kernel_errors() {
1895        assert!(connect_result_allows_peer_classification(&Ok(0)));
1896        assert!(connect_result_allows_peer_classification(&Err(
1897            Error::Errno(Errno::EINPROGRESS)
1898        )));
1899        assert!(!connect_result_allows_peer_classification(&Err(
1900            Error::Errno(Errno::EALREADY)
1901        )));
1902        assert!(!connect_result_allows_peer_classification(&Err(
1903            Error::Errno(Errno::EBADF)
1904        )));
1905        assert!(!connect_result_allows_peer_classification(&Err(
1906            Error::Errno(Errno::EFAULT)
1907        )));
1908    }
1909
1910    #[test]
1911    fn ppoll_timeout_uses_timespec_units() {
1912        assert_eq!(
1913            ppoll_timeout_duration(Timespec {
1914                tv_sec: 2,
1915                tv_nsec: 345_678_901,
1916            }),
1917            Ok(Duration::new(2, 345_678_901))
1918        );
1919        assert_eq!(
1920            timespec_from_duration(Duration::new(2, 345_678_901)),
1921            Timespec {
1922                tv_sec: 2,
1923                tv_nsec: 345_678_901,
1924            }
1925        );
1926    }
1927
1928    #[test]
1929    fn pselect6_fd_set_lengths_follow_the_raw_linux_abi() {
1930        assert_eq!(PSELECT6_INTERNAL_MAX_NFDS, 64);
1931        assert_eq!(pselect6_fd_set_len(-1), Err(Errno::EINVAL));
1932        assert_eq!(pselect6_fd_set_len(0), Ok(0));
1933        assert_eq!(pselect6_fd_set_len(1), Ok(8));
1934        assert_eq!(pselect6_fd_set_len(65), Ok(16));
1935        assert_eq!(pselect6_fd_set_len(libc::FD_SETSIZE as i32), Ok(128));
1936    }
1937
1938    #[test]
1939    fn pselect6_probe_result_maps_only_erestartsys_to_eintr() {
1940        assert_eq!(
1941            pselect6_probe_result(Err(Errno::ERESTARTSYS)),
1942            Err(Errno::EINTR)
1943        );
1944
1945        assert_eq!(pselect6_probe_result(Ok(0)), Ok(0));
1946        assert_eq!(pselect6_probe_result(Ok(1)), Ok(1));
1947        assert_eq!(pselect6_probe_result(Err(Errno::EBADF)), Err(Errno::EBADF));
1948        assert_eq!(
1949            pselect6_probe_result(Err(Errno::EINVAL)),
1950            Err(Errno::EINVAL)
1951        );
1952        assert_eq!(pselect6_probe_result(Err(Errno::EINTR)), Err(Errno::EINTR));
1953    }
1954
1955    #[test]
1956    fn ppoll_signal_mask_keeps_reverie_preemption_unblocked() {
1957        let preemption_bit = 1_u64 << ((reverie::PERF_EVENT_SIGNAL as usize) - 1);
1958        assert_eq!(sanitize_ppoll_signal_mask(u64::MAX), !preemption_bit);
1959    }
1960
1961    #[test]
1962    fn ppoll_record_replay_masked_waits_keep_fail_closed_probe() {
1963        assert!(!ppoll_uses_kernel_wait(true, true, true));
1964        assert!(ppoll_uses_kernel_wait(true, true, false));
1965        assert!(!ppoll_uses_kernel_wait(true, false, true));
1966        assert!(ppoll_uses_kernel_wait(false, true, true));
1967    }
1968
1969    #[test]
1970    fn ppoll_timeout_rejects_invalid_timespecs() {
1971        assert_eq!(
1972            ppoll_timeout_duration(Timespec {
1973                tv_sec: -1,
1974                tv_nsec: 0,
1975            }),
1976            Err(Errno::EINVAL)
1977        );
1978        assert_eq!(
1979            ppoll_timeout_duration(Timespec {
1980                tv_sec: 0,
1981                tv_nsec: 1_000_000_000,
1982            }),
1983            Err(Errno::EINVAL)
1984        );
1985    }
1986
1987    #[test]
1988    fn socket_timestamp_control_messages_use_logical_time() {
1989        let header_len = cmsg_align(std::mem::size_of::<libc::cmsghdr>());
1990        let timeval_len = std::mem::size_of::<libc::timeval>();
1991        let timespec_len = std::mem::size_of::<libc::timespec>();
1992        let first_len = header_len + timeval_len;
1993        let second_offset = cmsg_align(first_len);
1994        let second_len = header_len + timespec_len;
1995        let third_offset = second_offset + cmsg_align(second_len);
1996        let third_len = header_len + std::mem::size_of::<i32>();
1997        let mut control = vec![0; third_offset + cmsg_align(third_len)];
1998
1999        assert!(write_control_value(
2000            &mut control,
2001            libc::cmsghdr {
2002                cmsg_len: first_len,
2003                cmsg_level: libc::SOL_SOCKET,
2004                cmsg_type: SCM_TIMESTAMP_OLD,
2005            }
2006        ));
2007        assert!(write_control_value(
2008            &mut control[header_len..],
2009            libc::timeval {
2010                tv_sec: 99,
2011                tv_usec: 88,
2012            }
2013        ));
2014        assert!(write_control_value(
2015            &mut control[second_offset..],
2016            libc::cmsghdr {
2017                cmsg_len: second_len,
2018                cmsg_level: libc::SOL_SOCKET,
2019                cmsg_type: SCM_TIMESTAMPNS_OLD,
2020            }
2021        ));
2022        assert!(write_control_value(
2023            &mut control[second_offset + header_len..],
2024            libc::timespec {
2025                tv_sec: 77,
2026                tv_nsec: 66,
2027            }
2028        ));
2029        assert!(write_control_value(
2030            &mut control[third_offset..],
2031            libc::cmsghdr {
2032                cmsg_len: third_len,
2033                cmsg_level: libc::SOL_SOCKET,
2034                cmsg_type: libc::SCM_RIGHTS,
2035            }
2036        ));
2037        assert!(write_control_value(
2038            &mut control[third_offset + header_len..],
2039            42_i32
2040        ));
2041        let unrelated_message = control[third_offset..].to_vec();
2042
2043        assert_eq!(
2044            canonicalize_socket_timestamps(&mut control, LogicalTime::from_nanos(2_345_678_901)),
2045            2
2046        );
2047        let timeval = read_control_value::<libc::timeval>(&control[header_len..]).unwrap();
2048        assert_eq!(timeval.tv_sec, 2);
2049        assert_eq!(timeval.tv_usec, 345_678);
2050        let timespec =
2051            read_control_value::<libc::timespec>(&control[second_offset + header_len..]).unwrap();
2052        assert_eq!(timespec.tv_sec, 2);
2053        assert_eq!(timespec.tv_nsec, 345_678_901);
2054        assert_eq!(control[third_offset..], unrelated_message);
2055    }
2056
2057    #[test]
2058    fn truncated_timestamp_payload_prefix_is_rewritten() {
2059        let header_len = cmsg_align(std::mem::size_of::<libc::cmsghdr>());
2060        let full_len = header_len + std::mem::size_of::<libc::timeval>();
2061        let mut control = vec![0xaa; header_len + std::mem::size_of::<i32>()];
2062        assert!(write_control_value(
2063            &mut control,
2064            libc::cmsghdr {
2065                cmsg_len: full_len,
2066                cmsg_level: libc::SOL_SOCKET,
2067                cmsg_type: SCM_TIMESTAMP_OLD,
2068            }
2069        ));
2070
2071        assert_eq!(
2072            canonicalize_socket_timestamps(&mut control, LogicalTime::from_nanos(2_345_678_901)),
2073            1
2074        );
2075        assert_eq!(
2076            read_control_value::<i32>(&control[header_len..]),
2077            Some(2),
2078            "the visible timeval prefix must not retain host seconds"
2079        );
2080    }
2081
2082    #[test]
2083    fn timestamping_preserves_populated_source_slots() {
2084        let header_len = cmsg_align(std::mem::size_of::<libc::cmsghdr>());
2085        let timespec_len = std::mem::size_of::<libc::timespec>();
2086        let message_len = header_len + 3 * timespec_len;
2087        let mut control = vec![0; cmsg_align(message_len)];
2088        assert!(write_control_value(
2089            &mut control,
2090            libc::cmsghdr {
2091                cmsg_len: message_len,
2092                cmsg_level: libc::SOL_SOCKET,
2093                cmsg_type: SCM_TIMESTAMPING_OLD,
2094            }
2095        ));
2096        let zero = libc::timespec {
2097            tv_sec: 0,
2098            tv_nsec: 0,
2099        };
2100        let populated = libc::timespec {
2101            tv_sec: 99,
2102            tv_nsec: 88,
2103        };
2104        assert!(write_control_value(&mut control[header_len..], zero));
2105        assert!(write_control_value(
2106            &mut control[header_len + timespec_len..],
2107            populated
2108        ));
2109        assert!(write_control_value(
2110            &mut control[header_len + 2 * timespec_len..],
2111            populated
2112        ));
2113
2114        assert_eq!(
2115            canonicalize_socket_timestamps(&mut control, LogicalTime::from_nanos(2_345_678_901)),
2116            1
2117        );
2118        let first = read_control_value::<libc::timespec>(&control[header_len..]).unwrap();
2119        let second =
2120            read_control_value::<libc::timespec>(&control[header_len + timespec_len..]).unwrap();
2121        let third = read_control_value::<libc::timespec>(&control[header_len + 2 * timespec_len..])
2122            .unwrap();
2123        assert_eq!((first.tv_sec, first.tv_nsec), (0, 0));
2124        assert_eq!((second.tv_sec, second.tv_nsec), (2, 345_678_901));
2125        assert_eq!((third.tv_sec, third.tv_nsec), (2, 345_678_901));
2126    }
2127}