use std::sync::{Arc, Condvar, Mutex};
use serde::{Deserialize, Serialize};
use crate::start::Outcome;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum BootFrame {
UnitStarting {
project: String,
service: String,
},
Unit {
project: String,
service: String,
outcome: Outcome,
},
Done {
started: usize,
failed: usize,
},
}
impl BootFrame {
pub fn is_done(&self) -> bool {
matches!(self, BootFrame::Done { .. })
}
}
#[derive(Default)]
struct Inner {
frames: Vec<BootFrame>,
done: bool,
}
#[derive(Clone)]
pub struct BootJournal {
inner: Arc<(Mutex<Inner>, Condvar)>,
}
impl Default for BootJournal {
fn default() -> Self {
Self::new()
}
}
impl BootJournal {
pub fn new() -> Self {
Self {
inner: Arc::new((Mutex::new(Inner::default()), Condvar::new())),
}
}
pub fn push(&self, frame: BootFrame) {
let (lock, cvar) = &*self.inner;
let mut guard = lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if guard.done {
return;
}
if frame.is_done() {
guard.done = true;
}
guard.frames.push(frame);
cvar.notify_all();
}
pub fn record(&self, project: &str, service: &str, outcome: Outcome) {
self.push(BootFrame::Unit {
project: project.to_string(),
service: service.to_string(),
outcome,
});
}
pub fn is_done(&self) -> bool {
self.inner
.0
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.done
}
pub fn snapshot(&self) -> Vec<BootFrame> {
self.inner
.0
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.frames
.clone()
}
pub fn wait_from(&self, from: usize) -> Vec<BootFrame> {
let (lock, cvar) = &*self.inner;
let mut guard = lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
while guard.frames.len() <= from && !guard.done {
guard = cvar
.wait(guard)
.unwrap_or_else(std::sync::PoisonError::into_inner);
}
guard.frames.get(from..).unwrap_or(&[]).to_vec()
}
}
#[cfg(test)]
mod tests {
use std::thread;
use super::*;
use crate::start::Liveness;
fn up(service: &str) -> BootFrame {
BootFrame::Unit {
project: "p".into(),
service: service.into(),
outcome: Outcome::Up(Liveness { pid: 1 }),
}
}
#[test]
fn snapshot_replays_all_recorded_frames() {
let j = BootJournal::new();
j.push(up("a"));
j.push(up("b"));
j.push(BootFrame::Done {
started: 2,
failed: 0,
});
let snap = j.snapshot();
assert_eq!(snap.len(), 3);
assert!(snap[2].is_done());
assert!(j.is_done());
}
#[test]
fn push_after_done_is_ignored() {
let j = BootJournal::new();
j.push(BootFrame::Done {
started: 0,
failed: 0,
});
j.push(up("late"));
assert_eq!(j.snapshot().len(), 1);
}
#[test]
fn wait_from_blocks_until_a_new_frame_arrives() {
let j = BootJournal::new();
let producer = j.clone();
let handle = thread::spawn(move || {
let first = producer.clone();
first.push(up("a"));
first.push(BootFrame::Done {
started: 1,
failed: 0,
});
});
let mut seen = 0;
let mut all = Vec::new();
loop {
let batch = j.wait_from(seen);
seen += batch.len();
let done = batch.iter().any(BootFrame::is_done);
all.extend(batch);
if done {
break;
}
}
handle.join().unwrap();
assert!(
all.iter()
.any(|f| matches!(f, BootFrame::Unit { service, .. } if service == "a"))
);
assert!(all.last().unwrap().is_done());
}
}