1use 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
52fn print_poll(call: &syscalls::Poll) {
55 let len = call.nfds();
56 debug!("POLL: on {} fds, timeout {}", len, call.timeout());
57 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
69fn 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 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 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 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 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 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
328fn 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 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 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 Ok(self.handle_external_poll(guest, call).await?)
427 } else {
428 self.handle_internal_poll(guest, call).await
431 }
432 }
433
434 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 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 let sigmask_argument = match call.sigmask() {
484 Some(argument) => {
485 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 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 if call.nfds() > PSELECT6_INTERNAL_MAX_NFDS {
535 return self
536 .record_or_replay_blocking(guest, Syscall::Pselect6(call))
537 .await;
538 }
539
540 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 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 let probe_timeout = stack.reserve::<Timespec>();
614 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 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 if let Err(error) = guest.memory().write_value(timeout, &remaining) {
722 trace!(?error, "ignoring pselect6 timeout writeback failure");
724 }
725 }
726 Ok(())
727 }
728
729 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 return Ok(guest.inject(call).await?);
771 }
772
773 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 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 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 if let Err(error) = guest.memory().write_value(timeout, &remaining) {
919 trace!(?error, "ignoring select timeout writeback failure");
921 }
922 }
923 Ok(())
924 }
925
926 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 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 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 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 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 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, ×pec_from_duration(remaining))
1085 {
1086 trace!(?error, "ignoring ppoll timeout writeback failure");
1088 }
1089 Ok(())
1090 }
1091
1092 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 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?) } 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 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 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; let fd = self.record_or_replay(guest, call).await? as RawFd;
1151 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 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; Ok(self.record_or_replay(guest, call).await?)
1184 }
1185
1186 pub async fn handle_epoll_pwait<G: Guest<Self>>(
1188 &self,
1189 guest: &mut G,
1190 call: syscalls::EpollPwait,
1191 ) -> Result<i64, Error> {
1192 if call.sigmask().is_some() {
1218 let dettid = guest.thread_state().dettid;
1219 resource_request(guest, Resources::new(dettid)).await; return Ok(self.record_or_replay(guest, call).await?);
1221 }
1222 if self.cfg.recordreplay_modes && call.timeout() == 0 {
1223 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 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 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?) } 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 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; Ok(self.record_or_replay(guest, call).await?)
1292 }
1293
1294 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 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 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?) } 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 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 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 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 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 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 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 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 pub async fn handle_recvmsg<G: Guest<Self>>(
1469 &self,
1470 guest: &mut G,
1471 call: syscalls::Recvmsg,
1472 ) -> Result<i64, Error> {
1473 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 let geometry: Vec<Vec<(usize, usize)>> = {
1739 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 let counts: Vec<usize> = {
1768 let address: AddrMut<'_, libc::mmsghdr> =
1769 AddrMut::from_raw(base).ok_or(Errno::EFAULT)?;
1770 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 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 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 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}