use std::{
collections::BTreeSet,
sync::atomic::{AtomicBool, AtomicUsize, Ordering},
};
use gxhash::HashMap;
use parking_lot::Mutex;
use super::{
cleanup::vector_set_cleanup_work_channel::VectorSetCleanupWorkChannel,
disk_ann_service::{DiskANNService, DiskAnnInsertResult},
vector_manager__cleanup::{CleanupGate, CleanupRuntime},
vector_manager__context_metadata::ContextMetadata,
vector_manager__index::{INDEX_SIZE, Index},
vector_manager__locking::VectorSetLocks,
vector_manager__quantization::{QuantizationChannel, QuantizationState, QuantizationStep},
vector_manager__replication::ReplicationRuntime,
vector_manager_element_data::prepare_vector_data,
vector_types::{
VectorDistanceMetricType, VectorIdFormat, VectorQuantType, VectorSetFlags, VectorValueType,
},
};
pub const CONTEXT_STEP: u64 = 8;
pub const METADATA_NAMESPACE: u8 = 1;
pub const INDEX_SIZE_BYTES: usize = INDEX_SIZE;
pub const VADD_APPEND_LOG_ARG: i64 = i64::MIN; pub const RECREATE_INDEX_ARG: i64 = VADD_APPEND_LOG_ARG + 2;
pub const VREM_APPEND_LOG_ARG: i64 = RECREATE_INDEX_ARG + 1;
pub const MIGRATE_ELEMENT_KEY_LOG_ARG: i64 = VREM_APPEND_LOG_ARG + 1;
pub const MIGRATE_INDEX_KEY_LOG_ARG: i64 = MIGRATE_ELEMENT_KEY_LOG_ARG + 1;
pub const VADD_SET_FLAGS_ARG: i64 = MIGRATE_INDEX_KEY_LOG_ARG + 1;
pub const CREATE_INDEX_ARG: i64 = VADD_SET_FLAGS_ARG + 1;
pub const VSETATTR_APPEND_LOG_ARG: i64 = CREATE_INDEX_ARG + 1;
pub const RECORD_TYPE: u8 = 1;
pub const MINIMUM_SPACE_PER_ID: usize = 4 + 8;
pub const MAX_VECTOR_DIMENSIONS: u32 = 1 << 16;
pub const MAX_RETRIEVE_COUNT: usize = 100_000_000;
pub const MAX_FILTERING_SCALE_FACTOR: usize = 256;
pub const MAX_EXPLORATION_FACTOR: usize = 1_000_000;
pub const CONTEXT_METADATA_SIZE: usize = 4 * 8 + 64 * 2;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum VectorManagerResult {
#[default]
Invalid = 0,
OK,
BadParams,
Duplicate,
MissingElement,
}
#[derive(Debug, Clone, PartialEq)]
pub struct VectorOpError {
pub result: VectorManagerResult,
pub message: Vec<u8>,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct SimilarityOutput {
pub output_ids: Vec<u8>,
pub output_distances: Vec<f32>,
pub output_attributes: Vec<u8>,
pub filter_bitmap: Vec<u8>,
pub id_format: VectorIdFormat,
pub found: usize,
}
#[derive(Debug, Clone, Default)]
pub struct VectorManagerOptions {
pub is_enabled: bool,
pub quantization_task_count: usize,
}
pub struct VectorManager {
pub is_enabled: bool,
pub db_id: i32,
pub service: DiskANNService,
pub context_metadatas: Mutex<Vec<ContextMetadata>>,
pub dirty_context_metadatas: Mutex<BTreeSet<usize>>,
pub metadata_store: Mutex<HashMap<i32, [u8; CONTEXT_METADATA_SIZE]>>,
pub recovered_indexes: Mutex<HashMap<u64, u8>>,
pub recovered_metadata: Mutex<HashMap<i32, ContextMetadata>>,
pub requested_drops: Mutex<HashMap<Vec<u8>, u64>>,
pub potentially_deleted: Mutex<HashMap<Vec<u8>, u64>>,
pub(crate) key_index_registry: Mutex<HashMap<Vec<u8>, [u8; INDEX_SIZE_BYTES]>>,
pub vector_set_locks: VectorSetLocks,
pub cleanup_task_channel: VectorSetCleanupWorkChannel<u64>,
pub request_cleanup_task_channel: VectorSetCleanupWorkChannel<u64>,
pub request_drop_task_channel: VectorSetCleanupWorkChannel<()>,
pub quantization_channel: QuantizationChannel,
pub quantization_task_count: usize,
pub quantization_requests_processed: AtomicUsize,
pub quantization_backfills_processed: AtomicUsize,
pub cleanup_gate: CleanupGate,
pub cleanup_runtime: CleanupRuntime,
pub replication: ReplicationRuntime,
storage_session_attached: AtomicBool,
}
impl VectorManager {
pub fn new(options: VectorManagerOptions) -> Self {
let quantization_task_count = match options.quantization_task_count {
0 => 4, n => n.min(1024),
};
Self {
is_enabled: options.is_enabled,
db_id: 0,
service: DiskANNService::default(),
context_metadatas: Mutex::new(vec![Default::default()]),
dirty_context_metadatas: Mutex::new(BTreeSet::new()),
metadata_store: Mutex::new(HashMap::default()),
recovered_indexes: Mutex::new(HashMap::default()),
recovered_metadata: Mutex::new(HashMap::default()),
requested_drops: Mutex::new(HashMap::default()),
potentially_deleted: Mutex::new(HashMap::default()),
key_index_registry: Mutex::new(HashMap::default()),
vector_set_locks: VectorSetLocks::default(),
cleanup_task_channel: VectorSetCleanupWorkChannel::new(),
request_cleanup_task_channel: VectorSetCleanupWorkChannel::new(),
request_drop_task_channel: VectorSetCleanupWorkChannel::new(),
quantization_channel: QuantizationChannel::new(),
quantization_task_count,
quantization_requests_processed: AtomicUsize::new(0),
quantization_backfills_processed: AtomicUsize::new(0),
cleanup_gate: CleanupGate::new(),
cleanup_runtime: CleanupRuntime::new(),
replication: ReplicationRuntime::new(),
storage_session_attached: AtomicBool::new(true),
}
}
pub fn with_db_id(mut self, db_id: i32) -> Self {
self.db_id = db_id;
self
}
pub fn assert_have_storage_session(&self) {
debug_assert!(
self.storage_session_attached.load(Ordering::Relaxed),
"进入此路径前必须已挂载存储会话"
);
}
pub fn error_msg(result: VectorManagerResult) -> &'static [u8] {
match result {
VectorManagerResult::Invalid => b"ERR Invalid vector set operation",
VectorManagerResult::OK => b"",
VectorManagerResult::BadParams => b"ERR Invalid parameters for vector set operation",
VectorManagerResult::Duplicate => b"ERR Vector set element already exists",
VectorManagerResult::MissingElement => b"ERR Vector set element does not exist",
}
}
pub fn ensure_distance_buffer_size(buffer: &mut Vec<u8>, retrieve_count: usize) {
let size_bytes = retrieve_count * 4;
debug_assert!(retrieve_count <= MAX_RETRIEVE_COUNT, "结果数超出上限");
if buffer.len() < size_bytes {
buffer.resize(size_bytes, 0);
}
}
pub fn ensure_id_buffer_size(buffer: &mut Vec<u8>, retrieve_count: usize) {
let size_bytes = retrieve_count * MINIMUM_SPACE_PER_ID;
debug_assert!(retrieve_count <= MAX_RETRIEVE_COUNT, "结果数超出上限");
if buffer.len() < size_bytes {
buffer.resize(size_bytes, 0);
}
}
pub fn ensure_filter_bitmap_size(buffer: &mut Vec<u8>, result_count: usize) {
let size_bytes = (result_count + 7) >> 3;
if buffer.len() < size_bytes {
buffer.resize(size_bytes, 0);
}
}
#[allow(clippy::too_many_arguments)]
pub fn try_add(
&self,
key: &[u8],
index_value: &[u8],
element: &[u8],
value_type: VectorValueType,
values: &[u8],
attributes: &[u8],
provided_reduce_dims: u32,
provided_quant_type: VectorQuantType,
provided_num_links: u32,
provided_distance_metric: VectorDistanceMetricType,
) -> Result<VectorManagerResult, VectorOpError> {
self.assert_have_storage_session();
let err = |result: VectorManagerResult, message: &[u8]| {
Err(VectorOpError {
result,
message: message.to_vec(),
})
};
let Some(index) = Index::from_bytes(index_value) else {
return err(
VectorManagerResult::BadParams,
b"ERR Invalid vector set index",
);
};
if provided_reduce_dims != 0 && provided_reduce_dims != index.reduce_dims {
return err(
VectorManagerResult::BadParams,
b"ERR Provided REDUCE does not match Vector Set definition",
);
}
if provided_quant_type != VectorQuantType::Invalid && provided_quant_type != index.quant_type {
return err(
VectorManagerResult::BadParams,
b"ERR asked quantization mismatch with existing vector set",
);
}
if provided_distance_metric != index.distance_metric {
return err(
VectorManagerResult::BadParams,
format!(
"ERR Distance metric mismatch - got {} but set has {}",
provided_distance_metric.csharp_name(),
index.distance_metric.csharp_name()
)
.as_bytes(),
);
}
if provided_num_links != index.num_links {
return err(
VectorManagerResult::BadParams,
b"ERR asked M value mismatch with existing vector set",
);
}
let prepared = match prepare_vector_data(index.quant_type, value_type, values) {
Ok(p) => p,
Err(e) => return err(VectorManagerResult::BadParams, e.message()),
};
if prepared.element_count != index.dimensions as usize {
return err(
VectorManagerResult::BadParams,
format!(
"ERR Vector dimension mismatch - got {} but set has {}",
prepared.element_count, index.dimensions
)
.as_bytes(),
);
}
if provided_reduce_dims == 0 && index.reduce_dims != 0 {
return err(
VectorManagerResult::BadParams,
format!(
"ERR Vector dimension mismatch - got {} but set has {}",
prepared.element_count, index.reduce_dims
)
.as_bytes(),
);
}
match self
.service
.insert(index.context, element, &prepared.bytes, attributes)
{
DiskAnnInsertResult::True => Ok(VectorManagerResult::OK),
DiskAnnInsertResult::QuantizationRequested => {
let _ = self
.quantization_channel
.try_publish(QuantizationState::new(
key.to_vec(),
QuantizationStep::BuildQuantizationTable,
0,
));
Ok(VectorManagerResult::OK)
}
DiskAnnInsertResult::False => Ok(VectorManagerResult::Duplicate),
}
}
pub fn try_remove(&self, index_value: &[u8], element: &[u8]) -> VectorManagerResult {
self.assert_have_storage_session();
let Some(index) = Index::from_bytes(index_value) else {
return VectorManagerResult::Invalid;
};
if self.service.remove(index.context, element) {
VectorManagerResult::OK
} else {
VectorManagerResult::MissingElement
}
}
pub fn try_set_attribute(&self, index_value: &[u8], element: &[u8], attribute: &[u8]) -> bool {
self.assert_have_storage_session();
let Some(index) = Index::from_bytes(index_value) else {
return false;
};
self
.service
.set_attribute(index.context, element, attribute)
}
pub fn request_deletion(&self, value: &[u8]) {
if value.len() != INDEX_SIZE_BYTES {
log::warn!("Ignored Vector Set deletion due to size mismatch");
return;
}
let Some(index) = Index::from_bytes(value) else {
return;
};
if index.flags.contains(VectorSetFlags::SUPPRESS_CLEANUP) {
return;
}
if !self.request_cleanup_task_channel.try_publish(index.context) {
log::error!("Could not submit request for Vector Set cleanup, aborting delete");
return;
}
self.drop_index(value);
}
pub fn request_drop_in_memory_index(&self, key: &[u8], value: &[u8]) {
if value.len() != INDEX_SIZE_BYTES {
log::warn!("Ignored Vector Set drop index due to size mismatch");
return;
}
let Some(index) = Index::from_bytes(value) else {
return;
};
if index.flags.contains(VectorSetFlags::SUPPRESS_CLEANUP) {
return;
}
if index.index_ptr != 0 {
if self
.requested_drops
.lock()
.insert(key.to_vec(), index.context)
.is_some()
{
log::error!("Drop triggered multiple times for same index");
return;
}
let _ = self.request_drop_task_channel.try_publish(());
}
}
pub fn drop_in_memory_index(&self, value: &[u8]) {
if value.len() != INDEX_SIZE_BYTES {
log::warn!("Ignored Vector Set drop index due to size mismatch");
return;
}
self.drop_index(value);
}
fn drop_index(&self, value: &[u8]) {
let Some(index) = Index::from_bytes(value) else {
return;
};
if index.index_ptr == 0 {
return;
}
self.service.drop_index(index.context);
}
pub fn clear_index_pointer(value: &mut [u8]) {
if value.len() != INDEX_SIZE_BYTES {
return;
}
value[8..16].fill(0);
}
#[allow(clippy::too_many_arguments)]
pub fn value_similarity(
&self,
index_value: &[u8],
value_type: VectorValueType,
values: &[u8],
count: usize,
search_exploration_factor: usize,
filter: &[u8],
max_filtering_effort: usize,
delta: f32,
include_attributes: bool,
) -> Result<SimilarityOutput, VectorOpError> {
self.assert_have_storage_session();
let err = |result: VectorManagerResult, message: &[u8]| {
Err(VectorOpError {
result,
message: message.to_vec(),
})
};
let Some(index) = Index::from_bytes(index_value) else {
return err(
VectorManagerResult::BadParams,
b"ERR Invalid vector set index",
);
};
let effective_ef = effective_search_ef(
search_exploration_factor,
count,
filter,
max_filtering_effort,
);
let prepared = match prepare_vector_data(index.quant_type, value_type, values) {
Ok(p) => p,
Err(e) => return err(VectorManagerResult::BadParams, e.message()),
};
if prepared.element_count != index.dimensions as usize {
return err(
VectorManagerResult::BadParams,
b"ERR Dimensions provided do not match Vector Set dimensions",
);
}
let mut filter_state = match compile_filter_state(filter) {
Ok(state) => state,
Err(_) => {
return err(
VectorManagerResult::BadParams,
b"ERR Compiling filter failed",
);
}
};
let mut predicate = |external_id: &[u8]| -> bool {
match &mut filter_state {
None => true,
Some((program, state)) => {
self.evaluate_candidate_filter(index.context, external_id, program, filter, state)
}
}
};
let mut hits = self
.service
.search_vector(
index.context,
&prepared.bytes,
count,
effective_ef,
&mut predicate,
)
.map_err(|_| VectorOpError {
result: VectorManagerResult::BadParams,
message: b"ERR Error indicating response from vector service".to_vec(),
})?;
apply_delta_cutoff(&mut hits, delta);
self.build_similarity_output(index.context, hits, filter, include_attributes)
}
#[allow(clippy::too_many_arguments)]
pub fn element_similarity(
&self,
index_value: &[u8],
element: &[u8],
count: usize,
search_exploration_factor: usize,
filter: &[u8],
max_filtering_effort: usize,
delta: f32,
include_attributes: bool,
) -> Result<SimilarityOutput, VectorOpError> {
self.assert_have_storage_session();
let Some(index) = Index::from_bytes(index_value) else {
return Err(VectorOpError {
result: VectorManagerResult::BadParams,
message: b"ERR Invalid vector set index".to_vec(),
});
};
let effective_ef = effective_search_ef(
search_exploration_factor,
count,
filter,
max_filtering_effort,
);
if !self.service.check_external_id_valid(index.context, element) {
return Err(VectorOpError {
result: VectorManagerResult::MissingElement,
message: super::resp_server_session_vectors::ERR_ELEMENT_NOT_IN_SET.to_vec(),
});
}
let mut filter_state = compile_filter_state(filter).map_err(|_| VectorOpError {
result: VectorManagerResult::BadParams,
message: b"ERR Compiling filter failed".to_vec(),
})?;
let mut predicate = |external_id: &[u8]| -> bool {
match &mut filter_state {
None => true,
Some((program, state)) => {
self.evaluate_candidate_filter(index.context, external_id, program, filter, state)
}
}
};
let mut hits = self
.service
.search_element(index.context, element, count, effective_ef, &mut predicate)
.map_err(|_| VectorOpError {
result: VectorManagerResult::BadParams,
message: b"ERR Error indicating response from vector service".to_vec(),
})?;
apply_delta_cutoff(&mut hits, delta);
self.build_similarity_output(index.context, hits, filter, include_attributes)
}
fn build_similarity_output(
&self,
context: u64,
hits: Vec<super::disk_ann_service::SearchHit>,
filter: &[u8],
include_attributes: bool,
) -> Result<SimilarityOutput, VectorOpError> {
let found = hits.len();
let mut output = SimilarityOutput {
found,
id_format: VectorIdFormat::I32LengthPrefixed,
..Default::default()
};
for hit in &hits {
output
.output_ids
.extend_from_slice(&(hit.external_id.len() as i32).to_le_bytes());
output.output_ids.extend_from_slice(&hit.external_id);
output.output_distances.push(hit.distance);
}
if include_attributes || !filter.is_empty() {
for hit in &hits {
let attr = self
.service
.get_attribute(context, &hit.external_id)
.unwrap_or_default();
output
.output_attributes
.extend_from_slice(&(attr.len() as i32).to_le_bytes());
output.output_attributes.extend_from_slice(&attr);
}
}
if !filter.is_empty() {
Self::ensure_filter_bitmap_size(&mut output.filter_bitmap, found);
let view = AttributeView {
raw: &output.output_attributes,
};
let _ = super::vector_manager__filter::apply_post_filter(
filter,
found,
&view,
&mut output.filter_bitmap,
);
}
Ok(output)
}
pub fn fetch_single_vector_element_attributes(
&self,
index_value: &[u8],
element: &[u8],
) -> VectorManagerResult {
self.assert_have_storage_session();
let Some(index) = Index::from_bytes(index_value) else {
return VectorManagerResult::Invalid;
};
if self.service.get_attribute(index.context, element).is_some() {
VectorManagerResult::OK
} else {
VectorManagerResult::MissingElement
}
}
pub fn fetch_vector_element_attributes(&self, context: u64, ids: &[u8]) -> Vec<u8> {
let mut out = Vec::new();
for element in unpack_length_prefixed(ids) {
let attr = self
.service
.get_attribute(context, element)
.unwrap_or_default();
out.extend_from_slice(&(attr.len() as i32).to_le_bytes());
out.extend_from_slice(&attr);
}
out
}
pub fn try_get_embedding(&self, index_value: &[u8], element: &[u8]) -> Option<Vec<f32>> {
self.assert_have_storage_session();
let index = Index::from_bytes(index_value)?;
let embedding = self.service.embedding_of(index.context, element)?;
let internal = self.service.internal_id_of(index.context, element)?;
if !self
.service
.check_internal_id_valid(index.context, internal)
{
return None;
}
Some(embedding)
}
pub fn try_get_raw_embedding(
&self,
index_value: &[u8],
element: &[u8],
) -> Option<(Vec<u8>, VectorQuantType, f64, Option<f64>)> {
self.assert_have_storage_session();
let index = Index::from_bytes(index_value)?;
let quant = self.service.quant_of(index.context)?;
let bytes = self.service.get_full_vector(index.context, element)?;
let norm = 1.0;
let range = (quant == VectorQuantType::Q8).then_some(1.0);
Some((bytes, quant, norm, range))
}
pub fn is_member(&self, index_value: &[u8], element: &[u8]) -> bool {
let Some(index) = Index::from_bytes(index_value) else {
return false;
};
self.service.check_external_id_valid(index.context, element)
}
pub fn reconcile_recovered_state(&self, require_no_reserved_contexts: bool) -> bool {
if !self.is_enabled {
return true;
}
let mut needs_updated = false;
let mut metas = self.context_metadatas.lock();
if require_no_reserved_contexts {
for (i, meta) in metas.iter().enumerate() {
if !meta.is_empty() {
log::error!(
"Vector Set context reservation was not empty at index {i} when rebuilding after a full store replacement; expected the preceding flush to have cleared it"
);
return false;
}
}
}
{
let recovered = self.recovered_metadata.lock();
if !recovered.is_empty() {
let max_context = recovered.keys().copied().max().unwrap_or(0);
*metas = vec![Default::default(); (max_context + 1) as usize];
for i in 0..metas.len() {
if let Some(meta) = recovered.get(&(i as i32)) {
metas[i] = *meta;
}
}
}
}
for i in 0..metas.len() {
if let Some(abandoned) = metas[i].get_migrating() {
for ctx in abandoned {
metas[i].mark_migration_complete(i != 0, ctx, u16::MAX);
metas[i].mark_cleaning_up(i != 0, ctx);
}
self.dirty_context_metadatas.lock().insert(i);
needs_updated = true;
}
}
{
let recovered = self.recovered_indexes.lock();
for &context in recovered.keys() {
let (context_index, context_value) = Self::decompose_context(context);
if let Some(meta) = metas.get_mut(context_index) {
let allow_zero = context_index != 0;
if meta.is_cleaning_up(allow_zero, context_value) {
meta.clear_is_cleaning_up(allow_zero, context_value);
self.dirty_context_metadatas.lock().insert(context_index);
needs_updated = true;
}
}
}
}
for i in 0..metas.len() {
let offset = Self::offset_for_context_metadata(i);
for j in 0..64u64 {
let context = offset + j * CONTEXT_STEP;
if context == 0 {
continue;
}
if self.recovered_indexes.lock().contains_key(&context) {
continue;
}
let (_, context_value) = Self::decompose_context(context);
let allow_zero = i != 0;
if metas[i].is_in_use(allow_zero, context_value)
&& !metas[i].is_cleaning_up(allow_zero, context_value)
{
metas[i].mark_cleaning_up(allow_zero, context_value);
self.dirty_context_metadatas.lock().insert(i);
}
}
}
self.recovered_indexes.lock().clear();
if needs_updated {
self.update_context_metadata();
}
let _ = self.cleanup_task_channel.try_publish(0);
true
}
pub fn recovered_vector_set_index_key(&self, value: &[u8]) {
if value.len() != INDEX_SIZE_BYTES {
return;
}
let Some(index) = Index::from_bytes(value) else {
return;
};
self.recovered_indexes.lock().insert(index.context, 0);
}
pub fn recovered_context_metadata(&self, key: &[u8], value: &[u8]) -> bool {
if value.len() != CONTEXT_METADATA_SIZE || key.len() != 4 {
return true;
}
let record_index = i32::from_le_bytes(key.try_into().unwrap_or([0; 4]));
let Some(metadata) = ContextMetadata::from_bytes(value) else {
return true;
};
if metadata.is_empty() {
return true;
}
let mut recovered = self.recovered_metadata.lock();
if recovered.insert(record_index, metadata).is_some() {
log::error!("Recovered multiple instances of the same ContextMetadata: {record_index}");
return false;
}
true
}
pub fn sanitize_and_track_ingested_record_if_applicable(
&self,
tombstone: bool,
namespace_bytes: Option<&[u8]>,
record_type: u8,
key: &[u8],
value: &mut [u8],
) -> bool {
if tombstone {
return true;
}
if let Some(ns) = namespace_bytes {
if self.is_enabled && ns.len() == 1 && ns[0] == METADATA_NAMESPACE {
return self.recovered_context_metadata(key, value);
}
return true;
}
if record_type == RECORD_TYPE {
Self::clear_index_pointer(value);
if self.is_enabled {
self.recovered_vector_set_index_key(value);
}
}
true
}
}
fn effective_search_ef(
ef: usize,
count: usize,
filter: &[u8],
max_filtering_effort: usize,
) -> usize {
let base = ef.max(count);
if filter.is_empty() {
base
} else {
base.max(count.saturating_mul(max_filtering_effort.max(1)))
}
}
fn compile_filter_state(
filter: &[u8],
) -> Result<
Option<(
super::vector_filter_expression::ExprProgram,
super::vector_manager__filter::CandidateFilterState,
)>,
(),
> {
if filter.is_empty() {
return Ok(None);
}
match super::expr_compiler::try_compile(filter) {
Ok(program) => {
let state = super::vector_manager__filter::CandidateFilterState::new(&program, filter);
Ok(Some((program, state)))
}
Err(_) => Err(()),
}
}
fn apply_delta_cutoff(hits: &mut Vec<super::disk_ann_service::SearchHit>, delta: f32) {
if delta.is_finite() {
hits.retain(|h| h.distance <= delta);
}
}
pub fn unpack_length_prefixed(bytes: &[u8]) -> Vec<&[u8]> {
let mut out = Vec::new();
let mut rest = bytes;
while rest.len() >= 4 {
let len = i32::from_le_bytes(rest[..4].try_into().unwrap_or([0; 4])).max(0) as usize;
let total = 4 + len;
if rest.len() < total {
break;
}
out.push(&rest[4..total]);
rest = &rest[total..];
}
out
}
pub struct AttributeView<'a> {
pub raw: &'a [u8],
}
impl AttributeView<'_> {
pub fn segments(&self) -> Vec<&[u8]> {
unpack_length_prefixed(self.raw)
}
}
#[cfg(test)]
mod tests {
use super::{
super::{
vector_manager__index::Index,
vector_types::{VectorDistanceMetricType, VectorQuantType},
},
*,
};
fn manager() -> VectorManager {
VectorManager::new(VectorManagerOptions {
is_enabled: true,
..Default::default()
})
}
fn fresh_index(context: u64, dims: u32) -> Index {
Index {
context,
index_ptr: 1,
dimensions: dims,
reduce_dims: 0,
num_links: 8,
build_exploration_factor: 64,
quant_type: VectorQuantType::NoQuant,
distance_metric: VectorDistanceMetricType::L2,
flags: VectorSetFlags::NONE,
}
}
fn f32_bytes(vals: &[f32]) -> Vec<u8> {
vals.iter().flat_map(|v| v.to_le_bytes()).collect()
}
#[test]
fn buffer_size_helpers() {
let mut buf = Vec::new();
VectorManager::ensure_distance_buffer_size(&mut buf, 10);
assert_eq!(buf.len(), 40);
VectorManager::ensure_distance_buffer_size(&mut buf, 5);
assert_eq!(buf.len(), 40);
let mut ids = Vec::new();
VectorManager::ensure_id_buffer_size(&mut ids, 4);
assert_eq!(ids.len(), 4 * MINIMUM_SPACE_PER_ID);
let mut bits = Vec::new();
VectorManager::ensure_filter_bitmap_size(&mut bits, 65);
assert_eq!(bits.len(), 9);
VectorManager::ensure_filter_bitmap_size(&mut bits, 8);
assert_eq!(bits.len(), 9);
}
#[test]
fn log_args_are_contiguous() {
assert_eq!(VADD_APPEND_LOG_ARG, i64::MIN);
assert_eq!(RECREATE_INDEX_ARG, VADD_APPEND_LOG_ARG + 2);
assert_eq!(VREM_APPEND_LOG_ARG, RECREATE_INDEX_ARG + 1);
assert_eq!(MIGRATE_ELEMENT_KEY_LOG_ARG, VREM_APPEND_LOG_ARG + 1);
assert_eq!(MIGRATE_INDEX_KEY_LOG_ARG, MIGRATE_ELEMENT_KEY_LOG_ARG + 1);
assert_eq!(VADD_SET_FLAGS_ARG, MIGRATE_INDEX_KEY_LOG_ARG + 1);
assert_eq!(CREATE_INDEX_ARG, VADD_SET_FLAGS_ARG + 1);
assert_eq!(VSETATTR_APPEND_LOG_ARG, CREATE_INDEX_ARG + 1);
assert_eq!(RECORD_TYPE, 1);
assert_eq!(MAX_VECTOR_DIMENSIONS, 65_536);
assert_eq!(MAX_FILTERING_SCALE_FACTOR, 256);
assert_eq!(MAX_EXPLORATION_FACTOR, 1_000_000);
}
#[test]
fn error_msg_table() {
assert_eq!(VectorManager::error_msg(VectorManagerResult::OK), b"");
assert!(VectorManager::error_msg(VectorManagerResult::MissingElement).starts_with(b"ERR"));
assert_eq!(VectorManagerResult::default(), VectorManagerResult::Invalid);
}
#[test]
fn clear_index_pointer_only_for_index_sized() {
let mut value = fresh_index(8, 4).to_bytes();
value[8..16].copy_from_slice(&123u64.to_le_bytes());
VectorManager::clear_index_pointer(&mut value);
assert_eq!(u64::from_le_bytes(value[8..16].try_into().unwrap()), 0);
let mut junk = vec![1u8; 10];
VectorManager::clear_index_pointer(&mut junk);
assert!(junk.iter().all(|b| *b == 1));
}
#[test]
fn add_remove_and_query_lifecycle() {
let manager = manager();
manager.service.create_index(
16,
2,
0,
VectorQuantType::NoQuant,
64,
8,
VectorDistanceMetricType::L2,
);
let index = fresh_index(16, 2).to_bytes();
let res = manager.try_add(
b"set",
&index,
b"elem-a",
VectorValueType::FP32,
&f32_bytes(&[1.0, 1.0]),
b"{\"k\":1}",
0,
VectorQuantType::Invalid,
8,
VectorDistanceMetricType::L2,
);
assert_eq!(res.unwrap(), VectorManagerResult::OK);
let res = manager.try_add(
b"set",
&index,
b"elem-a",
VectorValueType::FP32,
&f32_bytes(&[2.0, 2.0]),
b"",
0,
VectorQuantType::Invalid,
8,
VectorDistanceMetricType::L2,
);
assert_eq!(res.unwrap(), VectorManagerResult::Duplicate);
let res = manager.try_add(
b"set",
&index,
b"elem-b",
VectorValueType::FP32,
&f32_bytes(&[1.0]),
b"",
0,
VectorQuantType::Invalid,
8,
VectorDistanceMetricType::L2,
);
assert_eq!(res.unwrap_err().result, VectorManagerResult::BadParams);
let res = manager.try_add(
b"set",
&index,
b"elem-b",
VectorValueType::FP32,
&f32_bytes(&[1.0, 2.0]),
b"",
0,
VectorQuantType::Invalid,
4,
VectorDistanceMetricType::L2,
);
assert_eq!(
res.unwrap_err().message,
b"ERR asked M value mismatch with existing vector set".to_vec()
);
assert!(manager.is_member(&index, b"elem-a"));
assert!(!manager.is_member(&index, b"elem-z"));
assert_eq!(
manager.fetch_single_vector_element_attributes(&index, b"elem-a"),
VectorManagerResult::OK
);
let emb = manager.try_get_embedding(&index, b"elem-a").unwrap();
assert_eq!(emb, vec![1.0, 1.0]);
let (raw, quant, norm, range) = manager.try_get_raw_embedding(&index, b"elem-a").unwrap();
assert_eq!(quant, VectorQuantType::NoQuant);
assert!((norm - 1.0).abs() < f64::EPSILON);
assert!(range.is_none());
assert_eq!(raw, f32_bytes(&[1.0, 1.0]));
assert_eq!(
manager.try_remove(&index, b"elem-a"),
VectorManagerResult::OK
);
assert_eq!(
manager.try_remove(&index, b"elem-a"),
VectorManagerResult::MissingElement
);
assert!(!manager.is_member(&index, b"elem-a"));
assert_eq!(manager.try_get_embedding(&index, b"elem-a"), None);
}
#[test]
fn value_similarity_with_post_filter() {
let manager = manager();
manager.service.create_index(
32,
2,
0,
VectorQuantType::NoQuant,
64,
8,
VectorDistanceMetricType::L2,
);
let index = fresh_index(32, 2).to_bytes();
for (name, v, attr) in [
("near", vec![1.0, 0.0], "{\"n\": 1}"),
("far", vec![5.0, 5.0], "{\"n\": 100}"),
("mid", vec![2.0, 0.0], "{\"n\": 5}"),
] {
manager
.try_add(
b"set",
&index,
name.as_bytes(),
VectorValueType::FP32,
&f32_bytes(&v),
attr.as_bytes(),
0,
VectorQuantType::Invalid,
8,
VectorDistanceMetricType::L2,
)
.unwrap();
}
let out = manager
.value_similarity(
&index,
VectorValueType::FP32,
&f32_bytes(&[1.0, 0.0]),
3,
32,
b"",
0,
f32::INFINITY,
true,
)
.unwrap();
assert_eq!(out.found, 3);
assert_eq!(out.id_format, VectorIdFormat::I32LengthPrefixed);
assert!(out.output_distances[0] < out.output_distances[2]);
let out = manager
.value_similarity(
&index,
VectorValueType::FP32,
&f32_bytes(&[1.0, 0.0]),
3,
32,
b"",
0,
2.0,
false,
)
.unwrap();
assert_eq!(out.found, 2);
let ids: Vec<&[u8]> = unpack_length_prefixed(&out.output_ids);
assert!(ids.contains(&b"near".as_slice()) && ids.contains(&b"mid".as_slice()));
let out = manager
.value_similarity(
&index,
VectorValueType::FP32,
&f32_bytes(&[1.0, 0.0]),
3,
32,
b".n > 1",
16,
f32::INFINITY,
true,
)
.unwrap();
assert_eq!(out.found, 2);
let ids: Vec<&[u8]> = unpack_length_prefixed(&out.output_ids);
assert!(ids.contains(&b"far".as_slice()) && ids.contains(&b"mid".as_slice()));
let passed = out
.filter_bitmap
.iter()
.map(|b| b.count_ones())
.sum::<u32>();
assert_eq!(passed, 2);
let out = manager.value_similarity(
&index,
VectorValueType::FP32,
&f32_bytes(&[1.0, 0.0]),
3,
32,
b".n > >",
16,
f32::INFINITY,
false,
);
assert_eq!(
out.unwrap_err().message,
b"ERR Compiling filter failed".to_vec()
);
}
#[test]
fn element_similarity_and_attributes_batch() {
let manager = manager();
manager.service.create_index(
48,
1,
0,
VectorQuantType::NoQuant,
64,
8,
VectorDistanceMetricType::L2,
);
let index = fresh_index(48, 1).to_bytes();
manager
.try_add(
b"set",
&index,
b"x",
VectorValueType::FP32,
&f32_bytes(&[0.0]),
b"{\"a\":1}",
0,
VectorQuantType::Invalid,
8,
VectorDistanceMetricType::L2,
)
.unwrap();
manager
.try_add(
b"set",
&index,
b"y",
VectorValueType::FP32,
&f32_bytes(&[10.0]),
b"",
0,
VectorQuantType::Invalid,
8,
VectorDistanceMetricType::L2,
)
.unwrap();
let out = manager
.element_similarity(&index, b"x", 2, 32, b"", 0, f32::INFINITY, false)
.unwrap();
assert_eq!(out.found, 2);
let err = manager
.element_similarity(&index, b"ghost", 2, 32, b"", 0, f32::INFINITY, false)
.unwrap_err();
assert_eq!(err.result, VectorManagerResult::MissingElement);
assert_eq!(err.message, b"Element not in Vector Set".to_vec());
let mut ids = Vec::new();
ids.extend_from_slice(&1i32.to_le_bytes());
ids.extend_from_slice(b"x");
ids.extend_from_slice(&1i32.to_le_bytes());
ids.extend_from_slice(b"y");
let attrs = manager.fetch_vector_element_attributes(48, &ids);
let parsed = unpack_length_prefixed(&attrs);
assert_eq!(parsed[0], b"{\"a\":1}".as_slice());
assert_eq!(parsed[1], b"".as_slice());
assert!(manager.try_set_attribute(&index, b"y", b"{\"b\":2}"));
assert!(!manager.try_set_attribute(&index, b"none", b"{}"));
}
#[test]
fn request_deletion_respects_suppress_cleanup() {
let manager = manager();
let mut index = fresh_index(64, 2);
manager.service.create_index(
64,
2,
0,
VectorQuantType::NoQuant,
16,
4,
VectorDistanceMetricType::L2,
);
assert_eq!(
manager
.service
.insert(64, b"e", &f32_bytes(&[0.0, 0.0]), b""),
DiskAnnInsertResult::True
);
manager.request_deletion(&index.to_bytes());
assert!(manager.request_cleanup_task_channel.has_pending());
assert_eq!(manager.request_cleanup_task_channel.try_read(), Some(64));
assert_eq!(manager.service.card(64), 0);
manager.service.create_index(
66,
2,
0,
VectorQuantType::NoQuant,
16,
4,
VectorDistanceMetricType::L2,
);
assert_eq!(
manager
.service
.insert(66, b"e", &f32_bytes(&[0.0, 0.0]), b""),
DiskAnnInsertResult::True
);
index.context = 66;
index.flags = VectorSetFlags::SUPPRESS_CLEANUP;
manager.request_deletion(&index.to_bytes());
assert!(!manager.request_cleanup_task_channel.has_pending());
assert_eq!(manager.service.card(66), 1);
manager.request_deletion(&[0u8; 10]);
}
#[test]
fn drop_in_memory_index_flow() {
let manager = manager();
manager.service.create_index(
80,
2,
0,
VectorQuantType::NoQuant,
16,
4,
VectorDistanceMetricType::L2,
);
assert_eq!(
manager
.service
.insert(80, b"e", &f32_bytes(&[0.0, 0.0]), b""),
DiskAnnInsertResult::True
);
let mut index = fresh_index(80, 2);
let key = b"dropkey".to_vec();
manager.request_drop_in_memory_index(&key, &index.to_bytes());
assert!(manager.requested_drops.lock().contains_key(&key));
assert!(manager.request_drop_task_channel.has_pending());
manager.request_drop_in_memory_index(&key, &index.to_bytes());
assert_eq!(manager.requested_drops.lock().len(), 1);
index.flags = VectorSetFlags::SUPPRESS_CLEANUP;
manager.request_drop_in_memory_index(b"other", &index.to_bytes());
assert!(
!manager
.requested_drops
.lock()
.contains_key(b"other".as_slice())
);
assert_eq!(manager.service.card(80), 1);
manager.drop_in_memory_index(&index.to_bytes());
assert_eq!(manager.service.card(80), 0);
}
#[test]
fn recovery_reconciliation_flow() {
let manager = manager();
let ctx = manager.next_vector_set_context(1).unwrap();
let record = fresh_index(ctx, 4).to_bytes();
manager.recovered_vector_set_index_key(&record);
let meta_bytes = manager.context_metadatas.lock()[0].to_bytes();
assert!(manager.recovered_context_metadata(&0i32.to_le_bytes(), &meta_bytes));
assert!(!manager.recovered_context_metadata(&0i32.to_le_bytes(), &meta_bytes));
assert!(manager.recovered_context_metadata(&0i32.to_le_bytes(), &[0u8; 8]));
assert!(manager.recovered_context_metadata(&[0u8; 8], &meta_bytes));
assert!(manager.reconcile_recovered_state(false));
manager.request_deletion(&[0u8; 3]);
manager.recovered_vector_set_index_key(&[0u8; 3]);
let manager2 = VectorManager::new(VectorManagerOptions {
is_enabled: true,
..Default::default()
});
let _ = manager2.next_vector_set_context(0);
assert!(!manager2.reconcile_recovered_state(true));
}
#[test]
fn sanitize_ingested_records() {
let manager = manager();
let mut meta = [0u8; CONTEXT_METADATA_SIZE];
assert!(manager.sanitize_and_track_ingested_record_if_applicable(
false,
Some(&[METADATA_NAMESPACE]),
0,
&0i32.to_le_bytes(),
&mut meta,
));
let mut value = fresh_index(96, 2).to_bytes();
assert!(manager.sanitize_and_track_ingested_record_if_applicable(
false,
None,
RECORD_TYPE,
b"k",
&mut value
));
assert!(manager.recovered_indexes.lock().contains_key(&96));
assert_eq!(u64::from_le_bytes(value[8..16].try_into().unwrap()), 0);
let mut value2 = fresh_index(98, 2).to_bytes();
assert!(manager.sanitize_and_track_ingested_record_if_applicable(
true,
None,
RECORD_TYPE,
b"k",
&mut value2
));
assert!(!manager.recovered_indexes.lock().contains_key(&98));
}
#[test]
fn storage_session_assertion() {
let manager = manager();
manager.assert_have_storage_session();
assert!(VectorManager::index_has_suppress_cleanup(&Index {
flags: VectorSetFlags::SUPPRESS_CLEANUP,
..Index::default()
}));
}
}