1mod exec_identity;
13mod parked;
14use std::cmp::Ordering;
15use std::collections::BTreeMap;
16use std::collections::BTreeSet;
17use std::collections::HashMap;
18use std::collections::HashSet;
19use std::collections::btree_map::Entry;
20use std::fmt::Debug;
21use std::fs;
22use std::fs::File;
23use std::io::Write;
24use std::num::NonZeroUsize;
25use std::os::fd::FromRawFd;
26use std::path::PathBuf;
27use std::sync::Arc;
28use std::sync::Mutex;
29use std::sync::atomic::AtomicBool;
30use std::sync::atomic::AtomicU16;
31use std::sync::atomic::Ordering::SeqCst;
32use std::task::Poll;
33use std::time::SystemTime;
34
35use anyhow::bail;
36use chrono::DateTime;
37use chrono::Utc;
38use detcore_model::procfs::mount_ids_are_ordered_subset;
39use detcore_model::summary::RunSummary;
40use detcore_model::summary::TimesliceStats;
41pub(crate) use exec_identity::reconnect_exec;
42pub(crate) use exec_identity::retire_exec;
43use nix::sys::signal;
44use nix::sys::signal::Signal;
45use nix::unistd::Pid;
46pub(crate) use parked::parked_wait_request;
47pub(crate) use parked::polled_read_request;
48pub(crate) use parked::signal_dequeued;
49use reverie::GlobalRPC;
50use reverie::GlobalTool;
51use reverie::Guest;
52use reverie::Tid;
53use reverie::syscalls::AddrMut;
54use reverie::syscalls::CloneFlags;
55use reverie::syscalls::MemoryAccess;
56use reverie::syscalls::Sysno;
57use serde::Deserialize;
58use serde::Serialize;
59use tracing::debug;
60use tracing::error;
61use tracing::info;
62use tracing::trace;
63use tracing::warn;
64
65use crate::config::Config;
66use crate::consts::ROOT_DETPID;
67use crate::ivar::Ivar;
68use crate::preemptions::PreemptionReader;
69use crate::preemptions::ThreadHistory;
70use crate::record_or_replay::RecordOrReplay;
71use crate::resources::ChaosEpochTransition;
72use crate::resources::Permission;
73use crate::resources::ResourceID;
74use crate::resources::Resources;
75use crate::scheduler::AdmitIntent;
76use crate::scheduler::AdmitSide;
77use crate::scheduler::ConsumeResult;
78use crate::scheduler::DEFAULT_PRIORITY;
79use crate::scheduler::ExecReconnect;
80use crate::scheduler::MaybePrintStack;
81use crate::scheduler::Priority;
82use crate::scheduler::SchedResponse;
83use crate::scheduler::SchedValue;
84use crate::scheduler::Scheduler;
85use crate::scheduler::ThreadNextTurn;
86use crate::scheduler::entropy_to_priority;
87use crate::scheduler::parked::*;
88use crate::scheduler::real_timer::DequeueAck;
89use crate::scheduler::real_timer::ItimerSnapshot;
90use crate::scheduler::real_timer::TimerFailure;
91use crate::scheduler::runqueue::FIRST_PRIORITY;
92use crate::scheduler::runqueue::LAST_PRIORITY;
93use crate::scheduler::runqueue::REPLAY_DEFERRED_PRIORITY;
94use crate::scheduler::runqueue::REPLAY_FOREGROUND_PRIORITY;
95use crate::scheduler::runqueue::is_ordinary_priority;
96use crate::scheduler::sched_loop;
97use crate::scheduler::sched_loop_external;
98use crate::tool_local::Detcore;
99use crate::tool_local::ExecFdBlockingOverrides;
100use crate::tool_local::RobustListWake;
101use crate::types::*;
102
103pub(crate) async fn yield_once() {
104 let mut yielded = false;
105 std::future::poll_fn(|context| {
106 if yielded {
107 Poll::Ready(())
108 } else {
109 yielded = true;
110 context.waker().wake_by_ref();
111 Poll::Pending
112 }
113 })
114 .await;
115}
116
117#[derive(Debug)]
118struct InodePool {
119 inodes: HashMap<RawInode, DetInode>,
121 detinodes_info: HashMap<DetInode, DetInodeInfo>,
122 next_inode: u64,
126}
127
128#[derive(Debug)]
130struct DetInodeInfo {
131 raw: RawInode,
132 mtime: Option<LogicalTime>,
139}
140
141pub const CANONICAL_FILE_MTIME_SECONDS: [i64; 2] = [0, 1];
161
162#[derive(PartialEq, Debug, Eq, Clone, Copy, Serialize, Deserialize)]
165pub enum ObservedMtime {
166 Unobserved,
168 HostSpecific,
170 Canonical(LogicalTime),
172}
173
174impl ObservedMtime {
175 pub fn from_host_mtime(tv_sec: i64, tv_nsec: i64) -> Self {
177 match u64::try_from(tv_sec) {
178 Ok(secs) if tv_nsec == 0 && CANONICAL_FILE_MTIME_SECONDS.contains(&tv_sec) => {
179 ObservedMtime::Canonical(LogicalTime::from_secs(secs))
180 }
181 _ => ObservedMtime::HostSpecific,
182 }
183 }
184
185 fn first_seen_mtime(self, epoch: LogicalTime) -> Option<LogicalTime> {
189 match self {
190 ObservedMtime::Unobserved => None,
191 ObservedMtime::HostSpecific => Some(epoch),
192 ObservedMtime::Canonical(mtime) => Some(mtime),
193 }
194 }
195}
196
197struct ChildRegistration {
202 parent_dettid: DetTid,
203 parent_detpid: DetPid,
204 child_dettid: DetTid,
205 child_tid_addr: usize,
206 flags: Option<CloneFlags>,
207 exit_signal: libc::c_int,
208 physical_ids: Option<(i32, i32)>,
209 maybe_priority: Option<Priority>,
210 parent_is_kernel_blocked: bool,
211}
212
213#[derive(Clone, Debug, PartialEq, Eq)]
214struct PendingExecState {
215 caller: DetTid,
216 process: DetPid,
217 mm: MmId,
218 fd_blocking: ExecFdBlockingOverrides,
219}
220
221pub struct BackendFailureCleanup {
223 pub scheduler: Result<(), tokio::task::JoinError>,
225 pub preemption_recording: Result<(), String>,
227}
228
229#[derive(Clone, Copy)]
230struct RpcIncarnation {
231 dettid: DetTid,
232 mm: MmId,
233}
234
235impl Default for InodePool {
236 fn default() -> Self {
237 InodePool::new()
238 }
239}
240
241impl InodePool {
242 fn new() -> Self {
243 InodePool {
244 inodes: HashMap::new(),
245 detinodes_info: HashMap::new(),
246 next_inode: 1,
247 }
248 }
249
250 fn add_inode(
259 &mut self,
260 raw_inode: RawInode,
261 observed: ObservedMtime,
262 epoch: LogicalTime,
263 ) -> (DetInode, LogicalTime) {
264 let dino = match self.inodes.get(&raw_inode) {
265 Some(dino) => *dino,
266 None => {
267 let new = DetInode::mint(self.next_inode);
272 self.next_inode += 1;
273 assert!(self.inodes.insert(raw_inode, new).is_none());
274 let prev = self.detinodes_info.insert(
275 new,
276 DetInodeInfo {
277 raw: raw_inode,
278 mtime: None,
279 },
280 );
281 assert!(prev.is_none()); new
283 }
284 };
285 let info = self
286 .detinodes_info
287 .get_mut(&dino)
288 .expect("Internal invariant broken, det_ino missing entry");
289 if info.mtime.is_none() {
290 info.mtime = observed.first_seen_mtime(epoch);
291 }
292 (dino, info.mtime.unwrap_or(epoch))
293 }
294
295 fn remove_inode(&mut self, det_inode: DetInode) {
297 if let Some(info) = self.detinodes_info.remove(&det_inode) {
298 self.inodes.remove(&info.raw);
299 }
300 }
301}
302
303#[derive(Debug)]
329struct DevicePool {
330 devices: HashMap<u64, u64>,
331 next_device: u64,
332}
333
334#[derive(Debug)]
344enum MountIdPool {
345 Uninitialized,
346 Invalid,
347 Ready {
348 mount_ids: BTreeMap<u64, u64>,
349 mountinfo_order: Vec<u64>,
350 allow_visible_subsets: bool,
351 unlisted_order: Vec<u64>,
352 next_mount_id: u64,
353 },
354}
355
356pub struct MountIdentityProvenance {
359 pub mountinfo_order: Vec<u64>,
360 pub unlisted_order: Vec<u64>,
361}
362
363impl MountIdPool {
364 fn from_config(mount_ids: &[u64], captured: bool, unlisted_ids: &[u64]) -> Self {
365 if !captured {
366 if !mount_ids.is_empty() || !unlisted_ids.is_empty() {
367 return Self::Invalid;
368 }
369 return Self::Uninitialized;
370 }
371 Self::from_orders(mount_ids, unlisted_ids, true).unwrap_or(Self::Invalid)
372 }
373
374 fn from_orders(
375 mount_ids: &[u64],
376 unlisted_ids: &[u64],
377 allow_visible_subsets: bool,
378 ) -> Option<Self> {
379 let mut seen = BTreeSet::new();
380 if !mount_ids.iter().all(|raw| seen.insert(*raw))
381 || !unlisted_ids
382 .iter()
383 .all(|raw| *raw != 0 && seen.insert(*raw))
384 {
385 return None;
386 }
387
388 let mut mappings = BTreeMap::new();
389 for (index, raw) in mount_ids.iter().chain(unlisted_ids).enumerate() {
390 mappings.insert(*raw, u64::try_from(index).ok()?.checked_add(1)?);
391 }
392 let next_mount_id = u64::try_from(mappings.len()).ok()?.checked_add(1)?;
393 Some(Self::Ready {
394 mount_ids: mappings,
395 mountinfo_order: mount_ids.to_vec(),
396 allow_visible_subsets,
397 unlisted_order: unlisted_ids.to_vec(),
398 next_mount_id,
399 })
400 }
401
402 fn validate_mountinfo_order(&mut self, mountinfo_order: &[u64]) -> bool {
403 if matches!(self, Self::Uninitialized) {
404 *self = Self::from_orders(mountinfo_order, &[], false).unwrap_or(Self::Invalid);
405 }
406 let Self::Ready {
407 mountinfo_order: expected,
408 allow_visible_subsets,
409 ..
410 } = self
411 else {
412 return false;
413 };
414 if *allow_visible_subsets {
415 mount_ids_are_ordered_subset(mountinfo_order, expected)
416 } else {
417 mountinfo_order == expected
418 }
419 }
420
421 fn determinize(&mut self, raw_mount_id: u64, mountinfo_order: Option<&[u64]>) -> Option<u64> {
422 if raw_mount_id == 0 {
425 return Some(0);
426 }
427 if let Some(order) = mountinfo_order {
428 if !self.validate_mountinfo_order(order) {
429 return None;
430 }
431 } else if matches!(self, Self::Uninitialized) {
432 return None;
433 }
434 let Self::Ready {
435 mount_ids,
436 unlisted_order,
437 next_mount_id,
438 ..
439 } = self
440 else {
441 return None;
442 };
443
444 if let Some(virtual_mount_id) = mount_ids.get(&raw_mount_id) {
445 return Some(*virtual_mount_id);
446 }
447 let virtual_mount_id = *next_mount_id;
448 *next_mount_id = next_mount_id.checked_add(1)?;
449 mount_ids.insert(raw_mount_id, virtual_mount_id);
450 unlisted_order.push(raw_mount_id);
451 Some(virtual_mount_id)
452 }
453
454 fn provenance(&self) -> Result<Option<MountIdentityProvenance>, &'static str> {
455 match self {
456 Self::Uninitialized => Ok(None),
457 Self::Invalid => Err("mount identity provenance is invalid"),
458 Self::Ready {
459 mountinfo_order,
460 unlisted_order,
461 ..
462 } => Ok(Some(MountIdentityProvenance {
463 mountinfo_order: mountinfo_order.clone(),
464 unlisted_order: unlisted_order.clone(),
465 })),
466 }
467 }
468}
469
470impl Default for DevicePool {
471 fn default() -> Self {
472 DevicePool::new()
473 }
474}
475
476impl DevicePool {
477 fn new() -> Self {
478 DevicePool {
481 devices: HashMap::new(),
482 next_device: 1,
483 }
484 }
485
486 fn determinize(&mut self, raw_device: u64) -> u64 {
489 match self.devices.get(&raw_device) {
490 Some(dev) => *dev,
491 None => {
492 let new = self.next_device;
493 self.next_device += 1;
494 self.devices.insert(raw_device, new);
495 new
496 }
497 }
498 }
499}
500
501#[derive(Debug)]
506pub struct GlobalState {
507 sched: Arc<Mutex<Scheduler>>,
508
509 inodes: Arc<Mutex<InodePool>>,
510
511 devices: Arc<Mutex<DevicePool>>,
514
515 mount_ids: Mutex<MountIdPool>,
517
518 next_port: AtomicU16,
520
521 used_ports: Mutex<HashSet<u16>>,
523
524 unsupported_syscalls: Mutex<BTreeSet<String>>,
526
527 unsupported_syscall_report_fd: Option<Mutex<File>>,
529
530 open_file_to_port: Mutex<HashMap<OpenFileId, u16>>,
532
533 port_start_range: AtomicU16,
534 port_end_range: AtomicU16,
535
536 past_first_execve: AtomicBool,
538
539 pending_exec_states: Mutex<BTreeMap<DetPid, PendingExecState>>,
544 completed_exec_transfers: Mutex<BTreeMap<DetPid, exec_identity::ExecTransferReceipt>>,
549
550 post_exec_fd_blocking: Mutex<BTreeMap<DetTid, ExecFdBlockingOverrides>>,
552
553 sched_handle: Option<tokio::task::JoinHandle<()>>,
554
555 global_time: Arc<Mutex<GlobalTime>>,
564
565 cfg: Config,
567
568 preemptions_to_replay: Option<PreemptionReader>,
570
571 realtime_start: SystemTime,
573}
574
575impl Default for GlobalState {
576 fn default() -> Self {
577 panic!("Detcore GlobalState Default impl should not be called");
580 }
581}
582
583impl Drop for GlobalState {
584 fn drop(&mut self) {
585 if let Some(message) =
587 format_unsupported_syscall_warning(&self.unsupported_syscalls.lock().unwrap())
588 {
589 warn!("{}", message);
590 }
591 info!("detcore shut down, destroying global state");
592 }
593}
594
595impl GlobalState {
596 async fn lock_rpc_scheduler(
601 &self,
602 consuming_cleanup: bool,
603 ) -> std::sync::MutexGuard<'_, Scheduler> {
604 std::future::poll_fn(|_| {
605 let sched = self.sched.lock().unwrap();
606 if !consuming_cleanup && sched.backend_failed() {
607 Poll::Pending
608 } else {
609 Poll::Ready(sched)
610 }
611 })
612 .await
613 }
614
615 pub fn mount_identity_provenance(
620 &self,
621 ) -> Result<Option<MountIdentityProvenance>, &'static str> {
622 self.mount_ids.lock().unwrap().provenance()
623 }
624
625 fn initialize(cfg: &Config, spawn_scheduler: bool) -> Self {
626 let sched = Arc::new(Mutex::new(Scheduler::new(cfg)));
627 let global_time = Arc::new(Mutex::new(GlobalTime::new(cfg)));
628 let handle = if cfg.sequentialize_threads && spawn_scheduler {
629 info!("[scheduler] daemon task starting up, waiting for guest thread start..");
636 Some(tokio::spawn(sched_loop(sched.clone(), global_time.clone())))
637 } else {
638 None
639 };
640
641 let preemptions_to_replay: Option<PreemptionReader> = cfg
642 .replay_preemptions_from
643 .as_ref()
644 .map(|path| PreemptionReader::new(path));
645 let range = Self::read_port_range();
646
647 let unsupported_syscall_report_fd = cfg.unsupported_syscall_report_fd.and_then(|fd| {
648 let duplicate = unsafe { libc::fcntl(fd, libc::F_DUPFD_CLOEXEC, fd) };
659 if duplicate == -1 {
660 warn!(
661 "failed to duplicate unsupported-syscall report fd {fd}: {}",
662 std::io::Error::last_os_error()
663 );
664 None
665 } else {
666 Some(Mutex::new(unsafe { File::from_raw_fd(duplicate) }))
668 }
669 });
670
671 Self {
672 sched,
673 next_port: AtomicU16::new(range[0]),
674 used_ports: Mutex::new(HashSet::new()),
675 unsupported_syscalls: Mutex::new(BTreeSet::new()),
676 unsupported_syscall_report_fd,
677 port_start_range: AtomicU16::new(range[0]),
678 port_end_range: AtomicU16::new(range[1]),
679 open_file_to_port: Mutex::new(HashMap::new()),
680 past_first_execve: AtomicBool::new(false),
681 pending_exec_states: Mutex::new(BTreeMap::new()),
682 completed_exec_transfers: Mutex::new(BTreeMap::new()),
683 post_exec_fd_blocking: Mutex::new(BTreeMap::new()),
684 inodes: Arc::new(Mutex::new(InodePool::new())),
685 devices: Arc::new(Mutex::new(DevicePool::new())),
688 mount_ids: Mutex::new(MountIdPool::from_config(
689 &cfg.mountinfo_mount_ids,
690 cfg.mountinfo_mount_ids_captured,
691 &cfg.fdinfo_unlisted_mount_ids,
692 )),
693 sched_handle: handle,
694 cfg: cfg.clone(),
695 realtime_start: SystemTime::now(),
696 global_time,
697 preemptions_to_replay,
698 }
699 }
700
701 pub fn init_for_external_scheduler(cfg: &Config) -> Self {
704 assert!(
705 cfg.sequentialize_threads,
706 "an external scheduler is only meaningful when threads are sequentialized"
707 );
708 Self::initialize(cfg, false)
709 }
710
711 pub fn run_external_scheduler(
719 &self,
720 observer: Arc<dyn Fn(&'static str) + Send + Sync>,
721 ) -> impl Future<Output = ()> + use<> {
722 info!("[scheduler] daemon task starting up, waiting for guest thread start..");
726 sched_loop_external(self.sched.clone(), self.global_time.clone(), observer)
727 }
728
729 pub fn complete_physical_process_exit(&self, raw_pid: i32) {
735 let detpid = DetPid::from_raw(raw_pid);
736 self.pending_exec_states.lock().unwrap().remove(&detpid);
737 self.completed_exec_transfers
738 .lock()
739 .unwrap()
740 .remove(&detpid);
741 self.post_exec_fd_blocking.lock().unwrap().remove(&detpid);
742 if self
743 .sched
744 .lock()
745 .unwrap()
746 .complete_physical_process_exit(detpid)
747 {
748 trace!(
749 "[detcore, dpid {}] backend completed final physical process exit",
750 detpid
751 );
752 }
753 }
754
755 pub fn release_all_physical_process_exits(&self) {
758 self.pending_exec_states.lock().unwrap().clear();
759 self.completed_exec_transfers.lock().unwrap().clear();
760 self.post_exec_fd_blocking.lock().unwrap().clear();
761 let released = self
762 .sched
763 .lock()
764 .unwrap()
765 .release_all_physical_process_exits();
766 if released != 0 {
767 trace!("released {released} final physical process-exit barrier(s)");
768 }
769 }
770
771 pub fn force_shutdown_with_error(&self) {
774 let start = std::time::Instant::now();
775 let sched = loop {
776 if start.elapsed().as_millis() > 1000 {
777 eprintln!(
778 "Could not acquire scheduler lock during forced shutdown (timeout)... proceeding anyway."
779 );
780 return;
781 }
782 match self.sched.try_lock() {
783 Ok(guard) => {
784 break guard;
785 }
786 Err(std::sync::TryLockError::WouldBlock) => {
787 std::thread::yield_now();
788 continue;
789 }
790 Err(e) => {
791 eprintln!(
792 "Could not acquire scheduler lock during forced shutdown ({})... proceeding anyway.",
793 e
794 );
795 return;
796 }
797 }
798 };
799 info!("Scheduler state at exit:\n{}", sched.full_summary());
800 }
801
802 pub async fn cancel_internal_scheduler(&mut self) {
808 if let Some(handle) = self.sched_handle.take() {
809 handle.abort();
810 match handle.await {
811 Ok(()) => {}
812 Err(error) if error.is_cancelled() => {}
813 Err(error) => panic!("cancelled scheduler task panicked: {error}"),
814 }
815 }
816 }
817
818 pub async fn clean_up_after_backend_failure(mut self) -> BackendFailureCleanup {
824 let scheduler = if let Some(handle) = self.sched_handle.take() {
825 handle.await
826 } else {
827 Ok(())
828 };
829 let writer = self
832 .sched
833 .lock()
834 .unwrap_or_else(std::sync::PoisonError::into_inner)
835 .preemption_writer
836 .take();
837 let preemption_recording = if self.cfg.record_preemptions_to.is_some() {
838 writer.map_or(Ok(()), |writer| writer.flush())
839 } else {
840 drop(writer);
843 Ok(())
844 };
845 BackendFailureCleanup {
846 scheduler,
847 preemption_recording,
848 }
849 }
850
851 pub async fn clean_up(self, to_stderr: bool, print_summary_to_json_file: &Option<PathBuf>) {
861 self.clean_up_with_dispatch_stats(to_stderr, print_summary_to_json_file, None)
862 .await
863 }
864
865 pub async fn clean_up_with_dispatch_stats(
869 mut self,
870 to_stderr: bool,
871 print_summary_to_json_file: &Option<PathBuf>,
872 dispatch_stats: Option<reverie::DispatchStats>,
873 ) {
874 if let Some(handle) = self.sched_handle.take() {
875 debug!("Global state cleanup, confirming scheduler has shut down...");
876 handle.await.expect("Global scheduler clean shutdown");
877 debug!("Global state cleanup, continuing...");
878 }
879 let banner =
880 " ------------------------------ hermit run report ------------------------------";
881 let recording_destination = self.cfg.record_preemptions_to.clone();
882 let (mut summary, info_reprio_descrip) = self.into_run_summary_for_log().unwrap();
883 summary.dispatch_stats = dispatch_stats;
884
885 if let Some(path) = print_summary_to_json_file {
887 let json = serde_json::to_string_pretty(&summary).unwrap();
888 fs::write(path, json + "\n").unwrap();
889 }
890
891 if to_stderr {
893 {
901 use std::io::Write;
902 let _ = write!(crate::util::RetryingStderr, "{}\n{}", banner, summary);
903 }
904 } else {
905 let rt = summary.realtime_elapsed.take();
907 log_run_summary(
908 banner,
909 &summary,
910 info_reprio_descrip.as_deref(),
911 recording_destination.as_deref(),
912 );
913 if let Some(x) = rt {
914 debug!("Nondeterministic realtime elapsed: {:?}", x);
915 }
916 }
917 }
918
919 #[cfg(test)]
920 fn into_run_summary(self) -> anyhow::Result<RunSummary> {
921 self.into_run_summary_for_log().map(|(summary, _)| summary)
922 }
923
924 fn into_run_summary_for_log(self) -> anyhow::Result<(RunSummary, Option<String>)> {
925 let (mut summary, info_reprio_descrip) = {
927 let mut sched = self.sched.lock().unwrap();
928 sched.generate_partial_run_summary_for_log(self.cfg.record_preemptions_to.as_ref())?
929 };
930 summary.realtime_elapsed = Some(self.realtime_start.elapsed()?);
936
937 if self.cfg.virtualize_time {
938 let final_time = self.global_time.lock().unwrap();
939 summary.virttime_final = final_time.as_nanos().as_nanos();
940 summary.virttime_elapsed = match final_time.elapsed_nanos() {
941 Some(elapsed) => elapsed.as_nanos(),
942 None => bail!(
943 "Internal invariant violated! Global time {} is before its epoch baseline",
944 final_time.as_nanos()
945 ),
946 };
947 }
948
949 Ok((summary, info_reprio_descrip))
950 }
951}
952
953fn log_run_summary(
954 banner: &str,
955 summary: &RunSummary,
956 info_reprio_descrip: Option<&str>,
957 recording_destination: Option<&std::path::Path>,
958) {
959 info!("\n{}\n{}", banner, summary.info(info_reprio_descrip));
960 debug!(
961 replayed_events = summary.schedevent_replayed,
962 ?recording_destination,
963 "Run recording/replay bookkeeping"
964 );
965}
966
967#[reverie::global_tool]
968impl GlobalTool for GlobalState {
969 type Config = Config;
970
971 type Request = (DetTime, MmId, GlobalRequest);
977
978 type Response = (Option<LogicalTime>, GlobalResponse);
986
987 async fn init_global_state(cfg: &Config) -> GlobalState {
989 GlobalState::initialize(cfg, true)
990 }
991
992 fn install_backend_signal_control(
993 &self,
994 control: Option<reverie::BackendSignalControl>,
995 ) -> Result<reverie::BackendSignalControlMode, reverie::Error> {
996 self.sched.lock().unwrap().install_signal_control(control)
997 }
998
999 fn authorize_backend_signal_boundary(
1000 &self,
1001 task: reverie::SignalTaskIdentity,
1002 ) -> Result<Option<reverie::SignalDeliveryPermit>, reverie::Error> {
1003 self.sched.lock().unwrap().authorize_signal_boundary(task)
1004 }
1005
1006 async fn on_backend_signal_boundary(
1007 &self,
1008 receipt: reverie::SignalBoundaryReceipt,
1009 ) -> Result<(), reverie::Error> {
1010 self.sched.lock().unwrap().consume_signal_boundary(receipt)
1011 }
1012
1013 fn on_backend_process_retired(
1014 &self,
1015 event: reverie::BackendProcessRetirement,
1016 ) -> Result<(), reverie::Error> {
1017 let (result, wakes) = {
1018 let mut sched = self.sched.lock().unwrap();
1019 let result = sched.consume_process_retirement(event);
1020 (result, sched.take_signal_failure_wakes())
1021 };
1022 for wake in wakes {
1023 let _ = wake.send(());
1024 }
1025 result
1026 }
1027
1028 fn report_backend_failure(&self, event: reverie::BackendFailure) {
1029 let (wake, deferred) = {
1030 let mut sched = self.sched.lock().unwrap();
1031 (
1032 sched.report_backend_failure(event),
1033 sched.take_signal_failure_wakes(),
1034 )
1035 };
1036 for wake in deferred {
1037 let _ = wake.send(());
1038 }
1039 if let Some(wake) = wake {
1040 let _ = wake.send(());
1043 }
1044 }
1045
1046 async fn wait_for_backend_failure(&self) {
1047 let wake = self.sched.lock().unwrap().backend_failure_waiter();
1048 wake.await
1049 .expect("GlobalState owns the failure sender until publication");
1050 }
1051
1052 async fn on_backend_child_wait_event(
1053 &self,
1054 event: reverie::BackendChildWaitEvent,
1055 ) -> Result<(), reverie::Error> {
1056 crate::scheduler::signal_control::ChildExitPublicationFuture::new(self.sched.clone(), event)
1062 .await
1063 }
1064
1065 async fn receive_rpc(&self, from: Tid, gr: Self::Request) -> Self::Response {
1066 type R = GlobalResponse;
1067 let dtid = DetTid::from_raw(from.into()); let (guest_time, request_mm, request) = gr;
1069 let time_from_guest = guest_time.as_nanos();
1070 match &request {
1074 GlobalRequest::ReconnectExec { former, process } => {
1075 return self
1076 .recv_exec_transfer(from, guest_time, request_mm, *former, *process)
1077 .await;
1078 }
1079 GlobalRequest::RetireExec { thread, signaled } => {
1080 return self
1081 .recv_retire_exec_transfer(
1082 from,
1083 guest_time,
1084 request_mm,
1085 thread.clone(),
1086 *signaled,
1087 )
1088 .await;
1089 }
1090 _ => {}
1091 }
1092 if let GlobalRequest::SignalDequeued {
1093 detpid,
1094 identity,
1095 dequeue,
1096 } = &request
1097 {
1098 return self
1099 .recv_signal_dequeued(dtid, request_mm, guest_time, *detpid, *identity, *dequeue)
1100 .await;
1101 }
1102
1103 if let GlobalRequest::CreateVforkChildThread(parent, process, child, _, flags, ..) =
1107 &request
1108 {
1109 let mut sched = self.lock_rpc_scheduler(false).await;
1110 if sched.transferred_exec_tid_requires_registration(dtid) {
1111 if *child != dtid
1112 || !(flags.contains(CloneFlags::CLONE_VFORK)
1113 || (self.cfg.backend_serializes_fork_children
1114 && !flags.contains(CloneFlags::CLONE_THREAD)))
1115 || !sched.pending_vfork_registration_matches(
1116 *parent,
1117 *process,
1118 *child,
1119 request_mm,
1120 flags.contains(CloneFlags::CLONE_VM),
1121 )
1122 {
1123 return (None, R::ThreadExited);
1124 }
1125 sched.register_reused_transferred_exec_tid(dtid, request_mm);
1126 }
1127 }
1128
1129 if matches!(&request, GlobalRequest::StartNewThread(child, ..) if *child == dtid) {
1135 loop {
1136 {
1137 let sched = self.lock_rpc_scheduler(false).await;
1138 if !sched.transferred_exec_tid_requires_registration(dtid)
1139 || sched.thread_is_logically_killed(dtid)
1140 {
1141 break;
1142 }
1143 }
1144 yield_once().await;
1145 }
1146 }
1147
1148 let is_deregister = matches!(&request, GlobalRequest::DeregisterThread(_));
1149 let consuming_cleanup =
1150 is_deregister || matches!(&request, GlobalRequest::RobustListWakes(_));
1151
1152 let (exec_reconnect, is_exec_caller_after_local_mm_swap) = {
1153 let pending = self.pending_exec_states.lock().unwrap();
1154 let reconnect = match &request {
1155 GlobalRequest::CreateChildThread(child, process, _, None, _, _, _)
1156 if *child == dtid && *child == *process =>
1157 {
1158 pending.get(process).cloned()
1159 }
1160 _ => None,
1161 };
1162 let is_exec_caller_after_local_mm_swap = pending.values().any(|state| {
1163 state.caller == dtid && state.mm.for_exec(state.process) == request_mm
1164 });
1165 (reconnect, is_exec_caller_after_local_mm_swap)
1166 };
1167
1168 let mut tombstoned_deregistration = None;
1173 {
1174 let sched = self.lock_rpc_scheduler(consuming_cleanup).await;
1175 if exec_reconnect.is_none()
1176 && !is_exec_caller_after_local_mm_swap
1177 && !sched.rpc_incarnation_matches(dtid, request_mm)
1178 {
1179 debug!(
1180 "[detcore, dtid {}] rejecting {:?} RPC from retired exec incarnation {:?}",
1181 dtid, request, request_mm,
1182 );
1183 return if is_deregister {
1184 (None, R::DeregisterThread(()))
1185 } else {
1186 (None, R::ThreadExited)
1187 };
1188 }
1189 if let GlobalRequest::DeregisterThread(owner) = &request {
1190 assert_eq!(
1191 owner.dettid, dtid,
1192 "deregistration must belong to its sender"
1193 );
1194 assert_eq!(owner.mm, request_mm, "deregistration must retain its MmId");
1195 if !sched.thread_was_registered(dtid) && !sched.thread_is_logically_killed(dtid) {
1198 assert!(
1199 sched.backend_failed() || !owner.thread_start_entered,
1200 "a started thread must have a scheduler registration before deregistration"
1201 );
1202 return (None, R::DeregisterThread(()));
1206 }
1207 }
1208 let child = match &request {
1209 GlobalRequest::CreateChildThread(child, ..)
1210 | GlobalRequest::CreateVforkChildThread(_, _, child, ..) => Some(*child),
1211 _ => None,
1212 };
1213 if sched.thread_is_logically_killed(dtid) && exec_reconnect.is_none() {
1214 trace!(
1215 "[detcore, dtid {}] rejecting RPC after permanent logical-thread removal",
1216 dtid
1217 );
1218 if let GlobalRequest::DeregisterThread(deregistration) = &request {
1219 tombstoned_deregistration = Some(deregistration.clone());
1220 } else {
1221 return (None, R::ThreadExited);
1222 }
1223 }
1224 if child.is_some_and(|child| sched.thread_is_logically_killed(child))
1225 && exec_reconnect.is_none()
1226 {
1227 trace!(
1228 "[detcore, dtid {}] rejecting registration that reuses a tombstoned child TID",
1229 dtid
1230 );
1231 return (None, R::ThreadExited);
1232 }
1233
1234 if let GlobalRequest::ResumeExec(process) = &request
1235 && (dtid != *process
1236 || sched.registered_process(dtid) != Some(*process)
1237 || !sched.next_turns.contains_key(&dtid))
1238 {
1239 return (None, R::ThreadExited);
1240 }
1241
1242 self.acknowledge_exec_transfer(dtid, request_mm);
1245
1246 let is_thread_reconnect = matches!(
1247 &request,
1248 GlobalRequest::StartNewThread(child_dettid, ..) if *child_dettid == dtid
1249 ) && self.global_time.lock().unwrap().contains_thread(dtid);
1250
1251 if tombstoned_deregistration.is_none()
1256 && exec_reconnect.is_none()
1257 && !is_thread_reconnect
1258 {
1259 self.global_time.lock().unwrap().update_global_time(
1260 dtid,
1261 time_from_guest,
1262 guest_time.inherited_nanos(),
1263 );
1264 }
1265 }
1266 if let Some(deregistration) = tombstoned_deregistration {
1267 self.recv_deregister_thread(from, deregistration).await;
1268 return (None, R::DeregisterThread(()));
1269 }
1270
1271 #[allow(clippy::unit_arg)]
1275 let resp = match request {
1276 GlobalRequest::ReconnectExec { .. } | GlobalRequest::RetireExec { .. } => {
1277 unreachable!("exec transfer handled before ordinary admission")
1278 }
1279 GlobalRequest::ResumeExec(process) => {
1280 let mut resources = Resources::new(dtid);
1281 resources.insert(ResourceID::MemAddrSpace(process), Permission::RW);
1282 let (response, duration) = self
1283 .recv_request_resources(from, process, resources, Some(request_mm))
1284 .await;
1285 match response {
1286 SchedulerRpcResult::Continue(status) => R::ResumeExec(status, duration),
1287 SchedulerRpcResult::ThreadExited => R::ThreadExited,
1288 }
1289 }
1290 GlobalRequest::SignalDequeued { .. } => {
1291 unreachable!("consuming path handled before ordinary cancellation")
1292 }
1293 GlobalRequest::ParkedRequest(rs, pid, capability) => {
1294 let (response, _) = self
1295 .recv_resources_with_origin(
1296 from,
1297 pid,
1298 rs,
1299 Some(request_mm),
1300 RpcOrigin::DirectRequestResources,
1301 capability,
1302 )
1303 .await;
1304 match response {
1305 SchedulerRpcResult::Continue(r) => R::ParkedRequest(r),
1306 SchedulerRpcResult::ThreadExited => R::ThreadExited,
1307 }
1308 }
1309 GlobalRequest::ResumeParkedRequest {
1310 ticket,
1311 current_site,
1312 } => {
1313 self.recv_resume_parked(from, request_mm, ticket, current_site)
1314 .await
1315 }
1316 GlobalRequest::FinishParkedObservation {
1317 wait,
1318 lease,
1319 site,
1320 finish,
1321 } => {
1322 let ack = Ivar::new();
1323 let intent = ControlIntent::Finish {
1324 wait,
1325 lease,
1326 site,
1327 finish,
1328 ack: ack.clone(),
1329 };
1330 let posted = self
1331 .sched
1332 .lock()
1333 .unwrap()
1334 .post_control(dtid, request_mm, intent);
1335 if let Err(error) = posted {
1336 self.sched.lock().unwrap().fail_parked(dtid, error);
1337 return (None, R::ThreadExited);
1338 }
1339 R::FinishParkedObservation(ack.await)
1340 }
1341 GlobalRequest::ParkedProtocolFailure(error) => {
1342 self.sched.lock().unwrap().fail_parked(dtid, error);
1343 return (None, R::ThreadExited);
1344 }
1345
1346 GlobalRequest::RequestResources(rs, pid) => {
1347 let (response, _endtime) = self
1348 .recv_request_resources(from, pid, rs, Some(request_mm))
1349 .await;
1350 match response {
1351 SchedulerRpcResult::Continue(response) => R::RequestResources(response),
1352 SchedulerRpcResult::ThreadExited => R::ThreadExited,
1353 }
1354 }
1355 GlobalRequest::ReleaseResources(rs) => {
1356 R::ReleaseResources(self.recv_release_resources(from, rs).await)
1357 }
1358 GlobalRequest::ReleaseAllResources => {
1359 R::ReleaseAllResources(self.recv_release_all_resources(from).await)
1360 }
1361 GlobalRequest::ReportUnsupportedSyscall(name) => {
1363 let _sched = self.lock_rpc_scheduler(false).await;
1364 let inserted = self
1365 .unsupported_syscalls
1366 .lock()
1367 .unwrap()
1368 .insert(name.clone());
1369 if inserted
1370 && let Some(report) = &self.unsupported_syscall_report_fd
1371 && let Err(error) = writeln!(report.lock().unwrap(), "{name}")
1372 {
1373 warn!("failed to append unsupported-syscall report: {error}");
1374 }
1375 R::ReportUnsupportedSyscall(())
1376 }
1377 GlobalRequest::PrepareExec(process, mm, fd_blocking) => {
1378 let mut sched = self.lock_rpc_scheduler(false).await;
1379 if mm != request_mm {
1380 return (None, R::ThreadExited);
1381 }
1382 if self.cfg.sequentialize_threads && !sched.prepare_exec_teardown(dtid, process, mm)
1383 {
1384 return (None, R::ThreadExited);
1385 }
1386 trace!(
1387 "[detcore, dtid {}] preparing exec for process {} with mm {:?} and logically blocking descriptors {:?}",
1388 dtid, process, mm, fd_blocking,
1389 );
1390 self.pending_exec_states.lock().unwrap().insert(
1391 process,
1392 PendingExecState {
1393 caller: dtid,
1394 process,
1395 mm,
1396 fd_blocking,
1397 },
1398 );
1399 R::PrepareExec(())
1400 }
1401 GlobalRequest::CancelExec(process) => {
1402 let mut sched = self.lock_rpc_scheduler(false).await;
1403 let mut pending = self.pending_exec_states.lock().unwrap();
1404 if pending
1405 .get(&process)
1406 .is_some_and(|state| state.caller == dtid)
1407 {
1408 let prepared = pending.remove(&process).unwrap();
1409 sched.finish_exec_teardown(prepared.caller, process, prepared.mm, false);
1410 }
1411 R::CancelExec(())
1412 }
1413 GlobalRequest::MarkPastFirstExecve(detpid, signal_identity) => {
1414 let mut sched = self.lock_rpc_scheduler(false).await;
1415 let prepared = self.pending_exec_states.lock().unwrap().get(&dtid).cloned();
1419 if self.cfg.kvm_shared_dequeue_timers {
1420 let result = (|| {
1421 let identity = signal_identity.ok_or(ProtocolFailure::Identity)?;
1422 let pid = sched
1423 .registered_process(dtid)
1424 .ok_or(ProtocolFailure::Identity)?;
1425 let mut pending = self.pending_exec_states.lock().unwrap();
1426 if let Some(prepared) = pending.get(&pid) {
1427 if prepared.caller != dtid || prepared.process != pid {
1428 return Err(ProtocolFailure::Identity);
1429 }
1430 sched.complete_signal_exec(
1431 pid,
1432 dtid,
1433 prepared.mm,
1434 request_mm,
1435 identity,
1436 )?;
1437 pending.remove(&pid);
1438 } else {
1439 sched
1443 .real_timers
1444 .validate_task(pid, dtid, request_mm, identity)?;
1445 }
1446 Ok(())
1447 })();
1448 if let Err(error) = result {
1449 sched.fail_parked(dtid, error);
1450 return (None, R::ThreadExited);
1451 }
1452 }
1453 sched.blocked.timed_waiters.remove_posix_timers(detpid);
1457 if let Some(prepared) = &prepared
1461 && self.cfg.sequentialize_threads
1462 && prepared.caller == dtid
1463 && prepared.process == dtid
1464 && prepared.mm.for_exec(dtid) == request_mm
1465 {
1466 sched.reconnect_after_exec(ExecReconnect {
1467 caller: dtid,
1468 new_leader: dtid,
1469 detpid: dtid,
1470 pre_exec_mm: prepared.mm,
1471 post_exec_mm: request_mm,
1472 child_tid_addr: 0,
1473 reconnect_priority: None,
1474 });
1475 self.pending_exec_states.lock().unwrap().remove(&dtid);
1476 }
1477 self.past_first_execve.store(true, SeqCst);
1478 let overrides = self
1479 .post_exec_fd_blocking
1480 .lock()
1481 .unwrap()
1482 .remove(&dtid)
1483 .unwrap_or_default();
1484 trace!(
1485 "[detcore, dtid {}] restoring logically blocking descriptors after exec: {:?}",
1486 dtid, overrides,
1487 );
1488 R::MarkPastFirstExecve(overrides)
1489 }
1490 GlobalRequest::CreateChildThread(
1492 dettid,
1493 parent_detpid,
1494 ctid,
1495 flags,
1496 exit_signal,
1497 physical_ids,
1498 priority,
1499 ) => {
1500 if let Some(prepared) = &exec_reconnect {
1501 let mut sched = self.lock_rpc_scheduler(false).await;
1502 let (pending, post_exec_mm) = {
1503 let mut states = self.pending_exec_states.lock().unwrap();
1504 let Some(pending) = states.remove(&parent_detpid) else {
1505 return (None, R::ThreadExited);
1506 };
1507 assert_eq!(&pending, prepared);
1508 let post_exec_mm = pending.mm.for_exec(pending.process);
1509 (pending, post_exec_mm)
1510 };
1511 assert_eq!(pending.process, parent_detpid);
1512 if let Some((physical_pid, physical_tid)) = physical_ids
1513 && let Err(open_error) = sched.register_physical_thread(
1514 dettid,
1515 post_exec_mm,
1516 physical_pid,
1517 physical_tid,
1518 )
1519 {
1520 error!(
1521 "[detcore, dtid {}] failed to register post-exec host process {} thread {}: {}",
1522 dettid, physical_pid, physical_tid, open_error,
1523 );
1524 return (None, R::ThreadExited);
1525 }
1526 let retired = sched.reconnect_after_exec(ExecReconnect {
1527 caller: pending.caller,
1528 new_leader: dettid,
1529 detpid: parent_detpid,
1530 pre_exec_mm: pending.mm,
1531 post_exec_mm,
1532 child_tid_addr: ctid,
1533 reconnect_priority: priority,
1534 });
1535 if pending.caller != dettid {
1536 self.global_time
1537 .lock()
1538 .unwrap()
1539 .reassign_thread(pending.caller, dettid);
1540 }
1541 if !pending.fd_blocking.is_empty() {
1542 self.post_exec_fd_blocking
1543 .lock()
1544 .unwrap()
1545 .insert(dettid, pending.fd_blocking);
1546 }
1547 debug!(
1548 "[detcore, dtid {}] reconciled successful exec from caller {}; retired prior identities {:?}",
1549 dtid, pending.caller, retired
1550 );
1551 R::CreateChildThread(Some(post_exec_mm))
1552 } else {
1553 match self
1554 .recv_create_child_thread(
1555 from,
1556 request_mm,
1557 ChildRegistration {
1558 parent_dettid: DetTid::from_raw(from.into()),
1559 parent_detpid,
1560 child_dettid: dettid,
1561 child_tid_addr: ctid,
1562 flags,
1563 exit_signal,
1564 physical_ids,
1565 maybe_priority: priority,
1566 parent_is_kernel_blocked: false,
1567 },
1568 )
1569 .await
1570 {
1571 SchedulerRpcResult::Continue(()) => R::CreateChildThread(None),
1572 SchedulerRpcResult::ThreadExited => R::ThreadExited,
1573 }
1574 }
1575 }
1576 GlobalRequest::CreateVforkChildThread(
1578 parent_dettid,
1579 parent_detpid,
1580 child_dettid,
1581 ctid,
1582 flags,
1583 exit_signal,
1584 priority,
1585 ) => match self
1586 .recv_create_child_thread(
1587 from,
1588 request_mm,
1589 ChildRegistration {
1590 parent_dettid,
1591 parent_detpid,
1592 child_dettid,
1593 child_tid_addr: ctid,
1594 flags: Some(flags),
1595 exit_signal,
1596 physical_ids: None,
1597 maybe_priority: priority,
1598 parent_is_kernel_blocked: true,
1599 },
1600 )
1601 .await
1602 {
1603 SchedulerRpcResult::Continue(()) => R::CreateChildThread(None),
1604 SchedulerRpcResult::ThreadExited => R::ThreadExited,
1605 },
1606 GlobalRequest::StartNewThread(dettid, detpid, physical_ids, signal_identity) => {
1608 match self
1609 .recv_start_new_thread(
1610 from,
1611 dettid,
1612 detpid,
1613 request_mm,
1614 physical_ids,
1615 signal_identity,
1616 )
1617 .await
1618 {
1619 SchedulerRpcResult::Continue(history) => R::StartNewThread(history),
1620 SchedulerRpcResult::ThreadExited => R::ThreadExited,
1621 }
1622 }
1623 GlobalRequest::DeregisterThread(deregistration) => {
1624 R::DeregisterThread(self.recv_deregister_thread(from, deregistration).await)
1625 }
1626 GlobalRequest::SetChildTidAddress(address) => {
1627 let updated = self
1628 .lock_rpc_scheduler(false)
1629 .await
1630 .set_child_tid_address(dtid, address);
1631 if updated {
1632 R::SetChildTidAddress(())
1633 } else {
1634 R::ThreadExited
1635 }
1636 }
1637 GlobalRequest::FutexAction(dettid, action, futexid, init_read, mask) => R::FutexAction(
1638 self.recv_futex_action(
1639 RpcIncarnation {
1640 dettid,
1641 mm: request_mm,
1642 },
1643 action,
1644 futexid,
1645 init_read,
1646 mask,
1647 )
1648 .await,
1649 ),
1650 GlobalRequest::RobustListWakes(wakes) => {
1651 R::RobustListWakes(self.recv_robust_list_wakes(wakes))
1652 }
1653 GlobalRequest::DeterminizeInode(ino, observed) => {
1654 R::DeterminizeInode(self.recv_determinize_inode(from, ino, observed).await)
1655 }
1656 GlobalRequest::DeterminizeDevice(dev) => {
1659 R::DeterminizeDevice(self.recv_determinize_device(from, dev).await)
1660 }
1661 GlobalRequest::DeterminizeMountId(raw_mount_id, fallback_order) => {
1662 R::DeterminizeMountId(
1663 self.recv_determinize_mount_id(from, raw_mount_id, fallback_order.as_deref())
1664 .await,
1665 )
1666 }
1667 GlobalRequest::ValidateMountIdOrder(mountinfo_order) => R::ValidateMountIdOrder(
1668 self.recv_validate_mount_id_order(from, &mountinfo_order)
1669 .await,
1670 ),
1671 GlobalRequest::UnlinkInode(d_ino) => {
1672 R::UnlinkInode(self.recv_unlink_inode(from, d_ino).await)
1673 }
1674 GlobalRequest::TouchFile(ino) => R::TouchFile(self.recv_touch_file(from, ino).await),
1675 GlobalRequest::SetFileMtime(ino, mtime) => {
1676 R::SetFileMtime(self.recv_set_file_mtime(from, ino, mtime).await)
1677 }
1678 GlobalRequest::GlobalTimeLowerBound => {
1679 let ns = self.global_time.lock().unwrap().as_nanos();
1680 R::GlobalTimeLowerBound(ns)
1681 }
1682 GlobalRequest::TraceSchedEvent(ev, detpid, command_bootstrap) => {
1683 match self
1684 .recv_trace_schedevent(ev, detpid, request_mm, command_bootstrap)
1685 .await
1686 {
1687 SchedulerRpcResult::Continue(response) => R::TraceSchedEvent(response),
1688 SchedulerRpcResult::ThreadExited => R::ThreadExited,
1689 }
1690 }
1691 GlobalRequest::RegisterAlarm(dpid, dtid, duration, interval, sig) => {
1695 let now = self.global_time.lock().unwrap().as_nanos();
1696 match self
1697 .recv_register_alarm(
1698 dpid,
1699 RpcIncarnation {
1700 dettid: dtid,
1701 mm: request_mm,
1702 },
1703 now,
1704 duration,
1705 interval,
1706 sig,
1707 )
1708 .await
1709 {
1710 SchedulerRpcResult::Continue(remaining) => R::RegisterAlarm(remaining),
1711 SchedulerRpcResult::ThreadExited => R::ThreadExited,
1712 }
1713 }
1714 GlobalRequest::AlarmRemaining(dpid) => {
1717 let now = self.global_time.lock().unwrap().as_nanos();
1718 let mut sched = self.lock_rpc_scheduler(false).await;
1719 match sched.itimer_snapshot(dpid, now) {
1720 Ok(snapshot) => R::AlarmRemaining(snapshot),
1721 Err(error) => {
1722 sched.fail_parked(dtid, ProtocolFailure::Timer(error));
1723 R::ThreadExited
1724 }
1725 }
1726 }
1727 GlobalRequest::RegisterPosixTimer(dpid, dtid, timer_id, deadline, interval, sig) => {
1730 match self
1731 .recv_register_posix_timer(
1732 dpid,
1733 RpcIncarnation {
1734 dettid: dtid,
1735 mm: request_mm,
1736 },
1737 timer_id,
1738 deadline,
1739 interval,
1740 sig,
1741 )
1742 .await
1743 {
1744 SchedulerRpcResult::Continue(()) => R::RegisterPosixTimer(()),
1745 SchedulerRpcResult::ThreadExited => R::ThreadExited,
1746 }
1747 }
1748 GlobalRequest::ResolveKillTargets(dpid) => R::ResolveKillTargets(
1751 self.lock_rpc_scheduler(false)
1752 .await
1753 .process_signal_targets(dpid),
1754 ),
1755 GlobalRequest::NotifySignalPending(dettid, SigWrapper(signal), target_process) => {
1756 let mut scheduler = self.lock_rpc_scheduler(false).await;
1757 scheduler.notify_signal_pending(dettid, SigWrapper(signal));
1758 if signal == libc::SIGKILL
1759 && let Some(detpid) = target_process
1760 {
1761 scheduler.note_process_sigkill(dettid, detpid);
1762 }
1763 R::NotifySignalPending(())
1764 }
1765 GlobalRequest::ThreadIsLive(dtid) => {
1766 R::ThreadIsLive(self.lock_rpc_scheduler(false).await.thread_is_live(dtid))
1767 }
1768 GlobalRequest::ExactChildWaitState(parent, child) => R::ExactChildWaitState(
1769 self.lock_rpc_scheduler(false)
1770 .await
1771 .exact_child_wait_state(parent, child),
1772 ),
1773 GlobalRequest::ReadyChildWait(parent, selector) => {
1774 let sched = self.lock_rpc_scheduler(false).await;
1775 R::ReadyChildWait((
1776 sched.ready_child_wait(parent, selector),
1777 sched.has_child_wait_target(parent, selector),
1778 ))
1779 }
1780 GlobalRequest::ConsumeChildWait(parent, child) => R::ConsumeChildWait(
1781 self.lock_rpc_scheduler(false)
1782 .await
1783 .consume_child_wait(parent, child),
1784 ),
1785 GlobalRequest::ProcessGroup(process) => R::ProcessGroup(
1786 self.lock_rpc_scheduler(false)
1787 .await
1788 .thread_tree
1789 .process_group(process),
1790 ),
1791 GlobalRequest::SetProcessGroup(process, group) => R::SetProcessGroup(
1792 self.lock_rpc_scheduler(false)
1793 .await
1794 .thread_tree
1795 .set_process_group(process, group),
1796 ),
1797 GlobalRequest::CreateSession(process) => R::CreateSession(
1798 self.lock_rpc_scheduler(false)
1799 .await
1800 .thread_tree
1801 .create_session(process),
1802 ),
1803 GlobalRequest::UnrecoverableShutdown => {
1804 self.force_shutdown_with_error();
1805 R::UnrecoverableShutdown(())
1806 }
1807 GlobalRequest::RequestPort(open_file_id) => {
1808 let _sched = self.lock_rpc_scheduler(false).await;
1809 let mut mut_used_ports = self.used_ports.lock().unwrap();
1810 self.update_port_range();
1811 let total_available =
1812 self.port_end_range.load(SeqCst) - self.port_start_range.load(SeqCst);
1813 let mut index = 0;
1814 while (*mut_used_ports).contains(&self.next_port.load(SeqCst))
1815 && index < total_available
1816 {
1817 self.next_port.fetch_add(1, SeqCst);
1818 if self.next_port.load(SeqCst) > self.port_end_range.load(SeqCst) {
1819 self.next_port
1820 .store(self.port_start_range.load(SeqCst), SeqCst);
1821 }
1822 index += 1;
1823 }
1824 if index == total_available {
1825 R::PortFull
1826 } else {
1827 (*mut_used_ports).insert(self.next_port.load(SeqCst));
1828 let mut open_file_to_port = self.open_file_to_port.lock().unwrap();
1829 open_file_to_port.insert(open_file_id, self.next_port.load(SeqCst));
1830 R::RequestPort(self.next_port.load(SeqCst))
1831 }
1832 }
1833 GlobalRequest::AddUsedPort(port, open_file_id) => {
1834 let _sched = self.lock_rpc_scheduler(false).await;
1835 let mut used_ports = self.used_ports.lock().unwrap();
1836 used_ports.insert(port);
1837 let mut open_file_to_port = self.open_file_to_port.lock().unwrap();
1838 open_file_to_port.insert(open_file_id, port);
1839 R::AddUsedPort
1840 }
1841 GlobalRequest::ReleasePort(open_file_id) => {
1842 let _sched = self.lock_rpc_scheduler(false).await;
1843 let mut used_ports = self.used_ports.lock().unwrap();
1844 let mut open_file_to_port = self.open_file_to_port.lock().unwrap();
1845 let port = open_file_to_port.remove(&open_file_id);
1846 if let Some(port) = port {
1847 used_ports.remove(&port);
1848 }
1849 R::ReleasePort(port)
1850 }
1851 };
1852
1853 let sender_became_terminal =
1856 if is_deregister || exec_reconnect.is_some() || is_exec_caller_after_local_mm_swap {
1857 false
1858 } else {
1859 let sched = self.lock_rpc_scheduler(consuming_cleanup).await;
1860 sched.thread_is_logically_killed(dtid)
1861 || !sched.rpc_incarnation_matches(dtid, request_mm)
1862 };
1863 if resp == R::ThreadExited || sender_became_terminal {
1864 return (None, R::ThreadExited);
1865 }
1866
1867 let time_from_sched = self.global_time.lock().unwrap().threads_time(dtid);
1868 let time_update = match time_from_sched.cmp(&time_from_guest) {
1869 Ordering::Equal => None,
1870 Ordering::Less => {
1871 panic!(
1872 "internal error: thread time should never go down, only monotonically up: time in sched {}, thread local time was {}",
1873 time_from_sched, time_from_guest
1874 )
1875 }
1876 Ordering::Greater => Some(time_from_sched),
1877 };
1878 (time_update, resp)
1879 }
1880}
1881
1882impl GlobalState {
1883 async fn recv_resources_with_origin(
1884 &self,
1885 from: Tid,
1886 detpid: DetPid,
1887 rs: Resources,
1888 request_mm: Option<MmId>,
1889 rpc: RpcOrigin,
1890 capability: ControlCapability,
1891 ) -> (SchedulerRpcResult<ResourceReply>, Option<LogicalTime>) {
1892 let dettid = DetTid::from_raw(from.into()); let resp2 = {
1895 let mut sched = self.lock_rpc_scheduler(false).await;
1896 if sched.thread_is_logically_killed(dettid)
1897 || request_mm.is_some_and(|mm| !sched.rpc_incarnation_matches(dettid, mm))
1898 {
1899 return (SchedulerRpcResult::ThreadExited, None);
1900 }
1901 let Some(nextturn) = sched.next_turns.get(&dettid).cloned() else {
1902 panic!(
1903 "Detcore internal error: no entry for dettid {} in next_turns during resource request.",
1904 dettid
1905 );
1906 };
1907 trace!(
1908 "[detcore, dtid {}] ResourceRequest, filling request into {}",
1909 &dettid, &nextturn.req
1910 );
1911 if let Some(mm) = request_mm
1912 && let Err(error) = sched.install_resource_origin(
1913 dettid,
1914 ResourceOrigin {
1915 rpc,
1916 mm,
1917 control: capability,
1918 },
1919 )
1920 {
1921 sched.fail_parked(dettid, error);
1922 return (SchedulerRpcResult::ThreadExited, None);
1923 }
1924 sched.request_put(&nextturn.req, rs.clone(), &self.global_time);
1925 nextturn.resp
1926 };
1927 trace!(
1928 "[detcore, dtid {}] waiting on {} for resources: {:?}",
1929 dettid, &resp2, rs
1930 );
1931 let answer = resp2.get().await; self.finish_resource_response(from, detpid, rs, request_mm, answer)
1933 .await
1934 }
1935
1936 async fn finish_resource_response(
1937 &self,
1938 from: Tid,
1939 detpid: DetPid,
1940 rs: Resources,
1941 request_mm: Option<MmId>,
1942 answer: SchedResponse,
1943 ) -> (SchedulerRpcResult<ResourceReply>, Option<LogicalTime>) {
1944 let dettid = DetTid::from_raw(from.as_raw());
1945 let request_became_stale = {
1946 let sched = self.lock_rpc_scheduler(false).await;
1947 sched.thread_is_logically_killed(dettid)
1948 || request_mm.is_some_and(|mm| !sched.rpc_incarnation_matches(dettid, mm))
1949 };
1950 if request_became_stale {
1951 trace!(
1957 "[detcore, dtid {}] terminating pending request after logical removal",
1958 dettid
1959 );
1960 return (SchedulerRpcResult::ThreadExited, None);
1961 }
1962 if let SchedResponse::ObserveSignal(control) = answer {
1965 return (
1966 SchedulerRpcResult::Continue(ResourceReply::ObserveSignal(control)),
1967 None,
1968 );
1969 }
1970 if let Some((true, process, mm)) = rs.exit_identity() {
1971 info!(
1972 "Scheduler authorized an exit-group scenario, from dettid {} / detpid {}",
1973 dettid, detpid
1974 );
1975 {
1982 let mut sched = self.lock_rpc_scheduler(false).await;
1983 if sched.thread_is_logically_killed(dettid)
1984 || request_mm.is_some_and(|mm| !sched.rpc_incarnation_matches(dettid, mm))
1985 {
1986 return (SchedulerRpcResult::ThreadExited, None);
1987 }
1988 for tid in sched.thread_tree.my_thread_group(&dettid) {
1989 if tid != dettid {
1992 sched.logically_kill_thread(&tid, &process, mm);
1993 }
1994 }
1995 }
1996 }
1997
1998 match answer {
1999 SchedResponse::ObserveSignal(_) => {
2000 self.sched
2001 .lock()
2002 .unwrap()
2003 .fail_parked(dettid, ProtocolFailure::UnexpectedControl);
2004 (SchedulerRpcResult::ThreadExited, None)
2005 }
2006 SchedResponse::Go(Some(schedval)) => {
2008 trace!(
2009 "[dtid {}] resources granted, resuming normally: {:?}",
2010 dettid, rs
2011 );
2012
2013 let endtime_update = match schedval {
2014 SchedValue::TimeOut => None,
2016 SchedValue::Value(timeslice) => Some(LogicalTime::from_nanos(timeslice)),
2017 };
2018 (
2019 SchedulerRpcResult::Continue(ResourceReply::Grant(ResumeStatus::Normal)),
2020 endtime_update,
2021 )
2022 }
2023 SchedResponse::Go(None) => {
2024 trace!(
2025 "[dtid {}] resources granted but no timeslice specified",
2026 dettid,
2027 );
2028 (
2029 SchedulerRpcResult::Continue(ResourceReply::Grant(ResumeStatus::Normal)),
2030 None,
2031 )
2032 }
2033 SchedResponse::Signaled(signal) => {
2034 trace!(
2035 "[dtid {}] resources granted but interrupted by signal",
2036 dettid,
2037 );
2038 (
2039 SchedulerRpcResult::Continue(ResourceReply::Grant(ResumeStatus::Signaled(
2040 signal,
2041 ))),
2042 None,
2043 )
2044 }
2045 }
2046 }
2047
2048 async fn recv_release_resources(&self, from: Tid, rs: Resources) {
2049 trace!("[detcore] Resources released to pid {}: {:?}", from, rs);
2051 }
2052
2053 async fn recv_release_all_resources(&self, from: Tid) {
2054 trace!("[detcore] All resources held by pid {} released", from);
2056 }
2057
2058 async fn recv_create_child_thread(
2062 &self,
2063 rpc_sender: Tid,
2064 request_mm: MmId,
2065 registration: ChildRegistration,
2066 ) -> SchedulerRpcResult<()> {
2067 let ChildRegistration {
2068 parent_dettid,
2069 parent_detpid,
2070 child_dettid,
2071 child_tid_addr: ctid,
2072 flags,
2073 exit_signal,
2074 physical_ids,
2075 maybe_priority,
2076 parent_is_kernel_blocked,
2077 } = registration;
2078 let initial_priority = if let Some(pr) = &self.preemptions_to_replay {
2079 assert!(maybe_priority.is_none());
2080 let prio = pr
2081 .thread_initial_priority(&child_dettid)
2082 .unwrap_or_else(|| {
2083 warn!(
2084 "Child thread {} not found in preemption history to replay",
2085 child_dettid
2086 );
2087 DEFAULT_PRIORITY
2088 });
2089 if !is_ordinary_priority(prio) {
2090 panic!(
2091 "Read a bad initial_prority from file: {}\nFull file: {}",
2092 prio,
2093 pr.load_all(),
2094 );
2095 }
2096 prio
2097 } else {
2098 let prio = maybe_priority.expect(
2099 "create_child_thread must take an initial priority unless replaying preemptions",
2100 );
2101 if !is_ordinary_priority(prio) {
2102 panic!(
2103 "recv_create_child_thread received a bad prority argument : {}",
2104 prio,
2105 );
2106 }
2107 prio
2108 };
2109
2110 {
2111 let mut sched = self.lock_rpc_scheduler(false).await;
2112 let sender = DetTid::from_raw(rpc_sender.into());
2113 if sched.thread_is_logically_killed(sender)
2114 || !sched.rpc_incarnation_matches(sender, request_mm)
2115 || sched.thread_is_logically_killed(child_dettid)
2116 {
2117 return SchedulerRpcResult::ThreadExited;
2118 }
2119
2120 if parent_is_kernel_blocked && self.cfg.sequentialize_threads {
2121 sched.complete_vfork_registration(parent_dettid, child_dettid);
2122 }
2123
2124 let child_mm = MmId::for_clone(
2125 request_mm,
2126 child_dettid,
2127 flags.is_some_and(|flags| flags.contains(CloneFlags::CLONE_VM)),
2128 );
2129 sched.register_reused_transferred_exec_tid(child_dettid, child_mm);
2130
2131 let _entry = sched
2133 .next_turns
2134 .entry(child_dettid)
2135 .or_insert_with(|| ThreadNextTurn {
2136 dettid: child_dettid,
2137 child_tid_addr: ctid,
2138 req: Ivar::new(),
2139 resp: Ivar::new(),
2140 protocol: Default::default(),
2141 });
2142
2143 {
2144 let is_group_leader = if let Some(f) = flags {
2145 !f.contains(CloneFlags::CLONE_THREAD)
2146 } else {
2147 true };
2149 sched.thread_tree.add_child_with_wait_metadata(
2150 parent_dettid,
2151 child_dettid,
2152 is_group_leader,
2153 flags.is_some_and(|flags| flags.contains(CloneFlags::CLONE_PARENT)),
2154 exit_signal,
2155 );
2156 }
2157
2158 if let Some((physical_pid, physical_tid)) = physical_ids {
2159 let child_is_thread =
2160 flags.is_some_and(|flags| flags.contains(CloneFlags::CLONE_THREAD));
2161 let child_detpid = if child_is_thread {
2162 parent_detpid
2163 } else {
2164 child_dettid
2165 };
2166 if let Err(open_error) = sched.register_physical_thread(
2167 child_dettid,
2168 child_mm,
2169 physical_pid,
2170 physical_tid,
2171 ) {
2172 error!(
2173 "[detcore, dtid {}] cannot bind host process {} thread {} during child registration: {}",
2174 child_dettid, physical_pid, physical_tid, open_error,
2175 );
2176 sched.logically_kill_thread(&child_dettid, &child_detpid, child_mm);
2177 return SchedulerRpcResult::ThreadExited;
2178 }
2179 }
2180
2181 sched.hb_note_spawn(child_dettid);
2184
2185 if self.cfg.replay_schedule_from.is_none() {
2186 let old_prio = sched.priorities.insert(child_dettid, initial_priority);
2188 assert!(old_prio.is_none());
2189 } else {
2190 if let std::collections::btree_map::Entry::Vacant(entry) =
2193 sched.priorities.entry(child_dettid)
2194 {
2195 assert_eq!(parent_detpid, ROOT_DETPID);
2196 entry.insert(initial_priority);
2197 }
2198 }
2199
2200 if let Some(pr) = &mut sched.preemption_writer {
2201 pr.register_thread(child_dettid, initial_priority);
2202 }
2203
2204 let intent = if self.cfg.sequentialize_threads && !parent_is_kernel_blocked {
2222 AdmitIntent::PostFork(self.cfg.runs_post_fork)
2223 } else {
2224 AdmitIntent::Fixed(AdmitSide::Back)
2225 };
2226 sched.admit_to_run_queue(child_dettid, intent);
2227 debug!(
2228 "[detcore] CreateChildThread with dtid {}: admit child via {:?}.",
2229 child_dettid, intent,
2230 );
2231 sched.started_up.try_put(());
2232 }
2233 if self.cfg.sequentialize_threads && !parent_is_kernel_blocked {
2238 let mut rs = Resources::new(parent_detpid);
2239 rs.insert(
2240 ResourceID::ParentContinue {
2241 parent: parent_dettid,
2242 child: child_dettid,
2243 },
2244 Permission::W,
2245 );
2246 if matches!(
2247 self.recv_grant_resources(
2248 rpc_sender,
2249 parent_detpid,
2250 rs,
2251 Some(request_mm),
2252 RpcOrigin::ParentContinue
2253 )
2254 .await
2255 .0,
2256 SchedulerRpcResult::ThreadExited
2257 ) {
2258 return SchedulerRpcResult::ThreadExited;
2259 }
2260 }
2261 SchedulerRpcResult::Continue(())
2262 }
2263
2264 async fn recv_start_new_thread(
2268 &self,
2269 from: Tid,
2270 dettid: DetTid,
2271 detpid: DetPid,
2272 request_mm: MmId,
2273 physical_ids: Option<(i32, i32)>,
2274 signal_identity: Option<reverie::SignalTaskIdentity>,
2275 ) -> SchedulerRpcResult<Option<ThreadHistory>> {
2276 let mut tries: u64 = 0;
2277 let response_ivar = loop {
2279 yield_once().await;
2280 let mut sched = self.lock_rpc_scheduler(false).await;
2281 if sched.thread_is_logically_killed(dettid)
2282 || !sched.rpc_incarnation_matches(dettid, request_mm)
2283 {
2284 return SchedulerRpcResult::ThreadExited;
2285 }
2286 if self.cfg.backend_requires_thread_directed_process_signals && physical_ids.is_none() {
2287 error!(
2288 "[detcore, dtid {}] backend requires a host thread ID at StartNewThread",
2289 dettid,
2290 );
2291 sched.logically_kill_thread(&dettid, &detpid, request_mm);
2292 return SchedulerRpcResult::ThreadExited;
2293 }
2294 let rsrcs = {
2296 let mut s = HashMap::new();
2297 s.insert(ResourceID::MemAddrSpace(detpid), Permission::RW); Resources {
2299 tid: dettid,
2300 resources: s,
2301 poll_attempt: 0,
2302 fyi: String::new(),
2303 signal_interrupt_errno: None,
2304 backend_runtime_bootstrap: false,
2305 }
2306 };
2307 let nextturn = match sched.next_turns.entry(dettid) {
2308 Entry::Vacant(_entry) => {
2309 if tries == 0 {
2315 trace!(
2316 "[detcore, dtid {}] thread showed up early, no queue entry yet. Waiting...",
2317 dettid
2318 );
2319 }
2320 tries += 1;
2321 continue;
2322 }
2323 Entry::Occupied(entry) => {
2324 trace!(
2325 "[detcore, dtid {}] handling StartNewThread rpc. Found next_turns entry (after {} tries)",
2326 from, tries
2327 );
2328 entry.get().clone()
2329 }
2330 };
2331 if let Some((physical_pid, physical_tid)) = physical_ids
2332 && let Err(open_error) =
2333 sched.register_physical_thread(dettid, request_mm, physical_pid, physical_tid)
2334 {
2335 error!(
2336 "[detcore, dtid {}] cannot bind host process {} thread {} for exact signal delivery: {}",
2337 dettid, physical_pid, physical_tid, open_error,
2338 );
2339 sched.logically_kill_thread(&dettid, &detpid, request_mm);
2340 return SchedulerRpcResult::ThreadExited;
2341 }
2342 if self.cfg.kvm_shared_dequeue_timers {
2343 let binding = signal_identity
2344 .ok_or(TimerFailure::Identity)
2345 .and_then(|identity| {
2346 if from.as_raw() != dettid.as_raw()
2347 || sched.registered_process(dettid) != Some(detpid)
2348 {
2349 return Err(TimerFailure::Identity);
2350 }
2351 sched.real_timers.bind(detpid, dettid, request_mm, identity)
2352 });
2353 if let Err(error) = binding {
2354 sched.fail_parked(dettid, ProtocolFailure::Timer(error));
2355 return SchedulerRpcResult::ThreadExited;
2356 }
2357 }
2358 if let Err(error) = sched.install_resource_origin(
2359 dettid,
2360 ResourceOrigin {
2361 rpc: RpcOrigin::ThreadStart,
2362 mm: request_mm,
2363 control: ControlCapability::None,
2364 },
2365 ) {
2366 sched.fail_parked(dettid, error);
2367 return SchedulerRpcResult::ThreadExited;
2368 }
2369 sched.request_put(&nextturn.req, rsrcs, &self.global_time);
2370 break nextturn.resp;
2371 };
2372 debug!(
2373 "[detcore, dtid {}] New thread will now wait for response on {}...",
2374 &dettid, &response_ivar
2375 );
2376 let answer = response_ivar.get().await;
2377 if matches!(answer, SchedResponse::ObserveSignal(_)) {
2378 self.sched
2379 .lock()
2380 .unwrap()
2381 .fail_parked(dettid, ProtocolFailure::UnexpectedControl);
2382 return SchedulerRpcResult::ThreadExited;
2383 }
2384 let request_became_stale = {
2385 let sched = self.lock_rpc_scheduler(false).await;
2386 sched.thread_is_logically_killed(dettid)
2387 || !sched.rpc_incarnation_matches(dettid, request_mm)
2388 };
2389 if request_became_stale {
2390 return SchedulerRpcResult::ThreadExited;
2391 }
2392 info!(
2393 "[detcore, dtid {}] New thread given go-ahead to proceed via {}",
2394 &dettid, &response_ivar
2395 );
2396 if let Some(pr) = &self.preemptions_to_replay {
2397 let (history, old_prio) = {
2398 let mut sched = self.lock_rpc_scheduler(false).await;
2399 if sched.thread_is_logically_killed(dettid)
2400 || !sched.rpc_incarnation_matches(dettid, request_mm)
2401 {
2402 return SchedulerRpcResult::ThreadExited;
2403 }
2404 let history = pr.extract_thread_record(&dettid).unwrap_or_else(|| {
2405 warn!(
2406 "Replaying preemptions, but no record found for thread {}",
2407 dettid
2408 );
2409 ThreadHistory::new()
2410 });
2411 let old_prio = sched.priorities.insert(dettid, history.initial_priority());
2412 (history, old_prio)
2413 };
2414 debug!(
2415 "[replay-preemption] Enqueing new thread at priority {:?} (changed from {:?})",
2416 history.initial_priority(),
2417 old_prio,
2418 );
2419 SchedulerRpcResult::Continue(Some(history))
2420 } else {
2421 SchedulerRpcResult::Continue(None)
2422 }
2423 }
2424
2425 async fn recv_deregister_thread(&self, _from: Tid, mut deregistration: ThreadDeregistration) {
2428 let mut sched = self.sched.lock().unwrap();
2430 let ThreadDeregistration {
2431 dettid, detpid, mm, ..
2432 } = &deregistration;
2433 let mut pending = self.pending_exec_states.lock().unwrap();
2438 if pending.get(detpid).is_some_and(|state| {
2439 state.caller == *dettid && (*mm == state.mm || *mm == state.mm.for_exec(state.process))
2440 }) {
2441 let prepared = pending.remove(detpid).unwrap();
2442 sched.finish_exec_teardown(prepared.caller, *detpid, prepared.mm, false);
2443 deregistration.mm = prepared.mm;
2448 }
2449 drop(pending);
2450
2451 assert!(self.cfg.sequentialize_threads);
2453 self.account_deregistered_thread(&mut sched, deregistration);
2454 }
2455
2456 fn account_deregistered_thread(
2459 &self,
2460 sched: &mut Scheduler,
2461 deregistration: ThreadDeregistration,
2462 ) {
2463 let ThreadDeregistration {
2464 dettid,
2465 detpid,
2466 mm,
2467 thread_start_entered: _,
2468 timeslice_stats,
2469 syscall_count,
2470 chaos_epochs,
2471 } = deregistration;
2472 if !sched.rpc_incarnation_matches(dettid, mm) {
2473 debug!(
2474 "[detcore, dtid {}] ignoring deregistration from retired exec incarnation {:?}",
2475 dettid, mm,
2476 );
2477 return;
2478 }
2479 self.post_exec_fd_blocking.lock().unwrap().remove(&dettid);
2480 if !sched.note_deregistration_accounted(dettid) {
2481 trace!(
2482 "[detcore, dtid {}] acknowledging already-accounted deregistration",
2483 dettid
2484 );
2485 return;
2486 }
2487 if let Some(writer) = &mut sched.preemption_writer {
2488 for transition in chaos_epochs {
2489 writer.insert_chaos_epoch(dettid, transition);
2490 }
2491 }
2492 sched.record_timeslice_stats(dettid, timeslice_stats);
2493 sched.record_syscall_count(dettid, syscall_count);
2494 if !sched.defer_exec_sibling_retirement(dettid, detpid, mm)
2495 && !sched.thread_is_logically_killed(dettid)
2496 {
2497 sched.logically_kill_thread(&dettid, &detpid, mm);
2498 }
2499 trace!(
2500 "[detcore, dtid {}] thread deregistered, removed from sched structures.",
2501 dettid
2502 );
2503 }
2504
2505 async fn recv_futex_action(
2506 &self,
2507 caller: RpcIncarnation,
2508 action: FutexAction,
2509 futexid: FutexID,
2510 init_read: i32,
2511 mask: u32,
2512 ) -> Option<SchedValue> {
2513 let RpcIncarnation { dettid, mm } = caller;
2514 trace!("[detcore, dtid {}] Futex action: {:?}", &dettid, action);
2515 let response_iv = {
2516 let mut sched = self.lock_rpc_scheduler(false).await;
2517 if sched.thread_is_logically_killed(dettid)
2518 || !sched.rpc_incarnation_matches(dettid, mm)
2519 {
2520 return Some(SchedValue::Value(nix::errno::Errno::EINTR as u64));
2521 }
2522 let Some(resp_iv) = sched
2523 .next_turns
2524 .get(&dettid)
2525 .map(|nextturn| nextturn.resp.clone())
2526 else {
2527 trace!(
2530 "[detcore, dtid {}] ignoring futex action after logical thread removal",
2531 dettid
2532 );
2533 return Some(SchedValue::Value(nix::errno::Errno::EINTR as u64));
2534 };
2535 match action {
2536 FutexAction::WaitRequest(maybe_timeout) => {
2537 if sched.child_tid_was_cleared(futexid, init_read) {
2538 trace!(
2539 "[detcore, dtid {}] late wait on cleared child-TID futex {:?}",
2540 dettid, futexid
2541 );
2542 return Some(SchedValue::Value(0));
2543 }
2544 if let Err(error) = sched.install_resource_origin(
2545 dettid,
2546 ResourceOrigin {
2547 rpc: RpcOrigin::FutexAction,
2548 mm,
2549 control: ControlCapability::None,
2550 },
2551 ) {
2552 sched.fail_parked(dettid, error);
2553 return Some(SchedValue::Value(nix::errno::Errno::EINTR as u64));
2554 }
2555 sched.sleep_futex_waiter(&dettid, futexid, maybe_timeout, mask);
2556 }
2558 FutexAction::WaitFinished => {
2559 return None;
2560 }
2561 FutexAction::WakeRequest(num_threads) => {
2562 let num = sched.wake_futex_waiters(dettid, futexid, num_threads, mask);
2563 return Some(SchedValue::Value(num));
2564 }
2565 FutexAction::WakeFinished(_num_threads) => {
2566 return None;
2567 }
2568 }
2569 assert!(sched.run_queue.remove_tid(dettid));
2571 resp_iv
2572 };
2573 match response_iv.get().await {
2575 SchedResponse::Go(answer) => {
2576 trace!(
2577 "[detcore, dtid {}] Unblocked from futex_wait! ({})",
2578 &dettid, &response_iv
2579 );
2580 answer
2581 }
2582 SchedResponse::Signaled(_) => Some(SchedValue::Value(nix::errno::Errno::EINTR as u64)),
2583 SchedResponse::ObserveSignal(_) => {
2584 self.sched
2585 .lock()
2586 .unwrap()
2587 .fail_parked(dettid, ProtocolFailure::UnexpectedControl);
2588 self.wait_for_backend_failure().await;
2589 futures::future::pending().await
2590 }
2591 }
2592 }
2593
2594 fn recv_robust_list_wakes(&self, wakes: Vec<(DetTid, FutexID)>) -> Vec<u64> {
2595 let mut sched = self.sched.lock().unwrap();
2596 sched.wake_futex_waiters_after_exit(&wakes)
2597 }
2598
2599 fn epoch_logical_time(&self) -> LogicalTime {
2602 let nanos = self
2603 .cfg
2604 .epoch
2605 .timestamp_nanos_opt()
2606 .expect("epoch cannot be represented in a timestamp with nanosecond precision")
2607 as u64;
2608 LogicalTime::from_nanos(nanos)
2609 }
2610
2611 async fn recv_determinize_inode(
2612 &self,
2613 from: Tid,
2614 ino: RawInode,
2615 observed: ObservedMtime,
2616 ) -> (DetInode, LogicalTime) {
2617 let _sched = self.lock_rpc_scheduler(false).await;
2618 let epoch = self.epoch_logical_time();
2622 let (dino, ns) = self.inodes.lock().unwrap().add_inode(ino, observed, epoch);
2623 trace!(
2624 "[detcore, dtid {}] resolved (raw) inode {:?} to {:?}, mtime {}",
2625 from, ino, dino, ns
2626 );
2627 (dino, ns)
2628 }
2629
2630 async fn recv_determinize_device(&self, from: Tid, raw_device: u64) -> u64 {
2633 let _sched = self.lock_rpc_scheduler(false).await;
2634 let det_device = self.devices.lock().unwrap().determinize(raw_device);
2635 trace!(
2636 "[detcore, dtid {}] resolved (raw) device {} to {}",
2637 from, raw_device, det_device
2638 );
2639 det_device
2640 }
2641
2642 async fn recv_determinize_mount_id(
2643 &self,
2644 from: Tid,
2645 raw_mount_id: u64,
2646 mountinfo_order: Option<&[u64]>,
2647 ) -> Option<u64> {
2648 let _sched = self.lock_rpc_scheduler(false).await;
2649 let virtual_mount_id = self
2650 .mount_ids
2651 .lock()
2652 .unwrap()
2653 .determinize(raw_mount_id, mountinfo_order);
2654 trace!(
2655 "[detcore, dtid {}] resolved fdinfo mount ID {} to {:?}",
2656 from, raw_mount_id, virtual_mount_id
2657 );
2658 virtual_mount_id
2659 }
2660
2661 async fn recv_validate_mount_id_order(&self, from: Tid, mountinfo_order: &[u64]) -> bool {
2662 let _sched = self.lock_rpc_scheduler(false).await;
2663 let valid = self
2664 .mount_ids
2665 .lock()
2666 .unwrap()
2667 .validate_mountinfo_order(mountinfo_order);
2668 trace!(
2669 "[detcore, dtid {}] validated mountinfo identity order: {}",
2670 from, valid
2671 );
2672 valid
2673 }
2674
2675 async fn recv_unlink_inode(&self, from: Tid, d_ino: DetInode) {
2676 let _sched = self.lock_rpc_scheduler(false).await;
2677 trace!("[detcore, dtid {}] unlink (det) inode {:?}", from, d_ino);
2678 self.inodes.lock().unwrap().remove_inode(d_ino);
2679 }
2680
2681 async fn recv_touch_file(&self, from: Tid, ino: RawInode) {
2682 let _sched = self.lock_rpc_scheduler(false).await;
2683 let mtime = if self.cfg.virtualize_time {
2684 self.global_time.lock().unwrap().as_nanos()
2685 } else {
2686 let dt: DateTime<Utc> = Utc::now();
2689 let nanos = dt.timestamp_nanos_opt().expect(
2690 "current time cannot be represented in a timestamp with nanosecond precision",
2691 ) as u64;
2692 LogicalTime::from_nanos(nanos)
2693 };
2694 trace!(
2695 "[dtid {}] bumping mtime on file (rawinode {:?}) to {}",
2696 from, ino, mtime,
2697 );
2698 self.set_inode_mtime(ino, mtime);
2699 }
2700
2701 async fn recv_set_file_mtime(&self, from: Tid, ino: RawInode, mtime: LogicalTime) {
2707 let _sched = self.lock_rpc_scheduler(false).await;
2708 trace!(
2709 "[dtid {}] setting mtime on file (rawinode {:?}) to {}",
2710 from, ino, mtime,
2711 );
2712 self.set_inode_mtime(ino, mtime);
2713 }
2714
2715 fn set_inode_mtime(&self, ino: RawInode, mtime: LogicalTime) {
2716 let epoch = self.epoch_logical_time();
2717 let mut mg = self.inodes.lock().unwrap();
2718 let (dino, _) = mg.add_inode(ino, ObservedMtime::Unobserved, epoch);
2722 let info = mg
2723 .detinodes_info
2724 .get_mut(&dino)
2725 .expect("Invariant violation: det inode missing from map.");
2727 info.mtime = Some(mtime);
2728 }
2729
2730 async fn recv_trace_schedevent(
2731 &self,
2732 ev: SchedEvent,
2733 detpid: DetPid,
2734 request_mm: MmId,
2735 command_bootstrap: bool,
2736 ) -> SchedulerRpcResult<TraceSchedEventResponse> {
2737 let ev = {
2738 let sched = self.lock_rpc_scheduler(false).await;
2739 if !sched.rpc_incarnation_matches(ev.dettid, request_mm) {
2740 return SchedulerRpcResult::ThreadExited;
2741 }
2742 let ev = {
2744 if self.past_first_execve.load(SeqCst) {
2745 ev
2746 } else {
2747 info!(
2748 "Warning: erasing rip of pre-execve sched event! {:?}",
2749 SchedEventForLog {
2750 event: &ev,
2751 command_bootstrap
2752 }
2753 );
2754 SchedEvent {
2755 end_rip: None,
2756 start_rip: None,
2757 ..ev
2758 }
2759 }
2760 };
2761 if ev.op == Op::Syscall(Sysno::execve, SyscallPhase::Prehook) {
2763 self.past_first_execve.store(true, SeqCst);
2764 }
2765 ev
2766 };
2767
2768 let result = if self.cfg.replay_schedule_from.is_some() {
2770 let (consumed, print_stack2) = {
2771 let mut sched = self.lock_rpc_scheduler(false).await;
2772 if sched.thread_is_logically_killed(ev.dettid)
2773 || !sched.rpc_incarnation_matches(ev.dettid, request_mm)
2774 {
2775 return SchedulerRpcResult::ThreadExited;
2776 }
2777 let consumed = sched.consume_schedevent(&ev);
2778 let print_stack2 = if self.cfg.record_preemptions {
2779 sched.record_event(&ev)
2780 } else {
2781 None
2782 };
2783 (consumed, print_stack2)
2784 };
2785 let ConsumeResult {
2786 keep_running,
2787 print_stack,
2788 event_ix: _,
2789 timeslice_remaining: mut end_of_timeslice,
2790 } = consumed;
2791 trace!(
2792 "keep_running :{}, end_of_timeslice: {:?}",
2793 keep_running, end_of_timeslice
2794 );
2795
2796 if !keep_running {
2797 trace!(
2798 "[detcore, dtid {}] Thread yielding to follow replay schedule",
2799 &ev.dettid,
2800 );
2801 let tid = reverie::Tid::from(ev.dettid.as_raw()); let mut rsrcs = Resources::new(ev.dettid);
2803 rsrcs.insert(ResourceID::TraceReplay, Permission::RW);
2804 let (response, timeslice) = self
2805 .recv_grant_resources(
2806 tid,
2807 detpid,
2808 rsrcs,
2809 Some(request_mm),
2810 RpcOrigin::TraceSchedEvent,
2811 )
2812 .await;
2813 if response == SchedulerRpcResult::ThreadExited {
2814 return SchedulerRpcResult::ThreadExited;
2815 }
2816 end_of_timeslice = timeslice;
2817 trace!(
2818 "[detcore, dtid {}] Thread reactivated after yielding for replay schedule",
2819 &ev.dettid,
2820 );
2821 }
2822
2823 TraceSchedEventResponse {
2824 print_stack_strace: print_stack.or(print_stack2),
2825 timeslice: end_of_timeslice,
2826 }
2827 } else {
2828 let print_stack_strace = {
2829 let mut sched = self.lock_rpc_scheduler(false).await;
2830 if sched.thread_is_logically_killed(ev.dettid)
2831 || !sched.rpc_incarnation_matches(ev.dettid, request_mm)
2832 {
2833 return SchedulerRpcResult::ThreadExited;
2834 }
2835 if self.cfg.record_preemptions {
2836 sched.record_event(&ev)
2837 } else {
2838 None
2839 }
2840 };
2841 TraceSchedEventResponse {
2842 print_stack_strace,
2843 timeslice: None,
2844 }
2845 };
2846
2847 if result.print_stack_strace.is_some()
2848 && let Some(sig) = &self.cfg.stacktrace_signal
2849 {
2850 let _sched = self.lock_rpc_scheduler(false).await;
2851 trace!(
2852 "[dtid {}] signaling thread with {} at the point of stack trace printing.",
2853 ev.dettid, sig.0
2854 );
2855 let tid = Pid::from_raw(ev.dettid.as_raw());
2856 match sig.signal() {
2863 Some(named) => signal::kill(tid, named).unwrap(),
2864 None => {
2865 let rc = unsafe { libc::kill(tid.as_raw(), sig.raw()) };
2868 assert_eq!(rc, 0, "raw kill of signal {} failed", sig.raw());
2869 }
2870 }
2871 }
2872
2873 SchedulerRpcResult::Continue(result)
2874 }
2875
2876 fn read_port_range() -> Vec<u16> {
2880 let contents = fs::read_to_string("/proc/sys/net/ipv4/ip_local_port_range")
2881 .expect("File should be present");
2882 let range: Vec<u16> = contents
2883 .split_whitespace()
2884 .filter_map(|number| number.parse().ok())
2885 .collect();
2886 range
2887 }
2888
2889 fn update_port_range(&self) {
2891 let range = Self::read_port_range();
2892 self.port_start_range.store(range[0], SeqCst);
2893 self.port_end_range.store(range[1], SeqCst);
2894 }
2895
2896 async fn recv_register_alarm(
2901 &self,
2902 detpid: DetPid,
2903 caller: RpcIncarnation,
2904 now: LogicalTime,
2905 duration: LogicalTime,
2906 interval: LogicalTime,
2907 sig: SigWrapper,
2908 ) -> SchedulerRpcResult<(LogicalTime, LogicalTime)> {
2909 let RpcIncarnation { dettid, mm } = caller;
2910 let mut sched = self.lock_rpc_scheduler(false).await;
2911 if sched.thread_is_logically_killed(dettid) || !sched.rpc_incarnation_matches(dettid, mm) {
2912 return SchedulerRpcResult::ThreadExited;
2913 }
2914 match sched.replace_real_timer(detpid, dettid, now, duration, interval, alarm_signal(sig)) {
2915 Ok(old) => SchedulerRpcResult::Continue(old),
2916 Err(error) => {
2917 sched.fail_parked(dettid, ProtocolFailure::Timer(error));
2918 SchedulerRpcResult::ThreadExited
2919 }
2920 }
2921 }
2922
2923 async fn recv_register_posix_timer(
2927 &self,
2928 detpid: DetPid,
2929 caller: RpcIncarnation,
2930 timer_id: i32,
2931 deadline: Option<LogicalTime>,
2932 interval: LogicalTime,
2933 sig: SigWrapper,
2934 ) -> SchedulerRpcResult<()> {
2935 let RpcIncarnation { dettid, mm } = caller;
2936 let mut sched = self.lock_rpc_scheduler(false).await;
2937 if sched.thread_is_logically_killed(dettid) || !sched.rpc_incarnation_matches(dettid, mm) {
2938 return SchedulerRpcResult::ThreadExited;
2939 }
2940 sched.register_posix_timer(
2941 detpid,
2942 dettid,
2943 timer_id,
2944 deadline,
2945 interval,
2946 alarm_signal(sig),
2947 );
2948 SchedulerRpcResult::Continue(())
2949 }
2950}
2951
2952#[derive(PartialEq, Debug, Eq, Clone, Serialize, Deserialize)]
2954pub struct ThreadDeregistration {
2955 pub(crate) dettid: DetTid,
2956 pub(crate) detpid: DetPid,
2957 pub(crate) mm: MmId,
2958 pub(crate) thread_start_entered: bool,
2961 pub(crate) timeslice_stats: TimesliceStats,
2962 pub(crate) syscall_count: u64,
2963 pub(crate) chaos_epochs: Vec<ChaosEpochTransition>,
2964}
2965
2966#[derive(PartialEq, Debug, Eq, Clone, Serialize, Deserialize)]
2971#[allow(clippy::enum_variant_names)]
2972pub enum GlobalRequest {
2973 RequestResources(Resources, DetPid),
2976 ParkedRequest(Resources, DetPid, ControlCapability),
2977 ResumeParkedRequest {
2978 ticket: ResumeTicket,
2979 current_site: reverie::CallbackSignalSite,
2980 },
2981 FinishParkedObservation {
2982 wait: ContinuationId,
2983 lease: reverie::ParkedObservationLease,
2984 site: reverie::CallbackSignalSite,
2985 finish: ObservationFinish,
2986 },
2987 ParkedProtocolFailure(ProtocolFailure),
2988 SignalDequeued {
2989 detpid: DetPid,
2990 identity: reverie::SignalTaskIdentity,
2991 dequeue: reverie::SignalDequeue,
2992 },
2993 ReleaseResources(Resources),
2995 ReleaseAllResources,
2997
2998 ReportUnsupportedSyscall(String),
3001
3002 PrepareExec(DetPid, MmId, ExecFdBlockingOverrides),
3007
3008 CancelExec(DetPid),
3010
3011 ReconnectExec {
3013 former: DetTid,
3014 process: DetPid,
3015 },
3016 ResumeExec(DetPid),
3018 RetireExec {
3022 thread: ThreadDeregistration,
3023 signaled: bool,
3024 },
3025
3026 MarkPastFirstExecve(DetPid, Option<reverie::SignalTaskIdentity>),
3029
3030 CreateChildThread(
3036 DetTid,
3037 DetPid,
3038 usize,
3039 Option<CloneFlags>,
3040 libc::c_int,
3041 Option<(i32, i32)>,
3042 Option<Priority>,
3043 ),
3044
3045 CreateVforkChildThread(
3050 DetTid,
3051 DetPid,
3052 DetTid,
3053 usize,
3054 CloneFlags,
3055 libc::c_int,
3056 Option<Priority>,
3057 ),
3058
3059 StartNewThread(
3062 DetTid,
3063 DetPid,
3064 Option<(i32, i32)>,
3065 Option<reverie::SignalTaskIdentity>,
3066 ),
3067
3068 DeregisterThread(ThreadDeregistration),
3072
3073 SetChildTidAddress(usize),
3076
3077 FutexAction(DetTid, FutexAction, FutexID, i32, u32),
3080
3081 DeterminizeInode(RawInode, ObservedMtime),
3084
3085 DeterminizeDevice(u64),
3090
3091 DeterminizeMountId(u64, Option<Vec<u64>>),
3095
3096 ValidateMountIdOrder(Vec<u64>),
3098
3099 UnlinkInode(DetInode),
3101
3102 TouchFile(RawInode),
3104
3105 SetFileMtime(RawInode, LogicalTime),
3107
3108 GlobalTimeLowerBound,
3110
3111 TraceSchedEvent(SchedEvent, DetPid, bool),
3114
3115 RegisterAlarm(DetPid, DetTid, LogicalTime, LogicalTime, SigWrapper),
3120
3121 RegisterPosixTimer(
3125 DetPid,
3126 DetTid,
3127 i32,
3128 Option<LogicalTime>,
3129 LogicalTime,
3130 SigWrapper,
3131 ),
3132
3133 AlarmRemaining(DetPid),
3137
3138 ReadyChildWait(DetPid, ChildWaitSpec),
3142 ConsumeChildWait(DetPid, DetPid),
3144 ProcessGroup(DetPid),
3146 SetProcessGroup(DetPid, DetPid),
3148 CreateSession(DetPid),
3150 ResolveKillTargets(DetPid),
3152 NotifySignalPending(DetTid, SigWrapper, Option<DetPid>),
3154 ThreadIsLive(DetTid),
3156 ExactChildWaitState(DetPid, DetPid),
3158
3159 UnrecoverableShutdown,
3161
3162 RequestPort(OpenFileId),
3164
3165 AddUsedPort(u16, OpenFileId),
3167
3168 ReleasePort(OpenFileId),
3170
3171 RobustListWakes(Vec<(DetTid, FutexID)>),
3174}
3175
3176#[allow(missing_docs, clippy::unit_arg)]
3178#[derive(PartialEq, Debug, Eq, Clone, Serialize, Deserialize)]
3179pub enum GlobalResponse {
3180 ThreadExited,
3183 RequestResources(ResumeStatus),
3184 ParkedRequest(ResourceReply),
3185 ResumeParkedRequest(ResourceReply),
3186 FinishParkedObservation(Result<FinishAck, ProtocolFailure>),
3187 SignalDequeued {
3188 ack: Result<DequeueAck, TimerFailure>,
3189 terminal: bool,
3190 },
3191 ReleaseResources(()),
3192 ReleaseAllResources(()),
3193 ReportUnsupportedSyscall(()),
3195 PrepareExec(()),
3196 CancelExec(()),
3197 ReconnectExec(bool),
3198 ResumeExec(ResumeStatus, Option<LogicalTime>),
3199 RetireExec(bool),
3200 MarkPastFirstExecve(ExecFdBlockingOverrides),
3201 CreateChildThread(Option<MmId>),
3202 StartNewThread(Option<ThreadHistory>),
3204 DeregisterThread(()),
3205 SetChildTidAddress(()),
3206 FutexAction(Option<SchedValue>),
3207 DeterminizeInode((DetInode, LogicalTime)),
3209 DeterminizeDevice(u64),
3212 DeterminizeMountId(Option<u64>),
3213 ValidateMountIdOrder(bool),
3214 UnlinkInode(()),
3215 TouchFile(()),
3216 SetFileMtime(()),
3217 GlobalTimeLowerBound(LogicalTime),
3218 TraceSchedEvent(TraceSchedEventResponse),
3219 RegisterAlarm((LogicalTime, LogicalTime)),
3223 RegisterPosixTimer(()),
3226 ReadyChildWait((Option<DetPid>, bool)),
3229 ConsumeChildWait(bool),
3230 ProcessGroup(Option<DetPid>),
3231 SetProcessGroup(bool),
3232 CreateSession(bool),
3233 AlarmRemaining(ItimerSnapshot),
3234 ResolveKillTargets(Vec<DetTid>),
3237 NotifySignalPending(()),
3238 ThreadIsLive(bool),
3239 ExactChildWaitState(ExactChildWaitState),
3240 UnrecoverableShutdown(()),
3242
3243 RequestPort(u16),
3244 AddUsedPort,
3245 ReleasePort(Option<u16>),
3246 PortFull,
3247 RobustListWakes(Vec<u64>),
3248}
3249
3250pub fn format_unsupported_syscall_warning(syscalls: &BTreeSet<String>) -> Option<String> {
3254 if syscalls.is_empty() {
3255 None
3256 } else {
3257 Some(format!(
3258 "syscalls {} used but not yet supported",
3259 syscalls.iter().cloned().collect::<Vec<_>>().join(",")
3260 ))
3261 }
3262}
3263
3264pub async fn prepare_exec<G, T>(guest: &mut G, mm: MmId, fd_blocking: ExecFdBlockingOverrides)
3272where
3273 G: Guest<Detcore<T>>,
3274 T: RecordOrReplay,
3275{
3276 let detpid = guest.thread_state().detpid.expect("detpid unset");
3277 let (_, response) =
3278 send_and_update_time(guest, GlobalRequest::PrepareExec(detpid, mm, fd_blocking)).await;
3279 assert_eq!(response, GlobalResponse::PrepareExec(()));
3280}
3281
3282pub async fn cancel_exec<G, T>(guest: &mut G)
3283where
3284 G: Guest<Detcore<T>>,
3285 T: RecordOrReplay,
3286{
3287 let detpid = guest.thread_state().detpid.expect("detpid unset");
3288 let (_, response) = send_and_update_time(guest, GlobalRequest::CancelExec(detpid)).await;
3289 assert_eq!(response, GlobalResponse::CancelExec(()));
3290}
3291
3292pub async fn mark_past_first_execve<G, T>(guest: &mut G)
3293where
3294 G: Guest<Detcore<T>>,
3295 T: RecordOrReplay,
3296{
3297 let signal_identity = guest
3298 .config()
3299 .kvm_shared_dequeue_timers
3300 .then(|| guest.signal_task_identity())
3301 .flatten();
3302 let detpid = guest.thread_state().detpid.expect("detpid unset");
3303 let (_, response) = send_and_update_time(
3304 guest,
3305 GlobalRequest::MarkPastFirstExecve(detpid, signal_identity),
3306 )
3307 .await;
3308 let overrides = match response {
3309 GlobalResponse::MarkPastFirstExecve(overrides) => overrides,
3310 _ => unreachable!(),
3311 };
3312 if guest.config().kvm_shared_dequeue_timers {
3313 guest.thread_state_mut().signal_task_identity = signal_identity;
3314 }
3315 if !overrides.is_empty() {
3316 let dettid = guest.thread_state().dettid;
3317 let metadata = Arc::clone(&guest.thread_state().file_metadata);
3318 metadata
3319 .lock()
3320 .unwrap()
3321 .apply_exec_blocking_overrides(dettid, overrides);
3322 }
3323}
3324
3325pub async fn report_unsupported_syscall<G, T>(guest: &mut G, sysno: Sysno)
3327where
3328 G: Guest<Detcore<T>>,
3329 T: RecordOrReplay,
3330{
3331 let (_, response) = send_and_update_time(
3332 guest,
3333 GlobalRequest::ReportUnsupportedSyscall(sysno.to_string()),
3334 )
3335 .await;
3336 assert_eq!(response, GlobalResponse::ReportUnsupportedSyscall(()));
3337}
3338
3339pub(crate) async fn set_child_tid_address<G, T>(guest: &mut G, address: usize)
3341where
3342 G: Guest<Detcore<T>>,
3343 T: RecordOrReplay,
3344{
3345 let (_, response) =
3346 send_and_update_time(guest, GlobalRequest::SetChildTidAddress(address)).await;
3347 assert_eq!(response, GlobalResponse::SetChildTidAddress(()));
3348}
3349
3350pub async fn send_and_update_time<G, T>(
3351 guest: &mut G,
3352 request: GlobalRequest,
3353) -> (Option<LogicalTime>, GlobalResponse)
3354where
3355 G: Guest<Detcore<T>>,
3356 T: RecordOrReplay,
3357{
3358 assert_eq!(
3364 DetTid::from_raw(guest.tid().as_raw()),
3365 guest.thread_state().dettid,
3366 "replacement image must reconnect before ordinary RPCs"
3367 );
3368 let mytime = guest.thread_state().thread_logical_time.clone();
3369 let mm = guest.thread_state().mm_id;
3370 let resp = guest.send_rpc((mytime, mm, request)).await;
3371 if resp.1 == GlobalResponse::ThreadExited {
3372 let dettid = guest.thread_state().dettid;
3373 trace!(
3374 "[detcore, dtid {}] exiting after terminal scheduler cancellation",
3375 dettid
3376 );
3377 guest.retire_current_thread().await
3381 }
3382 if let Some(time) = resp.0 {
3385 guest
3386 .thread_state_mut()
3387 .thread_logical_time
3388 .advance_to(time);
3389 }
3390 resp
3391}
3392
3393#[derive(PartialEq, Debug, Eq, Clone, Serialize, Deserialize)]
3397pub enum ResumeStatus {
3398 Normal,
3399 Signaled(Option<Vec<SigWrapper>>),
3400}
3401
3402#[derive(PartialEq, Debug, Eq, Clone)]
3405enum SchedulerRpcResult<T> {
3406 Continue(T),
3407 ThreadExited,
3408}
3409
3410pub async fn resource_request<G, T>(guest: &mut G, mut r: Resources) -> ResumeStatus
3414where
3415 G: Guest<Detcore<T>>,
3416 T: RecordOrReplay,
3417{
3418 r.backend_runtime_bootstrap |=
3423 guest.thread_state().in_uncharged_bootstrap_syscall && guest.is_backend_runtime_bootstrap();
3424 if guest.config().sequentialize_threads {
3425 if let Some(lease) = guest.signal_observation_lease() {
3426 if let Some(site) = guest.parked_signal_site() {
3427 return parked::capable_resource_request(
3428 guest,
3429 r,
3430 ControlCapability::PublishOnly { lease, site },
3431 )
3432 .await;
3433 }
3434 parked::terminate_protocol(guest, ProtocolFailure::Identity).await;
3435 }
3436 let dettid = guest.thread_state().dettid;
3437 let detpid = guest.thread_state().detpid.expect("detpid unset");
3438 trace!(
3439 "[detcore, dtid {}] BLOCKING on resource_request rpc... {:?}",
3440 &dettid, r
3441 );
3442 let resp =
3443 send_and_update_time(guest, GlobalRequest::RequestResources(r.clone(), detpid)).await;
3444 match resp.1 {
3445 GlobalResponse::RequestResources(x) => {
3446 trace!(
3447 "[detcore, dtid {}] UNBLOCKED, acquired resources: {:?}",
3448 &dettid, r
3449 );
3450 x
3451 }
3452 _ => unreachable!(),
3453 }
3454 } else {
3455 ResumeStatus::Normal
3456 }
3457}
3458
3459pub async fn resource_release_all<G, T>(guest: &mut G)
3464where
3465 G: Guest<Detcore<T>>,
3466 T: RecordOrReplay,
3467{
3468 if guest.config().sequentialize_threads {
3469 let resp = send_and_update_time(guest, GlobalRequest::ReleaseAllResources).await;
3470 match resp.1 {
3471 GlobalResponse::ReleaseAllResources(x) => x,
3472 _ => unreachable!(),
3473 }
3474 }
3475}
3476
3477pub async fn thread_start_request<G, T>(
3484 cfg: &Config,
3485 guest: &mut G,
3486 detpid: DetPid,
3487) -> Option<ThreadHistory>
3488where
3489 G: Guest<Detcore<T>>,
3490 T: RecordOrReplay,
3491{
3492 let dettid = guest.thread_state().dettid;
3493 if cfg.sequentialize_threads {
3494 trace!("[detcore, dtid {}] new thread BLOCKING on rpc...", &dettid);
3495 let physical_ids = guest
3496 .thread_state()
3497 .physical_tid
3498 .map(|tid| (guest.pid().as_raw(), tid));
3499 let signal_identity = guest.signal_task_identity();
3500 guest.thread_state_mut().signal_task_identity = signal_identity;
3501 let resp = send_and_update_time(
3502 guest,
3503 GlobalRequest::StartNewThread(dettid, detpid, physical_ids, signal_identity),
3504 )
3505 .await;
3506 match resp.1 {
3507 GlobalResponse::StartNewThread(preempts) => {
3508 trace!("[detcore, dtid {}] new thread UNBLOCKED (post-rpc)", dettid);
3509 preempts
3510 }
3511 _ => unreachable!(),
3512 }
3513 } else {
3514 None
3515 }
3516}
3517
3518pub(crate) fn child_tid_clear_address(flags: CloneFlags, address: usize) -> usize {
3521 if flags.contains(CloneFlags::CLONE_CHILD_CLEARTID) {
3522 address
3523 } else {
3524 0
3525 }
3526}
3527
3528pub async fn create_child_thread<G, T>(
3534 guest: &mut G,
3535 child_dettid: DetTid,
3536 ctid: usize,
3537 flags: Option<CloneFlags>,
3538 exit_signal: libc::c_int,
3539 physical_ids: Option<(i32, i32)>,
3540) -> Option<MmId>
3541where
3542 G: Guest<Detcore<T>>,
3543 T: RecordOrReplay,
3544{
3545 let starting_priority = if guest.config().replay_preemptions_from.is_some() {
3547 None
3550 } else if guest.config().replay_schedule_from.is_some() {
3551 if child_dettid <= DetTid::from_raw(3) {
3553 Some(REPLAY_FOREGROUND_PRIORITY)
3554 } else {
3555 Some(REPLAY_DEFERRED_PRIORITY)
3556 }
3557 } else if guest.config().chaos {
3558 let entropy = guest
3559 .thread_state_mut()
3560 .chaos_prng_next_u64("child_priority");
3561 if guest.config().chaos_target_races {
3562 if entropy.is_multiple_of(2) {
3568 Some(FIRST_PRIORITY)
3569 } else {
3570 Some(LAST_PRIORITY)
3571 }
3572 } else {
3573 Some(entropy_to_priority(entropy))
3574 }
3575 } else {
3576 Some(DEFAULT_PRIORITY)
3577 };
3578
3579 let detpid = guest.thread_state().detpid.expect("detpid unset");
3580
3581 let resp = send_and_update_time(
3582 guest,
3583 GlobalRequest::CreateChildThread(
3584 child_dettid,
3585 detpid,
3586 ctid,
3587 flags,
3588 exit_signal,
3589 physical_ids,
3590 starting_priority,
3591 ),
3592 )
3593 .await;
3594 match resp.1 {
3595 GlobalResponse::CreateChildThread(x) => x,
3596 _ => unreachable!(),
3597 }
3598}
3599
3600pub async fn create_vfork_child_thread<G, T>(
3608 guest: &mut G,
3609 child_dettid: DetTid,
3610 vfork: crate::tool_local::PendingVfork,
3611) where
3612 G: Guest<Detcore<T>>,
3613 T: RecordOrReplay,
3614{
3615 let starting_priority = if guest.config().replay_preemptions_from.is_some() {
3616 None
3617 } else if guest.config().replay_schedule_from.is_some() {
3618 Some(if child_dettid <= DetTid::from_raw(3) {
3619 REPLAY_FOREGROUND_PRIORITY
3620 } else {
3621 REPLAY_DEFERRED_PRIORITY
3622 })
3623 } else if guest.config().chaos {
3624 Some(entropy_to_priority(vfork.child_priority_entropy.expect(
3625 "vfork child priority entropy missing in chaos mode",
3626 )))
3627 } else {
3628 Some(DEFAULT_PRIORITY - 1)
3634 };
3635
3636 let resp = send_and_update_time(
3637 guest,
3638 GlobalRequest::CreateVforkChildThread(
3639 vfork.parent_dettid,
3640 vfork.parent_detpid,
3641 child_dettid,
3642 vfork.child_tid_addr,
3643 vfork.flags,
3644 vfork.exit_signal,
3645 starting_priority,
3646 ),
3647 )
3648 .await;
3649 match resp.1 {
3650 GlobalResponse::CreateChildThread(_) => (),
3651 _ => unreachable!(),
3652 }
3653}
3654
3655pub(crate) async fn deregister_thread<R>(
3660 threads_time: DetTime,
3661 cfg: &Config,
3662 reverie: &R,
3663 thread: ThreadDeregistration,
3664) where
3665 R: GlobalRPC<GlobalState>,
3667{
3668 if cfg.sequentialize_threads {
3669 let mm = thread.mm;
3670 let resp = reverie
3672 .send_rpc((threads_time, mm, GlobalRequest::DeregisterThread(thread)))
3673 .await;
3674 match resp.1 {
3676 GlobalResponse::DeregisterThread(x) => x,
3677 _ => unreachable!(),
3678 }
3679 }
3680}
3681
3682pub(crate) async fn acknowledge_robust_list_exit_time<R>(
3687 threads_time: DetTime,
3688 reverie: &R,
3689 mm: MmId,
3690) -> bool
3691where
3692 R: GlobalRPC<GlobalState>,
3693{
3694 let response = reverie
3695 .send_rpc((threads_time, mm, GlobalRequest::RobustListWakes(Vec::new())))
3696 .await;
3697 match response.1 {
3698 GlobalResponse::RobustListWakes(counts) => {
3699 assert!(
3700 counts.is_empty(),
3701 "an empty exit-clock acknowledgement woke a waiter"
3702 );
3703 true
3704 }
3705 GlobalResponse::ThreadExited => false,
3706 _ => unreachable!(),
3707 }
3708}
3709
3710pub(crate) async fn robust_list_wakes_after_exit<R>(
3714 threads_time: DetTime,
3715 reverie: &R,
3716 mm: MmId,
3717 wakes: Vec<(DetTid, RobustListWake)>,
3718) -> Vec<u64>
3719where
3720 R: GlobalRPC<GlobalState>,
3721{
3722 if wakes.is_empty() {
3723 return Vec::new();
3724 }
3725 let response = reverie
3726 .send_rpc((
3727 threads_time,
3728 mm,
3729 GlobalRequest::RobustListWakes(
3730 wakes
3731 .into_iter()
3732 .map(|(owner, wake)| (owner, wake.futex))
3733 .collect(),
3734 ),
3735 ))
3736 .await;
3737 match response.1 {
3738 GlobalResponse::RobustListWakes(counts) => counts,
3739 _ => unreachable!(),
3740 }
3741}
3742
3743#[derive(PartialEq, Debug, Eq, Clone, Copy, Serialize, Deserialize)]
3745pub enum FutexAction {
3746 WaitRequest(Option<LogicalTime>),
3748 WaitFinished,
3750 WakeRequest(i32),
3752 WakeFinished(i32),
3754}
3755
3756pub async fn futex_action<G, T>(
3759 guest: &mut G,
3760 futex_action: FutexAction,
3761 futexid: &FutexID,
3762 init_read: i32,
3763 mask: u32,
3764) -> Option<SchedValue>
3765where
3766 G: Guest<Detcore<T>>,
3767 T: RecordOrReplay,
3768{
3769 assert!(guest.config().sequentialize_threads);
3770 let dettid = guest.thread_state().dettid;
3771 let req = GlobalRequest::FutexAction(dettid, futex_action, *futexid, init_read, mask);
3772 trace!(
3773 "BLOCKING on futex_action: sending request to scheduler: {:?}",
3774 req
3775 );
3776 let resp = send_and_update_time(guest, req.clone()).await;
3778 match resp.1 {
3779 GlobalResponse::FutexAction(answer) => {
3780 trace!("UNBLOCKING after futex_action. Request was: {:?}", req);
3781 answer
3782 }
3783 _ => unreachable!(),
3784 }
3785}
3786
3787pub async fn determinize_inode<G, T>(guest: &mut G, inode: RawInode) -> (DetInode, LogicalTime)
3795where
3796 G: Guest<Detcore<T>>,
3797 T: RecordOrReplay,
3798{
3799 determinize_inode_observing_mtime(guest, inode, ObservedMtime::Unobserved).await
3800}
3801
3802pub async fn determinize_inode_observing_mtime<G, T>(
3807 guest: &mut G,
3808 inode: RawInode,
3809 observed: ObservedMtime,
3810) -> (DetInode, LogicalTime)
3811where
3812 G: Guest<Detcore<T>>,
3813 T: RecordOrReplay,
3814{
3815 let resp = send_and_update_time(guest, GlobalRequest::DeterminizeInode(inode, observed)).await;
3816 match resp.1 {
3817 GlobalResponse::DeterminizeInode(x) => x,
3818 _ => unreachable!(),
3819 }
3820}
3821
3822pub async fn determinize_device<G, T>(guest: &mut G, raw_device: u64) -> u64
3826where
3827 G: Guest<Detcore<T>>,
3828 T: RecordOrReplay,
3829{
3830 let resp = send_and_update_time(guest, GlobalRequest::DeterminizeDevice(raw_device)).await;
3831 match resp.1 {
3832 GlobalResponse::DeterminizeDevice(x) => x,
3833 _ => unreachable!(),
3834 }
3835}
3836
3837pub async fn determinize_mount_id<G, T>(
3839 guest: &mut G,
3840 raw_mount_id: u64,
3841 mountinfo_order: Option<Vec<u64>>,
3842) -> Option<u64>
3843where
3844 G: Guest<Detcore<T>>,
3845 T: RecordOrReplay,
3846{
3847 let resp = send_and_update_time(
3848 guest,
3849 GlobalRequest::DeterminizeMountId(raw_mount_id, mountinfo_order),
3850 )
3851 .await;
3852 match resp.1 {
3853 GlobalResponse::DeterminizeMountId(value) => value,
3854 _ => unreachable!(),
3855 }
3856}
3857
3858pub async fn validate_mountinfo_identity_order<G, T>(
3860 guest: &mut G,
3861 mountinfo_order: Vec<u64>,
3862) -> bool
3863where
3864 G: Guest<Detcore<T>>,
3865 T: RecordOrReplay,
3866{
3867 let resp =
3868 send_and_update_time(guest, GlobalRequest::ValidateMountIdOrder(mountinfo_order)).await;
3869 match resp.1 {
3870 GlobalResponse::ValidateMountIdOrder(valid) => valid,
3871 _ => unreachable!(),
3872 }
3873}
3874
3875#[allow(unused)]
3877pub async fn unlink_inode<G, T>(guest: &mut G, d_ino: DetInode)
3878where
3879 G: Guest<Detcore<T>>,
3880 T: RecordOrReplay,
3881{
3882 let resp = send_and_update_time(guest, GlobalRequest::UnlinkInode(d_ino)).await;
3883 match resp.1 {
3884 GlobalResponse::UnlinkInode(x) => x,
3885 _ => unreachable!(),
3886 }
3887}
3888
3889pub async fn touch_file<G, T>(guest: &mut G, inode: RawInode)
3892where
3893 G: Guest<Detcore<T>>,
3894 T: RecordOrReplay,
3895{
3896 let resp = send_and_update_time(guest, GlobalRequest::TouchFile(inode)).await;
3897 match resp.1 {
3898 GlobalResponse::TouchFile(x) => x,
3899 _ => unreachable!(),
3900 }
3901}
3902
3903pub async fn set_file_mtime<G, T>(guest: &mut G, inode: RawInode, mtime: LogicalTime)
3905where
3906 G: Guest<Detcore<T>>,
3907 T: RecordOrReplay,
3908{
3909 let resp = send_and_update_time(guest, GlobalRequest::SetFileMtime(inode, mtime)).await;
3910 match resp.1 {
3911 GlobalResponse::SetFileMtime(x) => x,
3912 _ => unreachable!(),
3913 }
3914}
3915
3916pub async fn global_time_lower_bound<G, T>(guest: &mut G) -> LogicalTime
3918where
3919 G: Guest<Detcore<T>>,
3920 T: RecordOrReplay,
3921{
3922 let resp = send_and_update_time(guest, GlobalRequest::GlobalTimeLowerBound).await;
3923 match resp.1 {
3924 GlobalResponse::GlobalTimeLowerBound(x) => x,
3925 _ => unreachable!(),
3926 }
3927}
3928
3929pub async fn thread_observe_time<G, T>(guest: &mut G) -> LogicalTime
3933where
3934 G: Guest<Detcore<T>>,
3935 T: RecordOrReplay,
3936{
3937 global_time_lower_bound(guest).await
3938}
3939
3940fn write_backtrace<G, T>(guest: &mut G, m_path: Option<&PathBuf>)
3942where
3943 G: Guest<Detcore<T>>,
3944 T: RecordOrReplay,
3945{
3946 if let Some(backtrace) = guest.backtrace() {
3947 if let Some(path) = m_path {
3948 let file = File::create(path).expect("Failed to open preemption stacktrace log file");
3949 serde_json::to_writer(file, &backtrace.force_pretty()).unwrap();
3950 } else {
3951 eprintln!("{}", backtrace.force_pretty());
3952 }
3953 } else {
3954 warn!("Could not read backtrace!");
3955 }
3956}
3957
3958#[derive(PartialEq, Debug, Eq, Clone, Serialize, Deserialize)]
3960pub struct TraceSchedEventResponse {
3961 print_stack_strace: MaybePrintStack,
3962 timeslice: Option<LogicalTime>,
3963}
3964
3965struct SchedEventForLog<'a> {
3966 event: &'a SchedEvent,
3967 command_bootstrap: bool,
3968}
3969
3970struct CommandBootstrapInstructionPointer(NonZeroUsize);
3971
3972impl std::fmt::Debug for CommandBootstrapInstructionPointer {
3973 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
3974 write!(f, "{}", crate::logdiff::host_addr(self.0.get()))
3975 }
3976}
3977
3978impl std::fmt::Debug for SchedEventForLog<'_> {
3979 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
3980 if !self.command_bootstrap {
3981 return std::fmt::Debug::fmt(self.event, f);
3982 }
3983 f.debug_struct("SchedEvent")
3984 .field("dettid", &self.event.dettid)
3985 .field("op", &self.event.op)
3986 .field("count", &self.event.count)
3987 .field(
3988 "start_rip",
3989 &self.event.start_rip.map(CommandBootstrapInstructionPointer),
3990 )
3991 .field(
3992 "end_rip",
3993 &self.event.end_rip.map(CommandBootstrapInstructionPointer),
3994 )
3995 .field("end_time", &self.event.end_time)
3996 .finish()
3997 }
3998}
3999
4000pub async fn trace_schedevent<G, T>(guest: &mut G, ev: SchedEvent, tag_end_rip: bool)
4007where
4008 G: Guest<Detcore<T>>,
4009 T: RecordOrReplay,
4010{
4011 assert!(guest.config().sequentialize_threads);
4012
4013 let ev = if tag_end_rip {
4015 let end_rip = if let Some(r) = ev.end_rip {
4016 r
4017 } else {
4018 let regs = guest.regs().await;
4019 NonZeroUsize::new(regs.rip.try_into().unwrap()).unwrap()
4020 };
4021 SchedEvent {
4022 end_rip: Some(end_rip),
4023 ..ev
4024 }
4025 } else {
4026 ev
4027 };
4028
4029 if tracing::enabled!(tracing::Level::TRACE)
4030 && let Some(rip) = ev.end_rip
4031 {
4032 let rip_addr = AddrMut::<u16>::from_raw(rip.into()).unwrap();
4033 let rip_contents = guest.memory().read_value(rip_addr);
4037 trace!(
4038 "Tracing sched event, after which rip is {}, next two instruction bytes {:?}",
4039 rip, rip_contents
4040 );
4041 }
4042
4043 let detpid = guest.thread_state().detpid.expect("detpid unset");
4044 let command_bootstrap = guest.is_command_bootstrap();
4045 let resp = send_and_update_time(
4046 guest,
4047 GlobalRequest::TraceSchedEvent(ev, detpid, command_bootstrap),
4048 )
4049 .await;
4050
4051 trace!("trace_schedevent result: {:?}", resp);
4052 match resp {
4053 (
4054 _,
4055 GlobalResponse::TraceSchedEvent(TraceSchedEventResponse {
4056 print_stack_strace,
4057 timeslice,
4058 }),
4059 ) => {
4060 if let Some(m_path) = print_stack_strace {
4061 trace!("[trace_schedevent] writing stacktrace via Reverie...");
4062 write_backtrace(guest, m_path.as_ref());
4063 }
4064
4065 if let Some(timeslice) = timeslice
4066 && guest.thread_state().past_global_first_execve
4067 {
4068 let end_of_timeslice =
4069 guest.thread_state().thread_logical_time.as_nanos() + timeslice;
4070 trace!(
4071 "[detcore][dettid {}] setting end_of_timeslice to {:?} as instructed by replayer",
4072 guest.thread_state().dettid,
4073 end_of_timeslice
4074 );
4075 guest.thread_state_mut().end_of_timeslice = Some(end_of_timeslice);
4076 if guest.config().max_timeslice.is_some() {
4077 guest.thread_state_mut().max_timeslice_end = Some(end_of_timeslice);
4078 }
4079 }
4080 }
4081 _ => {
4082 unreachable!()
4083 }
4084 }
4085}
4086
4087pub async fn register_alarm<G, T>(
4093 guest: &mut G,
4094 duration: LogicalTime,
4095 interval: LogicalTime,
4096 sig: Signal,
4097) -> (LogicalTime, LogicalTime)
4098where
4099 G: Guest<Detcore<T>>,
4100 T: RecordOrReplay,
4101{
4102 let dettid = guest.thread_state().dettid;
4103 let detpid = guest.thread_state().detpid.expect("detpid unset");
4104 let resp = send_and_update_time(
4105 guest,
4106 GlobalRequest::RegisterAlarm(detpid, dettid, duration, interval, SigWrapper::from(sig)),
4107 )
4108 .await;
4109 match resp.1 {
4110 GlobalResponse::RegisterAlarm(x) => x,
4111 _ => unreachable!(),
4112 }
4113}
4114
4115pub async fn alarm_remaining<G, T>(guest: &mut G) -> ItimerSnapshot
4119where
4120 G: Guest<Detcore<T>>,
4121 T: RecordOrReplay,
4122{
4123 let detpid = guest.thread_state().detpid.expect("detpid unset");
4124 let resp = send_and_update_time(guest, GlobalRequest::AlarmRemaining(detpid)).await;
4125 match resp.1 {
4126 GlobalResponse::AlarmRemaining(remaining) => remaining,
4127 _ => unreachable!(),
4128 }
4129}
4130
4131pub async fn register_posix_timer<G, T>(
4135 guest: &mut G,
4136 timer_id: i32,
4137 deadline: Option<LogicalTime>,
4138 interval: LogicalTime,
4139 sig: Signal,
4140) where
4141 G: Guest<Detcore<T>>,
4142 T: RecordOrReplay,
4143{
4144 let dettid = guest.thread_state().dettid;
4145 let detpid = guest.thread_state().detpid.expect("detpid unset");
4146 let resp = send_and_update_time(
4147 guest,
4148 GlobalRequest::RegisterPosixTimer(
4149 detpid,
4150 dettid,
4151 timer_id,
4152 deadline,
4153 interval,
4154 SigWrapper::from(sig),
4155 ),
4156 )
4157 .await;
4158 match resp.1 {
4159 GlobalResponse::RegisterPosixTimer(()) => {}
4160 _ => unreachable!(),
4161 }
4162}
4163
4164pub async fn thread_is_live<G, T>(guest: &mut G, dettid: DetTid) -> bool
4174where
4175 G: Guest<Detcore<T>>,
4176 T: RecordOrReplay,
4177{
4178 let response = send_and_update_time(guest, GlobalRequest::ThreadIsLive(dettid)).await;
4179 match response.1 {
4180 GlobalResponse::ThreadIsLive(live) => live,
4181 _ => unreachable!(),
4182 }
4183}
4184
4185pub async fn exact_child_wait_state<G, T>(guest: &mut G, child: DetPid) -> ExactChildWaitState
4187where
4188 G: Guest<Detcore<T>>,
4189 T: RecordOrReplay,
4190{
4191 let parent = guest.thread_state().detpid.expect("detpid unset");
4192 let response =
4193 send_and_update_time(guest, GlobalRequest::ExactChildWaitState(parent, child)).await;
4194 match response.1 {
4195 GlobalResponse::ExactChildWaitState(state) => state,
4196 _ => unreachable!(),
4197 }
4198}
4199
4200pub async fn await_exact_child_physical_exit<G, T>(
4202 guest: &mut G,
4203 child: DetPid,
4204) -> ExactChildWaitState
4205where
4206 G: Guest<Detcore<T>>,
4207 T: RecordOrReplay,
4208{
4209 let mut state = exact_child_wait_state(guest, child).await;
4210 if matches!(
4211 state,
4212 ExactChildWaitState::PhysicalExitPending | ExactChildWaitState::PhysicallyExited
4213 ) {
4214 let dettid = guest.thread_state().dettid;
4215 let mut resources = Resources::new(dettid);
4216 resources.insert(ResourceID::WaitPhysicalChild(child), Permission::R);
4217 resources.fyi("wait-child-physical-exit");
4218 let _ = resource_request(guest, resources).await;
4219 state = exact_child_wait_state(guest, child).await;
4220 }
4221 state
4222}
4223
4224pub async fn wait_for_child_lifecycle<G, T>(guest: &mut G, spec: ChildWaitSpec) -> ResumeStatus
4226where
4227 G: Guest<Detcore<T>>,
4228 T: RecordOrReplay,
4229{
4230 let dettid = guest.thread_state().dettid;
4231 let parent = guest.thread_state().detpid.expect("detpid unset");
4232 let mut resources = Resources::new(dettid);
4233 resources.insert(ResourceID::WaitChild { parent, spec }, Permission::R);
4234 resources.fyi("wait-child-lifecycle");
4235 resource_request(guest, resources).await
4236}
4237
4238pub async fn ready_child_wait<G, T>(guest: &mut G, spec: ChildWaitSpec) -> (Option<DetPid>, bool)
4239where
4240 G: Guest<Detcore<T>>,
4241 T: RecordOrReplay,
4242{
4243 let parent = guest.thread_state().detpid.expect("detpid unset");
4244 let response = send_and_update_time(guest, GlobalRequest::ReadyChildWait(parent, spec)).await;
4245 match response.1 {
4246 GlobalResponse::ReadyChildWait(snapshot) => snapshot,
4247 _ => unreachable!(),
4248 }
4249}
4250
4251pub async fn process_group<G, T>(guest: &mut G, process: DetPid) -> Option<DetPid>
4252where
4253 G: Guest<Detcore<T>>,
4254 T: RecordOrReplay,
4255{
4256 let response = send_and_update_time(guest, GlobalRequest::ProcessGroup(process)).await;
4257 match response.1 {
4258 GlobalResponse::ProcessGroup(group) => group,
4259 _ => unreachable!(),
4260 }
4261}
4262
4263pub async fn set_process_group<G, T>(guest: &mut G, process: DetPid, group: DetPid) -> bool
4264where
4265 G: Guest<Detcore<T>>,
4266 T: RecordOrReplay,
4267{
4268 let response =
4269 send_and_update_time(guest, GlobalRequest::SetProcessGroup(process, group)).await;
4270 match response.1 {
4271 GlobalResponse::SetProcessGroup(updated) => updated,
4272 _ => unreachable!(),
4273 }
4274}
4275
4276pub async fn create_session<G, T>(guest: &mut G, process: DetPid) -> bool
4277where
4278 G: Guest<Detcore<T>>,
4279 T: RecordOrReplay,
4280{
4281 let response = send_and_update_time(guest, GlobalRequest::CreateSession(process)).await;
4282 match response.1 {
4283 GlobalResponse::CreateSession(updated) => updated,
4284 _ => unreachable!(),
4285 }
4286}
4287
4288pub async fn consume_child_wait<G, T>(guest: &mut G, child: DetPid) -> bool
4289where
4290 G: Guest<Detcore<T>>,
4291 T: RecordOrReplay,
4292{
4293 let parent = guest.thread_state().detpid.expect("detpid unset");
4294 let response =
4295 send_and_update_time(guest, GlobalRequest::ConsumeChildWait(parent, child)).await;
4296 match response.1 {
4297 GlobalResponse::ConsumeChildWait(consumed) => consumed,
4298 _ => unreachable!(),
4299 }
4300}
4301
4302pub async fn resolve_kill_targets<G, T>(guest: &mut G, detpid: DetPid) -> Vec<DetTid>
4303where
4304 G: Guest<Detcore<T>>,
4305 T: RecordOrReplay,
4306{
4307 let response = send_and_update_time(guest, GlobalRequest::ResolveKillTargets(detpid)).await;
4308 match response.1 {
4309 GlobalResponse::ResolveKillTargets(targets) => targets,
4310 _ => unreachable!(),
4311 }
4312}
4313
4314fn alarm_signal(sig: SigWrapper) -> Signal {
4324 sig.signal().unwrap_or_else(|| {
4325 panic!(
4326 "timer registration received unnameable signal {}",
4327 sig.raw()
4328 )
4329 })
4330}
4331
4332pub async fn notify_signal_pending<G, T>(
4333 guest: &mut G,
4334 dettid: DetTid,
4335 signal: SigWrapper,
4336 target_process: Option<DetPid>,
4337) where
4338 G: Guest<Detcore<T>>,
4339 T: RecordOrReplay,
4340{
4341 let response = send_and_update_time(
4342 guest,
4343 GlobalRequest::NotifySignalPending(dettid, signal, target_process),
4344 )
4345 .await;
4346 match response.1 {
4347 GlobalResponse::NotifySignalPending(()) => {}
4348 _ => unreachable!(),
4349 }
4350}
4351
4352pub async fn unrecoverable_shutdown<G, T>(guest: &G, status: i32) -> !
4368where
4369 G: Guest<Detcore<T>>,
4370 T: RecordOrReplay,
4371{
4372 if cfg!(debug_assertions) {
4373 let mytime = guest.thread_state().thread_logical_time.clone();
4374 let mm = guest.thread_state().mm_id;
4375 let _ = guest
4377 .send_rpc((mytime, mm, GlobalRequest::UnrecoverableShutdown))
4378 .await;
4379 }
4380
4381 std::process::exit(status);
4400}
4401
4402#[cfg(test)]
4403mod tests {
4404 #[test]
4409 fn external_scheduler_announces_before_its_future_is_polled() {
4410 #[derive(Clone, Default)]
4411 struct Capture(std::sync::Arc<Mutex<Vec<String>>>);
4412 struct Visitor(String);
4413 impl tracing::field::Visit for Visitor {
4414 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
4415 if field.name() == "message" {
4416 use std::fmt::Write;
4417 write!(self.0, "{value:?}").unwrap();
4418 }
4419 }
4420 }
4421 impl tracing::Subscriber for Capture {
4422 fn enabled(&self, metadata: &tracing::Metadata<'_>) -> bool {
4423 *metadata.level() == tracing::Level::INFO
4424 }
4425 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
4426 tracing::span::Id::from_u64(1)
4427 }
4428 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
4429 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
4430 fn event(&self, event: &tracing::Event<'_>) {
4431 let mut visitor = Visitor(String::new());
4432 event.record(&mut visitor);
4433 self.0.lock().unwrap().push(visitor.0);
4434 }
4435 fn enter(&self, _: &tracing::span::Id) {}
4436 fn exit(&self, _: &tracing::span::Id) {}
4437 }
4438 let config = Config {
4439 sequentialize_threads: true,
4440 ..Config::default()
4441 };
4442 let state = GlobalState::init_for_external_scheduler(&config);
4443 let captured = Capture::default();
4444 let scheduler = tracing::subscriber::with_default(captured.clone(), || {
4445 state.run_external_scheduler(std::sync::Arc::new(|_| {}))
4446 });
4447 let announcement = "[scheduler] daemon task starting up, waiting for guest thread start..";
4448 assert_eq!(*captured.0.lock().unwrap(), [announcement]);
4449 drop(scheduler);
4450 }
4451
4452 #[test]
4453 fn schedule_event_host_markers_require_command_bootstrap_provenance() {
4454 let event = SchedEvent::branches(DetTid::from_raw(3), 223)
4455 .with_end_rip(std::num::NonZeroUsize::new(0x1234).unwrap())
4456 .with_time(LogicalTime::from_nanos(2230));
4457 let original = serde_json::to_string(&event).unwrap();
4458 let plain = format!(
4459 "{:?}",
4460 super::SchedEventForLog {
4461 event: &event,
4462 command_bootstrap: false
4463 }
4464 );
4465 assert_eq!(plain, format!("{event:?}"));
4466 let marked = format!(
4467 "{:?}",
4468 super::SchedEventForLog {
4469 event: &event,
4470 command_bootstrap: true
4471 }
4472 );
4473 assert_eq!(
4474 marked,
4475 "SchedEvent { dettid: DetPid(3), op: Branch, count: 223, start_rip: None, end_rip: Some(<hostaddr 0x1234>), end_time: Some(LogicalTime(2230)) }"
4476 );
4477 assert_eq!(serde_json::to_string(&event).unwrap(), original);
4478 }
4479
4480 #[test]
4481 fn summary_preemption_views_keep_counts_full_report_and_single_flush() {
4482 let (_config, state, tid, _) = cancellation_test_state();
4483 let directory = tempfile::tempdir().unwrap();
4484 let path = directory.path().join("recording");
4485 let mut writer = crate::preemptions::PreemptionWriter::new(Some(path.clone()));
4486 writer.register_thread(tid, DEFAULT_PRIORITY);
4487 writer.insert_reprioritization(tid, LogicalTime::from_nanos(100), 10, DEFAULT_PRIORITY, 20);
4488 let mut scheduler = state.sched.lock().unwrap();
4489 scheduler.preemption_writer = Some(writer);
4490 let (summary, info_description) = scheduler
4491 .generate_partial_run_summary_for_log(Some(&path))
4492 .unwrap();
4493 assert!(scheduler.preemption_writer.is_none());
4494 let recorded = std::fs::read(&path).unwrap();
4495 let parsed: serde_json::Value = serde_json::from_slice(&recorded).unwrap();
4496 assert_eq!(
4497 parsed["per_thread"][tid.to_string()]["prio_changes"],
4498 serde_json::json!([[100, DEFAULT_PRIORITY]])
4499 );
4500 let count_line = "Record of 1 preemption and reprioritization events:\n";
4501 assert_eq!(info_description.as_deref(), Some(count_line));
4502 assert_eq!(
4503 summary.reprio_descrip.as_deref(),
4504 Some(format!("{count_line} (Writing to file {path:?})\n").as_str())
4505 );
4506 let full = summary.to_string();
4507 let json = serde_json::to_vec(&summary).unwrap();
4508 let info = summary.info(info_description.as_deref()).to_string();
4509 assert!(info.contains(count_line));
4510 assert!(!info.contains("Writing to file"));
4511 assert!(full.contains(&format!("Writing to file {path:?}")));
4512 let _ = scheduler.generate_partial_run_summary(None).unwrap();
4514 assert_eq!(std::fs::read(&path).unwrap(), recorded);
4515 assert_eq!(summary.to_string(), full);
4516 assert_eq!(serde_json::to_vec(&summary).unwrap(), json);
4517 }
4518
4519 #[test]
4520 fn run_summary_info_keeps_semantics_and_debug_retains_bookkeeping() {
4521 #[derive(Clone)]
4522 struct Capture(std::sync::Arc<Mutex<Vec<(tracing::Level, String)>>>);
4523 struct Visitor(String);
4524 impl tracing::field::Visit for Visitor {
4525 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
4526 use std::fmt::Write;
4527 write!(self.0, "{}={value:?};", field.name()).unwrap();
4528 }
4529 }
4530 impl tracing::Subscriber for Capture {
4531 fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
4532 true
4533 }
4534 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
4535 tracing::span::Id::from_u64(1)
4536 }
4537 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
4538 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
4539 fn event(&self, event: &tracing::Event<'_>) {
4540 let mut visitor = Visitor(String::new());
4541 event.record(&mut visitor);
4542 self.0
4543 .lock()
4544 .unwrap()
4545 .push((*event.metadata().level(), visitor.0));
4546 }
4547 fn enter(&self, _: &tracing::span::Id) {}
4548 fn exit(&self, _: &tracing::span::Id) {}
4549 }
4550 let summary = super::RunSummary {
4551 sched_turns: 4,
4552 schedevent_recorded: 11,
4553 schedevent_replayed: 11,
4554 schedevent_desynced: 2,
4555 desync_descrip: Some("two real desyncs\n".into()),
4556 num_processes: 1,
4557 num_threads: 1,
4558 threads_descrip: "[3]".into(),
4559 syscalls: Some(3),
4560 virttime_elapsed: 2_512_380,
4561 virttime_final: 2_512_380,
4562 timeslice_stats: TimesliceStats {
4563 count: 1,
4564 sum_ns: 512_380,
4565 min_ns: 512_380,
4566 max_ns: 512_380,
4567 },
4568 ..Default::default()
4569 };
4570 let description = "Record of 7 preemption and reprioritization events:\n";
4571 let captured = Capture(Default::default());
4572 tracing::subscriber::with_default(captured.clone(), || {
4573 super::log_run_summary(
4574 "report",
4575 &summary,
4576 Some(description),
4577 Some(std::path::Path::new("/host/recording")),
4578 );
4579 });
4580 let events = captured.0.lock().unwrap();
4581 assert_eq!(events.len(), 2);
4582 assert_eq!(events[0].0, tracing::Level::INFO);
4583 assert_eq!(events[1].0, tracing::Level::DEBUG);
4584 let info = &events[0].1;
4585 for text in [
4586 "1 group leaders of 1 thread(s)",
4587 "3 syscalls",
4588 "4 turns, recorded 11 events (2 desynced)",
4589 "two real desyncs",
4590 "Record of 7 preemption",
4591 "2_512_380ns",
4592 "min=512380ns max=512380ns mean=512380ns count=1",
4593 ] {
4594 assert!(info.contains(text), "missing semantic field {text}: {info}");
4595 }
4596 assert!(!info.contains("/host/recording"));
4597 assert!(!info.contains("replayed"));
4598 assert!(events[1].1.contains("replayed_events=11"));
4599 assert!(events[1].1.contains("/host/recording"));
4600 let baseline = summary.info(Some(description)).to_string();
4601 let mutations: [fn(&mut super::RunSummary); 8] = [
4602 |s| s.sched_turns += 1,
4603 |s| s.schedevent_recorded += 1,
4604 |s| s.schedevent_desynced += 1,
4605 |s| s.num_threads += 1,
4606 |s| s.syscalls = Some(4),
4607 |s| s.virttime_elapsed += 1,
4608 |s| s.timeslice_stats.count += 1,
4609 |s| s.threads_descrip.push_str(",4"),
4610 ];
4611 for mutate in mutations {
4612 let mut changed = summary.clone();
4613 mutate(&mut changed);
4614 assert_ne!(changed.info(Some(description)).to_string(), baseline);
4615 }
4616 assert_ne!(
4617 summary
4618 .info(Some(
4619 "Record of 8 preemption and reprioritization events:\n"
4620 ))
4621 .to_string(),
4622 baseline
4623 );
4624 }
4625
4626 mod backend_failure_tests;
4627 mod exec_identity_tests;
4628 mod exec_teardown_tests;
4629 use std::collections::BTreeSet;
4630 use std::os::fd::AsRawFd;
4631 use std::os::fd::FromRawFd;
4632 use std::os::fd::OwnedFd;
4633 use std::sync::Mutex;
4634 use std::time::Duration;
4635
4636 use nix::sys::signal::Signal;
4637 use reverie::GlobalRPC;
4638 use reverie::GlobalTool;
4639 use reverie::Guest;
4640 use reverie::Tid;
4641 use reverie::syscalls::CloneFlags;
4642
4643 use super::FutexAction;
4644 use super::GlobalRequest;
4645 use super::GlobalResponse;
4646 use super::GlobalState;
4647 use super::MountIdPool;
4648 use super::PendingExecState;
4649 use super::ResumeStatus;
4650 use super::RpcIncarnation;
4651 use super::SchedulerRpcResult;
4652 use super::SigWrapper;
4653 use super::ThreadDeregistration;
4654 use super::TimesliceStats;
4655 use super::format_unsupported_syscall_warning;
4656 use crate::Detcore;
4657 use crate::config::Config;
4658 use crate::config::RunsPostFork;
4659 use crate::ivar::Ivar;
4660 use crate::preemptions::PreemptionRecord;
4661 use crate::resources::ExternalOpId;
4662 use crate::resources::Permission;
4663 use crate::resources::ResourceID;
4664 use crate::resources::Resources;
4665 use crate::scheduler::DEFAULT_PRIORITY;
4666 use crate::scheduler::SchedRequest;
4667 use crate::scheduler::SchedResponse;
4668 use crate::scheduler::SchedValue;
4669 use crate::scheduler::ThreadNextTurn;
4670 use crate::tool_local::ExecFdBlockingOverrides;
4671 use crate::types::DetPid;
4672 use crate::types::DetTid;
4673 use crate::types::DetTime;
4674 use crate::types::FutexID;
4675 use crate::types::LogicalTime;
4676 use crate::types::MmId;
4677 use crate::types::Op;
4678 use crate::types::SchedEvent;
4679
4680 #[test]
4681 fn fdinfo_mount_ids_preserve_raw_equivalence_and_distinctness() {
4682 let mut pool = MountIdPool::from_config(&[10, 20, 30], true, &[]);
4683 assert_eq!(pool.determinize(20, None), Some(2));
4684
4685 assert_eq!(pool.determinize(700, None), Some(4));
4688 assert_eq!(pool.determinize(701, None), Some(5));
4689 assert_eq!(pool.determinize(702, None), Some(6));
4690 assert_eq!(pool.determinize(701, None), Some(5));
4691 }
4692
4693 #[test]
4694 fn fdinfo_mount_ids_seed_from_low_level_snapshot_and_refuse_drift() {
4695 let mut pool = MountIdPool::from_config(&[], false, &[]);
4696 assert_eq!(pool.determinize(20, Some(&[10, 20, 30])), Some(2));
4697 assert_eq!(pool.determinize(700, Some(&[10, 20, 30])), Some(4));
4698 assert_eq!(pool.determinize(20, Some(&[10, 20, 30])), Some(2));
4699 assert_eq!(pool.determinize(20, Some(&[10, 99, 30])), None);
4700 assert_eq!(pool.determinize(20, Some(&[10, 20])), None);
4701 assert_eq!(pool.determinize(20, Some(&[10, 20, 30, 40])), None);
4702 }
4703
4704 #[test]
4705 fn mountinfo_reads_seed_once_and_refuse_later_table_changes() {
4706 let mut pool = MountIdPool::from_config(&[], false, &[]);
4707 assert!(pool.validate_mountinfo_order(&[10, 20, 30]));
4708 assert!(pool.validate_mountinfo_order(&[10, 20, 30]));
4709 assert!(!pool.validate_mountinfo_order(&[10, 99, 30]));
4710
4711 let mut empty = MountIdPool::from_config(&[], false, &[]);
4712 assert!(empty.validate_mountinfo_order(&[]));
4713 assert!(!empty.validate_mountinfo_order(&[10]));
4714 }
4715
4716 #[test]
4717 fn configured_mountinfo_order_accepts_only_ordered_known_subsets() {
4718 let mut pool = MountIdPool::from_config(&[10, 20, 30, 40], true, &[]);
4719 assert!(pool.validate_mountinfo_order(&[10, 30, 40]));
4720 assert_eq!(pool.determinize(30, None), Some(3));
4721 assert!(pool.validate_mountinfo_order(&[20, 40]));
4722 assert!(!pool.validate_mountinfo_order(&[30, 20]));
4723 assert!(!pool.validate_mountinfo_order(&[10, 10]));
4724 }
4725
4726 #[test]
4727 fn configured_mountinfo_order_refuses_new_namespace_ids() {
4728 let mut pool = MountIdPool::from_config(&[10, 20, 30, 40], true, &[]);
4729 assert!(pool.validate_mountinfo_order(&[10, 30, 40]));
4730 assert!(
4731 !pool.validate_mountinfo_order(&[10, 99, 40]),
4732 "a namespace-local mount ID absent from captured provenance must fail closed"
4733 );
4734 }
4735
4736 #[test]
4737 fn fdinfo_mount_ids_refuse_malformed_configured_provenance() {
4738 let mut pool = MountIdPool::from_config(&[10, 10], true, &[]);
4739 assert_eq!(pool.determinize(10, None), None);
4740 }
4741
4742 #[test]
4743 fn raw_zero_is_one_mount_identity_without_descriptor_type_partitioning() {
4744 let mut pool = MountIdPool::from_config(&[10, 0], true, &[]);
4745 assert_eq!(pool.determinize(0, None), Some(0));
4750 assert_eq!(pool.determinize(20, None), Some(3));
4751 assert_eq!(pool.determinize(0, None), Some(0));
4752 }
4753
4754 #[test]
4755 fn recorded_unlisted_mount_order_rebuilds_the_same_mapping() {
4756 let mut recording = MountIdPool::from_config(&[], false, &[]);
4757 assert!(recording.validate_mountinfo_order(&[10, 20]));
4758 assert_eq!(recording.determinize(700, None), Some(3));
4759 assert_eq!(recording.determinize(701, None), Some(4));
4760 let provenance = recording.provenance().unwrap().unwrap();
4761
4762 let mut replay = MountIdPool::from_config(
4763 &provenance.mountinfo_order,
4764 true,
4765 &provenance.unlisted_order,
4766 );
4767 assert_eq!(replay.determinize(701, None), Some(4));
4768 assert_eq!(replay.determinize(700, None), Some(3));
4769 }
4770
4771 fn cancellation_test_state() -> (Config, GlobalState, DetTid, DetPid) {
4772 let config = Config {
4773 sequentialize_threads: true,
4774 cancel_killed_thread_rpcs: true,
4775 ..Config::default()
4776 };
4777 let state = GlobalState::initialize(&config, false);
4778 let dettid = DetTid::from_raw(17);
4779 let detpid = DetPid::from_raw(17);
4780 state
4781 .sched
4782 .lock()
4783 .unwrap()
4784 .thread_tree
4785 .add_child(dettid, dettid, true);
4786 (config, state, dettid, detpid)
4787 }
4788
4789 fn install_test_registration(state: &GlobalState, dettid: DetTid, request: Ivar<SchedRequest>) {
4790 let mut scheduler = state.sched.lock().unwrap();
4791 assert!(!scheduler.thread_is_logically_killed(dettid));
4792 scheduler.next_turns.insert(
4793 dettid,
4794 ThreadNextTurn {
4795 dettid,
4796 child_tid_addr: 0,
4797 req: request,
4798 resp: Ivar::new(),
4799 protocol: Default::default(),
4800 },
4801 );
4802 scheduler.priorities.insert(dettid, DEFAULT_PRIORITY);
4803 scheduler.runqueue_push_back(dettid);
4804 }
4805
4806 struct ExternalRegistrationGuest<'a> {
4809 global: &'a GlobalState,
4810 config: &'a Config,
4811 thread: crate::ThreadState<()>,
4812 requests: Mutex<Vec<GlobalRequest>>,
4813 post_exec: bool,
4814 }
4815
4816 struct ExternalRegistrationStack;
4817 struct ExternalRegistrationStackGuard;
4818
4819 impl Drop for ExternalRegistrationStackGuard {
4820 fn drop(&mut self) {}
4821 }
4822
4823 impl reverie::Stack for ExternalRegistrationStack {
4824 type StackGuard = ExternalRegistrationStackGuard;
4825
4826 fn size(&self) -> usize {
4827 panic!("external registration must not use a guest stack")
4828 }
4829
4830 fn capacity(&self) -> usize {
4831 panic!("external registration must not use a guest stack")
4832 }
4833
4834 fn push<'stack, T>(&mut self, _value: T) -> reverie::syscalls::Addr<'stack, T> {
4835 panic!("external registration must not use a guest stack")
4836 }
4837
4838 fn reserve<'stack, T>(&mut self) -> reverie::syscalls::AddrMut<'stack, T> {
4839 panic!("external registration must not use a guest stack")
4840 }
4841
4842 fn commit(self) -> Result<Self::StackGuard, reverie::syscalls::Errno> {
4843 panic!("external registration must not use a guest stack")
4844 }
4845 }
4846
4847 #[reverie::tool]
4848 impl GlobalRPC<GlobalState> for ExternalRegistrationGuest<'_> {
4849 async fn send_rpc(
4850 &self,
4851 message: <GlobalState as GlobalTool>::Request,
4852 ) -> <GlobalState as GlobalTool>::Response {
4853 self.requests.lock().unwrap().push(message.2.clone());
4854 self.global
4855 .receive_rpc(Tid::from_raw(self.thread.dettid.as_raw()), message)
4856 .await
4857 }
4858
4859 fn config(&self) -> &Config {
4860 self.config
4861 }
4862 }
4863
4864 #[reverie::tool]
4865 impl Guest<Detcore> for ExternalRegistrationGuest<'_> {
4866 type Memory = reverie::syscalls::LocalMemory;
4867 type Stack = ExternalRegistrationStack;
4868
4869 fn tid(&self) -> reverie::Pid {
4870 reverie::Pid::from_raw(self.thread.dettid.as_raw())
4871 }
4872
4873 fn pid(&self) -> reverie::Pid {
4874 reverie::Pid::from_raw(self.thread.detpid.unwrap().as_raw())
4875 }
4876
4877 fn ppid(&self) -> Option<reverie::Pid> {
4878 None
4879 }
4880
4881 fn auxv(&self) -> reverie::Auxv {
4882 assert!(self.post_exec, "external registration must not read auxv");
4883 reverie::Auxv::from_entries([])
4884 }
4885
4886 fn memory(&self) -> Self::Memory {
4887 panic!("external registration must not access guest memory")
4888 }
4889
4890 fn thread_state_mut(&mut self) -> &mut crate::ThreadState<()> {
4891 &mut self.thread
4892 }
4893
4894 fn thread_state(&self) -> &crate::ThreadState<()> {
4895 &self.thread
4896 }
4897
4898 async fn regs(&mut self) -> libc::user_regs_struct {
4899 assert!(
4900 self.post_exec,
4901 "external registration must not read guest registers"
4902 );
4903 unsafe { std::mem::zeroed() }
4906 }
4907
4908 async fn stack(&mut self) -> Self::Stack {
4909 panic!("external registration must not use a guest stack")
4910 }
4911
4912 async fn daemonize(&mut self) {
4913 panic!("external registration must not daemonize")
4914 }
4915
4916 async fn inject<S: reverie::syscalls::SyscallInfo>(
4917 &mut self,
4918 _syscall: S,
4919 ) -> Result<i64, reverie::syscalls::Errno> {
4920 panic!("external registration must not inject a syscall")
4921 }
4922
4923 async fn tail_inject<S: reverie::syscalls::SyscallInfo>(
4924 &mut self,
4925 _syscall: S,
4926 ) -> reverie::Never {
4927 panic!("external registration unexpectedly retired its live parent")
4928 }
4929
4930 fn set_timer(&mut self, _schedule: reverie::TimerSchedule) -> Result<(), reverie::Error> {
4931 panic!("external registration must not set a timer")
4932 }
4933
4934 fn set_timer_precise(
4935 &mut self,
4936 _schedule: reverie::TimerSchedule,
4937 ) -> Result<(), reverie::Error> {
4938 panic!("external registration must not set a timer")
4939 }
4940
4941 fn read_clock(&mut self) -> Result<u64, reverie::Error> {
4942 panic!("external registration must not read a host clock")
4943 }
4944 }
4945
4946 struct RetirementGuest<'a> {
4947 global: &'a GlobalState,
4948 config: &'a Config,
4949 thread: crate::ThreadState<()>,
4950 requests: Mutex<Vec<GlobalRequest>>,
4951 retired: std::sync::atomic::AtomicBool,
4952 }
4953
4954 #[reverie::tool]
4955 impl GlobalRPC<GlobalState> for RetirementGuest<'_> {
4956 async fn send_rpc(
4957 &self,
4958 message: <GlobalState as GlobalTool>::Request,
4959 ) -> <GlobalState as GlobalTool>::Response {
4960 self.requests.lock().unwrap().push(message.2.clone());
4961 self.global
4962 .receive_rpc(Tid::from_raw(self.thread.dettid.as_raw()), message)
4963 .await
4964 }
4965
4966 fn config(&self) -> &Config {
4967 self.config
4968 }
4969 }
4970
4971 #[reverie::tool]
4972 impl Guest<Detcore> for RetirementGuest<'_> {
4973 type Memory = reverie::syscalls::LocalMemory;
4974 type Stack = ExternalRegistrationStack;
4975
4976 fn tid(&self) -> reverie::Pid {
4977 reverie::Pid::from_raw(self.thread.dettid.as_raw())
4978 }
4979
4980 fn pid(&self) -> reverie::Pid {
4981 reverie::Pid::from_raw(self.thread.detpid.unwrap().as_raw())
4982 }
4983
4984 fn ppid(&self) -> Option<reverie::Pid> {
4985 None
4986 }
4987
4988 fn memory(&self) -> Self::Memory {
4989 panic!("external registration must not access guest memory")
4990 }
4991
4992 fn thread_state_mut(&mut self) -> &mut crate::ThreadState<()> {
4993 &mut self.thread
4994 }
4995
4996 fn thread_state(&self) -> &crate::ThreadState<()> {
4997 &self.thread
4998 }
4999
5000 async fn regs(&mut self) -> libc::user_regs_struct {
5001 panic!("external registration must not read guest registers")
5002 }
5003
5004 async fn stack(&mut self) -> Self::Stack {
5005 panic!("external registration must not use a guest stack")
5006 }
5007
5008 async fn daemonize(&mut self) {
5009 panic!("external registration must not daemonize")
5010 }
5011
5012 async fn inject<S: reverie::syscalls::SyscallInfo>(
5013 &mut self,
5014 _syscall: S,
5015 ) -> Result<i64, reverie::syscalls::Errno> {
5016 panic!("external registration must not inject a syscall")
5017 }
5018
5019 async fn retire_current_thread(&mut self) -> reverie::Never {
5020 self.retired
5021 .store(true, std::sync::atomic::Ordering::SeqCst);
5022 futures::future::pending().await
5023 }
5024
5025 async fn cancel_current_thread(&mut self) -> reverie::Never {
5026 panic!("scheduler retirement must not request group cancellation")
5027 }
5028
5029 async fn tail_inject<S: reverie::syscalls::SyscallInfo>(
5030 &mut self,
5031 _syscall: S,
5032 ) -> reverie::Never {
5033 panic!("external registration unexpectedly retired its live parent")
5034 }
5035
5036 fn set_timer(&mut self, _schedule: reverie::TimerSchedule) -> Result<(), reverie::Error> {
5037 panic!("external registration must not set a timer")
5038 }
5039
5040 fn set_timer_precise(
5041 &mut self,
5042 _schedule: reverie::TimerSchedule,
5043 ) -> Result<(), reverie::Error> {
5044 panic!("external registration must not set a timer")
5045 }
5046
5047 fn read_clock(&mut self) -> Result<u64, reverie::Error> {
5048 panic!("external registration must not read a host clock")
5049 }
5050 }
5051
5052 #[tokio::test]
5053 async fn thread_exited_uses_natural_retirement_for_killed_stale_and_failed_rpcs() {
5054 use std::sync::atomic::Ordering;
5055
5056 use reverie::Tool;
5057
5058 for cause in ["killed", "stale image", "backend failure", "live"] {
5059 let (config, state, tid, pid) = cancellation_test_state();
5060 install_test_registration(&state, tid, Ivar::new());
5061 let tool: Detcore = Detcore::new(Tid::from_raw(tid.as_raw()), &config);
5062 let mut thread = tool.init_thread_state(Tid::from_raw(tid.as_raw()), None);
5063 thread.detpid = Some(pid);
5064 thread.thread_logical_time.add_syscall_with_cost(37);
5065 match cause {
5066 "killed" => {
5067 state
5068 .sched
5069 .lock()
5070 .unwrap()
5071 .logically_kill_thread(&tid, &pid, thread.mm_id)
5072 }
5073 "stale image" => {
5074 state
5075 .sched
5076 .lock()
5077 .unwrap()
5078 .install_test_exec_incarnation(tid, thread.mm_id);
5079 thread.mm_id = thread.mm_id.for_exec(pid);
5080 }
5081 "backend failure" => state.report_backend_failure(reverie::BackendFailure {
5082 pid: Tid::from_raw(pid.as_raw()),
5083 tid: Tid::from_raw(tid.as_raw()),
5084 phase: "natural retirement control",
5085 }),
5086 "live" => {}
5087 _ => unreachable!(),
5088 }
5089 let mut guest = RetirementGuest {
5090 global: &state,
5091 config: &config,
5092 thread,
5093 requests: Mutex::new(Vec::new()),
5094 retired: std::sync::atomic::AtomicBool::new(false),
5095 };
5096 let before = guest.thread.thread_logical_time.as_nanos();
5097 {
5098 let mut call = std::pin::pin!(super::send_and_update_time(
5099 &mut guest,
5100 GlobalRequest::GlobalTimeLowerBound
5101 ));
5102 if cause == "live" {
5103 assert!(matches!(
5104 futures::poll!(call.as_mut()),
5105 std::task::Poll::Ready((_, GlobalResponse::GlobalTimeLowerBound(_)))
5106 ));
5107 } else {
5108 assert!(futures::poll!(call.as_mut()).is_pending(), "{cause}");
5109 }
5110 }
5111 assert_eq!(
5115 guest.retired.load(Ordering::SeqCst),
5116 matches!(cause, "killed" | "stale image"),
5117 "{cause}"
5118 );
5119 assert_eq!(guest.requests.lock().unwrap().len(), 1);
5120 assert_eq!(guest.thread.thread_logical_time.as_nanos(), before);
5121 if cause == "backend failure" {
5122 assert!(state.sched.lock().unwrap().backend_failed());
5123 assert!(
5124 futures::poll!(std::pin::pin!(state.wait_for_backend_failure())).is_ready()
5125 );
5126 }
5127 }
5128 }
5129
5130 async fn check_external_child_tid_registration(
5131 flags: CloneFlags,
5132 supplied_address: usize,
5133 expected_address: usize,
5134 ) {
5135 use reverie::Tool;
5136
5137 let config = Config {
5138 sequentialize_threads: true,
5139 cancel_killed_thread_rpcs: true,
5140 runs_post_fork: crate::RunsPostFork::Parent,
5143 ..Config::default()
5144 };
5145 let state = GlobalState::initialize(&config, false);
5146 let parent = DetTid::from_raw(17);
5147 let parent_pid = DetPid::from_raw(17);
5148 state
5149 .sched
5150 .lock()
5151 .unwrap()
5152 .thread_tree
5153 .add_child(parent, parent, true);
5154 install_test_registration(&state, parent, Ivar::new());
5155 let tool = Detcore::new(reverie::Pid::from_raw(parent.as_raw()), &config);
5156 let mut thread = tool.init_thread_state(Tid::from_raw(parent.as_raw()), None);
5157 thread.detpid = Some(parent_pid);
5158 let mut guest = ExternalRegistrationGuest {
5159 global: &state,
5160 config: &config,
5161 thread,
5162 requests: Mutex::new(Vec::new()),
5163 post_exec: false,
5164 };
5165 let child = DetTid::from_raw(18);
5166 let exit_signal = if flags.contains(CloneFlags::CLONE_THREAD) {
5167 0
5168 } else {
5169 libc::SIGCHLD
5170 };
5171 tokio::time::timeout(Duration::from_secs(2), async {
5172 let mut registration = std::pin::pin!(tool.register_external_child(
5173 &mut guest,
5174 Tid::from_raw(child.as_raw()),
5175 supplied_address,
5176 flags,
5177 exit_signal,
5178 None,
5179 ));
5180 assert!(futures::poll!(registration.as_mut()).is_pending());
5181 let committed = crate::scheduler::do_a_turn_blocking(
5185 state.sched.clone(),
5186 state.global_time.clone(),
5187 &Err(crate::scheduler::SkipTurn),
5188 )
5189 .await
5190 .expect("the parent continuation must commit");
5191 assert_eq!(committed.tid, parent);
5192 assert_eq!(
5193 committed.resources,
5194 std::collections::HashMap::from([(
5195 crate::resources::ResourceID::ParentContinue { parent, child },
5196 crate::resources::Permission::W,
5197 )]),
5198 );
5199 registration.await;
5200 })
5201 .await
5202 .expect("external child registration must return without running a guest");
5203
5204 assert_eq!(guest.thread.clone_flags, None);
5205 assert_eq!(guest.thread.dettid, parent);
5206 let mut scheduler = state.sched.lock().unwrap();
5207 assert_eq!(scheduler.next_turns.len(), 2);
5208 assert!(scheduler.next_turns.contains_key(&parent));
5209 assert_eq!(
5210 scheduler.next_turns[&child].child_tid_addr,
5211 expected_address
5212 );
5213 assert_eq!(
5214 *guest.requests.lock().unwrap(),
5215 vec![GlobalRequest::CreateChildThread(
5216 child,
5217 parent_pid,
5218 expected_address,
5219 Some(flags),
5220 exit_signal,
5221 None,
5222 Some(DEFAULT_PRIORITY),
5223 )],
5224 "external registration must send one correctly gated real RPC"
5225 );
5226 let child_pid = if flags.contains(CloneFlags::CLONE_THREAD) {
5227 parent_pid
5228 } else {
5229 child
5230 };
5231 let child_mm = MmId::for_clone(
5232 MmId::initial(parent_pid),
5233 child,
5234 flags.contains(CloneFlags::CLONE_VM),
5235 );
5236 scheduler.logically_kill_thread(&child, &child_pid, child_mm);
5237 assert!(!scheduler.next_turns.contains_key(&child));
5238 assert!(scheduler.next_turns.contains_key(&parent));
5239 assert_eq!(
5240 scheduler.child_tid_was_cleared(
5241 FutexID::private(child_mm, supplied_address),
5242 child.as_raw(),
5243 ),
5244 expected_address != 0,
5245 "exit must use only the registered child-TID address"
5246 );
5247 assert!(!scheduler.child_tid_was_cleared(FutexID::private(child_mm, 0), child.as_raw(),));
5248 assert!(!scheduler.child_tid_was_cleared(
5249 FutexID::private(child_mm, supplied_address),
5250 parent.as_raw(),
5251 ));
5252 }
5253
5254 #[tokio::test]
5255 async fn external_registration_without_child_cleartid_disables_exit_wake() {
5256 for kind in [
5257 CloneFlags::empty(),
5258 CloneFlags::CLONE_THREAD | CloneFlags::CLONE_VM | CloneFlags::CLONE_SIGHAND,
5259 ] {
5260 for registration in [CloneFlags::empty(), CloneFlags::CLONE_CHILD_SETTID] {
5261 check_external_child_tid_registration(kind | registration, 0x1234, 0).await;
5262 }
5263 }
5264 }
5265
5266 #[tokio::test]
5267 async fn external_registration_with_child_cleartid_preserves_exact_exit_wake() {
5268 for kind in [
5269 CloneFlags::empty(),
5270 CloneFlags::CLONE_THREAD | CloneFlags::CLONE_VM | CloneFlags::CLONE_SIGHAND,
5271 ] {
5272 for registration in [
5273 CloneFlags::CLONE_CHILD_CLEARTID,
5274 CloneFlags::CLONE_CHILD_CLEARTID | CloneFlags::CLONE_CHILD_SETTID,
5275 ] {
5276 check_external_child_tid_registration(kind | registration, 0x1234, 0x1234).await;
5277 check_external_child_tid_registration(kind | registration, 0, 0).await;
5278 }
5279 }
5280 }
5281
5282 #[tokio::test]
5283 async fn set_child_tid_address_rpc_updates_and_resets_the_registration() {
5284 let (config, state, dettid, detpid) = cancellation_test_state();
5285 install_test_registration(&state, dettid, Ivar::new());
5286 state
5287 .sched
5288 .lock()
5289 .unwrap()
5290 .next_turns
5291 .get_mut(&dettid)
5292 .unwrap()
5293 .child_tid_addr = 0x1234;
5294
5295 let response = state
5296 .receive_rpc(
5297 reverie::Tid::from_raw(dettid.as_raw()),
5298 (
5299 DetTime::new(&config),
5300 MmId::initial(detpid),
5301 GlobalRequest::SetChildTidAddress(0),
5302 ),
5303 )
5304 .await;
5305
5306 assert_eq!(response.1, GlobalResponse::SetChildTidAddress(()));
5307 assert_eq!(
5308 state
5309 .sched
5310 .lock()
5311 .unwrap()
5312 .next_turns
5313 .get(&dettid)
5314 .unwrap()
5315 .child_tid_addr,
5316 0
5317 );
5318 }
5319
5320 async fn child_start_clock_trajectory(
5321 start_before_selection: bool,
5322 child_first: bool,
5323 ) -> (Vec<LogicalTime>, Vec<DetTid>, bool) {
5324 let config = Config {
5325 sequentialize_threads: true,
5326 runs_post_fork: if child_first {
5327 RunsPostFork::Child
5328 } else {
5329 RunsPostFork::Parent
5330 },
5331 ..Config::default()
5332 };
5333 let state = GlobalState::initialize(&config, false);
5334 let parent = DetTid::from_raw(17);
5335 let parent_pid = parent;
5336 state
5337 .sched
5338 .lock()
5339 .unwrap()
5340 .thread_tree
5341 .add_child(parent, parent, true);
5342 let child = DetTid::from_raw(parent.as_raw() + 1);
5343 let parent_mm = MmId::initial(parent_pid);
5344 let child_mm = MmId::for_clone(parent_mm, child, false);
5345 install_test_registration(&state, parent, Ivar::new());
5346 let epoch = DetTime::new(&config).as_nanos();
5347 let mut parent_clock = DetTime::new(&config);
5348 parent_clock.advance_to(epoch + LogicalTime::from_nanos(1_001));
5349 let child_clock = parent_clock.clone_for_child();
5350
5351 let mut registration = Box::pin(state.receive_rpc(
5354 Tid::from_raw(parent.as_raw()),
5355 (
5356 parent_clock.clone(),
5357 parent_mm,
5358 GlobalRequest::CreateChildThread(
5359 child,
5360 parent_pid,
5361 0,
5362 Some(CloneFlags::empty()),
5363 libc::SIGCHLD,
5364 None,
5365 Some(DEFAULT_PRIORITY),
5366 ),
5367 ),
5368 ));
5369 assert!(futures::poll!(&mut registration).is_pending());
5370 let child_request = state.sched.lock().unwrap().next_turns[&child].req.clone();
5371 let mut startup = Box::pin(state.receive_rpc(
5372 Tid::from_raw(child.as_raw()),
5373 (
5374 child_clock.clone(),
5375 child_mm,
5376 GlobalRequest::StartNewThread(child, child, None, None),
5377 ),
5378 ));
5379 if start_before_selection {
5380 assert!(futures::poll!(&mut startup).is_pending());
5382 assert!(futures::poll!(&mut startup).is_pending());
5383 assert!(child_request.try_read().is_some());
5384 }
5385
5386 let skipped = Err(crate::scheduler::SkipTurn);
5387 let mut turn = Box::pin(crate::scheduler::do_a_turn_blocking(
5388 state.sched.clone(),
5389 state.global_time.clone(),
5390 &skipped,
5391 ));
5392 if !start_before_selection && child_first {
5393 assert!(futures::poll!(&mut turn).is_pending());
5394 assert_eq!(child_request.to_string(), "<ivar HasWaiter>");
5395 assert!(futures::poll!(&mut startup).is_pending());
5396 assert!(futures::poll!(&mut startup).is_pending());
5397 assert!(child_request.try_read().is_some());
5398 }
5399 let first = turn.await.expect("first post-fork turn must commit");
5400 if !start_before_selection && !child_first {
5401 assert!(child_request.try_read().is_none());
5402 assert!(futures::poll!(&mut startup).is_pending());
5403 assert!(futures::poll!(&mut startup).is_pending());
5404 assert!(child_request.try_read().is_some());
5405 }
5406 let child_resource = ResourceID::MemAddrSpace(child);
5407 let parent_resource = ResourceID::ParentContinue { parent, child };
5408 let (first_tid, first_resource, first_permission, second_resource, second_permission) =
5409 if child_first {
5410 (
5411 child,
5412 &child_resource,
5413 Permission::RW,
5414 &parent_resource,
5415 Permission::W,
5416 )
5417 } else {
5418 (
5419 parent,
5420 &parent_resource,
5421 Permission::W,
5422 &child_resource,
5423 Permission::RW,
5424 )
5425 };
5426 assert_eq!(first.tid, first_tid);
5427 assert_eq!(first.resources.len(), 1);
5428 assert_eq!(first.resources.get(first_resource), Some(&first_permission));
5429 if child_first {
5430 assert_eq!(
5431 startup.as_mut().await,
5432 (None, GlobalResponse::StartNewThread(None))
5433 );
5434 } else {
5435 assert_eq!(
5436 registration.as_mut().await,
5437 (None, GlobalResponse::CreateChildThread(None))
5438 );
5439 }
5440 let first_time = state.sched.lock().unwrap().committed_time;
5441 assert_eq!(first_time, epoch + LogicalTime::from_nanos(1_001));
5442 assert_eq!(
5443 state.global_time.lock().unwrap().threads_time(child),
5444 child_clock.as_nanos()
5445 );
5446
5447 let (mut first_clock, first_mm) = if child_first {
5450 (child_clock, child_mm)
5451 } else {
5452 (parent_clock, parent_mm)
5453 };
5454 first_clock.advance_to(first_clock.as_nanos() + LogicalTime::from_nanos(1));
5455 let mut resources = Resources::new(first_tid);
5456 resources.insert(ResourceID::MemAddrSpace(first_tid), Permission::RW);
5457 let mut work = Box::pin(state.receive_rpc(
5458 Tid::from_raw(first_tid.as_raw()),
5459 (
5460 first_clock,
5461 first_mm,
5462 GlobalRequest::RequestResources(resources, first_tid),
5463 ),
5464 ));
5465 assert!(futures::poll!(&mut work).is_pending());
5466 let second = crate::scheduler::do_a_turn_blocking(
5467 state.sched.clone(),
5468 state.global_time.clone(),
5469 &Ok(first),
5470 )
5471 .await
5472 .expect("other thread's continuation must commit");
5473 assert_eq!(second.resources.len(), 1);
5474 assert_eq!(
5475 second.resources.get(second_resource),
5476 Some(&second_permission)
5477 );
5478 if child_first {
5479 assert_eq!(
5480 registration.await,
5481 (None, GlobalResponse::CreateChildThread(None))
5482 );
5483 } else {
5484 assert_eq!(startup.await, (None, GlobalResponse::StartNewThread(None)));
5485 }
5486 let mut scheduler = state.sched.lock().unwrap();
5487 let next_time = scheduler.committed_time;
5488 assert_eq!(next_time, epoch + LogicalTime::from_nanos(501_002));
5489 let queue = scheduler.run_queue.tids().copied().collect();
5490 let next_random = scheduler.child_runs_first_post_fork(RunsPostFork::Random);
5491 (vec![first_time, next_time], queue, next_random)
5492 }
5493
5494 #[tokio::test]
5495 async fn child_start_clock_is_independent_of_first_rpc_arrival() {
5496 for child_first in [true, false] {
5497 assert_eq!(
5498 child_start_clock_trajectory(true, child_first).await,
5499 child_start_clock_trajectory(false, child_first).await,
5500 );
5501 }
5502 }
5503
5504 #[tokio::test]
5505 async fn vfork_registration_does_not_charge_inherited_work_before_startup() {
5506 let (config, state, parent, parent_pid) = cancellation_test_state();
5507 let child = DetTid::from_raw(parent.as_raw() + 1);
5508 let mm = MmId::initial(parent_pid);
5509 install_test_registration(&state, parent, Ivar::new());
5510 let mut parent_clock = DetTime::new(&config);
5511 let epoch = parent_clock.as_nanos();
5512 parent_clock.advance_to(epoch + LogicalTime::from_nanos(1_001));
5513 let child_clock = parent_clock.clone_for_child();
5514 let mut resources = Resources::new(parent);
5515 resources.insert(
5516 ResourceID::BlockingVfork(ExternalOpId::new(parent, 1)),
5517 Permission::RW,
5518 );
5519 let mut blocking = Box::pin(state.receive_rpc(
5520 Tid::from_raw(parent.as_raw()),
5521 (
5522 parent_clock,
5523 mm,
5524 GlobalRequest::RequestResources(resources, parent_pid),
5525 ),
5526 ));
5527 assert!(futures::poll!(&mut blocking).is_pending());
5528 let background = crate::scheduler::do_a_turn_blocking(
5529 state.sched.clone(),
5530 state.global_time.clone(),
5531 &Err(crate::scheduler::SkipTurn),
5532 )
5533 .await;
5534 assert!(background.is_err());
5535 assert_eq!(
5536 blocking.await,
5537 (None, GlobalResponse::RequestResources(ResumeStatus::Normal))
5538 );
5539 let before_child = state.global_time.lock().unwrap().as_nanos();
5540 assert_eq!(before_child, epoch + LogicalTime::from_nanos(1_001));
5541
5542 let created = state
5544 .receive_rpc(
5545 Tid::from_raw(child.as_raw()),
5546 (
5547 child_clock.clone(),
5548 mm,
5549 GlobalRequest::CreateVforkChildThread(
5550 parent,
5551 parent_pid,
5552 child,
5553 0,
5554 CloneFlags::CLONE_VFORK | CloneFlags::CLONE_VM,
5555 libc::SIGCHLD,
5556 Some(DEFAULT_PRIORITY - 1),
5557 ),
5558 ),
5559 )
5560 .await;
5561 assert_eq!(created, (None, GlobalResponse::CreateChildThread(None)));
5562 assert_eq!(state.global_time.lock().unwrap().as_nanos(), before_child);
5563 assert_eq!(
5564 state.global_time.lock().unwrap().threads_time(child),
5565 child_clock.as_nanos()
5566 );
5567
5568 let mut startup = Box::pin(state.receive_rpc(
5569 Tid::from_raw(child.as_raw()),
5570 (
5571 child_clock,
5572 mm,
5573 GlobalRequest::StartNewThread(child, child, None, None),
5574 ),
5575 ));
5576 assert!(futures::poll!(&mut startup).is_pending());
5577 assert!(futures::poll!(&mut startup).is_pending());
5578 let first = crate::scheduler::do_a_turn_blocking(
5579 state.sched.clone(),
5580 state.global_time.clone(),
5581 &background,
5582 )
5583 .await
5584 .expect("vfork child must receive its first turn");
5585 assert_eq!(first.tid, child);
5586 assert_eq!(startup.await, (None, GlobalResponse::StartNewThread(None)));
5587 assert_eq!(state.global_time.lock().unwrap().as_nanos(), before_child);
5588 }
5589
5590 #[tokio::test]
5591 async fn first_nonstartup_rpc_counts_only_new_child_work() {
5592 let config = Config {
5593 sequentialize_threads: false,
5594 ..Config::default()
5595 };
5596 let state = GlobalState::initialize(&config, false);
5597 let parent = DetTid::from_raw(17);
5598 let child = DetTid::from_raw(18);
5599 let mut parent_clock = DetTime::new(&config);
5600 let epoch = parent_clock.as_nanos();
5601 parent_clock.advance_to(epoch + LogicalTime::from_nanos(101));
5602 let mut child_clock = parent_clock.clone_for_child();
5603 let _ = state
5604 .receive_rpc(
5605 Tid::from_raw(parent.as_raw()),
5606 (
5607 parent_clock,
5608 MmId::initial(parent),
5609 GlobalRequest::GlobalTimeLowerBound,
5610 ),
5611 )
5612 .await;
5613 child_clock.advance_to(child_clock.as_nanos() + LogicalTime::from_nanos(1));
5614 let observed = state
5615 .receive_rpc(
5616 Tid::from_raw(child.as_raw()),
5617 (
5618 child_clock,
5619 MmId::initial(child),
5620 GlobalRequest::GlobalTimeLowerBound,
5621 ),
5622 )
5623 .await;
5624 assert_eq!(
5625 observed,
5626 (
5627 None,
5628 GlobalResponse::GlobalTimeLowerBound(epoch + LogicalTime::from_nanos(102))
5629 )
5630 );
5631 }
5632
5633 #[tokio::test]
5634 async fn exec_reconnect_retains_inherited_work_accounting_across_local_reload() {
5635 let (config, state, leader, detpid) = cancellation_test_state();
5636 let ancestor = DetTid::from_raw(leader.as_raw() - 1);
5637 let worker = DetTid::from_raw(leader.as_raw() + 1);
5638 let old_mm = MmId::initial(detpid);
5639 install_test_registration(&state, leader, Ivar::new());
5640 state
5641 .sched
5642 .lock()
5643 .unwrap()
5644 .thread_tree
5645 .add_child(leader, worker, false);
5646 install_test_registration(&state, worker, Ivar::new());
5647
5648 let mut ancestor_clock = DetTime::new(&config);
5649 let epoch = ancestor_clock.as_nanos();
5650 ancestor_clock.advance_to(epoch + LogicalTime::from_nanos(1_000));
5651 let mut leader_clock = ancestor_clock.clone_for_child();
5652 leader_clock.advance_to(epoch + LogicalTime::from_nanos(1_100));
5653 let mut worker_clock = leader_clock.clone_for_child();
5654 worker_clock.advance_to(epoch + LogicalTime::from_nanos(1_350));
5655 for (tid, clock) in [
5656 (ancestor, ancestor_clock),
5657 (leader, leader_clock),
5658 (worker, worker_clock.clone()),
5659 ] {
5660 let _ = state
5661 .receive_rpc(
5662 Tid::from_raw(tid.as_raw()),
5663 (clock, old_mm, GlobalRequest::GlobalTimeLowerBound),
5664 )
5665 .await;
5666 }
5667 let total = state.global_time.lock().unwrap().as_nanos();
5668 assert_eq!(total, epoch + LogicalTime::from_nanos(1_350));
5669 let _ = state
5673 .receive_rpc(
5674 Tid::from_raw(worker.as_raw()),
5675 (
5676 worker_clock.clone(),
5677 old_mm,
5678 GlobalRequest::PrepareExec(detpid, old_mm, Default::default()),
5679 ),
5680 )
5681 .await;
5682 let cancelled = state
5683 .receive_rpc(
5684 Tid::from_raw(worker.as_raw()),
5685 (
5686 worker_clock.clone(),
5687 old_mm,
5688 GlobalRequest::CancelExec(detpid),
5689 ),
5690 )
5691 .await;
5692 assert_eq!(cancelled, (None, GlobalResponse::CancelExec(())));
5693 assert!(state.pending_exec_states.lock().unwrap().is_empty());
5694 assert_eq!(state.global_time.lock().unwrap().as_nanos(), total);
5695 assert_eq!(
5696 state.global_time.lock().unwrap().threads_time(worker),
5697 worker_clock.as_nanos()
5698 );
5699
5700 let prepared = state
5701 .receive_rpc(
5702 Tid::from_raw(worker.as_raw()),
5703 (
5704 worker_clock.clone(),
5705 old_mm,
5706 GlobalRequest::PrepareExec(detpid, old_mm, Default::default()),
5707 ),
5708 )
5709 .await;
5710 assert_eq!(prepared, (None, GlobalResponse::PrepareExec(())));
5711
5712 let mut fresh = DetTime::new(&config);
5713 let recreated = state
5714 .receive_rpc(
5715 Tid::from_raw(leader.as_raw()),
5716 (
5717 fresh.clone(),
5718 MmId::initial(leader),
5719 GlobalRequest::CreateChildThread(
5720 leader,
5721 detpid,
5722 0,
5723 None,
5724 libc::SIGCHLD,
5725 None,
5726 Some(DEFAULT_PRIORITY),
5727 ),
5728 ),
5729 )
5730 .await;
5731 assert_eq!(
5732 recreated,
5733 (
5734 Some(worker_clock.as_nanos()),
5735 GlobalResponse::CreateChildThread(Some(old_mm.for_exec(detpid)))
5736 )
5737 );
5738 let stale = state
5741 .receive_rpc(
5742 Tid::from_raw(leader.as_raw()),
5743 (
5744 worker_clock.clone(),
5745 old_mm,
5746 GlobalRequest::RequestResources(Resources::new(leader), detpid),
5747 ),
5748 )
5749 .await;
5750 assert_eq!(stale, (None, GlobalResponse::ThreadExited));
5751 assert_eq!(state.global_time.lock().unwrap().as_nanos(), total);
5752
5753 fresh.advance_to(recreated.0.unwrap());
5754 assert_eq!(fresh.inherited_nanos(), LogicalTime::ZERO);
5755 let mut startup = Box::pin(state.receive_rpc(
5756 Tid::from_raw(leader.as_raw()),
5757 (
5758 fresh.clone(),
5759 old_mm.for_exec(detpid),
5760 GlobalRequest::StartNewThread(leader, detpid, None, None),
5761 ),
5762 ));
5763 assert!(futures::poll!(&mut startup).is_pending());
5764 assert!(futures::poll!(&mut startup).is_pending());
5765 let first = crate::scheduler::do_a_turn_blocking(
5766 state.sched.clone(),
5767 state.global_time.clone(),
5768 &Err(crate::scheduler::SkipTurn),
5769 )
5770 .await
5771 .expect("replacement leader must run");
5772 assert_eq!(first.tid, leader);
5773 assert_eq!(startup.await, (None, GlobalResponse::StartNewThread(None)));
5774 assert_eq!(state.global_time.lock().unwrap().as_nanos(), total);
5775
5776 fresh.advance_to(fresh.as_nanos() + LogicalTime::from_nanos(1));
5777 let observed = state
5778 .receive_rpc(
5779 Tid::from_raw(leader.as_raw()),
5780 (
5781 fresh.clone(),
5782 old_mm.for_exec(detpid),
5783 GlobalRequest::GlobalTimeLowerBound,
5784 ),
5785 )
5786 .await;
5787 assert_eq!(
5788 observed,
5789 (
5790 None,
5791 GlobalResponse::GlobalTimeLowerBound(total + LogicalTime::from_nanos(1))
5792 )
5793 );
5794 let global = state.global_time.lock().unwrap();
5795 assert_eq!(global.threads_time(leader), fresh.as_nanos());
5796 assert!(!global.contains_thread(worker));
5797 }
5798
5799 #[test]
5800 fn live_registration_without_next_turn_is_not_terminal() {
5801 let (_, state, dettid, detpid) = cancellation_test_state();
5802 install_test_registration(&state, dettid, Ivar::new());
5803 state.sched.lock().unwrap().next_turns.remove(&dettid);
5804
5805 assert!(
5806 !state
5807 .sched
5808 .lock()
5809 .unwrap()
5810 .thread_is_logically_killed(dettid),
5811 "transient next-turn absence must not imply logical death"
5812 );
5813
5814 state
5815 .sched
5816 .lock()
5817 .unwrap()
5818 .logically_kill_thread(&dettid, &detpid, MmId::initial(detpid));
5819 assert!(
5820 state
5821 .sched
5822 .lock()
5823 .unwrap()
5824 .thread_is_logically_killed(dettid),
5825 "explicit logical death must install a permanent TID tombstone"
5826 );
5827 }
5828
5829 #[tokio::test]
5830 async fn exec_reconnect_retires_siblings_and_reuses_live_scheduler_and_clock_state() {
5831 let (config, state, dettid, detpid) = cancellation_test_state();
5832 let old_mm = MmId::initial(detpid).for_exec(detpid);
5833 install_test_registration(&state, dettid, Ivar::new());
5834 let sibling = DetTid::from_raw(dettid.as_raw() + 1);
5835 let sibling_request = Ivar::new();
5836 state
5837 .sched
5838 .lock()
5839 .unwrap()
5840 .thread_tree
5841 .add_child(dettid, sibling, false);
5842 install_test_registration(&state, sibling, sibling_request.clone());
5843 {
5844 let mut scheduler = state.sched.lock().unwrap();
5845 scheduler.next_turns.get_mut(&dettid).unwrap().resp =
5846 Ivar::full(SchedResponse::Go(None));
5847 scheduler
5848 .next_turns
5849 .get_mut(&sibling)
5850 .unwrap()
5851 .child_tid_addr = 0x1234;
5852 }
5853 let mut existing_time = DetTime::new(&config);
5854 existing_time.add_syscall();
5855 existing_time.add_syscall();
5856 state.global_time.lock().unwrap().update_global_time(
5857 dettid,
5858 existing_time.as_nanos(),
5859 LogicalTime::ZERO,
5860 );
5861 let (global_before, thread_before) = {
5862 let global_time = state.global_time.lock().unwrap();
5863 (global_time.as_nanos(), global_time.threads_time(dettid))
5864 };
5865 let fresh_local_time = DetTime::new(&config);
5866 let physical_pid = std::process::id() as i32;
5867 let physical_tid = unsafe { libc::syscall(libc::SYS_gettid) as i32 };
5868 let physical_ids = Some((physical_pid, physical_tid));
5869 state.pending_exec_states.lock().unwrap().insert(
5870 detpid,
5871 PendingExecState {
5872 caller: dettid,
5873 process: detpid,
5874 mm: old_mm,
5875 fd_blocking: Default::default(),
5876 },
5877 );
5878 state
5879 .sched
5880 .lock()
5881 .unwrap()
5882 .install_test_exec_incarnation(dettid, old_mm);
5883 let in_flight_exec_response = state
5884 .receive_rpc(
5885 reverie::Tid::from_raw(dettid.as_raw()),
5886 (
5887 existing_time.clone(),
5888 old_mm.for_exec(detpid),
5889 GlobalRequest::ReportUnsupportedSyscall("exec-in-flight".to_owned()),
5890 ),
5891 )
5892 .await;
5893 assert_eq!(
5894 in_flight_exec_response.1,
5895 GlobalResponse::ReportUnsupportedSyscall(())
5896 );
5897
5898 let create_response = state
5899 .receive_rpc(
5900 reverie::Tid::from_raw(dettid.as_raw()),
5901 (
5902 fresh_local_time.clone(),
5903 MmId::initial(dettid),
5904 GlobalRequest::CreateChildThread(
5905 dettid,
5906 detpid,
5907 0,
5908 None,
5909 libc::SIGCHLD,
5910 physical_ids,
5911 Some(DEFAULT_PRIORITY),
5912 ),
5913 ),
5914 )
5915 .await;
5916 assert_eq!(
5917 create_response,
5918 (
5919 Some(thread_before),
5920 GlobalResponse::CreateChildThread(Some(old_mm.for_exec(detpid)))
5921 )
5922 );
5923 assert_eq!(
5924 state.sched.lock().unwrap().physical_thread_identity(dettid),
5925 Some((old_mm.for_exec(detpid), physical_pid, physical_tid)),
5926 "post-exec host identity must be installed by CreateChildThread before admission"
5927 );
5928
5929 let start_response = state
5930 .receive_rpc(
5931 reverie::Tid::from_raw(dettid.as_raw()),
5932 (
5933 fresh_local_time,
5934 old_mm.for_exec(detpid),
5935 GlobalRequest::StartNewThread(dettid, detpid, physical_ids, None),
5936 ),
5937 )
5938 .await;
5939 assert_eq!(
5940 start_response,
5941 (Some(thread_before), GlobalResponse::StartNewThread(None))
5942 );
5943
5944 let scheduler = state.sched.lock().unwrap();
5945 assert!(!scheduler.thread_is_logically_killed(dettid));
5946 assert!(scheduler.thread_is_logically_killed(sibling));
5947 assert_eq!(scheduler.next_turns.len(), 1);
5948 assert!(matches!(sibling_request.try_read(), Some(Err(_))));
5949 assert!(
5950 scheduler.child_tid_was_cleared(FutexID::private(old_mm, 0x1234), sibling.as_raw())
5951 );
5952 drop(scheduler);
5953 let global_time = state.global_time.lock().unwrap();
5954 assert_eq!(global_time.as_nanos(), global_before);
5955 assert_eq!(global_time.threads_time(dettid), thread_before);
5956 }
5957
5958 #[tokio::test]
5959 async fn nonleader_exec_rebinds_caller_to_leader_and_preserves_its_clock() {
5960 let (config, state, leader, detpid) = cancellation_test_state();
5961 let worker = DetTid::from_raw(leader.as_raw() + 1);
5962 let sibling = DetTid::from_raw(leader.as_raw() + 2);
5963 let old_mm = MmId::initial(detpid).for_exec(detpid);
5964 let leader_request = Ivar::new();
5965 let worker_request = Ivar::new();
5966 let sibling_request = Ivar::new();
5967 install_test_registration(&state, leader, leader_request.clone());
5968 {
5969 let mut scheduler = state.sched.lock().unwrap();
5970 scheduler.thread_tree.add_child(leader, worker, false);
5971 scheduler.thread_tree.add_child(leader, sibling, false);
5972 }
5973 install_test_registration(&state, worker, worker_request.clone());
5974 install_test_registration(&state, sibling, sibling_request.clone());
5975 {
5976 let mut scheduler = state.sched.lock().unwrap();
5977 scheduler
5978 .next_turns
5979 .get_mut(&sibling)
5980 .unwrap()
5981 .child_tid_addr = 0x5678;
5982 scheduler
5983 .timeslices
5984 .insert(leader, Some(LogicalTime::from_nanos(99)));
5985 scheduler.install_test_vfork_barrier(leader, sibling);
5986 }
5987
5988 let mut leader_clock = DetTime::new(&config);
5989 leader_clock.add_syscall();
5990 let mut worker_clock = DetTime::new(&config);
5991 worker_clock.add_syscall();
5992 worker_clock.add_syscall();
5993 worker_clock.add_syscall();
5994 {
5995 let mut global_time = state.global_time.lock().unwrap();
5996 global_time.update_global_time(leader, leader_clock.as_nanos(), LogicalTime::ZERO);
5997 global_time.update_global_time(worker, worker_clock.as_nanos(), LogicalTime::ZERO);
5998 }
5999 let total_before = state.global_time.lock().unwrap().as_nanos();
6000 let fd_blocking: ExecFdBlockingOverrides = [42].into_iter().collect();
6001 state.pending_exec_states.lock().unwrap().insert(
6002 detpid,
6003 PendingExecState {
6004 caller: worker,
6005 process: detpid,
6006 mm: old_mm,
6007 fd_blocking: fd_blocking.clone(),
6008 },
6009 );
6010 let fresh_local_time = DetTime::new(&config);
6011
6012 let create_response = state
6013 .receive_rpc(
6014 reverie::Tid::from_raw(leader.as_raw()),
6015 (
6016 fresh_local_time.clone(),
6017 MmId::initial(leader),
6018 GlobalRequest::CreateChildThread(
6019 leader,
6020 detpid,
6021 0,
6022 None,
6023 libc::SIGCHLD,
6024 None,
6025 Some(DEFAULT_PRIORITY),
6026 ),
6027 ),
6028 )
6029 .await;
6030 assert_eq!(
6031 create_response,
6032 (
6033 Some(worker_clock.as_nanos()),
6034 GlobalResponse::CreateChildThread(Some(old_mm.for_exec(detpid)))
6035 )
6036 );
6037 assert!(state.pending_exec_states.lock().unwrap().is_empty());
6038 assert_eq!(
6039 state.post_exec_fd_blocking.lock().unwrap().get(&leader),
6040 Some(&fd_blocking)
6041 );
6042
6043 let late_old_request = state
6044 .receive_rpc(
6045 reverie::Tid::from_raw(leader.as_raw()),
6046 (
6047 leader_clock.clone(),
6048 old_mm,
6049 GlobalRequest::RequestResources(Resources::new(leader), detpid),
6050 ),
6051 )
6052 .await;
6053 assert_eq!(late_old_request, (None, GlobalResponse::ThreadExited));
6054 let admitted_before_fence = state
6055 .recv_request_resources(
6056 reverie::Tid::from_raw(leader.as_raw()),
6057 detpid,
6058 Resources::new(leader),
6059 Some(old_mm),
6060 )
6061 .await;
6062 assert_eq!(
6063 admitted_before_fence,
6064 (SchedulerRpcResult::ThreadExited, None)
6065 );
6066 let late_old_deregister = state
6067 .receive_rpc(
6068 reverie::Tid::from_raw(leader.as_raw()),
6069 (
6070 leader_clock.clone(),
6071 old_mm,
6072 GlobalRequest::DeregisterThread(ThreadDeregistration {
6073 thread_start_entered: true,
6074 dettid: leader,
6075 detpid,
6076 mm: old_mm,
6077 timeslice_stats: TimesliceStats::default(),
6078 syscall_count: 0,
6079 chaos_epochs: Vec::new(),
6080 }),
6081 ),
6082 )
6083 .await;
6084 assert_eq!(
6085 late_old_deregister,
6086 (None, GlobalResponse::DeregisterThread(()))
6087 );
6088 let duplicate_create = state
6089 .receive_rpc(
6090 reverie::Tid::from_raw(leader.as_raw()),
6091 (
6092 fresh_local_time.clone(),
6093 MmId::initial(leader),
6094 GlobalRequest::CreateChildThread(
6095 leader,
6096 detpid,
6097 0,
6098 None,
6099 libc::SIGCHLD,
6100 None,
6101 Some(DEFAULT_PRIORITY),
6102 ),
6103 ),
6104 )
6105 .await;
6106 assert_eq!(duplicate_create, (None, GlobalResponse::ThreadExited));
6107
6108 {
6109 let scheduler = state.sched.lock().unwrap();
6110 assert!(!scheduler.thread_is_logically_killed(leader));
6111 assert!(scheduler.thread_is_logically_killed(worker));
6112 assert!(scheduler.thread_is_logically_killed(sibling));
6113 assert_eq!(scheduler.next_turns.len(), 1);
6114 assert!(scheduler.next_turns.contains_key(&leader));
6115 assert!(!scheduler.timeslices.contains_key(&leader));
6116 assert!(!scheduler.vfork_barrier_mentions(leader));
6117 assert!(!scheduler.vfork_barrier_mentions(sibling));
6118 assert!(matches!(leader_request.try_read(), Some(Err(_))));
6119 assert!(matches!(worker_request.try_read(), Some(Err(_))));
6120 assert!(matches!(sibling_request.try_read(), Some(Err(_))));
6121 assert!(
6122 scheduler
6123 .child_tid_was_cleared(FutexID::private(old_mm, 0x5678), sibling.as_raw(),)
6124 );
6125 }
6126
6127 let turn_sched = state.sched.clone();
6132 let turn_time = state.global_time.clone();
6133 let turn = tokio::spawn(async move {
6134 let last: Result<Resources, crate::scheduler::SkipTurn> =
6135 Err(crate::scheduler::SkipTurn);
6136 crate::scheduler::do_a_turn_blocking(turn_sched, turn_time, &last).await
6137 });
6138 let start_response = state
6139 .receive_rpc(
6140 reverie::Tid::from_raw(leader.as_raw()),
6141 (
6142 fresh_local_time,
6143 old_mm.for_exec(detpid),
6144 GlobalRequest::StartNewThread(leader, detpid, None, None),
6145 ),
6146 )
6147 .await;
6148 assert!(
6149 turn.await
6150 .expect("replacement scheduler turn panicked")
6151 .is_ok(),
6152 "replacement leader did not survive the first step2 drain"
6153 );
6154 assert_eq!(
6155 start_response,
6156 (
6157 Some(worker_clock.as_nanos()),
6158 GlobalResponse::StartNewThread(None)
6159 )
6160 );
6161 {
6162 let scheduler = state.sched.lock().unwrap();
6163 assert_eq!(
6164 scheduler
6165 .run_queue
6166 .tids()
6167 .filter(|dettid| **dettid == leader)
6168 .count(),
6169 1
6170 );
6171 assert!(!scheduler.run_queue.contains_tid(worker));
6172 assert!(!scheduler.run_queue.contains_tid(sibling));
6173 }
6174 {
6175 let global_time = state.global_time.lock().unwrap();
6176 assert_eq!(global_time.as_nanos(), total_before);
6177 assert_eq!(global_time.threads_time(leader), worker_clock.as_nanos());
6178 assert!(!global_time.contains_thread(worker));
6179 }
6180
6181 let mark_response = state
6182 .receive_rpc(
6183 reverie::Tid::from_raw(leader.as_raw()),
6184 (
6185 worker_clock,
6186 old_mm.for_exec(detpid),
6187 GlobalRequest::MarkPastFirstExecve(detpid, None),
6188 ),
6189 )
6190 .await;
6191 assert_eq!(
6192 mark_response.1,
6193 GlobalResponse::MarkPastFirstExecve(fd_blocking)
6194 );
6195 assert!(state.post_exec_fd_blocking.lock().unwrap().is_empty());
6196 }
6197
6198 #[tokio::test]
6199 async fn post_exec_deletes_armed_disarmed_and_non_notifying_timers() {
6200 use reverie::Tool;
6201
6202 let config = Config {
6203 sequentialize_threads: true,
6204 cancel_killed_thread_rpcs: true,
6205 max_timeslice: None,
6208 ..Config::default()
6209 };
6210 let state = GlobalState::initialize(&config, false);
6211 let leader = DetTid::from_raw(17);
6212 let detpid = DetPid::from_raw(17);
6213 state
6214 .sched
6215 .lock()
6216 .unwrap()
6217 .thread_tree
6218 .add_child(leader, leader, true);
6219 let next_request = Ivar::new();
6220 install_test_registration(&state, leader, next_request.clone());
6221 let tool = Detcore::new(reverie::Pid::from_raw(leader.as_raw()), &config);
6222 let mut thread = tool.init_thread_state(Tid::from_raw(leader.as_raw()), None);
6223 thread.detpid = Some(detpid);
6224 thread.thread_logical_time.add_syscall_with_cost(137);
6225 let mut guest = ExternalRegistrationGuest {
6226 global: &state,
6227 config: &config,
6228 thread,
6229 requests: Mutex::new(Vec::new()),
6230 post_exec: true,
6231 };
6232 let t = LogicalTime::from_nanos;
6233 let [periodic, disarmed, silent] = {
6234 let mut timers = guest.thread.posix_timers.lock().unwrap();
6235 let periodic = timers.create(Some(libc::SIGUSR2));
6236 let disarmed = timers.create(Some(libc::SIGALRM));
6237 let silent = timers.create(None);
6238 timers.settime(periodic, 50, Some(t(100)), t(0));
6239 timers.settime(silent, 0, Some(t(200)), t(0));
6240 [periodic, disarmed, silent]
6241 };
6242
6243 for _ in 0..2 {
6247 let old_mm = guest.thread.mm_id;
6248 super::prepare_exec(&mut guest, old_mm, Default::default()).await;
6249 guest.thread.mm_id = old_mm.for_exec(detpid);
6250 let before = guest.thread.thread_logical_time.as_nanos();
6251 tokio::time::timeout(Duration::from_secs(2), tool.handle_post_exec(&mut guest))
6252 .await
6253 .expect("post-exec must finish within the caller's scheduler turn")
6254 .unwrap();
6255 assert_eq!(guest.thread.thread_logical_time.as_nanos(), before);
6256 assert!(
6257 next_request.try_read().is_none(),
6258 "post-exec unexpectedly yielded"
6259 );
6260 let mut timers = guest.thread.posix_timers.lock().unwrap();
6261 for id in [periodic, disarmed, silent] {
6262 assert!(!timers.contains(id), "post-exec retained POSIX timer {id}");
6263 assert_eq!(timers.gettime(id, t(10)), None);
6264 assert_eq!(timers.settime(id, 0, Some(t(300)), t(10)), None);
6265 assert_eq!(timers.signal(id), None);
6266 assert!(!timers.remove(id));
6267 }
6268 }
6269 let mut timers = guest.thread.posix_timers.lock().unwrap();
6270 let new = timers.create(Some(libc::SIGUSR1));
6271 assert_eq!(new, 3);
6272 assert_eq!(timers.settime(new, 0, Some(t(300)), t(10)), Some((0, 0)));
6273 assert_eq!(timers.gettime(new, t(20)), Some((280, 0)));
6274 }
6275
6276 #[tokio::test]
6277 async fn only_successful_exec_cancels_posix_deadlines() {
6278 let (config, state, leader, detpid) = cancellation_test_state();
6279 install_test_registration(&state, leader, Ivar::new());
6280 let clock = DetTime::new(&config);
6281 let mm = MmId::initial(detpid);
6282 let deadline = LogicalTime::from_nanos(1_000_000);
6283 {
6284 let mut sched = state.sched.lock().unwrap();
6285 sched.register_alarm(
6286 detpid,
6287 leader,
6288 LogicalTime::ZERO,
6289 deadline,
6290 LogicalTime::ZERO,
6291 Signal::SIGALRM,
6292 );
6293 sched.register_posix_timer(
6294 detpid,
6295 leader,
6296 0,
6297 Some(deadline),
6298 deadline,
6299 Signal::SIGUSR2,
6300 );
6301 }
6302 let before: Vec<_> = state
6303 .sched
6304 .lock()
6305 .unwrap()
6306 .blocked
6307 .timed_waiters
6308 .iter()
6309 .collect();
6310 for request in [
6311 GlobalRequest::PrepareExec(detpid, mm, Default::default()),
6312 GlobalRequest::CancelExec(detpid),
6313 GlobalRequest::PrepareExec(detpid, mm, Default::default()),
6314 ] {
6315 state
6316 .receive_rpc(
6317 reverie::Tid::from_raw(leader.as_raw()),
6318 (clock.clone(), mm, request),
6319 )
6320 .await;
6321 assert_eq!(
6322 state
6323 .sched
6324 .lock()
6325 .unwrap()
6326 .blocked
6327 .timed_waiters
6328 .iter()
6329 .collect::<Vec<_>>(),
6330 before
6331 );
6332 }
6333
6334 let response = state
6335 .receive_rpc(
6336 reverie::Tid::from_raw(leader.as_raw()),
6337 (
6338 clock,
6339 mm.for_exec(detpid),
6340 GlobalRequest::MarkPastFirstExecve(detpid, None),
6341 ),
6342 )
6343 .await;
6344 assert_eq!(
6345 response.1,
6346 GlobalResponse::MarkPastFirstExecve(Default::default())
6347 );
6348 let sched = state.sched.lock().unwrap();
6349 assert_eq!(
6350 sched.blocked.timed_waiters.alarm_state(detpid),
6351 Some((deadline, LogicalTime::ZERO))
6352 );
6353 assert_eq!(sched.blocked.timed_waiters.iter().count(), 1);
6354 }
6355
6356 #[tokio::test]
6357 async fn failed_exec_clears_prepared_state_without_retiring_siblings() {
6358 let (config, state, leader, detpid) = cancellation_test_state();
6359 let sibling = DetTid::from_raw(leader.as_raw() + 1);
6360 install_test_registration(&state, leader, Ivar::new());
6361 state
6362 .sched
6363 .lock()
6364 .unwrap()
6365 .thread_tree
6366 .add_child(leader, sibling, false);
6367 install_test_registration(&state, sibling, Ivar::new());
6368 let clock = DetTime::new(&config);
6369
6370 let prepared = state
6371 .receive_rpc(
6372 reverie::Tid::from_raw(leader.as_raw()),
6373 (
6374 clock.clone(),
6375 MmId::initial(leader),
6376 GlobalRequest::PrepareExec(detpid, MmId::initial(detpid), Default::default()),
6377 ),
6378 )
6379 .await;
6380 assert_eq!(prepared.1, GlobalResponse::PrepareExec(()));
6381 assert!(
6382 state
6383 .pending_exec_states
6384 .lock()
6385 .unwrap()
6386 .contains_key(&detpid)
6387 );
6388
6389 let cancelled = state
6390 .receive_rpc(
6391 reverie::Tid::from_raw(leader.as_raw()),
6392 (
6393 clock,
6394 MmId::initial(leader),
6395 GlobalRequest::CancelExec(detpid),
6396 ),
6397 )
6398 .await;
6399 assert_eq!(cancelled.1, GlobalResponse::CancelExec(()));
6400 assert!(state.pending_exec_states.lock().unwrap().is_empty());
6401 {
6402 let scheduler = state.sched.lock().unwrap();
6403 assert!(!scheduler.thread_is_logically_killed(leader));
6404 assert!(!scheduler.thread_is_logically_killed(sibling));
6405 assert_eq!(scheduler.next_turns.len(), 2);
6406 }
6407
6408 state.pending_exec_states.lock().unwrap().insert(
6409 detpid,
6410 PendingExecState {
6411 caller: leader,
6412 process: detpid,
6413 mm: MmId::initial(detpid),
6414 fd_blocking: Default::default(),
6415 },
6416 );
6417 state
6418 .post_exec_fd_blocking
6419 .lock()
6420 .unwrap()
6421 .insert(leader, [42].into_iter().collect());
6422 state
6423 .recv_deregister_thread(
6424 reverie::Tid::from_raw(leader.as_raw()),
6425 ThreadDeregistration {
6426 thread_start_entered: true,
6427 dettid: leader,
6428 detpid,
6429 mm: MmId::initial(detpid),
6430 timeslice_stats: TimesliceStats::default(),
6431 syscall_count: 0,
6432 chaos_epochs: Vec::new(),
6433 },
6434 )
6435 .await;
6436 assert!(state.pending_exec_states.lock().unwrap().is_empty());
6437 assert!(state.post_exec_fd_blocking.lock().unwrap().is_empty());
6438
6439 state.pending_exec_states.lock().unwrap().insert(
6440 detpid,
6441 PendingExecState {
6442 caller: leader,
6443 process: detpid,
6444 mm: MmId::initial(detpid),
6445 fd_blocking: Default::default(),
6446 },
6447 );
6448 state
6449 .recv_deregister_thread(
6450 reverie::Tid::from_raw(leader.as_raw()),
6451 ThreadDeregistration {
6452 thread_start_entered: true,
6453 dettid: leader,
6454 detpid,
6455 mm: MmId::initial(detpid).for_exec(detpid),
6456 timeslice_stats: TimesliceStats::default(),
6457 syscall_count: 0,
6458 chaos_epochs: Vec::new(),
6459 },
6460 )
6461 .await;
6462 assert!(state.pending_exec_states.lock().unwrap().is_empty());
6463
6464 state.pending_exec_states.lock().unwrap().insert(
6465 detpid,
6466 PendingExecState {
6467 caller: leader,
6468 process: detpid,
6469 mm: MmId::initial(detpid),
6470 fd_blocking: Default::default(),
6471 },
6472 );
6473 state
6474 .post_exec_fd_blocking
6475 .lock()
6476 .unwrap()
6477 .insert(leader, [42].into_iter().collect());
6478 state.complete_physical_process_exit(detpid.as_raw());
6479 assert!(state.pending_exec_states.lock().unwrap().is_empty());
6480 assert!(state.post_exec_fd_blocking.lock().unwrap().is_empty());
6481 }
6482
6483 #[test]
6484 fn unsupported_syscall_report_duplicate_is_close_on_exec() {
6485 let mut descriptors = [-1; 2];
6486 assert_eq!(
6487 unsafe { libc::pipe2(descriptors.as_mut_ptr(), libc::O_CLOEXEC) },
6488 0
6489 );
6490 let _reader = unsafe { OwnedFd::from_raw_fd(descriptors[0]) };
6492 let writer = unsafe { OwnedFd::from_raw_fd(descriptors[1]) };
6493 let config = Config {
6494 unsupported_syscall_report_fd: Some(writer.as_raw_fd()),
6495 ..Config::default()
6496 };
6497
6498 let state = GlobalState::initialize(&config, false);
6499 let duplicate = state
6500 .unsupported_syscall_report_fd
6501 .as_ref()
6502 .expect("report writer should be duplicated")
6503 .lock()
6504 .unwrap();
6505 let flags = unsafe { libc::fcntl(duplicate.as_raw_fd(), libc::F_GETFD) };
6506 assert_ne!(flags, -1);
6507 assert_ne!(flags & libc::FD_CLOEXEC, 0);
6508 }
6509
6510 #[test]
6513 fn device_pool_remaps_deterministically() {
6514 use super::DevicePool;
6515
6516 let raw_root = 0x20; let raw_proc_run1 = 3_145_792; let raw_proc_run2 = 3_145_788; let mut pool1 = DevicePool::new();
6524 let root1 = pool1.determinize(raw_root);
6525 let proc1 = pool1.determinize(raw_proc_run1);
6526 let mut pool2 = DevicePool::new();
6528 let root2 = pool2.determinize(raw_root);
6529 let proc2 = pool2.determinize(raw_proc_run2);
6530
6531 assert_eq!(root1, root2);
6534 assert_eq!(proc1, proc2);
6535
6536 assert_ne!(root1, proc1);
6538 assert_eq!(root1, 1);
6539 assert_eq!(proc1, 2);
6540
6541 assert_eq!(pool1.determinize(raw_root), root1);
6543 assert_eq!(pool1.determinize(raw_proc_run1), proc1);
6544 }
6545
6546 #[test]
6547 fn mountinfo_prepopulation_is_reused_by_later_stat_observations() {
6548 use super::DevicePool;
6549
6550 let mountinfo_devices = [libc::makedev(8, 1), libc::makedev(0, 44)];
6551 let mut pool = DevicePool::new();
6552 let rendered = mountinfo_devices
6553 .into_iter()
6554 .map(|raw| pool.determinize(raw))
6555 .collect::<Vec<_>>();
6556
6557 assert_eq!(rendered, [1, 2]);
6558 assert_eq!(pool.determinize(libc::makedev(0, 44)), rendered[1]);
6559 assert_eq!(pool.determinize(libc::makedev(8, 1)), rendered[0]);
6560 assert_eq!(pool.determinize(libc::makedev(259, 7)), 3);
6561 }
6562
6563 #[tokio::test]
6564 async fn late_futex_rpc_after_thread_removal_returns_eintr() {
6565 let config = Config {
6566 sequentialize_threads: true,
6567 ..Config::default()
6568 };
6569 let state = GlobalState::initialize(&config, false);
6570 let dettid = DetTid::from_raw(17);
6571 let detpid = DetPid::from_raw(17);
6572 let response = state
6573 .recv_futex_action(
6574 RpcIncarnation {
6575 dettid,
6576 mm: MmId::initial(detpid),
6577 },
6578 FutexAction::WaitRequest(None),
6579 FutexID::private(MmId::initial(detpid), 0x1000),
6580 0,
6581 u32::MAX,
6582 )
6583 .await;
6584
6585 assert!(matches!(
6586 response,
6587 Some(SchedValue::Value(value)) if value == nix::errno::Errno::EINTR as u64
6588 ));
6589 }
6590
6591 #[tokio::test]
6592 async fn late_child_tid_wait_after_exit_returns_spurious_wake() {
6593 let config = Config {
6594 sequentialize_threads: true,
6595 ..Config::default()
6596 };
6597 let state = GlobalState::initialize(&config, false);
6598 let detpid = DetPid::from_raw(17);
6599 let child = DetTid::from_raw(18);
6600 let futex = FutexID::private(MmId::initial(detpid), 0x1000);
6601 state.sched.lock().unwrap().next_turns.insert(
6602 detpid,
6603 ThreadNextTurn {
6604 dettid: detpid,
6605 child_tid_addr: 0,
6606 req: Ivar::new(),
6607 resp: Ivar::new(),
6608 protocol: Default::default(),
6609 },
6610 );
6611 state
6612 .sched
6613 .lock()
6614 .unwrap()
6615 .wake_futex_child_cleartid(futex, child);
6616 assert!(
6617 state
6618 .sched
6619 .lock()
6620 .unwrap()
6621 .child_tid_was_cleared(futex, child.as_raw())
6622 );
6623 assert!(
6624 !state
6625 .sched
6626 .lock()
6627 .unwrap()
6628 .child_tid_was_cleared(futex, child.as_raw() + 1)
6629 );
6630
6631 let response = state
6632 .recv_futex_action(
6633 RpcIncarnation {
6634 dettid: detpid,
6635 mm: MmId::initial(detpid),
6636 },
6637 FutexAction::WaitRequest(None),
6638 futex,
6639 child.as_raw(),
6640 u32::MAX,
6641 )
6642 .await;
6643
6644 assert!(matches!(response, Some(SchedValue::Value(0))));
6645 assert!(state.sched.lock().unwrap().blocked.futex_waiters.is_empty());
6646 }
6647
6648 #[tokio::test]
6649 async fn late_resource_request_after_logical_kill_is_cancelled() {
6650 let (config, state, dettid, detpid) = cancellation_test_state();
6651 install_test_registration(&state, dettid, Ivar::new());
6652 let mut current_time = DetTime::new(&config);
6653 current_time.add_syscall();
6654 state.global_time.lock().unwrap().update_global_time(
6655 dettid,
6656 current_time.as_nanos(),
6657 LogicalTime::ZERO,
6658 );
6659 state
6660 .sched
6661 .lock()
6662 .unwrap()
6663 .logically_kill_thread(&dettid, &detpid, MmId::initial(detpid));
6664 let (global_before, thread_before) = {
6665 let global_time = state.global_time.lock().unwrap();
6666 (global_time.as_nanos(), global_time.threads_time(dettid))
6667 };
6668 let mut late_time = current_time;
6669 late_time.add_syscall();
6670
6671 let response = state
6672 .receive_rpc(
6673 reverie::Tid::from_raw(dettid.as_raw()),
6674 (
6675 late_time,
6676 MmId::initial(dettid),
6677 GlobalRequest::RequestResources(Resources::new(dettid), detpid),
6678 ),
6679 )
6680 .await;
6681
6682 assert_eq!(response, (None, GlobalResponse::ThreadExited));
6683 assert!(!state.sched.lock().unwrap().next_turns.contains_key(&dettid));
6684 let global_time = state.global_time.lock().unwrap();
6685 assert_eq!(global_time.as_nanos(), global_before);
6686 assert_eq!(global_time.threads_time(dettid), thread_before);
6687 }
6688
6689 #[tokio::test]
6690 async fn duplicate_deregistration_is_acknowledged_without_clock_or_scheduler_mutation() {
6691 let (config, state, dettid, detpid) = cancellation_test_state();
6692 install_test_registration(&state, dettid, Ivar::new());
6693 let mut current_time = DetTime::new(&config);
6694 current_time.add_syscall();
6695 state.global_time.lock().unwrap().update_global_time(
6696 dettid,
6697 current_time.as_nanos(),
6698 LogicalTime::ZERO,
6699 );
6700 state
6701 .sched
6702 .lock()
6703 .unwrap()
6704 .logically_kill_thread(&dettid, &detpid, MmId::initial(detpid));
6705 let (global_before, thread_before) = {
6706 let global_time = state.global_time.lock().unwrap();
6707 (global_time.as_nanos(), global_time.threads_time(dettid))
6708 };
6709 let mut late_time = current_time;
6710 late_time.add_syscall();
6711 let mut final_stats = TimesliceStats::default();
6712 final_stats.record(7);
6713
6714 let first_response = state
6715 .receive_rpc(
6716 reverie::Tid::from_raw(dettid.as_raw()),
6717 (
6718 late_time.clone(),
6719 MmId::initial(dettid),
6720 GlobalRequest::DeregisterThread(ThreadDeregistration {
6721 thread_start_entered: true,
6722 dettid,
6723 detpid,
6724 mm: MmId::initial(detpid),
6725 timeslice_stats: final_stats,
6726 syscall_count: 17,
6727 chaos_epochs: Vec::new(),
6728 }),
6729 ),
6730 )
6731 .await;
6732 assert_eq!(first_response, (None, GlobalResponse::DeregisterThread(())));
6733 assert_eq!(
6734 state
6735 .sched
6736 .lock()
6737 .unwrap()
6738 .per_thread_timeslice
6739 .get(&dettid),
6740 Some(&final_stats)
6741 );
6742 assert_eq!(
6743 state.sched.lock().unwrap().per_thread_syscalls.get(&dettid),
6744 Some(&17)
6745 );
6746
6747 late_time.add_syscall();
6748 let duplicate_response = state
6749 .receive_rpc(
6750 reverie::Tid::from_raw(dettid.as_raw()),
6751 (
6752 late_time,
6753 MmId::initial(dettid),
6754 GlobalRequest::DeregisterThread(ThreadDeregistration {
6755 thread_start_entered: true,
6756 dettid,
6757 detpid,
6758 mm: MmId::initial(detpid),
6759 timeslice_stats: final_stats,
6760 syscall_count: 99,
6761 chaos_epochs: Vec::new(),
6762 }),
6763 ),
6764 )
6765 .await;
6766 assert_eq!(
6767 duplicate_response,
6768 (None, GlobalResponse::DeregisterThread(()))
6769 );
6770 assert_eq!(
6771 state
6772 .sched
6773 .lock()
6774 .unwrap()
6775 .per_thread_timeslice
6776 .get(&dettid),
6777 Some(&final_stats)
6778 );
6779 assert_eq!(
6780 state.sched.lock().unwrap().per_thread_syscalls.get(&dettid),
6781 Some(&17),
6782 "a duplicate deregistration must not double-count or replace final accounting"
6783 );
6784 let summary = state
6785 .sched
6786 .lock()
6787 .unwrap()
6788 .generate_partial_run_summary(None)
6789 .unwrap();
6790 assert_eq!(summary.syscalls, Some(17));
6791 assert!(!state.sched.lock().unwrap().next_turns.contains_key(&dettid));
6792 let global_time = state.global_time.lock().unwrap();
6793 assert_eq!(global_time.as_nanos(), global_before);
6794 assert_eq!(global_time.threads_time(dettid), thread_before);
6795 }
6796
6797 #[tokio::test]
6798 async fn child_registration_fails_closed_for_a_tombstoned_tid() {
6799 let (config, state, parent, detpid) = cancellation_test_state();
6800 install_test_registration(&state, parent, Ivar::new());
6801 let child = DetTid::from_raw(18);
6802 state
6803 .sched
6804 .lock()
6805 .unwrap()
6806 .thread_tree
6807 .add_child(parent, child, false);
6808 install_test_registration(&state, child, Ivar::new());
6809 state
6810 .sched
6811 .lock()
6812 .unwrap()
6813 .logically_kill_thread(&child, &detpid, MmId::initial(detpid));
6814
6815 let response = state
6816 .receive_rpc(
6817 reverie::Tid::from_raw(parent.as_raw()),
6818 (
6819 DetTime::new(&config),
6820 MmId::initial(parent),
6821 GlobalRequest::CreateChildThread(
6822 child,
6823 detpid,
6824 0,
6825 None,
6826 libc::SIGCHLD,
6827 None,
6828 Some(DEFAULT_PRIORITY),
6829 ),
6830 ),
6831 )
6832 .await;
6833
6834 assert_eq!(response, (None, GlobalResponse::ThreadExited));
6835 let scheduler = state.sched.lock().unwrap();
6836 assert!(scheduler.thread_is_logically_killed(child));
6837 assert!(!scheduler.next_turns.contains_key(&child));
6838 assert!(!scheduler.priorities.contains_key(&child));
6839 assert!(!state.global_time.lock().unwrap().contains_thread(parent));
6840 }
6841
6842 #[tokio::test]
6843 async fn pending_resource_request_woken_by_logical_kill_is_terminal() {
6844 let (_, state, dettid, detpid) = cancellation_test_state();
6845 let request_seen = Ivar::new();
6846 install_test_registration(&state, dettid, request_seen.clone());
6847
6848 let request = state.recv_request_resources(
6849 reverie::Tid::from_raw(dettid.as_raw()),
6850 detpid,
6851 Resources::new(dettid),
6852 None,
6853 );
6854 let kill_after_request = async {
6855 while request_seen.try_read().is_none() {
6856 tokio::task::yield_now().await;
6857 }
6858 state.sched.lock().unwrap().logically_kill_thread(
6859 &dettid,
6860 &detpid,
6861 MmId::initial(detpid),
6862 );
6863 };
6864
6865 let (response, ()) = tokio::join!(request, kill_after_request);
6866 assert_eq!(response, (SchedulerRpcResult::ThreadExited, None));
6867 }
6868
6869 #[tokio::test]
6870 async fn trace_replay_yield_propagates_terminal_scheduler_cancellation() {
6871 let dettid = DetTid::from_raw(17);
6872 let detpid = DetPid::from_raw(17);
6873 let next_tid = DetTid::from_raw(18);
6874 let event = SchedEvent {
6875 dettid,
6876 op: Op::OtherInstructions,
6877 count: 1,
6878 start_rip: None,
6879 end_rip: None,
6880 end_time: Some(LogicalTime::from_nanos(1)),
6881 };
6882 let next_event = SchedEvent {
6883 dettid: next_tid,
6884 end_time: Some(LogicalTime::from_nanos(2)),
6885 ..event.clone()
6886 };
6887 let trace_file = tempfile::NamedTempFile::new().unwrap();
6888 std::fs::write(
6889 trace_file.path(),
6890 PreemptionRecord::from_sched_events(vec![event.clone(), next_event]).to_string(),
6891 )
6892 .unwrap();
6893 let config = Config {
6894 sequentialize_threads: true,
6895 cancel_killed_thread_rpcs: true,
6896 replay_schedule_from: Some(trace_file.path().to_path_buf()),
6897 ..Config::default()
6898 };
6899 let state = GlobalState::initialize(&config, false);
6900 state
6901 .sched
6902 .lock()
6903 .unwrap()
6904 .thread_tree
6905 .add_child(dettid, dettid, true);
6906 state
6907 .sched
6908 .lock()
6909 .unwrap()
6910 .thread_tree
6911 .add_child(dettid, next_tid, false);
6912 let request_seen = Ivar::new();
6913 install_test_registration(&state, dettid, request_seen.clone());
6914 install_test_registration(&state, next_tid, Ivar::new());
6915
6916 let replay = state.recv_trace_schedevent(event, detpid, MmId::initial(detpid), false);
6917 let kill_after_replay_yield = async {
6918 while request_seen.try_read().is_none() {
6919 tokio::task::yield_now().await;
6920 }
6921 state.sched.lock().unwrap().logically_kill_thread(
6922 &dettid,
6923 &detpid,
6924 MmId::initial(detpid),
6925 );
6926 };
6927 let (response, ()) = tokio::time::timeout(Duration::from_secs(1), async {
6928 tokio::join!(replay, kill_after_replay_yield)
6929 })
6930 .await
6931 .expect("trace replay cancellation did not terminate the pending scheduler RPC");
6932
6933 assert_eq!(response, SchedulerRpcResult::ThreadExited);
6934 }
6935
6936 #[tokio::test]
6937 async fn pending_start_request_woken_by_logical_kill_is_terminal() {
6938 let (config, state, dettid, detpid) = cancellation_test_state();
6939 let request_seen = Ivar::new();
6940 install_test_registration(&state, dettid, request_seen.clone());
6941 let request = state.receive_rpc(
6942 reverie::Tid::from_raw(dettid.as_raw()),
6943 (
6944 DetTime::new(&config),
6945 MmId::initial(dettid),
6946 GlobalRequest::StartNewThread(dettid, detpid, None, None),
6947 ),
6948 );
6949 let kill_after_request = async {
6950 while request_seen.try_read().is_none() {
6951 tokio::task::yield_now().await;
6952 }
6953 state.sched.lock().unwrap().logically_kill_thread(
6954 &dettid,
6955 &detpid,
6956 MmId::initial(detpid),
6957 );
6958 };
6959
6960 let (response, ()) = tokio::join!(request, kill_after_request);
6961 assert_eq!(response, (None, GlobalResponse::ThreadExited));
6962 assert!(!state.sched.lock().unwrap().priorities.contains_key(&dettid));
6963 }
6964
6965 #[tokio::test]
6966 async fn required_physical_thread_id_missing_is_terminal() {
6967 let config = Config {
6968 sequentialize_threads: true,
6969 cancel_killed_thread_rpcs: true,
6970 backend_requires_thread_directed_process_signals: true,
6971 ..Config::default()
6972 };
6973 let state = GlobalState::initialize(&config, false);
6974 let dettid = DetTid::from_raw(17);
6975 let detpid = DetPid::from_raw(17);
6976 state
6977 .sched
6978 .lock()
6979 .unwrap()
6980 .thread_tree
6981 .add_child(dettid, dettid, true);
6982 install_test_registration(&state, dettid, Ivar::new());
6983
6984 let response = state
6985 .receive_rpc(
6986 reverie::Tid::from_raw(dettid.as_raw()),
6987 (
6988 DetTime::new(&config),
6989 MmId::initial(detpid),
6990 GlobalRequest::StartNewThread(dettid, detpid, None, None),
6991 ),
6992 )
6993 .await;
6994
6995 assert_eq!(response, (None, GlobalResponse::ThreadExited));
6996 let scheduler = state.sched.lock().unwrap();
6997 assert!(!scheduler.next_turns.contains_key(&dettid));
6998 }
6999
7000 #[tokio::test]
7001 async fn tombstoned_timer_registration_cannot_mutate_scheduler_state() {
7002 let (_, state, dettid, detpid) = cancellation_test_state();
7003 install_test_registration(&state, dettid, Ivar::new());
7004 state
7005 .sched
7006 .lock()
7007 .unwrap()
7008 .logically_kill_thread(&dettid, &detpid, MmId::initial(detpid));
7009 let now = LogicalTime::from_nanos(100);
7010
7011 let alarm = state
7012 .recv_register_alarm(
7013 detpid,
7014 RpcIncarnation {
7015 dettid,
7016 mm: MmId::initial(detpid),
7017 },
7018 now,
7019 LogicalTime::from_nanos(10),
7020 LogicalTime::ZERO,
7021 SigWrapper::from(Signal::SIGALRM),
7022 )
7023 .await;
7024 assert_eq!(alarm, SchedulerRpcResult::ThreadExited);
7025
7026 let posix = state
7027 .recv_register_posix_timer(
7028 detpid,
7029 RpcIncarnation {
7030 dettid,
7031 mm: MmId::initial(detpid),
7032 },
7033 1,
7034 Some(now + LogicalTime::from_nanos(10)),
7035 LogicalTime::ZERO,
7036 SigWrapper::from(Signal::SIGALRM),
7037 )
7038 .await;
7039 assert_eq!(posix, SchedulerRpcResult::ThreadExited);
7040 assert!(state.sched.lock().unwrap().blocked.timed_waiters.is_empty());
7041 }
7042
7043 #[tokio::test]
7044 async fn pending_futex_request_woken_by_logical_kill_is_terminal() {
7045 let (config, state, dettid, detpid) = cancellation_test_state();
7046 install_test_registration(&state, dettid, Ivar::new());
7047 let request = state.receive_rpc(
7048 reverie::Tid::from_raw(dettid.as_raw()),
7049 (
7050 DetTime::new(&config),
7051 MmId::initial(dettid),
7052 GlobalRequest::FutexAction(
7053 dettid,
7054 FutexAction::WaitRequest(None),
7055 FutexID::private(MmId::initial(detpid), 0x1000),
7056 0,
7057 u32::MAX,
7058 ),
7059 ),
7060 );
7061 let kill_after_wait = async {
7062 while state.sched.lock().unwrap().blocked.futex_waiters.is_empty() {
7063 tokio::task::yield_now().await;
7064 }
7065 state.sched.lock().unwrap().logically_kill_thread(
7066 &dettid,
7067 &detpid,
7068 MmId::initial(detpid),
7069 );
7070 };
7071
7072 let (response, ()) = tokio::time::timeout(Duration::from_secs(1), async {
7073 tokio::join!(request, kill_after_wait)
7074 })
7075 .await
7076 .expect("futex teardown did not wake the blocked RPC");
7077 assert_eq!(response, (None, GlobalResponse::ThreadExited));
7078 }
7079
7080 #[tokio::test]
7081 async fn parent_continue_propagates_terminal_scheduler_cancellation() {
7082 let (config, state, parent, detpid) = cancellation_test_state();
7083 let parent_request = Ivar::new();
7084 install_test_registration(&state, parent, parent_request.clone());
7085 let child = DetTid::from_raw(18);
7086 let physical_pid = std::process::id() as i32;
7087 let physical_tid = unsafe { libc::syscall(libc::SYS_gettid) as i32 };
7088 let request = state.receive_rpc(
7089 reverie::Tid::from_raw(parent.as_raw()),
7090 (
7091 DetTime::new(&config),
7092 MmId::initial(parent),
7093 GlobalRequest::CreateChildThread(
7094 child,
7095 detpid,
7096 0,
7097 None,
7098 libc::SIGCHLD,
7099 Some((physical_pid, physical_tid)),
7100 Some(DEFAULT_PRIORITY),
7101 ),
7102 ),
7103 );
7104 let kill_after_parent_parks = async {
7105 while parent_request.try_read().is_none() {
7106 tokio::task::yield_now().await;
7107 }
7108 let mut scheduler = state.sched.lock().unwrap();
7109 let (registered_mm, registered_pid, registered_tid) = scheduler
7110 .physical_thread_identity(child)
7111 .expect("parent registration must install the child pidfd before continuing");
7112 assert_eq!(
7113 registered_mm,
7114 MmId::for_clone(MmId::initial(parent), child, false)
7115 );
7116 assert_eq!(
7117 (registered_pid, registered_tid),
7118 (physical_pid, physical_tid)
7119 );
7120 scheduler.logically_kill_thread(&parent, &detpid, MmId::initial(detpid));
7121 };
7122
7123 let (response, ()) = tokio::join!(request, kill_after_parent_parks);
7124 assert_eq!(response, (None, GlobalResponse::ThreadExited));
7125 }
7126
7127 #[test]
7128 fn unsupported_syscall_warning_is_sorted_and_aggregated() {
7129 let syscalls = BTreeSet::from([
7130 "vmsplice".to_owned(),
7131 "getppid".to_owned(),
7132 "getppid".to_owned(),
7133 ]);
7134
7135 assert_eq!(
7136 format_unsupported_syscall_warning(&syscalls).as_deref(),
7137 Some("syscalls getppid,vmsplice used but not yet supported")
7138 );
7139 assert_eq!(format_unsupported_syscall_warning(&BTreeSet::new()), None);
7140 }
7141
7142 #[tokio::test]
7143 async fn abnormal_cleanup_cancels_an_unstarted_scheduler() {
7144 let config = Config {
7145 sequentialize_threads: true,
7146 ..Config::default()
7147 };
7148 let mut state = GlobalState::initialize(&config, true);
7149 state.cancel_internal_scheduler().await;
7150 let summary_path = None;
7151 let cleanup = state.clean_up(false, &summary_path);
7152
7153 assert!(
7154 tokio::time::timeout(Duration::from_millis(100), cleanup)
7155 .await
7156 .is_ok(),
7157 "cleanup waited for a scheduler whose guest never registered"
7158 );
7159 }
7160
7161 #[tokio::test]
7162 async fn abnormal_cleanup_cancels_a_registered_scheduler() {
7163 let config = Config {
7164 sequentialize_threads: true,
7165 ..Config::default()
7166 };
7167 let mut state = GlobalState::initialize(&config, true);
7168 let dettid = DetTid::from_raw(1);
7169 {
7170 let mut scheduler = state.sched.lock().unwrap();
7171 scheduler.priorities.insert(dettid, DEFAULT_PRIORITY);
7172 scheduler.next_turns.insert(
7173 dettid,
7174 ThreadNextTurn {
7175 dettid,
7176 child_tid_addr: 0,
7177 req: Ivar::new(),
7178 resp: Ivar::new(),
7179 protocol: Default::default(),
7180 },
7181 );
7182 scheduler.runqueue_push_back(dettid);
7183 scheduler.started_up.put(());
7184 }
7185 tokio::task::yield_now().await;
7186
7187 state.cancel_internal_scheduler().await;
7188 let summary_path = None;
7189 assert!(
7190 tokio::time::timeout(
7191 Duration::from_millis(100),
7192 state.clean_up(false, &summary_path),
7193 )
7194 .await
7195 .is_ok(),
7196 "cleanup waited after cancelling a registered scheduler"
7197 );
7198 }
7199
7200 #[test]
7206 fn det_inodes_are_minted_not_passed_through() {
7207 use crate::types::DetInode;
7208
7209 let mut pool = super::InodePool::new();
7210 let t = LogicalTime::from_nanos(0);
7211 let seen = super::ObservedMtime::Unobserved;
7212
7213 let host_a = 221_742_951; let host_b = 998_877_665;
7215 let (a, _) = pool.add_inode(host_a, seen, t);
7216 let (b, _) = pool.add_inode(host_b, seen, t);
7217
7218 assert_ne!(a.as_raw(), host_a, "det inode must not be the host inode");
7219 assert_ne!(b.as_raw(), host_b, "det inode must not be the host inode");
7220 assert_eq!(a, DetInode::mint(1), "minting starts at 1");
7221 assert_eq!(b, DetInode::mint(2), "minting is monotonic");
7222
7223 let (a_again, _) = pool.add_inode(host_a, seen, t);
7225 assert_eq!(a, a_again, "mapping must be stable per host inode");
7226 }
7227
7228 #[test]
7231 fn only_exact_canonical_host_mtimes_are_kept() {
7232 use super::ObservedMtime;
7233
7234 assert_eq!(
7235 ObservedMtime::from_host_mtime(1, 0),
7236 ObservedMtime::Canonical(LogicalTime::from_secs(1))
7237 );
7238 assert_eq!(
7239 ObservedMtime::from_host_mtime(0, 0),
7240 ObservedMtime::Canonical(LogicalTime::from_secs(0))
7241 );
7242 for (secs, nanos) in [(1, 5), (0, 1), (2, 0), (-1, 0), (1_600_000_000, 0)] {
7243 assert_eq!(
7244 ObservedMtime::from_host_mtime(secs, nanos),
7245 ObservedMtime::HostSpecific,
7246 "{secs}.{nanos:09} is a real timestamp, not a canonical one"
7247 );
7248 }
7249 }
7250
7251 #[test]
7254 fn first_seen_mtime_is_resolved_by_the_first_stat() {
7255 use super::ObservedMtime;
7256
7257 let epoch = LogicalTime::from_secs(1_798_761_600);
7258 let canonical = ObservedMtime::Canonical(LogicalTime::from_secs(1));
7259 let mut pool = super::InodePool::new();
7260
7261 let (_, mtime) = pool.add_inode(10, ObservedMtime::HostSpecific, epoch);
7263 assert_eq!(mtime, epoch);
7264 let (_, mtime) = pool.add_inode(11, canonical, epoch);
7266 assert_eq!(mtime, LogicalTime::from_secs(1));
7267
7268 let (_, mtime) = pool.add_inode(12, ObservedMtime::Unobserved, epoch);
7270 assert_eq!(mtime, epoch, "an unresolved mtime reads as the epoch");
7271 let (_, mtime) = pool.add_inode(12, canonical, epoch);
7272 assert_eq!(mtime, LogicalTime::from_secs(1));
7273
7274 let (_, mtime) = pool.add_inode(10, canonical, epoch);
7276 assert_eq!(mtime, epoch);
7277 let (_, mtime) = pool.add_inode(11, ObservedMtime::HostSpecific, epoch);
7278 assert_eq!(mtime, LogicalTime::from_secs(1));
7279 }
7280}
7281
7282#[cfg(test)]
7283mod robust_exit_clock_tests {
7284 use std::sync::Mutex;
7285 use std::task::Poll;
7286
7287 use nix::sys::signal::Signal;
7288 use reverie::ExitStatus;
7289 use reverie::GlobalRPC;
7290 use reverie::GlobalTool;
7291 use reverie::Tid;
7292 use reverie::Tool;
7293
7294 use super::GlobalRequest;
7295 use super::GlobalResponse;
7296 use super::GlobalState;
7297 use crate::Detcore;
7298 use crate::ThreadState;
7299 use crate::config::Config;
7300 use crate::ivar::Ivar;
7301 use crate::resources::Resources;
7302 use crate::scheduler::DEFAULT_PRIORITY;
7303 use crate::scheduler::SkipTurn;
7304 use crate::scheduler::ThreadNextTurn;
7305 use crate::tool_local::RobustListExit;
7306 use crate::tool_local::RobustListWake;
7307 use crate::types::*;
7308
7309 #[derive(Debug, PartialEq)]
7310 struct WakeObservation {
7311 wakes: Vec<(DetTid, FutexID)>,
7312 counts: Vec<u64>,
7313 clocks: serde_json::Value,
7314 turn: u64,
7315 }
7316
7317 #[derive(Debug)]
7318 struct RpcObservation {
7319 sender: DetTid,
7320 kind: &'static str,
7321 accepted: bool,
7322 clocks: serde_json::Value,
7323 queued: Vec<DetTid>,
7324 waiters: usize,
7325 turn: u64,
7326 }
7327
7328 struct ExitRpc<'a> {
7331 state: &'a GlobalState,
7332 sender: DetTid,
7333 observations: &'a Mutex<Vec<WakeObservation>>,
7334 rpc_observations: &'a Mutex<Vec<RpcObservation>>,
7335 }
7336
7337 #[reverie::tool]
7338 impl GlobalRPC<GlobalState> for ExitRpc<'_> {
7339 async fn send_rpc(
7340 &self,
7341 request: <GlobalState as GlobalTool>::Request,
7342 ) -> <GlobalState as GlobalTool>::Response {
7343 let kind = match &request.2 {
7344 GlobalRequest::RobustListWakes(wakes) if wakes.is_empty() => "empty-wake",
7345 GlobalRequest::RobustListWakes(_) => "wake",
7346 GlobalRequest::DeregisterThread(_) => "deregister",
7347 _ => panic!("unexpected exit RPC: {:?}", request.2),
7348 };
7349 let wakes = match &request.2 {
7350 GlobalRequest::RobustListWakes(wakes) if !wakes.is_empty() => Some(wakes.clone()),
7351 _ => None,
7352 };
7353 let response = self
7354 .state
7355 .receive_rpc(Tid::from_raw(self.sender.as_raw()), request)
7356 .await;
7357 let clocks = serde_json::to_value(&*self.state.global_time.lock().unwrap()).unwrap();
7358 {
7359 let sched = self.state.sched.lock().unwrap();
7360 self.rpc_observations.lock().unwrap().push(RpcObservation {
7361 sender: self.sender,
7362 kind,
7363 accepted: !matches!(response.1, GlobalResponse::ThreadExited),
7364 clocks: clocks.clone(),
7365 queued: sched.run_queue.tids().copied().collect(),
7366 waiters: sched
7367 .blocked
7368 .futex_waiters
7369 .values()
7370 .map(Vec::len)
7371 .sum::<usize>(),
7372 turn: sched.turn,
7373 });
7374 }
7375 if let Some(wakes) = wakes {
7376 let GlobalResponse::RobustListWakes(counts) = &response.1 else {
7377 panic!("a complete admitted exit batch was refused: {response:?}");
7378 };
7379 let clocks =
7380 serde_json::to_value(&*self.state.global_time.lock().unwrap()).unwrap();
7381 let turn = self.state.sched.lock().unwrap().turn;
7382 self.observations.lock().unwrap().push(WakeObservation {
7383 wakes,
7384 counts: counts.clone(),
7385 clocks,
7386 turn,
7387 });
7388 }
7389 response
7390 }
7391
7392 fn config(&self) -> &Config {
7393 &self.state.cfg
7394 }
7395 }
7396
7397 struct Fixture {
7398 state: GlobalState,
7399 tool: Detcore,
7400 owners: [ThreadState<()>; 2],
7401 initial_clocks: [DetTime; 2],
7402 waiters: [DetTid; 2],
7403 peer: DetTid,
7404 futexes: [FutexID; 2],
7405 observations: Mutex<Vec<WakeObservation>>,
7406 rpc_observations: Mutex<Vec<RpcObservation>>,
7407 }
7408
7409 impl Fixture {
7410 fn new(
7411 reason: RobustListExit,
7412 equal_clocks: bool,
7413 empty_owner: Option<usize>,
7414 cancel_killed_thread_rpcs: bool,
7415 ) -> Self {
7416 let config = Config {
7417 sequentialize_threads: true,
7418 cancel_killed_thread_rpcs,
7419 ..Config::default()
7420 };
7421 let state = GlobalState::initialize(&config, false);
7422 let parent = DetTid::from_raw(1);
7423 let leader = DetTid::from_raw(17);
7424 let worker = DetTid::from_raw(18);
7425 let waiters = [DetTid::from_raw(21), DetTid::from_raw(23)];
7426 let peer = DetTid::from_raw(25);
7427 let mm = MmId::initial(leader);
7428 let mut first = ThreadState::new(leader, &config, ());
7429 first.detpid = Some(leader);
7430 let mut second = first.clone();
7431 second.dettid = worker;
7432 let mut inherited = DetTime::new(&config);
7433 inherited.add_syscall_with_cost(1_000);
7434 let first_initial = inherited.clone_for_child();
7435 inherited.add_syscall_with_cost(370);
7436 let second_initial = inherited.clone_for_child();
7437 first.thread_logical_time = first_initial.clone();
7438 second.thread_logical_time = second_initial.clone();
7439 first
7440 .thread_logical_time
7441 .add_syscall_with_cost(if equal_clocks { 407 } else { 37 });
7442 second.thread_logical_time.add_syscall_with_cost(37);
7443 first.record_robust_list_head(Some(0x404100));
7444 second.record_robust_list_head(Some(0x404200));
7445 let object = SharedMemoryObjectId::Anonymous {
7446 origin: MmId::initial(parent),
7447 sequence: 1,
7448 };
7449 let futexes = [FutexID::shared(object, 0), FutexID::shared(object, 8)];
7450 first.stage_robust_list_wakes(
7451 reason,
7452 vec![
7453 (
7454 worker,
7455 if empty_owner == Some(1) {
7456 Vec::new()
7457 } else {
7458 vec![RobustListWake { futex: futexes[1] }]
7459 },
7460 ),
7461 (
7462 leader,
7463 if empty_owner == Some(0) {
7464 Vec::new()
7465 } else {
7466 vec![RobustListWake { futex: futexes[0] }]
7467 },
7468 ),
7469 ],
7470 );
7471 {
7472 let mut sched = state.sched.lock().unwrap();
7473 sched.thread_tree.add_child(parent, parent, true);
7474 sched.thread_tree.add_child(parent, leader, true);
7475 sched.thread_tree.add_child(leader, worker, false);
7476 for tid in [waiters[0], waiters[1], peer] {
7477 sched.thread_tree.add_child(parent, tid, true);
7478 }
7479 for tid in [leader, worker, waiters[0], waiters[1], peer] {
7480 sched.priorities.insert(tid, DEFAULT_PRIORITY);
7481 sched.next_turns.insert(
7482 tid,
7483 ThreadNextTurn {
7484 dettid: tid,
7485 child_tid_addr: 0,
7486 req: Ivar::new(),
7487 resp: Ivar::new(),
7488 protocol: Default::default(),
7489 },
7490 );
7491 sched.install_test_exec_incarnation(
7492 tid,
7493 if tid == leader || tid == worker {
7494 mm
7495 } else {
7496 MmId::initial(tid)
7497 },
7498 );
7499 }
7500 sched.runqueue_push_back(leader);
7504 sched.runqueue_push_back(worker);
7505 sched.runqueue_push_back(peer);
7506 sched.next_turns[&peer].req.put(Ok(Resources::new(peer)));
7507 for (waiter, futex) in waiters.into_iter().zip(futexes) {
7508 sched.sleep_futex_waiter(&waiter, futex, None, u32::MAX);
7509 }
7510 }
7511 {
7512 let mut time = state.global_time.lock().unwrap();
7513 for (tid, clock) in [(leader, &first_initial), (worker, &second_initial)] {
7514 time.update_global_time(tid, clock.as_nanos(), clock.inherited_nanos());
7515 }
7516 }
7517 let tool = Detcore::new(Tid::from_raw(leader.as_raw()), &config);
7518 Self {
7519 state,
7520 tool,
7521 owners: [first, second],
7522 initial_clocks: [first_initial, second_initial],
7523 waiters,
7524 peer,
7525 futexes,
7526 observations: Mutex::new(Vec::new()),
7527 rpc_observations: Mutex::new(Vec::new()),
7528 }
7529 }
7530
7531 fn assert_clocks(&self, completed: &[usize]) -> serde_json::Value {
7532 let time = self.state.global_time.lock().unwrap();
7533 let snapshot = serde_json::to_value(&*time).unwrap();
7534 let epoch = DetTime::new(&self.state.cfg).as_nanos();
7535 let mut expected = epoch;
7536 for (index, owner) in self.owners.iter().enumerate() {
7537 let clock = if completed.contains(&index) {
7538 &owner.thread_logical_time
7539 } else {
7540 &self.initial_clocks[index]
7541 };
7542 assert_eq!(time.threads_time(owner.dettid), clock.as_nanos());
7543 assert_eq!(
7544 snapshot["inherited_time"][owner.dettid.as_raw().to_string()],
7545 serde_json::to_value(clock.inherited_nanos()).unwrap()
7546 );
7547 expected = expected + (clock.as_nanos() - epoch - clock.inherited_nanos());
7548 }
7549 assert_eq!(
7550 time.as_nanos(),
7551 expected,
7552 "only each owner's own uninherited work contributes"
7553 );
7554 snapshot
7555 }
7556
7557 async fn exit(&self, index: usize, status: ExitStatus) {
7558 self.exit_thread(self.owners[index].clone(), status).await;
7559 }
7560
7561 async fn exit_thread(&self, thread: ThreadState<()>, status: ExitStatus) {
7562 let rpc = ExitRpc {
7563 state: &self.state,
7564 sender: thread.dettid,
7565 observations: &self.observations,
7566 rpc_observations: &self.rpc_observations,
7567 };
7568 let exit =
7569 self.tool
7570 .on_exit_thread(Tid::from_raw(rpc.sender.as_raw()), &rpc, thread, status);
7571 let mut exit = std::pin::pin!(exit);
7572 assert!(
7577 matches!(futures::poll!(exit.as_mut()), Poll::Ready(Ok(()))),
7578 "exit callback yielded before its accounting and cleanup completed"
7579 );
7580 }
7581
7582 async fn nonmember_exit(&self, raw_tid: i32) {
7583 let mut thread = self.owners[0].clone();
7584 thread.dettid = DetTid::from_raw(raw_tid);
7585 thread.thread_logical_time = DetTime::new(&self.state.cfg);
7586 let tid = thread.dettid;
7587 {
7588 let mut sched = self.state.sched.lock().unwrap();
7589 sched
7590 .thread_tree
7591 .add_child(self.owners[0].dettid, tid, false);
7592 sched.priorities.insert(tid, DEFAULT_PRIORITY);
7593 sched.next_turns.insert(
7594 tid,
7595 ThreadNextTurn {
7596 dettid: tid,
7597 child_tid_addr: 0,
7598 req: Ivar::new(),
7599 resp: Ivar::new(),
7600 protocol: Default::default(),
7601 },
7602 );
7603 sched.install_test_exec_incarnation(tid, thread.mm_id);
7604 sched.runqueue_push_back(tid);
7605 }
7606 let before = self.rpc_observations.lock().unwrap().len();
7607 self.exit_thread(thread, ExitStatus::Exited(0)).await;
7608 let observations = self.rpc_observations.lock().unwrap();
7609 assert_eq!(
7610 observations.len(),
7611 before + 1,
7612 "a nonmember sent an acknowledgement or wake RPC"
7613 );
7614 assert_eq!(observations[before].sender, tid);
7615 assert_eq!(observations[before].kind, "deregister");
7616 assert!(observations[before].accepted);
7617 }
7618 }
7619
7620 #[tokio::test]
7621 async fn robust_exit_acknowledgements_preserve_global_action_order_and_eligibility() {
7622 for order in [[0, 1], [1, 0]] {
7623 for empty_owner in [0, 1] {
7624 let f = Fixture::new(RobustListExit::ExitGroup, false, Some(empty_owner), false);
7625 f.nonmember_exit(30).await;
7626 assert!(f.observations.lock().unwrap().is_empty());
7627 let queue_before_first: Vec<_> = f
7628 .state
7629 .sched
7630 .lock()
7631 .unwrap()
7632 .run_queue
7633 .tids()
7634 .copied()
7635 .collect();
7636 let first_start = f.rpc_observations.lock().unwrap().len();
7637 f.exit(order[0], ExitStatus::Exited(0)).await;
7638 let first_clocks = f.assert_clocks(&[order[0]]);
7639 {
7640 let observations = f.rpc_observations.lock().unwrap();
7641 let actions = &observations[first_start..];
7642 assert_eq!(
7643 actions.iter().map(|a| a.kind).collect::<Vec<_>>(),
7644 ["empty-wake", "deregister"]
7645 );
7646 assert!(actions.iter().all(|a| a.accepted
7647 && a.sender == f.owners[order[0]].dettid
7648 && a.clocks == first_clocks
7649 && a.waiters == 2
7650 && a.turn == 0));
7651 assert_eq!(
7652 actions[0].queued, queue_before_first,
7653 "empty acknowledgement changed scheduler eligibility"
7654 );
7655 }
7656 f.nonmember_exit(31).await;
7657 f.assert_clocks(&[order[0]]);
7658 assert!(
7659 f.observations.lock().unwrap().is_empty(),
7660 "a nonmember completed the physical-exit barrier"
7661 );
7662 let queue_before_last: Vec<_> = f
7663 .state
7664 .sched
7665 .lock()
7666 .unwrap()
7667 .run_queue
7668 .tids()
7669 .copied()
7670 .collect();
7671 let last_start = f.rpc_observations.lock().unwrap().len();
7672 f.exit(order[1], ExitStatus::Exited(0)).await;
7673 let final_clocks = f.assert_clocks(&[0, 1]);
7674 {
7675 let observations = f.rpc_observations.lock().unwrap();
7676 let actions = &observations[last_start..];
7677 assert_eq!(
7678 actions.iter().map(|a| a.kind).collect::<Vec<_>>(),
7679 ["empty-wake", "wake", "deregister"]
7680 );
7681 assert!(actions.iter().all(|a| a.accepted
7682 && a.sender == f.owners[order[1]].dettid
7683 && a.clocks == final_clocks
7684 && a.turn == 0));
7685 assert_eq!(actions[0].waiters, 2);
7686 assert_eq!(actions[1].waiters, 1);
7687 assert_eq!(actions[0].queued, queue_before_last);
7688 assert_eq!(
7689 actions[1].queued, queue_before_last,
7690 "wake bypassed deferred admission"
7691 );
7692 }
7693 f.nonmember_exit(32).await;
7694 f.assert_clocks(&[0, 1]);
7695 assert_eq!(f.observations.lock().unwrap().len(), 1);
7696 let observations = f.rpc_observations.lock().unwrap();
7697 for owner in &f.owners {
7698 assert_eq!(
7699 observations
7700 .iter()
7701 .filter(|a| a.sender == owner.dettid && a.kind == "empty-wake")
7702 .count(),
7703 1,
7704 "each unique matching owner, including an empty-wake owner, must acknowledge before the batch clears"
7705 );
7706 }
7707 }
7708 }
7709 }
7710
7711 #[tokio::test]
7712 async fn robust_exit_callbacks_preserve_owner_clocks_across_arrival_orders() {
7713 for reason in [
7714 RobustListExit::ExitGroup,
7715 RobustListExit::Signal(libc::SIGTERM),
7716 ] {
7717 for equal_clocks in [false, true] {
7718 for empty_owner in [None, Some(0), Some(1)] {
7719 let mut results = Vec::new();
7720 for order in [[0, 1], [1, 0]] {
7721 let f = Fixture::new(reason, equal_clocks, empty_owner, false);
7722 let status = match reason {
7723 RobustListExit::ExitGroup => ExitStatus::Exited(0),
7724 RobustListExit::Signal(_) => {
7725 ExitStatus::Signaled(Signal::SIGTERM, false)
7726 }
7727 };
7728 f.exit(order[0], status).await;
7729 f.assert_clocks(&[order[0]]);
7730 assert!(f.observations.lock().unwrap().is_empty());
7731 {
7732 let sched = f.state.sched.lock().unwrap();
7733 assert_eq!(
7734 sched
7735 .blocked
7736 .futex_waiters
7737 .values()
7738 .map(Vec::len)
7739 .sum::<usize>(),
7740 2
7741 );
7742 assert_eq!(sched.turn, 0);
7743 }
7744 {
7747 let last = Err(SkipTurn);
7748 let turn = crate::scheduler::do_a_turn_blocking(
7749 f.state.sched.clone(),
7750 f.state.global_time.clone(),
7751 &last,
7752 );
7753 let mut turn = std::pin::pin!(turn);
7754 assert!(matches!(futures::poll!(turn.as_mut()), Poll::Pending));
7755 }
7756 f.exit(order[0], status).await;
7757 f.assert_clocks(&[order[0]]);
7758 assert!(
7759 f.observations.lock().unwrap().is_empty(),
7760 "duplicate physical exit released an incomplete group"
7761 );
7762 f.exit(order[1], status).await;
7763 let clocks = f.assert_clocks(&[0, 1]);
7764 {
7765 let observations = f.observations.lock().unwrap();
7766 assert_eq!(observations.len(), 1);
7767 let expected: Vec<_> = (0..2)
7768 .filter(|i| Some(*i) != empty_owner)
7769 .map(|i| (f.owners[i].dettid, f.futexes[i]))
7770 .collect();
7771 assert_eq!(observations[0].wakes, expected);
7772 assert_eq!(observations[0].counts, vec![1; expected.len()]);
7773 assert_eq!(
7774 observations[0].clocks, clocks,
7775 "all owner clocks must be accounted before wake admission"
7776 );
7777 assert_eq!(observations[0].turn, 0);
7778 }
7779 f.exit(order[1], status).await;
7780 assert_eq!(f.assert_clocks(&[0, 1]), clocks);
7781 assert_eq!(
7782 f.observations.lock().unwrap().len(),
7783 1,
7784 "repeated cleanup emitted a second batch"
7785 );
7786 {
7787 let sched = f.state.sched.lock().unwrap();
7788 assert!(!sched.run_queue.contains_tid(f.waiters[0]));
7789 assert!(!sched.run_queue.contains_tid(f.waiters[1]));
7790 }
7791 let result = crate::scheduler::do_a_turn_blocking(
7794 f.state.sched.clone(),
7795 f.state.global_time.clone(),
7796 &Err(SkipTurn),
7797 )
7798 .await;
7799 assert!(result.is_ok());
7800 let sched = f.state.sched.lock().unwrap();
7801 assert_eq!(sched.turn, 1);
7802 let queued: Vec<_> = sched.run_queue.tids().copied().collect();
7803 for (i, waiter) in f.waiters.iter().enumerate() {
7804 assert_eq!(queued.contains(waiter), Some(i) != empty_owner);
7805 }
7806 assert!(queued.contains(&f.peer));
7807 results.push((
7808 clocks,
7809 queued,
7810 std::mem::take(&mut *f.observations.lock().unwrap()),
7811 ));
7812 }
7813 assert_eq!(
7814 results[0], results[1],
7815 "callback order changed final clocks, typed wake observations or the actual scheduler drain"
7816 );
7817 }
7818 }
7819 }
7820 }
7821
7822 #[tokio::test]
7823 async fn robust_exit_callbacks_keep_rejected_or_mismatched_batches_incomplete() {
7824 for rejection in ["old-mm", "tombstone", "wrong-signal", "normal-exit"] {
7825 let f = Fixture::new(
7826 RobustListExit::Signal(libc::SIGTERM),
7827 false,
7828 None,
7829 rejection == "tombstone",
7830 );
7831 let rejected = f.owners[0].dettid;
7832 let mm = f.owners[0].mm_id;
7833 f.exit(1, ExitStatus::Signaled(Signal::SIGTERM, false))
7834 .await;
7835 let replacement_request = Ivar::new();
7836 if rejection == "old-mm" {
7837 let mut sched = f.state.sched.lock().unwrap();
7838 sched.install_test_exec_incarnation(rejected, mm.for_exec(rejected));
7839 sched.next_turns.insert(
7840 rejected,
7841 ThreadNextTurn {
7842 dettid: rejected,
7843 child_tid_addr: 0,
7844 req: replacement_request.clone(),
7845 resp: Ivar::new(),
7846 protocol: Default::default(),
7847 },
7848 );
7849 let mut replacement_time = f.owners[0].thread_logical_time.clone();
7850 replacement_time.add_syscall_with_cost(500);
7851 f.state.global_time.lock().unwrap().update_global_time(
7852 rejected,
7853 replacement_time.as_nanos(),
7854 replacement_time.inherited_nanos(),
7855 );
7856 } else if rejection == "tombstone" {
7857 let mut sched = f.state.sched.lock().unwrap();
7859 sched.logically_kill_thread(&rejected, &rejected, mm);
7860 }
7861 let before = serde_json::to_value(&*f.state.global_time.lock().unwrap()).unwrap();
7862 let status = match rejection {
7863 "wrong-signal" => ExitStatus::Signaled(Signal::SIGKILL, false),
7864 "normal-exit" => ExitStatus::Exited(0),
7865 _ => ExitStatus::Signaled(Signal::SIGTERM, false),
7866 };
7867 f.exit(0, status).await;
7868 assert!(
7869 f.observations.lock().unwrap().is_empty(),
7870 "{rejection} released a group"
7871 );
7872 if rejection == "old-mm" || rejection == "tombstone" {
7873 assert_eq!(
7874 serde_json::to_value(&*f.state.global_time.lock().unwrap()).unwrap(),
7875 before,
7876 "rejected acknowledgement changed clock state"
7877 );
7878 }
7879 let sched = f.state.sched.lock().unwrap();
7880 assert_eq!(
7881 sched
7882 .blocked
7883 .futex_waiters
7884 .values()
7885 .map(Vec::len)
7886 .sum::<usize>(),
7887 2,
7888 "rejected group lost a real waiter"
7889 );
7890 assert!(
7891 f.waiters
7892 .iter()
7893 .all(|tid| !sched.run_queue.contains_tid(*tid))
7894 );
7895 if rejection == "old-mm" {
7896 assert!(sched.rpc_incarnation_matches(rejected, mm.for_exec(rejected)));
7897 assert_eq!(
7898 sched.next_turns[&rejected].req, replacement_request,
7899 "old cleanup destroyed replacement registration"
7900 );
7901 }
7902 }
7903 }
7904
7905 #[tokio::test]
7906 #[should_panic(expected = "Attempted to update tid 17 time")]
7907 async fn robust_exit_clock_ack_still_refuses_a_backwards_owner_sample() {
7908 let f = Fixture::new(RobustListExit::ExitGroup, false, None, false);
7909 let owner = &f.owners[0];
7910 let mut later = owner.thread_logical_time.clone();
7911 later.add_syscall_with_cost(1);
7912 f.state.global_time.lock().unwrap().update_global_time(
7913 owner.dettid,
7914 later.as_nanos(),
7915 later.inherited_nanos(),
7916 );
7917 f.exit(0, ExitStatus::Exited(0)).await;
7918 }
7919
7920 #[tokio::test]
7921 async fn backend_failure_preserves_consuming_robust_exit_clock_accounting() {
7922 for order in [[0, 1], [1, 0]] {
7923 let f = Fixture::new(RobustListExit::ExitGroup, false, None, false);
7924 let selected = {
7925 let mut sched = f.state.sched.lock().unwrap();
7926 sched.next_turns.get_mut(&f.peer).unwrap().req = Ivar::new();
7927 sched.select_test_turn().unwrap()
7928 };
7929 let responses = {
7930 let sched = f.state.sched.lock().unwrap();
7931 f.waiters.map(|tid| sched.next_turns[&tid].resp.clone())
7932 };
7933 let mut turn = std::pin::pin!(crate::scheduler::finish_selected_turn(
7934 f.state.sched.clone(),
7935 f.state.global_time.clone(),
7936 selected.0,
7937 selected.1,
7938 selected.2,
7939 ));
7940 assert!(futures::poll!(turn.as_mut()).is_pending());
7941 f.state.report_backend_failure(reverie::BackendFailure {
7942 pid: Tid::from_raw(17),
7943 tid: Tid::from_raw(18),
7944 phase: "native robust cleanup control",
7945 });
7946 for index in order {
7947 f.exit(index, ExitStatus::Exited(0)).await;
7948 }
7949 assert!(matches!(futures::poll!(turn.as_mut()), Poll::Ready(Err(_))));
7950 f.assert_clocks(&[0, 1]);
7951 let observations = f.observations.lock().unwrap();
7952 assert_eq!(observations.len(), 1, "one complete batch");
7953 assert_eq!(observations[0].counts, vec![1, 1]);
7954 assert!(
7955 responses
7956 .iter()
7957 .all(|response| response.try_read().is_none())
7958 );
7959 let mut sched = f.state.sched.lock().unwrap();
7960 assert_eq!(sched.turn, 0);
7961 for owner in &f.owners {
7962 assert!(!sched.next_turns.contains_key(&owner.dettid));
7963 assert!(!sched.note_deregistration_accounted(owner.dettid));
7964 }
7965 }
7966 }
7967}