Skip to main content

puffin/puffin_manager/stager/
bounded_stager.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
51/// `BoundedStager` is a `Stager` that uses `moka` to manage staging area.
52pub struct BoundedStager<H> {
53    /// The base directory of the staging area.
54    base_dir: PathBuf,
55
56    /// The cache maintaining the cache key to the size of the file or directory.
57    cache: Cache<String, CacheValue>,
58
59    /// The recycle bin for the deleted files and directories.
60    recycle_bin: Cache<String, CacheValue>,
61
62    /// The delete queue for the cleanup task.
63    ///
64    /// The lifetime of a guard is:
65    ///   1. initially inserted into the cache
66    ///   2. moved to the recycle bin when evicted
67    ///      2.1 moved back to the cache when accessed
68    ///      2.2 deleted from the recycle bin after a certain period
69    ///   3. sent the delete task to the delete queue on drop
70    ///   4. background routine removes the file or directory
71    delete_queue: Sender<DeleteTask>,
72
73    /// Notifier for the stager.
74    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        // Dir is usually large.
306        // Runs pending tasks of the cache and recycle bin to free up space
307        // more quickly.
308        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(); // SAFETY: `support_invalidation_closures` is enabled
319        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        // To guarantee the atomicity of writing the file, we need to write
337        // the file to a temporary file first...
338        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        // ...then rename the temporary file to the target path
348        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        // To guarantee the atomicity of writing the directory, we need to write
356        // the directory to a temporary directory first...
357        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        // ...then rename the temporary directory to the target path
362        fs::rename(&tmp_base, target_path)
363            .await
364            .context(RenameSnafu)?;
365        Ok(size)
366    }
367
368    /// Recovers the staging area by iterating through the staging directory.
369    ///
370    /// Note: It can't recover the mapping between puffin files and keys, so TTL
371    ///       is configured to purge the dangling files and directories.
372    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                // Remove temporary or deleted files and directories
386                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                // Insert the guard of the file or directory to the cache
393                let meta = entry.metadata().await.context(MetadataSnafu)?;
394                let file_path = path.file_name().unwrap().to_string_lossy().into_owned();
395
396                // <key>.<uuid>
397                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                        // placeholder
416                        handle: String::new(),
417                    }));
418                    // A duplicate dir will be moved to the delete queue.
419                    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                        // placeholder
428                        handle: String::new(),
429                    }));
430                    // A duplicate file will be moved to the delete queue.
431                    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    /// Walks through the directory and calculate the total size of all files in the directory.
458    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                            // Remove the deleted directory if the rename fails and retry
505                            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                    // Purge recycle bin periodically to reclaim the space quickly.
525                    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/// `FsBlobGuard` is a `BlobGuard` for accessing the blob and
573/// automatically deleting the file on drop.
574#[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/// `FsDirGuard` is a `DirGuard` for accessing the directory and
606/// automatically deleting the directory on drop.
607#[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
635/// `MokaDirWriterProvider` implements `DirWriterProvider` for initializing a directory.
636struct 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        // initialize stager
955        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        // recover stager
1002        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, /* extremely small size */
1065            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        // First time to get the blob
1075        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        // The blob should be evicted
1093        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        // Second time to get the blob, get from recycle bin
1116        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        // The blob should be evicted
1129        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        // First time to get the directory
1161        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        // The directory should be evicted
1189        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        // Second time to get the directory
1208        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        // Still hold the guard
1226        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        // Third time to get the directory and all guards are dropped
1245        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        // Still hold the guard, so the directory should not be removed even if it's evicted
1257        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        // initialize stager
1374        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        // purge the stager
1421        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}