Skip to main content

kernel/
io.rs

1//! All operating-system contact lives here. Everything above this module is
2//! arithmetic over byte arrays and is platform-neutral by construction.
3//!
4//! SACRIFICE (Law 4): unbuffered I/O requires offset, length and buffer address
5//! to be page-aligned, and forgoes kernel readahead. Bought: the kernel does not
6//! keep a second copy of every page, so the cgroup charges only memory we chose.
7
8use crate::page::PAGE_SIZE;
9use crate::{Error, Result};
10use std::fs::{File, OpenOptions};
11use std::path::Path;
12
13#[cfg(unix)]
14type FileIdentity = (u64, u64);
15#[cfg(not(unix))]
16type FileIdentity = ();
17
18#[cfg(unix)]
19fn file_identity(file: &File) -> std::io::Result<FileIdentity> {
20    use std::os::unix::fs::MetadataExt;
21    let metadata = file.metadata()?;
22    Ok((metadata.dev(), metadata.ino()))
23}
24#[cfg(not(unix))]
25fn file_identity(_file: &File) -> std::io::Result<FileIdentity> { Ok(()) }
26
27fn writer_paths() -> &'static std::sync::Mutex<std::collections::HashMap<std::path::PathBuf, FileIdentity>> {
28    static PATHS: std::sync::OnceLock<
29        std::sync::Mutex<std::collections::HashMap<std::path::PathBuf, FileIdentity>>,
30    > = std::sync::OnceLock::new();
31    PATHS.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
32}
33
34pub(crate) fn writer_owned_by_this_process(path: &Path) -> bool {
35    let Ok(path) = std::fs::canonicalize(path) else { return false };
36    let paths = writer_paths().lock().unwrap();
37    let Some(owned) = paths.get(&path) else { return false };
38    let Ok(file) = File::open(&path) else { return false };
39    file_identity(&file).is_ok_and(|current| current == *owned)
40}
41
42#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43pub enum IoMode { Direct, Buffered }
44
45/// Which durability barrier a caller wants. Named rather than boolean,
46/// because `true` does not say what was promised.
47///
48/// There is deliberately no generic `sync()` on `FileIo`. One existed
49/// briefly and called `File::sync_data`, which issues `F_FULLFSYNC` on
50/// macOS and `fdatasync` on Linux -- the same call meaning two different
51/// promises depending on the platform, which is exactly the ambiguity
52/// `sync_data`/`sync_full` exist to remove. Leaving a `sync()` beside them
53/// preserves the bug in the one place nobody looks: the buffer pool's own
54/// checkpoint flush went through it, ungoverned by any `SyncMode`, even
55/// after `Store::commit` stopped being ambiguous.
56#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57pub enum Barrier { None, Data, Full }
58
59/// One file's I/O, counted. Per FILE: a global cursor interleaves the data file
60/// and the WAL, counting every alternation as a seek and summing both files'
61/// bytes as one -- both measured wrong before this existed.
62#[derive(Debug, Default)]
63pub struct IoStats {
64    pub writes: std::sync::atomic::AtomicU64,
65    pub write_bytes: std::sync::atomic::AtomicU64,
66    pub reads: std::sync::atomic::AtomicU64,
67}
68
69impl IoStats {
70    /// Snapshot and reset.
71    pub fn take(&self) -> (u64, u64, u64) {
72        use std::sync::atomic::Ordering::Relaxed;
73        (self.writes.swap(0, Relaxed), self.write_bytes.swap(0, Relaxed), self.reads.swap(0, Relaxed))
74    }
75}
76
77pub trait FileIo: Send + Sync {
78    /// Experimental transactional allocator; ordinary files retain pool reuse.
79    fn manages_free_pages(&self) -> bool { false }
80    fn pop_free_page(&self) -> Result<Option<u32>> { Ok(None) }
81    fn push_free_page(&self, _page: u32) -> Result<()> { unreachable!() }
82    /// This file's counters, if it keeps any. Defaulted so test doubles need no change.
83    fn stats(&self) -> Option<&IoStats> { None }
84
85    /// True when this file demands page-aligned offsets, lengths and buffers.
86    /// The WAL is byte-addressed and always opens Buffered, so it never does.
87    fn requires_alignment(&self) -> bool;
88    fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()>;
89    fn write_at(&self, buf: &[u8], off: u64) -> Result<()>;
90    /// A data barrier that does NOT force the drive's own write cache: the
91    /// write is durable against an OS crash, not necessarily against a loss
92    /// of power to the drive. `fdatasync` on Linux; plain `fsync` on macOS.
93    ///
94    /// Deliberately not `std::fs::File::sync_data`: libstd's implementation
95    /// calls `fcntl(F_FULLFSYNC)` on macOS (a safety choice in std, not a
96    /// bug), which is exactly the strong, ~65x-costlier barrier `sync_full`
97    /// exists to name separately. Going through it here would make `Normal`
98    /// and `Full` issue the identical primitive on macOS while differing on
99    /// Linux -- the SyncMode label would say one thing and the hardware would
100    /// hear another, and differently on different platforms.
101    fn sync_data(&self) -> Result<()>;
102    /// The strongest barrier this platform can issue: durable even against a
103    /// loss of power to the drive. `fcntl(F_FULLFSYNC)` on macOS (roughly 65x
104    /// the cost of `sync_data` on the same hardware); `File::sync_all`
105    /// (ordinary `fsync`) elsewhere.
106    fn sync_full(&self) -> Result<()>;
107    /// The exact primitive `sync_full` issues on this platform, so a
108    /// measurement can state what it did rather than imply it.
109    fn sync_full_primitive(&self) -> &'static str;
110    fn sync_dir(&self) -> Result<()>;
111    fn len(&self) -> Result<u64>;
112    fn set_len(&self, n: u64) -> Result<()>;
113}
114
115struct PosixFile {
116    f: File,
117    #[cfg(unix)]
118    dir: File,
119    mode: IoMode,
120    stats: IoStats,
121    /// Present only on the fd that owns this process's writer lock.
122    writer_path: Option<std::path::PathBuf>,
123    writer_identity: Option<FileIdentity>,
124}
125
126impl Drop for PosixFile {
127    fn drop(&mut self) {
128        if let (Some(path), Some(identity)) =
129            (self.writer_path.take(), self.writer_identity.take())
130        {
131            // This descriptor holds the writer's OS lock; release it by
132            // unlock, not by the close that follows (see `Locked`).
133            let _ = self.f.unlock();
134            let mut paths = writer_paths().lock().unwrap();
135            if paths.get(&path) == Some(&identity) { paths.remove(&path); }
136        }
137    }
138}
139
140impl FileIo for PosixFile {
141    fn stats(&self) -> Option<&IoStats> { Some(&self.stats) }
142
143    fn requires_alignment(&self) -> bool { self.mode == IoMode::Direct }
144    fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> {
145        self.stats.reads.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
146        if self.requires_alignment() {
147            debug_assert_eq!(buf.len() % PAGE_SIZE, 0, "unaligned length");
148            debug_assert_eq!(off as usize % PAGE_SIZE, 0, "unaligned offset");
149        }
150        #[cfg(unix)] {
151            use std::os::unix::fs::FileExt;
152            self.f.read_exact_at(buf, off)?;
153        }
154        #[cfg(windows)] {
155            use std::os::windows::fs::FileExt;
156            let mut done = 0usize;
157            while done < buf.len() {
158                let n = self.f.seek_read(&mut buf[done..], off + done as u64)?;
159                if n == 0 { return Err(std::io::Error::new(std::io::ErrorKind::UnexpectedEof, "seek_read hit EOF").into()); }
160                done += n;
161            }
162        }
163        Ok(())
164    }
165    fn write_at(&self, buf: &[u8], off: u64) -> Result<()> {
166        use std::sync::atomic::Ordering::Relaxed;
167        self.stats.writes.fetch_add(1, Relaxed);
168        self.stats.write_bytes.fetch_add(buf.len() as u64, Relaxed);
169        if self.requires_alignment() {
170            debug_assert_eq!(buf.len() % PAGE_SIZE, 0, "unaligned length");
171            debug_assert_eq!(off as usize % PAGE_SIZE, 0, "unaligned offset");
172        }
173        #[cfg(unix)] {
174            use std::os::unix::fs::FileExt;
175            self.f.write_all_at(buf, off)?;
176        }
177        #[cfg(windows)] {
178            // DuckDB/SQLite shape: positional WriteFile via OVERLAPPED --
179            // Rust's seek_write. Loop: seek_write may write short.
180            use std::os::windows::fs::FileExt;
181            let mut done = 0usize;
182            while done < buf.len() {
183                let n = self.f.seek_write(&buf[done..], off + done as u64)?;
184                if n == 0 { return Err(std::io::Error::new(std::io::ErrorKind::WriteZero, "seek_write wrote 0").into()); }
185                done += n;
186            }
187        }
188        Ok(())
189    }
190    fn sync_data(&self) -> Result<()> { sync_data_raw(&self.f) }
191    fn sync_full(&self) -> Result<()> { sync_full_raw(&self.f) }
192    fn sync_full_primitive(&self) -> &'static str { sync_full_primitive_name() }
193    fn sync_dir(&self) -> Result<()> {
194        // Unix: fsync the directory so the file's EXISTENCE is durable (the
195        // Law 3 checkpoint fix). Windows: opening a directory as a File is not
196        // a thing, and NTFS journals metadata -- SQLite's os_win.c fsyncs no
197        // directories either. No-op there, by design and stated.
198        #[cfg(unix)] { self.dir.sync_all()?; }
199        Ok(())
200    }
201    fn len(&self) -> Result<u64> { Ok(self.f.metadata()?.len()) }
202    fn set_len(&self, n: u64) -> Result<()> { self.f.set_len(n)?; Ok(()) }
203}
204
205#[cfg(target_os = "linux")]
206fn sync_data_raw(f: &File) -> Result<()> {
207    use std::os::unix::io::AsRawFd;
208    // SAFETY: fd is valid for the lifetime of `f`; fdatasync takes one int arg.
209    let rc = unsafe { libc::fdatasync(f.as_raw_fd()) };
210    if rc == -1 { return Err(std::io::Error::last_os_error().into()); }
211    Ok(())
212}
213
214#[cfg(target_os = "macos")]
215fn sync_data_raw(f: &File) -> Result<()> {
216    use std::os::unix::io::AsRawFd;
217    // Plain fsync -- NOT std::fs::File::sync_data, which libstd routes
218    // through fcntl(F_FULLFSYNC) on this platform. See the trait doc comment.
219    // SAFETY: fd is valid for the lifetime of `f`; fsync takes one int arg.
220    let rc = unsafe { libc::fsync(f.as_raw_fd()) };
221    if rc == -1 { return Err(std::io::Error::last_os_error().into()); }
222    Ok(())
223}
224
225// Windows: FlushFileBuffers for BOTH sync levels -- what SQLite's winSync and
226// DuckDB's FileSync call. There is no fdatasync distinction to honour, so
227// Normal and Full collapse to the same (strong) barrier: conservative, stated.
228// Android and other unixes: sync_all -- also conservative, also stated.
229#[cfg(not(any(target_os = "linux", target_os = "macos")))]
230fn sync_data_raw(f: &File) -> Result<()> { f.sync_all()?; Ok(()) }
231
232#[cfg(target_os = "macos")]
233fn sync_full_raw(f: &File) -> Result<()> {
234    use std::os::unix::io::AsRawFd;
235    // SAFETY: fd is valid for the lifetime of `f`; F_FULLFSYNC takes no argument.
236    let rc = unsafe { libc::fcntl(f.as_raw_fd(), libc::F_FULLFSYNC) };
237    if rc == -1 { return Err(std::io::Error::last_os_error().into()); }
238    Ok(())
239}
240
241#[cfg(not(target_os = "macos"))]
242fn sync_full_raw(f: &File) -> Result<()> { f.sync_all()?; Ok(()) }
243
244fn sync_full_primitive_name() -> &'static str {
245    #[cfg(target_os = "macos")] { "F_FULLFSYNC" }
246    #[cfg(not(target_os = "macos"))] { "fsync (sync_all)" }
247}
248
249/// Open `path`, asking for `want`. Returns the mode ACTUALLY obtained.
250///
251/// A platform that cannot give unbuffered I/O gets buffered I/O and says so.
252/// It must never silently behave like a different engine — that failure mode
253/// (an mmap path returning None on Windows, every index quietly falling back to
254/// resident) is exactly what this return value exists to prevent.
255pub fn open_file(path: &Path, want: IoMode) -> Result<(Box<dyn FileIo>, IoMode)> {
256    open_file_impl(path, want, false)
257}
258
259/// Open the database data file for its sole writer and retain an advisory
260/// whole-file lock for exactly as long as the returned `FileIo` owns the fd.
261/// Readers deliberately use `open_file_readonly` and never contend here.
262pub fn open_file_writer(path: &Path, want: IoMode) -> Result<(Box<dyn FileIo>, IoMode)> {
263    open_file_impl(path, want, true)
264}
265
266fn open_file_impl(path: &Path, want: IoMode, writer: bool) -> Result<(Box<dyn FileIo>, IoMode)> {
267    let parent = path.parent().unwrap_or(Path::new("."));
268    std::fs::create_dir_all(parent)?;
269    #[cfg(unix)]
270    let dir = File::open(parent)?;
271
272    // A FRESH builder per attempt. `custom_flags` mutates the builder in place
273    // and has no reset, so reusing one would carry O_DIRECT into the fallback
274    // and make the retry fail identically — turning "degrade and report" into a
275    // hard error on every filesystem that refuses unbuffered I/O. tmpfs and
276    // overlayfs both do, and this engine is measured inside a container.
277    let base = || {
278        let mut o = OpenOptions::new();
279        o.read(true).write(true).create(true);
280        o
281    };
282
283    let (f, got) = match want {
284        IoMode::Direct => match open_unbuffered(base(), path) {
285            Ok(f) => (f, IoMode::Direct),
286            Err(_) => (base().open(path)?, IoMode::Buffered),
287        },
288        IoMode::Buffered => (base().open(path)?, IoMode::Buffered),
289    };
290    let (writer_path, writer_identity) = if writer {
291        let path = std::fs::canonicalize(path)?;
292        let identity = file_identity(&f)?;
293        let mut paths = writer_paths().lock().unwrap();
294        if paths.get(&path) == Some(&identity) {
295            // A few crash tests deliberately `forget` a Store and reopen it
296            // in the same test process. A real process death closes the fd;
297            // forgetting Rust ownership cannot. Test-support builds allow
298            // that legacy simulation to reach recovery, while production
299            // refuses the second in-process writer here. Other processes do
300            // not share this registry and still contend on the OS lock below.
301            #[cfg(feature = "test-support")]
302            { (None, None) }
303            #[cfg(not(feature = "test-support"))]
304            { return Err(Error::WriterLocked); }
305        } else {
306            if !try_lock_exclusive(&f)? { return Err(Error::WriterLocked); }
307            paths.insert(path.clone(), identity);
308            (Some(path), Some(identity))
309        }
310    } else {
311        (None, None)
312    };
313    Ok((Box::new(PosixFile {
314        f,
315        #[cfg(unix)]
316        dir,
317        mode: got,
318        stats: IoStats::default(),
319        writer_path,
320        writer_identity,
321    }), got))
322}
323
324/// 2f: open for a SNAPSHOT READER -- read-only at the OS level, so the
325/// reader cannot write even by bug, and always Buffered (a reader shares
326/// the file with a live writer; O_DIRECT's alignment contract buys nothing
327/// on a cache the writer is also warming).
328pub fn open_file_readonly(path: &Path) -> Result<Box<dyn FileIo>> {
329    let parent = path.parent().unwrap_or(Path::new("."));
330    #[cfg(unix)]
331    let dir = File::open(parent)?;
332    let f = OpenOptions::new().read(true).open(path)?;
333    Ok(Box::new(PosixFile {
334        f,
335        #[cfg(unix)]
336        dir,
337        mode: IoMode::Buffered,
338        stats: IoStats::default(),
339        writer_path: None,
340        writer_identity: None,
341    }))
342}
343
344/// Recovery reads through an OS read-only descriptor while excluding writers.
345/// No create semantics: a missing source must stay missing.
346pub fn open_recovery_source(path: &Path) -> Result<Box<dyn FileIo>> {
347    let parent = path.parent().unwrap_or(Path::new("."));
348    #[cfg(unix)]
349    let dir = File::open(parent)?;
350    let f = OpenOptions::new().read(true).open(path)?;
351    let path = std::fs::canonicalize(path)?;
352    let identity = file_identity(&f)?;
353    let mut paths = writer_paths().lock().unwrap();
354    if paths.get(&path) == Some(&identity) || !try_lock_exclusive(&f)? {
355        return Err(Error::WriterLocked);
356    }
357    paths.insert(path.clone(), identity);
358    Ok(Box::new(PosixFile {
359        f,
360        #[cfg(unix)]
361        dir,
362        mode: IoMode::Buffered,
363        stats: IoStats::default(),
364        writer_path: Some(path),
365        writer_identity: Some(identity),
366    }))
367}
368
369/// Make a newly created directory's entry durable using the same portability
370/// contract as FileIo::sync_dir (Windows relies on the metadata journal).
371pub fn sync_directory(path: &Path) -> Result<()> {
372    #[cfg(unix)]
373    { File::open(path)?.sync_all()?; }
374    #[cfg(windows)]
375    { let _ = path; }
376    Ok(())
377}
378
379#[cfg(target_os = "linux")]
380fn open_unbuffered(mut opts: OpenOptions, path: &Path) -> std::io::Result<File> {
381    use std::os::unix::fs::OpenOptionsExt;
382    opts.custom_flags(libc::O_DIRECT).open(path)
383}
384
385#[cfg(target_os = "macos")]
386fn open_unbuffered(_opts: OpenOptions, _path: &Path) -> std::io::Result<File> {
387    // R1: on the tested external APFS volume, concurrent F_NOCACHE writes /
388    // reads returned unrelated bytes, also in a standalone C probe with aligned
389    // buffers (23/1600 mismatches; buffered control 0/1600). We cannot identify
390    // all affected OS/device combinations from an open() success. Until that
391    // boundary is proven, use the explicit Buffered fallback on macOS.
392    // Sacrifice: macOS Direct requests use the OS cache. P1/P2 already used
393    // Buffered. See docs/RECOVERY_R1.md in the E4 root for retained evidence.
394    Err(std::io::Error::new(std::io::ErrorKind::Unsupported,
395        "macOS uncached I/O disabled after failed data-isolation probe"))
396}
397
398// Windows upgrade path: FILE_FLAG_NO_BUFFERING via OpenOptionsExt (needs
399// sector-aligned buffers; AlignedRegion already is). Basic posture for now:
400// refuse, so open_file degrades to Buffered and REPORTS it -- same as any
401// filesystem that cannot do O_DIRECT.
402#[cfg(not(any(target_os = "linux", target_os = "macos")))]
403fn open_unbuffered(_opts: OpenOptions, _path: &Path) -> std::io::Result<File> {
404    Err(std::io::Error::new(std::io::ErrorKind::Unsupported, "no unbuffered mode"))
405}
406
407/// A page-aligned heap region. The buffer pool owns exactly one of these, which
408/// satisfies the alignment requirement of unbuffered I/O and the "allocate once,
409/// never free" requirement that keeps allocator retention out of the picture.
410pub struct AlignedRegion { ptr: *mut u8, len: usize }
411
412// SAFETY: the region is owned exclusively and never aliased across threads
413// except behind the pool's own synchronisation.
414unsafe impl Send for AlignedRegion {}
415unsafe impl Sync for AlignedRegion {}
416
417impl AlignedRegion {
418    pub fn new(len: usize) -> Result<Self> {
419        // A zero-size layout is undefined behaviour in alloc_zeroed, and
420        // `0 % PAGE_SIZE == 0` passes the alignment check, so it must be
421        // rejected explicitly rather than assumed away in a comment.
422        assert!(len > 0, "an AlignedRegion of zero bytes is a zero-size allocation");
423        assert_eq!(len % PAGE_SIZE, 0);
424        let layout = std::alloc::Layout::from_size_align(len, PAGE_SIZE).unwrap();
425        // SAFETY: layout has non-zero size and a valid power-of-two alignment.
426        let ptr = unsafe { std::alloc::alloc_zeroed(layout) };
427        if ptr.is_null() { return Err(Error::OutOfBudget); }
428        Ok(AlignedRegion { ptr, len })
429    }
430    fn frame_start(&self, i: usize) -> usize {
431        // checked, because in a release build `i + 1` on a pathological index
432        // wraps to 0 and defeats the bound entirely.
433        let start = i.checked_mul(PAGE_SIZE).expect("frame index overflow");
434        let end = start.checked_add(PAGE_SIZE).expect("frame index overflow");
435        assert!(end <= self.len, "frame {i} is outside the region");
436        start
437    }
438
439    /// # Safety
440    /// No `&mut` to frame `i` may be live for the lifetime of the returned
441    /// slice. The buffer pool guarantees this with its pin count; this type
442    /// cannot.
443    pub unsafe fn page(&self, i: usize) -> &[u8] {
444        let start = self.frame_start(i);
445        // SAFETY: bounds checked above; exclusivity is the caller's obligation.
446        unsafe { std::slice::from_raw_parts(self.ptr.add(start), PAGE_SIZE) }
447    }
448
449    /// # Safety
450    /// No other reference to frame `i` — shared or mutable — may be live for
451    /// the lifetime of the returned slice.
452    ///
453    /// This is deliberately an `unsafe fn`. As a safe fn, ordinary safe code
454    /// could call it twice and hold two `&mut [u8]` over the same bytes, which
455    /// is undefined behaviour however careful the pool is. A comment cannot
456    /// make an unenforceable invariant sound; moving the obligation to the
457    /// caller, in the type system, can.
458    #[allow(clippy::mut_from_ref)]
459    pub unsafe fn page_mut(&self, i: usize) -> &mut [u8] {
460        let start = self.frame_start(i);
461        // SAFETY: bounds checked above; exclusivity is the caller's obligation.
462        unsafe { std::slice::from_raw_parts_mut(self.ptr.add(start), PAGE_SIZE) }
463    }
464
465    /// A page-multiple prefix for one aligned multi-page I/O. The caller must
466    /// enforce the same aliasing discipline as [`Self::page`].
467    pub unsafe fn prefix(&self, len: usize) -> &[u8] {
468        assert!(len > 0 && len <= self.len && len % PAGE_SIZE == 0);
469        unsafe { std::slice::from_raw_parts(self.ptr, len) }
470    }
471}
472
473impl Drop for AlignedRegion {
474    fn drop(&mut self) {
475        let layout = std::alloc::Layout::from_size_align(self.len, PAGE_SIZE).unwrap();
476        // SAFETY: ptr came from alloc_zeroed with this exact layout.
477        unsafe { std::alloc::dealloc(self.ptr, layout) }
478    }
479}
480
481#[cfg(test)]
482mod tests {
483    use super::*;
484    use crate::page::PAGE_SIZE;
485
486    // A Direct request can succeed on Linux. Ordinary Vec allocations do not
487    // satisfy its buffer-address alignment contract, even at page-sized length.
488    #[repr(align(4096))]
489    struct TestPage([u8; PAGE_SIZE]);
490
491    #[test]
492    fn pages_round_trip_through_the_file() {
493        let dir = tempfile::tempdir().unwrap();
494        let path = dir.path().join("t.db");
495        let (f, _) = open_file(&path, IoMode::Buffered).unwrap();
496
497        let mut w = vec![0u8; PAGE_SIZE];
498        for (i, b) in w.iter_mut().enumerate() { *b = (i % 251) as u8; }
499        f.write_at(&w, (PAGE_SIZE * 3) as u64).unwrap();
500        f.sync_data().unwrap();
501
502        let mut r = vec![0u8; PAGE_SIZE];
503        f.read_at(&mut r, (PAGE_SIZE * 3) as u64).unwrap();
504        assert_eq!(r, w);
505        assert_eq!(f.len().unwrap(), (PAGE_SIZE * 4) as u64);
506    }
507
508    #[test]
509    fn a_read_past_the_end_is_an_error_not_a_short_buffer() {
510        let dir = tempfile::tempdir().unwrap();
511        let (f, _) = open_file(&dir.path().join("t.db"), IoMode::Buffered).unwrap();
512        let mut r = vec![0u8; PAGE_SIZE];
513        assert!(f.read_at(&mut r, 0).is_err());
514    }
515
516    /// Requesting Direct must NEVER fail merely because Direct is unavailable.
517    /// It must degrade to Buffered and say so. Filesystems that refuse
518    /// unbuffered I/O (tmpfs, overlayfs, many container mounts) are ordinary,
519    /// and this engine is measured inside a container.
520    ///
521    /// `assert!(matches!(got, Direct | Buffered))` would be tautological — the
522    /// enum has exactly those variants — so it asserts nothing. This asserts
523    /// the two things that can actually be wrong: that the open succeeded, and
524    /// that `requires_alignment` follows the mode obtained rather than the one
525    /// requested.
526    #[test]
527    fn requesting_direct_degrades_rather_than_failing() {
528        let dir = tempfile::tempdir().unwrap();
529        let (f, got) = open_file(&dir.path().join("t.db"), IoMode::Direct)
530            .expect("requesting Direct must degrade, never error");
531        assert_eq!(f.requires_alignment(), got == IoMode::Direct);
532        #[cfg(target_os = "macos")]
533        assert_eq!(got, IoMode::Buffered, "unproven uncached mode must stay disabled");
534
535        // Usable either way.
536        let w = TestPage([0u8; PAGE_SIZE]);
537        f.write_at(&w.0, 0).unwrap();
538        f.sync_data().unwrap();
539        let mut r = TestPage([0u8; PAGE_SIZE]);
540        f.read_at(&mut r.0, 0).unwrap();
541        assert_eq!(r.0, w.0);
542    }
543
544    #[test]
545    fn concurrent_direct_and_buffered_files_keep_their_own_bytes() {
546        // R1 observed one direct read return another file's freelist header.
547        // Exercise both zero and distinctive payloads under directory churn;
548        // retain actual bytes and the path if the symptom recurs.
549        std::thread::scope(|scope| {
550            for worker in 0..8u8 {
551                scope.spawn(move || {
552                    for round in 0..32u8 {
553                        let dir = tempfile::tempdir().unwrap();
554                        let mode = if worker % 2 == 0 { IoMode::Direct } else { IoMode::Buffered };
555                        let (f, _) = open_file(&dir.path().join("roundtrip"), mode).unwrap();
556                        let mut w = TestPage([0u8; PAGE_SIZE]);
557                        if round % 2 != 0 {
558                            for (i, byte) in w.0.iter_mut().enumerate() { *byte = worker.wrapping_add(round).wrapping_add(i as u8); }
559                        }
560                        f.write_at(&w.0, 0).unwrap();
561                        f.sync_data().unwrap();
562                        let mut r = TestPage([0u8; PAGE_SIZE]);
563                        f.read_at(&mut r.0, 0).unwrap();
564                        if r.0 != w.0 {
565                            let retained = dir.keep();
566                            std::fs::write(retained.join("expected"), &w.0).unwrap();
567                            std::fs::write(retained.join("observed"), &r.0).unwrap();
568                            panic!("I/O isolation failed: worker {worker}, round {round}, evidence {}", retained.display());
569                        }
570                    }
571                });
572            }
573        });
574    }
575}
576
577/// Advisory whole-file lock, non-blocking (2n reader table). Ok(true) =
578/// acquired; Ok(false) = held by a live process. The lock dies with the fd
579/// -- crash-safe by construction. Platform code lives HERE (the 2d rule).
580pub fn try_lock_exclusive(f: &std::fs::File) -> std::io::Result<bool> {
581    // std's file lock rather than raw flock: same contract on Unix, and it is
582    // the only form that EXISTS on Windows (LockFileEx underneath). The hand-
583    // rolled flock made the whole kernel un-compilable off Unix, which the
584    // reader table cannot afford -- without it the writer must assume a reader
585    // at generation 0 and page recycling stops for good.
586    match f.try_lock() {
587        Ok(()) => Ok(true),
588        Err(std::fs::TryLockError::WouldBlock) => Ok(false),
589        Err(std::fs::TryLockError::Error(e)) => Err(e),
590    }
591}
592
593/// Shared counterpart of `try_lock_exclusive`: Ok(false) when an exclusive
594/// holder is alive. Read-only handles suffice on every supported platform.
595/// Used by the page-WAL reader admission (ownership probe, admission gate).
596pub fn try_lock_shared(f: &std::fs::File) -> std::io::Result<bool> {
597    match f.try_lock_shared() {
598        Ok(()) => Ok(true),
599        Err(std::fs::TryLockError::WouldBlock) => Ok(false),
600        Err(std::fs::TryLockError::Error(e)) => Err(e),
601    }
602}
603
604/// Blocking shared lock: waits only for an exclusive holder's critical
605/// section (the page-WAL writer's tail truncation or checkpoint).
606pub fn lock_shared(f: &std::fs::File) -> std::io::Result<()> { f.lock_shared() }
607
608/// Blocking exclusive lock: waits only for admissions in flight.
609pub fn lock_exclusive(f: &std::fs::File) -> std::io::Result<()> { f.lock() }
610
611/// A file whose OS lock is RELEASED BY AN EXPLICIT UNLOCK when this is
612/// dropped, not by closing the descriptor.
613///
614/// An OS file lock belongs to the open file, and a child process started by
615/// ANY thread of this process inherits a copy of every descriptor for the
616/// instant between fork and exec. Closing our descriptor leaves the lock held
617/// by that copy until the child execs, so a database closed and reopened
618/// while anything spawns a subprocess could be refused as "already has an
619/// active writer", and a checkpoint could find a reader slot "held". Measured
620/// with the standard library alone -- one thread locking, closing and
621/// relocking a file 20,000 times while another starts child processes -- 45
622/// to 116 relocks per run were refused when the lock was released by close,
623/// and 0 when it was released by unlock first: an unlock acts on the open
624/// file itself, which the inherited copy shares. So every lock this engine
625/// holds past the statement that took it is held through this type.
626pub struct Locked(std::fs::File);
627
628impl Locked {
629    /// Hold `f`, which the caller has just locked (shared or exclusive).
630    pub fn held(f: std::fs::File) -> Self {
631        Self(f)
632    }
633}
634
635impl std::ops::Deref for Locked {
636    type Target = std::fs::File;
637    fn deref(&self) -> &std::fs::File {
638        &self.0
639    }
640}
641
642impl Drop for Locked {
643    fn drop(&mut self) {
644        let _ = self.0.unlock();
645    }
646}