pub mod aggregates;
pub mod bloom;
pub mod commit;
pub mod disk_cache;
pub mod encoding;
pub mod hll;
pub mod list;
pub mod list_prune;
pub mod options_hash;
pub mod part;
pub mod partition;
pub mod term_range;
use std::{
cmp::Ordering,
collections::{BTreeMap, HashMap, HashSet},
f32::consts::PI,
fmt,
ops::Deref,
sync::{Arc, OnceLock},
};
use arrow::compute::kernels::aggregate as agg;
use arrow_array::*;
use arrow_schema::{DataType, TimeUnit};
use bytes::Bytes;
use dashmap::DashMap;
use futures::future;
pub use list::{FtsSummaryAgg, GlobalVectorIndex, RoutingRef, ScalarStatsAgg};
use rayon::{ThreadPool, prelude::*};
use tokio::{sync::OnceCell, task::spawn_blocking};
use uuid::Uuid;
use xxhash_rust::xxh3::xxh3_64;
use super::options::SupertableOptions;
use crate::{
storage::{StorageError, StorageProvider},
superfile::{
builder::VectorConfig,
vector::{
distance::{
COSINE_DISTANCE_BASE, L2_CROSS_TERM_COEFF, Metric, all_centroid_scores_transposed,
distance, dot, insert_ranked, nearest_k_centroids_transposed,
transpose_centroids_cluster_major,
},
layout::VectorLayout,
quant::BitQuantizer,
rotation::RandomRotation,
},
},
supertable::{
CommitError,
error::ManifestError,
manifest::{
commit::{
EncodedPart, PointerFile, frame_content_size, part_uri, read_pointer,
translate_contention, write_manifest, write_part_bytes, write_pointer,
},
disk_cache::ManifestDiskCache,
encoding::SummaryWireMode,
list::{
FORMAT_VERSION as LIST_FORMAT_VERSION, Manifest, ManifestPartEntry,
PartitionStrategy,
},
part::{ContentHash, ManifestPart, PartId},
partition::{assign_partition, encode_partition_key},
},
query::{hierarchical_iter, prune::PruneLeaf},
slow_vector_state,
},
};
pub(crate) const SUPERFILE_DATA_DIR: &str = "data";
pub(crate) const DEFAULT_VECTOR_INDEX_PREFIX: &str = "_vector_index";
#[derive(Debug, Clone)]
pub struct SuperfileList {
pub manifest_id: u64,
pub options: Arc<SupertableOptions>,
pub superfiles: Vec<Arc<SuperfileEntry>>,
pub(crate) vector_index_storage_prefix: Option<String>,
}
impl SuperfileList {
pub fn empty(options: Arc<SupertableOptions>) -> Self {
Self {
manifest_id: 0,
options,
superfiles: Vec::new(),
vector_index_storage_prefix: None,
}
}
pub(crate) fn empty_with_vector_index_prefix(
options: Arc<SupertableOptions>,
vector_index_storage_prefix: Option<String>,
) -> Self {
Self {
manifest_id: 0,
options,
superfiles: Vec::new(),
vector_index_storage_prefix,
}
}
pub fn with_appended(&self, new_entries: Vec<Arc<SuperfileEntry>>) -> Self {
let mut superfiles = self.superfiles.clone();
superfiles.extend(new_entries);
Self {
manifest_id: self.manifest_id + 1,
options: self.options.clone(),
superfiles,
vector_index_storage_prefix: self.vector_index_storage_prefix.clone(),
}
}
pub fn n_docs_total(&self) -> u64 {
self.superfiles.iter().map(|s| s.n_docs).sum()
}
}
pub struct ManifestSnapshot {
superfile_list: SuperfileList,
list: Option<Manifest>,
parts: DashMap<PartId, Arc<OnceCell<Arc<ManifestPart>>>>,
loader: Option<Arc<ManifestPartLoader>>,
stamped_partition_strategy: Option<PartitionStrategy>,
stamped_global_vector_index: Option<list::GlobalVectorIndex>,
stamped_drained_ranges: Option<list::DrainedVersionRanges>,
}
impl fmt::Debug for ManifestSnapshot {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ManifestSnapshot")
.field("manifest_id", &self.superfile_list.manifest_id)
.field("n_superfiles", &self.superfile_list.superfiles.len())
.field("has_list", &self.list.is_some())
.field(
"n_parts",
&self.list.as_ref().map(|l| l.parts.len()).unwrap_or(0),
)
.field("n_parts_loaded", &self.parts.len())
.field("has_loader", &self.loader.is_some())
.finish()
}
}
impl Deref for ManifestSnapshot {
type Target = SuperfileList;
fn deref(&self) -> &Self::Target {
&self.superfile_list
}
}
impl ManifestSnapshot {
pub fn new(
manifest_id: u64,
options: Arc<SupertableOptions>,
superfile_list: Vec<Arc<SuperfileEntry>>,
storage: Option<Arc<dyn StorageProvider>>,
list: Option<Manifest>,
) -> Self {
let superfile_list = SuperfileList {
manifest_id,
options,
superfiles: superfile_list,
vector_index_storage_prefix: None,
};
if let Some(storage) = storage
&& let Some(list) = list
{
let manifest_cache = superfile_list.options.manifest_disk_cache.clone();
let loader = Arc::new(ManifestPartLoader::new_with_cache(
Arc::clone(&storage),
&list,
manifest_cache,
));
Self {
superfile_list,
list: Some(list),
parts: DashMap::new(),
loader: Some(loader),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
}
} else {
Self {
superfile_list,
list: None,
parts: DashMap::new(),
loader: None,
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
}
}
}
#[cfg(test)]
pub fn new_from_superfiles(
opts: Arc<SupertableOptions>,
superfiles: Vec<Arc<SuperfileEntry>>,
) -> Self {
ManifestSnapshot::empty(opts).with_appended(superfiles)
}
pub fn empty(options: Arc<SupertableOptions>) -> Self {
Self {
superfile_list: SuperfileList::empty(options),
list: None,
parts: DashMap::new(),
loader: None,
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
}
}
pub(crate) fn empty_with_vector_index_prefix(
options: Arc<SupertableOptions>,
vector_index_storage_prefix: Option<String>,
) -> Self {
Self {
superfile_list: SuperfileList::empty_with_vector_index_prefix(
options,
vector_index_storage_prefix,
),
list: None,
parts: dashmap::DashMap::new(),
loader: None,
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
}
}
pub(crate) fn materialized_empty_with_vector_index_prefix(
options: Arc<SupertableOptions>,
vector_index_storage_prefix: Option<String>,
) -> Self {
let strategy = options.effective_partition_strategy();
let list = Self::build_list(
&options,
strategy,
0,
Vec::new(),
vector_index_storage_prefix.clone(),
BTreeMap::new(),
);
let loader = options.storage.as_ref().map(|storage| {
Arc::new(ManifestPartLoader::new_with_cache(
storage.clone(),
&list,
options.manifest_disk_cache.clone(),
))
});
Self {
superfile_list: SuperfileList::empty_with_vector_index_prefix(
options,
vector_index_storage_prefix,
),
list: Some(list),
parts: DashMap::new(),
loader,
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
}
}
fn build_list(
options: &SupertableOptions,
strategy: PartitionStrategy,
manifest_id: u64,
parts: Vec<ManifestPartEntry>,
vector_index_storage_prefix: Option<String>,
tombstone_seqs: BTreeMap<Uuid, u64>,
) -> Manifest {
Manifest {
format_version: LIST_FORMAT_VERSION.into(),
manifest_id,
options_hash: options_hash::compute_options_hash(options, &strategy),
schema: Vec::new(),
id_column: options.id_column.clone(),
fts_columns: options
.fts_columns
.iter()
.map(|f| list::FtsColumnInfo {
column: f.column.clone(),
})
.collect(),
vector_columns: options
.vector_columns
.iter()
.map(|v| list::VectorColumnInfo {
column: v.column.clone(),
dim: v.dim,
n_cent: v.n_cent,
rot_seed: v.rot_seed,
metric: format!("{:?}", v.metric).to_lowercase(),
})
.collect(),
partition_strategy: strategy,
vector_index_storage_prefix,
global_vector_index: None,
drained_ranges: Default::default(),
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts,
tombstone_seqs,
}
}
pub fn get_manifest_id(&self) -> u64 {
self.superfile_list.manifest_id
}
pub fn get_next_manifest_id(&self) -> u64 {
self.get_manifest_id() + 1
}
pub fn get_opts(&self) -> Arc<SupertableOptions> {
self.superfile_list.options.clone()
}
pub fn get_partition_strategy(&self) -> list::PartitionStrategy {
self.partition_strategy()
.cloned()
.unwrap_or_else(|| self.superfile_list.options.effective_partition_strategy())
}
pub(crate) fn partition_strategy(&self) -> Option<&list::PartitionStrategy> {
self.stamped_partition_strategy
.as_ref()
.or_else(|| self.list.as_ref().map(|l| &l.partition_strategy))
}
pub(crate) fn vector_cell_routing(&self) -> Option<list::CellRoutingParams> {
match self.partition_strategy() {
Some(list::PartitionStrategy::VectorCell { routing, .. }) => Some(*routing),
_ => None,
}
}
pub(crate) fn vector_cell_clusters(&self, column: &str) -> Option<&ClusterCentroids> {
match self.partition_strategy() {
Some(list::PartitionStrategy::VectorCell {
column: cell_column,
clusters,
..
}) if cell_column == column => Some(clusters),
_ => None,
}
}
pub fn get_global_vector_index(&self) -> Option<list::GlobalVectorIndex> {
self.global_vector_index().cloned()
}
pub(crate) fn global_vector_index(&self) -> Option<&list::GlobalVectorIndex> {
self.stamped_global_vector_index.as_ref().or_else(|| {
self.list
.as_ref()
.and_then(|l| l.global_vector_index.as_ref())
})
}
pub fn get_drained_ranges(&self) -> list::DrainedVersionRanges {
if let Some(d) = &self.stamped_drained_ranges {
return d.clone();
}
self.list
.as_ref()
.map(|l| l.drained_ranges.clone())
.unwrap_or_default()
}
pub fn get_num_parts(&self) -> usize {
self.list.as_ref().map(|l| l.parts.len()).unwrap_or(0)
}
pub fn get_num_parts_loaded(&self) -> usize {
self.parts.len()
}
pub fn is_in_process_only(&self) -> bool {
self.list.is_none()
}
pub(crate) fn complete_flat_superfiles(&self) -> Option<&[Arc<SuperfileEntry>]> {
match self.list.as_ref() {
None => Some(&self.superfile_list.superfiles),
Some(list) if list.slow_vector_state_uri.is_some() => {
Some(&self.superfile_list.superfiles)
}
Some(list) => {
let expected: u64 = list.parts.iter().map(|entry| entry.n_superfiles).sum();
(self.superfile_list.superfiles.len() as u64 == expected)
.then_some(self.superfile_list.superfiles.as_slice())
}
}
}
pub(crate) fn vector_index_storage_prefix(&self) -> Option<&str> {
if let Some(list) = self.list.as_ref()
&& let Some(prefix) = list.vector_index_storage_prefix.as_deref()
{
return Some(prefix);
}
self.superfile_list.vector_index_storage_prefix.as_deref()
}
fn stamp_vector_index_storage_prefix(
&self,
vector_columns: &[list::VectorColumnInfo],
) -> Option<String> {
if vector_columns.is_empty() {
return None;
}
if let Some(prefix) = self.vector_index_storage_prefix() {
return Some(prefix.to_string());
}
Some(DEFAULT_VECTOR_INDEX_PREFIX.to_string())
}
pub fn get_cached_part_by_id(&self, part_id: &PartId) -> Option<Arc<ManifestPart>> {
self.parts
.get(part_id)
.and_then(|cell| cell.value().get().cloned())
}
pub fn get_cached_part_by_list_idx(&self, idx: usize) -> Option<Arc<ManifestPart>> {
let Some(list) = &self.list else {
return None;
};
let part_id = list.parts[idx].part_id;
self.get_cached_part_by_id(&part_id)
}
pub(crate) async fn load(
current_manifest: Option<Arc<Self>>,
storage: Arc<dyn StorageProvider>,
options: Option<Arc<SupertableOptions>>,
) -> Result<Arc<Self>, ManifestLoadError> {
let (pointer, _) = match read_pointer(storage.as_ref()).await? {
Some(p) => p,
None => return Err(ManifestLoadError::PointerNotFound),
};
Self::load_with_pointer(current_manifest, storage, options, pointer).await
}
pub(crate) async fn load_with_pointer(
current_manifest: Option<Arc<Self>>,
storage: Arc<dyn StorageProvider>,
options: Option<Arc<SupertableOptions>>,
pointer: PointerFile,
) -> Result<Arc<Self>, ManifestLoadError> {
if let Some(current_manifest) = ¤t_manifest
&& current_manifest.superfile_list.manifest_id >= pointer.manifest_id
{
return Err(ManifestLoadError::AlreadyLoaded);
}
let (list_bytes, _) = storage
.get(&pointer.manifest_uri)
.await
.map_err(ManifestLoadError::Storage)?;
let list = list::decode(&list_bytes).map_err(ManifestLoadError::ListParse)?;
let options = if let Some(options) = options {
options
} else if let Some(current) = ¤t_manifest {
current.options.clone()
} else {
return Err(ManifestLoadError::ContentHashMismatch {
expected: "valid options".to_string(),
actual: "None options".to_string(),
});
};
let expected_hash = options_hash::compute_options_hash(&options, &list.partition_strategy);
if let Err(mismatch) = options_hash::verify_options_hash(expected_hash, list.options_hash) {
return Err(ManifestLoadError::ContentHashMismatch {
expected: mismatch.expected,
actual: mismatch.actual,
});
}
let loader = Arc::new(ManifestPartLoader::new_with_cache_and_mode(
Arc::clone(&storage),
&list,
options.manifest_disk_cache.clone(),
options.summary_centroids_from_superfiles,
));
let parts: DashMap<_, _> = DashMap::new();
let mut all_superfiles: Vec<Arc<SuperfileEntry>> = Vec::new();
let expected_n_superfiles: Option<u64> = if list.slow_vector_state_uri.is_some() {
None
} else {
Some(list.parts.iter().map(|e| e.n_superfiles).sum())
};
let reused: Option<Vec<Arc<SuperfileEntry>>> = match (
list.slow_vector_state_uri.as_deref(),
list.slow_vector_state_content_hash,
) {
(Some(uri), Some(hash)) => current_manifest.as_ref().and_then(|cur| {
let same_ref = cur.list.as_ref().is_some_and(|cl| {
cl.slow_vector_state_uri.as_deref() == Some(uri)
&& cl.slow_vector_state_content_hash == Some(hash)
});
let complete = expected_n_superfiles
.is_none_or(|expected| cur.superfile_list.superfiles.len() as u64 == expected);
(same_ref && complete).then(|| cur.superfile_list.superfiles.clone())
}),
_ => None,
};
let entries_reused = reused.is_some();
let hydrated: Option<Vec<Arc<SuperfileEntry>>> = match reused {
Some(entries) => Some(entries),
None => match (
list.slow_vector_state_uri.as_deref(),
list.slow_vector_state_content_hash,
) {
(Some(uri), Some(hash)) => {
let entries = slow_vector_state::load_state(storage.as_ref(), uri, &hash)
.await
.map_err(|e| ManifestLoadError::SlowStateHydration(e.to_string()))?;
if let Some(expected) = expected_n_superfiles
&& entries.len() as u64 != expected
{
return Err(ManifestLoadError::SlowStateHydration(format!(
"blob entry count {} != list total {expected}",
entries.len(),
)));
}
Some(entries)
}
_ => None,
},
};
if let Some(entries) = hydrated {
for entry in &list.parts {
let inherited = current_manifest
.as_ref()
.and_then(|cur| cur.parts.get(&entry.part_id).map(|kv| kv.value().clone()));
parts.insert(
entry.part_id,
inherited.unwrap_or_else(|| Arc::new(OnceCell::new())),
);
}
all_superfiles = entries;
} else if let Some(current_manifest) = ¤t_manifest {
let mut missing_part_ids = Vec::new();
for entry in &list.parts {
if let Some(existing) = current_manifest.parts.get(&entry.part_id) {
parts.insert(entry.part_id, existing.value().clone());
} else {
missing_part_ids.push(entry.part_id);
}
}
let threshold = options.eager_load_threshold_parts as usize;
let eager = list.parts.len() <= threshold;
if eager {
let load_futs = missing_part_ids
.iter()
.map(|id| {
let loader = Arc::clone(&loader);
let pid = *id;
async move { loader.load(pid).await }
})
.collect::<Vec<_>>();
let loaded = future::join_all(load_futs).await;
for (pid, result) in missing_part_ids.iter().zip(loaded) {
let part = result?;
let cell = OnceCell::new();
cell.set(part).expect("fresh cell");
parts.insert(*pid, Arc::new(cell));
}
for entry in &list.parts {
let cell = parts.get(&entry.part_id).expect("part inserted above");
let part = cell
.value()
.get()
.expect("eager-fetched or inherited; must be set");
all_superfiles.extend(part.superfiles.iter().cloned());
}
} else {
for pid in &missing_part_ids {
parts.insert(*pid, Arc::new(OnceCell::new()));
}
}
} else {
let n_parts = list.parts.len();
let threshold = options.eager_load_threshold_parts as usize;
let eager = n_parts <= threshold;
if eager {
let part_ids: Vec<_> = list.parts.iter().map(|p| p.part_id).collect();
let load_futs = part_ids
.iter()
.map(|id| {
let loader = Arc::clone(&loader);
let pid = *id;
async move { loader.load(pid).await }
})
.collect::<Vec<_>>();
let loaded = future::join_all(load_futs).await;
for (pid, result) in part_ids.iter().zip(loaded) {
let part = result?;
all_superfiles.extend(part.superfiles.iter().cloned());
let cell = OnceCell::new();
cell.set(part).expect("fresh OnceCell");
parts.insert(*pid, Arc::new(cell));
}
} else {
for entry in &list.parts {
parts.insert(entry.part_id, Arc::new(OnceCell::new()));
}
}
}
if !entries_reused {
let strip = matches!(
list.partition_strategy,
PartitionStrategy::VectorCell { .. }
);
let vector_columns = options.vector_columns.clone();
let pool = Arc::clone(&options.reader_pool);
let mut entries = all_superfiles;
let prewarm = move || {
if strip {
strip_summary_centroids(&mut entries, &vector_columns);
}
prewarm_summary_admit_slabs(&entries, &vector_columns, &pool);
entries
};
all_superfiles = match spawn_blocking(prewarm).await {
Ok(entries) => entries,
Err(join_error) => {
return Err(ManifestLoadError::SlowStateHydration(format!(
"admit slab prewarm task failed: {join_error}"
)));
}
};
}
let mut new_superfile_list = SuperfileList::empty(options.clone());
new_superfile_list.manifest_id = pointer.manifest_id;
new_superfile_list.superfiles = all_superfiles;
let new_manifest = ManifestSnapshot {
superfile_list: new_superfile_list,
list: Some(list),
parts,
loader: Some(loader),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
};
Ok(Arc::new(new_manifest))
}
pub async fn write(
&self,
storage: &dyn StorageProvider,
expected_prev_etag: Option<&str>,
parts_to_write: &[&[u8]],
) -> Result<(), CommitError> {
let Some(list_to_write) = self.list.as_ref() else {
return Ok(());
};
let list_fut = write_manifest(storage, list_to_write);
let part_futs = parts_to_write
.iter()
.map(|encoded| write_part_bytes(storage, encoded));
let part_join = future::join_all(part_futs);
let (list_res, part_results) = tokio::join!(list_fut, part_join);
let list_res = list_res.map_err(translate_contention)?;
for part_result in part_results {
part_result.map_err(translate_contention)?;
}
let pointer = PointerFile {
manifest_id: self.get_manifest_id(),
manifest_uri: list_res.uri,
content_hash: list_res.content_hash,
};
write_pointer(storage, &pointer, expected_prev_etag).await?;
Ok(())
}
pub fn get_all_superfiles(&self) -> &[Arc<SuperfileEntry>] {
&self.superfile_list.superfiles
}
pub(crate) async fn get_pruned_superfiles(
&self,
leaves: &[PruneLeaf],
) -> Result<Vec<Arc<SuperfileEntry>>, ManifestLoadError> {
match &self.list {
Some(list) => {
if list.slow_vector_state_uri.is_some() {
return Ok(self.superfile_list.superfiles.clone());
}
let expected: u64 = list.parts.iter().map(|e| e.n_superfiles).sum();
if self.superfile_list.superfiles.len() as u64 == expected {
return Ok(self.superfile_list.superfiles.clone());
}
let mut kept: Option<HashSet<PartId>> = None;
for leaf in leaves {
if let Some(part_ids) = leaf.keep_parts(list) {
let set: HashSet<PartId> = part_ids.into_iter().collect();
kept = Some(match kept {
None => set,
Some(existing) => existing.intersection(&set).copied().collect(),
});
}
}
let ordered: Vec<PartId> = match kept {
Some(set) => list
.parts
.iter()
.map(|p| p.part_id)
.filter(|id| set.contains(id))
.collect(),
None => list.parts.iter().map(|p| p.part_id).collect(),
};
hierarchical_iter::load_and_flatten(self, &ordered).await
}
None => Ok(hierarchical_iter::fallback_to_flat_superfiles(self)),
}
}
pub(crate) async fn get_all_superfiles_loaded(
&self,
) -> Result<Vec<Arc<SuperfileEntry>>, ManifestLoadError> {
match &self.list {
Some(list) => {
if list.slow_vector_state_uri.is_some() {
return Ok(self.superfile_list.superfiles.clone());
}
let expected: u64 = list.parts.iter().map(|e| e.n_superfiles).sum();
if self.superfile_list.superfiles.len() as u64 == expected {
return Ok(self.superfile_list.superfiles.clone());
}
let all: Vec<PartId> = list.parts.iter().map(|p| p.part_id).collect();
hierarchical_iter::load_and_flatten(self, &all).await
}
None => Ok(hierarchical_iter::fallback_to_flat_superfiles(self)),
}
}
pub(crate) async fn get_undrained_superfiles_loaded(
&self,
drained: &list::DrainedVersionRanges,
) -> Result<Vec<Arc<SuperfileEntry>>, ManifestLoadError> {
let Some(list) = &self.list else {
return Ok(self
.superfile_list
.superfiles
.iter()
.filter(|entry| !drained.contains(entry.birth_version))
.cloned()
.collect());
};
if list.slow_vector_state_uri.is_some() {
return Ok(self
.superfile_list
.superfiles
.iter()
.filter(|entry| !drained.contains(entry.birth_version))
.cloned()
.collect());
}
let expected: u64 = list.parts.iter().map(|entry| entry.n_superfiles).sum();
if self.superfile_list.superfiles.len() as u64 == expected {
return Ok(self
.superfile_list
.superfiles
.iter()
.filter(|entry| !drained.contains(entry.birth_version))
.cloned()
.collect());
}
let part_ids: Vec<PartId> = list
.parts
.iter()
.filter(|entry| {
entry
.birth_version_range()
.is_none_or(|(lo, hi)| !drained.covers(lo, hi))
})
.map(|entry| entry.part_id)
.collect();
let entries = hierarchical_iter::load_and_flatten(self, &part_ids).await?;
Ok(entries
.into_iter()
.filter(|entry| !drained.contains(entry.birth_version))
.collect())
}
pub fn get_all_list_entries(&self) -> &[ManifestPartEntry] {
match &self.list {
Some(list) => &list.parts,
None => &[],
}
}
pub fn with_appended(&self, new_entries: Vec<Arc<SuperfileEntry>>) -> Self {
Self {
superfile_list: self.superfile_list.with_appended(new_entries),
list: self.list.clone(),
parts: DashMap::new(),
loader: self.loader.clone(),
stamped_partition_strategy: self.stamped_partition_strategy.clone(),
stamped_global_vector_index: self.stamped_global_vector_index.clone(),
stamped_drained_ranges: self.stamped_drained_ranges.clone(),
}
}
pub(crate) fn deleted_user_ids_inline(&self) -> Option<&[u8]> {
self.list.as_ref()?.deleted_user_ids_inline.as_deref()
}
pub(crate) fn slow_vector_state_blob(&self) -> Option<(&str, part::ContentHash)> {
let list = self.list.as_ref()?;
Some((
list.slow_vector_state_uri.as_deref()?,
list.slow_vector_state_content_hash?,
))
}
pub(crate) fn slow_vector_state_centroids_blob(&self) -> Option<&RoutingRef> {
self.list.as_ref()?.slow_vector_state_centroids.as_ref()
}
pub fn with_deleted_user_ids(&self, encoded: Vec<u8>) -> Self {
let next_id = self.get_next_manifest_id();
let new_list = self.list.as_ref().map(|list| {
let mut list = list.clone();
list.manifest_id = next_id;
list.deleted_user_ids_inline = Some(encoded.clone());
list
});
Self {
superfile_list: SuperfileList {
manifest_id: next_id,
options: Arc::clone(&self.superfile_list.options),
superfiles: self.superfile_list.superfiles.clone(),
vector_index_storage_prefix: self
.superfile_list
.vector_index_storage_prefix
.clone(),
},
list: new_list,
parts: self.parts.clone(),
loader: self.loader.clone(),
stamped_partition_strategy: self.stamped_partition_strategy.clone(),
stamped_global_vector_index: self.stamped_global_vector_index.clone(),
stamped_drained_ranges: self.stamped_drained_ranges.clone(),
}
}
pub fn with_slow_vector_state(
&self,
uri: String,
hash: part::ContentHash,
centroids: RoutingRef,
) -> Self {
let next_id = self.get_next_manifest_id();
let new_list = self.list.as_ref().map(|list| {
let mut list = list.clone();
list.manifest_id = next_id;
list.slow_vector_state_uri = Some(uri);
list.slow_vector_state_content_hash = Some(hash);
list.slow_vector_state_centroids = Some(centroids);
list
});
Self {
superfile_list: SuperfileList {
manifest_id: next_id,
options: Arc::clone(&self.superfile_list.options),
superfiles: self.superfile_list.superfiles.clone(),
vector_index_storage_prefix: self
.superfile_list
.vector_index_storage_prefix
.clone(),
},
list: new_list,
parts: self.parts.clone(),
loader: self.loader.clone(),
stamped_partition_strategy: self.stamped_partition_strategy.clone(),
stamped_global_vector_index: self.stamped_global_vector_index.clone(),
stamped_drained_ranges: self.stamped_drained_ranges.clone(),
}
}
pub(crate) fn with_slow_vector_state_ref(
&self,
uri: String,
hash: part::ContentHash,
centroids: RoutingRef,
) -> Self {
let new_list = self.list.as_ref().map(|list| {
let mut list = list.clone();
list.slow_vector_state_uri = Some(uri);
list.slow_vector_state_content_hash = Some(hash);
list.slow_vector_state_centroids = Some(centroids);
list
});
Self {
superfile_list: SuperfileList {
manifest_id: self.superfile_list.manifest_id,
options: Arc::clone(&self.superfile_list.options),
superfiles: self.superfile_list.superfiles.clone(),
vector_index_storage_prefix: self
.superfile_list
.vector_index_storage_prefix
.clone(),
},
list: new_list,
parts: self.parts.clone(),
loader: self.loader.clone(),
stamped_partition_strategy: self.stamped_partition_strategy.clone(),
stamped_global_vector_index: self.stamped_global_vector_index.clone(),
stamped_drained_ranges: self.stamped_drained_ranges.clone(),
}
}
pub fn with_partition_strategy(&self, strategy: list::PartitionStrategy) -> Self {
let new_list = match self.list.as_ref() {
Some(list) => {
let mut list = list.clone();
list.partition_strategy = strategy.clone();
Some(list)
}
None => None,
};
Self {
superfile_list: SuperfileList {
manifest_id: self.manifest_id,
options: Arc::clone(&self.options),
superfiles: self.superfiles.clone(),
vector_index_storage_prefix: self.vector_index_storage_prefix.clone(),
},
list: new_list.or_else(|| self.list.clone()),
parts: self.parts.clone(),
loader: self.loader.clone(),
stamped_partition_strategy: Some(strategy),
stamped_global_vector_index: self.stamped_global_vector_index.clone(),
stamped_drained_ranges: self.stamped_drained_ranges.clone(),
}
}
pub fn with_global_vector_index(&self, index: list::GlobalVectorIndex) -> Self {
let new_list = self.list.as_ref().map(|list| {
let mut list = list.clone();
list.global_vector_index = Some(index.clone());
list
});
Self {
superfile_list: SuperfileList {
manifest_id: self.manifest_id,
options: Arc::clone(&self.options),
superfiles: self.superfiles.clone(),
vector_index_storage_prefix: self.vector_index_storage_prefix.clone(),
},
list: new_list.or_else(|| self.list.clone()),
parts: self.parts.clone(),
loader: self.loader.clone(),
stamped_partition_strategy: self.stamped_partition_strategy.clone(),
stamped_global_vector_index: Some(index),
stamped_drained_ranges: self.stamped_drained_ranges.clone(),
}
}
pub fn with_drained_ranges(&self, ranges: list::DrainedVersionRanges) -> Self {
let new_list = self.list.as_ref().map(|list| {
let mut list = list.clone();
list.drained_ranges = ranges.clone();
list
});
Self {
superfile_list: SuperfileList {
manifest_id: self.manifest_id,
options: Arc::clone(&self.options),
superfiles: self.superfiles.clone(),
vector_index_storage_prefix: self.vector_index_storage_prefix.clone(),
},
list: new_list.or_else(|| self.list.clone()),
parts: self.parts.clone(),
loader: self.loader.clone(),
stamped_partition_strategy: self.stamped_partition_strategy.clone(),
stamped_global_vector_index: self.stamped_global_vector_index.clone(),
stamped_drained_ranges: Some(ranges),
}
}
pub fn get_tombstone_seqs(&self) -> Option<&BTreeMap<Uuid, u64>> {
self.list.as_ref().map(|l| &l.tombstone_seqs)
}
pub(crate) fn with_tombstone_seqs_bumped(&self, touched: &[Uuid]) -> Option<Self> {
let list = self.list.as_ref()?;
let next_id = self.get_next_manifest_id();
let mut new_list = list.clone();
new_list.manifest_id = next_id;
for id in touched {
new_list.tombstone_seqs.insert(*id, next_id);
}
let mut superfile_list = self.superfile_list.clone();
superfile_list.manifest_id = next_id;
let parts = DashMap::new();
for kv in self.parts.iter() {
parts.insert(*kv.key(), Arc::clone(kv.value()));
}
Some(Self {
superfile_list,
list: Some(new_list),
parts,
loader: self.loader.clone(),
stamped_partition_strategy: self.stamped_partition_strategy.clone(),
stamped_global_vector_index: self.stamped_global_vector_index.clone(),
stamped_drained_ranges: self.stamped_drained_ranges.clone(),
})
}
pub async fn get_part_by_id(
&self,
part_id: PartId,
) -> Result<Arc<ManifestPart>, ManifestLoadError> {
let loader = self
.loader
.as_ref()
.ok_or(ManifestLoadError::NoLoaderAttached)?;
let cell = self
.parts
.entry(part_id)
.or_insert_with(|| Arc::new(OnceCell::new()))
.clone();
let loaded = cell.get_or_try_init(|| loader.load(part_id)).await?;
Ok(Arc::clone(loaded))
}
pub(crate) async fn user_centroids_for_rescore(&self) -> Option<Arc<UserCentroidCache>> {
let list = self.list.as_ref()?;
if matches!(
list.partition_strategy,
PartitionStrategy::VectorCell { .. }
) {
return None;
}
let loader = self.loader.as_ref()?;
if list.parts.is_empty() {
return None;
}
let manifest_id = self.get_manifest_id();
let slot = Arc::clone(&self.options.user_centroid_cache);
let mut guard = slot.lock().await;
if let Some(cache) = guard.as_ref()
&& cache.manifest_id == manifest_id
{
return Some(Arc::clone(cache));
}
let waves = list.parts.iter().map(|part| {
let loader = Arc::clone(loader);
let part_id = part.part_id;
async move { loader.load_full(part_id).await }
});
match future::try_join_all(waves).await {
Ok(parts) => {
let cache = Arc::new(UserCentroidCache::from_parts(manifest_id, &parts));
*guard = Some(Arc::clone(&cache));
Some(cache)
}
Err(error) => {
eprintln!(
"[supertable] full-part centroid hydration unavailable ({error}); falling \
back to per-superfile centroid reads"
);
None
}
}
}
pub(crate) async fn lookup_superfile_entry(
&self,
uri: SuperfileUri,
) -> Result<Option<Arc<SuperfileEntry>>, ManifestLoadError> {
if let Some(entry) = self.superfiles.iter().find(|e| e.uri == uri) {
return Ok(Some(Arc::clone(entry)));
}
let Some(list) = &self.list else {
return Ok(None);
};
for part_entry in &list.parts {
let part = self.get_part_by_id(part_entry.part_id).await?;
if let Some(entry) = part.superfiles.iter().find(|e| e.uri == uri) {
return Ok(Some(Arc::clone(entry)));
}
}
Ok(None)
}
pub async fn update(
&self,
new_entries: &[Arc<SuperfileEntry>],
entries_to_remove: &[Arc<SuperfileEntry>],
) -> Result<(ManifestSnapshot, Vec<EncodedPart>), ManifestError> {
self.update_inner(new_entries, entries_to_remove, false)
.await
}
pub(crate) async fn update_preserving_birth_versions(
&self,
new_entries: &[Arc<SuperfileEntry>],
entries_to_remove: &[Arc<SuperfileEntry>],
) -> Result<(ManifestSnapshot, Vec<EncodedPart>), ManifestError> {
self.update_inner(new_entries, entries_to_remove, true)
.await
}
async fn update_inner(
&self,
new_entries: &[Arc<SuperfileEntry>],
entries_to_remove: &[Arc<SuperfileEntry>],
preserve_birth_versions: bool,
) -> Result<(ManifestSnapshot, Vec<EncodedPart>), ManifestError> {
let opts = self.get_opts();
let strategy = self.get_partition_strategy();
let hidden_table = matches!(&strategy, PartitionStrategy::VectorCell { .. });
let birth_version = self.get_next_manifest_id();
let stamped_new_entries: Vec<Arc<SuperfileEntry>> = new_entries
.iter()
.map(|e| {
if !e.partition_key.is_empty() {
return Err(ManifestError::EntryAlreadyPartitioned {
detail: format!(
"superfile {} arrived with a partition_key already set",
e.superfile_id
),
});
}
let pk = assign_partition(e, &strategy)?;
let entry_birth_version = if preserve_birth_versions {
e.birth_version
} else {
birth_version
};
Ok(Arc::new(SuperfileEntry {
partition_key: encode_partition_key(&pk),
birth_version: entry_birth_version,
..(**e).clone()
}))
})
.collect::<Result<_, ManifestError>>()?;
let list_entries: Option<&[ManifestPartEntry]> =
if matches!(&strategy, PartitionStrategy::VectorCell { .. }) {
None
} else {
Some(self.get_all_list_entries())
};
let latest_idx = list_entries.and_then(|entries| entries.len().checked_sub(1));
let mut out_list_entries: Vec<ManifestPartEntry> = Vec::new();
let mut parts_to_write: Vec<EncodedPart> = Vec::new();
let mut pending_new = list_entries
.map(|_| stamped_new_entries.to_vec())
.unwrap_or_default();
for (i, entry) in list_entries.unwrap_or_default().iter().enumerate() {
if Some(i) != latest_idx || pending_new.is_empty() {
out_list_entries.push(entry.clone());
continue;
}
let new_for_part = std::mem::take(&mut pending_new);
let combined_n = entry.n_superfiles + new_for_part.len() as u64;
let latest_at_size_cap = entry.size_bytes_compressed
>= self.superfile_list.options.part_size_threshold_bytes;
if combined_n > self.superfile_list.options.target_superfiles_per_part
|| latest_at_size_cap
{
out_list_entries.push(entry.clone());
let (fresh_entry, fresh_encoded_part) =
rebuild_part_and_entry(vec![], new_for_part, None, hidden_table);
out_list_entries.push(fresh_entry);
parts_to_write.push(fresh_encoded_part);
} else {
let existing_part = self.get_part_by_id(entry.part_id).await?;
let (rebuilt_entry, rebuilt_encoded_part) = rebuild_part_and_entry(
existing_part.superfiles.clone(),
new_for_part,
Some(entry),
hidden_table,
);
out_list_entries.push(rebuilt_entry);
parts_to_write.push(rebuilt_encoded_part);
}
}
if !pending_new.is_empty() {
let (fresh_entry, fresh_encoded_part) =
rebuild_part_and_entry(vec![], pending_new, None, hidden_table);
out_list_entries.push(fresh_entry);
parts_to_write.push(fresh_encoded_part);
}
let mut out_list_entries_after_removal = Vec::new();
if entries_to_remove.is_empty() {
out_list_entries_after_removal = out_list_entries;
} else {
let removal_ids = entries_to_remove
.iter()
.map(|r| r.superfile_id)
.collect::<HashSet<_>>();
for entry in out_list_entries {
let (superfile_entries_in_part, existing_part_to_update) = if let Some(existing) =
parts_to_write
.iter_mut()
.find(|ep| ep.part.part_id == entry.part_id)
{
(existing.part.superfiles.clone(), Some(existing))
} else if let Ok(existing_part) = self.get_part_by_id(entry.part_id).await {
(existing_part.superfiles.clone(), None)
} else {
return Err(ManifestError::UnknownPartId(entry.part_id));
};
let final_superfile_entries = superfile_entries_in_part
.iter()
.filter(|s| !removal_ids.contains(&s.superfile_id))
.cloned()
.collect::<Vec<_>>();
if final_superfile_entries.len() == superfile_entries_in_part.len() {
out_list_entries_after_removal.push(entry);
continue;
}
let (fresh_entry, fresh_encoded_part) =
rebuild_part_and_entry(vec![], final_superfile_entries, None, hidden_table);
if let Some(existing) = existing_part_to_update {
*existing = fresh_encoded_part;
} else {
parts_to_write.push(fresh_encoded_part);
}
out_list_entries_after_removal.push(fresh_entry);
}
}
let ids_to_remove = entries_to_remove
.iter()
.map(|e| e.superfile_id)
.collect::<HashSet<_>>();
let mut tombstone_seqs = self
.list
.as_ref()
.map(|list| list.tombstone_seqs.clone())
.unwrap_or_default();
tombstone_seqs.retain(|id, _| !ids_to_remove.contains(id));
let opts_hash = options_hash::compute_options_hash(opts.as_ref(), &strategy);
let vector_columns: Vec<list::VectorColumnInfo> = opts
.vector_columns
.iter()
.map(|v| list::VectorColumnInfo {
column: v.column.clone(),
dim: v.dim,
n_cent: v.n_cent,
rot_seed: v.rot_seed,
metric: format!("{:?}", v.metric).to_lowercase(),
})
.collect();
let new_list = Manifest {
drained_ranges: self.get_drained_ranges(),
tombstone_seqs,
format_version: LIST_FORMAT_VERSION.into(),
manifest_id: self.get_next_manifest_id(),
options_hash: opts_hash,
schema: Vec::new(),
id_column: opts.id_column.clone(),
fts_columns: opts
.fts_columns
.iter()
.map(|f| list::FtsColumnInfo {
column: f.column.clone(),
})
.collect(),
vector_columns: opts
.vector_columns
.iter()
.map(|v| list::VectorColumnInfo {
column: v.column.clone(),
dim: v.dim,
n_cent: v.n_cent,
rot_seed: v.rot_seed,
metric: format!("{:?}", v.metric).to_lowercase(),
})
.collect(),
partition_strategy: strategy,
vector_index_storage_prefix: if matches!(
self.get_partition_strategy(),
list::PartitionStrategy::VectorCell { .. }
) {
None
} else {
self.stamp_vector_index_storage_prefix(&vector_columns)
},
global_vector_index: self.get_global_vector_index(),
deleted_user_ids_inline: self
.list
.as_ref()
.and_then(|l| l.deleted_user_ids_inline.clone()),
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: out_list_entries_after_removal,
};
let mut new_superfile_list = self
.get_all_superfiles()
.iter()
.chain(stamped_new_entries.iter())
.map(Arc::clone)
.collect::<Vec<_>>();
new_superfile_list.retain(|e| !ids_to_remove.contains(&e.superfile_id));
let new_superfile_list = SuperfileList {
manifest_id: self.get_next_manifest_id(),
options: self.get_opts(),
superfiles: new_superfile_list,
vector_index_storage_prefix: None,
};
let loader = opts.storage.as_ref().map(|storage| {
Arc::new(ManifestPartLoader::new_with_cache(
storage.clone(),
&new_list,
opts.manifest_disk_cache.clone(),
))
});
let live_part_ids: HashSet<_> = new_list.parts.iter().map(|e| e.part_id).collect();
let parts = DashMap::new();
for kv in self.parts.iter() {
if live_part_ids.contains(kv.key()) {
parts.insert(*kv.key(), kv.value().clone());
}
}
for part in parts_to_write.iter() {
let part = part.part.clone();
parts.insert(
part.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part)))),
);
}
let new_manifest = ManifestSnapshot {
superfile_list: new_superfile_list,
list: Some(new_list),
parts,
loader,
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
};
Ok((new_manifest, parts_to_write))
}
}
fn rebuild_part_and_entry(
old_superfiles: Vec<Arc<SuperfileEntry>>,
new_superfiles: Vec<Arc<SuperfileEntry>>,
base_part: Option<&ManifestPartEntry>,
hidden: bool,
) -> (ManifestPartEntry, EncodedPart) {
let aggregates = aggregates::compute(&new_superfiles, base_part);
let superfiles = old_superfiles
.into_iter()
.chain(new_superfiles)
.collect::<Vec<_>>();
let part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles,
};
let primary_mode = if hidden {
SummaryWireMode::RoutingOnly
} else {
SummaryWireMode::Full
};
let compressed = part::encode_with_mode(&part, primary_mode);
let size_compressed = compressed.len() as u64;
let content_hash = ContentHash::of(&compressed);
let size_uncompressed = frame_content_size(&compressed, size_compressed);
let routing = if hidden {
None
} else {
let routing_encoded = part::encode_with_mode(&part, SummaryWireMode::RoutingOnly);
let routing_hash = ContentHash::of(&routing_encoded);
Some((routing_encoded, routing_hash))
};
let entry = ManifestPartEntry {
part_id: part.part_id,
uri: part_uri(&content_hash),
n_superfiles: part.superfiles.len() as u64,
size_bytes_compressed: size_compressed,
size_bytes_uncompressed: size_uncompressed,
content_hash,
routing: routing.as_ref().map(|(_, hash)| RoutingRef {
uri: part_uri(hash),
content_hash: *hash,
}),
id_range: aggregates.id_range,
scalar_stats_agg: aggregates.scalar_stats_agg,
fts_summary_agg: aggregates.fts_summary_agg,
};
(
entry,
EncodedPart {
part,
encoded: compressed,
routing_encoded: routing.map(|(bytes, _)| bytes),
},
)
}
pub struct ManifestPartLoader {
storage: Arc<dyn StorageProvider>,
parts_index: HashMap<PartId, (ContentHash, String, Option<RoutingRef>)>,
manifest_disk_cache: Option<Arc<ManifestDiskCache>>,
prefer_routing: bool,
}
impl ManifestPartLoader {
pub fn new(storage: Arc<dyn StorageProvider>, list: &Manifest) -> Self {
Self::new_with_cache(storage, list, None)
}
pub fn new_with_cache(
storage: Arc<dyn StorageProvider>,
list: &Manifest,
manifest_disk_cache: Option<Arc<ManifestDiskCache>>,
) -> Self {
Self::new_with_cache_and_mode(storage, list, manifest_disk_cache, false)
}
pub fn new_with_cache_and_mode(
storage: Arc<dyn StorageProvider>,
list: &Manifest,
manifest_disk_cache: Option<Arc<ManifestDiskCache>>,
prefer_routing: bool,
) -> Self {
let mut idx = HashMap::with_capacity(list.parts.len());
for entry in &list.parts {
idx.insert(
entry.part_id,
(entry.content_hash, entry.uri.clone(), entry.routing.clone()),
);
}
Self {
storage,
parts_index: idx,
manifest_disk_cache,
prefer_routing,
}
}
pub async fn load(&self, part_id: PartId) -> Result<Arc<ManifestPart>, ManifestLoadError> {
self.load_with_form(part_id, self.prefer_routing).await
}
pub async fn load_full(&self, part_id: PartId) -> Result<Arc<ManifestPart>, ManifestLoadError> {
self.load_with_form(part_id, false).await
}
async fn load_with_form(
&self,
part_id: PartId,
prefer_routing: bool,
) -> Result<Arc<ManifestPart>, ManifestLoadError> {
let (full_hash, full_uri, routing) = self
.parts_index
.get(&part_id)
.ok_or(ManifestLoadError::PartNotInList { part_id })?;
let (expected_hash, uri) = match (prefer_routing, routing) {
(true, Some(routing)) => (&routing.content_hash, &routing.uri),
_ => (full_hash, full_uri),
};
if let Some(cache) = &self.manifest_disk_cache
&& let Some(bytes) = cache.get(expected_hash).await
{
let parsed = decode_part_off_thread(Bytes::from(bytes)).await?;
return Ok(Arc::new(parsed));
}
let (bytes, _) = self
.storage
.get(uri)
.await
.map_err(ManifestLoadError::Storage)?;
let parsed = verify_and_decode_part_off_thread(bytes.clone(), *expected_hash).await?;
if let Some(cache) = &self.manifest_disk_cache {
cache.put(*expected_hash, &bytes).await;
}
Ok(Arc::new(parsed))
}
}
pub(crate) struct UserCentroidCache {
pub(crate) manifest_id: u64,
pub(crate) cells: HashMap<(Uuid, String), Vec<(Option<u32>, Arc<Vec<f32>>)>>,
}
impl UserCentroidCache {
pub(crate) fn cell(
&self,
superfile_id: Uuid,
column: &str,
cell_id: Option<u32>,
) -> Option<Arc<Vec<f32>>> {
self.cells
.get(&(superfile_id, column.to_owned()))?
.iter()
.find(|(id, _)| *id == cell_id)
.map(|(_, fp32)| Arc::clone(fp32))
}
pub(crate) fn from_parts(manifest_id: u64, parts: &[Arc<ManifestPart>]) -> Self {
let mut cells: HashMap<(Uuid, String), Vec<(Option<u32>, Arc<Vec<f32>>)>> = HashMap::new();
for part in parts {
for entry in &part.superfiles {
for (column, summary) in &entry.vector_summary {
let list = cells
.entry((entry.superfile_id, column.clone()))
.or_default();
for cell in &summary.cells {
if cell.clusters.vectors_resident() && cell.clusters.n_cent > 0 {
list.push((cell.cell_id, Arc::new(cell.clusters.centroids.clone())));
}
}
}
}
}
Self { manifest_id, cells }
}
}
async fn decode_part_off_thread(bytes: Bytes) -> Result<ManifestPart, ManifestLoadError> {
match spawn_blocking(move || part::decode(&bytes)).await {
Ok(result) => Ok(result?),
Err(join_error) => Err(ManifestLoadError::Parse(part::PartParseError::Avro(
format!("part decode task failed: {join_error}"),
))),
}
}
async fn verify_and_decode_part_off_thread(
bytes: Bytes,
expected_hash: ContentHash,
) -> Result<ManifestPart, ManifestLoadError> {
let verify_then_decode = move || {
let actual_hash = ContentHash::of(&bytes);
if actual_hash != expected_hash {
return Err(ManifestLoadError::ContentHashMismatch {
expected: expected_hash.to_hex(),
actual: actual_hash.to_hex(),
});
}
part::decode(&bytes).map_err(ManifestLoadError::from)
};
match spawn_blocking(verify_then_decode).await {
Ok(result) => result,
Err(join_error) => Err(ManifestLoadError::Parse(part::PartParseError::Avro(
format!("part verify/decode task failed: {join_error}"),
))),
}
}
#[derive(Debug, thiserror::Error)]
pub enum ManifestLoadError {
#[error("pointer not found in storage")]
PointerNotFound,
#[error("already loaded")]
AlreadyLoaded,
#[error("pointer parse error: {0}")]
PointerParse(String),
#[error("no storage / loader attached to this manifest")]
NoLoaderAttached,
#[error("list parse error: {0}")]
ListParse(#[source] list::ListParseError),
#[error("part_id not in manifest list: {part_id}")]
PartNotInList { part_id: PartId },
#[error("storage error during part load: {0}")]
Storage(#[source] StorageError),
#[error("content-hash mismatch: expected {expected}, got {actual}")]
ContentHashMismatch { expected: String, actual: String },
#[error("part parse failed")]
Parse(#[from] part::PartParseError),
#[error("slow vector-state hydration failed: {0}")]
SlowStateHydration(String),
}
#[derive(Debug, Clone)]
pub struct SuperfileEntry {
pub superfile_id: Uuid,
pub uri: SuperfileUri,
pub n_docs: u64,
pub id_min: i128,
pub id_max: i128,
pub scalar_stats: HashMap<String, ScalarStatsAgg>,
pub fts_summary: HashMap<String, FtsSummaryAgg>,
pub vector_summary: HashMap<String, VectorSummary>,
pub partition_key: Vec<u8>,
pub partition_hint: Option<u32>,
pub subsection_offsets: Option<SubsectionOffsets>,
pub(crate) vector_layout: VectorLayout,
pub birth_version: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubsectionOffsets {
pub total_size: u64,
pub vec: Option<(u64, u64)>,
pub fts: Option<(u64, u64)>,
pub vec_open_ranges: Vec<(u64, u64)>,
pub fts_open_ranges: Vec<(u64, u64)>,
pub open_blob: Vec<(u64, Vec<u8>)>,
}
#[derive(Clone, Copy, Hash, Eq, PartialEq, Ord, PartialOrd, Debug)]
pub struct SuperfileUri(pub Uuid);
impl SuperfileUri {
pub fn new_v4() -> Self {
Self(Uuid::new_v4())
}
pub fn storage_path(self) -> String {
format!("{SUPERFILE_DATA_DIR}/seg-{}.sf.parquet", self.0)
}
pub fn cache_filename(self) -> String {
format!("seg-{}.sf.parquet", self.0)
}
pub fn cache_tmp_filename(self) -> String {
format!("seg-{}.sf.parquet.tmp", self.0)
}
pub fn from_cache_filename(name: &str) -> Option<Self> {
let body = name.strip_prefix("seg-")?.strip_suffix(".sf.parquet")?;
Uuid::parse_str(body).ok().map(SuperfileUri)
}
}
pub(crate) fn merge_min_max_arrays(
existing_min: &ArrayRef,
other_min: &ArrayRef,
existing_max: &ArrayRef,
other_max: &ArrayRef,
) -> Option<(ArrayRef, ArrayRef)> {
#[inline]
fn merge_opt<T: PartialOrd>(a: Option<T>, b: Option<T>, keep_min: bool) -> Option<T> {
match (a, b) {
(Some(a), Some(b)) => Some(if keep_min == (a <= b) { a } else { b }),
(Some(v), None) | (None, Some(v)) => Some(v),
(None, None) => None,
}
}
macro_rules! prim_merge {
($array_ty:ty) => {{
let exn = existing_min.as_any().downcast_ref::<$array_ty>()?;
let otn = other_min.as_any().downcast_ref::<$array_ty>()?;
let exx = existing_max.as_any().downcast_ref::<$array_ty>()?;
let otx = other_max.as_any().downcast_ref::<$array_ty>()?;
let at = |a: &$array_ty| (!a.is_null(0)).then(|| a.value(0));
Some((
Arc::new(<$array_ty>::from(vec![merge_opt(at(exn), at(otn), true)])) as ArrayRef,
Arc::new(<$array_ty>::from(vec![merge_opt(at(exx), at(otx), false)])) as ArrayRef,
))
}};
}
macro_rules! ts_merge {
($array_ty:ty, $tz:expr) => {{
let exn = existing_min.as_any().downcast_ref::<$array_ty>()?;
let otn = other_min.as_any().downcast_ref::<$array_ty>()?;
let exx = existing_max.as_any().downcast_ref::<$array_ty>()?;
let otx = other_max.as_any().downcast_ref::<$array_ty>()?;
let at = |a: &$array_ty| (!a.is_null(0)).then(|| a.value(0));
Some((
Arc::new(
<$array_ty>::from(vec![merge_opt(at(exn), at(otn), true)])
.with_timezone_opt($tz.clone()),
) as ArrayRef,
Arc::new(
<$array_ty>::from(vec![merge_opt(at(exx), at(otx), false)])
.with_timezone_opt($tz.clone()),
) as ArrayRef,
))
}};
}
match existing_min.data_type() {
DataType::UInt8 => prim_merge!(UInt8Array),
DataType::UInt16 => prim_merge!(UInt16Array),
DataType::UInt32 => prim_merge!(UInt32Array),
DataType::UInt64 => prim_merge!(UInt64Array),
DataType::Int8 => prim_merge!(Int8Array),
DataType::Int16 => prim_merge!(Int16Array),
DataType::Int32 => prim_merge!(Int32Array),
DataType::Int64 => prim_merge!(Int64Array),
DataType::Float32 => prim_merge!(Float32Array),
DataType::Float64 => prim_merge!(Float64Array),
DataType::Boolean => prim_merge!(BooleanArray),
DataType::Utf8 => {
let exn = existing_min.as_any().downcast_ref::<StringArray>()?;
let otn = other_min.as_any().downcast_ref::<StringArray>()?;
let exx = existing_max.as_any().downcast_ref::<StringArray>()?;
let otx = other_max.as_any().downcast_ref::<StringArray>()?;
let min = merge_opt(
(!exn.is_null(0)).then(|| exn.value(0)),
(!otn.is_null(0)).then(|| otn.value(0)),
true,
);
let max = merge_opt(
(!exx.is_null(0)).then(|| exx.value(0)),
(!otx.is_null(0)).then(|| otx.value(0)),
false,
);
Some((
Arc::new(StringArray::from(vec![min])),
Arc::new(StringArray::from(vec![max])),
))
}
DataType::LargeUtf8 => {
let exn = existing_min.as_any().downcast_ref::<LargeStringArray>()?;
let otn = other_min.as_any().downcast_ref::<LargeStringArray>()?;
let exx = existing_max.as_any().downcast_ref::<LargeStringArray>()?;
let otx = other_max.as_any().downcast_ref::<LargeStringArray>()?;
let min = merge_opt(
(!exn.is_null(0)).then(|| exn.value(0)),
(!otn.is_null(0)).then(|| otn.value(0)),
true,
);
let max = merge_opt(
(!exx.is_null(0)).then(|| exx.value(0)),
(!otx.is_null(0)).then(|| otx.value(0)),
false,
);
Some((
Arc::new(LargeStringArray::from(vec![min])),
Arc::new(LargeStringArray::from(vec![max])),
))
}
DataType::Decimal128(precision, scale) => {
let exn = existing_min.as_any().downcast_ref::<Decimal128Array>()?;
let otn = other_min.as_any().downcast_ref::<Decimal128Array>()?;
let exx = existing_max.as_any().downcast_ref::<Decimal128Array>()?;
let otx = other_max.as_any().downcast_ref::<Decimal128Array>()?;
let at = |a: &Decimal128Array| (!a.is_null(0)).then(|| a.value(0));
let min = merge_opt(at(exn), at(otn), true);
let max = merge_opt(at(exx), at(otx), false);
Some((
Arc::new(
Decimal128Array::from(vec![min])
.with_precision_and_scale(*precision, *scale)
.ok()?,
),
Arc::new(
Decimal128Array::from(vec![max])
.with_precision_and_scale(*precision, *scale)
.ok()?,
),
))
}
DataType::Date32 => prim_merge!(Date32Array),
DataType::Date64 => prim_merge!(Date64Array),
DataType::Time32(TimeUnit::Second) => prim_merge!(Time32SecondArray),
DataType::Time32(TimeUnit::Millisecond) => prim_merge!(Time32MillisecondArray),
DataType::Time64(TimeUnit::Microsecond) => prim_merge!(Time64MicrosecondArray),
DataType::Time64(TimeUnit::Nanosecond) => prim_merge!(Time64NanosecondArray),
DataType::Timestamp(TimeUnit::Second, tz) => ts_merge!(TimestampSecondArray, tz),
DataType::Timestamp(TimeUnit::Millisecond, tz) => ts_merge!(TimestampMillisecondArray, tz),
DataType::Timestamp(TimeUnit::Microsecond, tz) => ts_merge!(TimestampMicrosecondArray, tz),
DataType::Timestamp(TimeUnit::Nanosecond, tz) => ts_merge!(TimestampNanosecondArray, tz),
_ => None,
}
}
pub(crate) fn column_sum(col: &ArrayRef) -> Option<ArrayRef> {
macro_rules! signed {
($array_ty:ty) => {{
let a = col.as_any().downcast_ref::<$array_ty>()?;
let total: i128 = a.iter().flatten().map(i128::from).sum();
let v = i64::try_from(total).ok()?;
Some(Arc::new(Int64Array::from(vec![v])) as ArrayRef)
}};
}
macro_rules! unsigned {
($array_ty:ty) => {{
let a = col.as_any().downcast_ref::<$array_ty>()?;
let total: u128 = a.iter().flatten().map(u128::from).sum();
let v = u64::try_from(total).ok()?;
Some(Arc::new(UInt64Array::from(vec![v])) as ArrayRef)
}};
}
macro_rules! float {
($array_ty:ty) => {{
let a = col.as_any().downcast_ref::<$array_ty>()?;
let total: f64 = a.iter().flatten().map(f64::from).sum();
Some(Arc::new(Float64Array::from(vec![total])) as ArrayRef)
}};
}
match col.data_type() {
DataType::Int8 => signed!(Int8Array),
DataType::Int16 => signed!(Int16Array),
DataType::Int32 => signed!(Int32Array),
DataType::Int64 => signed!(Int64Array),
DataType::UInt8 => unsigned!(UInt8Array),
DataType::UInt16 => unsigned!(UInt16Array),
DataType::UInt32 => unsigned!(UInt32Array),
DataType::UInt64 => unsigned!(UInt64Array),
DataType::Float32 => float!(Float32Array),
DataType::Float64 => float!(Float64Array),
_ => None,
}
}
pub(crate) fn add_sum_arrays(a: &ArrayRef, b: &ArrayRef) -> Option<ArrayRef> {
match (a.data_type(), b.data_type()) {
(DataType::Int64, DataType::Int64) => {
let x = a.as_any().downcast_ref::<Int64Array>()?.value(0);
let y = b.as_any().downcast_ref::<Int64Array>()?.value(0);
Some(Arc::new(Int64Array::from(vec![x.checked_add(y)?])) as ArrayRef)
}
(DataType::UInt64, DataType::UInt64) => {
let x = a.as_any().downcast_ref::<UInt64Array>()?.value(0);
let y = b.as_any().downcast_ref::<UInt64Array>()?.value(0);
Some(Arc::new(UInt64Array::from(vec![x.checked_add(y)?])) as ArrayRef)
}
(DataType::Float64, DataType::Float64) => {
let x = a.as_any().downcast_ref::<Float64Array>()?.value(0);
let y = b.as_any().downcast_ref::<Float64Array>()?.value(0);
Some(Arc::new(Float64Array::from(vec![x + y])) as ArrayRef)
}
_ => None,
}
}
pub(crate) fn column_hll(col: &ArrayRef) -> Option<hll::HllSketch> {
let mut sketch = hll::HllSketch::new();
macro_rules! ints {
($array_ty:ty) => {{
let a = col.as_any().downcast_ref::<$array_ty>()?;
for v in a.iter().flatten() {
sketch.insert_hash(xxh3_64(&v.to_le_bytes()));
}
}};
}
match col.data_type() {
DataType::Int8 => ints!(Int8Array),
DataType::Int16 => ints!(Int16Array),
DataType::Int32 => ints!(Int32Array),
DataType::Int64 => ints!(Int64Array),
DataType::UInt8 => ints!(UInt8Array),
DataType::UInt16 => ints!(UInt16Array),
DataType::UInt32 => ints!(UInt32Array),
DataType::UInt64 => ints!(UInt64Array),
DataType::Float32 => {
let a = col.as_any().downcast_ref::<Float32Array>()?;
for v in a.iter().flatten() {
sketch.insert_hash(xxh3_64(&v.to_bits().to_le_bytes()));
}
}
DataType::Float64 => {
let a = col.as_any().downcast_ref::<Float64Array>()?;
for v in a.iter().flatten() {
sketch.insert_hash(xxh3_64(&v.to_bits().to_le_bytes()));
}
}
DataType::Utf8 => {
let a = col.as_any().downcast_ref::<StringArray>()?;
for v in a.iter().flatten() {
sketch.insert_hash(xxh3_64(v.as_bytes()));
}
}
DataType::LargeUtf8 => {
let a = col.as_any().downcast_ref::<LargeStringArray>()?;
for v in a.iter().flatten() {
sketch.insert_hash(xxh3_64(v.as_bytes()));
}
}
_ => return None,
}
Some(sketch)
}
pub(crate) fn column_min_max(col: &ArrayRef) -> Option<(ArrayRef, ArrayRef)> {
macro_rules! prim {
($array_ty:ty) => {{
let a = col.as_any().downcast_ref::<$array_ty>()?;
let mn_arr: ArrayRef = Arc::new(<$array_ty>::from(vec![agg::min(a)]));
let mx_arr: ArrayRef = Arc::new(<$array_ty>::from(vec![agg::max(a)]));
Some((mn_arr, mx_arr))
}};
}
macro_rules! ts {
($array_ty:ty, $tz:expr) => {{
let a = col.as_any().downcast_ref::<$array_ty>()?;
let mn_arr: ArrayRef =
Arc::new(<$array_ty>::from(vec![agg::min(a)]).with_timezone_opt($tz.clone()));
let mx_arr: ArrayRef =
Arc::new(<$array_ty>::from(vec![agg::max(a)]).with_timezone_opt($tz.clone()));
Some((mn_arr, mx_arr))
}};
}
match col.data_type() {
DataType::UInt8 => prim!(UInt8Array),
DataType::UInt16 => prim!(UInt16Array),
DataType::UInt32 => prim!(UInt32Array),
DataType::UInt64 => prim!(UInt64Array),
DataType::Int8 => prim!(Int8Array),
DataType::Int16 => prim!(Int16Array),
DataType::Int32 => prim!(Int32Array),
DataType::Int64 => prim!(Int64Array),
DataType::Float32 => prim!(Float32Array),
DataType::Float64 => prim!(Float64Array),
DataType::Boolean => {
let a = col.as_any().downcast_ref::<BooleanArray>()?;
Some((
Arc::new(BooleanArray::from(vec![agg::min_boolean(a)])),
Arc::new(BooleanArray::from(vec![agg::max_boolean(a)])),
))
}
DataType::Utf8 => {
let a = col.as_any().downcast_ref::<StringArray>()?;
Some((
Arc::new(StringArray::from(vec![agg::min_string(a)])),
Arc::new(StringArray::from(vec![agg::max_string(a)])),
))
}
DataType::LargeUtf8 => {
let a = col.as_any().downcast_ref::<LargeStringArray>()?;
Some((
Arc::new(LargeStringArray::from(vec![agg::min_string(a)])),
Arc::new(LargeStringArray::from(vec![agg::max_string(a)])),
))
}
DataType::Decimal128(precision, scale) => {
let a = col.as_any().downcast_ref::<Decimal128Array>()?;
Some((
Arc::new(
Decimal128Array::from(vec![agg::min(a)])
.with_precision_and_scale(*precision, *scale)
.ok()?,
),
Arc::new(
Decimal128Array::from(vec![agg::max(a)])
.with_precision_and_scale(*precision, *scale)
.ok()?,
),
))
}
DataType::Date32 => prim!(Date32Array),
DataType::Date64 => prim!(Date64Array),
DataType::Time32(TimeUnit::Second) => prim!(Time32SecondArray),
DataType::Time32(TimeUnit::Millisecond) => prim!(Time32MillisecondArray),
DataType::Time64(TimeUnit::Microsecond) => prim!(Time64MicrosecondArray),
DataType::Time64(TimeUnit::Nanosecond) => prim!(Time64NanosecondArray),
DataType::Timestamp(TimeUnit::Second, tz) => ts!(TimestampSecondArray, tz),
DataType::Timestamp(TimeUnit::Millisecond, tz) => ts!(TimestampMillisecondArray, tz),
DataType::Timestamp(TimeUnit::Microsecond, tz) => ts!(TimestampMicrosecondArray, tz),
DataType::Timestamp(TimeUnit::Nanosecond, tz) => ts!(TimestampNanosecondArray, tz),
_ => None,
}
}
#[derive(Debug, Clone)]
pub struct VectorSummary {
pub centroid: Vec<f32>,
pub cells: Vec<CellVectorSummary>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CellVectorSummary {
pub cell_id: Option<u32>,
pub clusters: ClusterCentroids,
}
pub(crate) const ADMIT_CODE_WORD_BITS: usize = 64;
pub(crate) const RABITQ_ADMIT_CELL_SHORTLIST_FRACTION: f64 = 0.20;
pub(crate) const RABITQ_ADMIT_CELL_SHORTLIST_MIN: usize = 48;
#[derive(Debug, Clone)]
pub(crate) struct RabitqAdmitContext {
rot_seed: u64,
dim: usize,
rotation: Arc<RandomRotation>,
quant: BitQuantizer,
cos_table: Arc<Vec<f32>>,
}
impl RabitqAdmitContext {
pub(crate) fn new(dim: usize, rot_seed: u64) -> Self {
Self {
rot_seed,
dim,
rotation: Arc::new(RandomRotation::new(dim, rot_seed)),
quant: BitQuantizer::new(dim),
cos_table: Arc::new(
(0..=dim)
.map(|h| (PI * h as f32 / dim as f32).cos())
.collect(),
),
}
}
pub(crate) fn encode(&self, vector: &[f32]) -> RabitqAdmitQuery {
debug_assert_eq!(vector.len(), self.dim);
let mut rotated = vec![0.0f32; self.dim];
self.rotation.apply(vector, &mut rotated);
let mut code = vec![0u8; self.quant.code_bytes()];
self.quant.encode_rotated_into(&rotated, &mut code);
let q_l2sq = dot(vector, vector);
RabitqAdmitQuery {
rot_seed: self.rot_seed,
rotation: Arc::clone(&self.rotation),
quant: self.quant.clone(),
q_words: pack_code_bytes_to_words(&code),
q_norm: q_l2sq.sqrt(),
q_l2sq,
cos_table: Arc::clone(&self.cos_table),
}
}
}
#[derive(Debug)]
pub(crate) struct RabitqAdmitQuery {
rot_seed: u64,
rotation: Arc<RandomRotation>,
quant: BitQuantizer,
q_words: Vec<u64>,
q_norm: f32,
q_l2sq: f32,
cos_table: Arc<Vec<f32>>,
}
impl RabitqAdmitQuery {
pub(crate) fn new(dim: usize, rot_seed: u64, query: &[f32]) -> Self {
RabitqAdmitContext::new(dim, rot_seed).encode(query)
}
}
fn admit_encoders(
vector_columns: &[VectorConfig],
) -> HashMap<&str, (RandomRotation, BitQuantizer, u64)> {
vector_columns
.iter()
.map(|vc| {
(
vc.column.as_str(),
(
RandomRotation::new(vc.dim, vc.rot_seed),
BitQuantizer::new(vc.dim),
vc.rot_seed,
),
)
})
.collect()
}
fn strip_summary_centroids(
superfiles: &mut [Arc<SuperfileEntry>],
vector_columns: &[VectorConfig],
) {
let encoders = admit_encoders(vector_columns);
if encoders.is_empty() {
return;
}
for entry in superfiles.iter_mut() {
let Some(entry) = Arc::get_mut(entry) else {
continue;
};
for (column, summary) in entry.vector_summary.iter_mut() {
let Some((rotation, quant, rot_seed)) = encoders.get(column.as_str()) else {
continue;
};
for cell in &mut summary.cells {
if cell.clusters.dim as usize != quant.dim {
continue;
}
cell.clusters
.strip_centroids_after_slab(rotation, quant, *rot_seed);
}
}
}
}
fn prewarm_summary_admit_slabs(
superfiles: &[Arc<SuperfileEntry>],
vector_columns: &[VectorConfig],
pool: &ThreadPool,
) {
let encoders = admit_encoders(vector_columns);
if encoders.is_empty() {
return;
}
pool.install(|| {
superfiles.par_iter().for_each(|entry| {
for (column, summary) in &entry.vector_summary {
let Some((rotation, quant, rot_seed)) = encoders.get(column.as_str()) else {
continue;
};
for cell in &summary.cells {
if cell.clusters.dim as usize != quant.dim || !cell.clusters.vectors_resident()
{
continue;
}
cell.clusters
.prewarm_admit_codes(rotation, quant, *rot_seed);
}
}
});
});
}
fn pack_code_bytes_to_words(code: &[u8]) -> Vec<u64> {
let bytes_per_word = ADMIT_CODE_WORD_BITS / 8;
code.chunks(bytes_per_word)
.map(|chunk| {
let mut word = [0u8; 8];
word[..chunk.len()].copy_from_slice(chunk);
u64::from_le_bytes(word)
})
.collect()
}
#[inline]
fn hamming_words(a: &[u64], b: &[u64]) -> u32 {
debug_assert_eq!(a.len(), b.len());
a.iter().zip(b).map(|(x, y)| (x ^ y).count_ones()).sum()
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) struct RabitqAdmitCodes {
pub(crate) rot_seed: u64,
pub(crate) words_per_code: usize,
pub(crate) codes: Vec<u64>,
pub(crate) norms: Vec<f32>,
}
#[derive(Debug, Default)]
pub struct ClusterCentroids {
pub n_cent: u32,
pub dim: u32,
pub centroids: Vec<f32>,
pub counts: Vec<u32>,
transposed: OnceLock<Vec<f32>>,
admit_codes: OnceLock<RabitqAdmitCodes>,
}
impl Clone for ClusterCentroids {
fn clone(&self) -> Self {
let transposed = OnceLock::new();
if let Some(cache) = self.transposed.get() {
let _ = transposed.set(cache.clone());
}
let admit_codes = OnceLock::new();
if let Some(cache) = self.admit_codes.get() {
let _ = admit_codes.set(cache.clone());
}
Self {
n_cent: self.n_cent,
dim: self.dim,
centroids: self.centroids.clone(),
counts: self.counts.clone(),
transposed,
admit_codes,
}
}
}
impl PartialEq for ClusterCentroids {
fn eq(&self, other: &Self) -> bool {
self.n_cent == other.n_cent
&& self.dim == other.dim
&& self.centroids == other.centroids
&& self.counts == other.counts
}
}
impl Eq for ClusterCentroids {}
impl ClusterCentroids {
pub fn empty() -> Self {
Self::default()
}
pub fn is_empty(&self) -> bool {
self.n_cent == 0
}
pub fn centroid(&self, c: usize) -> &[f32] {
assert!(
self.vectors_resident(),
"centroid() on a stripped summary — fp32 centroids are not resident"
);
let d = self.dim as usize;
let base = c * d;
&self.centroids[base..base + d]
}
pub fn to_fp32(&self) -> Vec<f32> {
self.centroids.clone()
}
pub fn from_fp32(n_cent: u32, dim: u32, centroids: &[f32], counts: Vec<u32>) -> Self {
let stored: Vec<f32> = centroids
.iter()
.map(|v| if v.is_finite() { *v } else { 0.0 })
.collect();
Self::from_decoded(n_cent, dim, stored, counts)
}
pub(crate) fn from_decoded(
n_cent: u32,
dim: u32,
centroids: Vec<f32>,
counts: Vec<u32>,
) -> Self {
debug_assert_eq!(
counts.len(),
n_cent as usize,
"cluster grid counts ({}) must match n_cent ({n_cent})",
counts.len()
);
debug_assert_eq!(
centroids.len(),
n_cent as usize * dim as usize,
"cluster grid centroids ({}) must be n_cent*dim ({}*{})",
centroids.len(),
n_cent,
dim
);
Self {
n_cent,
dim,
centroids,
counts,
transposed: OnceLock::new(),
admit_codes: OnceLock::new(),
}
}
pub(crate) fn transposed(&self) -> &[f32] {
assert!(
self.vectors_resident(),
"fp32 centroids were dropped (summary_centroids_from_superfiles); \
exact scans must read the superfile centroid regions"
);
self.transposed.get_or_init(|| {
transpose_centroids_cluster_major(
&self.centroids,
self.n_cent as usize,
self.dim as usize,
)
})
}
#[allow(dead_code)]
pub(crate) fn invalidate_transposed(&mut self) {
self.transposed = OnceLock::new();
self.admit_codes = OnceLock::new();
}
pub(crate) fn vectors_resident(&self) -> bool {
self.n_cent == 0 || !self.centroids.is_empty()
}
fn build_admit_codes(
&self,
rotation: &RandomRotation,
quant: &BitQuantizer,
rot_seed: u64,
) -> RabitqAdmitCodes {
assert!(
self.vectors_resident(),
"admit codes need resident fp32 centroids; this summary was stripped"
);
let dim = self.dim as usize;
let n_cent = self.n_cent as usize;
let words_per_code = dim.div_ceil(ADMIT_CODE_WORD_BITS);
let mut codes = vec![0u64; n_cent.saturating_mul(words_per_code)];
let mut norms = vec![0.0f32; n_cent];
let mut rotated = vec![0.0f32; dim];
let mut byte_code = vec![0u8; quant.code_bytes()];
for c in 0..n_cent {
let centroid = self.centroid(c);
norms[c] = dot(centroid, centroid).sqrt();
rotation.apply(centroid, &mut rotated);
quant.encode_rotated_into(&rotated, &mut byte_code);
codes[c * words_per_code..(c + 1) * words_per_code]
.copy_from_slice(&pack_code_bytes_to_words(&byte_code));
}
RabitqAdmitCodes {
rot_seed,
words_per_code,
codes,
norms,
}
}
pub(crate) fn strip_centroids_after_slab(
&mut self,
rotation: &RandomRotation,
quant: &BitQuantizer,
rot_seed: u64,
) {
if self.n_cent == 0 || self.centroids.is_empty() {
return;
}
let codes = self.build_admit_codes(rotation, quant, rot_seed);
self.admit_codes = OnceLock::new();
let _ = self.admit_codes.set(codes);
self.centroids = Vec::new();
self.transposed = OnceLock::new();
}
fn admit_codes(&self, admit: &RabitqAdmitQuery) -> &RabitqAdmitCodes {
let cache = self.admit_codes.get_or_init(|| {
self.build_admit_codes(admit.rotation.as_ref(), &admit.quant, admit.rot_seed)
});
assert_eq!(
cache.rot_seed, admit.rot_seed,
"admit codes built with a different rot_seed"
);
cache
}
pub(crate) fn prewarm_admit_codes(
&self,
rotation: &RandomRotation,
quant: &BitQuantizer,
rot_seed: u64,
) {
let _ = self
.admit_codes
.get_or_init(|| self.build_admit_codes(rotation, quant, rot_seed));
}
pub(crate) fn admit_codes_built(&self) -> Option<&RabitqAdmitCodes> {
self.admit_codes.get()
}
pub(crate) fn from_decoded_routing(
n_cent: u32,
dim: u32,
counts: Vec<u32>,
admit: RabitqAdmitCodes,
) -> Self {
debug_assert_eq!(
counts.len(),
n_cent as usize,
"routing cluster counts ({}) must match n_cent ({n_cent})",
counts.len()
);
Self {
n_cent,
dim,
centroids: Vec::new(),
counts,
transposed: OnceLock::new(),
admit_codes: {
let lock = OnceLock::new();
let _ = lock.set(admit);
lock
},
}
}
pub(crate) fn estimate_min_admit_score(
&self,
metric: Metric,
admit: &RabitqAdmitQuery,
) -> Option<f32> {
let mut best: Option<f32> = None;
self.estimate_admit_scores_into(metric, admit, |_, score| {
best = Some(best.map_or(score, |b: f32| b.min(score)));
});
best
}
pub(crate) fn admit_shortlist(
&self,
metric: Metric,
admit: &RabitqAdmitQuery,
window: usize,
) -> Vec<(u32, f32)> {
let mut top: Vec<(u32, f32)> = Vec::with_capacity(window.saturating_add(1));
self.estimate_admit_scores_into(metric, admit, |c, score| {
insert_ranked(&mut top, window, c, score);
});
top
}
pub(crate) fn estimate_admit_scores_into(
&self,
metric: Metric,
admit: &RabitqAdmitQuery,
mut emit: impl FnMut(u32, f32),
) {
debug_assert_eq!(admit.cos_table.len(), self.dim as usize + 1);
let cache = self.admit_codes(admit);
let w = cache.words_per_code;
for c in 0..self.n_cent as usize {
if self.counts[c] == 0 {
continue;
}
let code = &cache.codes[c * w..(c + 1) * w];
let h = hamming_words(&admit.q_words, code) as usize;
let est_dot = admit.cos_table[h] * admit.q_norm * cache.norms[c];
let score = match metric {
Metric::Cosine => COSINE_DISTANCE_BASE - est_dot,
Metric::NegDot => -est_dot,
Metric::L2Sq => {
let c_norm = cache.norms[c];
admit.q_l2sq + c_norm * c_norm - L2_CROSS_TERM_COEFF * est_dot
}
};
emit(c as u32, score);
}
}
pub fn score_one(&self, metric: Metric, c: usize, query: &[f32]) -> f32 {
debug_assert_eq!(query.len(), self.dim as usize);
distance(metric, query, self.centroid(c))
}
pub fn score_clusters_into(
&self,
metric: Metric,
query: &[f32],
mut emit: impl FnMut(u32, f32),
) {
debug_assert_eq!(query.len(), self.dim as usize);
let n_cent = self.n_cent as usize;
let scores = all_centroid_scores_transposed(
metric,
query,
self.transposed(),
n_cent,
self.dim as usize,
);
for (c, &score) in scores.iter().enumerate() {
if self.counts[c] == 0 {
continue;
}
emit(c as u32, score);
}
}
pub fn rank_cells(&self, metric: Metric, query: &[f32]) -> Vec<(u32, f32)> {
debug_assert_eq!(query.len(), self.dim as usize);
let n_cent = self.n_cent as usize;
let scores = all_centroid_scores_transposed(
metric,
query,
self.transposed(),
n_cent,
self.dim as usize,
);
let mut ranked: Vec<(u32, f32)> = scores
.into_iter()
.enumerate()
.map(|(c, score)| (c as u32, score))
.collect();
ranked.sort_unstable_by(|a, b| {
a.1.partial_cmp(&b.1)
.unwrap_or(Ordering::Equal)
.then_with(|| a.0.cmp(&b.0))
});
ranked
}
pub fn nearest_cell(&self, metric: Metric, query: &[f32]) -> u32 {
debug_assert_eq!(query.len(), self.dim as usize);
nearest_k_centroids_transposed(
metric,
query,
self.transposed(),
self.n_cent as usize,
self.dim as usize,
Some(&self.counts),
1,
)
.first()
.map(|&(cell, _)| cell)
.unwrap_or(0)
}
pub fn assign_rows(&self, metric: Metric, vectors: &[f32], assignments: &mut [u32]) {
let dim = self.dim as usize;
assert_eq!(vectors.len() % dim, 0, "assign_rows: vectors len mismatch");
let n = vectors.len() / dim;
assert_eq!(
assignments.len(),
n,
"assign_rows: assignments len mismatch"
);
if n == 0 {
return;
}
assignments
.par_iter_mut()
.enumerate()
.for_each(|(d, slot)| {
*slot = self.nearest_cell(metric, &vectors[d * dim..(d + 1) * dim]);
});
}
}
#[cfg(test)]
mod tests {
use std::{hint::black_box, slice::from_ref, sync::Arc, time::Instant};
use arrow_array::{
Array, Date32Array, Date64Array, Int64Array, Time64MicrosecondArray,
TimestampMicrosecondArray,
};
use arrow_schema::{DataType, Field, Schema, TimeUnit};
use dashmap::DashMap;
use datafusion::scalar::ScalarValue;
use tempfile::TempDir;
use tokio::sync::OnceCell;
use super::*;
use crate::{
storage::LocalFsStorageProvider,
superfile::{builder::FtsConfig, vector::distance::distance},
supertable::manifest::{
commit::{PartWriteResult, write_manifest_part},
list::{Manifest, PartitionStrategy},
},
test_helpers::default_tokenizer,
};
fn synth_clusters(n_cent: u32, dim: u32, seed: u64) -> (ClusterCentroids, Vec<f32>) {
let (nc, d) = (n_cent as usize, dim as usize);
let mut centroids = vec![0f32; nc * d];
for c in 0..nc {
for j in 0..d {
let v = ((seed + (c * d + j) as u64 * 2_654_435_761) % 1000) as f32 / 250.0 - 2.0
+ c as f32 * 0.1;
centroids[c * d + j] = v;
}
}
let counts: Vec<u32> = (0..nc).map(|c| if c == nc / 2 { 0 } else { 10 }).collect();
let cc = ClusterCentroids::from_fp32(n_cent, dim, ¢roids, counts);
(cc, centroids)
}
#[test]
fn min_max_stats_record_and_fold_all_null_columns() {
let arr = |vals: Vec<Option<i64>>| Arc::new(Int64Array::from(vals)) as ArrayRef;
let scalar = |a: &ArrayRef| ScalarValue::try_from_array(a, 0).expect("decode");
let (mn, mx) = column_min_max(&arr(vec![None, None])).expect("all-null stat");
assert!(mn.is_null(0) && mx.is_null(0));
let (mn, mx) = column_min_max(&arr(vec![Some(5), Some(2), None])).expect("stat");
assert_eq!(scalar(&mn), ScalarValue::Int64(Some(2)));
assert_eq!(scalar(&mx), ScalarValue::Int64(Some(5)));
let null1 = arr(vec![None]);
let (mn, mx) =
merge_min_max_arrays(&null1, &arr(vec![Some(2)]), &null1, &arr(vec![Some(9)]))
.expect("merge real over null");
assert_eq!(scalar(&mn), ScalarValue::Int64(Some(2)));
assert_eq!(scalar(&mx), ScalarValue::Int64(Some(9)));
let (mn, mx) = merge_min_max_arrays(&null1, &null1, &null1, &null1).expect("merge null");
assert!(mn.is_null(0) && mx.is_null(0));
}
#[test]
fn min_max_stats_cover_temporal_columns() {
let scalar = |a: &ArrayRef| ScalarValue::try_from_array(a, 0).expect("decode");
let d: ArrayRef = Arc::new(Date32Array::from(vec![Some(20100), Some(19000), None]));
let (mn, mx) = column_min_max(&d).expect("date stat");
assert_eq!(scalar(&mn), ScalarValue::Date32(Some(19000)));
assert_eq!(scalar(&mx), ScalarValue::Date32(Some(20100)));
let lo: ArrayRef = Arc::new(Date32Array::from(vec![Some(19000)]));
let hi: ArrayRef = Arc::new(Date32Array::from(vec![Some(20100)]));
let olo: ArrayRef = Arc::new(Date32Array::from(vec![Some(18000)]));
let ohi: ArrayRef = Arc::new(Date32Array::from(vec![Some(21000)]));
let (mn, mx) = merge_min_max_arrays(&lo, &olo, &hi, &ohi).expect("date merge");
assert_eq!(scalar(&mn), ScalarValue::Date32(Some(18000)));
assert_eq!(scalar(&mx), ScalarValue::Date32(Some(21000)));
let tz = "+05:30";
let ts: ArrayRef = Arc::new(
TimestampMicrosecondArray::from(vec![Some(200i64), Some(100)]).with_timezone(tz),
);
let zoned = DataType::Timestamp(TimeUnit::Microsecond, Some(tz.into()));
let (mn, mx) = column_min_max(&ts).expect("ts stat");
assert_eq!(mn.data_type(), &zoned);
assert_eq!(
scalar(&mn),
ScalarValue::TimestampMicrosecond(Some(100), Some(tz.into()))
);
assert_eq!(
scalar(&mx),
ScalarValue::TimestampMicrosecond(Some(200), Some(tz.into()))
);
let (mmn, _mmx) = merge_min_max_arrays(&mn, &mn, &mx, &mx).expect("ts merge keeps tz");
assert_eq!(mmn.data_type(), &zoned);
let d64: ArrayRef = Arc::new(Date64Array::from(vec![Some(9i64), Some(2)]));
assert_eq!(
scalar(&column_min_max(&d64).expect("date64 stat").0),
ScalarValue::Date64(Some(2))
);
let t64: ArrayRef = Arc::new(Time64MicrosecondArray::from(vec![Some(7i64), Some(3)]));
assert_eq!(
scalar(&column_min_max(&t64).expect("time64 stat").0),
ScalarValue::Time64Microsecond(Some(3))
);
}
#[test]
fn merge_min_max_folds_every_supported_column_type() {
let scalar = |a: &ArrayRef| ScalarValue::try_from_array(a, 0).expect("decode");
let fold =
|lo: ArrayRef, hi: ArrayRef| merge_min_max_arrays(&hi, &lo, &lo, &hi).expect("merge");
macro_rules! prim_case {
($arr:ty, $sv:path, $lo:expr, $hi:expr) => {{
let lo: ArrayRef = Arc::new(<$arr>::from(vec![Some($lo)]));
let hi: ArrayRef = Arc::new(<$arr>::from(vec![Some($hi)]));
let (mn, mx) = fold(lo, hi);
assert_eq!(
scalar(&mn),
$sv(Some($lo)),
concat!(stringify!($arr), " min")
);
assert_eq!(
scalar(&mx),
$sv(Some($hi)),
concat!(stringify!($arr), " max")
);
}};
}
prim_case!(UInt8Array, ScalarValue::UInt8, 1u8, 9u8);
prim_case!(UInt16Array, ScalarValue::UInt16, 1u16, 9u16);
prim_case!(UInt32Array, ScalarValue::UInt32, 1u32, 9u32);
prim_case!(UInt64Array, ScalarValue::UInt64, 1u64, 9u64);
prim_case!(Int8Array, ScalarValue::Int8, -3i8, 4i8);
prim_case!(Int16Array, ScalarValue::Int16, -3i16, 4i16);
prim_case!(Int32Array, ScalarValue::Int32, -3i32, 4i32);
prim_case!(Float32Array, ScalarValue::Float32, -1.5f32, 2.5f32);
prim_case!(Float64Array, ScalarValue::Float64, -1.5f64, 2.5f64);
prim_case!(Date64Array, ScalarValue::Date64, 2i64, 9i64);
prim_case!(Time32SecondArray, ScalarValue::Time32Second, 2i32, 9i32);
prim_case!(
Time32MillisecondArray,
ScalarValue::Time32Millisecond,
2i32,
9i32
);
prim_case!(
Time64NanosecondArray,
ScalarValue::Time64Nanosecond,
2i64,
9i64
);
let (mn, mx) = fold(
Arc::new(BooleanArray::from(vec![Some(false)])),
Arc::new(BooleanArray::from(vec![Some(true)])),
);
assert_eq!(scalar(&mn), ScalarValue::Boolean(Some(false)), "bool min");
assert_eq!(scalar(&mx), ScalarValue::Boolean(Some(true)), "bool max");
let (mn, mx) = fold(
Arc::new(StringArray::from(vec![Some("apple")])),
Arc::new(StringArray::from(vec![Some("pear")])),
);
assert_eq!(
scalar(&mn),
ScalarValue::Utf8(Some("apple".into())),
"utf8 min"
);
assert_eq!(
scalar(&mx),
ScalarValue::Utf8(Some("pear".into())),
"utf8 max"
);
let (mn, mx) = fold(
Arc::new(LargeStringArray::from(vec![Some("apple")])),
Arc::new(LargeStringArray::from(vec![Some("pear")])),
);
assert_eq!(
scalar(&mn),
ScalarValue::LargeUtf8(Some("apple".into())),
"largeutf8 min"
);
assert_eq!(
scalar(&mx),
ScalarValue::LargeUtf8(Some("pear".into())),
"largeutf8 max"
);
let dec = |v: i128| -> ArrayRef {
Arc::new(
Decimal128Array::from(vec![Some(v)])
.with_precision_and_scale(10, 2)
.expect("decimal"),
)
};
let (mn, mx) = fold(dec(125), dec(999));
assert_eq!(
scalar(&mn),
ScalarValue::Decimal128(Some(125), 10, 2),
"dec min"
);
assert_eq!(
scalar(&mx),
ScalarValue::Decimal128(Some(999), 10, 2),
"dec max"
);
let (mn, mx) = fold(
Arc::new(TimestampSecondArray::from(vec![Some(100i64)])),
Arc::new(TimestampSecondArray::from(vec![Some(200i64)])),
);
assert_eq!(
scalar(&mn),
ScalarValue::TimestampSecond(Some(100), None),
"ts-sec min"
);
assert_eq!(
scalar(&mx),
ScalarValue::TimestampSecond(Some(200), None),
"ts-sec max"
);
let (mn, _mx) = fold(
Arc::new(TimestampNanosecondArray::from(vec![Some(100i64)])),
Arc::new(TimestampNanosecondArray::from(vec![Some(200i64)])),
);
assert_eq!(
scalar(&mn),
ScalarValue::TimestampNanosecond(Some(100), None),
"ts-nano min"
);
}
#[test]
fn clone_preserves_warm_transposed_cache() {
let (cc, _) = synth_clusters(32, 128, 3);
let warm = cc.transposed().as_ptr();
assert!(cc.transposed.get().is_some());
let cloned = cc.clone();
assert!(
cloned.transposed.get().is_some(),
"clone must carry a warm transposed cache"
);
assert_eq!(cloned.transposed().len(), cc.transposed().len());
assert_ne!(
cloned.transposed().as_ptr(),
warm,
"clone owns its own cache buffer"
);
let mut mutated = cc.clone();
mutated.centroids[0] += 1.0;
mutated.invalidate_transposed();
assert!(mutated.transposed.get().is_none());
let _ = mutated.transposed(); assert!(mutated.transposed.get().is_some());
}
#[test]
fn admit_estimate_prefers_matching_centroid_instance() {
const DIM: usize = 128;
const ROT_SEED: u64 = 7;
let mut near = vec![0.0f32; DIM];
near[0] = 1.0;
let mut far = vec![0.0f32; DIM];
far[5] = 1.0;
let a = ClusterCentroids::from_fp32(1, DIM as u32, &near, vec![1]);
let b = ClusterCentroids::from_fp32(1, DIM as u32, &far, vec![1]);
let mut query = near.clone();
query[1] = 0.05;
let admit = RabitqAdmitQuery::new(DIM, ROT_SEED, &query);
for metric in [Metric::Cosine, Metric::L2Sq, Metric::NegDot] {
let ea = a
.estimate_min_admit_score(metric, &admit)
.expect("a populated");
let eb = b
.estimate_min_admit_score(metric, &admit)
.expect("b populated");
assert!(
ea < eb,
"{metric:?}: matching instance must rank first ({ea} vs {eb})"
);
}
let unpopulated = ClusterCentroids::from_fp32(1, DIM as u32, &near, vec![0]);
assert!(
unpopulated
.estimate_min_admit_score(Metric::Cosine, &admit)
.is_none()
);
let cloned = a.clone();
assert!(cloned.admit_codes.get().is_some());
}
#[test]
fn strip_centroids_keeps_slab_and_counts() {
const DIM: usize = 64;
const ROT_SEED: u64 = 7;
let mut flat = vec![0.0f32; 2 * DIM];
flat[0] = 1.0;
flat[DIM + 5] = 1.0;
let mut cc = ClusterCentroids::from_fp32(2, DIM as u32, &flat, vec![3, 4]);
let rotation = RandomRotation::new(DIM, ROT_SEED);
let quant = BitQuantizer::new(DIM);
cc.strip_centroids_after_slab(&rotation, &quant, ROT_SEED);
assert!(!cc.vectors_resident());
assert_eq!(cc.n_cent, 2);
assert_eq!(cc.counts, vec![3, 4]);
assert!(cc.centroids.is_empty());
let mut query = vec![0.0f32; DIM];
query[0] = 1.0;
let admit = RabitqAdmitQuery::new(DIM, ROT_SEED, &query);
assert!(
cc.estimate_min_admit_score(Metric::Cosine, &admit)
.is_some(),
"estimates must keep serving from the pre-built slab"
);
cc.strip_centroids_after_slab(&rotation, &quant, ROT_SEED);
assert!(!cc.vectors_resident());
let cloned = cc.clone();
assert!(!cloned.vectors_resident());
assert!(cloned.admit_codes.get().is_some());
}
#[test]
#[should_panic(expected = "fp32 centroids were dropped")]
fn transposed_on_stripped_summary_panics() {
const DIM: usize = 64;
const ROT_SEED: u64 = 7;
let mut flat = vec![0.0f32; DIM];
flat[0] = 1.0;
let mut cc = ClusterCentroids::from_fp32(1, DIM as u32, &flat, vec![1]);
cc.strip_centroids_after_slab(
&RandomRotation::new(DIM, ROT_SEED),
&BitQuantizer::new(DIM),
ROT_SEED,
);
let _ = cc.transposed();
}
#[test]
fn score_clusters_into_matches_centroid_distance() {
let (n_cent, dim) = (17u32, 96u32);
let (cc, centroids) = synth_clusters(n_cent, dim, 7);
let query: Vec<f32> = (0..dim)
.map(|j| ((j as u64 * 40_503 + 11) % 997) as f32 / 500.0 - 1.0)
.collect();
for metric in [Metric::Cosine, Metric::L2Sq, Metric::NegDot] {
let mut scored: Vec<(u32, f32)> = Vec::new();
cc.score_clusters_into(metric, &query, |c, s| {
scored.push((c, s));
});
let mut reference: Vec<(u32, f32)> = Vec::new();
for c in 0..n_cent as usize {
if cc.counts[c] == 0 {
continue;
}
let d = distance(metric, &query, cc.centroid(c));
assert_eq!(
cc.score_one(metric, c, &query),
d,
"{metric:?}: score_one c{c}"
);
reference.push((c as u32, d));
}
assert_eq!(
scored.len(),
reference.len(),
"{metric:?}: cluster sets differ (count-0 skip)"
);
for ((sc, ss), (rc, rs)) in scored.iter().zip(&reference) {
assert_eq!(sc, rc, "{metric:?}: cluster order");
assert!(
(ss - rs).abs() <= 1e-5 * (1.0 + rs.abs()),
"{metric:?} cluster {sc}: {ss} vs {rs}"
);
}
}
let roundtrip = cc.to_fp32();
for (i, (&got, &want)) in roundtrip.iter().zip(centroids.iter()).enumerate() {
assert_eq!(got, want, "roundtrip[{i}]: {got} vs {want}");
}
}
#[test]
#[ignore = "perf microbench, not a correctness gate"]
fn score_clusters_microbench() {
let (n_cent, dim) = (4096u32, 384u32);
let iters = 50usize;
let (cc, _) = synth_clusters(n_cent, dim, 99);
let query: Vec<f32> = (0..dim).map(|j| (j as f32).sin()).collect();
for metric in [Metric::Cosine, Metric::L2Sq] {
let t0 = Instant::now();
for _ in 0..iters {
let mut acc = 0f32;
cc.score_clusters_into(metric, &query, |_, s| acc += s);
black_box(acc);
}
let us = t0.elapsed().as_micros() as f64 / iters as f64;
println!("score_clusters {metric:?}: {us:.0} µs/query");
}
}
fn schema() -> Arc<Schema> {
Arc::new(Schema::new(vec![Field::new(
"title",
DataType::LargeUtf8,
false,
)]))
}
fn opts() -> Arc<SupertableOptions> {
let tk = default_tokenizer();
Arc::new(
SupertableOptions::new(
schema(),
vec![FtsConfig {
column: "title".into(),
positions: false,
}],
vec![],
Some(tk),
)
.expect("valid options"),
)
}
fn seg_entry(uuid: Uuid, n_docs: u64) -> Arc<SuperfileEntry> {
Arc::new(SuperfileEntry {
birth_version: 0,
superfile_id: uuid,
uri: SuperfileUri(uuid),
n_docs,
id_min: 0,
id_max: n_docs.saturating_sub(1) as i128,
scalar_stats: HashMap::new(),
fts_summary: HashMap::new(),
vector_summary: HashMap::new(),
partition_key: Vec::new(),
partition_hint: None,
vector_layout: VectorLayout::Ivf,
subsection_offsets: None,
})
}
#[test]
fn empty_manifest_starts_at_zero() {
let m = ManifestSnapshot::empty(opts());
assert_eq!(m.manifest_id, 0);
assert_eq!(m.superfiles.len(), 0);
assert_eq!(m.n_docs_total(), 0);
}
#[test]
fn with_appended_increments_manifest_id_and_extends_superfiles() {
let m0 = ManifestSnapshot::empty(opts());
let entry = seg_entry(Uuid::new_v4(), 100);
let m1 = m0.with_appended(vec![entry.clone()]);
assert_eq!(m1.manifest_id, 1);
assert_eq!(m1.superfiles.len(), 1);
assert_eq!(m1.n_docs_total(), 100);
assert_eq!(m0.manifest_id, 0);
assert_eq!(m0.superfiles.len(), 0);
assert_eq!(m0.n_docs_total(), 0);
}
#[test]
fn with_appended_chains_to_higher_manifest_ids() {
let m0 = ManifestSnapshot::empty(opts());
let m1 = m0.with_appended(vec![seg_entry(Uuid::new_v4(), 50)]);
let m2 = m1.with_appended(vec![seg_entry(Uuid::new_v4(), 75)]);
assert_eq!(m0.manifest_id, 0);
assert_eq!(m1.manifest_id, 1);
assert_eq!(m2.manifest_id, 2);
assert_eq!(m0.superfiles.len(), 0);
assert_eq!(m1.superfiles.len(), 1);
assert_eq!(m2.superfiles.len(), 2);
assert_eq!(m2.n_docs_total(), 50 + 75);
}
#[test]
fn with_appended_shares_old_superfiles_via_arc() {
let entry = seg_entry(Uuid::new_v4(), 1);
let m0 = ManifestSnapshot::empty(opts()).with_appended(vec![entry.clone()]);
let m1 = m0.with_appended(vec![seg_entry(Uuid::new_v4(), 2)]);
assert!(Arc::ptr_eq(&m0.superfiles[0], &m1.superfiles[0]));
}
#[test]
fn with_appended_empty_input_still_bumps_manifest_id() {
let m0 = ManifestSnapshot::empty(opts());
let m1 = m0.with_appended(vec![]);
assert_eq!(m1.manifest_id, 1);
assert_eq!(m1.superfiles.len(), 0);
}
#[test]
fn new_from_superfiles_builds_manifest_at_id_one_with_entries() {
let a = seg_entry(Uuid::new_v4(), 10);
let b = seg_entry(Uuid::new_v4(), 20);
let m = ManifestSnapshot::new_from_superfiles(opts(), vec![a.clone(), b.clone()]);
assert_eq!(m.manifest_id, 1);
assert_eq!(m.superfiles.len(), 2);
assert_eq!(m.n_docs_total(), 30);
assert!(Arc::ptr_eq(&m.superfiles[0], &a));
assert!(Arc::ptr_eq(&m.superfiles[1], &b));
assert!(m.is_in_process_only());
}
#[test]
fn new_from_superfiles_with_empty_input_is_empty_at_id_one() {
let m = ManifestSnapshot::new_from_superfiles(opts(), vec![]);
assert_eq!(m.manifest_id, 1);
assert_eq!(m.superfiles.len(), 0);
assert_eq!(m.n_docs_total(), 0);
}
#[test]
fn get_next_manifest_id_is_current_plus_one() {
let m0 = ManifestSnapshot::empty(opts());
assert_eq!(m0.get_manifest_id(), 0);
assert_eq!(m0.get_next_manifest_id(), 1);
let m1 = m0.with_appended(vec![seg_entry(Uuid::new_v4(), 1)]);
assert_eq!(m1.get_manifest_id(), 1);
assert_eq!(m1.get_next_manifest_id(), 2);
}
#[test]
fn get_next_manifest_id_is_a_pure_read() {
let m = ManifestSnapshot::empty(opts());
let _ = m.get_next_manifest_id();
assert_eq!(m.get_manifest_id(), 0, "current id unchanged");
assert_eq!(m.get_next_manifest_id(), m.get_next_manifest_id());
}
#[test]
fn superfile_uri_is_distinct_per_call() {
let a = SuperfileUri::new_v4();
let b = SuperfileUri::new_v4();
assert_ne!(a, b);
}
mod lazy_load {
use std::{
collections::HashMap,
error::Error,
ops::Range,
slice::from_ref,
sync::{
Arc,
atomic::{AtomicUsize, Ordering},
},
time::SystemTime,
};
use arrow_schema::{DataType, Field, Schema};
use async_trait::async_trait;
use bytes::Bytes;
use dashmap::DashMap;
use tokio::spawn;
use uuid::Uuid;
use super::super::*;
use crate::{
storage::{ObjectMeta, StorageError, StorageProvider},
supertable::{
SupertableOptions,
manifest::{
list::{FORMAT_VERSION as LIST_FORMAT_VERSION, Manifest, PartitionStrategy},
part::{self as part_mod, ContentHash, ManifestPart, PartId},
},
},
};
#[derive(Debug)]
struct CountingMockStorage {
objects: HashMap<String, Bytes>,
get_calls: AtomicUsize,
}
impl CountingMockStorage {
fn new(objects: HashMap<String, Bytes>) -> Self {
Self {
objects,
get_calls: AtomicUsize::new(0),
}
}
fn get_call_count(&self) -> usize {
self.get_calls.load(Ordering::Acquire)
}
}
#[async_trait]
impl StorageProvider for CountingMockStorage {
async fn head(&self, uri: &str) -> Result<ObjectMeta, StorageError> {
match self.objects.get(uri) {
Some(b) => Ok(ObjectMeta {
size: b.len() as u64,
etag: Some("mock-etag".into()),
last_modified: SystemTime::UNIX_EPOCH,
}),
None => Err(StorageError::NotFound { uri: uri.into() }),
}
}
async fn get(&self, uri: &str) -> Result<(Bytes, ObjectMeta), StorageError> {
self.get_calls.fetch_add(1, Ordering::AcqRel);
match self.objects.get(uri) {
Some(b) => Ok((
b.clone(),
ObjectMeta {
size: b.len() as u64,
etag: Some("mock-etag".into()),
last_modified: SystemTime::UNIX_EPOCH,
},
)),
None => Err(StorageError::NotFound { uri: uri.into() }),
}
}
async fn get_range(
&self,
uri: &str,
_range: Range<u64>,
) -> Result<Bytes, StorageError> {
Err(permanent(uri, "get_range unimplemented for mock"))
}
async fn put_atomic(
&self,
uri: &str,
_bytes: Bytes,
) -> Result<Option<String>, StorageError> {
Err(permanent(uri, "put_atomic unimplemented for mock"))
}
async fn put_if_match(
&self,
uri: &str,
_bytes: Bytes,
_expected_etag: Option<&str>,
) -> Result<Option<String>, StorageError> {
Err(permanent(uri, "put_if_match unimplemented for mock"))
}
async fn put_multipart(
&self,
uri: &str,
) -> Result<Box<dyn object_store::MultipartUpload>, StorageError> {
Err(permanent(uri, "put_multipart unimplemented for mock"))
}
async fn delete(&self, _uri: &str) -> Result<(), StorageError> {
Ok(())
}
}
fn permanent(uri: &str, msg: &'static str) -> StorageError {
let boxed: Box<dyn Error + Send + Sync> = msg.into();
StorageError::Permanent {
uri: uri.into(),
source: boxed,
}
}
fn make_test_part(seed: u8) -> ManifestPart {
ManifestPart {
format_version: part_mod::FORMAT_VERSION.into(),
part_id: PartId(Uuid::from_bytes([seed; 16])),
superfiles: vec![],
}
}
fn encode_and_index(
parts: &[ManifestPart],
) -> (HashMap<String, Bytes>, Vec<ManifestPartEntry>) {
let mut objects = HashMap::new();
let mut entries = Vec::new();
for p in parts {
let bytes = part_mod::encode(p);
let hash = ContentHash::of(&bytes);
let uri = format!("manifests/part-{}.avro.zst", hash.to_hex());
let size_compressed = bytes.len() as u64;
objects.insert(uri.clone(), Bytes::from(bytes));
entries.push(ManifestPartEntry {
part_id: p.part_id,
uri,
n_superfiles: p.superfiles.len() as u64,
size_bytes_compressed: size_compressed,
size_bytes_uncompressed: size_compressed,
content_hash: hash,
routing: None,
id_range: (0, 0),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
});
}
(objects, entries)
}
fn fresh_list(entries: Vec<ManifestPartEntry>) -> Manifest {
Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: LIST_FORMAT_VERSION.into(),
manifest_id: 1,
options_hash: ContentHash([0u8; 32]),
schema: Vec::new(),
id_column: "doc_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "doc_id".into(),
n_buckets: 64,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: entries,
}
}
fn options_for_test() -> Arc<SupertableOptions> {
let s = Arc::new(Schema::new(vec![Field::new(
"title",
DataType::LargeUtf8,
false,
)]));
Arc::new(SupertableOptions::new(s, vec![], vec![], None).expect("opts"))
}
fn build_manifest_with_loader(
list: Manifest,
storage: Arc<dyn StorageProvider>,
) -> ManifestSnapshot {
let loader = Arc::new(ManifestPartLoader::new(Arc::clone(&storage), &list));
ManifestSnapshot {
superfile_list: SuperfileList::empty(options_for_test()),
list: Some(list),
parts: DashMap::new(),
loader: Some(loader),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
}
}
#[tokio::test]
async fn part_first_touch_loads_and_caches() {
let part = make_test_part(7);
let (objects, entries) = encode_and_index(from_ref(&part));
let storage = Arc::new(CountingMockStorage::new(objects));
let list = fresh_list(entries);
let manifest =
build_manifest_with_loader(list, Arc::clone(&storage) as Arc<dyn StorageProvider>);
let loaded = manifest.get_part_by_id(part.part_id).await.expect("load");
assert_eq!(loaded.part_id, part.part_id);
assert_eq!(storage.get_call_count(), 1, "exactly one storage.get");
}
#[tokio::test]
async fn second_touch_hits_cache_zero_additional_gets() {
let part = make_test_part(11);
let (objects, entries) = encode_and_index(from_ref(&part));
let storage = Arc::new(CountingMockStorage::new(objects));
let list = fresh_list(entries);
let manifest =
build_manifest_with_loader(list, Arc::clone(&storage) as Arc<dyn StorageProvider>);
let a = manifest
.get_part_by_id(part.part_id)
.await
.expect("first load");
let b = manifest
.get_part_by_id(part.part_id)
.await
.expect("second load");
assert!(Arc::ptr_eq(&a, &b), "second touch must return cached Arc");
assert_eq!(storage.get_call_count(), 1, "cache hit ⇒ no extra get");
}
#[tokio::test]
async fn concurrent_loaders_coalesce_to_one_get() {
let part = make_test_part(13);
let (objects, entries) = encode_and_index(from_ref(&part));
let storage = Arc::new(CountingMockStorage::new(objects));
let list = fresh_list(entries);
let manifest = Arc::new(build_manifest_with_loader(
list,
Arc::clone(&storage) as Arc<dyn StorageProvider>,
));
let mut handles = Vec::with_capacity(100);
for _ in 0..100 {
let m = Arc::clone(&manifest);
let pid = part.part_id;
handles.push(spawn(async move { m.get_part_by_id(pid).await }));
}
let mut first: Option<Arc<ManifestPart>> = None;
for h in handles {
let p = h.await.expect("join").expect("load");
match &first {
None => first = Some(p),
Some(f) => assert!(
Arc::ptr_eq(f, &p),
"all concurrent loaders must share the same Arc"
),
}
}
assert_eq!(
storage.get_call_count(),
1,
"100 concurrent loaders on cold cell ⇒ exactly one storage.get"
);
}
#[tokio::test]
async fn content_hash_mismatch_surfaces_typed_error_without_refetch() {
let part = make_test_part(17);
let (mut objects, entries) = encode_and_index(from_ref(&part));
let bytes = objects.values().next().expect("one obj").clone();
let mut tampered = bytes.to_vec();
let last = tampered.len() - 1;
tampered[last] ^= 0xff;
let uri = entries[0].uri.clone();
objects.insert(uri, Bytes::from(tampered));
let (_, fresh_entries) = encode_and_index(from_ref(&part));
let list = fresh_list(fresh_entries);
let storage = Arc::new(CountingMockStorage::new(objects));
let manifest =
build_manifest_with_loader(list, Arc::clone(&storage) as Arc<dyn StorageProvider>);
let err = manifest
.get_part_by_id(part.part_id)
.await
.expect_err("must reject tampered bytes");
assert!(
matches!(err, ManifestLoadError::ContentHashMismatch { .. }),
"expected ContentHashMismatch, got {err:?}"
);
let _pre = storage.get_call_count();
let err2 = manifest
.get_part_by_id(part.part_id)
.await
.expect_err("must reject on retry too");
assert!(matches!(
err2,
ManifestLoadError::ContentHashMismatch { .. }
));
}
#[tokio::test]
async fn part_id_not_in_list_surfaces_typed_error() {
let part = make_test_part(19);
let (objects, entries) = encode_and_index(&[part]);
let storage = Arc::new(CountingMockStorage::new(objects));
let list = fresh_list(entries);
let manifest =
build_manifest_with_loader(list, Arc::clone(&storage) as Arc<dyn StorageProvider>);
let stranger = PartId(Uuid::from_bytes([0xff; 16]));
let err = manifest
.get_part_by_id(stranger)
.await
.expect_err("must reject");
assert!(
matches!(err, ManifestLoadError::PartNotInList { .. }),
"expected PartNotInList, got {err:?}"
);
assert_eq!(
storage.get_call_count(),
0,
"missing-id check happens before any storage.get"
);
}
#[tokio::test]
async fn disk_cache_hit_serves_second_loader_without_storage_get() {
let part = make_test_part(23);
let (objects, entries) = encode_and_index(from_ref(&part));
let storage = Arc::new(CountingMockStorage::new(objects));
let storage_dyn = Arc::clone(&storage) as Arc<dyn StorageProvider>;
let list = fresh_list(entries);
let cache_root = std::env::temp_dir()
.join("infino-manifest-cache-loader-test-disk_cache_hit_second_loader");
let _ = std::fs::remove_dir_all(&cache_root);
let cache = ManifestDiskCache::new(cache_root.clone(), 1 << 20).expect("cache");
let loader_a = ManifestPartLoader::new_with_cache(
Arc::clone(&storage_dyn),
&list,
Some(Arc::clone(&cache)),
);
let a = loader_a.load(part.part_id).await.expect("first load");
assert_eq!(a.part_id, part.part_id);
assert_eq!(storage.get_call_count(), 1, "first loader fetches once");
assert_eq!(cache.stats().n_entries, 1, "part bytes cached on disk");
let loader_b = ManifestPartLoader::new_with_cache(
Arc::clone(&storage_dyn),
&list,
Some(Arc::clone(&cache)),
);
let b = loader_b.load(part.part_id).await.expect("second load");
assert_eq!(b.part_id, part.part_id);
assert_eq!(
storage.get_call_count(),
1,
"disk-cache hit ⇒ no additional storage.get"
);
assert!(cache.stats().n_hits >= 1, "recorded a cache hit");
let _ = std::fs::remove_dir_all(&cache_root);
}
#[tokio::test]
async fn loader_without_cache_always_hits_storage() {
let part = make_test_part(29);
let (objects, entries) = encode_and_index(from_ref(&part));
let storage = Arc::new(CountingMockStorage::new(objects));
let storage_dyn = Arc::clone(&storage) as Arc<dyn StorageProvider>;
let list = fresh_list(entries);
let loader = ManifestPartLoader::new(Arc::clone(&storage_dyn), &list);
loader.load(part.part_id).await.expect("load 1");
loader.load(part.part_id).await.expect("load 2");
assert_eq!(
storage.get_call_count(),
2,
"no cache ⇒ every load round-trips to storage"
);
}
#[tokio::test]
async fn no_loader_attached_surfaces_typed_error() {
let manifest = ManifestSnapshot::empty(options_for_test());
let err = manifest
.get_part_by_id(PartId(Uuid::nil()))
.await
.expect_err("must error");
assert!(
matches!(err, ManifestLoadError::NoLoaderAttached),
"expected NoLoaderAttached, got {err:?}"
);
}
}
#[test]
fn superfile_uri_path_helpers_share_the_same_uuid() {
let uri = SuperfileUri(Uuid::from_u128(0x1234_5678));
let id = uri.0;
assert_eq!(uri.storage_path(), format!("data/seg-{id}.sf.parquet"));
assert_eq!(uri.cache_filename(), format!("seg-{id}.sf.parquet"));
assert_eq!(uri.cache_tmp_filename(), format!("seg-{id}.sf.parquet.tmp"));
}
#[test]
fn manifest_debug_reports_counts() {
let m = ManifestSnapshot::empty(opts()).with_appended(vec![seg_entry(Uuid::new_v4(), 3)]);
let dbg = format!("{m:?}");
assert!(dbg.contains("ManifestSnapshot"));
assert!(dbg.contains("manifest_id"));
assert!(dbg.contains("n_superfiles"));
assert!(dbg.contains("has_loader"));
}
#[test]
fn manifest_debug_with_list_reports_part_count() {
use list::{Manifest, PartitionStrategy};
let entry = part::PartId::new_v4();
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 1,
options_hash: part::ContentHash([0u8; 32]),
schema: Vec::new(),
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![list::ManifestPartEntry {
part_id: entry,
uri: "manifests/part-x".into(),
n_superfiles: 0,
size_bytes_compressed: 0,
size_bytes_uncompressed: 0,
content_hash: part::ContentHash([0u8; 32]),
routing: None,
id_range: (0, 0),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let m = ManifestSnapshot {
superfile_list: SuperfileList::empty(opts()),
list: Some(list),
parts: DashMap::new(),
loader: None,
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
};
let dbg = format!("{m:?}");
assert!(dbg.contains("n_parts: 1"), "{dbg}");
assert!(dbg.contains("has_list: true"), "{dbg}");
}
#[test]
fn cluster_centroids_empty_is_empty_and_default_matches() {
let cc = ClusterCentroids::empty();
assert!(cc.is_empty());
assert_eq!(cc.n_cent, 0);
let cc = ClusterCentroids::from_fp32(2, 4, &[0.0; 8], vec![1, 1]);
assert!(!cc.is_empty());
assert_eq!(cc.n_cent, 2);
assert_eq!(cc.dim, 4);
}
#[test]
fn add_sum_arrays_handles_each_type_and_overflow() {
use arrow_array::{Float64Array, Int64Array, UInt64Array};
let r = add_sum_arrays(
&(Arc::new(Int64Array::from(vec![3])) as ArrayRef),
&(Arc::new(Int64Array::from(vec![4])) as ArrayRef),
)
.expect("int sum");
assert_eq!(
r.as_any()
.downcast_ref::<Int64Array>()
.expect("test")
.value(0),
7
);
let r = add_sum_arrays(
&(Arc::new(UInt64Array::from(vec![3u64])) as ArrayRef),
&(Arc::new(UInt64Array::from(vec![4u64])) as ArrayRef),
)
.expect("uint sum");
assert_eq!(
r.as_any()
.downcast_ref::<UInt64Array>()
.expect("test")
.value(0),
7
);
let r = add_sum_arrays(
&(Arc::new(Float64Array::from(vec![1.5])) as ArrayRef),
&(Arc::new(Float64Array::from(vec![2.5])) as ArrayRef),
)
.expect("float sum");
assert!(
(r.as_any()
.downcast_ref::<Float64Array>()
.expect("test")
.value(0)
- 4.0)
.abs()
< 1e-9
);
let r = add_sum_arrays(
&(Arc::new(Int64Array::from(vec![i64::MAX])) as ArrayRef),
&(Arc::new(Int64Array::from(vec![1])) as ArrayRef),
);
assert!(r.is_none(), "i64 overflow drops the stat");
let r = add_sum_arrays(
&(Arc::new(Int64Array::from(vec![1])) as ArrayRef),
&(Arc::new(UInt64Array::from(vec![1u64])) as ArrayRef),
);
assert!(r.is_none(), "type mismatch drops the stat");
}
fn make_entry(docs: u64, pk: Vec<u8>, hint: Option<u32>) -> Arc<SuperfileEntry> {
Arc::new(SuperfileEntry {
birth_version: 0,
superfile_id: uuid::Uuid::new_v4(),
uri: SuperfileUri::new_v4(),
n_docs: docs,
id_min: 0,
id_max: docs as i128 - 1,
scalar_stats: Default::default(),
fts_summary: Default::default(),
vector_summary: Default::default(),
partition_key: pk,
partition_hint: hint,
vector_layout: VectorLayout::Ivf,
subsection_offsets: None,
})
}
fn make_superfile_entry(docs: u64, pk: Vec<u8>) -> Arc<SuperfileEntry> {
make_entry(docs, pk, None)
}
fn make_new_entry(docs: u64) -> Arc<SuperfileEntry> {
make_entry(docs, Vec::new(), None)
}
fn hash_bucket_0_pk() -> Vec<u8> {
vec![0, 0, 0, 0]
}
fn simple_schema() -> Arc<Schema> {
Arc::new(Schema::new(vec![Field::new(
"text",
DataType::LargeUtf8,
false,
)]))
}
fn make_opts() -> Arc<SupertableOptions> {
SupertableOptions::new(simple_schema(), vec![], vec![], None)
.map(Arc::new)
.expect("valid options")
}
fn empty_manifest(opts: &Arc<SupertableOptions>) -> Arc<ManifestSnapshot> {
Arc::new(ManifestSnapshot {
superfile_list: SuperfileList::empty(opts.clone()),
list: Some(Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![],
}),
parts: DashMap::new(),
loader: None,
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
})
}
#[tokio::test]
async fn slow_vector_state_ref_stamp_clear_and_preserve() {
let opts = make_opts();
let manifest = empty_manifest(&opts);
assert!(manifest.slow_vector_state_blob().is_none());
let hash = ContentHash([3u8; 32]);
let centroids = RoutingRef {
uri: "slow-vector-state/state-c.bin".into(),
content_hash: ContentHash([5u8; 32]),
};
let stamped = manifest.with_slow_vector_state(
"slow-vector-state/state-x.bin".into(),
hash,
centroids.clone(),
);
let (uri, got_hash) = stamped.slow_vector_state_blob().expect("ref stamped");
assert_eq!(uri, "slow-vector-state/state-x.bin");
assert_eq!(got_hash, hash);
assert_eq!(
stamped.slow_vector_state_centroids_blob(),
Some(¢roids),
"centroid-section sibling stamped with the state ref"
);
assert_eq!(stamped.get_manifest_id(), manifest.get_next_manifest_id());
assert_eq!(
stamped.get_all_superfiles().len(),
manifest.get_all_superfiles().len(),
"stamp must not change membership"
);
let deleted_stamped = stamped.with_deleted_user_ids(Vec::new());
assert!(
deleted_stamped.slow_vector_state_blob().is_some(),
"deleted-ids stamp must preserve the slow-state ref"
);
assert!(
deleted_stamped.slow_vector_state_centroids_blob().is_some(),
"deleted-ids stamp must preserve the centroid-section ref"
);
let new_entry = make_new_entry(100);
let (updated, _parts) = stamped
.update(from_ref(&new_entry), &[])
.await
.expect("update");
assert!(
updated.slow_vector_state_blob().is_none(),
"membership change must clear the slow-state ref"
);
assert!(
updated.slow_vector_state_centroids_blob().is_none(),
"membership change must clear the centroid-section ref"
);
}
#[tokio::test]
async fn undrained_load_trusts_resident_view_with_slow_state_ref() {
let opts = make_opts();
let entries = vec![
make_superfile_entry(100, hash_bucket_0_pk()),
make_superfile_entry(50, hash_bucket_0_pk()),
];
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 1,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: Some("slow-vector-state/state-abc.bin".into()),
slow_vector_state_content_hash: Some(ContentHash([7u8; 32])),
slow_vector_state_centroids: None,
parts: Vec::new(),
};
let (_dir, storage) = local_storage();
let manifest = ManifestSnapshot::new(1, opts, entries.clone(), Some(storage), Some(list));
let drained = list::DrainedVersionRanges::default();
let got = manifest
.get_undrained_superfiles_loaded(&drained)
.await
.expect("undrained load");
assert_eq!(
got.len(),
entries.len(),
"resident blob-hydrated membership must be returned, not the empty part fan"
);
}
#[tokio::test]
async fn update_fresh_start_cold_partition_should_create_entry() {
let opts = make_opts();
let old_manifest = empty_manifest(&opts);
let new_entry = make_new_entry(100);
let new_entries = vec![new_entry];
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 1);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, 1);
assert_eq!(parts[0].part.superfiles.len(), 1);
assert_eq!(parts[0].part.superfiles[0].n_docs, 100);
}
#[tokio::test]
async fn update_fresh_start_multiple_cold_partitions_should_create_entries() {
let opts = make_opts();
let old_manifest = empty_manifest(&opts);
let entry1 = make_new_entry(100);
let entry2 = make_new_entry(200);
let new_entries = vec![entry1, entry2];
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 1);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, 2);
assert_eq!(parts[0].part.superfiles.len(), 2);
let total_docs: u64 = parts[0].part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(total_docs, 300);
}
fn local_storage() -> (TempDir, Arc<dyn StorageProvider>) {
let dir = TempDir::new().expect("tempdir");
let store: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("local"));
(dir, store)
}
fn n_parts_initialized(m: &ManifestSnapshot) -> usize {
m.parts
.iter()
.filter(|kv| kv.value().get().is_some())
.count()
}
async fn persist_two_entry_table(
storage: &Arc<dyn StorageProvider>,
slow_ref: Option<(String, ContentHash)>,
) -> Vec<Arc<SuperfileEntry>> {
let entries = vec![
make_superfile_entry(100, hash_bucket_0_pk()),
make_superfile_entry(50, hash_bucket_0_pk()),
];
let part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: entries.clone(),
};
let pw = write_manifest_part(storage.as_ref(), &part)
.await
.expect("write part");
let (slow_uri, slow_hash) = match slow_ref {
Some((u, h)) => (Some(u), Some(h)),
None => (None, None),
};
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 1,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: slow_uri,
slow_vector_state_content_hash: slow_hash,
slow_vector_state_centroids: None,
parts: vec![ManifestPartEntry {
part_id: pw.part_id,
uri: pw.uri,
content_hash: pw.content_hash,
routing: None,
size_bytes_compressed: pw.size_bytes_compressed,
size_bytes_uncompressed: pw.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 99),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let lw = write_manifest(storage.as_ref(), &list)
.await
.expect("write list");
write_pointer(
storage.as_ref(),
&PointerFile {
manifest_id: 1,
manifest_uri: lw.uri,
content_hash: lw.content_hash,
},
None,
)
.await
.expect("write pointer");
entries
}
#[tokio::test]
async fn load_hydrates_flat_view_from_slow_state_blob() {
let opts = make_opts();
let (_dir, storage) = local_storage();
let entries = vec![
make_superfile_entry(100, hash_bucket_0_pk()),
make_superfile_entry(50, hash_bucket_0_pk()),
];
let published = slow_vector_state::write_state(storage.as_ref(), &entries, None)
.await
.expect("write blob");
let (blob_uri, blob_hash) = (published.uri, published.content_hash);
let (_dir2, storage2) = local_storage();
let _ = _dir2;
drop(storage2); let persisted = persist_two_entry_table(&storage, Some((blob_uri, blob_hash))).await;
let loaded = ManifestSnapshot::load(None, Arc::clone(&storage), Some(opts))
.await
.expect("load");
assert_eq!(loaded.superfiles.len(), 2);
let want: HashSet<Uuid> = persisted.iter().map(|e| e.superfile_id).collect();
let got: HashSet<Uuid> = loaded.superfiles.iter().map(|e| e.superfile_id).collect();
assert_eq!(
got.len(),
2,
"blob-hydrated flat view must carry both entries"
);
assert_eq!(want.len(), 2);
assert_eq!(
n_parts_initialized(&loaded),
0,
"hydration must not fetch any manifest part"
);
assert!(loaded.slow_vector_state_blob().is_some());
}
const ROUTING_TEST_DIM: usize = 16;
const ROUTING_TEST_ROT_SEED: u64 = 5;
fn make_summary_entry(docs: u64) -> Arc<SuperfileEntry> {
let base = make_superfile_entry(docs, hash_bucket_0_pk());
let mut entry = (*base).clone();
let mut flat = vec![0.0f32; 2 * ROUTING_TEST_DIM];
flat[0] = 1.0;
flat[ROUTING_TEST_DIM + 1] = -1.0;
let clusters = ClusterCentroids::from_fp32(2, ROUTING_TEST_DIM as u32, &flat, vec![3, 4]);
clusters.prewarm_admit_codes(
&RandomRotation::new(ROUTING_TEST_DIM, ROUTING_TEST_ROT_SEED),
&BitQuantizer::new(ROUTING_TEST_DIM),
ROUTING_TEST_ROT_SEED,
);
entry.vector_summary.insert(
"emb".into(),
VectorSummary {
centroid: vec![0.5; ROUTING_TEST_DIM],
cells: vec![CellVectorSummary {
cell_id: Some(0),
clusters,
}],
},
);
Arc::new(entry)
}
#[tokio::test]
async fn state_blob_hydrates_stripped_for_all_consumers() {
let (_dir, storage) = local_storage();
let entries = vec![make_summary_entry(100), make_summary_entry(50)];
let published = slow_vector_state::write_state(storage.as_ref(), &entries, None)
.await
.expect("write blobs");
persist_two_entry_table(
&storage,
Some((published.uri.clone(), published.content_hash)),
)
.await;
let consumer_opts = |knob: bool| {
Arc::new(
SupertableOptions::new(simple_schema(), vec![], vec![], None)
.expect("valid options")
.with_summary_centroids_from_superfiles(knob),
)
};
let resident = |m: &ManifestSnapshot| {
m.superfiles
.iter()
.map(|e| {
let clusters = &e.vector_summary["emb"].cells[0].clusters;
assert!(clusters.admit_codes_built().is_some(), "slab always rides");
clusters.vectors_resident()
})
.collect::<Vec<bool>>()
};
let knob_on = ManifestSnapshot::load(None, Arc::clone(&storage), Some(consumer_opts(true)))
.await
.expect("knob-on load");
assert_eq!(
resident(&knob_on),
vec![false, false],
"the routing-shaped state blob hydrates stripped entries"
);
let knob_off =
ManifestSnapshot::load(None, Arc::clone(&storage), Some(consumer_opts(false)))
.await
.expect("knob-off load");
assert_eq!(
resident(&knob_off),
vec![false, false],
"knob-off hydrates the same routing-shaped blob — fp32 lives in the section"
);
}
async fn persist_summary_part_table(storage: &Arc<dyn StorageProvider>, with_routing: bool) {
let entries = vec![make_summary_entry(100), make_summary_entry(50)];
let part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: entries,
};
let full = part::encode(&part);
let full_hash = ContentHash::of(&full);
let routing = part::encode_with_mode(&part, SummaryWireMode::RoutingOnly);
let routing_hash = ContentHash::of(&routing);
assert!(
routing.len() < full.len(),
"routing part must shed the fp32 payload ({} vs {} bytes)",
routing.len(),
full.len()
);
write_part_bytes(storage.as_ref(), &full)
.await
.expect("put full part");
write_part_bytes(storage.as_ref(), &routing)
.await
.expect("put routing part");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 1,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![ManifestPartEntry {
part_id: part.part_id,
uri: part_uri(&full_hash),
content_hash: full_hash,
routing: with_routing.then(|| RoutingRef {
uri: part_uri(&routing_hash),
content_hash: routing_hash,
}),
size_bytes_compressed: full.len() as u64,
size_bytes_uncompressed: full.len() as u64,
n_superfiles: 2,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let lw = write_manifest(storage.as_ref(), &list)
.await
.expect("write list");
write_pointer(
storage.as_ref(),
&PointerFile {
manifest_id: 1,
manifest_uri: lw.uri,
content_hash: lw.content_hash,
},
None,
)
.await
.expect("write pointer");
}
#[tokio::test]
async fn consumer_mode_hydrates_parts_from_routing_sibling() {
let consumer_opts = |knob: bool| {
Arc::new(
SupertableOptions::new(simple_schema(), vec![], vec![], None)
.expect("valid options")
.with_summary_centroids_from_superfiles(knob),
)
};
let shape = |m: &ManifestSnapshot| {
m.superfiles
.iter()
.map(|e| {
let clusters = &e.vector_summary["emb"].cells[0].clusters;
(
clusters.vectors_resident(),
clusters.admit_codes_built().is_some(),
)
})
.collect::<Vec<(bool, bool)>>()
};
let (_dir, storage) = local_storage();
persist_summary_part_table(&storage, true).await;
let knob_on = ManifestSnapshot::load(None, Arc::clone(&storage), Some(consumer_opts(true)))
.await
.expect("knob-on load");
assert_eq!(
shape(&knob_on),
vec![(false, true), (false, true)],
"knob-on consumer must hydrate stripped entries (slab riding) from the routing part"
);
let knob_off =
ManifestSnapshot::load(None, Arc::clone(&storage), Some(consumer_opts(false)))
.await
.expect("knob-off load");
assert_eq!(
shape(&knob_off),
vec![(true, false), (true, false)],
"knob-off load must keep the full part's resident fp32 (slab rebuilt by prewarm \
on real tables)"
);
let (_dir2, storage2) = local_storage();
persist_summary_part_table(&storage2, false).await;
let fallback =
ManifestSnapshot::load(None, Arc::clone(&storage2), Some(consumer_opts(true)))
.await
.expect("fallback load");
assert_eq!(
shape(&fallback),
vec![(true, false), (true, false)],
"knob-on without a routing ref must fall back to the full part"
);
}
#[tokio::test]
async fn rebuild_part_and_entry_stamps_routing_sibling() {
let (entry, encoded_part) =
rebuild_part_and_entry(vec![], vec![make_summary_entry(10)], None, false);
let routing = entry.routing.expect("sibling stamped");
let routing_encoded = encoded_part
.routing_encoded
.as_ref()
.expect("user part carries sibling bytes");
assert_eq!(
routing.content_hash,
ContentHash::of(routing_encoded),
"entry ref must address the returned sibling bytes"
);
assert_eq!(routing.uri, part_uri(&routing.content_hash));
assert_eq!(entry.content_hash, ContentHash::of(&encoded_part.encoded));
assert!(
routing_encoded.len() < encoded_part.encoded.len(),
"sibling must shed the fp32 payload"
);
let (hidden_entry, hidden_encoded) =
rebuild_part_and_entry(vec![], vec![make_summary_entry(10)], None, true);
assert!(
hidden_entry.routing.is_none(),
"hidden part must not stamp a sibling — its primary is routing-shaped"
);
assert!(
hidden_encoded.routing_encoded.is_none(),
"hidden part must not carry sibling bytes"
);
assert!(
hidden_encoded.encoded.len() < encoded_part.encoded.len(),
"hidden primary must shed the fp32 payload ({} vs {} bytes)",
hidden_encoded.encoded.len(),
encoded_part.encoded.len()
);
}
#[tokio::test]
async fn refresh_with_unchanged_slow_ref_reuses_entries() {
let opts = make_opts();
let (_dir, storage) = local_storage();
let entries = vec![
make_superfile_entry(100, hash_bucket_0_pk()),
make_superfile_entry(50, hash_bucket_0_pk()),
];
let published = slow_vector_state::write_state(storage.as_ref(), &entries, None)
.await
.expect("write blob");
let (blob_uri, blob_hash) = (published.uri, published.content_hash);
persist_two_entry_table(&storage, Some((blob_uri, blob_hash))).await;
let a = ManifestSnapshot::load(None, Arc::clone(&storage), Some(Arc::clone(&opts)))
.await
.expect("load A");
let (_, meta) = read_pointer(storage.as_ref())
.await
.expect("read pointer")
.expect("pointer present");
let etag = meta.etag.expect("localfs pointer etag");
let stamped = a.with_deleted_user_ids(Vec::new());
stamped
.write(storage.as_ref(), Some(etag.as_str()), &[])
.await
.expect("stamp publish");
let b = ManifestSnapshot::load(Some(Arc::clone(&a)), Arc::clone(&storage), None)
.await
.expect("refresh");
assert_eq!(b.get_manifest_id(), a.get_manifest_id() + 1);
assert!(b.slow_vector_state_blob().is_some(), "ref preserved");
assert_eq!(b.superfiles.len(), a.superfiles.len());
for (be, ae) in b.superfiles.iter().zip(a.superfiles.iter()) {
assert!(
Arc::ptr_eq(be, ae),
"unchanged ref must reuse the SAME decoded entries — \
the centroid state never leaves memory on list-only churn"
);
}
assert_eq!(
n_parts_initialized(&b),
0,
"refresh with unchanged ref must not fetch parts"
);
}
#[tokio::test]
async fn load_with_corrupt_slow_ref_raises_hydration_error() {
let opts = make_opts();
let (_dir, storage) = local_storage();
let bogus = (
"slow-vector-state/state-missing.bin".to_string(),
ContentHash([9u8; 32]),
);
persist_two_entry_table(&storage, Some(bogus)).await;
let err = ManifestSnapshot::load(None, Arc::clone(&storage), Some(opts))
.await
.expect_err("corrupt slow-state ref must fail the load loudly");
assert!(
matches!(err, ManifestLoadError::SlowStateHydration(_)),
"expected SlowStateHydration, got: {err:?}"
);
}
#[tokio::test]
async fn update_add_to_existing_partition_rewrites_part() {
let opts = make_opts();
let (_dir, storage) = local_storage();
let old_superfile = make_superfile_entry(100, hash_bucket_0_pk());
let existing_part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![old_superfile.clone()],
};
let pw = write_manifest_part(storage.as_ref(), &existing_part)
.await
.expect("write part");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![ManifestPartEntry {
part_id: pw.part_id,
uri: pw.uri,
content_hash: pw.content_hash,
routing: None,
size_bytes_compressed: pw.size_bytes_compressed,
size_bytes_uncompressed: pw.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 99),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts = DashMap::new();
parts.insert(
pw.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(existing_part)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![old_superfile],
vector_index_storage_prefix: None,
},
list: Some(list),
parts,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let new_entry = make_new_entry(50);
let new_entries = vec![new_entry];
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 1);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, 2);
assert_eq!(parts[0].part.superfiles.len(), 2);
let total_docs: u64 = parts[0].part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(total_docs, 150);
}
#[tokio::test]
async fn update_leaves_unchanged_parts_untouched() {
const SUPERFILES_PER_PART: u64 = 2;
const TARGET_SUPERFILES_PER_PART: u64 = 3;
let (_dir, storage) = local_storage();
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = TARGET_SUPERFILES_PER_PART;
let opts = Arc::new(base_opts.with_storage(storage.clone()));
let pk_a = hash2_pk(0);
let pk_b = hash2_pk(1);
async fn two_superfile_part(
storage: &dyn StorageProvider,
pk: &[u8],
hint: u32,
docs: [u64; 2],
) -> (ManifestPart, PartWriteResult) {
let part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![
make_superfile_entry_hinted(docs[0], pk.to_vec(), hint),
make_superfile_entry_hinted(docs[1], pk.to_vec(), hint),
],
};
let pw = write_manifest_part(storage, &part)
.await
.expect("write part");
(part, pw)
}
let (part_a_old, pw_a_old) =
two_superfile_part(storage.as_ref(), &pk_a, 0, [100, 110]).await;
let (_part_a_latest, pw_a_latest) =
two_superfile_part(storage.as_ref(), &pk_a, 0, [120, 130]).await;
let (part_b, pw_b) = two_superfile_part(storage.as_ref(), &pk_b, 1, [200, 210]).await;
let entry_for = |pw: &PartWriteResult| -> ManifestPartEntry {
ManifestPartEntry {
part_id: pw.part_id,
uri: pw.uri.clone(),
content_hash: pw.content_hash,
routing: None,
size_bytes_compressed: pw.size_bytes_compressed,
size_bytes_uncompressed: pw.size_bytes_uncompressed,
n_superfiles: SUPERFILES_PER_PART,
id_range: (0, 0),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}
};
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 2,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![
entry_for(&pw_a_old),
entry_for(&pw_a_latest),
entry_for(&pw_b),
],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
part_b.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_b.clone())))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: part_a_old
.superfiles
.iter()
.chain(part_b.superfiles.iter())
.cloned()
.collect(),
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let new_entry = make_new_entry_hinted(140, 0);
let (new_manifest, parts_to_write) = old_manifest
.update(from_ref(&new_entry), &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 3, "list entry count");
assert_eq!(
parts_to_write.len(),
1,
"only the rewritten last part should be re-emitted; \
unchanged parts must not be re-encoded/PUT",
);
let find = |part_id: PartId| {
list_entries
.iter()
.find(|e| e.part_id == part_id)
.unwrap_or_else(|| panic!("entry for part {part_id:?} missing after update"))
};
let part0_after = find(pw_a_old.part_id);
assert_eq!(part0_after.uri, pw_a_old.uri, "carried-over part_0 uri");
assert_eq!(
part0_after.content_hash, pw_a_old.content_hash,
"carried-over part_0 content_hash",
);
assert_eq!(part0_after.n_superfiles, SUPERFILES_PER_PART);
let part1_after = find(pw_a_latest.part_id);
assert_eq!(part1_after.uri, pw_a_latest.uri, "carried-over part_1 uri");
assert_eq!(
part1_after.content_hash, pw_a_latest.content_hash,
"carried-over part_1 content_hash",
);
assert_eq!(part1_after.n_superfiles, SUPERFILES_PER_PART);
assert_eq!(
parts_to_write[0].part.superfiles.len(),
(SUPERFILES_PER_PART + 1) as usize,
"rewritten last part should hold its 2 superfiles + the new one",
);
assert!(
!list_entries.iter().any(|e| e.part_id == pw_b.part_id),
"the rewritten last part is replaced, so its old part_id must not survive",
);
let rewritten_after = list_entries
.iter()
.find(|e| e.part_id != pw_a_old.part_id && e.part_id != pw_a_latest.part_id)
.expect("rewritten last entry present after the add");
assert_eq!(rewritten_after.n_superfiles, SUPERFILES_PER_PART + 1);
let rewritten_v1_part_id = rewritten_after.part_id;
let (after_removal, removal_parts) = new_manifest
.update(&[], from_ref(&new_entry))
.await
.expect("update removal");
let entries_after = after_removal.get_all_list_entries();
assert_eq!(entries_after.len(), 3, "list entry count after removal");
assert!(
!entries_after
.iter()
.any(|e| e.part_id == rewritten_v1_part_id),
"the part we removed a superfile from must be rebuilt (new part_id)",
);
let part0_after_removal = entries_after
.iter()
.find(|e| e.part_id == pw_a_old.part_id)
.expect("part_0 must survive the removal unchanged");
assert_eq!(
part0_after_removal.uri, pw_a_old.uri,
"part_0 uri after removal",
);
assert_eq!(
part0_after_removal.content_hash, pw_a_old.content_hash,
"part_0 content_hash after removal",
);
let part1_after_removal = entries_after
.iter()
.find(|e| e.part_id == pw_a_latest.part_id)
.expect("part_1 must survive the removal unchanged");
assert_eq!(
part1_after_removal.uri, pw_a_latest.uri,
"part_1 uri after removal",
);
assert_eq!(
part1_after_removal.content_hash, pw_a_latest.content_hash,
"part_1 content_hash after removal",
);
assert_eq!(
removal_parts.len(),
1,
"only the part we removed from should be rewritten; unchanged parts \
must not be re-encoded/PUT",
);
}
#[tokio::test]
async fn update_rewrite_partition_within_target() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 3;
let opts = Arc::new(base_opts);
let (_dir, storage) = local_storage();
let sf1 = make_superfile_entry(100, hash_bucket_0_pk());
let sf2 = make_superfile_entry(150, hash_bucket_0_pk());
let existing_part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf1.clone(), sf2.clone()],
};
let pw = write_manifest_part(storage.as_ref(), &existing_part)
.await
.expect("write part");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![ManifestPartEntry {
part_id: pw.part_id,
uri: pw.uri,
content_hash: pw.content_hash,
routing: None,
size_bytes_compressed: pw.size_bytes_compressed,
size_bytes_uncompressed: pw.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts = DashMap::new();
parts.insert(
pw.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(existing_part)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf1, sf2],
vector_index_storage_prefix: None,
},
list: Some(list),
parts,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let new_entry = make_new_entry(75);
let new_entries = vec![new_entry];
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 1);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, 3);
let part = &parts[0];
assert_eq!(part.part.superfiles.len(), 3);
let total_docs: u64 = part.part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(total_docs, 325); }
#[tokio::test]
async fn update_split_partition_exceeds_target() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 2;
let opts = Arc::new(base_opts);
let (_dir, storage) = local_storage();
let sf1 = make_superfile_entry(100, hash_bucket_0_pk());
let sf2 = make_superfile_entry(150, hash_bucket_0_pk());
let existing_part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf1.clone(), sf2.clone()],
};
let pw = write_manifest_part(storage.as_ref(), &existing_part)
.await
.expect("write part");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![ManifestPartEntry {
part_id: pw.part_id,
uri: pw.uri,
content_hash: pw.content_hash,
routing: None,
size_bytes_compressed: pw.size_bytes_compressed,
size_bytes_uncompressed: pw.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts = DashMap::new();
parts.insert(
pw.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(existing_part)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf1, sf2],
vector_index_storage_prefix: None,
},
list: Some(list),
parts,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let new_entry1 = make_new_entry(75);
let new_entry2 = make_new_entry(80);
let new_entries = vec![new_entry1, new_entry2];
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 2);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, 2);
assert_eq!(list_entries[1].n_superfiles, 2);
let part = &parts[0];
assert_eq!(part.part.superfiles.len(), 2);
let total_docs: u64 = part.part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(total_docs, 155); }
fn make_superfile_entry_hinted(docs: u64, pk: Vec<u8>, hint: u32) -> Arc<SuperfileEntry> {
make_entry(docs, pk, Some(hint))
}
fn make_new_entry_hinted(docs: u64, hint: u32) -> Arc<SuperfileEntry> {
make_entry(docs, Vec::new(), Some(hint))
}
fn hash2_pk(bucket: u32) -> Vec<u8> {
bucket.to_le_bytes().to_vec()
}
#[tokio::test]
async fn update_split_partition_exceeds_size_threshold() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 10_000;
base_opts.part_size_threshold_bytes = 1;
let opts = Arc::new(base_opts);
let (_dir, storage) = local_storage();
let sf1 = make_superfile_entry(100, hash_bucket_0_pk());
let sf2 = make_superfile_entry(150, hash_bucket_0_pk());
let existing_part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf1.clone(), sf2.clone()],
};
let pw = write_manifest_part(storage.as_ref(), &existing_part)
.await
.expect("write part");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![ManifestPartEntry {
part_id: pw.part_id,
uri: pw.uri.clone(),
content_hash: pw.content_hash,
routing: None,
size_bytes_compressed: pw.size_bytes_compressed,
size_bytes_uncompressed: pw.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let frozen_part_id = pw.part_id;
let frozen_hash = pw.content_hash;
let loader = ManifestPartLoader::new(storage, &list);
let parts = DashMap::new();
parts.insert(
pw.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(existing_part)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf1, sf2],
vector_index_storage_prefix: None,
},
list: Some(list),
parts,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let new_entries = vec![make_new_entry(75)];
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(
list_entries.len(),
2,
"size-capped latest part must freeze; new entries go to a fresh part"
);
assert_eq!(parts.len(), 1, "only the fresh part is encoded + written");
assert_eq!(list_entries[0].part_id, frozen_part_id);
assert_eq!(list_entries[0].content_hash, frozen_hash);
assert_eq!(list_entries[0].n_superfiles, 2);
assert_eq!(list_entries[1].n_superfiles, 1);
assert_eq!(parts[0].part.superfiles.len(), 1);
assert_eq!(parts[0].part.superfiles[0].n_docs, 75);
}
#[tokio::test]
async fn update_older_entry_preserved_when_latest_rewritten() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 2;
let opts = Arc::new(base_opts);
let (_dir, storage) = local_storage();
let sf_old = make_superfile_entry(100, hash_bucket_0_pk());
let sf_latest = make_superfile_entry(150, hash_bucket_0_pk());
let part_old = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_old.clone()],
};
let pw_old = write_manifest_part(storage.as_ref(), &part_old)
.await
.expect("write part_old");
let part_latest = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_latest.clone()],
};
let pw_latest = write_manifest_part(storage.as_ref(), &part_latest)
.await
.expect("write part_latest");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![
ManifestPartEntry {
part_id: pw_old.part_id,
uri: pw_old.uri.clone(),
content_hash: pw_old.content_hash,
routing: None,
size_bytes_compressed: pw_old.size_bytes_compressed,
size_bytes_uncompressed: pw_old.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 99),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
ManifestPartEntry {
part_id: pw_latest.part_id,
uri: pw_latest.uri,
content_hash: pw_latest.content_hash,
routing: None,
size_bytes_compressed: pw_latest.size_bytes_compressed,
size_bytes_uncompressed: pw_latest.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts = DashMap::new();
parts.insert(
part_latest.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_latest)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf_old, sf_latest],
vector_index_storage_prefix: None,
},
list: Some(list),
parts,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let new_entries = vec![make_new_entry(75)];
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 2);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, 1);
assert_eq!(list_entries[0].uri, pw_old.uri);
assert_eq!(list_entries[1].n_superfiles, 2);
assert_eq!(parts[0].part.superfiles.len(), 2);
let total_docs: u64 = parts[0].part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(total_docs, 225); }
#[tokio::test]
async fn update_two_partitions_both_touched() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 3;
let opts = Arc::new(base_opts);
let (_dir, storage) = local_storage();
let sf_a = make_superfile_entry_hinted(100, hash2_pk(0), 0);
let part_a = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_a.clone()],
};
let pw_a = write_manifest_part(storage.as_ref(), &part_a)
.await
.expect("write part_a");
let sf_b = make_superfile_entry_hinted(200, hash2_pk(1), 1);
let part_b = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_b.clone()],
};
let pw_b = write_manifest_part(storage.as_ref(), &part_b)
.await
.expect("write part_b");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 2,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![
ManifestPartEntry {
part_id: pw_a.part_id,
uri: pw_a.uri,
content_hash: pw_a.content_hash,
routing: None,
size_bytes_compressed: pw_a.size_bytes_compressed,
size_bytes_uncompressed: pw_a.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 99),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
ManifestPartEntry {
part_id: pw_b.part_id,
uri: pw_b.uri,
content_hash: pw_b.content_hash,
routing: None,
size_bytes_compressed: pw_b.size_bytes_compressed,
size_bytes_uncompressed: pw_b.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 199),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
part_a.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_a)))),
);
parts_map.insert(
part_b.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_b)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf_a, sf_b],
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let new_entries = vec![make_new_entry_hinted(50, 0), make_new_entry_hinted(80, 1)];
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 2);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].part_id, pw_a.part_id);
assert_eq!(list_entries[0].n_superfiles, 1);
assert_eq!(list_entries[0].content_hash, pw_a.content_hash);
assert_eq!(list_entries[1].n_superfiles, 3);
assert_eq!(parts[0].part.superfiles.len(), 3);
let total_docs: u64 = parts[0].part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(total_docs, 330);
let hints: Vec<_> = parts[0]
.part
.superfiles
.iter()
.map(|s| s.partition_hint)
.collect();
assert!(hints.contains(&Some(0)), "hint-0 new entry preserved");
assert!(hints.contains(&Some(1)), "hint-1 new entry preserved");
}
#[tokio::test]
async fn update_two_partitions_one_touched_exact_carry_over() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 3;
let opts = Arc::new(base_opts);
let (_dir, storage) = local_storage();
let sf_a = make_superfile_entry_hinted(100, hash2_pk(0), 0);
let part_a = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_a.clone()],
};
let pw_a = write_manifest_part(storage.as_ref(), &part_a)
.await
.expect("write part_a");
let sf_b = make_superfile_entry_hinted(200, hash2_pk(1), 1);
let part_b = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_b.clone()],
};
let pw_b = write_manifest_part(storage.as_ref(), &part_b)
.await
.expect("write part_b");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 2,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![
ManifestPartEntry {
part_id: pw_a.part_id,
uri: pw_a.uri.clone(),
content_hash: pw_a.content_hash,
routing: None,
size_bytes_compressed: pw_a.size_bytes_compressed,
size_bytes_uncompressed: pw_a.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 99),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
ManifestPartEntry {
part_id: pw_b.part_id,
uri: pw_b.uri,
content_hash: pw_b.content_hash,
routing: None,
size_bytes_compressed: pw_b.size_bytes_compressed,
size_bytes_uncompressed: pw_b.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 199),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
part_a.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_a)))),
);
parts_map.insert(
part_b.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_b)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf_a, sf_b],
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let new_entries = vec![make_new_entry_hinted(50, 0)];
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 2);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].part_id, pw_a.part_id);
assert_eq!(list_entries[0].n_superfiles, 1);
assert_eq!(list_entries[0].content_hash, pw_a.content_hash);
assert_eq!(list_entries[1].n_superfiles, 2);
assert_eq!(parts[0].part.superfiles.len(), 2);
let docs: u64 = parts[0].part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(docs, 250); }
#[tokio::test]
async fn update_two_partitions_each_with_prior_split() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 2;
let opts = Arc::new(base_opts);
let (_dir, storage) = local_storage();
let sf_a_old = make_superfile_entry_hinted(100, hash2_pk(0), 0);
let part_a_old = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_a_old.clone()],
};
let pw_a_old = write_manifest_part(storage.as_ref(), &part_a_old)
.await
.expect("write part_a_old");
let sf_a_latest = make_superfile_entry_hinted(150, hash2_pk(0), 0);
let part_a_latest = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_a_latest.clone()],
};
let pw_a_latest = write_manifest_part(storage.as_ref(), &part_a_latest)
.await
.expect("write part_a_latest");
let sf_b_old = make_superfile_entry_hinted(200, hash2_pk(1), 1);
let part_b_old = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_b_old.clone()],
};
let pw_b_old = write_manifest_part(storage.as_ref(), &part_b_old)
.await
.expect("write part_b_old");
let sf_b_latest = make_superfile_entry_hinted(250, hash2_pk(1), 1);
let part_b_latest = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_b_latest.clone()],
};
let pw_b_latest = write_manifest_part(storage.as_ref(), &part_b_latest)
.await
.expect("write part_b_latest");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 2,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![
ManifestPartEntry {
part_id: pw_a_old.part_id,
uri: pw_a_old.uri.clone(),
content_hash: pw_a_old.content_hash,
routing: None,
size_bytes_compressed: pw_a_old.size_bytes_compressed,
size_bytes_uncompressed: pw_a_old.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 99),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
ManifestPartEntry {
part_id: pw_a_latest.part_id,
uri: pw_a_latest.uri.clone(),
content_hash: pw_a_latest.content_hash,
routing: None,
size_bytes_compressed: pw_a_latest.size_bytes_compressed,
size_bytes_uncompressed: pw_a_latest.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
ManifestPartEntry {
part_id: pw_b_old.part_id,
uri: pw_b_old.uri.clone(),
content_hash: pw_b_old.content_hash,
routing: None,
size_bytes_compressed: pw_b_old.size_bytes_compressed,
size_bytes_uncompressed: pw_b_old.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 199),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
ManifestPartEntry {
part_id: pw_b_latest.part_id,
uri: pw_b_latest.uri.clone(),
content_hash: pw_b_latest.content_hash,
routing: None,
size_bytes_compressed: pw_b_latest.size_bytes_compressed,
size_bytes_uncompressed: pw_b_latest.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 249),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
part_a_latest.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_a_latest)))),
);
parts_map.insert(
part_b_latest.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_b_latest)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf_a_old, sf_a_latest, sf_b_old, sf_b_latest],
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let new_entries = vec![make_new_entry_hinted(75, 0), make_new_entry_hinted(90, 1)];
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 5);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].part_id, pw_a_old.part_id);
assert_eq!(list_entries[0].uri, pw_a_old.uri);
assert_eq!(list_entries[0].content_hash, pw_a_old.content_hash);
assert_eq!(list_entries[1].part_id, pw_a_latest.part_id);
assert_eq!(list_entries[2].part_id, pw_b_old.part_id);
assert_eq!(list_entries[2].uri, pw_b_old.uri);
assert_eq!(list_entries[2].content_hash, pw_b_old.content_hash);
assert_eq!(list_entries[3].part_id, pw_b_latest.part_id);
for e in &list_entries[0..4] {
assert_eq!(e.n_superfiles, 1);
}
assert_eq!(list_entries[4].n_superfiles, 2);
assert_eq!(parts[0].part.superfiles.len(), 2);
let docs: u64 = parts[0].part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(docs, 165);
let hints: Vec<_> = parts[0]
.part
.superfiles
.iter()
.map(|s| s.partition_hint)
.collect();
assert!(hints.contains(&Some(0)), "hint-0 new entry preserved");
assert!(hints.contains(&Some(1)), "hint-1 new entry preserved");
}
#[tokio::test]
async fn update_multiple_partitions_land_in_one_lineage() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 8;
let opts = Arc::new(base_opts);
let old_manifest = empty_manifest(&opts);
let hints = [0u32, 1, 2, 3];
let new_entries: Vec<_> = hints
.iter()
.enumerate()
.map(|(i, &h)| make_new_entry_hinted(100 + i as u64, h))
.collect();
let (new_manifest, parts) = old_manifest
.update(&new_entries, &[])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 1);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, hints.len() as u64);
assert_eq!(parts[0].part.superfiles.len(), hints.len());
for (&h, expected) in hints.iter().zip(new_entries.iter()) {
let landed = parts[0]
.part
.superfiles
.iter()
.find(|s| s.superfile_id == expected.superfile_id)
.unwrap_or_else(|| panic!("superfile with hint {h} landed in the part"));
assert_eq!(landed.partition_hint, Some(h), "partition_hint preserved");
assert_eq!(
landed.partition_key,
hash2_pk(0),
"single-bucket Hash stamps bucket 0 on every entry"
);
}
}
#[tokio::test]
async fn update_remove_one_superfile_from_partition() {
let opts = make_opts();
let (_dir, storage) = local_storage();
let sf_keep = make_superfile_entry(100, hash_bucket_0_pk());
let sf_remove = make_superfile_entry(150, hash_bucket_0_pk());
let existing_part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_keep.clone(), sf_remove.clone()],
};
let pw = write_manifest_part(storage.as_ref(), &existing_part)
.await
.expect("write part");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![ManifestPartEntry {
part_id: pw.part_id,
uri: pw.uri,
content_hash: pw.content_hash,
routing: None,
size_bytes_compressed: pw.size_bytes_compressed,
size_bytes_uncompressed: pw.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
existing_part.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(existing_part)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf_keep.clone(), sf_remove.clone()],
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let (new_manifest, parts) = old_manifest
.update(&[], from_ref(&sf_remove))
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 1);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, 1);
assert_eq!(parts[0].part.superfiles.len(), 1);
assert_eq!(
parts[0].part.superfiles[0].superfile_id,
sf_keep.superfile_id
);
let total_docs: u64 = parts[0].part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(total_docs, 100);
}
#[tokio::test]
async fn update_add_and_remove_in_same_partition() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 3;
let opts = Arc::new(base_opts);
let (_dir, storage) = local_storage();
let sf_keep = make_superfile_entry(100, hash_bucket_0_pk());
let sf_remove = make_superfile_entry(150, hash_bucket_0_pk());
let existing_part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_keep.clone(), sf_remove.clone()],
};
let pw = write_manifest_part(storage.as_ref(), &existing_part)
.await
.expect("write part");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![ManifestPartEntry {
part_id: pw.part_id,
uri: pw.uri,
content_hash: pw.content_hash,
routing: None,
size_bytes_compressed: pw.size_bytes_compressed,
size_bytes_uncompressed: pw.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
existing_part.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(existing_part)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf_keep.clone(), sf_remove.clone()],
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let sf_new = make_new_entry(75);
let new_entries = vec![sf_new.clone()];
let (new_manifest, parts) = old_manifest
.update(&new_entries, from_ref(&sf_remove))
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 1);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, 2);
assert_eq!(parts[0].part.superfiles.len(), 2);
let ids: Vec<_> = parts[0]
.part
.superfiles
.iter()
.map(|s| s.superfile_id)
.collect();
assert!(ids.contains(&sf_keep.superfile_id));
assert!(ids.contains(&sf_new.superfile_id));
assert!(!ids.contains(&sf_remove.superfile_id));
let total_docs: u64 = parts[0].part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(total_docs, 175); }
#[tokio::test]
async fn update_remove_from_one_partition_other_carried_over_exactly() {
let opts = make_opts();
let (_dir, storage) = local_storage();
let sf_a_keep = make_superfile_entry_hinted(100, hash2_pk(0), 0);
let sf_a_remove = make_superfile_entry_hinted(150, hash2_pk(0), 0);
let part_a = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_a_keep.clone(), sf_a_remove.clone()],
};
let pw_a = write_manifest_part(storage.as_ref(), &part_a)
.await
.expect("write part_a");
let sf_b = make_superfile_entry_hinted(200, hash2_pk(1), 1);
let part_b = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_b.clone()],
};
let pw_b = write_manifest_part(storage.as_ref(), &part_b)
.await
.expect("write part_b");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 2,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![
ManifestPartEntry {
part_id: pw_a.part_id,
uri: pw_a.uri,
content_hash: pw_a.content_hash,
routing: None,
size_bytes_compressed: pw_a.size_bytes_compressed,
size_bytes_uncompressed: pw_a.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
ManifestPartEntry {
part_id: pw_b.part_id,
uri: pw_b.uri.clone(),
content_hash: pw_b.content_hash,
routing: None,
size_bytes_compressed: pw_b.size_bytes_compressed,
size_bytes_uncompressed: pw_b.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 199),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
part_a.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_a)))),
);
parts_map.insert(
part_b.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_b)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf_a_keep.clone(), sf_a_remove.clone(), sf_b.clone()],
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let (new_manifest, parts) = old_manifest
.update(&[], from_ref(&sf_a_remove))
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 2);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, 1);
assert_eq!(parts[0].part.superfiles.len(), 1);
assert_eq!(
parts[0].part.superfiles[0].superfile_id,
sf_a_keep.superfile_id
);
let docs_a: u64 = parts[0].part.superfiles.iter().map(|s| s.n_docs).sum();
assert_eq!(docs_a, 100);
assert_eq!(list_entries[1].n_superfiles, 1);
assert_eq!(list_entries[1].uri, pw_b.uri);
assert_eq!(list_entries[1].content_hash, pw_b.content_hash);
}
#[tokio::test]
async fn update_remove_from_latest_part_in_split_partition() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 2;
let opts = Arc::new(base_opts);
let (_dir, storage) = local_storage();
let sf_a_old = make_superfile_entry(100, hash_bucket_0_pk());
let part_a_old = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_a_old.clone()],
};
let pw_a_old = write_manifest_part(storage.as_ref(), &part_a_old)
.await
.expect("write part_a_old");
let sf_a_latest_keep = make_superfile_entry(150, hash_bucket_0_pk());
let sf_a_latest_remove = make_superfile_entry(200, hash_bucket_0_pk());
let part_a_latest = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_a_latest_keep.clone(), sf_a_latest_remove.clone()],
};
let pw_a_latest = write_manifest_part(storage.as_ref(), &part_a_latest)
.await
.expect("write part_a_latest");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![
ManifestPartEntry {
part_id: pw_a_old.part_id,
uri: pw_a_old.uri.clone(),
content_hash: pw_a_old.content_hash,
routing: None,
size_bytes_compressed: pw_a_old.size_bytes_compressed,
size_bytes_uncompressed: pw_a_old.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 99),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
ManifestPartEntry {
part_id: pw_a_latest.part_id,
uri: pw_a_latest.uri.clone(),
content_hash: pw_a_latest.content_hash,
routing: None,
size_bytes_compressed: pw_a_latest.size_bytes_compressed,
size_bytes_uncompressed: pw_a_latest.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 199),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
part_a_old.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_a_old)))),
);
parts_map.insert(
part_a_latest.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_a_latest)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![
sf_a_old.clone(),
sf_a_latest_keep.clone(),
sf_a_latest_remove.clone(),
],
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let (new_manifest, parts_to_write) = old_manifest
.update(&[], from_ref(&sf_a_latest_remove))
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 2);
assert_eq!(parts_to_write.len(), 1);
let all_ids: Vec<_> = parts_to_write
.iter()
.flat_map(|ep| ep.part.superfiles.iter())
.map(|s| s.superfile_id)
.collect();
assert!(
all_ids.contains(&sf_a_latest_keep.superfile_id),
"sf_a_latest_keep must survive"
);
assert!(
!all_ids.contains(&sf_a_latest_remove.superfile_id),
"sf_a_latest_remove must be absent"
);
assert_eq!(list_entries[0].n_superfiles, 1);
assert_eq!(list_entries[1].n_superfiles, 1);
}
#[tokio::test]
async fn update_remove_all_superfiles_empties_partition() {
let opts = make_opts();
let (_dir, storage) = local_storage();
let sf1 = make_superfile_entry(100, hash_bucket_0_pk());
let sf2 = make_superfile_entry(150, hash_bucket_0_pk());
let existing_part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf1.clone(), sf2.clone()],
};
let pw = write_manifest_part(storage.as_ref(), &existing_part)
.await
.expect("write part");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![ManifestPartEntry {
part_id: pw.part_id,
uri: pw.uri,
content_hash: pw.content_hash,
routing: None,
size_bytes_compressed: pw.size_bytes_compressed,
size_bytes_uncompressed: pw.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
existing_part.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(existing_part)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf1.clone(), sf2.clone()],
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let (new_manifest, parts) = old_manifest
.update(&[], &[sf1.clone(), sf2.clone()])
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 1);
assert_eq!(parts.len(), 1);
assert_eq!(list_entries[0].n_superfiles, 0);
assert_eq!(parts[0].part.superfiles.len(), 0);
}
#[tokio::test]
async fn update_remove_nonexistent_superfile_id_is_noop() {
let opts = make_opts();
let (_dir, storage) = local_storage();
let sf1 = make_superfile_entry(100, hash_bucket_0_pk());
let sf2 = make_superfile_entry(150, hash_bucket_0_pk());
let existing_part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf1.clone(), sf2.clone()],
};
let pw = write_manifest_part(storage.as_ref(), &existing_part)
.await
.expect("write part");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![ManifestPartEntry {
part_id: pw.part_id,
uri: pw.uri,
content_hash: pw.content_hash,
routing: None,
size_bytes_compressed: pw.size_bytes_compressed,
size_bytes_uncompressed: pw.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
}],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
existing_part.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(existing_part)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![sf1.clone(), sf2.clone()],
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let sf_ghost = make_superfile_entry(50, hash_bucket_0_pk());
let (new_manifest, parts_to_write) = old_manifest
.update(&[], from_ref(&sf_ghost))
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 1);
assert_eq!(parts_to_write.len(), 0);
assert_eq!(list_entries[0].n_superfiles, 2);
}
#[tokio::test]
async fn update_remove_from_older_frozen_part_in_split_partition() {
let mut base_opts =
SupertableOptions::new(simple_schema(), vec![], vec![], None).expect("valid options");
base_opts.target_superfiles_per_part = 2;
let opts = Arc::new(base_opts);
let (_dir, storage) = local_storage();
let sf_a_old_keep = make_superfile_entry(100, hash_bucket_0_pk());
let sf_a_old_remove = make_superfile_entry(150, hash_bucket_0_pk());
let part_a_old = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_a_old_keep.clone(), sf_a_old_remove.clone()],
};
let pw_a_old = write_manifest_part(storage.as_ref(), &part_a_old)
.await
.expect("write part_a_old");
let sf_a_latest = make_superfile_entry(200, hash_bucket_0_pk());
let part_a_latest = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![sf_a_latest.clone()],
};
let pw_a_latest = write_manifest_part(storage.as_ref(), &part_a_latest)
.await
.expect("write part_a_latest");
let list = Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 0,
options_hash: ContentHash([0u8; 32]),
schema: vec![],
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts: vec![
ManifestPartEntry {
part_id: pw_a_old.part_id,
uri: pw_a_old.uri,
content_hash: pw_a_old.content_hash,
routing: None,
size_bytes_compressed: pw_a_old.size_bytes_compressed,
size_bytes_uncompressed: pw_a_old.size_bytes_uncompressed,
n_superfiles: 2,
id_range: (0, 149),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
ManifestPartEntry {
part_id: pw_a_latest.part_id,
uri: pw_a_latest.uri,
content_hash: pw_a_latest.content_hash,
routing: None,
size_bytes_compressed: pw_a_latest.size_bytes_compressed,
size_bytes_uncompressed: pw_a_latest.size_bytes_uncompressed,
n_superfiles: 1,
id_range: (0, 199),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
},
],
};
let loader = ManifestPartLoader::new(storage, &list);
let parts_map = DashMap::new();
parts_map.insert(
part_a_old.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_a_old)))),
);
parts_map.insert(
part_a_latest.part_id,
Arc::new(OnceCell::new_with(Some(Arc::new(part_a_latest)))),
);
let old_manifest = Arc::new(ManifestSnapshot {
superfile_list: SuperfileList {
manifest_id: 0,
options: opts.clone(),
superfiles: vec![
sf_a_old_keep.clone(),
sf_a_old_remove.clone(),
sf_a_latest.clone(),
],
vector_index_storage_prefix: None,
},
list: Some(list),
parts: parts_map,
loader: Some(Arc::new(loader)),
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
});
let (new_manifest, parts_to_write) = old_manifest
.update(&[], from_ref(&sf_a_old_remove))
.await
.expect("update");
let list_entries = new_manifest.get_all_list_entries();
assert_eq!(list_entries.len(), 2);
assert_eq!(parts_to_write.len(), 1);
let all_ids: Vec<_> = parts_to_write
.iter()
.flat_map(|ep| ep.part.superfiles.iter())
.map(|s| s.superfile_id)
.collect();
assert!(
all_ids.contains(&sf_a_old_keep.superfile_id),
"sf_a_old_keep must survive"
);
assert!(
!all_ids.contains(&sf_a_old_remove.superfile_id),
"sf_a_old_remove must be absent"
);
assert_eq!(list_entries[0].n_superfiles, 1);
assert_eq!(list_entries[1].n_superfiles, 1);
}
fn list_with_parts(n_parts: usize) -> list::Manifest {
use list::{Manifest, ManifestPartEntry, PartitionStrategy};
let parts = (0..n_parts)
.map(|i| ManifestPartEntry {
part_id: part::PartId(Uuid::from_u128(i as u128 + 1)),
uri: format!("manifests/part-{i}"),
n_superfiles: 0,
size_bytes_compressed: 0,
size_bytes_uncompressed: 0,
content_hash: part::ContentHash([0u8; 32]),
routing: None,
id_range: (0, 0),
scalar_stats_agg: Default::default(),
fts_summary_agg: Default::default(),
})
.collect();
Manifest {
drained_ranges: Default::default(),
global_vector_index: None,
tombstone_seqs: Default::default(),
format_version: list::FORMAT_VERSION.into(),
manifest_id: 1,
options_hash: part::ContentHash([0u8; 32]),
schema: Vec::new(),
id_column: "_id".into(),
fts_columns: vec![],
vector_columns: vec![],
partition_strategy: PartitionStrategy::Hash {
column: "_id".into(),
n_buckets: 1,
},
vector_index_storage_prefix: None,
deleted_user_ids_inline: None,
slow_vector_state_uri: None,
slow_vector_state_content_hash: None,
slow_vector_state_centroids: None,
parts,
}
}
fn manifest_with_list(list: list::Manifest) -> ManifestSnapshot {
ManifestSnapshot {
superfile_list: SuperfileList::empty(opts()),
list: Some(list),
parts: DashMap::new(),
loader: None,
stamped_partition_strategy: None,
stamped_global_vector_index: None,
stamped_drained_ranges: None,
}
}
#[test]
fn list_accessors_read_from_attached_list() {
let m = manifest_with_list(list_with_parts(3));
assert_eq!(m.get_num_parts(), 3);
assert_eq!(m.get_all_list_entries().len(), 3);
assert_eq!(m.get_num_parts_loaded(), 0, "nothing eagerly loaded");
assert!(!m.is_in_process_only(), "a list is attached");
let empty = ManifestSnapshot::empty(opts());
assert_eq!(empty.get_num_parts(), 0);
assert!(empty.get_all_list_entries().is_empty());
assert!(empty.is_in_process_only());
}
#[test]
fn complete_flat_superfiles_rejects_partial_part_view() {
let mut list = list_with_parts(1);
list.parts[0].n_superfiles = 1;
let mut manifest = manifest_with_list(list);
assert!(
manifest.complete_flat_superfiles().is_none(),
"empty resident view cannot represent one listed superfile"
);
manifest
.superfile_list
.superfiles
.push(seg_entry(Uuid::new_v4(), 4));
let complete = manifest
.complete_flat_superfiles()
.expect("resident count now matches list");
assert_eq!(complete.len(), 1);
}
#[test]
fn cached_part_lookups_miss_before_load() {
let m = manifest_with_list(list_with_parts(2));
let known_id = part::PartId(Uuid::from_u128(1));
assert!(m.get_cached_part_by_id(&known_id).is_none());
assert!(m.get_cached_part_by_list_idx(0).is_none());
assert!(m.get_cached_part_by_list_idx(1).is_none());
let empty = ManifestSnapshot::empty(opts());
assert!(empty.get_cached_part_by_list_idx(0).is_none());
}
#[test]
fn manifest_new_without_storage_is_in_process_only() {
let m = ManifestSnapshot::new(7, opts(), vec![seg_entry(Uuid::new_v4(), 4)], None, None);
assert_eq!(m.get_manifest_id(), 7);
assert!(m.is_in_process_only());
assert_eq!(m.get_num_parts(), 0);
assert_eq!(m.superfiles.len(), 1);
}
#[test]
fn from_fp32_handles_non_finite_components() {
let centroids = [f32::INFINITY, f32::NEG_INFINITY, 0.0, 1.0];
let cc = ClusterCentroids::from_fp32(1, 4, ¢roids, vec![1]);
let out = cc.centroid(0);
assert!(out.iter().all(|v| v.is_finite()));
assert_eq!(out[0], 0.0);
assert_eq!(out[1], 0.0);
assert_eq!(out[2], 0.0);
assert_eq!(out[3], 1.0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn decode_part_off_thread_roundtrips_and_rejects_garbage() {
let id = Uuid::new_v4();
let seg = Arc::new(SuperfileEntry {
birth_version: 0,
superfile_id: id,
uri: SuperfileUri(id),
n_docs: 3,
id_min: -5,
id_max: 7,
scalar_stats: HashMap::new(),
fts_summary: HashMap::new(),
vector_summary: HashMap::new(),
partition_key: Vec::new(),
partition_hint: None,
vector_layout: VectorLayout::Ivf,
subsection_offsets: None,
});
let part = ManifestPart {
format_version: part::FORMAT_VERSION.into(),
part_id: PartId::new_v4(),
superfiles: vec![seg],
};
let bytes = part::encode(&part);
let decoded = decode_part_off_thread(Bytes::from(bytes))
.await
.expect("valid part decodes off-thread");
assert_eq!(decoded.part_id, part.part_id, "part_id round-trips");
assert_eq!(decoded.superfiles.len(), 1);
assert_eq!(decoded.superfiles[0].superfile_id, id);
assert_eq!(decoded.superfiles[0].id_min, -5);
assert_eq!(decoded.superfiles[0].id_max, 7);
let err = decode_part_off_thread(Bytes::from_static(b"not-a-valid-part-blob"))
.await
.expect_err("garbage bytes must fail to decode");
assert!(
matches!(err, ManifestLoadError::Parse(_)),
"expected a parse error, got {err:?}"
);
}
}