use super::cache::MetadataTableCache;
use super::error::ManifestLoadError;
use super::load::{
load_manifest_segment_rows_in_key_range_with_cache, load_segment_filter, SegmentKeyRangeBlocks,
SessionBlockMemo,
};
use super::runs::{MetadataRunManifest, MetadataTableManifest, CHECKPOINT_TABLE_FAMILIES};
#[cfg(test)]
use crate::metadata::MetadataState;
use futures::future::try_join_all;
use loonfs_api::wire::manifest::{
MetadataFileRef, MetadataRow, MetadataTableFamily, NamespaceManifestEnvelope,
};
use loonfs_api::wire::sst_blocks::string_prefix_upper_bound;
use loonfs_objectstore::ObjectStore;
#[cfg(test)]
use serde::{Deserialize, Serialize};
use std::sync::Arc;
pub(super) const MAX_MATERIALIZED_TABLE_LOADS: usize = 16;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Readahead {
Enabled,
Disabled,
}
#[cfg(test)]
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct ManifestMaterializationForInspection {
pub(crate) manifest: NamespaceManifestEnvelope,
pub(crate) metadata_state: MetadataState,
}
pub(crate) struct VerifiedMetadataTables<'a, S: ObjectStore + ?Sized> {
pub(super) store: &'a S,
pub(super) table_cache: Option<&'a MetadataTableCache>,
pub(super) manifest_object_key: String,
pub(super) manifest: Arc<NamespaceManifestEnvelope>,
pub(super) scan_runs: Arc<Vec<MetadataRunManifest>>,
pub(super) block_memo: SessionBlockMemo,
}
impl<'a, S: ObjectStore + ?Sized> VerifiedMetadataTables<'a, S> {
pub(crate) fn synthesized(store: &'a S, manifest: NamespaceManifestEnvelope) -> Self {
debug_assert!(
manifest.payload.metadata_files.is_empty(),
"a synthesized manifest must name no durable metadata files"
);
let scan_runs = Arc::new(Vec::new());
Self {
store,
table_cache: None,
manifest_object_key: String::new(),
manifest: Arc::new(manifest),
scan_runs,
block_memo: SessionBlockMemo::default(),
}
}
}
impl<S: ObjectStore + ?Sized> VerifiedMetadataTables<'_, S> {
pub(crate) fn manifest(&self) -> &NamespaceManifestEnvelope {
self.manifest.as_ref()
}
pub(crate) async fn get_for_lookup(
&self,
family: MetadataTableFamily,
key: &str,
filter_probe: &str,
) -> Result<Option<MetadataRow>, ManifestLoadError> {
Ok(self
.scan_prefix_rows(family, key, Some(filter_probe), Readahead::Enabled)
.await?
.into_iter()
.find(|(row_key, _)| row_key == key)
.map(|(_, row)| row))
}
#[cfg(test)]
pub(crate) async fn scan_prefix(
&self,
family: MetadataTableFamily,
prefix: &str,
) -> Result<Vec<MetadataRow>, ManifestLoadError> {
Ok(strip_row_keys(
self.scan_prefix_rows(family, prefix, None, Readahead::Enabled)
.await?,
))
}
pub(super) async fn scan_prefix_in_runs(
&self,
runs: &[MetadataRunManifest],
family: MetadataTableFamily,
prefix: &str,
) -> Result<Vec<MetadataRow>, ManifestLoadError> {
Ok(strip_row_keys(
self.scan_prefix_rows_in_runs(runs, family, prefix, None, Readahead::Disabled)
.await?,
))
}
pub(crate) async fn scan_prefix_for_lookup(
&self,
family: MetadataTableFamily,
prefix: &str,
filter_probe: &str,
readahead: Readahead,
) -> Result<Vec<MetadataRow>, ManifestLoadError> {
Ok(strip_row_keys(
self.scan_prefix_rows(family, prefix, Some(filter_probe), readahead)
.await?,
))
}
pub(crate) async fn scan_range_page(
&self,
family: MetadataTableFamily,
lower_bound: &str,
upper_bound: Option<&str>,
limit: usize,
) -> Result<Vec<MetadataRow>, ManifestLoadError> {
Ok(strip_row_keys(
self.scan_range_page_with_keys(family, lower_bound, upper_bound, limit)
.await?,
))
}
pub(crate) async fn scan_range_page_with_keys(
&self,
family: MetadataTableFamily,
lower_bound: &str,
upper_bound: Option<&str>,
limit: usize,
) -> Result<Vec<(String, MetadataRow)>, ManifestLoadError> {
if limit == 0 {
return Ok(Vec::new());
}
self.scan_range_page_rows(
self.scan_runs.as_ref(),
family,
lower_bound,
upper_bound,
limit,
None,
Readahead::Enabled,
)
.await
}
pub(crate) async fn scan_range_page_for_lookup(
&self,
family: MetadataTableFamily,
lower_bound: &str,
upper_bound: Option<&str>,
limit: usize,
filter_probe: &str,
) -> Result<Vec<MetadataRow>, ManifestLoadError> {
if limit == 0 {
return Ok(Vec::new());
}
Ok(strip_row_keys(
self.scan_range_page_rows(
self.scan_runs.as_ref(),
family,
lower_bound,
upper_bound,
limit,
Some(filter_probe),
Readahead::Enabled,
)
.await?,
))
}
async fn scan_prefix_rows(
&self,
family: MetadataTableFamily,
prefix: &str,
filter_probe: Option<&str>,
readahead: Readahead,
) -> Result<Vec<(String, MetadataRow)>, ManifestLoadError> {
self.scan_prefix_rows_in_runs(
self.scan_runs.as_ref(),
family,
prefix,
filter_probe,
readahead,
)
.await
}
async fn scan_prefix_rows_in_runs(
&self,
runs: &[MetadataRunManifest],
family: MetadataTableFamily,
prefix: &str,
filter_probe: Option<&str>,
readahead: Readahead,
) -> Result<Vec<(String, MetadataRow)>, ManifestLoadError> {
let upper_bound = string_prefix_upper_bound(prefix);
self.scan_range_page_rows(
runs,
family,
prefix,
upper_bound.as_deref(),
usize::MAX,
filter_probe,
readahead,
)
.await
}
async fn segment_filter_admits(
&self,
descriptor: &MetadataFileRef,
filter_probe: &str,
) -> Result<bool, ManifestLoadError> {
let filter =
load_segment_filter(self.store, self.table_cache, &self.block_memo, descriptor).await?;
if filter.may_contain(filter_probe) {
return Ok(true);
}
if let Some(cache) = self.table_cache {
cache.record_filter_skip();
}
Ok(false)
}
fn record_filter_false_positive_if_empty(&self, matched_rows: usize) {
if matched_rows == 0 {
if let Some(cache) = self.table_cache {
cache.record_filter_false_positive();
}
}
}
#[allow(clippy::too_many_arguments)]
async fn scan_range_page_rows(
&self,
runs: &[MetadataRunManifest],
family: MetadataTableFamily,
lower_bound: &str,
upper_bound: Option<&str>,
limit: usize,
filter_probe: Option<&str>,
readahead: Readahead,
) -> Result<Vec<(String, MetadataRow)>, ManifestLoadError> {
let mut candidates = Vec::new();
for run in runs {
let table = manifest_table_for_family(&self.manifest_object_key, &run.tables, family)?;
candidates.extend(table.segments.iter().filter(|descriptor| {
descriptor_may_intersect_range(descriptor, lower_bound, upper_bound)
}));
}
let mut matching_descriptors = self
.filter_admitted_descriptors(candidates, filter_probe)
.await?;
matching_descriptors.sort_by(|left, right| {
left.min_key
.cmp(&right.min_key)
.then(left.max_key.cmp(&right.max_key))
.then(left.object_key.cmp(&right.object_key))
});
let mut rows = Vec::<(String, MetadataRow)>::new();
let mut next_descriptor_index = 0;
while next_descriptor_index < matching_descriptors.len() {
let should_load_next = if rows.len() < limit {
true
} else {
let boundary_key = &rows[limit - 1].0;
matching_descriptors[next_descriptor_index].min_key <= *boundary_key
};
if !should_load_next {
break;
}
let chunk_end = (next_descriptor_index + MAX_MATERIALIZED_TABLE_LOADS)
.min(matching_descriptors.len());
let loaded_segments = try_join_all(
matching_descriptors[next_descriptor_index..chunk_end]
.iter()
.map(|descriptor| {
self.segment_rows(family, descriptor, lower_bound, upper_bound, readahead)
}),
)
.await?;
next_descriptor_index = chunk_end;
for segment_rows in loaded_segments {
let matched_before = rows.len();
rows.extend(
segment_rows
.rows_in_key_range(lower_bound, upper_bound)
.take(limit)
.map(|(row_key, row)| (row_key.to_owned(), row.clone())),
);
if filter_probe.is_some() {
self.record_filter_false_positive_if_empty(rows.len() - matched_before);
}
}
rows.sort_by(|(left_key, _), (right_key, _)| left_key.cmp(right_key));
}
rows.truncate(limit);
Ok(rows)
}
async fn filter_admitted_descriptors<'d>(
&self,
descriptors: Vec<&'d MetadataFileRef>,
filter_probe: Option<&str>,
) -> Result<Vec<&'d MetadataFileRef>, ManifestLoadError> {
let Some(filter_probe) = filter_probe else {
return Ok(descriptors);
};
let mut admitted = Vec::with_capacity(descriptors.len());
for chunk in descriptors.chunks(MAX_MATERIALIZED_TABLE_LOADS) {
let checks = try_join_all(
chunk
.iter()
.map(|descriptor| self.segment_filter_admits(descriptor, filter_probe)),
)
.await?;
admitted.extend(
chunk
.iter()
.zip(checks)
.filter(|(_, admits)| *admits)
.map(|(descriptor, _)| *descriptor),
);
}
Ok(admitted)
}
async fn segment_rows(
&self,
family: MetadataTableFamily,
descriptor: &MetadataFileRef,
lower_bound: &str,
upper_bound: Option<&str>,
readahead: Readahead,
) -> Result<SegmentKeyRangeBlocks, ManifestLoadError> {
load_manifest_segment_rows_in_key_range_with_cache(
self.store,
self.table_cache,
&self.block_memo,
family,
descriptor,
lower_bound,
upper_bound,
readahead,
)
.await
}
}
fn strip_row_keys(rows: Vec<(String, MetadataRow)>) -> Vec<MetadataRow> {
rows.into_iter().map(|(_, row)| row).collect()
}
pub(super) fn ordered_manifest_tables<'a>(
manifest_object_key: &str,
tables: &'a [MetadataTableManifest],
) -> Result<Vec<&'a MetadataTableManifest>, ManifestLoadError> {
let mut ordered = Vec::with_capacity(CHECKPOINT_TABLE_FAMILIES.len());
for family in CHECKPOINT_TABLE_FAMILIES {
let mut matching = tables.iter().filter(|table| table.family == family);
let Some(table) = matching.next() else {
return Err(ManifestLoadError::MissingTableFamily {
object_key: manifest_object_key.to_owned(),
family,
});
};
if matching.next().is_some() {
return Err(ManifestLoadError::DuplicateTableFamily {
object_key: manifest_object_key.to_owned(),
family,
});
}
ordered.push(table);
}
Ok(ordered)
}
pub(super) fn manifest_table_for_family<'a>(
manifest_object_key: &str,
tables: &'a [MetadataTableManifest],
family: MetadataTableFamily,
) -> Result<&'a MetadataTableManifest, ManifestLoadError> {
tables.iter().find(|table| table.family == family).ok_or(
ManifestLoadError::MissingTableFamily {
object_key: manifest_object_key.to_owned(),
family,
},
)
}
pub(super) fn descriptor_may_intersect_range(
descriptor: &MetadataFileRef,
lower_bound: &str,
upper_bound: Option<&str>,
) -> bool {
if descriptor.row_count == 0 {
return false;
}
if descriptor.max_key.as_str() < lower_bound {
return false;
}
if upper_bound
.map(|upper_bound| descriptor.min_key.as_str() >= upper_bound)
.unwrap_or(false)
{
return false;
}
true
}