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}