ic_backup/ops/persistence/attempt_journal/
mod.rs1use super::{
4 BackupLayoutGuard, ExecutionProgressPersistenceError, JournalLock, JournalLockError,
5 PersistenceError, create_json_durable, read_json, write_json_durable,
6};
7use crate::model::artifacts::ArtifactChecksumRecord;
8use crate::model::attempt_journal::{
9 AttemptAuthorityRecord, AttemptJournalRecord, AttemptJournalRecordError,
10 MAX_ATTEMPT_JOURNAL_BYTES, MutationReceiptRequest, ObservationReceiptRequest,
11};
12use std::path::{Path, PathBuf};
13use thiserror::Error;
14
15#[derive(Debug)]
20pub struct AttemptJournalGuard<'a> {
21 layout: &'a BackupLayoutGuard,
22 _lock: JournalLock,
23 record: AttemptJournalRecord,
24 usable: bool,
25}
26
27impl<'a> AttemptJournalGuard<'a> {
28 pub fn create(
33 layout: &'a BackupLayoutGuard,
34 authority: AttemptAuthorityRecord,
35 ) -> Result<Self, AttemptJournalError> {
36 layout.check_root()?;
37 let path = journal_path(layout, &authority);
38 let lock = JournalLock::acquire(&path)?;
39 let record = AttemptJournalRecord::new(authority);
40 check_size(&record)?;
41 create_json_durable(&path, &record)?;
42 Ok(Self {
43 layout,
44 _lock: lock,
45 record,
46 usable: true,
47 })
48 }
49 pub fn open(
54 layout: &'a BackupLayoutGuard,
55 expected: &AttemptAuthorityRecord,
56 ) -> Result<Self, AttemptJournalError> {
57 layout.check_root()?;
58 let path = journal_path(layout, expected);
59 let lock = JournalLock::acquire(&path)?;
60 let record: AttemptJournalRecord = read_json(&path, MAX_ATTEMPT_JOURNAL_BYTES)?;
61 check_size(&record)?;
62 if record.authority() != expected {
63 return Err(AttemptJournalError::AuthorityMismatch);
64 }
65 Ok(Self {
66 layout,
67 _lock: lock,
68 record,
69 usable: true,
70 })
71 }
72 pub fn record(&self) -> Result<&AttemptJournalRecord, AttemptJournalError> {
77 self.check_usable()?;
78 Ok(&self.record)
79 }
80 #[must_use]
82 pub fn path(&self) -> PathBuf {
83 journal_path(self.layout, self.record.authority())
84 }
85 pub fn reserve_mutation(&mut self) -> Result<u32, AttemptJournalError> {
90 self.reserve_with(AttemptJournalRecord::reserve_mutation, write_json_durable)
91 }
92 pub fn reserve_planned_mutation(
103 &mut self,
104 expected_plan: &ArtifactChecksumRecord,
105 ) -> Result<u32, ExecutionProgressPersistenceError> {
106 let view = super::execution_progress::read_with_current(
107 self.layout,
108 expected_plan,
109 Some(self.record()?),
110 )?;
111 let sequence = self.record.authority().binding().operation_sequence();
112 if view.operations.iter().any(|operation| {
113 operation.operation_sequence == sequence
114 && operation.state
115 == crate::policy::execution_progress::OperationProgressState::AwaitingDependencies
116 }) {
117 return Err(ExecutionProgressPersistenceError::DependenciesUnapplied(sequence));
118 }
119 Ok(self.reserve_mutation()?)
120 }
121 pub fn reserve_observation(
127 &mut self,
128 mutation: u32,
129 request: &str,
130 ) -> Result<u32, AttemptJournalError> {
131 self.reserve_with(
132 |record| record.reserve_observation(mutation, request),
133 write_json_durable,
134 )
135 }
136 pub fn reserve_planned_observation(
146 &mut self,
147 expected_plan: &ArtifactChecksumRecord,
148 mutation: u32,
149 request: &str,
150 ) -> Result<u32, ExecutionProgressPersistenceError> {
151 super::execution_progress::read_with_current(
152 self.layout,
153 expected_plan,
154 Some(self.record()?),
155 )?;
156 Ok(self.reserve_observation(mutation, request)?)
157 }
158 pub fn record_mutation(
163 &mut self,
164 receipt: MutationReceiptRequest,
165 ) -> Result<(), AttemptJournalError> {
166 self.reserve_with(|record| record.record_mutation(receipt), write_json_durable)
167 }
168 pub fn record_observation(
176 &mut self,
177 receipt: ObservationReceiptRequest,
178 ) -> Result<(), AttemptJournalError> {
179 self.reserve_with(
180 |record| record.record_observation(receipt),
181 write_json_durable,
182 )
183 }
184
185 fn reserve_with<T>(
186 &mut self,
187 transition: impl FnOnce(&mut AttemptJournalRecord) -> Result<T, AttemptJournalRecordError>,
188 write: impl FnOnce(&Path, &AttemptJournalRecord) -> Result<(), PersistenceError>,
189 ) -> Result<T, AttemptJournalError> {
190 self.check_usable()?;
191 let mut next = self.record.clone();
192 let result = transition(&mut next)?;
193 check_size(&next)?;
194 self.usable = false;
195 write(&self.path(), &next)?;
196 self.record = next;
197 self.usable = true;
198 Ok(result)
199 }
200 fn check_usable(&self) -> Result<(), AttemptJournalError> {
201 if !self.usable {
202 return Err(AttemptJournalError::IndeterminateWrite);
203 }
204 self.layout.check_root()?;
205 Ok(())
206 }
207}
208
209pub(super) fn read_original_journals(
210 layout: &BackupLayoutGuard,
211 authorities: &[AttemptAuthorityRecord],
212 current: Option<&AttemptJournalRecord>,
213) -> Result<Vec<AttemptJournalRecord>, AttemptJournalError> {
214 layout.check_root()?;
215 if let Some(current) = current
216 && !authorities.contains(current.authority())
217 {
218 return Err(AttemptJournalError::AuthorityMismatch);
219 }
220 let mut journals = Vec::with_capacity(authorities.len());
221 for authority in authorities {
222 if let Some(current) = current
223 && current.authority() == authority
224 {
225 journals.push(current.clone());
226 } else {
227 let guard = AttemptJournalGuard::open(layout, authority)?;
228 journals.push(guard.record()?.clone());
229 }
230 }
231 layout.check_root()?;
232 Ok(journals)
233}
234
235fn journal_path(layout: &BackupLayoutGuard, authority: &AttemptAuthorityRecord) -> PathBuf {
236 layout.root().join(format!(
237 "attempt-{}.json",
238 authority.binding().operation_sequence()
239 ))
240}
241fn check_size(record: &AttemptJournalRecord) -> Result<(), PersistenceError> {
242 super::json::check_json_size(record, MAX_ATTEMPT_JOURNAL_BYTES)
243}
244
245#[derive(Debug, Error)]
247pub enum AttemptJournalError {
248 #[error("attempt journal authority mismatch")]
250 AuthorityMismatch,
251 #[error("attempt journal write outcome indeterminate; reopen retained evidence")]
253 IndeterminateWrite,
254 #[error(transparent)]
256 Record(#[from] AttemptJournalRecordError),
257 #[error(transparent)]
259 Lock(#[from] JournalLockError),
260 #[error(transparent)]
262 Persistence(#[from] PersistenceError),
263}
264
265#[cfg(all(test, unix))]
266mod tests;