Skip to main content

foyer_storage/engine/block/
engine.rs

1// Copyright 2026 foyer Project Authors
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::{
16    fmt::Debug,
17    future::Future,
18    marker::PhantomData,
19    sync::{
20        Arc,
21        atomic::{AtomicBool, AtomicUsize, Ordering},
22    },
23    time::Instant,
24};
25
26use asyncband::mpsc::UnboundedReceiver;
27#[cfg(feature = "tracing")]
28use fastrace::prelude::*;
29use foyer_common::{
30    bits,
31    code::{StorageKey, StorageValue},
32    error::{Error, ErrorKind, Result},
33    metrics::Metrics,
34    properties::{Age, Properties},
35    spawn::Spawner,
36};
37use futures_core::future::BoxFuture;
38use futures_util::{
39    FutureExt,
40    future::{join_all, try_join_all},
41};
42use itertools::Itertools;
43
44use super::{
45    flusher::{Flusher, InvalidStats, Submission},
46    indexer::Indexer,
47    recover::RecoverRunner,
48};
49#[cfg(any(test, feature = "test_utils"))]
50use crate::test_utils::*;
51use crate::{
52    Device, Load, RejectAll, StorageFilter, StorageFilterResult,
53    compress::Compression,
54    engine::{
55        Engine, EngineBuildContext, EngineConfig, Populated,
56        block::{
57            eviction::{EvictionPicker, FifoPicker, InvalidRatioPicker},
58            manager::{BlockId, BlockManager},
59            reclaimer::{BlockCleaner, Reclaimer, ReclaimerTrait},
60            serde::{AtomicSequence, EntryHeader},
61            tombstone::{Tombstone, TombstoneLog},
62        },
63    },
64    filter::conditions::IoThrottle,
65    io::{PAGE, bytes::IoSliceMut},
66    keeper::PieceRef,
67    serde::EntryDeserializer,
68};
69
70/// Config for the block-based disk cache engine.
71///
72/// The block-based disk cache engine is suitable for general cache entries with size from 2K to hundreds of MiBs.
73///
74/// Each cache entry will be aligned to a multiplier of 4K on disk, hence too small cache entries will lead to heavy
75/// internal fragmentation.
76///
77/// The disk cache evicts cache entries in block unit.
78pub struct BlockEngineConfig<K, V, P>
79where
80    K: StorageKey,
81    V: StorageValue,
82    P: Properties,
83{
84    device: Arc<dyn Device>,
85    block_size: usize,
86    compression: Compression,
87    indexer_shards: usize,
88    recover_concurrency: usize,
89    flushers: usize,
90    reclaimers: usize,
91    buffer_pool_size: usize,
92    blob_index_size: usize,
93    submit_queue_size_threshold: usize,
94    clean_block_threshold: usize,
95    eviction_pickers: Vec<Box<dyn EvictionPicker>>,
96    admission_filter: StorageFilter,
97    reinsertion_filter: StorageFilter,
98    enable_tombstone_log: bool,
99    #[cfg(any(test, feature = "test_utils"))]
100    flush_switch: Switch,
101    #[cfg(any(test, feature = "test_utils"))]
102    load_holder: Holder,
103    marker: PhantomData<(K, V, P)>,
104}
105
106impl<K, V, P> Debug for BlockEngineConfig<K, V, P>
107where
108    K: StorageKey,
109    V: StorageValue,
110    P: Properties,
111{
112    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
113        f.debug_struct("BlockEngineConfig")
114            .field("device", &self.device)
115            .field("block_size", &self.block_size)
116            .field("compression", &self.compression)
117            .field("indexer_shards", &self.indexer_shards)
118            .field("recover_concurrency", &self.recover_concurrency)
119            .field("flushers", &self.flushers)
120            .field("reclaimers", &self.reclaimers)
121            .field("buffer_pool_size", &self.buffer_pool_size)
122            .field("blob_index_size", &self.blob_index_size)
123            .field("submit_queue_size_threshold", &self.submit_queue_size_threshold)
124            .field("clean_block_threshold", &self.clean_block_threshold)
125            .field("eviction_pickers", &self.eviction_pickers)
126            .field("admission_filter", &self.admission_filter)
127            .field("reinsertion_filter", &self.reinsertion_filter)
128            .field("enable_tombstone_log", &self.enable_tombstone_log)
129            .finish()
130    }
131}
132
133impl<K, V, P> BlockEngineConfig<K, V, P>
134where
135    K: StorageKey,
136    V: StorageValue,
137    P: Properties,
138{
139    /// Create a new block-based disk cache engine builder with default configurations.
140    pub fn new(device: Arc<dyn Device>) -> Self {
141        Self {
142            device,
143            block_size: 16 * 1024 * 1024, // 16 MiB
144            compression: Compression::default(),
145            indexer_shards: 64,
146            recover_concurrency: 8,
147            flushers: 1,
148            reclaimers: 1,
149            buffer_pool_size: 16 * 1024 * 1024,            // 16 MiB
150            blob_index_size: 4 * 1024,                     // 4 KiB
151            submit_queue_size_threshold: 16 * 1024 * 1024, // 16 MiB
152            clean_block_threshold: 1,
153            eviction_pickers: vec![Box::new(InvalidRatioPicker::new(0.8)), Box::<FifoPicker>::default()],
154            admission_filter: StorageFilter::new(),
155            reinsertion_filter: StorageFilter::new().with_condition(RejectAll),
156            enable_tombstone_log: false,
157            #[cfg(any(test, feature = "test_utils"))]
158            flush_switch: Switch::default(),
159            #[cfg(any(test, feature = "test_utils"))]
160            load_holder: Holder::default(),
161            marker: PhantomData,
162        }
163    }
164
165    /// Set the block size for the block-based disk cache engine.
166    ///
167    /// Block is the minimal cache eviction unit for the block-based disk cache,
168    /// its size also limits the max cacheable entry size.
169    ///
170    /// The block size must be 4K-aligned. the given value is not 4K-aligned, it will be automatically aligned up.
171    ///
172    /// Default: `16 MiB`.
173    pub fn with_block_size(mut self, block_size: usize) -> Self {
174        self.block_size = bits::align_up(PAGE, block_size);
175        self
176    }
177
178    /// Set the shard num of the indexer. Each shard has its own lock.
179    ///
180    /// Default: `64`.
181    pub fn with_indexer_shards(mut self, indexer_shards: usize) -> Self {
182        self.indexer_shards = indexer_shards;
183        self
184    }
185
186    /// Set the recover concurrency for the disk cache store.
187    ///
188    /// Default: `8`.
189    pub fn with_recover_concurrency(mut self, recover_concurrency: usize) -> Self {
190        self.recover_concurrency = recover_concurrency;
191        self
192    }
193
194    /// Set the flusher count for the disk cache store.
195    ///
196    /// The flusher count limits how many blocks can be concurrently written.
197    ///
198    /// Default: `1`.
199    pub fn with_flushers(mut self, flushers: usize) -> Self {
200        self.flushers = flushers;
201        self
202    }
203
204    /// Set the admission filter for th disk cache store.
205    ///
206    /// The admission filter is used to pick the entries that can be inserted into the disk cache store.
207    ///
208    /// Default: Admit all.
209    pub fn with_admission_filter(mut self, filter: StorageFilter) -> Self {
210        self.admission_filter = filter;
211        self
212    }
213
214    /// Set the reclaimer count for the disk cache store.
215    ///
216    /// The reclaimer count limits how many blocks can be concurrently reclaimed.
217    ///
218    /// Default: `1`.
219    pub fn with_reclaimers(mut self, reclaimers: usize) -> Self {
220        self.reclaimers = reclaimers;
221        self
222    }
223
224    /// Set the total flush buffer pool size.
225    ///
226    /// Each flusher shares a volume at `threshold / flushers`.
227    ///
228    /// If the buffer of the flush queue exceeds the threshold, the further entries will be ignored.
229    ///
230    /// Default: 16 MiB.
231    pub fn with_buffer_pool_size(mut self, buffer_pool_size: usize) -> Self {
232        self.buffer_pool_size = buffer_pool_size;
233        self
234    }
235
236    /// Set the blob index size for each blob.
237    ///
238    /// A larger blob index size can hold more blob entries, but it will also increase the io size of each blob part
239    /// write.
240    ///
241    /// NOTE: The size will be aligned up to a multiplier of 4K.
242    ///
243    /// Default: 4 KiB
244    pub fn with_blob_index_size(mut self, blob_index_size: usize) -> Self {
245        let blob_index_size = bits::align_up(PAGE, blob_index_size);
246        self.blob_index_size = blob_index_size;
247        self
248    }
249
250    /// Set the submit queue size threshold.
251    ///
252    /// If the total entry estimated size in the submit queue exceeds the threshold, the further entries will be
253    /// ignored.
254    ///
255    /// Default: `buffer_pool_size` * 2.
256    pub fn with_submit_queue_size_threshold(mut self, submit_queue_size_threshold: usize) -> Self {
257        self.submit_queue_size_threshold = submit_queue_size_threshold;
258        self
259    }
260
261    /// Set the clean block threshold for the disk cache store.
262    ///
263    /// The reclaimers only work when the clean block count is equal to or lower than the clean block threshold.
264    ///
265    /// Default: the same value as the `reclaimers`.
266    pub fn with_clean_block_threshold(mut self, clean_block_threshold: usize) -> Self {
267        self.clean_block_threshold = clean_block_threshold;
268        self
269    }
270
271    /// Set the eviction pickers for th disk cache store.
272    ///
273    /// The eviction picker is used to pick the block to reclaim.
274    ///
275    /// The eviction pickers are applied in order. If the previous eviction picker doesn't pick any block, the next one
276    /// will be applied.
277    ///
278    /// If no eviction picker picks a block, a block will be picked randomly.
279    ///
280    /// Default: [ invalid ratio picker { threshold = 0.8 }, fifo picker ]
281    pub fn with_eviction_pickers(mut self, eviction_pickers: Vec<Box<dyn EvictionPicker>>) -> Self {
282        self.eviction_pickers = eviction_pickers;
283        self
284    }
285
286    /// Set the reinsertion filter for th disk cache store.
287    ///
288    /// The reinsertion filter is used to pick the entries that can be reinsertion into the disk cache store while
289    /// reclaiming.
290    ///
291    /// Note: Only extremely important entries should be picked. If too many entries are picked, both insertion and
292    /// reinsertion will be stuck.
293    ///
294    /// Default: Reject all.
295    pub fn with_reinsertion_filter(mut self, filter: StorageFilter) -> Self {
296        self.reinsertion_filter = filter;
297        self
298    }
299
300    /// Enable the tombstone log.
301    ///
302    /// For updatable cache, either the tombstone log or [`crate::engine::RecoverMode::None`] must be enabled to prevent
303    /// from the phantom entries after reopen.
304    pub fn with_tombstone_log(mut self, enable: bool) -> Self {
305        self.enable_tombstone_log = enable;
306        self
307    }
308
309    /// Pass the flush holder for test.
310    #[cfg(any(test, feature = "test_utils"))]
311    pub fn with_flush_switch(mut self, flush_switch: Switch) -> Self {
312        self.flush_switch = flush_switch;
313        self
314    }
315
316    /// Pass the load holder for test.
317    #[cfg(any(test, feature = "test_utils"))]
318    pub fn with_load_holder(mut self, load_holder: Holder) -> Self {
319        self.load_holder = load_holder;
320        self
321    }
322
323    /// Build the block-based disk cache engine with the given configurations.
324    pub async fn build(
325        self: Box<Self>,
326        EngineBuildContext {
327            io_engine,
328            metrics,
329            spawner: runtime,
330            recover_mode,
331        }: EngineBuildContext,
332    ) -> Result<Arc<BlockEngine<K, V, P>>> {
333        let device = self.device;
334        let block_size = self.block_size;
335
336        let mut tombstones = vec![];
337
338        let tombstone_log = if self.enable_tombstone_log {
339            // TODO(MrCroxx): The tombstone log support multiples partitions for multiple device support.
340            let mut partitions = vec![];
341
342            let max_entries = device.capacity() / PAGE;
343            let pages = max_entries / TombstoneLog::SLOTS_PER_PAGE
344                + if max_entries.is_multiple_of(TombstoneLog::SLOTS_PER_PAGE) {
345                    0
346                } else {
347                    1
348                };
349            let partition = device.create_partition(pages * PAGE)?;
350            partitions.push(partition);
351
352            let tombstone_log = TombstoneLog::open(partitions, io_engine.clone(), &mut tombstones).await?;
353            Some(tombstone_log)
354        } else {
355            None
356        };
357
358        let indexer = Indexer::new(self.indexer_shards);
359        let submit_queue_size = Arc::<AtomicUsize>::default();
360
361        #[expect(clippy::type_complexity)]
362        let (flushers, rxs): (Vec<Flusher<K, V, P>>, Vec<UnboundedReceiver<Submission<K, V, P>>>) = (0..self.flushers)
363            .map(|id| Flusher::<K, V, P>::new(id, submit_queue_size.clone(), metrics.clone()))
364            .unzip();
365
366        let reclaimer = Reclaimer::new(
367            indexer.clone(),
368            flushers.clone(),
369            Arc::new(self.reinsertion_filter),
370            self.blob_index_size,
371            device.statistics().clone(),
372            runtime.clone(),
373        );
374        let reclaimer: Arc<dyn ReclaimerTrait> = Arc::new(reclaimer);
375
376        let block_manager = BlockManager::open(
377            device.clone(),
378            io_engine,
379            block_size,
380            self.eviction_pickers,
381            reclaimer,
382            self.reclaimers,
383            self.clean_block_threshold,
384            metrics.clone(),
385            runtime.clone(),
386        )?;
387        let blocks = block_manager.blocks();
388
389        if self.flushers + self.clean_block_threshold > blocks / 2 {
390            tracing::warn!(
391                "[block engine]: block-based object disk cache stable blocks count is too small, flusher [{flushers}] + clean block threshold [{clean_block_threshold}] (default = reclaimers) is supposed to be much larger than the block count [{blocks}]",
392                flushers = self.flushers,
393                clean_block_threshold = self.clean_block_threshold,
394            );
395        }
396
397        let sequence = AtomicSequence::default();
398
399        RecoverRunner::run(
400            self.recover_concurrency,
401            recover_mode,
402            self.blob_index_size,
403            (0..blocks as BlockId).collect_vec(),
404            &sequence,
405            &indexer,
406            &block_manager,
407            &tombstones,
408            runtime.clone(),
409            metrics.clone(),
410        )
411        .await?;
412
413        let io_buffer_size = self.buffer_pool_size / self.flushers;
414        for (flusher, rx) in flushers.iter().zip(rxs) {
415            flusher.run(
416                rx,
417                block_size,
418                io_buffer_size,
419                self.blob_index_size,
420                self.compression,
421                indexer.clone(),
422                block_manager.clone(),
423                tombstone_log.clone(),
424                metrics.clone(),
425                &runtime,
426                #[cfg(any(test, feature = "test_utils"))]
427                self.flush_switch.clone(),
428            )?;
429        }
430
431        let admission_filter = self.admission_filter.with_condition(IoThrottle);
432
433        let inner = BlockEngineInner {
434            admission_filter,
435            device,
436            indexer,
437            block_manager,
438            flushers,
439            submit_queue_size,
440            submit_queue_size_threshold: self.submit_queue_size_threshold,
441            sequence,
442            _spawner: runtime,
443            active: AtomicBool::new(true),
444            metrics,
445            #[cfg(any(test, feature = "test_utils"))]
446            flush_switch: self.flush_switch,
447            #[cfg(any(test, feature = "test_utils"))]
448            load_holder: self.load_holder,
449        };
450        let inner = Arc::new(inner);
451        let engine = BlockEngine { inner };
452        let engine = Arc::new(engine);
453        Ok(engine)
454    }
455}
456
457impl<K, V, P> EngineConfig<K, V, P> for BlockEngineConfig<K, V, P>
458where
459    K: StorageKey,
460    V: StorageValue,
461    P: Properties,
462{
463    fn build(self: Box<Self>, ctx: EngineBuildContext) -> BoxFuture<'static, Result<Arc<dyn Engine<K, V, P>>>> {
464        async move { self.build(ctx).await.map(|e| e as Arc<dyn Engine<K, V, P>>) }.boxed()
465    }
466}
467
468impl<K, V, P> From<BlockEngineConfig<K, V, P>> for Box<dyn EngineConfig<K, V, P>>
469where
470    K: StorageKey,
471    V: StorageValue,
472    P: Properties,
473{
474    fn from(builder: BlockEngineConfig<K, V, P>) -> Self {
475        builder.boxed()
476    }
477}
478
479/// Block-based disk cache engine.
480pub struct BlockEngine<K, V, P>
481where
482    K: StorageKey,
483    V: StorageValue,
484    P: Properties,
485{
486    inner: Arc<BlockEngineInner<K, V, P>>,
487}
488
489impl<K, V, P> Debug for BlockEngine<K, V, P>
490where
491    K: StorageKey,
492    V: StorageValue,
493    P: Properties,
494{
495    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
496        f.debug_struct("BlockEngine").finish()
497    }
498}
499
500struct BlockEngineInner<K, V, P>
501where
502    K: StorageKey,
503    V: StorageValue,
504    P: Properties,
505{
506    admission_filter: StorageFilter,
507
508    device: Arc<dyn Device>,
509
510    indexer: Indexer,
511    block_manager: BlockManager,
512
513    flushers: Vec<Flusher<K, V, P>>,
514
515    submit_queue_size: Arc<AtomicUsize>,
516    submit_queue_size_threshold: usize,
517
518    sequence: AtomicSequence,
519
520    _spawner: Spawner,
521
522    active: AtomicBool,
523
524    metrics: Arc<Metrics>,
525
526    #[cfg(any(test, feature = "test_utils"))]
527    flush_switch: Switch,
528
529    #[cfg(any(test, feature = "test_utils"))]
530    load_holder: Holder,
531}
532
533impl<K, V, P> Clone for BlockEngine<K, V, P>
534where
535    K: StorageKey,
536    V: StorageValue,
537    P: Properties,
538{
539    fn clone(&self) -> Self {
540        Self {
541            inner: self.inner.clone(),
542        }
543    }
544}
545
546impl<K, V, P> BlockEngine<K, V, P>
547where
548    K: StorageKey,
549    V: StorageValue,
550    P: Properties,
551{
552    fn wait(&self) -> impl Future<Output = ()> + Send + 'static {
553        let flushers = self.inner.flushers.clone();
554        let block_manager = self.inner.block_manager.clone();
555        async move {
556            join_all(flushers.iter().map(|flusher| flusher.wait())).await;
557            block_manager.wait_reclaim().await;
558        }
559    }
560
561    fn close(&self) -> BoxFuture<'static, Result<()>> {
562        let this = self.clone();
563        async move {
564            this.inner.active.store(false, Ordering::Relaxed);
565            this.wait().await;
566            Ok(())
567        }
568        .boxed()
569    }
570
571    #[cfg_attr(feature = "tracing", trace(name = "foyer::storage::engine::block::engine::enqueue"))]
572    fn enqueue(&self, piece: PieceRef<K, V, P>, estimated_size: usize) {
573        if !self.inner.active.load(Ordering::Relaxed) {
574            tracing::warn!("cannot enqueue new entry after closed");
575            return;
576        }
577
578        tracing::trace!(
579            hash = piece.hash(),
580            age = ?piece.properties().age().unwrap_or_default(),
581            "[block engine]: enqueue"
582        );
583        match piece.properties().age().unwrap_or_default() {
584            Age::Fresh | Age::Old => {}
585            Age::Young => {
586                // skip write block engine if the entry is still young
587                self.inner.metrics.storage_block_engine_enqueue_skip.increase(1);
588                return;
589            }
590        }
591
592        if self.inner.submit_queue_size.load(Ordering::Relaxed) > self.inner.submit_queue_size_threshold {
593            self.inner.metrics.storage_queue_channel_overflow.increase(1);
594            return;
595        }
596
597        let sequence = self.inner.sequence.fetch_add(1, Ordering::Relaxed);
598
599        self.inner.flushers[piece.hash() as usize % self.inner.flushers.len()].submit(Submission::CacheEntry {
600            piece,
601            estimated_size,
602            sequence,
603        });
604    }
605
606    fn load(&self, hash: u64) -> impl Future<Output = Result<Load<K, V, P>>> + Send + 'static {
607        tracing::trace!(hash, "[block engine]: load");
608
609        #[cfg(any(test, feature = "test_utils"))]
610        let load_holer = self.inner.load_holder.wait();
611
612        let indexer = self.inner.indexer.clone();
613        let metrics = self.inner.metrics.clone();
614        let block_manager = self.inner.block_manager.clone();
615
616        let load = async move {
617            #[cfg(any(test, feature = "test_utils"))]
618            load_holer.await;
619
620            let addr = match indexer.get(hash) {
621                Some(addr) => addr,
622                None => {
623                    return Ok(Load::Miss);
624                }
625            };
626
627            tracing::trace!(hash, ?addr, "[block engine]: load");
628
629            let block = block_manager.block(addr.block);
630            if block.partition().statistics().is_read_throttled() {
631                return Ok(Load::Throttled);
632            }
633
634            let buf = IoSliceMut::new(bits::align_up(PAGE, addr.len as _));
635            let (buf, res) = block.read(Box::new(buf), addr.offset as _).await;
636            match res {
637                Ok(_) => {}
638                Err(e) => {
639                    tracing::error!(hash, ?addr, ?e, "[block engine load]: load error");
640                    return Err(e);
641                }
642            }
643
644            let header = match EntryHeader::read(&buf[..EntryHeader::serialized_len()]) {
645                Ok(header) => header,
646                Err(e) => {
647                    return match e.kind() {
648                        ErrorKind::Parse
649                        | ErrorKind::MagicMismatch
650                        | ErrorKind::ChecksumMismatch
651                        | ErrorKind::OutOfRange => {
652                            tracing::warn!(
653                                hash,
654                                ?addr,
655                                ?e,
656                                "[block engine load]: deserialize read buffer raise error, remove this entry and skip"
657                            );
658                            indexer.remove(hash);
659                            Ok(Load::Miss)
660                        }
661                        _ => {
662                            tracing::error!(hash, ?addr, ?e, "[block engine load]: load error");
663                            Err(e)
664                        }
665                    };
666                }
667            };
668
669            let (key, value) = {
670                let now = Instant::now();
671                let res = match EntryDeserializer::deserialize::<K, V>(
672                    &buf[EntryHeader::serialized_len()..],
673                    header.key_len as _,
674                    header.value_len as _,
675                    header.compression,
676                    Some(header.checksum),
677                ) {
678                    Ok(res) => res,
679                    Err(e) => {
680                        return match e.kind() {
681                            ErrorKind::MagicMismatch | ErrorKind::ChecksumMismatch | ErrorKind::OutOfRange => {
682                                tracing::warn!(
683                                    hash,
684                                    ?addr,
685                                    ?header,
686                                    ?e,
687                                    "[block engine load]: deserialize read buffer raise error, remove this entry and skip"
688                                );
689                                indexer.remove(hash);
690                                Ok(Load::Miss)
691                            }
692                            _ => {
693                                tracing::error!(hash, ?addr, ?header, ?e, "[block engine load]: load error");
694                                Err(e)
695                            }
696                        };
697                    }
698                };
699                metrics
700                    .storage_entry_deserialize_duration
701                    .record(now.elapsed().as_secs_f64());
702                res
703            };
704
705            let age = match block.statistics().probation.load(Ordering::Relaxed) {
706                true => Age::Old,
707                false => Age::Young,
708            };
709
710            Ok(Load::Entry {
711                key,
712                value,
713                populated: Populated { age },
714            })
715        };
716        #[cfg(feature = "tracing")]
717        let load = load.in_span(Span::enter_with_local_parent(
718            "foyer::storage::engine::block::engine::load",
719        ));
720        load
721    }
722
723    fn delete(&self, hash: u64) {
724        if !self.inner.active.load(Ordering::Relaxed) {
725            tracing::warn!("cannot delete entry after closed");
726            return;
727        }
728
729        let sequence = self.inner.sequence.fetch_add(1, Ordering::Relaxed);
730        let stats = self
731            .inner
732            .indexer
733            .insert_tombstone(hash, sequence)
734            .map(|addr| InvalidStats {
735                block: addr.block,
736                size: bits::align_up(PAGE, addr.len as usize),
737            });
738
739        let this = self.clone();
740
741        this.inner.flushers[hash as usize % this.inner.flushers.len()].submit(Submission::Tombstone {
742            tombstone: Tombstone { hash, sequence },
743            stats,
744        });
745    }
746
747    fn may_contains(&self, hash: u64) -> bool {
748        self.inner.indexer.get(hash).is_some()
749    }
750
751    fn destroy(&self) -> BoxFuture<'static, Result<()>> {
752        let this = self.clone();
753        async move {
754            if !this.inner.active.load(Ordering::Relaxed) {
755                return Err(Error::new(ErrorKind::Closed, "cannot delete entry after closed"));
756            }
757
758            // Write a tombstone to clear tombstone log by increase the max sequence.
759            let sequence = this.inner.sequence.fetch_add(1, Ordering::Relaxed);
760
761            this.inner.flushers[0].submit(Submission::Tombstone {
762                tombstone: Tombstone { hash: 0, sequence },
763                stats: None,
764            });
765            this.wait().await;
766
767            // Clear indices.
768            //
769            // This step must perform after the latest writer finished,
770            // otherwise the indices of the latest batch cannot be cleared.
771            this.inner.indexer.clear();
772
773            // Clean blocks.
774            try_join_all((0..this.inner.block_manager.blocks() as BlockId).map(|id| {
775                let block = this.inner.block_manager.block(id).clone();
776                async move {
777                    let res = BlockCleaner::clean(&block).await;
778                    block.statistics().reset();
779                    res
780                }
781            }))
782            .await?;
783
784            Ok(())
785        }
786        .boxed()
787    }
788
789    #[cfg(any(test, feature = "test_utils"))]
790    pub fn hold_flush(&self) {
791        self.inner.flush_switch.on();
792    }
793
794    #[cfg(any(test, feature = "test_utils"))]
795    pub fn unhold_flush(&self) {
796        self.inner.flush_switch.off();
797    }
798}
799
800impl<K, V, P> Engine<K, V, P> for BlockEngine<K, V, P>
801where
802    K: StorageKey,
803    V: StorageValue,
804    P: Properties,
805{
806    fn device(&self) -> &Arc<dyn Device> {
807        &self.inner.device
808    }
809
810    fn filter(&self, hash: u64, estimated_size: usize) -> StorageFilterResult {
811        self.inner
812            .admission_filter
813            .filter(self.inner.device.statistics(), hash, estimated_size)
814    }
815
816    fn enqueue(&self, piece: PieceRef<K, V, P>, estimated_size: usize) {
817        self.enqueue(piece, estimated_size);
818    }
819
820    fn load(&self, hash: u64) -> BoxFuture<'static, Result<Load<K, V, P>>> {
821        // TODO(MrCroxx): refactor this.
822        self.load(hash).boxed()
823    }
824
825    fn delete(&self, hash: u64) {
826        self.delete(hash);
827    }
828
829    fn may_contains(&self, hash: u64) -> bool {
830        self.may_contains(hash)
831    }
832
833    fn destroy(&self) -> BoxFuture<'static, Result<()>> {
834        self.destroy()
835    }
836
837    fn wait(&self) -> BoxFuture<'static, ()> {
838        // TODO(MrCroxx): refactor this.
839        self.wait().boxed()
840    }
841
842    fn close(&self) -> BoxFuture<'static, Result<()>> {
843        self.close()
844    }
845}
846
847#[cfg(test)]
848mod tests {
849
850    use std::{fs::File, path::Path};
851
852    use bytesize::ByteSize;
853    use foyer_common::hasher::ModHasher;
854    use foyer_memory::{Cache, CacheBuilder, CacheEntry, FifoConfig, TestProperties};
855    use itertools::Itertools;
856
857    use super::*;
858    use crate::{
859        PsyncIoEngineConfig, RejectAll,
860        engine::RecoverMode,
861        io::{
862            device::{DeviceBuilder, combined::CombinedDeviceBuilder, fs::FsDeviceBuilder},
863            engine::{IoEngine, IoEngineBuildContext, IoEngineConfig},
864        },
865        serde::EntrySerializer,
866        test_utils::Biased,
867    };
868
869    const KB: usize = 1024;
870
871    fn cache_for_test() -> Cache<u64, Vec<u8>, ModHasher, TestProperties> {
872        CacheBuilder::new(10)
873            .with_shards(1)
874            .with_eviction_config(FifoConfig::default())
875            .with_hash_builder(ModHasher::default())
876            .build()
877    }
878
879    async fn io_engine_for_test(spawner: Spawner) -> Arc<dyn IoEngine> {
880        // TODO(MrCroxx): Test with other io engines.
881        PsyncIoEngineConfig::new()
882            .boxed()
883            .build(IoEngineBuildContext { spawner })
884            .await
885            .unwrap()
886    }
887
888    /// 4 files, fifo eviction, 16 KiB block, 64 KiB capacity.
889    async fn engine_for_test(dir: impl AsRef<Path>) -> Arc<BlockEngine<u64, Vec<u8>, TestProperties>> {
890        store_for_test_with_reinsertion_filter(dir, StorageFilter::new().with_condition(RejectAll)).await
891    }
892
893    async fn store_for_test_with_reinsertion_filter(
894        dir: impl AsRef<Path>,
895        reinsertion_filter: StorageFilter,
896    ) -> Arc<BlockEngine<u64, Vec<u8>, TestProperties>> {
897        let device = FsDeviceBuilder::new(dir)
898            .with_capacity(ByteSize::kib(64).as_u64() as _)
899            .build()
900            .unwrap();
901        let spawner = Spawner::current();
902        let io_engine = io_engine_for_test(spawner.clone()).await;
903        let metrics = Arc::new(Metrics::noop());
904        let builder = BlockEngineConfig {
905            device,
906            block_size: 16 * 1024,
907            compression: Compression::None,
908            indexer_shards: 4,
909            recover_concurrency: 2,
910            flushers: 1,
911            reclaimers: 1,
912            clean_block_threshold: 1,
913            admission_filter: StorageFilter::new(),
914            eviction_pickers: vec![Box::<FifoPicker>::default()],
915            reinsertion_filter,
916            enable_tombstone_log: false,
917            buffer_pool_size: 16 * 1024 * 1024,
918            blob_index_size: 4 * 1024,
919            submit_queue_size_threshold: 16 * 1024 * 1024 * 2,
920            flush_switch: Switch::default(),
921            load_holder: Holder::default(),
922            marker: PhantomData,
923        };
924
925        let builder = Box::new(builder);
926        builder
927            .build(EngineBuildContext {
928                io_engine,
929                metrics,
930                spawner,
931                recover_mode: RecoverMode::Strict,
932            })
933            .await
934            .unwrap()
935    }
936
937    async fn store_for_test_with_tombstone_log(
938        dir: impl AsRef<Path>,
939    ) -> Arc<BlockEngine<u64, Vec<u8>, TestProperties>> {
940        let device = FsDeviceBuilder::new(dir)
941            .with_capacity(ByteSize::kib(64).as_u64() as usize + ByteSize::kib(4).as_u64() as usize)
942            .build()
943            .unwrap();
944        let spawner = Spawner::current();
945        let io_engine = io_engine_for_test(spawner.clone()).await;
946        let metrics = Arc::new(Metrics::noop());
947        let builder = BlockEngineConfig {
948            device,
949            block_size: 16 * 1024,
950            compression: Compression::None,
951            indexer_shards: 4,
952            recover_concurrency: 2,
953            flushers: 1,
954            reclaimers: 1,
955            clean_block_threshold: 1,
956            eviction_pickers: vec![Box::<FifoPicker>::default()],
957            admission_filter: StorageFilter::new(),
958            reinsertion_filter: StorageFilter::new().with_condition(RejectAll),
959            enable_tombstone_log: true,
960            buffer_pool_size: 16 * 1024 * 1024,
961            blob_index_size: 4 * 1024,
962            submit_queue_size_threshold: 16 * 1024 * 1024 * 2,
963            flush_switch: Switch::default(),
964            load_holder: Holder::default(),
965            marker: PhantomData,
966        };
967        let builder = Box::new(builder);
968        builder
969            .build(EngineBuildContext {
970                io_engine,
971                metrics,
972                spawner,
973                recover_mode: RecoverMode::Strict,
974            })
975            .await
976            .unwrap()
977    }
978
979    fn enqueue(
980        store: &BlockEngine<u64, Vec<u8>, TestProperties>,
981        entry: CacheEntry<u64, Vec<u8>, ModHasher, TestProperties>,
982    ) {
983        let estimated_size = EntrySerializer::estimated_size(entry.key(), entry.value());
984        store.enqueue(entry.piece().into(), estimated_size);
985    }
986
987    #[test_log::test(tokio::test)]
988    async fn test_store_enqueue_lookup_recovery() {
989        let dir = tempfile::tempdir().unwrap();
990
991        let memory = cache_for_test();
992        let store = engine_for_test(dir.path()).await;
993
994        // [ [e1, e2], [], [], [] ]
995        store.hold_flush();
996        let e1 = memory.insert(1, vec![1; 7 * KB]);
997        let e2 = memory.insert(2, vec![2; 3 * KB]);
998        enqueue(&store, e1.clone());
999        enqueue(&store, e2);
1000        store.unhold_flush();
1001        store.wait().await;
1002
1003        let r1 = store.load(memory.hash(&1)).await.unwrap().kv().unwrap();
1004        assert_eq!(r1, (1, vec![1; 7 * KB]));
1005        let r2 = store.load(memory.hash(&2)).await.unwrap().kv().unwrap();
1006        assert_eq!(r2, (2, vec![2; 3 * KB]));
1007
1008        // [ [e1, e2], [e3, e4], [], [] ]
1009        store.hold_flush();
1010        let e3 = memory.insert(3, vec![3; 7 * KB]);
1011        let e4 = memory.insert(4, vec![4; 2 * KB]);
1012        enqueue(&store, e3);
1013        enqueue(&store, e4);
1014        store.unhold_flush();
1015        store.wait().await;
1016
1017        let r1 = store.load(memory.hash(&1)).await.unwrap().kv().unwrap();
1018        assert_eq!(r1, (1, vec![1; 7 * KB]));
1019        let r2 = store.load(memory.hash(&2)).await.unwrap().kv().unwrap();
1020        assert_eq!(r2, (2, vec![2; 3 * KB]));
1021        let r3 = store.load(memory.hash(&3)).await.unwrap().kv().unwrap();
1022        assert_eq!(r3, (3, vec![3; 7 * KB]));
1023        let r4 = store.load(memory.hash(&4)).await.unwrap().kv().unwrap();
1024        assert_eq!(r4, (4, vec![4; 2 * KB]));
1025
1026        // [ [e1, e2], [e3, e4], [e5], [] ]
1027        let e5 = memory.insert(5, vec![5; 11 * KB]);
1028        enqueue(&store, e5);
1029        store.wait().await;
1030
1031        let r1 = store.load(memory.hash(&1)).await.unwrap().kv().unwrap();
1032        assert_eq!(r1, (1, vec![1; 7 * KB]));
1033        let r2 = store.load(memory.hash(&2)).await.unwrap().kv().unwrap();
1034        assert_eq!(r2, (2, vec![2; 3 * KB]));
1035        let r3 = store.load(memory.hash(&3)).await.unwrap().kv().unwrap();
1036        assert_eq!(r3, (3, vec![3; 7 * KB]));
1037        let r4 = store.load(memory.hash(&4)).await.unwrap().kv().unwrap();
1038        assert_eq!(r4, (4, vec![4; 2 * KB]));
1039        let r5 = store.load(memory.hash(&5)).await.unwrap().kv().unwrap();
1040        assert_eq!(r5, (5, vec![5; 11 * KB]));
1041
1042        // [ [], [e3, e4], [e5], [e6, e4*] ]
1043        store.hold_flush();
1044        let e6 = memory.insert(6, vec![6; 7 * KB]);
1045        let e4v2 = memory.insert(4, vec![!4; 3 * KB]);
1046        enqueue(&store, e6);
1047        enqueue(&store, e4v2);
1048        store.unhold_flush();
1049        store.wait().await;
1050
1051        assert!(store.load(memory.hash(&1)).await.unwrap().kv().is_none());
1052        assert!(store.load(memory.hash(&2)).await.unwrap().kv().is_none());
1053        let r3 = store.load(memory.hash(&3)).await.unwrap().kv().unwrap();
1054        assert_eq!(r3, (3, vec![3; 7 * KB]));
1055        let r4v2 = store.load(memory.hash(&4)).await.unwrap().kv().unwrap();
1056        assert_eq!(r4v2, (4, vec![!4; 3 * KB]));
1057        let r5 = store.load(memory.hash(&5)).await.unwrap().kv().unwrap();
1058        assert_eq!(r5, (5, vec![5; 11 * KB]));
1059        let r6 = store.load(memory.hash(&6)).await.unwrap().kv().unwrap();
1060        assert_eq!(r6, (6, vec![6; 7 * KB]));
1061
1062        store.close().await.unwrap();
1063        enqueue(&store, e1);
1064        store.wait().await;
1065
1066        drop(store);
1067
1068        let store = engine_for_test(dir.path()).await;
1069
1070        assert!(store.load(memory.hash(&1)).await.unwrap().kv().is_none());
1071        assert!(store.load(memory.hash(&2)).await.unwrap().kv().is_none());
1072        let r3 = store.load(memory.hash(&3)).await.unwrap().kv().unwrap();
1073        assert_eq!(r3, (3, vec![3; 7 * KB]));
1074        let r4v2 = store.load(memory.hash(&4)).await.unwrap().kv().unwrap();
1075        assert_eq!(r4v2, (4, vec![!4; 3 * KB]));
1076        let r5 = store.load(memory.hash(&5)).await.unwrap().kv().unwrap();
1077        assert_eq!(r5, (5, vec![5; 11 * KB]));
1078        let r6 = store.load(memory.hash(&6)).await.unwrap().kv().unwrap();
1079        assert_eq!(r6, (6, vec![6; 7 * KB]));
1080    }
1081
1082    #[test_log::test(tokio::test)]
1083    async fn test_store_delete_recovery() {
1084        let dir = tempfile::tempdir().unwrap();
1085
1086        let memory = cache_for_test();
1087        let store = store_for_test_with_tombstone_log(dir.path()).await;
1088
1089        let es = (0..10).map(|i| memory.insert(i, vec![i as u8; 3 * KB])).collect_vec();
1090
1091        // [[0, 1, 2], [3, 4, 5], [6, 7, 8], []]
1092        for e in es.iter().take(9) {
1093            enqueue(&store, e.clone());
1094        }
1095        store.wait().await;
1096
1097        for i in 0..9 {
1098            assert_eq!(
1099                store.load(memory.hash(&i)).await.unwrap().kv(),
1100                Some((i, vec![i as u8; 3 * KB]))
1101            );
1102        }
1103
1104        store.delete(memory.hash(&3));
1105        store.wait().await;
1106        assert_eq!(store.load(memory.hash(&3)).await.unwrap().kv(), None);
1107
1108        store.close().await.unwrap();
1109        drop(store);
1110
1111        let store = store_for_test_with_tombstone_log(dir.path()).await;
1112        for i in 0..9 {
1113            if i != 3 {
1114                assert_eq!(
1115                    store.load(memory.hash(&i)).await.unwrap().kv(),
1116                    Some((i, vec![i as u8; 3 * KB]))
1117                );
1118            } else {
1119                assert_eq!(store.load(memory.hash(&3)).await.unwrap().kv(), None);
1120            }
1121        }
1122
1123        enqueue(&store, es[3].clone());
1124        store.wait().await;
1125        assert_eq!(
1126            store.load(memory.hash(&3)).await.unwrap().kv(),
1127            Some((3, vec![3; 3 * KB]))
1128        );
1129
1130        store.close().await.unwrap();
1131        drop(store);
1132
1133        let store = store_for_test_with_tombstone_log(dir.path()).await;
1134
1135        assert_eq!(
1136            store.load(memory.hash(&3)).await.unwrap().kv(),
1137            Some((3, vec![3; 3 * KB]))
1138        );
1139    }
1140
1141    #[test_log::test(tokio::test)]
1142    async fn test_store_destroy_recovery() {
1143        let dir = tempfile::tempdir().unwrap();
1144
1145        let memory = cache_for_test();
1146        let store = store_for_test_with_tombstone_log(dir.path()).await;
1147
1148        let es = (0..10).map(|i| memory.insert(i, vec![i as u8; 3 * KB])).collect_vec();
1149
1150        // [[0, 1, 2], [3, 4, 5], [6, 7, 8], []]
1151        store.hold_flush();
1152        for e in es.iter().take(9) {
1153            enqueue(&store, e.clone());
1154        }
1155        store.unhold_flush();
1156        store.wait().await;
1157
1158        for i in 0..9 {
1159            assert_eq!(
1160                store.load(memory.hash(&i)).await.unwrap().kv(),
1161                Some((i, vec![i as u8; 3 * KB]))
1162            );
1163        }
1164
1165        store.delete(memory.hash(&3));
1166        store.wait().await;
1167        assert_eq!(store.load(memory.hash(&3)).await.unwrap().kv(), None);
1168
1169        store.destroy().await.unwrap();
1170
1171        store.close().await.unwrap();
1172        drop(store);
1173
1174        let store = store_for_test_with_tombstone_log(dir.path()).await;
1175        for i in 0..9 {
1176            assert_eq!(store.load(memory.hash(&i)).await.unwrap().kv(), None);
1177        }
1178
1179        enqueue(&store, es[3].clone());
1180        store.wait().await;
1181        assert_eq!(
1182            store.load(memory.hash(&3)).await.unwrap().kv(),
1183            Some((3, vec![3; 3 * KB]))
1184        );
1185
1186        store.close().await.unwrap();
1187        drop(store);
1188
1189        let store = store_for_test_with_tombstone_log(dir.path()).await;
1190
1191        assert_eq!(
1192            store.load(memory.hash(&3)).await.unwrap().kv(),
1193            Some((3, vec![3; 3 * KB]))
1194        );
1195    }
1196
1197    // FIXME(MrCroxx): Move the admission test to store level.
1198    // #[test_log::test(tokio::test)]
1199    // async fn test_store_admission() {
1200    //     let dir = tempfile::tempdir().unwrap();
1201
1202    //     let memory = cache_for_test();
1203    //     let store = store_for_test_with_admission_picker(&memory, dir.path(),
1204    // Arc::new(BiasedPicker::new([1]))).await;
1205
1206    //     let e1 = memory.insert(1, vec![1; 7 * KB]);
1207    //     let e2 = memory.insert(2, vec![2; 7 * KB]);
1208
1209    //     assert!(enqueue(&store, e1.clone(),).await.unwrap());
1210    //     assert!(!enqueue(&store, e2,).await.unwrap());
1211
1212    //     let r1 = store.load(&1).await.unwrap().unwrap();
1213    //     assert_eq!(r1, (1, vec![1; 7 * KB]));
1214    //     assert!(store.load(&2).await.unwrap().is_none());
1215    // }
1216
1217    #[test_log::test(tokio::test)]
1218    async fn test_store_reinsertion() {
1219        let dir = tempfile::tempdir().unwrap();
1220
1221        let memory = cache_for_test();
1222        let store = store_for_test_with_reinsertion_filter(
1223            dir.path(),
1224            StorageFilter::new().with_condition(Biased::new(vec![1, 3, 5, 7, 9, 11, 13, 15, 17, 19])),
1225        )
1226        .await;
1227
1228        let es = (0..15).map(|i| memory.insert(i, vec![i as u8; 3 * KB])).collect_vec();
1229
1230        // [[(0), (1), (2)], [(3), (4), (5)], [(6), (7), (8)], []]
1231        for e in es.iter().take(9).cloned() {
1232            enqueue(&store, e);
1233            store.wait().await;
1234        }
1235
1236        for i in 0..9 {
1237            let r = store.load(memory.hash(&i)).await.unwrap().kv().unwrap();
1238            assert_eq!(r, (i, vec![i as u8; 3 * KB]));
1239        }
1240
1241        // [[], [(3), (4), (5)], [(6), (7), (8)], [(9), (10), (1)]]
1242        enqueue(&store, es[9].clone());
1243        enqueue(&store, es[10].clone());
1244        store.wait().await;
1245        let mut res = vec![];
1246        for i in 0..11 {
1247            res.push(store.load(memory.hash(&i)).await.unwrap().kv());
1248        }
1249        assert_eq!(
1250            res,
1251            vec![
1252                None,
1253                Some((1, vec![1; 3 * KB])),
1254                None,
1255                Some((3, vec![3; 3 * KB])),
1256                Some((4, vec![4; 3 * KB])),
1257                Some((5, vec![5; 3 * KB])),
1258                Some((6, vec![6; 3 * KB])),
1259                Some((7, vec![7; 3 * KB])),
1260                Some((8, vec![8; 3 * KB])),
1261                Some((9, vec![9; 3 * KB])),
1262                Some((10, vec![10; 3 * KB])),
1263            ]
1264        );
1265
1266        // [[(11), (3), (5)], [], [(6), (7), (8)], [(9), (10), (1)]]
1267        enqueue(&store, es[11].clone());
1268        store.wait().await;
1269        let mut res = vec![];
1270        for i in 0..12 {
1271            res.push(store.load(memory.hash(&i)).await.unwrap().kv());
1272        }
1273        assert_eq!(
1274            res,
1275            vec![
1276                None,
1277                Some((1, vec![1; 3 * KB])),
1278                None,
1279                Some((3, vec![3; 3 * KB])),
1280                None,
1281                Some((5, vec![5; 3 * KB])),
1282                Some((6, vec![6; 3 * KB])),
1283                Some((7, vec![7; 3 * KB])),
1284                Some((8, vec![8; 3 * KB])),
1285                Some((9, vec![9; 3 * KB])),
1286                Some((10, vec![10; 3 * KB])),
1287                Some((11, vec![11; 3 * KB])),
1288            ]
1289        );
1290
1291        // [[(11), (3), (5)], [(12), (13), (14)], [], [(9), (10), (1)]]
1292        store.delete(memory.hash(&7));
1293        store.wait().await;
1294        enqueue(&store, es[12].clone());
1295        store.wait().await;
1296        enqueue(&store, es[13].clone());
1297        store.wait().await;
1298        enqueue(&store, es[14].clone());
1299        store.wait().await;
1300        let mut res = vec![];
1301        for i in 0..15 {
1302            res.push(store.load(memory.hash(&i)).await.unwrap().kv());
1303        }
1304        assert_eq!(
1305            res,
1306            vec![
1307                None,
1308                Some((1, vec![1; 3 * KB])),
1309                None,
1310                Some((3, vec![3; 3 * KB])),
1311                None,
1312                Some((5, vec![5; 3 * KB])),
1313                None,
1314                None,
1315                None,
1316                Some((9, vec![9; 3 * KB])),
1317                Some((10, vec![10; 3 * KB])),
1318                Some((11, vec![11; 3 * KB])),
1319                Some((12, vec![12; 3 * KB])),
1320                Some((13, vec![13; 3 * KB])),
1321                Some((14, vec![14; 3 * KB])),
1322            ]
1323        );
1324    }
1325
1326    #[test_log::test(tokio::test)]
1327    async fn test_store_magic_checksum_mismatch() {
1328        let dir = tempfile::tempdir().unwrap();
1329
1330        let memory = cache_for_test();
1331        let store = engine_for_test(dir.path()).await;
1332
1333        // write entry 1
1334        let e1 = memory.insert(1, vec![1; 7 * KB]);
1335        enqueue(&store, e1);
1336        store.wait().await;
1337
1338        // check entry 1
1339        let r1 = store.load(memory.hash(&1)).await.unwrap().kv().unwrap();
1340        assert_eq!(r1, (1, vec![1; 7 * KB]));
1341
1342        // corrupt entry and header
1343        for entry in std::fs::read_dir(dir.path()).unwrap() {
1344            let entry = entry.unwrap();
1345            if !entry.metadata().unwrap().is_file() {
1346                continue;
1347            }
1348
1349            let file = File::options().write(true).open(entry.path()).unwrap();
1350            #[cfg(target_family = "unix")]
1351            {
1352                use std::os::unix::fs::FileExt;
1353                file.write_all_at(&[b'x'; 42], 5 * 1024).unwrap();
1354            }
1355            #[cfg(target_family = "windows")]
1356            {
1357                use std::os::windows::fs::FileExt;
1358                file.seek_write(&[b'x'; 42], 5 * 1024).unwrap();
1359            }
1360        }
1361
1362        assert!(store.load(memory.hash(&1)).await.unwrap().kv().is_none());
1363    }
1364
1365    #[test_log::test(tokio::test)]
1366    async fn test_aggregated_device() {
1367        let dir = tempfile::tempdir().unwrap();
1368
1369        const KB: usize = 1024;
1370        const MB: usize = 1024 * 1024;
1371
1372        let spawner = Spawner::current();
1373        let io_engine = io_engine_for_test(spawner.clone()).await;
1374
1375        let d1 = FsDeviceBuilder::new(dir.path().join("dev1"))
1376            .with_capacity(MB)
1377            .build()
1378            .unwrap();
1379        let d2 = FsDeviceBuilder::new(dir.path().join("dev2"))
1380            .with_capacity(2 * MB)
1381            .build()
1382            .unwrap();
1383        let d3 = FsDeviceBuilder::new(dir.path().join("dev3"))
1384            .with_capacity(4 * MB)
1385            .build()
1386            .unwrap();
1387        let device = CombinedDeviceBuilder::new()
1388            .with_device(d1)
1389            .with_device(d2)
1390            .with_device(d3)
1391            .build()
1392            .unwrap();
1393        let engine = BlockEngineConfig::<u64, Vec<u8>, TestProperties>::new(device)
1394            .with_block_size(64 * KB)
1395            .boxed()
1396            .build(EngineBuildContext {
1397                io_engine,
1398                metrics: Arc::new(Metrics::noop()),
1399                spawner,
1400                recover_mode: RecoverMode::None,
1401            })
1402            .await
1403            .unwrap();
1404        assert_eq!(engine.inner.block_manager.blocks(), (1 + 2 + 4) * MB / (64 * KB));
1405    }
1406}