1use 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
70pub 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 pub fn new(device: Arc<dyn Device>) -> Self {
141 Self {
142 device,
143 block_size: 16 * 1024 * 1024, compression: Compression::default(),
145 indexer_shards: 64,
146 recover_concurrency: 8,
147 flushers: 1,
148 reclaimers: 1,
149 buffer_pool_size: 16 * 1024 * 1024, blob_index_size: 4 * 1024, submit_queue_size_threshold: 16 * 1024 * 1024, 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 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 pub fn with_indexer_shards(mut self, indexer_shards: usize) -> Self {
182 self.indexer_shards = indexer_shards;
183 self
184 }
185
186 pub fn with_recover_concurrency(mut self, recover_concurrency: usize) -> Self {
190 self.recover_concurrency = recover_concurrency;
191 self
192 }
193
194 pub fn with_flushers(mut self, flushers: usize) -> Self {
200 self.flushers = flushers;
201 self
202 }
203
204 pub fn with_admission_filter(mut self, filter: StorageFilter) -> Self {
210 self.admission_filter = filter;
211 self
212 }
213
214 pub fn with_reclaimers(mut self, reclaimers: usize) -> Self {
220 self.reclaimers = reclaimers;
221 self
222 }
223
224 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 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 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 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 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 pub fn with_reinsertion_filter(mut self, filter: StorageFilter) -> Self {
296 self.reinsertion_filter = filter;
297 self
298 }
299
300 pub fn with_tombstone_log(mut self, enable: bool) -> Self {
305 self.enable_tombstone_log = enable;
306 self
307 }
308
309 #[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 #[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 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 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
479pub 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 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 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 this.inner.indexer.clear();
772
773 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 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 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 PsyncIoEngineConfig::new()
882 .boxed()
883 .build(IoEngineBuildContext { spawner })
884 .await
885 .unwrap()
886 }
887
888 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 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 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 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 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 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 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 #[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 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 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 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 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 let e1 = memory.insert(1, vec![1; 7 * KB]);
1335 enqueue(&store, e1);
1336 store.wait().await;
1337
1338 let r1 = store.load(memory.hash(&1)).await.unwrap().kv().unwrap();
1340 assert_eq!(r1, (1, vec![1; 7 * KB]));
1341
1342 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}