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