From 385c30392f8698a2badf30ff6c9081afa41b4318 Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Mon, 28 Sep 2026 05:13:29 +0800 Subject: [PATCH 1/2] fix(vfs): share advisory locks across VFS instances by lock file Each native backend (Tokio, io_uring, GCD, IOCP) kept its own in-process lock table, so distinct VFS instances rooted at the same directory could each believe they held the OS lock. GCD and IOCP also each carried their own bespoke fcntl/LockFileEx implementation. Replace the per-VFS tables with one process-wide table in oslock, keyed by the resolved (canonicalized parent + file name) lock path, so every instance over one directory contends correctly regardless of path spelling. A per-domain gate serializes acquiring the OS lock, and the single OS lock descriptor is now owned by the shared entry instead of each handle, so an early handle close can no longer drop the F_SETLK lock out from under a remaining in-process holder. GCD and IOCP now route through the same implementation instead of their own copies. --- src/vfs/gcd/vfs.rs | 150 +-------------------- src/vfs/iocp/vfs.rs | 186 ++------------------------ src/vfs/iouring/vfs.rs | 10 +- src/vfs/oslock.rs | 279 +++++++++++++++++++++++---------------- src/vfs/tokio_backend.rs | 12 +- 5 files changed, 185 insertions(+), 452 deletions(-) diff --git a/src/vfs/gcd/vfs.rs b/src/vfs/gcd/vfs.rs index 45dc378..7327990 100644 --- a/src/vfs/gcd/vfs.rs +++ b/src/vfs/gcd/vfs.rs @@ -1,20 +1,12 @@ //! `GcdVfs`: macOS / iOS / iPadOS VFS rooted at a directory, using Grand //! Central Dispatch I/O for per-file reads and writes. Advisory path locking -//! uses the same in-process state machine + POSIX `flock` protocol as the -//! Tokio fallback. +//! is the shared `oslock` implementation. //! //! `dispatch_io` covers reads and writes only. Path operations and directory -//! sync are plain blocking syscalls, so they run on the blocking pool; only -//! path validation and the in-process lock table stay on the executor. -#![allow(unsafe_code)] - -use std::collections::BTreeMap; -use std::os::unix::io::AsRawFd; +//! sync are plain blocking syscalls, so they run on the blocking pool. use std::path::PathBuf; use std::sync::Arc; -use parking_lot::Mutex; - use dispatch2::{DispatchQoS, DispatchQueue, DispatchRetained, GlobalQueueIdentifier}; use crate::Result; @@ -22,96 +14,15 @@ use crate::errors::PagedbError; use super::file::GcdFile; use crate::vfs::blocking::offload; +use crate::vfs::oslock::LockKind; use crate::vfs::traits::{Vfs, canonical_native_path, resolve_native_path}; use crate::vfs::types::OpenMode; -#[derive(Debug, Clone, Copy)] -enum LockState { - Free, - Exclusive, - Shared(u32), -} - -#[derive(Debug, Clone, Copy)] -enum LockKind { - Exclusive, - Shared, -} - -struct InProcLockEntry { - state: Mutex, -} - -struct OsFcntlHandle { - _file: std::fs::File, -} - -impl OsFcntlHandle { - fn try_acquire(path: &std::path::Path, kind: LockKind) -> Result { - let file = std::fs::OpenOptions::new() - .create(true) - .truncate(false) - .read(true) - .write(true) - .open(path) - .map_err(PagedbError::Io)?; - - let fd = file.as_raw_fd(); - #[allow(clippy::cast_possible_truncation)] - let l_type = match kind { - LockKind::Exclusive => libc::F_WRLCK as libc::c_short, - LockKind::Shared => libc::F_RDLCK as libc::c_short, - }; - #[allow(clippy::cast_possible_truncation)] - let flock = libc::flock { - l_type, - l_whence: libc::SEEK_SET as libc::c_short, - l_start: 0, - l_len: 0, - l_pid: 0, - }; - // SAFETY: `fd` valid (owned by `file`); `flock` fully initialised; - // F_SETLK is non-blocking. - let rc = unsafe { libc::fcntl(fd, libc::F_SETLK, &flock) }; - if rc == -1 { - let err = std::io::Error::last_os_error(); - let raw = err.raw_os_error().unwrap_or(0); - if raw == libc::EAGAIN || raw == libc::EACCES { - return Err(PagedbError::AlreadyLocked); - } - return Err(PagedbError::Io(err)); - } - Ok(Self { _file: file }) - } -} - -// SAFETY: fd is owned exclusively; struct moves whole-cloth. -unsafe impl Send for OsFcntlHandle {} - -pub struct GcdLockHandle { - lock_ref: Arc, - kind: LockKind, - _os_lock: OsFcntlHandle, -} - -impl Drop for GcdLockHandle { - fn drop(&mut self) { - let mut s = self.lock_ref.state.lock(); - match (self.kind, *s) { - (LockKind::Exclusive, LockState::Exclusive) - | (LockKind::Shared, LockState::Shared(1)) => *s = LockState::Free, - (LockKind::Shared, LockState::Shared(n)) if n > 1 => { - *s = LockState::Shared(n - 1); - } - _ => {} - } - } -} +pub use crate::vfs::oslock::NativeLockHandle as GcdLockHandle; struct GcdInner { root: PathBuf, queue: DispatchRetained, - locks: Mutex>>, } #[derive(Clone)] @@ -128,7 +39,6 @@ impl GcdVfs { inner: Arc::new(GcdInner { root: root.into(), queue, - locks: Mutex::new(BTreeMap::new()), }), } } @@ -137,57 +47,9 @@ impl GcdVfs { resolve_native_path(&self.inner.root, path) } - fn lookup_or_create_entry(&self, path: &str) -> Arc { - let mut locks = self.inner.locks.lock(); - locks - .entry(path.to_string()) - .or_insert_with(|| { - Arc::new(InProcLockEntry { - state: Mutex::new(LockState::Free), - }) - }) - .clone() - } - async fn do_lock(&self, path: &str, kind: LockKind) -> Result { - let logical_path = canonical_native_path(path)?; - let entry = self.lookup_or_create_entry(&logical_path); - { - let mut s = entry.state.lock(); - match (kind, *s) { - (LockKind::Exclusive, LockState::Free) => *s = LockState::Exclusive, - (LockKind::Shared, LockState::Free) => *s = LockState::Shared(1), - (LockKind::Shared, LockState::Shared(n)) => *s = LockState::Shared(n + 1), - _ => return Err(PagedbError::AlreadyLocked), - } - } - let lock_path = self.resolve(&logical_path)?; - // `F_SETLK` never waits on a conflict, but creating the sentinel file - // can still stall on the filesystem, so the pair goes to the pool. - let acquired = offload(move || { - if let Some(parent) = lock_path.parent() { - std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; - } - OsFcntlHandle::try_acquire(&lock_path, kind) - }) - .await; - match acquired { - Ok(os_lock) => Ok(GcdLockHandle { - lock_ref: entry, - kind, - _os_lock: os_lock, - }), - Err(e) => { - let mut s = entry.state.lock(); - match (kind, *s) { - (LockKind::Exclusive, LockState::Exclusive) - | (LockKind::Shared, LockState::Shared(1)) => *s = LockState::Free, - (LockKind::Shared, LockState::Shared(n)) => *s = LockState::Shared(n - 1), - _ => {} - } - Err(e) - } - } + let lock_path = self.resolve(&canonical_native_path(path)?)?; + crate::vfs::oslock::acquire(lock_path, kind).await } } diff --git a/src/vfs/iocp/vfs.rs b/src/vfs/iocp/vfs.rs index d507f55..2716050 100644 --- a/src/vfs/iocp/vfs.rs +++ b/src/vfs/iocp/vfs.rs @@ -1,150 +1,32 @@ //! `IocpVfs`: Windows IOCP-backed VFS rooted at a directory. Advisory path -//! locking uses an in-process state machine backed by `LockFileEx` for -//! cross-process exclusion — same protocol as the Tokio fallback. Segment +//! locking is the shared `oslock` implementation. Segment //! files open with `FILE_SHARE_DELETE` so tombstone-rename protocols succeed //! against held handles. //! //! Path operations have no overlapped form: `CreateFile`, `MoveFileEx`, //! directory enumeration and friends all park the calling thread. They run on -//! the blocking pool; only path validation and the in-process lock table stay -//! on the executor. -#![allow(unsafe_code)] +//! the blocking pool. -use std::collections::BTreeMap; use std::os::windows::fs::OpenOptionsExt; use std::os::windows::io::AsRawHandle; use std::path::PathBuf; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; -use parking_lot::Mutex; - use crate::Result; use crate::errors::PagedbError; use super::file::IocpFile; use super::port::Port; use crate::vfs::blocking::offload; +use crate::vfs::oslock::LockKind; use crate::vfs::traits::{Vfs, canonical_native_path, resolve_native_path}; use crate::vfs::types::OpenMode; -use windows_sys::Win32::Foundation::{ - ERROR_IO_PENDING, ERROR_LOCK_VIOLATION, GetLastError, HANDLE, -}; -use windows_sys::Win32::Storage::FileSystem::{ - FILE_FLAG_OVERLAPPED, LOCKFILE_EXCLUSIVE_LOCK, LOCKFILE_FAIL_IMMEDIATELY, LockFileEx, - UnlockFileEx, -}; -use windows_sys::Win32::System::IO::OVERLAPPED; - -// --------------------------------------------------------------------------- -// In-process lock state machine -// --------------------------------------------------------------------------- - -#[derive(Debug, Clone, Copy)] -enum LockState { - Free, - Exclusive, - Shared(u32), -} - -#[derive(Debug, Clone, Copy)] -enum LockKind { - Exclusive, - Shared, -} - -struct InProcLockEntry { - state: Mutex, -} - -// --------------------------------------------------------------------------- -// Cross-process lock via LockFileEx -// --------------------------------------------------------------------------- - -struct OsLockFileExHandle { - file: std::fs::File, -} - -impl OsLockFileExHandle { - fn try_acquire(path: &std::path::Path, kind: LockKind) -> Result { - // FILE_SHARE_READ | FILE_SHARE_WRITE: multiple processes must be able - // to open the sentinel file and contend on the lock. - const FILE_SHARE_READ_WRITE: u32 = 0x0000_0003; - - let file = std::fs::OpenOptions::new() - .create(true) - .truncate(false) - .read(true) - .write(true) - .share_mode(FILE_SHARE_READ_WRITE) - .open(path) - .map_err(PagedbError::Io)?; - - let handle = file.as_raw_handle() as HANDLE; - let flags = match kind { - LockKind::Exclusive => LOCKFILE_EXCLUSIVE_LOCK | LOCKFILE_FAIL_IMMEDIATELY, - LockKind::Shared => LOCKFILE_FAIL_IMMEDIATELY, - }; - - // SAFETY: `handle` is valid for the duration of this call (owned by - // `file` which is alive). Zero-initialised OVERLAPPED is the - // documented input for synchronous `LockFileEx` use. We lock the - // whole [0, u64::MAX) byte range. - let mut overlapped: OVERLAPPED = unsafe { std::mem::zeroed() }; - let rc = unsafe { LockFileEx(handle, flags, 0, u32::MAX, u32::MAX, &mut overlapped) }; - - if rc == 0 { - // SAFETY: documented pattern. - let err = unsafe { GetLastError() }; - if err == ERROR_LOCK_VIOLATION || err == ERROR_IO_PENDING { - return Err(PagedbError::AlreadyLocked); - } - return Err(PagedbError::Io(std::io::Error::last_os_error())); - } - - Ok(Self { file }) - } -} - -impl Drop for OsLockFileExHandle { - fn drop(&mut self) { - let handle = self.file.as_raw_handle() as HANDLE; - // SAFETY: `handle` valid until `file` drops at end of this method. - // Errors are ignored in Drop; closing the handle releases the lock - // regardless. - let mut overlapped: OVERLAPPED = unsafe { std::mem::zeroed() }; - let _ = unsafe { UnlockFileEx(handle, 0, u32::MAX, u32::MAX, &mut overlapped) }; - } -} +use windows_sys::Win32::Foundation::HANDLE; +use windows_sys::Win32::Storage::FileSystem::FILE_FLAG_OVERLAPPED; -// SAFETY: HANDLE is process-owned and not aliased between threads — the -// struct moves as a whole. -unsafe impl Send for OsLockFileExHandle {} - -// --------------------------------------------------------------------------- -// Public lock handle -// --------------------------------------------------------------------------- - -pub struct IocpLockHandle { - lock_ref: Arc, - kind: LockKind, - _os_lock: OsLockFileExHandle, -} - -impl Drop for IocpLockHandle { - fn drop(&mut self) { - let mut s = self.lock_ref.state.lock(); - match (self.kind, *s) { - (LockKind::Exclusive, LockState::Exclusive) - | (LockKind::Shared, LockState::Shared(1)) => *s = LockState::Free, - (LockKind::Shared, LockState::Shared(n)) if n > 1 => { - *s = LockState::Shared(n - 1); - } - _ => {} - } - } -} +pub use crate::vfs::oslock::NativeLockHandle as IocpLockHandle; // --------------------------------------------------------------------------- // IocpVfs @@ -157,7 +39,6 @@ struct IocpInner { /// means keys are not strictly required to disambiguate completions, but /// they are useful for diagnostics and future relaxation of serialisation. next_key: AtomicUsize, - locks: Mutex>>, } #[derive(Clone)] @@ -173,7 +54,6 @@ impl IocpVfs { root: root.into(), port, next_key: AtomicUsize::new(1), - locks: Mutex::new(BTreeMap::new()), }), }) } @@ -182,59 +62,9 @@ impl IocpVfs { resolve_native_path(&self.inner.root, path) } - fn lookup_or_create_entry(&self, path: &str) -> Arc { - let mut locks = self.inner.locks.lock(); - locks - .entry(path.to_string()) - .or_insert_with(|| { - Arc::new(InProcLockEntry { - state: Mutex::new(LockState::Free), - }) - }) - .clone() - } - async fn do_lock(&self, path: &str, kind: LockKind) -> Result { - let logical_path = canonical_native_path(path)?; - let entry = self.lookup_or_create_entry(&logical_path); - { - let mut s = entry.state.lock(); - match (kind, *s) { - (LockKind::Exclusive, LockState::Free) => *s = LockState::Exclusive, - (LockKind::Shared, LockState::Free) => *s = LockState::Shared(1), - (LockKind::Shared, LockState::Shared(n)) => *s = LockState::Shared(n + 1), - _ => return Err(PagedbError::AlreadyLocked), - } - } - let lock_path = self.resolve(&logical_path)?; - // Creating the sentinel file and taking the OS lock are both blocking - // syscalls. `LockFileEx` itself is `LOCKFILE_FAIL_IMMEDIATELY`, so it - // never waits on a conflict — but opening the file can still stall on - // the filesystem, so the pair goes to the pool together. - let acquired = offload(move || { - if let Some(parent) = lock_path.parent() { - std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; - } - OsLockFileExHandle::try_acquire(&lock_path, kind) - }) - .await; - match acquired { - Ok(os_lock) => Ok(IocpLockHandle { - lock_ref: entry, - kind, - _os_lock: os_lock, - }), - Err(e) => { - let mut s = entry.state.lock(); - match (kind, *s) { - (LockKind::Exclusive, LockState::Exclusive) - | (LockKind::Shared, LockState::Shared(1)) => *s = LockState::Free, - (LockKind::Shared, LockState::Shared(n)) => *s = LockState::Shared(n - 1), - _ => {} - } - Err(e) - } - } + let lock_path = self.resolve(&canonical_native_path(path)?)?; + crate::vfs::oslock::acquire(lock_path, kind).await } } diff --git a/src/vfs/iouring/vfs.rs b/src/vfs/iouring/vfs.rs index c75f279..f3ecc0a 100644 --- a/src/vfs/iouring/vfs.rs +++ b/src/vfs/iouring/vfs.rs @@ -17,7 +17,7 @@ use crate::errors::PagedbError; use super::file::IouringFile; use super::ring::Ring; use crate::vfs::blocking::offload; -use crate::vfs::oslock::{LockKind, LockTable}; +use crate::vfs::oslock::LockKind; use crate::vfs::traits::{Vfs, canonical_native_path, resolve_native_path}; use crate::vfs::types::OpenMode; @@ -33,12 +33,11 @@ pub use crate::vfs::oslock::NativeLockHandle as IouringLockHandle; struct IouringInner { root: PathBuf, ring: Ring, - locks: LockTable, } /// VFS rooted at a directory, using `io_uring` for file I/O and `std::fs` / -/// libc syscalls for path-level operations. Cloning shares the same root, -/// ring, and lock table. +/// libc syscalls for path-level operations. Cloning shares the same root and +/// ring. Locks are process-wide, so instances over one directory contend. #[derive(Clone)] pub struct IouringVfs { inner: Arc, @@ -54,7 +53,6 @@ impl IouringVfs { inner: Arc::new(IouringInner { root: root.into(), ring, - locks: LockTable::new(), }), }) } @@ -66,7 +64,7 @@ impl IouringVfs { async fn do_lock(&self, path: &str, kind: LockKind) -> Result { let logical_path = canonical_native_path(path)?; let lock_path = self.resolve(&logical_path)?; - crate::vfs::oslock::acquire(&self.inner.locks, &logical_path, lock_path, kind).await + crate::vfs::oslock::acquire(lock_path, kind).await } } diff --git a/src/vfs/oslock.rs b/src/vfs/oslock.rs index 318ab6e..c1bd24c 100644 --- a/src/vfs/oslock.rs +++ b/src/vfs/oslock.rs @@ -1,11 +1,17 @@ //! Advisory path locking shared by every filesystem-backed native VFS. //! -//! One protocol, one implementation. A lock is two layers: an in-process state -//! machine that gives fast single-process exclusion, and an OS-level lock that +//! One protocol, one implementation. A lock is two layers: a process-wide +//! table that excludes every handle in this process, and an OS-level lock that //! excludes other processes — `fcntl` OFD locks (`F_OFD_SETLK`) on Linux, //! classic `F_SETLK` on other Unix, and `LockFileEx` on Windows. On targets //! with neither, only the in-process layer applies. //! +//! The table is keyed by resolved lock-file path, so all VFS instances share +//! it. macOS `F_SETLK` never conflicts within one process. This table enforces +//! the conflict instead. +//! Each entry owns the process's single OS lock on its file. One descriptor +//! means an early close can never drop an `F_SETLK` lock. +//! //! Both layers live here rather than in each backend deliberately. The store's //! single-writer guarantee rests on the `.writer.lock` sentinel, and two //! backends may hold that sentinel at the same time — a process that got an @@ -48,66 +54,84 @@ pub(crate) enum LockKind { Shared, } +/// One lock domain: its holders in this process and the OS lock they share. +struct EntryState { + mode: LockState, + #[cfg(any(unix, windows))] + os: Option, +} + struct InProcLockEntry { - state: Mutex, + /// Serializes taking the OS lock for a free domain. + gate: tokio::sync::Mutex<()>, + state: Mutex, } impl InProcLockEntry { - /// Take the in-process layer, or report the conflict without touching it. - fn try_enter(&self, kind: LockKind) -> Result<()> { - let mut state = self.state.lock(); - match (kind, *state) { - (LockKind::Exclusive, LockState::Free) => *state = LockState::Exclusive, - (LockKind::Shared, LockState::Free) => *state = LockState::Shared(1), - (LockKind::Shared, LockState::Shared(n)) => *state = LockState::Shared(n + 1), - _ => return Err(PagedbError::AlreadyLocked), + fn new() -> Self { + Self { + gate: tokio::sync::Mutex::new(()), + state: Mutex::new(EntryState { + mode: LockState::Free, + #[cfg(any(unix, windows))] + os: None, + }), } - Ok(()) } - /// Give the in-process layer back. Used both when a handle drops and to - /// roll back after the OS layer refused, so a failed acquisition leaves no - /// trace. - fn leave(&self, kind: LockKind) { + /// Join a held domain. `false` means it is free and needs the OS lock. + fn try_join(&self, kind: LockKind) -> Result { let mut state = self.state.lock(); - match (kind, *state) { - (LockKind::Exclusive, LockState::Exclusive) - | (LockKind::Shared, LockState::Shared(1)) => *state = LockState::Free, - (LockKind::Shared, LockState::Shared(n)) if n > 1 => *state = LockState::Shared(n - 1), - _ => {} + match (kind, state.mode) { + (_, LockState::Free) => Ok(false), + (LockKind::Shared, LockState::Shared(n)) => { + state.mode = LockState::Shared(n + 1); + Ok(true) + } + _ => Err(PagedbError::AlreadyLocked), } } -} - -/// Per-VFS table of in-process lock entries, keyed by canonical logical path. -/// -/// Each distinct canonical path is its own lock domain. The table only ever -/// grows an entry per path that has been locked at least once; entries are -/// cheap and keeping them avoids racing a concurrent acquirer against removal. -pub(crate) struct LockTable { - entries: Mutex>>, -} -impl LockTable { - pub(crate) fn new() -> Self { - Self { - entries: Mutex::new(BTreeMap::new()), + /// Record the first holder of a free domain. + fn install(&self, kind: LockKind, #[cfg(any(unix, windows))] os: OsLock) { + let mut state = self.state.lock(); + state.mode = match kind { + LockKind::Exclusive => LockState::Exclusive, + LockKind::Shared => LockState::Shared(1), + }; + #[cfg(any(unix, windows))] + { + state.os = Some(os); } } - fn entry(&self, path: &str) -> Arc { - let mut entries = self.entries.lock(); - entries - .entry(path.to_string()) - .or_insert_with(|| { - Arc::new(InProcLockEntry { - state: Mutex::new(LockState::Free), - }) - }) - .clone() + /// Give one hold back. The last holder releases the OS lock. + fn leave(&self, kind: LockKind) { + let mut state = self.state.lock(); + state.mode = match (kind, state.mode) { + (LockKind::Shared, LockState::Shared(n)) if n > 1 => LockState::Shared(n - 1), + _ => LockState::Free, + }; + #[cfg(any(unix, windows))] + if matches!(state.mode, LockState::Free) { + state.os = None; + } } } +/// Every lock domain in the process, keyed by resolved lock-file path. +/// Entries are never removed, so no acquirer races a removal. +static LOCK_TABLE: std::sync::LazyLock>>> = + std::sync::LazyLock::new(|| Mutex::new(BTreeMap::new())); + +fn entry(key: PathBuf) -> Arc { + LOCK_TABLE + .lock() + .entry(key) + .or_insert_with(|| Arc::new(InProcLockEntry::new())) + .clone() +} + // --------------------------------------------------------------------------- // Unix cross-process lock via fcntl. // --------------------------------------------------------------------------- @@ -307,119 +331,142 @@ type OsLock = OsLockFileExHandle; // Public lock handle. // --------------------------------------------------------------------------- -/// RAII advisory lock handle returned by every native backend's -/// `lock_exclusive` / `lock_shared`. Holds the in-process state guard and, on -/// targets that have one, the OS-level lock. Dropping it releases both. +/// Advisory lock returned by every native backend's `lock_exclusive` and +/// `lock_shared`. The last holder of a file releases the OS lock on drop. pub struct NativeLockHandle { entry: Arc, kind: LockKind, - /// On Unix: holds the fcntl-locked file open (an OFD lock on Linux, a - /// process `F_SETLK` lock elsewhere). On Windows: holds the - /// `LockFileEx`-locked file open. Dropped together with this handle. - #[cfg(any(unix, windows))] - _os_lock: OsLock, } impl Drop for NativeLockHandle { fn drop(&mut self) { self.entry.leave(self.kind); - // `_os_lock` is dropped automatically after this, releasing the OS lock. } } -/// Acquire an advisory lock on one canonical logical path. -/// -/// `logical_path` names the lock domain in the in-process table; -/// `lock_path` is the sentinel file on disk that carries the OS-level lock. -pub(crate) async fn acquire( - table: &LockTable, - logical_path: &str, - lock_path: PathBuf, - kind: LockKind, -) -> Result { - let entry = table.entry(logical_path); - // In-process guard first: fast fail if this process already holds a - // conflicting lock on the path, without touching the filesystem. - entry.try_enter(kind)?; - +/// Acquire an advisory lock on the sentinel file at `lock_path`, creating the +/// file and its directory when absent. +pub(crate) async fn acquire(lock_path: PathBuf, kind: LockKind) -> Result { #[cfg(any(unix, windows))] { - // Neither `*_SETLK` nor `LOCKFILE_FAIL_IMMEDIATELY` waits on a - // conflict, but creating the sentinel file can still stall on the - // filesystem, so the pair goes to the blocking pool together. - let acquired = offload(move || { - if let Some(parent) = lock_path.parent() { - std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; - } - OsLock::try_acquire(&lock_path, kind) - }) - .await; - match acquired { - Ok(os_lock) => Ok(NativeLockHandle { - entry, - kind, - _os_lock: os_lock, - }), - Err(error) => { - // The OS layer refused, so the in-process layer must not stay - // taken — a rejected acquisition leaves no trace. - entry.leave(kind); - Err(error) + // Filesystem calls, so off the async thread. + let key = offload(move || resolve_lock_key(&lock_path)).await?; + let entry = entry(key.clone()); + { + let _gate = entry.gate.lock().await; + if !entry.try_join(kind)? { + // Creating the sentinel file can stall on the filesystem. + let os = offload(move || OsLock::try_acquire(&key, kind)).await?; + entry.install(kind, os); } } + Ok(NativeLockHandle { entry, kind }) } #[cfg(not(any(unix, windows)))] { - let _ = lock_path; + let entry = entry(lock_path); + { + let _gate = entry.gate.lock().await; + if !entry.try_join(kind)? { + entry.install(kind); + } + } Ok(NativeLockHandle { entry, kind }) } } +/// Canonical directory plus file name, so two spellings of one path match. +#[cfg(any(unix, windows))] +fn resolve_lock_key(lock_path: &std::path::Path) -> Result { + let parent = lock_path + .parent() + .filter(|parent| !parent.as_os_str().is_empty()) + .unwrap_or_else(|| std::path::Path::new(".")); + std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; + let parent = std::fs::canonicalize(parent).map_err(PagedbError::Io)?; + let name = lock_path.file_name().ok_or_else(|| { + PagedbError::Io(std::io::Error::other(format!( + "lock path has no file name: {}", + lock_path.display() + ))) + })?; + Ok(parent.join(name)) +} + #[cfg(test)] mod tests { use super::*; - #[test] - fn an_exclusive_entry_excludes_every_other_kind() { - let table = LockTable::new(); - let entry = table.entry("/db"); - entry.try_enter(LockKind::Exclusive).unwrap(); + fn lock_file(dir: &tempfile::TempDir, name: &str) -> PathBuf { + dir.path().join(name) + } + + #[tokio::test(flavor = "current_thread")] + async fn an_exclusive_lock_excludes_every_other_kind() { + let dir = tempfile::tempdir().unwrap(); + let _held = acquire(lock_file(&dir, "a.lock"), LockKind::Exclusive) + .await + .unwrap(); assert!(matches!( - entry.try_enter(LockKind::Exclusive), + acquire(lock_file(&dir, "a.lock"), LockKind::Exclusive).await, Err(PagedbError::AlreadyLocked) )); assert!(matches!( - entry.try_enter(LockKind::Shared), + acquire(lock_file(&dir, "a.lock"), LockKind::Shared).await, Err(PagedbError::AlreadyLocked) )); } - #[test] - fn shared_entries_stack_and_only_the_last_release_frees_the_domain() { - let table = LockTable::new(); - let entry = table.entry("/db"); - entry.try_enter(LockKind::Shared).unwrap(); - entry.try_enter(LockKind::Shared).unwrap(); - - entry.leave(LockKind::Shared); + #[tokio::test(flavor = "current_thread")] + async fn shared_holds_stack_and_only_the_last_release_frees_the_file() { + let dir = tempfile::tempdir().unwrap(); + let first = acquire(lock_file(&dir, "a.lock"), LockKind::Shared) + .await + .unwrap(); + let second = acquire(lock_file(&dir, "a.lock"), LockKind::Shared) + .await + .unwrap(); + + drop(first); assert!( matches!( - entry.try_enter(LockKind::Exclusive), + acquire(lock_file(&dir, "a.lock"), LockKind::Exclusive).await, Err(PagedbError::AlreadyLocked) ), "one shared holder remains, so exclusive must still be refused" ); - entry.leave(LockKind::Shared); - entry.try_enter(LockKind::Exclusive).unwrap(); + drop(second); + acquire(lock_file(&dir, "a.lock"), LockKind::Exclusive) + .await + .unwrap(); + } + + #[tokio::test(flavor = "current_thread")] + async fn two_spellings_of_one_lock_file_share_one_domain() { + let dir = tempfile::tempdir().unwrap(); + let _held = acquire(lock_file(&dir, "a.lock"), LockKind::Exclusive) + .await + .unwrap(); + let respelled = dir.path().join("sub").join("..").join("a.lock"); + std::fs::create_dir_all(dir.path().join("sub")).unwrap(); + assert!( + matches!( + acquire(respelled, LockKind::Exclusive).await, + Err(PagedbError::AlreadyLocked) + ), + "a second spelling of the path must not bypass the holder" + ); } - #[test] - fn the_same_logical_path_always_maps_to_one_entry() { - let table = LockTable::new(); - let first = table.entry("/db"); - let second = table.entry("/db"); - assert!(Arc::ptr_eq(&first, &second)); - assert!(!Arc::ptr_eq(&first, &table.entry("/other"))); + #[tokio::test(flavor = "current_thread")] + async fn distinct_lock_files_are_distinct_domains() { + let dir = tempfile::tempdir().unwrap(); + let _a = acquire(lock_file(&dir, "a.lock"), LockKind::Exclusive) + .await + .unwrap(); + acquire(lock_file(&dir, "b.lock"), LockKind::Exclusive) + .await + .unwrap(); } } diff --git a/src/vfs/tokio_backend.rs b/src/vfs/tokio_backend.rs index 8d15844..30e5116 100644 --- a/src/vfs/tokio_backend.rs +++ b/src/vfs/tokio_backend.rs @@ -13,7 +13,7 @@ use crate::Result; use crate::errors::PagedbError; use super::blocking::offload; -use super::oslock::{LockKind, LockTable}; +use super::oslock::LockKind; use super::traits::{Vfs, VfsFile, canonical_native_path, resolve_native_path}; use super::types::{OpenMode, ReadReq, WriteReq}; @@ -28,7 +28,7 @@ pub use super::oslock::NativeLockHandle as TokioLockHandle; /// VFS rooted at a directory. Paths supplied to all methods are resolved /// relative to this root; a leading `/` is stripped. Cloning shares the same -/// root and lock table. +/// root. Locks are process-wide, so instances over one directory contend. #[derive(Clone)] pub struct TokioVfs { inner: Arc, @@ -36,7 +36,6 @@ pub struct TokioVfs { struct TokioInner { root: PathBuf, - locks: LockTable, } impl TokioVfs { @@ -44,10 +43,7 @@ impl TokioVfs { /// exist or be created before the first `open` call. pub fn new(root: impl Into) -> Self { Self { - inner: Arc::new(TokioInner { - root: root.into(), - locks: LockTable::new(), - }), + inner: Arc::new(TokioInner { root: root.into() }), } } @@ -68,7 +64,7 @@ impl TokioVfs { async fn do_lock(&self, path: &str, kind: LockKind) -> Result { let logical_path = Self::canonical_logical_path(path)?; let lock_path = self.resolve(&logical_path)?; - super::oslock::acquire(&self.inner.locks, &logical_path, lock_path, kind).await + super::oslock::acquire(lock_path, kind).await } } From 7badf9ae7f6a9999adb401fc95255ae8ccfb0735 Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Mon, 28 Sep 2026 05:14:06 +0800 Subject: [PATCH 2/2] feat(open): fork a restored directory into an independent writer A snapshot directory and a directory `restore_from` fills copy `main.db` byte for byte, so they share the source's DEK and nonce space with it. Opening one Standalone previously proceeded anyway, letting independent writes on both directories repeat nonces under one key. Record a restore mode (STANDALONE, READ_ONLY, FOLLOWER) in the main.db header. A Standalone open now refuses any non-STANDALONE mode with RestoredNotPromoted. A restored directory can still open ReadOnly, or promote to Follower and track its source. `rekey_into_writer` forks a ReadOnly or Follower handle into a Standalone writer under a fresh identity: it rekeys the tree and quota catalogs and rewrites segments into a fork directory the open-time orphan scan skips, then adopts those segments into staging and publishes a fresh header once the fork takes the writer sentinel. Thread the KEK-changing-rekey resume path through open_with_mode via an optional counterpart KEK, and remove_if_present to make crash cleanup of scratch files idempotent. --- CHANGELOG.md | 2 +- README.md | 5 + src/compaction/full.rs | 4 +- src/compaction/helpers.rs | 17 +- src/pager/header.rs | 58 ++++- src/recovery/fork.rs | 70 ++++++ src/recovery/journal.rs | 4 +- src/recovery/mod.rs | 1 + src/recovery/quarantine/preserve.rs | 9 +- src/recovery/reconcile.rs | 12 +- src/segment/writer.rs | 90 ++++++- src/snapshot/tests/basic.rs | 4 +- src/txn/db/apply_journal.rs | 39 +-- src/txn/db/catalog/quotas.rs | 29 ++- src/txn/db/core.rs | 62 ++++- src/txn/db/mod.rs | 1 + src/txn/db/open/create.rs | 29 ++- src/txn/db/open/existing.rs | 83 +++---- src/txn/db/open/mod.rs | 1 + src/txn/db/open/modes.rs | 123 ++++------ src/txn/db/open/promote.rs | 137 +++++++++++ src/txn/db/open/recovery.rs | 22 +- src/txn/db/rekey/fork/entry.rs | 105 ++++++++ src/txn/db/rekey/fork/mod.rs | 6 + src/txn/db/rekey/fork/segments.rs | 67 ++++++ src/txn/db/rekey/fork/trees.rs | 199 ++++++++++++++++ src/txn/db/rekey/mod.rs | 1 + src/txn/db/restore_mode/mod.rs | 20 ++ src/txn/db/restore_mode/stamp.rs | 87 +++++++ src/txn/db/restore_mode/values.rs | 9 + src/txn/db/segment.rs | 10 +- src/txn/db/snapshot.rs | 50 +++- src/txn/db/util.rs | 70 +----- src/txn/write/commit.rs | 4 +- src/vfs/mod.rs | 2 +- src/vfs/traits.rs | 12 + tests/common/mod.rs | 69 ++++++ tests/open_mode_discipline.rs | 349 +++++++++++++++++++++++++++ tests/restored_store_fork.rs | 356 ++++++++++++++++++++++++++++ tests/restored_store_modes.rs | 178 ++++++++++++++ 40 files changed, 2083 insertions(+), 313 deletions(-) create mode 100644 src/recovery/fork.rs create mode 100644 src/txn/db/open/promote.rs create mode 100644 src/txn/db/rekey/fork/entry.rs create mode 100644 src/txn/db/rekey/fork/mod.rs create mode 100644 src/txn/db/rekey/fork/segments.rs create mode 100644 src/txn/db/rekey/fork/trees.rs create mode 100644 src/txn/db/restore_mode/mod.rs create mode 100644 src/txn/db/restore_mode/stamp.rs create mode 100644 src/txn/db/restore_mode/values.rs create mode 100644 tests/common/mod.rs create mode 100644 tests/open_mode_discipline.rs create mode 100644 tests/restored_store_fork.rs create mode 100644 tests/restored_store_modes.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index f670844..598bee9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -17,7 +17,7 @@ No version has been released yet. Pre-releases are published as `0.1.0-beta.N`; - **Snapshots** — `snapshot_to`, `restore_from`, and incremental apply, each authenticated against the state its manifest describes. Destinations must be empty; malformed or incomplete artifacts fail closed. - **Recovery** — open-flow GC, apply-journal replay, deep-walk `fsck`, and the `pagedb-fsck` binary. - **Online rekey** — rekey under a new key with mixed-cipher and mixed-epoch page coexistence; no full-file migration. A rotation whose source epoch is still pinned by a reader defers retiring it and completes the retirement once the reader set drains — including when the last reader leaves during the deferral itself, so the superseded master key never stays leasable behind an `Ok(())`. A retirement that cannot be taken yet stays queued for the next attempt without holding up the others. -- **Handle modes** — `Standalone`, `Follower`, `ReadOnly`, and `Observer`. +- **Handle modes** — `Standalone`, `Follower`, `ReadOnly`, and `Observer`, each holding its own sentinel. A snapshot or restored directory shares its source's identity, so it opens only as `ReadOnly` or as a `Follower` (`open_follower`, `promote_to_follower`), and a `Standalone` open reports `RestoredNotPromoted`. `rekey_into_writer` forks it into an independent `Standalone` writer under a fresh identity. - **Open refusals name the parameter, not the store** — `KeyMismatch`, `PageSizeMismatch`, and `RealmMismatch`, each decided before anything is read or written, and none reported as corruption. - **Failures report themselves** — an unreadable free-list chain, main file, or segment catalog fails `stats()` instead of reporting zero; compaction never skips a catalog entry whose file it cannot open; segment open distinguishes a missing file from a permission or backend error; and only genuine contention is reported as contention. Persisted named-counter rows are validated at open, and commit-history keys are rejected unless exactly eight bytes. diff --git a/README.md b/README.md index 4c65622..467a83f 100644 --- a/README.md +++ b/README.md @@ -185,6 +185,11 @@ defend against, stated plainly so nobody has to infer it. capability bit this build does not understand reports `HeaderCapabilityUnsupported`: a newer build wrote the store, and it is untouched. Open it with a build that understands the capability. +- **A restored directory is not a writer.** A snapshot directory, and a + directory `restore_from` fills, shares the source's key and nonce space. A + Standalone open of one reports `RestoredNotPromoted`. Open it with + `open_read_only`, track the source with `open_follower`, or call + `rekey_into_writer` to fork an independent writer under a fresh identity. - **`main.db` is not reconstructible from `seg/`.** Segment files are identity-keyed and the mapping from embedder name to segment id lives only in the catalog inside `main.db`. Losing `main.db` while `seg/` survives is diff --git a/src/compaction/full.rs b/src/compaction/full.rs index abd5c47..f6b51bb 100644 --- a/src/compaction/full.rs +++ b/src/compaction/full.rs @@ -11,7 +11,7 @@ use crate::catalog::codec::{Catalog, CatalogRowKind, SegmentMeta}; use crate::errors::PagedbError; use crate::segment::reader::SegmentReader; use crate::segment::types::SegmentPageKind; -use crate::segment::writer::SegmentWriter; +use crate::segment::writer::{STAGING_DIR, SegmentWriter}; use crate::txn::db::{Db, WriterState}; use crate::vfs::{Vfs, VfsFile}; @@ -149,7 +149,7 @@ async fn repack_one_segment( mmap_limit, ) .await?; - db.vfs.mkdir_all("seg/.staging").await?; + db.vfs.mkdir_all(STAGING_DIR).await?; let new_segment_id = crate::crypto::random::segment_id()?; let mut writer = SegmentWriter::create_internal( db.pager.clone(), diff --git a/src/compaction/helpers.rs b/src/compaction/helpers.rs index 8dd40c0..637fc36 100644 --- a/src/compaction/helpers.rs +++ b/src/compaction/helpers.rs @@ -187,7 +187,7 @@ pub(super) async fn replace_segment_compact( format_version: crate::pager::structural_header::MAIN_FORMAT_VERSION, cipher_id: db.cipher_id.as_byte(), page_size_log2: page_size_log2(db.page_size)?, - flags: 0, + flags: db.header_flags, file_id: db.file_id, kek_salt: db.kek_salt, mk_epoch: db.mk_epoch.load(std::sync::atomic::Ordering::SeqCst), @@ -202,10 +202,10 @@ pub(super) async fn replace_segment_compact( apply_journal_root_version: 0, commit_history_root_page_id: 0, commit_history_root_version: 0, - restore_mode: 0, + restore_mode: state.restore_mode, next_page_id: new_next, - commit_retain_policy_tag: 0, - commit_retain_policy_value: 0, + commit_retain_policy_tag: state.commit_retain_policy_tag, + commit_retain_policy_value: state.commit_retain_policy_value, realm_id: db.realm_id, }; @@ -272,7 +272,6 @@ pub(super) fn make_header_fields( new_next: u64, free_list_root_page_id: u64, ) -> MainDbHeaderFields { - let _ = state; let mut catalog_root_bytes = [0u8; 16]; catalog_root_bytes[..8].copy_from_slice(&new_cat_root.to_le_bytes()); catalog_root_bytes[8..].copy_from_slice(&new_commit_id.to_le_bytes()); @@ -280,7 +279,7 @@ pub(super) fn make_header_fields( format_version: crate::pager::structural_header::MAIN_FORMAT_VERSION, cipher_id: db.cipher_id.as_byte(), page_size_log2: page_size_log2(db.page_size).unwrap_or(12), - flags: 0, + flags: db.header_flags, file_id: db.file_id, kek_salt: db.kek_salt, mk_epoch: db.mk_epoch.load(std::sync::atomic::Ordering::SeqCst), @@ -295,10 +294,10 @@ pub(super) fn make_header_fields( apply_journal_root_version: 0, commit_history_root_page_id: 0, commit_history_root_version: 0, - restore_mode: 0, + restore_mode: state.restore_mode, next_page_id: new_next, - commit_retain_policy_tag: 0, - commit_retain_policy_value: 0, + commit_retain_policy_tag: state.commit_retain_policy_tag, + commit_retain_policy_value: state.commit_retain_policy_value, realm_id: db.realm_id, } } diff --git a/src/pager/header.rs b/src/pager/header.rs index ebd2322..55f3b7f 100644 --- a/src/pager/header.rs +++ b/src/pager/header.rs @@ -5,16 +5,17 @@ //! `seq` (HK-MAC-verified). A torn write to one slot leaves the other intact. use crate::Result; +use crate::crypto::SecretKey; +use crate::crypto::kdf::{derive_hk, derive_mk}; use crate::crypto::keys::DerivedKey; -// `decode_main_db_header` and the corruption detail it reports are needed only -// by `open_header`, which is test-only: production openers decode per slot -// because the HK has to be derived from each slot's own salt and epoch. +// `CorruptionDetail` is used only by `open_header`, which is test-only. +// Production openers authenticate per slot because each slot's HK derives from its own salt and epoch. #[cfg(test)] use crate::errors::CorruptionDetail; use crate::errors::PagedbError; -#[cfg(test)] -use crate::pager::format::structural_header::decode_main_db_header; -use crate::pager::format::structural_header::{MainDbHeaderFields, encode_main_db_header}; +use crate::pager::format::structural_header::{ + MainDbHeaderFields, decode_main_db_header, encode_main_db_header, +}; use crate::vfs::types::OpenMode; use crate::vfs::{Vfs, VfsFile, read_exact_at, write_all_at}; @@ -109,6 +110,51 @@ pub(crate) async fn read_header_slot( } } +/// Authenticate one header slot under `hk`. +/// +/// Returns `None` when the slot fails to verify. A wrong key, a torn write, +/// and damage look identical. +/// +/// Returns `Err` only for `HeaderCapabilityUnsupported`. The key already +/// authenticated the slot, so no other key or slot revokes its capability. +pub(crate) fn authenticate_slot( + slot: &[u8], + hk: &DerivedKey, + page_size: usize, +) -> Result> { + match decode_main_db_header(slot, hk, page_size) { + Ok(fields) => Ok(Some(fields)), + Err(error @ PagedbError::HeaderCapabilityUnsupported { .. }) => Err(error), + Err(_) => Ok(None), + } +} + +/// Authenticate one header slot under `kek`, as [`authenticate_slot`] does. +/// +/// The HK derives from the KEK salt and MK epoch, stored unencrypted at +/// bytes `32..48` and `48..56`. Returns the HK with the fields so the +/// caller can rewrite the slot. +pub(crate) fn authenticate_slot_with_kek( + slot: &[u8], + kek: &SecretKey, + page_size: usize, +) -> Result> { + if slot.len() < 56 { + return Ok(None); + } + let mut salt = [0u8; 16]; + salt.copy_from_slice(&slot[32..48]); + let mut epoch = [0u8; 8]; + epoch.copy_from_slice(&slot[48..56]); + let Ok(mk) = derive_mk(kek.as_bytes(), &salt, u64::from_le_bytes(epoch)) else { + return Ok(None); + }; + let Ok(hk) = derive_hk(&mk) else { + return Ok(None); + }; + Ok(authenticate_slot(slot, &hk, page_size)?.map(|fields| (fields, hk))) +} + /// Reads both slots; verifies each via HK-MAC; picks the one with the /// greater `seq`. If only one verifies, it wins. If neither verifies, returns /// `Corruption(HeaderUnverifiable)` — unrecoverable from inside the header diff --git a/src/recovery/fork.rs b/src/recovery/fork.rs new file mode 100644 index 0000000..b7e4251 --- /dev/null +++ b/src/recovery/fork.rs @@ -0,0 +1,70 @@ +//! Segments a `rekey_into_writer` fork wrote before its `main.db` was live. +//! +//! The fork writes its segments in [`FORK_DIR`], which the open-time orphan +//! scan skips. A fork that never publishes leaves the restored directory +//! openable. Once the fork's `main.db` is live, its catalog names those +//! segments. Moving them into staging hands them to catalog repair, which +//! publishes each one the catalog names and sweeps the rest. + +use crate::Result; +use crate::errors::PagedbError; +use crate::segment::writer::{FORK_DIR, STAGING_DIR}; +use crate::vfs::Vfs; + +/// Move every fork segment into staging. Returns the number moved. +/// +/// Callers must hold write authority over a store whose `main.db` is the +/// fork's. Callers must run catalog repair afterwards. +pub(crate) async fn adopt_forked_segments(vfs: &V) -> Result { + let names = forked_segment_names(vfs).await?; + if names.is_empty() { + return Ok(0); + } + vfs.mkdir_all(STAGING_DIR).await?; + let mut count: u64 = 0; + for name in names { + vfs.rename( + &format!("{FORK_DIR}/{name}"), + &format!("{STAGING_DIR}/{name}"), + ) + .await?; + count += 1; + } + vfs.sync_dir(STAGING_DIR).await?; + vfs.sync_dir(FORK_DIR).await?; + Ok(count) +} + +/// Remove every fork segment. Returns the number removed. +/// +/// A fork calls this before it starts, to drop what an earlier attempt +/// wrote without publishing. +pub(crate) async fn discard_forked_segments(vfs: &V) -> Result { + let names = forked_segment_names(vfs).await?; + if names.is_empty() { + return Ok(0); + } + let mut count: u64 = 0; + for name in names { + vfs.remove(&format!("{FORK_DIR}/{name}")).await?; + count += 1; + } + vfs.sync_dir(FORK_DIR).await?; + Ok(count) +} + +/// Names in the fork directory that this crate writes: segment identities. +/// Anything else is left alone. +async fn forked_segment_names(vfs: &V) -> Result> { + let entries = match vfs.list_dir(FORK_DIR).await { + Ok(entries) => entries, + Err(PagedbError::Io(error)) if error.kind() == std::io::ErrorKind::NotFound => { + return Ok(Vec::new()); + } + Err(error) => return Err(error), + }; + Ok(entries + .into_iter() + .filter(|name| crate::hex::parse_hex::<16>(name).is_some()) + .collect()) +} diff --git a/src/recovery/journal.rs b/src/recovery/journal.rs index 24d8750..88fdec6 100644 --- a/src/recovery/journal.rs +++ b/src/recovery/journal.rs @@ -310,7 +310,7 @@ pub async fn replay_apply_journal( #[cfg(test)] pub async fn execute_journal_actions(vfs: &V, actions: &[JournalAction]) -> Result<()> { vfs.mkdir_all("seg").await?; - vfs.mkdir_all("seg/.staging").await?; + vfs.mkdir_all(crate::segment::writer::STAGING_DIR).await?; vfs.mkdir_all("seg/.tombstone").await?; for action in actions { match action { @@ -338,7 +338,7 @@ pub async fn execute_journal_actions(vfs: &V, actions: &[JournalAction]) } } vfs.sync_dir("seg").await?; - vfs.sync_dir("seg/.staging").await?; + vfs.sync_dir(crate::segment::writer::STAGING_DIR).await?; vfs.sync_dir("seg/.tombstone").await } diff --git a/src/recovery/mod.rs b/src/recovery/mod.rs index d50ff78..54a0a10 100644 --- a/src/recovery/mod.rs +++ b/src/recovery/mod.rs @@ -2,6 +2,7 @@ //! tombstone GC, spill-scratch reclamation. pub(crate) mod deep_walk; +pub(crate) mod fork; pub(crate) mod gc; pub(crate) mod journal; pub(crate) mod provenance; diff --git a/src/recovery/quarantine/preserve.rs b/src/recovery/quarantine/preserve.rs index a5b5abd..4a9c6a9 100644 --- a/src/recovery/quarantine/preserve.rs +++ b/src/recovery/quarantine/preserve.rs @@ -18,6 +18,7 @@ use crate::Result; use crate::errors::PagedbError; +use crate::segment::writer::STAGING_DIR; use crate::vfs::Vfs; /// Where a quarantined store went, and what moved. @@ -43,13 +44,7 @@ pub struct QuarantineReport { /// So the layout this crate writes is enumerated rather than discovered. /// Anything a backend does report is still moved; this only ensures the /// directories pagedb itself creates are never missed. -const STORE_SUBDIRS: &[&str] = &[ - "seg", - "seg/.staging", - "seg/.tombstone", - "applyjournal", - "tmp", -]; +const STORE_SUBDIRS: &[&str] = &["seg", STAGING_DIR, "seg/.tombstone", "applyjournal", "tmp"]; /// Move every file of the store rooted at `store_dir` into /// `/quarantine/