Skip to main content

camel_core/cache/
disk_offload.rs

1//! Disk payload store for the [`crate::cache::offload::OffloadRepository`]
2//! decorator.
3//!
4//! [`DiskPayloadStore`] holds content-addressed payload blob files under a
5//! dedicated directory, addressed by the blob names the decorator
6//! computes. A background sweeper reclaims dead blobs by their
7//! name-encoded death epoch and stale `.tmp` leftovers by age, exiting
8//! when the context-owned shutdown token fires.
9//!
10//! # Blob lifecycle
11//!
12//! Blob names are `{blake3-128hex(key)}.{death_epoch_secs}.{blake3-128hex(
13//! bytes || content_type-discriminant)}.blob`. The death epoch is encoded
14//! in the name so the sweeper can reclaim dead blobs by file name alone,
15//! without consulting the index.
16//!
17//! # Failure policy
18//!
19//! - Writes are tmp-then-rename so a partially written blob is never
20//!   visible under its final name; write failures surface to the
21//!   decorator, whose inline fallback applies.
22//! - A vanished blob reads as `Ok(None)`; a blob that exists but cannot
23//!   be read (e.g. `PermissionDenied`) surfaces as `Err` per ADR-0023
24//!   Contract C1.
25//! - `unlink` counts an already-absent blob as success; `clear` never
26//!   fails — per-blob unlink failures WARN and the scan continues.
27
28use 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
43/// Max attempts to open a unique tmp file before giving up on a name.
44const TMP_NAME_ATTEMPTS: u32 = 8;
45
46/// [`PayloadStore`] that keeps offloaded payload blobs as files under
47/// `dir`, with a background death-epoch sweeper.
48///
49/// The sweeper is spawned on construction and aborted on [`Drop`]; its
50/// shutdown token is owned by the sweeper task and never cancelled by
51/// the store. All sweeping runs on the real clock — it must observe
52/// actual file ages (the decorator's injectable test clock never
53/// reaches the store).
54pub struct DiskPayloadStore {
55    /// Directory holding offloaded payload blobs.
56    dir: PathBuf,
57    /// Background payload sweeper; aborted on Drop.
58    sweep_handle: Mutex<Option<tokio::task::JoinHandle<()>>>,
59}
60
61impl DiskPayloadStore {
62    /// Create a store writing blobs into `dir` (production clock).
63    ///
64    /// The `shutdown_token` stops the background payload sweeper; it is
65    /// owned by the sweeper task and never cancelled by the store.
66    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    /// Write `bytes` to their content-addressed blob file `dest_name`.
75    ///
76    /// Tmp-then-rename so a partially written blob is never visible under
77    /// its final name. All I/O is async `tokio::fs` (the house file-I/O
78    /// style, matching `camel-file`'s `atomic_write`).
79    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        // Best-effort tmp cleanup on failure: a leaked `.tmp` would never
85        // be reclaimed by the epoch sweeper.
86        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    /// Open a unique exclusive tmp file next to `dest_name`, retrying name
103    /// collisions with a fresh nonce (bounded by [`TMP_NAME_ATTEMPTS`]).
104    ///
105    /// The nonce hashes `dest_name || clock_nanos || attempt_counter`, so
106    /// retries still produce fresh names under a frozen clock.
107    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    /// Best-effort fsync of the blob directory so the rename itself is
141    /// durable. Failures are WARNed and ignored: the blob is already
142    /// renamed, and a directory-fsync failure must not fail the write.
143    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    /// Best-effort unlink of every entry of the payload dir.
158    ///
159    /// Per-file `NotFound` is success (a concurrent sweeper or replica may
160    /// have reclaimed the blob already); any other per-file error WARNs
161    /// and iteration continues. [`Self::clear`] must never surface its
162    /// own unlink failures as `Err`.
163    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            // No dir = nothing was ever offloaded; nothing to unlink.
167            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                // NotFound = a concurrent sweeper or replica won the race.
194                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        // Trait contract (ADR-0065): a payload name is a bare single path
217        // component. The same sanitize guard `read`/`unlink` apply keeps a
218        // separator or `..` from ever escaping the payload dir through a
219        // write.
220        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        // The death epoch is already encoded in the blob name; the
226        // filename-epoch sweeper owns reclamation, so the deadline needs
227        // no separate storage on the disk tier.
228        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            // NotFound = a concurrent sweeper or replica won the race.
270            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    /// Reclaim payload space now: best-effort unlink of every entry of
279    /// the payload dir. Unlink failures never turn `clear` into `Err` —
280    /// each failure WARNs and the rest of the dir is still attempted.
281    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        // Abort ONLY the sweep task. Never cancel the context-owned token —
298        // that would shut down the entire context when one store drops.
299        if let Some(handle) = self.sweep_handle.lock().take() {
300            handle.abort();
301        }
302    }
303}
304
305// ── Payload sweeper ─────────────────────────────────────────────────────────
306
307/// Unlink one payload-dir file if it is dead: `.blob` files by their
308/// name-encoded death epoch, `.tmp` leftovers by age.
309///
310/// `Ok(true)` = unlinked here; `Ok(false)` = kept (still live, a foreign
311/// name without a parseable epoch, or vanished between listing and unlink —
312/// the ENOENT race counts as reclaimed-by-someone-else, never an error).
313/// Any other error is returned for the sweep loop to WARN over. All filesystem
314/// access is async (`tokio::fs`), keeping the sweeper off blocked
315/// runtime workers.
316async 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    // Clamp a pre-epoch clock to the Unix epoch, matching `set`'s
325    // death-epoch math.
326    let now_secs = now.duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
327    let dead = if name.ends_with(".blob") {
328        // Strictly-before: a blob dying exactly `now` survives this pass
329        // (the filename epoch is whole seconds; the next tick reclaims).
330        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/// One sweep pass over `dir`: reclaim dead blobs by their filename
353/// death epoch and stale `.tmp` leftovers by age.
354///
355/// Per-file `NotFound` (a concurrent sweeper or replica won the race)
356/// counts as success; other per-file errors WARN and the scan
357/// continues. A missing dir is not an error — nothing was ever
358/// offloaded. Returns `(blobs_unlinked, tmps_unlinked)`.
359#[derive(Debug, Default, PartialEq, Eq, Clone, Copy)]
360struct SweepStats {
361    /// Dead blobs unlinked this pass.
362    blobs_unlinked: u64,
363    /// Bytes reclaimed with those dead blobs.
364    blob_bytes_reclaimed: u64,
365    /// Stale tmp files unlinked this pass.
366    tmps_unlinked: u64,
367    /// Blobs still on disk after the pass (live, orphan pre-epoch, or
368    /// foreign names — anything the sweep kept; a blob that vanishes
369    /// mid-pass via the ENOENT race is counted here until the next pass).
370    live_blobs: u64,
371    /// Total bytes of those surviving blobs.
372    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
434/// Spawn the background payload sweeper for `dir`.
435///
436/// Mirrors the redb sweep loop: tick every `sweep_interval`, reclaim
437/// dead blobs and stale tmp files, exit when `shutdown_token` fires.
438/// The sweep always runs on the REAL clock (`SystemTime::now`), never
439/// an injected decorator clock — it must observe actual file ages.
440fn 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                    // Per-pass volume observability (bd rc-h3dp): live
452                    // bytes are the high-water baseline operators compare
453                    // against the eager-reclaim trigger; reclaimed bytes
454                    // show the pass's cleanup.
455                    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;