use std::collections::HashMap;
use std::fmt;
use std::fs::{self, File, OpenOptions};
use std::io::{self, ErrorKind, Read as _, Seek as _, SeekFrom, Write as _};
use std::path::{Path, PathBuf};
use std::pin::Pin;
use std::sync::{Arc, Mutex, Weak};
use std::task::{Context, Poll};
use std::time::{Duration, SystemTime};
use bytes::{Bytes, BytesMut};
use futures_core::Stream;
use mkit_core::hash::{Hash, Hasher};
use mkit_transport_file::{create_dir_all_durably, sync_dir, temp_path};
use super::{io_error, unavailable};
use crate::store::{
BlobBody, BlobKey, BlobMeta, BlobStore, ByteRange, CommitOutcome, MAX_BLOB_PIECE_BYTES,
PackSink, StoreError,
};
pub(super) const READ_BLOCK: usize = 64 * 1024;
type MultipartLocks = Arc<Mutex<HashMap<[u8; 32], Weak<tokio::sync::Mutex<()>>>>>;
#[derive(Debug, Clone)]
pub struct FsBlobStore {
pub(super) root: PathBuf,
keyspace: &'static str,
pub(super) multipart_locks: MultipartLocks,
}
impl FsBlobStore {
#[must_use]
pub fn new(root: impl Into<PathBuf>) -> Self {
Self::with_keyspace(root, "packs")
}
#[must_use]
pub fn with_keyspace(root: impl Into<PathBuf>, keyspace: &'static str) -> Self {
assert!(
!keyspace.is_empty() && !keyspace.starts_with('.') && !keyspace.contains(['/', '\\']),
"a keyspace is one plain path component: {keyspace:?}"
);
assert!(
!crate::store::is_reserved_pack_keyspace(keyspace),
"a keyspace must not alias a sibling namespace: {keyspace:?}"
);
Self {
root: root.into(),
keyspace,
multipart_locks: Arc::new(Mutex::new(HashMap::new())),
}
}
#[must_use]
pub fn root(&self) -> &Path {
&self.root
}
#[must_use]
pub fn keyspace(&self) -> &'static str {
self.keyspace
}
fn dir(&self) -> PathBuf {
self.root.join(self.keyspace)
}
fn path(&self, key: &BlobKey) -> Result<PathBuf, StoreError> {
Ok(self.root.join(key.relative_path(self.keyspace)?))
}
pub fn sweep_stale_uploads(&self, min_age: Duration) -> io::Result<usize> {
let now = SystemTime::now();
let mut removed = 0;
for dir in [
self.dir(),
self.root.join("upload-markers/v1"),
self.root.join("objects"),
self.root.join("object-offsets/v1"),
] {
let entries = match fs::read_dir(dir) {
Ok(entries) => entries,
Err(e) if e.kind() == ErrorKind::NotFound => continue,
Err(e) => return Err(e),
};
for entry in entries {
let Ok(entry) = entry else { continue };
let name = entry.file_name();
if !name.to_str().is_some_and(is_upload_temp_name) {
continue;
}
let Ok(meta) = entry.metadata() else { continue };
let stale = meta.is_file()
&& meta
.modified()
.ok()
.and_then(|m| now.duration_since(m).ok())
.is_some_and(|age| age >= min_age);
if stale && fs::remove_file(entry.path()).is_ok() {
removed += 1;
}
}
}
Ok(removed + super::multipart::sweep_sessions(&self.root, now)?)
}
}
pub(super) fn is_upload_temp_name(name: &str) -> bool {
let decimal = |s: &str, max: usize| {
!s.is_empty() && s.len() <= max && s.bytes().all(|b| b.is_ascii_digit())
};
let Some(rest) = name.strip_prefix('.') else {
return false;
};
let Some((hex, rest)) = rest.split_at_checked(64) else {
return false;
};
let Some((pid, seq)) = rest.strip_prefix(".tmp.").and_then(|r| r.split_once('.')) else {
return false;
};
hex.bytes().all(|b| matches!(b, b'0'..=b'9' | b'a'..=b'f'))
&& decimal(pid, 10)
&& decimal(seq, 20)
}
pub struct FsPackSink {
file: Option<File>,
tmp: Option<PathBuf>,
dest: PathBuf,
key: BlobKey,
declared: u64,
written: u64,
hasher: Hasher,
}
impl fmt::Debug for FsPackSink {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("FsPackSink")
.field("tmp", &self.tmp)
.field("dest", &self.dest)
.field("declared", &self.declared)
.field("written", &self.written)
.finish_non_exhaustive()
}
}
impl FsPackSink {
fn publish(&mut self, root: Option<Hash>) -> Result<bool, StoreError> {
let expected = self.key.expected_root(root)?;
if self.written != self.declared {
return Err(StoreError::Invalid("blob length does not match".into()));
}
if self.hasher.finalize() != expected {
return Err(StoreError::Invalid(
"blob hash does not match its key".into(),
));
}
let file = self.file.take().ok_or_else(closed)?;
file.sync_all().map_err(io_error)?;
drop(file);
let tmp = self.tmp.as_ref().ok_or_else(closed)?;
let existed = self.dest.exists();
fs::rename(tmp, &self.dest).map_err(io_error)?;
self.tmp = None;
if let Some(dir) = self.dest.parent() {
sync_dir(dir).map_err(io_error)?;
}
Ok(existed)
}
}
fn outcome(existed: bool) -> CommitOutcome {
if existed {
CommitOutcome::AlreadyPresent
} else {
CommitOutcome::Created
}
}
fn closed() -> StoreError {
StoreError::unavailable(io::Error::other("blob upload already closed"))
}
impl Drop for FsPackSink {
fn drop(&mut self) {
self.file = None;
if let Some(tmp) = self.tmp.take() {
let _ = fs::remove_file(tmp);
}
}
}
impl PackSink for FsPackSink {
async fn write(&mut self, chunk: Bytes) -> Result<(), StoreError> {
let total = self.written.checked_add(chunk.len() as u64);
let Some(total) = total.filter(|t| *t <= self.declared) else {
return Err(StoreError::Invalid("blob is longer than declared".into()));
};
let file = self.file.as_mut().ok_or_else(closed)?;
file.write_all(&chunk).map_err(io_error)?;
self.hasher.update(&chunk);
self.written = total;
Ok(())
}
async fn commit(mut self) -> Result<CommitOutcome, StoreError> {
Ok(outcome(self.publish(None)?))
}
async fn commit_with_root(mut self, content_root: Hash) -> Result<CommitOutcome, StoreError> {
Ok(outcome(self.publish(Some(content_root))?))
}
async fn abort(self) {}
}
struct Blocks {
file: File,
remaining: u64,
}
impl Stream for Blocks {
type Item = Result<Bytes, StoreError>;
fn poll_next(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let this = self.get_mut();
if this.remaining == 0 {
return Poll::Ready(None);
}
let n = usize::try_from(this.remaining).map_or(READ_BLOCK, |r| r.min(READ_BLOCK));
let mut piece = BytesMut::zeroed(n);
match this.file.read_exact(&mut piece) {
Ok(()) => {
this.remaining -= n as u64;
Poll::Ready(Some(Ok(piece.freeze())))
}
Err(e) => {
this.remaining = 0;
Poll::Ready(Some(Err(io_error(e))))
}
}
}
}
const _: () = assert!(READ_BLOCK <= MAX_BLOB_PIECE_BYTES);
impl BlobStore for FsBlobStore {
type Sink = FsPackSink;
async fn begin(&self, key: BlobKey, len: u64) -> Result<FsPackSink, StoreError> {
let dest = self.path(&key)?;
let dir = dest
.parent()
.ok_or_else(|| StoreError::Invalid("blob path has no directory".into()))?;
create_dir_all_durably(dir).map_err(io_error)?;
let tmp = temp_path(&dest).map_err(io_error)?;
let file = OpenOptions::new()
.write(true)
.create_new(true)
.open(&tmp)
.map_err(io_error)?;
Ok(FsPackSink {
file: Some(file),
tmp: Some(tmp),
dest,
key,
declared: len,
written: 0,
hasher: Hasher::new(),
})
}
async fn get(
&self,
key: &BlobKey,
range: Option<ByteRange>,
) -> Result<Option<BlobBody>, StoreError> {
let mut file = match File::open(self.path(key)?) {
Ok(file) => file,
Err(e) if e.kind() == ErrorKind::NotFound => return Ok(None),
Err(e) => return Err(io_error(e)),
};
let len = file.metadata().map_err(io_error)?.len();
let span = match range {
Some(range) => range.resolve(len)?,
None => 0..len,
};
if span.start > 0 {
file.seek(SeekFrom::Start(span.start)).map_err(io_error)?;
}
let n = span.end - span.start;
if let Ok(whole) = usize::try_from(n)
&& whole <= MAX_BLOB_PIECE_BYTES
{
let mut buf = vec![0; whole];
file.read_exact(&mut buf).map_err(io_error)?;
return Ok(Some(BlobBody::Bytes(Bytes::from(buf))));
}
Ok(Some(BlobBody::Stream {
len: n,
stream: Box::pin(Blocks { file, remaining: n }),
}))
}
async fn head(&self, key: &BlobKey) -> Result<Option<BlobMeta>, StoreError> {
match fs::metadata(self.path(key)?) {
Ok(meta) => Ok(Some(BlobMeta { len: meta.len() })),
Err(e) if e.kind() == ErrorKind::NotFound => Ok(None),
Err(e) => Err(io_error(e)),
}
}
async fn probe(&self) -> Result<(), StoreError> {
let meta = fs::metadata(&self.root).map_err(io_error)?;
if meta.is_dir() {
Ok(())
} else {
Err(unavailable(io::Error::other(
"blob root is not a directory",
)))
}
}
async fn delete(&self, key: &BlobKey) -> Result<bool, StoreError> {
let path = self.path(key)?;
match fs::remove_file(&path) {
Ok(()) => {}
Err(e) if e.kind() == ErrorKind::NotFound => return Ok(false),
Err(e) => return Err(io_error(e)),
}
if let Some(dir) = path.parent() {
sync_dir(dir).map_err(io_error)?;
}
Ok(true)
}
}