Skip to main content

detcore/
preemptions.rs

1/*
2 * Copyright (c) Meta Platforms, Inc. and affiliates.
3 * All rights reserved.
4 *
5 * This source code is licensed under the BSD-style license found in the
6 * LICENSE file in the root directory of this source tree.
7 */
8
9//! A datatype to abstract a record of thread preemptions, as generated during chaos mode execution.
10
11use std::collections::BTreeMap;
12use std::fs::File;
13use std::iter::FromIterator;
14use std::path::Path;
15use std::path::PathBuf;
16
17use chrono::DateTime;
18use chrono::Utc;
19use detcore_model::collections::ReplayCursor;
20use serde::Deserialize;
21use serde::Serialize;
22use tracing::trace;
23
24use crate::resources::ChaosEpochTransition;
25use crate::scheduler::Priority;
26use crate::scheduler::runqueue::DEFAULT_PRIORITY;
27use crate::scheduler::runqueue::FIRST_PRIORITY;
28use crate::scheduler::runqueue::LAST_PRIORITY;
29use crate::scheduler::runqueue::is_ordinary_priority;
30use crate::types::DetTid;
31use crate::types::LogicalTime;
32use crate::types::SchedEvent;
33
34/// A record of all the preemptions and other scheduling events that occur during execution.
35#[derive(PartialEq, Default, Debug, Eq, Clone, Hash, Serialize, Deserialize)]
36pub struct PreemptionRecord {
37    // TODO(T106838933): switch the keys from DetTid to Pedigree.
38    /// A sorted list of end-of-timeslice preemption times for each thread.
39    per_thread: BTreeMap<DetTid, ThreadHistory>,
40    global: Vec<SchedEvent>,
41    /// The virtual-time epoch of the run that produced this record.
42    ///
43    /// Every `LogicalTime` above is an absolute virtual instant, i.e. an offset
44    /// from this epoch, so the record only describes a run whose clock starts
45    /// here. A replay started from any other epoch would see every recorded
46    /// timeslice end shifted by the difference. `None` for records written
47    /// before the epoch was stored, and for records synthesized in memory.
48    #[serde(default, skip_serializing_if = "Option::is_none")]
49    epoch: Option<DateTime<Utc>>,
50}
51
52impl std::fmt::Display for PreemptionRecord {
53    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
54        let str = serde_json::to_string(&self).unwrap();
55        write!(f, "{}", str)
56    }
57}
58
59impl PreemptionRecord {
60    /// TODO: This is not ideal because PreemptionRecord shouldn't be doing double duty as an eventlog.
61    pub fn from_sched_events(events: Vec<SchedEvent>) -> Self {
62        Self {
63            per_thread: Default::default(),
64            global: events,
65            epoch: None,
66        }
67    }
68
69    /// The virtual-time epoch this record was produced under, if it was stored.
70    pub fn epoch(&self) -> Option<DateTime<Utc>> {
71        self.epoch
72    }
73
74    /// Make a copy of everything in the `PreemptionRecord` in a public format.
75    pub fn extract_all(&self) -> BTreeMap<DetTid, ThreadHistory> {
76        self.per_thread.clone()
77    }
78
79    /// Inverse of from_vecs.
80    pub fn as_vecs(&self) -> BTreeMap<DetTid, Vec<(LogicalTime, Priority)>> {
81        let mut bt = BTreeMap::new();
82        for (tid, th) in &self.per_thread {
83            bt.insert(*tid, th.as_vec());
84        }
85        bt
86    }
87    /// TODO: provide a better way to create an empty PreemptionRecord but with entries for a given set of tids.
88    pub fn strip_contents(mut self) -> Self {
89        for history in self.per_thread.values_mut() {
90            history.prio_changes = Vec::new();
91            history.preemption_rcbs = Vec::new();
92            history.chaos_epochs = Vec::new();
93            history.final_prio = 1000;
94        }
95        self.global = Vec::new();
96        self
97    }
98
99    /// Iterates over each sched event and applies function event -> *event yielding resulting iterator as new events
100    /// This effectively moves content of global allowing to split or delete events based on the function provided
101    pub fn split_map<F, R>(&mut self, splitter: F)
102    where
103        R: IntoIterator<Item = SchedEvent>,
104        F: Fn(SchedEvent, &ReplayCursor<SchedEvent>) -> R,
105    {
106        let mut result = Vec::new();
107        let global = std::mem::take(&mut self.global);
108        let mut cursor = ReplayCursor::from_iter(global);
109
110        while let Some(event) = cursor.next() {
111            for new_event in splitter(event, &cursor) {
112                result.push(new_event);
113            }
114        }
115        self.global = result;
116    }
117
118    /// Iterate over the schedevents.
119    pub fn schedevents_iter_mut(&mut self) -> std::slice::IterMut<'_, SchedEvent> {
120        self.global.iter_mut()
121    }
122
123    /// Leave only the per-thread preemption records, not the global schedule.
124    pub fn preemptions_only(&mut self) {
125        self.global.clear();
126    }
127
128    /// Make a copy just of the preemptions (cheap) rather than the whole schedule trace (expensive)
129    pub fn clone_preemptions_only(&self) -> Self {
130        PreemptionRecord {
131            per_thread: self.per_thread.clone(),
132            global: Vec::new(),
133            epoch: self.epoch,
134        }
135    }
136
137    /// Read-only access to the raw event schedule.
138    pub fn schedevents(&self) -> &Vec<SchedEvent> {
139        &self.global
140    }
141
142    /// Returns true if the schedule trace (SchedEvent) record is nonempty.
143    pub fn contains_schedevents(&self) -> bool {
144        !self.global.is_empty()
145    }
146
147    /// Convert from a flat vector representation (Time,Priority_AFTER_Time) into the internal representation.
148    /// The very first timeslice should have a zero timestamp, but it is ignored if nonzero.
149    pub fn from_vecs(bt: &BTreeMap<DetTid, Vec<(LogicalTime, Priority)>>) -> Self {
150        let mut bt2 = BTreeMap::new();
151        for (tid, vec) in bt {
152            let th = if vec.is_empty() {
153                ThreadHistory {
154                    final_prio: DEFAULT_PRIORITY,
155                    prio_changes: Vec::new(),
156                    preemption_rcbs: Vec::new(),
157                    chaos_epochs: Vec::new(),
158                }
159            } else {
160                // We could insist on this invariant, but it's a little more flexible not to:
161                // assert!(vec.get(0).unwrap().per_thread.is_zero());
162                let (_final_end, final_prio) = vec.last().unwrap();
163                let final_prio = *final_prio;
164
165                let mut prio_changes = Vec::new();
166                let mut it = vec.iter().peekable();
167                while let Some((_this_ns, this_p)) = it.next() {
168                    if let Some((next_ns, _next_p)) = it.peek() {
169                        prio_changes.push((*next_ns, *this_p));
170                    } else {
171                        // Don't need the last one, because we've already used it's
172                        // timestamp, and retrieved final_prio.
173                        break;
174                    }
175                }
176                assert!(final_prio >= FIRST_PRIORITY);
177                assert!(final_prio <= LAST_PRIORITY);
178                ThreadHistory {
179                    final_prio,
180                    prio_changes,
181                    preemption_rcbs: Vec::new(),
182                    chaos_epochs: Vec::new(),
183                }
184            };
185            bt2.insert(*tid, th);
186        }
187        PreemptionRecord {
188            per_thread: bt2,
189            global: Vec::new(),
190            epoch: None,
191        }
192    }
193
194    /// Extracts the inner global events.
195    pub fn into_global(self) -> Vec<SchedEvent> {
196        self.global
197    }
198
199    /// Save to disk.
200    pub fn write_to_disk(&self, path: &Path) -> Result<(), String> {
201        let mut str: String = self.to_string();
202        str.push('\n');
203        match File::create(path) {
204            Ok(mut file) => match std::io::Write::write_all(&mut file, str.as_bytes()) {
205                Ok(_) => Ok(()),
206                Err(err) => Err(format!(
207                    "Failed to write preemption record to file {:?}, error: {}",
208                    path, err
209                )),
210            },
211            Err(err) => Err(format!(
212                "Failed to create file for preemption record {:?}, error: {}",
213                path, err
214            )),
215        }
216    }
217
218    /// Perform internal invariant checks on the PreemptionRecord and return an error if
219    /// it is not well formed.
220    pub fn validate(&self) -> Result<(), String> {
221        for (tid, history) in &self.per_thread {
222            if !is_ordinary_priority(history.final_prio) {
223                return Err(format!(
224                    "final priority for thread {} invalid: {}",
225                    tid, history.final_prio
226                ));
227            }
228            {
229                let mut time_last = None;
230                for (count, (ns, prio)) in history.prio_changes.iter().enumerate() {
231                    if let Some(last) = time_last {
232                        if !is_ordinary_priority(*prio) {
233                            return Err(format!(
234                                "preemption priority #{} for thread {} invalid: {}",
235                                count, tid, history.final_prio
236                            ));
237                        }
238                        if *ns <= last {
239                            return Err(format!(
240                                "Timestamps failed to monotonically increase ({}), in series:\n {:?}",
241                                ns, history.prio_changes
242                            ));
243                        }
244                    }
245                    time_last = Some(*ns);
246                }
247            }
248            if !history.preemption_rcbs.is_empty()
249                && history.preemption_rcbs.len() != history.prio_changes.len()
250            {
251                return Err(format!(
252                    "thread {} has {} preemption times but {} RCB targets",
253                    tid,
254                    history.prio_changes.len(),
255                    history.preemption_rcbs.len()
256                ));
257            }
258            if history
259                .preemption_rcbs
260                .windows(2)
261                .any(|pair| pair[1] < pair[0])
262            {
263                return Err(format!(
264                    "preemption RCB targets failed to increase for thread {}",
265                    tid
266                ));
267            }
268            let mut transition_last = None;
269            let mut epoch_last = None;
270            for transition in &history.chaos_epochs {
271                if transition.factor.as_f64() <= 0.0 {
272                    return Err(format!(
273                        "chaos epoch factor for thread {} must be positive",
274                        tid
275                    ));
276                }
277                if transition_last.is_some_and(|last| transition.logical_time <= last) {
278                    return Err(format!(
279                        "chaos epoch transition times failed to increase for thread {}",
280                        tid
281                    ));
282                }
283                if epoch_last.is_some_and(|last| transition.epoch <= last) {
284                    return Err(format!(
285                        "chaos epoch numbers failed to increase for thread {}",
286                        tid
287                    ));
288                }
289                transition_last = Some(transition.logical_time);
290                epoch_last = Some(transition.epoch);
291            }
292        }
293        Ok(())
294    }
295
296    /// Normalize priorities for readibility.
297    pub fn normalize(&self) -> PreemptionRecord {
298        let mut clone = self.clone();
299
300        let mut priomap: BTreeMap<Priority, Priority> = BTreeMap::new();
301        for history in clone.per_thread.values_mut() {
302            let _ = priomap.insert(history.final_prio, 0);
303            for (_ns, prio) in &history.prio_changes {
304                let _ = priomap.insert(*prio, 0);
305            }
306        }
307
308        for (cur_prio, val) in (DEFAULT_PRIORITY..).zip(priomap.values_mut()) {
309            assert!(cur_prio <= LAST_PRIORITY);
310            *val = cur_prio;
311        }
312        for history in clone.per_thread.values_mut() {
313            history.final_prio = *priomap.get(&history.final_prio).unwrap();
314            for (_ns, prio) in &mut history.prio_changes {
315                *prio = *priomap.get(prio).unwrap();
316            }
317        }
318
319        // One more pass to remove anything that just has default priority its whole lifetime.
320        let mut finalmap = BTreeMap::new();
321        for (tid, history) in clone.per_thread.into_iter() {
322            if history.final_prio != DEFAULT_PRIORITY
323                || !history.prio_changes.is_empty()
324                || !history.chaos_epochs.is_empty()
325            {
326                assert!(finalmap.insert(tid, history).is_none());
327            }
328        }
329        clone.per_thread = finalmap;
330        clone
331    }
332
333    /// Remove whichever is the latest preemption across all thread histories. If
334    /// there are no priority changes in any thread history, the entire last thread
335    /// (highest thread number) is completely removed.
336    /// Note this function assumes the last priority change in the list has the latest
337    /// preempt time.
338    pub fn with_latest_preempt_removed(&self) -> PreemptionRecord {
339        let mut preempts_latest_prio_changes: Vec<(DetTid, LogicalTime)> = self
340            .per_thread
341            .clone()
342            .into_iter()
343            .map(|(tid, th)| {
344                (
345                    tid,
346                    th.prio_changes
347                        .last() // Assume last prio change in the list has the latest preempt time
348                        .map_or(LogicalTime::ZERO, |prio_change| prio_change.0),
349                )
350            })
351            .collect();
352        preempts_latest_prio_changes.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap());
353
354        let tid_with_preempt_to_drop = preempts_latest_prio_changes
355            .into_iter()
356            .map(|(tid, _th)| tid)
357            .next();
358
359        let mut clone = self.clone();
360        if let Some(tid) = tid_with_preempt_to_drop {
361            let thread_history = &mut clone.per_thread.get_mut(&tid).unwrap();
362            if !thread_history.prio_changes.is_empty() {
363                // There were prio changes for at least one of the threads, remove that
364                // preemption as the one to drop
365                let removed_prio = thread_history.prio_changes.pop().unwrap().1;
366                if !thread_history.preemption_rcbs.is_empty() {
367                    thread_history.preemption_rcbs.pop();
368                }
369                thread_history.final_prio = removed_prio;
370            } else {
371                // There were no prio changes for any threads, so just remove the last
372                // thread completely as the drop
373                clone.per_thread.pop_last();
374            }
375        }
376        clone
377    }
378}
379
380/// The record of priorities and preemptions for a particular thread.
381#[derive(
382    PartialEq, // Silly protection from rustfmt disagreements.
383    Debug,
384    Eq,
385    Clone,
386    Hash,
387    Serialize,
388    Deserialize,
389)]
390pub struct ThreadHistory {
391    /// The priority of the last timeslice, after the last preemption.
392    /// (This is also the initial priority if there are no preemptions.)
393    pub final_prio: Priority,
394
395    /// The preemption points and priority just BEFORE the preemption occurred.  It's
396    /// private, and thus you must go through the iterator, as this will change in the
397    /// future.
398    prio_changes: Vec<(LogicalTime, Priority)>,
399
400    // AUTONOMOUS-BOT-IMPLEMENTED
401    // TODO-HUMAN-REVIEW(PR-1151)
402    /// Exact per-thread PMU RCB target for each matching priority change. Empty
403    /// in older JSON artifacts, which fall back to logical-time replay.
404    // `ThreadHistory` crosses a non-self-describing bincode RPC; never skip it.
405    #[serde(default)]
406    preemption_rcbs: Vec<u64>,
407
408    // AUTONOMOUS-BOT-IMPLEMENTED
409    // TODO-HUMAN-REVIEW(PR-1151)
410    /// Exact per-thread slowdown transitions. Empty artifacts from older Hermit
411    /// versions remain wire-compatible when read through Serde's default.
412    // `ThreadHistory` also crosses Reverie's non-self-describing bincode RPC,
413    // so this field must never be skipped when it is empty.
414    #[serde(default)]
415    chaos_epochs: Vec<ChaosEpochTransition>,
416}
417
418impl ThreadHistory {
419    /// An empty thread history, with default priority only.
420    pub fn new() -> Self {
421        ThreadHistory {
422            final_prio: DEFAULT_PRIORITY,
423            prio_changes: Vec::new(),
424            preemption_rcbs: Vec::new(),
425            chaos_epochs: Vec::new(),
426        }
427    }
428
429    /// Turns the thread history into an iterator that produces a stream of preemption events.
430    #[allow(clippy::should_implement_trait)]
431    pub fn into_iter(self) -> ThreadHistoryIterator {
432        ThreadHistoryIterator {
433            full_history: self,
434            ix: 0,
435            chaos_epoch_ix: 0,
436        }
437    }
438
439    #[cfg(test)]
440    pub(crate) fn with_chaos_epochs(mut self, chaos_epochs: Vec<ChaosEpochTransition>) -> Self {
441        self.chaos_epochs = chaos_epochs;
442        self
443    }
444
445    #[cfg(test)]
446    pub(crate) fn with_prio_changes(mut self, prio_changes: Vec<(LogicalTime, Priority)>) -> Self {
447        self.prio_changes = prio_changes;
448        self
449    }
450
451    #[cfg(test)]
452    pub(crate) fn with_preemption_rcbs(mut self, preemption_rcbs: Vec<u64>) -> Self {
453        self.preemption_rcbs = preemption_rcbs;
454        self
455    }
456
457    /// Flatten the time series into a uniform representation.
458    /// This is a `(Time, Priority_AFTER_Time)` representation.
459    /// The very first timeslice will have a zero timestamp.
460    // Note that this is different than the internal representation which groups together
461    // `Priority_UNTIL_Time, Time`, i.e. shifted by one.
462    pub fn as_vec(&self) -> Vec<(LogicalTime, Priority)> {
463        let mut vec = Vec::new();
464        let mut time0 = LogicalTime::from_nanos(0);
465
466        for (time1, prio) in &self.prio_changes {
467            vec.push((time0, *prio));
468            time0 = *time1;
469        }
470        vec.push((time0, self.final_prio));
471        vec
472    }
473
474    /// Return the priority of the initial timeslice after the thread starts running.
475    pub fn initial_priority(&self) -> Priority {
476        if let Some((_ns, pr)) = self.prio_changes.first() {
477            *pr
478        } else {
479            self.final_prio
480        }
481    }
482}
483
484impl Default for ThreadHistory {
485    fn default() -> Self {
486        Self::new()
487    }
488}
489
490/// An iterator over the preemption history of a thread.  Does not include the final
491/// timeslice, which does not have an ending time.
492// TODO: in the future this will perform "lazy IO" and pull from disk in batches.
493#[derive(PartialEq, Eq, Clone, Serialize, Deserialize)]
494pub struct ThreadHistoryIterator {
495    /// For now we store the full history in memory:
496    full_history: ThreadHistory,
497    /// Index that tracks our position.
498    ix: usize,
499    // AUTONOMOUS-BOT-IMPLEMENTED
500    // TODO-HUMAN-REVIEW(PR-1151)
501    #[serde(default)]
502    chaos_epoch_ix: usize,
503}
504
505impl ThreadHistoryIterator {
506    /// Return the priority of the initial slice after the thread starts running.
507    pub fn initial_priority(&self) -> Priority {
508        self.full_history.initial_priority()
509    }
510
511    /// Return the priority of the last slice after the last preemption.
512    pub fn final_priority(&self) -> Priority {
513        self.full_history.final_prio
514    }
515
516    // AUTONOMOUS-BOT-IMPLEMENTED
517    // TODO-HUMAN-REVIEW(PR-1151)
518    /// Consume recorded epoch transitions up to `current_time`, returning a
519    /// transition only when the active recorded factor changed.
520    pub fn advance_chaos_epoch(
521        &mut self,
522        current_time: LogicalTime,
523    ) -> Option<ChaosEpochTransition> {
524        let mut changed = None;
525        while let Some(transition) = self.full_history.chaos_epochs.get(self.chaos_epoch_ix)
526            && transition.logical_time <= current_time
527        {
528            changed = Some(*transition);
529            self.chaos_epoch_ix += 1;
530        }
531        changed
532    }
533
534    /// Whether this artifact carries recorded slowdown epochs.
535    pub fn has_chaos_epochs(&self) -> bool {
536        !self.full_history.chaos_epochs.is_empty()
537    }
538
539    // AUTONOMOUS-BOT-IMPLEMENTED
540    // TODO-HUMAN-REVIEW(PR-1151)
541    /// Return the next recorded logical boundary and its exact RCB target when
542    /// the artifact was produced by a version that records one.
543    pub fn next_with_rcbs(&mut self) -> Option<(LogicalTime, Priority, Option<u64>)> {
544        let ix = self.ix;
545        self.next().map(|(time, priority)| {
546            (
547                time,
548                priority,
549                self.full_history.preemption_rcbs.get(ix).copied(),
550            )
551        })
552    }
553}
554
555impl Iterator for ThreadHistoryIterator {
556    type Item = (LogicalTime, Priority);
557
558    fn next(&mut self) -> Option<Self::Item> {
559        let vec = &self.full_history.prio_changes;
560        if vec.len() > self.ix {
561            let elt = vec[self.ix];
562            self.ix += 1;
563            Some(elt)
564        } else {
565            None
566        }
567    }
568}
569
570#[cfg(test)]
571mod tests {
572    use detcore_model::schedule::Op;
573    use pretty_assertions::assert_eq;
574    use test_case::test_case;
575
576    use super::*;
577    use crate::types::RcbTimeMultiplier;
578
579    // AUTONOMOUS-BOT-IMPLEMENTED
580    // TODO-HUMAN-REVIEW(PR-1151)
581    #[test]
582    fn chaos_epoch_transitions_round_trip_and_replay_exact_factors() {
583        let tid = DetTid::from_raw(2);
584        let first = ChaosEpochTransition {
585            logical_time: LogicalTime::from_nanos(100),
586            epoch: 0,
587            factor: RcbTimeMultiplier::from_f64(2.5),
588        };
589        let second = ChaosEpochTransition {
590            logical_time: LogicalTime::from_nanos(500),
591            epoch: 1,
592            factor: RcbTimeMultiplier::from_f64(0.75),
593        };
594
595        let mut writer = PreemptionWriter::new(None);
596        writer.register_thread(tid, DEFAULT_PRIORITY);
597        writer.insert_chaos_epoch(tid, first);
598        writer.insert_chaos_epoch(tid, second);
599        let encoded = writer.into_string();
600        assert!(encoded.contains("chaos_epochs"));
601
602        let decoded: PreemptionRecord = serde_json::from_str(&encoded).unwrap();
603        decoded.validate().unwrap();
604        let mut history = decoded.extract_all().remove(&tid).unwrap().into_iter();
605        assert_eq!(
606            history.advance_chaos_epoch(LogicalTime::from_nanos(99)),
607            None
608        );
609        assert_eq!(
610            history.advance_chaos_epoch(LogicalTime::from_nanos(100)),
611            Some(first)
612        );
613        assert_eq!(
614            history.advance_chaos_epoch(LogicalTime::from_nanos(499)),
615            None
616        );
617        assert_eq!(
618            history.advance_chaos_epoch(LogicalTime::from_nanos(500)),
619            Some(second)
620        );
621        assert_eq!(history.advance_chaos_epoch(LogicalTime::MAX), None);
622    }
623
624    /// <https://github.com/rrnewton/hermit/issues/3411>: the epoch a record was
625    /// made under survives serialization, the preemptions-only copy, and the
626    /// non-panicking CLI reader; records written before it was stored read as
627    /// `None`, and a malformed file is an error rather than a panic.
628    #[test]
629    fn recorded_epoch_round_trips_and_legacy_records_have_none() {
630        let epoch: DateTime<Utc> = "2000-12-31T23:59:59.123456789Z".parse().unwrap();
631        let tid = DetTid::from_raw(3);
632        let mut writer = PreemptionWriter::new(None).with_epoch(epoch);
633        writer.register_thread(tid, DEFAULT_PRIORITY);
634        writer.insert_reprioritization(
635            tid,
636            LogicalTime::from_nanos(978_307_199_223_456_789),
637            7,
638            DEFAULT_PRIORITY,
639            5,
640        );
641        let encoded = writer.into_string();
642
643        let decoded: PreemptionRecord = serde_json::from_str(&encoded).unwrap();
644        decoded.validate().unwrap();
645        assert_eq!(decoded.epoch(), Some(epoch));
646        assert_eq!(decoded.clone_preemptions_only().epoch(), Some(epoch));
647
648        let directory = tempfile::tempdir().unwrap();
649        let current = directory.path().join("current.json");
650        std::fs::write(&current, &encoded).unwrap();
651        assert_eq!(read_recorded_epoch(&current), Ok(Some(epoch)));
652
653        let mut legacy_value: serde_json::Value = serde_json::from_str(&encoded).unwrap();
654        legacy_value
655            .as_object_mut()
656            .unwrap()
657            .remove("epoch")
658            .unwrap();
659        let legacy = directory.path().join("legacy.json");
660        std::fs::write(&legacy, legacy_value.to_string()).unwrap();
661        assert_eq!(read_recorded_epoch(&legacy), Ok(None));
662        let legacy_record: PreemptionRecord = serde_json::from_value(legacy_value.clone()).unwrap();
663        assert_eq!(legacy_record.epoch(), None);
664        assert!(
665            !serde_json::to_string(&legacy_record)
666                .unwrap()
667                .contains("\"epoch\"")
668        );
669
670        let malformed = directory.path().join("malformed.json");
671        std::fs::write(&malformed, "{\"epoch\": 7}").unwrap();
672        assert!(read_recorded_epoch(&malformed).is_err());
673        assert!(read_recorded_epoch(&directory.path().join("absent.json")).is_err());
674    }
675
676    #[test]
677    fn print_preemptionrecord() {
678        let (file, path) = tempfile::NamedTempFile::new().unwrap().keep().unwrap();
679        drop(file);
680        let mut pw = PreemptionWriter::new(Some(path.clone()));
681        let tid1 = DetTid::from_raw(2);
682        let tid2 = DetTid::from_raw(4);
683        pw.register_thread(tid1, 1000);
684        pw.register_thread(tid2, 1000);
685
686        pw.insert_reprioritization(tid1, LogicalTime::from_nanos(3), 3, 1000, 3);
687        pw.insert_reprioritization(tid1, LogicalTime::from_nanos(30), 30, 3, 30);
688        pw.insert_reprioritization(tid1, LogicalTime::from_nanos(300), 300, 30, 300);
689        pw.insert_reprioritization(tid2, LogicalTime::from_nanos(2), 2, 1000, 2);
690        pw.insert_reprioritization(tid2, LogicalTime::from_nanos(20), 20, 2, 20);
691        pw.insert_reprioritization(tid2, LogicalTime::from_nanos(200), 200, 20, 200);
692
693        // First, round trip through a pretty-printed string
694        // ------------------------------------------------
695        let str: String = serde_json::to_string_pretty(&pw.inner).unwrap();
696        eprintln!("{}", str);
697        let pr2: PreemptionRecord = serde_json::from_str(&str).unwrap();
698        eprintln!("Round trip {:?}", pr2);
699        assert_eq!(pw.inner, pr2);
700
701        // Second, actually write it out to disk
702        // -------------------------------------
703        pw.flush().unwrap();
704        let reader = PreemptionReader::new(&path);
705        let th1 = reader.extract_thread_record(&tid1).unwrap();
706        let th2 = reader.extract_thread_record(&tid2).unwrap();
707        assert_eq!(th1.final_prio, 300);
708        assert_eq!(th2.final_prio, 200);
709
710        let mut exact = th1.clone().into_iter();
711        assert_eq!(
712            exact.next_with_rcbs(),
713            Some((LogicalTime::from_nanos(3), 1000, Some(3)))
714        );
715        assert_eq!(
716            exact.next_with_rcbs(),
717            Some((LogicalTime::from_nanos(30), 3, Some(30)))
718        );
719
720        let it1 = th1.into_iter();
721        let it2 = th2.into_iter();
722
723        assert_eq!(it1.initial_priority(), 1000);
724        assert_eq!(it2.initial_priority(), 1000);
725
726        let v1: Vec<(LogicalTime, Priority)> = it1.collect();
727        let v2: Vec<(LogicalTime, Priority)> = it2.collect();
728        assert_eq!(
729            v1,
730            vec![
731                (LogicalTime::from_nanos(3), 1000),
732                (LogicalTime::from_nanos(30), 3),
733                (LogicalTime::from_nanos(300), 30)
734            ]
735        );
736        assert_eq!(
737            v2,
738            vec![
739                (LogicalTime::from_nanos(2), 1000),
740                (LogicalTime::from_nanos(20), 2),
741                (LogicalTime::from_nanos(200), 20)
742            ]
743        );
744        std::fs::remove_file(path).unwrap();
745    }
746
747    #[test]
748    fn round_trip_vec_representations() {
749        let str = r#"{"per_thread":{"2":{"final_prio":1716,"prio_changes":[[946684799000013020,7301],[946684799000034020,9081],[946684799000041600,9238],[946684799000054790,865],
750[946684799000057440,751],[946684799000061970,275],[946684799000062730,5135],[946684799000069530,6339],
751[946684799000082850,1123],[946684799000101140,7875],[946684799000140625,4203],[946684799000171780,8611],
752[946684799000183550,6306],[946684799000184440,7958],[946684799000195750,8919],[946684799000226150,69],
753[946684799000236380,5915],[946684799000278180,3514],[946684799000320050,30],[946684799000334630,4629],
754[946684799000344650,2926],[946684799000355020,710],[946684799000365030,3513],[946684799000386350,4881],
755[946684799000396360,4852],[946684799000406840,4935],[946684799000426980,6672],[946684799000437970,7727],
756[946684799000452410,7017],[946684799000462430,1572],[946684799000546210,6395],[946684799000548120,3726],
757[946684799000562700,846],[946684799000583090,7838],[946684799000603310,8291],[946684799000655180,210],
758[946684799000666230,4576],[946684799000680910,4974],[946684799000723020,9160],[946684799000776780,439],
759[946684799000777080,6791],[946684799000787220,3015],[946684799000809090,7489],[946684799000840870,7165],
760[946684799000852855,4326],[946684799000854355,358],[946684799000866575,4448],[946684799000903415,6848],
761[946684799000913455,1899],[946684799000923630,2117],[946684799000963980,7705],[946684799001007090,8683],
762[946684799001016950,3317],[946684799001017980,5261],[946684799001027540,3478],[946684799001029990,6474],
763[946684799001053545,4823],[946684799001068395,4508],[946684799001073095,194],[946684799001121035,5944],
764[946684799001171285,8408],[946684799001171295,4493],[946684799001192845,3481]]}},"global":[]}"#;
765        let pr: PreemptionRecord = serde_json::from_str(str).unwrap();
766        let vecs = pr.as_vecs();
767        let pr2 = PreemptionRecord::from_vecs(&vecs);
768        assert_eq!(pr, pr2);
769    }
770
771    #[test]
772    fn normalize_preemption_record() {
773        let str = r#"{"per_thread":{
774"3":{"final_prio":1000,"prio_changes":[]},
775"5":{"final_prio":1000,"prio_changes":[]},
776"7":{"final_prio":9722,"prio_changes":[]},
777"9":{"final_prio":9982,"prio_changes":[[946684799006227400,7839]]},
778"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#;
779        let pr: PreemptionRecord = serde_json::from_str(str).unwrap();
780        let pr2 = pr.normalize();
781        pr2.validate().unwrap();
782        let bmap = pr2.as_vecs();
783        // Normalization may remove the useless default-prio entry:
784        if let Some(x) = bmap.get(&DetTid::from_raw(3)) {
785            assert_eq!(x, &vec![(LogicalTime::from_nanos(0), 1000)]);
786        }
787        if let Some(x) = bmap.get(&DetTid::from_raw(5)) {
788            assert_eq!(x, &vec![(LogicalTime::from_nanos(0), 1000)]);
789        }
790        assert_eq!(
791            bmap.get(&DetTid::from_raw(7)).unwrap(),
792            &vec![(LogicalTime::from_nanos(0), 1002)]
793        );
794        assert_eq!(
795            bmap.get(&DetTid::from_raw(9)).unwrap(),
796            &vec![
797                (LogicalTime::from_nanos(0), 1001),
798                (LogicalTime::from_nanos(946684799006227400), 1003)
799            ]
800        );
801        if let Some(x) = bmap.get(&DetTid::from_raw(11)) {
802            assert_eq!(x, &vec![(LogicalTime::from_nanos(0), 1000)]);
803        }
804    }
805
806    #[test_case(
807        r#"{"per_thread":{
808"3":{"final_prio":1000,"prio_changes":[]},
809"5":{"final_prio":1002,"prio_changes":[[946684799006227410,1000],[946684799006227415,1001]]},
810"9":{"final_prio":1002,"prio_changes":[[946684799006227400,1000],[946684799006227405,1001]]},
811"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#,
812        r#"{"per_thread":{
813"3":{"final_prio":1000,"prio_changes":[]},
814"5":{"final_prio":1001,"prio_changes":[[946684799006227410,1000]]},
815"9":{"final_prio":1002,"prio_changes":[[946684799006227400,1000],[946684799006227405,1001]]},
816"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#
817        ; "removes latest priority change and coalesces into final priority"
818    )]
819    #[test_case(
820        r#"{"per_thread":{
821"3":{"final_prio":1000,"prio_changes":[]},
822"5":{"final_prio":1000,"prio_changes":[]},
823"9":{"final_prio":1001,"prio_changes":[[946684799006227400,1000]]},
824"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#,
825        r#"{"per_thread":{
826"3":{"final_prio":1000,"prio_changes":[]},
827"5":{"final_prio":1000,"prio_changes":[]},
828"9":{"final_prio":1000,"prio_changes":[]},
829"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#
830        ; "removes only priority change and coalesces into final priority"
831    )]
832    #[test_case(
833        r#"{"per_thread":{
834"3":{"final_prio":1000,"prio_changes":[]},
835"5":{"final_prio":1000,"prio_changes":[]},
836"9":{"final_prio":1000,"prio_changes":[]},
837"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#,
838        r#"{"per_thread":{
839"3":{"final_prio":1000,"prio_changes":[]},
840"5":{"final_prio":1000,"prio_changes":[]},
841"9":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#
842        ; "removes entire history for last tid if no non-default priority changes"
843    )]
844    fn with_latest_preempt_removed(pr_json: &str, expected_pr_json: &str) {
845        let pr: PreemptionRecord = serde_json::from_str(pr_json).unwrap();
846        let expected_pr: PreemptionRecord = serde_json::from_str(expected_pr_json).unwrap();
847
848        let pr_with_latest_removed = pr.with_latest_preempt_removed();
849
850        pr_with_latest_removed.validate().unwrap();
851        self::assert_eq!(pr_with_latest_removed, expected_pr);
852    }
853
854    #[test]
855    fn test_split_map() {
856        let mut pr: PreemptionRecord = serde_json::from_str(
857            r#"
858            {
859                "per_thread" : {
860                },
861                "global" : [
862                    {
863                        "dettid": 3,
864                        "op": "OtherInstructions",
865                        "count": 1,
866                        "start_rip": null,
867                        "end_rip": null,
868                        "end_time": 946684799000000000
869                    },
870                    {
871                        "dettid": 3,
872                        "op": "Branch",
873                        "count": 311,
874                        "start_rip": null,
875                        "end_rip": null,
876                        "end_time": 946684799000003110
877                    },
878                    {
879                        "dettid": 3,
880                        "op": "OtherInstructions",
881                        "count": 1,
882                        "start_rip": null,
883                        "end_rip": null,
884                        "end_time": 946684799000003110
885                    }
886                ]
887            }
888        "#,
889        )
890        .unwrap();
891        let original = pr.global.clone();
892        pr.split_map(|e, _| vec![e]); // shouldn't change contents with this
893        assert_eq!(original, pr.global);
894
895        pr.split_map(|e, _| match &e {
896            // this will split event with Op::Branch into two
897            SchedEvent { op: Op::Branch, .. } => vec![e.clone(), e],
898            _ => vec![e],
899        });
900
901        assert_eq!(
902            pr.global.iter().map(|e| e.op).collect::<Vec<_>>(),
903            vec![
904                Op::OtherInstructions,
905                Op::Branch,
906                Op::Branch,
907                Op::OtherInstructions
908            ]
909        );
910
911        pr.split_map(|_, _| vec![]); // this will remove all entries
912        assert_eq!(pr.global, Vec::new());
913    }
914}
915
916// Reader and Writer implementations:
917// =============================================================================
918
919/// A writer for a stream of preemption events.
920#[derive(Debug)]
921pub struct PreemptionWriter {
922    inner: PreemptionRecord,
923    dest: Option<PathBuf>,
924    flushed: bool,
925}
926
927impl PreemptionWriter {
928    /// A new, empty record of preemptions, with an optional location on disk that it will be
929    /// written to.  If no path is supplied, the results will accumulate in memory only.
930    pub fn new(path: Option<PathBuf>) -> Self {
931        PreemptionWriter {
932            inner: Default::default(),
933            dest: path,
934            flushed: false,
935        }
936    }
937
938    /// Store the virtual-time epoch the recorded logical times are measured
939    /// from, so a replay can start its clock at the same instant.
940    pub fn with_epoch(mut self, epoch: DateTime<Utc>) -> Self {
941        self.inner.epoch = Some(epoch);
942        self
943    }
944
945    /// Does the record have zero entries?
946    pub fn is_empty(&self) -> bool {
947        self.inner.per_thread.is_empty()
948    }
949
950    /// Total number of elements in the record, across all threads.
951    pub fn len(&self) -> usize {
952        let mut count = 0;
953        for v in self.inner.per_thread.values() {
954            count += v.prio_changes.len() + v.chaos_epochs.len();
955        }
956        count
957    }
958
959    /// Register a thread with its default/final priority.
960    pub fn register_thread(&mut self, tid: DetTid, prio: Priority) {
961        if self
962            .inner
963            .per_thread
964            .insert(
965                tid,
966                ThreadHistory {
967                    final_prio: prio,
968                    prio_changes: Vec::new(),
969                    preemption_rcbs: Vec::new(),
970                    chaos_epochs: Vec::new(),
971                },
972            )
973            .is_some()
974        {
975            panic!(
976                "PreemptionRecord: error, cannot re-register thread id already registered: {}",
977                tid
978            )
979        }
980    }
981
982    /// Insert a new preemption point for the given thread.
983    /// It must monotonically increase in time.
984    ///
985    /// For now it (redundantly) takes both the priority of the timeslice just finished,
986    /// and the next one coming up.  This serves as a sanity check.
987    pub fn insert_reprioritization(
988        &mut self,
989        tid: DetTid,
990        time: LogicalTime,
991        rcbs: u64,
992        prior_prio: Priority,
993        next_prio: Priority,
994    ) {
995        let history = self.inner.per_thread.get_mut(&tid).unwrap_or_else(|| {
996            panic!(
997                "PreemptionRecord: Cannot insert a preemption before registering thread {}",
998                tid
999            )
1000        });
1001        assert_eq!(history.final_prio, prior_prio);
1002
1003        if let Some((last, _prio)) = history.prio_changes.last() {
1004            assert!(&time > last);
1005        }
1006        history.prio_changes.push((time, prior_prio));
1007        history.preemption_rcbs.push(rcbs);
1008        history.final_prio = next_prio;
1009    }
1010
1011    // AUTONOMOUS-BOT-IMPLEMENTED
1012    // TODO-HUMAN-REVIEW(PR-1151)
1013    /// Record a deterministic slowdown epoch transition for replay.
1014    pub fn insert_chaos_epoch(&mut self, tid: DetTid, transition: ChaosEpochTransition) {
1015        let history = self.inner.per_thread.get_mut(&tid).unwrap_or_else(|| {
1016            panic!(
1017                "PreemptionRecord: Cannot insert a chaos epoch before registering thread {}",
1018                tid
1019            )
1020        });
1021        if let Some(last) = history.chaos_epochs.last() {
1022            assert!(transition.logical_time > last.logical_time);
1023            assert!(transition.epoch > last.epoch);
1024        }
1025        history.chaos_epochs.push(transition);
1026    }
1027
1028    /// Add a SchedEvent to the global log of thread behavior.
1029    pub fn insert_schedevent(&mut self, ev: SchedEvent) {
1030        if ev.count > 0 {
1031            self.inner.global.push(ev)
1032            // TODO, aggregation: possibly check if the last event was the SAME, and combine them,
1033            // increasing the count.
1034        } else {
1035            trace!("NOT recording scheduled event with zero count!");
1036        }
1037    }
1038
1039    /// Set the priority of the current thread after a change.
1040    pub fn set_current(&mut self, tid: DetTid, new_prio: Priority) {
1041        let history = self.inner.per_thread.get_mut(&tid).unwrap_or_else(|| {
1042            panic!(
1043                "PreemptionRecord: Cannot set current priority before registering thread {}",
1044                tid
1045            )
1046        });
1047        history.final_prio = new_prio;
1048    }
1049
1050    /// Abort writing to disk and instead gather the output thusfar into a string.
1051    pub fn into_string(mut self) -> String {
1052        self.flushed = true;
1053        self.inner.to_string()
1054    }
1055
1056    /// Finish any IO necessary to flush the record to disk.
1057    /// This will also happen automatically on drop.
1058    pub fn flush(mut self) -> Result<(), String> {
1059        self.flushed = true;
1060        self.write_to_disk()
1061    }
1062
1063    fn write_to_disk(&mut self) -> Result<(), String> {
1064        if let Some(path) = &self.dest {
1065            self.inner.write_to_disk(path)
1066        } else {
1067            Err(
1068                "Cannot write_to_disk because this PreemptionWriter was created without a backing file.".to_string()
1069            )
1070        }
1071    }
1072}
1073
1074impl Drop for PreemptionWriter {
1075    fn drop(&mut self) {
1076        if !self.flushed
1077            && self.dest.is_some()
1078            && let Err(e) = self.write_to_disk()
1079        {
1080            panic!("Error while dropping PreemptionWriter: {}", e);
1081        }
1082    }
1083}
1084
1085/// A reader for a stream of preemption events.
1086#[derive(Debug)]
1087pub struct PreemptionReader {
1088    inner: PreemptionRecord,
1089}
1090
1091// Deprecated: use schedevents()
1092/// Read a full trace from disk.  Panic if it doesn't load.
1093pub fn read_trace(path: &Path) -> Vec<SchedEvent> {
1094    let pr = read_preemption_record(path);
1095    pr.global
1096}
1097
1098// TODO: we should implement streaming and not read this all at once.
1099fn read_preemption_record(path: &Path) -> PreemptionRecord {
1100    let string = std::fs::read_to_string(path)
1101        .unwrap_or_else(|e| panic!("Error reading file {:?}:\n {}", path, e));
1102    let pr: PreemptionRecord = serde_json::from_str(&string).unwrap_or_else(|e| {
1103        panic!(
1104            "Error parsing PreemptionRecord from JSON: {}\nJSON contents:\n{}",
1105            e, string
1106        )
1107    });
1108    if let Err(e) = pr.validate() {
1109        panic!(
1110            "Invalid PreemptionRecord when loading from path {}. Error:\n {}",
1111            path.display(),
1112            e
1113        );
1114    }
1115    pr
1116}
1117
1118/// Read only the virtual-time epoch stored in a preemption record file.
1119///
1120/// Unlike [`PreemptionReader::new`], this reports an unreadable or malformed
1121/// file as an error instead of panicking, because the CLI consults it before
1122/// any guest starts. `Ok(None)` means the record predates stored epochs.
1123pub fn read_recorded_epoch(path: &Path) -> Result<Option<DateTime<Utc>>, String> {
1124    #[derive(Deserialize)]
1125    struct EpochOnly {
1126        #[serde(default)]
1127        epoch: Option<DateTime<Utc>>,
1128    }
1129    let file = File::open(path)
1130        .map_err(|e| format!("cannot read preemption record {}: {}", path.display(), e))?;
1131    let record: EpochOnly = serde_json::from_reader(std::io::BufReader::new(file))
1132        .map_err(|e| format!("cannot parse preemption record {}: {}", path.display(), e))?;
1133    Ok(record.epoch)
1134}
1135
1136// TODO: eventually this should do streaming IO, and abstract it behind this interface.
1137impl PreemptionReader {
1138    /// Access the stored `PreemptionRecord`, and eagerly or lazily load its data from disk.
1139    pub fn new(path: &Path) -> Self {
1140        let pr = read_preemption_record(path);
1141        PreemptionReader { inner: pr }
1142    }
1143
1144    /// Gets the inner `PreemptionRecord`.
1145    pub fn into_inner(self) -> PreemptionRecord {
1146        self.inner
1147    }
1148
1149    /// Extract a copy of all the entries for a particular thread.  Note that this returns None if
1150    /// there is NO record for the thread (not registered).  Otherwise, it returns a copy of the
1151    /// `ThreadHistory`, which may be read into memory or not.
1152    pub fn extract_thread_record(&self, tid: &DetTid) -> Option<ThreadHistory> {
1153        self.inner.per_thread.get(tid).cloned()
1154    }
1155
1156    /// Return the initial priority for a thread, if it is recorded in the record.
1157    pub fn thread_initial_priority(&self, tid: &DetTid) -> Option<Priority> {
1158        self.inner.per_thread.get(tid).map(|x| x.initial_priority())
1159    }
1160
1161    /// Return all the threads in this record.
1162    pub fn all_threads(&self) -> Vec<DetTid> {
1163        // TODO(T110956298): have this return an iterator.
1164        self.inner.per_thread.keys().copied().collect()
1165    }
1166
1167    /// Load the full record of preemptions into memory.
1168    pub fn load_all(&self) -> PreemptionRecord {
1169        self.inner.clone()
1170    }
1171
1172    /// Size in number of preemptions, across all threads.
1173    pub fn size(&self) -> usize {
1174        let mut sum = 0;
1175        for th in self.inner.per_thread.values() {
1176            sum += th.prio_changes.len()
1177        }
1178        sum
1179    }
1180}
1181
1182/// Take a file containing a PreemptionRecord and strip the times from all the events.
1183/// Optionally take a destination path, otherwise create a destination file in the same directory
1184/// with our own naming convention.
1185pub fn strip_times_from_events_file(
1186    sched_path: &Path,
1187    dest: Option<PathBuf>,
1188) -> anyhow::Result<PathBuf> {
1189    let new_path = dest.unwrap_or_else(|| sched_path.with_extension("notimes"));
1190    let mut preemptions = PreemptionReader::new(sched_path).into_inner();
1191    for se in preemptions.schedevents_iter_mut() {
1192        se.end_time = None;
1193    }
1194    preemptions
1195        .write_to_disk(new_path.as_ref())
1196        .map_err(anyhow::Error::msg)?;
1197    Ok(new_path)
1198}