1use std::collections::HashMap;
16use std::path::PathBuf;
17use std::sync::Arc;
18use std::time::{Duration, Instant};
19
20use async_trait::async_trait;
21use async_walkdir::{Filtering, WalkDir};
22use base64::Engine;
23use base64::prelude::BASE64_URL_SAFE;
24use common_base::range_read::FileReader;
25use common_runtime::runtime::RuntimeTrait;
26use common_telemetry::{info, warn};
27use futures::{FutureExt, StreamExt};
28use moka::future::Cache;
29use moka::policy::EvictionPolicy;
30use sha2::{Digest, Sha256};
31use snafu::ResultExt;
32use tokio::fs;
33use tokio::sync::mpsc::error::TrySendError;
34use tokio::sync::mpsc::{Receiver, Sender};
35use tokio_util::compat::TokioAsyncWriteCompatExt;
36
37use crate::error::{
38 CacheGetSnafu, CreateSnafu, MetadataSnafu, OpenSnafu, ReadSnafu, RemoveSnafu, RenameSnafu,
39 Result, WalkDirSnafu,
40};
41use crate::puffin_manager::stager::{
42 BoxWriter, DirWriterProvider, InitBlobFn, InitDirFn, Stager, StagerNotifier,
43};
44use crate::puffin_manager::{BlobGuard, DirGuard, DirMetrics};
45
46const DELETE_QUEUE_SIZE: usize = 10240;
47const TMP_EXTENSION: &str = "tmp";
48const DELETED_EXTENSION: &str = "deleted";
49const RECYCLE_BIN_TTL: Duration = Duration::from_secs(60);
50
51pub struct BoundedStager<H> {
53 base_dir: PathBuf,
55
56 cache: Cache<String, CacheValue>,
58
59 recycle_bin: Cache<String, CacheValue>,
61
62 delete_queue: Sender<DeleteTask>,
72
73 notifier: Option<Arc<dyn StagerNotifier>>,
75
76 _phantom: std::marker::PhantomData<H>,
77}
78
79impl<H: 'static> BoundedStager<H> {
80 pub async fn new(
81 base_dir: PathBuf,
82 capacity: u64,
83 notifier: Option<Arc<dyn StagerNotifier>>,
84 cache_ttl: Option<Duration>,
85 ) -> Result<Self> {
86 tokio::fs::create_dir_all(&base_dir)
87 .await
88 .context(CreateSnafu)?;
89
90 let recycle_bin = Cache::builder().time_to_idle(RECYCLE_BIN_TTL).build();
91 let recycle_bin_cloned = recycle_bin.clone();
92 let notifier_cloned = notifier.clone();
93
94 let mut cache_builder = Cache::builder()
95 .max_capacity(capacity)
96 .weigher(|_: &String, v: &CacheValue| v.weight())
97 .eviction_policy(EvictionPolicy::lru())
98 .support_invalidation_closures()
99 .async_eviction_listener(move |k, v, _| {
100 let recycle_bin = recycle_bin_cloned.clone();
101 if let Some(notifier) = notifier_cloned.as_ref() {
102 notifier.on_cache_evict(v.size());
103 notifier.on_recycle_insert(v.size());
104 }
105 async move {
106 recycle_bin.insert(k.as_str().to_string(), v).await;
107 }
108 .boxed()
109 });
110 if let Some(ttl) = cache_ttl
111 && !ttl.is_zero()
112 {
113 cache_builder = cache_builder.time_to_live(ttl);
114 }
115 let cache = cache_builder.build();
116
117 let (delete_queue, rx) = tokio::sync::mpsc::channel(DELETE_QUEUE_SIZE);
118 let notifier_cloned = notifier.clone();
119 common_runtime::global_runtime().spawn(Self::delete_routine(
120 rx,
121 recycle_bin.clone(),
122 notifier_cloned,
123 ));
124 let stager = Self {
125 cache,
126 base_dir,
127 delete_queue,
128 recycle_bin,
129 notifier,
130 _phantom: std::marker::PhantomData,
131 };
132
133 stager.recover().await?;
134
135 Ok(stager)
136 }
137}
138
139#[async_trait]
140impl<H: ToString + Clone + Send + Sync> Stager for BoundedStager<H> {
141 type Blob = Arc<FsBlobGuard>;
142 type Dir = Arc<FsDirGuard>;
143 type FileHandle = H;
144
145 async fn get_blob<'a>(
146 &self,
147 handle: &Self::FileHandle,
148 key: &str,
149 init_fn: Box<InitBlobFn<'a>>,
150 ) -> Result<Self::Blob> {
151 let handle_str = handle.to_string();
152 let cache_key = Self::encode_cache_key(&handle_str, key);
153
154 let mut miss = false;
155 let v = self
156 .cache
157 .try_get_with_by_ref(&cache_key, async {
158 if let Some(v) = self.recycle_bin.remove(&cache_key).await {
159 if let Some(notifier) = self.notifier.as_ref() {
160 let size = v.size();
161 notifier.on_cache_insert(size);
162 notifier.on_recycle_clear(size);
163 }
164 return Ok(v);
165 }
166
167 miss = true;
168 let timer = Instant::now();
169 let file_name = format!("{}.{}", cache_key, uuid::Uuid::new_v4());
170 let path = self.base_dir.join(&file_name);
171
172 let size = Self::write_blob(&path, init_fn).await?;
173 if let Some(notifier) = self.notifier.as_ref() {
174 notifier.on_cache_insert(size);
175 notifier.on_load_blob(timer.elapsed());
176 }
177 let guard = Arc::new(FsBlobGuard {
178 handle: handle_str,
179 path,
180 delete_queue: self.delete_queue.clone(),
181 size,
182 });
183 Ok(CacheValue::File(guard))
184 })
185 .await
186 .context(CacheGetSnafu)?;
187
188 if let Some(notifier) = self.notifier.as_ref() {
189 if miss {
190 notifier.on_cache_miss(v.size());
191 } else {
192 notifier.on_cache_hit(v.size());
193 }
194 }
195 match v {
196 CacheValue::File(guard) => Ok(guard),
197 _ => unreachable!(),
198 }
199 }
200
201 async fn get_dir<'a>(
202 &self,
203 handle: &Self::FileHandle,
204 key: &str,
205 init_fn: Box<InitDirFn<'a>>,
206 ) -> Result<(Self::Dir, DirMetrics)> {
207 let handle_str = handle.to_string();
208
209 let cache_key = Self::encode_cache_key(&handle_str, key);
210
211 let mut miss = false;
212 let v = self
213 .cache
214 .try_get_with_by_ref(&cache_key, async {
215 if let Some(v) = self.recycle_bin.remove(&cache_key).await {
216 if let Some(notifier) = self.notifier.as_ref() {
217 let size = v.size();
218 notifier.on_cache_insert(size);
219 notifier.on_recycle_clear(size);
220 }
221 return Ok(v);
222 }
223
224 miss = true;
225 let timer = Instant::now();
226 let dir_name = format!("{}.{}", cache_key, uuid::Uuid::new_v4());
227 let path = self.base_dir.join(&dir_name);
228
229 let size = Self::write_dir(&path, init_fn).await?;
230 if let Some(notifier) = self.notifier.as_ref() {
231 notifier.on_cache_insert(size);
232 notifier.on_load_dir(timer.elapsed());
233 }
234 let guard = Arc::new(FsDirGuard {
235 handle: handle_str,
236 path,
237 size,
238 delete_queue: self.delete_queue.clone(),
239 });
240 Ok(CacheValue::Dir(guard))
241 })
242 .await
243 .context(CacheGetSnafu)?;
244
245 let dir_size = v.size();
246 if let Some(notifier) = self.notifier.as_ref() {
247 if miss {
248 notifier.on_cache_miss(dir_size);
249 } else {
250 notifier.on_cache_hit(dir_size);
251 }
252 }
253
254 let metrics = DirMetrics {
255 cache_hit: !miss,
256 dir_size,
257 };
258
259 match v {
260 CacheValue::Dir(guard) => Ok((guard, metrics)),
261 _ => unreachable!(),
262 }
263 }
264
265 async fn put_dir(
266 &self,
267 handle: &Self::FileHandle,
268 key: &str,
269 dir_path: PathBuf,
270 size: u64,
271 ) -> Result<()> {
272 let handle_str = handle.to_string();
273 let cache_key = Self::encode_cache_key(&handle_str, key);
274
275 self.cache
276 .try_get_with(cache_key.clone(), async move {
277 if let Some(v) = self.recycle_bin.remove(&cache_key).await {
278 if let Some(notifier) = self.notifier.as_ref() {
279 let size = v.size();
280 notifier.on_cache_insert(size);
281 notifier.on_recycle_clear(size);
282 }
283 return Ok(v);
284 }
285
286 let dir_name = format!("{}.{}", cache_key, uuid::Uuid::new_v4());
287 let path = self.base_dir.join(&dir_name);
288
289 fs::rename(&dir_path, &path).await.context(RenameSnafu)?;
290 if let Some(notifier) = self.notifier.as_ref() {
291 notifier.on_cache_insert(size);
292 }
293 let guard = Arc::new(FsDirGuard {
294 handle: handle_str,
295 path,
296 size,
297 delete_queue: self.delete_queue.clone(),
298 });
299 Ok(CacheValue::Dir(guard))
300 })
301 .await
302 .map(|_| ())
303 .context(CacheGetSnafu)?;
304
305 self.cache.run_pending_tasks().await;
309 self.recycle_bin.run_pending_tasks().await;
310
311 Ok(())
312 }
313
314 async fn purge(&self, handle: &Self::FileHandle) -> Result<()> {
315 let handle_str = handle.to_string();
316 self.cache
317 .invalidate_entries_if(move |_k, v| v.handle() == handle_str)
318 .unwrap(); self.cache.run_pending_tasks().await;
320 Ok(())
321 }
322}
323
324impl<H> BoundedStager<H> {
325 fn encode_cache_key(puffin_file_name: &str, key: &str) -> String {
326 let mut hasher = Sha256::new();
327 hasher.update(puffin_file_name);
328 hasher.update(key);
329 hasher.update(puffin_file_name);
330 let hash = hasher.finalize();
331
332 BASE64_URL_SAFE.encode(hash)
333 }
334
335 async fn write_blob(target_path: &PathBuf, init_fn: Box<InitBlobFn<'_>>) -> Result<u64> {
336 let tmp_path = target_path.with_extension(TMP_EXTENSION);
339 let writer = Box::new(
340 fs::File::create(&tmp_path)
341 .await
342 .context(CreateSnafu)?
343 .compat_write(),
344 );
345 let size = init_fn(writer).await?;
346
347 fs::rename(tmp_path, target_path)
349 .await
350 .context(RenameSnafu)?;
351 Ok(size)
352 }
353
354 async fn write_dir(target_path: &PathBuf, init_fn: Box<InitDirFn<'_>>) -> Result<u64> {
355 let tmp_base = target_path.with_extension(TMP_EXTENSION);
358 let writer_provider = Box::new(MokaDirWriterProvider(tmp_base.clone()));
359 let size = init_fn(writer_provider).await?;
360
361 fs::rename(&tmp_base, target_path)
363 .await
364 .context(RenameSnafu)?;
365 Ok(size)
366 }
367
368 async fn recover(&self) -> Result<()> {
373 let timer = std::time::Instant::now();
374 info!("Recovering the staging area, base_dir: {:?}", self.base_dir);
375
376 let mut read_dir = fs::read_dir(&self.base_dir).await.context(ReadSnafu)?;
377
378 let mut elems = HashMap::new();
379 while let Some(entry) = read_dir.next_entry().await.context(ReadSnafu)? {
380 let path = entry.path();
381
382 if path.extension() == Some(TMP_EXTENSION.as_ref())
383 || path.extension() == Some(DELETED_EXTENSION.as_ref())
384 {
385 if entry.metadata().await.context(MetadataSnafu)?.is_dir() {
387 fs::remove_dir_all(path).await.context(RemoveSnafu)?;
388 } else {
389 fs::remove_file(path).await.context(RemoveSnafu)?;
390 }
391 } else {
392 let meta = entry.metadata().await.context(MetadataSnafu)?;
394 let file_path = path.file_name().unwrap().to_string_lossy().into_owned();
395
396 let key = match file_path.split('.').next() {
398 Some(key) => key.to_string(),
399 None => {
400 warn!(
401 "Invalid staging file name: {}, expected format: <key>.<uuid>",
402 file_path
403 );
404 continue;
405 }
406 };
407
408 if meta.is_dir() {
409 let size = Self::get_dir_size(&path).await?;
410 let v = CacheValue::Dir(Arc::new(FsDirGuard {
411 path,
412 size,
413 delete_queue: self.delete_queue.clone(),
414
415 handle: String::new(),
417 }));
418 let _dup_dir = elems.insert(key, v);
420 } else {
421 let size = meta.len();
422 let v = CacheValue::File(Arc::new(FsBlobGuard {
423 path,
424 size,
425 delete_queue: self.delete_queue.clone(),
426
427 handle: String::new(),
429 }));
430 let _dup_file = elems.insert(key, v);
432 }
433 }
434 }
435
436 let mut size = 0;
437 let num_elems = elems.len();
438 for (key, value) in elems {
439 size += value.size();
440 self.cache.insert(key, value).await;
441 }
442 if let Some(notifier) = self.notifier.as_ref() {
443 notifier.on_cache_insert(size);
444 }
445
446 self.cache.run_pending_tasks().await;
447
448 info!(
449 "Recovered the staging area, num_entries: {}, num_bytes: {}, cost: {:?}",
450 num_elems,
451 size,
452 timer.elapsed()
453 );
454 Ok(())
455 }
456
457 async fn get_dir_size(path: &PathBuf) -> Result<u64> {
459 let mut size = 0;
460 let mut wd = WalkDir::new(path).filter(|entry| async move {
461 match entry.file_type().await {
462 Ok(ft) if ft.is_dir() => Filtering::Ignore,
463 _ => Filtering::Continue,
464 }
465 });
466
467 while let Some(entry) = wd.next().await {
468 let entry = entry.context(WalkDirSnafu)?;
469 size += entry.metadata().await.context(MetadataSnafu)?.len();
470 }
471
472 Ok(size)
473 }
474
475 async fn delete_routine(
476 mut receiver: Receiver<DeleteTask>,
477 recycle_bin: Cache<String, CacheValue>,
478 notifier: Option<Arc<dyn StagerNotifier>>,
479 ) {
480 loop {
481 match tokio::time::timeout(RECYCLE_BIN_TTL, receiver.recv()).await {
482 Ok(Some(task)) => match task {
483 DeleteTask::File(path, size) => {
484 if let Err(err) = fs::remove_file(&path).await {
485 if err.kind() == std::io::ErrorKind::NotFound {
486 continue;
487 }
488
489 warn!(err; "Failed to remove the file.");
490 }
491
492 if let Some(notifier) = notifier.as_ref() {
493 notifier.on_recycle_clear(size);
494 }
495 }
496
497 DeleteTask::Dir(path, size) => {
498 let deleted_path = path.with_extension(DELETED_EXTENSION);
499 if let Err(err) = fs::rename(&path, &deleted_path).await {
500 if err.kind() == std::io::ErrorKind::NotFound {
501 continue;
502 }
503
504 let _ = fs::remove_dir_all(&deleted_path).await;
506 if let Err(err) = fs::rename(&path, &deleted_path).await {
507 warn!(err; "Failed to rename the dangling directory to deleted path.");
508 continue;
509 }
510 }
511 if let Err(err) = fs::remove_dir_all(&deleted_path).await {
512 warn!(err; "Failed to remove the dangling directory.");
513 }
514 if let Some(notifier) = notifier.as_ref() {
515 notifier.on_recycle_clear(size);
516 }
517 }
518 DeleteTask::Terminate => {
519 break;
520 }
521 },
522 Ok(None) => break,
523 Err(_) => {
524 recycle_bin.run_pending_tasks().await;
526 }
527 }
528 }
529
530 info!("The delete routine for the bounded stager is terminated.");
531 }
532}
533
534impl<H> Drop for BoundedStager<H> {
535 fn drop(&mut self) {
536 let _ = self.delete_queue.try_send(DeleteTask::Terminate);
537 }
538}
539
540#[derive(Debug, Clone)]
541enum CacheValue {
542 File(Arc<FsBlobGuard>),
543 Dir(Arc<FsDirGuard>),
544}
545
546impl CacheValue {
547 fn size(&self) -> u64 {
548 match self {
549 CacheValue::File(guard) => guard.size,
550 CacheValue::Dir(guard) => guard.size,
551 }
552 }
553
554 fn weight(&self) -> u32 {
555 self.size().try_into().unwrap_or(u32::MAX)
556 }
557
558 fn handle(&self) -> &str {
559 match self {
560 CacheValue::File(guard) => &guard.handle,
561 CacheValue::Dir(guard) => &guard.handle,
562 }
563 }
564}
565
566enum DeleteTask {
567 File(PathBuf, u64),
568 Dir(PathBuf, u64),
569 Terminate,
570}
571
572#[derive(Debug)]
575pub struct FsBlobGuard {
576 handle: String,
577 path: PathBuf,
578 size: u64,
579 delete_queue: Sender<DeleteTask>,
580}
581
582#[async_trait]
583impl BlobGuard for FsBlobGuard {
584 type Reader = FileReader;
585
586 async fn reader(&self) -> Result<Self::Reader> {
587 FileReader::new(&self.path).await.context(OpenSnafu)
588 }
589}
590
591impl Drop for FsBlobGuard {
592 fn drop(&mut self) {
593 if let Err(err) = self
594 .delete_queue
595 .try_send(DeleteTask::File(self.path.clone(), self.size))
596 {
597 if matches!(err, TrySendError::Closed(_)) {
598 return;
599 }
600 warn!(err; "Failed to send the delete task for the file.");
601 }
602 }
603}
604
605#[derive(Debug)]
608pub struct FsDirGuard {
609 handle: String,
610 path: PathBuf,
611 size: u64,
612 delete_queue: Sender<DeleteTask>,
613}
614
615impl DirGuard for FsDirGuard {
616 fn path(&self) -> &PathBuf {
617 &self.path
618 }
619}
620
621impl Drop for FsDirGuard {
622 fn drop(&mut self) {
623 if let Err(err) = self
624 .delete_queue
625 .try_send(DeleteTask::Dir(self.path.clone(), self.size))
626 {
627 if matches!(err, TrySendError::Closed(_)) {
628 return;
629 }
630 warn!(err; "Failed to send the delete task for the directory.");
631 }
632 }
633}
634
635struct MokaDirWriterProvider(PathBuf);
637
638#[async_trait]
639impl DirWriterProvider for MokaDirWriterProvider {
640 async fn writer(&self, rel_path: &str) -> Result<BoxWriter> {
641 let full_path = if cfg!(windows) {
642 self.0.join(rel_path.replace('/', "\\"))
643 } else {
644 self.0.join(rel_path)
645 };
646 if let Some(parent) = full_path.parent() {
647 fs::create_dir_all(parent).await.context(CreateSnafu)?;
648 }
649 Ok(Box::new(
650 fs::File::create(full_path)
651 .await
652 .context(CreateSnafu)?
653 .compat_write(),
654 ) as BoxWriter)
655 }
656}
657
658#[cfg(test)]
659impl<H> BoundedStager<H> {
660 pub async fn must_get_file(&self, puffin_file_name: &str, key: &str) -> fs::File {
661 let cache_key = Self::encode_cache_key(puffin_file_name, key);
662 let value = self.cache.get(&cache_key).await.unwrap();
663 let path = match &value {
664 CacheValue::File(guard) => &guard.path,
665 _ => panic!("Expected a file, but got a directory."),
666 };
667 fs::File::open(path).await.unwrap()
668 }
669
670 pub async fn must_get_dir(&self, puffin_file_name: &str, key: &str) -> PathBuf {
671 let cache_key = Self::encode_cache_key(puffin_file_name, key);
672 let value = self.cache.get(&cache_key).await.unwrap();
673 let path = match &value {
674 CacheValue::Dir(guard) => &guard.path,
675 _ => panic!("Expected a directory, but got a file."),
676 };
677 path.clone()
678 }
679
680 pub fn in_cache(&self, puffin_file_name: &str, key: &str) -> bool {
681 let cache_key = Self::encode_cache_key(puffin_file_name, key);
682 self.cache.contains_key(&cache_key)
683 }
684}
685
686#[cfg(test)]
687mod tests {
688 use std::sync::atomic::AtomicU64;
689
690 use common_base::range_read::RangeReader;
691 use common_test_util::temp_dir::create_temp_dir;
692 use futures::AsyncWriteExt;
693 use tokio::io::AsyncReadExt as _;
694
695 use super::*;
696 use crate::error::BlobNotFoundSnafu;
697 use crate::puffin_manager::stager::Stager;
698
699 struct MockNotifier {
700 cache_insert_size: AtomicU64,
701 cache_evict_size: AtomicU64,
702 cache_hit_count: AtomicU64,
703 cache_hit_size: AtomicU64,
704 cache_miss_count: AtomicU64,
705 cache_miss_size: AtomicU64,
706 recycle_insert_size: AtomicU64,
707 recycle_clear_size: AtomicU64,
708 }
709
710 #[derive(Debug, PartialEq, Eq)]
711 struct Stats {
712 cache_insert_size: u64,
713 cache_evict_size: u64,
714 cache_hit_count: u64,
715 cache_hit_size: u64,
716 cache_miss_count: u64,
717 cache_miss_size: u64,
718 recycle_insert_size: u64,
719 recycle_clear_size: u64,
720 }
721
722 impl MockNotifier {
723 fn build() -> Arc<MockNotifier> {
724 Arc::new(Self {
725 cache_insert_size: AtomicU64::new(0),
726 cache_evict_size: AtomicU64::new(0),
727 cache_hit_count: AtomicU64::new(0),
728 cache_hit_size: AtomicU64::new(0),
729 cache_miss_count: AtomicU64::new(0),
730 cache_miss_size: AtomicU64::new(0),
731 recycle_insert_size: AtomicU64::new(0),
732 recycle_clear_size: AtomicU64::new(0),
733 })
734 }
735
736 fn stats(&self) -> Stats {
737 Stats {
738 cache_insert_size: self
739 .cache_insert_size
740 .load(std::sync::atomic::Ordering::Relaxed),
741 cache_evict_size: self
742 .cache_evict_size
743 .load(std::sync::atomic::Ordering::Relaxed),
744 cache_hit_count: self
745 .cache_hit_count
746 .load(std::sync::atomic::Ordering::Relaxed),
747 cache_hit_size: self
748 .cache_hit_size
749 .load(std::sync::atomic::Ordering::Relaxed),
750 cache_miss_count: self
751 .cache_miss_count
752 .load(std::sync::atomic::Ordering::Relaxed),
753 cache_miss_size: self
754 .cache_miss_size
755 .load(std::sync::atomic::Ordering::Relaxed),
756 recycle_insert_size: self
757 .recycle_insert_size
758 .load(std::sync::atomic::Ordering::Relaxed),
759 recycle_clear_size: self
760 .recycle_clear_size
761 .load(std::sync::atomic::Ordering::Relaxed),
762 }
763 }
764 }
765
766 impl StagerNotifier for MockNotifier {
767 fn on_cache_insert(&self, size: u64) {
768 self.cache_insert_size
769 .fetch_add(size, std::sync::atomic::Ordering::Relaxed);
770 }
771
772 fn on_cache_evict(&self, size: u64) {
773 self.cache_evict_size
774 .fetch_add(size, std::sync::atomic::Ordering::Relaxed);
775 }
776
777 fn on_cache_hit(&self, size: u64) {
778 self.cache_hit_count
779 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
780 self.cache_hit_size
781 .fetch_add(size, std::sync::atomic::Ordering::Relaxed);
782 }
783
784 fn on_cache_miss(&self, size: u64) {
785 self.cache_miss_count
786 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
787 self.cache_miss_size
788 .fetch_add(size, std::sync::atomic::Ordering::Relaxed);
789 }
790
791 fn on_recycle_insert(&self, size: u64) {
792 self.recycle_insert_size
793 .fetch_add(size, std::sync::atomic::Ordering::Relaxed);
794 }
795
796 fn on_recycle_clear(&self, size: u64) {
797 self.recycle_clear_size
798 .fetch_add(size, std::sync::atomic::Ordering::Relaxed);
799 }
800
801 fn on_load_blob(&self, _duration: Duration) {}
802
803 fn on_load_dir(&self, _duration: Duration) {}
804 }
805
806 #[tokio::test]
807 async fn test_get_blob() {
808 let tempdir = create_temp_dir("test_get_blob_");
809 let notifier = MockNotifier::build();
810 let stager = BoundedStager::new(
811 tempdir.path().to_path_buf(),
812 u64::MAX,
813 Some(notifier.clone()),
814 None,
815 )
816 .await
817 .unwrap();
818
819 let puffin_file_name = "test_get_blob".to_string();
820 let key = "key";
821 let reader = stager
822 .get_blob(
823 &puffin_file_name,
824 key,
825 Box::new(|mut writer| {
826 Box::pin(async move {
827 writer.write_all(b"hello world").await.unwrap();
828 Ok(11)
829 })
830 }),
831 )
832 .await
833 .unwrap()
834 .reader()
835 .await
836 .unwrap();
837
838 let m = reader.metadata().await.unwrap();
839 let buf = reader.read(0..m.content_length).await.unwrap();
840 assert_eq!(&*buf, b"hello world");
841
842 let mut file = stager.must_get_file(&puffin_file_name, key).await;
843 let mut buf = Vec::new();
844 file.read_to_end(&mut buf).await.unwrap();
845 assert_eq!(buf, b"hello world");
846
847 let stats = notifier.stats();
848 assert_eq!(
849 stats,
850 Stats {
851 cache_insert_size: 11,
852 cache_evict_size: 0,
853 cache_hit_count: 0,
854 cache_hit_size: 0,
855 cache_miss_count: 1,
856 cache_miss_size: 11,
857 recycle_insert_size: 0,
858 recycle_clear_size: 0,
859 }
860 );
861 }
862
863 #[tokio::test]
864 async fn test_get_dir() {
865 let tempdir = create_temp_dir("test_get_dir_");
866 let notifier = MockNotifier::build();
867 let stager = BoundedStager::new(
868 tempdir.path().to_path_buf(),
869 u64::MAX,
870 Some(notifier.clone()),
871 None,
872 )
873 .await
874 .unwrap();
875
876 let files_in_dir = [
877 ("file_a", "Hello, world!".as_bytes()),
878 ("file_b", "Hello, Rust!".as_bytes()),
879 ("file_c", "你好,世界!".as_bytes()),
880 ("subdir/file_d", "Hello, Puffin!".as_bytes()),
881 ("subdir/subsubdir/file_e", "¡Hola mundo!".as_bytes()),
882 ];
883
884 let puffin_file_name = "test_get_dir".to_string();
885 let key = "key";
886 let (dir_path, metrics) = stager
887 .get_dir(
888 &puffin_file_name,
889 key,
890 Box::new(|writer_provider| {
891 Box::pin(async move {
892 let mut size = 0;
893 for (rel_path, content) in &files_in_dir {
894 size += content.len();
895 let mut writer = writer_provider.writer(rel_path).await.unwrap();
896 writer.write_all(content).await.unwrap();
897 }
898 Ok(size as _)
899 })
900 }),
901 )
902 .await
903 .unwrap();
904
905 assert!(!metrics.cache_hit);
906 assert!(metrics.dir_size > 0);
907
908 for (rel_path, content) in &files_in_dir {
909 let file_path = dir_path.path().join(rel_path);
910 let mut file = tokio::fs::File::open(&file_path).await.unwrap();
911 let mut buf = Vec::new();
912 file.read_to_end(&mut buf).await.unwrap();
913 assert_eq!(buf, *content);
914 }
915
916 let dir_path = stager.must_get_dir(&puffin_file_name, key).await;
917 for (rel_path, content) in &files_in_dir {
918 let file_path = dir_path.join(rel_path);
919 let mut file = tokio::fs::File::open(&file_path).await.unwrap();
920 let mut buf = Vec::new();
921 file.read_to_end(&mut buf).await.unwrap();
922 assert_eq!(buf, *content);
923 }
924
925 let stats = notifier.stats();
926 assert_eq!(
927 stats,
928 Stats {
929 cache_insert_size: 70,
930 cache_evict_size: 0,
931 cache_hit_count: 0,
932 cache_hit_size: 0,
933 cache_miss_count: 1,
934 cache_miss_size: 70,
935 recycle_insert_size: 0,
936 recycle_clear_size: 0
937 }
938 );
939 }
940
941 #[tokio::test]
942 async fn test_recover() {
943 let tempdir = create_temp_dir("test_recover_");
944 let notifier = MockNotifier::build();
945 let stager = BoundedStager::new(
946 tempdir.path().to_path_buf(),
947 u64::MAX,
948 Some(notifier.clone()),
949 None,
950 )
951 .await
952 .unwrap();
953
954 let puffin_file_name = "test_recover".to_string();
956 let blob_key = "blob_key";
957 let guard = stager
958 .get_blob(
959 &puffin_file_name,
960 blob_key,
961 Box::new(|mut writer| {
962 Box::pin(async move {
963 writer.write_all(b"hello world").await.unwrap();
964 Ok(11)
965 })
966 }),
967 )
968 .await
969 .unwrap();
970 drop(guard);
971
972 let files_in_dir = [
973 ("file_a", "Hello, world!".as_bytes()),
974 ("file_b", "Hello, Rust!".as_bytes()),
975 ("file_c", "你好,世界!".as_bytes()),
976 ("subdir/file_d", "Hello, Puffin!".as_bytes()),
977 ("subdir/subsubdir/file_e", "¡Hola mundo!".as_bytes()),
978 ];
979
980 let dir_key = "dir_key";
981 let (guard, _metrics) = stager
982 .get_dir(
983 &puffin_file_name,
984 dir_key,
985 Box::new(|writer_provider| {
986 Box::pin(async move {
987 let mut size = 0;
988 for (rel_path, content) in &files_in_dir {
989 size += content.len();
990 let mut writer = writer_provider.writer(rel_path).await.unwrap();
991 writer.write_all(content).await.unwrap();
992 }
993 Ok(size as _)
994 })
995 }),
996 )
997 .await
998 .unwrap();
999 drop(guard);
1000
1001 drop(stager);
1003 let stager = BoundedStager::new(tempdir.path().to_path_buf(), u64::MAX, None, None)
1004 .await
1005 .unwrap();
1006
1007 let reader = stager
1008 .get_blob(
1009 &puffin_file_name,
1010 blob_key,
1011 Box::new(|_| Box::pin(async { Ok(0) })),
1012 )
1013 .await
1014 .unwrap()
1015 .reader()
1016 .await
1017 .unwrap();
1018
1019 let m = reader.metadata().await.unwrap();
1020 let buf = reader.read(0..m.content_length).await.unwrap();
1021 assert_eq!(&*buf, b"hello world");
1022
1023 let (dir_path, metrics) = stager
1024 .get_dir(
1025 &puffin_file_name,
1026 dir_key,
1027 Box::new(|_| Box::pin(async { Ok(0) })),
1028 )
1029 .await
1030 .unwrap();
1031
1032 assert!(metrics.cache_hit);
1033 assert!(metrics.dir_size > 0);
1034 for (rel_path, content) in &files_in_dir {
1035 let file_path = dir_path.path().join(rel_path);
1036 let mut file = tokio::fs::File::open(&file_path).await.unwrap();
1037 let mut buf = Vec::new();
1038 file.read_to_end(&mut buf).await.unwrap();
1039 assert_eq!(buf, *content);
1040 }
1041
1042 let stats = notifier.stats();
1043 assert_eq!(
1044 stats,
1045 Stats {
1046 cache_insert_size: 81,
1047 cache_evict_size: 0,
1048 cache_hit_count: 0,
1049 cache_hit_size: 0,
1050 cache_miss_count: 2,
1051 cache_miss_size: 81,
1052 recycle_insert_size: 0,
1053 recycle_clear_size: 0
1054 }
1055 );
1056 }
1057
1058 #[tokio::test]
1059 async fn test_eviction() {
1060 let tempdir = create_temp_dir("test_eviction_");
1061 let notifier = MockNotifier::build();
1062 let stager = BoundedStager::new(
1063 tempdir.path().to_path_buf(),
1064 1, Some(notifier.clone()),
1066 None,
1067 )
1068 .await
1069 .unwrap();
1070
1071 let puffin_file_name = "test_eviction".to_string();
1072 let blob_key = "blob_key";
1073
1074 let reader = stager
1076 .get_blob(
1077 &puffin_file_name,
1078 blob_key,
1079 Box::new(|mut writer| {
1080 Box::pin(async move {
1081 writer.write_all(b"Hello world").await.unwrap();
1082 Ok(11)
1083 })
1084 }),
1085 )
1086 .await
1087 .unwrap()
1088 .reader()
1089 .await
1090 .unwrap();
1091
1092 stager.cache.run_pending_tasks().await;
1094 assert!(!stager.in_cache(&puffin_file_name, blob_key));
1095
1096 let stats = notifier.stats();
1097 assert_eq!(
1098 stats,
1099 Stats {
1100 cache_insert_size: 11,
1101 cache_evict_size: 11,
1102 cache_hit_count: 0,
1103 cache_hit_size: 0,
1104 cache_miss_count: 1,
1105 cache_miss_size: 11,
1106 recycle_insert_size: 11,
1107 recycle_clear_size: 0
1108 }
1109 );
1110
1111 let m = reader.metadata().await.unwrap();
1112 let buf = reader.read(0..m.content_length).await.unwrap();
1113 assert_eq!(&*buf, b"Hello world");
1114
1115 let reader = stager
1117 .get_blob(
1118 &puffin_file_name,
1119 blob_key,
1120 Box::new(|_| async { Ok(0) }.boxed()),
1121 )
1122 .await
1123 .unwrap()
1124 .reader()
1125 .await
1126 .unwrap();
1127
1128 stager.cache.run_pending_tasks().await;
1130 assert!(!stager.in_cache(&puffin_file_name, blob_key));
1131
1132 let stats = notifier.stats();
1133 assert_eq!(
1134 stats,
1135 Stats {
1136 cache_insert_size: 22,
1137 cache_evict_size: 22,
1138 cache_hit_count: 1,
1139 cache_hit_size: 11,
1140 cache_miss_count: 1,
1141 cache_miss_size: 11,
1142 recycle_insert_size: 22,
1143 recycle_clear_size: 11
1144 }
1145 );
1146
1147 let m = reader.metadata().await.unwrap();
1148 let buf = reader.read(0..m.content_length).await.unwrap();
1149 assert_eq!(&*buf, b"Hello world");
1150
1151 let dir_key = "dir_key";
1152 let files_in_dir = [
1153 ("file_a", "Hello, world!".as_bytes()),
1154 ("file_b", "Hello, Rust!".as_bytes()),
1155 ("file_c", "你好,世界!".as_bytes()),
1156 ("subdir/file_d", "Hello, Puffin!".as_bytes()),
1157 ("subdir/subsubdir/file_e", "¡Hola mundo!".as_bytes()),
1158 ];
1159
1160 let (guard_0, _metrics) = stager
1162 .get_dir(
1163 &puffin_file_name,
1164 dir_key,
1165 Box::new(|writer_provider| {
1166 Box::pin(async move {
1167 let mut size = 0;
1168 for (rel_path, content) in &files_in_dir {
1169 let mut writer = writer_provider.writer(rel_path).await.unwrap();
1170 writer.write_all(content).await.unwrap();
1171 size += content.len() as u64;
1172 }
1173 Ok(size)
1174 })
1175 }),
1176 )
1177 .await
1178 .unwrap();
1179
1180 for (rel_path, content) in &files_in_dir {
1181 let file_path = guard_0.path().join(rel_path);
1182 let mut file = tokio::fs::File::open(&file_path).await.unwrap();
1183 let mut buf = Vec::new();
1184 file.read_to_end(&mut buf).await.unwrap();
1185 assert_eq!(buf, *content);
1186 }
1187
1188 stager.cache.run_pending_tasks().await;
1190 assert!(!stager.in_cache(&puffin_file_name, dir_key));
1191
1192 let stats = notifier.stats();
1193 assert_eq!(
1194 stats,
1195 Stats {
1196 cache_insert_size: 92,
1197 cache_evict_size: 92,
1198 cache_hit_count: 1,
1199 cache_hit_size: 11,
1200 cache_miss_count: 2,
1201 cache_miss_size: 81,
1202 recycle_insert_size: 92,
1203 recycle_clear_size: 11
1204 }
1205 );
1206
1207 let (guard_1, _metrics) = stager
1209 .get_dir(
1210 &puffin_file_name,
1211 dir_key,
1212 Box::new(|_| async { Ok(0) }.boxed()),
1213 )
1214 .await
1215 .unwrap();
1216
1217 for (rel_path, content) in &files_in_dir {
1218 let file_path = guard_1.path().join(rel_path);
1219 let mut file = tokio::fs::File::open(&file_path).await.unwrap();
1220 let mut buf = Vec::new();
1221 file.read_to_end(&mut buf).await.unwrap();
1222 assert_eq!(buf, *content);
1223 }
1224
1225 stager.cache.run_pending_tasks().await;
1227 assert!(!stager.in_cache(&puffin_file_name, dir_key));
1228
1229 let stats = notifier.stats();
1230 assert_eq!(
1231 stats,
1232 Stats {
1233 cache_insert_size: 162,
1234 cache_evict_size: 162,
1235 cache_hit_count: 2,
1236 cache_hit_size: 81,
1237 cache_miss_count: 2,
1238 cache_miss_size: 81,
1239 recycle_insert_size: 162,
1240 recycle_clear_size: 81
1241 }
1242 );
1243
1244 drop(guard_0);
1246 drop(guard_1);
1247 let (guard_2, _metrics) = stager
1248 .get_dir(
1249 &puffin_file_name,
1250 dir_key,
1251 Box::new(|_| Box::pin(async move { Ok(0) })),
1252 )
1253 .await
1254 .unwrap();
1255
1256 stager.cache.run_pending_tasks().await;
1258 assert!(!stager.in_cache(&puffin_file_name, blob_key));
1259
1260 for (rel_path, content) in &files_in_dir {
1261 let file_path = guard_2.path().join(rel_path);
1262 let mut file = tokio::fs::File::open(&file_path).await.unwrap();
1263 let mut buf = Vec::new();
1264 file.read_to_end(&mut buf).await.unwrap();
1265 assert_eq!(buf, *content);
1266 }
1267
1268 let stats = notifier.stats();
1269 assert_eq!(
1270 stats,
1271 Stats {
1272 cache_insert_size: 232,
1273 cache_evict_size: 232,
1274 cache_hit_count: 3,
1275 cache_hit_size: 151,
1276 cache_miss_count: 2,
1277 cache_miss_size: 81,
1278 recycle_insert_size: 232,
1279 recycle_clear_size: 151
1280 }
1281 );
1282 }
1283
1284 #[tokio::test]
1285 async fn test_get_blob_concurrency_on_fail() {
1286 let tempdir = create_temp_dir("test_get_blob_concurrency_on_fail_");
1287 let stager = BoundedStager::new(tempdir.path().to_path_buf(), u64::MAX, None, None)
1288 .await
1289 .unwrap();
1290
1291 let puffin_file_name = "test_get_blob_concurrency_on_fail".to_string();
1292 let key = "key";
1293
1294 let stager = Arc::new(stager);
1295 let handles = (0..10)
1296 .map(|_| {
1297 let stager = stager.clone();
1298 let puffin_file_name = puffin_file_name.clone();
1299 let task = async move {
1300 let failed_init = Box::new(|_| {
1301 async {
1302 tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
1303 BlobNotFoundSnafu { blob: "whatever" }.fail()
1304 }
1305 .boxed()
1306 });
1307 stager.get_blob(&puffin_file_name, key, failed_init).await
1308 };
1309
1310 tokio::spawn(task)
1311 })
1312 .collect::<Vec<_>>();
1313
1314 for handle in handles {
1315 let r = handle.await.unwrap();
1316 assert!(r.is_err());
1317 }
1318
1319 assert!(!stager.in_cache(&puffin_file_name, key));
1320 }
1321
1322 #[tokio::test]
1323 async fn test_get_dir_concurrency_on_fail() {
1324 let tempdir = create_temp_dir("test_get_dir_concurrency_on_fail_");
1325 let stager = BoundedStager::new(tempdir.path().to_path_buf(), u64::MAX, None, None)
1326 .await
1327 .unwrap();
1328
1329 let puffin_file_name = "test_get_dir_concurrency_on_fail".to_string();
1330 let key = "key";
1331
1332 let stager = Arc::new(stager);
1333 let handles = (0..10)
1334 .map(|_| {
1335 let stager = stager.clone();
1336 let puffin_file_name = puffin_file_name.clone();
1337 let task = async move {
1338 let failed_init = Box::new(|_| {
1339 async {
1340 tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
1341 BlobNotFoundSnafu { blob: "whatever" }.fail()
1342 }
1343 .boxed()
1344 });
1345 stager.get_dir(&puffin_file_name, key, failed_init).await
1346 };
1347
1348 tokio::spawn(task)
1349 })
1350 .collect::<Vec<_>>();
1351
1352 for handle in handles {
1353 let r = handle.await.unwrap();
1354 assert!(r.is_err());
1355 }
1356
1357 assert!(!stager.in_cache(&puffin_file_name, key));
1358 }
1359
1360 #[tokio::test]
1361 async fn test_purge() {
1362 let tempdir = create_temp_dir("test_purge_");
1363 let notifier = MockNotifier::build();
1364 let stager = BoundedStager::new(
1365 tempdir.path().to_path_buf(),
1366 u64::MAX,
1367 Some(notifier.clone()),
1368 None,
1369 )
1370 .await
1371 .unwrap();
1372
1373 let puffin_file_name = "test_purge".to_string();
1375 let blob_key = "blob_key";
1376 let guard = stager
1377 .get_blob(
1378 &puffin_file_name,
1379 blob_key,
1380 Box::new(|mut writer| {
1381 Box::pin(async move {
1382 writer.write_all(b"hello world").await.unwrap();
1383 Ok(11)
1384 })
1385 }),
1386 )
1387 .await
1388 .unwrap();
1389 drop(guard);
1390
1391 let files_in_dir = [
1392 ("file_a", "Hello, world!".as_bytes()),
1393 ("file_b", "Hello, Rust!".as_bytes()),
1394 ("file_c", "你好,世界!".as_bytes()),
1395 ("subdir/file_d", "Hello, Puffin!".as_bytes()),
1396 ("subdir/subsubdir/file_e", "¡Hola mundo!".as_bytes()),
1397 ];
1398
1399 let dir_key = "dir_key";
1400 let (guard, _metrics) = stager
1401 .get_dir(
1402 &puffin_file_name,
1403 dir_key,
1404 Box::new(|writer_provider| {
1405 Box::pin(async move {
1406 let mut size = 0;
1407 for (rel_path, content) in &files_in_dir {
1408 size += content.len();
1409 let mut writer = writer_provider.writer(rel_path).await.unwrap();
1410 writer.write_all(content).await.unwrap();
1411 }
1412 Ok(size as _)
1413 })
1414 }),
1415 )
1416 .await
1417 .unwrap();
1418 drop(guard);
1419
1420 stager.purge(&puffin_file_name).await.unwrap();
1422
1423 let stats = notifier.stats();
1424 assert_eq!(
1425 stats,
1426 Stats {
1427 cache_insert_size: 81,
1428 cache_evict_size: 81,
1429 cache_hit_count: 0,
1430 cache_hit_size: 0,
1431 cache_miss_count: 2,
1432 cache_miss_size: 81,
1433 recycle_insert_size: 81,
1434 recycle_clear_size: 0
1435 }
1436 );
1437 }
1438}