1use crate::drop_event::DropEvent;
59use crate::errors::FxfsError;
60use crate::filesystem::MAX_BLOCK_SIZE;
61use crate::log::*;
62use crate::lsm_tree::bloom_filter::{BloomFilterReader, BloomFilterStats, BloomFilterWriter};
63use crate::lsm_tree::types::{
64 BoxedLayerIterator, Existence, FuzzyHash, Item, ItemRef, Key, Layer, LayerIterator, LayerValue,
65 LayerWriter, MaybeContainsKey,
66};
67use crate::object_handle::{LayerObject, ObjectHandle, ReadObjectHandle, WriteBytes};
68use crate::object_store::caching_object_handle::{CHUNK_SIZE, CachedChunk, CachingObjectHandle};
69use crate::object_store::extent::MIN_BLOCK_SIZE;
70use crate::object_store::{DataObjectHandle, HandleOptions, ObjectStore};
71use crate::serialized_types::serialized_key::{KeyDeserializer, compare_keys};
72use crate::serialized_types::{
73 LATEST_VERSION, OLD_KEY_SERIALIZATION_VERSION, Version, Versioned, VersionedLatest,
74};
75use anyhow::{Context, Error, anyhow, ensure};
76use async_trait::async_trait;
77use byteorder::{ByteOrder, LittleEndian, ReadBytesExt, WriteBytesExt};
78use fprint::TypeFingerprint;
79use fuchsia_sync::Mutex;
80use futures::future::BoxFuture;
81use futures::stream::{FuturesUnordered, TryStreamExt};
82use fxfs_crypto::{Crypt, UnwrappedKey, WrappedKey};
83use serde::{Deserialize, Serialize};
84use static_assertions::const_assert;
85use std::cmp::Ordering;
86use std::io::{Read as _, Write as _};
87use std::marker::PhantomData;
88use std::ops::Bound;
89use std::sync::Arc;
90use storage_units::BlockSize;
91
92const PERSISTENT_LAYER_MAGIC: &[u8; 8] = b"FxfsLayr";
93
94pub type LayerHeader = LayerHeaderV39;
96
97#[derive(Debug, Serialize, Deserialize, TypeFingerprint, Versioned)]
98pub struct LayerHeaderV39 {
99 magic: [u8; 8],
101 block_size: u64,
107}
108
109pub type LayerInfo = LayerInfoV39;
111
112#[derive(Debug, Serialize, Deserialize, TypeFingerprint, Versioned)]
113pub struct LayerInfoV39 {
114 num_items: usize,
117 num_data_blocks: u64,
119 bloom_filter_size_bytes: usize,
121 bloom_filter_seed: u64,
123 bloom_filter_num_hashes: usize,
125}
126
127struct LayerData<K> {
128 object_id: u64,
129 version: Version,
130 block_size: BlockSize,
131 data_size: u64,
132 seek_table: Vec<u64>,
133 num_items: usize,
134 bloom_filter: Option<BloomFilterReader<K>>,
135 bloom_filter_stats: Option<BloomFilterStats>,
136 close_event: Mutex<Option<Arc<DropEvent>>>,
137}
138
139impl<K> LayerData<K> {
140 fn data_offset(&self) -> u64 {
141 NUM_HEADER_BLOCKS * self.block_size
142 }
143}
144
145pub struct PersistentLayer<K, V> {
147 object_handle: CachingObjectHandle<Arc<dyn LayerObject>>,
148 data: LayerData<K>,
149 _value_type: PhantomData<V>,
150}
151
152pub struct SyncPersistentLayer<K, V> {
154 object_handle: Arc<dyn LayerObject>,
155 data: LayerData<K>,
156 _value_type: PhantomData<V>,
157}
158
159struct BufferCursor<B> {
160 buffer: B,
161 pos: usize,
162}
163
164impl<B: LayerBuffer> BufferCursor<B> {
165 fn as_bytes(&self) -> &[u8] {
166 self.buffer.as_bytes_from(self.pos)
167 }
168}
169
170impl<B: LayerBuffer> std::io::Read for BufferCursor<B> {
171 fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
172 let to_read = self.buffer.read_at(self.pos, buf);
173 self.pos += to_read;
174 Ok(to_read)
175 }
176}
177
178trait LayerBuffer {
179 fn read_at(&self, pos: usize, buf: &mut [u8]) -> usize;
180
181 fn as_bytes_from(&self, pos: usize) -> &[u8];
182
183 fn has_io_error(&self) -> bool {
184 false
185 }
186}
187
188struct ChunkBuffer<'iter> {
189 handle: &'iter CachingObjectHandle<Arc<dyn LayerObject>>,
190 chunk: Option<(usize, CachedChunk)>,
192}
193
194impl LayerBuffer for ChunkBuffer<'_> {
195 fn read_at(&self, pos: usize, buf: &mut [u8]) -> usize {
196 let Some((_, chunk)) = &self.chunk else {
197 return 0;
198 };
199 let to_read = std::cmp::min(buf.len(), chunk.len().saturating_sub(pos));
200 if to_read > 0 {
201 buf[..to_read].copy_from_slice(&chunk[pos..pos + to_read]);
202 }
203 to_read
204 }
205
206 fn as_bytes_from(&self, pos: usize) -> &[u8] {
207 self.chunk.as_ref().and_then(|(_, c)| c.get(pos..)).unwrap_or(&[])
208 }
209}
210
211#[derive(Clone, Copy)]
212struct SliceBuffer<'iter> {
213 handle: &'iter dyn LayerObject,
214 slice: &'iter [u8],
215}
216
217impl LayerBuffer for SliceBuffer<'_> {
218 fn read_at(&self, pos: usize, buf: &mut [u8]) -> usize {
219 let to_read = std::cmp::min(buf.len(), self.slice.len().saturating_sub(pos));
220 if to_read > 0 {
221 buf[..to_read].copy_from_slice(&self.slice[pos..pos + to_read]);
222 }
223 to_read
224 }
225
226 fn as_bytes_from(&self, pos: usize) -> &[u8] {
227 self.slice.get(pos..).unwrap_or(&[])
228 }
229
230 fn has_io_error(&self) -> bool {
231 self.handle.has_io_error()
232 }
233}
234
235const MINIMUM_DATA_BLOCKS_FOR_BLOOM_FILTER: usize = 4;
237
238const NUM_HEADER_BLOCKS: u64 = 1;
240
241const MINIMUM_LAYER_FILE_BLOCKS: u64 = 2;
244
245const MAX_BLOOM_FILTER_SIZE: usize = 64 * 1024 * 1024;
248const MAX_SEEK_TABLE_SIZE: usize = 64 * 1024 * 1024;
249
250const PER_DATA_BLOCK_HEADER_SIZE: usize = 2;
252const PER_DATA_BLOCK_SEEK_ENTRY_SIZE: usize = 2;
253
254enum KeyState<K> {
255 None,
256 Deserialized(K),
257 InPlace,
258}
259
260struct KeyOnlyIterator<'iter, K: Key, V: LayerValue, B> {
262 buffer: BufferCursor<B>,
264
265 layer: &'iter LayerData<K>,
266
267 pos: u64,
269
270 item_index: u16,
272
273 item_count: u16,
275
276 key: KeyState<K>,
278
279 current_block_base_u64: u64,
281
282 _value_type: PhantomData<V>,
283}
284
285impl<K: Key, V: LayerValue, B: LayerBuffer> KeyOnlyIterator<'_, K, V, B> {
286 fn corruption_error(&self, err: impl std::fmt::Display + Send + Sync + 'static) -> Error {
287 if self.buffer.buffer.has_io_error() {
288 anyhow!(zx_status::Status::IO).context(err)
289 } else {
290 anyhow!(FxfsError::Inconsistent).context(err)
291 }
292 }
293
294 fn seek_to_block_item(&mut self, index: u16) -> Result<(), Error> {
298 ensure!(index < self.item_count, FxfsError::OutOfRange);
299 if index == self.item_index && matches!(self.key, KeyState::None) {
300 return Ok(());
303 }
304 let block_start =
305 self.layer.block_size.align_down((self.buffer.pos as u64).saturating_sub(1)) as usize;
306 let offset_in_block = if index == 0 {
307 PER_DATA_BLOCK_HEADER_SIZE
310 } else {
311 let old_buffer_pos = self.buffer.pos;
312 let seek_entry_pos = (block_start + self.layer.block_size.get() as usize)
313 .checked_sub(PER_DATA_BLOCK_SEEK_ENTRY_SIZE * usize::from(self.item_count - index))
314 .ok_or_else(|| {
315 self.corruption_error(format!(
316 "Invalid item count {} for index {index}",
317 self.item_count
318 ))
319 })?;
320 if seek_entry_pos < block_start + PER_DATA_BLOCK_HEADER_SIZE {
321 return Err(self.corruption_error(format!(
322 "Invalid item count {} for index {index}",
323 self.item_count
324 )));
325 }
326 self.buffer.pos = seek_entry_pos;
327 let res = self.buffer.read_u16::<LittleEndian>();
328 self.buffer.pos = old_buffer_pos;
329 let offset_in_block = res
330 .map_err(|e| self.corruption_error(e))
331 .context("Failed to read offset")? as usize;
332 if offset_in_block >= self.layer.block_size.get() as usize
333 || offset_in_block <= PER_DATA_BLOCK_HEADER_SIZE
334 {
335 return Err(self
336 .corruption_error(format!("Offset {offset_in_block} is out of valid range.")));
337 }
338 offset_in_block
339 };
340 self.item_index = index;
341 self.key = KeyState::None;
342 self.buffer.pos = block_start + offset_in_block;
343 Ok(())
344 }
345
346 fn take_item(&mut self) -> Result<Option<Item<K, V>>, Error> {
347 let key = match std::mem::replace(&mut self.key, KeyState::None) {
348 KeyState::None => return Ok(None),
349 KeyState::InPlace => {
350 let (mut deserializer, key_len) =
351 KeyDeserializer::new(self.buffer.as_bytes(), Some(self.current_block_base_u64))
352 .map_err(|e| self.corruption_error(e))
353 .context("Corrupt layer (key format)")?;
354 let key = K::deserialize_key_from(&mut deserializer)
355 .map_err(|e| self.corruption_error(e))
356 .context("Corrupt layer (key)")?;
357 if !deserializer.is_empty() {
358 return Err(self.corruption_error("Trailing bytes in serialized key"));
359 }
360 self.buffer.pos += key_len;
361 key
362 }
363 KeyState::Deserialized(key) => key,
364 };
365 let value = V::deserialize_from_version(self.buffer.by_ref(), self.layer.version)
366 .map_err(|e| self.corruption_error(e))
367 .context("Corrupt layer (value)")?;
368 Ok(Some(Item { key, value }))
369 }
370
371 fn read_block_header(&mut self) -> Result<(), Error> {
372 self.item_count =
373 self.buffer.read_u16::<LittleEndian>().map_err(|e| self.corruption_error(e))?;
374 if self.item_count == 0 {
375 return Err(self.corruption_error(format!(
376 "Read block with zero item count (object: {}, offset: {})",
377 self.layer.object_id, self.pos
378 )));
379 }
380 if PER_DATA_BLOCK_HEADER_SIZE
381 + (usize::from(self.item_count) - 1) * PER_DATA_BLOCK_SEEK_ENTRY_SIZE
382 >= self.layer.block_size.get() as usize
383 {
384 return Err(self.corruption_error("Block seek table overlaps header"));
385 }
386 debug!(
387 pos = self.pos,
388 object_size = self.layer.data_offset() + self.layer.data_size,
389 oid = self.layer.object_id;
390 ""
391 );
392 if self.layer.version > OLD_KEY_SERIALIZATION_VERSION {
393 let block_index = (self.pos - self.layer.data_offset()) / self.layer.block_size;
394 self.current_block_base_u64 = *self
395 .layer
396 .seek_table
397 .get(block_index as usize)
398 .ok_or_else(|| self.corruption_error("Block index out of bounds"))?;
399 }
400 self.pos += self.layer.block_size;
401 self.item_index = 0;
402 self.key = KeyState::None;
403 Ok(())
404 }
405
406 fn deserialize_current_key(&mut self) -> Result<(), Error> {
407 self.seek_to_block_item(self.item_index)?;
408 self.key = if self.layer.version <= OLD_KEY_SERIALIZATION_VERSION {
409 KeyState::Deserialized(
410 K::deserialize_from_version(self.buffer.by_ref(), self.layer.version)
411 .map_err(|e| self.corruption_error(e))
412 .context("Corrupt layer (key)")?,
413 )
414 } else {
415 KeyState::InPlace
416 };
417 self.item_index += 1;
418 Ok(())
419 }
420
421 fn binary_search_block(
428 &mut self,
429 search_key: &mut SearchKey<'_, K>,
430 ) -> Result<Option<bool>, Error> {
431 if self.layer.version > OLD_KEY_SERIALIZATION_VERSION {
432 let target = search_key
433 .serialized_with_base(self.current_block_base_u64)
434 .ok_or_else(|| self.corruption_error("Search key precedes block"))?;
435 return self.binary_search_block_in_place(target);
436 }
437
438 let mut left_index = 0;
439 let mut right_index = self.item_count;
440 while left_index < (right_index - 1) {
441 let mid_index = left_index + ((right_index - left_index) / 2);
442 self.seek_to_block_item(mid_index).context("Read index offset for binary search")?;
443 self.deserialize_current_key()?;
444 match search_key.compare(self)?.context("Unexpected EOF")? {
445 Ordering::Greater => right_index = mid_index,
446 Ordering::Equal => return Ok(Some(true)),
447 Ordering::Less => left_index = mid_index,
448 }
449 }
450 if right_index < self.item_count {
451 self.seek_to_block_item(right_index)
452 .context("Read index for offset of right pointer")?;
453 self.deserialize_current_key()?;
454 return Ok(Some(false));
455 }
456 self.key = KeyState::None;
457 Ok(None)
458 }
459
460 fn binary_search_block_in_place(&mut self, target: &[u8]) -> Result<Option<bool>, Error> {
463 let mut left_index = 0;
464 let mut right_index = self.item_count;
465 if right_index > 1 {
466 let block_size = self.layer.block_size.get() as usize;
467 let block_start =
468 self.layer.block_size.align_down((self.buffer.pos as u64).saturating_sub(1))
469 as usize;
470 let block_bytes = self
471 .buffer
472 .buffer
473 .as_bytes_from(block_start)
474 .get(..block_size)
475 .ok_or_else(|| self.corruption_error("Short block"))?;
476 let seek_table_offset =
478 block_size - usize::from(right_index - 1) * PER_DATA_BLOCK_SEEK_ENTRY_SIZE;
479 let (data_bytes, seek_entries) = block_bytes.split_at(seek_table_offset);
480 let mut found = None;
481 while left_index < (right_index - 1) {
482 let mid_index = left_index + ((right_index - left_index) / 2);
483 let entry_pos = usize::from(mid_index - 1) * PER_DATA_BLOCK_SEEK_ENTRY_SIZE;
484 let offset = LittleEndian::read_u16(
485 &seek_entries[entry_pos..entry_pos + PER_DATA_BLOCK_SEEK_ENTRY_SIZE],
486 ) as usize;
487 if offset >= seek_table_offset || offset <= PER_DATA_BLOCK_HEADER_SIZE {
488 return Err(
489 self.corruption_error(format!("Offset {offset} is out of valid range."))
490 );
491 }
492 match compare_keys(&data_bytes[offset..], target)
493 .map_err(|e| self.corruption_error(e))?
494 {
495 Ordering::Greater => {
496 right_index = mid_index;
497 found = Some((mid_index, offset, false));
498 }
499 Ordering::Equal => {
500 found = Some((mid_index, offset, true));
501 break;
502 }
503 Ordering::Less => left_index = mid_index,
504 }
505 }
506 if let Some((index, offset, exact)) = found {
507 self.item_index = index + 1;
508 self.key = KeyState::InPlace;
509 self.buffer.pos = block_start + offset;
510 return Ok(Some(exact));
511 }
512 }
513 self.item_index = self.item_count;
514 self.key = KeyState::None;
515 Ok(None)
516 }
517}
518
519impl<'iter, K: Key, V: LayerValue, B> KeyOnlyIterator<'iter, K, V, B> {
520 fn new(layer: &'iter LayerData<K>, buffer: B, buffer_pos: usize, pos: u64) -> Self {
523 assert!(layer.block_size.is_aligned(pos));
524 Self {
525 layer,
526 buffer: BufferCursor { buffer, pos: buffer_pos },
527 pos,
528 item_index: 0,
529 item_count: 0,
530 key: KeyState::None,
531 current_block_base_u64: 0,
532 _value_type: PhantomData,
533 }
534 }
535}
536
537impl<'iter, K: Key, V: LayerValue> KeyOnlyIterator<'iter, K, V, ChunkBuffer<'iter>> {
538 fn new_async(layer: &'iter PersistentLayer<K, V>, pos: u64) -> Self {
539 Self::new(
540 &layer.data,
541 ChunkBuffer { handle: &layer.object_handle, chunk: None },
542 (pos % CHUNK_SIZE) as usize,
543 pos,
544 )
545 }
546
547 fn reposition(&self, pos: u64, other: Option<&Self>) -> Self {
550 let chunk_num = (pos / CHUNK_SIZE) as usize;
551 let chunk =
552 std::iter::once(self).chain(other).find_map(|iter| match &iter.buffer.buffer.chunk {
553 Some((n, chunk)) if *n == chunk_num => Some((chunk_num, chunk.clone())),
554 _ => None,
555 });
556 Self::new(
557 self.layer,
558 ChunkBuffer { handle: self.buffer.buffer.handle, chunk },
559 (pos % CHUNK_SIZE) as usize,
560 pos,
561 )
562 }
563
564 async fn advance(&mut self) -> Result<(), Error> {
565 if !self.try_advance()? {
566 let chunk_num = (self.pos / CHUNK_SIZE) as usize;
567 let chunk = self
568 .buffer
569 .buffer
570 .handle
571 .read(self.pos as usize)
572 .await
573 .context("Reading during advance")?;
574 self.buffer.buffer.chunk = Some((chunk_num, chunk));
575 self.buffer.pos = (self.pos % CHUNK_SIZE) as usize;
576 self.read_block_header()?;
577 self.deserialize_current_key()?;
578 }
579 Ok(())
580 }
581
582 fn try_advance(&mut self) -> Result<bool, Error> {
585 if self.item_index >= self.item_count {
586 if self.pos >= self.layer.data_offset() + self.layer.data_size {
587 self.key = KeyState::None;
588 return Ok(true);
589 }
590 let chunk_num = (self.pos / CHUNK_SIZE) as usize;
591 if !matches!(&self.buffer.buffer.chunk, Some((n, _)) if *n == chunk_num) {
592 let Some(chunk) = self.buffer.buffer.handle.try_read(self.pos as usize) else {
593 return Ok(false);
594 };
595 self.buffer.buffer.chunk = Some((chunk_num, chunk));
596 }
597 self.buffer.pos = (self.pos % CHUNK_SIZE) as usize;
598 self.read_block_header()?;
599 }
600 self.deserialize_current_key()?;
601 Ok(true)
602 }
603}
604
605impl<'iter, K: Key, V: LayerValue> KeyOnlyIterator<'iter, K, V, SliceBuffer<'iter>> {
606 fn new_sync(layer: &'iter SyncPersistentLayer<K, V>, pos: u64) -> Self {
607 let slice = layer.object_handle.as_slice().expect("slice must be present");
608 Self::new(
609 &layer.data,
610 SliceBuffer { handle: layer.object_handle.as_ref(), slice },
611 pos as usize,
612 pos,
613 )
614 }
615
616 fn reposition(&self, pos: u64, _other: Option<&Self>) -> Self {
619 Self::new(self.layer, self.buffer.buffer, pos as usize, pos)
620 }
621
622 fn advance(&mut self) -> Result<(), Error> {
623 if self.item_index >= self.item_count {
624 if self.pos >= self.layer.data_offset() + self.layer.data_size {
625 self.key = KeyState::None;
626 return Ok(());
627 }
628 self.buffer.pos = self.pos as usize;
629 self.read_block_header()?;
630 }
631 self.deserialize_current_key()
632 }
633}
634
635struct Iterator<'iter, K: Key, V: LayerValue, B> {
636 inner: KeyOnlyIterator<'iter, K, V, B>,
637 item: Option<Item<K, V>>,
639}
640
641impl<'iter, K: Key, V: LayerValue, B: LayerBuffer> Iterator<'iter, K, V, B> {
642 fn new(mut seek_iterator: KeyOnlyIterator<'iter, K, V, B>) -> Result<Self, Error> {
643 let item = seek_iterator.take_item()?;
644 Ok(Self { inner: seek_iterator, item })
645 }
646}
647
648impl<'iter, K: Key, V: LayerValue> LayerIterator<K, V>
649 for Iterator<'iter, K, V, ChunkBuffer<'iter>>
650{
651 async fn advance(&mut self) -> Result<(), Error> {
652 self.inner.advance().await?;
653 self.item = self.inner.take_item()?;
654 Ok(())
655 }
656
657 fn advance_dyn<'a>(&'a mut self) -> Result<Option<BoxFuture<'a, Result<(), Error>>>, Error> {
658 if self.inner.try_advance()? {
659 self.item = self.inner.take_item()?;
660 Ok(None)
661 } else {
662 Ok(Some(Box::pin(self.advance())))
663 }
664 }
665
666 fn get(&self) -> Option<ItemRef<'_, K, V>> {
667 self.item.as_ref().map(<&Item<K, V>>::into)
668 }
669}
670
671impl<'iter, K: Key, V: LayerValue> LayerIterator<K, V>
672 for Iterator<'iter, K, V, SliceBuffer<'iter>>
673{
674 async fn advance(&mut self) -> Result<(), Error> {
675 self.inner.advance()?;
676 self.item = self.inner.take_item()?;
677 Ok(())
678 }
679
680 fn advance_dyn<'a>(&'a mut self) -> Result<Option<BoxFuture<'a, Result<(), Error>>>, Error> {
681 self.inner.advance()?;
682 self.item = self.inner.take_item()?;
683 Ok(None)
684 }
685
686 fn get(&self) -> Option<ItemRef<'_, K, V>> {
687 self.item.as_ref().map(<&Item<K, V>>::into)
688 }
689}
690
691async fn load_seek_table(
692 object_handle: &(impl ReadObjectHandle + 'static),
693 seek_table_offset: u64,
694 num_data_blocks: u64,
695 version: Version,
696) -> Result<Vec<u64>, Error> {
697 if num_data_blocks == 0 || version <= OLD_KEY_SERIALIZATION_VERSION {
698 return Ok(Vec::new());
701 }
702
703 let seek_table_size = (num_data_blocks as usize) * std::mem::size_of::<u64>();
704 if seek_table_size > MAX_SEEK_TABLE_SIZE {
705 return Err(anyhow!(FxfsError::NotSupported)).context("Seek table too large");
706 }
707 let aligned_size =
708 object_handle.block_size().align_up(seek_table_size as u64).ok_or(FxfsError::TooBig)?
709 as usize;
710 let mut buffer = object_handle.allocate_buffer(aligned_size).await;
711 let bytes_read = object_handle
712 .read_aligned(seek_table_offset, buffer.as_mut())
713 .await
714 .context("Reading seek table blocks")?;
715 ensure!(bytes_read >= seek_table_size, "Short read");
716
717 let mut seek_table = Vec::with_capacity(num_data_blocks as usize);
718 let mut prev = 0;
719 for chunk in buffer.subslice(0..seek_table_size).as_ptr_slice().iter_as::<[u8; 8]>() {
720 let next = u64::from_le_bytes(chunk);
721 if prev > next {
724 return Err(anyhow!(FxfsError::Inconsistent))
725 .context(format!("Seek table entry out of order, {prev:?} > {next:?}"));
726 }
727 prev = next;
728 seek_table.push(next);
729 }
730 Ok(seek_table)
731}
732
733const BLOOM_FILTER_READ_CHUNK_SIZE: usize = 1024 * 1024;
734
735async fn load_bloom_filter<K: FuzzyHash>(
736 handle: &(impl ReadObjectHandle + 'static),
737 bloom_filter_offset: u64,
738 layer_info: &LayerInfo,
739) -> Result<Option<BloomFilterReader<K>>, Error> {
740 if layer_info.bloom_filter_size_bytes == 0 {
741 return Ok(None);
742 }
743 if layer_info.bloom_filter_size_bytes > MAX_BLOOM_FILTER_SIZE {
744 return Err(anyhow!(FxfsError::NotSupported)).context("Bloom filter too large");
745 }
746 let aligned_size = handle
747 .block_size()
748 .align_up(layer_info.bloom_filter_size_bytes as u64)
749 .ok_or(FxfsError::TooBig)? as usize;
750 let mut buffer = handle.allocate_buffer(aligned_size).await;
751 let reads = FuturesUnordered::new();
752 let mut offset = bloom_filter_offset;
753 for chunk in buffer.as_mut().chunks_mut(BLOOM_FILTER_READ_CHUNK_SIZE) {
754 let chunk_len = chunk.len() as u64;
755 reads.push(async move {
756 handle.read_aligned(offset, chunk).await.context("Failed to read")?;
757 Ok::<(), Error>(())
758 });
759 offset += chunk_len;
760 }
761 reads.try_collect::<()>().await?;
762 Ok(Some(BloomFilterReader::read(
763 buffer.subslice(0..layer_info.bloom_filter_size_bytes).as_ptr_slice(),
764 layer_info.bloom_filter_seed,
765 layer_info.bloom_filter_num_hashes,
766 )?))
767}
768
769struct SearchKey<'a, K: Key> {
770 key: &'a K,
771 buf: Vec<u8>,
772 cached_base: Option<u64>,
773}
774
775impl<'a, K: Key> SearchKey<'a, K> {
776 fn new(key: &'a K) -> Self {
777 Self { key, buf: Vec::new(), cached_base: None }
778 }
779
780 fn serialized_with_base(&mut self, base: u64) -> Option<&[u8]> {
783 if self.key.get_leading_u64() < base {
784 return None;
785 }
786 if self.cached_base != Some(base) {
787 self.buf.clear();
788 self.key.serialize_key_with_base_into(&mut self.buf, base);
789 self.cached_base = Some(base);
790 }
791 Some(&self.buf)
792 }
793
794 fn compare<V: LayerValue, B: LayerBuffer>(
795 &mut self,
796 iter: &KeyOnlyIterator<'_, K, V, B>,
797 ) -> Result<Option<Ordering>, Error> {
798 match &iter.key {
799 KeyState::None => Ok(None),
800 KeyState::Deserialized(k) => Ok(Some(k.cmp_upper_bound(self.key))),
801 KeyState::InPlace => {
802 let Some(target) = self.serialized_with_base(iter.current_block_base_u64) else {
803 return Ok(Some(Ordering::Greater));
804 };
805 Ok(Some(
806 compare_keys(iter.buffer.as_bytes(), target)
807 .map_err(|e| iter.corruption_error(e))?,
808 ))
809 }
810 }
811 }
812}
813
814impl<K: FuzzyHash> LayerData<K> {
815 async fn open(handle: &Arc<dyn LayerObject>) -> Result<Self, Error> {
816 let handle_block_size = handle.block_size();
817 let mut buffer = handle.allocate_buffer(handle_block_size.get() as usize).await;
818 handle.read_aligned(0, buffer.as_mut()).await.context("Failed to read first block")?;
819 let mut reader = buffer.as_ptr_slice();
820 let version = Version::deserialize_from(&mut reader)?;
821
822 ensure!(version <= LATEST_VERSION, FxfsError::InvalidVersion);
823 let header = LayerHeader::deserialize_from_version(&mut reader, version)
824 .context("Failed to deserialize header")?;
825 if &header.magic != PERSISTENT_LAYER_MAGIC {
826 return Err(anyhow!(FxfsError::Inconsistent).context("Invalid layer file magic"));
827 }
828 let block_size = BlockSize::from_u64(header.block_size).ok_or_else(|| {
829 anyhow!(FxfsError::Inconsistent)
830 .context(format!("Invalid block size {}", header.block_size))
831 })?;
832 ensure!(block_size <= MAX_BLOCK_SIZE, FxfsError::NotSupported);
833 ensure!(block_size >= MIN_BLOCK_SIZE, FxfsError::NotSupported);
834 if !handle_block_size.is_aligned(block_size.get()) {
835 return Err(anyhow!(FxfsError::Inconsistent)).context(format!(
836 "{block_size} not a multiple of handle block size {handle_block_size}"
837 ));
838 }
839
840 if handle.get_size() < MINIMUM_LAYER_FILE_BLOCKS * block_size {
841 return Err(anyhow!(FxfsError::Inconsistent).context("Layer file too short"));
842 }
843
844 let bs = block_size.get() as usize;
845 let layer_info = {
846 let last_block_offset = handle
847 .get_size()
848 .checked_sub(block_size.get())
849 .ok_or(FxfsError::Inconsistent)
850 .context("Layer file unexpectedly short")?;
851 handle
852 .read_aligned(last_block_offset, buffer.subslice_mut(0..bs))
853 .await
854 .context("Failed to read layer info")?;
855 let layer_info_len =
856 u64::from_le_bytes(buffer.subslice(bs - 8..bs).as_ptr_slice().read().unwrap());
857 let layer_info_offset = bs
858 .checked_sub(std::mem::size_of::<u64>() + layer_info_len as usize)
859 .ok_or(FxfsError::Inconsistent)
860 .context("Invalid layer info length")?;
861 let mut reader = buffer.subslice(layer_info_offset..).as_ptr_slice();
862 LayerInfo::deserialize_from_version(&mut reader, version)
863 .context("Failed to deserialize LayerInfo")?
864 };
865 std::mem::drop(buffer);
866 if layer_info.num_items == 0 && layer_info.num_data_blocks > 0 {
867 return Err(anyhow!(FxfsError::Inconsistent))
868 .context("Invalid num_items/num_data_blocks");
869 }
870 let total_blocks = handle.get_size() / block_size;
871 let bloom_filter_blocks =
872 block_size.align_up_to_blocks(layer_info.bloom_filter_size_bytes as u64);
873 if layer_info.num_data_blocks + bloom_filter_blocks
874 > total_blocks - MINIMUM_LAYER_FILE_BLOCKS
875 {
876 return Err(anyhow!(FxfsError::Inconsistent)).context("Invalid number of blocks");
877 }
878
879 let bloom_filter_offset = block_size * (NUM_HEADER_BLOCKS + layer_info.num_data_blocks);
880 let bloom_filter = if version == LATEST_VERSION {
881 load_bloom_filter(handle, bloom_filter_offset, &layer_info)
882 .await
883 .context("Failed to load bloom filter")?
884 } else {
885 None
889 };
890 let bloom_filter_stats = bloom_filter.as_ref().map(|b| b.stats());
891
892 let seek_offset =
893 block_size * (NUM_HEADER_BLOCKS + layer_info.num_data_blocks + bloom_filter_blocks);
894 let seek_table = load_seek_table(handle, seek_offset, layer_info.num_data_blocks, version)
895 .await
896 .context("Failed to load seek table")?;
897
898 Ok(Self {
899 object_id: handle.object_id(),
900 version,
901 block_size,
902 data_size: block_size * layer_info.num_data_blocks,
903 seek_table,
904 num_items: layer_info.num_items,
905 bloom_filter,
906 bloom_filter_stats,
907 close_event: Mutex::new(Some(Arc::new(DropEvent::new()))),
908 })
909 }
910}
911
912macro_rules! seek_impl {
913 ($self:ident, $bound:ident $(, $await:tt)?) => {{
914 let (key, excluded) = match $bound {
915 Bound::Unbounded => {
916 let mut iterator = $self.key_only_iterator($self.data_offset());
917 iterator.advance()$(.$await)? .context("Unbounded seek advance")?;
918 return Ok(Iterator::new(iterator)?);
919 }
920 Bound::Included(k) => (k, false),
921 Bound::Excluded(k) => (k, true),
922 };
923 let first_data_block_index = $self.data_offset() / $self.data.block_size;
924
925 let (mut left_offset, mut right_offset) = if $self.data.seek_table.is_empty() {
926 ($self.data_offset(), $self.data_offset() + $self.data.data_size)
927 } else {
928 let target = key.get_leading_u64();
934 let right_index =
935 $self.data.seek_table.as_slice().partition_point(|&x| x <= target) as u64;
936 if right_index == 0 {
937 let mut iterator = $self.key_only_iterator($self.data_offset());
938 iterator.advance()$(.$await)? .context("Initial seek advance")?;
939 return Ok(Iterator::new(iterator)?);
940 }
941 let left_index = $self.data.seek_table.as_slice()[..right_index as usize]
944 .partition_point(|&x| x < target)
945 .saturating_sub(1) as u64;
946
947 (
948 (left_index + first_data_block_index) * $self.data.block_size,
949 (right_index + first_data_block_index) * $self.data.block_size,
950 )
951 };
952 let mut left = $self.key_only_iterator(left_offset);
953 left.advance()$(.$await)? .context("Initial seek advance")?;
954 let mut search_key = SearchKey::new(key);
955 match search_key.compare(&left)? {
956 Some(Ordering::Less) => {}
957 Some(Ordering::Equal) if excluded => {
958 left.advance()$(.$await)??;
959 return Ok(Iterator::new(left)?);
960 }
961 _ => return Ok(Iterator::new(left)?),
962 }
963 let mut right = None;
964 while right_offset - left_offset > $self.data.block_size {
965 let mid_offset =
967 $self.data.block_size.align_down(left_offset + (right_offset - left_offset) / 2);
968 let mut iterator = left.reposition(mid_offset, right.as_ref());
969 iterator.advance()$(.$await)??;
970 match search_key.compare(&iterator)?.context("Unexpected EOF")? {
971 Ordering::Greater => {
972 right_offset = mid_offset;
973 right = Some(iterator);
974 }
975 Ordering::Equal => {
976 if excluded {
977 iterator.advance()$(.$await)??;
978 }
979 return Ok(Iterator::new(iterator)?);
980 }
981 Ordering::Less => {
982 left_offset = mid_offset;
983 left = iterator;
984 }
985 }
986 }
987
988 match left.binary_search_block(&mut search_key)? {
990 Some(true) if excluded => left.advance()$(.$await)??,
991 Some(_) => {}
992 None => match right {
997 Some(right) => return Ok(Iterator::new(right)?),
998 None => left.advance()$(.$await)??,
999 },
1000 }
1001 Ok(Iterator::new(left)?)
1002 }};
1003}
1004
1005fn page_size() -> BlockSize {
1006 #[cfg(target_os = "fuchsia")]
1007 {
1008 storage_units::page_size().into()
1009 }
1010 #[cfg(not(target_os = "fuchsia"))]
1011 {
1012 BlockSize::SIZE_4KIB
1013 }
1014}
1015
1016impl<K: Key, V: LayerValue> PersistentLayer<K, V> {
1017 pub async fn open(handle: impl LayerObject + 'static) -> Result<Arc<dyn Layer<K, V>>, Error> {
1018 Self::open_layer(Arc::new(handle)).await
1019 }
1020
1021 pub async fn open_layer(handle: Arc<dyn LayerObject>) -> Result<Arc<dyn Layer<K, V>>, Error> {
1022 let data = LayerData::open(&handle).await?;
1023 if handle.as_slice().is_some() && data.block_size <= page_size() {
1024 Ok(Arc::new(SyncPersistentLayer {
1025 object_handle: handle,
1026 data,
1027 _value_type: PhantomData,
1028 }) as Arc<dyn Layer<K, V>>)
1029 } else {
1030 Ok(Arc::new(PersistentLayer {
1031 object_handle: CachingObjectHandle::new(handle),
1032 data,
1033 _value_type: PhantomData,
1034 }) as Arc<dyn Layer<K, V>>)
1035 }
1036 }
1037
1038 pub async fn open_async(handle: Arc<dyn LayerObject>) -> Result<Arc<Self>, Error> {
1039 let data = LayerData::open(&handle).await?;
1040 Ok(Arc::new(PersistentLayer {
1041 object_handle: CachingObjectHandle::new(handle),
1042 data,
1043 _value_type: PhantomData,
1044 }))
1045 }
1046
1047 pub async fn open_handle(
1050 handle: DataObjectHandle<ObjectStore>,
1051 unwrapped_key: Option<UnwrappedKey>,
1052 ) -> Result<Arc<dyn Layer<K, V>>, Error> {
1053 let layer_object: Arc<dyn LayerObject> =
1054 if let Some(layer_pager) = handle.store().filesystem().layer_pager() {
1055 layer_pager.open_layer(handle, unwrapped_key).await?
1056 } else {
1057 Arc::new(handle)
1058 };
1059 Self::open_layer(layer_object).await
1060 }
1061
1062 fn data_offset(&self) -> u64 {
1063 self.data.data_offset()
1064 }
1065
1066 fn key_only_iterator(&self, pos: u64) -> KeyOnlyIterator<'_, K, V, ChunkBuffer<'_>> {
1067 KeyOnlyIterator::new_async(self, pos)
1068 }
1069
1070 async fn seek<'a>(
1071 &'a self,
1072 bound: Bound<&K>,
1073 ) -> Result<Iterator<'a, K, V, ChunkBuffer<'a>>, Error> {
1074 seek_impl!(self, bound, await)
1075 }
1076}
1077
1078impl<K: Key, V: LayerValue> SyncPersistentLayer<K, V> {
1079 pub async fn open(handle: Arc<dyn LayerObject>) -> Result<Arc<Self>, Error> {
1080 ensure!(handle.as_slice().is_some(), FxfsError::InvalidArgs);
1081 let data = LayerData::open(&handle).await?;
1082 ensure!(data.block_size <= page_size(), FxfsError::NotSupported);
1083 Ok(Arc::new(Self { object_handle: handle, data, _value_type: PhantomData }))
1084 }
1085
1086 fn data_offset(&self) -> u64 {
1087 self.data.data_offset()
1088 }
1089
1090 fn key_only_iterator(&self, pos: u64) -> KeyOnlyIterator<'_, K, V, SliceBuffer<'_>> {
1091 KeyOnlyIterator::new_sync(self, pos)
1092 }
1093
1094 fn seek<'a>(&'a self, bound: Bound<&K>) -> Result<Iterator<'a, K, V, SliceBuffer<'a>>, Error> {
1095 seek_impl!(self, bound)
1096 }
1097}
1098
1099async fn unwrap_layer_key(
1101 store: &ObjectStore,
1102 object_id: u64,
1103 crypt: &dyn Crypt,
1104) -> Result<Option<UnwrappedKey>, Error> {
1105 if store.filesystem().layer_pager().is_none() {
1106 return Ok(None);
1107 }
1108 let keys = store
1109 .get_keys(object_id)
1110 .await
1111 .with_context(|| format!("Failed to get keys for layer file {object_id}"))?;
1112 let (_, key) = keys
1113 .first()
1114 .ok_or_else(|| anyhow!(FxfsError::Inconsistent))
1115 .with_context(|| format!("Missing key for encrypted layer file {object_id}"))?;
1116 let wrapped_key = WrappedKey::from(key.clone());
1117 let unwrapped = crypt
1118 .unwrap_key(&wrapped_key, object_id)
1119 .await
1120 .with_context(|| format!("Failed to unwrap key for layer file {object_id}"))?;
1121 Ok(Some(unwrapped))
1122}
1123
1124pub async fn layer_from_handle<K: Key, V: LayerValue>(
1129 handle: DataObjectHandle<ObjectStore>,
1130 unwrapped_key: Option<UnwrappedKey>,
1131) -> Result<Arc<dyn Layer<K, V>>, Error> {
1132 PersistentLayer::open_handle(handle, unwrapped_key).await
1133}
1134
1135pub async fn open_layers<K: Key, V: LayerValue>(
1140 store: &Arc<ObjectStore>,
1141 object_ids: impl IntoIterator<Item = u64>,
1142 crypt: Option<Arc<dyn Crypt>>,
1143) -> Result<(Vec<Arc<dyn Layer<K, V>>>, u64), Error> {
1144 let mut layers = Vec::new();
1145 let mut total_size = 0;
1146 for object_id in object_ids {
1147 let handle =
1148 ObjectStore::open_object(store, object_id, HandleOptions::default(), crypt.clone())
1149 .await
1150 .with_context(|| format!("Failed to open layer file {object_id}"))?;
1151 total_size += handle.get_size();
1152 let unwrapped_key = if let Some(crypt) = &crypt {
1153 unwrap_layer_key(store, object_id, crypt.as_ref()).await?
1154 } else {
1155 None
1156 };
1157 layers.push(PersistentLayer::open_handle(handle, unwrapped_key).await?);
1158 }
1159 Ok((layers, total_size))
1160}
1161
1162#[async_trait]
1163impl<K: Key, V: LayerValue> Layer<K, V> for PersistentLayer<K, V> {
1164 fn handle(&self) -> Option<&dyn ReadObjectHandle> {
1165 Some(self.object_handle.source())
1166 }
1167
1168 fn purge_cached_data(&self) {
1169 self.object_handle.purge();
1170 }
1171
1172 fn clear_cached_data(&self) {
1173 self.object_handle.clear();
1174 }
1175
1176 async fn seek<'a>(&'a self, bound: Bound<&K>) -> Result<BoxedLayerIterator<'a, K, V>, Error> {
1177 Ok(Box::new(PersistentLayer::seek(self, bound).await?))
1178 }
1179
1180 fn len(&self) -> usize {
1181 self.data.num_items
1182 }
1183
1184 fn maybe_contains_key(&self, key: &K) -> MaybeContainsKey {
1185 self.data.bloom_filter.as_ref().map_or(MaybeContainsKey::Maybe, |f| f.maybe_contains(key))
1186 }
1187
1188 fn has_bloom_filter(&self) -> bool {
1189 self.data.bloom_filter.is_some()
1190 }
1191
1192 async fn key_exists(&self, key: &K) -> Result<Existence, Error> {
1193 match &self.data.bloom_filter {
1194 Some(filter) => Ok(match filter.maybe_contains(key) {
1195 MaybeContainsKey::False => Existence::Missing,
1196 MaybeContainsKey::Maybe | MaybeContainsKey::RangeKeyTooLarge => {
1197 Existence::MaybeExists
1198 }
1199 }),
1200 None => {
1201 let iter = self.seek(Bound::Included(key)).await?;
1202 Ok(iter.get().map_or(Existence::Missing, |i| {
1203 if i.key.cmp_upper_bound(key).is_eq() {
1204 Existence::Exists
1205 } else {
1206 Existence::Missing
1207 }
1208 }))
1209 }
1210 }
1211 }
1212
1213 fn lock(&self) -> Option<Arc<DropEvent>> {
1214 self.data.close_event.lock().clone()
1215 }
1216
1217 async fn close(&self) {
1218 let listener = self.data.close_event.lock().take().expect("close already called").listen();
1219 listener.await;
1220 self.object_handle.source().close().await;
1221 }
1222
1223 fn get_version(&self) -> Version {
1224 return self.data.version;
1225 }
1226
1227 fn record_inspect_data(self: Arc<Self>, node: &fuchsia_inspect::Node) {
1228 node.record_uint("num_items", self.data.num_items as u64);
1229 node.record_bool("persistent", true);
1230 node.record_uint("size", self.object_handle.source().get_size());
1231 if let Some(stats) = self.data.bloom_filter_stats.as_ref() {
1232 node.record_child("bloom_filter", move |node| {
1233 node.record_uint("size", stats.size as u64);
1234 node.record_uint("num_hashes", stats.num_hashes as u64);
1235 node.record_uint("fill_percentage", stats.fill_percentage as u64);
1236 });
1237 }
1238 }
1239}
1240
1241#[async_trait]
1242impl<K: Key, V: LayerValue> Layer<K, V> for SyncPersistentLayer<K, V> {
1243 fn handle(&self) -> Option<&dyn ReadObjectHandle> {
1244 Some(&self.object_handle)
1245 }
1246
1247 fn purge_cached_data(&self) {
1248 self.object_handle.purge_cached_data();
1249 }
1250
1251 fn clear_cached_data(&self) {
1252 self.object_handle.purge_cached_data();
1253 }
1254
1255 async fn seek<'a>(&'a self, bound: Bound<&K>) -> Result<BoxedLayerIterator<'a, K, V>, Error> {
1256 Ok(Box::new(SyncPersistentLayer::seek(self, bound)?))
1257 }
1258
1259 fn len(&self) -> usize {
1260 self.data.num_items
1261 }
1262
1263 fn maybe_contains_key(&self, key: &K) -> MaybeContainsKey {
1264 self.data.bloom_filter.as_ref().map_or(MaybeContainsKey::Maybe, |f| f.maybe_contains(key))
1265 }
1266
1267 fn has_bloom_filter(&self) -> bool {
1268 self.data.bloom_filter.is_some()
1269 }
1270
1271 async fn key_exists(&self, key: &K) -> Result<Existence, Error> {
1272 match &self.data.bloom_filter {
1273 Some(filter) => Ok(match filter.maybe_contains(key) {
1274 MaybeContainsKey::False => Existence::Missing,
1275 MaybeContainsKey::Maybe | MaybeContainsKey::RangeKeyTooLarge => {
1276 Existence::MaybeExists
1277 }
1278 }),
1279 None => {
1280 let iter = self.seek(Bound::Included(key))?;
1281 Ok(iter.get().map_or(Existence::Missing, |i| {
1282 if i.key.cmp_upper_bound(key).is_eq() {
1283 Existence::Exists
1284 } else {
1285 Existence::Missing
1286 }
1287 }))
1288 }
1289 }
1290 }
1291
1292 fn lock(&self) -> Option<Arc<DropEvent>> {
1293 self.data.close_event.lock().clone()
1294 }
1295
1296 async fn close(&self) {
1297 let listener = self.data.close_event.lock().take().expect("close already called").listen();
1298 listener.await;
1299 self.object_handle.close().await;
1300 }
1301
1302 fn get_version(&self) -> Version {
1303 return self.data.version;
1304 }
1305
1306 fn record_inspect_data(self: Arc<Self>, node: &fuchsia_inspect::Node) {
1307 node.record_uint("num_items", self.data.num_items as u64);
1308 node.record_bool("persistent", true);
1309 node.record_uint("size", self.object_handle.get_size());
1310 if let Some(stats) = self.data.bloom_filter_stats.as_ref() {
1311 node.record_child("bloom_filter", move |node| {
1312 node.record_uint("size", stats.size as u64);
1313 node.record_uint("num_hashes", stats.num_hashes as u64);
1314 node.record_uint("fill_percentage", stats.fill_percentage as u64);
1315 });
1316 }
1317 }
1318}
1319
1320const_assert!(MAX_BLOCK_SIZE.size() <= u16::MAX as u64 + 1);
1322
1323pub struct PersistentLayerWriter<W: WriteBytes, K: Key, V: LayerValue> {
1326 writer: W,
1327 version: Version,
1328 block_size: BlockSize,
1329 buf: Vec<u8>,
1330 buf_item_count: LayerWriterBufItemCount,
1331 item_count: usize,
1332 block_offsets: Vec<u16>,
1333 block_keys: Vec<u64>,
1334 bloom_filter: BloomFilterWriter<K>,
1335 _value: PhantomData<V>,
1336}
1337
1338impl<W: WriteBytes, K: Key, V: LayerValue> PersistentLayerWriter<W, K, V> {
1339 pub async fn new(writer: W, num_items: usize, block_size: BlockSize) -> Result<Self, Error> {
1341 Self::new_with_version(writer, num_items, block_size, LATEST_VERSION).await
1342 }
1343
1344 pub(crate) async fn new_with_version(
1345 mut writer: W,
1346 num_items: usize,
1347 block_size: BlockSize,
1348 version: Version,
1349 ) -> Result<Self, Error> {
1350 ensure!(block_size <= MAX_BLOCK_SIZE, FxfsError::NotSupported);
1351 ensure!(block_size >= MIN_BLOCK_SIZE, FxfsError::NotSupported);
1352
1353 let header =
1355 LayerHeader { magic: PERSISTENT_LAYER_MAGIC.clone(), block_size: block_size.get() };
1356 let mut buf = vec![0u8; block_size.get() as usize];
1357 {
1358 let mut cursor = std::io::Cursor::new(&mut buf[..]);
1359 version.serialize_into(&mut cursor)?;
1360 header.serialize_into(&mut cursor)?;
1361 }
1362 writer.write_bytes(&buf[..]).await?;
1363
1364 let seed: u64 = rand::random();
1365 Ok(Self {
1366 writer,
1367 version,
1368 block_size,
1369 buf: Vec::new(),
1370 buf_item_count: LayerWriterBufItemCount(0),
1371 item_count: 0,
1372 block_offsets: Vec::new(),
1373 block_keys: Vec::new(),
1374 bloom_filter: BloomFilterWriter::new(seed, num_items),
1375 _value: PhantomData,
1376 })
1377 }
1378
1379 async fn write_block(&mut self) -> Result<(), Error> {
1384 if *self.buf_item_count == 0 {
1385 return Ok(());
1386 }
1387 let seek_table_size = self.block_offsets.len() * PER_DATA_BLOCK_SEEK_ENTRY_SIZE;
1388 assert!(
1389 PER_DATA_BLOCK_HEADER_SIZE + seek_table_size + self.buf.len()
1390 <= self.block_size.get() as usize
1391 );
1392 let mut cursor = std::io::Cursor::new(vec![0u8; self.block_size.get() as usize]);
1393 cursor.write_u16::<LittleEndian>(*self.buf_item_count)?;
1394 cursor.write_all(&self.buf)?;
1395 cursor.set_position(self.block_size - seek_table_size as u64);
1396 for &offset in &self.block_offsets {
1398 cursor.write_u16::<LittleEndian>(offset)?;
1399 }
1400 self.writer.write_bytes(cursor.get_ref()).await?;
1401 debug!(item_count = *self.buf_item_count, byte_count = self.buf.len(); "wrote items");
1402 self.buf.clear();
1403 *self.buf_item_count = 0;
1404 self.block_offsets.clear();
1405 Ok(())
1406 }
1407
1408 async fn write_seek_table(&mut self) -> Result<usize, Error> {
1413 let keys = if self.version <= OLD_KEY_SERIALIZATION_VERSION {
1414 self.block_keys.get(1..).unwrap_or(&[])
1415 } else {
1416 &self.block_keys
1417 };
1418 if keys.len() == 0 {
1419 return Ok(0);
1420 }
1421 let size = keys.len() * std::mem::size_of::<u64>();
1422 self.buf.resize(size, 0);
1423 let mut len = 0;
1424 for key in keys {
1425 LittleEndian::write_u64(&mut self.buf[len..len + std::mem::size_of::<u64>()], *key);
1426 len += std::mem::size_of::<u64>();
1427 }
1428 self.writer.write_bytes(&self.buf).await?;
1429 Ok(size)
1430 }
1431
1432 async fn write_info(
1435 &mut self,
1436 num_data_blocks: u64,
1437 bloom_filter_size_bytes: usize,
1438 seek_table_len: usize,
1439 ) -> Result<(), Error> {
1440 let block_size = self.writer.block_size().get() as usize;
1441 let layer_info = LayerInfo {
1442 num_items: self.item_count,
1443 num_data_blocks,
1444 bloom_filter_size_bytes,
1445 bloom_filter_seed: self.bloom_filter.seed(),
1446 bloom_filter_num_hashes: self.bloom_filter.num_hashes(),
1447 };
1448 self.buf.clear();
1449 layer_info.serialize_into(&mut self.buf)?;
1450 let layer_info_len = self.buf.len() as u64;
1451 self.buf.write_u64::<LittleEndian>(layer_info_len)?;
1452 let actual_len = self.buf.len();
1453
1454 let avail_in_block =
1457 block_size - (seek_table_len as u64 % self.writer.block_size()) as usize;
1458 let to_skip = if avail_in_block < actual_len {
1459 block_size + avail_in_block - actual_len
1460 } else {
1461 avail_in_block - actual_len
1462 };
1463 self.buf.resize(to_skip + actual_len, 0);
1464 self.buf.copy_within(0..actual_len, to_skip);
1465 self.buf[..to_skip].fill(0);
1466 self.writer.write_bytes(&self.buf).await?;
1467 Ok(())
1468 }
1469
1470 async fn write_bloom_filter(&mut self) -> Result<usize, Error> {
1473 if self.data_blocks() < MINIMUM_DATA_BLOCKS_FOR_BLOOM_FILTER {
1474 return Ok(0);
1475 }
1476 let size =
1478 self.block_size.align_up(self.bloom_filter.serialized_size() as u64).unwrap() as usize;
1479 self.buf.resize(size, 0);
1480 let mut cursor = std::io::Cursor::new(&mut self.buf);
1481 self.bloom_filter.write(&mut cursor)?;
1482 self.writer.write_bytes(&self.buf).await?;
1483 Ok(self.bloom_filter.serialized_size())
1484 }
1485
1486 #[cfg(test)]
1489 pub(crate) fn bloom_filter(&mut self) -> &mut BloomFilterWriter<K> {
1490 &mut self.bloom_filter
1491 }
1492
1493 fn serialize_item(&mut self, item: ItemRef<'_, K, V>) -> Result<(), Error> {
1494 if self.version <= OLD_KEY_SERIALIZATION_VERSION {
1495 item.key.serialize_into(&mut self.buf)?;
1496 } else {
1497 item.key.serialize_key_with_base_into(&mut self.buf, *self.block_keys.last().unwrap());
1498 }
1499 item.value.serialize_into(&mut self.buf)?;
1500 Ok(())
1501 }
1502
1503 fn data_blocks(&self) -> usize {
1504 self.block_keys.len()
1505 }
1506}
1507
1508impl<W: WriteBytes + Send, K: Key, V: LayerValue> LayerWriter<K, V>
1509 for PersistentLayerWriter<W, K, V>
1510{
1511 async fn write(&mut self, item: ItemRef<'_, K, V>) -> Result<(), Error> {
1512 let len = self.buf.len();
1514 if self.block_keys.is_empty() {
1518 self.block_keys.push(item.key.get_leading_u64());
1519 }
1520 self.serialize_item(item)?;
1521
1522 let mut added_offset = false;
1523 if *self.buf_item_count > 0 {
1525 self.block_offsets.push(u16::try_from(len + PER_DATA_BLOCK_HEADER_SIZE).unwrap());
1526 added_offset = true;
1527 }
1528
1529 if PER_DATA_BLOCK_HEADER_SIZE
1532 + self.buf.len()
1533 + (self.block_offsets.len() * PER_DATA_BLOCK_SEEK_ENTRY_SIZE)
1534 > self.block_size.get() as usize - 1
1535 {
1536 if added_offset {
1537 self.block_offsets.pop();
1540 }
1541 self.buf.truncate(len);
1542 self.write_block().await?;
1543
1544 self.block_keys.push(item.key.get_leading_u64());
1546 self.serialize_item(item)?;
1547 }
1548
1549 self.bloom_filter.insert(&item.key);
1550 *self.buf_item_count += 1;
1551 self.item_count += 1;
1552 Ok(())
1553 }
1554
1555 async fn complete(mut self) -> Result<u64, Error> {
1556 self.write_block().await?;
1557 let data_blocks = self.data_blocks() as u64;
1558 let bloom_filter_len = self.write_bloom_filter().await?;
1559 let seek_table_len = self.write_seek_table().await?;
1560 self.write_info(data_blocks, bloom_filter_len, seek_table_len).await?;
1561 self.writer.complete().await
1562 }
1563}
1564
1565#[repr(transparent)]
1567struct LayerWriterBufItemCount(u16);
1568
1569impl Drop for LayerWriterBufItemCount {
1570 fn drop(&mut self) {
1571 debug_assert!(self.0 == 0, "Dropping unwritten items; did you forget to call complete?");
1572 if self.0 > 0 {
1573 warn!("Dropping unwritten items; did you forget to call complete?");
1574 }
1575 }
1576}
1577
1578impl std::ops::Deref for LayerWriterBufItemCount {
1579 type Target = u16;
1580 fn deref(&self) -> &u16 {
1581 &self.0
1582 }
1583}
1584
1585impl std::ops::DerefMut for LayerWriterBufItemCount {
1586 fn deref_mut(&mut self) -> &mut u16 {
1587 &mut self.0
1588 }
1589}
1590
1591#[cfg(test)]
1592mod tests {
1593 use super::{
1594 BlockSize, FxfsError, PersistentLayer, PersistentLayerWriter, SyncPersistentLayer,
1595 };
1596 use crate::filesystem::MAX_BLOCK_SIZE;
1597 use crate::lsm_tree::LayerIterator;
1598 use crate::lsm_tree::persistent_layer::MINIMUM_DATA_BLOCKS_FOR_BLOOM_FILTER;
1599 use crate::lsm_tree::testing::TestKey;
1600 use crate::lsm_tree::types::{
1601 Existence, Item, ItemRef, Layer, LayerWriter, MaybeContainsKey, OrdUpperBound,
1602 };
1603 use crate::object_handle::{
1604 LayerObject, ObjectHandle, ReadObjectHandle, WriteBytes, WriteObjectHandle,
1605 };
1606 use crate::object_store::AttributeId;
1607 use crate::object_store::allocator::AllocatorKey;
1608 use crate::object_store::extent::{Extent, MIN_BLOCK_SIZE};
1609 use crate::object_store::object_record::ObjectKey;
1610 use crate::round::round_up;
1611 use crate::serialized_types::OLD_KEY_SERIALIZATION_VERSION;
1612 use crate::testing::fake_object::{FakeObject, FakeObjectHandle};
1613 use crate::testing::writer::Writer;
1614 use anyhow::Error;
1615 use async_trait::async_trait;
1616 use std::fmt::Debug;
1617 use std::ops::{Bound, Range};
1618 use std::sync::Arc;
1619 use std::sync::atomic::{AtomicBool, Ordering};
1620 use storage_device::buffer::{BufferFuture, MutableBufferRef};
1621
1622 impl<W: WriteBytes> Debug for PersistentLayerWriter<W, i32, i32> {
1623 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> Result<(), std::fmt::Error> {
1624 f.debug_struct("rPersistentLayerWriter")
1625 .field("block_size", &self.block_size)
1626 .field("item_count", &*self.buf_item_count)
1627 .finish()
1628 }
1629 }
1630
1631 #[fuchsia::test]
1632 async fn test_iterate_after_write() {
1633 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
1634 const ITEM_COUNT: i32 = 10000;
1635
1636 let handle = FakeObjectHandle::new(Arc::new(FakeObject::new()));
1637 {
1638 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
1639 Writer::new(&handle).await,
1640 ITEM_COUNT as usize * 4,
1641 BLOCK_SIZE,
1642 )
1643 .await
1644 .expect("writer new");
1645 for i in 0..ITEM_COUNT {
1646 writer.write(Item::new(i, i).as_item_ref()).await.expect("write failed");
1647 }
1648 writer.complete().await.expect("flush failed");
1649 }
1650 let layer = PersistentLayer::<i32, i32>::open(handle).await.expect("new failed");
1651 let mut iterator = layer.seek(Bound::Unbounded).await.expect("seek failed");
1652 for i in 0..ITEM_COUNT {
1653 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1654 assert_eq!((key, value), (&i, &i));
1655 iterator.advance().await.expect("failed to advance");
1656 }
1657 assert!(iterator.get().is_none());
1658 }
1659
1660 #[fuchsia::test]
1661 async fn test_seek_after_write() {
1662 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
1663 const ITEM_COUNT: i32 = 5000;
1664
1665 let handle = FakeObjectHandle::new(Arc::new(FakeObject::new()));
1666 {
1667 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
1668 Writer::new(&handle).await,
1669 ITEM_COUNT as usize * 18,
1670 BLOCK_SIZE,
1671 )
1672 .await
1673 .expect("writer new");
1674 for i in 0..ITEM_COUNT {
1675 writer.write(Item::new(i * 2, i * 2).as_item_ref()).await.expect("write failed");
1677 }
1678 writer.complete().await.expect("flush failed");
1679 }
1680 let layer = PersistentLayer::<i32, i32>::open(handle).await.expect("new failed");
1681 for i in 0..ITEM_COUNT * 2 {
1683 let expected = round_up(i, 2).unwrap();
1686 let mut iterator = layer.seek(Bound::Included(&i)).await.expect("failed to seek");
1687 if i >= (ITEM_COUNT * 2) - 1 {
1690 assert!(iterator.get().is_none());
1691 } else {
1692 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1693 assert_eq!((key, value), (&expected, &expected));
1694 }
1695
1696 iterator.advance().await.expect("failed to advance");
1698 if i >= (ITEM_COUNT * 2) - 3 {
1702 assert!(iterator.get().is_none());
1703 } else {
1704 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1705 let next = expected + 2;
1706 assert_eq!((key, value), (&next, &next));
1707 }
1708 }
1709 }
1710
1711 #[fuchsia::test]
1712 async fn test_seek_unbounded() {
1713 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
1714 const ITEM_COUNT: i32 = 1000;
1715
1716 let handle = FakeObjectHandle::new(Arc::new(FakeObject::new()));
1717 {
1718 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
1719 Writer::new(&handle).await,
1720 ITEM_COUNT as usize * 18,
1721 BLOCK_SIZE,
1722 )
1723 .await
1724 .expect("writer new");
1725 for i in 0..ITEM_COUNT {
1726 writer.write(Item::new(i, i).as_item_ref()).await.expect("write failed");
1727 }
1728 writer.complete().await.expect("flush failed");
1729 }
1730 let layer = PersistentLayer::<i32, i32>::open(handle).await.expect("new failed");
1731 let mut iterator = layer.seek(Bound::Unbounded).await.expect("failed to seek");
1732 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1733 assert_eq!((key, value), (&0, &0));
1734
1735 iterator.advance().await.expect("failed to advance");
1737 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1738 assert_eq!((key, value), (&1, &1));
1739 }
1740
1741 #[fuchsia::test]
1742 async fn test_zero_items() {
1743 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
1744
1745 let handle = FakeObjectHandle::new(Arc::new(FakeObject::new()));
1746 {
1747 let writer = PersistentLayerWriter::<_, i32, i32>::new(
1748 Writer::new(&handle).await,
1749 0,
1750 BLOCK_SIZE,
1751 )
1752 .await
1753 .expect("writer new");
1754 writer.complete().await.expect("flush failed");
1755 }
1756
1757 let layer = PersistentLayer::<i32, i32>::open(handle).await.expect("new failed");
1758 let iterator = (layer.as_ref() as &dyn Layer<i32, i32>)
1759 .seek(Bound::Unbounded)
1760 .await
1761 .expect("seek failed");
1762 assert!(iterator.get().is_none())
1763 }
1764
1765 #[fuchsia::test]
1766 async fn test_one_item() {
1767 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
1768
1769 let handle = FakeObjectHandle::new(Arc::new(FakeObject::new()));
1770 {
1771 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
1772 Writer::new(&handle).await,
1773 1,
1774 BLOCK_SIZE,
1775 )
1776 .await
1777 .expect("writer new");
1778 writer.write(Item::new(42, 42).as_item_ref()).await.expect("write failed");
1779 writer.complete().await.expect("flush failed");
1780 }
1781
1782 let layer = PersistentLayer::<i32, i32>::open(handle).await.expect("new failed");
1783 {
1784 let mut iterator = (layer.as_ref() as &dyn Layer<i32, i32>)
1785 .seek(Bound::Unbounded)
1786 .await
1787 .expect("seek failed");
1788 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1789 assert_eq!((key, value), (&42, &42));
1790 iterator.advance().await.expect("failed to advance");
1791 assert!(iterator.get().is_none())
1792 }
1793 {
1794 let mut iterator = (layer.as_ref() as &dyn Layer<i32, i32>)
1795 .seek(Bound::Included(&30))
1796 .await
1797 .expect("seek failed");
1798 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1799 assert_eq!((key, value), (&42, &42));
1800 iterator.advance().await.expect("failed to advance");
1801 assert!(iterator.get().is_none())
1802 }
1803 {
1804 let mut iterator = (layer.as_ref() as &dyn Layer<i32, i32>)
1805 .seek(Bound::Included(&42))
1806 .await
1807 .expect("seek failed");
1808 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1809 assert_eq!((key, value), (&42, &42));
1810 iterator.advance().await.expect("failed to advance");
1811 assert!(iterator.get().is_none())
1812 }
1813 {
1814 let iterator = (layer.as_ref() as &dyn Layer<i32, i32>)
1815 .seek(Bound::Included(&43))
1816 .await
1817 .expect("seek failed");
1818 assert!(iterator.get().is_none())
1819 }
1820 }
1821
1822 #[fuchsia::test]
1823 async fn test_large_block_size() {
1824 const BLOCK_SIZE: BlockSize = MAX_BLOCK_SIZE;
1826 let item_count: i32 = ((BLOCK_SIZE.get() as i32) / 18) * 3;
1828
1829 let handle = FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
1830 {
1831 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
1832 Writer::new(&handle).await,
1833 item_count as usize * 18,
1834 BLOCK_SIZE,
1835 )
1836 .await
1837 .expect("writer new");
1838 for i in 2000000000..(2000000000 + item_count) {
1840 writer.write(Item::new(i, i).as_item_ref()).await.expect("write failed");
1841 }
1842 writer.complete().await.expect("flush failed");
1843 }
1844
1845 let layer = PersistentLayer::<i32, i32>::open(handle).await.expect("new failed");
1846 let mut iterator = layer.seek(Bound::Unbounded).await.expect("seek failed");
1847 for i in 2000000000..(2000000000 + item_count) {
1848 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1849 assert_eq!((key, value), (&i, &i));
1850 iterator.advance().await.expect("failed to advance");
1851 }
1852 assert!(iterator.get().is_none());
1853 }
1854
1855 #[fuchsia::test]
1856 async fn test_overlarge_block_size() {
1857 const BLOCK_SIZE: BlockSize = BlockSize::from_u64(MAX_BLOCK_SIZE.size() * 2).unwrap();
1859
1860 let handle = FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
1861 PersistentLayerWriter::<_, i32, i32>::new(Writer::new(&handle).await, 0, BLOCK_SIZE)
1862 .await
1863 .expect_err("Creating writer with overlarge block size.");
1864 }
1865
1866 #[fuchsia::test]
1867 async fn test_seek_bound_excluded() {
1868 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
1869 const ITEM_COUNT: i32 = 10000;
1870
1871 let handle = FakeObjectHandle::new(Arc::new(FakeObject::new()));
1872 {
1873 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
1874 Writer::new(&handle).await,
1875 ITEM_COUNT as usize * 18,
1876 BLOCK_SIZE,
1877 )
1878 .await
1879 .expect("writer new");
1880 for i in 0..ITEM_COUNT {
1881 writer.write(Item::new(i, i).as_item_ref()).await.expect("write failed");
1882 }
1883 writer.complete().await.expect("flush failed");
1884 }
1885 let layer = PersistentLayer::<i32, i32>::open(handle).await.expect("new failed");
1886
1887 for i in 9982..ITEM_COUNT {
1888 let mut iterator = layer.seek(Bound::Excluded(&i)).await.expect("failed to seek");
1889 let i_plus_one = i + 1;
1890 if i_plus_one < ITEM_COUNT {
1891 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1892
1893 assert_eq!((key, value), (&i_plus_one, &i_plus_one));
1894
1895 iterator.advance().await.expect("failed to advance");
1897 let i_plus_two = i + 2;
1898 if i_plus_two < ITEM_COUNT {
1899 let ItemRef { key, value, .. } = iterator.get().expect("missing item");
1900 assert_eq!((key, value), (&i_plus_two, &i_plus_two));
1901 } else {
1902 assert!(iterator.get().is_none());
1903 }
1904 } else {
1905 assert!(iterator.get().is_none());
1906 }
1907 }
1908 }
1909
1910 fn generate_extents(
1914 object_id: u64,
1915 base_offset: u64,
1916 count: u64,
1917 ) -> (Vec<Item<ObjectKey, u64>>, u64) {
1918 let mut items = Vec::new();
1919 for i in 0..count {
1920 items.push(Item::new(
1921 ObjectKey::extent(
1922 object_id,
1923 AttributeId::TEST_ID,
1924 (base_offset + i) * MIN_BLOCK_SIZE..(base_offset + i + 1) * MIN_BLOCK_SIZE,
1925 ),
1926 object_id,
1927 ));
1928 }
1929 (items, object_id + 1)
1930 }
1931
1932 fn generate_objects(object_id_range: Range<u64>) -> (Vec<Item<ObjectKey, u64>>, u64) {
1936 let mut items = Vec::new();
1937 let end = object_id_range.end;
1938 for object_id in object_id_range {
1939 items.push(Item::new(ObjectKey::object(object_id), object_id));
1940 }
1941 (items, end)
1942 }
1943
1944 #[fuchsia::test]
1947 async fn test_block_seek_duplicate_leading_u64() {
1948 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
1950 const ITEMS_PER_PHASE: u64 = 50;
1951
1952 let mut to_find = Vec::new();
1953
1954 let handle = FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
1955 {
1956 let mut items = Vec::new();
1957 let mut object_id = u32::MAX as u64 + 1;
1959
1960 {
1963 let base_extent_offset = 0;
1964 let (mut generated, next_object_id) =
1965 generate_extents(object_id, base_extent_offset, ITEMS_PER_PHASE * 3);
1966 items.append(&mut generated);
1967 let count = ITEMS_PER_PHASE * 3;
1968 to_find.push(ObjectKey::extent(
1969 object_id,
1970 AttributeId::TEST_ID,
1971 base_extent_offset * MIN_BLOCK_SIZE..(base_extent_offset + 1) * MIN_BLOCK_SIZE,
1972 ));
1973 to_find.push(ObjectKey::extent(
1974 object_id,
1975 AttributeId::TEST_ID,
1976 (base_extent_offset + count / 2) * MIN_BLOCK_SIZE
1977 ..(base_extent_offset + count / 2 + 1) * MIN_BLOCK_SIZE,
1978 ));
1979 to_find.push(ObjectKey::extent(
1980 object_id,
1981 AttributeId::TEST_ID,
1982 (base_extent_offset + count - 1) * MIN_BLOCK_SIZE
1983 ..(base_extent_offset + count) * MIN_BLOCK_SIZE,
1984 ));
1985 object_id = next_object_id;
1986 }
1987
1988 {
1990 let (mut generated, next_object_id) =
1991 generate_objects(object_id..object_id + ITEMS_PER_PHASE * 3);
1992 items.append(&mut generated);
1993 object_id = next_object_id;
1994 }
1995
1996 {
1999 let base_extent_offset = 1000;
2000 let (mut generated, next_object_id) =
2001 generate_extents(object_id, base_extent_offset, ITEMS_PER_PHASE * 3);
2002 items.append(&mut generated);
2003 let count = ITEMS_PER_PHASE * 3;
2004 to_find.push(ObjectKey::extent(
2005 object_id,
2006 AttributeId::TEST_ID,
2007 base_extent_offset * MIN_BLOCK_SIZE..(base_extent_offset + 1) * MIN_BLOCK_SIZE,
2008 ));
2009 to_find.push(ObjectKey::extent(
2010 object_id,
2011 AttributeId::TEST_ID,
2012 (base_extent_offset + count / 2) * MIN_BLOCK_SIZE
2013 ..(base_extent_offset + count / 2 + 1) * MIN_BLOCK_SIZE,
2014 ));
2015 to_find.push(ObjectKey::extent(
2016 object_id,
2017 AttributeId::TEST_ID,
2018 (base_extent_offset + count - 1) * MIN_BLOCK_SIZE
2019 ..(base_extent_offset + count) * MIN_BLOCK_SIZE,
2020 ));
2021 object_id = next_object_id;
2022 }
2023
2024 {
2026 let (mut generated, next_object_id) =
2027 generate_objects(object_id..object_id + ITEMS_PER_PHASE * 3);
2028 items.append(&mut generated);
2029 object_id = next_object_id;
2030 }
2031
2032 {
2035 let base_extent_offset = 2000;
2036 let (mut generated, _) =
2037 generate_extents(object_id, base_extent_offset, ITEMS_PER_PHASE * 3);
2038 items.append(&mut generated);
2039 let count = ITEMS_PER_PHASE * 3;
2040 to_find.push(ObjectKey::extent(
2041 object_id,
2042 AttributeId::TEST_ID,
2043 base_extent_offset * MIN_BLOCK_SIZE..(base_extent_offset + 1) * MIN_BLOCK_SIZE,
2044 ));
2045 to_find.push(ObjectKey::extent(
2046 object_id,
2047 AttributeId::TEST_ID,
2048 (base_extent_offset + count / 2) * MIN_BLOCK_SIZE
2049 ..(base_extent_offset + count / 2 + 1) * MIN_BLOCK_SIZE,
2050 ));
2051 to_find.push(ObjectKey::extent(
2052 object_id,
2053 AttributeId::TEST_ID,
2054 (base_extent_offset + count - 1) * MIN_BLOCK_SIZE
2055 ..(base_extent_offset + count) * MIN_BLOCK_SIZE,
2056 ));
2057 }
2058
2059 items.sort_by(|a, b| a.key.cmp_upper_bound(&b.key));
2061
2062 let mut writer = PersistentLayerWriter::<_, ObjectKey, u64>::new(
2063 Writer::new(&handle).await,
2064 (3 * BLOCK_SIZE) as usize,
2065 BLOCK_SIZE,
2066 )
2067 .await
2068 .expect("writer new");
2069
2070 for item in items {
2071 writer.write(item.as_item_ref()).await.expect("write failed");
2072 }
2073
2074 writer.complete().await.expect("flush failed");
2075 }
2076
2077 let layer = PersistentLayer::<ObjectKey, u64>::open(handle).await.expect("new failed");
2078 for target in to_find {
2079 let iterator = layer.seek(Bound::Included(&target)).await.expect("failed to seek");
2080 let ItemRef { key, .. } = iterator.get().expect("missing item");
2081 assert_eq!(&target, key);
2082 }
2083 }
2084
2085 #[fuchsia::test]
2086 async fn test_two_seek_blocks() {
2087 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
2089 const ITEMS_PER_PHASE: u64 = 50;
2090 const ITEM_COUNT: u64 = ITEMS_PER_PHASE * ((BLOCK_SIZE.size() / 8) + 2);
2091
2092 let mut to_find = Vec::new();
2093
2094 let handle = FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
2095 {
2096 let mut writer = PersistentLayerWriter::<_, TestKey, u64>::new(
2097 Writer::new(&handle).await,
2098 ITEM_COUNT as usize * 18,
2099 BLOCK_SIZE,
2100 )
2101 .await
2102 .expect("writer new");
2103
2104 let initial_value = u32::MAX as u64 + 1;
2106 for i in 0..ITEM_COUNT {
2107 writer
2108 .write(
2109 Item::new(TestKey(initial_value + i..initial_value + i), initial_value)
2110 .as_item_ref(),
2111 )
2112 .await
2113 .expect("write failed");
2114 }
2115 to_find.push(TestKey(initial_value..initial_value));
2117 let middle = initial_value + ITEM_COUNT / 2;
2118 to_find.push(TestKey(middle..middle));
2119 let end = initial_value + ITEM_COUNT - 1;
2120 to_find.push(TestKey(end..end));
2121
2122 writer.complete().await.expect("flush failed");
2123 }
2124
2125 let layer = PersistentLayer::<TestKey, u64>::open(handle).await.expect("new failed");
2126 for target in to_find {
2127 let iterator = layer.seek(Bound::Included(&target)).await.expect("failed to seek");
2128 let ItemRef { key, .. } = iterator.get().expect("missing item");
2129 assert_eq!(&target, key);
2130 }
2131 }
2132
2133 #[fuchsia::test]
2136 async fn test_full_seek_block() {
2137 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
2138 const ITEMS_PER_PHASE: u64 = 50;
2139
2140 const SEEK_TABLE_ENTRIES: u64 = BLOCK_SIZE.size() / 8;
2142
2143 const START_ENTRIES_COUNT: u64 = ITEMS_PER_PHASE * SEEK_TABLE_ENTRIES;
2147
2148 for entries in START_ENTRIES_COUNT..START_ENTRIES_COUNT + (ITEMS_PER_PHASE * 2) {
2149 let handle =
2150 FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
2151 {
2152 let mut writer = PersistentLayerWriter::<_, TestKey, u64>::new(
2153 Writer::new(&handle).await,
2154 entries as usize,
2155 BLOCK_SIZE,
2156 )
2157 .await
2158 .expect("writer new");
2159
2160 let initial_value = u32::MAX as u64 + 1;
2162 for i in 0..entries {
2163 writer
2164 .write(
2165 Item::new(TestKey(initial_value + i..initial_value + i), initial_value)
2166 .as_item_ref(),
2167 )
2168 .await
2169 .expect("write failed");
2170 }
2171
2172 writer.complete().await.expect("flush failed");
2173 }
2174 PersistentLayer::<TestKey, u64>::open(handle).await.expect("new failed");
2175 }
2176 }
2177
2178 #[fuchsia::test]
2179 async fn test_ignore_bloom_filter_on_older_versions() {
2180 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
2181 const ITEMS_PER_PHASE: u64 = 50;
2182 const ITEM_COUNT: u64 = (1 + MINIMUM_DATA_BLOCKS_FOR_BLOOM_FILTER as u64) * ITEMS_PER_PHASE;
2184
2185 let old_version_handle =
2186 FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
2187 let current_version_handle =
2188 FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
2189 let initial_value = u32::MAX as u64 + 1;
2191 {
2192 let mut old_version_writer =
2193 PersistentLayerWriter::<_, TestKey, u64>::new_with_version(
2194 Writer::new(&old_version_handle).await,
2195 ITEM_COUNT as usize,
2196 BLOCK_SIZE,
2197 OLD_KEY_SERIALIZATION_VERSION,
2198 )
2199 .await
2200 .expect("writer new");
2201 let mut current_version_writer = PersistentLayerWriter::<_, TestKey, u64>::new(
2202 Writer::new(¤t_version_handle).await,
2203 ITEM_COUNT as usize,
2204 BLOCK_SIZE,
2205 )
2206 .await
2207 .expect("writer new");
2208
2209 for i in 0..ITEM_COUNT {
2210 old_version_writer
2211 .write(
2212 Item::new(TestKey(initial_value + i..initial_value + i), initial_value)
2213 .as_item_ref(),
2214 )
2215 .await
2216 .expect("write failed");
2217 current_version_writer
2218 .write(
2219 Item::new(TestKey(initial_value + i..initial_value + i), initial_value)
2220 .as_item_ref(),
2221 )
2222 .await
2223 .expect("write failed");
2224 }
2225
2226 old_version_writer.complete().await.expect("flush failed");
2227 current_version_writer.complete().await.expect("flush failed");
2228 }
2229
2230 let old_layer =
2231 PersistentLayer::<TestKey, u64>::open(old_version_handle).await.expect("open failed");
2232 let current_layer = PersistentLayer::<TestKey, u64>::open(current_version_handle)
2233 .await
2234 .expect("open failed");
2235 assert!(!old_layer.has_bloom_filter());
2236 assert!(current_layer.has_bloom_filter());
2237
2238 let iter = old_layer.seek(Bound::Included(&TestKey(0..0))).await.expect("seek failed");
2240 let item = iter.get().expect("missing item");
2241 assert_eq!(item.key.0.start, initial_value);
2242
2243 let iter = current_layer.seek(Bound::Included(&TestKey(0..0))).await.expect("seek failed");
2244 let item = iter.get().expect("missing item");
2245 assert_eq!(item.key.0.start, initial_value);
2246
2247 let iter = old_layer.seek(Bound::Unbounded).await.expect("seek failed");
2248 let item = iter.get().expect("missing item");
2249 assert_eq!(item.key.0.start, initial_value);
2250
2251 let iter = current_layer.seek(Bound::Unbounded).await.expect("seek failed");
2252 let item = iter.get().expect("missing item");
2253 assert_eq!(item.key.0.start, initial_value);
2254
2255 let target_val = initial_value + ITEM_COUNT / 2;
2256 let iter = old_layer
2257 .seek(Bound::Included(&TestKey(target_val..target_val)))
2258 .await
2259 .expect("seek failed");
2260 let item = iter.get().expect("missing item");
2261 assert_eq!(item.key.0.start, target_val);
2262
2263 let iter = current_layer
2264 .seek(Bound::Included(&TestKey(target_val..target_val)))
2265 .await
2266 .expect("seek failed");
2267 let item = iter.get().expect("missing item");
2268 assert_eq!(item.key.0.start, target_val);
2269 }
2270
2271 #[fuchsia::test]
2272 async fn test_allocator_key_older_version_seek() {
2273 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
2274 const ITEM_COUNT: u64 = 100;
2275
2276 let old_version_handle =
2277 FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
2278 let current_version_handle =
2279 FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
2280 let step = MIN_BLOCK_SIZE.get();
2281 {
2282 let mut old_version_writer =
2283 PersistentLayerWriter::<_, AllocatorKey, i64>::new_with_version(
2284 Writer::new(&old_version_handle).await,
2285 ITEM_COUNT as usize,
2286 BLOCK_SIZE,
2287 OLD_KEY_SERIALIZATION_VERSION,
2288 )
2289 .await
2290 .expect("writer new");
2291 let mut current_version_writer = PersistentLayerWriter::<_, AllocatorKey, i64>::new(
2292 Writer::new(¤t_version_handle).await,
2293 ITEM_COUNT as usize,
2294 BLOCK_SIZE,
2295 )
2296 .await
2297 .expect("writer new");
2298
2299 for i in 0..ITEM_COUNT {
2300 let key = AllocatorKey { device_range: Extent(i * step..(i + 1) * step) };
2301 old_version_writer
2302 .write(Item::new(key.clone(), i as i64).as_item_ref())
2303 .await
2304 .expect("write failed");
2305 current_version_writer
2306 .write(Item::new(key, i as i64).as_item_ref())
2307 .await
2308 .expect("write failed");
2309 }
2310
2311 old_version_writer.complete().await.expect("flush failed");
2312 current_version_writer.complete().await.expect("flush failed");
2313 }
2314
2315 let old_layer = PersistentLayer::<AllocatorKey, i64>::open(old_version_handle)
2316 .await
2317 .expect("open failed");
2318 let current_layer = PersistentLayer::<AllocatorKey, i64>::open(current_version_handle)
2319 .await
2320 .expect("open failed");
2321
2322 let target_idx = ITEM_COUNT / 2;
2323 let target_key =
2324 AllocatorKey { device_range: Extent(target_idx * step..(target_idx + 1) * step) };
2325
2326 let iter = old_layer.seek(Bound::Included(&target_key)).await.expect("seek failed");
2327 let item = iter.get().expect("missing item");
2328 assert_eq!(item.key, &target_key);
2329
2330 let iter = current_layer.seek(Bound::Included(&target_key)).await.expect("seek failed");
2331 let item = iter.get().expect("missing item");
2332 assert_eq!(item.key, &target_key);
2333 }
2334
2335 #[fuchsia::test]
2336 async fn test_key_exists_no_bloom_filter() {
2337 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_8KIB;
2338 const ITEM_COUNT: i32 = 100;
2340
2341 let handle = FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
2342 {
2343 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
2344 Writer::new(&handle).await,
2345 ITEM_COUNT as usize,
2346 BLOCK_SIZE,
2347 )
2348 .await
2349 .expect("writer new");
2350 for i in 1..ITEM_COUNT {
2351 writer.write(Item::new(i * 2, i * 2).as_item_ref()).await.expect("write failed");
2352 }
2353 writer.complete().await.expect("flush failed");
2354 }
2355 let layer = PersistentLayer::<i32, i32>::open(handle).await.expect("new failed");
2356 assert!(!layer.has_bloom_filter());
2357
2358 assert_eq!(layer.key_exists(&0).await.expect("key_exists failed"), Existence::Missing);
2359 assert_eq!(layer.key_exists(&1).await.expect("key_exists failed"), Existence::Missing);
2360 for i in 1..ITEM_COUNT {
2361 assert_eq!(
2362 layer.key_exists(&(i * 2)).await.expect("key_exists failed"),
2363 Existence::Exists
2364 );
2365 assert_eq!(
2366 layer.key_exists(&(i * 2 + 1)).await.expect("key_exists failed"),
2367 Existence::Missing
2368 );
2369 }
2370 }
2371
2372 #[fuchsia::test]
2373 async fn test_key_exists_with_bloom_filter() {
2374 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
2375 const ITEM_COUNT: i32 = 10000;
2377
2378 let handle = FakeObjectHandle::new(Arc::new(FakeObject::new()));
2379 {
2380 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
2381 Writer::new(&handle).await,
2382 ITEM_COUNT as usize,
2383 BLOCK_SIZE,
2384 )
2385 .await
2386 .expect("writer new");
2387 for i in 0..ITEM_COUNT {
2388 writer.write(Item::new(i * 2, i * 2).as_item_ref()).await.expect("write failed");
2389 }
2390 writer.complete().await.expect("flush failed");
2391 }
2392 let layer = PersistentLayer::<i32, i32>::open(handle).await.expect("new failed");
2393 assert!(layer.has_bloom_filter());
2394
2395 for i in 0..ITEM_COUNT {
2396 assert_eq!(
2398 layer.key_exists(&(i * 2)).await.expect("key_exists failed"),
2399 Existence::MaybeExists
2400 );
2401 }
2402
2403 let mut missing_count = 0;
2406 for i in 0..ITEM_COUNT {
2407 let result = layer.key_exists(&(i * 2 + 1)).await.expect("key_exists failed");
2408 assert_ne!(result, Existence::Exists);
2409 if result == Existence::Missing {
2410 missing_count += 1;
2411 }
2412 }
2413 assert!(missing_count > ITEM_COUNT / 2);
2415 }
2416
2417 #[fuchsia::test]
2418 async fn test_load_large_bloom_filter_multi_chunk() {
2419 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
2420 const ESTIMATED_ITEMS: usize = 600_000;
2422 const WRITTEN_ITEMS: i32 = 2000;
2423
2424 let handle = FakeObjectHandle::new(Arc::new(FakeObject::new()));
2425 {
2426 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
2427 Writer::new(&handle).await,
2428 ESTIMATED_ITEMS,
2429 BLOCK_SIZE,
2430 )
2431 .await
2432 .expect("writer new");
2433 for i in 0..WRITTEN_ITEMS {
2434 writer.write(Item::new(i * 2, i * 2).as_item_ref()).await.expect("write failed");
2435 }
2436 writer.complete().await.expect("flush failed");
2437 }
2438 let layer = PersistentLayer::<i32, i32>::open(handle).await.expect("open failed");
2439 assert!(layer.has_bloom_filter());
2440
2441 for i in 0..WRITTEN_ITEMS {
2442 assert_eq!(layer.maybe_contains_key(&(i * 2)), MaybeContainsKey::Maybe);
2443 }
2444 let mut false_count = 0;
2445 for i in 0..WRITTEN_ITEMS {
2446 if layer.maybe_contains_key(&(i * 2 + 1)) == MaybeContainsKey::False {
2447 false_count += 1;
2448 }
2449 }
2450 assert!(false_count > WRITTEN_ITEMS / 2);
2451 }
2452
2453 #[fuchsia::test]
2454 async fn test_clear_cached_data() {
2455 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
2456 let handle = FakeObjectHandle::new(Arc::new(FakeObject::new()));
2457 {
2458 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
2459 Writer::new(&handle).await,
2460 100,
2461 BLOCK_SIZE,
2462 )
2463 .await
2464 .expect("writer new");
2465 writer.write(Item::new(1, 1).as_item_ref()).await.expect("write failed");
2466 writer.complete().await.expect("flush failed");
2467 }
2468 let layer =
2469 PersistentLayer::<i32, i32>::open_async(Arc::new(handle)).await.expect("open failed");
2470 let iter = layer.seek(Bound::Unbounded).await.expect("seek failed");
2471 assert_eq!(iter.get().map(|i| (*i.key, *i.value)), Some((1, 1)));
2472 drop(iter);
2473
2474 assert!(layer.object_handle.try_read(BLOCK_SIZE.get() as usize).is_some());
2475
2476 layer.clear_cached_data();
2477
2478 assert!(layer.object_handle.try_read(BLOCK_SIZE.get() as usize).is_none());
2479 }
2480
2481 struct SliceLayerObject {
2482 handle: FakeObjectHandle,
2483 slice: Vec<u8>,
2484 io_error: AtomicBool,
2485 purged: AtomicBool,
2486 closed: AtomicBool,
2487 }
2488
2489 impl ObjectHandle for SliceLayerObject {
2490 fn object_id(&self) -> u64 {
2491 self.handle.object_id()
2492 }
2493 fn block_size(&self) -> BlockSize {
2494 self.handle.block_size()
2495 }
2496 fn allocate_buffer(&self, size: usize) -> BufferFuture<'_> {
2497 self.handle.allocate_buffer(size)
2498 }
2499 }
2500
2501 #[async_trait]
2502 impl ReadObjectHandle for SliceLayerObject {
2503 async fn read_aligned(
2504 &self,
2505 offset: u64,
2506 buf: MutableBufferRef<'_>,
2507 ) -> Result<usize, Error> {
2508 self.handle.read_aligned(offset, buf).await
2509 }
2510 fn get_size(&self) -> u64 {
2511 self.handle.get_size()
2512 }
2513 }
2514
2515 #[async_trait]
2516 impl LayerObject for SliceLayerObject {
2517 fn as_slice(&self) -> Option<&[u8]> {
2518 Some(&self.slice)
2519 }
2520 fn has_io_error(&self) -> bool {
2521 self.io_error.load(Ordering::SeqCst)
2522 }
2523 fn purge_cached_data(&self) {
2524 self.purged.store(true, Ordering::SeqCst);
2525 }
2526 async fn close(&self) {
2527 self.closed.store(true, Ordering::SeqCst);
2528 }
2529 }
2530
2531 #[fuchsia::test]
2532 async fn test_sync_persistent_layer() {
2533 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_512B;
2534 const ITEM_COUNT: i32 = 1000;
2535
2536 let handle = FakeObjectHandle::new(Arc::new(FakeObject::new()));
2537 {
2538 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
2539 Writer::new(&handle).await,
2540 ITEM_COUNT as usize,
2541 BLOCK_SIZE,
2542 )
2543 .await
2544 .expect("writer new");
2545 for i in 0..ITEM_COUNT {
2546 writer.write(Item::new(i * 2, i * 2).as_item_ref()).await.expect("write failed");
2547 }
2548 writer.complete().await.expect("flush failed");
2549 }
2550
2551 let size = handle.get_size() as usize;
2552 let mut buf = handle.allocate_buffer(size).await;
2553 handle.read_aligned(0, buf.as_mut()).await.expect("read failed");
2554 let slice = buf.subslice(..).to_vec();
2555 drop(buf);
2556 let slice_obj = Arc::new(SliceLayerObject {
2557 handle,
2558 slice: slice.clone(),
2559 io_error: AtomicBool::new(false),
2560 purged: AtomicBool::new(false),
2561 closed: AtomicBool::new(false),
2562 });
2563
2564 let layer =
2565 PersistentLayer::<i32, i32>::open_layer(slice_obj.clone()).await.expect("open failed");
2566 assert!(layer.has_bloom_filter());
2567 assert_eq!(layer.len(), ITEM_COUNT as usize);
2568
2569 let mut iter = layer.seek(Bound::Unbounded).await.expect("seek failed");
2571 for i in 0..ITEM_COUNT {
2572 let item = iter.get().expect("expected item");
2573 assert_eq!(*item.key, i * 2);
2574 assert_eq!(*item.value, i * 2);
2575 let fut = iter.advance_dyn().expect("advance_dyn failed");
2576 assert!(fut.is_none(), "SyncPersistentLayer iterator must advance synchronously");
2577 }
2578 assert!(iter.get().is_none());
2579
2580 let iter = layer.seek(Bound::Included(&500)).await.expect("seek included failed");
2582 assert_eq!(*iter.get().expect("expected item").key, 500);
2583
2584 let iter = layer.seek(Bound::Excluded(&500)).await.expect("seek excluded failed");
2585 assert_eq!(*iter.get().expect("expected item").key, 502);
2586
2587 let iter = layer.seek(Bound::Included(&501)).await.expect("seek non-existent failed");
2588 assert_eq!(*iter.get().expect("expected item").key, 502);
2589
2590 assert!(!slice_obj.purged.load(Ordering::SeqCst));
2592 layer.purge_cached_data();
2593 assert!(slice_obj.purged.load(Ordering::SeqCst));
2594
2595 assert!(!slice_obj.closed.load(Ordering::SeqCst));
2597 layer.close().await;
2598 assert!(slice_obj.closed.load(Ordering::SeqCst));
2599
2600 let mut zeroed_slice = slice;
2602 let bs = BLOCK_SIZE.get() as usize;
2603 zeroed_slice[bs..bs * 2].fill(0);
2604 let zeroed_obj = Arc::new(SliceLayerObject {
2605 handle: FakeObjectHandle::new(Arc::new(FakeObject::new())),
2606 slice: zeroed_slice,
2607 io_error: AtomicBool::new(false),
2608 purged: AtomicBool::new(false),
2609 closed: AtomicBool::new(false),
2610 });
2611 let mut wbuf = zeroed_obj.handle.allocate_buffer(zeroed_obj.slice.len()).await;
2614 wbuf.as_mut_ptr_slice().copy_from_slice(&zeroed_obj.slice);
2615 zeroed_obj.handle.write_or_append(Some(0), wbuf.as_ref()).await.unwrap();
2616 drop(wbuf);
2617 let zeroed_layer =
2618 PersistentLayer::<i32, i32>::open_layer(zeroed_obj.clone()).await.expect("open failed");
2619
2620 let err = zeroed_layer.seek(Bound::Unbounded).await.err().expect("seek should fail");
2623 assert!(FxfsError::Inconsistent.matches(&err), "expected Inconsistent, got {err:?}");
2624
2625 zeroed_obj.io_error.store(true, Ordering::SeqCst);
2628 let err = zeroed_layer.seek(Bound::Unbounded).await.err().expect("seek should fail");
2629 assert_eq!(
2630 err.root_cause().downcast_ref::<zx_status::Status>(),
2631 Some(&zx_status::Status::IO)
2632 );
2633 }
2634
2635 #[fuchsia::test]
2636 async fn test_layer_block_size_gt_4k_falls_back_to_async_persistent_layer() {
2637 const BLOCK_SIZE: BlockSize = BlockSize::SIZE_8KIB;
2638 let handle = FakeObjectHandle::new_with_block_size(Arc::new(FakeObject::new()), BLOCK_SIZE);
2639 {
2640 let mut writer = PersistentLayerWriter::<_, i32, i32>::new(
2641 Writer::new(&handle).await,
2642 10,
2643 BLOCK_SIZE,
2644 )
2645 .await
2646 .expect("writer new");
2647 for i in 0..10 {
2648 writer.write(Item::new(i, i).as_item_ref()).await.expect("write failed");
2649 }
2650 writer.complete().await.expect("flush failed");
2651 }
2652
2653 let size = handle.get_size() as usize;
2654 let mut buf = handle.allocate_buffer(size).await;
2655 handle.read_aligned(0, buf.as_mut()).await.expect("read failed");
2656 let slice = buf.subslice(..).to_vec();
2657 drop(buf);
2658 let slice_obj = Arc::new(SliceLayerObject {
2659 handle,
2660 slice,
2661 io_error: AtomicBool::new(false),
2662 purged: AtomicBool::new(false),
2663 closed: AtomicBool::new(false),
2664 });
2665
2666 assert!(SyncPersistentLayer::<i32, i32>::open(slice_obj.clone()).await.is_err());
2668
2669 let layer = PersistentLayer::<i32, i32>::open_layer(slice_obj).await.expect("open_layer");
2671 let mut iter = layer.seek(Bound::Unbounded).await.expect("seek");
2672 for i in 0..10 {
2673 assert_eq!(*iter.get().expect("item").key, i);
2674 iter.advance().await.expect("advance");
2675 }
2676 assert!(iter.get().is_none());
2677 }
2678}