Skip to main content

fxfs/lsm_tree/
persistent_layer.rs

1// Copyright 2024 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5// PersistentLayer object format
6//
7// The layer is made up of 1 or more "blocks" whose size are some multiple of the block size used
8// by the underlying handle.
9//
10// The persistent layer has 4 types of blocks:
11//  - Header block
12//  - Data block
13//  - BloomFilter block
14//  - Seek block (+LayerInfo)
15//
16// The structure of the file is as follows:
17//
18// blk#     contents
19// 0        [Header]
20// 1        [Data]
21// 2        [Data]
22// ...      [Data]
23// L        [BloomFilter]
24// L + 1    [BloomFilter]
25// ...      [BloomFilter]
26// M        [Seek]
27// M + 1    [Seek]
28// ...      [Seek]
29// N        [Seek/LayerInfo]
30//
31// Generally, there will be an order of magnitude more Data blocks than Seek/BloomFilter blocks.
32//
33// Header contains a Version-prefixed LayerHeader struct.  This version is used for everything in
34// the layer file.
35//
36// Data blocks contain a little endian encoded u16 item count at the start, then a series of
37// serialized items, and a list of little endian u16 offsets within the block for where
38// serialized items start, excluding the first item (since it is at a known offset). The list of
39// offsets ends at the end of the block and since the items are of variable length, there may be
40// space between the two sections if the next item and its offset cannot fit into the block.
41//
42// |item_count|item|item|item|item|item|item|dead space|offset|offset|offset|offset|offset|
43//
44// BloomFilter blocks contain a bitmap which is used to probabilistically determine if a given key
45// might exist in the layer file.   See `BloomFilter` for details on this structure.  Note that this
46// can be absent from the file for small layer files.
47//
48// Seek/LayerInfo blocks contain both the seek table, and a single LayerInfo struct at the tail of
49// the last block, with the LayerInfo's length written as a little-endian u64 at the very end.  The
50// padding between the two structs is ignored but nominally is zeroed. They share blocks to avoid
51// wasting padding bytes.  Note that the seek table can be absent from the file for small layer
52// files (but there will always be one block for the LayerInfo).
53//
54// The seek table contains a little-endian u64 for each data block, recording the leading u64 of
55// that block's first key. Entries are in sorted order (duplicates are possible when an object
56// spans multiple blocks).
57
58use 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
94/// LayerHeader is stored in the first block of the persistent layer.
95pub type LayerHeader = LayerHeaderV39;
96
97#[derive(Debug, Serialize, Deserialize, TypeFingerprint, Versioned)]
98pub struct LayerHeaderV39 {
99    /// 'FxfsLayr'
100    magic: [u8; 8],
101    /// The block size used within this layer file. This is typically set at compaction time to the
102    /// same block size as the underlying object handle.
103    ///
104    /// (Each block starts with a 2 byte item count so there is a 64k item limit per block,
105    /// regardless of block size).
106    block_size: u64,
107}
108
109/// The last block of each layer contains metadata for the rest of the layer.
110pub type LayerInfo = LayerInfoV39;
111
112#[derive(Debug, Serialize, Deserialize, TypeFingerprint, Versioned)]
113pub struct LayerInfoV39 {
114    /// How many items are in the layer file.  Mainly used for sizing bloom filters during
115    /// compaction.
116    num_items: usize,
117    /// The number of data blocks in the layer file.
118    num_data_blocks: u64,
119    /// The size of the bloom filter in the layer file.  Not necessarily block-aligned.
120    bloom_filter_size_bytes: usize,
121    /// The seed for the nonces used in the bloom filter.
122    bloom_filter_seed: u64,
123    /// How many nonces to use for bloom filter hashing.
124    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
145/// A handle to an asynchronous, chunk-cached persistent layer.
146pub struct PersistentLayer<K, V> {
147    object_handle: CachingObjectHandle<Arc<dyn LayerObject>>,
148    data: LayerData<K>,
149    _value_type: PhantomData<V>,
150}
151
152/// A handle to a synchronous, slice-backed persistent layer.
153pub 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    // The currently loaded chunk and its index.
191    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
235// For small layer files, don't bother with the bloom filter.  Arbitrarily chosen.
236const MINIMUM_DATA_BLOCKS_FOR_BLOOM_FILTER: usize = 4;
237
238// How many blocks we reserve for the header.  Data blocks start at this offset.
239const NUM_HEADER_BLOCKS: u64 = 1;
240
241/// The smallest possible (empty) layer file is always 2 blocks, one for the header and one for
242/// LayerInfo.
243const MINIMUM_LAYER_FILE_BLOCKS: u64 = 2;
244
245// Put safety rails on the size of the bloom filter and seek table to avoid OOMing the system.
246// It's more likely that tampering has occurred in these cases.
247const MAX_BLOOM_FILTER_SIZE: usize = 64 * 1024 * 1024;
248const MAX_SEEK_TABLE_SIZE: usize = 64 * 1024 * 1024;
249
250// The following constants refer to sizes of metadata in the data blocks.
251const 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
260// A key-only iterator, used while seeking through the tree.
261struct KeyOnlyIterator<'iter, K: Key, V: LayerValue, B> {
262    // Allocated out of |layer|.
263    buffer: BufferCursor<B>,
264
265    layer: &'iter LayerData<K>,
266
267    // The position of the _next_ block to be read.
268    pos: u64,
269
270    // The item index in the current block.
271    item_index: u16,
272
273    // The number of items in the current block.
274    item_count: u16,
275
276    // The current key state.
277    key: KeyState<K>,
278
279    // The base value used for decoding keys in the current block.
280    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    // Repositions the iterator to point to the `index`'th item in the current block.
295    // Returns an error if the index is out of range or the resulting offset contains an obviously
296    // invalid value.
297    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            // Fast-path when we are seeking in a linear manner, as is the case when advancing a
301            // wrapping iterator that also deserializes the values.
302            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            // First entry isn't actually recorded, it is at the start of the block after the item
308            // count.
309            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    // Binary searches the current block for `search_key`.  The iterator must be positioned at the
422    // first item in the block, and that item must be less than `search_key`.  If the block
423    // contains a key that is greater than or equal to `search_key`, the iterator is positioned at
424    // the first such key and `Some(true)` is returned for an exact match, or `Some(false)`
425    // otherwise.  If every key in the block is less than `search_key`, `None` is returned and the
426    // caller should advance the iterator to the next block.
427    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    // Like `binary_search_block`, but compares the serialized `target` key directly against the
461    // keys in the block without repositioning the iterator for each probe.
462    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            // `read_block_header` has checked that the seek entries don't overlap the header.
477            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    // Returns an iterator which reads the block at `pos` (the file offset of a key block) out of
521    // `buffer`, starting at `buffer_pos`.
522    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    // Returns a new iterator for the block at `pos`, reusing the chunk held by `self` or `other`
548    // if either contains `pos`, which saves a cache lookup.
549    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    // Advances without blocking.  Returns false, leaving the iterator unchanged, if the next block
583    // needs to be read from the underlying handle.
584    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    // Returns a new iterator for the block at `pos`.  `_other` exists for parity with the
617    // chunk-backed version; the whole layer is always available here.
618    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    // The current item.
638    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        // Ignore the seek table on older versions because AllocatorKey::get_leading_u64()
699        // changed from bytes to blocks.
700        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        // Should be in strict ascending order, otherwise something's broken, or we've gone off
722        // the end and we're reading zeroes.
723        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    // Returns the key serialized relative to `base`, or None if the key sorts before any key that
781    // can be encoded relative to `base`.
782    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            // Ignore the bloom filter for layer files in outdated versions.  We don't know whether
886            // keys have changed formats or not (and therefore have different hash values), so we
887            // must ignore the bloom filter and always query the layer.
888            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            // We are searching for a range here, as multiple items can have the same value in
929            // this approximate search. Since the values in the seek table represent the first
930            // key in each block (ordered by cmp_upper_bound), if the value equals the target
931            // we must also search the block before it. The goal is for
932            // table[left] < target < table[right].
933            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            // Since partition_point will find the index of the first place where the predicate
942            // is false, we subtract 1 to get the index where it was last true.
943            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            // Pick a block midway.
966            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        // Finish the binary search on the block pointed to by `left`.
989        match left.binary_search_block(&mut search_key)? {
990            Some(true) if excluded => left.advance()$(.$await)??,
991            Some(_) => {}
992            // When we don't find an entry that is greater than or equal to the target in `left`,
993            // we need to return with the first entry of the next block, which might already be
994            // pointed to by `right`.  Otherwise, advance to the next block (or the end of the
995            // layer).
996            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    /// Opens a persistent layer backed by `handle`. If the layer is encrypted and the platform
1048    /// uses a pager, `unwrapped_key` is registered with the pager and dropped immediately.
1049    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
1099/// Unwraps the encryption key for `object_id` if a `LayerPager` is active on `store`.
1100async 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
1124/// Opens a persistent layer from a newly created or existing `DataObjectHandle`.
1125///
1126/// If the layer is encrypted and the platform uses a pager, `unwrapped_key` is registered with the
1127/// pager and dropped immediately.
1128pub 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
1135/// Opens persistent layers for `object_ids` from `store`.
1136///
1137/// Returns `(layers, total_size)` where `layers` is the vector of opened `Layer` trait objects and
1138/// `total_size` is the sum of the sizes in bytes of all opened layer objects.
1139pub 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
1320// This ensures that item_count can't be overflowed below.
1321const_assert!(MAX_BLOCK_SIZE.size() <= u16::MAX as u64 + 1);
1322
1323// -- Writer support --
1324
1325pub 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    /// Creates a new writer that will serialize items to the object accessible via |object_handle|
1340    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        // Write the header block.
1354        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    /// Writes `self.buf` out as a block.
1380    ///
1381    /// Blocks are fixed size, consisting of a 16-bit item count, data, zero padding
1382    /// and seek table at the end.
1383    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        // Write the seek table. Entries are 2 bytes each and items are always at least 10.
1397        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    // Assumes the writer is positioned to a new block.
1409    // Returns the size, in bytes, of the seek table.
1410    // Note that the writer will be positioned to exactly the end of the seek table, not to the end
1411    // of a block.
1412    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    // Assumes the writer is positioned to exactly the end of the seek table, which was
1433    // `seek_table_len` bytes.
1434    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        // We want the LayerInfo to be at the end of the last block.  That might require creating a
1455        // new block if we don't have enough room.
1456        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    // Assumes the writer is positioned to a new block.
1471    // Returns the size of the bloom filter, in bytes.
1472    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        // TODO(https://fxbug.dev/323571978): Avoid bounce-buffering.
1477        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    // Returns the bloom filter writer. Intended to be used for testing purposes, e.g., gain access
1487    // to the bloom filter to then corrupt it.
1488    #[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        // Note the length before we write this item.
1513        let len = self.buf.len();
1514        // Each data block's keys are delta-encoded relative to the leading u64 of that block's
1515        // first key (stored in `block_keys`). Record the base key for the first block here before
1516        // serializing; for subsequent blocks, `block_keys` is updated below when a block overflows.
1517        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        // Never record the first item. The offset is always the same.
1524        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 writing the item took us over a block, flush the bytes in the buffer prior to this
1530        // item.
1531        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                // Drop the recently added offset from the list. The latest item will be the first
1538                // on the next block and have a known offset there.
1539                self.block_offsets.pop();
1540            }
1541            self.buf.truncate(len);
1542            self.write_block().await?;
1543
1544            // Start a new block with `item` as its first entry and re-serialize using the new base.
1545            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/// Logs a warning if this object is dropped and the contained value isn't 0.
1566#[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                // Populate every other value as an item.
1676                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        // Search for all values to check the in-between values.
1682        for i in 0..ITEM_COUNT * 2 {
1683            // We've written every other value, we expect to get either the exact value searched
1684            // for, or the next one after it. So round up to the nearest multiple of 2.
1685            let expected = round_up(i, 2).unwrap();
1686            let mut iterator = layer.seek(Bound::Included(&i)).await.expect("failed to seek");
1687            // We've written values up to (N-1)*2=2*N-2, so when looking for 2*N-1 we'll go off the
1688            // end of the layer and get back no item.
1689            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            // Check that we can advance to the next item.
1697            iterator.advance().await.expect("failed to advance");
1698            // The highest value is 2*N-2, searching for 2*N-3 will find the last value, and
1699            // advancing will go off the end of the layer and return no item. If there was
1700            // previously no item, then it will latch and always return no item.
1701            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        // Check that we can advance to the next item.
1736        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        // At the upper end of the supported size.
1825        const BLOCK_SIZE: BlockSize = MAX_BLOCK_SIZE;
1826        // Items will be 18 bytes, so fill up a few pages.
1827        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            // Use large values to force varint encoding to use consistent space.
1839            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        // At the upper end of the supported size.
1858        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                // Check that we can advance to the next item.
1896                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    /// Generates extent records for a given object_id (of size 1).
1911    /// This produces a series of records with the same leading_u64.
1912    /// Returns the generated items and the next available object_id.
1913    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    /// Generates objects object_ids over a range.
1933    /// This produced a series of records with unique leading_u64.
1934    /// Returns the generated items and the next available value for sequencing.
1935    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    // Create a large spread of data across several blocks to ensure that no part of the range is
1945    // lost by the partial search using the layer seek table.
1946    #[fuchsia::test]
1947    async fn test_block_seek_duplicate_leading_u64() {
1948        // At the upper end of the supported size.
1949        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            // Make all values take up maximum space for varint encoding.
1958            let mut object_id = u32::MAX as u64 + 1;
1959
1960            // First fill the front with duplicate object IDs, then look at the start,
1961            // middle and end of the range.
1962            {
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            // Add some filler of all different leading u64.
1989            {
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            // Fill the middle with duplicate object IDs, then look at the start,
1997            // middle and end of the range.
1998            {
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            // Add some filler of all different leading u64.
2025            {
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            // Fill the end with duplicate object IDs, then look at the start,
2033            // middle and end of the range.
2034            {
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            // Sort items by cmp_upper_bound!
2060            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        // At the upper end of the supported size.
2088        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            // Make all values take up maximum space for varint encoding.
2105            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            // Look at the start middle and end.
2116            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    // Verifies behaviour around creating full seek blocks, to ensure that it is able to be opened
2134    // and parsed afterward.
2135    #[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        // How many entries there are in a seek table block.
2141        const SEEK_TABLE_ENTRIES: u64 = BLOCK_SIZE.size() / 8;
2142
2143        // Number of entries to fill a seek block would need one more block of entries, but we're
2144        // starting low here on purpose to do a range and make sure we hit the size we are
2145        // interested in.
2146        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                // Make all values take up maximum space for varint encoding.
2161                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        // Add enough items to create enough blocks for a bloom filter to be necessary.
2183        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        // Make all values take up maximum space for varint encoding.
2190        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(&current_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        // Verify seeking works in both layers (including before the first block's key).
2239        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(&current_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        // Not enough items to trigger a bloom filter.
2339        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        // Enough items to trigger a bloom filter.
2376        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            // With a bloom filter, we expect MaybeExists for present keys.
2397            assert_eq!(
2398                layer.key_exists(&(i * 2)).await.expect("key_exists failed"),
2399                Existence::MaybeExists
2400            );
2401        }
2402
2403        // For missing keys, we expect Missing, but might get MaybeExists due to false positives.
2404        // We can at least assert it's NOT Exists.
2405        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        // We expect mostly Missing.
2414        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        // Sizing for 600_000 items creates a 2 MiB bloom filter (exceeding 1 MiB chunk size).
2421        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        // Unbounded iteration should complete synchronously (advance_dyn returns Ok(None)).
2570        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        // Seek Included and Excluded.
2581        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        // Verify purge_cached_data delegates to LayerObject::purge_cached_data.
2591        assert!(!slice_obj.purged.load(Ordering::SeqCst));
2592        layer.purge_cached_data();
2593        assert!(slice_obj.purged.load(Ordering::SeqCst));
2594
2595        // Verify close delegates to LayerObject::close.
2596        assert!(!slice_obj.closed.load(Ordering::SeqCst));
2597        layer.close().await;
2598        assert!(slice_obj.closed.load(Ordering::SeqCst));
2599
2600        // Simulate zero-filled data block on disk corruption vs pager I/O error.
2601        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        // Write valid header/footer into the underlying FakeObjectHandle so LayerData::open
2612        // succeeds.
2613        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        // When `has_io_error()` is false, corruption returns `FxfsError::Inconsistent`
2621        // (`ZX_ERR_IO_DATA_INTEGRITY`).
2622        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        // When `has_io_error()` is true (signaled by `SUPPLY_ZEROES_ON_ERROR`), it returns
2626        // `zx_status::Status::IO` (`ZX_ERR_IO`).
2627        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        // SyncPersistentLayer::open must reject layers with block_size > 4KiB.
2667        assert!(SyncPersistentLayer::<i32, i32>::open(slice_obj.clone()).await.is_err());
2668
2669        // PersistentLayer::open_layer must fall back to async PersistentLayer and succeed.
2670        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}