Skip to main content

zentinel_proxy/
disk_cache.rs

1//! Disk-based cache storage backend
2//!
3//! Implements Pingora's `Storage` trait using the local filesystem. Each cached
4//! response is stored as a pair of files (`.meta` + `.body`) distributed across
5//! sharded subdirectories to keep per-directory inode counts manageable.
6//!
7//! # Directory layout
8//!
9//! ```text
10//! <base_path>/
11//!   shard-00/
12//!     <2-char-hex-prefix>/
13//!       <combined-hex-hash>.meta
14//!       <combined-hex-hash>.body
15//!     tmp/
16//!   shard-01/
17//!     ...
18//! ```
19
20use async_trait::async_trait;
21use bytes::Bytes;
22use dashmap::DashMap;
23use pingora_cache::eviction::EvictionManager;
24use pingora_cache::key::{CacheHashKey, CacheKey, CompactCacheKey};
25use pingora_cache::meta::CacheMeta;
26use pingora_cache::storage::{
27    HandleHit, HandleMiss, HitHandler, MissFinishType, MissHandler, PurgeType, Storage,
28};
29use pingora_cache::trace::SpanHandle;
30use pingora_core::{Error, ErrorType, Result};
31use std::any::Any;
32use std::collections::HashSet;
33use std::path::{Path, PathBuf};
34use std::sync::atomic::{AtomicU64, Ordering};
35use std::sync::Arc;
36use tracing::{debug, error, info, warn};
37
38// ============================================================================
39// DiskCacheStorage
40// ============================================================================
41
42/// Disk-based cache storage backend implementing Pingora's `Storage` trait.
43///
44/// All disk I/O is performed via `tokio::task::spawn_blocking` to avoid
45/// blocking the async runtime.
46pub struct DiskCacheStorage {
47    base_path: PathBuf,
48    num_shards: u32,
49    #[allow(dead_code)]
50    max_size_bytes: usize,
51    /// Tracks in-flight writes: combined_hash -> set of temp_ids.
52    /// Lock-free so `DiskMissHandler::drop` can always clean up its entry
53    /// (an async lock's `try_write` could silently leak under contention).
54    inflight: Arc<DashMap<String, HashSet<u64>>>,
55    next_temp_id: AtomicU64,
56}
57
58impl DiskCacheStorage {
59    /// Create a new `DiskCacheStorage`.
60    ///
61    /// Creates the shard directory structure and cleans up any orphaned `.tmp`
62    /// files left behind by interrupted writes.
63    pub fn new(path: &Path, shards: u32, max_size: usize) -> Self {
64        let base = path.to_path_buf();
65
66        // Probe writability once up front so a misconfigured cache directory
67        // produces a single actionable error instead of one per subdirectory
68        // (and later, one per failed cache write).
69        if let Err(e) = std::fs::create_dir_all(&base).and_then(|()| {
70            let probe = base.join(".zentinel-write-probe");
71            std::fs::write(&probe, b"probe")?;
72            std::fs::remove_file(&probe)
73        }) {
74            error!(
75                path = %base.display(),
76                error = %e,
77                "Disk cache directory is not writable; cache writes WILL fail \
78                 (requests still proxied, uncached). Fix permissions or change \
79                 cache disk-path."
80            );
81        }
82
83        // Create shard dirs, hex-prefix subdirs, and tmp dirs
84        for shard in 0..shards {
85            let shard_dir = base.join(format!("shard-{:02}", shard));
86
87            // Create all 256 hex-prefix subdirs
88            for prefix in 0..=255u8 {
89                let prefix_dir = shard_dir.join(format!("{:02x}", prefix));
90                if let Err(e) = std::fs::create_dir_all(&prefix_dir) {
91                    error!(path = %prefix_dir.display(), error = %e, "Failed to create prefix dir");
92                }
93            }
94
95            // Create tmp dir and clean orphaned files
96            let tmp_dir = shard_dir.join("tmp");
97            if let Err(e) = std::fs::create_dir_all(&tmp_dir) {
98                error!(path = %tmp_dir.display(), error = %e, "Failed to create tmp dir");
99            } else {
100                Self::clean_orphaned_tmp(&tmp_dir);
101            }
102        }
103
104        info!(
105            path = %base.display(),
106            shards,
107            max_size_mb = max_size / 1024 / 1024,
108            "Disk cache storage initialized"
109        );
110
111        Self {
112            base_path: base,
113            num_shards: shards,
114            max_size_bytes: max_size,
115            inflight: Arc::new(DashMap::new()),
116            next_temp_id: AtomicU64::new(1),
117        }
118    }
119
120    /// Remove orphaned .tmp files from a tmp directory.
121    fn clean_orphaned_tmp(tmp_dir: &Path) {
122        let entries = match std::fs::read_dir(tmp_dir) {
123            Ok(e) => e,
124            Err(_) => return,
125        };
126        let mut cleaned = 0u64;
127        for entry in entries.flatten() {
128            let path = entry.path();
129            if path.extension().and_then(|e| e.to_str()) == Some("tmp") {
130                if let Err(e) = std::fs::remove_file(&path) {
131                    warn!(path = %path.display(), error = %e, "Failed to clean orphaned tmp file");
132                } else {
133                    cleaned += 1;
134                }
135            }
136        }
137        if cleaned > 0 {
138            info!(dir = %tmp_dir.display(), cleaned, "Cleaned orphaned tmp files");
139        }
140    }
141
142    // ========================================================================
143    // Path helpers
144    // ========================================================================
145
146    /// Compute the shard index for a combined hex hash string.
147    fn shard_for_key(combined: &str, num_shards: u32) -> u32 {
148        // Use first two hex chars (one byte) to determine shard
149        let byte = u8::from_str_radix(&combined[..2], 16).unwrap_or(0);
150        (byte as u32) % num_shards
151    }
152
153    /// Compute the 2-char hex prefix subdirectory for a combined hex hash.
154    fn prefix_for_key(combined: &str) -> &str {
155        // Use chars 2..4 (second byte) as prefix subdir
156        &combined[2..4]
157    }
158
159    /// Full path to the `.meta` file for a given combined hex hash.
160    fn meta_path(&self, combined: &str) -> PathBuf {
161        let shard = Self::shard_for_key(combined, self.num_shards);
162        let prefix = Self::prefix_for_key(combined);
163        self.base_path
164            .join(format!("shard-{:02}", shard))
165            .join(prefix)
166            .join(format!("{}.meta", combined))
167    }
168
169    /// Full path to the `.body` file for a given combined hex hash.
170    fn body_path(&self, combined: &str) -> PathBuf {
171        let shard = Self::shard_for_key(combined, self.num_shards);
172        let prefix = Self::prefix_for_key(combined);
173        self.base_path
174            .join(format!("shard-{:02}", shard))
175            .join(prefix)
176            .join(format!("{}.body", combined))
177    }
178
179    /// Path to the tmp directory for a given combined hex hash (shard-local).
180    fn tmp_dir_for_key(&self, combined: &str) -> PathBuf {
181        let shard = Self::shard_for_key(combined, self.num_shards);
182        self.base_path
183            .join(format!("shard-{:02}", shard))
184            .join("tmp")
185    }
186}
187
188// ============================================================================
189// Meta file serialization helpers
190// ============================================================================
191
192/// Serialize CacheMeta to the on-disk format.
193///
194/// Format: `[4 bytes: internal_meta_len as u32 LE][internal_meta bytes][header bytes]`
195fn serialize_meta_to_disk(meta: &CacheMeta) -> Result<Vec<u8>> {
196    let (internal, header) = meta.serialize()?;
197    let internal_len = internal.len() as u32;
198    let mut buf = Vec::with_capacity(4 + internal.len() + header.len());
199    buf.extend_from_slice(&internal_len.to_le_bytes());
200    buf.extend_from_slice(&internal);
201    buf.extend_from_slice(&header);
202    Ok(buf)
203}
204
205/// Deserialize CacheMeta from the on-disk format.
206fn deserialize_meta_from_disk(data: &[u8]) -> Result<CacheMeta> {
207    if data.len() < 4 {
208        return Error::e_explain(ErrorType::FileReadError, "meta file too short");
209    }
210    let internal_len = u32::from_le_bytes([data[0], data[1], data[2], data[3]]) as usize;
211    if data.len() < 4 + internal_len {
212        return Error::e_explain(ErrorType::FileReadError, "meta file truncated");
213    }
214    let internal = &data[4..4 + internal_len];
215    let header = &data[4 + internal_len..];
216    CacheMeta::deserialize(internal, header)
217}
218
219// ============================================================================
220// DiskHitHandler
221// ============================================================================
222
223/// Hit handler for disk cache lookups.
224///
225/// Loads the full body into memory for seekable access.
226pub struct DiskHitHandler {
227    body: Vec<u8>,
228    meta_size: usize,
229    done: bool,
230    range_start: usize,
231    range_end: usize,
232}
233
234#[async_trait]
235impl HandleHit for DiskHitHandler {
236    async fn read_body(&mut self) -> Result<Option<Bytes>> {
237        if self.done {
238            return Ok(None);
239        }
240        self.done = true;
241        Ok(Some(Bytes::copy_from_slice(
242            &self.body[self.range_start..self.range_end],
243        )))
244    }
245
246    async fn finish(
247        self: Box<Self>,
248        _storage: &'static (dyn Storage + Sync),
249        _key: &CacheKey,
250        _trace: &SpanHandle,
251    ) -> Result<()> {
252        Ok(())
253    }
254
255    fn can_seek(&self) -> bool {
256        true
257    }
258
259    fn seek(&mut self, start: usize, end: Option<usize>) -> Result<()> {
260        if start >= self.body.len() {
261            return Error::e_explain(
262                ErrorType::InternalError,
263                format!("seek start out of range {} >= {}", start, self.body.len()),
264            );
265        }
266        self.range_start = start;
267        if let Some(end) = end {
268            self.range_end = std::cmp::min(self.body.len(), end);
269        }
270        self.done = false;
271        Ok(())
272    }
273
274    fn get_eviction_weight(&self) -> usize {
275        self.meta_size + self.body.len()
276    }
277
278    fn as_any(&self) -> &(dyn Any + Send + Sync) {
279        self
280    }
281
282    fn as_any_mut(&mut self) -> &mut (dyn Any + Send + Sync) {
283        self
284    }
285}
286
287// ============================================================================
288// DiskMissHandler
289// ============================================================================
290
291/// Miss handler for disk cache writes.
292///
293/// Accumulates the response body in memory, then atomically writes both
294/// `.meta` and `.body` files to disk via temp-file + rename.
295pub struct DiskMissHandler {
296    body_buffer: Vec<u8>,
297    serialized_meta: Vec<u8>,
298    combined: String,
299    meta_path: PathBuf,
300    body_path: PathBuf,
301    tmp_dir: PathBuf,
302    temp_id: u64,
303    inflight: Arc<DashMap<String, HashSet<u64>>>,
304    finished: bool,
305}
306
307impl DiskMissHandler {
308    /// Remove this handler's temp_id from inflight tracking.
309    fn release_inflight(&self) {
310        if let Some(mut set) = self.inflight.get_mut(&self.combined) {
311            set.remove(&self.temp_id);
312            if set.is_empty() {
313                drop(set);
314                self.inflight.remove_if(&self.combined, |_, s| s.is_empty());
315            }
316        }
317    }
318}
319
320#[async_trait]
321impl HandleMiss for DiskMissHandler {
322    async fn write_body(&mut self, data: Bytes, _eof: bool) -> Result<()> {
323        self.body_buffer.extend_from_slice(&data);
324        Ok(())
325    }
326
327    async fn finish(mut self: Box<Self>) -> Result<MissFinishType> {
328        self.finished = true;
329        let body = std::mem::take(&mut self.body_buffer);
330        let meta = self.serialized_meta.clone();
331        let meta_path = self.meta_path.clone();
332        let body_path = self.body_path.clone();
333        let tmp_dir = self.tmp_dir.clone();
334        let temp_id = self.temp_id;
335
336        let size = meta.len() + body.len();
337
338        // Write to disk via spawn_blocking
339        tokio::task::spawn_blocking(move || {
340            let tmp_meta = tmp_dir.join(format!("{}.meta.tmp", temp_id));
341            let tmp_body = tmp_dir.join(format!("{}.body.tmp", temp_id));
342
343            // Write meta temp file
344            if let Err(e) = std::fs::write(&tmp_meta, &meta) {
345                error!(path = %tmp_meta.display(), error = %e, "Failed to write tmp meta");
346                let _ = std::fs::remove_file(&tmp_meta);
347                return Err(Error::explain(
348                    ErrorType::WriteError,
349                    format!("failed to write meta: {}", e),
350                ));
351            }
352
353            // Write body temp file
354            if let Err(e) = std::fs::write(&tmp_body, &body) {
355                error!(path = %tmp_body.display(), error = %e, "Failed to write tmp body");
356                let _ = std::fs::remove_file(&tmp_meta);
357                let _ = std::fs::remove_file(&tmp_body);
358                return Err(Error::explain(
359                    ErrorType::WriteError,
360                    format!("failed to write body: {}", e),
361                ));
362            }
363
364            // Atomic rename meta
365            if let Err(e) = std::fs::rename(&tmp_meta, &meta_path) {
366                error!(error = %e, "Failed to rename tmp meta to final path");
367                let _ = std::fs::remove_file(&tmp_meta);
368                let _ = std::fs::remove_file(&tmp_body);
369                return Err(Error::explain(
370                    ErrorType::WriteError,
371                    format!("failed to rename meta: {}", e),
372                ));
373            }
374
375            // Atomic rename body
376            if let Err(e) = std::fs::rename(&tmp_body, &body_path) {
377                error!(error = %e, "Failed to rename tmp body to final path");
378                // Meta already renamed; remove it to stay consistent
379                let _ = std::fs::remove_file(&meta_path);
380                return Err(Error::explain(
381                    ErrorType::WriteError,
382                    format!("failed to rename body: {}", e),
383                ));
384            }
385
386            Ok(())
387        })
388        .await
389        .map_err(|e| {
390            Error::explain(
391                ErrorType::InternalError,
392                format!("spawn_blocking join error: {}", e),
393            )
394        })??;
395
396        // Remove from inflight tracking
397        self.release_inflight();
398
399        debug!(combined = %self.combined, size, "Disk cache entry written");
400        Ok(MissFinishType::Created(size))
401    }
402}
403
404impl Drop for DiskMissHandler {
405    fn drop(&mut self) {
406        if !self.finished {
407            // Clean up inflight tracking if finish() was never called
408            // (aborted/cancelled write). DashMap removal is lock-free, so
409            // unlike an async lock this can never silently skip cleanup.
410            self.release_inflight();
411        }
412    }
413}
414
415// ============================================================================
416// Storage trait implementation
417// ============================================================================
418
419#[async_trait]
420impl Storage for DiskCacheStorage {
421    async fn lookup(
422        &'static self,
423        key: &CacheKey,
424        _trace: &SpanHandle,
425    ) -> Result<Option<(CacheMeta, HitHandler)>> {
426        let combined = key.combined();
427        let meta_path = self.meta_path(&combined);
428        let body_path = self.body_path(&combined);
429
430        let result =
431            tokio::task::spawn_blocking(move || -> Result<Option<(CacheMeta, HitHandler)>> {
432                let meta_data = match std::fs::read(&meta_path) {
433                    Ok(d) => d,
434                    Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
435                    Err(e) => {
436                        debug!(error = %e, "Failed to read cache meta file");
437                        return Ok(None);
438                    }
439                };
440
441                let body_data = match std::fs::read(&body_path) {
442                    Ok(d) => d,
443                    Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
444                    Err(e) => {
445                        debug!(error = %e, "Failed to read cache body file");
446                        return Ok(None);
447                    }
448                };
449
450                let meta = match deserialize_meta_from_disk(&meta_data) {
451                    Ok(m) => m,
452                    Err(e) => {
453                        warn!(error = %e, "Corrupted cache meta, removing entry");
454                        let _ = std::fs::remove_file(&meta_path);
455                        let _ = std::fs::remove_file(&body_path);
456                        return Ok(None);
457                    }
458                };
459
460                let body_len = body_data.len();
461                let hit_handler = DiskHitHandler {
462                    body: body_data,
463                    meta_size: meta_data.len(),
464                    done: false,
465                    range_start: 0,
466                    range_end: body_len,
467                };
468
469                Ok(Some((meta, Box::new(hit_handler) as HitHandler)))
470            })
471            .await
472            .map_err(|e| {
473                Error::explain(
474                    ErrorType::InternalError,
475                    format!("spawn_blocking join error: {}", e),
476                )
477            })??;
478
479        Ok(result)
480    }
481
482    async fn get_miss_handler(
483        &'static self,
484        key: &CacheKey,
485        meta: &CacheMeta,
486        _trace: &SpanHandle,
487    ) -> Result<MissHandler> {
488        let combined = key.combined();
489        let serialized_meta = serialize_meta_to_disk(meta)?;
490        let meta_path = self.meta_path(&combined);
491        let body_path = self.body_path(&combined);
492        let tmp_dir = self.tmp_dir_for_key(&combined);
493        let temp_id = self.next_temp_id.fetch_add(1, Ordering::Relaxed);
494
495        // Register in inflight tracking
496        self.inflight
497            .entry(combined.clone())
498            .or_default()
499            .insert(temp_id);
500
501        Ok(Box::new(DiskMissHandler {
502            body_buffer: Vec::new(),
503            serialized_meta,
504            combined,
505            meta_path,
506            body_path,
507            tmp_dir,
508            temp_id,
509            inflight: self.inflight.clone(),
510            finished: false,
511        }))
512    }
513
514    async fn purge(
515        &'static self,
516        key: &CompactCacheKey,
517        _purge_type: PurgeType,
518        _trace: &SpanHandle,
519    ) -> Result<bool> {
520        let combined = key.combined();
521        let meta_path = self.meta_path(&combined);
522        let body_path = self.body_path(&combined);
523
524        let removed = tokio::task::spawn_blocking(move || {
525            let meta_removed = std::fs::remove_file(&meta_path).is_ok();
526            let body_removed = std::fs::remove_file(&body_path).is_ok();
527            meta_removed || body_removed
528        })
529        .await
530        .map_err(|e| {
531            Error::explain(
532                ErrorType::InternalError,
533                format!("spawn_blocking join error: {}", e),
534            )
535        })?;
536
537        // Also remove from inflight tracking
538        self.inflight.remove(&combined);
539
540        Ok(removed)
541    }
542
543    async fn update_meta(
544        &'static self,
545        key: &CacheKey,
546        meta: &CacheMeta,
547        _trace: &SpanHandle,
548    ) -> Result<bool> {
549        let combined = key.combined();
550        let serialized = serialize_meta_to_disk(meta)?;
551        let meta_path = self.meta_path(&combined);
552        let tmp_dir = self.tmp_dir_for_key(&combined);
553
554        tokio::task::spawn_blocking(move || {
555            // Atomic rewrite: write to tmp, rename over existing
556            let tmp_path = tmp_dir.join(format!("{}.meta.update.tmp", combined));
557            std::fs::write(&tmp_path, &serialized).map_err(|e| {
558                Error::explain(
559                    ErrorType::WriteError,
560                    format!("failed to write updated meta: {}", e),
561                )
562            })?;
563            std::fs::rename(&tmp_path, &meta_path).map_err(|e| {
564                let _ = std::fs::remove_file(&tmp_path);
565                Error::explain(
566                    ErrorType::WriteError,
567                    format!("failed to rename updated meta: {}", e),
568                )
569            })?;
570            Ok(true)
571        })
572        .await
573        .map_err(|e| {
574            Error::explain(
575                ErrorType::InternalError,
576                format!("spawn_blocking join error: {}", e),
577            )
578        })?
579    }
580
581    fn support_streaming_partial_write(&self) -> bool {
582        false
583    }
584
585    fn as_any(&self) -> &(dyn Any + Send + Sync + 'static) {
586        self
587    }
588}
589
590// ============================================================================
591// Eviction state rebuild
592// ============================================================================
593
594/// Scan disk entries and register them with the eviction manager.
595///
596/// This is called at startup to rebuild the LRU eviction state from the
597/// files on disk.
598pub async fn rebuild_eviction_state(
599    base_path: &Path,
600    num_shards: u32,
601    eviction: &'static pingora_cache::eviction::simple_lru::Manager,
602) {
603    let base = base_path.to_path_buf();
604    let result = tokio::task::spawn_blocking(move || {
605        let mut count = 0usize;
606        let mut total_size = 0usize;
607
608        for shard in 0..num_shards {
609            let shard_dir = base.join(format!("shard-{:02}", shard));
610
611            for prefix in 0..=255u8 {
612                let prefix_dir = shard_dir.join(format!("{:02x}", prefix));
613                let entries = match std::fs::read_dir(&prefix_dir) {
614                    Ok(e) => e,
615                    Err(_) => continue,
616                };
617
618                for entry in entries.flatten() {
619                    let path = entry.path();
620                    let ext = path.extension().and_then(|e| e.to_str());
621                    if ext != Some("body") {
622                        continue;
623                    }
624
625                    // Extract combined hash from filename
626                    let stem = match path.file_stem().and_then(|s| s.to_str()) {
627                        Some(s) => s.to_string(),
628                        None => continue,
629                    };
630
631                    // Get file size for weight
632                    let body_size = match std::fs::metadata(&path) {
633                        Ok(m) => m.len() as usize,
634                        Err(_) => continue,
635                    };
636
637                    // Also add meta size
638                    let meta_path = prefix_dir.join(format!("{}.meta", stem));
639                    let meta_size = std::fs::metadata(&meta_path)
640                        .map(|m| m.len() as usize)
641                        .unwrap_or(0);
642
643                    let size = body_size + meta_size;
644
645                    // Reconstruct CompactCacheKey from the combined hex hash
646                    if let Some(primary) = pingora_cache::key::str2hex(&stem) {
647                        let compact = CompactCacheKey {
648                            primary,
649                            variance: None,
650                            user_tag: "".into(),
651                        };
652
653                        // Admit to eviction manager (use epoch as fresh_until since
654                        // we don't know the actual TTL without parsing meta)
655                        let _ = eviction.admit(
656                            compact,
657                            size,
658                            std::time::SystemTime::now() + std::time::Duration::from_secs(3600),
659                        );
660
661                        count += 1;
662                        total_size += size;
663                    }
664                }
665            }
666        }
667
668        (count, total_size)
669    })
670    .await;
671
672    match result {
673        Ok((count, total_size)) => {
674            info!(
675                entries = count,
676                total_size_mb = total_size / 1024 / 1024,
677                "Rebuilt disk cache eviction state"
678            );
679        }
680        Err(e) => {
681            error!(error = %e, "Failed to rebuild disk cache eviction state");
682        }
683    }
684}
685
686// ============================================================================
687// Tests
688// ============================================================================
689
690#[cfg(test)]
691mod tests {
692    use super::*;
693    use once_cell::sync::Lazy;
694    use pingora_cache::trace::Span;
695    use pingora_http::ResponseHeader;
696    use std::time::SystemTime;
697    use tempfile::TempDir;
698
699    fn create_test_meta() -> CacheMeta {
700        let mut header = ResponseHeader::build(200, None).unwrap();
701        header.append_header("content-type", "text/plain").unwrap();
702        header.append_header("x-test", "disk-cache").unwrap();
703        CacheMeta::new(
704            SystemTime::now() + std::time::Duration::from_secs(3600),
705            SystemTime::now(),
706            60,
707            300,
708            header,
709        )
710    }
711
712    fn span() -> SpanHandle {
713        Span::inactive().handle()
714    }
715
716    #[test]
717    fn test_directory_creation() {
718        let tmp = TempDir::new().unwrap();
719        let _storage = DiskCacheStorage::new(tmp.path(), 4, 100 * 1024 * 1024);
720
721        // Verify shard dirs exist
722        for shard in 0..4u32 {
723            let shard_dir = tmp.path().join(format!("shard-{:02}", shard));
724            assert!(shard_dir.is_dir(), "shard dir should exist");
725
726            // Verify some prefix dirs exist
727            assert!(shard_dir.join("00").is_dir());
728            assert!(shard_dir.join("ff").is_dir());
729            assert!(shard_dir.join("a5").is_dir());
730
731            // Verify tmp dir exists
732            assert!(shard_dir.join("tmp").is_dir());
733        }
734
735        // Shard-04 should not exist
736        assert!(!tmp.path().join("shard-04").exists());
737    }
738
739    #[test]
740    fn test_path_helpers() {
741        let tmp = TempDir::new().unwrap();
742        let storage = DiskCacheStorage::new(tmp.path(), 16, 100 * 1024 * 1024);
743
744        // "ab" prefix -> shard = 0xab % 16 = 11, prefix subdir = second byte
745        let combined = "abcd1234567890abcdef1234567890ab";
746
747        let shard = DiskCacheStorage::shard_for_key(combined, 16);
748        assert_eq!(shard, 0xab % 16); // 171 % 16 = 11
749
750        let prefix = DiskCacheStorage::prefix_for_key(combined);
751        assert_eq!(prefix, "cd");
752
753        let meta = storage.meta_path(combined);
754        assert!(meta.to_str().unwrap().contains("shard-11"));
755        assert!(meta.to_str().unwrap().contains("/cd/"));
756        assert!(meta.to_str().unwrap().ends_with(".meta"));
757
758        let body = storage.body_path(combined);
759        assert!(body.to_str().unwrap().contains("shard-11"));
760        assert!(body.to_str().unwrap().contains("/cd/"));
761        assert!(body.to_str().unwrap().ends_with(".body"));
762    }
763
764    #[tokio::test]
765    async fn test_write_and_read() {
766        static STORAGE: Lazy<DiskCacheStorage> = Lazy::new(|| {
767            let path = std::env::temp_dir().join("zentinel-disk-cache-test-write-read");
768            let _ = std::fs::remove_dir_all(&path);
769            DiskCacheStorage::new(&path, 4, 100 * 1024 * 1024)
770        });
771        let trace = &span();
772
773        let key = CacheKey::new("", "test-write-read", "1");
774        let meta = create_test_meta();
775
776        // Lookup should return None initially
777        let result = STORAGE.lookup(&key, trace).await.unwrap();
778        assert!(result.is_none());
779
780        // Write via miss handler
781        let mut miss_handler = STORAGE.get_miss_handler(&key, &meta, trace).await.unwrap();
782        miss_handler
783            .write_body(b"hello "[..].into(), false)
784            .await
785            .unwrap();
786        miss_handler
787            .write_body(b"world"[..].into(), true)
788            .await
789            .unwrap();
790        let finish_result = miss_handler.finish().await.unwrap();
791        assert!(matches!(finish_result, MissFinishType::Created(_)));
792
793        // Lookup should now return the cached entry
794        let (read_meta, mut hit_handler) = STORAGE.lookup(&key, trace).await.unwrap().unwrap();
795        assert_eq!(read_meta.response_header().status.as_u16(), 200);
796
797        let body = hit_handler.read_body().await.unwrap().unwrap();
798        assert_eq!(body.as_ref(), b"hello world");
799
800        // Second read should return None
801        let body2 = hit_handler.read_body().await.unwrap();
802        assert!(body2.is_none());
803
804        // Cleanup
805        let _ = std::fs::remove_dir_all(
806            std::env::temp_dir().join("zentinel-disk-cache-test-write-read"),
807        );
808    }
809
810    #[tokio::test]
811    async fn test_purge() {
812        static STORAGE: Lazy<DiskCacheStorage> = Lazy::new(|| {
813            let path = std::env::temp_dir().join("zentinel-disk-cache-test-purge");
814            let _ = std::fs::remove_dir_all(&path);
815            DiskCacheStorage::new(&path, 4, 100 * 1024 * 1024)
816        });
817        let trace = &span();
818
819        let key = CacheKey::new("", "test-purge", "1");
820        let meta = create_test_meta();
821
822        // Write an entry
823        let mut miss_handler = STORAGE.get_miss_handler(&key, &meta, trace).await.unwrap();
824        miss_handler
825            .write_body(b"purge-me"[..].into(), true)
826            .await
827            .unwrap();
828        miss_handler.finish().await.unwrap();
829
830        // Verify it's there
831        assert!(STORAGE.lookup(&key, trace).await.unwrap().is_some());
832
833        // Purge it
834        let compact = key.to_compact();
835        let purged = STORAGE
836            .purge(&compact, PurgeType::Invalidation, trace)
837            .await
838            .unwrap();
839        assert!(purged);
840
841        // Verify it's gone
842        assert!(STORAGE.lookup(&key, trace).await.unwrap().is_none());
843
844        // Cleanup
845        let _ =
846            std::fs::remove_dir_all(std::env::temp_dir().join("zentinel-disk-cache-test-purge"));
847    }
848
849    #[tokio::test]
850    async fn test_update_meta() {
851        static STORAGE: Lazy<DiskCacheStorage> = Lazy::new(|| {
852            let path = std::env::temp_dir().join("zentinel-disk-cache-test-update-meta");
853            let _ = std::fs::remove_dir_all(&path);
854            DiskCacheStorage::new(&path, 4, 100 * 1024 * 1024)
855        });
856        let trace = &span();
857
858        let key = CacheKey::new("", "test-update-meta", "1");
859        let meta = create_test_meta();
860
861        // Write an entry
862        let mut miss_handler = STORAGE.get_miss_handler(&key, &meta, trace).await.unwrap();
863        miss_handler
864            .write_body(b"body-data"[..].into(), true)
865            .await
866            .unwrap();
867        miss_handler.finish().await.unwrap();
868
869        // Create updated meta with different header
870        let mut new_header = ResponseHeader::build(200, None).unwrap();
871        new_header
872            .append_header("content-type", "application/json")
873            .unwrap();
874        new_header.append_header("x-updated", "true").unwrap();
875        let new_meta = CacheMeta::new(
876            SystemTime::now() + std::time::Duration::from_secs(7200),
877            SystemTime::now(),
878            120,
879            600,
880            new_header,
881        );
882
883        // Update meta
884        let updated = STORAGE.update_meta(&key, &new_meta, trace).await.unwrap();
885        assert!(updated);
886
887        // Verify updated meta
888        let (read_meta, _hit) = STORAGE.lookup(&key, trace).await.unwrap().unwrap();
889        let headers = read_meta.response_header().headers.clone();
890        assert_eq!(headers.get("x-updated").unwrap().to_str().unwrap(), "true");
891
892        // Cleanup
893        let _ = std::fs::remove_dir_all(
894            std::env::temp_dir().join("zentinel-disk-cache-test-update-meta"),
895        );
896    }
897
898    #[tokio::test]
899    async fn test_miss_handler_drop() {
900        static STORAGE: Lazy<DiskCacheStorage> = Lazy::new(|| {
901            let path = std::env::temp_dir().join("zentinel-disk-cache-test-miss-drop");
902            let _ = std::fs::remove_dir_all(&path);
903            DiskCacheStorage::new(&path, 4, 100 * 1024 * 1024)
904        });
905        let trace = &span();
906
907        let key = CacheKey::new("", "test-miss-drop", "1");
908        let meta = create_test_meta();
909
910        // Create miss handler and write some data but don't finish
911        {
912            let mut miss_handler = STORAGE.get_miss_handler(&key, &meta, trace).await.unwrap();
913            miss_handler
914                .write_body(b"incomplete"[..].into(), false)
915                .await
916                .unwrap();
917            // Drop without finish
918        }
919
920        // Verify no files were written
921        assert!(STORAGE.lookup(&key, trace).await.unwrap().is_none());
922
923        // Verify inflight tracking was cleaned up
924        assert!(STORAGE.inflight.is_empty());
925
926        // Cleanup
927        let _ = std::fs::remove_dir_all(
928            std::env::temp_dir().join("zentinel-disk-cache-test-miss-drop"),
929        );
930    }
931
932    #[tokio::test]
933    async fn test_corrupted_meta() {
934        static STORAGE: Lazy<DiskCacheStorage> = Lazy::new(|| {
935            let path = std::env::temp_dir().join("zentinel-disk-cache-test-corrupted");
936            let _ = std::fs::remove_dir_all(&path);
937            DiskCacheStorage::new(&path, 4, 100 * 1024 * 1024)
938        });
939        let trace = &span();
940
941        let key = CacheKey::new("", "test-corrupted", "1");
942        let combined = key.combined();
943
944        // Write garbage to the meta file
945        let meta_path = STORAGE.meta_path(&combined);
946        let body_path = STORAGE.body_path(&combined);
947        std::fs::write(&meta_path, b"not-valid-meta-data").unwrap();
948        std::fs::write(&body_path, b"some-body").unwrap();
949
950        // Lookup should gracefully return None
951        let result = STORAGE.lookup(&key, trace).await.unwrap();
952        assert!(result.is_none());
953
954        // Corrupted files should have been cleaned up
955        assert!(!meta_path.exists());
956        assert!(!body_path.exists());
957
958        // Cleanup
959        let _ = std::fs::remove_dir_all(
960            std::env::temp_dir().join("zentinel-disk-cache-test-corrupted"),
961        );
962    }
963
964    #[test]
965    fn test_orphan_cleanup() {
966        let tmp = TempDir::new().unwrap();
967
968        // Pre-create a shard with tmp dir and orphaned files
969        let shard_tmp = tmp.path().join("shard-00").join("tmp");
970        std::fs::create_dir_all(&shard_tmp).unwrap();
971        std::fs::write(shard_tmp.join("orphan1.tmp"), b"data1").unwrap();
972        std::fs::write(shard_tmp.join("orphan2.tmp"), b"data2").unwrap();
973        // Non-tmp file should be left alone
974        std::fs::write(shard_tmp.join("keep.txt"), b"keep").unwrap();
975
976        assert!(shard_tmp.join("orphan1.tmp").exists());
977        assert!(shard_tmp.join("orphan2.tmp").exists());
978
979        // Creating storage should clean orphaned .tmp files
980        let _storage = DiskCacheStorage::new(tmp.path(), 4, 100 * 1024 * 1024);
981
982        assert!(!shard_tmp.join("orphan1.tmp").exists());
983        assert!(!shard_tmp.join("orphan2.tmp").exists());
984        assert!(shard_tmp.join("keep.txt").exists());
985    }
986
987    #[test]
988    fn test_meta_serialization_roundtrip() {
989        let meta = create_test_meta();
990        let serialized = serialize_meta_to_disk(&meta).unwrap();
991        let deserialized = deserialize_meta_from_disk(&serialized).unwrap();
992
993        assert_eq!(
994            meta.response_header().status.as_u16(),
995            deserialized.response_header().status.as_u16(),
996        );
997    }
998}