use crate::model::{
artifacts::ArtifactChecksumRecord,
attempt_journal::{
AttemptAuthorityRecord, AttemptJournalRecord, AttemptJournalView, OperationBindingRecord,
},
effect_graph::MAX_EFFECT_OPERATIONS,
operation_plan::{OperationPlanRecord, PlanContextRecord, PlannedOperationRecord},
};
use serde::Serialize;
use std::collections::{BTreeMap, BTreeSet};
use thiserror::Error;
#[derive(Clone, Debug)]
pub struct ExecutionProgressRequest<'a> {
pub plan: &'a OperationPlanRecord,
pub journals: &'a [&'a AttemptJournalRecord],
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum OperationProgressState {
AwaitingDependencies,
MutationAvailable,
MutationExhausted,
MutationUnresolved,
ObservationUnresolved,
ReconciliationExhausted,
Applied,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
pub struct OperationProgressView {
pub operation_sequence: u64,
pub state: OperationProgressState,
pub attempts: AttemptJournalView,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
pub struct ExecutionAttemptTotalsView {
pub mutations_used: u32,
pub observations_used: u32,
pub mutations_remaining: u32,
pub observations_remaining: u32,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
pub struct ExecutionProgressView {
pub intent: ArtifactChecksumRecord,
pub graph: ArtifactChecksumRecord,
pub applied_operations: usize,
pub attempts: ExecutionAttemptTotalsView,
pub operations: Vec<OperationProgressView>,
}
pub fn progress(
request: &ExecutionProgressRequest<'_>,
) -> Result<ExecutionProgressView, ExecutionProgressError> {
if request.journals.len() > MAX_EFFECT_OPERATIONS {
return Err(ExecutionProgressError::TooManyJournals);
}
let plan = request.plan;
let intent = plan.digest();
let mut journals = BTreeMap::new();
for journal in request.journals {
let sequence = journal.authority().binding().operation_sequence();
let operation = plan
.operation(sequence)
.map_err(|_| ExecutionProgressError::UnknownOperation(sequence))?;
if journals.contains_key(&sequence) {
return Err(ExecutionProgressError::DuplicateJournal(sequence));
}
let original = OriginalJournalBinding {
intent: &intent,
context: plan.context(),
operation,
};
if !original.matches(journal.authority()) {
return Err(ExecutionProgressError::AuthorityMismatch(sequence));
}
journals.insert(sequence, journal.view());
}
for operation in plan.operations() {
if !journals.contains_key(&operation.operation_sequence()) {
return Err(ExecutionProgressError::MissingJournal(
operation.operation_sequence(),
));
}
}
let applied: BTreeSet<_> = journals
.iter()
.filter_map(|(sequence, attempts)| attempts.applied.then_some(*sequence))
.collect();
let mut totals = ExecutionAttemptTotalsView {
mutations_used: 0,
observations_used: 0,
mutations_remaining: 0,
observations_remaining: 0,
};
let mut operations = Vec::with_capacity(journals.len());
for node in plan.graph().ordered_nodes() {
let sequence = node.operation_sequence();
let attempts = journals
.remove(&sequence)
.ok_or(ExecutionProgressError::MissingJournal(sequence))?;
let unmet = node
.depends_on()
.iter()
.find(|dependency| !applied.contains(dependency));
if attempts.mutations_used != 0
&& let Some(prerequisite) = unmet
{
return Err(ExecutionProgressError::PrematureAttempt {
operation_sequence: sequence,
prerequisite: *prerequisite,
});
}
totals.add(&attempts)?;
operations.push(OperationProgressView {
operation_sequence: sequence,
state: condition(&attempts, unmet.is_some()),
attempts,
});
}
Ok(ExecutionProgressView {
intent,
graph: plan.graph().digest(),
applied_operations: applied.len(),
attempts: totals,
operations,
})
}
struct OriginalJournalBinding<'a> {
intent: &'a ArtifactChecksumRecord,
context: &'a PlanContextRecord,
operation: &'a PlannedOperationRecord,
}
impl OriginalJournalBinding<'_> {
fn matches(&self, authority: &AttemptAuthorityRecord) -> bool {
self.identity_matches(authority.binding())
&& self.context_matches(authority.binding())
&& authority.budget() == self.operation.budget()
}
fn identity_matches(&self, binding: &OperationBindingRecord) -> bool {
binding.intent() == self.intent.hash()
&& binding.operation_sequence() == self.operation.operation_sequence()
&& binding.target() == self.operation.target()
&& binding.request() == self.operation.request()
}
fn context_matches(&self, binding: &OperationBindingRecord) -> bool {
binding.network() == self.context.network()
&& binding.caller() == self.context.caller()
&& binding.release() == self.context.release()
}
}
impl ExecutionAttemptTotalsView {
fn add(&mut self, attempts: &AttemptJournalView) -> Result<(), ExecutionProgressError> {
self.mutations_used = sum(self.mutations_used, attempts.mutations_used)?;
self.observations_used = sum(self.observations_used, attempts.observations_used)?;
self.mutations_remaining = sum(self.mutations_remaining, attempts.mutations_remaining)?;
self.observations_remaining =
sum(self.observations_remaining, attempts.observations_remaining)?;
Ok(())
}
}
fn sum(left: u32, right: u32) -> Result<u32, ExecutionProgressError> {
left.checked_add(right)
.ok_or(ExecutionProgressError::AccountingOverflow)
}
fn condition(attempts: &AttemptJournalView, unmet_dependencies: bool) -> OperationProgressState {
if attempts.applied {
OperationProgressState::Applied
} else if attempts.pending_observation.is_some() {
OperationProgressState::ObservationUnresolved
} else if attempts.pending_mutation.is_some() {
if attempts.observations_remaining == 0 {
OperationProgressState::ReconciliationExhausted
} else {
OperationProgressState::MutationUnresolved
}
} else if unmet_dependencies {
OperationProgressState::AwaitingDependencies
} else if attempts.mutations_remaining == 0 {
OperationProgressState::MutationExhausted
} else {
OperationProgressState::MutationAvailable
}
}
#[derive(Debug, Error, Eq, PartialEq)]
pub enum ExecutionProgressError {
#[error("execution progress exceeds {MAX_EFFECT_OPERATIONS} journals")]
TooManyJournals,
#[error("duplicate journal for operation {0}")]
DuplicateJournal(u64),
#[error("journal operation {0} is absent from the original plan")]
UnknownOperation(u64),
#[error("missing original journal for operation {0}")]
MissingJournal(u64),
#[error("original journal authority mismatch for operation {0}")]
AuthorityMismatch(u64),
#[error(
"operation {operation_sequence} was attempted without applied prerequisite {prerequisite}"
)]
PrematureAttempt {
operation_sequence: u64,
prerequisite: u64,
},
#[error("execution attempt accounting overflow")]
AccountingOverflow,
}
#[cfg(test)]
mod tests;