From 0b5ee5abe9924d64a1c9b0200f1da91e9bf218cc Mon Sep 17 00:00:00 2001 From: Andrew Briscoe Date: Sun, 26 Jul 2026 12:36:04 -0600 Subject: [PATCH 1/2] fix(vfs): harden completion and path boundaries --- src/vfs/gcd/file.rs | 88 ++++-- src/vfs/gcd/vfs.rs | 25 +- src/vfs/iocp/file.rs | 166 ++++++---- src/vfs/iocp/vfs.rs | 26 +- src/vfs/iouring/file.rs | 127 ++++++-- src/vfs/iouring/vfs.rs | 25 +- src/vfs/memory.rs | 46 ++- src/vfs/opfs/handle.rs | 29 +- src/vfs/opfs/opfs_worker.js | 9 +- src/vfs/opfs/protocol.rs | 10 +- src/vfs/tokio_backend.rs | 117 +++++-- src/vfs/traits.rs | 592 ++++++++++++++++++++++++++++++++++-- tests/vfs_gcd.rs | 210 +++++++++++++ tests/vfs_iocp.rs | 117 +++++++ tests/vfs_iouring.rs | 135 ++++++++ tests/vfs_memory.rs | 25 ++ tests/vfs_tokio.rs | 127 +++++++- 17 files changed, 1662 insertions(+), 212 deletions(-) create mode 100644 tests/vfs_gcd.rs create mode 100644 tests/vfs_iocp.rs diff --git a/src/vfs/gcd/file.rs b/src/vfs/gcd/file.rs index 78ca888..428855e 100644 --- a/src/vfs/gcd/file.rs +++ b/src/vfs/gcd/file.rs @@ -1,7 +1,7 @@ //! `GcdFile`: per-file I/O against a `DispatchIO` channel. //! //! The channel is created in `DISPATCH_IO_RANDOM` mode and owns a duplicated -//! file descriptor (so dispatch_io can asynchronously close its own copy when +//! file descriptor (so `dispatch_io` can asynchronously close its own copy when //! the channel is released, independently of the `std::fs::File` we hold). //! Reads and writes call `dispatch_io_read` / `dispatch_io_write` with block //! handlers; the handlers accumulate chunks and, on `done`, send the result @@ -14,12 +14,12 @@ use std::sync::Arc; use parking_lot::Mutex; -use block2::{Block, DynBlock, RcBlock}; +use block2::{DynBlock, RcBlock}; use dispatch2::{DispatchData, DispatchIO, DispatchIOCloseFlags, DispatchQueue, DispatchRetained}; use crate::Result; use crate::errors::PagedbError; -use crate::vfs::traits::VfsFile; +use crate::vfs::traits::{VfsFile, checked_signed_file_len}; use crate::vfs::types::{ReadReq, WriteReq}; pub struct GcdFile { @@ -81,6 +81,37 @@ impl GcdFile { fn fd(&self) -> std::os::unix::io::RawFd { self.file.as_raw_fd() } + + fn checked_dispatch_offset(offset: u64, len: usize) -> Result { + let last = if len > 0 { + let last_delta = u64::try_from(len - 1).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "buffer length does not fit in u64", + )) + })?; + offset.checked_add(last_delta).ok_or_else(|| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "dispatch_io offset range overflow", + )) + })? + } else { + offset + }; + libc::off_t::try_from(last).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "dispatch_io offset range does not fit into libc::off_t", + )) + })?; + offset.try_into().map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "dispatch_io offset does not fit into libc::off_t", + )) + }) + } } impl Drop for GcdFile { @@ -100,7 +131,7 @@ impl Drop for GcdFile { fn submit_read( channel: &DispatchIO, queue: &DispatchQueue, - offset: u64, + offset: libc::off_t, len: usize, tx: tokio::sync::oneshot::Sender>>, ) { @@ -135,22 +166,15 @@ fn submit_read( } }); - // SAFETY: transmute is from the stand-in block signature to the typedef - // declared by dispatch2 — the ABI is identical because `bool` and `u8` - // share the same one-byte ABI, and a `*mut DispatchData` is bit-identical - // to a `*mut c_void`. - let handler_ptr: *mut DynBlock = unsafe { - std::mem::transmute::< - *mut Block, - *mut DynBlock, - >(RcBlock::as_ptr(&handler)) - }; + // SAFETY: ABI-compatible pointer cast to the typedef declared by dispatch2. + // `bool` and `u8` are one byte, while both data arguments are pointers. + let handler_ptr: *mut DynBlock = + RcBlock::as_ptr(&handler).cast::>(); // SAFETY: channel and queue are valid (owned by the caller's `GcdFile`); // the handler block is retained by libdispatch on submission. unsafe { - #[allow(clippy::cast_possible_wrap)] - channel.read(offset as libc::off_t, len, queue, handler_ptr); + channel.read(offset, len, queue, handler_ptr); } // libdispatch has retained the block internally; we can drop our Rc. drop(handler); @@ -162,7 +186,7 @@ fn submit_read( fn submit_write( channel: &DispatchIO, queue: &DispatchQueue, - offset: u64, + offset: libc::off_t, buf: &[u8], tx: tokio::sync::oneshot::Sender>, ) { @@ -185,19 +209,14 @@ fn submit_write( }, ); - // SAFETY: ABI-compatible transmute; see `submit_read`. - let handler_ptr: *mut DynBlock = unsafe { - std::mem::transmute::< - *mut Block, - *mut DynBlock, - >(RcBlock::as_ptr(&handler)) - }; + // SAFETY: ABI-compatible pointer cast; see `submit_read`. + let handler_ptr: *mut DynBlock = + RcBlock::as_ptr(&handler).cast::>(); // SAFETY: channel/queue/data all valid; libdispatch retains both the data // object and the handler block for the operation's duration. unsafe { - #[allow(clippy::cast_possible_wrap)] - channel.write(offset as libc::off_t, &data, queue, handler_ptr); + channel.write(offset, &data, queue, handler_ptr); } drop(handler); drop(data); @@ -209,6 +228,7 @@ impl VfsFile for GcdFile { return Ok(0); } let len = buf.len(); + let offset = Self::checked_dispatch_offset(offset, len)?; let (tx, rx) = tokio::sync::oneshot::channel::>>(); // The handler block and its raw block pointer are `!Send`; submitting // from a synchronous helper keeps them out of this future's state @@ -227,6 +247,10 @@ impl VfsFile for GcdFile { } async fn read_at_vectored(&self, reqs: &mut [ReadReq<'_>]) -> Result<()> { + for req in reqs.iter() { + Self::checked_dispatch_offset(req.offset, req.buf.len())?; + } + // dispatch_io operations are intrinsically sequential per channel // (the channel serialises ops in submission order); issuing them one // at a time matches that and keeps the bridge simple. @@ -247,6 +271,7 @@ impl VfsFile for GcdFile { return Ok(0); } let len = buf.len(); + let offset = Self::checked_dispatch_offset(offset, len)?; let (tx, rx) = tokio::sync::oneshot::channel::>(); // `!Send` handler/data confined to a synchronous helper; see `read_at`. submit_write(&self.channel, &self.queue, offset, buf, tx); @@ -261,6 +286,9 @@ impl VfsFile for GcdFile { if !self.writable { return Err(PagedbError::ReadOnly); } + for req in reqs { + Self::checked_dispatch_offset(req.offset, req.buf.len())?; + } for req in reqs { self.write_at(req.offset, req.buf).await?; } @@ -282,10 +310,10 @@ impl VfsFile for GcdFile { if !self.writable { return Err(PagedbError::ReadOnly); } - // SAFETY: `fd()` valid as above; `len` fits in `libc::off_t` on - // 64-bit Apple targets (`off_t` is `i64` on macOS / iOS). - #[allow(clippy::cast_possible_wrap)] - let rc = unsafe { libc::ftruncate(self.fd(), len as libc::off_t) }; + let len = checked_signed_file_len(len, "ftruncate")?; + // SAFETY: `fd()` is valid as above and `len` was checked to fit the + // signed native file-offset type. + let rc = unsafe { libc::ftruncate(self.fd(), len) }; if rc != 0 { return Err(PagedbError::Io(std::io::Error::last_os_error())); } diff --git a/src/vfs/gcd/vfs.rs b/src/vfs/gcd/vfs.rs index b93f1ce..16a5ae0 100644 --- a/src/vfs/gcd/vfs.rs +++ b/src/vfs/gcd/vfs.rs @@ -17,7 +17,7 @@ use crate::Result; use crate::errors::PagedbError; use super::file::GcdFile; -use crate::vfs::traits::Vfs; +use crate::vfs::traits::{Vfs, canonical_native_path, resolve_native_path}; use crate::vfs::types::OpenMode; #[derive(Debug, Clone, Copy)] @@ -128,8 +128,8 @@ impl GcdVfs { } } - fn resolve(&self, p: &str) -> PathBuf { - self.inner.root.join(p.trim_start_matches('/')) + fn resolve(&self, path: &str) -> Result { + resolve_native_path(&self.inner.root, path) } fn lookup_or_create_entry(&self, path: &str) -> Arc { @@ -145,7 +145,8 @@ impl GcdVfs { } fn do_lock(&self, path: &str, kind: LockKind) -> Result { - let entry = self.lookup_or_create_entry(path); + 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) { @@ -155,7 +156,7 @@ impl GcdVfs { _ => return Err(PagedbError::AlreadyLocked), } } - let lock_path = self.resolve(path); + let lock_path = self.resolve(&logical_path)?; if let Some(parent) = lock_path.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; } @@ -184,7 +185,7 @@ impl Vfs for GcdVfs { type LockHandle = GcdLockHandle; async fn open(&self, path: &str, mode: OpenMode) -> Result { - let p = self.resolve(path); + let p = self.resolve(path)?; if matches!(mode, OpenMode::CreateNew | OpenMode::CreateOrOpen) { if let Some(parent) = p.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; @@ -231,7 +232,7 @@ impl Vfs for GcdVfs { } async fn remove(&self, path: &str) -> Result<()> { - let p = self.resolve(path); + let p = self.resolve(path)?; match std::fs::remove_file(&p) { Ok(()) => Ok(()), Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()), @@ -240,8 +241,8 @@ impl Vfs for GcdVfs { } async fn rename(&self, from: &str, to: &str) -> Result<()> { - let f = self.resolve(from); - let t = self.resolve(to); + let f = self.resolve(from)?; + let t = self.resolve(to)?; if let Some(parent) = t.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; } @@ -249,7 +250,7 @@ impl Vfs for GcdVfs { } async fn list_dir(&self, path: &str) -> Result> { - let p = self.resolve(path); + let p = self.resolve(path)?; let iter = match std::fs::read_dir(&p) { Ok(it) => it, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()), @@ -267,13 +268,13 @@ impl Vfs for GcdVfs { } async fn mkdir_all(&self, path: &str) -> Result<()> { - let p = self.resolve(path); + let p = self.resolve(path)?; std::fs::create_dir_all(&p).map_err(PagedbError::Io) } async fn sync_dir(&self, path: &str) -> Result<()> { // POSIX fsync on the directory fd; HFS+/APFS honor it. - let p = self.resolve(path); + let p = self.resolve(path)?; let dir = match std::fs::File::open(&p) { Ok(d) => d, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()), diff --git a/src/vfs/iocp/file.rs b/src/vfs/iocp/file.rs index c29e0d5..a4883d4 100644 --- a/src/vfs/iocp/file.rs +++ b/src/vfs/iocp/file.rs @@ -9,7 +9,10 @@ use std::sync::Arc; use crate::Result; use crate::errors::PagedbError; -use crate::vfs::traits::VfsFile; +use crate::vfs::traits::{ + OverlappedStart, VfsFile, checked_overlapped_read_start, checked_overlapped_start, + checked_read_count, checked_readfile_len, checked_signed_file_len, checked_write_progress, +}; use crate::vfs::types::{ReadReq, WriteReq}; use super::port::{Port, PortInner}; @@ -51,6 +54,14 @@ impl IocpFile { self.file.as_raw_handle() as HANDLE } + fn writefile_len(len: usize) -> Result { + u32::try_from(len).map_err(|_| { + PagedbError::Io(std::io::Error::other( + "buffer too large for u32 in WriteFile", + )) + }) + } + /// Submit one overlapped read or write at `offset`. The closure performs /// the `ReadFile` / `WriteFile` syscall against `&mut OVERLAPPED`. Caller /// must hold the port mutex. @@ -61,6 +72,34 @@ impl IocpFile { /// must return `0` (failure) or non-zero (immediate completion); on /// `ERROR_IO_PENDING` we wait via `Port::dequeue`. unsafe fn submit_overlapped(port: &Port, offset: u64, op: F) -> std::io::Result + where + F: FnOnce(&mut OVERLAPPED) -> i32, + { + // SAFETY: caller forwards the same safety contract as + // `submit_overlapped_impl`. + unsafe { Self::submit_overlapped_impl(port, offset, false, op) } + } + + /// Submit one overlapped read at `offset`; EOF is reported as zero bytes. + /// + /// # Safety + /// + /// Same contract as [`Self::submit_overlapped`]. + unsafe fn submit_read_overlapped(port: &Port, offset: u64, op: F) -> std::io::Result + where + F: FnOnce(&mut OVERLAPPED) -> i32, + { + // SAFETY: caller forwards the same safety contract as + // `submit_overlapped_impl`. + unsafe { Self::submit_overlapped_impl(port, offset, true, op) } + } + + unsafe fn submit_overlapped_impl( + port: &Port, + offset: u64, + eof_is_empty_read: bool, + op: F, + ) -> std::io::Result where F: FnOnce(&mut OVERLAPPED) -> i32, { @@ -76,30 +115,46 @@ impl IocpFile { } let rc = op(&mut overlapped); - if rc != 0 { - // Immediate synchronous completion. `OVERLAPPED.InternalHigh` - // holds the byte count on success. - #[allow(clippy::cast_possible_truncation)] - return Ok(overlapped.InternalHigh as u32); + let err_code = if rc == 0 { + // SAFETY: GetLastError immediately after a failed Win32 call is + // the documented pattern. + unsafe { windows_sys::Win32::Foundation::GetLastError() } + } else { + 0 + }; + let start = if eof_is_empty_read { + checked_overlapped_read_start(rc, err_code, ERROR_IO_PENDING, ERROR_HANDLE_EOF) + } else { + checked_overlapped_start(rc, err_code, ERROR_IO_PENDING) } - // SAFETY: GetLastError immediately after a failed Win32 call is the - // documented pattern. - let err_code = unsafe { windows_sys::Win32::Foundation::GetLastError() }; - if err_code != ERROR_IO_PENDING { - return Err(std::io::Error::from_raw_os_error(err_code as i32)); + .map_err(|error| match error { + PagedbError::Io(io) => io, + other => std::io::Error::other(other.to_string()), + })?; + if start == OverlappedStart::EmptyRead { + return Ok(0); } - // I/O is in flight against the port — wait for its completion. + + // Associated handles queue a completion packet even when the + // overlapped call succeeds immediately, so every non-EOF start must + // drain its matching packet before the stack-allocated OVERLAPPED dies. // SAFETY: caller holds the port lock; `overlapped` lives on this // stack frame and is the only outstanding request. - let (bytes, _key, _ov) = unsafe { port.dequeue() }.or_else(|e| { + let (bytes, _key, completed) = unsafe { port.dequeue() }.or_else(|error| { // ERROR_HANDLE_EOF on read is "read past EOF" — surface as // 0 bytes transferred, like POSIX read() at EOF. - if e.raw_os_error() == Some(ERROR_HANDLE_EOF as i32) { + if error.raw_os_error() == Some(ERROR_HANDLE_EOF as i32) { Ok((0u32, 0usize, std::ptr::null_mut::())) } else { - Err(e) + Err(error) } })?; + let expected = (&raw mut overlapped).cast::(); + if !completed.is_null() && completed != expected { + return Err(std::io::Error::other( + "IOCP completion did not match the submitted OVERLAPPED", + )); + } Ok(bytes) } } @@ -116,11 +171,7 @@ impl VfsFile for IocpFile { return Ok(0); } let handle = self.handle(); - let len = u32::try_from(buf.len()).map_err(|_| { - PagedbError::Io(std::io::Error::other( - "buffer too large for u32 in ReadFile", - )) - })?; + let len = checked_readfile_len(buf.len())?; let port = Port { inner: Arc::clone(&self.port), }; @@ -129,13 +180,12 @@ impl VfsFile for IocpFile { // SAFETY: `buf` slice outlives this fn frame; the closure is invoked // synchronously before `submit_overlapped` returns. let bytes = unsafe { - IocpFile::submit_overlapped(&port, offset, |ov| { - let mut bytes_read: u32 = 0; - ReadFile(handle, buf_ptr.cast(), len, &mut bytes_read, ov) + IocpFile::submit_read_overlapped(&port, offset, |ov| { + ReadFile(handle, buf_ptr.cast(), len, std::ptr::null_mut(), ov) }) } .map_err(PagedbError::Io)?; - Ok(bytes as usize) + checked_read_count(bytes as usize, buf.len()) } async fn read_at_vectored(&self, reqs: &mut [ReadReq<'_>]) -> Result<()> { @@ -147,27 +197,26 @@ impl VfsFile for IocpFile { // Issue ops sequentially under one lock acquisition so the batch is // observed atomically with respect to other concurrent users of the // port. + let lengths: Vec = reqs + .iter() + .map(|req| checked_readfile_len(req.buf.len())) + .collect::>()?; let handle = self.handle(); let port = Port { inner: Arc::clone(&self.port), }; let _guard = port.lock(); - for req in reqs.iter_mut() { - let len = u32::try_from(req.buf.len()).map_err(|_| { - PagedbError::Io(std::io::Error::other( - "buffer too large for u32 in ReadFile", - )) - })?; + for (req, len) in reqs.iter_mut().zip(lengths) { let buf_ptr = req.buf.as_mut_ptr(); // SAFETY: `req.buf` outlives this frame; closure is synchronous // within `submit_overlapped`. let bytes = unsafe { - IocpFile::submit_overlapped(&port, req.offset, |ov| { - let mut bytes_read: u32 = 0; - ReadFile(handle, buf_ptr.cast(), len, &mut bytes_read, ov) + IocpFile::submit_read_overlapped(&port, req.offset, |ov| { + ReadFile(handle, buf_ptr.cast(), len, std::ptr::null_mut(), ov) }) } .map_err(PagedbError::Io)? as usize; + let bytes = checked_read_count(bytes, req.buf.len())?; // Zero tail past EOF — matches MemVfs / TokioVfs / Iouring. for b in &mut req.buf[bytes..] { *b = 0; @@ -184,11 +233,7 @@ impl VfsFile for IocpFile { return Ok(0); } let handle = self.handle(); - let len = u32::try_from(buf.len()).map_err(|_| { - PagedbError::Io(std::io::Error::other( - "buffer too large for u32 in WriteFile", - )) - })?; + let len = Self::writefile_len(buf.len())?; let port = Port { inner: Arc::clone(&self.port), }; @@ -197,8 +242,7 @@ impl VfsFile for IocpFile { // SAFETY: `buf` slice outlives this frame; closure runs synchronously. let bytes = unsafe { IocpFile::submit_overlapped(&port, offset, |ov| { - let mut bytes_written: u32 = 0; - WriteFile(handle, buf_ptr, len, &mut bytes_written, ov) + WriteFile(handle, buf_ptr, len, std::ptr::null_mut(), ov) }) } .map_err(PagedbError::Io)?; @@ -212,26 +256,36 @@ impl VfsFile for IocpFile { if reqs.is_empty() { return Ok(()); } + for req in reqs { + Self::writefile_len(req.buf.len())?; + } let handle = self.handle(); let port = Port { inner: Arc::clone(&self.port), }; let _guard = port.lock(); for req in reqs { - let len = u32::try_from(req.buf.len()).map_err(|_| { - PagedbError::Io(std::io::Error::other( - "buffer too large for u32 in WriteFile", - )) - })?; - let buf_ptr = req.buf.as_ptr(); - // SAFETY: `req.buf` outlives this frame; closure is synchronous. - unsafe { - IocpFile::submit_overlapped(&port, req.offset, |ov| { - let mut bytes_written: u32 = 0; - WriteFile(handle, buf_ptr, len, &mut bytes_written, ov) - }) + let mut offset = req.offset; + let mut remaining = req.buf; + while !remaining.is_empty() { + let len = Self::writefile_len(remaining.len())?; + let buf_ptr = remaining.as_ptr(); + // SAFETY: `remaining` is a subslice of `req.buf` and outlives + // this synchronous submission/completion wait. + let written = unsafe { + IocpFile::submit_overlapped(&port, offset, |ov| { + WriteFile(handle, buf_ptr, len, std::ptr::null_mut(), ov) + }) + } + .map_err(PagedbError::Io)?; + let written = usize::try_from(written).map_err(|_| { + PagedbError::Io(std::io::Error::other( + "WriteFile reported byte count that does not fit in usize", + )) + })?; + let consumed = checked_write_progress(&mut offset, written, remaining.len())?; + remaining = &remaining[consumed..]; } - .map_err(PagedbError::Io)?; } Ok(()) } @@ -255,10 +309,10 @@ impl VfsFile for IocpFile { // SetFilePointerEx then SetEndOfFile. We don't care about the file // pointer for subsequent I/O (all our ops are positional/overlapped), // so we just seek to `len` and set EOF there. + let len = checked_signed_file_len(len, "SetFilePointerEx")?; // SAFETY: `handle` is valid; output ptr is null because we don't need - // the new position back. - #[allow(clippy::cast_possible_wrap)] - let rc = unsafe { SetFilePointerEx(handle, len as i64, std::ptr::null_mut(), FILE_BEGIN) }; + // the new position back. `len` was validated before the syscall. + let rc = unsafe { SetFilePointerEx(handle, len, std::ptr::null_mut(), FILE_BEGIN) }; if rc == 0 { return Err(PagedbError::Io(std::io::Error::last_os_error())); } diff --git a/src/vfs/iocp/vfs.rs b/src/vfs/iocp/vfs.rs index 11da654..f74c625 100644 --- a/src/vfs/iocp/vfs.rs +++ b/src/vfs/iocp/vfs.rs @@ -19,7 +19,7 @@ use crate::errors::PagedbError; use super::file::IocpFile; use super::port::Port; -use crate::vfs::traits::Vfs; +use crate::vfs::traits::{Vfs, canonical_native_path, resolve_native_path}; use crate::vfs::types::OpenMode; use windows_sys::Win32::Foundation::{ @@ -172,8 +172,8 @@ impl IocpVfs { }) } - fn resolve(&self, p: &str) -> PathBuf { - self.inner.root.join(p.trim_start_matches('/')) + fn resolve(&self, path: &str) -> Result { + resolve_native_path(&self.inner.root, path) } fn lookup_or_create_entry(&self, path: &str) -> Arc { @@ -189,7 +189,8 @@ impl IocpVfs { } fn do_lock(&self, path: &str, kind: LockKind) -> Result { - let entry = self.lookup_or_create_entry(path); + 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) { @@ -199,7 +200,7 @@ impl IocpVfs { _ => return Err(PagedbError::AlreadyLocked), } } - let lock_path = self.resolve(path); + let lock_path = self.resolve(&logical_path)?; if let Some(parent) = lock_path.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; } @@ -228,7 +229,7 @@ impl Vfs for IocpVfs { type LockHandle = IocpLockHandle; async fn open(&self, path: &str, mode: OpenMode) -> Result { - let p = self.resolve(path); + let p = self.resolve(path)?; if matches!(mode, OpenMode::CreateNew | OpenMode::CreateOrOpen) { if let Some(parent) = p.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; @@ -294,7 +295,7 @@ impl Vfs for IocpVfs { } async fn remove(&self, path: &str) -> Result<()> { - let p = self.resolve(path); + let p = self.resolve(path)?; match std::fs::remove_file(&p) { Ok(()) => Ok(()), Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()), @@ -303,8 +304,8 @@ impl Vfs for IocpVfs { } async fn rename(&self, from: &str, to: &str) -> Result<()> { - let f = self.resolve(from); - let t = self.resolve(to); + let f = self.resolve(from)?; + let t = self.resolve(to)?; if let Some(parent) = t.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; } @@ -314,7 +315,7 @@ impl Vfs for IocpVfs { } async fn list_dir(&self, path: &str) -> Result> { - let p = self.resolve(path); + let p = self.resolve(path)?; let iter = match std::fs::read_dir(&p) { Ok(it) => it, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()), @@ -332,11 +333,12 @@ impl Vfs for IocpVfs { } async fn mkdir_all(&self, path: &str) -> Result<()> { - let p = self.resolve(path); + let p = self.resolve(path)?; std::fs::create_dir_all(&p).map_err(PagedbError::Io) } - async fn sync_dir(&self, _path: &str) -> Result<()> { + async fn sync_dir(&self, path: &str) -> Result<()> { + canonical_native_path(path)?; // NTFS folds rename durability into its metadata journal, and // `FlushFileBuffers` on a directory handle is not generally available // through `std::fs`. Best-effort no-op on Windows; rename + the diff --git a/src/vfs/iouring/file.rs b/src/vfs/iouring/file.rs index ef3e14a..af110ed 100644 --- a/src/vfs/iouring/file.rs +++ b/src/vfs/iouring/file.rs @@ -14,7 +14,10 @@ use parking_lot::Mutex; use crate::Result; use crate::errors::PagedbError; -use crate::vfs::traits::VfsFile; +use crate::vfs::traits::{ + VfsFile, checked_indexed_completion, checked_iouring_positioned_offset, checked_read_count, + checked_signed_file_len, write_all_at, +}; use crate::vfs::types::{ReadReq, WriteReq}; /// Per-file handle backed by an `std::fs::File` fd and the shared `io_uring`. @@ -33,6 +36,22 @@ impl IouringFile { } } + fn check_write_range(offset: u64, len: usize) -> Result<()> { + let len = u64::try_from(len).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "buffer length does not fit in u64", + )) + })?; + offset.checked_add(len).ok_or_else(|| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "write offset overflow", + )) + })?; + Ok(()) + } + /// Submit a single SQE, wait for exactly one CQE with matching /// `user_data`, and return the CQE result. /// @@ -117,16 +136,18 @@ impl IouringFile { } } ring.submit_and_wait(chunk_len)?; + let mut chunk_results = vec![None; chunk_len]; let mut found = 0usize; { let mut cq = ring.completion(); cq.sync(); for cqe in cq.by_ref() { - let ud = cqe.user_data(); - if ud < chunk_len as u64 { - #[allow(clippy::cast_possible_truncation)] - let idx = base + ud as usize; - results[idx] = cqe.result(); + if checked_indexed_completion(&mut chunk_results, cqe.user_data(), cqe.result()) + .map_err(|error| match error { + PagedbError::Io(io) => io, + other => std::io::Error::other(other.to_string()), + })? + { found += 1; } if found == chunk_len { @@ -139,6 +160,10 @@ impl IouringFile { "io_uring: fewer CQEs returned than submitted", )); } + for (index, result) in chunk_results.into_iter().enumerate() { + results[base + index] = result + .ok_or_else(|| std::io::Error::other("io_uring: missing indexed CQE result"))?; + } base = end; } Ok(results) @@ -158,6 +183,7 @@ impl VfsFile for IouringFile { if buf.is_empty() { return Ok(0); } + checked_iouring_positioned_offset(offset, buf.len())?; let fd = Fd(self.file.as_raw_fd()); let len = u32::try_from(buf.len()) .map_err(|_| PagedbError::Io(std::io::Error::other("buffer too large for u32")))?; @@ -172,7 +198,7 @@ impl VfsFile for IouringFile { let n = unsafe { Self::submit_one(&mut ring, &entry, 0) }.map_err(PagedbError::Io)?; // n >= 0 guaranteed by submit_one (negative becomes Err). #[allow(clippy::cast_sign_loss)] - Ok(n as usize) + checked_read_count(n as usize, buf.len()) } async fn read_at_vectored(&self, reqs: &mut [ReadReq<'_>]) -> Result<()> { @@ -183,6 +209,7 @@ impl VfsFile for IouringFile { // Build one Read SQE per request; each gets its index as user_data. let mut entries: Vec = Vec::with_capacity(reqs.len()); for (i, req) in reqs.iter_mut().enumerate() { + checked_iouring_positioned_offset(req.offset, req.buf.len())?; let len = u32::try_from(req.buf.len()) .map_err(|_| PagedbError::Io(std::io::Error::other("buffer too large for u32")))?; entries.push( @@ -208,7 +235,7 @@ impl VfsFile for IouringFile { } // res >= 0 guaranteed above. #[allow(clippy::cast_sign_loss)] - let nread = res as usize; + let nread = checked_read_count(res as usize, req.buf.len())?; for b in &mut req.buf[nread..] { *b = 0; } @@ -223,6 +250,7 @@ impl VfsFile for IouringFile { if buf.is_empty() { return Ok(0); } + Self::check_write_range(offset, buf.len())?; let fd = Fd(self.file.as_raw_fd()); let len = u32::try_from(buf.len()) .map_err(|_| PagedbError::Io(std::io::Error::other("buffer too large for u32")))?; @@ -246,30 +274,85 @@ impl VfsFile for IouringFile { if reqs.is_empty() { return Ok(()); } + for req in reqs { + Self::check_write_range(req.offset, req.buf.len())?; + } let fd = Fd(self.file.as_raw_fd()); - // Build one Write SQE per request; each gets its index as user_data. + // Empty requests are already complete. Skipping them also ensures that + // a zero-byte CQE always represents impossible progress on real data. let mut entries: Vec = Vec::with_capacity(reqs.len()); + let mut entry_to_request = Vec::with_capacity(reqs.len()); for (i, req) in reqs.iter().enumerate() { + if req.buf.is_empty() { + continue; + } let len = u32::try_from(req.buf.len()) .map_err(|_| PagedbError::Io(std::io::Error::other("buffer too large for u32")))?; + let user_data = u64::try_from(entry_to_request.len()).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "too many vectored write requests", + )) + })?; entries.push( opcode::Write::new(fd, req.buf.as_ptr(), len) .offset(req.offset) .build() - .user_data(i as u64), + .user_data(user_data), ); + entry_to_request.push(i); + } + if entries.is_empty() { + return Ok(()); } - let mut ring = self.ring.lock(); - // SAFETY: `req.buf` slices are tied to the `reqs` argument's `'_` - // lifetime. The ring lock is held across submit+drain. - let results = - unsafe { Self::submit_batch(&mut ring, &entries) }.map_err(PagedbError::Io)?; + let results = { + let mut ring = self.ring.lock(); + // SAFETY: `req.buf` slices are tied to the `reqs` argument's `'_` + // lifetime. The ring lock is held across submit+drain. + unsafe { Self::submit_batch(&mut ring, &entries) }.map_err(PagedbError::Io)? + }; + drop(entries); - for &res in &results { + let mut short_writes = Vec::new(); + for (entry_index, &res) in results.iter().enumerate() { if res < 0 { return Err(PagedbError::Io(std::io::Error::from_raw_os_error(-res))); } + let written = usize::try_from(res) + .map_err(|_| PagedbError::Io(std::io::Error::other("negative write result")))?; + let request_index = entry_to_request[entry_index]; + let request = &reqs[request_index]; + if written > request.buf.len() { + return Err(PagedbError::Io(std::io::Error::other( + "io_uring write overreported bytes", + ))); + } + if written == 0 { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::WriteZero, + ))); + } + if written < request.buf.len() { + short_writes.push((request_index, written)); + } + } + + for (request_index, written) in short_writes { + let request = &reqs[request_index]; + let written_u64 = u64::try_from(written).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "write count does not fit in u64", + )) + })?; + let offset = request.offset.checked_add(written_u64).ok_or_else(|| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "write offset overflow", + )) + })?; + write_all_at(self, offset, &request.buf[written..]).await?; } Ok(()) } @@ -291,14 +374,10 @@ impl VfsFile for IouringFile { // v0.7. Use the syscall directly via libc; for regular files this is // synchronous and does not trigger disk I/O in the common path. // - // SAFETY: `self.file.as_raw_fd()` is a valid open fd for the lifetime - // of this method call (`self` keeps `file` alive). `len` fits in - // `libc::off_t` (i64) on any 64-bit Linux target; on 32-bit targets - // `off_t` is 32-bit but `_FILE_OFFSET_BITS=64` is standard, so the - // cast is safe in practice. We allow the truncation lint here because - // we are on Linux where `off_t` is always i64 in practice. - #[allow(clippy::cast_possible_wrap)] - let rc = unsafe { libc::ftruncate(self.file.as_raw_fd(), len as libc::off_t) }; + let len = checked_signed_file_len(len, "ftruncate")?; + // SAFETY: `self.file.as_raw_fd()` is valid for this method call and + // `len` was checked before entering the signed native syscall. + let rc = unsafe { libc::ftruncate(self.file.as_raw_fd(), len) }; if rc != 0 { return Err(PagedbError::Io(std::io::Error::last_os_error())); } diff --git a/src/vfs/iouring/vfs.rs b/src/vfs/iouring/vfs.rs index e73a127..67ca515 100644 --- a/src/vfs/iouring/vfs.rs +++ b/src/vfs/iouring/vfs.rs @@ -16,7 +16,7 @@ use crate::errors::PagedbError; use super::file::IouringFile; use super::ring::Ring; -use crate::vfs::traits::Vfs; +use crate::vfs::traits::{Vfs, canonical_native_path, resolve_native_path}; use crate::vfs::types::OpenMode; // --------------------------------------------------------------------------- @@ -159,8 +159,8 @@ impl IouringVfs { }) } - fn resolve(&self, p: &str) -> PathBuf { - self.inner.root.join(p.trim_start_matches('/')) + fn resolve(&self, path: &str) -> Result { + resolve_native_path(&self.inner.root, path) } fn lookup_or_create_entry(&self, path: &str) -> Arc { @@ -176,7 +176,8 @@ impl IouringVfs { } fn do_lock(&self, path: &str, kind: LockKind) -> Result { - let entry = self.lookup_or_create_entry(path); + let logical_path = canonical_native_path(path)?; + let entry = self.lookup_or_create_entry(&logical_path); // In-process guard first. { let mut s = entry.state.lock(); @@ -187,7 +188,7 @@ impl IouringVfs { _ => return Err(PagedbError::AlreadyLocked), } } - let lock_path = self.resolve(path); + let lock_path = self.resolve(&logical_path)?; if let Some(parent) = lock_path.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; } @@ -217,7 +218,7 @@ impl Vfs for IouringVfs { type LockHandle = IouringLockHandle; async fn open(&self, path: &str, mode: OpenMode) -> Result { - let p = self.resolve(path); + let p = self.resolve(path)?; if matches!(mode, OpenMode::CreateNew | OpenMode::CreateOrOpen) { if let Some(parent) = p.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; @@ -267,7 +268,7 @@ impl Vfs for IouringVfs { } async fn remove(&self, path: &str) -> Result<()> { - let p = self.resolve(path); + let p = self.resolve(path)?; match std::fs::remove_file(&p) { Ok(()) => Ok(()), Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()), @@ -276,8 +277,8 @@ impl Vfs for IouringVfs { } async fn rename(&self, from: &str, to: &str) -> Result<()> { - let f = self.resolve(from); - let t = self.resolve(to); + let f = self.resolve(from)?; + let t = self.resolve(to)?; if let Some(parent) = t.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; } @@ -285,7 +286,7 @@ impl Vfs for IouringVfs { } async fn list_dir(&self, path: &str) -> Result> { - let p = self.resolve(path); + let p = self.resolve(path)?; let iter = match std::fs::read_dir(&p) { Ok(it) => it, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()), @@ -303,12 +304,12 @@ impl Vfs for IouringVfs { } async fn mkdir_all(&self, path: &str) -> Result<()> { - let p = self.resolve(path); + let p = self.resolve(path)?; std::fs::create_dir_all(&p).map_err(PagedbError::Io) } async fn sync_dir(&self, path: &str) -> Result<()> { - let p = self.resolve(path); + let p = self.resolve(path)?; // Open the directory with O_RDONLY|O_DIRECTORY and fsync the fd. let dir = match std::fs::File::open(&p) { Ok(d) => d, diff --git a/src/vfs/memory.rs b/src/vfs/memory.rs index a6cd233..6e00b39 100644 --- a/src/vfs/memory.rs +++ b/src/vfs/memory.rs @@ -9,7 +9,7 @@ use parking_lot::Mutex; use crate::Result; use crate::errors::PagedbError; -use super::traits::{Vfs, VfsFile}; +use super::traits::{Vfs, VfsFile, canonical_native_path}; use super::types::{OpenMode, ReadReq, WriteReq}; /// Shared inode storage. A path entry in `MemVfs::files` points to one of @@ -107,6 +107,7 @@ impl Vfs for MemVfs { type LockHandle = MemLockHandle; async fn open(&self, path: &str, mode: OpenMode) -> Result { + let path = &canonical_native_path(path)?; let mut files = self.inner.files.lock(); let exists = files.contains_key(path); let (inode, writable) = match (mode, exists) { @@ -132,23 +133,27 @@ impl Vfs for MemVfs { } async fn remove(&self, path: &str) -> Result<()> { + let path = &canonical_native_path(path)?; let mut files = self.inner.files.lock(); files.remove(path); Ok(()) } async fn rename(&self, from: &str, to: &str) -> Result<()> { + let from = &canonical_native_path(from)?; + let to = canonical_native_path(to)?; let mut files = self.inner.files.lock(); let Some(inode) = files.remove(from) else { return Err(PagedbError::Io(std::io::Error::from( std::io::ErrorKind::NotFound, ))); }; - files.insert(to.to_string(), inode); + files.insert(to, inode); Ok(()) } async fn list_dir(&self, path: &str) -> Result> { + let path = canonical_native_path(path)?; let prefix = if path.ends_with('/') { path.to_string() } else { @@ -165,18 +170,22 @@ impl Vfs for MemVfs { Ok(out) } - async fn mkdir_all(&self, _path: &str) -> Result<()> { - // In-memory backend has no persistent directory entries; mkdir is a no-op. + async fn mkdir_all(&self, path: &str) -> Result<()> { + // No persistent directory entries to create, but the path still has to + // be a legal one — silently accepting what every other backend rejects + // is how a test suite stops catching path bugs. + canonical_native_path(path)?; Ok(()) } - async fn sync_dir(&self, _path: &str) -> Result<()> { - // In-memory backend has no durability semantics; sync_dir is a no-op. + async fn sync_dir(&self, path: &str) -> Result<()> { + // No durability semantics, but the same path contract applies. + canonical_native_path(path)?; Ok(()) } async fn lock_exclusive(&self, path: &str) -> Result { - let lock_ref = self.lookup_or_create_lock(path); + let lock_ref = self.lookup_or_create_lock(&canonical_native_path(path)?); let mut state = lock_ref.state.lock(); match *state { LockState::Free => { @@ -192,7 +201,7 @@ impl Vfs for MemVfs { } async fn lock_shared(&self, path: &str) -> Result { - let lock_ref = self.lookup_or_create_lock(path); + let lock_ref = self.lookup_or_create_lock(&canonical_native_path(path)?); let mut state = lock_ref.state.lock(); let next = match *state { LockState::Free => LockState::Shared(1), @@ -251,7 +260,9 @@ impl VfsFile for MemFile { let mut inode = self.inode.lock(); let offset = usize::try_from(offset) .map_err(|_| PagedbError::Io(std::io::Error::from(std::io::ErrorKind::InvalidInput)))?; - let end = offset + buf.len(); + let end = offset.checked_add(buf.len()).ok_or_else(|| { + PagedbError::Io(std::io::Error::from(std::io::ErrorKind::InvalidInput)) + })?; if end > inode.data.len() { inode.data.resize(end, 0); } @@ -264,14 +275,23 @@ impl VfsFile for MemFile { return Err(PagedbError::ReadOnly); } let mut inode = self.inode.lock(); + let mut ranges = Vec::with_capacity(reqs.len()); for req in reqs { let offset = usize::try_from(req.offset).map_err(|_| { PagedbError::Io(std::io::Error::from(std::io::ErrorKind::InvalidInput)) })?; - let end = offset + req.buf.len(); - if end > inode.data.len() { - inode.data.resize(end, 0); - } + let end = offset.checked_add(req.buf.len()).ok_or_else(|| { + PagedbError::Io(std::io::Error::from(std::io::ErrorKind::InvalidInput)) + })?; + ranges.push((offset, end)); + } + + let max_end = ranges.iter().map(|(_, end)| *end).max().unwrap_or(0); + if max_end > inode.data.len() { + inode.data.resize(max_end, 0); + } + + for (req, (offset, end)) in reqs.iter().zip(ranges) { inode.data[offset..end].copy_from_slice(req.buf); } Ok(()) diff --git a/src/vfs/opfs/handle.rs b/src/vfs/opfs/handle.rs index 2ac6c86..720303f 100644 --- a/src/vfs/opfs/handle.rs +++ b/src/vfs/opfs/handle.rs @@ -10,7 +10,9 @@ use std::sync::Arc; use crate::Result; use crate::errors::PagedbError; -use crate::vfs::traits::VfsFile; +use crate::vfs::traits::{ + VfsFile, checked_opfs_byte_count, checked_opfs_file_size, checked_opfs_js_range, write_all_at, +}; use crate::vfs::types::{ReadReq, WriteReq}; use super::protocol::{ErrKind, OpfsOp, OpfsResult}; @@ -48,6 +50,7 @@ impl Drop for OpfsFile { impl VfsFile for OpfsFile { async fn read_at(&self, offset: u64, buf: &mut [u8]) -> Result { let len = buf.len(); + checked_opfs_js_range(offset, len)?; let result = self .vfs .dispatch(OpfsOp::Read { @@ -57,8 +60,14 @@ impl VfsFile for OpfsFile { }) .await?; match result { - OpfsResult::Data { bytes } => { - let n = bytes.len().min(buf.len()); + OpfsResult::Data { bytes, count } => { + let n = checked_opfs_byte_count("read", count, buf.len())?; + if bytes.len() != n { + return Err(PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "OPFS read payload length did not match its byte count", + ))); + } buf[..n].copy_from_slice(&bytes[..n]); Ok(n) } @@ -68,6 +77,9 @@ impl VfsFile for OpfsFile { } async fn read_at_vectored(&self, reqs: &mut [ReadReq<'_>]) -> Result<()> { + for req in reqs.iter() { + checked_opfs_js_range(req.offset, req.buf.len())?; + } for req in reqs.iter_mut() { let n = self.read_at(req.offset, req.buf).await?; // Zero-fill past the returned bytes (all-or-nothing contract). @@ -84,6 +96,7 @@ impl VfsFile for OpfsFile { } let data = buf.to_vec(); let len = data.len(); + checked_opfs_js_range(offset, len)?; let result = self .vfs .dispatch(OpfsOp::Write { @@ -93,7 +106,7 @@ impl VfsFile for OpfsFile { }) .await?; match result { - OpfsResult::Ok => Ok(len), + OpfsResult::Written { count } => checked_opfs_byte_count("write", count, len), OpfsResult::Err { reason, kind } => Err(map_err(&reason, kind)), _ => Err(PagedbError::Unsupported), } @@ -104,7 +117,10 @@ impl VfsFile for OpfsFile { return Err(PagedbError::ReadOnly); } for req in reqs { - self.write_at(req.offset, req.buf).await?; + checked_opfs_js_range(req.offset, req.buf.len())?; + } + for req in reqs { + write_all_at(self, req.offset, req.buf).await?; } Ok(()) } @@ -127,6 +143,7 @@ impl VfsFile for OpfsFile { if self.read_only { return Err(PagedbError::ReadOnly); } + checked_opfs_js_range(len, 0)?; let result = self .vfs .dispatch(OpfsOp::Truncate { @@ -149,7 +166,7 @@ impl VfsFile for OpfsFile { }) .await?; match result { - OpfsResult::Size { len } => Ok(len), + OpfsResult::Size { len } => checked_opfs_file_size(len), OpfsResult::Err { reason, kind } => Err(map_err(&reason, kind)), _ => Err(PagedbError::Unsupported), } diff --git a/src/vfs/opfs/opfs_worker.js b/src/vfs/opfs/opfs_worker.js index 76f4d9a..5350e45 100644 --- a/src/vfs/opfs/opfs_worker.js +++ b/src/vfs/opfs/opfs_worker.js @@ -128,8 +128,11 @@ function opRead({ handle_id, offset, len }) { const buf = new ArrayBuffer(len); const view = new Uint8Array(buf); const read = h.read(view, { at: Number(offset) }); + if (!Number.isSafeInteger(read) || read < 0 || read > len) { + return errResult("OPFS read returned an invalid byte count", "io"); + } const bytes = Array.from(new Uint8Array(buf, 0, read)); - return { type: "data", bytes }; + return { type: "data", bytes, count: read }; } catch (e) { return errResult(e.message || String(e), classifyError(e)); } @@ -140,8 +143,8 @@ function opWrite({ handle_id, offset, data }) { if (!h) return errResult("handle not found: " + handle_id, "notFound"); try { const view = new Uint8Array(data); - h.write(view, { at: Number(offset) }); - return ok(); + const written = h.write(view, { at: Number(offset) }); + return { type: "written", count: written }; } catch (e) { return errResult(e.message || String(e), classifyError(e)); } diff --git a/src/vfs/opfs/protocol.rs b/src/vfs/opfs/protocol.rs index dcb3b6a..4937f0f 100644 --- a/src/vfs/opfs/protocol.rs +++ b/src/vfs/opfs/protocol.rs @@ -94,10 +94,12 @@ pub enum OpfsResult { Locked { lock_id: u32 }, /// Generic success (no payload). Ok, - /// Byte data returned from a read. - Data { bytes: Vec }, - /// File size in bytes. - Size { len: u64 }, + /// Byte data and the raw JavaScript count returned from a read. + Data { bytes: Vec, count: f64 }, + /// Raw JavaScript count returned from a write. + Written { count: f64 }, + /// Raw JavaScript file size. + Size { len: f64 }, /// Directory entries (names only, not full paths). Entries { names: Vec }, /// Operation failed. diff --git a/src/vfs/tokio_backend.rs b/src/vfs/tokio_backend.rs index a6ecf6c..dabc25f 100644 --- a/src/vfs/tokio_backend.rs +++ b/src/vfs/tokio_backend.rs @@ -20,7 +20,7 @@ use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt}; use crate::Result; use crate::errors::PagedbError; -use super::traits::{Vfs, VfsFile}; +use super::traits::{Vfs, VfsFile, canonical_native_path, resolve_native_path}; use super::types::{OpenMode, ReadReq, WriteReq}; // --------------------------------------------------------------------------- @@ -298,8 +298,12 @@ impl TokioVfs { } } - fn resolve(&self, p: &str) -> PathBuf { - self.inner.root.join(p.trim_start_matches('/')) + fn canonical_logical_path(path: &str) -> Result { + canonical_native_path(path) + } + + fn resolve(&self, path: &str) -> Result { + resolve_native_path(&self.inner.root, path) } /// Return the filesystem root directory of this VFS instance. @@ -343,7 +347,7 @@ impl Vfs for TokioVfs { type LockHandle = TokioLockHandle; async fn open(&self, path: &str, mode: OpenMode) -> Result { - let p = self.resolve(path); + let p = self.resolve(path)?; // Ensure parent directories exist for create modes. if matches!(mode, OpenMode::CreateNew | OpenMode::CreateOrOpen) { if let Some(parent) = p.parent() { @@ -397,7 +401,7 @@ impl Vfs for TokioVfs { } async fn remove(&self, path: &str) -> Result<()> { - let p = self.resolve(path); + let p = self.resolve(path)?; match fs::remove_file(&p).await { Ok(()) => Ok(()), Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()), @@ -406,8 +410,8 @@ impl Vfs for TokioVfs { } async fn rename(&self, from: &str, to: &str) -> Result<()> { - let f = self.resolve(from); - let t = self.resolve(to); + let f = self.resolve(from)?; + let t = self.resolve(to)?; if let Some(parent) = t.parent() { fs::create_dir_all(parent).await.map_err(PagedbError::Io)?; } @@ -415,7 +419,7 @@ impl Vfs for TokioVfs { } async fn list_dir(&self, path: &str) -> Result> { - let p = self.resolve(path); + let p = self.resolve(path)?; let mut entries = match fs::read_dir(&p).await { Ok(e) => e, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()), @@ -432,12 +436,12 @@ impl Vfs for TokioVfs { } async fn mkdir_all(&self, path: &str) -> Result<()> { - let p = self.resolve(path); + let p = self.resolve(path)?; fs::create_dir_all(&p).await.map_err(PagedbError::Io) } async fn sync_dir(&self, path: &str) -> Result<()> { - let p = self.resolve(path); + let p = self.resolve(path)?; // Open the directory with std::fs (synchronous) to call sync_all. // On platforms where opening a directory handle is unsupported or // syncing it returns Unsupported/PermissionDenied (e.g., some Windows @@ -475,7 +479,8 @@ impl Vfs for TokioVfs { } async fn lock_exclusive(&self, path: &str) -> Result { - let entry = self.lookup_or_create_entry(path); + let logical_path = Self::canonical_logical_path(path)?; + let entry = self.lookup_or_create_entry(&logical_path); // In-process guard first: fast fail if the same process already holds // any lock on this path. { @@ -489,7 +494,7 @@ impl Vfs for TokioVfs { // so another process opening the same directory is also excluded. #[cfg(unix)] { - let lock_path = self.resolve(path); + let lock_path = self.resolve(&logical_path)?; if let Some(parent) = lock_path.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; } @@ -511,7 +516,7 @@ impl Vfs for TokioVfs { // lock so another process is also excluded. #[cfg(windows)] { - let lock_path = self.resolve(path); + let lock_path = self.resolve(&logical_path)?; if let Some(parent) = lock_path.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; } @@ -543,7 +548,8 @@ impl Vfs for TokioVfs { } async fn lock_shared(&self, path: &str) -> Result { - let entry = self.lookup_or_create_entry(path); + let logical_path = Self::canonical_logical_path(path)?; + let entry = self.lookup_or_create_entry(&logical_path); { let mut s = entry.state.lock(); match *s { @@ -554,7 +560,7 @@ impl Vfs for TokioVfs { } #[cfg(unix)] { - let lock_path = self.resolve(path); + let lock_path = self.resolve(&logical_path)?; if let Some(parent) = lock_path.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; } @@ -577,7 +583,7 @@ impl Vfs for TokioVfs { } #[cfg(windows)] { - let lock_path = self.resolve(path); + let lock_path = self.resolve(&logical_path)?; if let Some(parent) = lock_path.parent() { std::fs::create_dir_all(parent).map_err(PagedbError::Io)?; } @@ -659,6 +665,11 @@ impl VfsFile for TokioFile { if !self.writable { return Err(PagedbError::ReadOnly); } + let buf_len = u64::try_from(buf.len()) + .map_err(|_| PagedbError::Io(std::io::Error::from(std::io::ErrorKind::InvalidInput)))?; + offset.checked_add(buf_len).ok_or_else(|| { + PagedbError::Io(std::io::Error::from(std::io::ErrorKind::InvalidInput)) + })?; let mut f = self.inner.lock().await; f.seek(std::io::SeekFrom::Start(offset)) .await @@ -671,6 +682,14 @@ impl VfsFile for TokioFile { if !self.writable { return Err(PagedbError::ReadOnly); } + for req in reqs { + let buf_len = u64::try_from(req.buf.len()).map_err(|_| { + PagedbError::Io(std::io::Error::from(std::io::ErrorKind::InvalidInput)) + })?; + req.offset.checked_add(buf_len).ok_or_else(|| { + PagedbError::Io(std::io::Error::from(std::io::ErrorKind::InvalidInput)) + })?; + } let mut f = self.inner.lock().await; for req in reqs { f.seek(std::io::SeekFrom::Start(req.offset)) @@ -723,19 +742,77 @@ mod tests { use super::*; fn tempdir() -> PathBuf { + use std::sync::atomic::{AtomicU64, Ordering}; + + static NEXT: AtomicU64 = AtomicU64::new(0); let mut p = std::env::temp_dir(); let nanos = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map_or(0, |d| d.as_nanos()); + let sequence = NEXT.fetch_add(1, Ordering::Relaxed); p.push(format!( - "pagedb-tokio-unit-{}-{}", - std::process::id(), - nanos + "pagedb-tokio-unit-{}-{nanos}-{sequence}", + std::process::id() )); - std::fs::create_dir_all(&p).unwrap(); + std::fs::create_dir(&p).unwrap(); p } + #[test] + fn tempdir_helper_allocates_unique_roots() { + let dirs: Vec<_> = (0..128).map(|_| tempdir()).collect(); + let mut paths = dirs.clone(); + paths.sort(); + paths.dedup(); + assert_eq!(paths.len(), dirs.len()); + + for dir in dirs { + std::fs::remove_dir_all(dir).ok(); + } + } + + #[test] + fn canonical_logical_paths_have_one_spelling() { + assert_eq!( + TokioVfs::canonical_logical_path("/main.db").unwrap(), + "/main.db" + ); + assert_eq!( + TokioVfs::canonical_logical_path("main.db").unwrap(), + "/main.db" + ); + assert_eq!( + TokioVfs::canonical_logical_path("/seg/file/").unwrap(), + "/seg/file" + ); + assert_eq!( + TokioVfs::canonical_logical_path(r"seg\file").unwrap(), + "/seg/file" + ); + assert_eq!(TokioVfs::canonical_logical_path("/").unwrap(), "/"); + } + + #[test] + fn rejects_parent_current_and_empty_path_components() { + let vfs = TokioVfs::new("/tmp/pagedb-root"); + + for path in [ + "../escape", + "/seg/../escape", + "./main.db", + "seg/./x", + "seg//x", + ] { + let error = vfs.resolve(path).unwrap_err(); + match error { + PagedbError::Io(error) => { + assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput); + } + other => panic!("expected InvalidInput, got {other:?}"), + } + } + } + #[tokio::test(flavor = "current_thread")] async fn write_and_read_round_trip() { let dir = tempdir(); diff --git a/src/vfs/traits.rs b/src/vfs/traits.rs index 104f9c9..8c52323 100644 --- a/src/vfs/traits.rs +++ b/src/vfs/traits.rs @@ -10,8 +10,14 @@ use super::types::{OpenMode, ReadReq, WriteReq}; /// Abstract file-system surface used by the pager. Implementations provide /// platform-appropriate I/O and advisory path locking. /// -/// Each `path` passed to a lock method is its own lock domain — locks on -/// distinct paths never conflict. +/// Root-backed implementations treat paths as logical, root-relative names. +/// Leading separators are optional, while parent/current-directory and empty +/// interior components are rejected so no operation can escape the configured +/// root or create an aliased lock domain. +/// +/// Each distinct logical path passed to a lock method is its own lock domain. +/// Equivalent spellings normalize to the same domain before lookup; distinct +/// canonical paths never conflict. /// /// All async methods return futures that are `Send`, matching the `Send + Sync` /// supertrait bounds on `Vfs` itself. This allows `Db` futures to be @@ -62,9 +68,11 @@ pub trait Vfs: Send + Sync { } } -/// Per-file I/O surface. Vectored ops are all-or-nothing — either every -/// request is satisfied or the call returns an error and no partial state is -/// observable to the caller. +/// Per-file I/O surface. Vectored operations validate every request before +/// performing any I/O, so invalid deterministic input cannot leave a valid +/// prefix of the batch applied. A runtime device or filesystem failure part-way +/// through is still reported as it happens — the guarantee is about request +/// validation, not rollback. /// /// All async methods return `Send` futures so that `VfsFile` values can be /// held across await points inside `Send` futures. @@ -137,19 +145,8 @@ pub(crate) async fn write_all_at( ) -> Result<()> { while !buf.is_empty() { let written = file.write_at(offset, buf).await?; - if written == 0 { - return Err(PagedbError::Io(std::io::Error::from( - std::io::ErrorKind::WriteZero, - ))); - } - checked_transfer_progress( - &mut offset, - written, - buf.len(), - "write_at", - "positional write offset", - )?; - buf = &buf[written..]; + let consumed = checked_write_progress(&mut offset, written, buf.len())?; + buf = &buf[consumed..]; } Ok(()) } @@ -210,6 +207,311 @@ fn checked_transfer_progress( Ok(()) } +/// Canonicalize a logical VFS path and reject components that could escape the +/// configured root. +pub(crate) fn canonical_native_path(path: &str) -> Result { + let trimmed = path.trim_matches(['/', '\\']); + if trimmed.is_empty() { + return Ok("/".to_string()); + } + + // Split first, then judge each component on its own. Parsing the path as a + // whole is not enough: a platform prefix is only recognised in leading + // position, so a drive letter deeper in the path (`a/C:/b`) parses as an + // ordinary name — and `PathBuf::push` would later let it replace the root + // outright. Every component has to survive parsing in isolation. + let parts: Vec<_> = trimmed.split(['/', '\\']).collect(); + for part in &parts { + if !is_plain_path_component(part) { + return Err(invalid_native_path(path)); + } + } + + Ok(format!("/{}", parts.join("/"))) +} + +/// Whether one path component is an ordinary name that can only ever extend a +/// path — never re-root it. +/// +/// A colon is rejected on every target, not just the one that gives it meaning. +/// `a/C:/b` naming a file under the root on Linux and a different drive on +/// Windows would make the same logical path resolve to two different places, +/// which is exactly the portability the store is supposed to guarantee. +fn is_plain_path_component(part: &str) -> bool { + use std::path::Component; + if part.contains(':') { + return false; + } + let mut components = std::path::Path::new(part).components(); + let Some(Component::Normal(name)) = components.next() else { + return false; + }; + components.next().is_none() && name == std::ffi::OsStr::new(part) +} + +/// Resolve a canonical logical path beneath a native VFS root. +#[cfg(any(test, not(target_arch = "wasm32")))] +pub(crate) fn resolve_native_path( + root: &std::path::Path, + path: &str, +) -> Result { + let logical = canonical_native_path(path)?; + let mut resolved = root.to_path_buf(); + for part in logical.trim_start_matches('/').split('/') { + if !part.is_empty() { + resolved.push(part); + } + } + Ok(resolved) +} + +fn invalid_native_path(path: &str) -> PagedbError { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + format!("VFS path escapes root: {path:?}"), + )) +} + +/// Validate a buffer length before passing it to Win32 `ReadFile`. +#[cfg(any(test, target_os = "windows"))] +pub(crate) fn checked_readfile_len(len: usize) -> Result { + u32::try_from(len).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "buffer too large for u32 in ReadFile", + )) + }) +} + +/// Validate one reported read completion before using it as a slice boundary. +/// +/// The offset-advancing form is [`checked_read_progress`]; this is for backends +/// that already know where the read landed and only need the count checked. A +/// short read is legal here — only a count the caller never asked for is not. +#[cfg(any( + test, + target_os = "windows", + target_os = "linux", + all(target_os = "android", not(target_arch = "arm")), +))] +pub(crate) fn checked_read_count(read: usize, requested: usize) -> Result { + if read > requested { + return Err(PagedbError::vfs_contract_violated( + "read_at", + "reported more bytes than the caller requested", + )); + } + Ok(read) +} + +/// Record one indexed batch completion, rejecting duplicate in-range keys. +#[cfg(any( + test, + target_os = "linux", + all(target_os = "android", not(target_arch = "arm")), +))] +pub(crate) fn checked_indexed_completion( + slots: &mut [Option], + user_data: u64, + result: i32, +) -> Result { + let Ok(index) = usize::try_from(user_data) else { + return Ok(false); + }; + let Some(slot) = slots.get_mut(index) else { + return Ok(false); + }; + if slot.is_some() { + return Err(PagedbError::Io(std::io::Error::other( + "duplicate indexed completion", + ))); + } + *slot = Some(result); + Ok(true) +} + +/// Reject `io_uring`'s `offset == -1` sentinel for positioned I/O. +#[cfg(any( + test, + target_os = "linux", + all(target_os = "android", not(target_arch = "arm")), +))] +pub(crate) fn checked_iouring_positioned_offset(offset: u64, len: usize) -> Result<()> { + if len > 0 && offset == u64::MAX { + return Err(PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "io_uring offset -1 uses the current file position", + ))); + } + Ok(()) +} + +/// Largest exactly representable integer in a JavaScript `number`. +#[cfg(any(test, all(target_arch = "wasm32", feature = "opfs")))] +pub(crate) const OPFS_MAX_SAFE_INTEGER: u64 = 9_007_199_254_740_991; + +#[cfg(any(test, all(target_arch = "wasm32", feature = "opfs")))] +const OPFS_MAX_SAFE_INTEGER_F64: f64 = 9_007_199_254_740_991.0; + +/// Validate an OPFS offset and length before crossing the JavaScript number +/// boundary. +#[cfg(any(test, all(target_arch = "wasm32", feature = "opfs")))] +pub(crate) fn checked_opfs_js_range(offset: u64, len: usize) -> Result<()> { + let len = u64::try_from(len).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "OPFS length does not fit u64", + )) + })?; + let end = offset.checked_add(len).ok_or_else(|| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "OPFS offset range overflow", + )) + })?; + if end > OPFS_MAX_SAFE_INTEGER { + return Err(PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + format!("offset {offset} + len {len} exceeds JS safe integer range"), + ))); + } + Ok(()) +} + +/// Validate a byte count returned through an OPFS JavaScript `number`. +#[cfg(any(test, all(target_arch = "wasm32", feature = "opfs")))] +pub(crate) fn checked_opfs_byte_count( + kind: &'static str, + n: f64, + requested: usize, +) -> Result { + if !n.is_finite() || n < 0.0 || n.fract() != 0.0 || n > OPFS_MAX_SAFE_INTEGER_F64 { + return Err(PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!("OPFS {kind} returned an invalid byte count"), + ))); + } + #[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)] + let count = n as u64; + let requested = u64::try_from(requested).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "OPFS request length does not fit u64", + )) + })?; + if count > requested { + return Err(PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!("OPFS {kind} overreported bytes"), + ))); + } + usize::try_from(count).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!("OPFS {kind} byte count does not fit usize"), + )) + }) +} + +/// Validate a file size returned through an OPFS JavaScript `number`. +#[cfg(any(test, all(target_arch = "wasm32", feature = "opfs")))] +pub(crate) fn checked_opfs_file_size(size: f64) -> Result { + if !size.is_finite() || size < 0.0 || size.fract() != 0.0 || size > OPFS_MAX_SAFE_INTEGER_F64 { + return Err(PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "OPFS returned an invalid file size", + ))); + } + #[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)] + Ok(size as u64) +} + +/// Validate a `u64` file length before passing it to signed native syscalls. +#[cfg(any( + test, + target_os = "windows", + target_os = "linux", + all(target_os = "android", not(target_arch = "arm")), + target_os = "macos", + target_os = "ios", +))] +pub(crate) fn checked_signed_file_len(len: u64, syscall: &'static str) -> Result { + i64::try_from(len).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + format!("{syscall} length does not fit signed file offset"), + )) + }) +} + +/// Classified result of starting a Windows overlapped I/O request. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[cfg(any(test, target_os = "windows"))] +pub(crate) enum OverlappedStart { + CompletionQueued, + EmptyRead, +} + +/// Decide whether a Win32 overlapped I/O call produced a completion to drain. +#[cfg(any(test, target_os = "windows"))] +pub(crate) fn checked_overlapped_start( + rc: i32, + last_error: u32, + error_io_pending: u32, +) -> Result { + if rc != 0 || last_error == error_io_pending { + return Ok(OverlappedStart::CompletionQueued); + } + let code = i32::try_from(last_error).map_err(|_| { + PagedbError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "Win32 error code does not fit i32", + )) + })?; + Err(PagedbError::Io(std::io::Error::from_raw_os_error(code))) +} + +/// Decide whether a Win32 overlapped read produced data, EOF, or an error. +#[cfg(any(test, target_os = "windows"))] +pub(crate) fn checked_overlapped_read_start( + rc: i32, + last_error: u32, + error_io_pending: u32, + error_handle_eof: u32, +) -> Result { + if rc == 0 && last_error == error_handle_eof { + return Ok(OverlappedStart::EmptyRead); + } + checked_overlapped_start(rc, last_error, error_io_pending) +} + +/// Validate one reported write completion, advance the positioned offset, and +/// return how many bytes the caller may consume. +/// +/// The mirror of [`checked_read_progress`], and the single write-side rule: +/// [`write_all_at`] and every backend that completes a short write use it, so +/// there is one definition of what progress a backend is allowed to claim. +#[inline] +pub(crate) fn checked_write_progress( + offset: &mut u64, + written: usize, + remaining: usize, +) -> Result { + if written == 0 { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::WriteZero, + ))); + } + checked_transfer_progress( + offset, + written, + remaining, + "write_at", + "positional write offset", + )?; + Ok(written) +} + #[cfg(test)] mod tests { use super::*; @@ -257,4 +559,258 @@ mod tests { "expected ArithmeticOverflow, got {err:?}" ); } + + use std::sync::{Arc, Mutex}; + + use super::{ + OverlappedStart, VfsFile, checked_indexed_completion, checked_iouring_positioned_offset, + checked_opfs_byte_count, checked_opfs_file_size, checked_opfs_js_range, + checked_overlapped_read_start, checked_overlapped_start, checked_read_count, + checked_readfile_len, checked_signed_file_len, checked_write_progress, write_all_at, + }; + use crate::vfs::types::{ReadReq, WriteReq}; + + #[derive(Clone, Default)] + struct ScriptedWriteFile { + state: Arc>, + } + + #[derive(Default)] + struct ScriptedWriteState { + completions: Vec, + calls: Vec<(u64, usize)>, + } + + impl ScriptedWriteFile { + fn with_completions(completions: Vec) -> Self { + Self { + state: Arc::new(Mutex::new(ScriptedWriteState { + completions, + calls: Vec::new(), + })), + } + } + + fn calls(&self) -> Vec<(u64, usize)> { + self.state.lock().unwrap().calls.clone() + } + } + + impl VfsFile for ScriptedWriteFile { + async fn read_at(&self, _offset: u64, _buf: &mut [u8]) -> crate::Result { + unimplemented!("read_at is not used by write_all_at tests") + } + + async fn read_at_vectored(&self, _reqs: &mut [ReadReq<'_>]) -> crate::Result<()> { + unimplemented!("read_at_vectored is not used by write_all_at tests") + } + + async fn write_at(&mut self, offset: u64, buf: &[u8]) -> crate::Result { + let mut state = self.state.lock().unwrap(); + state.calls.push((offset, buf.len())); + Ok(state.completions.remove(0)) + } + + async fn write_at_vectored(&mut self, _reqs: &[WriteReq<'_>]) -> crate::Result<()> { + unimplemented!("write_at_vectored is not used by write_all_at tests") + } + + async fn sync(&mut self) -> crate::Result<()> { + Ok(()) + } + + async fn truncate(&mut self, _len: u64) -> crate::Result<()> { + Ok(()) + } + + async fn len(&self) -> crate::Result { + Ok(0) + } + + async fn is_empty(&self) -> crate::Result { + Ok(true) + } + + fn supports_direct_io(&self) -> bool { + false + } + } + + #[test] + fn readfile_len_accepts_u32_max() { + assert_eq!(checked_readfile_len(u32::MAX as usize).unwrap(), u32::MAX); + } + + #[cfg(target_pointer_width = "64")] + #[test] + fn readfile_len_rejects_above_u32_max() { + let err = checked_readfile_len(u32::MAX as usize + 1).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + } + + #[test] + fn read_count_accepts_short_completion() { + assert_eq!(checked_read_count(3, 10).unwrap(), 3); + } + + #[test] + fn read_count_rejects_overreported_completion() { + let err = checked_read_count(11, 10).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + } + + #[test] + fn indexed_completion_records_unique_user_data() { + let mut slots = vec![None, None]; + assert!(checked_indexed_completion(&mut slots, 1, -5).unwrap()); + assert_eq!(slots, vec![None, Some(-5)]); + } + + #[test] + fn indexed_completion_ignores_out_of_range_user_data() { + let mut slots = vec![None, None]; + assert!(!checked_indexed_completion(&mut slots, 2, 9).unwrap()); + assert_eq!(slots, vec![None, None]); + } + + #[test] + fn indexed_completion_rejects_duplicate_user_data() { + let mut slots = vec![None, None]; + assert!(checked_indexed_completion(&mut slots, 0, 7).unwrap()); + let err = checked_indexed_completion(&mut slots, 0, 9).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + assert_eq!(slots, vec![Some(7), None]); + } + + #[test] + fn iouring_positioned_offset_rejects_u64_max_for_non_empty_io() { + let err = checked_iouring_positioned_offset(u64::MAX, 1).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + } + + #[test] + fn iouring_positioned_offset_allows_u64_max_for_empty_io() { + checked_iouring_positioned_offset(u64::MAX, 0).unwrap(); + } + + #[test] + fn opfs_js_range_rejects_end_above_safe_integer() { + let err = checked_opfs_js_range(9_007_199_254_740_991, 1).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + } + + #[test] + fn opfs_js_range_accepts_exact_safe_integer_end() { + checked_opfs_js_range(9_007_199_254_740_990, 1).unwrap(); + } + + #[test] + fn opfs_byte_count_rejects_nan() { + let err = checked_opfs_byte_count("read", f64::NAN, 8).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + } + + #[test] + fn opfs_byte_count_rejects_fractional_count() { + let err = checked_opfs_byte_count("read", 1.5, 8).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + } + + #[test] + fn opfs_byte_count_rejects_overreported_count() { + let err = checked_opfs_byte_count("write", 9.0, 8).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + } + + #[test] + fn opfs_byte_count_accepts_exact_count() { + assert_eq!(checked_opfs_byte_count("write", 8.0, 8).unwrap(), 8); + } + + #[test] + fn opfs_file_size_rejects_unsafe_integer() { + let error = checked_opfs_file_size(9_007_199_254_740_992.0).unwrap_err(); + assert!(matches!(error, crate::errors::PagedbError::Io(_))); + } + + #[test] + fn signed_file_len_accepts_i64_max() { + assert_eq!( + checked_signed_file_len(i64::MAX as u64, "truncate").unwrap(), + i64::MAX + ); + } + + #[test] + fn signed_file_len_rejects_above_i64_max() { + let err = checked_signed_file_len(i64::MAX as u64 + 1, "truncate").unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + } + + #[test] + fn overlapped_start_treats_immediate_success_as_queued_completion() { + assert_eq!( + checked_overlapped_start(1, 0, 997).unwrap(), + OverlappedStart::CompletionQueued + ); + } + + #[test] + fn overlapped_start_treats_pending_as_queued_completion() { + assert_eq!( + checked_overlapped_start(0, 997, 997).unwrap(), + OverlappedStart::CompletionQueued + ); + } + + #[test] + fn overlapped_start_rejects_immediate_error_without_completion() { + let err = checked_overlapped_start(0, 5, 997).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + } + + #[test] + fn overlapped_read_start_maps_immediate_eof_to_empty_read() { + assert_eq!( + checked_overlapped_read_start(0, 38, 997, 38).unwrap(), + OverlappedStart::EmptyRead + ); + } + + #[test] + fn write_progress_advances_offset_for_short_write() { + let mut offset = 7; + assert_eq!(checked_write_progress(&mut offset, 3, 10).unwrap(), 3); + assert_eq!(offset, 10); + } + + #[test] + fn write_progress_rejects_zero_for_non_empty_write() { + let mut offset = 7; + let err = checked_write_progress(&mut offset, 0, 10).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + assert_eq!(offset, 7, "failed progress must not advance offset"); + } + + #[test] + fn write_progress_rejects_overreported_count() { + let mut offset = 7; + let err = checked_write_progress(&mut offset, 11, 10).unwrap_err(); + assert!(matches!(err, crate::errors::PagedbError::Io(_))); + assert_eq!(offset, 7, "failed progress must not advance offset"); + } + + #[tokio::test(flavor = "current_thread")] + async fn write_all_at_returns_after_one_full_write() { + let mut file = ScriptedWriteFile::with_completions(vec![4]); + write_all_at(&mut file, 11, b"page").await.unwrap(); + assert_eq!(file.calls(), vec![(11, 4)]); + } + + #[tokio::test(flavor = "current_thread")] + async fn write_all_at_retries_after_short_write() { + let mut file = ScriptedWriteFile::with_completions(vec![2, 2]); + write_all_at(&mut file, 11, b"page").await.unwrap(); + assert_eq!(file.calls(), vec![(11, 4), (13, 2)]); + } } diff --git a/tests/vfs_gcd.rs b/tests/vfs_gcd.rs new file mode 100644 index 0000000..5a0a8be --- /dev/null +++ b/tests/vfs_gcd.rs @@ -0,0 +1,210 @@ +//! Integration tests for the Apple-platform `GcdVfs` backend. + +#![cfg(any(target_os = "macos", target_os = "ios"))] + +use pagedb::vfs::{GcdVfs, OpenMode, ReadReq, Vfs, VfsFile, WriteReq}; + +fn tempdir() -> tempfile::TempDir { + tempfile::Builder::new() + .prefix("pagedb-vfs-gcd-") + .tempdir() + .unwrap() +} + +#[tokio::test(flavor = "current_thread")] +async fn vectored_round_trip_reopens_and_zero_fills_past_eof() { + let dir = tempdir(); + let vfs = GcdVfs::new(dir.path()); + let mut writer = vfs.open("/round-trip", OpenMode::CreateNew).await.unwrap(); + writer + .write_at_vectored(&[ + WriteReq { + offset: 0, + buf: b"abc", + }, + WriteReq { + offset: 8, + buf: b"xyz", + }, + ]) + .await + .unwrap(); + writer.sync().await.unwrap(); + drop(writer); + + let reader = vfs.open("/round-trip", OpenMode::Read).await.unwrap(); + let mut first = [0xff; 3]; + let mut tail = [0xff; 5]; + reader + .read_at_vectored(&mut [ + ReadReq { + offset: 0, + buf: &mut first, + }, + ReadReq { + offset: 8, + buf: &mut tail, + }, + ]) + .await + .unwrap(); + + assert_eq!(&first, b"abc"); + assert_eq!(tail, [b'x', b'y', b'z', 0, 0]); +} + +#[tokio::test(flavor = "current_thread")] +async fn truncate_len_and_read_only_contract() { + let dir = tempdir(); + let vfs = GcdVfs::new(dir.path()); + let mut writer = vfs.open("/truncate", OpenMode::CreateNew).await.unwrap(); + writer.write_at(0, b"abcdefgh").await.unwrap(); + assert_eq!(writer.len().await.unwrap(), 8); + writer.truncate(4).await.unwrap(); + assert_eq!(writer.len().await.unwrap(), 4); + drop(writer); + + let mut reader = vfs.open("/truncate", OpenMode::Read).await.unwrap(); + assert!(matches!( + reader.write_at(0, b"x").await, + Err(pagedb::errors::PagedbError::ReadOnly) + )); + assert!(matches!( + reader.truncate(0).await, + Err(pagedb::errors::PagedbError::ReadOnly) + )); + assert!(!reader.supports_direct_io()); +} + +#[tokio::test(flavor = "current_thread")] +async fn rejects_parent_directory_escape() { + let dir = tempdir(); + let root_name = dir.path().file_name().unwrap().to_string_lossy(); + let escaped_file = dir + .path() + .parent() + .unwrap() + .join(format!("{root_name}-escaped-file")); + let vfs = GcdVfs::new(dir.path()); + + let open_result = vfs.open("../escaped-file", OpenMode::CreateNew).await; + let escaped_file_was_created = escaped_file.exists(); + std::fs::remove_file(&escaped_file).ok(); + let error = match open_result { + Ok(_) => panic!("open must reject paths outside the configured root"), + Err(error) => error, + }; + match error { + pagedb::errors::PagedbError::Io(error) => { + assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput); + } + other => panic!("expected InvalidInput, got {other:?}"), + } + assert!(!escaped_file_was_created); +} + +#[tokio::test(flavor = "current_thread")] +async fn equivalent_lock_paths_share_one_domain() { + let dir = tempdir(); + let vfs = GcdVfs::new(dir.path()); + let _held = vfs.lock_exclusive("/db.lock").await.unwrap(); + + let error = match vfs.lock_exclusive("db.lock").await { + Ok(_) => panic!("equivalent logical paths must share one lock domain"), + Err(error) => error, + }; + assert!(matches!(error, pagedb::errors::PagedbError::AlreadyLocked)); +} + +#[tokio::test(flavor = "current_thread")] +async fn vectored_write_rejects_unrepresentable_offset_without_partial_write() { + let dir = tempdir(); + let vfs = GcdVfs::new(dir.path()); + let mut file = vfs + .open("/vec-overflow", OpenMode::CreateNew) + .await + .unwrap(); + let writes = [ + WriteReq { + offset: 0, + buf: b"kept-out", + }, + WriteReq { + offset: u64::MAX, + buf: b"x", + }, + ]; + + let error = file + .write_at_vectored(&writes) + .await + .expect_err("an invalid offset must reject the whole vectored write"); + assert!(matches!(error, pagedb::errors::PagedbError::Io(_))); + + let mut actual = [0xff; 8]; + assert_eq!(file.read_at(0, &mut actual).await.unwrap(), 0); +} + +#[tokio::test(flavor = "current_thread")] +async fn vectored_write_rejects_signed_offset_range_without_partial_write() { + let dir = tempdir(); + let vfs = GcdVfs::new(dir.path()); + let mut file = vfs + .open("/vec-signed-overflow", OpenMode::CreateNew) + .await + .unwrap(); + let writes = [ + WriteReq { + offset: 0, + buf: b"kept-out", + }, + WriteReq { + offset: i64::MAX as u64, + buf: b"xx", + }, + ]; + + let error = file + .write_at_vectored(&writes) + .await + .expect_err("a range past off_t::MAX must reject the whole vectored write"); + assert!(matches!(error, pagedb::errors::PagedbError::Io(_))); + + let mut actual = [0xff; 8]; + assert_eq!(file.read_at(0, &mut actual).await.unwrap(), 0); +} + +#[tokio::test(flavor = "current_thread")] +async fn vectored_read_rejects_unrepresentable_offset_without_touching_buffers() { + let dir = tempdir(); + let vfs = GcdVfs::new(dir.path()); + let mut writer = vfs + .open("/read-overflow", OpenMode::CreateNew) + .await + .unwrap(); + writer.write_at(0, b"visible!").await.unwrap(); + writer.sync().await.unwrap(); + drop(writer); + + let reader = vfs.open("/read-overflow", OpenMode::Read).await.unwrap(); + let mut first = [0x55; 8]; + let mut second = [0x66; 1]; + let mut reads = [ + ReadReq { + offset: 0, + buf: &mut first, + }, + ReadReq { + offset: u64::MAX, + buf: &mut second, + }, + ]; + + let error = reader + .read_at_vectored(&mut reads) + .await + .expect_err("an invalid offset must reject before filling buffers"); + assert!(matches!(error, pagedb::errors::PagedbError::Io(_))); + assert_eq!(first, [0x55; 8]); + assert_eq!(second, [0x66; 1]); +} diff --git a/tests/vfs_iocp.rs b/tests/vfs_iocp.rs new file mode 100644 index 0000000..64cf476 --- /dev/null +++ b/tests/vfs_iocp.rs @@ -0,0 +1,117 @@ +//! Integration tests for the Windows `IocpVfs` backend. + +#![cfg(target_os = "windows")] + +use pagedb::vfs::{IocpVfs, OpenMode, ReadReq, Vfs, VfsFile, WriteReq}; + +fn tempdir() -> tempfile::TempDir { + tempfile::Builder::new() + .prefix("pagedb-vfs-iocp-") + .tempdir() + .unwrap() +} + +#[tokio::test(flavor = "current_thread")] +async fn vectored_round_trip_reopens_and_zero_fills_past_eof() { + let dir = tempdir(); + let vfs = IocpVfs::new(dir.path()).unwrap(); + let mut writer = vfs.open("/round-trip", OpenMode::CreateNew).await.unwrap(); + writer + .write_at_vectored(&[ + WriteReq { + offset: 0, + buf: b"abc", + }, + WriteReq { + offset: 8, + buf: b"xyz", + }, + ]) + .await + .unwrap(); + writer.sync().await.unwrap(); + drop(writer); + + let reader = vfs.open("/round-trip", OpenMode::Read).await.unwrap(); + let mut first = [0xff; 3]; + let mut tail = [0xff; 5]; + reader + .read_at_vectored(&mut [ + ReadReq { + offset: 0, + buf: &mut first, + }, + ReadReq { + offset: 8, + buf: &mut tail, + }, + ]) + .await + .unwrap(); + + assert_eq!(&first, b"abc"); + assert_eq!(tail, [b'x', b'y', b'z', 0, 0]); +} + +#[tokio::test(flavor = "current_thread")] +async fn truncate_len_and_read_only_contract() { + let dir = tempdir(); + let vfs = IocpVfs::new(dir.path()).unwrap(); + let mut writer = vfs.open("/truncate", OpenMode::CreateNew).await.unwrap(); + writer.write_at(0, b"abcdefgh").await.unwrap(); + assert_eq!(writer.len().await.unwrap(), 8); + writer.truncate(4).await.unwrap(); + assert_eq!(writer.len().await.unwrap(), 4); + drop(writer); + + let mut reader = vfs.open("/truncate", OpenMode::Read).await.unwrap(); + assert!(matches!( + reader.write_at(0, b"x").await, + Err(pagedb::errors::PagedbError::ReadOnly) + )); + assert!(matches!( + reader.truncate(0).await, + Err(pagedb::errors::PagedbError::ReadOnly) + )); + assert!(!reader.supports_direct_io()); +} + +#[tokio::test(flavor = "current_thread")] +async fn rejects_parent_directory_escape() { + let dir = tempdir(); + let root_name = dir.path().file_name().unwrap().to_string_lossy(); + let escaped_file = dir + .path() + .parent() + .unwrap() + .join(format!("{root_name}-escaped-file")); + let vfs = IocpVfs::new(dir.path()).unwrap(); + + let open_result = vfs.open("../escaped-file", OpenMode::CreateNew).await; + let escaped_file_was_created = escaped_file.exists(); + std::fs::remove_file(&escaped_file).ok(); + let error = match open_result { + Ok(_) => panic!("open must reject paths outside the configured root"), + Err(error) => error, + }; + match error { + pagedb::errors::PagedbError::Io(error) => { + assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput); + } + other => panic!("expected InvalidInput, got {other:?}"), + } + assert!(!escaped_file_was_created); +} + +#[tokio::test(flavor = "current_thread")] +async fn equivalent_lock_paths_share_one_domain() { + let dir = tempdir(); + let vfs = IocpVfs::new(dir.path()).unwrap(); + let _held = vfs.lock_exclusive("/db.lock").await.unwrap(); + + let error = match vfs.lock_exclusive("db.lock").await { + Ok(_) => panic!("equivalent logical paths must share one lock domain"), + Err(error) => error, + }; + assert!(matches!(error, pagedb::errors::PagedbError::AlreadyLocked)); +} diff --git a/tests/vfs_iouring.rs b/tests/vfs_iouring.rs index 730bde5..c881aeb 100644 --- a/tests/vfs_iouring.rs +++ b/tests/vfs_iouring.rs @@ -81,6 +81,100 @@ async fn vectored_write_and_read() { std::fs::remove_dir_all(&dir).ok(); } +#[tokio::test(flavor = "current_thread")] +async fn vectored_write_rejects_invalid_offset_without_partial_write() { + let dir = tempdir("vec_write_invalid_offset"); + let vfs = IouringVfs::new(&dir).unwrap(); + let mut file = vfs + .open("/vec-write-invalid-offset", OpenMode::CreateNew) + .await + .unwrap(); + let writes = [ + WriteReq { + offset: 0, + buf: b"kept-out", + }, + WriteReq { + offset: u64::MAX, + buf: b"x", + }, + ]; + + let error = file + .write_at_vectored(&writes) + .await + .expect_err("an invalid io_uring offset must reject the entire batch"); + assert!(matches!(error, pagedb::errors::PagedbError::Io(_))); + + let mut actual = [0xff; 8]; + assert_eq!(file.read_at(0, &mut actual).await.unwrap(), 0); + std::fs::remove_dir_all(&dir).ok(); +} + +#[tokio::test(flavor = "current_thread")] +async fn vectored_read_rejects_invalid_offset_without_touching_buffers() { + let dir = tempdir("vec_read_invalid_offset"); + let vfs = IouringVfs::new(&dir).unwrap(); + let mut writer = vfs + .open("/vec-read-invalid-offset", OpenMode::CreateNew) + .await + .unwrap(); + writer.write_at(0, b"visible!").await.unwrap(); + writer.sync().await.unwrap(); + drop(writer); + + let reader = vfs + .open("/vec-read-invalid-offset", OpenMode::Read) + .await + .unwrap(); + let mut first = [0x55; 8]; + let mut second = [0x66; 1]; + let mut reads = [ + ReadReq { + offset: 0, + buf: &mut first, + }, + ReadReq { + offset: u64::MAX, + buf: &mut second, + }, + ]; + + let error = reader + .read_at_vectored(&mut reads) + .await + .expect_err("an invalid io_uring offset must reject before reading"); + assert!(matches!(error, pagedb::errors::PagedbError::Io(_))); + assert_eq!(first, [0x55; 8]); + assert_eq!(second, [0x66; 1]); + std::fs::remove_dir_all(&dir).ok(); +} + +#[tokio::test(flavor = "current_thread")] +async fn vectored_read_zero_fills_past_eof() { + let dir = tempdir("vec_read_eof"); + let vfs = IouringVfs::new(&dir).unwrap(); + let mut writer = vfs + .open("/vec-read-eof", OpenMode::CreateNew) + .await + .unwrap(); + writer.write_at(0, b"abc").await.unwrap(); + writer.sync().await.unwrap(); + drop(writer); + + let reader = vfs.open("/vec-read-eof", OpenMode::Read).await.unwrap(); + let mut actual = [0xff; 5]; + reader + .read_at_vectored(&mut [ReadReq { + offset: 0, + buf: &mut actual, + }]) + .await + .unwrap(); + assert_eq!(actual, [b'a', b'b', b'c', 0, 0]); + std::fs::remove_dir_all(&dir).ok(); +} + #[tokio::test(flavor = "current_thread")] async fn truncate_and_len() { let dir = tempdir("trunc"); @@ -95,6 +189,47 @@ async fn truncate_and_len() { std::fs::remove_dir_all(&dir).ok(); } +#[tokio::test(flavor = "current_thread")] +async fn rejects_parent_directory_escape() { + let dir = tempdir("root_escape"); + let root_name = dir.file_name().unwrap().to_string_lossy(); + let escaped_file = dir + .parent() + .unwrap() + .join(format!("{root_name}-escaped-file")); + let vfs = IouringVfs::new(&dir).unwrap(); + + let open_result = vfs.open("../escaped-file", OpenMode::CreateNew).await; + let escaped_file_was_created = escaped_file.exists(); + std::fs::remove_file(&escaped_file).ok(); + let error = match open_result { + Ok(_) => panic!("open must reject paths outside the configured root"), + Err(error) => error, + }; + match error { + pagedb::errors::PagedbError::Io(error) => { + assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput); + } + other => panic!("expected InvalidInput, got {other:?}"), + } + assert!(!escaped_file_was_created); + std::fs::remove_dir_all(&dir).ok(); +} + +#[tokio::test(flavor = "current_thread")] +async fn equivalent_lock_paths_share_one_domain() { + let dir = tempdir("lock_alias"); + let vfs = IouringVfs::new(&dir).unwrap(); + let _held = vfs.lock_exclusive("/db.lock").await.unwrap(); + + let error = match vfs.lock_exclusive("db.lock").await { + Ok(_) => panic!("equivalent logical paths must share one lock domain"), + Err(error) => error, + }; + assert!(matches!(error, pagedb::errors::PagedbError::AlreadyLocked)); + std::fs::remove_dir_all(&dir).ok(); +} + #[tokio::test(flavor = "current_thread")] async fn sync_dir_smoke() { let dir = tempdir("syncdir"); diff --git a/tests/vfs_memory.rs b/tests/vfs_memory.rs index d7812b5..3ef0d8b 100644 --- a/tests/vfs_memory.rs +++ b/tests/vfs_memory.rs @@ -49,6 +49,31 @@ async fn vectored_read_write_round_trip() { assert_eq!(&b, b"BBBB"); } +#[tokio::test(flavor = "current_thread")] +async fn vectored_write_rejects_offset_overflow_without_partial_write() { + let vfs = MemVfs::new(); + let mut file = vfs.open("/x", OpenMode::CreateNew).await.unwrap(); + let requests = [ + WriteReq { + offset: 0, + buf: b"kept-out", + }, + WriteReq { + offset: u64::MAX, + buf: b"x", + }, + ]; + + let error = file + .write_at_vectored(&requests) + .await + .expect_err("the invalid request must reject the entire vector"); + assert!(matches!(error, PagedbError::Io(_))); + + let mut actual = [0xff; 8]; + assert_eq!(file.read_at(0, &mut actual).await.unwrap(), 0); +} + #[tokio::test(flavor = "current_thread")] async fn vectored_read_zero_fills_past_eof() { let vfs = MemVfs::new(); diff --git a/tests/vfs_tokio.rs b/tests/vfs_tokio.rs index 16cf138..f46bbdf 100644 --- a/tests/vfs_tokio.rs +++ b/tests/vfs_tokio.rs @@ -1,19 +1,42 @@ +#![cfg(not(target_arch = "wasm32"))] + //! Integration tests for `TokioVfs` — real disk I/O under a temporary directory. +use pagedb::errors::PagedbError; use pagedb::vfs::tokio_backend::TokioVfs; use pagedb::vfs::{OpenMode, ReadReq, Vfs, VfsFile, WriteReq}; fn tempdir() -> std::path::PathBuf { + use std::sync::atomic::{AtomicU64, Ordering}; + + static NEXT: AtomicU64 = AtomicU64::new(0); let mut p = std::env::temp_dir(); let nanos = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map(|d| d.as_nanos()) .unwrap_or(0); - p.push(format!("pagedb-vfs-tokio-{}-{}", std::process::id(), nanos)); - std::fs::create_dir_all(&p).unwrap(); + let sequence = NEXT.fetch_add(1, Ordering::Relaxed); + p.push(format!( + "pagedb-vfs-tokio-{}-{nanos}-{sequence}", + std::process::id() + )); + std::fs::create_dir(&p).unwrap(); p } +#[test] +fn tempdir_helper_allocates_unique_roots() { + let dirs: Vec<_> = (0..128).map(|_| tempdir()).collect(); + let mut paths = dirs.clone(); + paths.sort(); + paths.dedup(); + assert_eq!(paths.len(), dirs.len()); + + for dir in dirs { + std::fs::remove_dir_all(dir).ok(); + } +} + #[tokio::test(flavor = "current_thread")] async fn round_trip_file() { let dir = tempdir(); @@ -77,6 +100,36 @@ async fn vectored_write_then_read() { std::fs::remove_dir_all(&dir).ok(); } +#[tokio::test(flavor = "current_thread")] +async fn vectored_write_rejects_offset_overflow_without_partial_write() { + let dir = tempdir(); + let vfs = TokioVfs::new(&dir); + let mut file = vfs + .open("/vec-overflow", OpenMode::CreateNew) + .await + .unwrap(); + let requests = [ + WriteReq { + offset: 0, + buf: b"kept-out", + }, + WriteReq { + offset: u64::MAX, + buf: b"x", + }, + ]; + + let error = file + .write_at_vectored(&requests) + .await + .expect_err("the invalid request must reject the entire vector"); + assert!(matches!(error, PagedbError::Io(_))); + + let mut actual = [0xff; 8]; + assert_eq!(file.read_at(0, &mut actual).await.unwrap(), 0); + std::fs::remove_dir_all(&dir).ok(); +} + #[tokio::test(flavor = "current_thread")] async fn sync_dir_succeeds() { let dir = tempdir(); @@ -87,6 +140,76 @@ async fn sync_dir_succeeds() { std::fs::remove_dir_all(&dir).ok(); } +#[tokio::test(flavor = "current_thread")] +async fn rejects_parent_directory_escape() { + let dir = tempdir(); + let root_name = dir.file_name().unwrap().to_string_lossy(); + let escaped_file = dir + .parent() + .unwrap() + .join(format!("{root_name}-escaped-file")); + let escaped_dir = dir + .parent() + .unwrap() + .join(format!("{root_name}-escaped-dir")); + let vfs = TokioVfs::new(&dir); + + let open_result = vfs.open("../escaped-file", OpenMode::CreateNew).await; + let escaped_file_was_created = escaped_file.exists(); + std::fs::remove_file(&escaped_file).ok(); + let open_error = match open_result { + Ok(_) => panic!("open must reject paths outside the configured root"), + Err(error) => error, + }; + match open_error { + PagedbError::Io(error) => { + assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput); + } + other => panic!("expected InvalidInput, got {other:?}"), + } + assert!( + !escaped_file_was_created, + "an invalid open path must not create a file outside the VFS root" + ); + + let mkdir_result = vfs.mkdir_all("../escaped-dir").await; + let escaped_dir_was_created = escaped_dir.exists(); + std::fs::remove_dir_all(&escaped_dir).ok(); + let mkdir_error = mkdir_result.expect_err("mkdir_all must reject a root escape"); + match mkdir_error { + PagedbError::Io(error) => { + assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput); + } + other => panic!("expected InvalidInput, got {other:?}"), + } + assert!( + !escaped_dir_was_created, + "an invalid mkdir path must not create a directory outside the VFS root" + ); + + std::fs::remove_dir_all(&dir).ok(); +} + +#[tokio::test(flavor = "current_thread")] +async fn lock_paths_are_normalized_before_conflict_check() { + let dir = tempdir(); + let vfs = TokioVfs::new(&dir); + let _held = vfs.lock_exclusive("/db.lock").await.unwrap(); + + let error = match vfs.lock_exclusive("db.lock").await { + Ok(_) => panic!("equivalent logical paths must share one lock domain"), + Err(error) => error, + }; + assert!(matches!(error, PagedbError::AlreadyLocked)); + + let error = match vfs.lock_shared("db.lock").await { + Ok(_) => panic!("equivalent logical paths must share one lock domain"), + Err(error) => error, + }; + assert!(matches!(error, PagedbError::AlreadyLocked)); + std::fs::remove_dir_all(&dir).ok(); +} + #[tokio::test(flavor = "current_thread")] async fn list_dir_returns_sorted_entries() { let dir = tempdir(); From 937de160c46b45bb276d7bb44c8ff55d059f92cf Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Mon, 27 Jul 2026 07:04:55 +0800 Subject: [PATCH 2/2] test(vfs): cover path canonicalization and lock aliasing on the memory backend Adds tests asserting canonical_native_path/resolve_native_path reject escaping and drive-letter components, and that the in-memory VFS treats equivalent path spellings as one lock domain and one file, matching the native backends' contract. Also switches two String clones in MemVfs away from unnecessary to_string() calls. --- src/vfs/memory.rs | 4 +-- src/vfs/traits.rs | 66 +++++++++++++++++++++++++++++++++++++++++++-- tests/vfs_memory.rs | 48 +++++++++++++++++++++++++++++++++ 3 files changed, 114 insertions(+), 4 deletions(-) diff --git a/src/vfs/memory.rs b/src/vfs/memory.rs index 6e00b39..62bd702 100644 --- a/src/vfs/memory.rs +++ b/src/vfs/memory.rs @@ -125,7 +125,7 @@ impl Vfs for MemVfs { } (OpenMode::CreateNew | OpenMode::CreateOrOpen, false) => { let inode = Arc::new(Mutex::new(MemInode { data: Vec::new() })); - files.insert(path.to_string(), inode.clone()); + files.insert(path.clone(), inode.clone()); (inode, true) } }; @@ -155,7 +155,7 @@ impl Vfs for MemVfs { async fn list_dir(&self, path: &str) -> Result> { let path = canonical_native_path(path)?; let prefix = if path.ends_with('/') { - path.to_string() + path.clone() } else { format!("{path}/") }; diff --git a/src/vfs/traits.rs b/src/vfs/traits.rs index 8c52323..b60be08 100644 --- a/src/vfs/traits.rs +++ b/src/vfs/traits.rs @@ -636,6 +636,50 @@ mod tests { } } + #[test] + fn a_logical_path_normalizes_its_separators_and_leading_slash() { + for spelling in ["/seg/abc", "seg/abc", "seg\\abc", "/seg/abc/"] { + assert_eq!(canonical_native_path(spelling).unwrap(), "/seg/abc"); + } + assert_eq!(canonical_native_path("").unwrap(), "/"); + assert_eq!(canonical_native_path("/").unwrap(), "/"); + } + + #[test] + fn a_parent_or_current_directory_component_never_resolves() { + for escape in ["../escaped", "a/../../escaped", "./a", "a//b", "a/./b"] { + assert!( + canonical_native_path(escape).is_err(), + "{escape:?} must not canonicalize" + ); + } + } + + #[test] + fn a_drive_letter_component_is_rejected_wherever_it_appears() { + // Rejected on every target, not just the one that gives it meaning. + // A platform prefix is only recognised in leading position, so a drive + // letter deeper in the path parses as an ordinary name — and pushing it + // would replace the root outright rather than extend it. + for drive in ["C:", "a/C:/b", "C:/escaped", "a/C:"] { + assert!( + canonical_native_path(drive).is_err(), + "{drive:?} must not canonicalize" + ); + } + } + + #[test] + fn a_resolved_path_always_stays_under_its_root() { + let root = std::path::Path::new("/srv/pagedb"); + assert_eq!( + resolve_native_path(root, "seg/abc").unwrap(), + root.join("seg").join("abc") + ); + assert!(resolve_native_path(root, "../escaped").is_err()); + assert!(resolve_native_path(root, "a/C:/b").is_err()); + } + #[test] fn readfile_len_accepts_u32_max() { assert_eq!(checked_readfile_len(u32::MAX as usize).unwrap(), u32::MAX); @@ -656,7 +700,16 @@ mod tests { #[test] fn read_count_rejects_overreported_completion() { let err = checked_read_count(11, 10).unwrap_err(); - assert!(matches!(err, crate::errors::PagedbError::Io(_))); + assert!( + matches!( + err, + crate::errors::PagedbError::VfsContractViolated { + operation: "read_at", + .. + } + ), + "expected VfsContractViolated, got {err:?}" + ); } #[test] @@ -796,7 +849,16 @@ mod tests { fn write_progress_rejects_overreported_count() { let mut offset = 7; let err = checked_write_progress(&mut offset, 11, 10).unwrap_err(); - assert!(matches!(err, crate::errors::PagedbError::Io(_))); + assert!( + matches!( + err, + crate::errors::PagedbError::VfsContractViolated { + operation: "write_at", + .. + } + ), + "expected VfsContractViolated, got {err:?}" + ); assert_eq!(offset, 7, "failed progress must not advance offset"); } diff --git a/tests/vfs_memory.rs b/tests/vfs_memory.rs index 3ef0d8b..8a82311 100644 --- a/tests/vfs_memory.rs +++ b/tests/vfs_memory.rs @@ -208,3 +208,51 @@ async fn read_mode_handle_cannot_write() { let err = g.write_at(0, b"nope").await.err().unwrap(); assert!(matches!(err, PagedbError::ReadOnly)); } + +/// The in-memory backend is the one nearly every test runs on, so it has to +/// honour the same logical-path contract as the native backends. Spellings that +/// name one file must occupy one lock domain here too, or the suite stops being +/// able to catch path aliasing at all. +#[tokio::test(flavor = "current_thread")] +async fn equivalent_spellings_share_one_lock_domain() { + let vfs = MemVfs::new(); + let held = vfs.lock_exclusive("/db.lock").await.unwrap(); + for alias in ["db.lock", "/db.lock", "db.lock/"] { + assert!( + matches!( + vfs.lock_exclusive(alias).await, + Err(pagedb::PagedbError::AlreadyLocked) + ), + "{alias:?} must contend with the held lock" + ); + } + drop(held); + vfs.lock_exclusive("db.lock").await.unwrap(); +} + +/// Equivalent spellings must also name one file. +#[tokio::test(flavor = "current_thread")] +async fn equivalent_spellings_name_one_file() { + let vfs = MemVfs::new(); + let mut f = vfs.open("seg/one", OpenMode::CreateNew).await.unwrap(); + f.write_at(0, b"payload").await.unwrap(); + drop(f); + + let mut buf = [0u8; 7]; + let g = vfs.open("/seg/one", OpenMode::Read).await.unwrap(); + assert_eq!(g.read_at(0, &mut buf).await.unwrap(), 7); + assert_eq!(&buf, b"payload"); + assert_eq!(vfs.list_dir("/seg").await.unwrap(), vec!["one".to_string()]); +} + +/// A path that escapes the root is rejected by the in-memory backend as well. +#[tokio::test(flavor = "current_thread")] +async fn a_path_that_escapes_the_root_is_rejected() { + let vfs = MemVfs::new(); + for escape in ["../escaped", "a/../../escaped", "a/C:/b"] { + assert!( + vfs.open(escape, OpenMode::CreateNew).await.is_err(), + "{escape:?} must not open" + ); + } +}