use std::fs::File;
use std::io::BufReader;
use std::os::unix::fs::FileExt;
use std::path::{Path, PathBuf};
use thiserror::Error;
use crate::paths;
use crate::types::BlobHash;
#[derive(Debug, Clone)]
pub enum CasIoOperation {
ReadContent,
ReadMetadata,
OpenBuffered,
OpenRangeRead,
ReadRange,
CreateSubdir,
MoveStaged,
RemoveStaged,
RemoveFile,
}
#[derive(Error, Debug)]
pub enum CasManagerError {
#[error("CAS IO error during {operation:?} for path {path:?}")]
FileOperation {
operation: CasIoOperation,
path: PathBuf,
#[source]
source: std::io::Error,
},
#[error("Invalid range: start ({start}) > end ({end})")]
InvalidRangeStartEnd { start: u64, end: u64 },
}
pub struct CasManager {
paths: paths::DbPaths,
dir_tree_is_pre_created: bool,
}
impl CasManager {
pub fn new(paths: paths::DbPaths, dir_tree_is_pre_created: bool) -> Self {
Self { paths, dir_tree_is_pre_created }
}
pub fn read_blob(&self, blob_hash: &BlobHash) -> Result<bytes::Bytes, CasManagerError> {
let cas_path = self.paths.cas_file_path(blob_hash);
let bytes = std::fs::read(&cas_path).map_err(|e| CasManagerError::FileOperation {
operation: CasIoOperation::ReadContent,
path: cas_path.clone(),
source: e,
})?;
Ok(bytes::Bytes::from(bytes))
}
pub fn blob_bufreader(&self, blob_hash: &BlobHash) -> Result<BufReader<File>, CasManagerError> {
let cas_path = self.paths.cas_file_path(blob_hash);
let file = File::open(&cas_path).map_err(|e| CasManagerError::FileOperation {
operation: CasIoOperation::OpenBuffered,
path: cas_path.clone(),
source: e,
})?;
Ok(BufReader::new(file))
}
pub fn read_blob_range(
&self,
blob_hash: &BlobHash,
range_start: u64,
range_end: u64,
) -> Result<bytes::Bytes, CasManagerError> {
if range_start > range_end {
return Err(CasManagerError::InvalidRangeStartEnd {
start: range_start,
end: range_end,
});
}
let cas_path = self.paths.cas_file_path(blob_hash);
let file = File::open(&cas_path).map_err(|e| CasManagerError::FileOperation {
operation: CasIoOperation::OpenRangeRead,
path: cas_path.clone(),
source: e,
})?;
let read_len = range_end - range_start;
if read_len == 0 {
return Ok(bytes::Bytes::new());
}
let mut buff = Vec::with_capacity(read_len as usize);
let mut total_bytes_read = 0;
let mut current_offset = range_start;
while total_bytes_read < read_len as usize {
let spare = buff.spare_capacity_mut();
let remaining = std::cmp::min(spare.len(), read_len as usize - total_bytes_read);
let bytes_read = unsafe {
let ptr = spare.as_mut_ptr().cast::<u8>();
let slice = std::slice::from_raw_parts_mut(ptr, remaining);
file.read_at(slice, current_offset).map_err(|e| CasManagerError::FileOperation {
operation: CasIoOperation::ReadRange,
path: cas_path.clone(),
source: e,
})?
};
if bytes_read == 0 {
break;
}
unsafe {
let new_len = buff.len() + bytes_read;
buff.set_len(new_len);
}
total_bytes_read += bytes_read;
current_offset += bytes_read as u64;
}
Ok(bytes::Bytes::from(buff))
}
pub fn commit_blob(
&self,
staging_path: &Path,
blob_hash: &BlobHash,
) -> Result<PathBuf, CasManagerError> {
let final_cas_path = self.paths.cas_file_path(blob_hash);
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(()) => {
tracing::debug!(
"Moved blob {} to CAS path '{}'",
blob_hash,
final_cas_path.display()
);
Ok(final_cas_path)
}
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
std::fs::remove_file(staging_path).map_err(|e| CasManagerError::FileOperation {
operation: CasIoOperation::RemoveStaged,
path: staging_path.to_path_buf(),
source: e,
})?;
tracing::debug!(
"Blob {} already exists at '{}', removed staging",
blob_hash,
final_cas_path.display()
);
Ok(final_cas_path)
}
Err(e) => Err(CasManagerError::FileOperation {
operation: CasIoOperation::MoveStaged,
path: final_cas_path,
source: e,
}),
}
}
pub fn delete_blobs(&self, hashes: &[BlobHash]) -> Result<(), CasManagerError> {
for hash in hashes {
let file_path = self.paths.cas_file_path(hash);
match std::fs::remove_file(&file_path) {
Ok(_) => {
tracing::debug!(
"Successfully deleted unreferenced CAS file: {}",
file_path.display()
);
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
tracing::warn!(
"CAS file '{}' for unreferenced hash {} not found during deletion, skipping.",
file_path.display(),
hash
);
}
Err(e) => {
return Err(CasManagerError::FileOperation {
operation: CasIoOperation::RemoveFile,
path: file_path,
source: e,
});
}
}
}
Ok(())
}
}