use std::cmp::Ordering;
use std::collections::{BTreeMap, BTreeSet, BinaryHeap};
use std::path::PathBuf;
use std::sync::Arc;
use frankensearch_core::{BoundQueryEmbedding, DocId};
use frankensearch_index::{
FsviAdmissionError, FsviV2IdentityBinding, FsviV2Witness, ValidatedFsviBytes,
};
use super::{NativeAnnIndex, NativeRetrievalMode, checkpoint, invalid};
use crate::{Cx, Embedder, SearchResult};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct NativeShardRow {
pub shard: usize,
pub physical_row: u32,
}
#[derive(Debug, Clone, PartialEq)]
pub struct NativeShardHit {
pub doc_id: DocId,
pub score: f32,
pub row: NativeShardRow,
}
impl NativeShardHit {
#[must_use]
pub fn cmp_rank(&self, other: &Self) -> Ordering {
other
.score
.total_cmp(&self.score)
.then_with(|| self.doc_id.cmp(&other.doc_id))
.then_with(|| self.row.shard.cmp(&other.row.shard))
.then_with(|| self.row.physical_row.cmp(&other.row.physical_row))
}
}
#[derive(Debug, Clone)]
pub struct NativeShardSet {
shards: Vec<Arc<NativeAnnIndex>>,
live_count: usize,
physical_count: usize,
}
impl NativeShardSet {
pub fn try_replace(
&mut self,
cx: &Cx,
expected_current: &[FsviV2Witness],
candidate: Self,
) -> SearchResult<()> {
self.validate_replacement(cx, expected_current, candidate.owner_witnesses().next())?;
checkpoint(cx, "native_ann.shards.replace_commit")?;
*self = candidate;
Ok(())
}
pub fn try_replace_published(
&mut self,
cx: &Cx,
expected_current: &[FsviV2Witness],
expected_next: &[FsviV2Witness],
artifacts: &[(PathBuf, FsviV2IdentityBinding, Option<PathBuf>)],
) -> Result<(), FsviAdmissionError> {
self.validate_replacement(cx, expected_current, expected_next.first())?;
let candidate = Self::open_published(cx, expected_next, artifacts)?;
self.try_replace(cx, expected_current, candidate)?;
Ok(())
}
fn validate_replacement(
&self,
cx: &Cx,
expected_current: &[FsviV2Witness],
next: Option<&FsviV2Witness>,
) -> SearchResult<()> {
checkpoint(cx, "native_ann.shards.replace_admission")?;
if expected_current.len() != self.shards.len() {
return Err(invalid(
"shards.replace.expected_current",
"cardinality",
"replacement requires the complete current shard inventory",
));
}
for (shard, expected) in self.shards.iter().zip(expected_current) {
checkpoint(cx, "native_ann.shards.replace_current")?;
if shard.owner_witness() != expected {
return Err(invalid(
"shards.replace.expected_current",
"stale",
"the current exact ordered inventory differs from the refresh's expectation",
));
}
}
let next = next.ok_or_else(|| invalid(
"shards.inventory", "empty-successor",
"a successor must retain at least one identity-bearing shard, even for an empty corpus",
))?;
let current = self.shards[0].owner_witness();
if next.generation.sequence <= current.generation.sequence {
return Err(invalid(
"shards.replace.generation",
"not-newer",
"replacement requires a strictly higher generation sequence; same-sequence nonce changes are not successors",
));
}
if next.space_fingerprint != current.space_fingerprint
|| next.producer_fingerprint != current.producer_fingerprint
|| next.input_fingerprint != current.input_fingerprint
|| next.dimension != current.dimension
{
return Err(invalid(
"shards.replace.identity",
"changed",
"live replacement must preserve space, producer, input and dimension; open a separate handle for a model migration",
));
}
Ok(())
}
pub fn open_published(
cx: &Cx,
expected: &[FsviV2Witness],
artifacts: &[(PathBuf, FsviV2IdentityBinding, Option<PathBuf>)],
) -> Result<Self, FsviAdmissionError> {
checkpoint(cx, "native_ann.shards.open")?;
if expected.is_empty() || expected.len() != artifacts.len() {
return Err(invalid(
"shards.inventory",
"cardinality",
"a nonempty expected inventory must match every declared artifact",
)
.into());
}
let reference = &expected[0];
let mut images = BTreeSet::new();
for (witness, (path, binding, graph)) in expected.iter().zip(artifacts) {
checkpoint(cx, "native_ann.shards.open_spec")?;
if !path.is_absolute() || graph.as_ref().is_some_and(|path| !path.is_absolute()) {
return Err(invalid(
"shards.paths",
"relative",
"vector and optional graph paths must be explicit absolute paths",
)
.into());
}
if !images.insert(witness.whole_image_sha256) {
return Err(invalid(
"shards.inventory",
"duplicate-image",
"the selected inventory must not repeat a physical shard image",
)
.into());
}
if witness.generation != reference.generation
|| binding.generation() != witness.generation
{
return Err(invalid(
"shards.generation",
"mismatch",
"each declared binding and expected shard must name the same generation",
)
.into());
}
if witness.space_fingerprint != reference.space_fingerprint
|| witness.producer_fingerprint != reference.producer_fingerprint
|| witness.input_fingerprint != reference.input_fingerprint
|| witness.dimension != reference.dimension
{
return Err(invalid(
"shards.identity",
"mismatch",
"expected partitions must agree on space, producer, input and dimension",
)
.into());
}
}
let mut shards = Vec::with_capacity(expected.len());
for (witness, (path, binding, _)) in expected.iter().zip(artifacts) {
checkpoint(cx, "native_ann.shards.open_vector")?;
let opened = ValidatedFsviBytes::reopen_exact(path, binding, witness);
checkpoint(cx, "native_ann.shards.open_vector_complete")?;
let owner = Arc::new(opened?);
shards.push(Arc::new(NativeAnnIndex::exact(cx, owner)?));
}
let mut admitted = Self::admit(cx, expected, shards)?;
for (shard, (_, _, graph)) in admitted.shards.iter_mut().zip(artifacts) {
if let Some(path) = graph {
checkpoint(cx, "native_ann.shards.open_graph")?;
let opened = NativeAnnIndex::load_or_exact(cx, Arc::clone(&shard.owner), path)?;
*shard = Arc::new(opened);
}
}
checkpoint(cx, "native_ann.shards.open_complete")?;
Ok(admitted)
}
pub fn owner_witnesses(&self) -> impl ExactSizeIterator<Item = &FsviV2Witness> + '_ {
self.shards.iter().map(|shard| shard.owner_witness())
}
pub fn admit(
cx: &Cx,
expected: &[FsviV2Witness],
shards: Vec<Arc<NativeAnnIndex>>,
) -> SearchResult<Self> {
checkpoint(cx, "native_ann.shards.admission")?;
if shards.is_empty() || shards.len() != expected.len() {
return Err(invalid(
"shards.inventory",
"cardinality",
"a nonempty exact ordered shard inventory is required",
));
}
let reference = shards[0].owner_witness();
let mut live_count = 0_usize;
let mut physical_count = 0_usize;
let mut seen_images = BTreeSet::new();
for (shard, expected) in shards.iter().zip(expected) {
checkpoint(cx, "native_ann.shards.identity")?;
let actual = shard.owner_witness();
if actual != expected {
return Err(invalid(
"shards.inventory",
"witness-mismatch",
"each shard must match the expected whole-image witness at its exact position",
));
}
if !seen_images.insert(actual.whole_image_sha256) {
return Err(invalid(
"shards.inventory",
"duplicate-image",
"the selected inventory must not repeat an identical physical shard image",
));
}
if actual.generation != reference.generation {
return Err(invalid(
"shards.generation",
"mismatch",
"all partitions must belong to the same complete artifact generation",
));
}
if actual.space_fingerprint != reference.space_fingerprint
|| actual.producer_fingerprint != reference.producer_fingerprint
|| actual.input_fingerprint != reference.input_fingerprint
|| actual.dimension != reference.dimension
{
return Err(invalid(
"shards.identity",
"mismatch",
"partitions must agree on space, producer, input contract and dimension",
));
}
live_count = live_count.checked_add(shard.live_count()).ok_or_else(|| {
invalid(
"shards.live_count",
"overflow",
"live row count must fit usize",
)
})?;
physical_count = physical_count.checked_add(shard.len()).ok_or_else(|| {
invalid(
"shards.physical_count",
"overflow",
"physical row count must fit usize",
)
})?;
}
let mut partition_by_document: BTreeMap<String, usize> = BTreeMap::new();
for (ordinal, shard) in shards.iter().enumerate() {
for physical in 0..shard.len() {
checkpoint(cx, "native_ann.shards.membership")?;
let row = shard.owner.row(physical)?;
if partition_by_document
.insert(row.doc_id().to_owned(), ordinal)
.is_some_and(|previous| previous != ordinal)
{
return Err(invalid(
"shards.membership",
"overlap",
"physical document identities must not overlap across shard partitions",
));
}
}
}
checkpoint(cx, "native_ann.shards.admitted")?;
Ok(Self {
shards,
live_count,
physical_count,
})
}
#[must_use]
pub fn shard_count(&self) -> usize {
self.shards.len()
}
#[must_use]
pub const fn live_count(&self) -> usize {
self.live_count
}
#[must_use]
pub const fn physical_count(&self) -> usize {
self.physical_count
}
#[must_use]
pub fn shard(&self, ordinal: usize) -> Option<&NativeAnnIndex> {
self.shards.get(ordinal).map(AsRef::as_ref)
}
pub fn retrieval_modes(&self) -> impl ExactSizeIterator<Item = NativeRetrievalMode> + '_ {
self.shards.iter().map(|shard| shard.retrieval_mode())
}
pub fn search(
&self,
cx: &Cx,
query: &BoundQueryEmbedding,
k: usize,
ef: Option<usize>,
) -> SearchResult<Vec<NativeShardHit>> {
self.search_filtered(cx, query, k, ef, |_| true)
}
pub fn search_filtered<F>(
&self,
cx: &Cx,
query: &BoundQueryEmbedding,
k: usize,
ef: Option<usize>,
accept: F,
) -> SearchResult<Vec<NativeShardHit>>
where
F: Fn(&str) -> bool,
{
checkpoint(cx, "native_ann.shards.search")?;
self.shards[0].admit_identity(query.identity())?;
let target = k.min(self.live_count);
if target == 0 {
return Ok(Vec::new());
}
let mut winners: BinaryHeap<RankedHit> = BinaryHeap::new();
for (ordinal, shard) in self.shards.iter().enumerate() {
checkpoint(cx, "native_ann.shards.partition")?;
let candidates = shard.search_filtered(cx, query, target, ef, &accept)?;
for candidate in candidates {
checkpoint(cx, "native_ann.shards.merge")?;
let hit = NativeShardHit {
doc_id: candidate.doc_id,
score: candidate.score,
row: NativeShardRow {
shard: ordinal,
physical_row: candidate.index,
},
};
if winners.len() == target {
if winners
.peek()
.is_some_and(|worst| hit.cmp_rank(&worst.0) != Ordering::Less)
{
continue;
}
let _ = winners.pop();
}
winners.push(RankedHit(hit));
}
}
let mut hits: Vec<_> = winners.into_iter().map(|hit| hit.0).collect();
hits.sort_unstable_by(NativeShardHit::cmp_rank);
checkpoint(cx, "native_ann.shards.complete")?;
Ok(hits)
}
pub async fn search_text(
&self,
cx: &Cx,
embedder: &dyn Embedder,
text: &str,
k: usize,
ef: Option<usize>,
) -> SearchResult<Vec<NativeShardHit>> {
self.search_text_filtered(cx, embedder, text, k, ef, |_| true)
.await
}
pub async fn search_text_filtered<F>(
&self,
cx: &Cx,
embedder: &dyn Embedder,
text: &str,
k: usize,
ef: Option<usize>,
accept: F,
) -> SearchResult<Vec<NativeShardHit>>
where
F: Fn(&str) -> bool + Send,
{
checkpoint(cx, "native_ann.shards.text")?;
let reference = &self.shards[0];
reference.admit_identity(embedder.identity()?)?;
if k == 0 || self.live_count == 0 {
return Ok(Vec::new());
}
let query = reference.embed_query(cx, embedder, text).await?;
self.search_filtered(cx, &query, k, ef, accept)
}
}
struct RankedHit(NativeShardHit);
impl PartialEq for RankedHit {
fn eq(&self, other: &Self) -> bool {
self.cmp(other) == Ordering::Equal
}
}
impl Eq for RankedHit {}
impl PartialOrd for RankedHit {
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
impl Ord for RankedHit {
fn cmp(&self, other: &Self) -> Ordering {
self.0.cmp_rank(&other.0)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::cell::Cell;
use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
use crate::SearchError;
use frankensearch_core::generation::{
ArtifactGenerationIdentityV1, EmbeddingIdentityBundleV1, QuantizationFormat,
};
use frankensearch_core::traits::{IdentityBoundEmbedding, ModelCategory, SearchFuture};
use frankensearch_index::native_hnsw::HnswParams;
use frankensearch_index::{FsviV2IdentityBinding, ValidatedFsviBytes, VectorIndex};
fn identity() -> EmbeddingIdentityBundleV1 {
EmbeddingIdentityBundleV1::explicit_test_model("native-shards", 2)
}
fn partition(
cx: &Cx,
identity: &EmbeddingIdentityBundleV1,
rows: &[(&str, [f32; 2], bool)],
generation: u64,
format: QuantizationFormat,
ann: bool,
) -> Arc<NativeAnnIndex> {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("partition.fsvi");
let mut bundle = identity.clone();
bundle.storage.format = "fsvi-v2".to_owned();
bundle.storage.quantization = format;
bundle.storage.endianness = "little-endian".to_owned();
let binding = FsviV2IdentityBinding::new(
ArtifactGenerationIdentityV1::new(generation, [0x55; 16]).unwrap(),
bundle.freeze().unwrap(),
)
.unwrap();
let mut writer = VectorIndex::create_v2(&path, binding.clone()).unwrap();
for &(id, vector, live) in rows {
if live {
writer.write_record(id, &vector).unwrap();
} else {
writer.write_tombstone_record(id, &vector).unwrap();
}
}
writer.finish().unwrap();
let bytes: Arc<[u8]> = std::fs::read(path).unwrap().into();
let owner = Arc::new(ValidatedFsviBytes::from_arc(bytes, &binding).unwrap());
Arc::new(if ann {
NativeAnnIndex::build(cx, owner, HnswParams::default(), 7).unwrap()
} else {
NativeAnnIndex::exact(cx, owner).unwrap()
})
}
fn admit(cx: &Cx, shards: Vec<Arc<NativeAnnIndex>>) -> NativeShardSet {
let expected: Vec<_> = shards
.iter()
.map(|shard| shard.owner_witness().clone())
.collect();
NativeShardSet::admit(cx, &expected, shards).unwrap()
}
fn query() -> BoundQueryEmbedding {
BoundQueryEmbedding::new(vec![1.0, 0.0], identity()).unwrap()
}
struct Provider {
identity: EmbeddingIdentityBundleV1,
calls: AtomicUsize,
cancel: bool,
}
impl Provider {
fn new() -> Self {
Self {
identity: identity(),
calls: AtomicUsize::new(0),
cancel: false,
}
}
}
impl Embedder for Provider {
fn embed<'a>(&'a self, _cx: &'a Cx, _text: &'a str) -> SearchFuture<'a, Vec<f32>> {
Box::pin(async { Ok(vec![1.0, 0.0]) })
}
fn embed_bound<'a>(
&'a self,
_cx: &'a Cx,
_text: &'a str,
) -> SearchFuture<'a, IdentityBoundEmbedding> {
Box::pin(async move {
self.calls.fetch_add(1, AtomicOrdering::SeqCst);
if self.cancel {
return Err(SearchError::Cancelled {
phase: "test.provider".to_owned(),
reason: "cancelled".to_owned(),
});
}
Ok(IdentityBoundEmbedding {
values: vec![1.0, 0.0],
identity: self.identity.clone(),
})
})
}
fn identity(&self) -> SearchResult<&EmbeddingIdentityBundleV1> {
Ok(&self.identity)
}
fn id(&self) -> &'static str {
"native-shards-test"
}
fn model_name(&self) -> &str {
self.id()
}
fn dimension(&self) -> usize {
2
}
fn is_ready(&self) -> bool {
true
}
fn is_semantic(&self) -> bool {
false
}
fn category(&self) -> ModelCategory {
ModelCategory::HashEmbedder
}
}
#[test]
fn global_top_k_matches_the_full_union_for_mixed_storage_and_backends() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
for ann in [false, true] {
let a = partition(
&cx,
&identity(),
&[
("dead", [1.0, 0.0], false),
("a", [0.7, 0.3], true),
("c", [0.2, 0.8], true),
],
1,
QuantizationFormat::F16,
ann,
);
let b = partition(
&cx,
&identity(),
&[("b", [0.9, 0.1], true), ("d", [-0.1, 0.9], true)],
1,
QuantizationFormat::F32,
false,
);
let set = admit(&cx, vec![a, b]);
let mut exact = Vec::new();
for (ordinal, shard) in set.shards.iter().enumerate() {
for hit in shard.search(&cx, &query(), usize::MAX, Some(100)).unwrap() {
exact.push(NativeShardHit {
doc_id: hit.doc_id,
score: hit.score,
row: NativeShardRow {
shard: ordinal,
physical_row: hit.index,
},
});
}
}
exact.sort_unstable_by(NativeShardHit::cmp_rank);
for k in [0, 1, 2, 3, 4, usize::MAX] {
assert_eq!(
set.search(&cx, &query(), k, Some(100)).unwrap(),
exact[..k.min(exact.len())]
);
}
let hits = set.search(&cx, &query(), 2, None).unwrap();
assert_eq!(set.shards[1].owner.doc_id_at(1).unwrap(), "b");
assert_eq!(
hits[0].row,
NativeShardRow {
shard: 1,
physical_row: 1
}
);
assert_eq!(
hits[1].row,
NativeShardRow {
shard: 0,
physical_row: 1
}
);
for hit in &hits {
let owner = &set.shards[hit.row.shard].owner;
assert_eq!(
owner.doc_id_at(hit.row.physical_row as usize).unwrap(),
hit.doc_id
);
}
assert_eq!(
(set.shard_count(), set.live_count(), set.physical_count()),
(2, 4, 5)
);
assert!(set.shard(2).is_none());
assert_eq!(set.retrieval_modes().len(), 2);
}
});
}
#[test]
fn filters_expand_each_partition_before_global_selection() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let a = partition(
&cx,
&identity(),
&[("reject-a", [1.0, 0.0], true), ("keep-a", [0.2, 0.8], true)],
1,
QuantizationFormat::F32,
true,
);
let b = partition(
&cx,
&identity(),
&[("reject-b", [0.9, 0.1], true), ("keep-b", [0.1, 0.9], true)],
1,
QuantizationFormat::F32,
true,
);
let set = admit(&cx, vec![a, b]);
let calls = Cell::new(0); let hits = set
.search_filtered(&cx, &query(), 2, Some(1), |id| {
calls.set(calls.get() + 1);
id.starts_with("keep-")
})
.unwrap();
assert_eq!(
hits.iter()
.map(|hit| hit.doc_id.as_str())
.collect::<Vec<_>>(),
["keep-a", "keep-b"]
);
assert!(calls.get() > 0);
assert!(
set.search_filtered(&cx, &query(), 2, None, |_| false)
.unwrap()
.is_empty()
);
});
}
#[test]
fn equal_scores_use_document_order_not_partition_order() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let z = partition(
&cx,
&identity(),
&[("z", [1.0, 0.0], true)],
1,
QuantizationFormat::F32,
false,
);
let a = partition(
&cx,
&identity(),
&[("a", [1.0, 0.0], true)],
1,
QuantizationFormat::F32,
false,
);
let set = admit(&cx, vec![z, a]);
assert_eq!(set.search(&cx, &query(), 1, None).unwrap()[0].doc_id, "a");
});
}
#[test]
fn exact_inventory_rejects_missing_reordered_and_substituted_shards() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let a = partition(
&cx,
&identity(),
&[("a", [1.0, 0.0], true)],
1,
QuantizationFormat::F32,
false,
);
let b = partition(
&cx,
&identity(),
&[("b", [0.0, 1.0], true)],
1,
QuantizationFormat::F32,
false,
);
let replacement = partition(
&cx,
&identity(),
&[("a", [0.0, 1.0], true)],
1,
QuantizationFormat::F32,
false,
);
let expected = [a.owner_witness().clone(), b.owner_witness().clone()];
for shards in [
vec![Arc::clone(&a)],
vec![Arc::clone(&b), Arc::clone(&a)],
vec![replacement, Arc::clone(&b)],
vec![Arc::clone(&a), Arc::clone(&a)],
] {
assert!(
matches!(NativeShardSet::admit(&cx, &expected, shards), Err(SearchError::InvalidConfig { ref field, .. }) if field == "native_ann.shards.inventory")
);
}
assert!(NativeShardSet::admit(&cx, &[], Vec::new()).is_err());
assert_eq!(a.search(&cx, &query(), 1, None).unwrap()[0].doc_id, "a");
});
}
#[test]
fn independently_valid_shards_must_share_generation_and_producer() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let a = partition(
&cx,
&identity(),
&[("a", [1.0, 0.0], true)],
1,
QuantizationFormat::F32,
false,
);
let mut foreign = identity();
foreign.producer.backend = "different-producer".to_owned();
for (identity, generation, field) in [
(identity(), 2, "native_ann.shards.generation"),
(foreign, 1, "native_ann.shards.identity"),
(
EmbeddingIdentityBundleV1::explicit_test_model("other-space", 2),
1,
"native_ann.shards.identity",
),
] {
let b = partition(
&cx,
&identity,
&[("b", [0.0, 1.0], true)],
generation,
QuantizationFormat::F32,
false,
);
let expected = [a.owner_witness().clone(), b.owner_witness().clone()];
assert!(
matches!(NativeShardSet::admit(&cx, &expected, vec![Arc::clone(&a), b]), Err(SearchError::InvalidConfig { field: actual, .. }) if actual == field)
);
}
});
}
#[test]
fn overlapping_live_or_tombstoned_partitions_are_not_deduplicated_or_resurrected() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let a = partition(
&cx,
&identity(),
&[("same", [1.0, 0.0], true)],
1,
QuantizationFormat::F32,
false,
);
for live in [true, false] {
let b = partition(
&cx,
&identity(),
&[("same", [0.0, 1.0], live)],
1,
QuantizationFormat::F32,
false,
);
let expected = [a.owner_witness().clone(), b.owner_witness().clone()];
assert!(
matches!(NativeShardSet::admit(&cx, &expected, vec![Arc::clone(&a), b]), Err(SearchError::InvalidConfig { ref field, .. }) if field == "native_ann.shards.membership")
);
}
});
}
#[test]
fn text_embeds_once_even_when_the_first_partition_is_empty() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let a = partition(&cx, &identity(), &[], 1, QuantizationFormat::F32, false);
let b = partition(
&cx,
&identity(),
&[("b", [1.0, 0.0], true)],
1,
QuantizationFormat::F32,
true,
);
let set = admit(&cx, vec![a, b]);
let provider = Provider::new();
let hits = set
.search_text(&cx, &provider, "query", 1, None)
.await
.unwrap();
assert_eq!(hits[0].doc_id, "b");
assert_eq!(hits[0].row.shard, 1);
assert_eq!(provider.calls.load(AtomicOrdering::SeqCst), 1);
});
}
#[test]
fn zero_work_queries_still_refuse_foreign_identities_without_filter_calls() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
for rows in [vec![], vec![("a", [1.0, 0.0], true)]] {
let set = admit(
&cx,
vec![partition(
&cx,
&identity(),
&rows,
1,
QuantizationFormat::F32,
false,
)],
);
let mut foreign = identity();
foreign.producer.backend = "foreign".to_owned();
let query = BoundQueryEmbedding::new(vec![1.0, 0.0], foreign).unwrap();
let calls = Cell::new(0);
for k in [0, 1] {
assert!(
set.search_filtered(&cx, &query, k, None, |_| {
calls.set(calls.get() + 1);
true
})
.is_err()
);
}
assert_eq!(calls.get(), 0);
let provider = Provider::new();
assert!(
set.search_text(&cx, &provider, "query", 0, None)
.await
.unwrap()
.is_empty()
);
assert_eq!(provider.calls.load(AtomicOrdering::SeqCst), 0);
}
});
}
#[test]
fn selected_owners_survive_caller_drop_and_no_fallback_masks_provider_cancellation() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let a = partition(
&cx,
&identity(),
&[("a", [1.0, 0.0], true)],
1,
QuantizationFormat::F32,
false,
);
let weak = Arc::downgrade(&a);
let set = admit(&cx, vec![a]);
assert!(weak.upgrade().is_some());
let mut provider = Provider::new();
provider.cancel = true;
assert!(
matches!(set.search_text(&cx, &provider, "query", 1, None).await, Err(SearchError::Cancelled { ref phase, .. }) if phase == "test.provider")
);
assert_eq!(provider.calls.load(AtomicOrdering::SeqCst), 1);
assert_eq!(set.search(&cx, &query(), 1, None).unwrap()[0].doc_id, "a");
drop(set);
assert!(weak.upgrade().is_none());
});
}
}
#[cfg(test)]
mod merge_tests {
use super::*;
#[test]
fn bounded_heap_matches_sorted_union_across_orders_limits_and_signed_zero() {
let values = [-2.0_f32, -0.0, 0.0, 1.0, 1.0, 3.0];
for rotation in 0..values.len() {
let input: Vec<_> = (0..values.len())
.map(|offset| {
let i = (rotation + offset) % values.len();
NativeShardHit {
doc_id: format!("doc-{i}").into(),
score: values[i],
row: NativeShardRow {
shard: i % 3,
physical_row: u32::try_from(i / 3).unwrap(),
},
}
})
.collect();
let mut sorted = input.clone();
sorted.sort_unstable_by(NativeShardHit::cmp_rank);
for target in 1..=values.len() {
let mut heap: BinaryHeap<RankedHit> = BinaryHeap::new();
for hit in input.iter().cloned() {
if heap.len() == target {
if heap
.peek()
.is_some_and(|worst| hit.cmp_rank(&worst.0) != Ordering::Less)
{
continue;
}
let _ = heap.pop();
}
heap.push(RankedHit(hit));
}
let mut actual: Vec<_> = heap.into_iter().map(|hit| hit.0).collect();
actual.sort_unstable_by(NativeShardHit::cmp_rank);
assert_eq!(actual, sorted[..target]);
}
}
}
}
#[cfg(all(test, any(target_os = "linux", target_os = "android")))]
mod published_tests {
use super::*;
use crate::SearchError;
use crate::native_ann::NativeExactReason;
use frankensearch_core::generation::{
ArtifactGenerationIdentityV1, EmbeddingIdentityBundleV1, QuantizationFormat,
};
use frankensearch_index::native_hnsw::HnswParams;
use frankensearch_index::{FsviSnapshotRejectionReason, VectorIndex};
type Artifact = (PathBuf, FsviV2IdentityBinding, Option<PathBuf>);
fn fixture(
cx: &Cx,
root: &std::path::Path,
name: &str,
generation: u64,
format: QuantizationFormat,
rows: &[(&str, [f32; 2], bool)],
ann: bool,
) -> (FsviV2Witness, Artifact) {
let mut identity = EmbeddingIdentityBundleV1::explicit_test_model("reopen-native", 2);
"fsvi-v2".clone_into(&mut identity.storage.format);
identity.storage.quantization = format;
"little-endian".clone_into(&mut identity.storage.endianness);
let binding = FsviV2IdentityBinding::new(
ArtifactGenerationIdentityV1::new(generation, [0x73; 16]).unwrap(),
identity.freeze().unwrap(),
)
.unwrap();
let path = root.join(format!("{name}.fsvi"));
let mut writer = VectorIndex::create_v2(&path, binding.clone()).unwrap();
for &(id, vector, live) in rows {
if live {
writer.write_record(id, &vector).unwrap();
} else {
writer.write_tombstone_record(id, &vector).unwrap();
}
}
writer.finish().unwrap();
let bytes: Arc<[u8]> = std::fs::read(&path).unwrap().into();
let owner = Arc::new(ValidatedFsviBytes::from_arc(bytes, &binding).unwrap());
let witness = owner.witness().clone();
let graph = ann.then(|| {
let graph = root.join(format!("{name}.fshnsw"));
NativeAnnIndex::build(cx, owner, HnswParams::default(), 7)
.unwrap()
.save(cx, &graph)
.unwrap();
graph
});
(witness, (path, binding, graph))
}
fn query() -> BoundQueryEmbedding {
BoundQueryEmbedding::new(
vec![1.0, 0.0],
EmbeddingIdentityBundleV1::explicit_test_model("reopen-native", 2),
)
.unwrap()
}
#[test]
fn restart_reopens_mixed_storage_and_backends_without_rebuilding() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (a, first) = fixture(
&cx,
&root,
"fast-a",
1,
QuantizationFormat::F16,
&[("dead", [1.0, 0.0], false), ("a", [0.5, 0.5], true)],
true,
);
let (b, second) = fixture(
&cx,
&root,
"fast-b",
1,
QuantizationFormat::F32,
&[("b", [0.75, 0.25], true)],
false,
);
let expected = [a, b];
let artifacts = [first, second];
let before: Vec<_> = artifacts
.iter()
.map(|item| std::fs::read(&item.0).unwrap())
.collect();
let opened = NativeShardSet::open_published(&cx, &expected, &artifacts).unwrap();
let hits = opened.search(&cx, &query(), 2, None).unwrap();
assert_eq!(
hits[0].row,
NativeShardRow {
shard: 1,
physical_row: 0
}
);
assert_eq!(
hits[1].row,
NativeShardRow {
shard: 0,
physical_row: 1
}
);
assert_eq!(
opened.owner_witnesses().cloned().collect::<Vec<_>>(),
expected
);
assert_eq!(
opened.retrieval_modes().collect::<Vec<_>>(),
[
NativeRetrievalMode::Ann,
NativeRetrievalMode::Exact {
reason: NativeExactReason::Requested
},
]
);
drop(opened);
let reopened = NativeShardSet::open_published(&cx, &expected, &artifacts).unwrap();
assert_eq!(reopened.search(&cx, &query(), 2, None).unwrap(), hits);
for (artifact, bytes) in artifacts.iter().zip(before) {
assert_eq!(std::fs::read(&artifact.0).unwrap(), bytes);
}
assert_eq!((reopened.live_count(), reopened.physical_count()), (2, 3));
});
}
#[test]
fn absent_optional_graph_is_explicit_exact_and_creates_no_sidecar() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (witness, mut artifact) = fixture(
&cx,
&root,
"one",
1,
QuantizationFormat::F32,
&[("a", [1.0, 0.0], true)],
false,
);
let graph = root.join("missing.fshnsw");
artifact.2 = Some(graph.clone());
let set = NativeShardSet::open_published(&cx, &[witness], &[artifact]).unwrap();
assert_eq!(
set.retrieval_modes().collect::<Vec<_>>(),
[NativeRetrievalMode::Exact {
reason: NativeExactReason::SidecarMissing
},]
);
assert_eq!(set.search(&cx, &query(), 1, None).unwrap()[0].doc_id, "a");
assert!(!graph.exists());
assert!(!root.join("missing.fshnsw.receipt").exists());
});
}
#[test]
fn stale_graph_does_not_rebind_its_old_vectors_or_mask_mandatory_drift() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (_, old) = fixture(
&cx,
&root,
"old",
1,
QuantizationFormat::F32,
&[("old", [1.0, 0.0], true)],
true,
);
let (witness, mut new) = fixture(
&cx,
&root,
"new",
2,
QuantizationFormat::F32,
&[("new", [0.5, 0.5], true)],
false,
);
new.2 = old.2;
let graph = new.2.as_ref().unwrap().clone();
let before = std::fs::read(&graph).unwrap();
let set = NativeShardSet::open_published(
&cx,
std::slice::from_ref(&witness),
std::slice::from_ref(&new),
)
.unwrap();
assert_eq!(
set.retrieval_modes().collect::<Vec<_>>(),
[NativeRetrievalMode::Exact {
reason: NativeExactReason::SidecarRejected
},]
);
assert_eq!(set.search(&cx, &query(), 1, None).unwrap()[0].doc_id, "new");
assert_eq!(std::fs::read(&graph).unwrap(), before);
let (_, substitute) = fixture(
&cx,
&root,
"substitute",
2,
QuantizationFormat::F32,
&[("new", [0.0, 1.0], true)],
false,
);
new.0 = substitute.0;
assert!(
matches!(NativeShardSet::open_published(&cx, &[witness], &[new]),
Err(FsviAdmissionError::SnapshotRejected(rejected))
if rejected.reason == FsviSnapshotRejectionReason::WitnessMismatch)
);
assert_eq!(set.search(&cx, &query(), 1, None).unwrap()[0].score, 0.5);
});
}
#[test]
fn every_mandatory_partition_is_admitted_before_any_graph() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (a, mut first) = fixture(
&cx,
&root,
"a",
1,
QuantizationFormat::F32,
&[("a", [1.0, 0.0], true)],
false,
);
let (b, mut second) = fixture(
&cx,
&root,
"b",
1,
QuantizationFormat::F32,
&[("b", [0.0, 1.0], true)],
false,
);
first.2 = Some(root.join("wrong-extension.txt"));
let (_, substituted) = fixture(
&cx,
&root,
"b-substitute",
1,
QuantizationFormat::F32,
&[("b", [1.0, 0.0], true)],
false,
);
second.0 = substituted.0;
assert!(
matches!(NativeShardSet::open_published(&cx, &[a, b], &[first, second]),
Err(FsviAdmissionError::SnapshotRejected(rejected))
if rejected.reason == FsviSnapshotRejectionReason::WitnessMismatch)
);
});
}
#[test]
fn adjacent_wal_is_not_recovered_as_an_optional_graph_failure() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (witness, mut artifact) = fixture(
&cx,
&root,
"wal",
1,
QuantizationFormat::F32,
&[("a", [1.0, 0.0], true)],
false,
);
let wal = frankensearch_index::wal_path_for(&artifact.0);
std::fs::write(&wal, []).unwrap();
artifact.2 = Some(root.join("missing.fshnsw"));
assert!(
matches!(NativeShardSet::open_published(&cx, &[witness], &[artifact]),
Err(FsviAdmissionError::SnapshotRejected(rejected))
if rejected.reason == FsviSnapshotRejectionReason::PublishedWalPresent)
);
assert!(wal.exists());
});
}
#[test]
fn physical_overlap_is_rejected_before_optional_graph_loading() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (a, mut first) = fixture(
&cx,
&root,
"live",
1,
QuantizationFormat::F32,
&[("same", [1.0, 0.0], true)],
false,
);
let (b, second) = fixture(
&cx,
&root,
"dead",
1,
QuantizationFormat::F32,
&[("same", [0.0, 1.0], false)],
false,
);
first.2 = Some(root.join("invalid.txt"));
assert!(
matches!(NativeShardSet::open_published(&cx, &[a, b], &[first, second]),
Err(FsviAdmissionError::Index(SearchError::InvalidConfig { ref field, .. }))
if field == "native_ann.shards.membership")
);
});
}
#[test]
fn empty_corpus_reopens_with_its_identity_and_no_query_work() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (witness, artifact) =
fixture(&cx, &root, "empty", 1, QuantizationFormat::F32, &[], false);
let set = NativeShardSet::open_published(&cx, &[witness], &[artifact]).unwrap();
assert_eq!(set.live_count(), 0);
assert!(set.search(&cx, &query(), 1, None).unwrap().is_empty());
let foreign = BoundQueryEmbedding::new(
vec![1.0, 0.0],
EmbeddingIdentityBundleV1::explicit_test_model("foreign", 2),
)
.unwrap();
assert!(set.search(&cx, &foreign, 0, None).is_err());
});
}
#[test]
fn missing_partition_and_relative_path_are_not_silently_dropped() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (a, first) = fixture(
&cx,
&root,
"a",
1,
QuantizationFormat::F32,
&[("a", [1.0, 0.0], true)],
false,
);
let (b, mut second) = fixture(
&cx,
&root,
"b",
1,
QuantizationFormat::F32,
&[("b", [0.0, 1.0], true)],
false,
);
assert!(
matches!(NativeShardSet::open_published(&cx, &[a.clone(), b.clone()],
std::slice::from_ref(&first)),
Err(FsviAdmissionError::Index(SearchError::InvalidConfig { ref field, .. }))
if field == "native_ann.shards.inventory")
);
second.0 = PathBuf::from("relative.fsvi");
assert!(
matches!(NativeShardSet::open_published(&cx, std::slice::from_ref(&b),
std::slice::from_ref(&second)),
Err(FsviAdmissionError::Index(SearchError::InvalidConfig { ref field, .. }))
if field == "native_ann.shards.paths")
);
second.0 = root.join("absent.fsvi");
assert!(NativeShardSet::open_published(&cx, &[a, b], &[first, second]).is_err());
});
}
#[test]
fn replacement_switches_the_whole_set_while_cloned_readers_keep_old_owners() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (a, first) = fixture(
&cx,
&root,
"old-a",
1,
QuantizationFormat::F16,
&[("a", [1.0, 0.0], true)],
true,
);
let (b, second) = fixture(
&cx,
&root,
"old-b",
1,
QuantizationFormat::F32,
&[("b", [0.5, 0.5], true)],
false,
);
let old_expected = [a, b];
let mut current =
NativeShardSet::open_published(&cx, &old_expected, &[first, second]).unwrap();
let retained = current.clone();
assert!(Arc::ptr_eq(&retained.shards[0], ¤t.shards[0]));
let before = retained.search(&cx, &query(), 2, None).unwrap();
let (next, artifact) = fixture(
&cx,
&root,
"successor",
2,
QuantizationFormat::F32,
&[("new", [0.75, 0.25], true)],
false,
);
current
.try_replace_published(
&cx,
&old_expected,
std::slice::from_ref(&next),
std::slice::from_ref(&artifact),
)
.unwrap();
assert_eq!((current.shard_count(), current.live_count()), (1, 1));
assert_eq!(
current.search(&cx, &query(), 2, None).unwrap()[0].doc_id,
"new"
);
assert_eq!(retained.search(&cx, &query(), 2, None).unwrap(), before);
assert_eq!(
retained.owner_witnesses().cloned().collect::<Vec<_>>(),
old_expected
);
assert!(!Arc::ptr_eq(&retained.shards[0], ¤t.shards[0]));
let reopened = NativeShardSet::open_published(&cx, &[next], &[artifact]).unwrap();
assert_eq!(
current.search(&cx, &query(), 2, None).unwrap(),
reopened.search(&cx, &query(), 2, None).unwrap()
);
});
}
#[test]
fn stale_expected_current_checks_every_partition_before_candidate_io() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (a, first) = fixture(
&cx,
&root,
"a",
1,
QuantizationFormat::F32,
&[("a", [1.0, 0.0], true)],
false,
);
let (b, second) = fixture(
&cx,
&root,
"b",
1,
QuantizationFormat::F32,
&[("b", [0.5, 0.5], true)],
false,
);
let expected = [a, b];
let mut current =
NativeShardSet::open_published(&cx, &expected, &[first, second]).unwrap();
let before = current.search(&cx, &query(), 2, None).unwrap();
let mut stale = expected.clone();
stale[1].whole_image_sha256[0] ^= 1;
assert!(
matches!(current.try_replace_published(&cx, &stale, &[], &[]),
Err(FsviAdmissionError::Index(SearchError::InvalidConfig { ref field, ref value, .. }))
if field == "native_ann.shards.replace.expected_current" && value == "stale")
);
assert!(
matches!(current.try_replace_published(&cx, &expected[..1], &[], &[]),
Err(FsviAdmissionError::Index(SearchError::InvalidConfig { ref field, ref value, .. }))
if field == "native_ann.shards.replace.expected_current" && value == "cardinality")
);
assert_eq!(current.search(&cx, &query(), 2, None).unwrap(), before);
assert_eq!(
current.owner_witnesses().cloned().collect::<Vec<_>>(),
expected
);
});
}
#[test]
fn rollback_and_same_sequence_nonce_twins_refuse_before_opening_files() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (witness, artifact) = fixture(
&cx,
&root,
"current",
2,
QuantizationFormat::F32,
&[("a", [1.0, 0.0], true)],
false,
);
let expected = [witness];
let mut current = NativeShardSet::open_published(&cx, &expected, &[artifact]).unwrap();
for generation in [
ArtifactGenerationIdentityV1::new(1, [0x73; 16]).unwrap(),
ArtifactGenerationIdentityV1::new(2, [0x74; 16]).unwrap(),
expected[0].generation,
] {
let mut next = expected[0].clone();
next.generation = generation;
assert!(
matches!(current.try_replace_published(&cx, &expected, &[next], &[]),
Err(FsviAdmissionError::Index(SearchError::InvalidConfig { ref field, .. }))
if field == "native_ann.shards.replace.generation")
);
}
assert_eq!(
current.owner_witnesses().cloned().collect::<Vec<_>>(),
expected
);
assert_eq!(
current.search(&cx, &query(), 1, None).unwrap()[0].doc_id,
"a"
);
});
}
#[test]
fn a_bad_later_vector_or_graph_never_partially_replaces_the_live_inventory() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (old, artifact) = fixture(
&cx,
&root,
"old",
1,
QuantizationFormat::F32,
&[("old", [1.0, 0.0], true)],
true,
);
let expected = [old];
let mut current = NativeShardSet::open_published(&cx, &expected, &[artifact]).unwrap();
let retained = current.clone();
let (a, first) = fixture(
&cx,
&root,
"next-a",
2,
QuantizationFormat::F32,
&[("a", [0.75, 0.25], true)],
true,
);
let (b, mut second) = fixture(
&cx,
&root,
"next-b",
2,
QuantizationFormat::F32,
&[("b", [0.5, 0.5], true)],
false,
);
second.2 = Some(root.join("invalid-graph.txt"));
let next = [a, b];
let mut artifacts = [first, second];
assert!(
current
.try_replace_published(&cx, &expected, &next, &artifacts)
.is_err()
);
assert!(Arc::ptr_eq(&retained.shards[0], ¤t.shards[0]));
artifacts[1].2 = None;
std::fs::write(&artifacts[1].0, b"corrupted mandatory second vector").unwrap();
assert!(
current
.try_replace_published(&cx, &expected, &next, &artifacts)
.is_err()
);
assert!(Arc::ptr_eq(&retained.shards[0], ¤t.shards[0]));
assert_eq!(
current.owner_witnesses().cloned().collect::<Vec<_>>(),
expected
);
assert_eq!(
current.search(&cx, &query(), 1, None).unwrap()[0].doc_id,
"old"
);
});
}
#[test]
fn a_prepared_candidate_cannot_overwrite_an_intervening_refresh() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (old, first) = fixture(
&cx,
&root,
"gen-1",
1,
QuantizationFormat::F32,
&[("old", [1.0, 0.0], true)],
false,
);
let expected_old = [old];
let mut current = NativeShardSet::open_published(&cx, &expected_old, &[first]).unwrap();
let (second, a) = fixture(
&cx,
&root,
"gen-2",
2,
QuantizationFormat::F32,
&[("second", [0.75, 0.25], true)],
false,
);
let (third, b) = fixture(
&cx,
&root,
"gen-3",
3,
QuantizationFormat::F32,
&[("third", [0.5, 0.5], true)],
false,
);
let expected_second = [second];
let prepared_second =
NativeShardSet::open_published(&cx, &expected_second, &[a]).unwrap();
let prepared_third = NativeShardSet::open_published(&cx, &[third], &[b]).unwrap();
current
.try_replace(&cx, &expected_old, prepared_second)
.unwrap();
assert!(
matches!(current.try_replace(&cx, &expected_old, prepared_third.clone()),
Err(SearchError::InvalidConfig { ref field, .. })
if field == "native_ann.shards.replace.expected_current")
);
assert_eq!(
current.search(&cx, &query(), 1, None).unwrap()[0].doc_id,
"second"
);
current
.try_replace(&cx, &expected_second, prepared_third)
.unwrap();
assert_eq!(
current.search(&cx, &query(), 1, None).unwrap()[0].doc_id,
"third"
);
});
}
#[test]
fn independently_valid_new_producer_requires_a_separate_search_handle() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (old, artifact) = fixture(
&cx,
&root,
"old",
1,
QuantizationFormat::F32,
&[("old", [1.0, 0.0], true)],
false,
);
let mut identity = artifact.1.frozen_identity().identity.clone();
identity.producer.backend = "new-producer".to_owned();
identity.validate().unwrap();
let mut query_identity = identity.clone();
"in-memory-f32-v1".clone_into(&mut query_identity.storage.format);
"native-f32-values".clone_into(&mut query_identity.storage.endianness);
let new_query = BoundQueryEmbedding::new(vec![1.0, 0.0], query_identity).unwrap();
let binding = FsviV2IdentityBinding::new(
ArtifactGenerationIdentityV1::new(2, [0x73; 16]).unwrap(),
identity.freeze().unwrap(),
)
.unwrap();
let path = root.join("different-producer.fsvi");
let mut writer = VectorIndex::create_v2(&path, binding.clone()).unwrap();
writer.write_record("new", &[1.0, 0.0]).unwrap();
writer.finish().unwrap();
let bytes: Arc<[u8]> = std::fs::read(&path).unwrap().into();
let witness = ValidatedFsviBytes::from_arc(bytes, &binding)
.unwrap()
.witness()
.clone();
let candidate =
NativeShardSet::open_published(&cx, &[witness], &[(path, binding, None)]).unwrap();
assert_eq!(
candidate.search(&cx, &new_query, 1, None).unwrap()[0].doc_id,
"new"
);
let expected = [old];
let mut current = NativeShardSet::open_published(&cx, &expected, &[artifact]).unwrap();
assert!(matches!(current.try_replace(&cx, &expected, candidate),
Err(SearchError::InvalidConfig { ref field, .. })
if field == "native_ann.shards.replace.identity"));
assert_eq!(
current.search(&cx, &query(), 1, None).unwrap()[0].doc_id,
"old"
);
});
}
#[test]
fn explicit_delete_all_installs_a_new_identity_bearing_empty_generation() {
asupersync::test_utils::run_test_with_cx(|cx| async move {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
let (old, first) = fixture(
&cx,
&root,
"live",
1,
QuantizationFormat::F32,
&[("old", [1.0, 0.0], true)],
false,
);
let expected = [old];
let mut current = NativeShardSet::open_published(&cx, &expected, &[first]).unwrap();
let retained = current.clone();
let (empty, next) =
fixture(&cx, &root, "empty", 2, QuantizationFormat::F16, &[], false);
current
.try_replace_published(&cx, &expected, &[empty], &[next])
.unwrap();
assert_eq!((current.shard_count(), current.live_count()), (1, 0));
assert!(current.search(&cx, &query(), 1, None).unwrap().is_empty());
assert_eq!(
retained.search(&cx, &query(), 1, None).unwrap()[0].doc_id,
"old"
);
let foreign = BoundQueryEmbedding::new(
vec![1.0, 0.0],
EmbeddingIdentityBundleV1::explicit_test_model("foreign", 2),
)
.unwrap();
assert!(current.search(&cx, &foreign, 0, None).is_err());
});
}
}