use std::{
borrow::Cow,
cmp::Ordering,
collections::{BTreeMap, BTreeSet, BinaryHeap, HashMap, HashSet},
future::Future,
mem,
sync::{
Arc, Mutex, PoisonError,
atomic::{self, AtomicU64},
},
time::Instant,
};
use arrow::record_batch::RecordBatch;
use arrow_array::{Array, Decimal128Array};
use arrow_schema::Schema;
use futures::future::try_join_all;
use roaring::RoaringBitmap;
use tokio::{join, sync::OnceCell};
use uuid::Uuid;
use super::{
SuperfileHit,
candidate::CandidatePlan,
dispatch,
exec::common::{SCORE_COLUMN, id_score_batch, resolve_hits_named, take_rows_byte_source},
prune::{PruneLeaf, select_superfiles},
};
pub use crate::superfile::reader::VectorSearchOptions;
#[cfg(feature = "test-helpers")]
use crate::test_helpers::{admit_trace, served_shortlist_probe};
use crate::{
config,
runtime_bridge::run_on_pool,
runtime_metrics::op_stats::{self, OpStatsCollector},
storage::io_counters,
superfile::{
SuperfileReader,
error::ReadError,
fts::reader::BoolMode,
vector::{
distance::{Metric, distance, normalize, relative_score_window},
layout::VectorLayout,
reader::ScanCandidate,
},
},
supertable::{
error::QueryError,
handle::{Supertable, SupertableReader},
manifest::{
ManifestSnapshot, RABITQ_ADMIT_CELL_SHORTLIST_FRACTION,
RABITQ_ADMIT_CELL_SHORTLIST_MIN, RabitqAdmitQuery, SuperfileEntry, SuperfileUri,
VectorSummary,
list::{CellRoutingParams, PartitionStrategy},
},
opann::REPLICA_CLOSURE_DISTANCE_RATIO,
options::{GappedPlacementCell, GappedPlacementIndex},
slow_vector_state::{CentroidSection, fetch_centroid_section},
tombstones::SidecarCache,
},
};
const LAW_WIDTH_WITHIN_DEFAULT: usize = 1;
const WIDTH_BUDGET_OVERSAMPLE: usize = 4;
const DELETE_REFILL_GROWTH_FACTOR: usize = 2;
test_visible! {
const USER_COARSE_CELLS: usize = 16;
}
test_visible! {
const USER_FINE_RUNS_PER_FRAGMENT: usize = 8;
}
const FILTERED_USER_CELL_NPROBE: usize = 4;
const UNION_FINE_PICKS_MIN: usize = 2;
const FILTERED_HIDDEN_CELL_NPROBE: usize = 256;
const FILTERED_HIDDEN_FINE_NPROBE: usize = 16;
fn gate_fine_candidates_by_fragment(
candidates: Vec<(usize, u32, f32, Option<u32>, u64)>,
selected: &HashSet<u32>,
selected_ordered: &[u32],
keep_floor: usize,
keep_pct: f64,
gated_target: u64,
candidate_counts: &HashMap<(usize, u32), u64>,
scored: &mut Vec<(usize, u32, f32)>,
generation_of: Option<&[u64]>,
extension_depth: Option<(&HashSet<u32>, usize)>,
) -> Vec<(usize, u32, f32)> {
let fine_keep = |available: usize, floor: usize| -> usize {
((keep_pct * available as f64).floor() as usize)
.max(floor)
.max(1)
.min(available)
};
let refill = |mut gated: Vec<(usize, u32, f32)>,
mut remaining: Vec<(usize, u32, f32, u64)>|
-> Vec<(usize, u32, f32)> {
let mut postings: u64 = gated
.iter()
.map(|(si, cluster, _)| candidate_counts.get(&(*si, *cluster)).copied().unwrap_or(0))
.sum();
if postings < gated_target {
remaining.sort_unstable_by(|a, b| {
a.2.partial_cmp(&b.2)
.unwrap_or(Ordering::Equal)
.then_with(|| (a.0, a.1).cmp(&(b.0, b.1)))
});
for (si, cluster, score, count) in remaining {
gated.push((si, cluster, score));
postings += count;
if postings >= gated_target {
break;
}
}
}
gated
};
if let Some(gen_of) = generation_of {
let mut fine_by_generation: HashMap<u64, Vec<(usize, u32, f32, u64)>> = HashMap::new();
for (si, cluster, score, cell, count) in candidates {
match cell {
Some(cell) if selected.contains(&cell) => {
let generation = gen_of.get(si).copied().unwrap_or(0);
fine_by_generation
.entry(generation)
.or_default()
.push((si, cluster, score, count));
}
Some(_) => {}
None => scored.push((si, cluster, score)),
}
}
let mut gated = Vec::new();
let mut remaining = Vec::new();
let mut generations: Vec<u64> = fine_by_generation.keys().copied().collect();
generations.sort_unstable();
for generation in generations {
let Some(mut fine) = fine_by_generation.remove(&generation) else {
continue;
};
fine.sort_unstable_by(|a, b| {
a.2.partial_cmp(&b.2)
.unwrap_or(Ordering::Equal)
.then_with(|| (a.0, a.1).cmp(&(b.0, b.1)))
});
let keep = fine_keep(fine.len(), keep_floor);
let tail = fine.split_off(keep);
gated.extend(
fine.into_iter()
.map(|(si, cluster, score, _)| (si, cluster, score)),
);
remaining.extend(tail);
}
return refill(gated, remaining);
}
let mut fine_by_fragment: HashMap<(u32, usize), Vec<(u32, f32, u64)>> = HashMap::new();
for (si, cluster, score, cell, count) in candidates {
match cell {
Some(cell) if selected.contains(&cell) => fine_by_fragment
.entry((cell, si))
.or_default()
.push((cluster, score, count)),
Some(_) => {}
None => scored.push((si, cluster, score)),
}
}
let mut gated = Vec::new();
let mut remaining = Vec::new();
for &cell in selected_ordered {
let cell_floor = match extension_depth {
Some((extension, depth)) if extension.contains(&cell) => depth,
_ => keep_floor,
};
let mut fragment_ids: Vec<usize> = fine_by_fragment
.keys()
.filter_map(|(candidate_cell, si)| (*candidate_cell == cell).then_some(*si))
.collect();
fragment_ids.sort_unstable();
for si in fragment_ids {
let Some(mut fine) = fine_by_fragment.remove(&(cell, si)) else {
continue;
};
fine.sort_unstable_by(|a, b| {
a.1.partial_cmp(&b.1)
.unwrap_or(Ordering::Equal)
.then_with(|| a.0.cmp(&b.0))
});
let keep = fine_keep(fine.len(), cell_floor);
let tail = fine.split_off(keep);
gated.extend(
fine.into_iter()
.map(|(cluster, score, _)| (si, cluster, score)),
);
remaining.extend(
tail.into_iter()
.map(|(cluster, score, count)| (si, cluster, score, count)),
);
}
}
refill(gated, remaining)
}
fn cells_ranked_by_fine_score(
candidates: &[(usize, u32, f32, Option<u32>, u64)],
) -> Vec<(u32, f32)> {
let mut best: HashMap<u32, f32> = HashMap::new();
for &(_, _, score, cell, _) in candidates {
if let Some(cell) = cell {
best.entry(cell)
.and_modify(|s| *s = s.min(score))
.or_insert(score);
}
}
let mut ranked: Vec<(u32, f32)> = best.into_iter().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
}
fn postings_by_cell_from_summaries(
superfiles: &[Arc<SuperfileEntry>],
column: &str,
allow: Option<&HashMap<SuperfileUri, Arc<RoaringBitmap>>>,
superseded: &BTreeMap<Uuid, BTreeSet<u32>>,
) -> (HashMap<u32, u64>, bool) {
let mut postings: HashMap<u32, u64> = HashMap::new();
let mut any_tagged = false;
for entry in superfiles {
if allow.is_some_and(|m| !m.contains_key(&entry.uri)) {
continue;
}
let Some(vs) = entry.vector_summary.get(column) else {
continue;
};
for cell in &vs.cells {
let Some(cell_id) = cell.cell_id else {
continue;
};
if superseded
.get(&entry.superfile_id)
.is_some_and(|s| s.contains(&cell_id))
{
continue;
}
any_tagged = true;
let n: u64 = cell.clusters.counts.iter().map(|&c| u64::from(c)).sum();
*postings.entry(cell_id).or_default() += n;
}
}
(postings, any_tagged)
}
fn admit_shortlist_window(ranked_cells: usize) -> usize {
let scaled = (ranked_cells as f64 * RABITQ_ADMIT_CELL_SHORTLIST_FRACTION).ceil() as usize;
scaled.max(RABITQ_ADMIT_CELL_SHORTLIST_MIN)
}
fn admit_extension_round(
admit_ranking: &[(u32, f32)],
admitted: &HashSet<u32>,
exact_best_by_cell: &HashMap<u32, f32>,
serve_threshold: f32,
) -> Vec<u32> {
let mut residual_floor = f32::INFINITY;
for (cell, estimate) in admit_ranking {
if let Some(exact) = exact_best_by_cell.get(cell) {
residual_floor = residual_floor.min(exact - estimate);
}
}
if !residual_floor.is_finite() {
return Vec::new();
}
admit_ranking
.iter()
.filter(|(cell, _)| !admitted.contains(cell))
.filter(|(_, estimate)| estimate + residual_floor <= serve_threshold)
.map(|(cell, _)| *cell)
.collect()
}
fn law_floor_serve_selection(
fine_ranked: &[(u32, f32)],
grid_cells: &[u32],
fine_base: usize,
serve_threshold: f32,
) -> (Vec<u32>, HashSet<u32>) {
let fine_cells: Vec<u32> = fine_ranked
.iter()
.enumerate()
.take_while(|(rank, (_, score))| *rank < fine_base || *score <= serve_threshold)
.map(|(_, (cell, _))| *cell)
.collect();
let extension: HashSet<u32> = fine_cells
.iter()
.skip(fine_base)
.copied()
.filter(|cell| !grid_cells.contains(cell))
.collect();
(union_cell_selection(grid_cells, &fine_cells), extension)
}
type FineCandidate = (usize, u32, f32, Option<u32>, u64);
struct DeferredCellRescore {
si: usize,
cell_id: Option<u32>,
flat_base: u32,
}
fn eligible_summary<'e>(
entry: &'e SuperfileEntry,
column: &str,
query_dim: usize,
) -> Result<&'e VectorSummary, QueryError> {
match entry.vector_summary.get(column) {
Some(vs) if !vs.cells.is_empty() => {
for cell in &vs.cells {
if cell.clusters.dim as usize != query_dim {
return Err(QueryError::Execute(format!(
"vector summary dimension {} for column `{column}` on superfile {} \
does not match query dimension {query_dim}",
cell.clusters.dim, entry.superfile_id,
)));
}
}
Ok(vs)
}
Some(_) => Err(QueryError::Execute(format!(
"superfile {} has no cluster centroids in its vector summary for \
column `{column}` — malformed build; refusing to degrade to a \
blind per-superfile probe",
entry.superfile_id
))),
None => Err(QueryError::Execute(format!(
"superfile {} has no vector summary for column `{column}` — \
malformed build; refusing to degrade to a blind per-superfile \
probe",
entry.superfile_id
))),
}
}
fn estimate_admit_ranking(
superfiles: &[Arc<SuperfileEntry>],
column: &str,
query_len: usize,
metric: Metric,
admit_q: &RabitqAdmitQuery,
allow: Option<&HashMap<SuperfileUri, Arc<RoaringBitmap>>>,
superseded: &BTreeMap<Uuid, BTreeSet<u32>>,
) -> Result<Vec<(u32, f32)>, QueryError> {
let eligible = |entry: &Arc<SuperfileEntry>| allow.is_none_or(|m| m.contains_key(&entry.uri));
let is_superseded = |entry: &Arc<SuperfileEntry>, cell_id: u32| {
superseded
.get(&entry.superfile_id)
.is_some_and(|s| s.contains(&cell_id))
};
let mut cell_best: HashMap<u32, f32> = HashMap::new();
for entry in superfiles.iter().filter(|e| eligible(e)) {
let vs = eligible_summary(entry, column, query_len)?;
for cell in &vs.cells {
let Some(cell_id) = cell.cell_id else {
continue;
};
if is_superseded(entry, cell_id) {
continue;
}
let Some(est) = cell.clusters.estimate_min_admit_score(metric, admit_q) else {
continue;
};
cell_best
.entry(cell_id)
.and_modify(|best| {
if est < *best {
*best = est;
}
})
.or_insert(est);
}
}
let mut ranked: Vec<(u32, f32)> = cell_best.into_iter().collect();
ranked.sort_unstable_by(|a, b| {
a.1.partial_cmp(&b.1)
.unwrap_or(Ordering::Equal)
.then_with(|| a.0.cmp(&b.0))
});
Ok(ranked)
}
fn score_fine_candidates(
superfiles: &[Arc<SuperfileEntry>],
column: &str,
query: &[f32],
metric: Metric,
admit: Option<&HashSet<u32>>,
include_untagged: bool,
allow: Option<&HashMap<SuperfileUri, Arc<RoaringBitmap>>>,
superseded: &BTreeMap<Uuid, BTreeSet<u32>>,
) -> Result<(Vec<FineCandidate>, Vec<DeferredCellRescore>), QueryError> {
let eligible = |entry: &Arc<SuperfileEntry>| allow.is_none_or(|m| m.contains_key(&entry.uri));
let is_superseded = |entry: &Arc<SuperfileEntry>, cell_id: u32| {
superseded
.get(&entry.superfile_id)
.is_some_and(|s| s.contains(&cell_id))
};
let shortlist: Option<&HashSet<u32>> = admit;
let mut candidates: Vec<FineCandidate> = Vec::new();
let mut deferred: Vec<DeferredCellRescore> = Vec::new();
for (si, entry) in superfiles.iter().enumerate() {
if !eligible(entry) {
continue;
}
let vs = eligible_summary(entry, column, query.len())?;
let mut flat_base = 0u32;
for cell in &vs.cells {
let skipped = cell.cell_id.is_some_and(|cid| is_superseded(entry, cid))
|| (cell.cell_id.is_none() && !include_untagged)
|| shortlist
.is_some_and(|keep| cell.cell_id.is_some_and(|cid| !keep.contains(&cid)));
if !skipped {
if cell.clusters.vectors_resident() {
cell.clusters
.score_clusters_into(metric, query, |local, score| {
let count = cell
.clusters
.counts
.get(local as usize)
.copied()
.unwrap_or(0) as u64;
candidates.push((si, flat_base + local, score, cell.cell_id, count));
});
} else {
deferred.push(DeferredCellRescore {
si,
cell_id: cell.cell_id,
flat_base,
});
}
}
flat_base = flat_base.saturating_add(cell.clusters.n_cent);
}
}
Ok((candidates, deferred))
}
fn union_cell_selection(grid: &[u32], fine: &[u32]) -> Vec<u32> {
let mut selected: Vec<u32> = Vec::with_capacity(grid.len() + fine.len());
for &cell in grid.iter().chain(fine) {
if !selected.contains(&cell) {
selected.push(cell);
}
}
selected
}
fn cover_k_cell_cutoff(
cutoff: usize,
ranked: &[(u32, f32)],
postings_by_cell: &HashMap<u32, u64>,
k: usize,
) -> usize {
let cell_rows = |cell: u32| postings_by_cell.get(&cell).copied().unwrap_or(0);
let mut covered: u64 = ranked[..cutoff].iter().map(|(c, _)| cell_rows(*c)).sum();
let mut cut = cutoff;
while cut < ranked.len() && covered < k as u64 {
covered += cell_rows(ranked[cut].0);
cut += 1;
}
cut
}
fn fine_first_cell_selection(fine_ranked: &[(u32, f32)], grid_top: Option<u32>) -> Vec<u32> {
let Some(&(fine_top, fine_top_score)) = fine_ranked.first() else {
return grid_top.into_iter().collect();
};
let mut cells = vec![fine_top];
if let Some(grid_top) = grid_top
&& grid_top != fine_top
{
let tie_threshold =
relative_score_window(fine_top_score, REPLICA_CLOSURE_DISTANCE_RATIO - 1.0);
let grid_top_fine_score = fine_ranked
.iter()
.find(|(cell, _)| *cell == grid_top)
.map(|(_, score)| *score);
if grid_top_fine_score.is_some_and(|score| score <= tie_threshold) {
cells.push(grid_top);
}
}
cells
}
fn vector_read_query_error(e: ReadError) -> QueryError {
if let Some(msg) = e.over_budget() {
return QueryError::OverBudget(msg.to_string());
}
QueryError::Parquet(e.to_string())
}
pub struct VectorFilter<'a> {
pub column: &'a str,
pub query: &'a str,
pub mode: BoolMode,
}
#[derive(Clone)]
#[cfg(feature = "test-helpers")]
pub struct PreparedGlobalAllow {
use_hidden_index: bool,
allow_by_uri: HashMap<SuperfileUri, Arc<RoaringBitmap>>,
}
#[derive(Clone)]
#[cfg(not(feature = "test-helpers"))]
pub(crate) struct PreparedGlobalAllow {
use_hidden_index: bool,
allow_by_uri: HashMap<SuperfileUri, Arc<RoaringBitmap>>,
}
async fn lookup_user_placements_by_id(
manifest: &ManifestSnapshot,
user_row_ids: &[i128],
op_stats: &Option<Arc<OpStatsCollector>>,
) -> Result<Vec<(Arc<SuperfileEntry>, u32)>, QueryError> {
if user_row_ids.is_empty() {
return Ok(Vec::new());
}
let id_column = manifest.options.id_column.as_str();
let entries = manifest
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?;
let mut placements: Vec<Option<(Arc<SuperfileEntry>, u32)>> = vec![None; user_row_ids.len()];
let mut gapped = Vec::new();
let mut live_uris: HashSet<SuperfileUri> = HashSet::new();
for entry in entries {
live_uris.insert(entry.uri);
let matching: Vec<usize> = user_row_ids
.iter()
.enumerate()
.filter_map(|(index, &id)| {
(placements[index].is_none() && id >= entry.id_min && id <= entry.id_max)
.then_some(index)
})
.collect();
if matching.is_empty() {
continue;
}
if row_id_from_manifest_entry(&entry, 0).is_some() {
for index in matching {
let local = u32::try_from(user_row_ids[index] - entry.id_min).map_err(|_| {
QueryError::Execute(format!(
"local_doc_id out of range for id {}",
user_row_ids[index]
))
})?;
placements[index] = Some((Arc::clone(&entry), local));
}
} else {
gapped.push(entry);
}
}
if !gapped.is_empty() {
let slot = Arc::clone(&manifest.options.gapped_id_placement_cache);
let version = manifest.get_manifest_id();
let cells: Vec<(Arc<SuperfileEntry>, GappedPlacementCell)> = {
let mut cache = slot.lock().await;
if version > cache.pruned_through {
cache.entries.retain(|uri, _| live_uris.contains(uri));
cache.pruned_through = version;
}
gapped
.iter()
.map(|e| {
let cell = Arc::clone(
cache
.entries
.entry(e.uri)
.or_insert_with(|| Arc::new(OnceCell::new())),
);
(Arc::clone(e), cell)
})
.collect()
};
let built: Vec<(Arc<SuperfileEntry>, Arc<GappedPlacementIndex>)> =
try_join_all(cells.into_iter().map(|(entry, cell)| async move {
let index = cell
.get_or_try_init(|| {
build_gapped_placement_index(manifest, &entry, id_column, op_stats)
})
.await?;
Ok::<_, QueryError>((entry, Arc::clone(index)))
}))
.await?;
for (entry, index) in &built {
for (i, &id) in user_row_ids.iter().enumerate() {
if placements[i].is_none()
&& id >= entry.id_min
&& id <= entry.id_max
&& let Some(local) = index.local_for(id)
{
placements[i] = Some((Arc::clone(entry), local));
}
}
}
}
placements
.into_iter()
.enumerate()
.map(|(index, placement)| {
placement.ok_or_else(|| {
QueryError::Execute(format!("no user superfile owns id {}", user_row_ids[index]))
})
})
.collect()
}
fn id_values_from_batch(batch: &RecordBatch) -> Result<Vec<i128>, QueryError> {
batch
.column(0)
.as_any()
.downcast_ref::<Decimal128Array>()
.map(|a| a.values().to_vec())
.ok_or_else(|| QueryError::Execute("_id column missing".into()))
}
async fn build_gapped_placement_index(
manifest: &ManifestSnapshot,
entry: &SuperfileEntry,
id_column: &str,
op_stats: &Option<Arc<OpStatsCollector>>,
) -> Result<Arc<GappedPlacementIndex>, QueryError> {
let locals = Arc::new((0..entry.n_docs as u32).collect::<Vec<u32>>());
let ids = read_ids_for_locals(manifest, entry, &locals, id_column, true, op_stats).await?;
let scratch_bytes = ids.len()
* (std::mem::size_of::<i128>()
+ std::mem::size_of::<u32>()
+ std::mem::size_of::<(i128, u32)>());
let scratch = manifest
.options
.connection_memory_budget
.try_reserve(scratch_bytes)
.ok();
let pool = Arc::clone(&manifest.options.reader_pool);
let index = run_on_pool(
Some(&pool),
"gapped id-map: reader pool dropped result",
move || {
let mut pairs: Vec<(i128, u32)> = ids.into_iter().zip(locals.iter().copied()).collect();
pairs.sort_unstable_by(|a, b| a.0.cmp(&b.0).then_with(|| a.1.cmp(&b.1)));
pairs.dedup_by_key(|p| p.0);
let mut sorted_ids = Vec::with_capacity(pairs.len());
let mut sorted_locals = Vec::with_capacity(pairs.len());
for (id, local) in pairs {
sorted_ids.push(id);
sorted_locals.push(local);
}
drop(scratch);
GappedPlacementIndex::new(sorted_ids, sorted_locals)
},
)
.await
.map_err(|e| QueryError::Execute(e.to_string()))?;
Ok(Arc::new(index))
}
pub(crate) fn row_id_from_manifest_entry(
entry: &SuperfileEntry,
local_doc_id: u32,
) -> Option<i128> {
if entry.vector_layout == VectorLayout::MultiCellIvf {
return None;
}
let n_docs = i128::from(entry.n_docs);
let span = entry.id_max.checked_sub(entry.id_min)?.checked_add(1)?;
if n_docs == 0 || span != n_docs {
return None;
}
Some(entry.id_min + i128::from(local_doc_id))
}
pub(crate) async fn stable_ids_by_local_for_routing(
manifest: &ManifestSnapshot,
entry: &SuperfileEntry,
reader: &SuperfileReader,
op_stats: &Option<Arc<OpStatsCollector>>,
) -> Result<Vec<i128>, QueryError> {
if row_id_from_manifest_entry(entry, 0).is_some() {
return Ok((0..entry.n_docs as u32)
.map(|local| entry.id_min + i128::from(local))
.collect());
}
let locals = Arc::new((0..reader.n_docs() as u32).collect::<Vec<u32>>());
if let Some(ids) = reader
.vec()
.and_then(|v| v.inline_stable_ids_for_locals(&locals))
{
return Ok(ids);
}
let id_column = reader.id_column();
if reader.parquet_bytes().is_some() {
let batch = reader
.take_by_local_doc_ids(&locals, &[id_column])
.map_err(|e| QueryError::Execute(e.to_string()))?;
return id_values_from_batch(&batch);
}
read_ids_for_locals(manifest, entry, &locals, id_column, true, op_stats).await
}
async fn read_ids_for_locals(
manifest: &ManifestSnapshot,
entry: &SuperfileEntry,
local_ids: &Arc<Vec<u32>>,
id_column: &str,
allow_inline_region: bool,
op_stats: &Option<Arc<OpStatsCollector>>,
) -> Result<Vec<i128>, QueryError> {
let storage = manifest.options.storage.as_ref();
let store = Arc::clone(&manifest.options.store);
let disk_cache = manifest.options.disk_cache.as_ref();
let reader = dispatch::open_reader(&store, disk_cache, storage, entry, false).await?;
let inline_is_parquet_ordered = reader.vec().is_none_or(|v| v.n_docs() == reader.n_docs());
if allow_inline_region && inline_is_parquet_ordered {
let resident = {
let reader = Arc::clone(&reader);
let locals = Arc::clone(local_ids);
run_on_pool(
Some(&manifest.options.reader_pool),
"inline stable-id decode: reader pool dropped result",
move || {
reader
.vec()
.and_then(|v| v.inline_stable_ids_for_locals(&locals))
},
)
.await
.map_err(|e| QueryError::Execute(e.to_string()))?
};
if let Some(ids) = resident {
if let Some(stats) = op_stats {
stats.add_planned_read_ranges(1);
}
return Ok(ids);
}
if let Some(v) = reader.vec()
&& let Some(ids) = v
.inline_stable_ids_for_locals_async(local_ids)
.await
.map_err(|e| QueryError::Execute(e.to_string()))?
{
if let Some(stats) = op_stats {
stats.add_planned_read_ranges(1);
}
return Ok(ids);
}
}
if reader.parquet_bytes().is_some() {
let (batch, decode_ns) = {
let reader = Arc::clone(&reader);
let locals = Arc::clone(local_ids);
let id_column = id_column.to_string();
run_on_pool(
Some(&manifest.options.reader_pool),
"scalar id decode: reader pool dropped result",
move || {
op_stats::timed_section(|| {
reader
.take_by_local_doc_ids(&locals, &[id_column.as_str()])
.map_err(|e| QueryError::Execute(e.to_string()))
})
},
)
.await
.map_err(|e| QueryError::Execute(e.to_string()))?
};
if let Some(stats) = op_stats {
stats.add_kernel_cpu_ns(decode_ns);
}
return id_values_from_batch(&batch?);
}
let batch = take_rows_byte_source(&reader, local_ids, &[id_column])
.await
.map_err(|error| QueryError::Execute(error.to_string()))?;
id_values_from_batch(&batch)
}
async fn hidden_hits_user_ids(
hidden_manifest: &ManifestSnapshot,
hidden_hits: &[SuperfileHit],
id_column: &str,
op_stats: &Option<Arc<OpStatsCollector>>,
) -> Result<Vec<i128>, QueryError> {
let mut ids = vec![0i128; hidden_hits.len()];
let mut by_superfile: HashMap<SuperfileUri, Vec<usize>> = HashMap::new();
for (i, hit) in hidden_hits.iter().enumerate() {
if let Some(id) = hit.stable_id {
ids[i] = id;
continue;
}
by_superfile.entry(hit.superfile).or_default().push(i);
}
for (uri, idxs) in by_superfile {
let entry = hidden_manifest
.lookup_superfile_entry(uri)
.await
.map_err(QueryError::ManifestLoad)?
.ok_or_else(|| {
QueryError::Execute(format!("hidden superfile {uri:?} missing from manifest"))
})?;
if row_id_from_manifest_entry(&entry, 0).is_some() {
for &i in &idxs {
ids[i] = entry.id_min + i128::from(hidden_hits[i].local_doc_id);
}
continue;
}
let locals = Arc::new(
idxs.iter()
.map(|&i| hidden_hits[i].local_doc_id)
.collect::<Vec<u32>>(),
);
let vals = read_ids_for_locals(hidden_manifest, &entry, &locals, id_column, true, op_stats)
.await?;
for (j, &i) in idxs.iter().enumerate() {
ids[i] = vals[j];
}
}
Ok(ids)
}
const ID_SCORE_ID_COLUMN: usize = 0;
const ID_SCORE_SCORE_COLUMN: usize = 1;
pub(crate) fn free_column_slot(name: &str, id_column: &str) -> Option<usize> {
if id_column == SCORE_COLUMN {
return None;
}
match name {
n if n == id_column => Some(ID_SCORE_ID_COLUMN),
n if n == SCORE_COLUMN => Some(ID_SCORE_SCORE_COLUMN),
_ => None,
}
}
pub(crate) fn free_columns_unambiguous(user_schema: &Schema, id_column: &str) -> bool {
id_column != SCORE_COLUMN
&& !user_schema
.fields()
.iter()
.any(|f| f.name() == SCORE_COLUMN || f.name() == id_column)
}
pub(crate) fn id_score_projection_indices(
projection: Option<&[&str]>,
id_column: &str,
) -> Option<Vec<usize>> {
let Some(names) = projection else {
return Some(vec![ID_SCORE_ID_COLUMN, ID_SCORE_SCORE_COLUMN]);
};
if names.is_empty() {
return None;
}
names
.iter()
.map(|name| free_column_slot(name, id_column))
.collect()
}
fn is_hidden_vector_manifest(manifest: &ManifestSnapshot) -> bool {
matches!(
manifest.partition_strategy(),
Some(PartitionStrategy::VectorCell { .. })
)
}
pub(crate) fn hits_id_score_batch(
user_reader: &SupertableReader,
hits: &[SuperfileHit],
) -> Result<RecordBatch, QueryError> {
let mut ids = Vec::with_capacity(hits.len());
let mut scores = Vec::with_capacity(hits.len());
for hit in hits {
let id = hit.stable_id.ok_or_else(|| {
QueryError::Execute(format!(
"hit {:?}/{} missing stable _id before output materialization",
hit.superfile, hit.local_doc_id
))
})?;
ids.push(id);
scores.push(hit.score);
}
id_score_batch(user_reader, &ids, &scores).map_err(|e| QueryError::Execute(e.to_string()))
}
pub(crate) async fn user_placement_for_scalar_resolve(
user_reader: &SupertableReader,
hits: &[SuperfileHit],
) -> Result<Vec<SuperfileHit>, QueryError> {
if hits.is_empty() {
return Ok(Vec::new());
}
let user_manifest = user_reader.manifest();
let id_column = user_reader.options().id_column.as_str();
let hidden_manifest = user_reader.vector_index_table().map(|vit| {
Arc::clone(
vit.pinned_reader_with(user_reader.op_stats.clone())
.manifest(),
)
});
let deleted = user_reader.vector_index_table().and_then(|vit| {
vit.pinned_reader_with(user_reader.op_stats.clone())
.hidden_deleted_ids()
.ok()
});
let mut out: Vec<Option<SuperfileHit>> = vec![None; hits.len()];
let mut placement_requests: Vec<(usize, i128)> = Vec::new();
for (i, hit) in hits.iter().enumerate() {
if let Some(user_entry) = user_manifest
.lookup_superfile_entry(hit.superfile)
.await
.map_err(QueryError::ManifestLoad)?
&& !(user_entry.vector_layout == VectorLayout::MultiCellIvf && hit.stable_id.is_some())
{
out[i] = Some(*hit);
continue;
}
let user_row_id = if let Some(id) = hit.stable_id {
id
} else if let Some(ref hm) = hidden_manifest {
hidden_hits_user_ids(
hm,
std::slice::from_ref(hit),
id_column,
&user_reader.op_stats,
)
.await?[0]
} else {
return Err(QueryError::Execute(format!(
"hit superfile {:?} missing from manifests",
hit.superfile
)));
};
if deleted
.as_ref()
.is_some_and(|d| d.binary_search(&user_row_id).is_ok())
{
continue;
}
placement_requests.push((i, user_row_id));
}
let requested_ids: Vec<i128> = placement_requests.iter().map(|(_, id)| *id).collect();
let placements =
lookup_user_placements_by_id(user_manifest, &requested_ids, &user_reader.op_stats).await?;
for ((index, stable_id), (entry, local_doc_id)) in
placement_requests.into_iter().zip(placements)
{
out[index] = Some(SuperfileHit {
superfile: entry.uri,
local_doc_id,
score: hits[index].score,
stable_id: Some(stable_id),
});
}
Ok(out.into_iter().flatten().collect())
}
fn score_cell_fp32(
superfiles: &[Arc<SuperfileEntry>],
column: &str,
d: &DeferredCellRescore,
fp32: &[f32],
query: &[f32],
metric: Metric,
candidates: &mut Vec<FineCandidate>,
) -> bool {
let entry = &superfiles[d.si];
let Some(cell) = entry
.vector_summary
.get(column)
.and_then(|vs| vs.cells.iter().find(|cell| cell.cell_id == d.cell_id))
else {
return false;
};
let dim = cell.clusters.dim as usize;
if dim == 0 || fp32.len() != cell.clusters.n_cent as usize * dim {
return false;
}
for (local, centroid) in fp32.chunks_exact(dim).enumerate() {
let count = cell.clusters.counts.get(local).copied().unwrap_or(0) as u64;
if count == 0 {
continue;
}
let score = distance(metric, query, centroid);
candidates.push((d.si, d.flat_base + local as u32, score, d.cell_id, count));
}
true
}
impl SupertableReader {
async fn centroid_section(&self) -> Option<Arc<CentroidSection>> {
let manifest = self.manifest();
let reference = manifest.slow_vector_state_centroids_blob()?.clone();
let storage = manifest.options.storage.as_ref()?;
let slot = Arc::clone(&manifest.options.centroid_section_cache);
let mut guard = slot.lock().await;
if let Some(section) = guard.as_ref()
&& section.uri() == reference.uri
{
return Some(Arc::clone(section));
}
let entries = manifest.get_all_superfiles();
match fetch_centroid_section(storage.as_ref(), &reference, entries).await {
Ok(section) => {
let section = Arc::new(section);
*guard = Some(Arc::clone(§ion));
Some(section)
}
Err(error) => {
eprintln!(
"[supertable] centroid section {} unavailable ({error}); deferred rescores \
will fail unless the parts cache covers their cells",
reference.uri
);
None
}
}
}
async fn rescore_deferred_cells(
&self,
superfiles: &[Arc<SuperfileEntry>],
column: &str,
query: &[f32],
metric: Metric,
candidates: &mut Vec<FineCandidate>,
deferred: Vec<DeferredCellRescore>,
) -> Result<(), QueryError> {
let deferred = if deferred.is_empty() {
deferred
} else if let Some(section) = self.centroid_section().await {
let mut cells_read = 0u64;
let (leftovers, rescore_ns) = op_stats::timed_section(|| {
let mut leftovers = Vec::new();
for d in deferred {
let entry = &superfiles[d.si];
let read = section
.read_cell(entry.superfile_id, column, d.cell_id)
.map_err(|e| {
QueryError::Execute(format!("centroid section spill read: {e}"))
})?;
let Some(fp32) = read else {
leftovers.push(d);
continue;
};
cells_read += 1;
if !score_cell_fp32(superfiles, column, &d, &fp32, query, metric, candidates) {
leftovers.push(d);
}
}
Ok::<_, QueryError>(leftovers)
});
if let Some(stats) = &self.op_stats {
stats.add_planned_read_ranges(cells_read);
stats.add_kernel_cpu_ns(rescore_ns);
}
leftovers?
} else {
deferred
};
let deferred = if deferred.is_empty() {
deferred
} else if let Some(cache) = self.manifest().user_centroids_for_rescore().await {
let (leftovers, rescore_ns) = op_stats::timed_section(|| {
let mut leftovers = Vec::new();
for d in deferred {
let entry = &superfiles[d.si];
let Some(fp32) = cache.cell(entry.superfile_id, column, d.cell_id) else {
leftovers.push(d);
continue;
};
if !score_cell_fp32(
superfiles,
column,
&d,
fp32.as_slice(),
query,
metric,
candidates,
) {
leftovers.push(d);
}
}
leftovers
});
if let Some(stats) = &self.op_stats {
stats.add_kernel_cpu_ns(rescore_ns);
}
leftovers
} else {
deferred
};
if let Some(d) = deferred.first() {
let entry = &superfiles[d.si];
return Err(QueryError::Execute(format!(
"deferred admit rescore: no manifest-published fp32 covers superfile {} column \
{column} cell {:?} ({} cell(s) uncovered) — the centroid section / full parts \
must cover every stripped summary cell",
entry.superfile_id,
d.cell_id,
deferred.len(),
)));
}
Ok(())
}
async fn fanout_vector_clusters(
&self,
superfiles: &[Arc<SuperfileEntry>],
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
) -> Result<Vec<SuperfileHit>, QueryError> {
if superfiles.is_empty() {
return Ok(Vec::new());
}
self.vector_fanout_over_superfiles(superfiles.to_vec(), column, query, k, options, None)
.await
}
async fn vector_fanout_over_superfiles(
&self,
superfiles: Vec<Arc<SuperfileEntry>>,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
allow: Option<HashMap<SuperfileUri, Arc<RoaringBitmap>>>,
) -> Result<Vec<SuperfileHit>, QueryError> {
let filtered = allow.is_some();
let (resolved_nprobe, _) = options.resolve(filtered);
let manifest = self.manifest();
let hidden_vector_index = is_hidden_vector_manifest(manifest);
let hidden_routing = manifest.vector_cell_routing();
let law_rerank = rerank_mult_from_law(
hidden_vector_index,
filtered,
options.rerank_mult(),
hidden_routing.as_ref(),
k,
);
let law_rerank_served = law_rerank.is_some();
let options = match law_rerank {
Some(mult) => options.with_rerank_mult(mult),
None => options,
};
let nprobe = if !hidden_vector_index && !filtered && options.nprobe.is_none() {
USER_COARSE_CELLS
} else {
resolved_nprobe
};
let (metric, rot_seed) = manifest
.options
.vector_columns
.iter()
.find(|vc| vc.column == column)
.map(|vc| (vc.metric, vc.rot_seed))
.ok_or_else(|| QueryError::Execute(format!("unknown vector column `{column}`")))?;
let grid = manifest
.global_vector_index()
.filter(|g| g.column == column)
.map(|g| g.user_grid())
.filter(|grid| grid.n_cent > 0 && grid.dim as usize == query.len())
.or_else(|| {
manifest
.vector_cell_clusters(column)
.filter(|clusters| clusters.n_cent > 0 && clusters.dim as usize == query.len())
});
let admit_t0 = io_counters::phase_start();
let ranked_cells_scored: Option<Vec<(u32, f32)>> =
op_stats::timed_kernel(&self.op_stats, || {
grid.map(|grid| grid.rank_cells(metric, query))
});
let ranked_cells: Option<Vec<u32>> = ranked_cells_scored
.as_ref()
.map(|cells| cells.iter().map(|(cell, _)| *cell).collect());
let grid_cell_cutoff = |ranked: &[(u32, f32)], routing: &CellRoutingParams| -> usize {
if ranked.is_empty() {
return 0;
}
let mut cutoff = routing.nprobe_min.max(1).min(ranked.len());
let max_cells = routing.nprobe_max.max(routing.nprobe_min).min(ranked.len());
let threshold = relative_score_window(ranked[0].1, routing.slack);
while cutoff < max_cells && ranked[cutoff].1 <= threshold {
cutoff += 1;
}
cutoff
};
let birth_versions: Vec<u64> = superfiles.iter().map(|e| e.birth_version).collect();
let gated_target = (k as f64
* f64::from(config::global().vector.drain_replica_target_factor.max(1.0)))
.ceil() as u64;
let allow_ref = allow.as_ref();
let empty_superseded = BTreeMap::new();
let superseded = manifest.get_superseded_cells().unwrap_or(&empty_superseded);
let (postings_by_cell, any_tagged) =
postings_by_cell_from_summaries(&superfiles, column, allow_ref, superseded);
let mut gated = Vec::new();
let mut scored = Vec::new();
let mut sweep_width: Option<usize> = None;
let candidate_counts: HashMap<(usize, u32), u64>;
let mut served_cells_over_width: (usize, usize) = (1, 1);
if let (Some(ranked_scored), true) = (&ranked_cells_scored, any_tagged) {
let mut cell_routing = if hidden_vector_index {
let base = hidden_routing.ok_or_else(|| {
QueryError::Execute("hidden manifest missing cell routing".into())
})?;
if filtered {
CellRoutingParams {
nprobe_min: base.nprobe_min.max(FILTERED_HIDDEN_CELL_NPROBE),
nprobe_max: base.nprobe_max.max(FILTERED_HIDDEN_CELL_NPROBE),
fine_nprobe: base.fine_nprobe.max(FILTERED_HIDDEN_FINE_NPROBE),
..base
}
} else {
base
}
} else if filtered {
CellRoutingParams {
nprobe_min: FILTERED_USER_CELL_NPROBE,
nprobe_max: FILTERED_USER_CELL_NPROBE,
..CellRoutingParams::default()
}
} else {
CellRoutingParams::default()
};
if hidden_vector_index
&& !filtered
&& let Some(fine) = hidden_routing.and_then(|r| r.fine_for_k_at(k))
{
cell_routing.fine_nprobe = cell_routing.fine_nprobe.max(fine);
}
let law_width: Option<usize> =
if hidden_vector_index && !filtered && options.nprobe.is_none() {
hidden_routing
.and_then(|r| r.width_for_k_at(k))
.filter(|w| *w > LAW_WIDTH_WITHIN_DEFAULT)
} else {
None
};
let prepin_fine_depth = cell_routing.fine_nprobe;
let populated_cells = postings_by_cell.len().max(1);
sweep_width = apply_width_pin(
&mut cell_routing,
options.nprobe.map(|n| n.max(1)),
law_width,
filtered,
populated_cells,
);
let fine_nprobe_pct = if filtered {
0.0
} else {
config::global().vector.fine_nprobe_pct
};
let serve_near_tie_slack = config::global().vector.serve_near_tie_slack;
let ranked_for_beam: Vec<(u32, f32)> = ranked_scored
.iter()
.filter(|(cell, _)| postings_by_cell.contains_key(cell))
.copied()
.collect();
if ranked_for_beam.is_empty() {
return Err(QueryError::Execute(
"vector candidates name no cell present in the grid — \
malformed cell tags"
.into(),
));
}
let cutoff = grid_cell_cutoff(&ranked_for_beam, &cell_routing);
let cutoff = if options.nprobe.is_none() {
cover_k_cell_cutoff(cutoff, &ranked_for_beam, &postings_by_cell, k)
} else {
cutoff
};
let admit_q = RabitqAdmitQuery::new(query.len(), rot_seed, query);
let must_include: Vec<u32> = ranked_for_beam[..cutoff]
.iter()
.map(|(cell, _)| *cell)
.collect();
let admit_ranking = estimate_admit_ranking(
&superfiles,
column,
query.len(),
metric,
&admit_q,
allow_ref,
superseded,
)?;
let mut admitted: HashSet<u32> = admit_ranking
.iter()
.take(admit_shortlist_window(admit_ranking.len()))
.map(|(cell, _)| *cell)
.collect();
admitted.extend(must_include.iter().copied());
let (mut candidates, deferred) = score_fine_candidates(
&superfiles,
column,
query,
metric,
Some(&admitted),
true,
allow_ref,
superseded,
)?;
if !deferred.is_empty() {
self.rescore_deferred_cells(
&superfiles,
column,
query,
metric,
&mut candidates,
deferred,
)
.await?;
}
let law_default = !filtered && options.nprobe.is_none() && law_width.is_some();
if law_default {
loop {
let fine_ranked_now = cells_ranked_by_fine_score(&candidates);
let Some(&(_, best_exact)) = fine_ranked_now.first() else {
break;
};
let serve_threshold = relative_score_window(best_exact, serve_near_tie_slack);
let exact_best_by_cell: HashMap<u32, f32> =
fine_ranked_now.into_iter().collect();
let round = admit_extension_round(
&admit_ranking,
&admitted,
&exact_best_by_cell,
serve_threshold,
);
if round.is_empty() {
break;
}
let delta: HashSet<u32> = round.into_iter().collect();
admitted.extend(delta.iter().copied());
let (mut delta_candidates, delta_deferred) = score_fine_candidates(
&superfiles,
column,
query,
metric,
Some(&delta),
false,
allow_ref,
superseded,
)?;
if !delta_deferred.is_empty() {
self.rescore_deferred_cells(
&superfiles,
column,
query,
metric,
&mut delta_candidates,
delta_deferred,
)
.await?;
}
candidates.extend(delta_candidates);
}
}
#[cfg(feature = "test-helpers")]
admit_trace::record_admit(admitted.iter().copied().collect());
candidate_counts = candidates
.iter()
.map(|(si, cluster, _, _, count)| ((*si, *cluster), *count))
.collect();
let ranked = ranked_cells
.as_ref()
.expect("ranked cell ids exist with scored ranking");
if hidden_vector_index {
let fine_ranked = cells_ranked_by_fine_score(&candidates);
#[cfg(feature = "test-helpers")]
admit_trace::record_fine(fine_ranked.clone());
let default_p1 = !filtered && options.nprobe.is_none() && cutoff == 1;
let (selected_cells_ordered, extension_cells): (Vec<u32>, HashSet<u32>) =
if default_p1 {
(
fine_first_cell_selection(
&fine_ranked,
ranked_for_beam.first().map(|(cell, _)| *cell),
),
HashSet::new(),
)
} else if law_default {
let grid_cells: Vec<u32> = ranked_for_beam[..cutoff]
.iter()
.map(|(cell, _)| *cell)
.collect();
match fine_ranked
.first()
.map(|(_, score)| relative_score_window(*score, serve_near_tie_slack))
{
Some(threshold) => law_floor_serve_selection(
&fine_ranked,
&grid_cells,
cutoff.max(UNION_FINE_PICKS_MIN),
threshold,
),
None => (grid_cells, HashSet::new()),
}
} else {
let grid_cells: Vec<u32> = ranked_for_beam[..cutoff]
.iter()
.map(|(cell, _)| *cell)
.collect();
let fine_cells: Vec<u32> = fine_ranked
.iter()
.take(cutoff.max(UNION_FINE_PICKS_MIN))
.map(|(cell, _)| *cell)
.collect();
(
union_cell_selection(&grid_cells, &fine_cells),
HashSet::new(),
)
};
if law_default {
served_cells_over_width = (
selected_cells_ordered.len().max(1),
sweep_width.unwrap_or(1).max(1),
);
}
let selected_cells: HashSet<u32> = selected_cells_ordered.iter().copied().collect();
let generation_of = if sweep_width.is_some() {
None
} else {
Some(birth_versions.as_slice())
};
gated = gate_fine_candidates_by_fragment(
candidates,
&selected_cells,
&selected_cells_ordered,
cell_routing.fine_nprobe,
fine_nprobe_pct,
gated_target,
&candidate_counts,
&mut scored,
generation_of,
(!extension_cells.is_empty()).then_some((&extension_cells, prepin_fine_depth)),
);
} else {
let fine_ranked = cells_ranked_by_fine_score(&candidates);
let default_p1 = !filtered && options.nprobe.is_none() && cutoff == 1;
let mut selected_cells: Vec<u32> = if default_p1 && !fine_ranked.is_empty() {
fine_first_cell_selection(&fine_ranked, ranked.first().copied())
} else {
let grid_cells: Vec<u32> = ranked[..cutoff].to_vec();
let fine_cells: Vec<u32> = fine_ranked
.iter()
.take(cutoff)
.map(|(cell, _)| *cell)
.collect();
union_cell_selection(&grid_cells, &fine_cells)
};
let mut covered: u64 = selected_cells
.iter()
.map(|cell| postings_by_cell.get(cell).copied().unwrap_or(0))
.sum();
for cell in ranked.iter().copied() {
if covered >= gated_target {
break;
}
if selected_cells.contains(&cell) {
continue;
}
covered += postings_by_cell.get(&cell).copied().unwrap_or(0);
selected_cells.push(cell);
}
let selected: HashSet<u32> = selected_cells.iter().copied().collect();
gated = gate_fine_candidates_by_fragment(
candidates,
&selected,
&selected_cells,
USER_FINE_RUNS_PER_FRAGMENT,
0.0, gated_target,
&candidate_counts,
&mut scored,
None,
None,
);
}
} else {
let (mut candidates, deferred) = score_fine_candidates(
&superfiles,
column,
query,
metric,
None,
true,
allow_ref,
superseded,
)?;
if !deferred.is_empty() {
self.rescore_deferred_cells(
&superfiles,
column,
query,
metric,
&mut candidates,
deferred,
)
.await?;
}
candidate_counts = candidates
.iter()
.map(|(si, cluster, _, _, count)| ((*si, *cluster), *count))
.collect();
scored = candidates
.into_iter()
.map(|(si, cluster, score, _, _)| (si, cluster, score))
.collect();
}
let n_eligible = {
let mut segs: Vec<usize> = scored
.iter()
.chain(gated.iter())
.map(|&(si, _, _)| si)
.collect();
segs.sort_unstable();
segs.dedup();
segs.len()
};
let scaled_budget = nprobe.saturating_mul(n_eligible.max(1)).max(nprobe);
let default_budget = if hidden_vector_index {
hidden_routing
.expect("hidden manifest carries routing")
.fine_nprobe
.max(1)
} else {
scaled_budget
};
let budget = if hidden_vector_index {
default_budget
} else {
config::global()
.vector
.inner_budget
.map(|value| value.max(1))
.unwrap_or(default_budget)
};
let cluster_count = |&(si, cluster, _): &(usize, u32, f32)| -> u64 {
candidate_counts.get(&(si, cluster)).copied().unwrap_or(0)
};
let gated_postings: u64 = gated.iter().map(cluster_count).sum();
if scored.len() > budget {
scored.sort_unstable_by(|a, b| {
a.2.partial_cmp(&b.2)
.unwrap_or(Ordering::Equal)
.then_with(|| (a.0, a.1).cmp(&(b.0, b.1)))
});
let mut kept = budget;
let mut postings =
gated_postings + scored[..kept].iter().map(cluster_count).sum::<u64>();
while kept < scored.len() && postings < k as u64 {
postings += cluster_count(&scored[kept]);
kept += 1;
}
scored.truncate(kept);
}
let mut per_seg: HashMap<usize, Vec<u32>> = HashMap::new();
for (si, c, _) in scored.into_iter().chain(gated) {
per_seg.entry(si).or_default().push(c);
}
let mut units: Vec<(
Arc<SuperfileEntry>,
(usize, Vec<u32>, Option<Arc<RoaringBitmap>>),
)> = Vec::new();
for (si, entry) in superfiles.iter().enumerate() {
let Some(ids) = per_seg.remove(&si) else {
continue;
};
let bitmap = match allow.as_ref() {
Some(m) => match m.get(&entry.uri) {
Some(bm) => Some(Arc::clone(bm)),
None => continue,
},
None => None,
};
units.push((Arc::clone(entry), (si, ids, bitmap)));
}
if units.is_empty() {
if let Some(t0) = admit_t0 {
io_counters::phase_record("vec.admit", t0.elapsed().as_micros() as u64);
}
return Ok(Vec::new());
}
if let Some(t0) = admit_t0 {
io_counters::phase_record("vec.admit", t0.elapsed().as_micros() as u64);
}
let global_shortlist_width = if hidden_vector_index {
sweep_width.filter(|w| *w > 1)
} else {
None
};
let (_, plan_rerank_mult) = options.resolve(filtered);
let mut cold_rerank_mult = 0;
let options = match sweep_width {
Some(w) if w > 1 && options.nprobe.is_some() => {
if global_shortlist_width.is_some() {
let (_, rerank_mult) = options.resolve(filtered);
cold_rerank_mult = rerank_mult;
}
options
}
Some(w) if w > 1 => {
let (_, rerank_mult) = options.resolve(filtered);
let divided = rerank_mult
.saturating_mul(WIDTH_BUDGET_OVERSAMPLE)
.div_ceil(w)
.max(1);
if global_shortlist_width.is_some() {
cold_rerank_mult = divided;
options
} else {
options.with_rerank_mult(divided)
}
}
_ => options,
};
let column_arc = Arc::new(column.to_owned());
let query_arc = Arc::new(query.to_vec());
let column_arc2 = Arc::clone(&column_arc);
let query_arc2 = Arc::clone(&query_arc);
let op_stats_scan = self.op_stats.clone();
let reader_pool = Arc::clone(&manifest.options.reader_pool);
let budget = Some(Arc::clone(&manifest.options.connection_memory_budget));
let storage = manifest.options.storage.as_ref().map(Arc::clone);
let scan_pool: Arc<Mutex<Vec<(usize, u64, usize, Vec<ScanCandidate>)>>> =
Arc::new(Mutex::new(Vec::new()));
let scan_pool_body = Arc::clone(&scan_pool);
let max_replica_overhead = Arc::new(AtomicU64::new(0));
let max_replica_overhead_body = Arc::clone(&max_replica_overhead);
let body =
move |reader: Arc<SuperfileReader>,
entry: Arc<SuperfileEntry>,
tombstone_cache: Option<Arc<SidecarCache>>,
now: Instant,
(si, ids, bitmap): (usize, Vec<u32>, Option<Arc<RoaringBitmap>>)| {
let column = Arc::clone(&column_arc);
let query = Arc::clone(&query_arc);
let reader_pool = Arc::clone(&reader_pool);
let budget = budget.clone();
let storage = storage.clone();
let scan_pool = Arc::clone(&scan_pool_body);
let max_replica_overhead = Arc::clone(&max_replica_overhead_body);
let op_stats = op_stats_scan.clone();
async move {
let deny_pushdown = !hidden_vector_index
&& bitmap.is_none()
&& entry.vector_layout != VectorLayout::MultiCellIvf;
let deny = match tombstone_cache.as_ref() {
Some(cache) if deny_pushdown => {
dispatch::tombstone_deny_set(cache, entry.superfile_id, now)?
}
_ => None,
};
let pool = Some(Arc::clone(&reader_pool));
let replica_overhead = reader
.vec()
.map(|v| (v.n_docs() as usize).saturating_sub(reader.n_docs() as usize))
.unwrap_or(0);
let k_fetch = k.saturating_add(replica_overhead);
let reader_for_ids = Arc::clone(&reader);
let hits = if global_shortlist_width.is_some() {
let scan = reader
.vector_scan_clusters_filtered(
&column,
&query,
k_fetch,
&ids,
options,
cold_rerank_mult,
bitmap,
deny,
pool,
budget,
)
.await
.map_err(vector_read_query_error)?;
if let Some(stats) = &op_stats {
stats.add_vector_scan(scan.cells_scanned, scan.candidates_scanned);
stats.add_planned_read_ranges(scan.ranges_requested);
stats.add_vector_rows_reranked(scan.rows_reranked);
stats.add_kernel_cpu_ns(scan.kernel_cpu_ns);
}
max_replica_overhead
.fetch_max(replica_overhead as u64, atomic::Ordering::Relaxed);
if !scan.candidates.is_empty() {
scan_pool
.lock()
.unwrap_or_else(PoisonError::into_inner)
.push((si, scan.rot_seed, replica_overhead, scan.candidates));
}
scan.hits
} else {
let (hits, tally) = reader
.vector_search_clusters_filtered(
&column, &query, k_fetch, &ids, options, bitmap, deny, pool, budget,
)
.await
.map_err(vector_read_query_error)?;
if let Some(stats) = &op_stats {
stats.add_vector_scan(tally.cells_scanned, tally.candidates_scanned);
stats.add_planned_read_ranges(tally.ranges_requested);
stats.add_vector_rows_reranked(tally.rows_reranked);
stats.add_kernel_cpu_ns(tally.kernel_cpu_ns);
}
hits
};
let mut tagged = dispatch::tag_hits(&entry, hits);
io_counters::phase_timed_async("vec.stable_id", async {
dispatch::attach_stable_ids(
&reader_for_ids,
&entry,
&mut tagged,
false,
&op_stats,
)
.await
})
.await?;
if !hidden_vector_index && !deny_pushdown {
dispatch::apply_resolved_tombstone_filter(
&reader_for_ids,
storage.as_ref(),
tombstone_cache.as_ref(),
&entry,
&mut tagged,
now,
&op_stats,
)
.await?;
}
Ok::<Vec<SuperfileHit>, QueryError>(tagged)
}
};
let fanout_t0 = io_counters::phase_start();
let mut per_superfile = if allow.is_some() {
let fanout_width = manifest.options.reader_pool.current_num_threads().max(1);
let mut collected = Vec::new();
while !units.is_empty() {
let n = fanout_width.min(units.len());
let wave: Vec<_> = units.drain(..n).collect();
collected.extend(
dispatch::fanout_with(self, wave, !hidden_vector_index, false, body.clone())
.await?,
);
}
collected
} else {
dispatch::fanout_with(self, units, !hidden_vector_index, false, body).await?
};
if global_shortlist_width.is_some() {
let pooled = {
let mut guard = scan_pool.lock().unwrap_or_else(PoisonError::into_inner);
mem::take(&mut *guard)
};
if !pooled.is_empty() {
if pooled.windows(2).any(|w| w[0].1 != w[1].1) {
return Err(QueryError::Execute(
"pooled 1-bit estimates require one rotation seed per column".into(),
));
}
let replica_overhead =
usize::try_from(max_replica_overhead.load(atomic::Ordering::Relaxed))
.unwrap_or(0);
let mut flat: Vec<(usize, ScanCandidate)> = pooled
.into_iter()
.flat_map(|(si, _, _, cands)| cands.into_iter().map(move |c| (si, c)))
.collect();
let cell_floor = if options.nprobe.is_some() {
k.saturating_mul(plan_rerank_mult)
} else {
0
};
let shortlist_limit = deferred_shortlist_limit(
k,
replica_overhead,
plan_rerank_mult,
law_rerank_served,
options.nprobe.is_some(),
served_cells_over_width,
);
#[cfg(feature = "test-helpers")]
served_shortlist_probe::record(shortlist_limit, cell_floor);
flat = select_global_shortlist(flat, shortlist_limit, cell_floor);
let mut winners_by_seg: HashMap<usize, Vec<ScanCandidate>> = HashMap::new();
for (si, cand) in flat {
winners_by_seg.entry(si).or_default().push(cand);
}
let rerank_units: Vec<(Arc<SuperfileEntry>, Vec<ScanCandidate>)> = superfiles
.iter()
.enumerate()
.filter_map(|(si, entry)| {
winners_by_seg
.remove(&si)
.map(|sel| (Arc::clone(entry), sel))
})
.collect();
if let Some(stats) = &self.op_stats {
let rows: u64 = rerank_units.iter().map(|(_, sel)| sel.len() as u64).sum();
stats.add_vector_rows_reranked(rows);
}
let column = Arc::clone(&column_arc2);
let query = Arc::clone(&query_arc2);
let reader_pool = Arc::clone(&manifest.options.reader_pool);
let op_stats_c = self.op_stats.clone();
let body_c = move |reader: Arc<SuperfileReader>,
entry: Arc<SuperfileEntry>,
_tombstone_cache: Option<Arc<SidecarCache>>,
_now: Instant,
selected: Vec<ScanCandidate>| {
let column = Arc::clone(&column);
let query = Arc::clone(&query);
let reader_pool = Arc::clone(&reader_pool);
let op_stats = op_stats_c.clone();
async move {
let replica_overhead = reader
.vec()
.map(|v| (v.n_docs() as usize).saturating_sub(reader.n_docs() as usize))
.unwrap_or(0);
let k_fetch = k.saturating_add(replica_overhead);
let reader_for_ids = Arc::clone(&reader);
let (hits, rerank_kernel_ns) = reader
.vector_rerank_selected(
&column,
&query,
k_fetch,
selected,
Some(reader_pool),
)
.await
.map_err(vector_read_query_error)?;
if let Some(stats) = &op_stats {
stats.add_kernel_cpu_ns(rerank_kernel_ns);
}
let mut tagged = dispatch::tag_hits(&entry, hits);
io_counters::phase_timed_async("vec.stable_id", async {
dispatch::attach_stable_ids(
&reader_for_ids,
&entry,
&mut tagged,
false,
&op_stats,
)
.await
})
.await?;
Ok::<Vec<SuperfileHit>, QueryError>(tagged)
}
};
per_superfile
.extend(dispatch::fanout_with(self, rerank_units, false, false, body_c).await?);
}
}
if let Some(t0) = fanout_t0 {
io_counters::phase_record("vec.fanout_wall", t0.elapsed().as_micros() as u64);
}
Ok(top_k_ascending(per_superfile, k))
}
#[cfg_attr(
feature = "detailed-tracing",
tracing::instrument(skip_all, fields(column = column, k = k, dim = query.len()))
)]
pub(crate) async fn vector_hits_filtered_async(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
filter: VectorFilter<'_>,
) -> Result<Vec<SuperfileHit>, QueryError> {
if k == 0 {
return Ok(Vec::new());
}
let manifest = self.manifest();
let Some(tokenizer) = manifest.options.tokenizer.as_ref() else {
return Ok(Vec::new());
};
let tokens: Vec<String> = tokenizer.tokenize(filter.query).collect();
if tokens.is_empty() {
return Ok(Vec::new());
}
let prune_leaves = [PruneLeaf::TermPresence {
column: filter.column.to_owned(),
terms: tokens.clone(),
mode: filter.mode,
}];
let surviving: HashSet<u128> = select_superfiles(manifest, &prune_leaves)
.await?
.iter()
.map(|e| e.superfile_id.as_u128())
.collect();
if surviving.is_empty() {
return Ok(Vec::new());
}
let superfiles = self
.vector_pruned_superfiles_intersect(manifest, &surviving)
.await?;
if superfiles.is_empty() {
return Ok(Vec::new());
}
let allow = self
.candidate_bitmaps(&superfiles, filter.column, &tokens, filter.mode)
.await?;
if allow.is_empty() {
return Ok(Vec::new());
}
self.route_filtered_vector_hits_async(superfiles, allow, column, query, k, options)
.await
}
async fn vector_pruned_superfiles_intersect(
&self,
manifest: &ManifestSnapshot,
surviving: &HashSet<u128>,
) -> Result<Vec<Arc<SuperfileEntry>>, QueryError> {
Ok(manifest
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?
.into_iter()
.filter(|e| surviving.contains(&e.superfile_id.as_u128()))
.collect())
}
async fn candidate_bitmaps(
&self,
superfiles: &[Arc<SuperfileEntry>],
filter_col: &str,
tokens: &[String],
mode: BoolMode,
) -> Result<HashMap<SuperfileUri, Arc<RoaringBitmap>>, QueryError> {
let filter_col_arc = Arc::new(filter_col.to_owned());
let tokens_arc: Arc<Vec<String>> = Arc::new(tokens.to_vec());
let op_stats = self.op_stats.clone();
self.fanout_candidate_bitmaps(superfiles, move |r, _entry| {
let filter_col_arc = Arc::clone(&filter_col_arc);
let tokens_arc = Arc::clone(&tokens_arc);
let op_stats = op_stats.clone();
async move {
let refs: Vec<&str> = tokens_arc.iter().map(String::as_str).collect();
let (docs, work) = r
.token_match(&filter_col_arc, &refs, mode)
.await
.map_err(|e| QueryError::Parquet(e.to_string()))?;
if let Some(stats) = &op_stats {
stats.add_fts_postings_bytes(work.postings_bytes);
stats.add_planned_read_ranges(work.planned_ranges);
stats.add_kernel_cpu_ns(work.kernel_cpu_ns);
}
Ok(docs.into_iter().collect::<RoaringBitmap>())
}
})
.await
}
pub(crate) async fn vector_hits_filtered_by_plan(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
plan: &CandidatePlan,
) -> Result<Vec<SuperfileHit>, QueryError> {
if k == 0 {
return Ok(Vec::new());
}
let query = calibrated_query(self, column, query);
let query: &[f32] = &query;
let manifest = self.manifest();
let superfiles = match plan.surviving_superfile_ids(manifest).await? {
None => manifest
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?,
Some(surviving) if surviving.is_empty() => return Ok(Vec::new()),
Some(surviving) => {
self.vector_pruned_superfiles_intersect(manifest, &surviving)
.await?
}
};
if superfiles.is_empty() {
return Ok(Vec::new());
}
let allow = self.candidate_bitmaps_from_plan(&superfiles, plan).await?;
if allow.is_empty() {
return Ok(Vec::new());
}
self.route_filtered_vector_hits_async(superfiles, allow, column, query, k, options)
.await
}
async fn stable_ids_from_user_allow_async(
&self,
user_allow: &HashMap<SuperfileUri, Arc<RoaringBitmap>>,
) -> Result<Vec<i128>, QueryError> {
let mut out: HashSet<i128> = HashSet::new();
let manifest = self.manifest();
let id_column = self.options().id_column.as_str();
for (uri, bm) in user_allow {
let entry = manifest
.lookup_superfile_entry(*uri)
.await
.map_err(QueryError::ManifestLoad)?
.ok_or_else(|| {
QueryError::Execute(format!("user superfile {uri:?} missing from manifest"))
})?;
if row_id_from_manifest_entry(&entry, 0).is_some() {
for local in bm.iter() {
out.insert(entry.id_min + i128::from(local));
}
continue;
}
let locals = Arc::new(bm.iter().collect::<Vec<u32>>());
let ids =
read_ids_for_locals(manifest, &entry, &locals, id_column, false, &self.op_stats)
.await?;
out.extend(ids);
}
Ok(out.into_iter().collect())
}
async fn route_filtered_vector_hits_async(
&self,
user_superfiles: Vec<Arc<SuperfileEntry>>,
user_allow: HashMap<SuperfileUri, Arc<RoaringBitmap>>,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
) -> Result<Vec<SuperfileHit>, QueryError> {
if user_allow.is_empty() {
return Ok(Vec::new());
}
let drained = self
.vector_index_table()
.map(|hidden| {
hidden
.pinned_reader_with(self.op_stats.clone())
.manifest()
.get_drained_ranges()
})
.unwrap_or_default();
let mut drained_allow = HashMap::new();
let mut undrained_user = Vec::new();
for entry in user_superfiles {
if drained.contains(entry.birth_version) {
if let Some(bitmap) = user_allow.get(&entry.uri) {
drained_allow.insert(entry.uri, Arc::clone(bitmap));
}
} else {
undrained_user.push(entry);
}
}
let user_hits = if undrained_user.is_empty() {
Vec::new()
} else {
self.vector_fanout_over_superfiles(
undrained_user,
column,
query,
k,
options,
Some(user_allow.clone()),
)
.await?
};
let stable_ids = self
.stable_ids_from_user_allow_async(&drained_allow)
.await?;
let hidden_hits = if stable_ids.is_empty() {
Vec::new()
} else {
let prepared = self
.prepare_vector_stable_allow_async(Arc::new(stable_ids))
.await?;
if !prepared.use_hidden_index {
return Err(QueryError::Execute(
"drained filtered-vector ids resolved to a user allow-set instead of the \
hidden index"
.into(),
));
}
self.vector_hits_prepared_global_allow_async(column, query, k, options, &prepared)
.await?
};
Ok(top_k_ascending(vec![hidden_hits, user_hits], k))
}
#[cfg(feature = "test-helpers")]
pub async fn vector_hits_global_allow_async(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
allow_global: Arc<RoaringBitmap>,
) -> Result<Vec<SuperfileHit>, QueryError> {
let prepared = self.prepare_vector_global_allow_async(allow_global).await?;
self.vector_hits_prepared_global_allow_async(column, query, k, options, &prepared)
.await
}
#[cfg(feature = "test-helpers")]
pub async fn prepare_vector_global_allow_async(
&self,
allow_global: Arc<RoaringBitmap>,
) -> Result<PreparedGlobalAllow, QueryError> {
if allow_global.is_empty() {
return Ok(PreparedGlobalAllow {
use_hidden_index: self.vector_index_table().is_some(),
allow_by_uri: HashMap::new(),
});
}
if let Some(vit) = self.vector_index_table() {
let hidden_reader = vit.pinned_reader_with(self.op_stats.clone());
let hidden_manifest = Arc::clone(hidden_reader.manifest());
let drained = hidden_manifest.get_drained_ranges();
let superfiles = hidden_manifest
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?;
if !superfiles.is_empty() {
let allow_for_cell = Arc::clone(&allow_global);
let manifest_for_ids = Arc::clone(&hidden_manifest);
let routing_stats = hidden_reader.op_stats.clone();
let allow_by_uri = hidden_reader
.fanout_candidate_bitmaps(&superfiles, move |r, entry| {
let allow_for_cell = Arc::clone(&allow_for_cell);
let manifest_for_ids = Arc::clone(&manifest_for_ids);
let routing_stats = routing_stats.clone();
async move {
let stable_ids = stable_ids_by_local_for_routing(
&manifest_for_ids,
&entry,
&r,
&routing_stats,
)
.await?;
let mut local = RoaringBitmap::new();
for (local_doc_id, stable_id) in stable_ids.into_iter().enumerate() {
if let Ok(global_id) = u32::try_from(stable_id)
&& allow_for_cell.contains(global_id)
{
local.insert(local_doc_id as u32);
}
}
Ok(local)
}
})
.await?;
if allow_by_uri.is_empty() {
return Err(QueryError::Execute(
"global allow ids for drained filtered-vector rows did not map to any \
hidden superfile"
.into(),
));
}
return Ok(PreparedGlobalAllow {
use_hidden_index: true,
allow_by_uri,
});
}
if !drained.is_empty() {
return Err(QueryError::Execute(
"hidden vector manifest has drained ranges but no hidden superfiles".into(),
));
}
}
let manifest = self.manifest();
let superfiles = manifest
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?;
let mut allow_by_uri: HashMap<SuperfileUri, RoaringBitmap> = HashMap::new();
let mut allowed = allow_global.iter().peekable();
let mut base = 0u64;
for entry in &superfiles {
let end = base.saturating_add(entry.n_docs);
while allowed.peek().is_some_and(|&id| (id as u64) < base) {
allowed.next();
}
let mut local = RoaringBitmap::new();
while let Some(id) = allowed.peek().copied() {
let id = id as u64;
if id >= end {
break;
}
local.insert((id - base) as u32);
allowed.next();
}
if !local.is_empty() {
allow_by_uri.insert(entry.uri, local);
}
base = end;
}
Ok(PreparedGlobalAllow {
use_hidden_index: false,
allow_by_uri: allow_by_uri
.into_iter()
.map(|(uri, bm)| (uri, Arc::new(bm)))
.collect(),
})
}
#[cfg(feature = "test-helpers")]
pub async fn prepare_vector_stable_allow_async(
&self,
allow_stable_ids: Arc<Vec<i128>>,
) -> Result<PreparedGlobalAllow, QueryError> {
self.prepare_vector_stable_allow_inner(allow_stable_ids)
.await
}
#[cfg(not(feature = "test-helpers"))]
pub(crate) async fn prepare_vector_stable_allow_async(
&self,
allow_stable_ids: Arc<Vec<i128>>,
) -> Result<PreparedGlobalAllow, QueryError> {
self.prepare_vector_stable_allow_inner(allow_stable_ids)
.await
}
#[cfg(any(test, feature = "test-helpers"))]
pub async fn diag_hidden_stable_cell_map(
&self,
column: &str,
) -> Result<HashMap<i128, u32>, QueryError> {
let Some(vit) = self.vector_index_table() else {
return Ok(HashMap::new());
};
let hidden_reader = vit.pinned_reader_with(self.op_stats.clone());
let hidden_manifest = Arc::clone(hidden_reader.manifest());
let superfiles = hidden_manifest
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?;
let map = Arc::new(Mutex::new(HashMap::new()));
let column_owned = column.to_string();
let map_for_fanout = Arc::clone(&map);
let manifest_for_ids = Arc::clone(&hidden_manifest);
let routing_stats = hidden_reader.op_stats.clone();
let _ = hidden_reader
.fanout_candidate_bitmaps(&superfiles, move |r, entry| {
let map = Arc::clone(&map_for_fanout);
let manifest_for_ids = Arc::clone(&manifest_for_ids);
let column = column_owned.clone();
let routing_stats = routing_stats.clone();
async move {
let stable_ids = stable_ids_by_local_for_routing(
&manifest_for_ids,
&entry,
&r,
&routing_stats,
)
.await?;
if let Some(vs) = entry.vector_summary.get(&column) {
let mut idx = 0usize;
let mut guard = map.lock().expect("diag cell-map lock");
for cell in &vs.cells {
let n: u64 = cell.clusters.counts.iter().map(|&c| u64::from(c)).sum();
for _ in 0..n {
if idx >= stable_ids.len() {
break;
}
if let Some(cid) = cell.cell_id {
guard.insert(stable_ids[idx], cid);
}
idx += 1;
}
}
}
Ok(RoaringBitmap::new())
}
})
.await?;
let map = Arc::try_unwrap(map)
.map(|m| m.into_inner().expect("diag cell-map lock"))
.unwrap_or_default();
Ok(map)
}
#[cfg(any(test, feature = "test-helpers"))]
pub fn diag_hidden_probe_laws(&self) -> Option<(Vec<u32>, Vec<u32>, Vec<u32>)> {
let vit = self.vector_index_table()?;
match vit
.pinned_reader_with(self.op_stats.clone())
.manifest()
.get_partition_strategy()
{
PartitionStrategy::VectorCell { routing, .. } => Some((
routing.width_for_k.to_vec(),
routing.fine_for_k.to_vec(),
routing.rerank_for_k.to_vec(),
)),
_ => None,
}
}
async fn prepare_vector_stable_allow_inner(
&self,
allow_stable_ids: Arc<Vec<i128>>,
) -> Result<PreparedGlobalAllow, QueryError> {
if allow_stable_ids.is_empty() {
return Ok(PreparedGlobalAllow {
use_hidden_index: false,
allow_by_uri: HashMap::new(),
});
}
let allow_set: Arc<HashSet<i128>> =
Arc::new(allow_stable_ids.iter().copied().collect::<HashSet<i128>>());
if let Some(vit) = self.vector_index_table() {
let hidden_reader = vit.pinned_reader_with(self.op_stats.clone());
let hidden_manifest = Arc::clone(hidden_reader.manifest());
let drained = hidden_manifest.get_drained_ranges();
let superfiles = hidden_manifest
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?;
if !superfiles.is_empty() {
let allow_for_cell = Arc::clone(&allow_set);
let manifest_for_ids = Arc::clone(&hidden_manifest);
let routing_stats = hidden_reader.op_stats.clone();
let allow_by_uri = hidden_reader
.fanout_candidate_bitmaps(&superfiles, move |r, entry| {
let allow_for_cell = Arc::clone(&allow_for_cell);
let manifest_for_ids = Arc::clone(&manifest_for_ids);
let routing_stats = routing_stats.clone();
async move {
let stable_ids = stable_ids_by_local_for_routing(
&manifest_for_ids,
&entry,
&r,
&routing_stats,
)
.await?;
let mut local = RoaringBitmap::new();
for (local_doc_id, stable_id) in stable_ids.into_iter().enumerate() {
if allow_for_cell.contains(&stable_id) {
local.insert(local_doc_id as u32);
}
}
Ok(local)
}
})
.await?;
if allow_by_uri.is_empty() {
return Err(QueryError::Execute(
"stable ids for drained filtered-vector rows did not map to any hidden \
superfile"
.into(),
));
}
return Ok(PreparedGlobalAllow {
use_hidden_index: true,
allow_by_uri,
});
}
if !drained.is_empty() {
return Err(QueryError::Execute(
"hidden vector manifest has drained ranges but no hidden superfiles".into(),
));
}
}
let manifest = self.manifest();
let superfiles = manifest
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?;
if superfiles.is_empty() {
return Ok(PreparedGlobalAllow {
use_hidden_index: false,
allow_by_uri: HashMap::new(),
});
}
let allow_for_user = Arc::clone(&allow_set);
let manifest_for_ids = Arc::clone(manifest);
let routing_stats = self.op_stats.clone();
let allow_by_uri = self
.fanout_candidate_bitmaps(&superfiles, move |r, entry| {
let allow_for_user = Arc::clone(&allow_for_user);
let manifest_for_ids = Arc::clone(&manifest_for_ids);
let routing_stats = routing_stats.clone();
async move {
let stable_ids = stable_ids_by_local_for_routing(
&manifest_for_ids,
&entry,
&r,
&routing_stats,
)
.await?;
let mut local = RoaringBitmap::new();
for (local_doc_id, stable_id) in stable_ids.into_iter().enumerate() {
if allow_for_user.contains(&stable_id) {
local.insert(local_doc_id as u32);
}
}
Ok(local)
}
})
.await?;
Ok(PreparedGlobalAllow {
use_hidden_index: false,
allow_by_uri,
})
}
#[cfg(feature = "test-helpers")]
pub async fn vector_hits_prepared_global_allow_async(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
prepared: &PreparedGlobalAllow,
) -> Result<Vec<SuperfileHit>, QueryError> {
let query = calibrated_query(self, column, query);
self.vector_hits_prepared_global_allow_inner(column, &query, k, options, prepared)
.await
}
#[cfg(not(feature = "test-helpers"))]
pub(crate) async fn vector_hits_prepared_global_allow_async(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
prepared: &PreparedGlobalAllow,
) -> Result<Vec<SuperfileHit>, QueryError> {
self.vector_hits_prepared_global_allow_inner(column, query, k, options, prepared)
.await
}
async fn vector_hits_prepared_global_allow_inner(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
prepared: &PreparedGlobalAllow,
) -> Result<Vec<SuperfileHit>, QueryError> {
if k == 0 || prepared.allow_by_uri.is_empty() {
return Ok(Vec::new());
}
if prepared.use_hidden_index {
let vit = self.vector_index_table().ok_or_else(|| {
QueryError::Execute("prepared hidden allow-set but no hidden index table".into())
})?;
let hidden_reader = vit.pinned_reader_with(self.op_stats.clone());
let superfiles = hidden_reader
.manifest()
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?;
if superfiles.is_empty() {
return Ok(Vec::new());
}
return hidden_reader
.vector_fanout_over_superfiles(
superfiles,
column,
query,
k,
options,
Some(prepared.allow_by_uri.clone()),
)
.await;
}
let superfiles = self
.manifest()
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?;
if superfiles.is_empty() {
return Ok(Vec::new());
}
self.vector_fanout_over_superfiles(
superfiles,
column,
query,
k,
options,
Some(prepared.allow_by_uri.clone()),
)
.await
}
async fn candidate_bitmaps_from_plan(
&self,
superfiles: &[Arc<SuperfileEntry>],
plan: &CandidatePlan,
) -> Result<HashMap<SuperfileUri, Arc<RoaringBitmap>>, QueryError> {
let plan_arc = Arc::new(plan.clone());
let op_stats = self.op_stats.clone();
self.fanout_candidate_bitmaps(superfiles, move |r, _entry| {
let plan = Arc::clone(&plan_arc);
let op_stats = op_stats.clone();
async move {
let (bitmap, work) = plan
.evaluate(r.as_ref())
.await
.map_err(|e| QueryError::Parquet(e.to_string()))?;
if let Some(stats) = &op_stats {
stats.add_fts_postings_bytes(work.postings_bytes);
stats.add_planned_read_ranges(work.planned_ranges);
stats.add_kernel_cpu_ns(work.kernel_cpu_ns);
}
bitmap.ok_or_else(|| {
QueryError::Execute(
"bounded CandidatePlan evaluated to Unbounded — planner bug".into(),
)
})
}
})
.await
}
async fn fanout_candidate_bitmaps<F, Fut>(
&self,
superfiles: &[Arc<SuperfileEntry>],
doc_ids: F,
) -> Result<HashMap<SuperfileUri, Arc<RoaringBitmap>>, QueryError>
where
F: Fn(Arc<SuperfileReader>, Arc<SuperfileEntry>) -> Fut + Send + Sync + Clone + 'static,
Fut: Future<Output = Result<RoaringBitmap, QueryError>> + Send,
{
let units: Vec<(Arc<SuperfileEntry>, ())> =
superfiles.iter().map(|e| (Arc::clone(e), ())).collect();
let body = move |r: Arc<SuperfileReader>,
entry: Arc<SuperfileEntry>,
tombstone_cache: Option<Arc<SidecarCache>>,
now: Instant,
_: ()| {
let doc_ids = doc_ids.clone();
async move {
let mut bm = doc_ids(r, Arc::clone(&entry)).await?;
subtract_tombstones(&mut bm, &entry, tombstone_cache.as_deref(), now)?;
Ok((entry.uri, bm))
}
};
let pairs: Vec<(SuperfileUri, RoaringBitmap)> =
dispatch::fanout_with(self, units, true, false, body).await?;
Ok(pairs
.into_iter()
.filter(|(_, bm)| !bm.is_empty())
.map(|(uri, bm)| (uri, Arc::new(bm)))
.collect())
}
pub(crate) async fn vector_search_user_table_async(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
) -> Result<Vec<SuperfileHit>, QueryError> {
if k == 0 {
return Ok(Vec::new());
}
let manifest = self.manifest();
let superfiles = manifest
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?;
if superfiles.is_empty() {
return Ok(Vec::new());
}
self.fanout_vector_clusters(&superfiles, column, query, k, options)
.await
}
pub(crate) async fn vector_search_global_index_async(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
) -> Result<Vec<SuperfileHit>, QueryError> {
if k == 0 {
return Ok(Vec::new());
}
let Some(vit) = self.vector_index_table() else {
if let Some(reason) = self.hidden_index_open_error() {
return Err(QueryError::Execute(format!(
"hidden vector index present but failed to open: {reason}"
)));
}
return self
.vector_search_user_table_async(column, query, k, options)
.await;
};
let hidden_reader = vit.pinned_reader_with(self.op_stats.clone());
let hidden_manifest = Arc::clone(hidden_reader.manifest());
let drained = hidden_manifest.get_drained_ranges();
let hidden_entries = hidden_manifest
.get_all_superfiles_loaded()
.await
.map_err(QueryError::ManifestLoad)?;
let hidden_search = async {
if hidden_entries.is_empty() {
Ok(Vec::new())
} else {
hidden_reader
.fanout_vector_clusters(&hidden_entries, column, query, k, options)
.await
}
};
let fast_state = async {
vit.ensure_fresh_async().await;
vit.pinned_reader_with(self.op_stats.clone())
.hidden_deleted_ids()
.map_err(|error| QueryError::Execute(error.to_string()))
};
let user_parts = self.manifest().get_undrained_superfiles_loaded(&drained);
let (hidden_hits, deleted, user_entries) = join!(hidden_search, fast_state, user_parts);
let mut hidden_hits = hidden_hits?;
let deleted = deleted?;
let user_entries = user_entries.map_err(QueryError::ManifestLoad)?;
let mut user_hits = if user_entries.is_empty() {
Vec::new()
} else {
self.fanout_vector_clusters(&user_entries, column, query, k, options)
.await?
};
let refill_cap = k.saturating_add(deleted.len()).max(k);
let mut requested = k;
loop {
let mut combined = top_k_ascending(vec![hidden_hits, user_hits], requested);
if let Some(hit) = combined.iter().find(|hit| hit.stable_id.is_none()) {
return Err(QueryError::Execute(format!(
"hit {:?}/{} missing stable _id before combined delete filtering",
hit.superfile, hit.local_doc_id
)));
}
let live = combined
.iter()
.filter(|hit| {
hit.stable_id
.is_some_and(|id| deleted.binary_search(&id).is_err())
})
.count();
let deleted_occupies_top_k = live < k && !deleted.is_empty();
if !deleted_occupies_top_k || requested >= refill_cap {
combined.retain(|hit| {
hit.stable_id
.is_some_and(|id| deleted.binary_search(&id).is_err())
});
combined.truncate(k);
return Ok(combined);
}
let next = requested
.saturating_mul(DELETE_REFILL_GROWTH_FACTOR)
.min(refill_cap);
if next == requested {
combined.retain(|hit| {
hit.stable_id
.is_some_and(|id| deleted.binary_search(&id).is_err())
});
combined.truncate(k);
return Ok(combined);
}
requested = next;
let hidden_retry = async {
if hidden_entries.is_empty() {
Ok(Vec::new())
} else {
hidden_reader
.fanout_vector_clusters(&hidden_entries, column, query, requested, options)
.await
}
};
let user_retry = async {
if user_entries.is_empty() {
Ok(Vec::new())
} else {
self.fanout_vector_clusters(&user_entries, column, query, requested, options)
.await
}
};
let (next_hidden, next_user) = join!(hidden_retry, user_retry);
hidden_hits = next_hidden?;
user_hits = next_user?;
}
}
pub(crate) async fn vector_search_async(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
) -> Result<Vec<SuperfileHit>, QueryError> {
let query = calibrated_query(self, column, query);
self.vector_search_global_index_async(column, &query, k, options)
.await
}
}
pub(crate) fn calibrated_query<'q>(
reader: &SupertableReader,
column: &str,
query: &'q [f32],
) -> Cow<'q, [f32]> {
let cosine = reader
.options()
.vector_columns
.iter()
.any(|c| c.column == column && c.metric == Metric::Cosine);
calibrated_query_for(cosine, query)
}
fn calibrated_query_for(cosine: bool, query: &[f32]) -> Cow<'_, [f32]> {
if cosine {
let mut q = query.to_vec();
normalize(&mut q);
Cow::Owned(q)
} else {
Cow::Borrowed(query)
}
}
impl SupertableReader {
pub fn vector_search(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
filter: Option<VectorFilter<'_>>,
projection: Option<&[&str]>,
) -> Result<Vec<RecordBatch>, QueryError> {
let query = calibrated_query(self, column, query);
let query: &[f32] = &query;
let _fg = crate::supertable::reader_cache::disk::ForegroundQueryGuard::enter();
self.block_on(async {
let hits = match filter {
None => {
self.vector_search_global_index_async(column, query, k, options)
.await?
}
Some(f) => {
self.vector_hits_filtered_async(column, query, k, options, f)
.await?
}
};
let id_column = self.options().id_column.as_str();
if free_columns_unambiguous(&self.options().schema, id_column)
&& let Some(indices) = id_score_projection_indices(projection, id_column)
&& hits.iter().all(|hit| hit.stable_id.is_some())
{
let batch = hits_id_score_batch(self, &hits)?
.project(&indices)
.map_err(|e| QueryError::Execute(e.to_string()))?;
return Ok(vec![batch]);
}
let hits = user_placement_for_scalar_resolve(self, &hits).await?;
let batch = resolve_hits_named(self, &hits, projection, "vector_search")
.await
.map_err(|e| QueryError::Execute(e.to_string()))?;
Ok(vec![batch])
})
}
pub fn vector_hits(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
filter: Option<VectorFilter<'_>>,
) -> Result<Vec<SuperfileHit>, QueryError> {
let query = calibrated_query(self, column, query);
let query: &[f32] = &query;
let _fg = crate::supertable::reader_cache::disk::ForegroundQueryGuard::enter();
match filter {
None => self.block_on(self.vector_search_global_index_async(column, query, k, options)),
Some(f) => self.block_on(self.vector_hits_filtered_async(column, query, k, options, f)),
}
}
}
fn subtract_tombstones(
bm: &mut RoaringBitmap,
entry: &SuperfileEntry,
tombstone_cache: Option<&SidecarCache>,
now: Instant,
) -> Result<(), QueryError> {
if let Some(cache) = tombstone_cache {
let deleted = cache
.bitmap_for(entry.superfile_id, now)
.map_err(|e| QueryError::Store(format!("tombstone cache: {e}")))?;
if !deleted.is_empty() {
*bm -= &*deleted;
}
}
Ok(())
}
fn rerank_mult_from_law(
hidden_vector_index: bool,
filtered: bool,
caller_rerank_mult: Option<usize>,
hidden_routing: Option<&CellRoutingParams>,
k: usize,
) -> Option<usize> {
if !hidden_vector_index || filtered || caller_rerank_mult.is_some() {
return None;
}
hidden_routing
.and_then(|r| r.rerank_for_k_at(k))
.map(|n| n.div_ceil(k.max(1)).max(1))
}
fn apply_width_pin(
routing: &mut CellRoutingParams,
caller_nprobe: Option<usize>,
law_width: Option<usize>,
filtered: bool,
populated_cells: usize,
) -> Option<usize> {
if let Some(nprobe) = caller_nprobe {
routing.nprobe_min = nprobe;
routing.nprobe_max = nprobe;
if filtered {
return None;
}
routing.fine_nprobe = usize::MAX;
Some(nprobe.clamp(1, populated_cells))
} else if let Some(width) = law_width {
routing.nprobe_min = width;
routing.nprobe_max = width;
routing.fine_nprobe = usize::MAX;
Some(width.min(populated_cells))
} else {
None
}
}
fn deferred_shortlist_limit(
k: usize,
replica_overhead: usize,
rerank_mult: usize,
law_rerank_served: bool,
caller_nprobe: bool,
served_cells_over_width: (usize, usize),
) -> usize {
let base = k
.saturating_add(replica_overhead)
.saturating_mul(rerank_mult);
if law_rerank_served && !caller_nprobe {
let (served, stamped_width) = served_cells_over_width;
base.saturating_mul(served).div_ceil(stamped_width.max(1))
} else {
base
}
}
fn select_global_shortlist(
mut pooled: Vec<(usize, ScanCandidate)>,
limit: usize,
cell_floor: usize,
) -> Vec<(usize, ScanCandidate)> {
let cmp = |a: &(usize, ScanCandidate), b: &(usize, ScanCandidate)| {
b.1.estimate.total_cmp(&a.1.estimate).then_with(|| {
(a.0, a.1.cell_idx, a.1.pos, a.1.did).cmp(&(b.0, b.1.cell_idx, b.1.pos, b.1.did))
})
};
if pooled.len() <= limit {
return pooled;
}
let mut floor_keep: Vec<(usize, ScanCandidate)> = Vec::new();
if cell_floor > 0 {
let mut start = 0;
while start < pooled.len() {
let (si, cell) = (pooled[start].0, pooled[start].1.cell_idx);
let mut end = start + 1;
while end < pooled.len() && pooled[end].0 == si && pooled[end].1.cell_idx == cell {
end += 1;
}
let group = &mut pooled[start..end];
if group.len() > cell_floor {
group.select_nth_unstable_by(cell_floor, cmp);
floor_keep.extend_from_slice(&group[..cell_floor]);
} else {
floor_keep.extend_from_slice(group);
}
start = end;
}
}
pooled.select_nth_unstable_by(limit, cmp);
pooled.truncate(limit);
if floor_keep.is_empty() {
return pooled;
}
let kept: HashSet<(usize, usize, u32, u32)> = pooled
.iter()
.map(|(si, c)| (*si, c.cell_idx, c.pos, c.did))
.collect();
for (si, cand) in floor_keep {
if !kept.contains(&(si, cand.cell_idx, cand.pos, cand.did)) {
pooled.push((si, cand));
}
}
pooled
}
fn top_k_ascending(per_superfile: Vec<Vec<SuperfileHit>>, k: usize) -> Vec<SuperfileHit> {
fn hit_order(a: &SuperfileHit, b: &SuperfileHit) -> Ordering {
a.score
.partial_cmp(&b.score)
.unwrap_or(Ordering::Equal)
.then_with(|| a.superfile.cmp(&b.superfile))
.then_with(|| a.local_doc_id.cmp(&b.local_doc_id))
}
#[derive(PartialEq)]
struct MaxByScore(SuperfileHit);
impl Eq for MaxByScore {}
impl PartialOrd for MaxByScore {
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
impl Ord for MaxByScore {
fn cmp(&self, other: &Self) -> Ordering {
hit_order(&self.0, &other.0)
}
}
let mut best_by_id: HashMap<i128, SuperfileHit> = HashMap::new();
let mut passthrough = Vec::new();
for hit in per_superfile.into_iter().flatten() {
if let Some(id) = hit.stable_id {
best_by_id
.entry(id)
.and_modify(|existing| {
if hit_order(&hit, existing) == Ordering::Less {
*existing = hit;
}
})
.or_insert(hit);
} else {
passthrough.push(hit);
}
}
let mut heap = BinaryHeap::with_capacity(k + 1);
for hit in best_by_id.into_values().chain(passthrough) {
if heap.len() < k {
heap.push(MaxByScore(hit));
} else if let Some(worst) = heap.peek()
&& hit_order(&hit, &worst.0) == Ordering::Less
{
heap.pop();
heap.push(MaxByScore(hit));
}
}
let mut result: Vec<SuperfileHit> = heap.into_iter().map(|m| m.0).collect();
result.sort_unstable_by(hit_order);
result
}
impl Supertable {
#[cfg_attr(
feature = "detailed-tracing",
tracing::instrument(skip_all, fields(column = column, k = k, dim = query.len()))
)]
pub fn vector_search(
&self,
column: &str,
query: &[f32],
k: usize,
options: VectorSearchOptions,
filter: Option<VectorFilter<'_>>,
projection: Option<&[&str]>,
) -> Result<Vec<RecordBatch>, crate::InfinoError> {
self.reader()?
.vector_search(column, query, k, options, filter, projection)
.map_err(crate::InfinoError::from)
.map_err(|e| e.with_context("vector_search", None))
}
}
#[cfg(test)]
mod tests {
use std::{
borrow::Cow,
collections::{BTreeMap, BTreeSet, HashMap, HashSet},
sync::Arc,
};
use arrow::array::Array;
#[test]
fn calibrated_query_normalizes_cosine_only() {
let q = [3.0f32, 4.0];
let cos = calibrated_query_for(true, &q);
assert!((cos[0] - 0.6).abs() < 1e-6);
assert!((cos[1] - 0.8).abs() < 1e-6);
let non_cos = calibrated_query_for(false, &q);
assert!(matches!(non_cos, Cow::Borrowed(_)));
assert_eq!(&*non_cos, &q[..]);
}
use arrow_array::{
Decimal128Array, FixedSizeListArray, Float32Array, LargeStringArray, RecordBatch,
};
use arrow_schema::{DataType, Field, Schema};
use super::{
RABITQ_ADMIT_CELL_SHORTLIST_MIN, SCORE_COLUMN, ScanCandidate, VectorFilter,
VectorSearchOptions, admit_extension_round, admit_shortlist_window, apply_width_pin,
calibrated_query_for, cells_ranked_by_fine_score, free_column_slot,
free_columns_unambiguous, gate_fine_candidates_by_fragment, hidden_hits_user_ids,
id_score_projection_indices, is_hidden_vector_manifest, law_floor_serve_selection,
postings_by_cell_from_summaries, rerank_mult_from_law, score_fine_candidates,
select_global_shortlist, union_cell_selection, vector_read_query_error,
};
use crate::{
BoolMode, InfinoError,
superfile::{
SuperfileReader,
builder::{BuilderOptions, FtsConfig, SuperfileBuilder, VectorConfig},
error::{ReadError, VectorError},
fts::reader::Bm25Stats,
vector::{distance::Metric, rerank_codec::RerankCodec},
},
supertable::{
Supertable, SupertableOptions,
error::QueryError,
manifest::{
ClusterCentroids, ManifestSnapshot,
list::{CellRoutingParams, PartitionStrategy},
},
writer::{recalibrate_probe_laws, split_overflow_cell},
},
test_helpers::default_tokenizer as tok,
};
fn block_on<F: std::future::Future>(fut: F) -> F::Output {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("test runtime")
.block_on(fut)
}
#[test]
fn a_user_column_shadowing_a_free_name_declines_the_fast_path() {
let id = "doc_id";
let clean = Schema::new(vec![
Field::new("title", DataType::LargeUtf8, false),
Field::new("body", DataType::LargeUtf8, false),
]);
assert!(free_columns_unambiguous(&clean, id));
let shadows_score = Schema::new(vec![
Field::new("title", DataType::LargeUtf8, false),
Field::new(SCORE_COLUMN, DataType::Float32, false),
]);
assert!(
!free_columns_unambiguous(&shadows_score, id),
"a visible `score` column must decline the fast path"
);
let shadows_id = Schema::new(vec![Field::new(id, DataType::Int64, false)]);
assert!(
!free_columns_unambiguous(&shadows_id, id),
"a user column matching the id column must decline the fast path"
);
assert_eq!(free_column_slot(SCORE_COLUMN, SCORE_COLUMN), None);
assert!(!free_columns_unambiguous(&clean, SCORE_COLUMN));
}
#[test]
fn tvf_index_classification_matches_the_name_path() {
let id = "doc_id";
let names = [id, "title", "body", SCORE_COLUMN];
let by_index = |requested: &[usize]| -> Option<Vec<usize>> {
requested
.iter()
.map(|&i| names.get(i).and_then(|n| free_column_slot(n, id)))
.collect()
};
for requested in [
vec![0, 3], vec![3, 0], vec![0], vec![3], vec![0, 1, 3], vec![1], ] {
let as_names: Vec<&str> = requested.iter().map(|&i| names[i]).collect();
assert_eq!(
by_index(&requested),
id_score_projection_indices(Some(&as_names), id),
"index and name classification disagree for {requested:?}"
);
}
}
#[test]
fn id_score_projection_admits_any_subset_of_id_and_score() {
let id = "doc_id";
assert_eq!(id_score_projection_indices(None, id), Some(vec![0, 1]));
assert_eq!(
id_score_projection_indices(Some(&[id, SCORE_COLUMN]), id),
Some(vec![0, 1])
);
assert_eq!(
id_score_projection_indices(Some(&[SCORE_COLUMN, id]), id),
Some(vec![1, 0])
);
assert_eq!(id_score_projection_indices(Some(&[id]), id), Some(vec![0]));
assert_eq!(
id_score_projection_indices(Some(&[SCORE_COLUMN]), id),
Some(vec![1])
);
assert_eq!(
id_score_projection_indices(Some(&[id, SCORE_COLUMN, id]), id),
Some(vec![0, 1, 0])
);
assert_eq!(
id_score_projection_indices(Some(&["other", SCORE_COLUMN]), id),
None
);
assert_eq!(id_score_projection_indices(Some(&[]), id), None);
}
#[test]
fn admit_shortlist_window_scales_with_cell_population() {
assert_eq!(admit_shortlist_window(0), RABITQ_ADMIT_CELL_SHORTLIST_MIN);
assert_eq!(admit_shortlist_window(64), RABITQ_ADMIT_CELL_SHORTLIST_MIN);
assert_eq!(admit_shortlist_window(240), RABITQ_ADMIT_CELL_SHORTLIST_MIN);
assert_eq!(admit_shortlist_window(256), 52);
assert_eq!(admit_shortlist_window(512), 103);
assert_eq!(admit_shortlist_window(1024), 205);
assert_eq!(admit_shortlist_window(241), 49);
}
#[test]
fn admit_extension_round_follows_evidence() {
let ranking = vec![
(1u32, 0.10f32),
(2, 0.12),
(3, 0.14),
(4, 0.20),
(5, 0.24),
(6, 0.80),
];
let admitted: HashSet<u32> = [1, 2, 3].into_iter().collect();
let exacts: HashMap<u32, f32> = [(1u32, 0.15f32), (2, 0.17), (3, 0.19)].into();
assert_eq!(
admit_extension_round(&ranking, &admitted, &exacts, 0.30),
vec![4, 5]
);
assert!(admit_extension_round(&ranking, &admitted, &exacts, 0.16).is_empty());
assert!(admit_extension_round(&ranking, &admitted, &HashMap::new(), 0.30).is_empty());
}
#[test]
fn law_floor_serve_selection_extends_past_the_served_floor() {
let cliff = vec![(7u32, 0.10f32), (8, 0.90), (9, 0.95)];
let (cells, ext) = law_floor_serve_selection(&cliff, &[7, 8], 2, 0.20);
assert_eq!(cells, vec![7, 8]);
assert!(ext.is_empty());
let flat = vec![(7u32, 0.10f32), (8, 0.11), (9, 0.12)];
let (cells, ext) = law_floor_serve_selection(&flat, &[7, 8], 2, 0.20);
assert_eq!(cells, vec![7, 8, 9]);
assert_eq!(ext, [9u32].into_iter().collect());
let (cells, ext) = law_floor_serve_selection(&flat, &[9], 2, 0.20);
assert_eq!(cells, vec![9, 7, 8]);
assert!(ext.is_empty());
}
#[test]
fn union_cell_selection_dedups_with_grid_priority() {
assert_eq!(union_cell_selection(&[4], &[9]), vec![4, 9]);
assert_eq!(union_cell_selection(&[4], &[4]), vec![4]);
assert_eq!(union_cell_selection(&[4, 9], &[9, 1]), vec![4, 9, 1]);
assert_eq!(union_cell_selection(&[], &[2]), vec![2]);
}
#[test]
fn cells_ranked_by_fine_score_takes_min_per_cell_in_order() {
let candidates: Vec<(usize, u32, f32, Option<u32>, u64)> = vec![
(0, 0, 0.9, Some(7), 10),
(0, 1, 0.2, Some(7), 10), (1, 2, 0.5, Some(3), 10), (1, 3, 0.5, Some(2), 10), (0, 4, 0.1, None, 10), ];
let ranked = cells_ranked_by_fine_score(&candidates);
assert_eq!(ranked.len(), 3);
assert_eq!(ranked[0], (7, 0.2));
assert_eq!(ranked[1].0, 2, "score tie broken by lower cell id");
assert_eq!(ranked[2].0, 3);
}
#[test]
fn hidden_hits_user_ids_uses_inline_stable_id_fast_path() {
let dim = 16;
let table = Supertable::create(options_one_superfile_per_commit(dim)).expect("create");
let reader = table.reader().expect("reader");
let manifest = reader.manifest();
let mk = |sid: i128| SuperfileHit {
superfile: SuperfileUri(uuid::Uuid::new_v4()),
local_doc_id: 0,
score: 0.0,
stable_id: Some(sid),
};
let hits = [mk(42)];
let ids =
block_on(hidden_hits_user_ids(manifest, &hits, "_id", &None)).expect("resolve one id");
assert_eq!(ids, vec![42], "single inline stable id returned verbatim");
let hits = [mk(42), mk(7)];
let ids =
block_on(hidden_hits_user_ids(manifest, &hits, "_id", &None)).expect("resolve two ids");
assert_eq!(ids, vec![42, 7], "inline stable ids returned in hit order");
}
#[test]
fn hidden_classification_uses_manifest_strategy_when_options_are_unstamped() {
let dim = 16;
let table = Supertable::create(options_one_superfile_per_commit(dim)).expect("create");
assert!(table.options().partition_strategy.is_none());
let mut centroid = vec![0.0; dim];
centroid[0] = 1.0;
let manifest = table
.reader()
.expect("reader")
.manifest()
.with_partition_strategy(PartitionStrategy::VectorCell {
column: "emb".into(),
clusters: ClusterCentroids::from_fp32(1, dim as u32, ¢roid, vec![1]),
routing: Default::default(),
});
assert!(is_hidden_vector_manifest(&manifest));
}
#[test]
fn per_fragment_keep_probes_small_fragment_in_shared_cell() {
let candidates = vec![
(0usize, 10u32, 0.10f32, Some(0u32), 5u64),
(0, 11, 0.11, Some(0), 5),
(0, 12, 0.12, Some(0), 5),
(1, 20, 0.30, Some(0), 5),
];
let selected: HashSet<u32> = [0].into_iter().collect();
let selected_ordered = [0u32];
let candidate_counts: HashMap<(usize, u32), u64> = candidates
.iter()
.map(|(si, cluster, _, _, count)| ((*si, *cluster), *count))
.collect();
let mut scored = Vec::new();
let gated = gate_fine_candidates_by_fragment(
candidates,
&selected,
&selected_ordered,
2, 0.0, 1, &candidate_counts,
&mut scored,
None,
None,
);
assert!(
gated.iter().any(|(si, _, _)| *si == 1),
"small fragment starved from the probe set: {gated:?}"
);
assert_eq!(gated.iter().filter(|(si, _, _)| *si == 0).count(), 2);
}
#[test]
fn extension_cells_keep_bounded_depth_under_pin() {
let candidates = vec![
(0usize, 10u32, 0.10f32, Some(0u32), 5u64),
(0, 11, 0.11, Some(0), 5),
(0, 12, 0.12, Some(0), 5),
(0, 20, 0.20, Some(1), 5),
(0, 21, 0.21, Some(1), 5),
(0, 22, 0.22, Some(1), 5),
];
let selected: HashSet<u32> = [0, 1].into_iter().collect();
let selected_ordered = [0u32, 1];
let extension: HashSet<u32> = [1].into_iter().collect();
let candidate_counts: HashMap<(usize, u32), u64> = candidates
.iter()
.map(|(si, cluster, _, _, count)| ((*si, *cluster), *count))
.collect();
let mut scored = Vec::new();
let gated = gate_fine_candidates_by_fragment(
candidates,
&selected,
&selected_ordered,
usize::MAX, 0.0, 1, &candidate_counts,
&mut scored,
None,
Some((&extension, 1)), );
let pinned: HashSet<u32> = gated
.iter()
.filter(|(_, c, _)| *c < 20)
.map(|(_, c, _)| *c)
.collect();
assert_eq!(pinned, [10u32, 11, 12].into_iter().collect());
let extended: Vec<u32> = gated
.iter()
.filter(|(_, c, _)| *c >= 20)
.map(|(_, c, _)| *c)
.collect();
assert_eq!(
extended,
vec![20],
"extension cell read past its bounded depth: {gated:?}"
);
}
#[test]
fn per_generation_keep_bounds_across_probed_cells() {
let candidates = vec![
(0usize, 10u32, 0.10f32, Some(0u32), 5u64),
(0, 11, 0.11, Some(0), 5),
(0, 12, 0.12, Some(0), 5),
(0, 20, 0.13, Some(1), 5),
(0, 21, 0.14, Some(1), 5),
(0, 22, 0.15, Some(1), 5),
];
let selected: HashSet<u32> = [0, 1].into_iter().collect();
let selected_ordered = [0u32, 1];
let birth_versions = [100u64];
let candidate_counts: HashMap<(usize, u32), u64> = candidates
.iter()
.map(|(si, cluster, _, _, count)| ((*si, *cluster), *count))
.collect();
let mut scored = Vec::new();
let gated = gate_fine_candidates_by_fragment(
candidates,
&selected,
&selected_ordered,
2, 0.0, 1, &candidate_counts,
&mut scored,
Some(&birth_versions),
None,
);
assert_eq!(
gated.len(),
2,
"per-wave keep multiplied by cells: {gated:?}"
);
let kept: HashSet<u32> = gated.iter().map(|(_, c, _)| *c).collect();
assert_eq!(kept, [10u32, 11].into_iter().collect());
}
#[test]
fn per_generation_keep_scales_with_fraction() {
let candidates = vec![
(0usize, 10u32, 0.10f32, Some(0u32), 5u64),
(0, 11, 0.11, Some(0), 5),
(0, 12, 0.12, Some(0), 5),
(0, 20, 0.13, Some(1), 5),
(0, 21, 0.14, Some(1), 5),
(0, 22, 0.15, Some(1), 5),
];
let selected: HashSet<u32> = [0, 1].into_iter().collect();
let selected_ordered = [0u32, 1];
let birth_versions = [100u64];
let candidate_counts: HashMap<(usize, u32), u64> = candidates
.iter()
.map(|(si, cluster, _, _, count)| ((*si, *cluster), *count))
.collect();
let mut scored = Vec::new();
let gated = gate_fine_candidates_by_fragment(
candidates,
&selected,
&selected_ordered,
1, 0.5, 1, &candidate_counts,
&mut scored,
Some(&birth_versions),
None,
);
assert_eq!(gated.len(), 3, "fraction keep miscounted: {gated:?}");
let kept: HashSet<u32> = gated.iter().map(|(_, c, _)| *c).collect();
assert_eq!(kept, [10u32, 11, 12].into_iter().collect());
}
#[test]
fn per_generation_keep_probes_small_delta_wave() {
let candidates = vec![
(0usize, 10u32, 0.10f32, Some(0u32), 5u64),
(0, 11, 0.11, Some(0), 5),
(0, 12, 0.12, Some(0), 5),
(1, 20, 0.30, Some(0), 5),
];
let selected: HashSet<u32> = [0].into_iter().collect();
let selected_ordered = [0u32];
let birth_versions = [100u64, 200];
let candidate_counts: HashMap<(usize, u32), u64> = candidates
.iter()
.map(|(si, cluster, _, _, count)| ((*si, *cluster), *count))
.collect();
let mut scored = Vec::new();
let gated = gate_fine_candidates_by_fragment(
candidates,
&selected,
&selected_ordered,
2, 0.0, 1, &candidate_counts,
&mut scored,
Some(&birth_versions),
None,
);
assert!(
gated.iter().any(|(si, _, _)| *si == 1),
"small delta wave starved from the probe set: {gated:?}"
);
assert_eq!(gated.iter().filter(|(si, _, _)| *si == 0).count(), 2);
}
#[test]
fn over_budget_vector_error_surfaces_as_infino_over_budget() {
let read_err = ReadError::Vector(Box::new(VectorError::OverBudget("gate".into())));
let q = vector_read_query_error(read_err);
assert!(matches!(q, QueryError::OverBudget(_)), "got {q:?}");
assert!(matches!(
InfinoError::from(QueryError::OverBudget("x".into())),
InfinoError::OverBudget(_)
));
assert!(matches!(
vector_read_query_error(ReadError::MissingKv("k")),
QueryError::Parquet(_)
));
}
fn fixed_list_f32(dim: usize) -> DataType {
DataType::FixedSizeList(
Arc::new(Field::new("item", DataType::Float32, true)),
dim as i32,
)
}
fn schema_with_vector(dim: usize) -> Arc<Schema> {
Arc::new(Schema::new(vec![
Field::new("title", DataType::LargeUtf8, false),
Field::new("emb", fixed_list_f32(dim), false),
]))
}
fn options_one_superfile_per_commit(dim: usize) -> SupertableOptions {
let pool = Arc::new(
rayon::ThreadPoolBuilder::new()
.num_threads(1)
.build()
.expect("pool"),
);
SupertableOptions::new(
schema_with_vector(dim),
vec![FtsConfig {
column: "title".into(),
positions: false,
}],
vec![VectorConfig {
column: "emb".into(),
dim,
rot_seed: 7,
metric: Metric::Cosine,
rerank_codec: RerankCodec::Fp32,
provided_centroids: None,
}],
Some(tok()),
)
.expect("valid options")
.with_writer_pool(pool)
}
fn build_vector_batch(start: u64, n: usize, dim: usize, schema: Arc<Schema>) -> RecordBatch {
let titles = LargeStringArray::from((0..n).map(|i| format!("doc {i}")).collect::<Vec<_>>());
let mut flat = Vec::<f32>::with_capacity(n * dim);
for i in 0..n {
let global = (start as usize) + i;
for d in 0..dim {
flat.push(if d == global % dim { 1.0 } else { 0.0 });
}
}
let item_field = Arc::new(Field::new("item", DataType::Float32, true));
let values = Float32Array::from(flat);
let fsl = FixedSizeListArray::try_new(
item_field,
dim as i32,
Arc::new(values) as Arc<dyn Array>,
None,
)
.expect("FSL");
RecordBatch::try_new(schema, vec![Arc::new(titles), Arc::new(fsl)]).expect("batch")
}
fn build_oracle_superfile(n_total: usize, dim: usize) -> Arc<SuperfileReader> {
let scalar_schema = Arc::new(Schema::new(vec![
Field::new(
"_id",
DataType::Decimal128(
crate::supertable::options::DECIMAL128_PRECISION,
crate::supertable::options::DECIMAL128_SCALE,
),
false,
),
Field::new("title", DataType::LargeUtf8, false),
]));
let opts = BuilderOptions::new(
scalar_schema.clone(),
"_id",
vec![FtsConfig {
column: "title".into(),
positions: false,
}],
vec![VectorConfig {
column: "emb".into(),
dim,
rot_seed: 7,
metric: Metric::Cosine,
rerank_codec: RerankCodec::Fp32,
provided_centroids: None,
}],
Some(tok()),
);
let mut b = SuperfileBuilder::new(opts).expect("builder");
let ids = arrow_array::Decimal128Array::from((0..n_total as i128).collect::<Vec<_>>())
.with_precision_and_scale(
crate::supertable::options::DECIMAL128_PRECISION,
crate::supertable::options::DECIMAL128_SCALE,
)
.expect("decimal128");
let titles =
LargeStringArray::from((0..n_total).map(|i| format!("doc {i}")).collect::<Vec<_>>());
let scalar_batch =
RecordBatch::try_new(scalar_schema, vec![Arc::new(ids), Arc::new(titles)])
.expect("scalar batch");
let mut flat = Vec::<f32>::with_capacity(n_total * dim);
for i in 0..n_total {
for d in 0..dim {
flat.push(if d == i % dim { 1.0 } else { 0.0 });
}
}
b.add_batch(&scalar_batch, &[flat.as_slice()])
.expect("add_batch");
let bytes = bytes::Bytes::from(b.finish().expect("finish"));
Arc::new(SuperfileReader::open(bytes).expect("open"))
}
#[test]
fn vector_search_empty_supertable_returns_empty() {
let st = Supertable::create(options_one_superfile_per_commit(16)).expect("create");
let r = st.reader().expect("reader");
let q = vec![0.1f32; 16];
let hits = r
.vector_hits("emb", &q, 5, VectorSearchOptions::new(), None)
.expect("query");
assert!(hits.is_empty());
}
fn superfile_bytes_with_ids(ids: &[i128], dim: usize) -> bytes::Bytes {
let scalar_schema = Arc::new(Schema::new(vec![
Field::new(
"_id",
DataType::Decimal128(
crate::supertable::options::DECIMAL128_PRECISION,
crate::supertable::options::DECIMAL128_SCALE,
),
false,
),
Field::new("title", DataType::LargeUtf8, false),
]));
let opts = BuilderOptions::new(
scalar_schema.clone(),
"_id",
vec![FtsConfig {
column: "title".into(),
positions: false,
}],
vec![VectorConfig {
column: "emb".into(),
dim,
rot_seed: 7,
metric: Metric::Cosine,
rerank_codec: RerankCodec::Fp32,
provided_centroids: None,
}],
Some(tok()),
);
let mut b = SuperfileBuilder::new(opts).expect("builder");
let id_arr = Decimal128Array::from(ids.to_vec())
.with_precision_and_scale(
crate::supertable::options::DECIMAL128_PRECISION,
crate::supertable::options::DECIMAL128_SCALE,
)
.expect("decimal128");
let titles = LargeStringArray::from(
(0..ids.len())
.map(|i| format!("doc {i}"))
.collect::<Vec<_>>(),
);
let scalar_batch =
RecordBatch::try_new(scalar_schema, vec![Arc::new(id_arr), Arc::new(titles)])
.expect("scalar batch");
let mut flat = Vec::<f32>::with_capacity(ids.len() * dim);
for i in 0..ids.len() {
for d in 0..dim {
flat.push(if d == i % dim { 1.0 } else { 0.0 });
}
}
b.add_batch(&scalar_batch, &[flat.as_slice()])
.expect("add_batch");
bytes::Bytes::from(b.finish().expect("finish"))
}
fn insert_gapped_entry(
ids: &[i128],
dim: usize,
store: &Arc<dyn crate::supertable::reader_cache::SuperfileReaderCache>,
seed: u128,
) -> Arc<SuperfileEntry> {
let id = Uuid::from_u128(seed);
let uri = SuperfileUri(id);
store
.insert(uri, superfile_bytes_with_ids(ids, dim))
.expect("insert superfile bytes");
Arc::new(SuperfileEntry {
birth_version: 0,
superfile_id: id,
uri,
n_docs: ids.len() as u64,
id_min: *ids.iter().min().expect("nonempty ids"),
id_max: *ids.iter().max().expect("nonempty ids"),
scalar_stats: std::collections::HashMap::new(),
fts_summary: std::collections::HashMap::new(),
vector_summary: std::collections::HashMap::new(),
partition_key: Vec::new(),
partition_hint: None,
vector_layout: crate::superfile::vector::layout::VectorLayout::Ivf,
subsection_offsets: None,
})
}
fn contiguous_entry(id_min: i128, n_docs: u64, seed: u128) -> Arc<SuperfileEntry> {
let id = Uuid::from_u128(seed);
Arc::new(SuperfileEntry {
birth_version: 0,
superfile_id: id,
uri: SuperfileUri(id),
n_docs,
id_min,
id_max: id_min + n_docs as i128 - 1,
scalar_stats: std::collections::HashMap::new(),
fts_summary: std::collections::HashMap::new(),
vector_summary: std::collections::HashMap::new(),
partition_key: Vec::new(),
partition_hint: None,
vector_layout: crate::superfile::vector::layout::VectorLayout::Ivf,
subsection_offsets: None,
})
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn gapped_id_placement_cache_places_reuses_and_prunes_monotonically() {
let dim = 16;
let opts = options_one_superfile_per_commit(dim);
let store = Arc::clone(&opts.store);
let opts = Arc::new(opts);
let gapped = insert_gapped_entry(&[0, 5, 10], dim, &store, 0xA1);
let contig = contiguous_entry(100, 3, 0xC0);
let m = ManifestSnapshot::new(
1,
Arc::clone(&opts),
vec![Arc::clone(&gapped), Arc::clone(&contig)],
None,
None,
);
let got = super::lookup_user_placements_by_id(&m, &[0, 5, 10, 100], &None)
.await
.expect("placements");
assert_eq!((got[0].0.uri, got[0].1), (gapped.uri, 0));
assert_eq!((got[1].0.uri, got[1].1), (gapped.uri, 1));
assert_eq!((got[2].0.uri, got[2].1), (gapped.uri, 2));
assert_eq!((got[3].0.uri, got[3].1), (contig.uri, 0));
let first_index = {
let cache = opts.gapped_id_placement_cache.lock().await;
assert_eq!(
cache.entries.len(),
1,
"only the gapped superfile is cached"
);
assert!(cache.entries.contains_key(&gapped.uri));
assert!(
!cache.entries.contains_key(&contig.uri),
"contiguous entries resolve by arithmetic and are never cached"
);
Arc::clone(
cache
.entries
.get(&gapped.uri)
.expect("cell")
.get()
.expect("index built"),
)
};
super::lookup_user_placements_by_id(&m, &[5], &None)
.await
.expect("2nd query");
{
let cache = opts.gapped_id_placement_cache.lock().await;
let second_index = Arc::clone(
cache
.entries
.get(&gapped.uri)
.expect("cell")
.get()
.expect("index built"),
);
assert!(
Arc::ptr_eq(&first_index, &second_index),
"cached index reused, not rebuilt"
);
}
let gapped2 = insert_gapped_entry(&[200, 205, 210], dim, &store, 0xB2);
let m_next =
ManifestSnapshot::new(2, Arc::clone(&opts), vec![Arc::clone(&gapped2)], None, None);
super::lookup_user_placements_by_id(&m_next, &[200, 210], &None)
.await
.expect("next-gen query");
{
let cache = opts.gapped_id_placement_cache.lock().await;
assert!(
!cache.entries.contains_key(&gapped.uri),
"superseded gapped superfile is evicted by the newer generation"
);
assert!(
cache.entries.contains_key(&gapped2.uri),
"the live gapped superfile is cached"
);
assert_eq!(cache.entries.len(), 1);
assert_eq!(cache.pruned_through, 2);
}
super::lookup_user_placements_by_id(&m, &[0], &None)
.await
.expect("older-gen query");
{
let cache = opts.gapped_id_placement_cache.lock().await;
assert!(
cache.entries.contains_key(&gapped2.uri),
"older generation did not evict the newer generation's map"
);
assert!(
cache.entries.contains_key(&gapped.uri),
"older generation rebuilt its own map"
);
assert_eq!(cache.entries.len(), 2);
assert_eq!(
cache.pruned_through, 2,
"prune watermark stays at the newest generation seen"
);
}
}
#[test]
fn vector_search_k_zero_short_circuits() {
let st = Supertable::create(options_one_superfile_per_commit(16)).expect("create");
let mut w = st.writer().expect("writer");
let schema = st.options().schema.clone();
w.append(&build_vector_batch(0, 8, 16, schema)).expect("a");
w.commit().expect("c");
let r = st.reader().expect("reader");
let q = vec![0.1f32; 16];
let hits = r
.vector_hits("emb", &q, 0, VectorSearchOptions::new(), None)
.expect("query");
assert!(hits.is_empty());
}
#[test]
fn vector_search_returns_ascending_distance_order() {
let dim = 16;
let st = Supertable::create(options_one_superfile_per_commit(dim)).expect("create");
let mut w = st.writer().expect("writer");
let schema = st.options().schema.clone();
w.append(&build_vector_batch(0, 8, dim, schema)).expect("a");
w.commit().expect("c");
let r = st.reader().expect("reader");
let mut q = vec![0.0f32; dim];
for (d, x) in q.iter_mut().enumerate() {
*x = (d as f32) / 100.0 + 0.001;
}
let hits = r
.vector_hits("emb", &q, 5, VectorSearchOptions::new(), None)
.expect("query");
assert!(!hits.is_empty());
for w in hits.windows(2) {
assert!(
w[0].score <= w[1].score,
"expected ascending: {:?} then {:?}",
w[0],
w[1]
);
}
}
#[test]
fn vector_search_top_k_caps_at_k() {
let dim = 16;
let st = Supertable::create(options_one_superfile_per_commit(dim)).expect("create");
let mut w = st.writer().expect("writer");
let schema = st.options().schema.clone();
for chunk in 0..3u64 {
w.append(&build_vector_batch(chunk * 8, 8, dim, schema.clone()))
.expect("a");
w.commit().expect("c");
}
let r = st.reader().expect("reader");
let q = vec![0.1f32; dim];
let hits = r
.vector_hits("emb", &q, 7, VectorSearchOptions::new(), None)
.expect("query");
assert_eq!(hits.len(), 7);
}
#[test]
fn vector_search_global_selection_recovers_neighbors_under_low_budget() {
let dim = 16;
let st = Supertable::create(options_one_superfile_per_commit(dim)).expect("create");
let mut w = st.writer().expect("writer");
let schema = st.options().schema.clone();
let n_seg = 10u64;
for chunk in 0..n_seg {
w.append(&build_vector_batch(chunk * 16, 16, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
}
assert_eq!(st.reader().expect("reader").n_superfiles(), n_seg as usize);
let mut q = vec![0f32; dim];
q[0] = 1.0;
let opts = VectorSearchOptions::new().with_nprobe(1);
let hits = st
.reader()
.expect("reader")
.vector_hits("emb", &q, 10, opts, None)
.expect("query");
let exact_neighbors = hits.iter().filter(|h| h.score < 1e-3).count();
assert!(
exact_neighbors >= 9,
"recall@10 ≥ 0.90 under aggressive global cluster pruning; \
recovered {exact_neighbors}/10 exact neighbors"
);
}
#[test]
fn vector_search_carries_superfile_uris_for_multi_superfile_results() {
let dim = 16;
let st = Supertable::create(options_one_superfile_per_commit(dim)).expect("create");
let mut w = st.writer().expect("writer");
let schema = st.options().schema.clone();
for chunk in 0..3u64 {
w.append(&build_vector_batch(chunk * 8, 8, dim, schema.clone()))
.expect("a");
w.commit().expect("c");
}
let r = st.reader().expect("reader");
let q = vec![0.1f32; dim];
let hits = r
.vector_hits("emb", &q, 24, VectorSearchOptions::new(), None)
.expect("query");
let superfile_uris: HashSet<_> = hits.iter().map(|h| h.superfile).collect();
assert_eq!(superfile_uris.len(), 3);
}
#[test]
fn vector_search_oracle_top_k_set_matches_single_superfile() {
let dim = 16;
let st = Supertable::create(options_one_superfile_per_commit(dim)).expect("create");
let mut w = st.writer().expect("writer");
let schema = st.options().schema.clone();
for chunk in 0..3u64 {
w.append(&build_vector_batch(chunk * 8, 8, dim, schema.clone()))
.expect("a");
w.commit().expect("c");
}
let oracle = build_oracle_superfile(24, dim);
let opts = VectorSearchOptions::new().with_nprobe(4);
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let oracle_hits =
block_on(oracle.vector_hits_async("emb", &q, 2, opts)).expect("oracle query");
let oracle_globals: HashSet<u32> = oracle_hits.iter().map(|(d, _)| *d).collect();
assert_eq!(oracle_globals, [0u32, 16].iter().copied().collect());
let st_reader = st.reader().expect("reader");
let st_hits = st_reader
.vector_hits("emb", &q, 2, opts, None)
.expect("supertable query");
let manifest = st_reader.manifest();
let st_globals: HashSet<u32> = st_hits
.iter()
.map(|h| {
let seg_idx = manifest
.superfiles
.iter()
.position(|e| e.uri == h.superfile)
.expect("superfile in manifest");
(seg_idx as u32) * 8 + h.local_doc_id
})
.collect();
assert_eq!(st_hits.len(), oracle_hits.len());
assert_eq!(st_globals, oracle_globals);
}
#[test]
fn vector_search_unknown_column_errors() {
let dim = 16;
let st = Supertable::create(options_one_superfile_per_commit(dim)).expect("create");
let mut w = st.writer().expect("writer");
let schema = st.options().schema.clone();
w.append(&build_vector_batch(0, 8, dim, schema)).expect("a");
w.commit().expect("c");
let r = st.reader().expect("reader");
let q = vec![0.1f32; dim];
let err = r
.vector_hits("nope", &q, 5, VectorSearchOptions::new(), None)
.expect_err("expected error");
assert!(
matches!(&err, QueryError::Execute(m) if m.contains("unknown vector column")),
"got {err:?}"
);
}
use tempfile::TempDir;
use uuid::Uuid;
use crate::{
storage::{LocalFsStorageProvider, StorageProvider},
supertable::{
manifest::{SuperfileEntry, SuperfileUri},
query::SuperfileHit,
tombstones::{SidecarCache, TombstoneSeqView, cache::DEFAULT_SEAL_TTL},
wal::{WalStore, tombstones_codec::TombstonesSidecar},
},
};
fn synthetic_entry(superfile_id: Uuid) -> SuperfileEntry {
SuperfileEntry {
birth_version: 0,
superfile_id,
uri: SuperfileUri(superfile_id),
n_docs: 100,
id_min: 0,
id_max: 99,
scalar_stats: std::collections::HashMap::new(),
fts_summary: std::collections::HashMap::new(),
vector_summary: std::collections::HashMap::new(),
partition_key: Vec::new(),
partition_hint: None,
vector_layout: crate::superfile::vector::layout::VectorLayout::Ivf,
subsection_offsets: None,
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn apply_tombstone_filter_drops_set_bits() {
let dir = TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("provider"));
let ws = WalStore::new(Arc::clone(&storage));
let sf_id = Uuid::from_u128(0xFEEDFACE);
let cache = Arc::new(SidecarCache::new(
ws.clone(),
DEFAULT_SEAL_TTL,
Arc::new(TombstoneSeqView {
manifest_id: 1,
seqs: [(sf_id, 1u64)].into_iter().collect(),
}),
));
let mut bitmap = roaring::RoaringBitmap::new();
bitmap.insert(1);
bitmap.insert(3);
bitmap.insert(5);
ws.put_tombstones(sf_id, None, &TombstonesSidecar { seal: None, bitmap })
.await
.expect("put sidecar");
let entry = synthetic_entry(sf_id);
let mut hits: Vec<SuperfileHit> = (0..8u32)
.map(|d| SuperfileHit {
superfile: entry.uri,
local_doc_id: d,
score: d as f32,
stable_id: None,
})
.collect();
crate::supertable::query::dispatch::apply_tombstone_filter(
Some(&cache),
&entry,
&mut hits,
std::time::Instant::now(),
)
.expect("filter");
let remaining: Vec<u32> = hits.iter().map(|h| h.local_doc_id).collect();
assert_eq!(remaining, vec![0u32, 2, 4, 6, 7]);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn apply_tombstone_filter_is_no_op_without_cache() {
let entry = synthetic_entry(Uuid::from_u128(0xABCD));
let mut hits: Vec<SuperfileHit> = (0..4u32)
.map(|d| SuperfileHit {
superfile: entry.uri,
local_doc_id: d,
score: 0.0,
stable_id: None,
})
.collect();
let original = hits.clone();
crate::supertable::query::dispatch::apply_tombstone_filter(
None,
&entry,
&mut hits,
std::time::Instant::now(),
)
.expect("no-cache");
assert_eq!(hits, original);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn apply_tombstone_filter_short_circuits_on_empty_bitmap() {
let dir = TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("provider"));
let ws = WalStore::new(Arc::clone(&storage));
let cache = Arc::new(SidecarCache::new(
ws,
DEFAULT_SEAL_TTL,
Arc::new(TombstoneSeqView::default()),
));
let entry = synthetic_entry(Uuid::from_u128(0x1111));
let mut hits: Vec<SuperfileHit> = (0..4u32)
.map(|d| SuperfileHit {
superfile: entry.uri,
local_doc_id: d,
score: 0.0,
stable_id: None,
})
.collect();
let original = hits.clone();
crate::supertable::query::dispatch::apply_tombstone_filter(
Some(&cache),
&entry,
&mut hits,
std::time::Instant::now(),
)
.expect("filter");
assert_eq!(hits, original);
}
#[test]
fn postings_by_cell_from_summaries_skips_superseded_cells() {
use crate::supertable::manifest::{CellVectorSummary, ClusterCentroids, VectorSummary};
const DIM: u32 = 4;
let column = "emb";
let cell = |cell_id: u32, count: u32| CellVectorSummary {
cell_id: Some(cell_id),
clusters: ClusterCentroids::from_fp32(
1,
DIM,
&vec![cell_id as f32; DIM as usize],
vec![count],
),
};
let sf_id = Uuid::from_u128(0xC0FFEE);
let mut entry = synthetic_entry(sf_id);
entry.vector_summary.insert(
column.into(),
VectorSummary {
centroid: vec![0.0; DIM as usize],
cells: vec![cell(1, 10), cell(2, 20), cell(3, 30)],
},
);
let entries = vec![Arc::new(entry)];
let empty = BTreeMap::new();
let (postings, any_tagged) =
postings_by_cell_from_summaries(&entries, column, None, &empty);
assert!(any_tagged);
assert_eq!(postings.get(&1), Some(&10));
assert_eq!(postings.get(&2), Some(&20));
assert_eq!(postings.get(&3), Some(&30));
let mut superseded = BTreeMap::new();
superseded.insert(sf_id, BTreeSet::from([2u32]));
let (postings, any_tagged) =
postings_by_cell_from_summaries(&entries, column, None, &superseded);
assert!(any_tagged, "surviving cells still tag");
assert!(!postings.contains_key(&2), "superseded cell is skipped");
assert_eq!(postings.get(&1), Some(&10));
assert_eq!(postings.get(&3), Some(&30));
let mut other = BTreeMap::new();
other.insert(Uuid::from_u128(0xDEAD), BTreeSet::from([1u32]));
let (postings, _) = postings_by_cell_from_summaries(&entries, column, None, &other);
assert_eq!(postings.get(&1), Some(&10));
assert_eq!(postings.get(&2), Some(&20));
assert_eq!(postings.get(&3), Some(&30));
}
#[test]
fn score_fine_candidates_skips_superseded_cells() {
use crate::supertable::manifest::{CellVectorSummary, ClusterCentroids, VectorSummary};
const DIM: u32 = 4;
let column = "emb";
let cell = |cell_id: u32, count: u32| CellVectorSummary {
cell_id: Some(cell_id),
clusters: ClusterCentroids::from_fp32(
1,
DIM,
&vec![cell_id as f32; DIM as usize],
vec![count],
),
};
let sf_id = Uuid::from_u128(0xC0FFEE);
let mut entry = synthetic_entry(sf_id);
entry.vector_summary.insert(
column.into(),
VectorSummary {
centroid: vec![0.0; DIM as usize],
cells: vec![cell(1, 10), cell(2, 20), cell(3, 30)],
},
);
let entries = vec![Arc::new(entry)];
let query = vec![0.0f32; DIM as usize];
let touched = |superseded: &BTreeMap<Uuid, BTreeSet<u32>>| -> HashSet<u32> {
let (cands, deferred) = score_fine_candidates(
&entries,
column,
&query,
Metric::L2Sq,
None,
true,
None,
superseded,
)
.expect("score");
cands
.iter()
.filter_map(|(_, _, _, cid, _)| *cid)
.chain(deferred.iter().filter_map(|d| d.cell_id))
.collect()
};
let empty = BTreeMap::new();
assert_eq!(touched(&empty), HashSet::from([1, 2, 3]));
let mut superseded = BTreeMap::new();
superseded.insert(sf_id, BTreeSet::from([2u32]));
assert_eq!(
touched(&superseded),
HashSet::from([1, 3]),
"superseded cell is not fine-scored or fetched"
);
}
#[test]
fn hybrid_vector_leg_uses_user_superfiles_not_hidden() {
let dim = 16usize;
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let opts = opts.with_storage(storage);
let st = Supertable::create(opts).expect("create");
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(0, 32, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
let reader = st.reader().expect("reader");
let user_uris: HashSet<_> = reader.manifest().superfiles.iter().map(|e| e.uri).collect();
assert!(
reader.vector_index_table().is_some(),
"hidden index must exist"
);
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let hits = reader
.hybrid_search(
"title",
"doc",
crate::superfile::fts::reader::BoolMode::Or,
"emb",
&q,
VectorSearchOptions::new(),
5,
)
.expect("hybrid");
assert!(!hits.is_empty());
for hit in &hits {
assert!(
user_uris.contains(&hit.superfile),
"hybrid vector leg must fan out on user superfiles, got {:?}",
hit.superfile
);
}
}
#[test]
fn vector_search_row_return_resolves_through_hidden_index() {
let dim = 16usize;
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let opts = opts.with_storage(storage);
let st = Supertable::create(opts).expect("create");
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(0, 16, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let batches = st
.reader()
.expect("reader")
.vector_search(
"emb",
&q,
5,
VectorSearchOptions::new(),
None,
Some(&["_id", "score"]),
)
.expect("vector_search rows");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert!(
rows >= 1,
"row-returning vector_search must resolve user rows"
);
}
#[test]
fn compaction_after_drain_preserves_search() {
use datafusion::prelude::{col, lit};
use crate::config::OptimizeOptions;
let dim = 16usize;
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
for c in 0..4u64 {
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(c * 32, 32, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
}
st.drain_vectors_to_cells_sync().expect("drain");
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let search = |st: &Supertable| {
st.reader()
.expect("reader")
.vector_hits(
"emb",
&q,
20,
VectorSearchOptions::new().with_nprobe(4),
None,
)
.expect("search")
.len()
};
assert!(
search(&st) >= 8,
"e_0's exact matches present pre-compaction"
);
st.optimize(&OptimizeOptions::default()).expect("optimize");
assert!(
search(&st) >= 8,
"compaction must preserve the exact-match docs"
);
let stats = st.delete(col("title").eq(lit("doc 0"))).expect("delete");
assert!(stats.n_tombstoned() >= 1, "delete tombstones matching docs");
st.optimize(&OptimizeOptions::default())
.expect("optimize after delete");
assert!(
!st.reader()
.expect("reader")
.vector_hits(
"emb",
&q,
20,
VectorSearchOptions::new().with_nprobe(4),
None
)
.expect("search after delete+compact")
.is_empty(),
"search still returns hits after delete + compaction"
);
}
#[test]
fn vector_cold_arm_matches_warm_top_k() {
use crate::test_helpers::lazy_foreground_disk_cache;
let dim = 16usize;
let schema = schema_with_vector(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let warm_hits = {
let st = Supertable::create(
options_one_superfile_per_commit(dim).with_storage(Arc::clone(&storage)),
)
.expect("create");
for c in 0..4u64 {
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(c * 32, 32, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
}
st.drain_vectors_to_cells_sync().expect("drain");
st.reader()
.expect("reader")
.vector_hits(
"emb",
&q,
20,
VectorSearchOptions::new().with_nprobe(4),
None,
)
.expect("warm search")
};
let cache_dir = tempfile::TempDir::new().expect("cache dir");
let cache = lazy_foreground_disk_cache(Arc::clone(&storage), cache_dir.path());
let st_cold = Supertable::open(
options_one_superfile_per_commit(dim)
.with_storage(Arc::clone(&storage))
.with_disk_cache(Arc::clone(&cache)),
)
.expect("open cold");
let cold_hits = st_cold
.reader()
.expect("reader")
.vector_hits(
"emb",
&q,
20,
VectorSearchOptions::new().with_nprobe(4),
None,
)
.expect("cold search");
assert!(
cache.stats().n_cold_fetches >= 1,
"the cold handle must actually read through the lazy cache \
(otherwise this fixture proves nothing)"
);
assert_eq!(
warm_hits, cold_hits,
"cache state must not change the top-k (warm deferred-global vs \
cold divided in-probe, both exact at an untruncated budget)"
);
}
#[test]
fn vector_search_multi_cell_rerank_over_larger_corpus() {
use crate::superfile::fts::reader::BoolMode;
let dim = 16usize;
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
for c in 0..4u64 {
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(c * 32, 32, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
}
st.drain_vectors_to_cells_sync().expect("drain");
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let hits = st
.reader()
.expect("reader")
.vector_hits(
"emb",
&q,
20,
VectorSearchOptions::new().with_nprobe(4),
None,
)
.expect("wide search");
let unfiltered_exact = near_count(&hits);
assert!(
unfiltered_exact >= 8,
"e_0 has 8 exact matches across commits; wide search must find \
them, got {unfiltered_exact}"
);
let filtered = st
.reader()
.expect("reader")
.vector_hits(
"emb",
&q,
20,
VectorSearchOptions::new().with_nprobe(4),
Some(VectorFilter {
column: "title",
query: "doc",
mode: BoolMode::Or,
}),
)
.expect("filtered wide search");
assert_eq!(
near_count(&filtered),
unfiltered_exact,
"filtered+nprobe must find the same exact matches as the \
unfiltered sweep"
);
}
const ORTHOGONAL_SCORE: f32 = 0.9;
const FIXTURE_DIM: usize = 16;
const FIXTURE_COMMITS: u64 = 4;
const FIXTURE_ROWS_PER_COMMIT: usize = 32;
const DOCS_PER_DIRECTION: usize =
FIXTURE_COMMITS as usize * FIXTURE_ROWS_PER_COMMIT / FIXTURE_DIM;
fn drained_three_direction_fixture() -> (tempfile::TempDir, Supertable, Vec<f32>, usize) {
let dim = FIXTURE_DIM;
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
for c in 0..FIXTURE_COMMITS {
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(
c * FIXTURE_ROWS_PER_COMMIT as u64,
FIXTURE_ROWS_PER_COMMIT,
dim,
schema.clone(),
))
.expect("append");
w.commit().expect("commit");
}
st.drain_vectors_to_cells_sync().expect("drain");
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
q[1] = 0.9;
q[2] = 0.8;
(dir, st, q, 3 * DOCS_PER_DIRECTION)
}
fn near_count(hits: &[SuperfileHit]) -> usize {
hits.iter().filter(|h| h.score < ORTHOGONAL_SCORE).count()
}
#[test]
fn caller_nprobe_widens_unfiltered_post_drain_sweep() {
let (_dir, st, q, k) = drained_three_direction_fixture();
let narrow_hits = st
.reader()
.expect("reader")
.vector_hits(
"emb",
&q,
k,
VectorSearchOptions::new().with_nprobe(1),
None,
)
.expect("narrow search");
let narrow = near_count(&narrow_hits);
assert!(
narrow >= DOCS_PER_DIRECTION && narrow < k,
"with_nprobe(1) must narrow the sweep below the law-widened \
default (got {narrow} of {k})"
);
let wide_hits = st
.reader()
.expect("reader")
.vector_hits(
"emb",
&q,
k,
VectorSearchOptions::new().with_nprobe(DOCS_PER_DIRECTION),
None,
)
.expect("wide search");
assert_eq!(
near_count(&wide_hits),
k,
"with_nprobe must widen the unfiltered post-drain sweep to all \
{k} exact neighbors across three cells — caller nprobe is an \
override, never discarded"
);
}
#[test]
fn width_pin_lifts_fine_depth_on_unfiltered_overrides() {
let base = CellRoutingParams {
nprobe_min: 1,
nprobe_max: 1,
fine_nprobe: 6,
..CellRoutingParams::default()
};
let mut r = base;
let sweep = apply_width_pin(&mut r, Some(64), None, false, 40);
assert_eq!((r.nprobe_min, r.nprobe_max), (64, 64));
assert_eq!(r.fine_nprobe, usize::MAX, "depth rides the pin");
assert_eq!(sweep, Some(40), "sweep width clamps to populated cells");
let mut r = base;
let sweep = apply_width_pin(&mut r, Some(64), None, true, 40);
assert_eq!((r.nprobe_min, r.nprobe_max), (64, 64));
assert_eq!(r.fine_nprobe, 6, "filtered floors untouched by the pin");
assert_eq!(sweep, None, "filtered sweeps keep pre-width budget");
let mut r = base;
let sweep = apply_width_pin(&mut r, None, Some(45), false, 256);
assert_eq!((r.nprobe_min, r.nprobe_max), (45, 45));
assert_eq!(r.fine_nprobe, usize::MAX, "depth rides the law");
assert_eq!(sweep, Some(45));
let mut r = base;
let sweep = apply_width_pin(&mut r, None, None, false, 256);
assert_eq!(r, base);
assert_eq!(sweep, None);
}
#[test]
fn global_shortlist_survives_minimal_rerank_budget() {
let (_dir, st, q, k) = drained_three_direction_fixture();
let hits = st
.reader()
.expect("reader")
.vector_hits(
"emb",
&q,
k,
VectorSearchOptions::new()
.with_nprobe(DOCS_PER_DIRECTION)
.with_rerank_mult(1),
None,
)
.expect("wide search, minimal budget");
assert_eq!(
near_count(&hits),
k,
"global shortlist selection at k x 1 must keep every planted \
neighbor across three cells"
);
}
#[test]
fn select_global_shortlist_is_deterministic() {
let cand = |est: f32, cell: usize, pos: u32, did: u32| ScanCandidate {
did,
estimate: est,
pos,
cluster_id: 0,
cell_idx: cell,
};
let a = vec![
(0usize, cand(0.9, 0, 1, 1)),
(1usize, cand(0.9, 0, 1, 1)),
(0usize, cand(0.5, 1, 2, 2)),
(1usize, cand(0.7, 0, 3, 3)),
];
let mut b = a.clone();
b.reverse();
let pick = |v: Vec<(usize, ScanCandidate)>| {
select_global_shortlist(v, 3, 0)
.into_iter()
.map(|(si, c)| (si, c.cell_idx, c.pos, c.did))
.collect::<Vec<_>>()
};
let mut from_a = pick(a);
let mut from_b = pick(b);
from_a.sort_unstable();
from_b.sort_unstable();
assert_eq!(
from_a, from_b,
"the kept SET must be arrival-order independent (its internal \
order is unspecified — the partition does not sort)"
);
assert_eq!(from_a.len(), 3, "limit is a hard truncation");
assert_eq!(from_a, vec![(0, 0, 1, 1), (1, 0, 1, 1), (1, 0, 3, 3)]);
}
#[test]
fn rerank_law_yields_to_caller_filter_and_uncalibrated_k() {
let routing = CellRoutingParams {
rerank_for_k: [40, 320, 2400, 0],
..CellRoutingParams::default()
};
let r = Some(&routing);
assert_eq!(rerank_mult_from_law(true, false, None, r, 10), Some(32));
assert_eq!(rerank_mult_from_law(true, false, Some(8), r, 10), None);
assert_eq!(rerank_mult_from_law(true, true, None, r, 10), None);
assert_eq!(rerank_mult_from_law(false, false, None, r, 10), None);
assert_eq!(rerank_mult_from_law(true, false, None, r, 1000), None);
assert_eq!(rerank_mult_from_law(true, false, None, None, 10), None);
}
#[test]
fn select_global_shortlist_cell_floor_rescues_swamped_cells() {
let cand = |est: f32, cell: usize, pos: u32, did: u32| ScanCandidate {
did,
estimate: est,
pos,
cluster_id: 0,
cell_idx: cell,
};
let mut pooled: Vec<(usize, ScanCandidate)> =
(0..10).map(|i| (0usize, cand(0.9, 0, i, i))).collect();
pooled.push((0, cand(0.30, 1, 100, 100)));
pooled.push((0, cand(0.20, 1, 101, 101)));
pooled.push((0, cand(0.10, 1, 102, 102)));
let kept = select_global_shortlist(pooled.clone(), 4, 0);
assert!(
kept.iter().all(|(_, c)| c.cell_idx == 0),
"without a floor the flooded cell evicts cell 1 entirely"
);
let kept = select_global_shortlist(pooled, 4, 2);
let cell1: Vec<u32> = kept
.iter()
.filter(|(_, c)| c.cell_idx == 1)
.map(|(_, c)| c.did)
.collect();
assert_eq!(
cell1,
vec![100, 101],
"the floor keeps cell 1's two best despite global eviction"
);
assert_eq!(
kept.iter().filter(|(_, c)| c.cell_idx == 0).count(),
4,
"the global prefix is untouched by the rescue"
);
}
#[test]
fn select_global_shortlist_widening_never_evicts_floored() {
let cand = |est: f32, cell: usize, pos: u32, did: u32| ScanCandidate {
did,
estimate: est,
pos,
cluster_id: 0,
cell_idx: cell,
};
const FLOOR: usize = 2;
const LIMIT: usize = 4;
let narrow: Vec<(usize, ScanCandidate)> = vec![
(0, cand(0.9, 0, 0, 0)),
(0, cand(0.8, 0, 1, 1)),
(0, cand(0.4, 1, 2, 2)),
(0, cand(0.3, 1, 3, 3)),
];
let mut wide = narrow.clone();
for i in 0..8u32 {
wide.push((0, cand(0.7, 2, 10 + i, 10 + i)));
}
let keep_ids = |v: Vec<(usize, ScanCandidate)>| {
let mut ids: Vec<u32> = select_global_shortlist(v, LIMIT, FLOOR)
.into_iter()
.map(|(_, c)| c.did)
.collect();
ids.sort_unstable();
ids
};
let narrow_kept = keep_ids(narrow);
let wide_kept = keep_ids(wide);
for id in &narrow_kept {
assert!(
wide_kept.contains(id),
"widening dropped candidate {id}: floored survivors must \
be immune to added cells (kept narrow {narrow_kept:?} vs \
wide {wide_kept:?})"
);
}
}
#[test]
fn select_global_shortlist_matches_naive_floor_definition() {
let mut state: u64 = 0x5eed_1234_abcd_9876;
let mut next = move || {
state = state
.wrapping_mul(6364136223846793005)
.wrapping_add(1442695040888963407);
state
};
const POOL: usize = 4_000;
const UNITS: usize = 3;
const CELLS: usize = 17;
const LIMIT: usize = 300;
const FLOOR: usize = 9;
let mut pooled: Vec<(usize, ScanCandidate)> = Vec::with_capacity(POOL);
for i in 0..POOL {
let r = next();
let est = ((r >> 32) % 64) as f32 / 64.0;
pooled.push((
(r % UNITS as u64) as usize,
ScanCandidate {
did: i as u32,
estimate: est,
pos: i as u32,
cluster_id: 0,
cell_idx: ((r >> 8) % CELLS as u64) as usize,
},
));
}
let cmp = |a: &(usize, ScanCandidate), b: &(usize, ScanCandidate)| {
b.1.estimate.total_cmp(&a.1.estimate).then_with(|| {
(a.0, a.1.cell_idx, a.1.pos, a.1.did).cmp(&(b.0, b.1.cell_idx, b.1.pos, b.1.did))
})
};
let mut sorted = pooled.clone();
sorted.sort_by(cmp);
let mut expect: HashSet<u32> = sorted[..LIMIT].iter().map(|(_, c)| c.did).collect();
let mut per_cell: HashMap<(usize, usize), usize> = HashMap::new();
for (si, cand) in &sorted {
let taken = per_cell.entry((*si, cand.cell_idx)).or_default();
if *taken < FLOOR {
expect.insert(cand.did);
*taken += 1;
}
}
let mut grouped = pooled.clone();
grouped.sort_by_key(|(si, c)| (*si, c.cell_idx));
let kept: HashSet<u32> = select_global_shortlist(grouped, LIMIT, FLOOR)
.into_iter()
.map(|(_, c)| c.did)
.collect();
assert_eq!(kept, expect, "kept set must match the naive definition");
let kept_ungrouped: HashSet<u32> = select_global_shortlist(pooled, LIMIT, FLOOR)
.into_iter()
.map(|(_, c)| c.did)
.collect();
assert!(
kept_ungrouped.is_superset(&expect),
"ungrouped pools must never lose a guaranteed survivor"
);
}
#[test]
fn drain_calibrated_width_law_widens_default_search() {
let (_dir, st, q, k) = drained_three_direction_fixture();
let hits = st
.reader()
.expect("reader")
.vector_hits("emb", &q, k, VectorSearchOptions::new(), None)
.expect("default search");
let near = near_count(&hits);
assert_eq!(
near, k,
"drain-calibrated width law must widen the default sweep to all \
{k} exact neighbors across three cells, got {near}"
);
}
#[test]
fn count_sums_matches_across_superfiles() {
let (_dir, st, _q, _k) = drained_three_direction_fixture();
let reader = st.reader().expect("reader");
assert_eq!(
reader
.count("title", "5", BoolMode::And)
.expect("sparse count"),
FIXTURE_COMMITS,
"one match per superfile"
);
assert_eq!(
reader
.count("title", "doc", BoolMode::And)
.expect("dense count"),
(FIXTURE_COMMITS as usize * FIXTURE_ROWS_PER_COMMIT) as u64,
"every row matches"
);
}
#[test]
fn bm25_global_stats_scores_across_superfiles() {
let (_dir, st, _q, _k) = drained_three_direction_fixture();
let reader = st.reader().expect("reader");
let batches = reader
.bm25_search("title", "5", 8, BoolMode::And, Bm25Stats::Global, None)
.expect("global-stats bm25");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(
rows, FIXTURE_COMMITS as usize,
"token \"5\" lives in one row per superfile; global-idf \
scoring must find all of them"
);
}
#[test]
fn raw_cosine_corpus_ranks_like_the_normalized_twin() {
const RAW_SCALE: f32 = 3.7;
let dim = FIXTURE_DIM;
let build = |scale: f32| {
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
let mut w = st.writer().expect("writer");
let base = build_vector_batch(0, FIXTURE_ROWS_PER_COMMIT, dim, schema.clone());
let emb = base
.column(1)
.as_any()
.downcast_ref::<FixedSizeListArray>()
.expect("fsl")
.clone();
let scaled: Vec<f32> = emb
.values()
.as_any()
.downcast_ref::<Float32Array>()
.expect("f32")
.values()
.iter()
.map(|v| v * scale)
.collect();
let fsl = FixedSizeListArray::try_new(
Arc::new(arrow_schema::Field::new(
"item",
arrow_schema::DataType::Float32,
true,
)),
dim as i32,
Arc::new(Float32Array::from(scaled)) as Arc<dyn Array>,
None,
)
.expect("fsl scaled");
let batch =
RecordBatch::try_new(schema.clone(), vec![base.column(0).clone(), Arc::new(fsl)])
.expect("batch");
w.append(&batch).expect("append");
w.commit().expect("commit");
drop(w);
st.drain_vectors_to_cells_sync().expect("drain");
(dir, st)
};
let (_d_unit, unit) = build(1.0);
let (_d_raw, raw) = build(RAW_SCALE);
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
q[1] = 0.7;
let scaled_q: Vec<f32> = q.iter().map(|v| v * RAW_SCALE).collect();
let unit_hits = unit
.reader()
.expect("reader")
.vector_hits("emb", &q, 8, VectorSearchOptions::new(), None)
.expect("unit search");
let raw_hits = raw
.reader()
.expect("reader")
.vector_hits("emb", &scaled_q, 8, VectorSearchOptions::new(), None)
.expect("raw search");
let positions = |hits: &[SuperfileHit]| -> Vec<i128> {
let ids: Vec<i128> = hits.iter().map(|h| h.stable_id.expect("id")).collect();
let base = *ids.iter().min().expect("hits");
ids.iter().map(|id| id - base).collect()
};
assert_eq!(
positions(&unit_hits),
positions(&raw_hits),
"raw corpus + scaled query must rank exactly like the unit twin"
);
for (u, r) in unit_hits.iter().zip(raw_hits.iter()) {
assert!(
(u.score - r.score).abs() < 1e-3,
"calibrated scores must match: {} vs {}",
u.score,
r.score
);
}
}
#[test]
fn vector_filter_restricts_hits_to_predicate_matches() {
let (_dir, st, q, _k) = drained_three_direction_fixture();
let reader = st.reader().expect("reader");
let matched = reader
.vector_hits(
"emb",
&q,
10,
VectorSearchOptions::new(),
Some(VectorFilter {
column: "title",
query: "5",
mode: BoolMode::And,
}),
)
.expect("sparse filtered search");
assert_eq!(
matched.len(),
FIXTURE_COMMITS as usize,
"the predicate matches one row per commit — nothing more"
);
let predicate_rows: HashSet<i128> = st
.exact_match("title", "doc 5", Some(&["_id"]))
.expect("exact-match oracle")
.iter()
.filter_map(|b| {
b.column_by_name("_id")
.and_then(|c| c.as_any().downcast_ref::<Decimal128Array>())
})
.flat_map(|c| c.values().iter().copied())
.collect();
assert_eq!(
predicate_rows.len(),
FIXTURE_COMMITS as usize,
"oracle sanity: exactly one \"doc 5\" row per commit"
);
assert!(
matched
.iter()
.all(|h| h.stable_id.is_some_and(|id| predicate_rows.contains(&id))),
"every hit is a predicate row, not a nearest neighbor: {matched:?}"
);
let all = reader
.vector_hits(
"emb",
&q,
10,
VectorSearchOptions::new(),
Some(VectorFilter {
column: "title",
query: "doc",
mode: BoolMode::And,
}),
)
.expect("match-all filtered search");
assert_eq!(
all.len(),
10,
"a predicate matching every row fills the full top-k"
);
}
#[test]
fn recalibrated_laws_serve_full_recall_at_default_search() {
let (_dir, st, q, k) = drained_three_direction_fixture();
let hidden = st
.reader()
.expect("reader")
.vector_index_table()
.expect("hidden index")
.clone();
let strategy = hidden
.reader()
.expect("hidden reader")
.manifest()
.get_partition_strategy();
let PartitionStrategy::VectorCell { clusters, .. } = strategy else {
panic!("hidden index must be VectorCell");
};
let mut direction = vec![0.0f32; FIXTURE_DIM];
direction[0] = 1.0;
let target_cell = clusters.nearest_cell(Metric::Cosine, &direction);
hidden
.block_on_query(split_overflow_cell(
hidden.inner().clone(),
target_cell,
0.0,
))
.expect("split")
.expect("populated cell must split");
let stamped = hidden
.block_on_query(recalibrate_probe_laws(hidden.inner()))
.expect("recalibrate");
assert!(stamped, "the reshaped grid must restamp the laws");
let hits = st
.reader()
.expect("reader")
.vector_hits("emb", &q, k, VectorSearchOptions::new(), None)
.expect("default search after recalibration");
let near = near_count(&hits);
assert_eq!(
near, k,
"recalibrated laws must serve the full {k} planted neighbors \
at default settings, got {near}"
);
}
#[test]
fn incremental_drain_never_narrows_the_width_law() {
let (_dir, st, q, k) = drained_three_direction_fixture();
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(
(FIXTURE_COMMITS * FIXTURE_ROWS_PER_COMMIT as u64) + 1,
DOCS_PER_DIRECTION,
FIXTURE_DIM,
schema_with_vector(FIXTURE_DIM),
))
.expect("append delta");
w.commit().expect("commit delta");
st.drain_vectors_to_cells_sync().expect("incremental drain");
let hits = st
.reader()
.expect("reader")
.vector_hits("emb", &q, k, VectorSearchOptions::new(), None)
.expect("default search after incremental drain");
let near = near_count(&hits);
assert_eq!(
near,
k,
"the delta-only calibration must not narrow the stamped law: \
default search lost {} of {k} planted neighbors",
k - near
);
}
#[test]
fn supertable_vector_search_wrapper_returns_rows() {
let dim = 16usize;
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(0, 16, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
drop(w);
st.drain_vectors_to_cells_sync().expect("drain");
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let batches = st
.vector_search(
"emb",
&q,
5,
VectorSearchOptions::new(),
None,
Some(&["_id"]),
)
.expect("handle-level vector_search");
assert!(
batches.iter().map(|b| b.num_rows()).sum::<usize>() >= 1,
"handle wrapper must return rows"
);
}
#[test]
fn vector_hits_global_allow_restricts_to_allowed_ids() {
use roaring::RoaringBitmap;
let dim = 16usize;
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(0, 16, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
drop(w);
let allow: Arc<RoaringBitmap> = Arc::new([0u32, 1, 2].into_iter().collect());
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let hits = block_on(st.reader().expect("reader").vector_hits_global_allow_async(
"emb",
&q,
16,
VectorSearchOptions::new().with_nprobe(32),
allow,
))
.expect("global-allow search");
assert!(!hits.is_empty(), "the e_0 doc is allowed and must be found");
assert!(
hits.len() <= 3,
"only the 3 allowed global rows may appear, got {}",
hits.len()
);
}
#[test]
fn prepare_vector_stable_allow_maps_valid_drained_id() {
use arrow_array::Decimal128Array;
let dim = 16usize;
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(0, 16, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
drop(w);
st.drain_vectors_to_cells_sync().expect("drain");
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let batches = st
.reader()
.expect("reader")
.vector_search(
"emb",
&q,
1,
VectorSearchOptions::new(),
None,
Some(&["_id"]),
)
.expect("row search");
let id = batches
.iter()
.find_map(|b| {
b.column_by_name("_id")
.and_then(|c| c.as_any().downcast_ref::<Decimal128Array>())
.filter(|c| !c.is_empty())
.map(|c| c.value(0))
})
.expect("a resolved _id");
let prepared = block_on(
st.reader()
.expect("reader")
.prepare_vector_stable_allow_async(Arc::new(vec![id])),
)
.expect("valid drained id must map");
assert!(
prepared.use_hidden_index,
"post-drain allow-set is keyed by the hidden index"
);
assert!(
!prepared.allow_by_uri.is_empty(),
"a valid id resolves to a non-empty hidden-cell allow-set"
);
}
#[test]
fn vector_search_rows_post_drain_resolve_hidden_ids() {
use arrow_array::Decimal128Array;
let dim = 16usize;
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(0, 16, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
drop(w);
st.drain_vectors_to_cells_sync().expect("drain");
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let batches = st
.reader()
.expect("reader")
.vector_search(
"emb",
&q,
5,
VectorSearchOptions::new(),
None,
Some(&["_id", "score"]),
)
.expect("post-drain row search");
let mut ids = Vec::new();
for b in &batches {
let col = b
.column_by_name("_id")
.expect("_id column")
.as_any()
.downcast_ref::<Decimal128Array>()
.expect("_id is decimal128");
for i in 0..col.len() {
ids.push(col.value(i));
}
}
assert_eq!(ids.len(), 5, "k=5 over 16 docs returns 5 rows");
assert_eq!(
ids[0],
*ids.iter().min().expect("ids is non-empty"),
"the exact-match doc must rank first, got {ids:?}"
);
}
#[test]
fn filtered_vector_search_row_return_fans_out_over_user_superfiles() {
use crate::superfile::fts::reader::BoolMode;
let dim = 16usize;
let schema = schema_with_vector(dim);
let opts = options_one_superfile_per_commit(dim);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
for start in [0u64, 16, 32] {
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(start, 16, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
}
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let batches = st
.reader()
.expect("reader")
.vector_search(
"emb",
&q,
10,
VectorSearchOptions::new(),
Some(VectorFilter {
column: "title",
query: "doc",
mode: BoolMode::Or,
}),
Some(&["_id", "score"]),
)
.expect("filtered row search");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert!(
rows >= 1,
"filtered fan-out must resolve rows across user superfiles"
);
}
#[test]
fn vector_hits_filtered_by_plan_returns_matching_docs() {
use std::collections::HashSet;
use datafusion::prelude::{col, lit};
use crate::{
superfile::vector::rerank_codec::RerankCodec,
supertable::query::candidate::CandidatePlan,
};
let dim = 16usize;
let schema = schema_with_vector(dim);
let pool = Arc::new(
rayon::ThreadPoolBuilder::new()
.num_threads(1)
.build()
.expect("pool"),
);
let opts = SupertableOptions::new(
schema.clone(),
vec![FtsConfig {
column: "title".into(),
positions: false,
}],
vec![VectorConfig {
column: "emb".into(),
dim,
rot_seed: 7,
metric: Metric::Cosine,
rerank_codec: RerankCodec::Sq8Residual,
provided_centroids: None,
}],
Some(tok()),
)
.expect("valid options")
.with_writer_pool(pool);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(0, 32, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
drop(w);
st.drain_vectors_to_cells_sync().expect("drain");
let reader = st.reader().expect("reader");
let manifest = reader.manifest();
let fts_cols: HashSet<&str> = HashSet::from(["title"]);
let filters = [col("title").eq(lit("doc"))];
let plan = CandidatePlan::from_filters(&filters, &fts_cols, &|col| {
manifest.options.fts_tokenizer_for(col)
});
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let hits = block_on(reader.vector_hits_filtered_by_plan(
"emb",
&q,
10,
VectorSearchOptions::new(),
&plan,
))
.expect("plan-filtered vector search");
assert!(
!hits.is_empty(),
"the title-token plan must admit docs for vector ranking"
);
}
#[test]
fn filtered_vector_search_post_drain_uses_hidden_index() {
use crate::superfile::vector::rerank_codec::RerankCodec;
let dim = 16usize;
let schema = schema_with_vector(dim);
let pool = Arc::new(
rayon::ThreadPoolBuilder::new()
.num_threads(1)
.build()
.expect("pool"),
);
let opts = SupertableOptions::new(
schema.clone(),
vec![FtsConfig {
column: "title".into(),
positions: false,
}],
vec![VectorConfig {
column: "emb".into(),
dim,
rot_seed: 7,
metric: Metric::Cosine,
rerank_codec: RerankCodec::Sq8Residual,
provided_centroids: None,
}],
Some(tok()),
)
.expect("valid options")
.with_writer_pool(pool);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let opts = opts.with_storage(storage);
let st = Supertable::create(opts).expect("create");
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(0, 32, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
st.drain_vectors_to_cells_sync().expect("drain");
let reader = st.reader().expect("reader");
let user_uris: HashSet<_> = reader.manifest().superfiles.iter().map(|e| e.uri).collect();
let hidden = reader
.vector_index_table()
.expect("hidden index must exist");
let hidden_uris: HashSet<_> = hidden
.reader()
.expect("reader")
.manifest()
.superfiles
.iter()
.map(|e| e.uri)
.collect();
assert!(
!hidden_uris.is_empty(),
"drain must publish at least one hidden superfile"
);
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let hits = reader
.vector_hits(
"emb",
&q,
5,
VectorSearchOptions::new(),
Some(VectorFilter {
column: "title",
query: "doc",
mode: crate::superfile::fts::reader::BoolMode::Or,
}),
)
.expect("filtered vector_hits");
assert!(!hits.is_empty(), "filtered search must return hits");
for hit in &hits {
assert!(
hidden_uris.contains(&hit.superfile),
"post-drain filtered hits must come from hidden superfiles, got {:?} \
(user={user_uris:?}, hidden={hidden_uris:?})",
hit.superfile
);
assert!(
!user_uris.contains(&hit.superfile),
"post-drain filtered hits must not come from user superfiles"
);
}
let mapping_error =
block_on(reader.prepare_vector_stable_allow_async(Arc::new(vec![i128::MAX])))
.err()
.expect("unknown drained id must fail hidden mapping");
assert!(
mapping_error
.to_string()
.contains("did not map to any hidden superfile"),
"unexpected mapping error: {mapping_error}"
);
}
#[test]
fn vector_search_post_drain_excludes_deleted() {
use datafusion::prelude::{col, lit};
use crate::superfile::vector::rerank_codec::RerankCodec;
let dim = 16usize;
let schema = schema_with_vector(dim);
let pool = Arc::new(
rayon::ThreadPoolBuilder::new()
.num_threads(1)
.build()
.expect("pool"),
);
let opts = SupertableOptions::new(
schema.clone(),
vec![FtsConfig {
column: "title".into(),
positions: false,
}],
vec![VectorConfig {
column: "emb".into(),
dim,
rot_seed: 7,
metric: Metric::Cosine,
rerank_codec: RerankCodec::Sq8Residual,
provided_centroids: None,
}],
Some(tok()),
)
.expect("valid options")
.with_writer_pool(pool);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
let mut w = st.writer().expect("writer");
w.append(&build_vector_batch(0, 32, dim, schema.clone()))
.expect("append");
w.commit().expect("commit");
drop(w); st.drain_vectors_to_cells_sync().expect("drain");
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let hits_before = st
.reader()
.expect("reader")
.vector_hits("emb", &q, 32, VectorSearchOptions::new(), None)
.expect("pre-delete search");
assert!(!hits_before.is_empty(), "docs retrievable pre-delete");
let stats = st.delete(col("title").eq(lit("doc 0"))).expect("delete");
assert_eq!(stats.n_tombstoned(), 1, "exactly one row tombstoned");
let hits_after = st
.reader()
.expect("reader")
.vector_hits("emb", &q, 32, VectorSearchOptions::new(), None)
.expect("post-delete search");
assert_eq!(
hits_after.len(),
hits_before.len() - 1,
"the deleted doc must drop out of the results"
);
}
#[test]
fn commit_user_superfiles_cell_packed_no_duplicate_parquet_rows() {
use crate::superfile::vector::layout::VectorLayout;
let dim = 16;
let st = Supertable::create(options_one_superfile_per_commit(dim)).expect("create");
let mut w = st.writer().expect("writer");
let schema = st.options().schema.clone();
let n = 200usize;
w.append(&build_vector_batch(0, n, dim, schema))
.expect("append");
w.commit().expect("commit");
let r = st.reader().expect("reader");
let manifest = r.manifest();
assert!(
!manifest.superfiles.is_empty(),
"commit must publish user superfiles"
);
let mut total_primary_rows = 0u64;
for entry in manifest.superfiles.iter() {
assert_eq!(
entry.vector_layout,
VectorLayout::MultiCellIvf,
"commit must write cell-packed MultiCellIvf user superfiles, got {:?}",
entry.vector_layout
);
total_primary_rows += entry.n_docs;
}
assert_eq!(
total_primary_rows, n as u64,
"each ingested row is a Parquet primary exactly once; boundary stubs \
must not add Parquet rows (got {total_primary_rows}, expected {n})"
);
}
#[test]
fn vector_search_dedups_and_resolves_with_stub_boundaries() {
use arrow_array::Decimal128Array;
use crate::superfile::vector::rerank_codec::RerankCodec;
let dim = 16usize;
let schema = schema_with_vector(dim);
let pool = Arc::new(
rayon::ThreadPoolBuilder::new()
.num_threads(1)
.build()
.expect("pool"),
);
let opts = SupertableOptions::new(
schema.clone(),
vec![FtsConfig {
column: "title".into(),
positions: false,
}],
vec![VectorConfig {
column: "emb".into(),
dim,
rot_seed: 7,
metric: Metric::Cosine,
rerank_codec: RerankCodec::Sq8Residual,
provided_centroids: None,
}],
Some(tok()),
)
.expect("valid options")
.with_writer_pool(pool);
let dir = tempfile::TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(crate::storage::LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(opts.with_storage(storage)).expect("create");
let mut w = st.writer().expect("writer");
let n = 200usize;
w.append(&build_vector_batch(0, n, dim, schema))
.expect("append");
w.commit().expect("commit");
let r = st.reader().expect("reader");
let mut q = vec![0.0f32; dim];
q[0] = 1.0;
let k = 20usize;
let batches = r
.vector_search(
"emb",
&q,
k,
VectorSearchOptions::new().with_nprobe(4),
None,
Some(&["_id", "title"]),
)
.expect("vector_search");
let mut seen: HashSet<i128> = HashSet::new();
let mut total = 0usize;
for b in &batches {
let ids = b
.column(0)
.as_any()
.downcast_ref::<Decimal128Array>()
.expect("_id column is Decimal128");
let titles = b
.column(1)
.as_any()
.downcast_ref::<LargeStringArray>()
.expect("title column is LargeString");
assert_eq!(titles.len(), ids.len());
for i in 0..ids.len() {
total += 1;
assert!(!titles.value(i).is_empty());
assert!(
seen.insert(ids.value(i)),
"duplicate _id {} in results — a boundary stub was not deduped \
against its primary",
ids.value(i)
);
}
}
assert_eq!(total, k, "search must return k distinct rows, got {total}");
}
#[test]
fn user_multicell_inline_ids_match_parquet_id_column() {
use arrow_array::Decimal128Array;
let dim = 16usize;
let schema = schema_with_vector(dim);
let dir = TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(
options_one_superfile_per_commit(dim).with_storage(Arc::clone(&storage)),
)
.expect("create");
let mut w = st.writer().expect("writer");
let n = 200usize;
w.append(&build_vector_batch(0, n, dim, schema))
.expect("append");
w.commit().expect("commit");
let r = st.reader().expect("reader");
let manifest = r.manifest();
let mut checked_files = 0usize;
for entry in manifest.superfiles.iter() {
let reader = manifest
.options
.store
.reader(&entry.uri)
.expect("writer-published reader");
let vec_reader = reader.vec().expect("vector reader");
let locals: Vec<u32> = (0..entry.n_docs as u32).collect();
let batch = reader
.take_by_local_doc_ids(&locals, &[reader.id_column()])
.expect("take _id column");
let truth = batch
.column(0)
.as_any()
.downcast_ref::<Decimal128Array>()
.expect("_id is Decimal128");
let Some(inline) = vec_reader.inline_stable_ids_for_locals(&locals) else {
panic!(
"inline stable-id lookup unavailable on user superfile {:?} \
(layout {:?}): stable_ids_for_tagged_hits would silently fall \
back to the _id page read",
entry.uri, entry.vector_layout
);
};
for (i, &local) in locals.iter().enumerate() {
assert_eq!(
inline[i],
truth.value(i),
"inline stable-id for parquet-local {local} in {:?} diverges \
from the _id column",
entry.uri
);
}
checked_files += 1;
}
assert!(checked_files > 0, "commit published no user superfiles");
}
#[test]
fn global_union_includes_undrained_user_delta() {
let dim = 16usize;
let schema = schema_with_vector(dim);
let dir = TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("storage"));
let st = Supertable::create(
options_one_superfile_per_commit(dim).with_storage(Arc::clone(&storage)),
)
.expect("create");
let mut writer = st.writer().expect("writer");
writer
.append(&build_vector_batch(0, 8, dim, Arc::clone(&schema)))
.expect("append base");
writer.commit().expect("commit base");
drop(writer);
st.drain_vectors_to_cells_sync().expect("drain base");
let mut writer = st.writer().expect("writer delta");
writer
.append(&build_vector_batch(15, 1, dim, schema))
.expect("append delta");
writer.commit().expect("commit delta");
drop(writer);
let reader = st.reader().expect("reader");
let hidden = reader.vector_index_table().expect("hidden index");
let drained = hidden
.reader()
.expect("reader")
.manifest()
.get_drained_ranges();
let undrained: Vec<_> = reader
.manifest()
.superfiles
.iter()
.filter(|entry| !drained.contains(entry.birth_version))
.collect();
assert_eq!(undrained.len(), 1);
let mut query = vec![0.0f32; dim];
query[15] = 1.0;
let hits = reader
.vector_hits("emb", &query, 1, VectorSearchOptions::new(), None)
.expect("global union search");
assert_eq!(hits.len(), 1);
assert_eq!(hits[0].superfile, undrained[0].uri);
}
}