mod codec;
mod dedupe;
mod rewrite;
mod staging;
mod sweep;
#[cfg(test)]
mod tests;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Instant;
use tokio::sync::Mutex;
use futures_util::{Stream, StreamExt};
use sha2::{Digest, Sha256};
use tokio::fs;
use tokio::io::AsyncWriteExt;
use crate::error::Error;
use crate::namespace::Namespace;
pub use dedupe::DedupeReport;
use dedupe::shares_bytes_with;
pub use rewrite::CompressReport;
enum Sink {
Raw(fs::File),
Framed(Box<codec::Writer>),
}
impl Sink {
async fn write(&mut self, chunk: &[u8]) -> Result<(), Error> {
match self {
Self::Raw(file) => Ok(file.write_all(chunk).await?),
Self::Framed(writer) => writer.push(chunk).await,
}
}
async fn finish(self) -> Result<(), Error> {
match self {
Self::Raw(mut file) => {
file.flush().await?;
Ok(file.sync_all().await?)
}
Self::Framed(writer) => writer.finish().await,
}
}
}
pub enum Object {
Raw { file: fs::File, size: u64 },
Framed(codec::Framed),
}
impl Object {
pub fn size(&self) -> u64 {
match self {
Self::Raw { size, .. } => *size,
Self::Framed(framed) => framed.plaintext(),
}
}
pub async fn stream(
self,
start: u64,
length: u64,
) -> Result<futures_util::stream::BoxStream<'static, Result<axum::body::Bytes, Error>>, Error>
{
use futures_util::StreamExt;
use tokio::io::AsyncSeekExt;
match self {
Self::Raw { mut file, .. } => {
file.seek(std::io::SeekFrom::Start(start)).await?;
let reader =
tokio_util::io::ReaderStream::new(tokio::io::AsyncReadExt::take(file, length));
Ok(reader.map(|chunk| chunk.map_err(Error::from)).boxed())
}
Self::Framed(framed) => Ok(framed.stream(start, length).boxed()),
}
}
}
pub use staging::{Reclaimed, reclaim};
pub use sweep::SweepReport;
#[derive(Debug, Clone, Copy)]
pub struct Budget {
pub used: u64,
pub limit: u64,
}
impl Budget {
pub fn exceeded_by(&self, arriving: u64) -> bool {
self.used + arriving > self.limit
}
pub fn refusal(&self) -> Error {
Error::OverQuota {
used: self.used,
limit: self.limit,
}
}
}
pub struct LocalStore {
root: PathBuf,
counter: AtomicU64,
usage: Mutex<Option<(Instant, u64, u64)>>,
per_namespace: Mutex<HashMap<String, (Instant, u64, u64)>>,
scans: AtomicU64,
max_object_size: Option<u64>,
compression: Option<i32>,
}
impl LocalStore {
pub fn new(root: impl Into<PathBuf>) -> Self {
Self {
root: root.into(),
counter: AtomicU64::new(0),
usage: Mutex::new(None),
per_namespace: Mutex::new(HashMap::new()),
scans: AtomicU64::new(0),
max_object_size: None,
compression: None,
}
}
pub fn with_compression(mut self, level: Option<i32>) -> Self {
self.compression = level;
self
}
pub fn with_max_object_size(mut self, limit: Option<u64>) -> Self {
self.max_object_size = limit;
self
}
pub fn validate_oid(oid: &str) -> Result<(), Error> {
let well_formed = oid.len() == 64
&& oid
.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b));
well_formed.then_some(()).ok_or(Error::MalformedOid)
}
fn object_path(&self, ns: &Namespace, oid: &str) -> PathBuf {
self.root
.join(ns.org())
.join(ns.repo())
.join(&oid[0..2])
.join(&oid[2..4])
.join(oid)
}
fn content_path(&self, oid: &str) -> PathBuf {
self.root
.join(".content")
.join(&oid[0..2])
.join(&oid[2..4])
.join(oid)
}
pub fn scans(&self) -> u64 {
self.scans.load(Ordering::Relaxed)
}
pub async fn writable(&self) -> Result<(), Error> {
fs::create_dir_all(&self.root).await?;
let ticket = self.counter.fetch_add(1, Ordering::Relaxed);
let probe = self.root.join(format!(".readiness.{ticket}"));
fs::write(&probe, b"").await?;
fs::remove_file(&probe).await?;
Ok(())
}
pub async fn exists(&self, ns: &Namespace, oid: &str) -> bool {
Self::validate_oid(oid).is_ok() && fs::metadata(self.object_path(ns, oid)).await.is_ok()
}
pub async fn open(&self, ns: &Namespace, oid: &str) -> Result<Object, Error> {
Self::validate_oid(oid)?;
let path = self.object_path(ns, oid);
let file = fs::File::open(&path).await.map_err(|_| Error::NotFound)?;
let on_disk = file.metadata().await?.len();
match codec::Framed::open(file, on_disk).await? {
Some(framed) => Ok(Object::Framed(framed)),
None => Ok(Object::Raw {
file: fs::File::open(&path).await.map_err(|_| Error::NotFound)?,
size: on_disk,
}),
}
}
pub async fn write<S, E>(
&self,
ns: &Namespace,
oid: &str,
expected_size: Option<u64>,
budget: Option<Budget>,
mut chunks: S,
) -> Result<u64, Error>
where
S: Stream<Item = Result<axum::body::Bytes, E>> + Unpin,
E: std::error::Error + Send + Sync + 'static,
{
Self::validate_oid(oid)?;
if let Some(limit) = self.max_object_size
&& expected_size.is_some_and(|declared| declared > limit)
{
return Err(Error::TooLarge { limit });
}
let path = self.object_path(ns, oid);
let parent = path.parent().expect("object paths always have a parent");
fs::create_dir_all(parent).await?;
let fresh = fs::metadata(&path).await.is_err();
let staged = self.staging_path(parent, oid);
let outcome = self.stream_to(&staged, budget, &mut chunks).await;
match outcome {
Ok((digest, written)) => {
self.finish(&staged, &path, oid, expected_size, &digest, written)
.await?;
if fresh {
self.stored(ns, written).await;
}
Ok(written)
}
Err(error) => {
let _ = fs::remove_file(&staged).await;
Err(error)
}
}
}
fn staging_path(&self, parent: &Path, oid: &str) -> PathBuf {
let ticket = self.counter.fetch_add(1, Ordering::Relaxed);
parent.join(format!("{oid}.{ticket}.part"))
}
async fn stream_to<S, E>(
&self,
staged: &Path,
budget: Option<Budget>,
chunks: &mut S,
) -> Result<(String, u64), Error>
where
S: Stream<Item = Result<axum::body::Bytes, E>> + Unpin,
E: std::error::Error + Send + Sync + 'static,
{
let file = fs::File::create(staged).await?;
let mut sink = match self.compression {
Some(level) => Sink::Framed(Box::new(codec::Writer::open(file, level).await?)),
None => Sink::Raw(file),
};
let mut hasher = Sha256::new();
let mut written = 0u64;
while let Some(chunk) = chunks.next().await {
let chunk = chunk.map_err(std::io::Error::other)?;
hasher.update(&chunk);
written += chunk.len() as u64;
if let Some(limit) = self.max_object_size.filter(|limit| written > *limit) {
return Err(Error::TooLarge { limit });
}
if let Some(budget) = budget.filter(|budget| budget.exceeded_by(written)) {
return Err(budget.refusal());
}
sink.write(&chunk).await?;
}
sink.finish().await?;
Ok((hex::encode(hasher.finalize()), written))
}
async fn finish(
&self,
staged: &Path,
final_path: &Path,
oid: &str,
expected_size: Option<u64>,
digest: &str,
written: u64,
) -> Result<(), Error> {
if let Some(declared) = expected_size.filter(|declared| *declared != written) {
let _ = fs::remove_file(staged).await;
return Err(Error::SizeMismatch {
declared,
actual: written,
});
}
if digest != oid {
let _ = fs::remove_file(staged).await;
return Err(Error::OidMismatch {
declared: oid.to_owned(),
actual: digest.to_owned(),
});
}
self.link_or_move(staged, final_path, oid).await
}
async fn link_or_move(&self, staged: &Path, final_path: &Path, oid: &str) -> Result<(), Error> {
let content = self.content_path(oid);
let parent = content.parent().expect("content paths have a parent");
fs::create_dir_all(parent).await?;
if fs::metadata(&content).await.is_err() {
fs::rename(staged, &content).await?;
}
match self.link(&content, final_path).await {
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
fs::rename(staged, &content).await?;
self.link(&content, final_path).await?;
}
outcome => outcome?,
}
let _ = fs::remove_file(staged).await;
Ok(())
}
async fn link(&self, content: &Path, final_path: &Path) -> Result<(), std::io::Error> {
let from = content.to_path_buf();
let to = final_path.to_path_buf();
let linked = tokio::task::spawn_blocking(move || std::fs::hard_link(&from, &to))
.await
.map_err(std::io::Error::other)?;
match linked {
Ok(()) => Ok(()),
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Err(error),
Err(_) => fs::copy(content, final_path).await.map(|_| ()),
}
}
}