Skip to main content

vsh_store/
persistent.rs

1use std::collections::BTreeMap;
2use std::fs::File;
3use std::io::{self, Read, Seek, SeekFrom, Write};
4use std::path::{Path, PathBuf};
5use std::sync::{Arc, Mutex, MutexGuard};
6
7use cap_std::fs::{Dir, OpenOptions};
8use vsh_types::{BlobId, PrincipalId, SnapshotId, TransactionId, TransactionState};
9
10use super::{
11    ApprovalGrant, CommitReservation, DataDirectory, TransactionRecord, TransactionStore,
12    TransactionStoreError, directory,
13};
14
15const LOG_MAGIC: &[u8; 8] = b"VSHST001";
16const CONTROL_MAGIC: &[u8; 8] = b"VSHCT001";
17const LOCK_FILE: &str = "transactions.lock";
18const LOG_FILE: &str = "transactions.vsh";
19const ALTERNATE_LOG_FILE: &str = "transactions.alt.vsh";
20const HEADER_BYTES: u64 = LOG_MAGIC.len() as u64;
21const CONTROL_HEADER_BYTES: u64 = CONTROL_MAGIC.len() as u64;
22const CONTROL_PAYLOAD_BYTES: usize = 8 + 1;
23const CONTROL_RECORD_BYTES: usize = CONTROL_PAYLOAD_BYTES + DIGEST_BYTES;
24const CONTROL_SLOT_COUNT: usize = 2;
25const CONTROL_BYTES: u64 =
26    CONTROL_HEADER_BYTES + (CONTROL_RECORD_BYTES * CONTROL_SLOT_COUNT) as u64;
27const DIGEST_BYTES: usize = 32;
28const OLD_MIN_PAYLOAD_BYTES: usize = 32 + 32 + 1 + 1;
29const MIN_PAYLOAD_BYTES: usize = 32 + 32 + 1 + 1 + 1;
30const ARTIFACT_PAYLOAD_BYTES: usize = 32;
31const APPROVAL_PAYLOAD_BYTES: usize = 32 + 8 + 8;
32const OLD_MAX_PAYLOAD_BYTES: usize = OLD_MIN_PAYLOAD_BYTES + APPROVAL_PAYLOAD_BYTES;
33const MAX_PAYLOAD_BYTES: usize =
34    MIN_PAYLOAD_BYTES + ARTIFACT_PAYLOAD_BYTES + APPROVAL_PAYLOAD_BYTES;
35
36/// Hard bounds for the dependency-free compacting transaction state store.
37#[derive(Clone, Copy, Debug, Eq, PartialEq)]
38pub struct FileStoreConfig {
39    /// Maximum active log bytes, including framing and checksums. When the next
40    /// append would cross this bound, the latest record for every transaction is
41    /// rewritten to the inactive slot and atomically selected by the control file.
42    pub max_log_bytes: u64,
43    /// Maximum unique transactions retained by one store.
44    pub max_records: usize,
45}
46
47impl Default for FileStoreConfig {
48    fn default() -> Self {
49        Self {
50            max_log_bytes: 256 * 1024 * 1024,
51            max_records: 1_000_000,
52        }
53    }
54}
55
56#[derive(Debug, Default)]
57struct PersistentState {
58    records: BTreeMap<TransactionId, TransactionRecord>,
59    offset: u64,
60    generation: u64,
61    slot: u8,
62}
63
64#[derive(Clone, Copy, Debug, Eq, PartialEq)]
65struct ControlState {
66    generation: u64,
67    slot: u8,
68}
69
70/// Durable cross-process transaction state using only the Rust standard library.
71///
72/// Each operation takes a short process-local mutex and an OS file lock, refreshes only
73/// unseen append-log bytes, validates checksums and lifecycle edges, then synchronizes
74/// at most one complete record. At the configured byte ceiling it compacts into the
75/// inactive log and switches a fixed-size, double-buffered control record. Monty
76/// execution, VFS work, and host commit operations never run while either lock is held.
77#[derive(Clone, Debug)]
78pub struct FileTransactionStore {
79    state: Arc<Mutex<PersistentState>>,
80    lock: Arc<File>,
81    data_directory: DataDirectory,
82    config: FileStoreConfig,
83}
84
85impl FileTransactionStore {
86    /// Open or initialize a durable transaction store below `data_directory`.
87    ///
88    /// # Errors
89    ///
90    /// Returns an error for filesystem/lock failures, invalid bounds, corruption, an
91    /// invalid persisted state edge, or a configured size ceiling.
92    pub fn open(
93        data_directory: impl AsRef<Path>,
94        config: FileStoreConfig,
95    ) -> Result<Self, TransactionStoreError> {
96        validate_config(config)?;
97        let data_directory = DataDirectory::open_trusted(data_directory).map_err(|source| {
98            persistent_io("open trusted data directory", io::Error::other(source))
99        })?;
100        Self::open_validated_in(&data_directory, config)
101    }
102
103    /// Open or initialize a durable transaction store below a pinned capability.
104    ///
105    /// # Errors
106    ///
107    /// Returns an error for invalid bounds, filesystem/lock failures, corruption,
108    /// an invalid persisted state edge, or a configured size ceiling.
109    pub fn open_in(
110        data_directory: &DataDirectory,
111        config: FileStoreConfig,
112    ) -> Result<Self, TransactionStoreError> {
113        validate_config(config)?;
114        Self::open_validated_in(data_directory, config)
115    }
116
117    fn open_validated_in(
118        data_directory: &DataDirectory,
119        config: FileStoreConfig,
120    ) -> Result<Self, TransactionStoreError> {
121        let mut lock_options = OpenOptions::new();
122        lock_options
123            .read(true)
124            .write(true)
125            .create(true)
126            .truncate(false);
127        let lock = directory::open_real_file(data_directory.directory(), LOCK_FILE, &lock_options)
128            .map_err(|source| persistent_io("open lock file", source))?;
129        let guard = FileLockGuard::exclusive(&lock)?;
130        let control = initialize_control(&lock)?;
131        let mut log = open_log(data_directory.directory(), control.slot)?;
132        initialize_header(&mut log)?;
133        let mut state = PersistentState {
134            records: BTreeMap::new(),
135            offset: HEADER_BYTES,
136            generation: control.generation,
137            slot: control.slot,
138        };
139        refresh(&mut state, &mut log, config, control)?;
140        directory::sync_directory(data_directory.directory())
141            .map_err(|source| persistent_io("sync data directory", source))?;
142        drop(guard);
143        Ok(Self {
144            state: Arc::new(Mutex::new(state)),
145            lock: Arc::new(lock),
146            data_directory: data_directory.clone(),
147            config,
148        })
149    }
150
151    /// Return the currently active state-log path for diagnostics and backup tooling.
152    ///
153    /// # Errors
154    ///
155    /// Returns an error if the cross-process lock or compact-log control journal
156    /// cannot be read safely.
157    pub fn active_log_path(&self) -> Result<PathBuf, TransactionStoreError> {
158        let _guard = FileLockGuard::exclusive(&self.lock)?;
159        let control = read_control(&self.lock)?;
160        Ok(log_path_for(self.data_directory.path(), control.slot))
161    }
162
163    fn state(&self) -> Result<MutexGuard<'_, PersistentState>, TransactionStoreError> {
164        self.state
165            .lock()
166            .map_err(|_| TransactionStoreError::Poisoned)
167    }
168
169    fn transact<R>(
170        &self,
171        operation: impl FnOnce(
172            &BTreeMap<TransactionId, TransactionRecord>,
173        ) -> Result<Mutation<R>, TransactionStoreError>,
174    ) -> Result<R, TransactionStoreError> {
175        let mut state = self.state()?;
176        let _guard = FileLockGuard::exclusive(&self.lock)?;
177        let control = read_control(&self.lock)?;
178        let mut log = open_log(self.data_directory.directory(), control.slot)?;
179        initialize_header(&mut log)?;
180        refresh(&mut state, &mut log, self.config, control)?;
181        let mutation = operation(&state.records)?;
182        if let Some(record) = mutation.persist {
183            if !state.records.contains_key(&record.id())
184                && state.records.len() >= self.config.max_records
185            {
186                return Err(TransactionStoreError::PersistentRecordLimit {
187                    observed: state.records.len().saturating_add(1),
188                    maximum: self.config.max_records,
189                });
190            }
191            match append_record(&mut log, &record, &mut state.offset, self.config) {
192                Ok(()) => {
193                    state.records.insert(record.id(), record);
194                }
195                Err(TransactionStoreError::PersistentLogLimit { .. }) => {
196                    self.compact(&mut state, control, record)?;
197                }
198                Err(error) => return Err(error),
199            }
200        }
201        mutation.result
202    }
203
204    fn compact(
205        &self,
206        state: &mut PersistentState,
207        control: ControlState,
208        record: TransactionRecord,
209    ) -> Result<(), TransactionStoreError> {
210        let generation =
211            control
212                .generation
213                .checked_add(1)
214                .ok_or(TransactionStoreError::PersistentCorrupt {
215                    offset: CONTROL_HEADER_BYTES,
216                    reason: "state-log control generation exhausted",
217                })?;
218        let next = ControlState {
219            generation,
220            slot: control.slot ^ 1,
221        };
222        let mut options = OpenOptions::new();
223        options.read(true).write(true).create(true).truncate(true);
224        let mut compacted = directory::open_real_file(
225            self.data_directory.directory(),
226            log_name_for(next.slot),
227            &options,
228        )
229        .map_err(|source| persistent_io("open compacted state log", source))?;
230        compacted
231            .write_all(LOG_MAGIC)
232            .map_err(|source| persistent_io("write compacted state header", source))?;
233        let mut offset = HEADER_BYTES;
234        let mut wrote_record = false;
235        for (id, existing) in &state.records {
236            if !wrote_record && record.id() < *id {
237                offset = write_record_frame(&mut compacted, &record, offset, self.config)?;
238                wrote_record = true;
239            }
240            if *id == record.id() {
241                offset = write_record_frame(&mut compacted, &record, offset, self.config)?;
242                wrote_record = true;
243            } else {
244                offset = write_record_frame(&mut compacted, existing, offset, self.config)?;
245            }
246        }
247        if !wrote_record {
248            offset = write_record_frame(&mut compacted, &record, offset, self.config)?;
249        }
250        compacted
251            .sync_all()
252            .map_err(|source| persistent_io("sync compacted state log", source))?;
253        directory::sync_directory(self.data_directory.directory())
254            .map_err(|source| persistent_io("sync compacted state directory", source))?;
255        append_control(&self.lock, control, next)?;
256
257        state.records.insert(record.id(), record);
258        state.offset = offset;
259        state.generation = next.generation;
260        state.slot = next.slot;
261        Ok(())
262    }
263}
264
265impl TransactionStore for FileTransactionStore {
266    fn create(&self, record: TransactionRecord) -> Result<(), TransactionStoreError> {
267        validate_record(&record, 0)?;
268        self.transact(|records| {
269            if records.contains_key(&record.id()) {
270                return Err(TransactionStoreError::Duplicate { id: record.id() });
271            }
272            Ok(Mutation::persist(record, ()))
273        })
274    }
275
276    fn get(&self, id: TransactionId) -> Result<TransactionRecord, TransactionStoreError> {
277        self.transact(|records| {
278            records
279                .get(&id)
280                .cloned()
281                .map(Mutation::return_value)
282                .ok_or(TransactionStoreError::NotFound { id })
283        })
284    }
285
286    fn compare_and_transition(
287        &self,
288        id: TransactionId,
289        expected: TransactionState,
290        next: TransactionState,
291    ) -> Result<TransactionRecord, TransactionStoreError> {
292        self.transact(|records| {
293            let mut record = records
294                .get(&id)
295                .cloned()
296                .ok_or(TransactionStoreError::NotFound { id })?;
297            if record.state() != expected {
298                return Err(TransactionStoreError::StateConflict {
299                    id,
300                    expected,
301                    actual: record.state(),
302                });
303            }
304            record
305                .transition(next)
306                .map_err(TransactionStoreError::Transition)?;
307            Ok(Mutation::persist(record.clone(), record))
308        })
309    }
310
311    fn approve(
312        &self,
313        id: TransactionId,
314        grant: ApprovalGrant,
315    ) -> Result<TransactionRecord, TransactionStoreError> {
316        if grant.transaction() != id {
317            return Err(TransactionStoreError::ApprovalBindingMismatch {
318                requested: id,
319                bound: grant.transaction(),
320            });
321        }
322        self.transact(|records| {
323            let mut record = records
324                .get(&id)
325                .cloned()
326                .ok_or(TransactionStoreError::NotFound { id })?;
327            if record.state() != TransactionState::PendingApproval {
328                return Err(TransactionStoreError::StateConflict {
329                    id,
330                    expected: TransactionState::PendingApproval,
331                    actual: record.state(),
332                });
333            }
334            record
335                .transition(TransactionState::Approved)
336                .map_err(TransactionStoreError::Transition)?;
337            record.approval = Some(grant);
338            Ok(Mutation::persist(record.clone(), record))
339        })
340    }
341
342    fn reserve(
343        &self,
344        id: TransactionId,
345        now_unix_ms: u64,
346    ) -> Result<CommitReservation, TransactionStoreError> {
347        self.transact(|records| {
348            let mut record = records
349                .get(&id)
350                .cloned()
351                .ok_or(TransactionStoreError::NotFound { id })?;
352            match record.state() {
353                TransactionState::AutoApproved => {}
354                TransactionState::Approved => {
355                    let grant = record
356                        .approval()
357                        .ok_or(TransactionStoreError::MissingApproval { id })?;
358                    if grant.is_expired_at(now_unix_ms) {
359                        record
360                            .transition(TransactionState::Expired)
361                            .map_err(TransactionStoreError::Transition)?;
362                        return Ok(Mutation::persist_error(
363                            record,
364                            TransactionStoreError::ApprovalExpired {
365                                id,
366                                expired_at_unix_ms: grant.expires_at_unix_ms(),
367                                observed_at_unix_ms: now_unix_ms,
368                            },
369                        ));
370                    }
371                }
372                actual => {
373                    return Err(TransactionStoreError::NotReservable { id, actual });
374                }
375            }
376            record
377                .transition(TransactionState::Reserved)
378                .map_err(TransactionStoreError::Transition)?;
379            let reservation = CommitReservation {
380                transaction: id,
381                base_snapshot: record.base_snapshot(),
382            };
383            Ok(Mutation::persist(record, reservation))
384        })
385    }
386}
387
388struct Mutation<R> {
389    persist: Option<TransactionRecord>,
390    result: Result<R, TransactionStoreError>,
391}
392
393impl<R> Mutation<R> {
394    fn return_value(value: R) -> Self {
395        Self {
396            persist: None,
397            result: Ok(value),
398        }
399    }
400
401    fn persist(record: TransactionRecord, value: R) -> Self {
402        Self {
403            persist: Some(record),
404            result: Ok(value),
405        }
406    }
407
408    fn persist_error(record: TransactionRecord, error: TransactionStoreError) -> Self {
409        Self {
410            persist: Some(record),
411            result: Err(error),
412        }
413    }
414}
415
416struct FileLockGuard<'a>(&'a File);
417
418impl<'a> FileLockGuard<'a> {
419    fn exclusive(file: &'a File) -> Result<Self, TransactionStoreError> {
420        File::lock(file).map_err(|source| persistent_io("acquire lock", source))?;
421        Ok(Self(file))
422    }
423}
424
425impl Drop for FileLockGuard<'_> {
426    fn drop(&mut self) {
427        let _ = File::unlock(self.0);
428    }
429}
430
431fn log_path_for(data_directory: &Path, slot: u8) -> PathBuf {
432    data_directory.join(log_name_for(slot))
433}
434
435fn log_name_for(slot: u8) -> &'static str {
436    match slot {
437        0 => LOG_FILE,
438        1 => ALTERNATE_LOG_FILE,
439        _ => unreachable!("validated control slot"),
440    }
441}
442
443fn initialize_control(lock: &File) -> Result<ControlState, TransactionStoreError> {
444    let mut control = lock
445        .try_clone()
446        .map_err(|source| persistent_io("clone state control file", source))?;
447    let length = control
448        .metadata()
449        .map_err(|source| persistent_io("inspect state control file", source))?
450        .len();
451    if length < CONTROL_HEADER_BYTES {
452        let mut prefix = vec![0_u8; usize::try_from(length).unwrap_or(0)];
453        control
454            .seek(SeekFrom::Start(0))
455            .map_err(|source| persistent_io("seek state control header", source))?;
456        control
457            .read_exact(&mut prefix)
458            .map_err(|source| persistent_io("read partial state control header", source))?;
459        if !CONTROL_MAGIC.starts_with(&prefix) {
460            return Err(TransactionStoreError::PersistentCorrupt {
461                offset: 0,
462                reason: "invalid state-control header",
463            });
464        }
465        initialize_control_slots(&mut control)?;
466        return read_control(lock);
467    }
468
469    control
470        .seek(SeekFrom::Start(0))
471        .map_err(|source| persistent_io("seek state control header", source))?;
472    let mut magic = [0_u8; CONTROL_MAGIC.len()];
473    control
474        .read_exact(&mut magic)
475        .map_err(|source| persistent_io("read state control header", source))?;
476    if &magic != CONTROL_MAGIC {
477        return Err(TransactionStoreError::PersistentCorrupt {
478            offset: 0,
479            reason: "invalid state-control header",
480        });
481    }
482    if length != CONTROL_BYTES {
483        control
484            .set_len(CONTROL_BYTES)
485            .map_err(|source| persistent_io("repair state control length", source))?;
486        control
487            .sync_all()
488            .map_err(|source| persistent_io("sync repaired state control length", source))?;
489    }
490    match read_control(lock) {
491        Ok(state) => Ok(state),
492        Err(TransactionStoreError::PersistentCorrupt {
493            reason: "state-control file has no valid slot",
494            ..
495        }) if length < CONTROL_BYTES => {
496            initialize_control_slots(&mut control)?;
497            read_control(lock)
498        }
499        Err(error) => Err(error),
500    }
501}
502
503fn read_control(lock: &File) -> Result<ControlState, TransactionStoreError> {
504    let mut control = lock
505        .try_clone()
506        .map_err(|source| persistent_io("clone state control file", source))?;
507    let length = control
508        .metadata()
509        .map_err(|source| persistent_io("inspect state control file", source))?
510        .len();
511    if length != CONTROL_BYTES {
512        return Err(TransactionStoreError::PersistentCorrupt {
513            offset: 0,
514            reason: "invalid state-control file length",
515        });
516    }
517    control
518        .seek(SeekFrom::Start(0))
519        .map_err(|source| persistent_io("seek state control header", source))?;
520    let mut magic = [0_u8; CONTROL_MAGIC.len()];
521    control
522        .read_exact(&mut magic)
523        .map_err(|source| persistent_io("read state control header", source))?;
524    if &magic != CONTROL_MAGIC {
525        return Err(TransactionStoreError::PersistentCorrupt {
526            offset: 0,
527            reason: "invalid state-control header",
528        });
529    }
530
531    let mut candidates = Vec::with_capacity(CONTROL_SLOT_COUNT);
532    for slot_index in 0..CONTROL_SLOT_COUNT {
533        let offset = control_slot_offset(slot_index);
534        control
535            .seek(SeekFrom::Start(offset))
536            .map_err(|source| persistent_io("seek state control slot", source))?;
537        let mut encoded = [0_u8; CONTROL_RECORD_BYTES];
538        control
539            .read_exact(&mut encoded)
540            .map_err(|source| persistent_io("read state control slot", source))?;
541        let (payload, digest) = encoded.split_at(CONTROL_PAYLOAD_BYTES);
542        if BlobId::digest(payload).as_bytes().as_slice() != digest {
543            continue;
544        }
545        let candidate = ControlState {
546            generation: u64::from_le_bytes(copy_array(&payload[..8])),
547            slot: payload[8],
548        };
549        if usize::from(candidate.slot) != slot_index {
550            return Err(TransactionStoreError::PersistentCorrupt {
551                offset,
552                reason: "invalid state-control slot",
553            });
554        }
555        if usize::try_from(candidate.generation & 1).unwrap_or(usize::MAX) != slot_index {
556            return Err(TransactionStoreError::PersistentCorrupt {
557                offset,
558                reason: "state-control generation is in the wrong slot",
559            });
560        }
561        candidates.push(candidate);
562    }
563    candidates.sort_unstable_by_key(|candidate| candidate.generation);
564    if let [older, newer] = candidates.as_slice()
565        && newer.generation != older.generation.saturating_add(1)
566    {
567        return Err(TransactionStoreError::PersistentCorrupt {
568            offset: CONTROL_HEADER_BYTES,
569            reason: "state-control generations are not consecutive",
570        });
571    }
572    candidates
573        .last()
574        .copied()
575        .ok_or(TransactionStoreError::PersistentCorrupt {
576            offset: CONTROL_HEADER_BYTES,
577            reason: "state-control file has no valid slot",
578        })
579}
580
581fn append_control(
582    lock: &File,
583    current: ControlState,
584    next: ControlState,
585) -> Result<(), TransactionStoreError> {
586    if next.generation != current.generation.saturating_add(1) || next.slot != (current.slot ^ 1) {
587        return Err(TransactionStoreError::PersistentCorrupt {
588            offset: CONTROL_HEADER_BYTES,
589            reason: "invalid state-control append",
590        });
591    }
592    if read_control(lock)? != current {
593        return Err(TransactionStoreError::PersistentCorrupt {
594            offset: CONTROL_HEADER_BYTES,
595            reason: "state-control changed while exclusively locked",
596        });
597    }
598    let mut control = lock
599        .try_clone()
600        .map_err(|source| persistent_io("clone state control file", source))?;
601    control
602        .seek(SeekFrom::Start(control_slot_offset(usize::from(next.slot))))
603        .map_err(|source| persistent_io("seek next state control slot", source))?;
604    write_control_record(&mut control, next)?;
605    control
606        .sync_all()
607        .map_err(|source| persistent_io("sync state control record", source))
608}
609
610fn initialize_control_slots(control: &mut File) -> Result<(), TransactionStoreError> {
611    control
612        .set_len(0)
613        .map_err(|source| persistent_io("reset state control file", source))?;
614    control
615        .seek(SeekFrom::Start(0))
616        .map_err(|source| persistent_io("seek initial state control", source))?;
617    control
618        .write_all(CONTROL_MAGIC)
619        .map_err(|source| persistent_io("write state control header", source))?;
620    write_control_record(
621        control,
622        ControlState {
623            generation: 0,
624            slot: 0,
625        },
626    )?;
627    control
628        .write_all(&[0_u8; CONTROL_RECORD_BYTES])
629        .map_err(|source| persistent_io("clear alternate state control slot", source))?;
630    control
631        .sync_all()
632        .map_err(|source| persistent_io("sync initial state control slots", source))
633}
634
635fn write_control_record(
636    control: &mut File,
637    state: ControlState,
638) -> Result<(), TransactionStoreError> {
639    let mut payload = [0_u8; CONTROL_PAYLOAD_BYTES];
640    payload[..8].copy_from_slice(&state.generation.to_le_bytes());
641    payload[8] = state.slot;
642    control
643        .write_all(&payload)
644        .and_then(|()| control.write_all(BlobId::digest(&payload).as_bytes()))
645        .map_err(|source| persistent_io("write state control record", source))
646}
647
648const fn control_slot_offset(slot: usize) -> u64 {
649    CONTROL_HEADER_BYTES + (slot * CONTROL_RECORD_BYTES) as u64
650}
651
652fn open_log(directory: &Dir, slot: u8) -> Result<File, TransactionStoreError> {
653    let mut options = OpenOptions::new();
654    options.read(true).write(true).create(true).truncate(false);
655    directory::open_real_file(directory, log_name_for(slot), &options)
656        .map_err(|source| persistent_io("open state log", source))
657}
658
659fn initialize_header(log: &mut File) -> Result<(), TransactionStoreError> {
660    let length = log
661        .metadata()
662        .map_err(|source| persistent_io("inspect state log", source))?
663        .len();
664    if length == 0 {
665        log.write_all(LOG_MAGIC)
666            .map_err(|source| persistent_io("write state header", source))?;
667        log.sync_all()
668            .map_err(|source| persistent_io("sync state header", source))?;
669        return Ok(());
670    }
671    if length < HEADER_BYTES {
672        let mut prefix = vec![0_u8; usize::try_from(length).unwrap_or(0)];
673        log.seek(SeekFrom::Start(0))
674            .map_err(|source| persistent_io("seek state header", source))?;
675        log.read_exact(&mut prefix)
676            .map_err(|source| persistent_io("read partial state header", source))?;
677        if !LOG_MAGIC.starts_with(&prefix) {
678            return Err(TransactionStoreError::PersistentCorrupt {
679                offset: 0,
680                reason: "invalid state-log header",
681            });
682        }
683        log.set_len(0)
684            .map_err(|source| persistent_io("repair partial state header", source))?;
685        log.seek(SeekFrom::Start(0))
686            .map_err(|source| persistent_io("seek repaired state header", source))?;
687        log.write_all(LOG_MAGIC)
688            .map_err(|source| persistent_io("write repaired state header", source))?;
689        log.sync_all()
690            .map_err(|source| persistent_io("sync repaired state header", source))?;
691        return Ok(());
692    }
693    let mut magic = [0_u8; LOG_MAGIC.len()];
694    log.seek(SeekFrom::Start(0))
695        .map_err(|source| persistent_io("seek state header", source))?;
696    log.read_exact(&mut magic)
697        .map_err(|source| persistent_io("read state header", source))?;
698    if &magic != LOG_MAGIC {
699        return Err(TransactionStoreError::PersistentCorrupt {
700            offset: 0,
701            reason: "invalid state-log header",
702        });
703    }
704    Ok(())
705}
706
707#[allow(clippy::too_many_lines)]
708fn refresh(
709    state: &mut PersistentState,
710    log: &mut File,
711    config: FileStoreConfig,
712    control: ControlState,
713) -> Result<(), TransactionStoreError> {
714    if state.generation != control.generation || state.slot != control.slot {
715        state.records.clear();
716        state.offset = HEADER_BYTES;
717        state.generation = control.generation;
718        state.slot = control.slot;
719    }
720    let log_length = log
721        .metadata()
722        .map_err(|source| persistent_io("inspect state log", source))?
723        .len();
724    if log_length > config.max_log_bytes {
725        return Err(TransactionStoreError::PersistentLogLimit {
726            observed: log_length,
727            maximum: config.max_log_bytes,
728        });
729    }
730    if log_length < state.offset || state.offset < HEADER_BYTES {
731        return Err(TransactionStoreError::PersistentCorrupt {
732            offset: state.offset,
733            reason: "state log moved backwards",
734        });
735    }
736    log.seek(SeekFrom::Start(state.offset))
737        .map_err(|source| persistent_io("seek state tail", source))?;
738
739    while state.offset < log_length {
740        let frame_start = state.offset;
741        let mut length_bytes = [0_u8; 4];
742        if !read_exact_or_torn(log, &mut length_bytes)? {
743            truncate_torn_tail(log, state, frame_start)?;
744            break;
745        }
746        let payload_length =
747            usize::try_from(u32::from_le_bytes(length_bytes)).unwrap_or(usize::MAX);
748        if !(OLD_MIN_PAYLOAD_BYTES..=MAX_PAYLOAD_BYTES).contains(&payload_length) {
749            return Err(TransactionStoreError::PersistentCorrupt {
750                offset: frame_start,
751                reason: "invalid state-frame length",
752            });
753        }
754        let mut frame = vec![0_u8; payload_length.saturating_add(DIGEST_BYTES)];
755        if !read_exact_or_torn(log, &mut frame)? {
756            truncate_torn_tail(log, state, frame_start)?;
757            break;
758        }
759        let (payload, encoded_digest) = frame.split_at(payload_length);
760        let actual_digest = BlobId::digest(payload);
761        if actual_digest.as_bytes().as_slice() != encoded_digest {
762            return Err(TransactionStoreError::PersistentCorrupt {
763                offset: frame_start,
764                reason: "state-frame checksum mismatch",
765            });
766        }
767        let record = decode_record(payload, frame_start)?;
768        apply_replayed_record(&mut state.records, record, frame_start, config.max_records)?;
769        state.offset = log
770            .stream_position()
771            .map_err(|source| persistent_io("inspect state tail", source))?;
772    }
773    Ok(())
774}
775
776fn read_exact_or_torn(log: &mut File, output: &mut [u8]) -> Result<bool, TransactionStoreError> {
777    match log.read_exact(output) {
778        Ok(()) => Ok(true),
779        Err(source) if source.kind() == io::ErrorKind::UnexpectedEof => Ok(false),
780        Err(source) => Err(persistent_io("read state frame", source)),
781    }
782}
783
784fn truncate_torn_tail(
785    log: &mut File,
786    state: &mut PersistentState,
787    frame_start: u64,
788) -> Result<(), TransactionStoreError> {
789    log.set_len(frame_start)
790        .map_err(|source| persistent_io("truncate torn state tail", source))?;
791    log.sync_all()
792        .map_err(|source| persistent_io("sync repaired state tail", source))?;
793    log.seek(SeekFrom::Start(frame_start))
794        .map_err(|source| persistent_io("seek repaired state tail", source))?;
795    state.offset = frame_start;
796    Ok(())
797}
798
799fn append_record(
800    log: &mut File,
801    record: &TransactionRecord,
802    offset: &mut u64,
803    config: FileStoreConfig,
804) -> Result<(), TransactionStoreError> {
805    let observed = write_record_frame(log, record, *offset, config)?;
806    log.sync_all()
807        .map_err(|source| persistent_io("sync state frame", source))?;
808    *offset = observed;
809    Ok(())
810}
811
812fn write_record_frame(
813    log: &mut File,
814    record: &TransactionRecord,
815    offset: u64,
816    config: FileStoreConfig,
817) -> Result<u64, TransactionStoreError> {
818    validate_record(record, offset)?;
819    let payload = encode_record(record);
820    let payload_length =
821        u32::try_from(payload.len()).map_err(|_| TransactionStoreError::PersistentCorrupt {
822            offset,
823            reason: "state record cannot be framed",
824        })?;
825    let frame_length = 4_u64
826        .saturating_add(u64::from(payload_length))
827        .saturating_add(DIGEST_BYTES as u64);
828    let observed = offset.saturating_add(frame_length);
829    if observed > config.max_log_bytes {
830        return Err(TransactionStoreError::PersistentLogLimit {
831            observed,
832            maximum: config.max_log_bytes,
833        });
834    }
835    log.seek(SeekFrom::Start(offset))
836        .map_err(|source| persistent_io("seek state append", source))?;
837    log.write_all(&payload_length.to_le_bytes())
838        .and_then(|()| log.write_all(&payload))
839        .and_then(|()| log.write_all(BlobId::digest(&payload).as_bytes()))
840        .map_err(|source| persistent_io("append state frame", source))?;
841    Ok(observed)
842}
843
844fn encode_record(record: &TransactionRecord) -> Vec<u8> {
845    let mut payload = Vec::with_capacity(MAX_PAYLOAD_BYTES);
846    payload.extend_from_slice(record.id().as_bytes());
847    payload.extend_from_slice(record.base_snapshot().as_bytes());
848    payload.push(state_tag(record.state()));
849    match record.artifact() {
850        None => payload.push(0),
851        Some(artifact) => {
852            payload.push(1);
853            payload.extend_from_slice(artifact.as_bytes());
854        }
855    }
856    match record.approval() {
857        None => payload.push(0),
858        Some(grant) => {
859            payload.push(1);
860            payload.extend_from_slice(grant.principal().as_bytes());
861            payload.extend_from_slice(&grant.issued_at_unix_ms().to_le_bytes());
862            payload.extend_from_slice(&grant.expires_at_unix_ms().to_le_bytes());
863        }
864    }
865    payload
866}
867
868fn decode_record(payload: &[u8], offset: u64) -> Result<TransactionRecord, TransactionStoreError> {
869    if payload.len() == OLD_MIN_PAYLOAD_BYTES || payload.len() == OLD_MAX_PAYLOAD_BYTES {
870        return decode_legacy_record(payload, offset);
871    }
872    let valid_sizes = [
873        MIN_PAYLOAD_BYTES,
874        MIN_PAYLOAD_BYTES + ARTIFACT_PAYLOAD_BYTES,
875        MIN_PAYLOAD_BYTES + APPROVAL_PAYLOAD_BYTES,
876        MAX_PAYLOAD_BYTES,
877    ];
878    if !valid_sizes.contains(&payload.len()) {
879        return Err(TransactionStoreError::PersistentCorrupt {
880            offset,
881            reason: "invalid transaction-record size",
882        });
883    }
884    let id = TransactionId::from_bytes(copy_digest(&payload[..32]));
885    let base_snapshot = SnapshotId::from_bytes(copy_digest(&payload[32..64]));
886    let state = decode_state(payload[64]).ok_or(TransactionStoreError::PersistentCorrupt {
887        offset,
888        reason: "unknown transaction state tag",
889    })?;
890    let mut cursor = 66;
891    let artifact = match payload[65] {
892        0 => None,
893        1 => {
894            let end = cursor + ARTIFACT_PAYLOAD_BYTES;
895            let bytes =
896                payload
897                    .get(cursor..end)
898                    .ok_or(TransactionStoreError::PersistentCorrupt {
899                        offset,
900                        reason: "truncated persisted artifact identity",
901                    })?;
902            cursor = end;
903            Some(BlobId::from_bytes(copy_digest(bytes)))
904        }
905        _ => {
906            return Err(TransactionStoreError::PersistentCorrupt {
907                offset,
908                reason: "invalid persisted artifact encoding",
909            });
910        }
911    };
912    let approval_tag = *payload
913        .get(cursor)
914        .ok_or(TransactionStoreError::PersistentCorrupt {
915            offset,
916            reason: "missing persisted approval encoding",
917        })?;
918    cursor += 1;
919    let approval = match approval_tag {
920        0 if cursor == payload.len() => None,
921        1 if cursor + APPROVAL_PAYLOAD_BYTES == payload.len() => {
922            let principal = PrincipalId::from_bytes(copy_digest(&payload[cursor..cursor + 32]));
923            cursor += 32;
924            let issued_at_unix_ms = u64::from_le_bytes(copy_array(&payload[cursor..cursor + 8]));
925            cursor += 8;
926            let expires_at_unix_ms = u64::from_le_bytes(copy_array(&payload[cursor..cursor + 8]));
927            Some(
928                ApprovalGrant::new(id, principal, issued_at_unix_ms, expires_at_unix_ms).map_err(
929                    |_| TransactionStoreError::PersistentCorrupt {
930                        offset,
931                        reason: "invalid persisted approval window",
932                    },
933                )?,
934            )
935        }
936        _ => {
937            return Err(TransactionStoreError::PersistentCorrupt {
938                offset,
939                reason: "invalid persisted approval encoding",
940            });
941        }
942    };
943    let record = TransactionRecord {
944        id,
945        base_snapshot,
946        state,
947        artifact,
948        approval,
949    };
950    validate_record(&record, offset)?;
951    Ok(record)
952}
953
954fn decode_legacy_record(
955    payload: &[u8],
956    offset: u64,
957) -> Result<TransactionRecord, TransactionStoreError> {
958    let id = TransactionId::from_bytes(copy_digest(&payload[..32]));
959    let base_snapshot = SnapshotId::from_bytes(copy_digest(&payload[32..64]));
960    let state = decode_state(payload[64]).ok_or(TransactionStoreError::PersistentCorrupt {
961        offset,
962        reason: "unknown legacy transaction state tag",
963    })?;
964    let approval = match payload[65] {
965        0 if payload.len() == OLD_MIN_PAYLOAD_BYTES => None,
966        1 if payload.len() == OLD_MAX_PAYLOAD_BYTES => {
967            let principal = PrincipalId::from_bytes(copy_digest(&payload[66..98]));
968            let issued_at_unix_ms = u64::from_le_bytes(copy_array(&payload[98..106]));
969            let expires_at_unix_ms = u64::from_le_bytes(copy_array(&payload[106..114]));
970            Some(
971                ApprovalGrant::new(id, principal, issued_at_unix_ms, expires_at_unix_ms).map_err(
972                    |_| TransactionStoreError::PersistentCorrupt {
973                        offset,
974                        reason: "invalid legacy approval window",
975                    },
976                )?,
977            )
978        }
979        _ => {
980            return Err(TransactionStoreError::PersistentCorrupt {
981                offset,
982                reason: "invalid legacy approval encoding",
983            });
984        }
985    };
986    let record = TransactionRecord {
987        id,
988        base_snapshot,
989        state,
990        artifact: None,
991        approval,
992    };
993    validate_record(&record, offset)?;
994    Ok(record)
995}
996
997fn validate_record(record: &TransactionRecord, offset: u64) -> Result<(), TransactionStoreError> {
998    if let Some(grant) = record.approval()
999        && grant.transaction() != record.id()
1000    {
1001        return Err(TransactionStoreError::PersistentCorrupt {
1002            offset,
1003            reason: "approval does not bind its transaction",
1004        });
1005    }
1006    let state_forbids_approval = matches!(
1007        record.state(),
1008        TransactionState::Created
1009            | TransactionState::Running
1010            | TransactionState::VirtualComplete
1011            | TransactionState::Denied
1012            | TransactionState::AutoApproved
1013            | TransactionState::PendingApproval
1014    );
1015    if state_forbids_approval && record.approval().is_some() {
1016        return Err(TransactionStoreError::PersistentCorrupt {
1017            offset,
1018            reason: "approval exists before the approved state",
1019        });
1020    }
1021    if matches!(
1022        record.state(),
1023        TransactionState::Approved | TransactionState::Expired
1024    ) && record.approval().is_none()
1025    {
1026        return Err(TransactionStoreError::PersistentCorrupt {
1027            offset,
1028            reason: "manual approval state lacks its grant",
1029        });
1030    }
1031    Ok(())
1032}
1033
1034fn apply_replayed_record(
1035    records: &mut BTreeMap<TransactionId, TransactionRecord>,
1036    record: TransactionRecord,
1037    offset: u64,
1038    max_records: usize,
1039) -> Result<(), TransactionStoreError> {
1040    match records.get(&record.id()) {
1041        None => {
1042            if records.len() >= max_records {
1043                return Err(TransactionStoreError::PersistentRecordLimit {
1044                    observed: records.len().saturating_add(1),
1045                    maximum: max_records,
1046                });
1047            }
1048        }
1049        Some(previous) => {
1050            if previous.base_snapshot() != record.base_snapshot() {
1051                return Err(TransactionStoreError::PersistentCorrupt {
1052                    offset,
1053                    reason: "transaction base snapshot changed",
1054                });
1055            }
1056            if previous.artifact() != record.artifact() {
1057                return Err(TransactionStoreError::PersistentCorrupt {
1058                    offset,
1059                    reason: "transaction artifact binding changed",
1060                });
1061            }
1062            if !previous.state().can_transition_to(record.state()) {
1063                return Err(TransactionStoreError::PersistentCorrupt {
1064                    offset,
1065                    reason: "invalid persisted transaction transition",
1066                });
1067            }
1068            match (previous.approval(), record.approval()) {
1069                (None, Some(_))
1070                    if previous.state() == TransactionState::PendingApproval
1071                        && record.state() == TransactionState::Approved => {}
1072                (Some(before), Some(after)) if before == after => {}
1073                (None, None) => {}
1074                _ => {
1075                    return Err(TransactionStoreError::PersistentCorrupt {
1076                        offset,
1077                        reason: "persisted approval binding changed",
1078                    });
1079                }
1080            }
1081        }
1082    }
1083    records.insert(record.id(), record);
1084    Ok(())
1085}
1086
1087const fn state_tag(state: TransactionState) -> u8 {
1088    match state {
1089        TransactionState::Created => 1,
1090        TransactionState::Running => 2,
1091        TransactionState::VirtualComplete => 3,
1092        TransactionState::Denied => 4,
1093        TransactionState::AutoApproved => 5,
1094        TransactionState::PendingApproval => 6,
1095        TransactionState::Approved => 7,
1096        TransactionState::Reserved => 8,
1097        TransactionState::Revalidating => 9,
1098        TransactionState::Committing => 10,
1099        TransactionState::Committed => 11,
1100        TransactionState::Stale => 12,
1101        TransactionState::Expired => 13,
1102        TransactionState::RecoveryRequired => 14,
1103        TransactionState::Failed => 15,
1104        TransactionState::Rejected => 16,
1105        _ => 0,
1106    }
1107}
1108
1109const fn decode_state(tag: u8) -> Option<TransactionState> {
1110    match tag {
1111        1 => Some(TransactionState::Created),
1112        2 => Some(TransactionState::Running),
1113        3 => Some(TransactionState::VirtualComplete),
1114        4 => Some(TransactionState::Denied),
1115        5 => Some(TransactionState::AutoApproved),
1116        6 => Some(TransactionState::PendingApproval),
1117        7 => Some(TransactionState::Approved),
1118        8 => Some(TransactionState::Reserved),
1119        9 => Some(TransactionState::Revalidating),
1120        10 => Some(TransactionState::Committing),
1121        11 => Some(TransactionState::Committed),
1122        12 => Some(TransactionState::Stale),
1123        13 => Some(TransactionState::Expired),
1124        14 => Some(TransactionState::RecoveryRequired),
1125        15 => Some(TransactionState::Failed),
1126        16 => Some(TransactionState::Rejected),
1127        _ => None,
1128    }
1129}
1130
1131fn copy_digest(bytes: &[u8]) -> [u8; 32] {
1132    let mut output = [0_u8; 32];
1133    output.copy_from_slice(bytes);
1134    output
1135}
1136
1137fn copy_array<const N: usize>(bytes: &[u8]) -> [u8; N] {
1138    let mut output = [0_u8; N];
1139    output.copy_from_slice(bytes);
1140    output
1141}
1142
1143fn validate_config(config: FileStoreConfig) -> Result<(), TransactionStoreError> {
1144    if config.max_log_bytes < HEADER_BYTES {
1145        return Err(TransactionStoreError::PersistentLogLimit {
1146            observed: HEADER_BYTES,
1147            maximum: config.max_log_bytes,
1148        });
1149    }
1150    if config.max_records == 0 {
1151        return Err(TransactionStoreError::PersistentRecordLimit {
1152            observed: 1,
1153            maximum: 0,
1154        });
1155    }
1156    Ok(())
1157}
1158
1159#[allow(
1160    clippy::needless_pass_by_value,
1161    reason = "map_err owns io::Error; the stable store error intentionally retains only ErrorKind"
1162)]
1163fn persistent_io(operation: &'static str, source: io::Error) -> TransactionStoreError {
1164    TransactionStoreError::PersistentIo {
1165        operation,
1166        kind: source.kind(),
1167    }
1168}
1169
1170#[cfg(test)]
1171mod tests {
1172    use std::fs::{self, OpenOptions};
1173    use std::io::Write;
1174    use std::path::{Path, PathBuf};
1175    use std::sync::atomic::{AtomicU64, Ordering};
1176    use std::sync::{Arc, Barrier};
1177    use std::thread;
1178
1179    use vsh_types::{BlobId, PrincipalId, SnapshotId, TransactionId, TransactionState};
1180
1181    use super::{FileStoreConfig, FileTransactionStore};
1182    use crate::{ApprovalGrant, TransactionRecord, TransactionStore, TransactionStoreError};
1183
1184    static TEST_SEQUENCE: AtomicU64 = AtomicU64::new(0);
1185
1186    struct TestDirectory(PathBuf);
1187
1188    impl TestDirectory {
1189        fn new(name: &str) -> Self {
1190            let sequence = TEST_SEQUENCE.fetch_add(1, Ordering::Relaxed);
1191            let path = std::env::temp_dir().join(format!(
1192                "vsh-file-store-{name}-{}-{sequence}",
1193                std::process::id()
1194            ));
1195            fs::create_dir(&path).unwrap();
1196            Self(path)
1197        }
1198
1199        fn path(&self) -> &Path {
1200            &self.0
1201        }
1202    }
1203
1204    impl Drop for TestDirectory {
1205        fn drop(&mut self) {
1206            let _ = fs::remove_dir_all(&self.0);
1207        }
1208    }
1209
1210    fn id(byte: u8) -> TransactionId {
1211        TransactionId::from_bytes([byte; 32])
1212    }
1213
1214    fn snapshot(byte: u8) -> SnapshotId {
1215        SnapshotId::from_bytes([byte; 32])
1216    }
1217
1218    const fn compacting_config() -> FileStoreConfig {
1219        FileStoreConfig {
1220            max_log_bytes: 180,
1221            max_records: 16,
1222        }
1223    }
1224
1225    #[test]
1226    fn lifecycle_and_approval_survive_reopen() {
1227        let directory = TestDirectory::new("reopen");
1228        let store =
1229            FileTransactionStore::open(directory.path(), FileStoreConfig::default()).unwrap();
1230        let artifact = BlobId::digest(b"pending artifact");
1231        let mut record = TransactionRecord::new(id(1), snapshot(2)).with_artifact(artifact);
1232        record.transition(TransactionState::Running).unwrap();
1233        record
1234            .transition(TransactionState::VirtualComplete)
1235            .unwrap();
1236        record
1237            .transition(TransactionState::PendingApproval)
1238            .unwrap();
1239        store.create(record).unwrap();
1240        let grant =
1241            ApprovalGrant::new(id(1), PrincipalId::digest_label("independent"), 10, 20).unwrap();
1242        store.approve(id(1), grant).unwrap();
1243
1244        let reopened =
1245            FileTransactionStore::open(directory.path(), FileStoreConfig::default()).unwrap();
1246        let loaded = reopened.get(id(1)).unwrap();
1247        assert_eq!(loaded.state(), TransactionState::Approved);
1248        assert_eq!(loaded.artifact(), Some(artifact));
1249        assert_eq!(loaded.approval(), Some(grant));
1250        assert_eq!(reopened.reserve(id(1), 11).unwrap().transaction(), id(1));
1251    }
1252
1253    #[test]
1254    fn independent_handles_have_one_cross_process_style_reservation_winner() {
1255        let directory = TestDirectory::new("reserve");
1256        let first =
1257            FileTransactionStore::open(directory.path(), FileStoreConfig::default()).unwrap();
1258        let second =
1259            FileTransactionStore::open(directory.path(), FileStoreConfig::default()).unwrap();
1260        let mut record = TransactionRecord::new(id(3), snapshot(4));
1261        record.transition(TransactionState::Running).unwrap();
1262        record
1263            .transition(TransactionState::VirtualComplete)
1264            .unwrap();
1265        record.transition(TransactionState::AutoApproved).unwrap();
1266        first.create(record).unwrap();
1267
1268        let barrier = Arc::new(Barrier::new(3));
1269        let handles = [first, second]
1270            .into_iter()
1271            .map(|store| {
1272                let barrier = Arc::clone(&barrier);
1273                thread::spawn(move || {
1274                    barrier.wait();
1275                    store.reserve(id(3), 0)
1276                })
1277            })
1278            .collect::<Vec<_>>();
1279        barrier.wait();
1280        let results = handles
1281            .into_iter()
1282            .map(|handle| handle.join().unwrap())
1283            .collect::<Vec<_>>();
1284        assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1);
1285        assert_eq!(
1286            results
1287                .iter()
1288                .filter(|result| matches!(result, Err(TransactionStoreError::NotReservable { .. })))
1289                .count(),
1290            1
1291        );
1292    }
1293
1294    #[test]
1295    fn torn_tail_is_truncated_to_last_checksummed_state() {
1296        let directory = TestDirectory::new("torn");
1297        let store =
1298            FileTransactionStore::open(directory.path(), FileStoreConfig::default()).unwrap();
1299        store
1300            .create(TransactionRecord::new(id(5), snapshot(6)))
1301            .unwrap();
1302        let log_path = store.active_log_path().unwrap();
1303        let valid_length = fs::metadata(&log_path).unwrap().len();
1304        let mut file = OpenOptions::new().append(true).open(&log_path).unwrap();
1305        file.write_all(&[114, 0, 0]).unwrap();
1306        file.sync_all().unwrap();
1307        drop(store);
1308
1309        let reopened =
1310            FileTransactionStore::open(directory.path(), FileStoreConfig::default()).unwrap();
1311        assert_eq!(
1312            reopened.get(id(5)).unwrap().state(),
1313            TransactionState::Created
1314        );
1315        assert_eq!(
1316            fs::metadata(reopened.active_log_path().unwrap())
1317                .unwrap()
1318                .len(),
1319            valid_length
1320        );
1321    }
1322
1323    #[test]
1324    fn complete_frame_with_bad_checksum_fails_closed_even_at_end_of_log() {
1325        let directory = TestDirectory::new("checksum");
1326        let store =
1327            FileTransactionStore::open(directory.path(), FileStoreConfig::default()).unwrap();
1328        store
1329            .create(TransactionRecord::new(id(6), snapshot(7)))
1330            .unwrap();
1331        let log_path = store.active_log_path().unwrap();
1332        drop(store);
1333
1334        let mut bytes = fs::read(&log_path).unwrap();
1335        let header_bytes = usize::try_from(super::HEADER_BYTES).unwrap();
1336        let payload_length =
1337            u32::from_le_bytes(bytes[header_bytes..header_bytes + 4].try_into().unwrap()) as usize;
1338        let checksum_start = header_bytes + 4 + payload_length;
1339        bytes[checksum_start] ^= 0x80;
1340        fs::write(&log_path, bytes).unwrap();
1341
1342        assert!(matches!(
1343            FileTransactionStore::open(directory.path(), FileStoreConfig::default()),
1344            Err(TransactionStoreError::PersistentCorrupt {
1345                reason: "state-frame checksum mismatch",
1346                ..
1347            })
1348        ));
1349    }
1350
1351    #[test]
1352    fn partial_matching_header_is_repaired_but_wrong_header_fails_closed() {
1353        let repairable = TestDirectory::new("partial-header");
1354        let repairable_log = repairable.path().join(super::LOG_FILE);
1355        fs::write(&repairable_log, &super::LOG_MAGIC[..4]).unwrap();
1356        let store =
1357            FileTransactionStore::open(repairable.path(), FileStoreConfig::default()).unwrap();
1358        drop(store);
1359        assert_eq!(
1360            &fs::read(repairable_log).unwrap()[..super::LOG_MAGIC.len()],
1361            super::LOG_MAGIC
1362        );
1363
1364        let corrupt = TestDirectory::new("wrong-header");
1365        fs::write(corrupt.path().join(super::LOG_FILE), b"NOTVSH01").unwrap();
1366        assert!(matches!(
1367            FileTransactionStore::open(corrupt.path(), FileStoreConfig::default()),
1368            Err(TransactionStoreError::PersistentCorrupt {
1369                reason: "invalid state-log header",
1370                ..
1371            })
1372        ));
1373    }
1374
1375    #[test]
1376    fn bounded_compaction_switches_logs_and_stale_handles_replay_the_new_generation() {
1377        let directory = TestDirectory::new("compact");
1378        let first = FileTransactionStore::open(directory.path(), compacting_config()).unwrap();
1379        let second = FileTransactionStore::open(directory.path(), compacting_config()).unwrap();
1380        first
1381            .create(TransactionRecord::new(id(7), snapshot(8)))
1382            .unwrap();
1383        let initial_log = first.active_log_path().unwrap();
1384
1385        first
1386            .compare_and_transition(id(7), TransactionState::Created, TransactionState::Running)
1387            .unwrap();
1388        let first_compacted_log = first.active_log_path().unwrap();
1389        assert_ne!(first_compacted_log, initial_log);
1390        assert!(
1391            fs::metadata(&first_compacted_log).unwrap().len() <= compacting_config().max_log_bytes
1392        );
1393        assert_eq!(
1394            second.get(id(7)).unwrap().state(),
1395            TransactionState::Running
1396        );
1397
1398        second
1399            .compare_and_transition(
1400                id(7),
1401                TransactionState::Running,
1402                TransactionState::VirtualComplete,
1403            )
1404            .unwrap();
1405        assert_eq!(second.active_log_path().unwrap(), initial_log);
1406        drop(first);
1407        drop(second);
1408
1409        let reopened = FileTransactionStore::open(directory.path(), compacting_config()).unwrap();
1410        assert_eq!(
1411            reopened.get(id(7)).unwrap().state(),
1412            TransactionState::VirtualComplete
1413        );
1414    }
1415
1416    #[test]
1417    fn torn_control_tail_and_inactive_log_never_replace_the_durable_generation() {
1418        let directory = TestDirectory::new("compact-torn-control");
1419        let store = FileTransactionStore::open(directory.path(), compacting_config()).unwrap();
1420        store
1421            .create(TransactionRecord::new(id(9), snapshot(10)))
1422            .unwrap();
1423        let inactive_log = store.active_log_path().unwrap();
1424        store
1425            .compare_and_transition(id(9), TransactionState::Created, TransactionState::Running)
1426            .unwrap();
1427        let active_log = store.active_log_path().unwrap();
1428        assert_ne!(active_log, inactive_log);
1429        drop(store);
1430
1431        fs::write(&inactive_log, b"uncommitted inactive generation").unwrap();
1432        let lock_path = directory.path().join(super::LOCK_FILE);
1433        let valid_control_length = fs::metadata(&lock_path).unwrap().len();
1434        let mut lock = OpenOptions::new().append(true).open(&lock_path).unwrap();
1435        lock.write_all(&[0xA5; 7]).unwrap();
1436        lock.sync_all().unwrap();
1437        drop(lock);
1438
1439        let reopened = FileTransactionStore::open(directory.path(), compacting_config()).unwrap();
1440        assert_eq!(reopened.active_log_path().unwrap(), active_log);
1441        assert_eq!(
1442            reopened.get(id(9)).unwrap().state(),
1443            TransactionState::Running
1444        );
1445        assert_eq!(fs::metadata(lock_path).unwrap().len(), valid_control_length);
1446    }
1447
1448    #[test]
1449    fn compaction_that_cannot_fit_all_latest_records_leaves_active_state_unchanged() {
1450        let directory = TestDirectory::new("compact-overflow");
1451        let store = FileTransactionStore::open(directory.path(), compacting_config()).unwrap();
1452        store
1453            .create(TransactionRecord::new(id(11), snapshot(12)))
1454            .unwrap();
1455        let active_log = store.active_log_path().unwrap();
1456
1457        assert!(matches!(
1458            store.create(TransactionRecord::new(id(13), snapshot(14))),
1459            Err(TransactionStoreError::PersistentLogLimit { .. })
1460        ));
1461        assert_eq!(store.active_log_path().unwrap(), active_log);
1462        assert_eq!(
1463            store.get(id(11)).unwrap().state(),
1464            TransactionState::Created
1465        );
1466        assert!(matches!(
1467            store.get(id(13)),
1468            Err(TransactionStoreError::NotFound { .. })
1469        ));
1470    }
1471
1472    #[test]
1473    fn invalid_store_bounds_fail_before_creating_files() {
1474        let directory = TestDirectory::new("invalid-bounds");
1475        let too_small = FileStoreConfig {
1476            max_log_bytes: 1,
1477            max_records: 1,
1478        };
1479        assert!(matches!(
1480            FileTransactionStore::open(directory.path(), too_small),
1481            Err(TransactionStoreError::PersistentLogLimit { .. })
1482        ));
1483        let zero_records = FileStoreConfig {
1484            max_log_bytes: 1024,
1485            max_records: 0,
1486        };
1487        assert!(matches!(
1488            FileTransactionStore::open(directory.path(), zero_records),
1489            Err(TransactionStoreError::PersistentRecordLimit { .. })
1490        ));
1491        assert_eq!(fs::read_dir(directory.path()).unwrap().count(), 0);
1492    }
1493
1494    #[cfg(unix)]
1495    #[test]
1496    fn internal_state_file_symlink_cannot_redirect_open() {
1497        use std::os::unix::fs::symlink;
1498
1499        let directory = TestDirectory::new("state-symlink");
1500        let outside = TestDirectory::new("state-symlink-outside");
1501        let target = outside.path().join("target");
1502        fs::write(&target, b"outside must remain unchanged").unwrap();
1503        symlink(&target, directory.path().join(super::LOCK_FILE)).unwrap();
1504
1505        assert!(matches!(
1506            FileTransactionStore::open(directory.path(), FileStoreConfig::default()),
1507            Err(TransactionStoreError::PersistentIo { .. })
1508        ));
1509        assert_eq!(fs::read(&target).unwrap(), b"outside must remain unchanged");
1510        assert_eq!(fs::read_dir(outside.path()).unwrap().count(), 1);
1511    }
1512}