1use crate::{
4 CacheKey, CachePolicy, CacheStorage, PutHandle, StoredEntry, fs_shims, policy::PolicyRepr,
5};
6use futures_lite::{AsyncRead, AsyncWrite, AsyncWriteExt};
7use moka::{notification::RemovalCause, sync::Cache};
8use sha2::{Digest, Sha256};
9use std::{
10 fmt::{self, Debug, Formatter, Write as _},
11 io,
12 path::{Path, PathBuf},
13 pin::Pin,
14 sync::{
15 Arc,
16 atomic::{AtomicU64, Ordering},
17 },
18 task::{Context, Poll},
19 time::Duration,
20};
21use trillium_http::{Body, BodySource, Headers};
22
23const META_SUFFIX: &str = ".meta";
24const BODY_SUFFIX: &str = ".body";
25
26const DEFAULT_MAX_CAPACITY_BYTES: u64 = 1024 * 1024 * 1024;
29
30static TEMP_COUNTER: AtomicU64 = AtomicU64::new(0);
33
34#[derive(Clone)]
106pub struct FileSystemStorage {
107 root: Arc<PathBuf>,
108 index: Cache<VariantId, u64>,
109 max_capacity_bytes: Option<u64>,
110 time_to_idle: Option<Duration>,
111 time_to_live: Option<Duration>,
112}
113
114impl Debug for FileSystemStorage {
115 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
116 f.debug_struct("FileSystemStorage")
117 .field("root", &self.root)
118 .field("weighted_size", &self.index.weighted_size())
119 .field("max_capacity_bytes", &self.max_capacity_bytes)
120 .field("time_to_idle", &self.time_to_idle)
121 .field("time_to_live", &self.time_to_live)
122 .finish()
123 }
124}
125
126impl FileSystemStorage {
127 pub fn new(root: impl Into<PathBuf>) -> Self {
131 let root = Arc::new(root.into());
132 let max_capacity_bytes = Some(DEFAULT_MAX_CAPACITY_BYTES);
133 let index = build_index(Arc::clone(&root), max_capacity_bytes, None, None);
134 scan_root(&root, &index);
135 Self {
136 root,
137 index,
138 max_capacity_bytes,
139 time_to_idle: None,
140 time_to_live: None,
141 }
142 }
143
144 pub fn with_max_capacity_bytes(mut self, bytes: u64) -> Self {
148 self.max_capacity_bytes = Some(bytes);
149 self.rebuild();
150 self
151 }
152
153 pub fn unbounded(mut self) -> Self {
157 self.max_capacity_bytes = None;
158 self.rebuild();
159 self
160 }
161
162 pub fn with_time_to_idle(mut self, duration: Duration) -> Self {
176 self.time_to_idle = Some(duration);
177 self.rebuild();
178 self
179 }
180
181 pub fn with_time_to_live(mut self, duration: Duration) -> Self {
192 self.time_to_live = Some(duration);
193 self.rebuild();
194 self
195 }
196
197 pub fn weighted_size(&self) -> u64 {
201 self.index.weighted_size()
202 }
203
204 pub fn entry_count(&self) -> u64 {
207 self.index.entry_count()
208 }
209
210 pub async fn run_pending_tasks(&self) {
214 self.index.run_pending_tasks();
215 }
216
217 fn rebuild(&mut self) {
221 self.index = build_index(
222 Arc::clone(&self.root),
223 self.max_capacity_bytes,
224 self.time_to_idle,
225 self.time_to_live,
226 );
227 scan_root(&self.root, &self.index);
228 }
229}
230
231#[derive(Clone, Hash, PartialEq, Eq)]
234struct VariantId {
235 key_hash: String,
236 variant_hash: String,
237}
238
239fn build_index(
243 root: Arc<PathBuf>,
244 max_capacity_bytes: Option<u64>,
245 time_to_idle: Option<Duration>,
246 time_to_live: Option<Duration>,
247) -> Cache<VariantId, u64> {
248 let mut builder = Cache::<VariantId, u64>::builder()
249 .weigher(|_key, &body_len| u32::try_from(body_len).unwrap_or(u32::MAX))
250 .eviction_listener(move |id: Arc<VariantId>, _body_len, cause: RemovalCause| {
251 if cause.was_evicted() {
252 let dir = root.join(&id.key_hash);
253 let _ = std::fs::remove_file(dir.join(format!("{}{META_SUFFIX}", id.variant_hash)));
254 let _ = std::fs::remove_file(dir.join(format!("{}{BODY_SUFFIX}", id.variant_hash)));
255 }
256 });
257 if let Some(cap) = max_capacity_bytes {
258 builder = builder.max_capacity(cap);
259 }
260 if let Some(tti) = time_to_idle {
261 builder = builder.time_to_idle(tti);
262 }
263 if let Some(ttl) = time_to_live {
264 builder = builder.time_to_live(ttl);
265 }
266 builder.build()
267}
268
269fn scan_root(root: &Path, index: &Cache<VariantId, u64>) {
273 let Ok(key_dirs) = std::fs::read_dir(root) else {
274 return;
275 };
276 for key_entry in key_dirs.flatten() {
277 let key_dir = key_entry.path();
278 let Some(key_hash) = file_stem_string(&key_dir) else {
279 continue;
280 };
281 let Ok(files) = std::fs::read_dir(&key_dir) else {
282 continue;
283 };
284 for file in files.flatten() {
285 let path = file.path();
286 let Some(variant_hash) = path
287 .file_name()
288 .and_then(|name| name.to_str())
289 .and_then(|name| name.strip_suffix(META_SUFFIX))
290 .map(str::to_string)
291 else {
292 continue;
293 };
294 let body = key_dir.join(format!("{variant_hash}{BODY_SUFFIX}"));
295 let Ok(metadata) = std::fs::metadata(&body) else {
296 continue;
297 };
298 index.insert(
299 VariantId {
300 key_hash: key_hash.clone(),
301 variant_hash,
302 },
303 metadata.len(),
304 );
305 }
306 }
307 index.run_pending_tasks();
308}
309
310fn file_stem_string(path: &Path) -> Option<String> {
311 path.file_name()
312 .and_then(|name| name.to_str())
313 .map(str::to_string)
314}
315
316#[derive(rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
319struct StoredMeta {
320 policy: PolicyRepr,
321 trailers: Option<Headers>,
322}
323
324impl CacheStorage for FileSystemStorage {
325 type PutHandle = FsPutHandle;
326 type StoredEntry = FsStoredEntry;
327
328 async fn get(&self, key: &CacheKey) -> Vec<Self::StoredEntry> {
329 let key_hash = key_hash(key);
330 let dir = self.root.join(&key_hash);
331 let Ok(paths) = fs_shims::read_dir_paths(&dir).await else {
332 return Vec::new();
333 };
334
335 let mut entries = Vec::new();
336 for path in paths {
337 let Some(variant_hash) = path
338 .file_name()
339 .and_then(|name| name.to_str())
340 .and_then(|name| name.strip_suffix(META_SUFFIX))
341 .map(str::to_string)
342 else {
343 continue;
344 };
345 let Ok(bytes) = fs_shims::read(&path).await else {
346 continue;
347 };
348 let Ok(meta) = deserialize_meta(&bytes) else {
349 continue;
350 };
351 self.index.get(&VariantId {
353 key_hash: key_hash.clone(),
354 variant_hash: variant_hash.clone(),
355 });
356 entries.push(FsStoredEntry {
357 meta_path: path,
358 body_path: dir.join(format!("{variant_hash}{BODY_SUFFIX}")),
359 policy: meta.policy.into(),
360 trailers: meta.trailers,
361 });
362 }
363 entries
364 }
365
366 async fn put(&self, key: CacheKey, policy: CachePolicy) -> io::Result<Self::PutHandle> {
367 let key_hash = key_hash(&key);
368 let dir = self.root.join(&key_hash);
369 fs_shims::create_dir_all(&dir).await?;
370
371 let variant_hash = variant_hash(&policy);
372 let n = TEMP_COUNTER.fetch_add(1, Ordering::Relaxed);
373 let body_tmp = dir.join(format!("{variant_hash}{BODY_SUFFIX}.tmp.{n}"));
374 let writer = fs_shims::create(&body_tmp).await?;
375
376 Ok(FsPutHandle {
377 writer,
378 body_tmp,
379 body_final: dir.join(format!("{variant_hash}{BODY_SUFFIX}")),
380 meta_tmp: dir.join(format!("{variant_hash}{META_SUFFIX}.tmp.{n}")),
381 meta_final: dir.join(format!("{variant_hash}{META_SUFFIX}")),
382 policy,
383 index: self.index.clone(),
384 variant_id: VariantId {
385 key_hash,
386 variant_hash,
387 },
388 written: 0,
389 committed: false,
390 })
391 }
392
393 async fn invalidate(&self, key: &CacheKey) {
394 let key_hash = key_hash(key);
395 let dir = self.root.join(&key_hash);
396 if let Ok(paths) = fs_shims::read_dir_paths(&dir).await {
399 for path in paths {
400 if let Some(variant_hash) = path
401 .file_name()
402 .and_then(|name| name.to_str())
403 .and_then(|name| name.strip_suffix(META_SUFFIX))
404 {
405 self.index.invalidate(&VariantId {
406 key_hash: key_hash.clone(),
407 variant_hash: variant_hash.to_string(),
408 });
409 }
410 }
411 }
412 let _ = fs_shims::remove_dir_all(&dir).await;
413 }
414}
415
416pub struct FsPutHandle {
422 writer: fs_shims::Writer,
423 body_tmp: PathBuf,
424 body_final: PathBuf,
425 meta_tmp: PathBuf,
426 meta_final: PathBuf,
427 policy: CachePolicy,
428 index: Cache<VariantId, u64>,
429 variant_id: VariantId,
430 written: u64,
431 committed: bool,
432}
433
434impl Debug for FsPutHandle {
435 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
436 f.debug_struct("FsPutHandle")
437 .field("body_final", &self.body_final)
438 .finish_non_exhaustive()
439 }
440}
441
442impl AsyncWrite for FsPutHandle {
443 fn poll_write(
444 self: Pin<&mut Self>,
445 cx: &mut Context<'_>,
446 buf: &[u8],
447 ) -> Poll<io::Result<usize>> {
448 let this = self.get_mut();
449 let poll = Pin::new(&mut this.writer).poll_write(cx, buf);
450 if let Poll::Ready(Ok(n)) = &poll {
451 this.written += *n as u64;
452 }
453 poll
454 }
455
456 fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
457 Pin::new(&mut self.get_mut().writer).poll_flush(cx)
458 }
459
460 fn poll_close(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
461 Pin::new(&mut self.get_mut().writer).poll_close(cx)
462 }
463}
464
465impl PutHandle for FsPutHandle {
466 async fn finalize(mut self, trailers: Option<Headers>) -> io::Result<()> {
467 self.writer.close().await?;
468 fs_shims::rename(&self.body_tmp, &self.body_final).await?;
469
470 let meta = StoredMeta {
471 policy: PolicyRepr::from(&self.policy),
472 trailers,
473 };
474 let bytes = serialize_meta(&meta)?;
475 fs_shims::write(&self.meta_tmp, &bytes).await?;
476 fs_shims::rename(&self.meta_tmp, &self.meta_final).await?;
477
478 self.index.insert(self.variant_id.clone(), self.written);
481
482 self.committed = true;
483 Ok(())
484 }
485}
486
487impl Drop for FsPutHandle {
488 fn drop(&mut self) {
489 if !self.committed {
490 let _ = std::fs::remove_file(&self.body_tmp);
491 }
492 }
493}
494
495#[derive(Clone)]
500pub struct FsStoredEntry {
501 meta_path: PathBuf,
502 body_path: PathBuf,
503 policy: CachePolicy,
504 trailers: Option<Headers>,
505}
506
507impl Debug for FsStoredEntry {
508 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
509 f.debug_struct("FsStoredEntry")
510 .field("body_path", &self.body_path)
511 .field("has_trailers", &self.trailers.is_some())
512 .finish_non_exhaustive()
513 }
514}
515
516impl StoredEntry for FsStoredEntry {
517 fn policy(&self) -> &CachePolicy {
518 &self.policy
519 }
520
521 async fn refresh_policy(&mut self, new_policy: CachePolicy) -> io::Result<()> {
522 let meta = StoredMeta {
523 policy: PolicyRepr::from(&new_policy),
524 trailers: self.trailers.clone(),
525 };
526 let bytes = serialize_meta(&meta)?;
527 let tmp = temp_sibling(&self.meta_path);
528 fs_shims::write(&tmp, &bytes).await?;
529 fs_shims::rename(&tmp, &self.meta_path).await?;
530
531 self.policy = new_policy;
532 Ok(())
533 }
534
535 async fn open(self) -> io::Result<Body> {
536 let len = fs_shims::metadata_len(&self.body_path).await?;
537 let reader = fs_shims::open(&self.body_path).await?;
538 let source = FsBodySource {
539 reader,
540 trailers: self.trailers,
541 };
542 Ok(Body::new_with_trailers(source, Some(len)))
543 }
544}
545
546struct FsBodySource {
549 reader: fs_shims::Reader,
550 trailers: Option<Headers>,
551}
552
553impl AsyncRead for FsBodySource {
554 fn poll_read(
555 self: Pin<&mut Self>,
556 cx: &mut Context<'_>,
557 buf: &mut [u8],
558 ) -> Poll<io::Result<usize>> {
559 Pin::new(&mut self.get_mut().reader).poll_read(cx, buf)
560 }
561}
562
563impl BodySource for FsBodySource {
564 fn trailers(self: Pin<&mut Self>) -> Option<Headers> {
565 self.get_mut().trailers.take()
566 }
567}
568
569fn hash_hex(bytes: &[u8]) -> String {
570 let mut hasher = Sha256::new();
571 hasher.update(bytes);
572 finalize_hex(hasher)
573}
574
575fn key_hash(key: &CacheKey) -> String {
576 hash_hex(key.to_string().as_bytes())
577}
578
579fn variant_hash(policy: &CachePolicy) -> String {
580 let mut hasher = Sha256::new();
581 for (name, value) in &policy.vary_snapshot {
582 hasher.update(name.as_bytes());
583 hasher.update([0]);
584 match value {
585 Some(value) => {
586 hasher.update([1]);
587 hasher.update(value.as_bytes());
588 }
589 None => hasher.update([0]),
590 }
591 hasher.update([0]);
592 }
593 finalize_hex(hasher)
594}
595
596fn finalize_hex(hasher: Sha256) -> String {
597 let digest = hasher.finalize();
598 let mut out = String::with_capacity(digest.len() * 2);
599 for byte in digest {
600 write!(out, "{byte:02x}").expect("writing to a String cannot fail");
601 }
602 out
603}
604
605fn temp_sibling(path: &Path) -> PathBuf {
607 let n = TEMP_COUNTER.fetch_add(1, Ordering::Relaxed);
608 let mut name = path.as_os_str().to_owned();
609 name.push(format!(".tmp.{n}"));
610 PathBuf::from(name)
611}
612
613fn serialize_meta(meta: &StoredMeta) -> io::Result<rkyv::util::AlignedVec> {
614 rkyv::to_bytes::<rkyv::rancor::Error>(meta)
615 .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))
616}
617
618fn deserialize_meta(bytes: &[u8]) -> io::Result<StoredMeta> {
619 let mut aligned = rkyv::util::AlignedVec::<16>::new();
622 aligned.extend_from_slice(bytes);
623 rkyv::from_bytes::<StoredMeta, rkyv::rancor::Error>(&aligned)
624 .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))
625}
626
627#[cfg(test)]
628mod tests {
629 use super::*;
630 use crate::test_helpers::*;
631 use futures_lite::{AsyncReadExt, AsyncWriteExt};
632 use std::time::{Duration, SystemTime};
633 use tempfile::TempDir;
634 use trillium_client::Conn;
635 use trillium_http::{KnownHeaderName::*, Method, Status};
636 use trillium_testing::{TestResult, harness, test};
637
638 fn key() -> CacheKey {
639 CacheKey::new(Method::Get, "http://example.com/".parse().unwrap())
640 }
641
642 fn new_storage() -> (TempDir, FileSystemStorage) {
643 let dir = tempfile::tempdir().unwrap();
644 let storage = FileSystemStorage::new(dir.path());
645 (dir, storage)
646 }
647
648 async fn store_at(storage: &FileSystemStorage, url: &str, body: &[u8]) {
649 let key = CacheKey::new(Method::Get, url.parse().unwrap());
650 let conn = exchange(
651 Method::Get,
652 &[],
653 Status::Ok,
654 &[(CacheControl, "max-age=600")],
655 );
656 let policy = policy_from(&conn, SystemTime::now(), private_cache());
657 let mut handle = storage.put(key, policy).await.unwrap();
658 handle.write_all(body).await.unwrap();
659 handle.finalize(None).await.unwrap();
660 }
661
662 async fn store(storage: &FileSystemStorage, key: CacheKey, conn: &Conn, body: &[u8]) {
663 let policy = policy_from(conn, SystemTime::now(), private_cache());
664 let mut handle = storage.put(key, policy).await.unwrap();
665 handle.write_all(body).await.unwrap();
666 handle.finalize(None).await.unwrap();
667 }
668
669 async fn read_body(entry: FsStoredEntry) -> Vec<u8> {
670 let mut body = entry.open().await.unwrap();
671 let mut buf = Vec::new();
672 body.read_to_end(&mut buf).await.unwrap();
673 buf
674 }
675
676 #[test(harness)]
677 async fn get_missing_key_returns_empty() -> TestResult {
678 let (_dir, storage) = new_storage();
679 assert!(storage.get(&key()).await.is_empty());
680 Ok(())
681 }
682
683 #[test(harness)]
684 async fn put_then_get_round_trips_through_disk() -> TestResult {
685 let (_dir, storage) = new_storage();
686 let conn = exchange(
687 Method::Get,
688 &[],
689 Status::Ok,
690 &[(CacheControl, "max-age=600")],
691 );
692 store(&storage, key(), &conn, b"hello").await;
693 let result = storage.get(&key()).await;
694 assert_eq!(result.len(), 1);
695 assert_eq!(read_body(result[0].clone()).await, b"hello");
696 Ok(())
697 }
698
699 #[test(harness)]
700 async fn put_with_same_vary_replaces() -> TestResult {
701 let (_dir, storage) = new_storage();
702 let conn = exchange(
703 Method::Get,
704 &[(AcceptEncoding, "gzip")],
705 Status::Ok,
706 &[(CacheControl, "max-age=600"), (Vary, "Accept-Encoding")],
707 );
708 store(&storage, key(), &conn, b"v1").await;
709 store(&storage, key(), &conn, b"v2").await;
710 let result = storage.get(&key()).await;
711 assert_eq!(result.len(), 1);
712 assert_eq!(read_body(result[0].clone()).await, b"v2");
713 Ok(())
714 }
715
716 #[test(harness)]
717 async fn put_with_different_vary_appends() -> TestResult {
718 let (_dir, storage) = new_storage();
719 let gzip = exchange(
720 Method::Get,
721 &[(AcceptEncoding, "gzip")],
722 Status::Ok,
723 &[(CacheControl, "max-age=600"), (Vary, "Accept-Encoding")],
724 );
725 let br = exchange(
726 Method::Get,
727 &[(AcceptEncoding, "br")],
728 Status::Ok,
729 &[(CacheControl, "max-age=600"), (Vary, "Accept-Encoding")],
730 );
731 store(&storage, key(), &gzip, b"gz").await;
732 store(&storage, key(), &br, b"br").await;
733 assert_eq!(storage.get(&key()).await.len(), 2);
734 Ok(())
735 }
736
737 #[test(harness)]
738 async fn invalidate_removes_all_entries_for_key() -> TestResult {
739 let (_dir, storage) = new_storage();
740 let conn = exchange(
741 Method::Get,
742 &[],
743 Status::Ok,
744 &[(CacheControl, "max-age=600")],
745 );
746 store(&storage, key(), &conn, b"x").await;
747 storage.invalidate(&key()).await;
748 assert!(storage.get(&key()).await.is_empty());
749 Ok(())
750 }
751
752 #[test(harness)]
753 async fn invalidate_does_not_touch_other_keys() -> TestResult {
754 let (_dir, storage) = new_storage();
755 let conn = exchange(
756 Method::Get,
757 &[],
758 Status::Ok,
759 &[(CacheControl, "max-age=600")],
760 );
761 let key_a = CacheKey::new(Method::Get, "http://a.example/".parse().unwrap());
762 let key_b = CacheKey::new(Method::Get, "http://b.example/".parse().unwrap());
763 store(&storage, key_a.clone(), &conn, b"a").await;
764 store(&storage, key_b.clone(), &conn, b"b").await;
765 storage.invalidate(&key_a).await;
766 assert!(storage.get(&key_a).await.is_empty());
767 assert_eq!(storage.get(&key_b).await.len(), 1);
768 Ok(())
769 }
770
771 #[test(harness)]
772 async fn drop_put_handle_without_finalize_discards() -> TestResult {
773 let (_dir, storage) = new_storage();
774 let conn = exchange(
775 Method::Get,
776 &[],
777 Status::Ok,
778 &[(CacheControl, "max-age=600")],
779 );
780 let policy = policy_from(&conn, SystemTime::now(), private_cache());
781 let mut handle = storage.put(key(), policy).await.unwrap();
782 handle.write_all(b"partial").await.unwrap();
783 drop(handle);
784 assert!(storage.get(&key()).await.is_empty());
785 Ok(())
786 }
787
788 #[test(harness)]
789 async fn refresh_policy_updates_meta_and_keeps_body() -> TestResult {
790 let (_dir, storage) = new_storage();
791 let conn = exchange(
792 Method::Get,
793 &[],
794 Status::Ok,
795 &[(CacheControl, "max-age=600")],
796 );
797 store(&storage, key(), &conn, b"body").await;
798
799 let mut entries = storage.get(&key()).await;
800 let original_time = entries[0].policy().response_time;
801 let refreshed = exchange(
802 Method::Get,
803 &[],
804 Status::Ok,
805 &[(CacheControl, "max-age=1200")],
806 );
807 let new_policy = policy_from(
808 &refreshed,
809 original_time + Duration::from_secs(100),
810 private_cache(),
811 );
812 entries[0].refresh_policy(new_policy).await.unwrap();
813
814 let fresh = storage.get(&key()).await;
815 assert_eq!(fresh.len(), 1);
816 assert_ne!(fresh[0].policy().response_time, original_time);
817 assert_eq!(read_body(fresh[0].clone()).await, b"body");
818 Ok(())
819 }
820
821 #[test(harness)]
822 async fn trailers_round_trip() -> TestResult {
823 let (_dir, storage) = new_storage();
824 let conn = exchange(
825 Method::Get,
826 &[],
827 Status::Ok,
828 &[(CacheControl, "max-age=600")],
829 );
830 let policy = policy_from(&conn, SystemTime::now(), private_cache());
831 let mut handle = storage.put(key(), policy).await.unwrap();
832 handle.write_all(b"data").await.unwrap();
833 let mut trailers = Headers::new();
834 trailers.insert("x-checksum", "abc123");
835 handle.finalize(Some(trailers)).await.unwrap();
836
837 let entry = storage.get(&key()).await.remove(0);
838 let mut body = entry.open().await.unwrap();
839 let mut buf = Vec::new();
840 body.read_to_end(&mut buf).await.unwrap();
841 assert_eq!(buf, b"data");
842 let trailers = body
843 .trailers()
844 .expect("stored trailers should surface after EOF");
845 assert_eq!(trailers.get_str("x-checksum"), Some("abc123"));
846 Ok(())
847 }
848
849 #[test(harness)]
850 async fn persists_across_new_storage_on_same_root() -> TestResult {
851 let dir = tempfile::tempdir().unwrap();
852 let conn = exchange(
853 Method::Get,
854 &[],
855 Status::Ok,
856 &[(CacheControl, "max-age=600")],
857 );
858 {
859 let storage = FileSystemStorage::new(dir.path());
860 store(&storage, key(), &conn, b"persisted").await;
861 }
862
863 let reopened = FileSystemStorage::new(dir.path());
865 let result = reopened.get(&key()).await;
866 assert_eq!(result.len(), 1);
867 assert_eq!(read_body(result[0].clone()).await, b"persisted");
868 Ok(())
869 }
870
871 #[test(harness)]
872 async fn size_cap_evicts_and_deletes_files() -> TestResult {
873 let dir = tempfile::tempdir().unwrap();
875 let storage = FileSystemStorage::new(dir.path()).with_max_capacity_bytes(1024);
876 let body = vec![b'x'; 600];
877 for i in 0..10 {
878 store_at(&storage, &format!("http://example.com/{i}"), &body).await;
879 }
880 storage.run_pending_tasks().await;
881 assert!(
882 storage.weighted_size() <= 1024,
883 "weighted size {} should be within cap of 1024",
884 storage.weighted_size()
885 );
886
887 let reopened = FileSystemStorage::new(dir.path()).unbounded();
891 assert!(
892 reopened.weighted_size() <= 1024,
893 "on-disk bytes {} should be within cap of 1024",
894 reopened.weighted_size()
895 );
896 Ok(())
897 }
898
899 #[test(harness)]
900 async fn rebuild_scan_trims_over_cap_directory() -> TestResult {
901 let dir = tempfile::tempdir().unwrap();
902 let body = vec![b'x'; 600];
903 {
904 let unbounded = FileSystemStorage::new(dir.path()).unbounded();
905 for i in 0..10 {
906 store_at(&unbounded, &format!("http://example.com/{i}"), &body).await;
907 }
908 unbounded.run_pending_tasks().await;
909 assert_eq!(unbounded.entry_count(), 10);
910 }
911
912 let capped = FileSystemStorage::new(dir.path()).with_max_capacity_bytes(1024);
914 assert!(
915 capped.weighted_size() <= 1024,
916 "weighted size {} should be within cap of 1024",
917 capped.weighted_size()
918 );
919 Ok(())
920 }
921
922 #[test(harness)]
923 async fn unbounded_keeps_all_entries() -> TestResult {
924 let dir = tempfile::tempdir().unwrap();
925 let storage = FileSystemStorage::new(dir.path()).unbounded();
926 let body = vec![b'x'; 600];
927 for i in 0..10 {
928 store_at(&storage, &format!("http://example.com/{i}"), &body).await;
929 }
930 storage.run_pending_tasks().await;
931 assert_eq!(storage.entry_count(), 10);
932 assert_eq!(storage.weighted_size(), 6000);
933 Ok(())
934 }
935
936 #[test(harness)]
937 async fn replacing_a_variant_does_not_double_count() -> TestResult {
938 let (_dir, storage) = new_storage();
939 store_at(&storage, "http://example.com/", &vec![b'x'; 600]).await;
940 store_at(&storage, "http://example.com/", &vec![b'y'; 300]).await;
941 storage.run_pending_tasks().await;
942 assert_eq!(storage.entry_count(), 1);
943 assert_eq!(storage.weighted_size(), 300);
944 Ok(())
945 }
946
947 #[test(harness)]
950 async fn time_to_live_evicts_and_deletes_files() -> TestResult {
951 let dir = tempfile::tempdir().unwrap();
952 let storage =
953 FileSystemStorage::new(dir.path()).with_time_to_live(Duration::from_millis(50));
954 store_at(&storage, "http://example.com/", b"x").await;
955 storage.run_pending_tasks().await;
956 assert_eq!(storage.entry_count(), 1);
957
958 std::thread::sleep(Duration::from_millis(120));
959 storage.run_pending_tasks().await;
960 assert_eq!(storage.entry_count(), 0);
961
962 let reopened = FileSystemStorage::new(dir.path()).unbounded();
965 assert_eq!(reopened.entry_count(), 0);
966 Ok(())
967 }
968
969 #[test(harness)]
970 async fn time_to_idle_evicts_unread_entries() -> TestResult {
971 let dir = tempfile::tempdir().unwrap();
972 let storage =
973 FileSystemStorage::new(dir.path()).with_time_to_idle(Duration::from_millis(50));
974 store_at(&storage, "http://example.com/", b"x").await;
975 storage.run_pending_tasks().await;
976 assert_eq!(storage.entry_count(), 1);
977
978 std::thread::sleep(Duration::from_millis(120));
980 storage.run_pending_tasks().await;
981 assert_eq!(storage.entry_count(), 0);
982 let reopened = FileSystemStorage::new(dir.path()).unbounded();
983 assert_eq!(reopened.entry_count(), 0);
984 Ok(())
985 }
986}