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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
38pub struct FileStoreConfig {
39 pub max_log_bytes: u64,
43 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#[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 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 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 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}