Skip to main content

lenso_module_management/
store.rs

1use crate::{
2    MODULE_OPERATION_JOURNAL_PROTOCOL, ModuleOperation, ModuleOperationJournal,
3    ModuleOperationJournalEvent, ModuleOperationLease, journal_event_digest,
4};
5use fs2::FileExt as _;
6use std::collections::BTreeMap;
7use std::fs::{self, File, OpenOptions};
8use std::io::{BufRead as _, BufReader, Write as _};
9use std::path::{Path, PathBuf};
10use std::sync::Mutex;
11use thiserror::Error;
12
13#[derive(Debug, Error)]
14pub enum ModuleOperationStoreError {
15    #[error("operation `{0}` was not found")]
16    NotFound(String),
17    #[error(
18        "operation `{operation_id}` revision changed: expected {expected}, observed {observed}"
19    )]
20    RevisionConflict {
21        operation_id: String,
22        expected: u64,
23        observed: u64,
24    },
25    #[error("operation `{0}` already exists")]
26    AlreadyExists(String),
27    #[error("operation journal is invalid: {0}")]
28    InvalidJournal(String),
29    #[error("operation store I/O failed: {0}")]
30    Io(#[from] std::io::Error),
31    #[error("operation store JSON failed: {0}")]
32    Json(#[from] serde_json::Error),
33}
34
35#[derive(Debug, Clone, PartialEq, Eq)]
36pub enum CreateOperationResult {
37    Created,
38    Existing(Box<ModuleOperation>),
39}
40
41pub trait ModuleOperationStore: std::fmt::Debug + Send + Sync {
42    fn create_idempotent(
43        &self,
44        operation: &ModuleOperation,
45        initial_event: &ModuleOperationJournalEvent,
46    ) -> Result<CreateOperationResult, ModuleOperationStoreError>;
47    fn load(&self, operation_id: &str) -> Result<ModuleOperation, ModuleOperationStoreError>;
48    fn journal(
49        &self,
50        operation_id: &str,
51    ) -> Result<ModuleOperationJournal, ModuleOperationStoreError>;
52    fn compare_and_append(
53        &self,
54        expected_revision: u64,
55        event: &ModuleOperationJournalEvent,
56    ) -> Result<(), ModuleOperationStoreError>;
57    fn find_by_idempotency_key(
58        &self,
59        application_id: &str,
60        idempotency_key: &str,
61    ) -> Result<Option<ModuleOperation>, ModuleOperationStoreError>;
62    fn load_lease(&self) -> Result<Option<ModuleOperationLease>, ModuleOperationStoreError>;
63    fn compare_and_set_lease(
64        &self,
65        expected_revision: Option<u64>,
66        lease: Option<&ModuleOperationLease>,
67    ) -> Result<(), ModuleOperationStoreError>;
68}
69
70#[derive(Debug, Default)]
71pub struct MemoryModuleOperationStore {
72    inner: Mutex<MemoryStoreState>,
73}
74
75#[derive(Debug, Default)]
76struct MemoryStoreState {
77    journals: BTreeMap<String, Vec<ModuleOperationJournalEvent>>,
78    lease: Option<ModuleOperationLease>,
79}
80
81impl ModuleOperationStore for MemoryModuleOperationStore {
82    fn create_idempotent(
83        &self,
84        operation: &ModuleOperation,
85        initial_event: &ModuleOperationJournalEvent,
86    ) -> Result<CreateOperationResult, ModuleOperationStoreError> {
87        let mut state = self.inner.lock().expect("memory operation store poisoned");
88        if let Some(existing) = state
89            .journals
90            .values()
91            .filter_map(|events| events.last())
92            .map(|event| &event.operation_after)
93            .find(|existing| {
94                existing.application_id == operation.application_id
95                    && existing.idempotency_key == operation.idempotency_key
96            })
97        {
98            return Ok(CreateOperationResult::Existing(Box::new(existing.clone())));
99        }
100        if state.journals.contains_key(&operation.operation_id) {
101            return Err(ModuleOperationStoreError::AlreadyExists(
102                operation.operation_id.clone(),
103            ));
104        }
105        state
106            .journals
107            .insert(operation.operation_id.clone(), vec![initial_event.clone()]);
108        Ok(CreateOperationResult::Created)
109    }
110
111    fn load(&self, operation_id: &str) -> Result<ModuleOperation, ModuleOperationStoreError> {
112        Ok(self
113            .journal(operation_id)?
114            .events
115            .last()
116            .expect("validated journal is non-empty")
117            .operation_after
118            .clone())
119    }
120
121    fn journal(
122        &self,
123        operation_id: &str,
124    ) -> Result<ModuleOperationJournal, ModuleOperationStoreError> {
125        let state = self.inner.lock().expect("memory operation store poisoned");
126        let events = state
127            .journals
128            .get(operation_id)
129            .cloned()
130            .ok_or_else(|| ModuleOperationStoreError::NotFound(operation_id.to_owned()))?;
131        validate_journal(operation_id, &events)?;
132        Ok(ModuleOperationJournal {
133            protocol: MODULE_OPERATION_JOURNAL_PROTOCOL.to_owned(),
134            operation_id: operation_id.to_owned(),
135            events,
136        })
137    }
138
139    fn compare_and_append(
140        &self,
141        expected_revision: u64,
142        event: &ModuleOperationJournalEvent,
143    ) -> Result<(), ModuleOperationStoreError> {
144        let mut state = self.inner.lock().expect("memory operation store poisoned");
145        let events = state
146            .journals
147            .get_mut(&event.operation_id)
148            .ok_or_else(|| ModuleOperationStoreError::NotFound(event.operation_id.clone()))?;
149        let observed = events
150            .last()
151            .expect("operation journal is non-empty")
152            .revision;
153        if observed != expected_revision {
154            return Err(ModuleOperationStoreError::RevisionConflict {
155                operation_id: event.operation_id.clone(),
156                expected: expected_revision,
157                observed,
158            });
159        }
160        events.push(event.clone());
161        Ok(())
162    }
163
164    fn find_by_idempotency_key(
165        &self,
166        application_id: &str,
167        idempotency_key: &str,
168    ) -> Result<Option<ModuleOperation>, ModuleOperationStoreError> {
169        let state = self.inner.lock().expect("memory operation store poisoned");
170        Ok(state
171            .journals
172            .values()
173            .filter_map(|events| events.last())
174            .map(|event| &event.operation_after)
175            .find(|operation| {
176                operation.application_id == application_id
177                    && operation.idempotency_key == idempotency_key
178            })
179            .cloned())
180    }
181
182    fn load_lease(&self) -> Result<Option<ModuleOperationLease>, ModuleOperationStoreError> {
183        Ok(self
184            .inner
185            .lock()
186            .expect("memory operation store poisoned")
187            .lease
188            .clone())
189    }
190
191    fn compare_and_set_lease(
192        &self,
193        expected_revision: Option<u64>,
194        lease: Option<&ModuleOperationLease>,
195    ) -> Result<(), ModuleOperationStoreError> {
196        let mut state = self.inner.lock().expect("memory operation store poisoned");
197        let observed = state.lease.as_ref().map(|lease| lease.revision);
198        if observed != expected_revision {
199            return Err(ModuleOperationStoreError::RevisionConflict {
200                operation_id: "composition-lease".to_owned(),
201                expected: expected_revision.unwrap_or(0),
202                observed: observed.unwrap_or(0),
203            });
204        }
205        state.lease = lease.cloned();
206        Ok(())
207    }
208}
209
210#[derive(Debug, Clone)]
211pub struct JsonFileModuleOperationStore {
212    root: PathBuf,
213}
214
215impl JsonFileModuleOperationStore {
216    pub fn new(root: impl Into<PathBuf>) -> Self {
217        Self { root: root.into() }
218    }
219
220    fn with_lock<T>(
221        &self,
222        operation: impl FnOnce() -> Result<T, ModuleOperationStoreError>,
223    ) -> Result<T, ModuleOperationStoreError> {
224        fs::create_dir_all(&self.root)?;
225        let lock = OpenOptions::new()
226            .create(true)
227            .truncate(false)
228            .read(true)
229            .write(true)
230            .open(self.root.join("store.lock"))?;
231        lock.lock_exclusive()?;
232        let result = operation();
233        fs2::FileExt::unlock(&lock)?;
234        result
235    }
236
237    fn operations_root(&self) -> PathBuf {
238        self.root.join("operations")
239    }
240
241    fn journal_path(&self, operation_id: &str) -> Result<PathBuf, ModuleOperationStoreError> {
242        validate_storage_identity(operation_id)?;
243        Ok(self
244            .operations_root()
245            .join(operation_id)
246            .join("journal.jsonl"))
247    }
248
249    fn read_events(
250        &self,
251        operation_id: &str,
252    ) -> Result<Vec<ModuleOperationJournalEvent>, ModuleOperationStoreError> {
253        let path = self.journal_path(operation_id)?;
254        let file = File::open(path).map_err(|error| match error.kind() {
255            std::io::ErrorKind::NotFound => {
256                ModuleOperationStoreError::NotFound(operation_id.to_owned())
257            }
258            _ => ModuleOperationStoreError::Io(error),
259        })?;
260        let events = BufReader::new(file)
261            .lines()
262            .filter_map(|line| match line {
263                Ok(line) if line.trim().is_empty() => None,
264                other => Some(other),
265            })
266            .map(|line| serde_json::from_str(&line?).map_err(ModuleOperationStoreError::from))
267            .collect::<Result<Vec<_>, _>>()?;
268        validate_journal(operation_id, &events)?;
269        Ok(events)
270    }
271
272    fn append_event(
273        path: &Path,
274        event: &ModuleOperationJournalEvent,
275    ) -> Result<(), ModuleOperationStoreError> {
276        let mut file = OpenOptions::new().create(true).append(true).open(path)?;
277        serde_json::to_writer(&mut file, event)?;
278        file.write_all(b"\n")?;
279        file.sync_data()?;
280        Ok(())
281    }
282}
283
284impl ModuleOperationStore for JsonFileModuleOperationStore {
285    fn create_idempotent(
286        &self,
287        operation: &ModuleOperation,
288        initial_event: &ModuleOperationJournalEvent,
289    ) -> Result<CreateOperationResult, ModuleOperationStoreError> {
290        self.with_lock(|| {
291            if self.operations_root().exists() {
292                for entry in fs::read_dir(self.operations_root())? {
293                    let entry = entry?;
294                    if !entry.file_type()?.is_dir() {
295                        continue;
296                    }
297                    let operation_id = entry.file_name().to_string_lossy().into_owned();
298                    let existing = self
299                        .read_events(&operation_id)?
300                        .last()
301                        .expect("validated journal is non-empty")
302                        .operation_after
303                        .clone();
304                    if existing.application_id == operation.application_id
305                        && existing.idempotency_key == operation.idempotency_key
306                    {
307                        return Ok(CreateOperationResult::Existing(Box::new(existing)));
308                    }
309                }
310            }
311            let path = self.journal_path(&operation.operation_id)?;
312            if path.exists() {
313                return Err(ModuleOperationStoreError::AlreadyExists(
314                    operation.operation_id.clone(),
315                ));
316            }
317            fs::create_dir_all(path.parent().expect("journal has parent"))?;
318            Self::append_event(&path, initial_event)?;
319            Ok(CreateOperationResult::Created)
320        })
321    }
322
323    fn load(&self, operation_id: &str) -> Result<ModuleOperation, ModuleOperationStoreError> {
324        Ok(self
325            .journal(operation_id)?
326            .events
327            .last()
328            .expect("validated journal is non-empty")
329            .operation_after
330            .clone())
331    }
332
333    fn journal(
334        &self,
335        operation_id: &str,
336    ) -> Result<ModuleOperationJournal, ModuleOperationStoreError> {
337        self.with_lock(|| {
338            let events = self.read_events(operation_id)?;
339            Ok(ModuleOperationJournal {
340                protocol: MODULE_OPERATION_JOURNAL_PROTOCOL.to_owned(),
341                operation_id: operation_id.to_owned(),
342                events,
343            })
344        })
345    }
346
347    fn compare_and_append(
348        &self,
349        expected_revision: u64,
350        event: &ModuleOperationJournalEvent,
351    ) -> Result<(), ModuleOperationStoreError> {
352        self.with_lock(|| {
353            let events = self.read_events(&event.operation_id)?;
354            let observed = events
355                .last()
356                .expect("validated journal is non-empty")
357                .revision;
358            if observed != expected_revision {
359                return Err(ModuleOperationStoreError::RevisionConflict {
360                    operation_id: event.operation_id.clone(),
361                    expected: expected_revision,
362                    observed,
363                });
364            }
365            Self::append_event(&self.journal_path(&event.operation_id)?, event)
366        })
367    }
368
369    fn find_by_idempotency_key(
370        &self,
371        application_id: &str,
372        idempotency_key: &str,
373    ) -> Result<Option<ModuleOperation>, ModuleOperationStoreError> {
374        self.with_lock(|| {
375            if !self.operations_root().exists() {
376                return Ok(None);
377            }
378            for entry in fs::read_dir(self.operations_root())? {
379                let entry = entry?;
380                if !entry.file_type()?.is_dir() {
381                    continue;
382                }
383                let operation_id = entry.file_name().to_string_lossy().into_owned();
384                let operation = self
385                    .read_events(&operation_id)?
386                    .last()
387                    .expect("validated journal is non-empty")
388                    .operation_after
389                    .clone();
390                if operation.application_id == application_id
391                    && operation.idempotency_key == idempotency_key
392                {
393                    return Ok(Some(operation));
394                }
395            }
396            Ok(None)
397        })
398    }
399
400    fn load_lease(&self) -> Result<Option<ModuleOperationLease>, ModuleOperationStoreError> {
401        let path = self.root.join("composition-lease.json");
402        match fs::read(path) {
403            Ok(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)),
404            Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
405            Err(error) => Err(error.into()),
406        }
407    }
408
409    fn compare_and_set_lease(
410        &self,
411        expected_revision: Option<u64>,
412        lease: Option<&ModuleOperationLease>,
413    ) -> Result<(), ModuleOperationStoreError> {
414        self.with_lock(|| {
415            let observed = self.load_lease()?.map(|lease| lease.revision);
416            if observed != expected_revision {
417                return Err(ModuleOperationStoreError::RevisionConflict {
418                    operation_id: "composition-lease".to_owned(),
419                    expected: expected_revision.unwrap_or(0),
420                    observed: observed.unwrap_or(0),
421                });
422            }
423            let path = self.root.join("composition-lease.json");
424            if let Some(lease) = lease {
425                let temporary = self.root.join("composition-lease.next.json");
426                fs::write(&temporary, serde_json::to_vec_pretty(lease)?)?;
427                File::open(&temporary)?.sync_all()?;
428                fs::rename(temporary, path)?;
429            } else if path.exists() {
430                fs::remove_file(path)?;
431            }
432            Ok(())
433        })
434    }
435}
436
437fn validate_storage_identity(value: &str) -> Result<(), ModuleOperationStoreError> {
438    if value.is_empty()
439        || !value
440            .bytes()
441            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_'))
442    {
443        return Err(ModuleOperationStoreError::InvalidJournal(
444            "operation identity is unsafe for target-owned storage".to_owned(),
445        ));
446    }
447    Ok(())
448}
449
450fn validate_journal(
451    operation_id: &str,
452    events: &[ModuleOperationJournalEvent],
453) -> Result<(), ModuleOperationStoreError> {
454    if events.is_empty() {
455        return Err(ModuleOperationStoreError::InvalidJournal(
456            "journal has no events".to_owned(),
457        ));
458    }
459    let mut prior_digest: Option<String> = None;
460    for (index, event) in events.iter().enumerate() {
461        let expected_revision = u64::try_from(index).expect("journal index fits u64");
462        if event.operation_id != operation_id
463            || event.revision != expected_revision
464            || event.operation_after.revision != event.revision
465            || event.prior_event_digest != prior_digest
466        {
467            return Err(ModuleOperationStoreError::InvalidJournal(format!(
468                "event {index} breaks identity, revision, or digest chaining"
469            )));
470        }
471        prior_digest = Some(journal_event_digest(event)?);
472    }
473    Ok(())
474}