From cf188b52924e8993e41e9b445cdd363443b9ce3f Mon Sep 17 00:00:00 2001 From: Vladimir Petrzhikovskii Date: Mon, 29 Sep 2025 14:05:28 +0200 Subject: [PATCH] fix test in ci, improve is_same_fs Apply new clippy lints --- Cargo.toml | 2 +- src/cas_manager.rs | 16 ++-- src/index/mod.rs | 16 ++-- src/index/state.rs | 8 +- src/io.rs | 206 +++++++++++++++++---------------------------- src/lib.rs | 24 +++--- src/orphan.rs | 8 +- src/wal/mod.rs | 26 +++--- src/wal/storage.rs | 26 +++--- src/wal/tests.rs | 6 +- 10 files changed, 138 insertions(+), 200 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 856d740..297486f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -2,7 +2,7 @@ name = "cassadilia" version = "0.4.2" edition = "2024" -rust-version = "1.87" +rust-version = "1.89.0" authors = ["0xdeafbeef"] description = "A content-addressable storage (CAS) system optimized for large blobs with read-mostly access patterns" license = "MIT OR Apache-2.0" diff --git a/src/cas_manager.rs b/src/cas_manager.rs index 7cd6010..8d226d5 100644 --- a/src/cas_manager.rs +++ b/src/cas_manager.rs @@ -139,14 +139,14 @@ impl CasManager { blob_hash: &BlobHash, ) -> Result { let final_cas_path = self.paths.cas_file_path(blob_hash); - if !self.dir_tree_is_pre_created { - if let Some(parent) = final_cas_path.parent() { - std::fs::create_dir_all(parent).map_err(|e| CasManagerError::FileOperation { - operation: CasIoOperation::CreateSubdir, - path: final_cas_path.clone(), - source: e, - })?; - } + if !self.dir_tree_is_pre_created + && let Some(parent) = final_cas_path.parent() + { + std::fs::create_dir_all(parent).map_err(|e| CasManagerError::FileOperation { + operation: CasIoOperation::CreateSubdir, + path: final_cas_path.clone(), + source: e, + })?; } match std::fs::rename(staging_path, &final_cas_path) { Ok(()) => { diff --git a/src/index/mod.rs b/src/index/mod.rs index bf9b5f4..77f0552 100644 --- a/src/index/mod.rs +++ b/src/index/mod.rs @@ -117,14 +117,14 @@ where // Revert: Remove our intent from pending_intents let mut intents = self.index.pending_intents.lock(); - if let Some(current_hash) = intents.get(&self.key) { - if *current_hash == self.hash { - intents.remove(&self.key); - - // If we had replaced an existing intent, restore it - if let Some(replaced_hash) = self.replaced_hash { - intents.insert(self.key.clone(), replaced_hash); - } + if let Some(current_hash) = intents.get(&self.key) + && *current_hash == self.hash + { + intents.remove(&self.key); + + // If we had replaced an existing intent, restore it + if let Some(replaced_hash) = self.replaced_hash { + intents.insert(self.key.clone(), replaced_hash); } } } diff --git a/src/index/state.rs b/src/index/state.rs index 85785ae..cfa9ca5 100644 --- a/src/index/state.rs +++ b/src/index/state.rs @@ -76,10 +76,10 @@ where WalOp::Remove { keys } => { // Remove mappings and decrement each old blob's refcount. for key in keys { - if let Some(item) = self.key_to_hash.remove(key) { - if let Some(h) = self.decrement_ref(&item.blob_hash)? { - unreferenced_hashes.push(h); - } + if let Some(item) = self.key_to_hash.remove(key) + && let Some(h) = self.decrement_ref(&item.blob_hash)? + { + unreferenced_hashes.push(h); } } } diff --git a/src/io.rs b/src/io.rs index 6f54785..fea3604 100644 --- a/src/io.rs +++ b/src/io.rs @@ -1,8 +1,6 @@ -use std::fs::{File, OpenOptions}; +use std::fs::OpenOptions; use std::io; use std::io::Write; -use std::mem::ManuallyDrop; -use std::os::fd::AsRawFd; use std::path::{Path, PathBuf}; use thiserror::Error; @@ -45,77 +43,6 @@ pub(crate) fn atomically_write_file_bytes( Ok(()) } -pub(crate) trait FileExt { - fn lock(self) -> io::Result; -} - -impl FileExt for File { - #[inline] - fn lock(self) -> io::Result { - LockedFile::new(self) - } -} - -#[derive(Debug)] -#[repr(transparent)] -pub(crate) struct LockedFile { - inner: File, -} - -impl LockedFile { - pub fn new(file: File) -> std::io::Result { - // SAFETY: An existing file descriptor is used. - if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) } != 0 { - return Err(io::Error::last_os_error()); - } - - Ok(Self { inner: file }) - } - - #[allow(unused)] - pub fn unlock(self) -> std::io::Result { - let res = self.unlock_impl(); - - // Avoid calling `Drop` on manual unlock. - let this = ManuallyDrop::new(self); - // SAFETY: `this` is a valid pointer. - let this = unsafe { std::ptr::read(&this.inner) }; - - res.map(|_| this) - } - - fn unlock_impl(&self) -> std::io::Result<()> { - // SAFETY: An existing file descriptor is used. - if unsafe { libc::flock(self.inner.as_raw_fd(), libc::LOCK_UN | libc::LOCK_NB) } != 0 { - Err(io::Error::last_os_error()) - } else { - Ok(()) - } - } -} - -impl Drop for LockedFile { - fn drop(&mut self) { - self.unlock_impl().ok(); - } -} - -impl std::ops::Deref for LockedFile { - type Target = File; - - #[inline] - fn deref(&self) -> &Self::Target { - &self.inner - } -} - -impl std::ops::DerefMut for LockedFile { - #[inline] - fn deref_mut(&mut self) -> &mut Self::Target { - &mut self.inner - } -} - #[derive(Debug, Clone)] pub enum AtomicWriteStep { CreateTemp, @@ -138,56 +65,25 @@ pub enum IoError { #[cfg(target_os = "linux")] pub fn is_same_fs(paths: &[&Path]) -> Result { - use std::ffi::CString; - use std::mem::MaybeUninit; - use std::os::unix::ffi::OsStrExt; - if paths.len() < 2 { return Ok(true); } - // Best-effort mount-id via statx (Linux usually ≥ 5.8; may be backported). - fn mount_id(p: &Path) -> io::Result> { - let c = CString::new(p.as_os_str().as_bytes()) - .map_err(|_e| io::Error::new(io::ErrorKind::InvalidInput, "path contains NUL"))?; - unsafe { - let mut stx = MaybeUninit::::uninit(); - let ret = libc::statx( - libc::AT_FDCWD, - c.as_ptr(), - libc::AT_SYMLINK_NOFOLLOW, - libc::STATX_MNT_ID | libc::STATX_TYPE, - stx.as_mut_ptr(), - ); - if ret == 0 { - let stx = stx.assume_init(); - // Only trust stx_mnt_id if the kernel set the bit. - if stx.stx_mask & libc::STATX_MNT_ID != 0 { - return Ok(Some(stx.stx_mnt_id)); - } - return Ok(None); - } - match io::Error::last_os_error().raw_os_error() { - Some(libc::ENOSYS | libc::EOPNOTSUPP | libc::EINVAL) => Ok(None), - _ => Err(io::Error::last_os_error()), - } - } - } - // Try mount-ids first. - if let Some(first_mid) = mount_id(paths[0])? { + if let Some(first_mid) = mount_id_statx(paths[0])? { for p in paths.iter().skip(1) { - match mount_id(p) { + match mount_id_statx(p) { Ok(Some(mid)) if mid == first_mid => {} Ok(Some(_)) => return Ok(false), Ok(None) => return same_fs_dev(paths), Err(e) => return Err(e), } } + Ok(true) + } else { + // Fallback: compare st_dev + same_fs_dev(paths) } - - // Fallback: compare st_dev - same_fs_dev(paths) } #[cfg(all(unix, not(target_os = "linux")))] @@ -213,43 +109,93 @@ pub fn same_fs_dev(paths: &[&Path]) -> Result { Ok(true) } -#[cfg(not(unix))] -fn is_same_fs(path: &[&Path]) -> Result { - compile_error!("plz send a pr if you want to use it on non-unix systems"); +// Best-effort mount-id via statx (Linux usually ≥ 5.8; may be backported). +#[cfg(target_os = "linux")] +fn mount_id_statx(p: &Path) -> io::Result> { + use std::ffi::CString; + use std::mem::MaybeUninit; + use std::os::unix::ffi::OsStrExt; + + let c = CString::new(p.as_os_str().as_bytes()) + .map_err(|_e| io::Error::new(io::ErrorKind::InvalidInput, "path contains NUL"))?; + unsafe { + let mut stx = MaybeUninit::::uninit(); + let ret = libc::statx( + libc::AT_FDCWD, + c.as_ptr(), + libc::AT_NO_AUTOMOUNT, + libc::STATX_MNT_ID | libc::STATX_TYPE, + stx.as_mut_ptr(), + ); + if ret == 0 { + let stx = stx.assume_init(); + // Only trust stx_mnt_id if the kernel set the bit. + if stx.stx_mask & libc::STATX_MNT_ID != 0 { + return Ok(Some(stx.stx_mnt_id)); + } + return Ok(None); + } + + let err = io::Error::last_os_error(); + match err.raw_os_error() { + Some(libc::ENOSYS | libc::EOPNOTSUPP | libc::EINVAL) => Ok(None), + _ => Err(err), + } + } } +#[cfg(not(unix))] +compile_error!("plz send a pr if you want to use it on non-unix systems"); + #[cfg(test)] mod tests { use super::*; + use crate::{Cas, LibError}; #[test] fn file_lock() { let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join("tempfile"); - let mut options = std::fs::OpenOptions::new(); - options.create(true).truncate(true).write(true); + let _cas: Cas = Cas::open(dir.path(), Default::default()).unwrap(); - let file = options.open(&path).unwrap().lock().unwrap(); + let cas2: Result, _> = Cas::::open(dir.path(), Default::default()); - options.open(&path).unwrap().lock().unwrap_err(); + assert!(matches!(cas2, Err(LibError::AlreadyOpened))); + } + + #[cfg(target_os = "linux")] + #[test] + fn same_fs_linux() -> anyhow::Result<()> { + let root = Path::new("/"); + let dev = Path::new("/dev"); + let shm = Path::new("/dev/shm"); + let null = Path::new("/dev/null"); - drop(file); + assert!(is_same_fs(&[dev, null])?); + assert!(!is_same_fs(&[root, dev])?); + assert!(!is_same_fs(&[shm, null])?); - options.open(&path).unwrap().lock().unwrap(); + assert_eq!(mount_id_statx(dev)?, mount_id_statx(null)?); + assert_ne!(mount_id_statx(root)?, mount_id_statx(dev)?); + assert_ne!(mount_id_statx(shm)?, mount_id_statx(null)?); + + assert!(same_fs_dev(&[dev, null])?); + assert!(!same_fs_dev(&[root, dev])?); + assert!(!same_fs_dev(&[shm, null])?); + + Ok(()) } + #[cfg(target_os = "macos")] #[test] - fn same_fs() { - let tmp = tempfile::tempdir().unwrap(); - let root = Path::new("/"); + fn same_fs_macos() -> anyhow::Result<()> { + let temp_root = tempfile::tempdir()?; + let nested = temp_root.path().join("nested"); + std::fs::create_dir_all(&nested)?; - assert!(!is_same_fs(&[tmp.path(), root]).unwrap()); - assert!(is_same_fs(&[tmp.path(), Path::new("/tmp")]).unwrap()); - assert!(is_same_fs(&[tmp.path(), tmp.path()]).unwrap()); + assert!(is_same_fs(&[temp_root.path(), nested.as_path()])?); - assert!(!same_fs_dev(&[tmp.path(), root]).unwrap()); - assert!(same_fs_dev(&[tmp.path(), Path::new("/tmp")]).unwrap()); - assert!(same_fs_dev(&[tmp.path(), tmp.path()]).unwrap()); + // I've no access to macOS, so gues it's enough to test. + Ok(()) } } diff --git a/src/lib.rs b/src/lib.rs index b2ba551..69e0f1f 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -36,8 +36,6 @@ mod transaction; use cas_manager::{CasManager, CasManagerError}; pub use transaction::Transaction; -use self::io::{FileExt, LockedFile}; - #[derive(Debug, Clone)] pub enum LibIoOperation { CreateLockFile, @@ -194,8 +192,7 @@ where } pub struct CasInner { - #[allow(unused)] - lockfile: LockedFile, + _lockfile: File, pub(crate) paths: paths::DbPaths, pub(crate) index: Index, datasync_channel: Option>, @@ -242,9 +239,9 @@ where operation: LibIoOperation::CreateLockFile, path: Some(paths.lockfile_path().to_path_buf()), source: e, - })? - .lock() - .map_err(|_e| LibError::AlreadyOpened)?; + })?; + + lockfile.try_lock().map_err(|_e| LibError::AlreadyOpened)?; // Load or create settings let settings_persister = SettingsPersister::new(paths.settings_path().to_path_buf()); @@ -297,7 +294,7 @@ where } }; - Ok(Self { lockfile, paths, index, datasync_channel, cas_manager }) + Ok(Self { _lockfile: lockfile, paths, index, datasync_channel, cas_manager }) } /// Returns a read-only view of the index state. @@ -434,13 +431,12 @@ where Err(cas_error) => { if let Some(io_err) = cas_error.source().and_then(|s| s.downcast_ref::()) + && io_err.kind() == std::io::ErrorKind::NotFound { - if io_err.kind() == std::io::ErrorKind::NotFound { - return Err(LibError::BlobDataMissing { - key: format!("{key:?}"), - hash: item.blob_hash, - }); - } + return Err(LibError::BlobDataMissing { + key: format!("{key:?}"), + hash: item.blob_hash, + }); } Err(LibError::Cas(cas_error)) } diff --git a/src/orphan.rs b/src/orphan.rs index dcf80e0..b298198 100644 --- a/src/orphan.rs +++ b/src/orphan.rs @@ -357,10 +357,10 @@ where source: e, })?; - if let Ok(metadata) = entry.metadata() { - if metadata.is_file() { - staging_files.push(entry.path()); - } + if let Ok(metadata) = entry.metadata() + && metadata.is_file() + { + staging_files.push(entry.path()); } } diff --git a/src/wal/mod.rs b/src/wal/mod.rs index 577bf4c..20e6edf 100644 --- a/src/wal/mod.rs +++ b/src/wal/mod.rs @@ -159,15 +159,15 @@ impl WalManager { last_checkpointed_version: CheckpointState, ) -> Result<(), WalError> { // Do not go backwards; allow equal (idempotent). - if let Some(last) = last_checkpointed_version { - if version <= last { - tracing::debug!( - last_checkpointed_version = last, - commit_version = version, - "Skipping checkpoint commit (idempotent/no-op).", - ); - return Ok(()); - } + if let Some(last) = last_checkpointed_version + && version <= last + { + tracing::debug!( + last_checkpointed_version = last, + commit_version = version, + "Skipping checkpoint commit (idempotent/no-op).", + ); + return Ok(()); } // Prune segments where all operations have version <= checkpoint_version @@ -266,10 +266,10 @@ impl WalManager { impl Drop for WalManager { fn drop(&mut self) { - if let Some(writer) = self.active_writer.take() { - if let Err(e) = writer.close() { - tracing::error!("Error closing WAL segment during drop: {:?}", e); - } + if let Some(writer) = self.active_writer.take() + && let Err(e) = writer.close() + { + tracing::error!("Error closing WAL segment during drop: {:?}", e); } } } diff --git a/src/wal/storage.rs b/src/wal/storage.rs index 6cdccc2..525ac50 100644 --- a/src/wal/storage.rs +++ b/src/wal/storage.rs @@ -280,20 +280,18 @@ impl SegmentStorage { source: e, })?; let path = entry.path(); - if path.is_file() { - if let Some(filename_str) = path.file_name().and_then(|name| name.to_str()) { - if filename_str.ends_with("_index.wal") { - if let Some(id_str) = filename_str.split('_').next() { - if let Ok(id) = id_str.parse::() { - segments.push(SegmentInfo::new(id, path.clone())); - } else { - tracing::warn!( - "Found WAL-like file with non-numeric segment ID: {}", - path.display() - ); - } - } - } + if path.is_file() + && let Some(filename_str) = path.file_name().and_then(|name| name.to_str()) + && filename_str.ends_with("_index.wal") + && let Some(id_str) = filename_str.split('_').next() + { + if let Ok(id) = id_str.parse::() { + segments.push(SegmentInfo::new(id, path.clone())); + } else { + tracing::warn!( + "Found WAL-like file with non-numeric segment ID: {}", + path.display() + ); } } } diff --git a/src/wal/tests.rs b/src/wal/tests.rs index a295921..ec54143 100644 --- a/src/wal/tests.rs +++ b/src/wal/tests.rs @@ -17,10 +17,8 @@ fn perform_checkpoint( return Ok(None); }; - if seal_current_segment { - if let Some(writer) = wal.active_writer.take() { - writer.seal()?; - } + if seal_current_segment && let Some(writer) = wal.active_writer.take() { + writer.seal()?; } wal.commit_checkpoint(version, last_checkpointed_version)?;