use super::block_fetch::load_segment_index;
use super::cache::{DecodedMetadataTableBlock, MetadataTableCache, MetadataTableCacheKey};
use super::data_block_load::load_segment_data_block_span;
use super::error::ManifestLoadError;
use super::scan::Readahead;
use super::validate::validate_manifest_row_seq_range;
use loonfs_api::wire::manifest::{MetadataFileRef, MetadataRow, MetadataTableFamily};
use loonfs_api::wire::sst_blocks::{index_blocks_for_key_range, DecodedDataBlock};
use loonfs_objectstore::ObjectStore;
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
pub(crate) struct SegmentKeyRangeBlocks {
blocks: Vec<Arc<DecodedDataBlock>>,
}
impl SegmentKeyRangeBlocks {
pub(super) fn rows_in_key_range<'a>(
&'a self,
lower_bound: &'a str,
upper_bound: Option<&'a str>,
) -> impl Iterator<Item = (&'a str, &'a MetadataRow)> + 'a {
self.blocks.iter().flat_map(move |block| {
let start = block
.row_keys
.partition_point(|key| key.as_str() < lower_bound);
let end = upper_bound.map_or(block.row_keys.len(), |upper_bound| {
block
.row_keys
.partition_point(|key| key.as_str() < upper_bound)
});
let range = start..end.max(start);
block.row_keys[range.clone()]
.iter()
.zip(&block.rows[range])
.map(|(key, row)| (key.as_str(), row))
})
}
#[cfg(test)]
pub(super) fn rows(&self) -> impl Iterator<Item = &MetadataRow> {
self.blocks.iter().flat_map(|block| block.rows.iter())
}
#[cfg(test)]
pub(super) fn row_keys(&self) -> impl Iterator<Item = &String> {
self.blocks.iter().flat_map(|block| block.row_keys.iter())
}
}
#[derive(Debug, Default)]
pub(super) struct SessionBlockMemo {
blocks: Mutex<HashMap<MetadataTableCacheKey, DecodedMetadataTableBlock>>,
}
impl SessionBlockMemo {
pub(super) fn get(
&self,
cache_key: &MetadataTableCacheKey,
) -> Option<DecodedMetadataTableBlock> {
self.blocks
.lock()
.expect("session block memo lock should not be poisoned")
.get(cache_key)
.cloned()
}
pub(super) fn record(
&self,
cache_key: &MetadataTableCacheKey,
block: &DecodedMetadataTableBlock,
) {
self.blocks
.lock()
.expect("session block memo lock should not be poisoned")
.insert(cache_key.clone(), block.clone());
}
}
const RANGE_SCAN_READAHEAD_BLOCKS: usize = 32;
#[allow(clippy::too_many_arguments)]
pub(super) async fn load_manifest_segment_rows_in_key_range_with_cache<S: ObjectStore + ?Sized>(
store: &S,
table_cache: Option<&MetadataTableCache>,
memo: &SessionBlockMemo,
family: MetadataTableFamily,
descriptor: &MetadataFileRef,
lower_bound: &str,
upper_bound: Option<&str>,
readahead: Readahead,
) -> Result<SegmentKeyRangeBlocks, ManifestLoadError> {
let index = load_segment_index(store, table_cache, memo, descriptor).await?;
let needed = index_blocks_for_key_range(&index, lower_bound, upper_bound);
let extended_end = if readahead == Readahead::Enabled {
needed
.start
.saturating_add(RANGE_SCAN_READAHEAD_BLOCKS)
.max(needed.end)
.min(index.len())
} else {
needed.end
};
let blocks: Vec<_> = load_segment_data_block_span(
store,
table_cache,
memo,
family,
descriptor,
&index[needed.start..extended_end],
)
.await?
.into_iter()
.take(needed.len())
.collect();
validate_manifest_row_seq_range(
&descriptor.object_key,
blocks.iter().flat_map(|block| block.rows.iter()),
descriptor.run_seq,
)?;
Ok(SegmentKeyRangeBlocks { blocks })
}