use core::pin::Pin;
use futures::{Stream, StreamExt};
use std::fs;
use std::io::Write;
#[cfg(unix)]
use std::os::fd::AsRawFd;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use tokio::sync::watch;
use tokio_stream::wrappers::WatchStream;
use tsoracle_consensus::{ConsensusDriver, ConsensusError, LeaderState};
use tsoracle_core::{Epoch, LeaseRecord, PHYSICAL_MS_MAX};
use crate::{dense_record, lease_record, record};
pub const DEFAULT_DENSE_CARDINALITY_CAP: u64 = 10_000;
#[derive(Debug)]
struct DenseState {
map: std::collections::BTreeMap<String, u64>,
cap: u64,
}
#[derive(Debug, thiserror::Error)]
pub enum FileDriverError {
#[error("io: {0}")]
Io(#[from] std::io::Error),
#[error("decode: {0}")]
Decode(#[from] record::RecordError),
#[error("physical_ms {0} exceeds 46-bit maximum")]
PhysicalMsOutOfRange(u64),
#[error("state directory {path} is already locked by another FileDriver: {source}")]
AlreadyLocked {
path: PathBuf,
#[source]
source: std::io::Error,
},
}
#[derive(Debug)]
pub struct FileDriver {
dir: PathBuf,
state: Arc<AtomicU64>,
write_lock: tokio::sync::Mutex<()>,
#[expect(
dead_code,
reason = "kept to hold the watch channel open for leader_rx consumers"
)]
leader_tx: watch::Sender<LeaderState>,
leader_rx: watch::Receiver<LeaderState>,
_lock: fs::File,
dense: tokio::sync::Mutex<DenseState>,
leases: tokio::sync::Mutex<Vec<LeaseRecord>>,
}
impl FileDriver {
pub fn open_or_init(dir: impl AsRef<Path>) -> Result<Arc<Self>, FileDriverError> {
let dir = dir.as_ref().to_path_buf();
fs::create_dir_all(&dir)?;
let lock_path = dir.join("LOCK");
let lock_file = fs::OpenOptions::new()
.create(true)
.read(true)
.write(true)
.truncate(false)
.open(&lock_path)?;
acquire_exclusive_lock(&lock_file, &lock_path)?;
let state_path = dir.join("state");
let current = if state_path.exists() {
let bytes = fs::read(&state_path)?;
let high_water = record::decode(&bytes)?;
if high_water > PHYSICAL_MS_MAX {
return Err(FileDriverError::PhysicalMsOutOfRange(high_water));
}
high_water
} else {
0
};
let dense_path = dir.join("dense");
let dense_state = if dense_path.exists() {
let bytes = fs::read(&dense_path)?;
let (map, cap) = dense_record::decode(&bytes)
.map_err(|e| FileDriverError::Io(std::io::Error::other(e)))?;
DenseState { map, cap }
} else {
DenseState {
map: std::collections::BTreeMap::new(),
cap: DEFAULT_DENSE_CARDINALITY_CAP,
}
};
let leases_path = dir.join("leases");
let leases = if leases_path.exists() {
let bytes = fs::read(&leases_path)?;
lease_record::decode(&bytes)
.map_err(|e| FileDriverError::Io(std::io::Error::other(e)))?
} else {
Vec::new()
};
let (tx, rx) = watch::channel(LeaderState::Leader { epoch: Epoch::ZERO });
Ok(Arc::new(FileDriver {
dir,
state: Arc::new(AtomicU64::new(current)),
write_lock: tokio::sync::Mutex::new(()),
leader_tx: tx,
leader_rx: rx,
_lock: lock_file,
dense: tokio::sync::Mutex::new(dense_state),
leases: tokio::sync::Mutex::new(leases),
}))
}
pub fn init_seeded(
dir: impl AsRef<Path>,
seed_physical_ms: u64,
) -> Result<(), FileDriverError> {
if seed_physical_ms > PHYSICAL_MS_MAX {
return Err(FileDriverError::PhysicalMsOutOfRange(seed_physical_ms));
}
let dir = dir.as_ref();
fs::create_dir_all(dir)?;
let state_path = dir.join("state");
if state_path.exists() {
return Err(FileDriverError::Io(std::io::Error::new(
std::io::ErrorKind::AlreadyExists,
"state file already exists; refusing to overwrite",
)));
}
write_record(dir, seed_physical_ms)?;
Ok(())
}
}
fn acquire_exclusive_lock(lock_file: &fs::File, lock_path: &Path) -> Result<(), FileDriverError> {
use fs2::FileExt;
match lock_file.try_lock_exclusive() {
Ok(()) => Ok(()),
Err(err) if err.raw_os_error() == fs2::lock_contended_error().raw_os_error() => {
Err(FileDriverError::AlreadyLocked {
path: lock_path.to_path_buf(),
source: err,
})
}
Err(err) => Err(FileDriverError::Io(err)),
}
}
fn write_record(dir: &Path, high_water: u64) -> Result<(), FileDriverError> {
tsoracle_failpoint::failpoint!(
"file_driver::before_write",
|arg: Option<String>| -> Result<(), FileDriverError> {
let _ = arg; Err(FileDriverError::Io(std::io::Error::other(
"failpoint: file_driver::before_write",
)))
}
);
let tmp = dir.join("state.tmp");
let final_path = dir.join("state");
let bytes = record::encode(high_water);
let mut file = fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.open(&tmp)?;
file.write_all(&bytes)?;
file.sync_all()?;
drop(file);
tsoracle_failpoint::failpoint!(
"file_driver::after_tmp_fsync_before_rename",
|arg: Option<String>| -> Result<(), FileDriverError> {
let _ = arg;
Err(FileDriverError::Io(std::io::Error::other(
"failpoint: file_driver::after_tmp_fsync_before_rename",
)))
}
);
fs::rename(&tmp, &final_path)?;
tsoracle_failpoint::failpoint!("file_driver::after_rename_before_dir_fsync");
#[cfg(unix)]
{
let dir_file = fs::File::open(dir)?;
let fd = dir_file.as_raw_fd();
let rc = unsafe { libc::fsync(fd) };
if rc != 0 {
return Err(FileDriverError::Io(std::io::Error::last_os_error()));
}
}
#[cfg(not(unix))]
{
let final_file = fs::OpenOptions::new().write(true).open(&final_path)?;
final_file.sync_all()?;
}
Ok(())
}
fn write_dense_record(
dir: &Path,
map: &std::collections::BTreeMap<String, u64>,
cap: u64,
) -> Result<(), FileDriverError> {
tsoracle_failpoint::failpoint!("file_driver::dense::before_write", |_arg: Option<
String,
>|
-> Result<
(),
FileDriverError,
> {
Err(FileDriverError::Io(std::io::Error::other(
"failpoint: file_driver::dense::before_write",
)))
});
let tmp = dir.join("dense.tmp");
let final_path = dir.join("dense");
let bytes = dense_record::encode(map, cap);
let mut file = fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.open(&tmp)?;
file.write_all(&bytes)?;
file.sync_all()?;
drop(file);
tsoracle_failpoint::failpoint!(
"file_driver::dense::after_tmp_fsync_before_rename",
|_arg: Option<String>| -> Result<(), FileDriverError> {
Err(FileDriverError::Io(std::io::Error::other(
"failpoint: file_driver::dense::after_tmp_fsync_before_rename",
)))
}
);
fs::rename(&tmp, &final_path)?;
tsoracle_failpoint::failpoint!("file_driver::dense::after_rename_before_dir_fsync");
#[cfg(unix)]
{
let dir_file = fs::File::open(dir)?;
let fd = dir_file.as_raw_fd();
let rc = unsafe { libc::fsync(fd) };
if rc != 0 {
return Err(FileDriverError::Io(std::io::Error::last_os_error()));
}
}
#[cfg(not(unix))]
{
let final_file = fs::OpenOptions::new().write(true).open(&final_path)?;
final_file.sync_all()?;
}
Ok(())
}
fn write_lease_record(dir: &Path, records: &[LeaseRecord]) -> Result<(), FileDriverError> {
tsoracle_failpoint::failpoint!("file_driver::leases::before_write", |_arg: Option<
String,
>|
-> Result<
(),
FileDriverError,
> {
Err(FileDriverError::Io(std::io::Error::other(
"failpoint: file_driver::leases::before_write",
)))
});
let tmp = dir.join("leases.tmp");
let final_path = dir.join("leases");
let bytes = lease_record::encode(records);
let mut file = fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.open(&tmp)?;
file.write_all(&bytes)?;
file.sync_all()?;
drop(file);
tsoracle_failpoint::failpoint!(
"file_driver::leases::after_tmp_fsync_before_rename",
|_arg: Option<String>| -> Result<(), FileDriverError> {
Err(FileDriverError::Io(std::io::Error::other(
"failpoint: file_driver::leases::after_tmp_fsync_before_rename",
)))
}
);
fs::rename(&tmp, &final_path)?;
tsoracle_failpoint::failpoint!("file_driver::leases::after_rename_before_dir_fsync");
#[cfg(unix)]
{
let dir_file = fs::File::open(dir)?;
let fd = dir_file.as_raw_fd();
let rc = unsafe { libc::fsync(fd) };
if rc != 0 {
return Err(FileDriverError::Io(std::io::Error::last_os_error()));
}
}
#[cfg(not(unix))]
{
let final_file = fs::OpenOptions::new().write(true).open(&final_path)?;
final_file.sync_all()?;
}
Ok(())
}
#[async_trait::async_trait]
impl ConsensusDriver for FileDriver {
fn leadership_events(&self) -> Pin<Box<dyn Stream<Item = LeaderState> + Send>> {
Box::pin(WatchStream::new(self.leader_rx.clone()).boxed())
}
async fn load_high_water(&self) -> Result<u64, ConsensusError> {
Ok(self.state.load(Ordering::Acquire))
}
async fn persist_high_water(
&self,
at_least: u64,
_epoch: Epoch,
) -> Result<u64, ConsensusError> {
tsoracle_consensus::reject_out_of_range_advance(at_least)?;
let _guard = self.write_lock.lock().await;
let current = self.state.load(Ordering::Acquire);
if at_least <= current {
return Ok(current);
}
let target = at_least;
let dir = self.dir.clone();
tokio::task::spawn_blocking(move || {
tsoracle_failpoint::failpoint!("file_driver::write_blocked");
write_record(&dir, target)
})
.await
.map_err(|e| ConsensusError::PermanentDriver(Box::new(std::io::Error::other(e))))?
.map_err(|e| ConsensusError::PermanentDriver(Box::new(e)))?;
self.state.store(target, Ordering::Release);
Ok(target)
}
async fn load_dense_seq(&self, key: &tsoracle_core::SeqKey) -> Result<u64, ConsensusError> {
let dense = self.dense.lock().await;
Ok(dense.map.get(key.as_str()).copied().unwrap_or(0))
}
async fn advance_dense(
&self,
key: &tsoracle_core::SeqKey,
count: u32,
_expected_epoch: Epoch,
) -> Result<u64, ConsensusError> {
let mut dense = self.dense.lock().await;
let present = dense.map.contains_key(key.as_str());
if !present && dense.map.len() as u64 >= dense.cap {
return Err(ConsensusError::SeqKeyCardinalityExceeded { cap: dense.cap });
}
let start = dense.map.get(key.as_str()).copied().unwrap_or(0);
let next = start
.checked_add(u64::from(count))
.ok_or(ConsensusError::SeqOverflow)?;
let mut new_map = dense.map.clone();
new_map.insert(key.as_str().to_string(), next);
let cap = dense.cap;
let dir = self.dir.clone();
let to_write = new_map.clone();
tokio::task::spawn_blocking(move || write_dense_record(&dir, &to_write, cap))
.await
.map_err(|e| ConsensusError::PermanentDriver(Box::new(std::io::Error::other(e))))?
.map_err(|e| ConsensusError::PermanentDriver(Box::new(e)))?;
dense.map = new_map;
Ok(start)
}
async fn advance_dense_batch(
&self,
entries: &[(tsoracle_core::SeqKey, u32)],
_expected_epoch: Epoch,
) -> Result<Vec<u64>, ConsensusError> {
if entries.is_empty() {
return Ok(Vec::new());
}
let mut dense = self.dense.lock().await;
let new_keys: std::collections::BTreeSet<&str> = entries
.iter()
.map(|(key, _)| key.as_str())
.filter(|k| !dense.map.contains_key(*k))
.collect();
if dense.map.len() as u64 + new_keys.len() as u64 > dense.cap {
return Err(ConsensusError::SeqKeyCardinalityExceeded { cap: dense.cap });
}
let mut scratch: std::collections::BTreeMap<String, u64> =
std::collections::BTreeMap::new();
let mut starts: Vec<u64> = Vec::with_capacity(entries.len());
for (key, count) in entries {
let key_str = key.as_str();
let running = scratch
.get(key_str)
.copied()
.or_else(|| dense.map.get(key_str).copied())
.unwrap_or(0);
starts.push(running);
let next = running
.checked_add(u64::from(*count))
.ok_or(ConsensusError::SeqOverflow)?;
scratch.insert(key_str.to_string(), next);
}
let mut new_map = dense.map.clone();
for (key_str, next) in &scratch {
new_map.insert(key_str.clone(), *next);
}
let cap = dense.cap;
let dir = self.dir.clone();
let to_write = new_map.clone();
tokio::task::spawn_blocking(move || write_dense_record(&dir, &to_write, cap))
.await
.map_err(|e| ConsensusError::PermanentDriver(Box::new(std::io::Error::other(e))))?
.map_err(|e| ConsensusError::PermanentDriver(Box::new(e)))?;
dense.map = new_map;
Ok(starts)
}
async fn load_leases(&self) -> Result<Vec<LeaseRecord>, ConsensusError> {
Ok(self.leases.lock().await.clone())
}
async fn persist_leases(
&self,
live: &[LeaseRecord],
_epoch: Epoch,
) -> Result<(), ConsensusError> {
let _guard = self.write_lock.lock().await;
let dir = self.dir.clone();
let to_write = live.to_vec();
tokio::task::spawn_blocking(move || write_lease_record(&dir, &to_write))
.await
.map_err(|e| ConsensusError::PermanentDriver(Box::new(std::io::Error::other(e))))?
.map_err(|e| ConsensusError::PermanentDriver(Box::new(e)))?;
*self.leases.lock().await = live.to_vec();
Ok(())
}
}
#[cfg(test)]
mod dense_tests {
use super::*;
use tsoracle_core::{Epoch, SeqKey};
fn key(s: &str) -> SeqKey {
SeqKey::try_new(s).unwrap()
}
#[tokio::test]
async fn advance_is_gapless_and_per_key() {
let dir = tempfile::tempdir().unwrap();
let d = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(
d.advance_dense(&key("orders"), 5, Epoch(1)).await.unwrap(),
0
);
assert_eq!(
d.advance_dense(&key("orders"), 3, Epoch(1)).await.unwrap(),
5
);
assert_eq!(
d.advance_dense(&key("users"), 1, Epoch(1)).await.unwrap(),
0
);
assert_eq!(d.load_dense_seq(&key("orders")).await.unwrap(), 8);
assert_eq!(d.load_dense_seq(&key("users")).await.unwrap(), 1);
assert_eq!(d.load_dense_seq(&key("absent")).await.unwrap(), 0);
}
#[tokio::test]
async fn counters_survive_reopen() {
let dir = tempfile::tempdir().unwrap();
{
let d = FileDriver::open_or_init(dir.path()).unwrap();
d.advance_dense(&key("orders"), 10, Epoch(1)).await.unwrap();
}
let d2 = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(
d2.advance_dense(&key("orders"), 1, Epoch(1)).await.unwrap(),
10
);
}
#[tokio::test]
async fn fresh_key_advance_succeeds() {
let dir = tempfile::tempdir().unwrap();
let d = FileDriver::open_or_init(dir.path()).unwrap();
assert!(d.advance_dense(&key("k"), 1, Epoch(1)).await.is_ok());
}
#[tokio::test]
async fn advance_past_u64_max_is_seq_overflow() {
use std::collections::BTreeMap;
let dir = tempfile::tempdir().unwrap();
let mut m = BTreeMap::new();
m.insert("k".to_string(), u64::MAX - 1);
let bytes = crate::dense_record::encode(&m, DEFAULT_DENSE_CARDINALITY_CAP);
std::fs::write(dir.path().join("dense"), bytes).unwrap();
let d = FileDriver::open_or_init(dir.path()).unwrap();
let err = d.advance_dense(&key("k"), 2, Epoch(1)).await;
assert!(matches!(
err,
Err(tsoracle_consensus::ConsensusError::SeqOverflow)
));
assert_eq!(
d.advance_dense(&key("k"), 1, Epoch(1)).await.unwrap(),
u64::MAX - 1
);
}
#[tokio::test]
async fn cardinality_cap_rejects_new_keys_when_full() {
use std::collections::BTreeMap;
let dir = tempfile::tempdir().unwrap();
let mut m = BTreeMap::new();
m.insert("a".to_string(), 1u64);
m.insert("b".to_string(), 1u64);
let bytes = crate::dense_record::encode(&m, 2); std::fs::write(dir.path().join("dense"), bytes).unwrap();
let d = FileDriver::open_or_init(dir.path()).unwrap();
assert!(d.advance_dense(&key("a"), 1, Epoch(1)).await.is_ok());
let err = d.advance_dense(&key("c"), 1, Epoch(1)).await;
assert!(matches!(
err,
Err(tsoracle_consensus::ConsensusError::SeqKeyCardinalityExceeded { cap: 2 })
));
}
#[tokio::test]
async fn batch_advance_is_gapless_atomic_and_one_durable_write() {
let dir = tempfile::tempdir().unwrap();
{
let d = FileDriver::open_or_init(dir.path()).unwrap();
let starts = d
.advance_dense_batch(&[(key("orders"), 5), (key("users"), 2)], Epoch(1))
.await
.unwrap();
assert_eq!(starts, vec![0, 0]);
}
let d2 = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(d2.load_dense_seq(&key("orders")).await.unwrap(), 5);
assert_eq!(d2.load_dense_seq(&key("users")).await.unwrap(), 2);
}
#[tokio::test]
async fn batch_advance_duplicate_key_yields_adjacent_starts() {
let dir = tempfile::tempdir().unwrap();
let d = FileDriver::open_or_init(dir.path()).unwrap();
let starts = d
.advance_dense_batch(&[(key("k"), 3), (key("k"), 5)], Epoch(1))
.await
.unwrap();
assert_eq!(starts, vec![0, 3]); assert_eq!(d.load_dense_seq(&key("k")).await.unwrap(), 8);
}
#[tokio::test]
async fn batch_advance_empty_is_noop() {
let dir = tempfile::tempdir().unwrap();
let d = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(d.advance_dense_batch(&[], Epoch(1)).await.unwrap(), vec![]);
}
#[tokio::test]
async fn batch_cardinality_is_atomic() {
use std::collections::BTreeMap;
let dir = tempfile::tempdir().unwrap();
let mut m = BTreeMap::new();
m.insert("a".to_string(), 1u64);
let bytes = crate::dense_record::encode(&m, 1); std::fs::write(dir.path().join("dense"), bytes).unwrap();
let d = FileDriver::open_or_init(dir.path()).unwrap();
let err = d
.advance_dense_batch(&[(key("a"), 1), (key("b"), 1)], Epoch(1))
.await;
assert!(matches!(
err,
Err(tsoracle_consensus::ConsensusError::SeqKeyCardinalityExceeded { cap: 1 })
));
assert_eq!(d.load_dense_seq(&key("a")).await.unwrap(), 1);
}
#[tokio::test]
async fn batch_overflow_is_atomic_and_accumulates_duplicates() {
use std::collections::BTreeMap;
let dir = tempfile::tempdir().unwrap();
let mut m = BTreeMap::new();
m.insert("k".to_string(), u64::MAX - 5);
let bytes = crate::dense_record::encode(&m, DEFAULT_DENSE_CARDINALITY_CAP);
std::fs::write(dir.path().join("dense"), bytes).unwrap();
let d = FileDriver::open_or_init(dir.path()).unwrap();
let err = d
.advance_dense_batch(&[(key("k"), 4), (key("k"), 4)], Epoch(1))
.await;
assert!(matches!(
err,
Err(tsoracle_consensus::ConsensusError::SeqOverflow)
));
assert_eq!(d.load_dense_seq(&key("k")).await.unwrap(), u64::MAX - 5);
}
}
#[cfg(test)]
mod lease_tests {
use super::*;
use tempfile::tempdir;
fn rec(lease_id: u64) -> LeaseRecord {
LeaseRecord {
lease_id,
holder: format!("holder-{lease_id}").into_bytes(),
holder_epoch: lease_id + 10,
ttl_ms: 10_000,
ts_upper_bound: lease_id * 100,
expires_at_ms: lease_id * 100 + 10_000,
superseded: lease_id % 2 == 0,
}
}
#[tokio::test]
async fn fresh_dir_has_empty_lease_set() {
let dir = tempdir().unwrap();
let driver = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(
driver.load_leases().await.unwrap(),
Vec::<LeaseRecord>::new()
);
}
#[tokio::test]
async fn persist_leases_then_load_leases_roundtrips() {
let dir = tempdir().unwrap();
let driver = FileDriver::open_or_init(dir.path()).unwrap();
let records = vec![rec(1), rec(2)];
driver.persist_leases(&records, Epoch(1)).await.unwrap();
assert_eq!(driver.load_leases().await.unwrap(), records);
}
#[tokio::test]
async fn leases_survive_reopen() {
let dir = tempdir().unwrap();
let records = vec![rec(1), rec(2)];
{
let driver = FileDriver::open_or_init(dir.path()).unwrap();
driver.persist_leases(&records, Epoch(1)).await.unwrap();
}
let reopened = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(reopened.load_leases().await.unwrap(), records);
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::tempdir;
#[tokio::test]
async fn fresh_init_starts_at_zero() {
let dir = tempdir().unwrap();
let driver = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(driver.load_high_water().await.unwrap(), 0);
}
#[tokio::test]
async fn persist_then_reload() {
let dir = tempdir().unwrap();
let driver = FileDriver::open_or_init(dir.path()).unwrap();
let actual = driver.persist_high_water(12345, Epoch::ZERO).await.unwrap();
assert_eq!(actual, 12345);
drop(driver);
let reopened = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(reopened.load_high_water().await.unwrap(), 12345);
}
#[tokio::test]
async fn persist_is_monotonic() {
let dir = tempdir().unwrap();
let driver = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(
driver.persist_high_water(100, Epoch::ZERO).await.unwrap(),
100
);
assert_eq!(
driver.persist_high_water(50, Epoch::ZERO).await.unwrap(),
100
);
assert_eq!(
driver.persist_high_water(200, Epoch::ZERO).await.unwrap(),
200
);
}
#[tokio::test]
async fn init_seeded_rejects_existing_state() {
let dir = tempdir().unwrap();
FileDriver::init_seeded(dir.path(), 1_700_000_000_000).unwrap();
let err = FileDriver::init_seeded(dir.path(), 1_700_000_000_000).unwrap_err();
match err {
FileDriverError::Io(e) => assert_eq!(e.kind(), std::io::ErrorKind::AlreadyExists),
_ => panic!("expected AlreadyExists"),
}
}
#[tokio::test]
async fn init_seeded_reloads_as_physical_ms() {
let dir = tempdir().unwrap();
let seed = 1_700_000_000_000u64;
FileDriver::init_seeded(dir.path(), seed).unwrap();
let driver = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(driver.load_high_water().await.unwrap(), seed);
assert!(seed < tsoracle_core::PHYSICAL_MS_MAX);
}
#[tokio::test]
async fn init_seeded_rejects_out_of_range_physical_ms() {
let dir = tempdir().unwrap();
let err = FileDriver::init_seeded(dir.path(), PHYSICAL_MS_MAX + 1).unwrap_err();
assert!(matches!(err, FileDriverError::PhysicalMsOutOfRange(_)));
}
#[tokio::test]
async fn persist_rejects_out_of_range_physical_ms() {
let dir = tempdir().unwrap();
let driver = FileDriver::open_or_init(dir.path()).unwrap();
let err = driver
.persist_high_water(PHYSICAL_MS_MAX + 1, Epoch::ZERO)
.await
.unwrap_err();
assert!(
matches!(err, ConsensusError::AdvanceOutOfRange(at_least) if at_least == PHYSICAL_MS_MAX + 1),
"out-of-range advance must surface as AdvanceOutOfRange carrying the offending value, got {err:?}"
);
}
#[tokio::test]
async fn open_or_init_rejects_out_of_range_state() {
let dir = tempdir().unwrap();
let state_path = dir.path().join("state");
let bytes = record::encode(PHYSICAL_MS_MAX + 1);
fs::write(&state_path, bytes).unwrap();
let err = FileDriver::open_or_init(dir.path()).unwrap_err();
assert!(
matches!(err, FileDriverError::PhysicalMsOutOfRange(v) if v == PHYSICAL_MS_MAX + 1)
);
}
#[tokio::test]
async fn leadership_events_emits_initial_leader_at_epoch_zero() {
let dir = tempdir().unwrap();
let driver = FileDriver::open_or_init(dir.path()).unwrap();
let mut stream = driver.leadership_events();
let first = tokio::time::timeout(std::time::Duration::from_secs(1), stream.next())
.await
.expect("stream emits initial state within the timeout")
.expect("stream is not closed");
assert_eq!(first, LeaderState::Leader { epoch: Epoch::ZERO });
}
}