use std::{
fs::{self, File},
io::{Error as IoError, ErrorKind, Read, Seek, SeekFrom},
path::{Path, PathBuf},
};
use shardline_protocol::ByteRange;
use thiserror::Error;
use crate::{
DeleteOutcome, ObjectBody, ObjectIntegrity, ObjectKey, ObjectMetadata, ObjectPrefix,
ObjectStore, PutOutcome,
local_fs::{
PutBytesIfAbsentOutcome, hard_link_file_if_absent, put_bytes_if_absent,
write_bytes_atomically,
},
};
use super::{
io::{
self, ensure_file_matches_bytes, ensure_files_match, open_existing_object_file,
verify_file_integrity,
},
metadata::{self, ensure_regular_file_metadata, object_file_metadata},
util::verify_integrity,
walk::{self, read_dir_if_exists},
};
#[derive(Debug, Clone)]
pub struct LocalObjectStore {
root: PathBuf,
}
impl LocalObjectStore {
#[must_use]
pub const fn open(root: PathBuf) -> Self {
Self { root }
}
pub fn new(root: PathBuf) -> Result<Self, LocalObjectStoreError> {
let store = Self::open(root);
metadata::ensure_parent_directories_are_not_symlinked(&store.root, &store.root)?;
fs::create_dir_all(&store.root)?;
Ok(store)
}
#[must_use]
pub fn root(&self) -> &Path {
&self.root
}
#[must_use]
pub fn path_for_key(&self, key: &ObjectKey) -> PathBuf {
self.key_path(key)
}
pub fn open_object_file(&self, key: &ObjectKey) -> Result<File, LocalObjectStoreError> {
let path = self.key_path(key);
metadata::ensure_parent_directories_are_not_symlinked(&self.root, &path)?;
open_existing_object_file(&path)
}
pub fn put_temporary_file_if_absent(
&self,
key: &ObjectKey,
temporary: &Path,
integrity: &ObjectIntegrity,
) -> Result<PutOutcome, LocalObjectStoreError> {
verify_file_integrity(temporary, integrity)?;
let path = self.key_path(key);
metadata::ensure_parent_directories_are_not_symlinked(&self.root, &path)?;
io::link_temporary_file_if_absent(&self.root, &path, temporary, integrity, None)
}
pub fn copy_object_if_absent(
&self,
source: &ObjectKey,
destination: &ObjectKey,
) -> Result<PutOutcome, LocalObjectStoreError> {
let source_path = self.key_path(source);
let destination_path = self.key_path(destination);
let _source = open_existing_object_file(&source_path)?;
metadata::ensure_parent_directories_are_not_symlinked(&self.root, &destination_path)?;
match hard_link_file_if_absent(&self.root, &destination_path, &source_path) {
Ok(()) => Ok(PutOutcome::Inserted),
Err(error) if error.kind() == ErrorKind::AlreadyExists => {
let existing = open_existing_object_file(&destination_path)?;
let source_file = open_existing_object_file(&source_path)?;
ensure_files_match(existing, source_file)?;
Ok(PutOutcome::AlreadyExists)
}
Err(error) => Err(LocalObjectStoreError::Io(error)),
}
}
pub fn put_overwrite(
&self,
key: &ObjectKey,
body: ObjectBody<'_>,
integrity: &ObjectIntegrity,
) -> Result<(), LocalObjectStoreError> {
let bytes = body.into_bytes();
verify_integrity(&bytes, integrity)?;
let path = self.key_path(key);
metadata::ensure_parent_directories_are_not_symlinked(&self.root, &path)?;
write_bytes_atomically(&self.root, &path, &bytes).map_err(LocalObjectStoreError::Io)
}
pub fn list_flat_namespace_page(
&self,
prefix: &ObjectPrefix,
start_after: Option<&ObjectKey>,
limit: usize,
) -> Result<Vec<ObjectMetadata>, LocalObjectStoreError> {
let prefix_str = prefix.as_str();
if let Some(start_after) = start_after
&& !start_after.as_str().starts_with(prefix_str)
{
return Err(LocalObjectStoreError::InvalidStartAfter);
}
let directory = self.root.join(prefix_str);
let Some(entries) = read_dir_if_exists(&directory)? else {
return Ok(Vec::new());
};
let mut children = Vec::new();
for entry in entries {
let entry = entry.map_err(LocalObjectStoreError::Io)?;
let file_type = entry.file_type().map_err(LocalObjectStoreError::Io)?;
if !file_type.is_file() {
continue;
}
let name = entry
.file_name()
.into_string()
.map_err(|_error| LocalObjectStoreError::InvalidStoredKey)?;
let key = ObjectKey::parse(&format!("{}{}", prefix.as_str(), name))
.map_err(|_error| LocalObjectStoreError::InvalidStoredKey)?;
if start_after.is_some_and(|offset| key.as_str() <= offset.as_str()) {
continue;
}
let metadata = fs::symlink_metadata(entry.path()).map_err(LocalObjectStoreError::Io)?;
ensure_regular_file_metadata(&metadata)?;
children.push(ObjectMetadata::new(key, metadata.len(), None));
}
children.sort_by(|left, right| left.key().as_str().cmp(right.key().as_str()));
children.truncate(limit);
Ok(children)
}
fn key_path(&self, key: &ObjectKey) -> PathBuf {
self.root.join(key.as_str())
}
}
impl ObjectStore for LocalObjectStore {
type Error = LocalObjectStoreError;
fn put_if_absent(
&self,
key: &ObjectKey,
body: ObjectBody<'_>,
integrity: &ObjectIntegrity,
) -> Result<PutOutcome, Self::Error> {
verify_integrity(body.as_slice(), integrity)?;
let path = self.key_path(key);
metadata::ensure_parent_directories_are_not_symlinked(&self.root, &path)?;
match open_existing_object_file(&path) {
Ok(file) => {
ensure_file_matches_bytes(file, body.as_slice())?;
return Ok(PutOutcome::AlreadyExists);
}
Err(LocalObjectStoreError::Io(error)) if error.kind() == ErrorKind::NotFound => {}
Err(error) => return Err(error),
}
match put_bytes_if_absent(&self.root, &path, body.as_slice())
.map_err(LocalObjectStoreError::Io)?
{
PutBytesIfAbsentOutcome::Inserted => Ok(PutOutcome::Inserted),
PutBytesIfAbsentOutcome::AlreadyExists => Ok(PutOutcome::AlreadyExists),
}
}
fn read_range(&self, key: &ObjectKey, range: ByteRange) -> Result<Vec<u8>, Self::Error> {
let path = self.key_path(key);
metadata::ensure_parent_directories_are_not_symlinked(&self.root, &path)?;
let mut file = open_existing_object_file(&path)?;
let Some(length_u64) = range.len() else {
return Err(LocalObjectStoreError::RangeOutOfBounds);
};
file.seek(SeekFrom::Start(range.start()))?;
let capacity = usize::try_from(length_u64)
.map_err(|_error| LocalObjectStoreError::RangeOutOfBounds)?;
let mut output = vec![0_u8; capacity];
if let Err(error) = file.read_exact(&mut output) {
if error.kind() == ErrorKind::UnexpectedEof {
return Err(LocalObjectStoreError::RangeOutOfBounds);
}
return Err(LocalObjectStoreError::Io(error));
}
Ok(output)
}
fn contains(&self, key: &ObjectKey) -> Result<bool, Self::Error> {
self.metadata(key).map(|metadata| metadata.is_some())
}
fn metadata(&self, key: &ObjectKey) -> Result<Option<ObjectMetadata>, Self::Error> {
let path = self.key_path(key);
metadata::ensure_parent_directories_are_not_symlinked(&self.root, &path)?;
let Some(metadata) = object_file_metadata(&path)? else {
return Ok(None);
};
Ok(Some(ObjectMetadata::new(key.clone(), metadata.len(), None)))
}
fn list_prefix(&self, prefix: &ObjectPrefix) -> Result<Vec<ObjectMetadata>, Self::Error> {
let mut listed = Vec::new();
let mut visitor = |metadata| {
listed.push(metadata);
Ok::<(), LocalObjectStoreError>(())
};
walk::collect_metadata_recursive(&self.root, &self.root, prefix.as_str(), &mut visitor)?;
listed.sort_by(|left, right| left.key().as_str().cmp(right.key().as_str()));
Ok(listed)
}
fn visit_prefix<Visitor, VisitorError>(
&self,
prefix: &ObjectPrefix,
mut visitor: Visitor,
) -> Result<(), VisitorError>
where
Self::Error: Into<VisitorError>,
Visitor: FnMut(ObjectMetadata) -> Result<(), VisitorError>,
{
walk::collect_metadata_recursive(&self.root, &self.root, prefix.as_str(), &mut visitor)
}
fn delete_if_present(&self, key: &ObjectKey) -> Result<DeleteOutcome, Self::Error> {
let path = self.key_path(key);
metadata::ensure_parent_directories_are_not_symlinked(&self.root, &path)?;
match fs::remove_file(&path) {
Ok(()) => {
walk::remove_empty_ancestors(&path, &self.root)?;
Ok(DeleteOutcome::Deleted)
}
Err(error) if error.kind() == ErrorKind::NotFound => Ok(DeleteOutcome::NotFound),
Err(error) => Err(LocalObjectStoreError::Io(error)),
}
}
}
#[derive(Debug, Error)]
pub enum LocalObjectStoreError {
#[error("local object store operation failed")]
Io(#[from] IoError),
#[error("object body length did not match expected integrity")]
IntegrityLengthMismatch,
#[error("object body hash did not match expected integrity")]
IntegrityHashMismatch,
#[error("object key already exists with conflicting bytes")]
ExistingObjectConflict,
#[error("requested byte range exceeded stored object length")]
RangeOutOfBounds,
#[error("stored object path could not be represented as a valid object key")]
InvalidStoredKey,
#[error("validated object key could not be mapped to a local path")]
InvalidObjectPath,
#[error("start_after key is outside the requested prefix")]
InvalidStartAfter,
}