1use std::collections::BTreeMap;
12use std::collections::BTreeSet;
13use std::collections::HashMap;
14use std::os::fd::BorrowedFd;
15use std::path::Path;
16use std::sync::Arc;
17use std::sync::Mutex;
18use std::sync::MutexGuard;
19use std::time::Duration;
20
21use detcore_model::pedigree::Pedigree;
22use detcore_model::summary::TimesliceStats;
23use nix::fcntl::OFlag;
24use nix::sys::stat;
25use nix::unistd::Pid;
26use rand::Rng as _;
27use rand::RngExt as _;
28use rand::SeedableRng;
29use rand_distr::Distribution;
30use rand_distr::Exp;
31use rand_pcg::Pcg64Mcg;
32use reverie::Errno;
33use reverie::Error;
34use reverie::Guest;
35use reverie::syscalls::CloneFlags;
36use reverie::syscalls::Syscall;
37use reverie::syscalls::Sysno;
38use serde::Deserialize;
39use serde::Serialize;
40use sha2::Digest as _;
41use sha2::Sha256;
42use tracing::debug;
43use tracing::info;
44
45use crate::config::Config;
46use crate::detlog;
47use crate::fd::*;
48use crate::memory::MemoryMetadata;
49use crate::preemptions::ThreadHistoryIterator;
50use crate::record_or_replay::NoopTool;
51use crate::record_or_replay::RecordOrReplay;
52use crate::resources::ChaosEpochTransition;
53use crate::resources::Device;
54use crate::resources::Permission;
55use crate::resources::ResourceID;
56use crate::resources::Resources;
57use crate::scheduler::Priority;
58use crate::stat::*;
59use crate::types::*;
60
61#[derive(Debug, Serialize, Deserialize)]
63pub struct Detcore<T = NoopTool> {
64 pub(crate) detpid: DetPid,
70
71 pub(crate) cfg: Config,
73
74 pub(crate) record_or_replay: T,
78}
79
80#[derive(Debug, Clone, Serialize, Deserialize)]
82pub struct FileMetadata {
83 pub(crate) files_id: FilesId,
85 next_open_file_sequence: u64,
87 #[serde(default)]
89 next_socket_open_file_sequence: u64,
90 pub(crate) file_handles: HashMap<RawFd, DetFd>,
93}
94
95#[derive(Debug)]
100pub(crate) struct CapturedDetFd {
101 files_id: FilesId,
102 detfd: DetFd,
103}
104
105#[derive(Debug)]
109pub(crate) struct CapturedDetFdInstallError {
110 pub(crate) expected_files_id: FilesId,
111 pub(crate) actual_files_id: FilesId,
112 returned_fd: RawFd,
113 pub(crate) captured: CapturedDetFd,
114}
115
116#[derive(Debug, PartialEq, Eq)]
119pub(crate) struct CapturedDetFdInstallCleanup {
120 pub(crate) close_fd: RawFd,
122 pub(crate) release_open_file: Option<OpenFileId>,
124}
125
126impl CapturedDetFdInstallError {
127 pub(crate) fn into_cleanup(self) -> CapturedDetFdInstallCleanup {
129 let release_open_file = (self.captured.detfd.open_file_alias_count() == 1)
130 .then(|| self.captured.detfd.open_file_id());
131 drop(self.captured);
132 CapturedDetFdInstallCleanup {
133 close_fd: self.returned_fd,
134 release_open_file,
135 }
136 }
137}
138
139pub(crate) fn pidfd_getfd_targets_calling_task(
147 target: Option<DetPid>,
148 current_tgid: DetPid,
149 current_tid: DetTid,
150) -> bool {
151 target == Some(current_tgid) && current_tid == current_tgid
152}
153
154pub type ExecFdBlockingOverrides = BTreeSet<RawFd>;
159
160#[derive(Debug, Clone, Serialize, Deserialize)]
167struct PosixTimer {
168 interval_ns: u64,
170 deadline: Option<LogicalTime>,
173 signal: Option<i32>,
176}
177
178#[derive(Debug, Clone, Default, Serialize, Deserialize)]
185pub struct PosixTimers {
186 next_id: i32,
189 timers: HashMap<i32, PosixTimer>,
190}
191
192impl PosixTimers {
193 pub(crate) fn create(&mut self, signal: Option<i32>) -> i32 {
197 let id = self.next_id;
198 self.next_id += 1;
199 self.timers.insert(
200 id,
201 PosixTimer {
202 interval_ns: 0,
203 deadline: None,
204 signal,
205 },
206 );
207 id
208 }
209
210 pub(crate) fn settime(
216 &mut self,
217 id: i32,
218 interval_ns: u64,
219 deadline: Option<LogicalTime>,
220 now: LogicalTime,
221 ) -> Option<(u64, u64)> {
222 let timer = self.timers.get_mut(&id)?;
223 let old = (
224 remaining_ns(timer.deadline, timer.interval_ns, now),
225 timer.interval_ns,
226 );
227 timer.interval_ns = interval_ns;
228 timer.deadline = deadline;
229 Some(old)
230 }
231
232 pub(crate) fn gettime(&self, id: i32, now: LogicalTime) -> Option<(u64, u64)> {
235 let timer = self.timers.get(&id)?;
236 Some((
237 remaining_ns(timer.deadline, timer.interval_ns, now),
238 timer.interval_ns,
239 ))
240 }
241
242 pub(crate) fn contains(&self, id: i32) -> bool {
244 self.timers.contains_key(&id)
245 }
246
247 pub(crate) fn signal(&self, id: i32) -> Option<Option<i32>> {
250 self.timers.get(&id).map(|timer| timer.signal)
251 }
252
253 pub(crate) fn remove(&mut self, id: i32) -> bool {
255 self.timers.remove(&id).is_some()
256 }
257
258 pub(crate) fn clear_for_exec(&mut self) {
261 self.timers.clear();
262 }
263}
264
265#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
267pub(crate) struct ResourceLimit {
268 pub(crate) current: u64,
269 pub(crate) maximum: u64,
270}
271
272#[derive(Debug, Clone, Serialize, Deserialize)]
274pub(crate) struct ResourceLimits {
275 limits: Vec<ResourceLimit>,
276}
277
278impl Default for ResourceLimits {
279 fn default() -> Self {
280 let unlimited = ResourceLimit {
281 current: libc::RLIM64_INFINITY,
282 maximum: libc::RLIM64_INFINITY,
283 };
284 let mut limits = vec![unlimited; libc::RLIMIT_RTTIME as usize + 1];
285 limits[libc::RLIMIT_STACK as usize] = ResourceLimit {
286 current: 8 * 1024 * 1024,
287 maximum: libc::RLIM64_INFINITY,
288 };
289 limits[libc::RLIMIT_NOFILE as usize] = ResourceLimit {
290 current: 1_048_576,
291 maximum: 1_048_576,
292 };
293 Self { limits }
294 }
295}
296
297impl ResourceLimits {
298 pub(crate) fn get(&self, resource: u32) -> Option<ResourceLimit> {
300 self.limits.get(resource as usize).copied()
301 }
302
303 pub(crate) fn set(&mut self, resource: u32, limit: ResourceLimit) {
305 self.limits[resource as usize] = limit;
306 }
307}
308
309fn remaining_ns(deadline: Option<LogicalTime>, interval_ns: u64, now: LogicalTime) -> u64 {
313 match deadline {
314 Some(d) if d > now => d.as_nanos() - now.as_nanos(),
315 Some(d) if interval_ns != 0 => {
316 let elapsed = now.as_nanos() - d.as_nanos();
317 interval_ns - (elapsed % interval_ns)
318 }
319 Some(_) => 0,
320 None => 0,
321 }
322}
323
324impl<T> Default for Detcore<T> {
325 fn default() -> Self {
326 panic!("Detcore Default impl should not be called");
330 }
331}
332
333impl<T: RecordOrReplay> AsRef<T> for Detcore<T> {
334 fn as_ref(&self) -> &T {
335 &self.record_or_replay
336 }
337}
338
339impl<T: RecordOrReplay> AsMut<T> for Detcore<T> {
340 fn as_mut(&mut self) -> &mut T {
341 &mut self.record_or_replay
342 }
343}
344
345impl<T: RecordOrReplay> Detcore<T> {
346 pub(crate) async fn record_or_replay_preserving_tool_errors<G, S>(
350 &self,
351 guest: &mut G,
352 syscall: S,
353 ) -> Result<i64, Error>
354 where
355 G: Guest<Self>,
356 S: Into<Syscall>,
357 {
358 preserve_record_or_replay_error(
359 self.record_or_replay
360 .handle_syscall_event(&mut guest.into_guest(), syscall.into())
361 .await,
362 )?
363 .map_err(Error::from)
364 }
365
366 pub(crate) async fn record_or_replay<G, S>(
383 &self,
384 guest: &mut G,
385 syscall: S,
386 ) -> Result<i64, Errno>
387 where
388 G: Guest<Self>,
389 S: Into<Syscall>,
390 {
391 self.record_or_replay
392 .handle_syscall_event(&mut guest.into_guest(), syscall.into())
393 .await
394 .map_err(|err| err.into_errno().unwrap())
396 }
397}
398
399fn preserve_record_or_replay_error(
400 result: Result<i64, Error>,
401) -> Result<Result<i64, Errno>, Error> {
402 match result {
403 Ok(value) => Ok(Ok(value)),
404 Err(Error::Errno(error)) => Ok(Err(error)),
405 Err(error) => Err(error),
406 }
407}
408
409pub(crate) fn finish_partial_record_or_replay_write(
414 written: i64,
415 error: Error,
416) -> Result<i64, Error> {
417 match error {
418 Error::Errno(_) if written > 0 => Ok(written),
419 error => Err(error),
420 }
421}
422
423#[cfg(test)]
424mod record_or_replay_error_tests {
425 use super::*;
426
427 #[test]
428 fn errno_remains_a_guest_syscall_result() {
429 let result = preserve_record_or_replay_error(Err(Error::Errno(Errno::EBADF)))
430 .expect("guest errno must not abort the tool");
431 assert_eq!(result, Err(Errno::EBADF));
432 }
433
434 #[test]
435 fn tool_failure_remains_an_outer_error() {
436 let failure =
437 preserve_record_or_replay_error(Err(Error::Tool(anyhow::anyhow!("restore failed"))))
438 .expect_err("tool failure must abort replay");
439 assert!(matches!(failure.into_errno(), Err(Error::Tool(_))));
440 }
441
442 #[test]
443 fn partial_write_suppresses_only_guest_errno() {
444 assert_eq!(
445 finish_partial_record_or_replay_write(7, Error::Errno(Errno::EPIPE)).unwrap(),
446 7
447 );
448
449 let io_failure = finish_partial_record_or_replay_write(
450 7,
451 Error::Io(std::io::Error::from_raw_os_error(libc::ENOSPC)),
452 )
453 .expect_err("host replay-output failure must not become a partial guest write");
454 assert!(matches!(io_failure.into_errno(), Err(Error::Io(_))));
455
456 let tool_failure = finish_partial_record_or_replay_write(
457 7,
458 Error::Tool(anyhow::anyhow!("replay endpoint disappeared")),
459 )
460 .expect_err("tool failure must not become a partial guest write");
461 assert!(matches!(tool_failure.into_errno(), Err(Error::Tool(_))));
462
463 let guest_failure = finish_partial_record_or_replay_write(0, Error::Errno(Errno::EPIPE))
464 .expect_err("zero-progress guest failure must remain guest-visible");
465 assert_eq!(guest_failure.into_errno().unwrap(), Errno::EPIPE);
466 }
467}
468
469impl FileMetadata {
470 fn new(owner: DetTid) -> Self {
472 FileMetadata {
473 files_id: FilesId::initial(owner),
474 next_open_file_sequence: 0,
475 next_socket_open_file_sequence: 0,
476 file_handles: HashMap::new(),
477 }
478 }
479
480 fn allocate_open_file_id(&mut self, creator: DetTid, ty: FdType) -> OpenFileId {
481 if ty == FdType::Socket {
482 let id = OpenFileId::new_socket(creator, self.next_socket_open_file_sequence);
483 self.next_socket_open_file_sequence += 1;
484 id
485 } else {
486 let id = OpenFileId::new(creator, self.next_open_file_sequence);
487 self.next_open_file_sequence += 1;
488 id
489 }
490 }
491
492 fn count_open_files_at_paths(&self, paths: &[&Path]) -> usize {
493 self.file_handles
494 .values()
495 .filter(|fd| {
496 fd.path()
497 .is_some_and(|path| paths.iter().any(|candidate| path == *candidate))
498 })
499 .map(DetFd::open_file_id)
500 .collect::<BTreeSet<_>>()
501 .len()
502 }
503
504 fn has_loopback_peer(&self) -> bool {
505 self.file_handles.values().any(DetFd::is_loopback_peer)
506 }
507
508 fn has_unsafe_vfork_flock_state(&self) -> bool {
525 self.file_handles.values().any(|detfd| {
526 matches!(detfd.known_flock_mode(), Some(Some(_))) || detfd.flock_mode_may_be_stale()
527 })
528 }
529
530 fn forget_flock_modes(&self) {
531 for detfd in self.file_handles.values() {
532 detfd.forget_flock_mode();
533 }
534 }
535
536 pub(crate) fn fork_for(&self, child: DetTid) -> Self {
563 self.forget_flock_modes();
564 Self {
565 files_id: FilesId::forked(child),
566 next_open_file_sequence: self.next_open_file_sequence,
567 next_socket_open_file_sequence: self.next_socket_open_file_sequence,
568 file_handles: self.file_handles.clone(),
569 }
570 }
571
572 pub(crate) fn for_exec(&self, task: DetTid) -> Self {
573 Self {
574 files_id: self.files_id.for_exec(task),
575 next_open_file_sequence: self.next_open_file_sequence,
576 next_socket_open_file_sequence: self.next_socket_open_file_sequence,
577 file_handles: self
578 .file_handles
579 .iter()
580 .filter_map(|(&fd, detfd)| (!detfd.is_cloexec()).then_some((fd, detfd.clone())))
581 .collect(),
582 }
583 }
584
585 pub(crate) fn exec_blocking_overrides(&self) -> ExecFdBlockingOverrides {
586 self.file_handles
587 .iter()
588 .filter_map(|(&fd, detfd)| {
589 (!detfd.is_cloexec() && detfd.physically_nonblocking() && !detfd.is_nonblocking())
590 .then_some(fd)
591 })
592 .collect()
593 }
594
595 pub(crate) fn apply_exec_blocking_overrides(
596 &mut self,
597 owner: DetTid,
598 overrides: ExecFdBlockingOverrides,
599 ) {
600 for fd in overrides {
601 tracing::trace!(
602 "[detcore, dtid {}] restoring descriptor {} as logically blocking after exec",
603 owner,
604 fd,
605 );
606 match self.discover_fd_from_current_process(owner, fd) {
607 Ok(()) => {
608 self.with_detfd(fd, |detfd| detfd.set_logical_nonblocking(false))
609 .expect("a just-discovered descriptor must remain registered");
610 }
611 Err(error) => {
612 tracing::warn!(
613 "[detcore, dtid {}] unable to restore inherited descriptor {} after exec: {}",
614 owner,
615 fd,
616 error,
617 );
618 }
619 }
620 }
621 }
622
623 pub(crate) fn open_files_closed_on_exec(&self, table_is_shared: bool) -> Vec<OpenFileId> {
624 if table_is_shared {
625 return Vec::new();
626 }
627
628 let mut open_files = HashMap::new();
629 for detfd in self.file_handles.values() {
630 let id = detfd.open_file_id();
631 let total_aliases = detfd.open_file_alias_count();
632 let entry = open_files.entry(id).or_insert((0, total_aliases, true));
633 debug_assert_eq!(entry.1, total_aliases);
634 entry.0 += 1;
635 entry.2 &= detfd.is_cloexec();
636 }
637
638 let mut closed: Vec<_> = open_files
639 .into_iter()
640 .filter_map(|(id, (table_aliases, total_aliases, all_cloexec))| {
641 (all_cloexec && table_aliases == total_aliases).then_some(id)
642 })
643 .collect();
644 closed.sort();
645 closed
646 }
647
648 fn setup_stdio(mut self, _pid: Pid, owner: DetTid) -> Self {
650 let stat: DetStat = stat::fstat(unsafe { BorrowedFd::borrow_raw(0) })
654 .unwrap()
655 .into();
656 let stdin = DetFd::new(
657 0,
658 OFlag::empty(),
659 FdType::Regular,
660 self.allocate_open_file_id(owner, FdType::Regular),
661 )
662 .with_stat(stat)
663 .with_resource(ResourceID::Device(Device::ContainerStdin));
664 let stdout = DetFd::new(
665 1,
666 OFlag::empty(),
667 FdType::Regular,
668 self.allocate_open_file_id(owner, FdType::Regular),
669 )
670 .with_stat(stat)
671 .with_resource(ResourceID::Device(Device::ContainerStdout));
672 let stderr = DetFd::new(
673 2,
674 OFlag::empty(),
675 FdType::Regular,
676 self.allocate_open_file_id(owner, FdType::Regular),
677 )
678 .with_stat(stat)
679 .with_resource(ResourceID::Device(Device::ContainerStderr));
680
681 stdin.mark_flock_mode_unobserved();
684 stdout.mark_flock_mode_unobserved();
685 stderr.mark_flock_mode_unobserved();
686
687 self.add_detfd(stdin);
688 self.add_detfd(stdout);
689 self.add_detfd(stderr);
690
691 self
692 }
693
694 fn discover_fd_from_current_process(&mut self, owner: DetTid, fd: RawFd) -> Result<(), Errno> {
697 if self.file_handles.contains_key(&fd) {
698 return Ok(());
699 }
700
701 let fd_flags = unsafe { libc::fcntl(fd, libc::F_GETFD) };
702 let status_flags = unsafe { libc::fcntl(fd, libc::F_GETFL) };
703 if fd_flags == -1 || status_flags == -1 {
704 return Err(Errno::last());
705 }
706 let raw_stat =
707 stat::fstat(unsafe { BorrowedFd::borrow_raw(fd) }).map_err(|_| Errno::last())?;
708 let file_type = stat::SFlag::from_bits_truncate(raw_stat.st_mode);
709 let ty = if file_type.contains(stat::SFlag::S_IFIFO) {
710 FdType::Pipe
711 } else if file_type.contains(stat::SFlag::S_IFSOCK) {
712 FdType::Socket
713 } else {
714 FdType::Regular
715 };
716 let mut flags = OFlag::from_bits_truncate(status_flags);
717 let physically_nonblocking = flags.contains(OFlag::O_NONBLOCK);
718 if fd_flags & libc::FD_CLOEXEC != 0 {
722 flags.insert(OFlag::O_CLOEXEC);
723 }
724 self.add_fd(owner, fd, flags, ty, Some(raw_stat.into()))?;
725 self.with_detfd(fd, |detfd| detfd.forget_flock_mode())?;
726 if let Some(resource) = stdio_resource(fd) {
727 self.with_detfd(fd, |detfd| detfd.set_resource(resource.clone()))?;
728 }
729 if ty == FdType::Pipe && physically_nonblocking {
730 self.with_detfd(fd, |detfd| detfd.set_physically_nonblocking())?;
731 }
732 Ok(())
733 }
734
735 fn with_detfd<F, U>(&mut self, fd: RawFd, mut f: F) -> Result<U, Errno>
737 where
738 F: FnMut(&mut DetFd) -> U,
739 {
740 let detfd = self.file_handles.get_mut(&fd).ok_or(Errno::EBADF)?;
741 Ok(f(detfd))
742 }
743
744 fn add_detfd(&mut self, detfd: DetFd) {
746 let fd = detfd.fd;
747 self.file_handles.insert(fd, detfd);
748 }
749
750 fn add_fd(
752 &mut self,
753 creator: DetTid,
754 fd: RawFd,
755 flags: OFlag,
756 ty: FdType,
757 stat: Option<DetStat>,
758 ) -> Result<(), Errno> {
759 let id = self.allocate_open_file_id(creator, ty);
760 let detfd = DetFd::new(fd, flags, ty, id).with_stat(stat);
761 self.add_detfd(detfd);
762 Ok(())
763 }
764
765 fn remove_fd(&mut self, fd: RawFd) -> Option<OpenFileId> {
767 let detfd = self.file_handles.remove(&fd)?;
768 (detfd.open_file_alias_count() == 1).then(|| detfd.open_file_id())
769 }
770
771 fn remove_fd_range(&mut self, first: u32, last: u32) -> Vec<OpenFileId> {
773 let mut descriptors: Vec<_> = self
774 .file_handles
775 .keys()
776 .copied()
777 .filter(|fd| *fd >= 0 && first <= *fd as u32 && *fd as u32 <= last)
778 .collect();
779 descriptors.sort_unstable();
780 descriptors
781 .into_iter()
782 .filter_map(|fd| self.remove_fd(fd))
783 .collect()
784 }
785
786 fn dup_fd(
788 &mut self,
789 oldfd: RawFd,
790 newfd: RawFd,
791 flags: OFlag,
792 ) -> Result<Option<OpenFileId>, Errno> {
793 if oldfd == newfd {
794 self.with_detfd(oldfd, |_| ())?;
795 return Ok(None);
796 }
797
798 let detfd = self.with_detfd(oldfd, |old_detfd| {
799 old_detfd.clone().with_fd(newfd).with_fd_flags(flags)
800 })?;
801 let replaced = self.file_handles.insert(newfd, detfd);
802 Ok(replaced
803 .and_then(|detfd| (detfd.open_file_alias_count() == 1).then(|| detfd.open_file_id())))
804 }
805
806 fn capture_fd(&mut self, fd: RawFd) -> Result<CapturedDetFd, Errno> {
808 let detfd = self.with_detfd(fd, |detfd| detfd.clone())?;
809 Ok(CapturedDetFd {
810 files_id: self.files_id,
811 detfd,
812 })
813 }
814
815 fn install_captured_fd(
821 &mut self,
822 captured: CapturedDetFd,
823 newfd: RawFd,
824 flags: OFlag,
825 ) -> Result<Option<OpenFileId>, CapturedDetFdInstallError> {
826 if captured.files_id != self.files_id {
827 return Err(CapturedDetFdInstallError {
828 expected_files_id: captured.files_id,
829 actual_files_id: self.files_id,
830 returned_fd: newfd,
831 captured,
832 });
833 }
834
835 let detfd = captured.detfd.with_fd(newfd).with_fd_flags(flags);
836 let replaced = self.file_handles.insert(newfd, detfd);
837 Ok(replaced
838 .and_then(|detfd| (detfd.open_file_alias_count() == 1).then(|| detfd.open_file_id())))
839 }
840
841 fn abandon_captured_fd(&mut self, captured: CapturedDetFd) -> Option<OpenFileId> {
844 let release =
845 (captured.detfd.open_file_alias_count() == 1).then(|| captured.detfd.open_file_id());
846 drop(captured);
847 release
848 }
849}
850
851fn stdio_resource(fd: RawFd) -> Option<ResourceID> {
852 match fd {
853 0 => Some(ResourceID::Device(Device::ContainerStdin)),
854 1 => Some(ResourceID::Device(Device::ContainerStdout)),
855 2 => Some(ResourceID::Device(Device::ContainerStderr)),
856 _ => None,
857 }
858}
859
860#[cfg(test)]
861mod posix_timers_tests {
862 use super::*;
863
864 fn t(ns: u64) -> LogicalTime {
865 LogicalTime::from_nanos(ns)
866 }
867
868 #[test]
869 fn ids_are_deterministic_and_sequential() {
870 let mut timers = PosixTimers::default();
871 assert_eq!(timers.create(None), 0);
872 assert_eq!(timers.create(Some(libc::SIGALRM)), 1);
873 assert_eq!(timers.create(None), 2);
874 }
875
876 #[test]
877 fn settime_reports_previous_arming_and_remaining_uses_virtual_clock() {
878 let mut timers = PosixTimers::default();
879 let id = timers.create(None);
880
881 let old = timers.settime(id, 0, Some(t(100)), t(0)).expect("known id");
884 assert_eq!(old, (0, 0));
885
886 assert_eq!(timers.gettime(id, t(40)), Some((60, 0)));
888 assert_eq!(timers.gettime(id, t(150)), Some((0, 0)));
890 }
891
892 #[test]
893 fn resetting_reports_old_remaining() {
894 let mut timers = PosixTimers::default();
895 let id = timers.create(None);
896 timers.settime(id, 0, Some(t(100)), t(0));
897 let old = timers
899 .settime(id, 50, Some(t(200)), t(30))
900 .expect("known id");
901 assert_eq!(old, (70, 0));
902 assert_eq!(timers.gettime(id, t(30)), Some((170, 50)));
903 }
904
905 #[test]
908 fn periodic_remaining_advances_past_each_deadline() {
909 let mut timers = PosixTimers::default();
910 let id = timers.create(Some(libc::SIGALRM));
911 timers.settime(id, 50, Some(t(100)), t(0));
912
913 assert_eq!(timers.gettime(id, t(100)), Some((50, 50)));
914 assert_eq!(timers.gettime(id, t(125)), Some((25, 50)));
915 assert_eq!(timers.gettime(id, t(150)), Some((50, 50)));
916 assert_eq!(timers.signal(id), Some(Some(libc::SIGALRM)));
917 }
918
919 #[test]
920 fn disarm_and_unknown_ids() {
921 let mut timers = PosixTimers::default();
922 let id = timers.create(None);
923 timers.settime(id, 0, Some(t(100)), t(0));
924 timers.settime(id, 0, None, t(10));
926 assert_eq!(timers.gettime(id, t(10)), Some((0, 0)));
927
928 assert_eq!(timers.settime(99, 0, Some(t(1)), t(0)), None);
930 assert_eq!(timers.gettime(99, t(0)), None);
931 assert!(!timers.contains(99));
932 }
933
934 #[test]
935 fn delete_removes_timer() {
936 let mut timers = PosixTimers::default();
937 let id = timers.create(None);
938 assert!(timers.contains(id));
939 assert!(timers.remove(id));
940 assert!(!timers.contains(id));
941 assert!(!timers.remove(id));
943 }
944}
945
946#[cfg(test)]
947mod resource_limits_tests {
948 use super::*;
949
950 #[test]
951 fn defaults_are_fixed_and_cover_linux_resources() {
952 let limits = ResourceLimits::default();
953 assert_eq!(
954 limits.get(libc::RLIMIT_STACK),
955 Some(ResourceLimit {
956 current: 8 * 1024 * 1024,
957 maximum: libc::RLIM64_INFINITY,
958 })
959 );
960 assert_eq!(
961 limits.get(libc::RLIMIT_NOFILE),
962 Some(ResourceLimit {
963 current: 1_048_576,
964 maximum: 1_048_576,
965 })
966 );
967 assert_eq!(
968 limits.get(libc::RLIMIT_CORE),
969 Some(ResourceLimit {
970 current: libc::RLIM64_INFINITY,
971 maximum: libc::RLIM64_INFINITY,
972 })
973 );
974 assert_eq!(limits.get(libc::RLIMIT_RTTIME + 1), None);
975 }
976
977 #[test]
978 fn cloned_process_state_changes_independently() {
979 let parent = ResourceLimits::default();
980 let mut child = parent.clone();
981 let lowered = ResourceLimit {
982 current: 1024,
983 maximum: 1_048_576,
984 };
985 child.set(libc::RLIMIT_NOFILE, lowered);
986
987 assert_eq!(child.get(libc::RLIMIT_NOFILE), Some(lowered));
988 assert_eq!(
989 parent.get(libc::RLIMIT_NOFILE),
990 Some(ResourceLimit {
991 current: 1_048_576,
992 maximum: 1_048_576,
993 })
994 );
995 }
996}
997
998#[cfg(test)]
999mod file_metadata_tests {
1000 use std::os::fd::AsRawFd;
1001
1002 use super::*;
1003
1004 #[test]
1005 fn on_demand_discovery_finds_a_live_descriptor() {
1006 let owner = DetTid::from_raw(9);
1007 let file = std::fs::File::open("/dev/null").expect("open test descriptor");
1008 let fd = file.as_raw_fd();
1009 let mut metadata = FileMetadata::new(owner);
1010
1011 assert_eq!(metadata.with_detfd(fd, |_| ()), Err(Errno::EBADF));
1012 metadata
1013 .discover_fd_from_current_process(owner, fd)
1014 .expect("live descriptor should be discovered");
1015 assert_eq!(
1016 metadata
1017 .with_detfd(fd, |detfd| detfd.ty())
1018 .expect("discovered descriptor should be tracked"),
1019 FdType::Regular
1020 );
1021 assert_eq!(
1022 metadata
1023 .with_detfd(fd, |detfd| detfd.known_flock_mode())
1024 .expect("discovered descriptor should be tracked"),
1025 None,
1026 "a live descriptor may already hold a flock that Detcore did not observe"
1027 );
1028 }
1029
1030 #[test]
1031 fn fork_makes_inherited_flock_state_unknown_in_parent_and_child() {
1032 let owner = DetTid::from_raw(9);
1033 let child = DetTid::from_raw(10);
1034 let mut metadata = FileMetadata::new(owner);
1035 metadata
1036 .add_fd(owner, 3, OFlag::empty(), FdType::Regular, None)
1037 .expect("register descriptor");
1038 metadata
1039 .with_detfd(3, |detfd| detfd.set_flock_mode(Some(libc::LOCK_SH)))
1040 .expect("registered descriptor");
1041
1042 let mut child_metadata = metadata.fork_for(child);
1043 assert_eq!(
1044 metadata
1045 .with_detfd(3, |detfd| detfd.known_flock_mode())
1046 .expect("parent descriptor should remain tracked"),
1047 None,
1048 "parent cache must not become stale when the child changes the shared open file description"
1049 );
1050 assert_eq!(
1051 child_metadata
1052 .with_detfd(3, |detfd| detfd.known_flock_mode())
1053 .expect("child descriptor should be inherited"),
1054 None,
1055 "child starts with the same deliberately unknown shared state"
1056 );
1057 }
1058
1059 #[test]
1060 fn exec_preserves_known_flock_mode_for_surviving_descriptors() {
1061 let owner = DetTid::from_raw(9);
1062 let mut metadata = FileMetadata::new(owner);
1063 metadata
1064 .add_fd(owner, 3, OFlag::empty(), FdType::Regular, None)
1065 .expect("register descriptor");
1066 metadata
1067 .with_detfd(3, |detfd| detfd.set_flock_mode(Some(libc::LOCK_SH)))
1068 .expect("registered descriptor");
1069
1070 let mut after_exec = metadata.for_exec(DetTid::from_raw(10));
1071 assert_eq!(
1072 after_exec
1073 .with_detfd(3, |detfd| detfd.known_flock_mode())
1074 .expect("non-cloexec descriptor should survive"),
1075 Some(Some(libc::LOCK_SH)),
1076 "exec must preserve known flock state for a surviving open file description"
1077 );
1078 }
1079
1080 #[test]
1081 fn discovered_stdio_uses_container_wide_resources() {
1082 let owner = DetTid::from_raw(9);
1083 let mut metadata = FileMetadata::new(owner);
1084
1085 metadata
1086 .discover_fd_from_current_process(owner, libc::STDOUT_FILENO)
1087 .expect("live stdout should be discovered");
1088
1089 assert_eq!(
1090 metadata
1091 .with_detfd(libc::STDOUT_FILENO, |detfd| detfd.resource())
1092 .expect("discovered stdout should be tracked"),
1093 Some(ResourceID::Device(Device::ContainerStdout))
1094 );
1095
1096 use crate::syscalls::deterministic_stdio_inode_for_resource;
1100
1101 let inherited_stat = metadata
1102 .with_detfd(libc::STDOUT_FILENO, |detfd| detfd.stat().unwrap())
1103 .unwrap();
1104 metadata.dup_fd(1, 7, OFlag::empty()).unwrap();
1105 assert_eq!(
1106 metadata.with_detfd(7, |detfd| detfd.resource()).unwrap(),
1107 Some(ResourceID::Device(Device::ContainerStdout))
1108 );
1109 assert_eq!(
1110 metadata
1111 .with_detfd(7, |detfd| {
1112 deterministic_stdio_inode_for_resource(7, detfd.resource())
1113 })
1114 .unwrap(),
1115 None
1116 );
1117 for fd in 0..=2 {
1118 metadata.remove_fd(fd);
1119 let mut ordinary = inherited_stat;
1120 ordinary.inode = 123 + fd as u64;
1121 metadata
1122 .add_fd(owner, fd, OFlag::empty(), FdType::Regular, Some(ordinary))
1123 .unwrap();
1124 metadata.dup_fd(fd, 8, OFlag::empty()).unwrap();
1125 for ordinary_fd in [fd, 8] {
1126 assert_eq!(
1127 metadata
1128 .with_detfd(ordinary_fd, |detfd| {
1129 deterministic_stdio_inode_for_resource(ordinary_fd, detfd.resource())
1130 })
1131 .unwrap(),
1132 None
1133 );
1134 }
1135 assert_eq!(
1136 metadata
1137 .with_detfd(fd, |detfd| detfd.open_file_id())
1138 .unwrap(),
1139 metadata
1140 .with_detfd(8, |detfd| detfd.open_file_id())
1141 .unwrap()
1142 );
1143 metadata.dup_fd(7, fd, OFlag::empty()).unwrap();
1144 assert_eq!(
1145 metadata
1146 .with_detfd(fd, |detfd| {
1147 deterministic_stdio_inode_for_resource(fd, detfd.resource())
1148 })
1149 .unwrap(),
1150 Some(DetInode::mint(1000 + fd as u64)),
1151 "inherited streams retain the existing numeric-slot outcome"
1152 );
1153 }
1154 metadata.dup_fd(7, 9, OFlag::O_CLOEXEC).unwrap();
1155 let child = metadata.fork_for(DetTid::from_raw(10));
1156 let after_exec = child.for_exec(DetTid::from_raw(10));
1157 assert!(!after_exec.file_handles.contains_key(&9));
1158 assert!(after_exec.file_handles.contains_key(&7));
1159 }
1160
1161 #[test]
1162 fn discovered_pipe_preserves_inherited_nonblocking() {
1163 let owner = DetTid::from_raw(9);
1164 let mut fds = [-1; 2];
1165 assert_eq!(
1166 unsafe { libc::pipe2(fds.as_mut_ptr(), libc::O_NONBLOCK) },
1167 0
1168 );
1169 let mut metadata = FileMetadata::new(owner);
1170
1171 metadata
1172 .discover_fd_from_current_process(owner, fds[0])
1173 .expect("live pipe should be discovered");
1174 let flags = metadata
1175 .with_detfd(fds[0], |detfd| {
1176 (detfd.is_nonblocking(), detfd.physically_nonblocking())
1177 })
1178 .expect("discovered pipe should be tracked");
1179
1180 assert_eq!(flags, (true, true));
1181 unsafe {
1182 libc::close(fds[0]);
1183 libc::close(fds[1]);
1184 }
1185 }
1186
1187 #[test]
1188 fn exec_handoff_restores_scheduler_pipe_as_logically_blocking() {
1189 let owner = DetTid::from_raw(9);
1190 let mut fds = [-1; 2];
1191 assert_eq!(
1192 unsafe { libc::pipe2(fds.as_mut_ptr(), libc::O_NONBLOCK) },
1193 0
1194 );
1195 let mut before_exec = FileMetadata::new(owner);
1196 before_exec
1197 .add_fd(owner, fds[0], OFlag::empty(), FdType::Pipe, None)
1198 .expect("scheduler pipe should be registered");
1199 before_exec
1200 .with_detfd(fds[0], |detfd| detfd.set_physically_nonblocking())
1201 .expect("scheduler pipe should remain registered");
1202
1203 let overrides = before_exec.exec_blocking_overrides();
1204 assert_eq!(overrides, BTreeSet::from([fds[0]]));
1205
1206 let mut after_exec = FileMetadata::new(owner);
1207 after_exec.apply_exec_blocking_overrides(owner, overrides);
1208 assert_eq!(
1209 after_exec
1210 .with_detfd(fds[0], |detfd| {
1211 (detfd.is_nonblocking(), detfd.physically_nonblocking())
1212 })
1213 .expect("inherited pipe should be rediscovered"),
1214 (false, true)
1215 );
1216
1217 unsafe {
1218 libc::close(fds[0]);
1219 libc::close(fds[1]);
1220 }
1221 }
1222
1223 #[test]
1224 fn fork_copies_slots_but_preserves_open_file_aliases() {
1225 let parent_tid = DetTid::from_raw(10);
1226 let child_tid = DetTid::from_raw(11);
1227 let mut parent = FileMetadata::new(parent_tid);
1228 parent
1229 .add_fd(parent_tid, 3, OFlag::O_NONBLOCK, FdType::Socket, None)
1230 .expect("parent fd should be inserted");
1231 parent
1232 .dup_fd(3, 4, OFlag::O_CLOEXEC)
1233 .expect("dup should succeed");
1234
1235 let parent_open = parent
1236 .with_detfd(3, |fd| fd.open_file_id())
1237 .expect("parent fd should exist");
1238 let duplicate_open = parent
1239 .with_detfd(4, |fd| fd.open_file_id())
1240 .expect("duplicate fd should exist");
1241 assert_eq!(parent_open, duplicate_open);
1242
1243 let initial_timestamp = LogicalTime::from_nanos(1_234_567_890);
1244 parent
1245 .with_detfd(3, |fd| fd.set_socket_receive_timestamp(initial_timestamp))
1246 .expect("parent socket should accept a receive timestamp");
1247
1248 let mut child = parent.fork_for(child_tid);
1249 assert_ne!(parent.files_id, child.files_id);
1250 assert_ne!(
1251 FdSlot {
1252 files: parent.files_id,
1253 fd: 3,
1254 },
1255 FdSlot {
1256 files: child.files_id,
1257 fd: 3,
1258 }
1259 );
1260 assert_eq!(
1261 parent_open,
1262 child
1263 .with_detfd(3, |fd| fd.open_file_id())
1264 .expect("forked fd should retain its open file identity")
1265 );
1266 assert_eq!(
1267 child
1268 .with_detfd(3, |fd| fd.socket_receive_timestamp())
1269 .expect("forked fd should retain its receive timestamp"),
1270 Some(initial_timestamp)
1271 );
1272 let child_timestamp = LogicalTime::from_nanos(2_345_678_901);
1273 child
1274 .with_detfd(3, |fd| fd.set_socket_receive_timestamp(child_timestamp))
1275 .expect("child socket should update the shared receive timestamp");
1276 assert_eq!(
1277 parent
1278 .with_detfd(4, |fd| fd.socket_receive_timestamp())
1279 .expect("parent duplicate should see the child update"),
1280 Some(child_timestamp)
1281 );
1282
1283 parent
1284 .add_fd(parent_tid, 5, OFlag::empty(), FdType::Regular, None)
1285 .expect("new parent fd should be inserted");
1286 child
1287 .add_fd(child_tid, 5, OFlag::empty(), FdType::Regular, None)
1288 .expect("new child fd should be inserted");
1289 assert_ne!(
1290 parent
1291 .with_detfd(5, |fd| fd.open_file_id())
1292 .expect("new parent fd should exist"),
1293 child
1294 .with_detfd(5, |fd| fd.open_file_id())
1295 .expect("new child fd should exist"),
1296 "separate opens after fork must not alias"
1297 );
1298 }
1299
1300 #[test]
1301 fn pidfd_getfd_target_requires_the_calling_leader() {
1302 let leader = DetPid::from_raw(31);
1303 assert!(pidfd_getfd_targets_calling_task(
1304 Some(leader),
1305 leader,
1306 DetTid::from_raw(31),
1307 ));
1308 assert!(!pidfd_getfd_targets_calling_task(
1309 Some(leader),
1310 leader,
1311 DetTid::from_raw(32),
1312 ));
1313 assert!(!pidfd_getfd_targets_calling_task(
1314 Some(DetPid::from_raw(32)),
1315 leader,
1316 DetTid::from_raw(31),
1317 ));
1318 }
1319
1320 #[test]
1321 fn captured_fd_installs_when_the_descriptor_table_is_unchanged() {
1322 let owner = DetTid::from_raw(32);
1323 let mut metadata = FileMetadata::new(owner);
1324 metadata
1325 .add_fd(owner, 3, OFlag::empty(), FdType::Regular, None)
1326 .expect("source should be inserted");
1327 let source_id = metadata
1328 .with_detfd(3, |fd| fd.open_file_id())
1329 .expect("source should exist");
1330 let captured = metadata.capture_fd(3).expect("source should be captured");
1331
1332 assert_eq!(
1333 metadata
1334 .install_captured_fd(captured, 4, OFlag::O_CLOEXEC)
1335 .expect("an unchanged table should accept the captured alias"),
1336 None
1337 );
1338 assert_eq!(
1339 metadata
1340 .with_detfd(4, |fd| (fd.open_file_id(), fd.is_cloexec()))
1341 .expect("captured alias should be installed"),
1342 (source_id, true)
1343 );
1344 }
1345
1346 #[test]
1347 fn failed_captured_fd_install_preserves_cleanup_obligations() {
1348 let owner = DetTid::from_raw(36);
1349 let mut original = FileMetadata::new(owner);
1350 original
1351 .add_fd(owner, 3, OFlag::empty(), FdType::Socket, None)
1352 .expect("source should be inserted");
1353 let source_id = original
1354 .with_detfd(3, |fd| fd.open_file_id())
1355 .expect("source should exist");
1356 let captured = original.capture_fd(3).expect("source should be captured");
1357 assert_eq!(
1358 original.remove_fd(3),
1359 None,
1360 "the capture defers final-OFD cleanup while the syscall is in flight"
1361 );
1362
1363 let mut replacement_table = FileMetadata::new(DetTid::from_raw(37));
1364 let failure = replacement_table
1365 .install_captured_fd(captured, 41, OFlag::O_CLOEXEC)
1366 .expect_err("a capture must not cross descriptor-table identity");
1367 assert_ne!(failure.expected_files_id, failure.actual_files_id);
1368 let cleanup = failure.into_cleanup();
1369 assert_eq!(cleanup.close_fd, 41);
1370 assert_eq!(cleanup.release_open_file, Some(source_id));
1371 }
1372
1373 #[test]
1374 fn equal_fd_dup_preserves_descriptor_flags() {
1375 let owner = DetTid::from_raw(20);
1376 let mut metadata = FileMetadata::new(owner);
1377 metadata
1378 .add_fd(owner, 3, OFlag::O_CLOEXEC, FdType::Regular, None)
1379 .expect("fd should be inserted");
1380
1381 assert_eq!(
1382 metadata
1383 .dup_fd(3, 3, OFlag::empty())
1384 .expect("equal-fd dup should validate the source"),
1385 None
1386 );
1387 assert!(
1388 metadata
1389 .with_detfd(3, |fd| fd.is_cloexec())
1390 .expect("fd should remain present"),
1391 "dup2(fd, fd) must not clear close-on-exec"
1392 );
1393 }
1394
1395 #[test]
1396 fn last_open_file_alias_survives_dup_and_fork() {
1397 let parent_tid = DetTid::from_raw(30);
1398 let child_tid = DetTid::from_raw(31);
1399 let mut parent = FileMetadata::new(parent_tid);
1400 parent
1401 .add_fd(parent_tid, 3, OFlag::empty(), FdType::Socket, None)
1402 .expect("socket should be inserted");
1403 let open_file_id = parent
1404 .with_detfd(3, |fd| fd.open_file_id())
1405 .expect("socket should exist");
1406 assert_eq!(
1407 parent
1408 .dup_fd(3, 4, OFlag::empty())
1409 .expect("dup should succeed"),
1410 None
1411 );
1412 assert_eq!(parent.remove_fd(3), None, "duplicate retains the OFD");
1413
1414 let mut child = parent.fork_for(child_tid);
1415 assert_eq!(parent.remove_fd(4), None, "forked child retains the OFD");
1416 assert_eq!(
1417 child.remove_fd(4),
1418 Some(open_file_id),
1419 "only the final alias releases the OFD"
1420 );
1421
1422 let mut replacement = FileMetadata::new(parent_tid);
1423 replacement
1424 .add_fd(parent_tid, 3, OFlag::empty(), FdType::Socket, None)
1425 .expect("source should be inserted");
1426 replacement
1427 .add_fd(parent_tid, 4, OFlag::empty(), FdType::Socket, None)
1428 .expect("target should be inserted");
1429 let target_id = replacement
1430 .with_detfd(4, |fd| fd.open_file_id())
1431 .expect("target should exist");
1432 assert_eq!(
1433 replacement
1434 .dup_fd(3, 4, OFlag::empty())
1435 .expect("dup replacement should succeed"),
1436 Some(target_id),
1437 "replacing the target must release its last OFD alias"
1438 );
1439 }
1440
1441 #[test]
1442 fn close_range_removes_selected_slots_and_releases_final_aliases() {
1443 let owner = DetTid::from_raw(35);
1444 let mut metadata = FileMetadata::new(owner);
1445 metadata
1446 .add_fd(owner, 3, OFlag::empty(), FdType::Regular, None)
1447 .expect("source should be inserted");
1448 metadata
1449 .dup_fd(3, 4, OFlag::empty())
1450 .expect("duplicate should be inserted");
1451 metadata
1452 .add_fd(owner, 100, OFlag::empty(), FdType::Regular, None)
1453 .expect("high fd should be inserted");
1454 let high_id = metadata
1455 .with_detfd(100, |fd| fd.open_file_id())
1456 .expect("high fd should exist");
1457
1458 assert_eq!(metadata.remove_fd_range(4, 100), [high_id]);
1459 assert!(metadata.with_detfd(3, |_| ()).is_ok());
1460 assert_eq!(metadata.with_detfd(4, |_| ()), Err(Errno::EBADF));
1461 assert_eq!(metadata.with_detfd(100, |_| ()), Err(Errno::EBADF));
1462 }
1463
1464 #[test]
1465 fn exec_reports_only_cloexec_open_files_with_no_other_aliases() {
1466 let owner = DetTid::from_raw(40);
1467 let child_tid = DetTid::from_raw(41);
1468 let mut metadata = FileMetadata::new(owner);
1469 metadata
1470 .add_fd(owner, 3, OFlag::O_CLOEXEC, FdType::Socket, None)
1471 .expect("socket should be inserted");
1472 let open_file_id = metadata
1473 .with_detfd(3, |fd| fd.open_file_id())
1474 .expect("socket should exist");
1475
1476 assert_eq!(metadata.open_files_closed_on_exec(false), [open_file_id]);
1477 assert!(
1478 metadata.open_files_closed_on_exec(true).is_empty(),
1479 "a shared descriptor table retains the original slot"
1480 );
1481
1482 let child = metadata.fork_for(child_tid);
1483 assert!(
1484 metadata.open_files_closed_on_exec(false).is_empty(),
1485 "a copied table retains an OFD alias"
1486 );
1487 drop(child);
1488
1489 metadata
1490 .dup_fd(3, 4, OFlag::empty())
1491 .expect("non-CLOEXEC alias should be created");
1492 assert!(
1493 metadata.open_files_closed_on_exec(false).is_empty(),
1494 "a non-CLOEXEC alias keeps the OFD live across exec"
1495 );
1496 }
1497}
1498
1499#[derive(Debug, Serialize, Deserialize, Clone, Default)]
1502pub struct ThreadStats {
1503 pub syscall_count: u64,
1505
1506 pub regs_sample_index: u64,
1516
1517 pub signal_count: u64,
1519
1520 pub timeslice_syscall_count: u64,
1522
1523 pub timeslice_signal_count: u64,
1525
1526 pub timeslice_count: u64,
1529
1530 pub last_recorded_slice: Option<u64>,
1533
1534 pub timeslice_stats: TimesliceStats,
1538
1539 pub timeslice_start_ns: Option<LogicalTime>,
1543}
1544
1545impl ThreadStats {
1546 pub fn new() -> Self {
1548 Default::default()
1549 }
1550
1551 pub fn count_syscall(&mut self) {
1554 self.syscall_count += 1;
1555 self.timeslice_syscall_count += 1;
1556 }
1557
1558 pub fn count_signal(&mut self) {
1560 self.signal_count += 1;
1561 self.timeslice_signal_count += 1;
1562 }
1563
1564 pub(crate) fn reset_timeslice(&mut self) {
1567 self.timeslice_syscall_count = 0;
1568 self.timeslice_signal_count = 0;
1569 self.timeslice_count += 1;
1570 }
1571
1572 pub fn close_final_timeslice(&mut self, now: LogicalTime) {
1577 if let Some(start) = self.timeslice_start_ns.take()
1578 && now >= start
1579 {
1580 self.timeslice_stats.record((now - start).as_nanos());
1581 }
1582 }
1583}
1584
1585#[derive(Debug, Clone, Serialize, Deserialize)]
1588pub struct PendingVfork {
1589 pub parent_dettid: DetTid,
1590 pub parent_detpid: DetPid,
1591 pub child_tid_addr: usize,
1592 pub flags: CloneFlags,
1593 pub exit_signal: libc::c_int,
1594 pub child_priority_entropy: Option<u64>,
1595}
1596
1597#[derive(Debug, Default, Clone, Copy, Serialize, Deserialize)]
1599pub(crate) struct ProcessCpuSnapshot {
1600 pub user: LogicalTime,
1601 pub system: LogicalTime,
1602 pub children_user: LogicalTime,
1603 pub children_system: LogicalTime,
1604}
1605
1606#[derive(Debug, Default, Clone, Serialize, Deserialize)]
1607pub(crate) struct ProcessCpuTime {
1608 snapshot: ProcessCpuSnapshot,
1609 exited_children: BTreeMap<DetPid, ProcessCpuSnapshot>,
1610}
1611
1612impl ProcessCpuTime {
1613 fn add_thread_delta(&mut self, user: LogicalTime, system: LogicalTime) {
1614 self.snapshot.user = self.snapshot.user + user;
1615 self.snapshot.system = self.snapshot.system + system;
1616 }
1617
1618 fn record_exited_child(&mut self, pid: DetPid, child: ProcessCpuSnapshot) {
1619 self.exited_children
1620 .entry(pid)
1621 .and_modify(|previous| {
1622 previous.user = previous.user.max(child.user);
1623 previous.system = previous.system.max(child.system);
1624 previous.children_user = previous.children_user.max(child.children_user);
1625 previous.children_system = previous.children_system.max(child.children_system);
1626 })
1627 .or_insert(child);
1628 }
1629
1630 fn reap_child(&mut self, pid: DetPid) {
1631 let Some(child) = self.exited_children.remove(&pid) else {
1632 return;
1633 };
1634 self.snapshot.children_user =
1635 self.snapshot.children_user + child.user + child.children_user;
1636 self.snapshot.children_system =
1637 self.snapshot.children_system + child.system + child.children_system;
1638 }
1639
1640 fn prepare_child(&mut self, pid: DetPid) {
1641 self.exited_children.remove(&pid);
1642 }
1643}
1644
1645#[derive(Debug, Default, Serialize, Deserialize)]
1653pub(crate) struct GuestClock {
1654 now: LogicalTime,
1655}
1656
1657#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1660pub(crate) struct RobustListWake {
1661 pub(crate) futex: FutexID,
1662}
1663
1664#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1667pub(crate) enum RobustListExit {
1668 ExitGroup,
1669 Signal(i32),
1670}
1671
1672#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1675pub(crate) struct PendingRobustListWakes {
1676 pub(crate) reason: RobustListExit,
1677 pub(crate) wakes: Vec<RobustListWake>,
1678}
1679
1680impl PendingRobustListWakes {
1681 fn matches_exit(&self, exit_signal: Option<i32>) -> bool {
1682 match self.reason {
1683 RobustListExit::ExitGroup => true,
1684 RobustListExit::Signal(expected) => exit_signal == Some(expected),
1685 }
1686 }
1687}
1688
1689#[derive(Debug, Clone, Default, Serialize, Deserialize)]
1693pub(crate) struct RobustListProcessState {
1694 heads: BTreeMap<DetTid, usize>,
1695 pending: BTreeMap<DetTid, PendingRobustListWakes>,
1696 matched_physical_exits: BTreeMap<DetTid, DetTime>,
1697}
1698
1699impl GuestClock {
1700 fn observe(&mut self, raw: LogicalTime) -> LogicalTime {
1701 self.now = self.now.max(raw);
1702 self.now
1703 }
1704}
1705
1706pub(crate) const DEFAULT_TIMER_SLACK_NS: u64 = 50_000;
1712
1713fn default_timer_slack_ns() -> u64 {
1714 DEFAULT_TIMER_SLACK_NS
1715}
1716
1717#[derive(Serialize, Deserialize, Clone)]
1719pub struct ThreadState<T> {
1720 pub dettid: DetTid,
1722 pub detpid: Option<DetTid>,
1724
1725 pub(crate) thread_start_entered: bool,
1728
1729 #[serde(default)]
1733 pub physical_tid: Option<i32>,
1734
1735 #[serde(default)]
1737 pub(crate) signal_task_identity: Option<reverie::SignalTaskIdentity>,
1738
1739 #[serde(default)]
1743 pub(crate) open_file_creator: Option<DetTid>,
1744
1745 pub mm_id: MmId,
1747
1748 pub(crate) memory_metadata: Arc<Mutex<MemoryMetadata>>,
1750
1751 pub pedigree: Pedigree,
1754
1755 pub stats: ThreadStats,
1757
1758 pub preemption_points: Option<ThreadHistoryIterator>,
1760
1761 pub interrupt_at: BTreeSet<u64>,
1763
1764 pub clone_flags: Option<CloneFlags>,
1772
1773 pub pending_vfork: Option<PendingVfork>,
1777
1778 pub file_metadata: Arc<Mutex<FileMetadata>>,
1781
1782 #[serde(default)]
1786 pub(crate) discover_live_file_metadata: bool,
1787
1788 #[serde(default = "default_timer_slack_ns")]
1795 pub(crate) timer_slack_ns: u64,
1796 #[serde(default = "default_timer_slack_ns")]
1799 pub(crate) default_timer_slack_ns: u64,
1800
1801 pub(crate) posix_timers: Arc<Mutex<PosixTimers>>,
1804
1805 pub(crate) resource_limits: Arc<Mutex<ResourceLimits>>,
1807
1808 pub(crate) process_cpu_time: Arc<Mutex<ProcessCpuTime>>,
1810
1811 #[serde(default)]
1813 pub(crate) guest_clock: Arc<Mutex<GuestClock>>,
1814
1815 pub(crate) parent_process_cpu_time: Option<Arc<Mutex<ProcessCpuTime>>>,
1817
1818 pub(crate) last_accounted_user_time: LogicalTime,
1820 pub(crate) last_accounted_system_time: LogicalTime,
1821
1822 pub(crate) thread_cpu_start_user_time: LogicalTime,
1826 pub(crate) thread_cpu_start_system_time: LogicalTime,
1827
1828 pub prng: Pcg64Mcg,
1830
1831 #[serde(default)]
1835 pub(crate) initialized_random_auxv: Option<crate::random::InitialImage>,
1836
1837 pub chaos_prng: Pcg64Mcg,
1839
1840 pub thread_logical_time: DetTime,
1842
1843 pub committed_clock_value: u64,
1845
1846 #[serde(default)]
1850 pub(crate) uncharged_bootstrap_syscalls: u32,
1851
1852 #[serde(default)]
1863 pub(crate) in_uncharged_bootstrap_syscall: bool,
1864
1865 pub record_or_replay: T,
1867
1868 pub end_of_timeslice: Option<LogicalTime>,
1877
1878 #[serde(default)]
1883 pub replay_rcb_end: Option<u64>,
1884
1885 #[serde(default = "chaos_epoch_sentinel")]
1892 pub chaos_epoch: u64,
1893
1894 #[serde(default)]
1898 pub chaos_slowdown_factor: RcbTimeMultiplier,
1899
1900 #[serde(default)]
1904 pub chaos_slowdown_active: bool,
1905
1906 #[serde(default)]
1910 pub pending_chaos_epochs: Vec<ChaosEpochTransition>,
1911
1912 pub max_timeslice_end: Option<LogicalTime>,
1915
1916 pub last_rcb_timer: Option<u64>,
1920
1921 #[serde(default)]
1923 pub last_rcb_timer_is_max: bool,
1924
1925 pub(crate) past_global_first_execve: bool,
1928
1929 #[serde(default)]
1938 pub(crate) robust_list_head: Option<usize>,
1939
1940 #[serde(default)]
1943 pub(crate) robust_list_process: Arc<Mutex<RobustListProcessState>>,
1944}
1945
1946impl<T> std::fmt::Debug for ThreadState<T> {
1949 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1950 f.debug_struct("ThreadState")
1951 .field("dettid", &self.dettid)
1952 .field("detpid", &self.detpid)
1953 .field("physical_tid", &self.physical_tid)
1954 .field("mm_id", &self.mm_id)
1955 .field("memory_metadata", &self.memory_metadata)
1956 .field("stats", &self.stats)
1957 .field("clone_flags", &self.clone_flags)
1958 .field("file_metadata", &self.file_metadata)
1959 .field("posix_timers", &self.posix_timers)
1960 .field("resource_limits", &self.resource_limits)
1961 .field("process_cpu_time", &self.process_cpu_time)
1962 .field("prng", &self.prng)
1963 .field("chaos_prng", &self.chaos_prng)
1964 .field("thread_logical_time", &self.thread_logical_time)
1965 .field("committed_clock_value", &self.committed_clock_value)
1966 .field("end_of_timeslice", &self.end_of_timeslice)
1967 .field("replay_rcb_end", &self.replay_rcb_end)
1968 .field("chaos_epoch", &self.chaos_epoch)
1969 .field("chaos_slowdown_factor", &self.chaos_slowdown_factor)
1970 .field("chaos_slowdown_active", &self.chaos_slowdown_active)
1971 .field("max_timeslice_end", &self.max_timeslice_end)
1972 .field("last_rcb_timer", &self.last_rcb_timer)
1973 .field("last_rcb_timer_is_max", &self.last_rcb_timer_is_max)
1974 .finish()
1975 }
1976}
1977
1978impl<T> Default for ThreadState<T> {
1979 fn default() -> Self {
1980 unreachable!()
1981 }
1982}
1983
1984pub(crate) fn chaos_epoch_sentinel() -> u64 {
1990 u64::MAX
1991}
1992
1993impl<T> AsRef<T> for ThreadState<T> {
1994 fn as_ref(&self) -> &T {
1995 &self.record_or_replay
1996 }
1997}
1998
1999impl<T> AsMut<T> for ThreadState<T> {
2000 fn as_mut(&mut self) -> &mut T {
2001 &mut self.record_or_replay
2002 }
2003}
2004
2005pub(crate) fn chaos_per_thread_slowdown_factor(
2034 sched_seed: u64,
2035 dettid: DetTid,
2036 epoch: u64,
2037 max_factor: f64,
2038) -> RcbTimeMultiplier {
2039 if max_factor <= 1.0 {
2042 return RcbTimeMultiplier::ONE;
2043 }
2044 const SLOWDOWN_SALT: u64 = 0x736c_6f77_646f_776e; const EPOCH_GOLDEN: u64 = 0xbf58_476d_1ce4_e5b9;
2054 let mixed = sched_seed
2055 ^ SLOWDOWN_SALT
2056 ^ ((dettid.as_raw() as u32 as u64).wrapping_mul(0x9e37_79b9_7f4a_7c15))
2057 ^ epoch.wrapping_mul(EPOCH_GOLDEN);
2058 let mut prng = Pcg64Mcg::seed_from_u64(mixed);
2059 let u: f64 = prng.random::<f64>();
2062 let exponent = 2.0 * u - 1.0;
2063 RcbTimeMultiplier::from_f64(max_factor.powf(exponent))
2064}
2065
2066impl<T> ThreadState<T> {
2067 pub(crate) fn observe_guest_clock(&self, raw: LogicalTime) -> LogicalTime {
2068 self.guest_clock
2069 .lock()
2070 .expect("guest clock mutex poisoned")
2071 .observe(raw)
2072 }
2073
2074 pub fn reseed_child_rngs(&mut self, parent: &Self, entropy: u128) {
2079 self.prng = thread_rng_from_parent_entropy("USER RAND", &parent.prng, entropy);
2080 self.chaos_prng = thread_rng_from_parent_entropy("CHAOSRAND", &parent.chaos_prng, entropy);
2081 }
2082
2083 pub(crate) fn recover_process_mm_id(&mut self, detpid: DetPid) -> bool {
2088 if self.dettid == detpid || self.mm_id != MmId::initial(self.dettid) {
2089 return false;
2090 }
2091
2092 self.mm_id = MmId::initial(detpid);
2093 true
2094 }
2095
2096 pub(crate) fn account_process_cpu_time(&mut self) {
2097 let user = self.thread_logical_time.user_cpu_time();
2098 let system = self.thread_logical_time.system_cpu_time();
2099 let user_delta = user - self.last_accounted_user_time;
2100 let system_delta = system - self.last_accounted_system_time;
2101 self.process_cpu_time
2102 .lock()
2103 .expect("process CPU time mutex poisoned")
2104 .add_thread_delta(user_delta, system_delta);
2105 self.last_accounted_user_time = user;
2106 self.last_accounted_system_time = system;
2107 }
2108
2109 pub(crate) fn charge_syscall_time(
2124 &mut self,
2125 in_backend_runtime_bootstrap: bool,
2126 sysno: Sysno,
2127 ) -> bool {
2128 if !in_backend_runtime_bootstrap {
2129 if self.uncharged_bootstrap_syscalls != 0 {
2130 debug!(
2131 "[dtid {}] backend runtime bootstrap window left {} syscalls uncharged",
2132 self.dettid, self.uncharged_bootstrap_syscalls
2133 );
2134 self.uncharged_bootstrap_syscalls = 0;
2135 }
2136 return true;
2137 }
2138 if crate::syscall_time::observes_virtual_time(sysno) {
2139 return true;
2140 }
2141 if self.uncharged_bootstrap_syscalls < crate::syscall_time::MAX_UNCHARGED_BOOTSTRAP_SYSCALLS
2142 {
2143 self.uncharged_bootstrap_syscalls += 1;
2144 if self.uncharged_bootstrap_syscalls
2145 == crate::syscall_time::MAX_UNCHARGED_BOOTSTRAP_SYSCALLS
2146 {
2147 info!(
2149 "[dtid {}] backend runtime bootstrap window reached its cap of {} uncharged syscalls; every further syscall in this window is charged",
2150 self.dettid,
2151 crate::syscall_time::MAX_UNCHARGED_BOOTSTRAP_SYSCALLS
2152 );
2153 }
2154 return false;
2155 }
2156 true
2157 }
2158
2159 pub(crate) fn process_cpu_time(&mut self) -> ProcessCpuSnapshot {
2160 self.account_process_cpu_time();
2161 self.process_cpu_time
2162 .lock()
2163 .expect("process CPU time mutex poisoned")
2164 .snapshot
2165 }
2166
2167 pub(crate) fn thread_cpu_time(&self) -> (LogicalTime, LogicalTime) {
2175 (
2176 self.thread_logical_time.user_cpu_time() - self.thread_cpu_start_user_time,
2177 self.thread_logical_time.system_cpu_time() - self.thread_cpu_start_system_time,
2178 )
2179 }
2180
2181 pub(crate) fn record_exited_child_process_cpu_time(&mut self, pid: DetPid) {
2182 self.account_process_cpu_time();
2183 let Some(parent) = &self.parent_process_cpu_time else {
2184 return;
2185 };
2186 let child = self
2187 .process_cpu_time
2188 .lock()
2189 .expect("process CPU time mutex poisoned")
2190 .snapshot;
2191 parent
2192 .lock()
2193 .expect("parent process CPU time mutex poisoned")
2194 .record_exited_child(pid, child);
2195 }
2196
2197 pub(crate) fn has_exited_child_process_cpu_time(&self, pid: DetPid) -> bool {
2198 self.process_cpu_time
2199 .lock()
2200 .expect("process CPU time mutex poisoned")
2201 .exited_children
2202 .contains_key(&pid)
2203 }
2204
2205 pub(crate) fn reap_child_process_cpu_time(&mut self, pid: DetPid) {
2206 self.account_process_cpu_time();
2207 self.process_cpu_time
2208 .lock()
2209 .expect("process CPU time mutex poisoned")
2210 .reap_child(pid);
2211 }
2212
2213 pub(crate) fn prepare_child_process_cpu_time(&self, pid: DetPid) {
2214 self.process_cpu_time
2215 .lock()
2216 .expect("process CPU time mutex poisoned")
2217 .prepare_child(pid);
2218 }
2219
2220 pub fn new(pid: DetPid, cfg: &Config, record_or_replay: T) -> Self {
2223 detlog!(
2224 "USER RAND: seeding PRNG for root thread with seed {}",
2225 cfg.rng_seed()
2226 );
2227 detlog!(
2228 "CHAOSRAND: seeding chaos scheduler with seed {}",
2229 cfg.sched_seed()
2230 );
2231 let thread_logical_time = DetTime::new(cfg);
2232 let last_accounted_user_time = thread_logical_time.user_cpu_time();
2233 let last_accounted_system_time = thread_logical_time.system_cpu_time();
2234 let file_metadata = if cfg.discover_live_file_metadata {
2235 let mut metadata = FileMetadata::new(pid);
2236 for fd in 0..=2 {
2237 metadata
2238 .discover_fd_from_current_process(pid, fd)
2239 .expect("SaBRe guest stdio must be open");
2240 }
2241 metadata
2242 } else {
2243 FileMetadata::new(pid).setup_stdio(pid.into(), pid)
2244 };
2245 ThreadState {
2246 dettid: pid,
2247 detpid: None, thread_start_entered: false,
2249 physical_tid: None,
2250 signal_task_identity: None,
2251 open_file_creator: None,
2252 mm_id: MmId::initial(pid),
2253 memory_metadata: Arc::new(Mutex::new(MemoryMetadata::new())),
2254 pedigree: Pedigree::new(), stats: ThreadStats::new(),
2256 file_metadata: Arc::new(Mutex::new(file_metadata)),
2257 discover_live_file_metadata: cfg.discover_live_file_metadata,
2258 timer_slack_ns: DEFAULT_TIMER_SLACK_NS,
2259 default_timer_slack_ns: DEFAULT_TIMER_SLACK_NS,
2260 posix_timers: Arc::new(Mutex::new(PosixTimers::default())),
2261 resource_limits: Arc::new(Mutex::new(ResourceLimits::default())),
2262 process_cpu_time: Arc::new(Mutex::new(ProcessCpuTime::default())),
2263 guest_clock: Arc::new(Mutex::new(GuestClock::default())),
2264 parent_process_cpu_time: None,
2265 last_accounted_user_time,
2266 last_accounted_system_time,
2267 thread_cpu_start_user_time: last_accounted_user_time,
2268 thread_cpu_start_system_time: last_accounted_system_time,
2269 clone_flags: None,
2270 pending_vfork: None,
2271 prng: crate::random::root_prng(cfg.rng_seed()),
2273 initialized_random_auxv: None,
2274 chaos_prng: Pcg64Mcg::seed_from_u64(cfg.sched_seed()),
2275 thread_logical_time,
2276 committed_clock_value: 0,
2277 uncharged_bootstrap_syscalls: 0,
2278 in_uncharged_bootstrap_syscall: false,
2279 end_of_timeslice: None, replay_rcb_end: None,
2281 chaos_epoch: chaos_epoch_sentinel(),
2284 chaos_slowdown_factor: RcbTimeMultiplier::ONE,
2287 chaos_slowdown_active: false,
2288 pending_chaos_epochs: Vec::new(),
2289 max_timeslice_end: None,
2290 last_rcb_timer: None,
2291 last_rcb_timer_is_max: false,
2292 record_or_replay,
2293 preemption_points: None,
2294 past_global_first_execve: false,
2295 interrupt_at: cfg.interrupts_for_thread(pid),
2296 robust_list_head: None,
2297 robust_list_process: Arc::new(Mutex::new(RobustListProcessState::default())),
2298 }
2299 }
2300
2301 pub fn apply_initial_random_state(
2304 &mut self,
2305 bytes: &[u8],
2306 config: &Config,
2307 image: crate::random::InitialImage,
2308 ) -> Result<(), Errno> {
2309 if self.initialized_random_auxv.is_some()
2310 || self.past_global_first_execve
2311 || self.dettid.as_raw() != image.pid
2312 || self.clone_flags.is_some()
2313 || self.pedigree.raw() != Pedigree::new().raw()
2314 {
2315 return Err(Errno::EPROTO);
2316 }
2317 let prng = crate::random::decode_initial_state(bytes, config, image)?;
2318 self.prng = prng;
2319 self.initialized_random_auxv = Some(image);
2320 Ok(())
2321 }
2322
2323 pub(crate) fn complete_initial_random_auxv(
2324 &mut self,
2325 pointer: Option<usize>,
2326 ) -> Result<bool, Errno> {
2327 let Some(image) = self.initialized_random_auxv else {
2328 return Ok(false);
2329 };
2330 if pointer != Some(image.at_random) || self.dettid.as_raw() != image.pid {
2331 return Err(Errno::EPROTO);
2332 }
2333 self.initialized_random_auxv = None;
2334 Ok(true)
2335 }
2336
2337 pub(crate) fn record_robust_list_head(&mut self, head: Option<usize>) {
2338 self.robust_list_head = head;
2339 let mut process = self
2340 .robust_list_process
2341 .lock()
2342 .expect("robust-list process state mutex poisoned");
2343 if let Some(head) = head {
2344 process.heads.insert(self.dettid, head);
2345 } else {
2346 process.heads.remove(&self.dettid);
2347 }
2348 }
2349
2350 pub(crate) fn robust_list_heads(&self) -> Vec<(DetTid, usize)> {
2351 self.robust_list_process
2352 .lock()
2353 .expect("robust-list process state mutex poisoned")
2354 .heads
2355 .iter()
2356 .map(|(&tid, &head)| (tid, head))
2357 .collect()
2358 }
2359
2360 pub(crate) fn stage_robust_list_wakes(
2361 &self,
2362 reason: RobustListExit,
2363 wakes: Vec<(DetTid, Vec<RobustListWake>)>,
2364 ) {
2365 let mut process = self
2366 .robust_list_process
2367 .lock()
2368 .expect("robust-list process state mutex poisoned");
2369 process.pending = wakes
2370 .into_iter()
2371 .map(|(owner, wakes)| (owner, PendingRobustListWakes { reason, wakes }))
2372 .collect();
2373 process.matched_physical_exits.clear();
2374 }
2375
2376 pub(crate) fn has_matching_robust_list_exit(&self, exit_signal: Option<i32>) -> bool {
2377 self.robust_list_process
2378 .lock()
2379 .expect("robust-list process state mutex poisoned")
2380 .pending
2381 .get(&self.dettid)
2382 .is_some_and(|pending| pending.matches_exit(exit_signal))
2383 }
2384
2385 pub(crate) fn take_robust_list_wakes_after_exit(
2386 &self,
2387 exit_signal: Option<i32>,
2388 exit_time: DetTime,
2389 ) -> Option<(DetTime, Vec<(DetTid, RobustListWake)>)> {
2390 let mut process = self
2391 .robust_list_process
2392 .lock()
2393 .expect("robust-list process state mutex poisoned");
2394 process.heads.remove(&self.dettid);
2395 let pending = process.pending.get(&self.dettid)?;
2396 if !pending.matches_exit(exit_signal) {
2397 return None;
2398 }
2399
2400 process
2401 .matched_physical_exits
2402 .insert(self.dettid, exit_time);
2403 if process
2404 .pending
2405 .keys()
2406 .any(|owner| !process.matched_physical_exits.contains_key(owner))
2407 {
2408 return None;
2409 }
2410
2411 let request_time = process
2412 .matched_physical_exits
2413 .values()
2414 .max()
2415 .expect("a completed robust-list exit group has a physical-exit time")
2416 .clone();
2417 process.matched_physical_exits.clear();
2418 let pending = std::mem::take(&mut process.pending);
2419 let mut ready = Vec::new();
2420 for (owner, mut pending) in pending {
2421 pending
2422 .wakes
2423 .sort_by_key(|wake| format!("{:?}", wake.futex));
2424 ready.extend(pending.wakes.into_iter().map(|wake| (owner, wake)));
2425 }
2426 Some((request_time, ready))
2427 }
2428
2429 pub(crate) fn take_robust_list_for_exec(&mut self) -> Option<usize> {
2438 let previous = self.robust_list_head.take();
2439 self.robust_list_process
2440 .lock()
2441 .expect("robust-list process state mutex poisoned")
2442 .heads
2443 .remove(&self.dettid);
2444 previous
2445 }
2446
2447 pub(crate) fn restore_robust_list_after_failed_exec(&mut self, previous: Option<usize>) {
2450 self.record_robust_list_head(previous);
2451 }
2452
2453 pub(crate) fn futex_id(&self, address: usize, is_private: bool) -> FutexID {
2455 if is_private {
2456 FutexID::private(self.mm_id, address)
2457 } else {
2458 self.memory_metadata
2459 .lock()
2460 .expect("memory metadata mutex poisoned")
2461 .futex_id(self.mm_id, address)
2462 }
2463 }
2464
2465 pub(crate) fn map_shared_anonymous(&self, start: usize, len: usize) {
2467 self.memory_metadata
2468 .lock()
2469 .expect("memory metadata mutex poisoned")
2470 .map_anonymous(self.mm_id, start, len);
2471 }
2472
2473 pub(crate) fn map_shared_object(
2475 &self,
2476 start: usize,
2477 len: usize,
2478 object: SharedMemoryObjectId,
2479 object_offset: u64,
2480 ) {
2481 self.memory_metadata
2482 .lock()
2483 .expect("memory metadata mutex poisoned")
2484 .map_object(start, len, object, object_offset);
2485 }
2486
2487 pub(crate) fn unmap_memory(&self, start: usize, len: usize) {
2489 self.memory_metadata
2490 .lock()
2491 .expect("memory metadata mutex poisoned")
2492 .unmap(start, len);
2493 }
2494
2495 pub(crate) fn remap_memory(
2497 &self,
2498 old_start: usize,
2499 old_len: usize,
2500 new_start: usize,
2501 new_len: usize,
2502 ) {
2503 self.memory_metadata
2504 .lock()
2505 .expect("memory metadata mutex poisoned")
2506 .remap(old_start, old_len, new_start, new_len);
2507 }
2508
2509 pub fn mk_request(&self, rid: ResourceID, perm: Permission) -> Resources {
2511 let mut resources = HashMap::new();
2512 resources.insert(rid, perm);
2513 Resources {
2514 tid: self.dettid,
2515 resources,
2516 poll_attempt: 0,
2517 fyi: String::new(),
2518 signal_interrupt_errno: None,
2519 backend_runtime_bootstrap: false,
2520 }
2521 }
2522
2523 pub fn chaos_prng_next_u64(&mut self, msg: &str) -> u64 {
2525 let r = self.chaos_prng.next_u64();
2526 detlog!("[dtid {}] CHAOSRAND({}): u64 => {}", self.dettid, msg, r);
2527 r
2528 }
2529
2530 fn metadata(&self) -> MutexGuard<'_, FileMetadata> {
2532 self.file_metadata.lock().unwrap()
2533 }
2534
2535 pub fn add_fd(
2550 &self,
2551 fd: RawFd,
2552 flags: OFlag,
2553 ty: FdType,
2554 stat: Option<DetStat>,
2555 ) -> Result<(), Errno> {
2556 self.metadata().add_fd(
2557 self.open_file_creator.unwrap_or(self.dettid),
2558 fd,
2559 flags,
2560 ty,
2561 stat,
2562 )
2563 }
2564
2565 pub fn set_open_file_creator(&mut self, creator: DetTid) {
2569 self.open_file_creator = Some(creator);
2570 }
2571
2572 pub fn with_detfd<F, U>(&self, fd: RawFd, f: F) -> Result<U, Errno>
2575 where
2576 F: FnMut(&mut DetFd) -> U,
2577 {
2578 let mut metadata = self.metadata();
2579 if self.discover_live_file_metadata {
2580 metadata.discover_fd_from_current_process(self.dettid, fd)?;
2581 }
2582 metadata.with_detfd(fd, f)
2583 }
2584
2585 pub(crate) fn scheduler_managed_pipe_fds(&self) -> Vec<RawFd> {
2595 let mut fds: Vec<RawFd> = self
2596 .metadata()
2597 .file_handles
2598 .iter()
2599 .filter(|(_, detfd)| {
2600 detfd.ty() == FdType::Pipe
2601 && detfd.physically_nonblocking()
2602 && !detfd.is_nonblocking()
2603 })
2604 .map(|(&fd, _)| fd)
2605 .collect();
2606 fds.sort_unstable();
2607 fds
2608 }
2609
2610 pub(crate) fn count_open_files_at_paths(&self, paths: &[&Path]) -> usize {
2611 self.metadata().count_open_files_at_paths(paths)
2612 }
2613
2614 pub(crate) fn has_loopback_peer(&self) -> bool {
2616 self.metadata().has_loopback_peer()
2617 }
2618
2619 pub fn has_unsafe_vfork_flock_state(&self) -> bool {
2621 self.metadata().has_unsafe_vfork_flock_state()
2622 }
2623
2624 pub fn forget_flock_modes(&self) {
2628 self.metadata().forget_flock_modes();
2629 }
2630
2631 pub fn remove_fd(&self, fd: RawFd) -> Option<OpenFileId> {
2633 self.metadata().remove_fd(fd)
2634 }
2635
2636 pub(crate) fn remove_fd_range(&self, first: u32, last: u32) -> Vec<OpenFileId> {
2638 self.metadata().remove_fd_range(first, last)
2639 }
2640
2641 pub fn dup_fd(
2643 &mut self,
2644 oldfd: RawFd,
2645 newfd: RawFd,
2646 flags: OFlag,
2647 ) -> Result<Option<OpenFileId>, Errno> {
2648 let mut metadata = self.metadata();
2649 if self.discover_live_file_metadata {
2650 metadata.discover_fd_from_current_process(self.dettid, oldfd)?;
2651 }
2652 metadata.dup_fd(oldfd, newfd, flags)
2653 }
2654
2655 pub(crate) fn capture_pidfd_getfd_source(
2661 &self,
2662 pidfd: RawFd,
2663 targetfd: RawFd,
2664 current_tgid: DetPid,
2665 current_tid: DetTid,
2666 ) -> Result<CapturedDetFd, Errno> {
2667 let mut metadata = self.metadata();
2668 if self.discover_live_file_metadata {
2669 metadata.discover_fd_from_current_process(self.dettid, pidfd)?;
2670 }
2671 let (is_pidfd, target) = metadata.with_detfd(pidfd, |detfd| {
2672 (matches!(detfd.ty(), FdType::Pidfd), detfd.pidfd_target())
2673 })?;
2674 if !is_pidfd {
2675 return Err(Errno::EBADF);
2676 }
2677 if !pidfd_getfd_targets_calling_task(target, current_tgid, current_tid) {
2678 return Err(Errno::EOPNOTSUPP);
2679 }
2680 if self.discover_live_file_metadata {
2681 metadata.discover_fd_from_current_process(self.dettid, targetfd)?;
2682 }
2683 metadata.capture_fd(targetfd)
2684 }
2685
2686 pub(crate) fn install_captured_fd(
2689 &mut self,
2690 captured: CapturedDetFd,
2691 newfd: RawFd,
2692 flags: OFlag,
2693 ) -> Result<Option<OpenFileId>, CapturedDetFdInstallError> {
2694 self.metadata().install_captured_fd(captured, newfd, flags)
2695 }
2696
2697 pub(crate) fn abandon_captured_fd(&self, captured: CapturedDetFd) -> Option<OpenFileId> {
2699 self.metadata().abandon_captured_fd(captured)
2700 }
2701
2702 pub fn thread_prng(&mut self) -> &mut Pcg64Mcg {
2705 &mut self.prng
2706 }
2707
2708 pub(crate) fn timeslice_expired(&self) -> bool {
2710 if let Some(replay_rcb_end) = self.replay_rcb_end {
2711 return self.committed_clock_value >= replay_rcb_end;
2712 }
2713 let current_time = self.thread_logical_time.as_nanos();
2714 self.end_of_timeslice
2715 .is_some_and(|end_of_timeslice| current_time >= end_of_timeslice)
2716 }
2717
2718 pub(crate) fn rcb_time_multiplier(&self) -> RcbTimeMultiplier {
2722 if self.chaos_slowdown_active {
2723 self.chaos_slowdown_factor
2724 } else {
2725 RcbTimeMultiplier::ONE
2726 }
2727 }
2728
2729 pub(crate) fn take_pending_chaos_epochs(&mut self) -> Vec<ChaosEpochTransition> {
2732 std::mem::take(&mut self.pending_chaos_epochs)
2733 }
2734
2735 fn install_chaos_epoch(&mut self, transition: ChaosEpochTransition, record: bool) {
2738 self.chaos_epoch = transition.epoch;
2739 self.chaos_slowdown_factor = transition.factor;
2740 self.chaos_slowdown_active = true;
2741 if record {
2742 self.pending_chaos_epochs.push(transition);
2743 }
2744 detlog!(
2745 "[dtid {}] CHAOSEPOCH => epoch = {}, factor = {}, logical_time = {}",
2746 self.dettid,
2747 transition.epoch,
2748 transition.factor.as_f64(),
2749 transition.logical_time
2750 );
2751 }
2752
2753 pub fn next_timeslice(&mut self, cfg: &Config) -> Option<Priority> {
2763 let logical_timeslice = cfg.target_timeslice.or(cfg.max_timeslice);
2764 if let Some(timeout_ns) = logical_timeslice {
2765 let current_ns = self.thread_logical_time.as_nanos();
2766
2767 let replay_has_epochs = self
2773 .preemption_points
2774 .as_ref()
2775 .is_some_and(ThreadHistoryIterator::has_chaos_epochs);
2776 if replay_has_epochs {
2777 let transition_time = self.end_of_timeslice.unwrap_or(current_ns);
2784 let transition = self
2785 .preemption_points
2786 .as_mut()
2787 .and_then(|history| history.advance_chaos_epoch(transition_time));
2788 if let Some(transition) = transition {
2789 self.install_chaos_epoch(transition, false);
2790 }
2791 } else if cfg.chaos && cfg.chaos_per_thread_slowdown {
2792 let elapsed_ns = self.thread_logical_time.without_starting().as_nanos();
2793 let epoch = elapsed_ns
2794 .checked_div(cfg.chaos_epoch_length_ns)
2795 .unwrap_or(0);
2796 let factor = chaos_per_thread_slowdown_factor(
2797 cfg.sched_seed(),
2798 self.dettid,
2799 epoch,
2800 cfg.chaos_slowdown_max_factor,
2801 );
2802 if epoch != self.chaos_epoch {
2803 self.install_chaos_epoch(
2804 ChaosEpochTransition {
2805 logical_time: current_ns,
2806 epoch,
2807 factor,
2808 },
2809 true,
2810 );
2811 } else {
2812 self.chaos_slowdown_factor = factor;
2813 self.chaos_slowdown_active = true;
2814 }
2815 } else {
2816 self.chaos_slowdown_factor = RcbTimeMultiplier::ONE;
2817 self.chaos_slowdown_active = false;
2818 self.pending_chaos_epochs.clear();
2819 }
2820
2821 let mut result = None;
2822 self.replay_rcb_end = None;
2823 let replay_controls_deadline =
2824 self.preemption_points.is_some() || cfg.replay_schedule_from.is_some();
2825
2826 if let Some(thi) = &mut self.preemption_points {
2828 if self.stats.last_recorded_slice.is_none() {
2829 if let Some((end_time, prio, replay_rcb_end)) = thi.next_with_rcbs() {
2831 debug!(
2832 "[dtid {}] next timeslice (T{}), set by recording to {:?} (current {}), priority {}",
2833 self.dettid,
2834 self.stats.timeslice_count + 1,
2835 end_time,
2836 current_ns,
2837 prio
2838 );
2839 let current_rcbs = self.committed_clock_value;
2840 if let Some(target_rcbs) = replay_rcb_end
2841 && target_rcbs < current_rcbs
2842 {
2843 panic!(
2844 "Cannot set RCB end of timeslice to {} for thread {}, when current RCB count is already {}.",
2845 target_rcbs, self.dettid, current_rcbs
2846 )
2847 }
2848 let exact_deadline_is_now =
2849 replay_rcb_end.is_some_and(|target| target == current_rcbs);
2850 if end_time <= current_ns && !exact_deadline_is_now {
2851 panic!(
2858 "Cannot set end of timeslice to {} for thread {}, when current thread logical time is already {}. \
2859 The replayed preemption record's timeslice ends are absolute virtual times; if it was \
2860 recorded under a different virtual-time epoch, replay it with that --epoch \
2861 (https://github.com/rrnewton/hermit/issues/3413).",
2862 end_time, self.dettid, current_ns
2863 )
2864 }
2865 self.end_of_timeslice = Some(end_time);
2866 self.replay_rcb_end = replay_rcb_end;
2867 result = Some(prio);
2868 } else {
2869 let max = LogicalTime::MAX;
2870 let prio = thi.final_priority();
2871 debug!(
2872 "[dtid {}] next timeslice (T{}) final slice after recorded preemption points... setting end_of_timeslice to max {}, final priority {}",
2873 self.dettid,
2874 self.stats.timeslice_count + 1,
2875 max,
2876 prio
2877 );
2878 self.stats.last_recorded_slice = Some(self.stats.timeslice_count);
2879 self.end_of_timeslice = Some(max);
2880 result = Some(prio)
2881 }
2882 } else {
2883 tracing::warn!(
2884 "[dtid {}] next timeslice: timer expired beyond the last recorded preemption. Not handled yet.",
2885 self.dettid
2886 );
2887 self.end_of_timeslice = Some(LogicalTime::MAX);
2888 result = Some(thi.final_priority())
2889 }
2890 } else if !cfg.chaos {
2891 if cfg.replay_schedule_from.is_some() {
2892 if cfg.no_rcb_time {
2893 let max_timeslice = cfg
2894 .max_timeslice
2895 .expect("schedule replay with PMU requires a maximum");
2896 self.end_of_timeslice =
2897 Some(current_ns + Duration::from_nanos(u64::from(max_timeslice)));
2898 } else {
2899 debug!(
2901 "[dtid {}] next timeslice (T{}), in replay mode setting timeslice to max (current time {})",
2902 self.dettid,
2903 self.stats.timeslice_count + 1,
2904 current_ns
2905 );
2906 self.end_of_timeslice = Some(LogicalTime::MAX);
2907 }
2908 } else {
2909 self.end_of_timeslice =
2913 Some(current_ns + Duration::from_nanos(u64::from(timeout_ns)));
2914 debug!(
2915 "[dtid {}] next timeslice (T{}), end of slice set to {} (current {})",
2916 self.dettid,
2917 self.stats.timeslice_count + 1,
2918 self.end_of_timeslice.unwrap(),
2919 current_ns,
2920 );
2921 }
2922 } else {
2923 let slowdown = self.rcb_time_multiplier().as_f64();
2929 let nanos_per_rcb = NANOS_PER_RCB * cfg.clock_multiplier.unwrap_or(1.0) * slowdown;
2930 let target_timeout_rcbs = u64::from(timeout_ns) as f64 / nanos_per_rcb;
2931 if self.chaos_slowdown_active {
2932 detlog!(
2933 "[dtid {}] CHAOSSLOWDOWN => factor = {}, virtual ns/rcb = {}",
2934 self.dettid,
2935 slowdown,
2936 nanos_per_rcb
2937 );
2938 }
2939 let next_rcbs: u64 = if cfg.chaos {
2940 let lambda = 1.0 / target_timeout_rcbs;
2942 let exp = Exp::new(lambda).unwrap();
2943 let rcbs = 1 + exp.sample(&mut self.chaos_prng) as u64;
2945 detlog!("[dtid {}] CHAOSRAND => next_rcbs = {}", self.dettid, rcbs);
2946 rcbs
2947 } else {
2948 target_timeout_rcbs as u64
2949 };
2950 assert!(next_rcbs > 0);
2951 self.last_rcb_timer = None;
2952 self.end_of_timeslice = Some(
2953 current_ns
2954 + Duration::from_nanos((next_rcbs as f64 * nanos_per_rcb).ceil() as u64),
2955 );
2956 debug!(
2957 "[dtid {}] next timeslice (T{}) chosen as {} rcbs, end of slice = {} (current {})",
2958 self.dettid,
2959 self.stats.timeslice_count + 1,
2960 next_rcbs,
2961 self.end_of_timeslice.unwrap(),
2962 current_ns
2963 );
2964 }
2965
2966 let configured_max_end = cfg
2967 .max_timeslice
2968 .map(|max_timeslice| current_ns + Duration::from_nanos(u64::from(max_timeslice)));
2969 self.max_timeslice_end = if replay_controls_deadline {
2970 if cfg.max_timeslice.is_some() {
2971 self.end_of_timeslice
2976 .filter(|end| *end != LogicalTime::MAX)
2977 .or(configured_max_end)
2978 } else {
2979 None
2980 }
2981 } else if cfg.target_timeslice.is_none() {
2982 match (self.end_of_timeslice, configured_max_end) {
2983 (Some(logical_end), Some(configured_end)) => {
2984 Some(logical_end.min(configured_end))
2985 }
2986 (_, configured_end) => configured_end,
2987 }
2988 } else {
2989 configured_max_end
2990 };
2991
2992 if let (Some(target_end), Some(max_end)) =
2993 (self.end_of_timeslice, self.max_timeslice_end)
2994 && target_end > max_end
2995 {
2996 self.end_of_timeslice = Some(max_end);
2997 }
2998
2999 self.last_rcb_timer = None;
3000 self.last_rcb_timer_is_max = false;
3001 self.reset_timeslice_stats(current_ns);
3002 result
3003 } else {
3004 self.end_of_timeslice = None;
3005 self.replay_rcb_end = None;
3006 self.max_timeslice_end = None;
3007 self.last_rcb_timer = None;
3008 self.last_rcb_timer_is_max = false;
3009 None
3010 }
3011 }
3012
3013 pub fn reset_timeslice_for_explicit_yield(&mut self) {
3017 let current_ns = self.thread_logical_time.as_nanos();
3018 self.reset_timeslice_stats(current_ns);
3019 }
3020
3021 fn reset_timeslice_stats(&mut self, current_ns: LogicalTime) {
3022 if let Some(start) = self.stats.timeslice_start_ns
3023 && current_ns >= start
3024 {
3025 self.stats
3026 .timeslice_stats
3027 .record((current_ns - start).as_nanos());
3028 }
3029 self.stats.timeslice_start_ns = Some(current_ns);
3030 self.stats.reset_timeslice();
3031 }
3032
3033 pub fn guest_past_first_execve(&self) -> bool {
3039 self.past_global_first_execve
3040 }
3041}
3042
3043#[cfg(test)]
3044mod timeslice_tests {
3045 use std::num::NonZeroU64;
3046
3047 use super::*;
3048 use crate::DEFAULT_PRIORITY;
3049 use crate::preemptions::ThreadHistory;
3050
3051 #[test]
3057 fn a_successful_exec_clears_the_robust_list_head() {
3058 let mut state = ThreadState::<()>::new(DetTid::from_raw(3), &Config::default(), ());
3059 assert_eq!(state.robust_list_head, None, "a fresh thread has no list");
3060
3061 state.record_robust_list_head(Some(0x7ffff7bff920));
3062 assert_eq!(
3063 state.take_robust_list_for_exec(),
3064 Some(0x7ffff7bff920),
3065 "the previous head is handed back for the failure path"
3066 );
3067 assert_eq!(
3068 state.robust_list_head, None,
3069 "the candidate exec image starts with no robust list, as after copy_process"
3070 );
3071 assert!(
3072 state.robust_list_heads().is_empty(),
3073 "the old address-space registration must also leave the shared index"
3074 );
3075 }
3076
3077 #[test]
3080 fn a_failed_exec_restores_the_robust_list_head() {
3081 let mut state = ThreadState::<()>::new(DetTid::from_raw(3), &Config::default(), ());
3082 state.record_robust_list_head(Some(0x404100));
3083
3084 let saved = state.take_robust_list_for_exec();
3085 state.restore_robust_list_after_failed_exec(saved);
3086
3087 assert_eq!(state.robust_list_head, Some(0x404100));
3088 assert_eq!(
3089 state.robust_list_heads(),
3090 vec![(DetTid::from_raw(3), 0x404100)],
3091 "a failed exec keeps the old address-space registration indexed"
3092 );
3093 }
3094
3095 fn two_owner_robust_list_state() -> (
3096 ThreadState<()>,
3097 ThreadState<()>,
3098 RobustListWake,
3099 RobustListWake,
3100 RobustListWake,
3101 ) {
3102 let first_owner = DetTid::from_raw(3);
3103 let second_owner = DetTid::from_raw(4);
3104 let mm = MmId::initial(first_owner);
3105 let lower = RobustListWake {
3106 futex: FutexID::private(mm, 0x4040),
3107 };
3108 let middle = RobustListWake {
3109 futex: FutexID::private(mm, 0x5050),
3110 };
3111 let higher = RobustListWake {
3112 futex: FutexID::private(mm, 0x6060),
3113 };
3114 let mut first = ThreadState::<()>::new(first_owner, &Config::default(), ());
3115 let mut second = first.clone();
3116 second.dettid = second_owner;
3117 first.record_robust_list_head(Some(0x404100));
3118 second.record_robust_list_head(Some(0x404200));
3119 first.stage_robust_list_wakes(
3120 RobustListExit::Signal(libc::SIGTERM),
3121 vec![
3122 (second_owner, vec![higher]),
3123 (first_owner, vec![middle, lower]),
3124 ],
3125 );
3126 (first, second, lower, middle, higher)
3127 }
3128
3129 #[test]
3130 fn group_robust_list_wakes_wait_for_every_owner_and_sort_the_request() {
3131 let (first, second, lower, middle, higher) = two_owner_robust_list_state();
3132 let mut first_time = DetTime::default();
3133 first_time.add_syscall();
3134 let mut second_time = first_time.clone();
3135 second_time.add_syscall();
3136
3137 assert_eq!(
3138 second.take_robust_list_wakes_after_exit(Some(libc::SIGTERM), second_time.clone()),
3139 None,
3140 "the first physical exit must not release part of the group"
3141 );
3142 assert_eq!(
3143 first.take_robust_list_wakes_after_exit(Some(libc::SIGTERM), first_time),
3144 Some((
3145 second_time,
3146 vec![
3147 (DetTid::from_raw(3), lower),
3148 (DetTid::from_raw(3), middle),
3149 (DetTid::from_raw(4), higher),
3150 ],
3151 )),
3152 "the final matching exit releases one owner/futex-sorted request"
3153 );
3154 }
3155
3156 #[test]
3157 fn group_robust_list_wake_request_is_independent_of_exit_arrival_order() {
3158 let release = |reverse: bool| {
3159 let (first, second, _, _, _) = two_owner_robust_list_state();
3160 let first_time = DetTime::default();
3161 let mut second_time = first_time.clone();
3162 second_time.add_syscall();
3163 if reverse {
3164 assert_eq!(
3165 second.take_robust_list_wakes_after_exit(
3166 Some(libc::SIGTERM),
3167 second_time.clone(),
3168 ),
3169 None
3170 );
3171 first.take_robust_list_wakes_after_exit(Some(libc::SIGTERM), first_time)
3172 } else {
3173 assert_eq!(
3174 first.take_robust_list_wakes_after_exit(Some(libc::SIGTERM), first_time),
3175 None
3176 );
3177 second.take_robust_list_wakes_after_exit(Some(libc::SIGTERM), second_time)
3178 }
3179 };
3180
3181 assert_eq!(release(false), release(true));
3182 }
3183
3184 #[test]
3185 fn mismatched_or_nonfatal_exit_does_not_release_group_robust_list_wakes() {
3186 let (first, second, _, _, _) = two_owner_robust_list_state();
3187 assert_eq!(
3188 first.take_robust_list_wakes_after_exit(Some(libc::SIGKILL), DetTime::default()),
3189 None
3190 );
3191 assert_eq!(
3192 second.take_robust_list_wakes_after_exit(Some(libc::SIGTERM), DetTime::default()),
3193 None,
3194 "one matching owner cannot release a group whose peer exited for another reason"
3195 );
3196
3197 let (first, second, _, _, _) = two_owner_robust_list_state();
3198 assert_eq!(
3199 first.take_robust_list_wakes_after_exit(None, DetTime::default()),
3200 None
3201 );
3202 assert_eq!(
3203 second.take_robust_list_wakes_after_exit(Some(libc::SIGTERM), DetTime::default()),
3204 None,
3205 "a normal exit must not satisfy a staged fatal-signal exit"
3206 );
3207 }
3208
3209 #[test]
3210 fn regular_file_opens_do_not_shift_socket_cookie_identity() {
3211 let owner = DetTid::from_raw(3);
3212 let mut files = FileMetadata::new(owner);
3213 let first = files.allocate_open_file_id(owner, FdType::Socket);
3214 for _ in 0..4 {
3215 files.allocate_open_file_id(owner, FdType::Regular);
3216 }
3217 let second = files.allocate_open_file_id(owner, FdType::Socket);
3218
3219 assert_eq!(first.deterministic_socket_cookie(), 3_u64 << 32);
3220 assert_eq!(second.deterministic_socket_cookie(), (3_u64 << 32) | 1);
3221 }
3222
3223 #[test]
3224 fn guest_clock_tracks_raw_logical_time_without_lag() {
3225 let epoch = LogicalTime::from_secs(1_000);
3226 let mut clock = GuestClock::default();
3227
3228 assert_eq!(
3229 clock.observe(epoch + Duration::from_nanos(41_000_000)),
3230 epoch + Duration::from_nanos(41_000_000)
3231 );
3232 assert_eq!(
3233 clock.observe(epoch + Duration::from_nanos(41_025_000)),
3234 epoch + Duration::from_nanos(41_025_000)
3235 );
3236 assert_eq!(
3238 clock.observe(epoch + Duration::from_nanos(41_010_000)),
3239 epoch + Duration::from_nanos(41_025_000)
3240 );
3241 }
3242
3243 #[test]
3244 fn guest_clock_absolute_deadline_stays_ahead_of_committed_time() {
3245 let committed_time = LogicalTime::from_secs(1_000) + Duration::from_millis(250);
3246 let mut clock = GuestClock::default();
3247 let guest_now = clock.observe(committed_time);
3248 let deadline = guest_now + Duration::from_millis(100);
3249
3250 assert_eq!(guest_now, committed_time);
3251 assert!(deadline > committed_time);
3252 }
3253
3254 #[test]
3255 fn guest_clock_process_tree_shares_one_monotonic_domain() {
3256 let epoch = LogicalTime::from_secs(1_000);
3257 let root = Arc::new(Mutex::new(GuestClock::default()));
3258 let forked_child = Arc::clone(&root);
3259
3260 assert!(Arc::ptr_eq(&root, &forked_child));
3261 assert_eq!(
3262 root.lock().unwrap().observe(epoch + Duration::from_secs(1)),
3263 epoch + Duration::from_secs(1)
3264 );
3265 assert_eq!(
3266 forked_child
3267 .lock()
3268 .unwrap()
3269 .observe(epoch + Duration::from_secs(2)),
3270 epoch + Duration::from_secs(2)
3271 );
3272
3273 let execed_child = Arc::clone(&forked_child);
3275 assert!(Arc::ptr_eq(&root, &execed_child));
3276 assert_eq!(
3277 execed_child
3278 .lock()
3279 .unwrap()
3280 .observe(epoch + Duration::from_secs(9)),
3281 epoch + Duration::from_secs(9)
3282 );
3283 }
3284
3285 #[test]
3286 fn unparented_thread_recovers_process_memory_identity() {
3287 let detpid = DetPid::from_raw(4);
3288 let dettid = DetTid::from_raw(7);
3289 let mut state = ThreadState::new(dettid, &Config::default(), ());
3290
3291 assert!(state.recover_process_mm_id(detpid));
3292 assert_eq!(state.mm_id, MmId::initial(detpid));
3293 assert!(!state.recover_process_mm_id(detpid));
3294 }
3295
3296 #[test]
3297 fn inherited_thread_keeps_existing_memory_identity() {
3298 let detpid = DetPid::from_raw(4);
3299 let dettid = DetTid::from_raw(7);
3300 let inherited_mm = MmId::initial(detpid).for_exec(detpid);
3301 let mut state = ThreadState::new(dettid, &Config::default(), ());
3302 state.mm_id = inherited_mm;
3303
3304 assert!(!state.recover_process_mm_id(detpid));
3305 assert_eq!(state.mm_id, inherited_mm);
3306 }
3307
3308 #[test]
3309 fn backend_can_override_open_file_creator_identity() {
3310 let host_tid = DetTid::from_raw(10_003);
3311 let virtual_tid = DetTid::from_raw(3);
3312 let mut state = ThreadState::new(host_tid, &Config::default(), ());
3313
3314 assert_eq!(state.open_file_creator, None);
3315 state.set_open_file_creator(virtual_tid);
3316 assert_eq!(state.open_file_creator, Some(virtual_tid));
3317 }
3318
3319 #[test]
3322 fn chaos_per_thread_slowdown_factor_is_stable_and_deterministic() {
3323 let seed = 0xdead_beef_u64;
3324 let max_factor: f64 = 10.0;
3325 for raw in 1..=64 {
3327 let tid = DetTid::from_raw(raw);
3328 let a = chaos_per_thread_slowdown_factor(seed, tid, 0, max_factor);
3329 let b = chaos_per_thread_slowdown_factor(seed, tid, 0, max_factor);
3330 assert_eq!(
3331 a, b,
3332 "factor must be a pure function of (seed, dettid, epoch)"
3333 );
3334 let a = a.as_f64();
3336 assert!(
3337 a >= 1.0 / max_factor - 1e-9 && a <= max_factor + 1e-9,
3338 "factor {} out of [1/{max_factor}, {max_factor}] for tid {raw}",
3339 a
3340 );
3341 }
3342 }
3343
3344 #[test]
3347 fn chaos_per_thread_slowdown_factor_varies_across_threads_and_seeds() {
3348 let max_factor = 10.0;
3349 let factors: Vec<f64> = (1..=32)
3351 .map(|raw| {
3352 chaos_per_thread_slowdown_factor(1234, DetTid::from_raw(raw), 0, max_factor)
3353 .as_f64()
3354 })
3355 .collect();
3356 let first = factors[0];
3357 assert!(
3358 factors.iter().any(|&f| (f - first).abs() > 1e-6),
3359 "per-thread factors should differ across threads"
3360 );
3361 let tid = DetTid::from_raw(7);
3363 let f_a = chaos_per_thread_slowdown_factor(1, tid, 0, max_factor).as_f64();
3364 let f_b = chaos_per_thread_slowdown_factor(2, tid, 0, max_factor).as_f64();
3365 assert!(
3366 (f_a - f_b).abs() > 1e-12,
3367 "different seeds should yield different factors for the same thread"
3368 );
3369 }
3370
3371 #[test]
3374 fn chaos_per_thread_slowdown_factor_disabled_when_max_factor_at_most_one() {
3375 for raw in 1..=16 {
3377 let tid = DetTid::from_raw(raw);
3378 assert_eq!(
3379 chaos_per_thread_slowdown_factor(99, tid, 0, 1.0),
3380 RcbTimeMultiplier::ONE
3381 );
3382 assert_eq!(
3383 chaos_per_thread_slowdown_factor(99, tid, 0, 0.5),
3384 RcbTimeMultiplier::ONE
3385 );
3386 }
3387 }
3388
3389 #[test]
3392 fn chaos_epoch_zero_reproduces_epochless_factor() {
3393 use rand::RngExt as _;
3397 use rand::SeedableRng as _;
3398 let max_factor: f64 = 10.0;
3399 for raw in 1..=64 {
3400 let tid = DetTid::from_raw(raw);
3401 for &seed in &[0u64, 1, 7, 0xdead_beef, u64::MAX] {
3402 let epochless = {
3403 const SLOWDOWN_SALT: u64 = 0x736c_6f77_646f_776e;
3406 let mixed = seed
3407 ^ SLOWDOWN_SALT
3408 ^ ((tid.as_raw() as u32 as u64).wrapping_mul(0x9e37_79b9_7f4a_7c15));
3409 let mut prng = Pcg64Mcg::seed_from_u64(mixed);
3410 let u: f64 = prng.random::<f64>();
3411 max_factor.powf(2.0 * u - 1.0)
3412 };
3413 assert_eq!(
3414 chaos_per_thread_slowdown_factor(seed, tid, 0, max_factor),
3415 RcbTimeMultiplier::from_f64(epochless),
3416 "epoch 0 must reproduce the epoch-less factor for seed {seed}, tid {raw}"
3417 );
3418 }
3419 }
3420 }
3421
3422 #[test]
3425 fn chaos_epoch_factor_varies_deterministically_across_epochs() {
3426 let max_factor = 10.0;
3427 let seed = 0x1234_5678_u64;
3428 let tid = DetTid::from_raw(3);
3429 let factors: Vec<f64> = (0..16)
3432 .map(|epoch| chaos_per_thread_slowdown_factor(seed, tid, epoch, max_factor).as_f64())
3433 .collect();
3434 for (epoch, &f) in factors.iter().enumerate() {
3436 assert_eq!(
3437 chaos_per_thread_slowdown_factor(seed, tid, epoch as u64, max_factor).as_f64(),
3438 f,
3439 "factor must be pure in epoch"
3440 );
3441 assert!(f >= 1.0 / max_factor - 1e-9 && f <= max_factor + 1e-9);
3443 }
3444 let first = factors[0];
3446 assert!(
3447 factors.iter().any(|&f| (f - first).abs() > 1e-6),
3448 "per-epoch factors should differ across epochs"
3449 );
3450 }
3451
3452 fn cpu_snapshot(
3453 user: u64,
3454 system: u64,
3455 children_user: u64,
3456 children_system: u64,
3457 ) -> ProcessCpuSnapshot {
3458 ProcessCpuSnapshot {
3459 user: LogicalTime::from_nanos(user),
3460 system: LogicalTime::from_nanos(system),
3461 children_user: LogicalTime::from_nanos(children_user),
3462 children_system: LogicalTime::from_nanos(children_system),
3463 }
3464 }
3465
3466 #[test]
3467 fn child_cpu_time_is_hidden_until_reap() {
3468 let pid = DetPid::from_raw(2);
3469 let mut parent = ProcessCpuTime::default();
3470 parent.record_exited_child(pid, cpu_snapshot(10, 20, 3, 4));
3471
3472 assert_eq!(parent.snapshot.children_user, LogicalTime::ZERO);
3473 assert_eq!(parent.snapshot.children_system, LogicalTime::ZERO);
3474
3475 parent.reap_child(pid);
3476 assert_eq!(parent.snapshot.children_user, LogicalTime::from_nanos(13));
3477 assert_eq!(parent.snapshot.children_system, LogicalTime::from_nanos(24));
3478 }
3479
3480 #[test]
3481 fn reaping_nonexited_child_does_not_change_accounting() {
3482 let pid = DetPid::from_raw(2);
3483 let mut parent = ProcessCpuTime::default();
3484
3485 parent.reap_child(pid);
3486 assert_eq!(parent.snapshot.children_user, LogicalTime::ZERO);
3487 assert_eq!(parent.snapshot.children_system, LogicalTime::ZERO);
3488 }
3489
3490 #[test]
3491 fn child_cpu_time_uses_final_thread_snapshot_and_drops_reaped_state() {
3492 let pid = DetPid::from_raw(2);
3493 let mut parent = ProcessCpuTime::default();
3494
3495 parent.record_exited_child(pid, cpu_snapshot(10, 20, 3, 4));
3496 parent.record_exited_child(pid, cpu_snapshot(12, 25, 4, 5));
3497
3498 parent.reap_child(pid);
3499 assert_eq!(parent.snapshot.children_user, LogicalTime::from_nanos(16));
3500 assert_eq!(parent.snapshot.children_system, LogicalTime::from_nanos(30));
3501 assert!(parent.exited_children.is_empty());
3502
3503 parent.reap_child(pid);
3504 assert_eq!(parent.snapshot.children_user, LogicalTime::from_nanos(16));
3505 assert_eq!(parent.snapshot.children_system, LogicalTime::from_nanos(30));
3506 }
3507
3508 fn nz(value: u64) -> Option<NonZeroU64> {
3509 NonZeroU64::new(value)
3510 }
3511
3512 #[test]
3515 fn constant_slowdown_is_the_single_epoch_case() {
3516 let config = Config {
3517 chaos: true,
3518 chaos_per_thread_slowdown: true,
3519 chaos_epoch_length_ns: 0,
3520 target_timeslice: nz(10_000),
3521 max_timeslice: nz(100_000),
3522 ..Default::default()
3523 };
3524 let mut state = ThreadState::new(DetPid::from_raw(3), &config, ());
3525 state.next_timeslice(&config);
3526 let first = state.take_pending_chaos_epochs().pop().unwrap();
3527 assert_eq!(first.epoch, 0);
3528 assert_eq!(state.chaos_epoch, 0);
3529
3530 state
3531 .thread_logical_time
3532 .add_rcbs_with_multiplier(10_000, first.factor);
3533 state.next_timeslice(&config);
3534 assert_eq!(state.chaos_epoch, 0);
3535 assert_eq!(state.chaos_slowdown_factor, first.factor);
3536 assert!(state.take_pending_chaos_epochs().is_empty());
3537 }
3538
3539 #[test]
3542 fn epoch_redraw_uses_elapsed_logical_time_at_commit_boundaries() {
3543 let config = Config {
3544 chaos: true,
3545 chaos_per_thread_slowdown: true,
3546 chaos_epoch_length_ns: 100,
3547 target_timeslice: nz(10_000),
3548 max_timeslice: nz(100_000),
3549 ..Default::default()
3550 };
3551 let tid = DetPid::from_raw(5);
3552 let mut state = ThreadState::new(tid, &config, ());
3553 state.next_timeslice(&config);
3554 let first = state.pending_chaos_epochs[0];
3555
3556 state
3557 .thread_logical_time
3558 .add_rcbs_with_multiplier(1_000, first.factor);
3559 let expected_epoch =
3560 state.thread_logical_time.without_starting().as_nanos() / config.chaos_epoch_length_ns;
3561 assert!(expected_epoch > 0);
3562
3563 state.next_timeslice(&config);
3564 let transitions = state.take_pending_chaos_epochs();
3565 assert_eq!(transitions.len(), 2);
3566 assert_eq!(transitions[0], first);
3567 let redraw = transitions[1];
3568 assert_eq!(redraw.epoch, expected_epoch);
3569 assert_eq!(
3570 redraw.factor,
3571 chaos_per_thread_slowdown_factor(
3572 config.sched_seed(),
3573 tid,
3574 expected_epoch,
3575 config.chaos_slowdown_max_factor,
3576 )
3577 );
3578 assert!(redraw.logical_time > first.logical_time);
3579 }
3580
3581 #[test]
3584 fn replay_installs_recorded_epoch_without_ambient_chaos_flags() {
3585 let config = Config {
3586 target_timeslice: nz(10_000),
3587 max_timeslice: nz(100_000),
3588 ..Default::default()
3589 };
3590 let transition = ChaosEpochTransition {
3591 logical_time: LogicalTime::ZERO,
3592 epoch: 7,
3593 factor: RcbTimeMultiplier::from_f64(3.25),
3594 };
3595 let mut state = ThreadState::new(DetPid::from_raw(5), &config, ());
3596 state.preemption_points = Some(
3597 ThreadHistory::new()
3598 .with_chaos_epochs(vec![transition])
3599 .into_iter(),
3600 );
3601
3602 state.next_timeslice(&config);
3603
3604 assert!(state.chaos_slowdown_active);
3605 assert_eq!(state.chaos_epoch, transition.epoch);
3606 assert_eq!(state.chaos_slowdown_factor, transition.factor);
3607 assert!(state.take_pending_chaos_epochs().is_empty());
3608 }
3609
3610 #[test]
3613 fn replay_installs_boundary_epoch_when_pmu_checks_in_early() {
3614 let config = Config {
3615 target_timeslice: nz(10_000),
3616 max_timeslice: nz(100_000),
3617 ..Default::default()
3618 };
3619 let mut state = ThreadState::new(DetPid::from_raw(5), &config, ());
3620 let now = state.thread_logical_time.as_nanos();
3621 let first = ChaosEpochTransition {
3622 logical_time: now,
3623 epoch: 0,
3624 factor: RcbTimeMultiplier::from_f64(2.0),
3625 };
3626 let second = ChaosEpochTransition {
3627 logical_time: now + Duration::from_nanos(100),
3628 epoch: 1,
3629 factor: RcbTimeMultiplier::from_f64(3.0),
3630 };
3631 let history = ThreadHistory::new()
3632 .with_prio_changes(vec![
3633 (now + Duration::from_nanos(100), DEFAULT_PRIORITY),
3634 (now + Duration::from_nanos(200), DEFAULT_PRIORITY),
3635 ])
3636 .with_preemption_rcbs(vec![4, 8])
3637 .with_chaos_epochs(vec![first, second]);
3638 state.preemption_points = Some(history.into_iter());
3639
3640 state.next_timeslice(&config);
3641 assert_eq!(state.chaos_slowdown_factor, first.factor);
3642 assert_eq!(
3643 state.end_of_timeslice,
3644 Some(now + Duration::from_nanos(100))
3645 );
3646
3647 state
3650 .thread_logical_time
3651 .add_rcbs_with_multiplier(4, first.factor);
3652 state.committed_clock_value = 4;
3653 assert!(state.timeslice_expired());
3654 state.next_timeslice(&config);
3655
3656 assert_eq!(
3657 state.thread_logical_time.as_nanos(),
3658 now + Duration::from_nanos(80)
3659 );
3660 assert_eq!(state.chaos_epoch, second.epoch);
3661 assert_eq!(state.chaos_slowdown_factor, second.factor);
3662 assert_eq!(
3663 state.end_of_timeslice,
3664 Some(now + Duration::from_nanos(200))
3665 );
3666 }
3667
3668 #[test]
3671 fn replay_preserves_adjacent_zero_rcb_slices() {
3672 let config = Config {
3673 target_timeslice: nz(10_000),
3674 max_timeslice: nz(100_000),
3675 ..Default::default()
3676 };
3677 let mut state = ThreadState::new(DetPid::from_raw(5), &config, ());
3678 let now = state.thread_logical_time.as_nanos();
3679 let history = ThreadHistory::new()
3680 .with_prio_changes(vec![
3681 (now + Duration::from_nanos(70), DEFAULT_PRIORITY),
3682 (now + Duration::from_nanos(80), DEFAULT_PRIORITY),
3683 ])
3684 .with_preemption_rcbs(vec![4, 4]);
3685 state.preemption_points = Some(history.into_iter());
3686
3687 state.next_timeslice(&config);
3688 state.thread_logical_time.add_rcbs(4);
3689 state.committed_clock_value = 4;
3690 assert!(state.timeslice_expired());
3691
3692 state.next_timeslice(&config);
3693 assert_eq!(state.replay_rcb_end, Some(4));
3694 assert!(state.timeslice_expired());
3695 }
3696
3697 #[test]
3698 fn target_and_pmu_deadlines_are_independent() {
3699 let config = Config {
3700 target_timeslice: nz(20_000),
3701 max_timeslice: nz(100_000),
3702 ..Default::default()
3703 };
3704 let mut state = ThreadState::new(DetPid::from_raw(1), &config, ());
3705 let now = state.thread_logical_time.as_nanos();
3706
3707 state.next_timeslice(&config);
3708
3709 assert_eq!(
3710 state.end_of_timeslice,
3711 Some(now + Duration::from_nanos(20_000))
3712 );
3713 assert_eq!(
3714 state.max_timeslice_end,
3715 Some(now + Duration::from_nanos(100_000))
3716 );
3717 }
3718
3719 #[test]
3720 fn target_only_mode_does_not_create_a_pmu_deadline() {
3721 let config = Config {
3722 target_timeslice: nz(20_000),
3723 max_timeslice: None,
3724 ..Default::default()
3725 };
3726 let mut state = ThreadState::new(DetPid::from_raw(1), &config, ());
3727 let now = state.thread_logical_time.as_nanos();
3728
3729 state.next_timeslice(&config);
3730
3731 assert_eq!(
3732 state.end_of_timeslice,
3733 Some(now + Duration::from_nanos(20_000))
3734 );
3735 assert_eq!(state.max_timeslice_end, None);
3736 }
3737
3738 #[test]
3739 fn max_timeslice_caps_a_larger_target() {
3740 let config = Config {
3741 target_timeslice: nz(100_000),
3742 max_timeslice: nz(20_000),
3743 ..Default::default()
3744 };
3745 let mut state = ThreadState::new(DetPid::from_raw(1), &config, ());
3746 let now = state.thread_logical_time.as_nanos();
3747
3748 state.next_timeslice(&config);
3749
3750 let max_end = now + Duration::from_nanos(20_000);
3751 assert_eq!(state.end_of_timeslice, Some(max_end));
3752 assert_eq!(state.max_timeslice_end, Some(max_end));
3753 }
3754
3755 #[test]
3756 fn chaos_without_target_caps_randomized_deadline_at_maximum() {
3757 let config = Config {
3758 chaos: true,
3759 target_timeslice: None,
3760 max_timeslice: nz(100_000),
3761 clock_multiplier: Some(1.05),
3762 ..Default::default()
3763 };
3764 let mut state = ThreadState::new(DetPid::from_raw(1), &config, ());
3765 let now = state.thread_logical_time.as_nanos();
3766 let configured_max = now + Duration::from_nanos(100_000);
3767 let minimum_progress = now + Duration::from_nanos(11);
3768
3769 state.next_timeslice(&config);
3770
3771 assert_eq!(state.max_timeslice_end, state.end_of_timeslice);
3772 assert!(state.max_timeslice_end.unwrap() <= configured_max);
3773 assert!(state.max_timeslice_end.unwrap() >= minimum_progress);
3774 }
3775
3776 #[test]
3777 fn schedule_replay_without_rcb_time_arms_pmu_maximum() {
3778 let config = Config {
3779 no_rcb_time: true,
3780 max_timeslice: nz(100_000),
3781 replay_schedule_from: Some(std::path::PathBuf::from("schedule.json")),
3782 ..Default::default()
3783 };
3784 let mut state = ThreadState::new(DetPid::from_raw(1), &config, ());
3785 let now = state.thread_logical_time.as_nanos();
3786
3787 state.next_timeslice(&config);
3788
3789 let expected = now + Duration::from_nanos(100_000);
3790 assert_eq!(state.end_of_timeslice, Some(expected));
3791 assert_eq!(state.max_timeslice_end, Some(expected));
3792 }
3793
3794 #[test]
3795 fn exhausted_preemption_replay_uses_bounded_pmu_maximum() {
3796 let config = Config {
3797 max_timeslice: nz(100_000),
3798 ..Default::default()
3799 };
3800 let mut state = ThreadState::new(DetPid::from_raw(3), &config, ());
3801 state.preemption_points = Some(ThreadHistory::new().into_iter());
3802 let now = state.thread_logical_time.as_nanos();
3803
3804 state.next_timeslice(&config);
3805
3806 let bounded_end = Some(now + Duration::from_nanos(100_000));
3807 assert_eq!(state.end_of_timeslice, bounded_end);
3808 assert_eq!(state.max_timeslice_end, bounded_end);
3809 }
3810
3811 #[test]
3812 fn timeslice_expiry_is_inclusive() {
3813 let config = Config::default();
3814 let mut state = ThreadState::new(DetPid::from_raw(1), &config, ());
3815 let now = state.thread_logical_time.as_nanos();
3816
3817 state.end_of_timeslice = Some(now + Duration::from_nanos(1));
3818 assert!(!state.timeslice_expired());
3819 state.end_of_timeslice = Some(now);
3820 assert!(state.timeslice_expired());
3821 }
3822
3823 #[test]
3824 fn child_rng_distinguishes_adjacent_thread_ids() {
3825 let parent = Pcg64Mcg::seed_from_u64(0);
3826 let mut even = thread_rng_from_parent("test", &parent, DetTid::from_raw(8));
3827 let mut odd = thread_rng_from_parent("test", &parent, DetTid::from_raw(9));
3828
3829 let even_values: [u64; 4] = std::array::from_fn(|_| even.next_u64());
3830 let odd_values: [u64; 4] = std::array::from_fn(|_| odd.next_u64());
3831 assert_ne!(even_values, odd_values);
3832 }
3833
3834 #[test]
3835 fn child_rng_uses_high_entropy_bits() {
3836 let parent = Pcg64Mcg::seed_from_u64(0);
3837 let mut low = thread_rng_from_parent_entropy("test", &parent, 1);
3838 let mut high = thread_rng_from_parent_entropy("test", &parent, (1_u128 << 64) | 1);
3839
3840 let low_values: [u64; 4] = std::array::from_fn(|_| low.next_u64());
3841 let high_values: [u64; 4] = std::array::from_fn(|_| high.next_u64());
3842 assert_ne!(low_values, high_values);
3843 }
3844
3845 #[test]
3846 fn child_rng_pedigree_uses_the_full_unbounded_path() {
3847 let parent_rng = Pcg64Mcg::seed_from_u64(0);
3848 let mut long_path = Pedigree::new();
3849 for index in 0..2_048 {
3850 let (parent, child) = long_path.fork();
3851 long_path = if index % 2 == 0 { child } else { parent };
3852 }
3853 let (left, right) = long_path.fork();
3854 let mut left_rng =
3855 thread_rng_from_parent_pedigree("test", &parent_rng, &left, ChildRngStream::User);
3856 let mut right_rng =
3857 thread_rng_from_parent_pedigree("test", &parent_rng, &right, ChildRngStream::User);
3858
3859 let left_values: [u64; 4] = std::array::from_fn(|_| left_rng.next_u64());
3860 let right_values: [u64; 4] = std::array::from_fn(|_| right_rng.next_u64());
3861 assert_ne!(left_values, right_values);
3862 }
3863
3864 #[test]
3865 fn child_rng_pedigree_separates_user_and_chaos_streams() {
3866 let parent_rng = Pcg64Mcg::seed_from_u64(0);
3867 let (_, child) = Pedigree::new().fork();
3868 let mut user_rng = thread_rng_from_parent_pedigree(
3869 "same logging label",
3870 &parent_rng,
3871 &child,
3872 ChildRngStream::User,
3873 );
3874 let mut chaos_rng = thread_rng_from_parent_pedigree(
3875 "same logging label",
3876 &parent_rng,
3877 &child,
3878 ChildRngStream::Chaos,
3879 );
3880
3881 let user_values: [u64; 4] = std::array::from_fn(|_| user_rng.next_u64());
3882 let chaos_values: [u64; 4] = std::array::from_fn(|_| chaos_rng.next_u64());
3883 assert_ne!(user_values, chaos_values);
3884 }
3885}
3886
3887pub fn thread_rng_from_parent(msg: &str, parent: &Pcg64Mcg, child: DetTid) -> Pcg64Mcg {
3892 thread_rng_from_parent_entropy_labeled(msg, parent, child.as_raw() as u32 as u128, "tid")
3893}
3894
3895#[derive(Clone, Copy, Debug)]
3896pub(crate) enum ChildRngStream {
3897 User,
3898 Chaos,
3899}
3900
3901impl ChildRngStream {
3902 fn domain(self) -> &'static [u8] {
3903 match self {
3904 Self::User => b"user",
3905 Self::Chaos => b"chaos",
3906 }
3907 }
3908}
3909
3910pub(crate) fn thread_rng_from_parent_pedigree(
3912 msg: &str,
3913 parent: &Pcg64Mcg,
3914 child: &Pedigree,
3915 stream: ChildRngStream,
3916) -> Pcg64Mcg {
3917 let bits = child.raw();
3918 let mut packed = Vec::with_capacity(bits.len().div_ceil(8));
3919 let mut byte = 0_u8;
3920 for (index, bit) in bits.iter().enumerate() {
3921 if *bit {
3922 byte |= 1 << (index % 8);
3923 }
3924 if index % 8 == 7 {
3925 packed.push(byte);
3926 byte = 0;
3927 }
3928 }
3929 if !bits.len().is_multiple_of(8) {
3930 packed.push(byte);
3931 }
3932
3933 let mut hasher = Sha256::new();
3940 hasher.update(b"hermit-child-rng-pedigree-v1\0");
3941 hasher.update(stream.domain());
3942 hasher.update([0]);
3943 hasher.update((bits.len() as u64).to_le_bytes());
3944 hasher.update(&packed);
3945 let digest = hasher.finalize();
3946
3947 let mut seed = <Pcg64Mcg as SeedableRng>::Seed::default();
3948 parent.clone().fill_bytes(seed.as_mut());
3949 for (seed_byte, pedigree_byte) in seed.iter_mut().zip(digest) {
3950 *seed_byte ^= pedigree_byte;
3951 }
3952 detlog!(
3953 "RNG {} seeding child {:?} pedigree {}: {:?} from parent {:?}",
3954 msg,
3955 stream,
3956 child,
3957 seed,
3958 parent
3959 );
3960 let mut rng = Pcg64Mcg::from_seed(seed);
3961 rng.next_u64();
3962 rng.next_u64();
3963 rng.next_u64();
3964 rng.next_u64();
3965 rng
3966}
3967
3968fn thread_rng_from_parent_entropy(msg: &str, parent: &Pcg64Mcg, entropy: u128) -> Pcg64Mcg {
3969 thread_rng_from_parent_entropy_labeled(msg, parent, entropy, "entropy")
3970}
3971
3972fn thread_rng_from_parent_entropy_labeled(
3973 msg: &str,
3974 parent: &Pcg64Mcg,
3975 entropy: u128,
3976 identity_kind: &str,
3977) -> Pcg64Mcg {
3978 let mut seed = <Pcg64Mcg as SeedableRng>::Seed::default();
3980 parent.clone().fill_bytes(seed.as_mut());
3982 detlog!("RNG {} Generated new seed {:?}", msg, seed);
3983 let entropy_bytes = entropy.to_le_bytes();
3987 for (seed_byte, entropy_byte) in seed[4..].iter_mut().zip(entropy_bytes) {
3988 *seed_byte ^= entropy_byte;
3989 }
3990 detlog!(
3991 "RNG {} seeding child {} {}: {:?} from parent {:?}",
3992 msg,
3993 identity_kind,
3994 entropy,
3995 seed,
3996 parent
3997 );
3998 let mut rng = Pcg64Mcg::from_seed(seed);
3999 rng.next_u64();
4002 rng.next_u64();
4003 rng.next_u64();
4004 rng.next_u64();
4005 rng
4006}