mod build;
mod cache;
mod diskann;
mod diskann_index;
mod flat;
mod fts;
mod hnsw;
mod ivf;
mod ordinal_map;
mod ordinals;
mod product_quantization;
mod quantization;
mod rabitq;
mod rabitq_index;
mod rebuild;
mod scalar;
mod vamana;
mod vector_query;
use crate::config::IoBackend;
use crate::doc::DocumentMap;
use crate::error::{Error, Result};
use crate::query::SearchQuery;
use crate::schema::{CollectionSchema, IndexParams};
use crate::stats::IndexStat;
use crate::types::{DataType, IndexType, MetricType};
use build::{build_vector_index, encode_vector};
use diskann_index::DiskannIndex;
use flat::FlatIndex;
use fts::FtsIndexRegistry;
use hnsw::HnswIndex;
use ivf::IvfIndex;
use ordinal_map::OrdinalMap;
pub(crate) use ordinals::OrdinalScores;
use ordinals::{OrdinalSet, OrdinalTable};
use quantization::{
dense_query_norm, score_dense_with_query_norm, score_with_query_norm, QuantizedVector,
};
use rabitq_index::{HnswRabitqIndex, IvfRabitqIndex};
use roaring::RoaringTreemap;
use scalar::{ScalarCandidates, ScalarIndexRegistry};
use std::collections::{BTreeMap, BTreeSet};
use std::sync::Arc;
use std::sync::OnceLock;
use vamana::VamanaIndex;
const MIN_DELTA_COMPACTION: usize = 64;
const MAX_DELTA_COMPACTION: usize = 2_048;
const DELTA_COMPACTION_DIVISOR: usize = 8;
const SCALAR_EXACT_PREFILTER_MIN: usize = 4_096;
const SCALAR_EXACT_PREFILTER_PER_RESULT: usize = 64;
#[derive(Debug, Clone, Default)]
pub(crate) struct IndexRegistry {
ordinals: OrdinalTable,
indexes: BTreeMap<String, VectorIndex>,
scalar_indexes: ScalarIndexRegistry,
fts_indexes: FtsIndexRegistry,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
struct VectorIndex {
#[serde(with = "cache::index_params_serde")]
params: IndexParams,
source_revision: u64,
base: Arc<VectorIndexBase>,
delta: BTreeMap<u64, QuantizedVector>,
delta_ordinals: RoaringTreemap,
tombstones: RoaringTreemap,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
struct VectorIndexBase {
vectors: OrdinalMap<QuantizedVector>,
vector_ordinals: RoaringTreemap,
kind: VectorIndexKind,
#[serde(skip)]
diskann: Option<Arc<diskann::FieldReader>>,
#[serde(skip)]
cosine_norms: OnceLock<Vec<f32>>,
#[serde(skip)]
exact_cosine_norms: OnceLock<Vec<f64>>,
#[serde(skip)]
exact_cosine_inv_norms: OnceLock<Vec<f64>>,
#[serde(skip)]
dense_f32: OnceLock<Option<DenseF32Base>>,
#[serde(skip)]
dense_f64: OnceLock<Option<DenseF64Base>>,
}
#[derive(Clone, Debug)]
struct DenseF32Base {
dimension: usize,
values: Vec<f32>,
}
#[derive(Clone, Debug)]
struct DenseF64Base {
dimension: usize,
values: Vec<f64>,
}
struct AnnSearchContext<'a> {
query: &'a SearchQuery,
vector: &'a [f32],
topk: usize,
metric: MetricType,
allowed: Option<&'a RoaringTreemap>,
eligible_count: usize,
ordinals: &'a OrdinalTable,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
enum VectorIndexKind {
Flat(FlatIndex),
Hnsw(HnswIndex),
HnswRabitq(HnswRabitqIndex),
Ivf(IvfIndex),
IvfRabitq(IvfRabitqIndex),
Diskann(DiskannIndex),
Vamana(VamanaIndex),
}
#[derive(Debug, Clone)]
pub(crate) struct CandidateSelection {
ids: OrdinalSet,
}
#[derive(Debug, Clone)]
pub(crate) struct CandidatePlan {
pub selection: Option<CandidateSelection>,
pub fts_scores: Option<OrdinalScores>,
pub used_ann: bool,
pub diskann_sector_reads: u64,
pub diskann_io_backend: Option<IoBackend>,
pub used_scalar: bool,
pub used_fts_index: bool,
}
struct AnnOrdinals {
ids: RoaringTreemap,
diskann_sector_reads: u64,
diskann_io_backend: Option<IoBackend>,
}
struct AnnCandidates {
selection: CandidateSelection,
diskann_sector_reads: u64,
diskann_io_backend: Option<IoBackend>,
}
impl CandidatePlan {
pub(crate) fn candidate_count(&self, document_count: usize) -> u64 {
self.fts_scores.as_ref().map_or_else(
|| {
self.selection.as_ref().map_or_else(
|| u64::try_from(document_count).unwrap_or(u64::MAX),
CandidateSelection::count,
)
},
OrdinalScores::candidate_count,
)
}
}
impl CandidateSelection {
pub(crate) fn count(&self) -> u64 {
u64::try_from(self.ids.len()).unwrap_or(u64::MAX)
}
pub(crate) fn ids(&self) -> impl Iterator<Item = &str> {
self.ids.ids()
}
pub(crate) fn iter_ordinals(&self) -> impl Iterator<Item = u64> + '_ {
self.ids.iter_ordinals()
}
pub(crate) fn id(&self, ordinal: u64) -> Option<&str> {
self.ids.id(ordinal)
}
pub(crate) fn contains(&self, id: &str) -> bool {
self.ids.contains(id)
}
}
impl IndexRegistry {
pub(crate) fn document_ordinal(&self, id: &str) -> Option<u64> {
self.ordinals.ordinal(id)
}
pub(crate) fn exact_unquantized_f32_score_at(
&self,
field: &str,
ordinal: u64,
query: &[f32],
query_norm: f64,
metric: MetricType,
) -> Option<f64> {
let index = self.indexes.get(field)?;
let coordinates = index.unquantized_f32(ordinal)?;
let score = score_dense_with_query_norm(query, coordinates, metric, query_norm);
score.is_finite().then_some(score)
}
pub(crate) fn build(
schema: &CollectionSchema,
docs: &DocumentMap,
source_revision: u64,
) -> Result<Self> {
let ordinals = OrdinalTable::build(docs)?;
let mut indexes = BTreeMap::new();
for field in &schema.vectors {
let Some(params) = field.index_params.as_ref() else {
continue;
};
if !builds_packed_vector_index(params.index_type, field.data_type) {
continue;
}
indexes.insert(
field.name.clone(),
build_vector_index(
docs,
&field.name,
field.dimension,
params,
source_revision,
&ordinals,
)?,
);
}
let scalar_indexes = ScalarIndexRegistry::build(schema, docs, source_revision, &ordinals)?;
let fts_indexes = FtsIndexRegistry::build(schema, docs, source_revision, &ordinals)?;
Ok(Self {
ordinals,
indexes,
scalar_indexes,
fts_indexes,
})
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn restore_cache(
bytes: &[u8],
diskann_file: Option<crate::storage::PositionedFile>,
io_backend: IoBackend,
schema: &CollectionSchema,
docs: &DocumentMap,
source_revision: u64,
source_identity: &str,
ceilings: crate::storage_ceilings::StorageCeilings,
) -> Option<Self> {
cache::restore(
bytes,
diskann_file,
io_backend,
schema,
docs,
source_revision,
source_identity,
ceilings.max_index_cache_bytes(),
ceilings.max_diskann_file_bytes(),
)
}
pub(crate) fn cache_bytes(
&self,
schema: &CollectionSchema,
source_revision: u64,
source_identity: &str,
ceilings: crate::storage_ceilings::StorageCeilings,
) -> Result<Vec<u8>> {
cache::encode_with_limit(
self,
schema,
source_revision,
source_identity,
ceilings.max_index_cache_bytes(),
)
}
pub(crate) fn diskann_bytes(
&self,
schema: &CollectionSchema,
source_revision: u64,
source_identity: &str,
ceilings: crate::storage_ceilings::StorageCeilings,
) -> Result<Option<Vec<u8>>> {
diskann::encode_with_limit(
self,
schema,
source_revision,
source_identity,
ceilings.max_diskann_file_bytes(),
)
}
pub(crate) fn has_cacheable_indexes(&self) -> bool {
!self.indexes.is_empty() || !self.scalar_indexes.is_empty() || !self.fts_indexes.is_empty()
}
pub(crate) fn accounted_payload_bytes(&self) -> u64 {
self.indexes.values().fold(0_u64, |total, index| {
total.saturating_add(index.estimated_payload_bytes())
})
}
pub(crate) fn apply_document_changes(
&self,
schema: &CollectionSchema,
previous_docs: &DocumentMap,
docs: &DocumentMap,
source_revision: u64,
changed_ids: &BTreeSet<String>,
) -> Result<Self> {
let mut ordinals = self.ordinals.clone();
for id in changed_ids {
match (previous_docs.contains_key(id), docs.contains_key(id)) {
(_, true) => {
ordinals.ensure_live(id)?;
}
(true, false) => ordinals.remove_live(id),
(false, false) => {}
}
}
if ordinals.should_compact() {
return Self::build(schema, docs, source_revision);
}
let mut indexes = BTreeMap::new();
for field in &schema.vectors {
let Some(params) = field.index_params.as_ref() else {
continue;
};
if !is_incremental_vector(params.index_type, field.data_type) {
continue;
}
let Some(current) = self
.indexes
.get(&field.name)
.filter(|index| index.params == *params)
else {
indexes.insert(
field.name.clone(),
build_vector_index(
docs,
&field.name,
field.dimension,
params,
source_revision,
&ordinals,
)?,
);
continue;
};
let mut next = current.clone();
next.source_revision = source_revision;
for id in changed_ids {
let previous = previous_docs
.get(id)
.and_then(|doc| doc.vector(&field.name));
let current = docs.get(id).and_then(|doc| doc.vector(&field.name));
if previous == current {
continue;
}
let ordinal = ordinals.ordinal(id).ok_or_else(|| {
Error::internal(format!("vector ordinal is missing for document '{id}'"))
})?;
next.delta.remove(&ordinal);
next.delta_ordinals.remove(ordinal);
if next.base.vectors.contains_key(ordinal) {
next.tombstones.insert(ordinal);
} else {
next.tombstones.remove(ordinal);
}
if let Some(vector) = current {
next.delta
.insert(ordinal, encode_vector(id, &field.name, params, vector)?);
next.delta_ordinals.insert(ordinal);
}
}
if next.should_compact() {
next = build_vector_index(
docs,
&field.name,
field.dimension,
params,
source_revision,
&ordinals,
)?;
}
indexes.insert(field.name.clone(), next);
}
let scalar_indexes = self.scalar_indexes.apply_document_changes(
schema,
previous_docs,
docs,
source_revision,
changed_ids,
&ordinals,
)?;
let fts_indexes = self.fts_indexes.apply_document_changes(
schema,
previous_docs,
docs,
source_revision,
changed_ids,
&ordinals,
)?;
Ok(Self {
ordinals,
indexes,
scalar_indexes,
fts_indexes,
})
}
fn candidates(
&self,
docs: &DocumentMap,
revision: u64,
query: &SearchQuery,
allowed: Option<&OrdinalSet>,
) -> Result<Option<AnnCandidates>> {
let Some(index) = self.indexes.get(&query.field_name) else {
return Ok(None);
};
if index.source_revision != revision
|| query
.params
.get("is_linear")
.and_then(serde_json::Value::as_bool)
.unwrap_or(false)
{
return Ok(None);
}
if requested_metric(query).is_some_and(|metric| metric != index.params.metric_type) {
return Ok(None);
}
let Some(query_vector) = query_vector(docs, query)? else {
return Ok(None);
};
let topk = usize::try_from(query.topk)
.map_err(|_| Error::invalid_argument("query topk must be positive"))?;
let metric = index.params.metric_type;
let eligible_count = allowed.map_or_else(
|| index.live_vector_count(),
|allowed| index.eligible_vector_count(allowed.bitmap()),
);
if eligible_count == 0 {
return Ok(Some(AnnCandidates {
selection: CandidateSelection {
ids: OrdinalSet::new(&self.ordinals, RoaringTreemap::new()),
},
diskann_sector_reads: 0,
diskann_io_backend: None,
}));
}
let search = AnnSearchContext {
query,
vector: &query_vector,
topk,
metric,
allowed: allowed.map(OrdinalSet::bitmap),
eligible_count,
ordinals: &self.ordinals,
};
let result =
match &index.base.kind {
VectorIndexKind::Flat(_) => index.flat_candidates(&search).map(|ids| AnnOrdinals {
ids,
diskann_sector_reads: 0,
diskann_io_backend: None,
}),
VectorIndexKind::Hnsw(hnsw) => {
index.hnsw_candidates(hnsw, &search).map(|ids| AnnOrdinals {
ids,
diskann_sector_reads: 0,
diskann_io_backend: None,
})
}
VectorIndexKind::HnswRabitq(hnsw) => index
.hnsw_rabitq_candidates(hnsw, &search)
.map(|ids| AnnOrdinals {
ids,
diskann_sector_reads: 0,
diskann_io_backend: None,
}),
VectorIndexKind::Ivf(ivf) => {
index.ivf_candidates(ivf, &search).map(|ids| AnnOrdinals {
ids,
diskann_sector_reads: 0,
diskann_io_backend: None,
})
}
VectorIndexKind::IvfRabitq(ivf) => {
index
.ivf_rabitq_candidates(ivf, &search)
.map(|ids| AnnOrdinals {
ids,
diskann_sector_reads: 0,
diskann_io_backend: None,
})
}
VectorIndexKind::Diskann(diskann) => index.diskann_candidates(diskann, &search),
VectorIndexKind::Vamana(vamana) => index.vamana_candidates(vamana, &search),
};
Ok(result.map(|result| AnnCandidates {
selection: CandidateSelection {
ids: OrdinalSet::new(&self.ordinals, result.ids),
},
diskann_sector_reads: result.diskann_sector_reads,
diskann_io_backend: result.diskann_io_backend,
}))
}
pub(crate) fn scalar_candidates(
&self,
source_revision: u64,
filter: &crate::filter::FilterExpr,
) -> Option<ScalarCandidates> {
self.scalar_indexes
.candidates(source_revision, filter, &self.ordinals)
}
fn fts_scores(
&self,
docs: &DocumentMap,
source_revision: u64,
query: &SearchQuery,
candidates: Option<&OrdinalSet>,
topk: Option<usize>,
) -> Result<Option<OrdinalScores>> {
self.fts_indexes
.search(source_revision, query, candidates, docs, &self.ordinals)?
.map(|scores| match topk {
Some(topk) => OrdinalScores::new_topk(&self.ordinals, scores, topk),
None => OrdinalScores::new(&self.ordinals, scores),
})
.transpose()
}
fn plan_fts_candidates(
&self,
docs: &DocumentMap,
source_revision: u64,
query: &SearchQuery,
scalar: Option<ScalarCandidates>,
filter_exact: bool,
) -> Result<CandidatePlan> {
let used_scalar = scalar.is_some();
let topk = filter_exact
.then(|| {
usize::try_from(query.topk)
.ok()
.filter(|topk| *topk > 0)
.ok_or_else(|| Error::invalid_argument("query topk must be positive"))
})
.transpose()?;
let fts_scores = self.fts_scores(
docs,
source_revision,
query,
scalar.as_ref().map(|value| &value.ids),
topk,
)?;
if let Some(fts_scores) = fts_scores {
return Ok(CandidatePlan {
selection: None,
fts_scores: Some(fts_scores),
used_ann: false,
diskann_sector_reads: 0,
diskann_io_backend: None,
used_scalar,
used_fts_index: true,
});
}
Ok(CandidatePlan {
selection: scalar.map(|selection| CandidateSelection {
ids: selection.into_ids(),
}),
fts_scores: None,
used_ann: false,
diskann_sector_reads: 0,
diskann_io_backend: None,
used_scalar,
used_fts_index: false,
})
}
pub(crate) fn plan_candidates(
&self,
docs: &DocumentMap,
source_revision: u64,
query: &SearchQuery,
filter: Option<&crate::filter::FilterExpr>,
) -> Result<CandidatePlan> {
let scalar = filter.and_then(|filter| self.scalar_candidates(source_revision, filter));
if query.fts.is_some() {
let filter_exact =
filter.is_none() || scalar.as_ref().is_some_and(|candidates| candidates.exact);
return self.plan_fts_candidates(docs, source_revision, query, scalar, filter_exact);
}
let used_scalar = scalar.is_some();
let Some(mut scalar) = scalar else {
let ann = self.candidates(docs, source_revision, query, None)?;
let used_ann = ann.as_ref().is_some_and(|_| {
self.indexes
.get(&query.field_name)
.is_some_and(|index| is_approximate_ann(index.params.index_type))
});
let diskann_sector_reads = ann.as_ref().map_or(0, |ann| ann.diskann_sector_reads);
let diskann_io_backend = ann.as_ref().and_then(|ann| ann.diskann_io_backend);
return Ok(CandidatePlan {
selection: ann.map(|ann| ann.selection),
fts_scores: None,
used_ann,
diskann_sector_reads,
diskann_io_backend,
used_scalar,
used_fts_index: false,
});
};
let topk = usize::try_from(query.topk)
.map_err(|_| Error::invalid_argument("query topk must be positive"))?;
let exact_limit = topk
.saturating_mul(SCALAR_EXACT_PREFILTER_PER_RESULT)
.max(SCALAR_EXACT_PREFILTER_MIN);
if scalar.len() <= exact_limit {
return Ok(CandidatePlan {
selection: Some(CandidateSelection {
ids: scalar.into_ids(),
}),
fts_scores: None,
used_ann: false,
diskann_sector_reads: 0,
diskann_io_backend: None,
used_scalar,
used_fts_index: false,
});
}
if !scalar.exact {
let Some(filter) = filter else {
return Err(Error::internal(
"scalar candidate refinement requires a parsed filter",
));
};
scalar.retain_ids(|id| docs.get(id).is_some_and(|doc| filter.matches(doc)));
scalar.exact = true;
}
if scalar.len() <= exact_limit {
return Ok(CandidatePlan {
selection: Some(CandidateSelection {
ids: scalar.into_ids(),
}),
fts_scores: None,
used_ann: false,
diskann_sector_reads: 0,
diskann_io_backend: None,
used_scalar,
used_fts_index: false,
});
}
let Some(ann) = self.candidates(docs, source_revision, query, Some(&scalar.ids))? else {
return Ok(CandidatePlan {
selection: Some(CandidateSelection {
ids: scalar.into_ids(),
}),
fts_scores: None,
used_ann: false,
diskann_sector_reads: 0,
diskann_io_backend: None,
used_scalar,
used_fts_index: false,
});
};
Ok(CandidatePlan {
selection: Some(ann.selection),
fts_scores: None,
used_ann: self
.indexes
.get(&query.field_name)
.is_some_and(|index| is_approximate_ann(index.params.index_type)),
diskann_sector_reads: ann.diskann_sector_reads,
diskann_io_backend: ann.diskann_io_backend,
used_scalar,
used_fts_index: false,
})
}
pub(crate) fn stats(
&self,
schema: &CollectionSchema,
docs: &DocumentMap,
source_revision: u64,
) -> Vec<IndexStat> {
let mut stats = Vec::new();
for field in &schema.vectors {
let Some(params) = field.index_params.as_ref() else {
continue;
};
if builds_packed_vector_index(params.index_type, field.data_type) {
if let Some(index) = self.indexes.get(&field.name) {
stats.push(IndexStat {
name: field.name.clone(),
index_type: index.params.index_type,
completeness: 1.0,
source_revision: index.source_revision,
document_count: u64::try_from(index.live_vector_count())
.unwrap_or(u64::MAX),
estimated_payload_bytes: Some(index.estimated_payload_bytes()),
state: "ready".into(),
});
} else if params.index_type == IndexType::Flat {
let document_count = u64::try_from(
docs.values()
.filter(|doc| doc.vector(&field.name).is_some())
.count(),
)
.unwrap_or(u64::MAX);
stats.push(IndexStat {
name: field.name.clone(),
index_type: IndexType::Flat,
completeness: 1.0,
source_revision,
document_count,
estimated_payload_bytes: None,
state: "ready".into(),
});
} else {
stats.push(IndexStat {
name: field.name.clone(),
index_type: params.index_type,
completeness: 0.0,
source_revision: 0,
document_count: 0,
estimated_payload_bytes: None,
state: "missing".into(),
});
}
continue;
}
if params.index_type == IndexType::Flat {
let document_count = u64::try_from(
docs.values()
.filter(|doc| doc.vector(&field.name).is_some())
.count(),
)
.unwrap_or(u64::MAX);
stats.push(IndexStat {
name: field.name.clone(),
index_type: IndexType::Flat,
completeness: 1.0,
source_revision,
document_count,
estimated_payload_bytes: None,
state: "ready".into(),
});
}
}
stats.extend(self.scalar_indexes.stats());
stats.extend(self.fts_indexes.stats());
stats
}
}
impl VectorIndex {
fn estimated_payload_bytes(&self) -> u64 {
let base_vectors = self
.base
.vectors
.values()
.fold(self.base.vectors.slot_count(), |total, vector| {
total.saturating_add(vector.encoded_bytes())
});
let vectors = self.delta.values().fold(base_vectors, |total, vector| {
total
.saturating_add(std::mem::size_of::<u64>())
.saturating_add(vector.encoded_bytes())
});
let membership = self
.base
.vector_ordinals
.serialized_size()
.saturating_add(self.delta_ordinals.serialized_size())
.saturating_add(self.tombstones.serialized_size());
let kind = match &self.base.kind {
VectorIndexKind::Flat(_) => FlatIndex::estimated_payload_bytes(),
VectorIndexKind::Hnsw(index) => index.estimated_payload_bytes(),
VectorIndexKind::HnswRabitq(index) => index.estimated_payload_bytes(),
VectorIndexKind::Ivf(index) => index.estimated_payload_bytes(),
VectorIndexKind::IvfRabitq(index) => index.estimated_payload_bytes(),
VectorIndexKind::Diskann(index) => index.estimated_payload_bytes(),
VectorIndexKind::Vamana(index) => index.estimated_payload_bytes(),
};
u64::try_from(vectors.saturating_add(membership).saturating_add(kind)).unwrap_or(u64::MAX)
}
fn live_vector_count(&self) -> usize {
let hidden_base =
bitmap_count_to_usize(self.tombstones.intersection_len(&self.base.vector_ordinals));
self.base
.vectors
.len()
.saturating_sub(hidden_base)
.saturating_add(self.delta.len())
}
fn eligible_vector_count(&self, allowed: &RoaringTreemap) -> usize {
let base = self
.base
.vector_ordinals
.intersection_len(allowed)
.saturating_sub(self.tombstones.intersection_len(allowed));
let delta = self.delta_ordinals.intersection_len(allowed);
bitmap_count_to_usize(base.saturating_add(delta))
}
fn base_eligible_vector_count(&self, allowed: &RoaringTreemap) -> usize {
bitmap_count_to_usize(
self.base
.vector_ordinals
.intersection_len(allowed)
.saturating_sub(self.tombstones.intersection_len(allowed)),
)
}
fn overlay_len(&self) -> usize {
let new_vectors = self
.delta
.keys()
.filter(|&&ordinal| !self.base.vectors.contains_key(ordinal))
.count();
bitmap_count_to_usize(self.tombstones.len()).saturating_add(new_vectors)
}
fn should_compact(&self) -> bool {
self.overlay_len() >= delta_compaction_limit(self.base.vectors.len())
}
fn merge_candidates(
&self,
mut base: RoaringTreemap,
query: &[f32],
limit: Option<usize>,
metric: MetricType,
allowed: Option<&RoaringTreemap>,
ordinals: &OrdinalTable,
) -> RoaringTreemap {
base -= &self.tombstones;
if let Some(allowed) = allowed {
base &= allowed;
}
let mut ids = base;
let mut delta = self.delta_ordinals.clone();
if let Some(allowed) = allowed {
delta &= allowed;
}
ids |= delta;
match limit {
Some(limit) if bitmap_count_to_usize(ids.len()) > limit => {
self.limit_candidates(&ids, query, limit, metric, ordinals)
}
Some(_) | None => ids,
}
}
fn limit_candidates(
&self,
ids: &RoaringTreemap,
query: &[f32],
limit: usize,
metric: MetricType,
ordinals: &OrdinalTable,
) -> RoaringTreemap {
let query_norm = if metric == MetricType::Cosine {
dense_query_norm(query)
} else {
0.0
};
let mut scored: Vec<(u64, f64)> = ids
.iter()
.filter_map(|ordinal| {
let vector = self
.delta
.get(&ordinal)
.or_else(|| self.base.vectors.get(ordinal))?;
Some((
ordinal,
score_with_query_norm(query, vector, metric, query_norm),
))
})
.collect();
scored.sort_by(|left, right| {
right.1.total_cmp(&left.1).then_with(|| {
ordinals
.id(left.0)
.unwrap_or_default()
.cmp(ordinals.id(right.0).unwrap_or_default())
.then_with(|| left.0.cmp(&right.0))
})
});
scored
.into_iter()
.take(limit)
.map(|(ordinal, _)| ordinal)
.collect()
}
}
fn candidate_set_is_sufficient(ids: &RoaringTreemap, search: &AnnSearchContext<'_>) -> bool {
search.allowed.is_none()
|| bitmap_count_to_usize(ids.len()) >= search.topk.min(search.eligible_count)
}
fn bitmap_count_to_usize(count: u64) -> usize {
usize::try_from(count).unwrap_or(usize::MAX)
}
fn proportional_candidate_limit(target: usize, population: usize, eligible: usize) -> usize {
if target == 0 || population == 0 || eligible == 0 {
return 0;
}
let target = u128::try_from(target).unwrap_or(u128::MAX);
let population_u128 = u128::try_from(population).unwrap_or(u128::MAX);
let eligible_u128 = u128::try_from(eligible).unwrap_or(u128::MAX);
let scaled = target
.saturating_mul(population_u128)
.saturating_add(eligible_u128.saturating_sub(1))
/ eligible_u128;
usize::try_from(scaled)
.unwrap_or(population)
.min(population)
}
fn delta_compaction_limit(base_len: usize) -> usize {
let fractional =
base_len.saturating_add(DELTA_COMPACTION_DIVISOR - 1) / DELTA_COMPACTION_DIVISOR;
fractional.clamp(MIN_DELTA_COMPACTION, MAX_DELTA_COMPACTION)
}
pub(super) fn is_in_memory_vector(index_type: IndexType) -> bool {
matches!(
index_type,
IndexType::Flat
| IndexType::Hnsw
| IndexType::HnswRabitq
| IndexType::Ivf
| IndexType::IvfRabitq
| IndexType::Diskann
| IndexType::Vamana
)
}
pub(super) fn builds_packed_vector_index(index_type: IndexType, data_type: DataType) -> bool {
is_in_memory_vector(index_type) && supports_f32_ann_kernel(data_type)
}
fn supports_f32_ann_kernel(data_type: DataType) -> bool {
matches!(
data_type,
DataType::VectorFp16
| DataType::VectorFp32
| DataType::VectorFp64
| DataType::VectorInt4
| DataType::VectorInt8
| DataType::VectorInt16
)
}
fn is_incremental_vector(index_type: IndexType, data_type: DataType) -> bool {
builds_packed_vector_index(index_type, data_type) && index_type != IndexType::Flat
}
fn is_approximate_ann(index_type: IndexType) -> bool {
is_in_memory_vector(index_type) && index_type != IndexType::Flat
}
fn query_vector(docs: &DocumentMap, query: &SearchQuery) -> Result<Option<Vec<f32>>> {
if let Some(vector) = &query.vector {
return Ok(Some(vector.clone()));
}
let Some(id) = query.id.as_deref() else {
return Ok(None);
};
let Some(vector) = docs.get(id).and_then(|doc| doc.vector(&query.field_name)) else {
return Ok(None);
};
vector.to_dense_f32().map(Some).ok_or_else(|| {
Error::resource_exhausted(format!(
"source document '{id}' cannot be represented by the f32 ANN kernel"
))
})
}
fn optional_positive_query_parameter(query: &SearchQuery, name: &str) -> Option<usize> {
query
.params
.get(name)
.and_then(serde_json::Value::as_u64)
.and_then(|value| (value > 0).then_some(value))
.and_then(|value| usize::try_from(value).ok())
}
#[allow(clippy::cast_possible_truncation)]
fn optional_f32_query_parameter(query: &SearchQuery, name: &str) -> Option<f32> {
query
.params
.get(name)
.and_then(serde_json::Value::as_f64)
.filter(|value| value.is_finite() && *value > 0.0 && *value <= f64::from(f32::MAX))
.map(|value| value as f32)
}
fn requested_metric(query: &SearchQuery) -> Option<MetricType> {
match query
.params
.get("metric")?
.as_str()?
.to_ascii_lowercase()
.as_str()
{
"l2" | "euclidean" => Some(MetricType::L2),
"ip" | "inner_product" | "dot" => Some(MetricType::Ip),
"cosine" => Some(MetricType::Cosine),
"mips_l2" | "mips-l2" => Some(MetricType::MipsL2),
_ => None,
}
}
#[cfg(test)]
mod tests;