Skip to main content

tsoracle_driver_file/
driver.rs

1//
2//  ░▀█▀░█▀▀░█▀█░█▀▄░█▀█░█▀▀░█░░░█▀▀
3//  ░░█░░▀▀█░█░█░█▀▄░█▀█░█░░░█░░░█▀▀
4//  ░░▀░░▀▀▀░▀▀▀░▀░▀░▀░▀░▀▀▀░▀▀▀░▀▀▀
5//
6//  tsoracle — Distributed Timestamp Oracle
7//  https://www.tsoracle.rs
8//
9//  Copyright (c) 2026 Prisma Risk
10//
11//  Licensed under the Apache License, Version 2.0 (the "License");
12//  you may not use this file except in compliance with the License.
13//  You may obtain a copy of the License at
14//
15//      https://www.apache.org/licenses/LICENSE-2.0
16//
17//  Unless required by applicable law or agreed to in writing, software
18//  distributed under the License is distributed on an "AS IS" BASIS,
19//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
20//  See the License for the specific language governing permissions and
21//  limitations under the License.
22//
23
24// #[PerformanceCriticalPath]
25
26use 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
42/// Default genesis cardinality cap for a freshly-initialized dense record.
43/// Immutable once written (Plan 1 has no reconfiguration path).
44pub 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    // Published high-water for readers. Writers are externally serialized by
72    // `write_lock`, so this is a publish-to-readers cell, not a mutual-exclusion
73    // lock. Reads (`load_high_water`) are wait-free; writers do disk I/O
74    // without holding any state lock and then publish via a Release store.
75    state: Arc<AtomicU64>,
76    write_lock: tokio::sync::Mutex<()>,
77    // Held to keep the watch channel open; FileDriver never sends after the
78    // initial Leader { epoch: 0 } published at construction. Dropping it would
79    // close the channel and terminate every `WatchStream::new(leader_rx.clone())`
80    // consumer prematurely.
81    #[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    // Holds the OS-level exclusive lock on `dir/LOCK` for the driver's
88    // lifetime. The kernel releases the flock when this file is closed —
89    // on graceful Drop, on panic unwind, and on hard process death — so
90    // there is no stale-lock cleanup path to maintain.
91    _lock: fs::File,
92    /// In-memory mirror of the on-disk dense map. The `Mutex` serializes
93    /// the read-modify-write fetch-add so concurrent callers don't race.
94    dense: tokio::sync::Mutex<DenseState>,
95    /// In-memory mirror of the on-disk lease set.
96    leases: tokio::sync::Mutex<Vec<LeaseRecord>>,
97}
98
99impl FileDriver {
100    /// Open the state directory. Creates it if missing. Reads and validates the
101    /// state file if present. Single-node deployments serve `Leader { epoch: 0 }`
102    /// continuously.
103    ///
104    /// Acquires an exclusive OS-level lock on the `LOCK` sentinel file under
105    /// `dir` before reading state, and holds it for the lifetime of the
106    /// returned driver. A second concurrent `open_or_init` against the same
107    /// directory returns [`FileDriverError::AlreadyLocked`] immediately —
108    /// `FileDriver` enforces its one-writer-per-directory precondition rather
109    /// than trusting the operator to. The lock is released by the kernel when
110    /// the driver is dropped or the process exits (including crash).
111    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        // Lock BEFORE reading state so the in-memory snapshot can't race a
116        // concurrent writer in another process. The sentinel is a stable
117        // inode — `write_record` replaces `state` via atomic rename, so a
118        // lock held on `state` itself would not cover the post-rename file.
119        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    /// Seed a fresh state directory with a high-water value. Used by the `init`
174    /// CLI subcommand for migrations. Fails if state already exists.
175    ///
176    /// The stored high-water is a physical_ms (the same units the allocator
177    /// uses for `committed_high_water`), NOT a packed `Timestamp`. The seed
178    /// argument is interpreted as the maximum physical_ms ever observed in the
179    /// prior system; on first serve, the failover fence will advance above it.
180    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
201/// Try-acquire an exclusive flock on `lock_file`. Classify the contended
202/// case (another live `FileDriver` holds it) as
203/// [`FileDriverError::AlreadyLocked`]; any other I/O error becomes
204/// [`FileDriverError::Io`].
205///
206/// We don't trust `io::Error::kind()` alone here: on Unix the contended
207/// errno is `EWOULDBLOCK` (mapped to `ErrorKind::WouldBlock`), but on
208/// Windows `LockFileEx` returns `ERROR_LOCK_VIOLATION`, which stdlib does
209/// not necessarily map to `WouldBlock`. `fs2::lock_contended_error()`
210/// returns the exact `io::Error` shape the platform uses, so we match on
211/// `raw_os_error()` for a portable check.
212fn 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; // currently only one action shape; future tags can match here
231            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    // Force the rename's metadata to durable media. The tmpfile `sync_all`
265    // above keeps the *data* durable on both platforms; this block adds the
266    // *metadata* barrier that makes the new directory entry survive a crash.
267    //
268    // Unix: open the parent directory and `fsync` its descriptor. This is
269    // the canonical POSIX barrier for a rename — it flushes the directory
270    // entry that names the new inode.
271    //
272    // Windows: there is no portable directory-level flush. `FlushFileBuffers`
273    // on a directory handle is undefined for most filesystems. NTFS journals
274    // `MoveFileEx` as a metadata transaction, but the `$LogFile` record is
275    // itself only durable after a checkpoint or an explicit
276    // `FlushFileBuffers` on a file on the same volume. Re-opening the
277    // renamed file with write access (required by `FlushFileBuffers`) and
278    // calling `sync_all` flushes the journal entry covering this rename.
279    // This is the pattern SQLite and RocksDB use on Windows.
280    #[cfg(unix)]
281    {
282        let dir_file = fs::File::open(dir)?;
283        let fd = dir_file.as_raw_fd();
284        // SAFETY: fd is a valid open directory descriptor for the duration of this call.
285        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        // `write(true)` is required: `FlushFileBuffers` rejects handles
293        // without `GENERIC_WRITE`. Default `truncate: false` leaves the
294        // file contents (the record we just renamed into place) intact.
295        let final_file = fs::OpenOptions::new().write(true).open(&final_path)?;
296        final_file.sync_all()?;
297    }
298    Ok(())
299}
300
301/// Atomically persist the dense map. Mirrors `write_record`'s tmp+fsync+rename+dir-fsync
302/// protocol (and its failpoint structure) for the dense `dense` file.
303fn 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        // SAFETY: fd is a valid open directory descriptor for the duration of this call.
351        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
364/// Atomically persist the lease set. Mirrors `write_record`'s
365/// tmp+fsync+rename+dir-fsync protocol for the `leases` file.
366fn 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        // SAFETY: fd is a valid open directory descriptor for the duration of this call.
410        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        // Wait-free read; pairs with the Release store in `persist_high_water`.
431        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        // Shared with the consensus backends so every driver rejects an
440        // out-of-range advance at the same bound before persisting it.
441        tsoracle_consensus::reject_out_of_range_advance(at_least)?;
442
443        // `write_lock` serializes writers — no two `persist_high_water` calls
444        // can race the disk write or the publish step below.
445        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        // spawn_blocking JoinError: the worker thread panicked. That is a
460        // bug, not a transient condition — fail permanently.
461        .map_err(|e| ConsensusError::PermanentDriver(Box::new(std::io::Error::other(e))))?
462        // FileDriverError covers the disk path: I/O failure, CRC/length
463        // checks, fsync failure. None of these are safely retried at this
464        // layer without operator visibility (a stuck disk does not clear
465        // itself). Classify as permanent.
466        .map_err(|e| ConsensusError::PermanentDriver(Box::new(e)))?;
467
468        // Publish only after the disk write is durable. Release pairs with
469        // the Acquire load in `load_high_water` and the snapshot above.
470        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        // Build the would-be-new map, persist it durably, THEN publish in memory.
497        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        // An empty batch is a no-op: return before locking or touching disk so a
517        // degenerate call never triggers a gratuitous durable rewrite.
518        if entries.is_empty() {
519            return Ok(Vec::new());
520        }
521
522        let mut dense = self.dense.lock().await;
523
524        // Phase 1: cardinality. Count distinct keys not already present;
525        // duplicates within the batch count only once.
526        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        // Phase 2: sequential fold into a scratch map seeded from current
536        // counters; each entry's start is the running value before the advance.
537        // Overflow against the accumulated value rejects the whole batch.
538        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        // Phase 3: build the would-be-new map, persist it durably ONCE, THEN
556        // publish. A crash before the publish leaves the prior map intact.
557        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        // First block for "orders": [0, 5).
609        assert_eq!(
610            d.advance_dense(&key("orders"), 5, Epoch(1)).await.unwrap(),
611            0
612        );
613        // Next block for "orders": [5, 8).
614        assert_eq!(
615            d.advance_dense(&key("orders"), 3, Epoch(1)).await.unwrap(),
616            5
617        );
618        // "users" is independent, starts at 0.
619        assert_eq!(
620            d.advance_dense(&key("users"), 1, Epoch(1)).await.unwrap(),
621            0
622        );
623        // load reflects committed counters.
624        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        // Next start is exactly the persisted counter — no gap, no rewind.
638        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        // Sanity: a fresh key with a small count succeeds.
649        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        // Seed a dense record with "k" already near the ceiling.
657        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        // count=2 would need u64::MAX-1 + 2 = overflow.
663        let err = d.advance_dense(&key("k"), 2, Epoch(1)).await;
664        assert!(matches!(
665            err,
666            Err(tsoracle_consensus::ConsensusError::SeqOverflow)
667        ));
668        // count=1 is exactly representable (lands at u64::MAX), so it succeeds.
669        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); // cap = 2, already full
683        std::fs::write(dir.path().join("dense"), bytes).unwrap();
684        let d = FileDriver::open_or_init(dir.path()).unwrap();
685        // Existing keys still advance.
686        assert!(d.advance_dense(&key("a"), 1, Epoch(1)).await.is_ok());
687        // A NEW key is rejected at the cap.
688        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        // FileDriver holds an EXCLUSIVE directory lock for its lifetime, so the
699        // first driver must be dropped (scope) before reopening — mirroring the
700        // existing `counters_survive_reopen` test.
701        {
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        // Reopen: the whole batch was one durable write.
710        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        // A duplicate key within one batch (rejected by the server pre-commit,
718        // but handled deterministically here so the driver is never the weak
719        // link) must produce ADJACENT, non-overlapping blocks: the second
720        // entry's start is the first entry's post-advance value. This pins the
721        // accumulation contract directly, since the file driver returns starts.
722        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]); // [0,3) then [3,8)
729        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        // An empty batch returns an empty start list without touching disk.
735        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); // cap 1, full
747        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        // "a" unchanged — nothing committed.
757        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        // The seed argument is a physical_ms; on reload the driver reports the
884        // same value (NOT shifted into a packed Timestamp). The allocator's
885        // bounds and the file driver's stored value must use identical units.
886        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        // Hand-write a state file whose encoded high_water exceeds the
918        // 46-bit physical_ms cap. open_or_init must refuse to load it
919        // rather than silently propagating an invariant violation into
920        // the allocator.
921        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        // FileDriver is single-node by design: every observer sees a single,
934        // permanent `Leader { epoch: 0 }` transition on subscription.
935        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}