use std::collections::BTreeMap;
use std::fs::File;
use std::iter::FromIterator;
use std::path::Path;
use std::path::PathBuf;
use chrono::DateTime;
use chrono::Utc;
use detcore_model::collections::ReplayCursor;
use serde::Deserialize;
use serde::Serialize;
use tracing::trace;
use crate::resources::ChaosEpochTransition;
use crate::scheduler::Priority;
use crate::scheduler::runqueue::DEFAULT_PRIORITY;
use crate::scheduler::runqueue::FIRST_PRIORITY;
use crate::scheduler::runqueue::LAST_PRIORITY;
use crate::scheduler::runqueue::is_ordinary_priority;
use crate::types::DetTid;
use crate::types::LogicalTime;
use crate::types::SchedEvent;
#[derive(PartialEq, Default, Debug, Eq, Clone, Hash, Serialize, Deserialize)]
pub struct PreemptionRecord {
per_thread: BTreeMap<DetTid, ThreadHistory>,
global: Vec<SchedEvent>,
#[serde(default, skip_serializing_if = "Option::is_none")]
epoch: Option<DateTime<Utc>>,
}
impl std::fmt::Display for PreemptionRecord {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let str = serde_json::to_string(&self).unwrap();
write!(f, "{}", str)
}
}
impl PreemptionRecord {
pub fn from_sched_events(events: Vec<SchedEvent>) -> Self {
Self {
per_thread: Default::default(),
global: events,
epoch: None,
}
}
pub fn epoch(&self) -> Option<DateTime<Utc>> {
self.epoch
}
pub fn extract_all(&self) -> BTreeMap<DetTid, ThreadHistory> {
self.per_thread.clone()
}
pub fn as_vecs(&self) -> BTreeMap<DetTid, Vec<(LogicalTime, Priority)>> {
let mut bt = BTreeMap::new();
for (tid, th) in &self.per_thread {
bt.insert(*tid, th.as_vec());
}
bt
}
pub fn strip_contents(mut self) -> Self {
for history in self.per_thread.values_mut() {
history.prio_changes = Vec::new();
history.preemption_rcbs = Vec::new();
history.chaos_epochs = Vec::new();
history.final_prio = 1000;
}
self.global = Vec::new();
self
}
pub fn split_map<F, R>(&mut self, splitter: F)
where
R: IntoIterator<Item = SchedEvent>,
F: Fn(SchedEvent, &ReplayCursor<SchedEvent>) -> R,
{
let mut result = Vec::new();
let global = std::mem::take(&mut self.global);
let mut cursor = ReplayCursor::from_iter(global);
while let Some(event) = cursor.next() {
for new_event in splitter(event, &cursor) {
result.push(new_event);
}
}
self.global = result;
}
pub fn schedevents_iter_mut(&mut self) -> std::slice::IterMut<'_, SchedEvent> {
self.global.iter_mut()
}
pub fn preemptions_only(&mut self) {
self.global.clear();
}
pub fn clone_preemptions_only(&self) -> Self {
PreemptionRecord {
per_thread: self.per_thread.clone(),
global: Vec::new(),
epoch: self.epoch,
}
}
pub fn schedevents(&self) -> &Vec<SchedEvent> {
&self.global
}
pub fn contains_schedevents(&self) -> bool {
!self.global.is_empty()
}
pub fn from_vecs(bt: &BTreeMap<DetTid, Vec<(LogicalTime, Priority)>>) -> Self {
let mut bt2 = BTreeMap::new();
for (tid, vec) in bt {
let th = if vec.is_empty() {
ThreadHistory {
final_prio: DEFAULT_PRIORITY,
prio_changes: Vec::new(),
preemption_rcbs: Vec::new(),
chaos_epochs: Vec::new(),
}
} else {
let (_final_end, final_prio) = vec.last().unwrap();
let final_prio = *final_prio;
let mut prio_changes = Vec::new();
let mut it = vec.iter().peekable();
while let Some((_this_ns, this_p)) = it.next() {
if let Some((next_ns, _next_p)) = it.peek() {
prio_changes.push((*next_ns, *this_p));
} else {
break;
}
}
assert!(final_prio >= FIRST_PRIORITY);
assert!(final_prio <= LAST_PRIORITY);
ThreadHistory {
final_prio,
prio_changes,
preemption_rcbs: Vec::new(),
chaos_epochs: Vec::new(),
}
};
bt2.insert(*tid, th);
}
PreemptionRecord {
per_thread: bt2,
global: Vec::new(),
epoch: None,
}
}
pub fn into_global(self) -> Vec<SchedEvent> {
self.global
}
pub fn write_to_disk(&self, path: &Path) -> Result<(), String> {
let mut str: String = self.to_string();
str.push('\n');
match File::create(path) {
Ok(mut file) => match std::io::Write::write_all(&mut file, str.as_bytes()) {
Ok(_) => Ok(()),
Err(err) => Err(format!(
"Failed to write preemption record to file {:?}, error: {}",
path, err
)),
},
Err(err) => Err(format!(
"Failed to create file for preemption record {:?}, error: {}",
path, err
)),
}
}
pub fn validate(&self) -> Result<(), String> {
for (tid, history) in &self.per_thread {
if !is_ordinary_priority(history.final_prio) {
return Err(format!(
"final priority for thread {} invalid: {}",
tid, history.final_prio
));
}
{
let mut time_last = None;
for (count, (ns, prio)) in history.prio_changes.iter().enumerate() {
if let Some(last) = time_last {
if !is_ordinary_priority(*prio) {
return Err(format!(
"preemption priority #{} for thread {} invalid: {}",
count, tid, history.final_prio
));
}
if *ns <= last {
return Err(format!(
"Timestamps failed to monotonically increase ({}), in series:\n {:?}",
ns, history.prio_changes
));
}
}
time_last = Some(*ns);
}
}
if !history.preemption_rcbs.is_empty()
&& history.preemption_rcbs.len() != history.prio_changes.len()
{
return Err(format!(
"thread {} has {} preemption times but {} RCB targets",
tid,
history.prio_changes.len(),
history.preemption_rcbs.len()
));
}
if history
.preemption_rcbs
.windows(2)
.any(|pair| pair[1] < pair[0])
{
return Err(format!(
"preemption RCB targets failed to increase for thread {}",
tid
));
}
let mut transition_last = None;
let mut epoch_last = None;
for transition in &history.chaos_epochs {
if transition.factor.as_f64() <= 0.0 {
return Err(format!(
"chaos epoch factor for thread {} must be positive",
tid
));
}
if transition_last.is_some_and(|last| transition.logical_time <= last) {
return Err(format!(
"chaos epoch transition times failed to increase for thread {}",
tid
));
}
if epoch_last.is_some_and(|last| transition.epoch <= last) {
return Err(format!(
"chaos epoch numbers failed to increase for thread {}",
tid
));
}
transition_last = Some(transition.logical_time);
epoch_last = Some(transition.epoch);
}
}
Ok(())
}
pub fn normalize(&self) -> PreemptionRecord {
let mut clone = self.clone();
let mut priomap: BTreeMap<Priority, Priority> = BTreeMap::new();
for history in clone.per_thread.values_mut() {
let _ = priomap.insert(history.final_prio, 0);
for (_ns, prio) in &history.prio_changes {
let _ = priomap.insert(*prio, 0);
}
}
for (cur_prio, val) in (DEFAULT_PRIORITY..).zip(priomap.values_mut()) {
assert!(cur_prio <= LAST_PRIORITY);
*val = cur_prio;
}
for history in clone.per_thread.values_mut() {
history.final_prio = *priomap.get(&history.final_prio).unwrap();
for (_ns, prio) in &mut history.prio_changes {
*prio = *priomap.get(prio).unwrap();
}
}
let mut finalmap = BTreeMap::new();
for (tid, history) in clone.per_thread.into_iter() {
if history.final_prio != DEFAULT_PRIORITY
|| !history.prio_changes.is_empty()
|| !history.chaos_epochs.is_empty()
{
assert!(finalmap.insert(tid, history).is_none());
}
}
clone.per_thread = finalmap;
clone
}
pub fn with_latest_preempt_removed(&self) -> PreemptionRecord {
let mut preempts_latest_prio_changes: Vec<(DetTid, LogicalTime)> = self
.per_thread
.clone()
.into_iter()
.map(|(tid, th)| {
(
tid,
th.prio_changes
.last() .map_or(LogicalTime::ZERO, |prio_change| prio_change.0),
)
})
.collect();
preempts_latest_prio_changes.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap());
let tid_with_preempt_to_drop = preempts_latest_prio_changes
.into_iter()
.map(|(tid, _th)| tid)
.next();
let mut clone = self.clone();
if let Some(tid) = tid_with_preempt_to_drop {
let thread_history = &mut clone.per_thread.get_mut(&tid).unwrap();
if !thread_history.prio_changes.is_empty() {
let removed_prio = thread_history.prio_changes.pop().unwrap().1;
if !thread_history.preemption_rcbs.is_empty() {
thread_history.preemption_rcbs.pop();
}
thread_history.final_prio = removed_prio;
} else {
clone.per_thread.pop_last();
}
}
clone
}
}
#[derive(
PartialEq, // Silly protection from rustfmt disagreements.
Debug,
Eq,
Clone,
Hash,
Serialize,
Deserialize,
)]
pub struct ThreadHistory {
pub final_prio: Priority,
prio_changes: Vec<(LogicalTime, Priority)>,
#[serde(default)]
preemption_rcbs: Vec<u64>,
#[serde(default)]
chaos_epochs: Vec<ChaosEpochTransition>,
}
impl ThreadHistory {
pub fn new() -> Self {
ThreadHistory {
final_prio: DEFAULT_PRIORITY,
prio_changes: Vec::new(),
preemption_rcbs: Vec::new(),
chaos_epochs: Vec::new(),
}
}
#[allow(clippy::should_implement_trait)]
pub fn into_iter(self) -> ThreadHistoryIterator {
ThreadHistoryIterator {
full_history: self,
ix: 0,
chaos_epoch_ix: 0,
}
}
#[cfg(test)]
pub(crate) fn with_chaos_epochs(mut self, chaos_epochs: Vec<ChaosEpochTransition>) -> Self {
self.chaos_epochs = chaos_epochs;
self
}
#[cfg(test)]
pub(crate) fn with_prio_changes(mut self, prio_changes: Vec<(LogicalTime, Priority)>) -> Self {
self.prio_changes = prio_changes;
self
}
#[cfg(test)]
pub(crate) fn with_preemption_rcbs(mut self, preemption_rcbs: Vec<u64>) -> Self {
self.preemption_rcbs = preemption_rcbs;
self
}
pub fn as_vec(&self) -> Vec<(LogicalTime, Priority)> {
let mut vec = Vec::new();
let mut time0 = LogicalTime::from_nanos(0);
for (time1, prio) in &self.prio_changes {
vec.push((time0, *prio));
time0 = *time1;
}
vec.push((time0, self.final_prio));
vec
}
pub fn initial_priority(&self) -> Priority {
if let Some((_ns, pr)) = self.prio_changes.first() {
*pr
} else {
self.final_prio
}
}
}
impl Default for ThreadHistory {
fn default() -> Self {
Self::new()
}
}
#[derive(PartialEq, Eq, Clone, Serialize, Deserialize)]
pub struct ThreadHistoryIterator {
full_history: ThreadHistory,
ix: usize,
#[serde(default)]
chaos_epoch_ix: usize,
}
impl ThreadHistoryIterator {
pub fn initial_priority(&self) -> Priority {
self.full_history.initial_priority()
}
pub fn final_priority(&self) -> Priority {
self.full_history.final_prio
}
pub fn advance_chaos_epoch(
&mut self,
current_time: LogicalTime,
) -> Option<ChaosEpochTransition> {
let mut changed = None;
while let Some(transition) = self.full_history.chaos_epochs.get(self.chaos_epoch_ix)
&& transition.logical_time <= current_time
{
changed = Some(*transition);
self.chaos_epoch_ix += 1;
}
changed
}
pub fn has_chaos_epochs(&self) -> bool {
!self.full_history.chaos_epochs.is_empty()
}
pub fn next_with_rcbs(&mut self) -> Option<(LogicalTime, Priority, Option<u64>)> {
let ix = self.ix;
self.next().map(|(time, priority)| {
(
time,
priority,
self.full_history.preemption_rcbs.get(ix).copied(),
)
})
}
}
impl Iterator for ThreadHistoryIterator {
type Item = (LogicalTime, Priority);
fn next(&mut self) -> Option<Self::Item> {
let vec = &self.full_history.prio_changes;
if vec.len() > self.ix {
let elt = vec[self.ix];
self.ix += 1;
Some(elt)
} else {
None
}
}
}
#[cfg(test)]
mod tests {
use detcore_model::schedule::Op;
use pretty_assertions::assert_eq;
use test_case::test_case;
use super::*;
use crate::types::RcbTimeMultiplier;
#[test]
fn chaos_epoch_transitions_round_trip_and_replay_exact_factors() {
let tid = DetTid::from_raw(2);
let first = ChaosEpochTransition {
logical_time: LogicalTime::from_nanos(100),
epoch: 0,
factor: RcbTimeMultiplier::from_f64(2.5),
};
let second = ChaosEpochTransition {
logical_time: LogicalTime::from_nanos(500),
epoch: 1,
factor: RcbTimeMultiplier::from_f64(0.75),
};
let mut writer = PreemptionWriter::new(None);
writer.register_thread(tid, DEFAULT_PRIORITY);
writer.insert_chaos_epoch(tid, first);
writer.insert_chaos_epoch(tid, second);
let encoded = writer.into_string();
assert!(encoded.contains("chaos_epochs"));
let decoded: PreemptionRecord = serde_json::from_str(&encoded).unwrap();
decoded.validate().unwrap();
let mut history = decoded.extract_all().remove(&tid).unwrap().into_iter();
assert_eq!(
history.advance_chaos_epoch(LogicalTime::from_nanos(99)),
None
);
assert_eq!(
history.advance_chaos_epoch(LogicalTime::from_nanos(100)),
Some(first)
);
assert_eq!(
history.advance_chaos_epoch(LogicalTime::from_nanos(499)),
None
);
assert_eq!(
history.advance_chaos_epoch(LogicalTime::from_nanos(500)),
Some(second)
);
assert_eq!(history.advance_chaos_epoch(LogicalTime::MAX), None);
}
#[test]
fn recorded_epoch_round_trips_and_legacy_records_have_none() {
let epoch: DateTime<Utc> = "2000-12-31T23:59:59.123456789Z".parse().unwrap();
let tid = DetTid::from_raw(3);
let mut writer = PreemptionWriter::new(None).with_epoch(epoch);
writer.register_thread(tid, DEFAULT_PRIORITY);
writer.insert_reprioritization(
tid,
LogicalTime::from_nanos(978_307_199_223_456_789),
7,
DEFAULT_PRIORITY,
5,
);
let encoded = writer.into_string();
let decoded: PreemptionRecord = serde_json::from_str(&encoded).unwrap();
decoded.validate().unwrap();
assert_eq!(decoded.epoch(), Some(epoch));
assert_eq!(decoded.clone_preemptions_only().epoch(), Some(epoch));
let directory = tempfile::tempdir().unwrap();
let current = directory.path().join("current.json");
std::fs::write(¤t, &encoded).unwrap();
assert_eq!(read_recorded_epoch(¤t), Ok(Some(epoch)));
let mut legacy_value: serde_json::Value = serde_json::from_str(&encoded).unwrap();
legacy_value
.as_object_mut()
.unwrap()
.remove("epoch")
.unwrap();
let legacy = directory.path().join("legacy.json");
std::fs::write(&legacy, legacy_value.to_string()).unwrap();
assert_eq!(read_recorded_epoch(&legacy), Ok(None));
let legacy_record: PreemptionRecord = serde_json::from_value(legacy_value.clone()).unwrap();
assert_eq!(legacy_record.epoch(), None);
assert!(
!serde_json::to_string(&legacy_record)
.unwrap()
.contains("\"epoch\"")
);
let malformed = directory.path().join("malformed.json");
std::fs::write(&malformed, "{\"epoch\": 7}").unwrap();
assert!(read_recorded_epoch(&malformed).is_err());
assert!(read_recorded_epoch(&directory.path().join("absent.json")).is_err());
}
#[test]
fn print_preemptionrecord() {
let (file, path) = tempfile::NamedTempFile::new().unwrap().keep().unwrap();
drop(file);
let mut pw = PreemptionWriter::new(Some(path.clone()));
let tid1 = DetTid::from_raw(2);
let tid2 = DetTid::from_raw(4);
pw.register_thread(tid1, 1000);
pw.register_thread(tid2, 1000);
pw.insert_reprioritization(tid1, LogicalTime::from_nanos(3), 3, 1000, 3);
pw.insert_reprioritization(tid1, LogicalTime::from_nanos(30), 30, 3, 30);
pw.insert_reprioritization(tid1, LogicalTime::from_nanos(300), 300, 30, 300);
pw.insert_reprioritization(tid2, LogicalTime::from_nanos(2), 2, 1000, 2);
pw.insert_reprioritization(tid2, LogicalTime::from_nanos(20), 20, 2, 20);
pw.insert_reprioritization(tid2, LogicalTime::from_nanos(200), 200, 20, 200);
let str: String = serde_json::to_string_pretty(&pw.inner).unwrap();
eprintln!("{}", str);
let pr2: PreemptionRecord = serde_json::from_str(&str).unwrap();
eprintln!("Round trip {:?}", pr2);
assert_eq!(pw.inner, pr2);
pw.flush().unwrap();
let reader = PreemptionReader::new(&path);
let th1 = reader.extract_thread_record(&tid1).unwrap();
let th2 = reader.extract_thread_record(&tid2).unwrap();
assert_eq!(th1.final_prio, 300);
assert_eq!(th2.final_prio, 200);
let mut exact = th1.clone().into_iter();
assert_eq!(
exact.next_with_rcbs(),
Some((LogicalTime::from_nanos(3), 1000, Some(3)))
);
assert_eq!(
exact.next_with_rcbs(),
Some((LogicalTime::from_nanos(30), 3, Some(30)))
);
let it1 = th1.into_iter();
let it2 = th2.into_iter();
assert_eq!(it1.initial_priority(), 1000);
assert_eq!(it2.initial_priority(), 1000);
let v1: Vec<(LogicalTime, Priority)> = it1.collect();
let v2: Vec<(LogicalTime, Priority)> = it2.collect();
assert_eq!(
v1,
vec![
(LogicalTime::from_nanos(3), 1000),
(LogicalTime::from_nanos(30), 3),
(LogicalTime::from_nanos(300), 30)
]
);
assert_eq!(
v2,
vec![
(LogicalTime::from_nanos(2), 1000),
(LogicalTime::from_nanos(20), 2),
(LogicalTime::from_nanos(200), 20)
]
);
std::fs::remove_file(path).unwrap();
}
#[test]
fn round_trip_vec_representations() {
let str = r#"{"per_thread":{"2":{"final_prio":1716,"prio_changes":[[946684799000013020,7301],[946684799000034020,9081],[946684799000041600,9238],[946684799000054790,865],
[946684799000057440,751],[946684799000061970,275],[946684799000062730,5135],[946684799000069530,6339],
[946684799000082850,1123],[946684799000101140,7875],[946684799000140625,4203],[946684799000171780,8611],
[946684799000183550,6306],[946684799000184440,7958],[946684799000195750,8919],[946684799000226150,69],
[946684799000236380,5915],[946684799000278180,3514],[946684799000320050,30],[946684799000334630,4629],
[946684799000344650,2926],[946684799000355020,710],[946684799000365030,3513],[946684799000386350,4881],
[946684799000396360,4852],[946684799000406840,4935],[946684799000426980,6672],[946684799000437970,7727],
[946684799000452410,7017],[946684799000462430,1572],[946684799000546210,6395],[946684799000548120,3726],
[946684799000562700,846],[946684799000583090,7838],[946684799000603310,8291],[946684799000655180,210],
[946684799000666230,4576],[946684799000680910,4974],[946684799000723020,9160],[946684799000776780,439],
[946684799000777080,6791],[946684799000787220,3015],[946684799000809090,7489],[946684799000840870,7165],
[946684799000852855,4326],[946684799000854355,358],[946684799000866575,4448],[946684799000903415,6848],
[946684799000913455,1899],[946684799000923630,2117],[946684799000963980,7705],[946684799001007090,8683],
[946684799001016950,3317],[946684799001017980,5261],[946684799001027540,3478],[946684799001029990,6474],
[946684799001053545,4823],[946684799001068395,4508],[946684799001073095,194],[946684799001121035,5944],
[946684799001171285,8408],[946684799001171295,4493],[946684799001192845,3481]]}},"global":[]}"#;
let pr: PreemptionRecord = serde_json::from_str(str).unwrap();
let vecs = pr.as_vecs();
let pr2 = PreemptionRecord::from_vecs(&vecs);
assert_eq!(pr, pr2);
}
#[test]
fn normalize_preemption_record() {
let str = r#"{"per_thread":{
"3":{"final_prio":1000,"prio_changes":[]},
"5":{"final_prio":1000,"prio_changes":[]},
"7":{"final_prio":9722,"prio_changes":[]},
"9":{"final_prio":9982,"prio_changes":[[946684799006227400,7839]]},
"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#;
let pr: PreemptionRecord = serde_json::from_str(str).unwrap();
let pr2 = pr.normalize();
pr2.validate().unwrap();
let bmap = pr2.as_vecs();
if let Some(x) = bmap.get(&DetTid::from_raw(3)) {
assert_eq!(x, &vec![(LogicalTime::from_nanos(0), 1000)]);
}
if let Some(x) = bmap.get(&DetTid::from_raw(5)) {
assert_eq!(x, &vec![(LogicalTime::from_nanos(0), 1000)]);
}
assert_eq!(
bmap.get(&DetTid::from_raw(7)).unwrap(),
&vec![(LogicalTime::from_nanos(0), 1002)]
);
assert_eq!(
bmap.get(&DetTid::from_raw(9)).unwrap(),
&vec![
(LogicalTime::from_nanos(0), 1001),
(LogicalTime::from_nanos(946684799006227400), 1003)
]
);
if let Some(x) = bmap.get(&DetTid::from_raw(11)) {
assert_eq!(x, &vec![(LogicalTime::from_nanos(0), 1000)]);
}
}
#[test_case(
r#"{"per_thread":{
"3":{"final_prio":1000,"prio_changes":[]},
"5":{"final_prio":1002,"prio_changes":[[946684799006227410,1000],[946684799006227415,1001]]},
"9":{"final_prio":1002,"prio_changes":[[946684799006227400,1000],[946684799006227405,1001]]},
"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#,
r#"{"per_thread":{
"3":{"final_prio":1000,"prio_changes":[]},
"5":{"final_prio":1001,"prio_changes":[[946684799006227410,1000]]},
"9":{"final_prio":1002,"prio_changes":[[946684799006227400,1000],[946684799006227405,1001]]},
"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#
; "removes latest priority change and coalesces into final priority"
)]
#[test_case(
r#"{"per_thread":{
"3":{"final_prio":1000,"prio_changes":[]},
"5":{"final_prio":1000,"prio_changes":[]},
"9":{"final_prio":1001,"prio_changes":[[946684799006227400,1000]]},
"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#,
r#"{"per_thread":{
"3":{"final_prio":1000,"prio_changes":[]},
"5":{"final_prio":1000,"prio_changes":[]},
"9":{"final_prio":1000,"prio_changes":[]},
"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#
; "removes only priority change and coalesces into final priority"
)]
#[test_case(
r#"{"per_thread":{
"3":{"final_prio":1000,"prio_changes":[]},
"5":{"final_prio":1000,"prio_changes":[]},
"9":{"final_prio":1000,"prio_changes":[]},
"11":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#,
r#"{"per_thread":{
"3":{"final_prio":1000,"prio_changes":[]},
"5":{"final_prio":1000,"prio_changes":[]},
"9":{"final_prio":1000,"prio_changes":[]}},"global":[]}"#
; "removes entire history for last tid if no non-default priority changes"
)]
fn with_latest_preempt_removed(pr_json: &str, expected_pr_json: &str) {
let pr: PreemptionRecord = serde_json::from_str(pr_json).unwrap();
let expected_pr: PreemptionRecord = serde_json::from_str(expected_pr_json).unwrap();
let pr_with_latest_removed = pr.with_latest_preempt_removed();
pr_with_latest_removed.validate().unwrap();
self::assert_eq!(pr_with_latest_removed, expected_pr);
}
#[test]
fn test_split_map() {
let mut pr: PreemptionRecord = serde_json::from_str(
r#"
{
"per_thread" : {
},
"global" : [
{
"dettid": 3,
"op": "OtherInstructions",
"count": 1,
"start_rip": null,
"end_rip": null,
"end_time": 946684799000000000
},
{
"dettid": 3,
"op": "Branch",
"count": 311,
"start_rip": null,
"end_rip": null,
"end_time": 946684799000003110
},
{
"dettid": 3,
"op": "OtherInstructions",
"count": 1,
"start_rip": null,
"end_rip": null,
"end_time": 946684799000003110
}
]
}
"#,
)
.unwrap();
let original = pr.global.clone();
pr.split_map(|e, _| vec![e]); assert_eq!(original, pr.global);
pr.split_map(|e, _| match &e {
SchedEvent { op: Op::Branch, .. } => vec![e.clone(), e],
_ => vec![e],
});
assert_eq!(
pr.global.iter().map(|e| e.op).collect::<Vec<_>>(),
vec![
Op::OtherInstructions,
Op::Branch,
Op::Branch,
Op::OtherInstructions
]
);
pr.split_map(|_, _| vec![]); assert_eq!(pr.global, Vec::new());
}
}
#[derive(Debug)]
pub struct PreemptionWriter {
inner: PreemptionRecord,
dest: Option<PathBuf>,
flushed: bool,
}
impl PreemptionWriter {
pub fn new(path: Option<PathBuf>) -> Self {
PreemptionWriter {
inner: Default::default(),
dest: path,
flushed: false,
}
}
pub fn with_epoch(mut self, epoch: DateTime<Utc>) -> Self {
self.inner.epoch = Some(epoch);
self
}
pub fn is_empty(&self) -> bool {
self.inner.per_thread.is_empty()
}
pub fn len(&self) -> usize {
let mut count = 0;
for v in self.inner.per_thread.values() {
count += v.prio_changes.len() + v.chaos_epochs.len();
}
count
}
pub fn register_thread(&mut self, tid: DetTid, prio: Priority) {
if self
.inner
.per_thread
.insert(
tid,
ThreadHistory {
final_prio: prio,
prio_changes: Vec::new(),
preemption_rcbs: Vec::new(),
chaos_epochs: Vec::new(),
},
)
.is_some()
{
panic!(
"PreemptionRecord: error, cannot re-register thread id already registered: {}",
tid
)
}
}
pub fn insert_reprioritization(
&mut self,
tid: DetTid,
time: LogicalTime,
rcbs: u64,
prior_prio: Priority,
next_prio: Priority,
) {
let history = self.inner.per_thread.get_mut(&tid).unwrap_or_else(|| {
panic!(
"PreemptionRecord: Cannot insert a preemption before registering thread {}",
tid
)
});
assert_eq!(history.final_prio, prior_prio);
if let Some((last, _prio)) = history.prio_changes.last() {
assert!(&time > last);
}
history.prio_changes.push((time, prior_prio));
history.preemption_rcbs.push(rcbs);
history.final_prio = next_prio;
}
pub fn insert_chaos_epoch(&mut self, tid: DetTid, transition: ChaosEpochTransition) {
let history = self.inner.per_thread.get_mut(&tid).unwrap_or_else(|| {
panic!(
"PreemptionRecord: Cannot insert a chaos epoch before registering thread {}",
tid
)
});
if let Some(last) = history.chaos_epochs.last() {
assert!(transition.logical_time > last.logical_time);
assert!(transition.epoch > last.epoch);
}
history.chaos_epochs.push(transition);
}
pub fn insert_schedevent(&mut self, ev: SchedEvent) {
if ev.count > 0 {
self.inner.global.push(ev)
} else {
trace!("NOT recording scheduled event with zero count!");
}
}
pub fn set_current(&mut self, tid: DetTid, new_prio: Priority) {
let history = self.inner.per_thread.get_mut(&tid).unwrap_or_else(|| {
panic!(
"PreemptionRecord: Cannot set current priority before registering thread {}",
tid
)
});
history.final_prio = new_prio;
}
pub fn into_string(mut self) -> String {
self.flushed = true;
self.inner.to_string()
}
pub fn flush(mut self) -> Result<(), String> {
self.flushed = true;
self.write_to_disk()
}
fn write_to_disk(&mut self) -> Result<(), String> {
if let Some(path) = &self.dest {
self.inner.write_to_disk(path)
} else {
Err(
"Cannot write_to_disk because this PreemptionWriter was created without a backing file.".to_string()
)
}
}
}
impl Drop for PreemptionWriter {
fn drop(&mut self) {
if !self.flushed
&& self.dest.is_some()
&& let Err(e) = self.write_to_disk()
{
panic!("Error while dropping PreemptionWriter: {}", e);
}
}
}
#[derive(Debug)]
pub struct PreemptionReader {
inner: PreemptionRecord,
}
pub fn read_trace(path: &Path) -> Vec<SchedEvent> {
let pr = read_preemption_record(path);
pr.global
}
fn read_preemption_record(path: &Path) -> PreemptionRecord {
let string = std::fs::read_to_string(path)
.unwrap_or_else(|e| panic!("Error reading file {:?}:\n {}", path, e));
let pr: PreemptionRecord = serde_json::from_str(&string).unwrap_or_else(|e| {
panic!(
"Error parsing PreemptionRecord from JSON: {}\nJSON contents:\n{}",
e, string
)
});
if let Err(e) = pr.validate() {
panic!(
"Invalid PreemptionRecord when loading from path {}. Error:\n {}",
path.display(),
e
);
}
pr
}
pub fn read_recorded_epoch(path: &Path) -> Result<Option<DateTime<Utc>>, String> {
#[derive(Deserialize)]
struct EpochOnly {
#[serde(default)]
epoch: Option<DateTime<Utc>>,
}
let file = File::open(path)
.map_err(|e| format!("cannot read preemption record {}: {}", path.display(), e))?;
let record: EpochOnly = serde_json::from_reader(std::io::BufReader::new(file))
.map_err(|e| format!("cannot parse preemption record {}: {}", path.display(), e))?;
Ok(record.epoch)
}
impl PreemptionReader {
pub fn new(path: &Path) -> Self {
let pr = read_preemption_record(path);
PreemptionReader { inner: pr }
}
pub fn into_inner(self) -> PreemptionRecord {
self.inner
}
pub fn extract_thread_record(&self, tid: &DetTid) -> Option<ThreadHistory> {
self.inner.per_thread.get(tid).cloned()
}
pub fn thread_initial_priority(&self, tid: &DetTid) -> Option<Priority> {
self.inner.per_thread.get(tid).map(|x| x.initial_priority())
}
pub fn all_threads(&self) -> Vec<DetTid> {
self.inner.per_thread.keys().copied().collect()
}
pub fn load_all(&self) -> PreemptionRecord {
self.inner.clone()
}
pub fn size(&self) -> usize {
let mut sum = 0;
for th in self.inner.per_thread.values() {
sum += th.prio_changes.len()
}
sum
}
}
pub fn strip_times_from_events_file(
sched_path: &Path,
dest: Option<PathBuf>,
) -> anyhow::Result<PathBuf> {
let new_path = dest.unwrap_or_else(|| sched_path.with_extension("notimes"));
let mut preemptions = PreemptionReader::new(sched_path).into_inner();
for se in preemptions.schedevents_iter_mut() {
se.end_time = None;
}
preemptions
.write_to_disk(new_path.as_ref())
.map_err(anyhow::Error::msg)?;
Ok(new_path)
}