ic_backup/policy/execution_progress/
mod.rs1use crate::model::{
4 artifacts::ArtifactChecksumRecord,
5 attempt_journal::{
6 AttemptAuthorityRecord, AttemptJournalRecord, AttemptJournalView, OperationBindingRecord,
7 },
8 effect_graph::MAX_EFFECT_OPERATIONS,
9 operation_plan::{OperationPlanRecord, PlanContextRecord, PlannedOperationRecord},
10};
11use serde::Serialize;
12use std::collections::{BTreeMap, BTreeSet};
13use thiserror::Error;
14
15#[derive(Clone, Debug)]
17pub struct ExecutionProgressRequest<'a> {
18 pub plan: &'a OperationPlanRecord,
20 pub journals: &'a [&'a AttemptJournalRecord],
22}
23
24#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
26#[serde(rename_all = "snake_case")]
27pub enum OperationProgressState {
28 AwaitingDependencies,
30 MutationAvailable,
32 MutationExhausted,
34 MutationUnresolved,
36 ObservationUnresolved,
38 ReconciliationExhausted,
40 Applied,
42}
43
44#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
46pub struct OperationProgressView {
47 pub operation_sequence: u64,
49 pub state: OperationProgressState,
51 pub attempts: AttemptJournalView,
53}
54
55#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
57pub struct ExecutionAttemptTotalsView {
58 pub mutations_used: u32,
60 pub observations_used: u32,
62 pub mutations_remaining: u32,
64 pub observations_remaining: u32,
66}
67
68#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
73pub struct ExecutionProgressView {
74 pub intent: ArtifactChecksumRecord,
76 pub graph: ArtifactChecksumRecord,
78 pub applied_operations: usize,
80 pub attempts: ExecutionAttemptTotalsView,
82 pub operations: Vec<OperationProgressView>,
84}
85
86pub fn progress(
97 request: &ExecutionProgressRequest<'_>,
98) -> Result<ExecutionProgressView, ExecutionProgressError> {
99 if request.journals.len() > MAX_EFFECT_OPERATIONS {
100 return Err(ExecutionProgressError::TooManyJournals);
101 }
102 let plan = request.plan;
103 let intent = plan.digest();
104 let mut journals = BTreeMap::new();
105 for journal in request.journals {
106 let sequence = journal.authority().binding().operation_sequence();
107 let operation = plan
108 .operation(sequence)
109 .map_err(|_| ExecutionProgressError::UnknownOperation(sequence))?;
110 if journals.contains_key(&sequence) {
111 return Err(ExecutionProgressError::DuplicateJournal(sequence));
112 }
113 let original = OriginalJournalBinding {
114 intent: &intent,
115 context: plan.context(),
116 operation,
117 };
118 if !original.matches(journal.authority()) {
119 return Err(ExecutionProgressError::AuthorityMismatch(sequence));
120 }
121 journals.insert(sequence, journal.view());
122 }
123 for operation in plan.operations() {
124 if !journals.contains_key(&operation.operation_sequence()) {
125 return Err(ExecutionProgressError::MissingJournal(
126 operation.operation_sequence(),
127 ));
128 }
129 }
130 let applied: BTreeSet<_> = journals
131 .iter()
132 .filter_map(|(sequence, attempts)| attempts.applied.then_some(*sequence))
133 .collect();
134 let mut totals = ExecutionAttemptTotalsView {
135 mutations_used: 0,
136 observations_used: 0,
137 mutations_remaining: 0,
138 observations_remaining: 0,
139 };
140 let mut operations = Vec::with_capacity(journals.len());
141 for node in plan.graph().ordered_nodes() {
142 let sequence = node.operation_sequence();
143 let attempts = journals
145 .remove(&sequence)
146 .ok_or(ExecutionProgressError::MissingJournal(sequence))?;
147 let unmet = node
148 .depends_on()
149 .iter()
150 .find(|dependency| !applied.contains(dependency));
151 if attempts.mutations_used != 0
152 && let Some(prerequisite) = unmet
153 {
154 return Err(ExecutionProgressError::PrematureAttempt {
155 operation_sequence: sequence,
156 prerequisite: *prerequisite,
157 });
158 }
159 totals.add(&attempts)?;
160 operations.push(OperationProgressView {
161 operation_sequence: sequence,
162 state: condition(&attempts, unmet.is_some()),
163 attempts,
164 });
165 }
166 Ok(ExecutionProgressView {
167 intent,
168 graph: plan.graph().digest(),
169 applied_operations: applied.len(),
170 attempts: totals,
171 operations,
172 })
173}
174
175struct OriginalJournalBinding<'a> {
176 intent: &'a ArtifactChecksumRecord,
177 context: &'a PlanContextRecord,
178 operation: &'a PlannedOperationRecord,
179}
180impl OriginalJournalBinding<'_> {
181 fn matches(&self, authority: &AttemptAuthorityRecord) -> bool {
182 self.identity_matches(authority.binding())
183 && self.context_matches(authority.binding())
184 && authority.budget() == self.operation.budget()
185 }
186 fn identity_matches(&self, binding: &OperationBindingRecord) -> bool {
187 binding.intent() == self.intent.hash()
188 && binding.operation_sequence() == self.operation.operation_sequence()
189 && binding.target() == self.operation.target()
190 && binding.request() == self.operation.request()
191 }
192 fn context_matches(&self, binding: &OperationBindingRecord) -> bool {
193 binding.network() == self.context.network()
194 && binding.caller() == self.context.caller()
195 && binding.release() == self.context.release()
196 }
197}
198impl ExecutionAttemptTotalsView {
199 fn add(&mut self, attempts: &AttemptJournalView) -> Result<(), ExecutionProgressError> {
200 self.mutations_used = sum(self.mutations_used, attempts.mutations_used)?;
201 self.observations_used = sum(self.observations_used, attempts.observations_used)?;
202 self.mutations_remaining = sum(self.mutations_remaining, attempts.mutations_remaining)?;
203 self.observations_remaining =
204 sum(self.observations_remaining, attempts.observations_remaining)?;
205 Ok(())
206 }
207}
208fn sum(left: u32, right: u32) -> Result<u32, ExecutionProgressError> {
209 left.checked_add(right)
210 .ok_or(ExecutionProgressError::AccountingOverflow)
211}
212fn condition(attempts: &AttemptJournalView, unmet_dependencies: bool) -> OperationProgressState {
213 if attempts.applied {
214 OperationProgressState::Applied
215 } else if attempts.pending_observation.is_some() {
216 OperationProgressState::ObservationUnresolved
217 } else if attempts.pending_mutation.is_some() {
218 if attempts.observations_remaining == 0 {
219 OperationProgressState::ReconciliationExhausted
220 } else {
221 OperationProgressState::MutationUnresolved
222 }
223 } else if unmet_dependencies {
224 OperationProgressState::AwaitingDependencies
225 } else if attempts.mutations_remaining == 0 {
226 OperationProgressState::MutationExhausted
227 } else {
228 OperationProgressState::MutationAvailable
229 }
230}
231
232#[derive(Debug, Error, Eq, PartialEq)]
234pub enum ExecutionProgressError {
235 #[error("execution progress exceeds {MAX_EFFECT_OPERATIONS} journals")]
237 TooManyJournals,
238 #[error("duplicate journal for operation {0}")]
240 DuplicateJournal(u64),
241 #[error("journal operation {0} is absent from the original plan")]
243 UnknownOperation(u64),
244 #[error("missing original journal for operation {0}")]
246 MissingJournal(u64),
247 #[error("original journal authority mismatch for operation {0}")]
249 AuthorityMismatch(u64),
250 #[error(
252 "operation {operation_sequence} was attempted without applied prerequisite {prerequisite}"
253 )]
254 PrematureAttempt {
255 operation_sequence: u64,
257 prerequisite: u64,
259 },
260 #[error("execution attempt accounting overflow")]
262 AccountingOverflow,
263}
264
265#[cfg(test)]
266mod tests;