1use core::pin::Pin;
27use futures::{Stream, StreamExt};
28use std::fs;
29use std::io::Write;
30#[cfg(unix)]
31use std::os::fd::AsRawFd;
32use std::path::{Path, PathBuf};
33use std::sync::Arc;
34use std::sync::atomic::{AtomicU64, Ordering};
35use tokio::sync::watch;
36use tokio_stream::wrappers::WatchStream;
37use tsoracle_consensus::{ConsensusDriver, ConsensusError, LeaderState};
38use tsoracle_core::{Epoch, LeaseRecord, PHYSICAL_MS_MAX};
39
40use crate::{dense_record, lease_record, record};
41
42pub const DEFAULT_DENSE_CARDINALITY_CAP: u64 = 10_000;
45
46#[derive(Debug)]
47struct DenseState {
48 map: std::collections::BTreeMap<String, u64>,
49 cap: u64,
50}
51
52#[derive(Debug, thiserror::Error)]
53pub enum FileDriverError {
54 #[error("io: {0}")]
55 Io(#[from] std::io::Error),
56 #[error("decode: {0}")]
57 Decode(#[from] record::RecordError),
58 #[error("physical_ms {0} exceeds 46-bit maximum")]
59 PhysicalMsOutOfRange(u64),
60 #[error("state directory {path} is already locked by another FileDriver: {source}")]
61 AlreadyLocked {
62 path: PathBuf,
63 #[source]
64 source: std::io::Error,
65 },
66}
67
68#[derive(Debug)]
69pub struct FileDriver {
70 dir: PathBuf,
71 state: Arc<AtomicU64>,
76 write_lock: tokio::sync::Mutex<()>,
77 #[expect(
82 dead_code,
83 reason = "kept to hold the watch channel open for leader_rx consumers"
84 )]
85 leader_tx: watch::Sender<LeaderState>,
86 leader_rx: watch::Receiver<LeaderState>,
87 _lock: fs::File,
92 dense: tokio::sync::Mutex<DenseState>,
95 leases: tokio::sync::Mutex<Vec<LeaseRecord>>,
97}
98
99impl FileDriver {
100 pub fn open_or_init(dir: impl AsRef<Path>) -> Result<Arc<Self>, FileDriverError> {
112 let dir = dir.as_ref().to_path_buf();
113 fs::create_dir_all(&dir)?;
114
115 let lock_path = dir.join("LOCK");
120 let lock_file = fs::OpenOptions::new()
121 .create(true)
122 .read(true)
123 .write(true)
124 .truncate(false)
125 .open(&lock_path)?;
126 acquire_exclusive_lock(&lock_file, &lock_path)?;
127
128 let state_path = dir.join("state");
129 let current = if state_path.exists() {
130 let bytes = fs::read(&state_path)?;
131 let high_water = record::decode(&bytes)?;
132 if high_water > PHYSICAL_MS_MAX {
133 return Err(FileDriverError::PhysicalMsOutOfRange(high_water));
134 }
135 high_water
136 } else {
137 0
138 };
139 let dense_path = dir.join("dense");
140 let dense_state = if dense_path.exists() {
141 let bytes = fs::read(&dense_path)?;
142 let (map, cap) = dense_record::decode(&bytes)
143 .map_err(|e| FileDriverError::Io(std::io::Error::other(e)))?;
144 DenseState { map, cap }
145 } else {
146 DenseState {
147 map: std::collections::BTreeMap::new(),
148 cap: DEFAULT_DENSE_CARDINALITY_CAP,
149 }
150 };
151 let leases_path = dir.join("leases");
152 let leases = if leases_path.exists() {
153 let bytes = fs::read(&leases_path)?;
154 lease_record::decode(&bytes)
155 .map_err(|e| FileDriverError::Io(std::io::Error::other(e)))?
156 } else {
157 Vec::new()
158 };
159
160 let (tx, rx) = watch::channel(LeaderState::Leader { epoch: Epoch::ZERO });
161 Ok(Arc::new(FileDriver {
162 dir,
163 state: Arc::new(AtomicU64::new(current)),
164 write_lock: tokio::sync::Mutex::new(()),
165 leader_tx: tx,
166 leader_rx: rx,
167 _lock: lock_file,
168 dense: tokio::sync::Mutex::new(dense_state),
169 leases: tokio::sync::Mutex::new(leases),
170 }))
171 }
172
173 pub fn init_seeded(
181 dir: impl AsRef<Path>,
182 seed_physical_ms: u64,
183 ) -> Result<(), FileDriverError> {
184 if seed_physical_ms > PHYSICAL_MS_MAX {
185 return Err(FileDriverError::PhysicalMsOutOfRange(seed_physical_ms));
186 }
187 let dir = dir.as_ref();
188 fs::create_dir_all(dir)?;
189 let state_path = dir.join("state");
190 if state_path.exists() {
191 return Err(FileDriverError::Io(std::io::Error::new(
192 std::io::ErrorKind::AlreadyExists,
193 "state file already exists; refusing to overwrite",
194 )));
195 }
196 write_record(dir, seed_physical_ms)?;
197 Ok(())
198 }
199}
200
201fn acquire_exclusive_lock(lock_file: &fs::File, lock_path: &Path) -> Result<(), FileDriverError> {
213 use fs2::FileExt;
214 match lock_file.try_lock_exclusive() {
215 Ok(()) => Ok(()),
216 Err(err) if err.raw_os_error() == fs2::lock_contended_error().raw_os_error() => {
217 Err(FileDriverError::AlreadyLocked {
218 path: lock_path.to_path_buf(),
219 source: err,
220 })
221 }
222 Err(err) => Err(FileDriverError::Io(err)),
223 }
224}
225
226fn write_record(dir: &Path, high_water: u64) -> Result<(), FileDriverError> {
227 tsoracle_failpoint::failpoint!(
228 "file_driver::before_write",
229 |arg: Option<String>| -> Result<(), FileDriverError> {
230 let _ = arg; Err(FileDriverError::Io(std::io::Error::other(
232 "failpoint: file_driver::before_write",
233 )))
234 }
235 );
236
237 let tmp = dir.join("state.tmp");
238 let final_path = dir.join("state");
239 let bytes = record::encode(high_water);
240
241 let mut file = fs::OpenOptions::new()
242 .create(true)
243 .write(true)
244 .truncate(true)
245 .open(&tmp)?;
246 file.write_all(&bytes)?;
247 file.sync_all()?;
248 drop(file);
249
250 tsoracle_failpoint::failpoint!(
251 "file_driver::after_tmp_fsync_before_rename",
252 |arg: Option<String>| -> Result<(), FileDriverError> {
253 let _ = arg;
254 Err(FileDriverError::Io(std::io::Error::other(
255 "failpoint: file_driver::after_tmp_fsync_before_rename",
256 )))
257 }
258 );
259
260 fs::rename(&tmp, &final_path)?;
261
262 tsoracle_failpoint::failpoint!("file_driver::after_rename_before_dir_fsync");
263
264 #[cfg(unix)]
281 {
282 let dir_file = fs::File::open(dir)?;
283 let fd = dir_file.as_raw_fd();
284 let rc = unsafe { libc::fsync(fd) };
286 if rc != 0 {
287 return Err(FileDriverError::Io(std::io::Error::last_os_error()));
288 }
289 }
290 #[cfg(not(unix))]
291 {
292 let final_file = fs::OpenOptions::new().write(true).open(&final_path)?;
296 final_file.sync_all()?;
297 }
298 Ok(())
299}
300
301fn write_dense_record(
304 dir: &Path,
305 map: &std::collections::BTreeMap<String, u64>,
306 cap: u64,
307) -> Result<(), FileDriverError> {
308 tsoracle_failpoint::failpoint!("file_driver::dense::before_write", |_arg: Option<
309 String,
310 >|
311 -> Result<
312 (),
313 FileDriverError,
314 > {
315 Err(FileDriverError::Io(std::io::Error::other(
316 "failpoint: file_driver::dense::before_write",
317 )))
318 });
319
320 let tmp = dir.join("dense.tmp");
321 let final_path = dir.join("dense");
322 let bytes = dense_record::encode(map, cap);
323
324 let mut file = fs::OpenOptions::new()
325 .create(true)
326 .write(true)
327 .truncate(true)
328 .open(&tmp)?;
329 file.write_all(&bytes)?;
330 file.sync_all()?;
331 drop(file);
332
333 tsoracle_failpoint::failpoint!(
334 "file_driver::dense::after_tmp_fsync_before_rename",
335 |_arg: Option<String>| -> Result<(), FileDriverError> {
336 Err(FileDriverError::Io(std::io::Error::other(
337 "failpoint: file_driver::dense::after_tmp_fsync_before_rename",
338 )))
339 }
340 );
341
342 fs::rename(&tmp, &final_path)?;
343
344 tsoracle_failpoint::failpoint!("file_driver::dense::after_rename_before_dir_fsync");
345
346 #[cfg(unix)]
347 {
348 let dir_file = fs::File::open(dir)?;
349 let fd = dir_file.as_raw_fd();
350 let rc = unsafe { libc::fsync(fd) };
352 if rc != 0 {
353 return Err(FileDriverError::Io(std::io::Error::last_os_error()));
354 }
355 }
356 #[cfg(not(unix))]
357 {
358 let final_file = fs::OpenOptions::new().write(true).open(&final_path)?;
359 final_file.sync_all()?;
360 }
361 Ok(())
362}
363
364fn write_lease_record(dir: &Path, records: &[LeaseRecord]) -> Result<(), FileDriverError> {
367 tsoracle_failpoint::failpoint!("file_driver::leases::before_write", |_arg: Option<
368 String,
369 >|
370 -> Result<
371 (),
372 FileDriverError,
373 > {
374 Err(FileDriverError::Io(std::io::Error::other(
375 "failpoint: file_driver::leases::before_write",
376 )))
377 });
378
379 let tmp = dir.join("leases.tmp");
380 let final_path = dir.join("leases");
381 let bytes = lease_record::encode(records);
382
383 let mut file = fs::OpenOptions::new()
384 .create(true)
385 .write(true)
386 .truncate(true)
387 .open(&tmp)?;
388 file.write_all(&bytes)?;
389 file.sync_all()?;
390 drop(file);
391
392 tsoracle_failpoint::failpoint!(
393 "file_driver::leases::after_tmp_fsync_before_rename",
394 |_arg: Option<String>| -> Result<(), FileDriverError> {
395 Err(FileDriverError::Io(std::io::Error::other(
396 "failpoint: file_driver::leases::after_tmp_fsync_before_rename",
397 )))
398 }
399 );
400
401 fs::rename(&tmp, &final_path)?;
402
403 tsoracle_failpoint::failpoint!("file_driver::leases::after_rename_before_dir_fsync");
404
405 #[cfg(unix)]
406 {
407 let dir_file = fs::File::open(dir)?;
408 let fd = dir_file.as_raw_fd();
409 let rc = unsafe { libc::fsync(fd) };
411 if rc != 0 {
412 return Err(FileDriverError::Io(std::io::Error::last_os_error()));
413 }
414 }
415 #[cfg(not(unix))]
416 {
417 let final_file = fs::OpenOptions::new().write(true).open(&final_path)?;
418 final_file.sync_all()?;
419 }
420 Ok(())
421}
422
423#[async_trait::async_trait]
424impl ConsensusDriver for FileDriver {
425 fn leadership_events(&self) -> Pin<Box<dyn Stream<Item = LeaderState> + Send>> {
426 Box::pin(WatchStream::new(self.leader_rx.clone()).boxed())
427 }
428
429 async fn load_high_water(&self) -> Result<u64, ConsensusError> {
430 Ok(self.state.load(Ordering::Acquire))
432 }
433
434 async fn persist_high_water(
435 &self,
436 at_least: u64,
437 _epoch: Epoch,
438 ) -> Result<u64, ConsensusError> {
439 tsoracle_consensus::reject_out_of_range_advance(at_least)?;
442
443 let _guard = self.write_lock.lock().await;
446
447 let current = self.state.load(Ordering::Acquire);
448 if at_least <= current {
449 return Ok(current);
450 }
451 let target = at_least;
452
453 let dir = self.dir.clone();
454 tokio::task::spawn_blocking(move || {
455 tsoracle_failpoint::failpoint!("file_driver::write_blocked");
456 write_record(&dir, target)
457 })
458 .await
459 .map_err(|e| ConsensusError::PermanentDriver(Box::new(std::io::Error::other(e))))?
462 .map_err(|e| ConsensusError::PermanentDriver(Box::new(e)))?;
467
468 self.state.store(target, Ordering::Release);
471 Ok(target)
472 }
473
474 async fn load_dense_seq(&self, key: &tsoracle_core::SeqKey) -> Result<u64, ConsensusError> {
475 let dense = self.dense.lock().await;
476 Ok(dense.map.get(key.as_str()).copied().unwrap_or(0))
477 }
478
479 async fn advance_dense(
480 &self,
481 key: &tsoracle_core::SeqKey,
482 count: u32,
483 _expected_epoch: Epoch,
484 ) -> Result<u64, ConsensusError> {
485 let mut dense = self.dense.lock().await;
486
487 let present = dense.map.contains_key(key.as_str());
488 if !present && dense.map.len() as u64 >= dense.cap {
489 return Err(ConsensusError::SeqKeyCardinalityExceeded { cap: dense.cap });
490 }
491 let start = dense.map.get(key.as_str()).copied().unwrap_or(0);
492 let next = start
493 .checked_add(u64::from(count))
494 .ok_or(ConsensusError::SeqOverflow)?;
495
496 let mut new_map = dense.map.clone();
498 new_map.insert(key.as_str().to_string(), next);
499 let cap = dense.cap;
500 let dir = self.dir.clone();
501 let to_write = new_map.clone();
502 tokio::task::spawn_blocking(move || write_dense_record(&dir, &to_write, cap))
503 .await
504 .map_err(|e| ConsensusError::PermanentDriver(Box::new(std::io::Error::other(e))))?
505 .map_err(|e| ConsensusError::PermanentDriver(Box::new(e)))?;
506
507 dense.map = new_map;
508 Ok(start)
509 }
510
511 async fn advance_dense_batch(
512 &self,
513 entries: &[(tsoracle_core::SeqKey, u32)],
514 _expected_epoch: Epoch,
515 ) -> Result<Vec<u64>, ConsensusError> {
516 if entries.is_empty() {
519 return Ok(Vec::new());
520 }
521
522 let mut dense = self.dense.lock().await;
523
524 let new_keys: std::collections::BTreeSet<&str> = entries
527 .iter()
528 .map(|(key, _)| key.as_str())
529 .filter(|k| !dense.map.contains_key(*k))
530 .collect();
531 if dense.map.len() as u64 + new_keys.len() as u64 > dense.cap {
532 return Err(ConsensusError::SeqKeyCardinalityExceeded { cap: dense.cap });
533 }
534
535 let mut scratch: std::collections::BTreeMap<String, u64> =
539 std::collections::BTreeMap::new();
540 let mut starts: Vec<u64> = Vec::with_capacity(entries.len());
541 for (key, count) in entries {
542 let key_str = key.as_str();
543 let running = scratch
544 .get(key_str)
545 .copied()
546 .or_else(|| dense.map.get(key_str).copied())
547 .unwrap_or(0);
548 starts.push(running);
549 let next = running
550 .checked_add(u64::from(*count))
551 .ok_or(ConsensusError::SeqOverflow)?;
552 scratch.insert(key_str.to_string(), next);
553 }
554
555 let mut new_map = dense.map.clone();
558 for (key_str, next) in &scratch {
559 new_map.insert(key_str.clone(), *next);
560 }
561 let cap = dense.cap;
562 let dir = self.dir.clone();
563 let to_write = new_map.clone();
564 tokio::task::spawn_blocking(move || write_dense_record(&dir, &to_write, cap))
565 .await
566 .map_err(|e| ConsensusError::PermanentDriver(Box::new(std::io::Error::other(e))))?
567 .map_err(|e| ConsensusError::PermanentDriver(Box::new(e)))?;
568
569 dense.map = new_map;
570 Ok(starts)
571 }
572
573 async fn load_leases(&self) -> Result<Vec<LeaseRecord>, ConsensusError> {
574 Ok(self.leases.lock().await.clone())
575 }
576
577 async fn persist_leases(
578 &self,
579 live: &[LeaseRecord],
580 _epoch: Epoch,
581 ) -> Result<(), ConsensusError> {
582 let _guard = self.write_lock.lock().await;
583 let dir = self.dir.clone();
584 let to_write = live.to_vec();
585 tokio::task::spawn_blocking(move || write_lease_record(&dir, &to_write))
586 .await
587 .map_err(|e| ConsensusError::PermanentDriver(Box::new(std::io::Error::other(e))))?
588 .map_err(|e| ConsensusError::PermanentDriver(Box::new(e)))?;
589 *self.leases.lock().await = live.to_vec();
590 Ok(())
591 }
592}
593
594#[cfg(test)]
595mod dense_tests {
596 use super::*;
597 use tsoracle_core::{Epoch, SeqKey};
598
599 fn key(s: &str) -> SeqKey {
600 SeqKey::try_new(s).unwrap()
601 }
602
603 #[tokio::test]
604 async fn advance_is_gapless_and_per_key() {
605 let dir = tempfile::tempdir().unwrap();
606 let d = FileDriver::open_or_init(dir.path()).unwrap();
607
608 assert_eq!(
610 d.advance_dense(&key("orders"), 5, Epoch(1)).await.unwrap(),
611 0
612 );
613 assert_eq!(
615 d.advance_dense(&key("orders"), 3, Epoch(1)).await.unwrap(),
616 5
617 );
618 assert_eq!(
620 d.advance_dense(&key("users"), 1, Epoch(1)).await.unwrap(),
621 0
622 );
623 assert_eq!(d.load_dense_seq(&key("orders")).await.unwrap(), 8);
625 assert_eq!(d.load_dense_seq(&key("users")).await.unwrap(), 1);
626 assert_eq!(d.load_dense_seq(&key("absent")).await.unwrap(), 0);
627 }
628
629 #[tokio::test]
630 async fn counters_survive_reopen() {
631 let dir = tempfile::tempdir().unwrap();
632 {
633 let d = FileDriver::open_or_init(dir.path()).unwrap();
634 d.advance_dense(&key("orders"), 10, Epoch(1)).await.unwrap();
635 }
636 let d2 = FileDriver::open_or_init(dir.path()).unwrap();
637 assert_eq!(
639 d2.advance_dense(&key("orders"), 1, Epoch(1)).await.unwrap(),
640 10
641 );
642 }
643
644 #[tokio::test]
645 async fn fresh_key_advance_succeeds() {
646 let dir = tempfile::tempdir().unwrap();
647 let d = FileDriver::open_or_init(dir.path()).unwrap();
648 assert!(d.advance_dense(&key("k"), 1, Epoch(1)).await.is_ok());
650 }
651
652 #[tokio::test]
653 async fn advance_past_u64_max_is_seq_overflow() {
654 use std::collections::BTreeMap;
655 let dir = tempfile::tempdir().unwrap();
656 let mut m = BTreeMap::new();
658 m.insert("k".to_string(), u64::MAX - 1);
659 let bytes = crate::dense_record::encode(&m, DEFAULT_DENSE_CARDINALITY_CAP);
660 std::fs::write(dir.path().join("dense"), bytes).unwrap();
661 let d = FileDriver::open_or_init(dir.path()).unwrap();
662 let err = d.advance_dense(&key("k"), 2, Epoch(1)).await;
664 assert!(matches!(
665 err,
666 Err(tsoracle_consensus::ConsensusError::SeqOverflow)
667 ));
668 assert_eq!(
670 d.advance_dense(&key("k"), 1, Epoch(1)).await.unwrap(),
671 u64::MAX - 1
672 );
673 }
674
675 #[tokio::test]
676 async fn cardinality_cap_rejects_new_keys_when_full() {
677 use std::collections::BTreeMap;
678 let dir = tempfile::tempdir().unwrap();
679 let mut m = BTreeMap::new();
680 m.insert("a".to_string(), 1u64);
681 m.insert("b".to_string(), 1u64);
682 let bytes = crate::dense_record::encode(&m, 2); std::fs::write(dir.path().join("dense"), bytes).unwrap();
684 let d = FileDriver::open_or_init(dir.path()).unwrap();
685 assert!(d.advance_dense(&key("a"), 1, Epoch(1)).await.is_ok());
687 let err = d.advance_dense(&key("c"), 1, Epoch(1)).await;
689 assert!(matches!(
690 err,
691 Err(tsoracle_consensus::ConsensusError::SeqKeyCardinalityExceeded { cap: 2 })
692 ));
693 }
694
695 #[tokio::test]
696 async fn batch_advance_is_gapless_atomic_and_one_durable_write() {
697 let dir = tempfile::tempdir().unwrap();
698 {
702 let d = FileDriver::open_or_init(dir.path()).unwrap();
703 let starts = d
704 .advance_dense_batch(&[(key("orders"), 5), (key("users"), 2)], Epoch(1))
705 .await
706 .unwrap();
707 assert_eq!(starts, vec![0, 0]);
708 }
709 let d2 = FileDriver::open_or_init(dir.path()).unwrap();
711 assert_eq!(d2.load_dense_seq(&key("orders")).await.unwrap(), 5);
712 assert_eq!(d2.load_dense_seq(&key("users")).await.unwrap(), 2);
713 }
714
715 #[tokio::test]
716 async fn batch_advance_duplicate_key_yields_adjacent_starts() {
717 let dir = tempfile::tempdir().unwrap();
723 let d = FileDriver::open_or_init(dir.path()).unwrap();
724 let starts = d
725 .advance_dense_batch(&[(key("k"), 3), (key("k"), 5)], Epoch(1))
726 .await
727 .unwrap();
728 assert_eq!(starts, vec![0, 3]); assert_eq!(d.load_dense_seq(&key("k")).await.unwrap(), 8);
730 }
731
732 #[tokio::test]
733 async fn batch_advance_empty_is_noop() {
734 let dir = tempfile::tempdir().unwrap();
736 let d = FileDriver::open_or_init(dir.path()).unwrap();
737 assert_eq!(d.advance_dense_batch(&[], Epoch(1)).await.unwrap(), vec![]);
738 }
739
740 #[tokio::test]
741 async fn batch_cardinality_is_atomic() {
742 use std::collections::BTreeMap;
743 let dir = tempfile::tempdir().unwrap();
744 let mut m = BTreeMap::new();
745 m.insert("a".to_string(), 1u64);
746 let bytes = crate::dense_record::encode(&m, 1); std::fs::write(dir.path().join("dense"), bytes).unwrap();
748 let d = FileDriver::open_or_init(dir.path()).unwrap();
749 let err = d
750 .advance_dense_batch(&[(key("a"), 1), (key("b"), 1)], Epoch(1))
751 .await;
752 assert!(matches!(
753 err,
754 Err(tsoracle_consensus::ConsensusError::SeqKeyCardinalityExceeded { cap: 1 })
755 ));
756 assert_eq!(d.load_dense_seq(&key("a")).await.unwrap(), 1);
758 }
759
760 #[tokio::test]
761 async fn batch_overflow_is_atomic_and_accumulates_duplicates() {
762 use std::collections::BTreeMap;
763 let dir = tempfile::tempdir().unwrap();
764 let mut m = BTreeMap::new();
765 m.insert("k".to_string(), u64::MAX - 5);
766 let bytes = crate::dense_record::encode(&m, DEFAULT_DENSE_CARDINALITY_CAP);
767 std::fs::write(dir.path().join("dense"), bytes).unwrap();
768 let d = FileDriver::open_or_init(dir.path()).unwrap();
769 let err = d
770 .advance_dense_batch(&[(key("k"), 4), (key("k"), 4)], Epoch(1))
771 .await;
772 assert!(matches!(
773 err,
774 Err(tsoracle_consensus::ConsensusError::SeqOverflow)
775 ));
776 assert_eq!(d.load_dense_seq(&key("k")).await.unwrap(), u64::MAX - 5);
777 }
778}
779
780#[cfg(test)]
781mod lease_tests {
782 use super::*;
783 use tempfile::tempdir;
784
785 fn rec(lease_id: u64) -> LeaseRecord {
786 LeaseRecord {
787 lease_id,
788 holder: format!("holder-{lease_id}").into_bytes(),
789 holder_epoch: lease_id + 10,
790 ttl_ms: 10_000,
791 ts_upper_bound: lease_id * 100,
792 expires_at_ms: lease_id * 100 + 10_000,
793 superseded: lease_id % 2 == 0,
794 }
795 }
796
797 #[tokio::test]
798 async fn fresh_dir_has_empty_lease_set() {
799 let dir = tempdir().unwrap();
800 let driver = FileDriver::open_or_init(dir.path()).unwrap();
801 assert_eq!(
802 driver.load_leases().await.unwrap(),
803 Vec::<LeaseRecord>::new()
804 );
805 }
806
807 #[tokio::test]
808 async fn persist_leases_then_load_leases_roundtrips() {
809 let dir = tempdir().unwrap();
810 let driver = FileDriver::open_or_init(dir.path()).unwrap();
811 let records = vec![rec(1), rec(2)];
812 driver.persist_leases(&records, Epoch(1)).await.unwrap();
813 assert_eq!(driver.load_leases().await.unwrap(), records);
814 }
815
816 #[tokio::test]
817 async fn leases_survive_reopen() {
818 let dir = tempdir().unwrap();
819 let records = vec![rec(1), rec(2)];
820 {
821 let driver = FileDriver::open_or_init(dir.path()).unwrap();
822 driver.persist_leases(&records, Epoch(1)).await.unwrap();
823 }
824 let reopened = FileDriver::open_or_init(dir.path()).unwrap();
825 assert_eq!(reopened.load_leases().await.unwrap(), records);
826 }
827}
828
829#[cfg(test)]
830mod tests {
831 use super::*;
832 use tempfile::tempdir;
833
834 #[tokio::test]
835 async fn fresh_init_starts_at_zero() {
836 let dir = tempdir().unwrap();
837 let driver = FileDriver::open_or_init(dir.path()).unwrap();
838 assert_eq!(driver.load_high_water().await.unwrap(), 0);
839 }
840
841 #[tokio::test]
842 async fn persist_then_reload() {
843 let dir = tempdir().unwrap();
844 let driver = FileDriver::open_or_init(dir.path()).unwrap();
845 let actual = driver.persist_high_water(12345, Epoch::ZERO).await.unwrap();
846 assert_eq!(actual, 12345);
847 drop(driver);
848 let reopened = FileDriver::open_or_init(dir.path()).unwrap();
849 assert_eq!(reopened.load_high_water().await.unwrap(), 12345);
850 }
851
852 #[tokio::test]
853 async fn persist_is_monotonic() {
854 let dir = tempdir().unwrap();
855 let driver = FileDriver::open_or_init(dir.path()).unwrap();
856 assert_eq!(
857 driver.persist_high_water(100, Epoch::ZERO).await.unwrap(),
858 100
859 );
860 assert_eq!(
861 driver.persist_high_water(50, Epoch::ZERO).await.unwrap(),
862 100
863 );
864 assert_eq!(
865 driver.persist_high_water(200, Epoch::ZERO).await.unwrap(),
866 200
867 );
868 }
869
870 #[tokio::test]
871 async fn init_seeded_rejects_existing_state() {
872 let dir = tempdir().unwrap();
873 FileDriver::init_seeded(dir.path(), 1_700_000_000_000).unwrap();
874 let err = FileDriver::init_seeded(dir.path(), 1_700_000_000_000).unwrap_err();
875 match err {
876 FileDriverError::Io(e) => assert_eq!(e.kind(), std::io::ErrorKind::AlreadyExists),
877 _ => panic!("expected AlreadyExists"),
878 }
879 }
880
881 #[tokio::test]
882 async fn init_seeded_reloads_as_physical_ms() {
883 let dir = tempdir().unwrap();
887 let seed = 1_700_000_000_000u64;
888 FileDriver::init_seeded(dir.path(), seed).unwrap();
889 let driver = FileDriver::open_or_init(dir.path()).unwrap();
890 assert_eq!(driver.load_high_water().await.unwrap(), seed);
891 assert!(seed < tsoracle_core::PHYSICAL_MS_MAX);
892 }
893
894 #[tokio::test]
895 async fn init_seeded_rejects_out_of_range_physical_ms() {
896 let dir = tempdir().unwrap();
897 let err = FileDriver::init_seeded(dir.path(), PHYSICAL_MS_MAX + 1).unwrap_err();
898 assert!(matches!(err, FileDriverError::PhysicalMsOutOfRange(_)));
899 }
900
901 #[tokio::test]
902 async fn persist_rejects_out_of_range_physical_ms() {
903 let dir = tempdir().unwrap();
904 let driver = FileDriver::open_or_init(dir.path()).unwrap();
905 let err = driver
906 .persist_high_water(PHYSICAL_MS_MAX + 1, Epoch::ZERO)
907 .await
908 .unwrap_err();
909 assert!(
910 matches!(err, ConsensusError::AdvanceOutOfRange(at_least) if at_least == PHYSICAL_MS_MAX + 1),
911 "out-of-range advance must surface as AdvanceOutOfRange carrying the offending value, got {err:?}"
912 );
913 }
914
915 #[tokio::test]
916 async fn open_or_init_rejects_out_of_range_state() {
917 let dir = tempdir().unwrap();
922 let state_path = dir.path().join("state");
923 let bytes = record::encode(PHYSICAL_MS_MAX + 1);
924 fs::write(&state_path, bytes).unwrap();
925 let err = FileDriver::open_or_init(dir.path()).unwrap_err();
926 assert!(
927 matches!(err, FileDriverError::PhysicalMsOutOfRange(v) if v == PHYSICAL_MS_MAX + 1)
928 );
929 }
930
931 #[tokio::test]
932 async fn leadership_events_emits_initial_leader_at_epoch_zero() {
933 let dir = tempdir().unwrap();
936 let driver = FileDriver::open_or_init(dir.path()).unwrap();
937 let mut stream = driver.leadership_events();
938 let first = tokio::time::timeout(std::time::Duration::from_secs(1), stream.next())
939 .await
940 .expect("stream emits initial state within the timeout")
941 .expect("stream is not closed");
942 assert_eq!(first, LeaderState::Leader { epoch: Epoch::ZERO });
943 }
944}