use super::door::require_dev_node;
use super::{
CandidateBytes, DevReply, DevRequest, DevStatusEvent, GenerationReport, NodeMode, ReloadStage,
StageRefusal,
};
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum StartFailure {
Load(String),
Tree(String),
}
pub trait NodeControl {
fn stop_component(&mut self) -> Result<(), String>;
fn stage(&mut self, candidate: &CandidateBytes) -> Result<(), String>;
fn start_component(&mut self) -> Result<(), StartFailure>;
fn witness_mailbox(&mut self) -> Result<(), String>;
fn witness_content(&mut self) -> Result<(), String>;
fn report(
&mut self,
build_generation: u64,
content_digest: [u8; 32],
) -> Result<GenerationReport, String>;
}
pub struct BarrierEngine {
mode: NodeMode,
last_good: CandidateBytes,
newest_seen: u64,
in_flight: bool,
}
impl BarrierEngine {
#[must_use]
pub fn new(mode: NodeMode, boot: CandidateBytes) -> Self {
Self {
mode,
last_good: boot,
newest_seen: 0,
in_flight: false,
}
}
pub fn handle(
&mut self,
request: DevRequest,
node: &mut dyn NodeControl,
status: &mut dyn FnMut(DevStatusEvent),
) -> DevReply {
let DevRequest::StageAndActivate(candidate) = request;
if let Err(refusal) = require_dev_node(self.mode) {
return DevReply::Refused(refusal);
}
if self.in_flight {
return DevReply::Refused(StageRefusal::ActivationInFlight);
}
if candidate.build_generation <= self.newest_seen {
return DevReply::Refused(StageRefusal::StaleGeneration {
offered: candidate.build_generation,
newest_seen: self.newest_seen,
});
}
self.newest_seen = candidate.build_generation;
self.in_flight = true;
let reply = self.activate(&candidate, node, status);
self.in_flight = false;
reply
}
fn activate(
&mut self,
candidate: &CandidateBytes,
node: &mut dyn NodeControl,
status: &mut dyn FnMut(DevStatusEvent),
) -> DevReply {
let generation = candidate.build_generation;
let mut stage = |at: ReloadStage| {
status(DevStatusEvent::Reloading {
build_generation: generation,
stage: at,
});
};
stage(ReloadStage::Drain);
stage(ReloadStage::StopOld);
if let Err(detail) = node.stop_component() {
return self.roll_back(ReloadStage::StopOld, detail, node, status);
}
stage(ReloadStage::PurgeOld);
if let Err(detail) = node.stage(candidate) {
return self.roll_back(ReloadStage::PurgeOld, detail, node, status);
}
stage(ReloadStage::StartFresh);
match node.start_component() {
Ok(()) => {}
Err(StartFailure::Load(detail)) => {
return self.roll_back(ReloadStage::LoadCandidate, detail, node, status);
}
Err(StartFailure::Tree(detail)) => {
return self.roll_back(ReloadStage::StartFresh, detail, node, status);
}
}
stage(ReloadStage::LivenessMailbox);
if let Err(detail) = node.witness_mailbox() {
return self.roll_back_live(ReloadStage::LivenessMailbox, detail, node, status);
}
stage(ReloadStage::LivenessContent);
if let Err(detail) = node.witness_content() {
return self.roll_back_live(ReloadStage::LivenessContent, detail, node, status);
}
self.last_good = candidate.clone();
match node.report(generation, candidate.content_digest) {
Ok(report) => {
status(DevStatusEvent::RunningCurrent(report.clone()));
DevReply::Activated(report)
}
Err(detail) => node_failed(ReloadStage::Report, detail, status),
}
}
fn roll_back_live(
&mut self,
failed: ReloadStage,
detail: String,
node: &mut dyn NodeControl,
status: &mut dyn FnMut(DevStatusEvent),
) -> DevReply {
if let Err(stop_detail) = node.stop_component() {
return node_failed(
failed,
format!(
"{detail}; stopping the unproven candidate tree also failed: {stop_detail}"
),
status,
);
}
self.roll_back(failed, detail, node, status)
}
fn roll_back(
&mut self,
failed: ReloadStage,
detail: String,
node: &mut dyn NodeControl,
status: &mut dyn FnMut(DevStatusEvent),
) -> DevReply {
let restore = self.restore_last_good(node);
match restore {
Ok(()) => {
let serving = match node.report(
self.last_good.build_generation,
self.last_good.content_digest,
) {
Ok(serving) => serving,
Err(report_detail) => {
return node_failed(
failed,
format!("{detail}; last-good witnesses unreadable: {report_detail}"),
status,
);
}
};
status(DevStatusEvent::ReloadFailed {
failed_generation: self.newest_seen,
failed,
detail: detail.clone(),
serving: serving.clone(),
});
DevReply::RolledBack {
failed,
detail,
serving,
}
}
Err(restore_detail) => node_failed(
failed,
format!("{detail}; restoring last good failed: {restore_detail}"),
status,
),
}
}
fn restore_last_good(&mut self, node: &mut dyn NodeControl) -> Result<(), String> {
node.stage(&self.last_good)
.map_err(|error| format!("restore stage: {error}"))?;
node.start_component().map_err(|error| match error {
StartFailure::Load(detail) => format!("restore load: {detail}"),
StartFailure::Tree(detail) => format!("restore start: {detail}"),
})?;
node.witness_mailbox()
.map_err(|error| format!("restore mailbox witness: {error}"))?;
node.witness_content()
.map_err(|error| format!("restore content witness: {error}"))
}
}
fn node_failed(
failed: ReloadStage,
detail: String,
status: &mut dyn FnMut(DevStatusEvent),
) -> DevReply {
status(DevStatusEvent::NodeFailed {
detail: detail.clone(),
});
DevReply::NodeFailed { failed, detail }
}
#[cfg(test)]
mod tests {
#![allow(
clippy::expect_used,
clippy::panic,
clippy::struct_excessive_bools,
clippy::cast_possible_truncation
)]
use super::{BarrierEngine, NodeControl, StartFailure};
use crate::dev::{
CandidateBytes, DevReply, DevRequest, DevStatusEvent, GenerationReport, NodeMode,
ReloadStage, StageRefusal,
};
#[derive(Default)]
struct FakeNode {
calls: Vec<String>,
fail_stop: bool,
fail_stage: bool,
fail_start: Option<StartFailure>,
fail_mailbox: bool,
fail_content: bool,
failures_left: usize,
}
impl FakeNode {
fn failing(field: fn(&mut FakeNode), failures: usize) -> Self {
let mut node = Self {
failures_left: failures,
..Self::default()
};
field(&mut node);
node
}
fn consume_failure(&mut self) -> bool {
if self.failures_left > 0 {
self.failures_left -= 1;
true
} else {
false
}
}
}
impl NodeControl for FakeNode {
fn stop_component(&mut self) -> Result<(), String> {
self.calls.push("stop".to_owned());
if self.fail_stop && self.consume_failure() {
return Err("stop refused".to_owned());
}
Ok(())
}
fn stage(&mut self, candidate: &CandidateBytes) -> Result<(), String> {
self.calls
.push(format!("stage:{}", candidate.build_generation));
if self.fail_stage && self.consume_failure() {
return Err("purge refused: still referenced".to_owned());
}
Ok(())
}
fn start_component(&mut self) -> Result<(), StartFailure> {
self.calls.push("start".to_owned());
if let Some(failure) = self.fail_start.clone()
&& self.consume_failure()
{
return Err(failure);
}
Ok(())
}
fn witness_mailbox(&mut self) -> Result<(), String> {
self.calls.push("mailbox".to_owned());
if self.fail_mailbox && self.consume_failure() {
return Err("mailbox witness failed".to_owned());
}
Ok(())
}
fn witness_content(&mut self) -> Result<(), String> {
self.calls.push("content".to_owned());
if self.fail_content && self.consume_failure() {
return Err("content witness failed".to_owned());
}
Ok(())
}
fn report(
&mut self,
build_generation: u64,
content_digest: [u8; 32],
) -> Result<GenerationReport, String> {
Ok(GenerationReport {
build_generation,
content_digest,
component_incarnation: 1,
module_generation: 1,
})
}
}
fn candidate(generation: u64) -> CandidateBytes {
CandidateBytes {
build_generation: generation,
content_digest: [generation as u8; 32],
component_beam: vec![generation as u8],
ffi_beam: vec![generation as u8, 0xff],
}
}
fn engine(mode: NodeMode) -> BarrierEngine {
BarrierEngine::new(mode, candidate(0))
}
fn drive(
engine: &mut BarrierEngine,
node: &mut FakeNode,
generation: u64,
) -> (DevReply, Vec<DevStatusEvent>) {
let mut events = Vec::new();
let reply = engine.handle(
DevRequest::StageAndActivate(candidate(generation)),
node,
&mut |event| events.push(event),
);
(reply, events)
}
#[test]
fn production_engine_refuses_with_zero_node_calls() {
let mut node = FakeNode::default();
let (reply, events) = drive(&mut engine(NodeMode::Production), &mut node, 1);
assert_eq!(reply, DevReply::Refused(StageRefusal::NotADevNode));
assert!(node.calls.is_empty(), "no path from request to the node");
assert!(
events.is_empty(),
"no status invented for a refused request"
);
}
#[test]
fn activation_runs_the_barrier_in_order() {
let mut node = FakeNode::default();
let mut engine = engine(NodeMode::Dev);
let (reply, events) = drive(&mut engine, &mut node, 1);
assert_eq!(
node.calls,
vec!["stop", "stage:1", "start", "mailbox", "content"]
);
match reply {
DevReply::Activated(report) => assert_eq!(report.build_generation, 1),
other => panic!("expected Activated, got {other:?}"),
}
assert!(matches!(
events.last(),
Some(DevStatusEvent::RunningCurrent(report)) if report.build_generation == 1
));
}
#[test]
fn stage_failures_roll_back_and_reprove_last_good() {
type Inject = fn(&mut FakeNode);
let scripted: [(Inject, ReloadStage); 3] = [
(|n| n.fail_stage = true, ReloadStage::PurgeOld),
(
|n| n.fail_start = Some(StartFailure::Load("bad beam".to_owned())),
ReloadStage::LoadCandidate,
),
(
|n| n.fail_start = Some(StartFailure::Tree("tree died".to_owned())),
ReloadStage::StartFresh,
),
];
for (inject, expected_stage) in scripted {
let mut node = FakeNode::failing(inject, 1);
let mut engine = engine(NodeMode::Dev);
let (reply, events) = drive(&mut engine, &mut node, 1);
match reply {
DevReply::RolledBack {
failed, serving, ..
} => {
assert_eq!(failed, expected_stage);
assert_eq!(serving.build_generation, 0, "last-good serves");
}
other => panic!("expected RolledBack at {expected_stage:?}, got {other:?}"),
}
let tail: Vec<_> = node.calls.iter().rev().take(4).rev().cloned().collect();
assert_eq!(
tail,
vec!["stage:0", "start", "mailbox", "content"],
"restore re-proves last good at {expected_stage:?}"
);
assert!(matches!(
events.last(),
Some(DevStatusEvent::ReloadFailed { serving, .. })
if serving.build_generation == 0
));
}
}
#[test]
fn witness_failure_stops_the_unproven_tree_before_restore() {
let mut node = FakeNode::failing(|n| n.fail_content = true, 1);
let mut engine = engine(NodeMode::Dev);
let (reply, _events) = drive(&mut engine, &mut node, 1);
match reply {
DevReply::RolledBack { failed, .. } => {
assert_eq!(failed, ReloadStage::LivenessContent);
}
other => panic!("expected RolledBack, got {other:?}"),
}
assert_eq!(
node.calls,
vec![
"stop", "stage:1", "start", "mailbox", "content", "stop", "stage:0", "start", "mailbox", "content" ]
);
}
#[test]
fn restore_failure_is_node_failed() {
let mut node = FakeNode::failing(|n| n.fail_stage = true, 2);
let mut engine = engine(NodeMode::Dev);
let (reply, events) = drive(&mut engine, &mut node, 1);
match reply {
DevReply::NodeFailed { failed, detail } => {
assert_eq!(failed, ReloadStage::PurgeOld);
assert!(detail.contains("restoring last good failed"), "{detail}");
}
other => panic!("expected NodeFailed, got {other:?}"),
}
assert!(matches!(
events.last(),
Some(DevStatusEvent::NodeFailed { .. })
));
}
#[test]
fn stale_generations_are_refused_by_identity() {
let mut node = FakeNode::failing(|n| n.fail_stage = true, 1);
let mut engine = engine(NodeMode::Dev);
let (reply, _) = drive(&mut engine, &mut node, 3);
assert!(matches!(reply, DevReply::RolledBack { .. }));
let calls_before = node.calls.len();
let (reply, _) = drive(&mut engine, &mut node, 3);
assert_eq!(
reply,
DevReply::Refused(StageRefusal::StaleGeneration {
offered: 3,
newest_seen: 3
})
);
assert_eq!(node.calls.len(), calls_before);
let (reply, _) = drive(&mut engine, &mut node, 4);
assert!(matches!(reply, DevReply::Activated(_)));
}
}