use std::path::Path;
use std::path::PathBuf;
use std::time::Duration;
use std::time::SystemTime;
use std::time::UNIX_EPOCH;
use async_trait::async_trait;
use camel_api::CamelError;
use parking_lot::Mutex;
use tokio::io::AsyncWriteExt;
use tokio_util::sync::CancellationToken;
use tracing::{info, warn};
use crate::cache::offload::{PayloadStore, hasher_128hex, parse_death_epoch, sanitize_blob_name};
const TMP_NAME_ATTEMPTS: u32 = 8;
pub struct DiskPayloadStore {
dir: PathBuf,
sweep_handle: Mutex<Option<tokio::task::JoinHandle<()>>>,
}
impl DiskPayloadStore {
pub fn new(dir: PathBuf, sweep_interval: Duration, shutdown_token: CancellationToken) -> Self {
let sweep_handle = spawn_sweeper(dir.clone(), sweep_interval, shutdown_token);
Self {
dir,
sweep_handle: Mutex::new(Some(sweep_handle)),
}
}
async fn write_blob(&self, dest_name: &str, bytes: &[u8]) -> std::io::Result<()> {
tokio::fs::create_dir_all(&self.dir).await?;
let dest_path = self.dir.join(dest_name);
let (mut file, tmp_path) = self.open_tmp_exclusive(dest_name).await?;
if let Err(e) = file.write_all(bytes).await {
let _ = tokio::fs::remove_file(&tmp_path).await;
return Err(e);
}
if let Err(e) = file.sync_all().await {
let _ = tokio::fs::remove_file(&tmp_path).await;
return Err(e);
}
if let Err(e) = tokio::fs::rename(&tmp_path, &dest_path).await {
let _ = tokio::fs::remove_file(&tmp_path).await;
return Err(e);
}
self.fsync_dir_best_effort().await;
Ok(())
}
async fn open_tmp_exclusive(
&self,
dest_name: &str,
) -> std::io::Result<(tokio::fs::File, PathBuf)> {
let clock_nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
let mut last_collision: Option<std::io::Error> = None;
for attempt in 0..TMP_NAME_ATTEMPTS {
let mut hasher = blake3::Hasher::new();
hasher.update(dest_name.as_bytes());
hasher.update(&clock_nanos.to_le_bytes());
hasher.update(&attempt.to_le_bytes());
let nonce = hasher_128hex(hasher);
let tmp_path = self.dir.join(format!("{dest_name}.{nonce}.tmp"));
match tokio::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&tmp_path)
.await
{
Ok(file) => return Ok((file, tmp_path)),
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
last_collision = Some(e);
}
Err(e) => return Err(e),
}
}
Err(last_collision
.unwrap_or_else(|| std::io::Error::other("tmp blob name collisions exhausted")))
}
async fn fsync_dir_best_effort(&self) {
let result = match tokio::fs::File::open(&self.dir).await {
Ok(dir_file) => dir_file.sync_all().await,
Err(e) => Err(e),
};
if let Err(e) = result {
warn!(
dir = %self.dir.display(),
error = %e,
"cache blob directory fsync failed (best-effort, ignored)"
);
}
}
async fn clear_payload_dir_best_effort(&self) {
let mut read_dir = match tokio::fs::read_dir(&self.dir).await {
Ok(read_dir) => read_dir,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return,
Err(e) => {
warn!(
dir = %self.dir.display(),
error = %e,
"cache payload dir read failed during clear (best-effort, skipped)"
);
return;
}
};
loop {
let entry = match read_dir.next_entry().await {
Ok(Some(entry)) => entry,
Ok(None) => return,
Err(e) => {
warn!(
dir = %self.dir.display(),
error = %e,
"cache payload dir iteration failed during clear (best-effort, stopped)"
);
return;
}
};
let path = entry.path();
match tokio::fs::remove_file(&path).await {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => {
warn!(
dir = %self.dir.display(),
blob = %path.display(),
error = %e,
"cache payload blob unlink failed during clear (best-effort, skipped)"
);
}
}
}
}
}
#[async_trait]
impl PayloadStore for DiskPayloadStore {
async fn put(
&self,
name: &str,
bytes: &[u8],
_death_epoch: SystemTime,
) -> Result<(), CamelError> {
let name = sanitize_blob_name(name).ok_or_else(|| {
CamelError::Config(format!(
"cache payload name must be a bare file name, got '{name}'"
))
})?;
self.write_blob(name, bytes).await.map_err(|e| {
CamelError::Io(format!(
"cache blob write '{}': {e}",
self.dir.join(name).display()
))
})
}
async fn read(&self, name: &str) -> Result<Option<Vec<u8>>, CamelError> {
let Some(name) = sanitize_blob_name(name) else {
return Err(CamelError::Io(format!(
"cache payload name must be a bare file name, got '{name}'"
)));
};
let blob_path = self.dir.join(name);
match tokio::fs::read(&blob_path).await {
Ok(bytes) => Ok(Some(bytes)),
Err(e)
if matches!(
e.kind(),
std::io::ErrorKind::NotFound | std::io::ErrorKind::NotADirectory
) =>
{
Ok(None)
}
Err(e) => Err(CamelError::Io(format!(
"cache payload blob read '{}': {e}",
blob_path.display()
))),
}
}
async fn unlink(&self, name: &str) -> Result<(), CamelError> {
let Some(name) = sanitize_blob_name(name) else {
return Err(CamelError::Io(format!(
"cache payload name must be a bare file name, got '{name}'"
)));
};
let blob_path = self.dir.join(name);
match tokio::fs::remove_file(&blob_path).await {
Ok(()) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(e) => Err(CamelError::Io(format!(
"cache payload blob unlink '{}': {e}",
blob_path.display()
))),
}
}
async fn clear(&self) {
self.clear_payload_dir_best_effort().await;
}
}
impl std::fmt::Debug for DiskPayloadStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("DiskPayloadStore")
.field("dir", &self.dir)
.field("sweep_attached", &self.sweep_handle.lock().is_some())
.finish()
}
}
impl Drop for DiskPayloadStore {
fn drop(&mut self) {
if let Some(handle) = self.sweep_handle.lock().take() {
handle.abort();
}
}
}
async fn unlink_payload_file(
path: &Path,
now: SystemTime,
sweep_interval: Duration,
) -> std::io::Result<bool> {
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
return Ok(false);
};
let now_secs = now.duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
let dead = if name.ends_with(".blob") {
parse_death_epoch(name).is_some_and(|death| death < now_secs)
} else if name.ends_with(".tmp") {
let threshold = now.checked_sub(sweep_interval).unwrap_or(UNIX_EPOCH);
let mtime = match tokio::fs::metadata(path).await {
Ok(meta) => meta.modified()?,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(false),
Err(e) => return Err(e),
};
mtime < threshold
} else {
return Ok(false);
};
if !dead {
return Ok(false);
}
match tokio::fs::remove_file(path).await {
Ok(()) => Ok(true),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false),
Err(e) => Err(e),
}
}
#[derive(Debug, Default, PartialEq, Eq, Clone, Copy)]
struct SweepStats {
blobs_unlinked: u64,
blob_bytes_reclaimed: u64,
tmps_unlinked: u64,
live_blobs: u64,
live_blob_bytes: u64,
}
async fn sweep_payload_dir(dir: &Path, now: SystemTime, sweep_interval: Duration) -> SweepStats {
let mut read_dir = match tokio::fs::read_dir(dir).await {
Ok(read_dir) => read_dir,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return SweepStats::default(),
Err(e) => {
warn!(
dir = %dir.display(),
error = %e,
"cache payload dir read failed during sweep (skipped)"
);
return SweepStats::default();
}
};
let mut stats = SweepStats::default();
loop {
let entry = match read_dir.next_entry().await {
Ok(Some(entry)) => entry,
Ok(None) => break,
Err(e) => {
warn!(
dir = %dir.display(),
error = %e,
"cache payload dir iteration failed during sweep (stopped)"
);
break;
}
};
let path = entry.path();
let is_tmp = path
.file_name()
.and_then(|n| n.to_str())
.is_some_and(|n| n.ends_with(".tmp"));
let size = entry.metadata().await.map(|m| m.len()).unwrap_or(0);
match unlink_payload_file(&path, now, sweep_interval).await {
Ok(true) => {
if is_tmp {
stats.tmps_unlinked += 1;
} else {
stats.blobs_unlinked += 1;
stats.blob_bytes_reclaimed += size;
}
}
Ok(false) => {
if !is_tmp {
stats.live_blobs += 1;
stats.live_blob_bytes += size;
}
}
Err(e) => warn!(
dir = %dir.display(),
file = %path.display(),
error = %e,
"cache payload file unlink failed during sweep (skipped)"
),
}
}
stats
}
fn spawn_sweeper(
dir: PathBuf,
sweep_interval: Duration,
shutdown_token: CancellationToken,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut ticker = tokio::time::interval(sweep_interval);
loop {
tokio::select! {
_ = ticker.tick() => {
let s = sweep_payload_dir(&dir, SystemTime::now(), sweep_interval).await;
info!(
dir = %dir.display(),
live_blobs = s.live_blobs,
live_blob_bytes = s.live_blob_bytes,
blobs_unlinked = s.blobs_unlinked,
blob_bytes_reclaimed = s.blob_bytes_reclaimed,
tmps_unlinked = s.tmps_unlinked,
"cache payload sweep pass"
);
}
_ = shutdown_token.cancelled() => break,
}
}
})
}
#[cfg(test)]
#[path = "disk_offload_tests.rs"]
mod disk_offload_tests;
#[cfg(test)]
#[path = "disk_offload_reclaim_tests.rs"]
mod disk_offload_reclaim_tests;