Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
16 changes: 8 additions & 8 deletions src/cas_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -139,14 +139,14 @@ impl CasManager {
blob_hash: &BlobHash,
) -> Result<PathBuf, CasManagerError> {
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(()) => {
Expand Down
16 changes: 8 additions & 8 deletions src/index/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
}
Expand Down
8 changes: 4 additions & 4 deletions src/index/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
}
Expand Down
206 changes: 76 additions & 130 deletions src/io.rs
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -45,77 +43,6 @@ pub(crate) fn atomically_write_file_bytes(
Ok(())
}

pub(crate) trait FileExt {
fn lock(self) -> io::Result<LockedFile>;
}

impl FileExt for File {
#[inline]
fn lock(self) -> io::Result<LockedFile> {
LockedFile::new(self)
}
}

#[derive(Debug)]
#[repr(transparent)]
pub(crate) struct LockedFile {
inner: File,
}

impl LockedFile {
pub fn new(file: File) -> std::io::Result<Self> {
// 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<File> {
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,
Expand All @@ -138,56 +65,25 @@ pub enum IoError {

#[cfg(target_os = "linux")]
pub fn is_same_fs(paths: &[&Path]) -> Result<bool, io::Error> {
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<Option<u64>> {
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::<libc::statx>::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")))]
Expand All @@ -213,43 +109,93 @@ pub fn same_fs_dev(paths: &[&Path]) -> Result<bool, io::Error> {
Ok(true)
}

#[cfg(not(unix))]
fn is_same_fs(path: &[&Path]) -> Result<bool, IsSameFsError> {
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<Option<u64>> {
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::<libc::statx>::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<String> = Cas::open(dir.path(), Default::default()).unwrap();

let file = options.open(&path).unwrap().lock().unwrap();
let cas2: Result<Cas<String>, _> = Cas::<String>::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(())
}
}
24 changes: 10 additions & 14 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -194,8 +192,7 @@ where
}

pub struct CasInner<K> {
#[allow(unused)]
lockfile: LockedFile,
_lockfile: File,
pub(crate) paths: paths::DbPaths,
pub(crate) index: Index<K>,
datasync_channel: Option<std::sync::mpsc::Sender<File>>,
Expand Down Expand Up @@ -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());
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -434,13 +431,12 @@ where
Err(cas_error) => {
if let Some(io_err) =
cas_error.source().and_then(|s| s.downcast_ref::<std::io::Error>())
&& 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))
}
Expand Down
Loading