1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57pub enum Barrier { None, Data, Full }
58
59#[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 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 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 fn stats(&self) -> Option<&IoStats> { None }
84
85 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 fn sync_data(&self) -> Result<()>;
102 fn sync_full(&self) -> Result<()>;
107 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 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 let _ = unlock(&self.f);
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 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 #[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 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 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#[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 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
249pub fn open_file(path: &Path, want: IoMode) -> Result<(Box<dyn FileIo>, IoMode)> {
256 open_file_impl(path, want, false)
257}
258
259pub 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 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 #[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
324pub 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
344pub 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
369pub 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 Err(std::io::Error::new(std::io::ErrorKind::Unsupported,
395 "macOS uncached I/O disabled after failed data-isolation probe"))
396}
397
398#[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
407pub struct AlignedRegion { ptr: *mut u8, len: usize }
411
412unsafe impl Send for AlignedRegion {}
415unsafe impl Sync for AlignedRegion {}
416
417impl AlignedRegion {
418 pub fn new(len: usize) -> Result<Self> {
419 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 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 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 pub unsafe fn page(&self, i: usize) -> &[u8] {
444 let start = self.frame_start(i);
445 unsafe { std::slice::from_raw_parts(self.ptr.add(start), PAGE_SIZE) }
447 }
448
449 #[allow(clippy::mut_from_ref)]
459 pub unsafe fn page_mut(&self, i: usize) -> &mut [u8] {
460 let start = self.frame_start(i);
461 unsafe { std::slice::from_raw_parts_mut(self.ptr.add(start), PAGE_SIZE) }
463 }
464
465 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 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 #[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 #[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 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 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
577pub fn try_lock_exclusive(f: &std::fs::File) -> std::io::Result<bool> {
581 os_lock::lock(f, true, false)
582}
583
584pub fn try_lock_shared(f: &std::fs::File) -> std::io::Result<bool> {
587 os_lock::lock(f, false, false)
588}
589
590pub fn lock_shared(f: &std::fs::File) -> std::io::Result<()> {
593 os_lock::lock(f, false, true).map(|_| ())
594}
595
596pub fn lock_exclusive(f: &std::fs::File) -> std::io::Result<()> {
598 os_lock::lock(f, true, true).map(|_| ())
599}
600
601pub fn unlock(f: &std::fs::File) -> std::io::Result<()> {
605 os_lock::unlock(f)
606}
607
608#[cfg(not(windows))]
625mod os_lock {
626 pub fn lock(f: &std::fs::File, exclusive: bool, wait: bool) -> std::io::Result<bool> {
627 let outcome = match (exclusive, wait) {
628 (true, true) => return f.lock().map(|_| true),
629 (false, true) => return f.lock_shared().map(|_| true),
630 (true, false) => f.try_lock(),
631 (false, false) => f.try_lock_shared(),
632 };
633 match outcome {
634 Ok(()) => Ok(true),
635 Err(std::fs::TryLockError::WouldBlock) => Ok(false),
636 Err(std::fs::TryLockError::Error(e)) => Err(e),
637 }
638 }
639 pub fn unlock(f: &std::fs::File) -> std::io::Result<()> {
640 f.unlock()
641 }
642}
643
644#[cfg(windows)]
645mod os_lock {
646 use std::os::windows::io::AsRawHandle;
647 use windows_sys::Win32::Foundation::ERROR_LOCK_VIOLATION;
648 use windows_sys::Win32::Storage::FileSystem::{
649 LockFileEx, UnlockFileEx, LOCKFILE_EXCLUSIVE_LOCK, LOCKFILE_FAIL_IMMEDIATELY,
650 };
651 use windows_sys::Win32::System::IO::OVERLAPPED;
652
653 const LOCK_BYTE: u64 = 0x7fff_ffff_0000_0000;
656
657 fn region() -> OVERLAPPED {
658 let mut o: OVERLAPPED = unsafe { std::mem::zeroed() };
660 o.Anonymous.Anonymous.Offset = LOCK_BYTE as u32;
661 o.Anonymous.Anonymous.OffsetHigh = (LOCK_BYTE >> 32) as u32;
662 o
663 }
664
665 pub fn lock(f: &std::fs::File, exclusive: bool, wait: bool) -> std::io::Result<bool> {
666 let mut flags = 0;
667 if exclusive {
668 flags |= LOCKFILE_EXCLUSIVE_LOCK;
669 }
670 if !wait {
671 flags |= LOCKFILE_FAIL_IMMEDIATELY;
672 }
673 let mut o = region();
674 let ok = unsafe { LockFileEx(f.as_raw_handle() as _, flags, 0, 1, 0, &mut o) };
677 if ok != 0 {
678 return Ok(true);
679 }
680 let error = std::io::Error::last_os_error();
681 if !wait && error.raw_os_error() == Some(ERROR_LOCK_VIOLATION as i32) {
682 return Ok(false);
683 }
684 Err(error)
685 }
686
687 pub fn unlock(f: &std::fs::File) -> std::io::Result<()> {
688 let mut o = region();
689 let ok = unsafe { UnlockFileEx(f.as_raw_handle() as _, 0, 1, 0, &mut o) };
691 if ok != 0 {
692 Ok(())
693 } else {
694 Err(std::io::Error::last_os_error())
695 }
696 }
697}
698
699pub struct Locked(std::fs::File);
715
716impl Locked {
717 pub fn held(f: std::fs::File) -> Self {
719 Self(f)
720 }
721}
722
723impl std::ops::Deref for Locked {
724 type Target = std::fs::File;
725 fn deref(&self) -> &std::fs::File {
726 &self.0
727 }
728}
729
730impl Drop for Locked {
731 fn drop(&mut self) {
732 let _ = unlock(&self.0);
733 }
734}