1use 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
38pub struct DiskCacheStorage {
47 base_path: PathBuf,
48 num_shards: u32,
49 #[allow(dead_code)]
50 max_size_bytes: usize,
51 inflight: Arc<DashMap<String, HashSet<u64>>>,
55 next_temp_id: AtomicU64,
56}
57
58impl DiskCacheStorage {
59 pub fn new(path: &Path, shards: u32, max_size: usize) -> Self {
64 let base = path.to_path_buf();
65
66 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 for shard in 0..shards {
85 let shard_dir = base.join(format!("shard-{:02}", shard));
86
87 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 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 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 fn shard_for_key(combined: &str, num_shards: u32) -> u32 {
148 let byte = u8::from_str_radix(&combined[..2], 16).unwrap_or(0);
150 (byte as u32) % num_shards
151 }
152
153 fn prefix_for_key(combined: &str) -> &str {
155 &combined[2..4]
157 }
158
159 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 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 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
188fn 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
205fn 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
219pub 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
287pub 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 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 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 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 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 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 if let Err(e) = std::fs::rename(&tmp_body, &body_path) {
377 error!(error = %e, "Failed to rename tmp body to final path");
378 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 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 self.release_inflight();
411 }
412 }
413}
414
415#[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 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 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 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
590pub 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 let stem = match path.file_stem().and_then(|s| s.to_str()) {
627 Some(s) => s.to_string(),
628 None => continue,
629 };
630
631 let body_size = match std::fs::metadata(&path) {
633 Ok(m) => m.len() as usize,
634 Err(_) => continue,
635 };
636
637 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 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 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#[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 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 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 assert!(shard_dir.join("tmp").is_dir());
733 }
734
735 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 let combined = "abcd1234567890abcdef1234567890ab";
746
747 let shard = DiskCacheStorage::shard_for_key(combined, 16);
748 assert_eq!(shard, 0xab % 16); 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 let result = STORAGE.lookup(&key, trace).await.unwrap();
778 assert!(result.is_none());
779
780 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 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 let body2 = hit_handler.read_body().await.unwrap();
802 assert!(body2.is_none());
803
804 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 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 assert!(STORAGE.lookup(&key, trace).await.unwrap().is_some());
832
833 let compact = key.to_compact();
835 let purged = STORAGE
836 .purge(&compact, PurgeType::Invalidation, trace)
837 .await
838 .unwrap();
839 assert!(purged);
840
841 assert!(STORAGE.lookup(&key, trace).await.unwrap().is_none());
843
844 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 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 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 let updated = STORAGE.update_meta(&key, &new_meta, trace).await.unwrap();
885 assert!(updated);
886
887 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 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 {
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 }
919
920 assert!(STORAGE.lookup(&key, trace).await.unwrap().is_none());
922
923 assert!(STORAGE.inflight.is_empty());
925
926 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 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 let result = STORAGE.lookup(&key, trace).await.unwrap();
952 assert!(result.is_none());
953
954 assert!(!meta_path.exists());
956 assert!(!body_path.exists());
957
958 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 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 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 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}