Skip to main content

core_storage/
fs.rs

1use std::fs::{File, OpenOptions};
2use std::io::{Read, Write};
3use std::path::PathBuf;
4
5#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
6pub enum FileId {
7    Wal,
8    Snapshot,
9    /// Backup of the previous snapshot, kept until the next clean open at
10    /// the current format version. Written before any migration to preserve
11    /// the original bytes if the migration step fails.
12    SnapshotBak,
13    /// RBAC role definitions sidecar. Written atomically by `apply_schema`
14    /// when roles change; loaded at open. Never part of WAL/snapshot format.
15    Roles,
16    /// Commit → wall-clock sidecar. Written atomically as commits land; loaded
17    /// at open. **Never part of WAL/snapshot format**, exactly as `Roles` is
18    /// not — a release that does not know this file does not read it, and opens
19    /// the store as it always did. See [`crate::commit_times`] for why the map
20    /// is not a WAL record.
21    CommitTimes,
22}
23
24impl FileId {
25    fn name(self) -> &'static str {
26        match self {
27            FileId::Wal => "wal.bin",
28            FileId::Snapshot => "snapshot.bin",
29            FileId::SnapshotBak => "snapshot.bin.bak",
30            FileId::Roles => "roles.json",
31            FileId::CommitTimes => "commit_times.bin",
32        }
33    }
34}
35
36pub trait Fs {
37    fn append(&mut self, file: FileId, data: &[u8]) -> std::io::Result<()>;
38    fn sync(&mut self, file: FileId) -> std::io::Result<()>;
39    fn read(&self, file: FileId) -> std::io::Result<Vec<u8>>;
40    fn write_atomic(&mut self, file: FileId, data: &[u8]) -> std::io::Result<()>;
41    /// Return the on-disk path of the snapshot file, if any.
42    ///
43    /// `Some` for `RealFs` (used by `MappedBase::map` for true file mmap).
44    /// `None` for `SimFs` and other in-memory implementations (falls back to
45    /// `MappedBase::from_bytes`).
46    fn snapshot_path(&self) -> Option<std::path::PathBuf> {
47        None
48    }
49
50    /// Return the on-disk path of the WAL file, if any.
51    ///
52    /// `Some` for `RealFs`. `None` for in-memory implementations.
53    /// Used by [`GraphDb::wal_size_bytes`] to read WAL file metadata.
54    fn wal_path(&self) -> Option<std::path::PathBuf> {
55        None
56    }
57    /// Read at most `n` bytes from the beginning of `file` without loading
58    /// the full contents.
59    ///
60    /// Used by the open path to sniff the 6-byte magic+version header before
61    /// deciding whether to mmap (V8) or full-read (legacy V5-V7).
62    ///
63    /// The default implementation calls `read()` and truncates; override in
64    /// `RealFs` for a true partial read.
65    fn read_prefix(&self, file: FileId, n: usize) -> std::io::Result<Vec<u8>> {
66        let mut bytes = self.read(file)?;
67        bytes.truncate(n);
68        Ok(bytes)
69    }
70
71    // ── Cross-process write lock ──────────────────────────────────────────────
72
73    /// Try to take the store's advisory exclusive write lock without blocking.
74    ///
75    /// Returns `Ok(true)` when the lock is now held by this handle, `Ok(false)`
76    /// when another handle holds it.  Taking a lock this handle already holds
77    /// is a successful no-op, so callers may re-acquire freely.
78    ///
79    /// The lock is advisory: it coordinates cooperating mushroomdb processes
80    /// and does not stop an unrelated program from writing the files.
81    ///
82    /// Takes `&self`, not `&mut self`, on purpose: a writer must be able to
83    /// poll for the lock without holding the in-process write guard, or a busy
84    /// peer in another process would stall every reader in this one.
85    ///
86    /// Default: `Ok(true)` — an in-memory store has no other process to
87    /// coordinate with.
88    fn try_lock_exclusive(&self) -> std::io::Result<bool> {
89        Ok(true)
90    }
91
92    /// Release the advisory write lock.  No-op when it is not held.
93    ///
94    /// Default: no-op.
95    fn unlock(&self) -> std::io::Result<()> {
96        Ok(())
97    }
98
99    // ── WAL tailing (refresh) ─────────────────────────────────────────────────
100
101    /// Current length of the WAL in bytes.
102    ///
103    /// Must not read file contents: this is the hot half of a staleness check
104    /// and runs on the read path.  The default implementation reads the WAL
105    /// because a generic `Fs` has no cheaper option; `RealFs` overrides it with
106    /// a metadata-only stat.
107    fn wal_len(&self) -> std::io::Result<u64> {
108        Ok(self.read(FileId::Wal)?.len() as u64)
109    }
110
111    /// Read `file` from byte offset `from` to the end.
112    ///
113    /// Returns an empty `Vec` when `from` is at or past the end of the file.
114    /// Used to decode the WAL tail another process appended since this handle
115    /// last consumed it.
116    ///
117    /// The default implementation reads the whole file and slices; `RealFs`
118    /// overrides it with a seek + read of the tail only.
119    fn read_range(&self, file: FileId, from: u64) -> std::io::Result<Vec<u8>> {
120        let bytes = self.read(file)?;
121        let from = from.min(bytes.len() as u64) as usize;
122        Ok(bytes[from..].to_vec())
123    }
124
125    /// Identity of the current snapshot file as `(len, mtime_nanos)`, or `None`
126    /// when no snapshot exists.
127    ///
128    /// A change in this pair means some process replaced the snapshot, so the
129    /// WAL this handle was tailing no longer continues its in-memory state and
130    /// a full reload is required.  Like [`wal_len`](Fs::wal_len) this must not
131    /// read file contents.
132    ///
133    /// Default: length-only identity derived from the snapshot bytes (in-memory
134    /// stores have no mtime); `RealFs` overrides it with a metadata-only stat.
135    fn snapshot_ident(&self) -> std::io::Result<Option<(u64, u64)>> {
136        let len = self.read(FileId::Snapshot)?.len() as u64;
137        Ok(if len == 0 { None } else { Some((len, 0)) })
138    }
139
140    // ── WAL archive methods (Task 4: history-preserving snapshots) ─────────────
141
142    /// List all WAL archive identifiers, sorted ascending (oldest first).
143    ///
144    /// Each archive created by [`archive_wal`] with commit-seq `N` appears as `N`.
145    /// Returns an empty list when no archives exist.
146    ///
147    /// Default: no archive support — returns empty.
148    fn list_archives(&self) -> std::io::Result<Vec<u64>> {
149        Ok(vec![])
150    }
151
152    /// Read the byte contents of WAL archive `n`.
153    ///
154    /// Returns an empty `Vec` when the archive does not exist.
155    ///
156    /// Default: no archive support — returns empty.
157    fn read_archive(&self, _n: u64) -> std::io::Result<Vec<u8>> {
158        Ok(vec![])
159    }
160
161    /// Atomically rename the current WAL to `wal.<n>.archive` (same directory,
162    /// same filesystem — the rename is guaranteed atomic at the OS level).
163    ///
164    /// The caller ensures the snapshot has been durably written before calling
165    /// this method.  After a successful rename the old WAL no longer exists as
166    /// `wal.bin`; a subsequent [`write_atomic`] on `FileId::Wal` creates a new
167    /// empty WAL.
168    ///
169    /// Returns `Err` if the operation is not supported or fails.
170    fn archive_wal(&mut self, _n: u64) -> std::io::Result<()> {
171        Err(std::io::Error::other(
172            "archive_wal not supported by this Fs implementation",
173        ))
174    }
175
176    /// Delete archive `n`.  No-op if it does not exist.
177    ///
178    /// Retention pruning (inside `snapshot_with`) is the only call site.
179    ///
180    /// Default: no-op.
181    fn delete_archive(&mut self, _n: u64) -> std::io::Result<()> {
182        Ok(())
183    }
184
185    /// Return the persisted horizon floor — the global frame index of the
186    /// first commit that is still reachable through surviving archives.
187    ///
188    /// Defaults to `0` (all history reachable / no pruning ever performed).
189    fn read_horizon_floor(&self) -> std::io::Result<u64> {
190        Ok(0)
191    }
192
193    /// Atomically persist `floor` so that a subsequent [`read_horizon_floor`]
194    /// after reopen returns the same value.
195    ///
196    /// Default: no-op (in-memory only; override in durable implementations).
197    fn write_horizon_floor(&mut self, _floor: u64) -> std::io::Result<()> {
198        Ok(())
199    }
200
201    /// Return `true` when the `wal.genesis` marker file is present.
202    ///
203    /// The marker signals that the surviving archive chain forms a complete,
204    /// uninterrupted WAL history starting from the store's first-ever commit
205    /// (the genesis chain).  When absent, archive-resident commits are not
206    /// safe to replay from empty state and `open_at` must refuse them.
207    ///
208    /// Default: `false` (no genesis chain / no archive support).
209    fn has_genesis_marker(&self) -> bool {
210        false
211    }
212
213    /// Durably create the `wal.genesis` marker file.
214    ///
215    /// Written exactly once, when the first WAL archive is taken from a store
216    /// that has never undergone a WAL-truncating snapshot.
217    ///
218    /// Default: no-op.
219    fn write_genesis_marker(&mut self) -> std::io::Result<()> {
220        Ok(())
221    }
222
223    /// Remove the `wal.genesis` marker file.
224    ///
225    /// Called when the genesis chain is broken: either by archive pruning
226    /// (floor advances past 0) or by a WAL-truncating snapshot taken after
227    /// archives already exist.  No-op if the marker is absent.
228    ///
229    /// Default: no-op.
230    fn delete_genesis_marker(&mut self) -> std::io::Result<()> {
231        Ok(())
232    }
233}
234
235pub trait FsIntrospect {
236    fn total_appended(&self) -> usize;
237    fn sync_count(&self) -> usize {
238        0
239    }
240}
241
242/// Name of the advisory cross-process write lock file.
243///
244/// Always empty: the file exists only to carry the OS lock. It is created on
245/// the first lock attempt and never removed, so the same inode backs the lock
246/// for every process that opens the store.
247pub const LOCK_FILE: &str = "LOCK";
248
249/// This handle's one open description of the store's `LOCK` file, plus whether
250/// it currently owns the lock.
251///
252/// Exactly one per `RealFs` on purpose: the OS lock is held per open file
253/// description, so two descriptions of `LOCK` inside one process would contend
254/// with each other. Opened lazily — a `RealFs` created for a one-shot file
255/// operation never touches the lock file at all.
256#[derive(Debug, Default)]
257struct LockState {
258    file: Option<File>,
259    held: bool,
260}
261
262#[derive(Debug)]
263pub struct RealFs {
264    dir: PathBuf,
265    /// Behind a `Mutex` so the lock can be taken and released through `&self`.
266    /// A writer polls for the cross-process lock *before* it takes the
267    /// in-process write guard, so that a peer holding the lock cannot stall
268    /// this process's readers.
269    lock: std::sync::Mutex<LockState>,
270}
271
272impl RealFs {
273    pub fn new(dir: &std::path::Path) -> std::io::Result<Self> {
274        std::fs::create_dir_all(dir)?;
275        Ok(Self {
276            dir: dir.to_path_buf(),
277            lock: std::sync::Mutex::new(LockState::default()),
278        })
279    }
280
281    /// The database directory this filesystem is rooted at.
282    pub fn dir(&self) -> &std::path::Path {
283        &self.dir
284    }
285
286    fn path(&self, file: FileId) -> PathBuf {
287        self.dir.join(file.name())
288    }
289}
290
291impl Fs for RealFs {
292    fn append(&mut self, file: FileId, data: &[u8]) -> std::io::Result<()> {
293        let mut f = OpenOptions::new()
294            .create(true)
295            .append(true)
296            .open(self.path(file))?;
297        f.write_all(data)
298    }
299
300    fn sync(&mut self, file: FileId) -> std::io::Result<()> {
301        let f = File::open(self.path(file))?;
302        full_sync(&f)
303    }
304
305    fn read(&self, file: FileId) -> std::io::Result<Vec<u8>> {
306        match File::open(self.path(file)) {
307            Ok(mut f) => {
308                let mut buf = Vec::new();
309                f.read_to_end(&mut buf)?;
310                Ok(buf)
311            }
312            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Vec::new()),
313            Err(e) => Err(e),
314        }
315    }
316
317    fn write_atomic(&mut self, file: FileId, data: &[u8]) -> std::io::Result<()> {
318        let tmp = self.dir.join(format!("{}.tmp", file.name()));
319        {
320            let mut f = File::create(&tmp)?;
321            f.write_all(data)?;
322            full_sync(&f)?;
323        }
324        std::fs::rename(&tmp, self.path(file))?;
325        sync_dir(&self.dir)
326    }
327
328    fn snapshot_path(&self) -> Option<std::path::PathBuf> {
329        Some(self.path(FileId::Snapshot))
330    }
331
332    fn wal_path(&self) -> Option<std::path::PathBuf> {
333        Some(self.path(FileId::Wal))
334    }
335
336    fn read_prefix(&self, file: FileId, n: usize) -> std::io::Result<Vec<u8>> {
337        use std::io::Read as _;
338        match File::open(self.path(file)) {
339            Ok(mut f) => {
340                let mut buf = vec![0u8; n];
341                let read = f.read(&mut buf)?;
342                buf.truncate(read);
343                Ok(buf)
344            }
345            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Vec::new()),
346            Err(e) => Err(e),
347        }
348    }
349
350    fn try_lock_exclusive(&self) -> std::io::Result<bool> {
351        let mut state = self.lock.lock().unwrap_or_else(|e| e.into_inner());
352        if state.held {
353            return Ok(true);
354        }
355        if state.file.is_none() {
356            state.file = Some(
357                OpenOptions::new()
358                    .create(true)
359                    .read(true)
360                    .write(true)
361                    .truncate(false)
362                    .open(self.dir.join(LOCK_FILE))?,
363            );
364        }
365        let f = state.file.as_ref().expect("lock file just opened");
366        match f.try_lock() {
367            Ok(()) => {
368                state.held = true;
369                Ok(true)
370            }
371            Err(std::fs::TryLockError::WouldBlock) => Ok(false),
372            Err(std::fs::TryLockError::Error(e)) => Err(e),
373        }
374    }
375
376    fn unlock(&self) -> std::io::Result<()> {
377        let mut state = self.lock.lock().unwrap_or_else(|e| e.into_inner());
378        if !state.held {
379            return Ok(());
380        }
381        // Clear the flag first: a failed unlock must not leave the handle
382        // believing it still owns a lock it may have lost.
383        state.held = false;
384        match state.file.as_ref() {
385            Some(f) => f.unlock(),
386            None => Ok(()),
387        }
388    }
389
390    fn wal_len(&self) -> std::io::Result<u64> {
391        match std::fs::metadata(self.path(FileId::Wal)) {
392            Ok(m) => Ok(m.len()),
393            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(0),
394            Err(e) => Err(e),
395        }
396    }
397
398    fn read_range(&self, file: FileId, from: u64) -> std::io::Result<Vec<u8>> {
399        use std::io::{Read as _, Seek as _, SeekFrom};
400        let mut f = match File::open(self.path(file)) {
401            Ok(f) => f,
402            Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
403            Err(e) => return Err(e),
404        };
405        let len = f.metadata()?.len();
406        if from >= len {
407            return Ok(Vec::new());
408        }
409        f.seek(SeekFrom::Start(from))?;
410        let mut buf = Vec::with_capacity((len - from) as usize);
411        f.read_to_end(&mut buf)?;
412        Ok(buf)
413    }
414
415    fn snapshot_ident(&self) -> std::io::Result<Option<(u64, u64)>> {
416        let m = match std::fs::metadata(self.path(FileId::Snapshot)) {
417            Ok(m) => m,
418            Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
419            Err(e) => return Err(e),
420        };
421        // mtime is a change hint, not a clock: an unreadable or pre-epoch
422        // timestamp degrades to 0, leaving length alone to detect the change.
423        let mtime_nanos = m
424            .modified()
425            .ok()
426            .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
427            .map(|d| d.as_nanos() as u64)
428            .unwrap_or(0);
429        Ok(Some((m.len(), mtime_nanos)))
430    }
431
432    fn list_archives(&self) -> std::io::Result<Vec<u64>> {
433        let mut ns = Vec::new();
434        for entry in std::fs::read_dir(&self.dir)? {
435            let entry = entry?;
436            let name = entry.file_name();
437            let s = name.to_string_lossy();
438            if let Some(mid) = s
439                .strip_prefix("wal.")
440                .and_then(|r| r.strip_suffix(".archive"))
441            {
442                if let Ok(n) = mid.parse::<u64>() {
443                    ns.push(n);
444                }
445            }
446        }
447        ns.sort_unstable();
448        Ok(ns)
449    }
450
451    fn read_archive(&self, n: u64) -> std::io::Result<Vec<u8>> {
452        let path = self.dir.join(format!("wal.{n}.archive"));
453        match std::fs::read(&path) {
454            Ok(b) => Ok(b),
455            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(vec![]),
456            Err(e) => Err(e),
457        }
458    }
459
460    fn archive_wal(&mut self, n: u64) -> std::io::Result<()> {
461        let wal_path = self.path(FileId::Wal);
462        let archive_path = self.dir.join(format!("wal.{n}.archive"));
463        std::fs::rename(&wal_path, &archive_path)?;
464        sync_dir(&self.dir)
465    }
466
467    fn delete_archive(&mut self, n: u64) -> std::io::Result<()> {
468        let path = self.dir.join(format!("wal.{n}.archive"));
469        match std::fs::remove_file(&path) {
470            Ok(()) => sync_dir(&self.dir),
471            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
472            Err(e) => Err(e),
473        }
474    }
475
476    fn read_horizon_floor(&self) -> std::io::Result<u64> {
477        let path = self.dir.join("wal.floor");
478        match std::fs::read(&path) {
479            Ok(b) if b.len() >= 8 => Ok(u64::from_le_bytes([
480                b[0], b[1], b[2], b[3], b[4], b[5], b[6], b[7],
481            ])),
482            Ok(_) => Ok(0),
483            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(0),
484            Err(e) => Err(e),
485        }
486    }
487
488    fn write_horizon_floor(&mut self, floor: u64) -> std::io::Result<()> {
489        let tmp = self.dir.join("wal.floor.tmp");
490        {
491            let mut f = File::create(&tmp)?;
492            f.write_all(&floor.to_le_bytes())?;
493            full_sync(&f)?;
494        }
495        std::fs::rename(&tmp, self.dir.join("wal.floor"))?;
496        sync_dir(&self.dir)
497    }
498
499    fn has_genesis_marker(&self) -> bool {
500        self.dir.join("wal.genesis").exists()
501    }
502
503    fn write_genesis_marker(&mut self) -> std::io::Result<()> {
504        let path = self.dir.join("wal.genesis");
505        {
506            let mut f = File::create(&path)?;
507            f.write_all(b"")?;
508            full_sync(&f)?;
509        }
510        sync_dir(&self.dir)
511    }
512
513    fn delete_genesis_marker(&mut self) -> std::io::Result<()> {
514        match std::fs::remove_file(self.dir.join("wal.genesis")) {
515            Ok(()) => sync_dir(&self.dir),
516            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
517            Err(e) => Err(e),
518        }
519    }
520}
521
522fn full_sync(file: &File) -> std::io::Result<()> {
523    #[cfg(target_os = "macos")]
524    {
525        use std::os::unix::io::AsRawFd;
526        let fd = file.as_raw_fd();
527        let rc = unsafe { libc::fcntl(fd, libc::F_FULLFSYNC) };
528        if rc == -1 {
529            return Err(std::io::Error::last_os_error());
530        }
531        Ok(())
532    }
533    #[cfg(not(target_os = "macos"))]
534    {
535        file.sync_all()
536    }
537}
538
539/// Sync the WAL file at `dir/wal.bin` to persistent storage without
540/// requiring a `&mut Fs`.  Used by the group-commit drain thread to fsync
541/// outside the exclusive write-lock window (reducing reader-visible latency).
542///
543/// On macOS, uses `F_FULLFSYNC` for true durability.  On other platforms,
544/// falls back to `fdatasync` / `fsync`.  Returns `Ok(())` if the WAL file
545/// does not exist (nothing to sync).
546pub fn sync_wal_at(dir: &std::path::Path) -> std::io::Result<()> {
547    let path = dir.join(FileId::Wal.name());
548    let f = match std::fs::File::open(&path) {
549        Ok(f) => f,
550        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
551        Err(e) => return Err(e),
552    };
553    full_sync(&f)
554}
555
556/// Truncate the WAL file at `dir/wal.bin` to exactly `len` bytes and fsync
557/// the truncation to persistent storage.
558///
559/// Used by the group-commit drain thread when a group fsync fails: truncating
560/// the WAL back to the last known-good synced offset removes the unsynced
561/// frames, ensuring a crash-then-replay cannot silently make the failed group
562/// durable via a later successful fsync flushing the whole inode.
563///
564/// Returns `Ok(())` if the file does not exist (nothing to truncate).
565pub fn truncate_wal_at(dir: &std::path::Path, len: u64) -> std::io::Result<()> {
566    let path = dir.join(FileId::Wal.name());
567    let f = match OpenOptions::new().write(true).open(&path) {
568        Ok(f) => f,
569        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
570        Err(e) => return Err(e),
571    };
572    f.set_len(len)?;
573    f.sync_all() // plain sync_all is sufficient for a truncation barrier
574}
575
576fn sync_dir(dir: &std::path::Path) -> std::io::Result<()> {
577    let d = File::open(dir)?;
578    d.sync_all()
579}
580
581#[cfg(test)]
582mod tests {
583    use super::*;
584
585    fn tmp() -> std::path::PathBuf {
586        let d = std::env::temp_dir().join(format!("graphdb-fs-{}", std::process::id()));
587        let _ = std::fs::remove_dir_all(&d);
588        d
589    }
590
591    #[test]
592    fn append_read_and_atomic_write() {
593        let mut fs = RealFs::new(&tmp()).unwrap();
594        assert_eq!(fs.read(FileId::Wal).unwrap(), Vec::<u8>::new()); // absent = empty
595        fs.append(FileId::Wal, b"ab").unwrap();
596        fs.append(FileId::Wal, b"cd").unwrap();
597        fs.sync(FileId::Wal).unwrap();
598        assert_eq!(fs.read(FileId::Wal).unwrap(), b"abcd");
599        fs.write_atomic(FileId::Snapshot, b"snap1").unwrap();
600        fs.write_atomic(FileId::Snapshot, b"snap2").unwrap(); // replaces
601        assert_eq!(fs.read(FileId::Snapshot).unwrap(), b"snap2");
602        fs.write_atomic(FileId::Wal, b"").unwrap(); // truncation path
603        assert_eq!(fs.read(FileId::Wal).unwrap(), Vec::<u8>::new());
604    }
605
606    #[test]
607    fn write_atomic_replaces_and_still_readable() {
608        // existing append_read_and_atomic_write already covers replace;
609        // keep it; dir-sync is best-effort observable only via crash tests.
610        // Do not fake F_FULLFSYNC in SimFs.
611        let d = std::env::temp_dir().join(format!("graphdb-fs-atomic-{}", std::process::id()));
612        let _ = std::fs::remove_dir_all(&d);
613        let mut fs = RealFs::new(&d).unwrap();
614        fs.write_atomic(FileId::Snapshot, b"snap1").unwrap();
615        fs.write_atomic(FileId::Snapshot, b"snap2").unwrap();
616        assert_eq!(fs.read(FileId::Snapshot).unwrap(), b"snap2");
617    }
618}