1use std::path::Path;
29use std::path::PathBuf;
30use std::time::Duration;
31use std::time::SystemTime;
32use std::time::UNIX_EPOCH;
33
34use async_trait::async_trait;
35use camel_api::CamelError;
36use parking_lot::Mutex;
37use tokio::io::AsyncWriteExt;
38use tokio_util::sync::CancellationToken;
39use tracing::{info, warn};
40
41use crate::cache::offload::{PayloadStore, hasher_128hex, parse_death_epoch, sanitize_blob_name};
42
43const TMP_NAME_ATTEMPTS: u32 = 8;
45
46pub struct DiskPayloadStore {
55 dir: PathBuf,
57 sweep_handle: Mutex<Option<tokio::task::JoinHandle<()>>>,
59}
60
61impl DiskPayloadStore {
62 pub fn new(dir: PathBuf, sweep_interval: Duration, shutdown_token: CancellationToken) -> Self {
67 let sweep_handle = spawn_sweeper(dir.clone(), sweep_interval, shutdown_token);
68 Self {
69 dir,
70 sweep_handle: Mutex::new(Some(sweep_handle)),
71 }
72 }
73
74 async fn write_blob(&self, dest_name: &str, bytes: &[u8]) -> std::io::Result<()> {
80 tokio::fs::create_dir_all(&self.dir).await?;
81 let dest_path = self.dir.join(dest_name);
82 let (mut file, tmp_path) = self.open_tmp_exclusive(dest_name).await?;
83
84 if let Err(e) = file.write_all(bytes).await {
87 let _ = tokio::fs::remove_file(&tmp_path).await;
88 return Err(e);
89 }
90 if let Err(e) = file.sync_all().await {
91 let _ = tokio::fs::remove_file(&tmp_path).await;
92 return Err(e);
93 }
94 if let Err(e) = tokio::fs::rename(&tmp_path, &dest_path).await {
95 let _ = tokio::fs::remove_file(&tmp_path).await;
96 return Err(e);
97 }
98 self.fsync_dir_best_effort().await;
99 Ok(())
100 }
101
102 async fn open_tmp_exclusive(
108 &self,
109 dest_name: &str,
110 ) -> std::io::Result<(tokio::fs::File, PathBuf)> {
111 let clock_nanos = SystemTime::now()
112 .duration_since(UNIX_EPOCH)
113 .map(|d| d.as_nanos())
114 .unwrap_or(0);
115 let mut last_collision: Option<std::io::Error> = None;
116 for attempt in 0..TMP_NAME_ATTEMPTS {
117 let mut hasher = blake3::Hasher::new();
118 hasher.update(dest_name.as_bytes());
119 hasher.update(&clock_nanos.to_le_bytes());
120 hasher.update(&attempt.to_le_bytes());
121 let nonce = hasher_128hex(hasher);
122 let tmp_path = self.dir.join(format!("{dest_name}.{nonce}.tmp"));
123 match tokio::fs::OpenOptions::new()
124 .write(true)
125 .create_new(true)
126 .open(&tmp_path)
127 .await
128 {
129 Ok(file) => return Ok((file, tmp_path)),
130 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
131 last_collision = Some(e);
132 }
133 Err(e) => return Err(e),
134 }
135 }
136 Err(last_collision
137 .unwrap_or_else(|| std::io::Error::other("tmp blob name collisions exhausted")))
138 }
139
140 async fn fsync_dir_best_effort(&self) {
144 let result = match tokio::fs::File::open(&self.dir).await {
145 Ok(dir_file) => dir_file.sync_all().await,
146 Err(e) => Err(e),
147 };
148 if let Err(e) = result {
149 warn!(
150 dir = %self.dir.display(),
151 error = %e,
152 "cache blob directory fsync failed (best-effort, ignored)"
153 );
154 }
155 }
156
157 async fn clear_payload_dir_best_effort(&self) {
164 let mut read_dir = match tokio::fs::read_dir(&self.dir).await {
165 Ok(read_dir) => read_dir,
166 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return,
168 Err(e) => {
169 warn!(
170 dir = %self.dir.display(),
171 error = %e,
172 "cache payload dir read failed during clear (best-effort, skipped)"
173 );
174 return;
175 }
176 };
177 loop {
178 let entry = match read_dir.next_entry().await {
179 Ok(Some(entry)) => entry,
180 Ok(None) => return,
181 Err(e) => {
182 warn!(
183 dir = %self.dir.display(),
184 error = %e,
185 "cache payload dir iteration failed during clear (best-effort, stopped)"
186 );
187 return;
188 }
189 };
190 let path = entry.path();
191 match tokio::fs::remove_file(&path).await {
192 Ok(()) => {}
193 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
195 Err(e) => {
196 warn!(
197 dir = %self.dir.display(),
198 blob = %path.display(),
199 error = %e,
200 "cache payload blob unlink failed during clear (best-effort, skipped)"
201 );
202 }
203 }
204 }
205 }
206}
207
208#[async_trait]
209impl PayloadStore for DiskPayloadStore {
210 async fn put(
211 &self,
212 name: &str,
213 bytes: &[u8],
214 _death_epoch: SystemTime,
215 ) -> Result<(), CamelError> {
216 let name = sanitize_blob_name(name).ok_or_else(|| {
221 CamelError::Config(format!(
222 "cache payload name must be a bare file name, got '{name}'"
223 ))
224 })?;
225 self.write_blob(name, bytes).await.map_err(|e| {
229 CamelError::Io(format!(
230 "cache blob write '{}': {e}",
231 self.dir.join(name).display()
232 ))
233 })
234 }
235
236 async fn read(&self, name: &str) -> Result<Option<Vec<u8>>, CamelError> {
237 let Some(name) = sanitize_blob_name(name) else {
238 return Err(CamelError::Io(format!(
239 "cache payload name must be a bare file name, got '{name}'"
240 )));
241 };
242 let blob_path = self.dir.join(name);
243 match tokio::fs::read(&blob_path).await {
244 Ok(bytes) => Ok(Some(bytes)),
245 Err(e)
246 if matches!(
247 e.kind(),
248 std::io::ErrorKind::NotFound | std::io::ErrorKind::NotADirectory
249 ) =>
250 {
251 Ok(None)
252 }
253 Err(e) => Err(CamelError::Io(format!(
254 "cache payload blob read '{}': {e}",
255 blob_path.display()
256 ))),
257 }
258 }
259
260 async fn unlink(&self, name: &str) -> Result<(), CamelError> {
261 let Some(name) = sanitize_blob_name(name) else {
262 return Err(CamelError::Io(format!(
263 "cache payload name must be a bare file name, got '{name}'"
264 )));
265 };
266 let blob_path = self.dir.join(name);
267 match tokio::fs::remove_file(&blob_path).await {
268 Ok(()) => Ok(()),
269 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
271 Err(e) => Err(CamelError::Io(format!(
272 "cache payload blob unlink '{}': {e}",
273 blob_path.display()
274 ))),
275 }
276 }
277
278 async fn clear(&self) {
282 self.clear_payload_dir_best_effort().await;
283 }
284}
285
286impl std::fmt::Debug for DiskPayloadStore {
287 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
288 f.debug_struct("DiskPayloadStore")
289 .field("dir", &self.dir)
290 .field("sweep_attached", &self.sweep_handle.lock().is_some())
291 .finish()
292 }
293}
294
295impl Drop for DiskPayloadStore {
296 fn drop(&mut self) {
297 if let Some(handle) = self.sweep_handle.lock().take() {
300 handle.abort();
301 }
302 }
303}
304
305async fn unlink_payload_file(
317 path: &Path,
318 now: SystemTime,
319 sweep_interval: Duration,
320) -> std::io::Result<bool> {
321 let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
322 return Ok(false);
323 };
324 let now_secs = now.duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
327 let dead = if name.ends_with(".blob") {
328 parse_death_epoch(name).is_some_and(|death| death < now_secs)
331 } else if name.ends_with(".tmp") {
332 let threshold = now.checked_sub(sweep_interval).unwrap_or(UNIX_EPOCH);
333 let mtime = match tokio::fs::metadata(path).await {
334 Ok(meta) => meta.modified()?,
335 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(false),
336 Err(e) => return Err(e),
337 };
338 mtime < threshold
339 } else {
340 return Ok(false);
341 };
342 if !dead {
343 return Ok(false);
344 }
345 match tokio::fs::remove_file(path).await {
346 Ok(()) => Ok(true),
347 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false),
348 Err(e) => Err(e),
349 }
350}
351
352#[derive(Debug, Default, PartialEq, Eq, Clone, Copy)]
360struct SweepStats {
361 blobs_unlinked: u64,
363 blob_bytes_reclaimed: u64,
365 tmps_unlinked: u64,
367 live_blobs: u64,
371 live_blob_bytes: u64,
373}
374
375async fn sweep_payload_dir(dir: &Path, now: SystemTime, sweep_interval: Duration) -> SweepStats {
376 let mut read_dir = match tokio::fs::read_dir(dir).await {
377 Ok(read_dir) => read_dir,
378 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return SweepStats::default(),
379 Err(e) => {
380 warn!(
381 dir = %dir.display(),
382 error = %e,
383 "cache payload dir read failed during sweep (skipped)"
384 );
385 return SweepStats::default();
386 }
387 };
388 let mut stats = SweepStats::default();
389 loop {
390 let entry = match read_dir.next_entry().await {
391 Ok(Some(entry)) => entry,
392 Ok(None) => break,
393 Err(e) => {
394 warn!(
395 dir = %dir.display(),
396 error = %e,
397 "cache payload dir iteration failed during sweep (stopped)"
398 );
399 break;
400 }
401 };
402 let path = entry.path();
403 let is_tmp = path
404 .file_name()
405 .and_then(|n| n.to_str())
406 .is_some_and(|n| n.ends_with(".tmp"));
407 let size = entry.metadata().await.map(|m| m.len()).unwrap_or(0);
408 match unlink_payload_file(&path, now, sweep_interval).await {
409 Ok(true) => {
410 if is_tmp {
411 stats.tmps_unlinked += 1;
412 } else {
413 stats.blobs_unlinked += 1;
414 stats.blob_bytes_reclaimed += size;
415 }
416 }
417 Ok(false) => {
418 if !is_tmp {
419 stats.live_blobs += 1;
420 stats.live_blob_bytes += size;
421 }
422 }
423 Err(e) => warn!(
424 dir = %dir.display(),
425 file = %path.display(),
426 error = %e,
427 "cache payload file unlink failed during sweep (skipped)"
428 ),
429 }
430 }
431 stats
432}
433
434fn spawn_sweeper(
441 dir: PathBuf,
442 sweep_interval: Duration,
443 shutdown_token: CancellationToken,
444) -> tokio::task::JoinHandle<()> {
445 tokio::spawn(async move {
446 let mut ticker = tokio::time::interval(sweep_interval);
447 loop {
448 tokio::select! {
449 _ = ticker.tick() => {
450 let s = sweep_payload_dir(&dir, SystemTime::now(), sweep_interval).await;
451 info!(
456 dir = %dir.display(),
457 live_blobs = s.live_blobs,
458 live_blob_bytes = s.live_blob_bytes,
459 blobs_unlinked = s.blobs_unlinked,
460 blob_bytes_reclaimed = s.blob_bytes_reclaimed,
461 tmps_unlinked = s.tmps_unlinked,
462 "cache payload sweep pass"
463 );
464 }
465 _ = shutdown_token.cancelled() => break,
466 }
467 }
468 })
469}
470
471#[cfg(test)]
472#[path = "disk_offload_tests.rs"]
473mod disk_offload_tests;
474
475#[cfg(test)]
476#[path = "disk_offload_reclaim_tests.rs"]
477mod disk_offload_reclaim_tests;