use super::{
AttemptJournalError, AttemptJournalGuard, BackupLayoutGuard, ExecutionProgressPersistenceError,
ExecutionSettlementCheckpointError, ExecutionSettlementPersistenceError, JournalLock,
JournalLockError, OperationPlanPersistenceError, PersistenceError,
checkpoint_execution_settlement, create_json_durable, create_operation_plan,
read_execution_progress, read_execution_settlement, read_json, read_operation_plan,
};
use crate::model::{
artifacts::ArtifactChecksumRecord,
execution_workflow::{
ExecutionStageBindingRecord, ExecutionStagePredecessorRecord, ExecutionWorkflowError,
ExecutionWorkflowRecord, MAX_EXECUTION_STAGE_BYTES, MAX_EXECUTION_WORKFLOW_BYTES,
},
operation_plan::{OperationPlanError, OperationPlanRecord},
};
use crate::policy::execution_progress::ExecutionProgressView;
use std::{
collections::{BTreeMap, btree_map::Entry},
fs,
path::PathBuf,
};
use thiserror::Error;
const WORKFLOW_FILE: &str = "execution-workflow.json";
const BINDING_FILE: &str = "stage-binding.json";
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum StagePreparationBarrier {
Plan,
Binding,
Journal(u64),
}
pub fn create_execution_workflow(
layout: &BackupLayoutGuard,
record: &ExecutionWorkflowRecord,
) -> Result<(), ExecutionWorkflowPersistenceError> {
layout.check_root()?;
super::json::check_json_size(record, MAX_EXECUTION_WORKFLOW_BYTES)?;
let path = layout.root().join(WORKFLOW_FILE);
let _lock = JournalLock::acquire(&path)?;
create_json_durable(&path, record)?;
Ok(())
}
pub fn read_execution_workflow(
layout: &BackupLayoutGuard,
expected: &ArtifactChecksumRecord,
) -> Result<ExecutionWorkflowRecord, ExecutionWorkflowPersistenceError> {
layout.check_root()?;
let path = layout.root().join(WORKFLOW_FILE);
let _lock = JournalLock::acquire(&path)?;
let record: ExecutionWorkflowRecord = read_json(&path, MAX_EXECUTION_WORKFLOW_BYTES)?;
super::json::check_json_size(&record, MAX_EXECUTION_WORKFLOW_BYTES)?;
if &record.digest() != expected {
return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
}
Ok(record)
}
#[derive(Debug)]
pub struct ExecutionStageGuard<'a> {
workflow_layout: &'a BackupLayoutGuard,
stage_layout: BackupLayoutGuard,
binding: ExecutionStageBindingRecord,
plan: OperationPlanRecord,
}
impl<'a> ExecutionStageGuard<'a> {
pub fn checkpoint(
&self,
learned_evidence: ArtifactChecksumRecord,
) -> Result<ExecutionStagePredecessorRecord, ExecutionStageCheckpointError> {
self.checkpoint_with(learned_evidence, || {})
}
fn checkpoint_with(
&self,
learned_evidence: ArtifactChecksumRecord,
after_publication: impl FnOnce(),
) -> Result<ExecutionStagePredecessorRecord, ExecutionStageCheckpointError> {
let settlement = checkpoint_execution_settlement(self.layout()?, &self.plan.digest())?;
after_publication();
self.layout()?;
Ok(ExecutionStagePredecessorRecord::new(
self.binding.stage_sequence(),
self.binding.digest(),
settlement.digest(),
learned_evidence,
))
}
pub fn resume(
workflow_layout: &'a BackupLayoutGuard,
expected_workflow: &ArtifactChecksumRecord,
sequence: u64,
expected_binding: &ArtifactChecksumRecord,
) -> Result<(Self, ExecutionProgressView), ExecutionStageResumeError> {
let stage = Self::open(
workflow_layout,
expected_workflow,
sequence,
expected_binding,
)?;
let progress = read_execution_progress(&stage.stage_layout, &stage.plan.digest())?;
stage.layout()?;
Ok((stage, progress))
}
pub fn prepare(
workflow_layout: &'a BackupLayoutGuard,
binding: ExecutionStageBindingRecord,
plan: OperationPlanRecord,
) -> Result<Self, ExecutionStagePreparationError> {
Self::prepare_with(workflow_layout, binding, plan, |_| {})
}
fn prepare_with(
workflow_layout: &'a BackupLayoutGuard,
binding: ExecutionStageBindingRecord,
plan: OperationPlanRecord,
mut barrier: impl FnMut(StagePreparationBarrier),
) -> Result<Self, ExecutionStagePreparationError> {
let authorities = plan.attempt_authorities()?;
let stage = Self::create_with(workflow_layout, binding, plan, &mut barrier)?;
let layout = stage.layout()?;
for authority in authorities {
let sequence = authority.binding().operation_sequence();
drop(AttemptJournalGuard::create(layout, authority)?);
barrier(StagePreparationBarrier::Journal(sequence));
}
stage.layout()?;
Ok(stage)
}
pub fn create(
workflow_layout: &'a BackupLayoutGuard,
binding: ExecutionStageBindingRecord,
plan: OperationPlanRecord,
) -> Result<Self, ExecutionWorkflowPersistenceError> {
Self::create_with(workflow_layout, binding, plan, |_| {})
}
fn create_with(
workflow_layout: &'a BackupLayoutGuard,
binding: ExecutionStageBindingRecord,
plan: OperationPlanRecord,
mut barrier: impl FnMut(StagePreparationBarrier),
) -> Result<Self, ExecutionWorkflowPersistenceError> {
let workflow = read_execution_workflow(workflow_layout, binding.workflow())?;
binding.validate(&workflow, &plan)?;
super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
super::json::check_json_size(
&plan,
crate::model::operation_plan::MAX_OPERATION_PLAN_BYTES,
)?;
validate_predecessors(workflow_layout, &workflow, &binding)?;
let path = stage_path(workflow_layout, binding.stage_sequence());
create_private_directory(&path)?;
fs::File::open(workflow_layout.root())
.and_then(|directory| directory.sync_all())
.map_err(PersistenceError::from)?;
let stage_layout = acquire_stage(workflow_layout, binding.stage_sequence())?;
create_operation_plan(&stage_layout, &plan)?;
barrier(StagePreparationBarrier::Plan);
create_json_durable(&stage_layout.root().join(BINDING_FILE), &binding)?;
barrier(StagePreparationBarrier::Binding);
workflow_layout.check_root()?;
Ok(Self {
workflow_layout,
stage_layout,
binding,
plan,
})
}
pub fn open(
workflow_layout: &'a BackupLayoutGuard,
expected_workflow: &ArtifactChecksumRecord,
sequence: u64,
expected_binding: &ArtifactChecksumRecord,
) -> Result<Self, ExecutionWorkflowPersistenceError> {
let workflow = read_execution_workflow(workflow_layout, expected_workflow)?;
workflow
.stage(sequence)
.map_err(|_| ExecutionWorkflowError::UnknownStage(sequence))?;
let stage_layout = acquire_stage(workflow_layout, sequence)?;
let (binding, plan) = read_binding(&stage_layout, &workflow, sequence, expected_binding)?;
validate_predecessors(workflow_layout, &workflow, &binding)?;
workflow_layout.check_root()?;
Ok(Self {
workflow_layout,
stage_layout,
binding,
plan,
})
}
#[must_use]
pub const fn plan(&self) -> &OperationPlanRecord {
&self.plan
}
#[must_use]
pub const fn binding(&self) -> &ExecutionStageBindingRecord {
&self.binding
}
pub fn layout(&self) -> Result<&BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
let workflow = read_execution_workflow(self.workflow_layout, self.binding.workflow())?;
read_binding(
&self.stage_layout,
&workflow,
self.binding.stage_sequence(),
&self.binding.digest(),
)?;
validate_predecessors(self.workflow_layout, &workflow, &self.binding)?;
Ok(&self.stage_layout)
}
}
#[derive(Debug, Error)]
pub enum ExecutionStageCheckpointError {
#[error(transparent)]
Stage(#[from] ExecutionWorkflowPersistenceError),
#[error(transparent)]
Settlement(#[from] ExecutionSettlementCheckpointError),
}
#[derive(Debug, Error)]
pub enum ExecutionStagePreparationError {
#[error(transparent)]
Stage(#[from] ExecutionWorkflowPersistenceError),
#[error(transparent)]
Authority(#[from] OperationPlanError),
#[error(transparent)]
Journal(#[from] AttemptJournalError),
}
#[derive(Debug, Error)]
pub enum ExecutionStageResumeError {
#[error(transparent)]
Stage(#[from] ExecutionWorkflowPersistenceError),
#[error(transparent)]
Progress(#[from] ExecutionProgressPersistenceError),
}
fn read_binding(
layout: &BackupLayoutGuard,
workflow: &ExecutionWorkflowRecord,
sequence: u64,
expected: &ArtifactChecksumRecord,
) -> Result<(ExecutionStageBindingRecord, OperationPlanRecord), ExecutionWorkflowPersistenceError> {
layout.check_root()?;
let binding: ExecutionStageBindingRecord =
read_json(&layout.root().join(BINDING_FILE), MAX_EXECUTION_STAGE_BYTES)?;
super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
if binding.stage_sequence() != sequence || &binding.digest() != expected {
return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
}
let plan = read_operation_plan(layout, binding.plan())?;
binding.validate(workflow, &plan)?;
Ok((binding, plan))
}
fn validate_predecessors(
layout: &BackupLayoutGuard,
workflow: &ExecutionWorkflowRecord,
binding: &ExecutionStageBindingRecord,
) -> Result<(), ExecutionWorkflowPersistenceError> {
let mut expected: BTreeMap<u64, ExecutionStagePredecessorRecord> = binding
.predecessors()
.iter()
.map(|row| (row.stage_sequence(), row.clone()))
.collect();
let mut pending: Vec<_> = expected.keys().copied().collect();
while let Some(sequence) = pending.pop() {
let row = &expected[&sequence];
let original = {
let predecessor = acquire_stage(layout, sequence)?;
let (original, _) = read_binding(&predecessor, workflow, sequence, row.binding())?;
read_execution_settlement(&predecessor, original.plan(), row.settlement())?;
original
};
for ancestor in original.predecessors() {
match expected.entry(ancestor.stage_sequence()) {
Entry::Vacant(entry) => {
pending.push(ancestor.stage_sequence());
entry.insert(ancestor.clone());
}
Entry::Occupied(entry) => {
if entry.get().binding() != ancestor.binding()
|| entry.get().settlement() != ancestor.settlement()
{
return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
}
}
}
}
}
Ok(())
}
fn stage_path(layout: &BackupLayoutGuard, sequence: u64) -> PathBuf {
layout.root().join(format!("execution-stage-{sequence}"))
}
fn acquire_stage(
layout: &BackupLayoutGuard,
sequence: u64,
) -> Result<BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
layout.check_root()?;
let path = stage_path(layout, sequence);
let metadata = fs::symlink_metadata(&path).map_err(PersistenceError::from)?;
if !metadata.is_dir() {
return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
}
let stage = BackupLayoutGuard::acquire(&path)?;
if stage.root() != path {
return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
}
layout.check_root()?;
Ok(stage)
}
fn create_private_directory(
path: &std::path::Path,
) -> Result<(), ExecutionWorkflowPersistenceError> {
#[cfg(unix)]
{
use std::os::unix::fs::DirBuilderExt;
fs::DirBuilder::new()
.mode(0o700)
.create(path)
.map_err(PersistenceError::from)?;
Ok(())
}
#[cfg(not(unix))]
{
let _ = path;
Err(PersistenceError::Io(std::io::Error::from(std::io::ErrorKind::Unsupported)).into())
}
}
#[derive(Debug, Error)]
pub enum ExecutionWorkflowPersistenceError {
#[error("execution workflow/stage digest mismatch")]
DigestMismatch,
#[error("unsafe execution stage directory")]
UnsafeStage,
#[error(transparent)]
Model(#[from] ExecutionWorkflowError),
#[error(transparent)]
Plan(#[from] OperationPlanPersistenceError),
#[error(transparent)]
Settlement(#[from] ExecutionSettlementPersistenceError),
#[error(transparent)]
Lock(#[from] JournalLockError),
#[error(transparent)]
Persistence(#[from] PersistenceError),
}
#[cfg(all(test, unix))]
mod tests;