#![expect(clippy::wildcard_enum_match_arm, reason = "test code")]
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use crate::common::flags::{FeatureFlags, init_feature_flags};
use crate::common::universal_io::{MmapFile, MmapFs};
use crate::segment::common::operation_error::OperationResult;
use crate::segment::data_types::vectors::{VectorInternal, VectorStructInternal};
use crate::segment::types::{Distance, ExtendedPointId, WithPayloadInterface, WithVector};
use crate::shard::files::{SEGMENTS_PATH, segment_manifest_path};
use crate::shard::operations::CollectionUpdateOperations::PointOperation;
use crate::shard::operations::point_ops::PointInsertOperationsInternal::PointsList;
use crate::shard::operations::point_ops::PointOperations::{DeletePoints, UpsertPoints};
use crate::shard::operations::point_ops::{PointStructPersisted, VectorStructPersisted};
use uuid::Uuid;
use crate::edge::config::vectors::EdgeVectorParams;
use crate::edge::edge_shard::scan_segment_dirs;
use crate::edge::read_only::{
LocalSegmentEnumerator, ManifestSegmentEnumerator, ReadOnlyEdgeShard, SegmentEnumerator,
};
use crate::edge::read_view::EdgeShardRead;
use crate::edge::{CountRequest, EdgeConfig, EdgeShard, RetrieveRequestBuilder, ScrollRequestBuilder};
const VECTOR_NAME: &str = "edge-ro-test-vector";
fn test_config() -> EdgeConfig {
EdgeConfig {
on_disk_payload: Some(false),
vectors: HashMap::from([(
VECTOR_NAME.to_string(),
EdgeVectorParams {
size: 1,
distance: Distance::Dot,
quantization_config: None,
multivector_config: None,
datatype: None,
on_disk: None,
hnsw_config: None,
},
)]),
sparse_vectors: HashMap::new(),
hnsw_config: None,
quantization_config: None,
optimizers: None,
wal_options: None,
max_search_threads: None,
search_pool_core: None,
}
}
fn point(id: u64) -> PointStructPersisted {
PointStructPersisted {
id: ExtendedPointId::NumId(id),
vector: VectorStructPersisted::from(VectorStructInternal::Named(HashMap::from([(
VECTOR_NAME.to_string(),
VectorInternal::from(vec![id as f32]),
)]))),
payload: None,
}
}
fn upsert(shard: &EdgeShard, ids: impl IntoIterator<Item = u64>) {
let points = ids.into_iter().map(point).collect::<Vec<_>>();
shard
.update(PointOperation(UpsertPoints(PointsList(points))))
.unwrap();
}
fn delete(shard: &EdgeShard, ids: impl IntoIterator<Item = u64>) {
let ids = ids.into_iter().map(ExtendedPointId::NumId).collect();
shard.update(PointOperation(DeletePoints { ids })).unwrap();
}
fn open_follower(path: &std::path::Path) -> ReadOnlyEdgeShard<MmapFile> {
ReadOnlyEdgeShard::<MmapFile>::open_with_enumerator(
MmapFs,
path,
LocalSegmentEnumerator::new(path),
None,
None,
)
.unwrap()
}
fn exact_count(follower: &ReadOnlyEdgeShard<MmapFile>) -> usize {
follower.count(CountRequest::new()).unwrap()
}
fn leader_exact_count(leader: &EdgeShard) -> usize {
leader.count(CountRequest::new()).unwrap()
}
fn scrolled_ids(follower: &ReadOnlyEdgeShard<MmapFile>) -> Vec<ExtendedPointId> {
let (records, _) = follower
.scroll(
ScrollRequestBuilder::new()
.limit(10_000)
.with_payload(WithPayloadInterface::Bool(false))
.build(),
)
.unwrap();
let mut ids = records.into_iter().map(|r| r.id).collect::<Vec<_>>();
ids.sort_unstable();
ids
}
fn assert_follower_vectors(follower: &ReadOnlyEdgeShard<MmapFile>, ids: &[u64]) {
let point_ids = ids
.iter()
.map(|id| ExtendedPointId::NumId(*id))
.collect::<Vec<_>>();
let results = follower
.retrieve(
RetrieveRequestBuilder::new(point_ids)
.with_payload(WithPayloadInterface::Bool(false))
.with_vector(WithVector::Bool(true))
.build(),
)
.unwrap();
assert_eq!(results.len(), ids.len(), "expected all ids retrievable");
for (result, &expected_id) in results.iter().zip(ids) {
assert_eq!(result.id, ExtendedPointId::NumId(expected_id));
let vectors = match result.vector.as_ref().expect("vector present") {
VectorStructInternal::Named(named) => named,
other => panic!("expected Named vectors, got {other:?}"),
};
let vec = match vectors.get(VECTOR_NAME).expect("vector name exists") {
VectorInternal::Dense(v) => v,
other => panic!("expected Dense vector, got {other:?}"),
};
assert_eq!(
vec,
&vec![expected_id as f32],
"vector mismatch for {expected_id}"
);
}
}
#[test]
fn follower_sees_flushed_data() {
let dir = tempfile::Builder::new()
.prefix("edge-ro-visibility")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=100);
leader.flush().unwrap();
let follower = open_follower(dir.path());
follower.refresh().unwrap();
assert_eq!(exact_count(&follower), 100);
assert_eq!(exact_count(&follower), leader_exact_count(&leader));
assert_eq!(scrolled_ids(&follower).len(), 100);
assert_follower_vectors(&follower, &[1, 50, 100]);
}
#[test]
fn follower_with_load_profile_serves_reads() {
let dir = tempfile::Builder::new()
.prefix("edge-ro-load-profile")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=100);
leader.flush().unwrap();
let scroll_request = ScrollRequestBuilder::new()
.limit(10_000)
.with_payload(WithPayloadInterface::Bool(false))
.build();
let follower = ReadOnlyEdgeShard::<MmapFile>::open_with_enumerator(
MmapFs,
dir.path(),
LocalSegmentEnumerator::new(dir.path()),
None,
Some(scroll_request.load_profile()),
)
.unwrap();
let (records, _) = follower.scroll(scroll_request).unwrap();
assert_eq!(records.len(), 100);
assert_eq!(exact_count(&follower), leader_exact_count(&leader));
assert_follower_vectors(&follower, &[1, 50, 100]);
upsert(&leader, 101..=150);
leader.flush().unwrap();
follower.refresh().unwrap();
assert_eq!(exact_count(&follower), 150);
}
#[test]
fn refresh_picks_up_incremental_writes() {
let dir = tempfile::Builder::new()
.prefix("edge-ro-incremental")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=50);
leader.flush().unwrap();
let follower = open_follower(dir.path());
follower.refresh().unwrap();
assert_eq!(exact_count(&follower), 50);
upsert(&leader, 51..=100);
leader.flush().unwrap();
follower.refresh().unwrap();
assert_eq!(exact_count(&follower), 100);
assert_eq!(exact_count(&follower), leader_exact_count(&leader));
assert_follower_vectors(&follower, &[1, 50, 51, 100]);
}
#[test]
fn follower_reflects_deletes() {
let dir = tempfile::Builder::new()
.prefix("edge-ro-deletes")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=100);
leader.flush().unwrap();
let follower = open_follower(dir.path());
follower.refresh().unwrap();
assert_eq!(exact_count(&follower), 100);
delete(&leader, 1..=40);
leader.flush().unwrap();
follower.refresh().unwrap();
assert_eq!(exact_count(&follower), 60);
assert_eq!(exact_count(&follower), leader_exact_count(&leader));
let ids = scrolled_ids(&follower);
assert_eq!(ids.len(), 60);
assert!(!ids.contains(&ExtendedPointId::NumId(1)));
assert!(ids.contains(&ExtendedPointId::NumId(100)));
}
#[cfg_attr(
target_os = "windows",
ignore = "leader can't delete a segment dir held open (mmap) by the follower"
)]
#[test]
fn follower_tracks_optimization_swap() {
let dir = tempfile::Builder::new()
.prefix("edge-ro-optimize")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=1000);
delete(&leader, 1..=300);
leader.flush().unwrap();
let follower = open_follower(dir.path());
follower.refresh().unwrap();
assert_eq!(exact_count(&follower), 700);
let optimized = leader.optimize().unwrap();
assert!(optimized, "expected a vacuum optimization to run");
leader.flush().unwrap();
follower.refresh().unwrap();
assert_eq!(exact_count(&follower), 700);
assert_eq!(exact_count(&follower), leader_exact_count(&leader));
assert_follower_vectors(&follower, &[301, 500, 1000]);
}
#[test]
fn refresh_on_unchanged_dir_is_noop() {
let dir = tempfile::Builder::new()
.prefix("edge-ro-noop")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=10);
leader.flush().unwrap();
let follower = open_follower(dir.path());
follower.refresh().unwrap();
let before = exact_count(&follower);
follower.refresh().unwrap();
follower.refresh().unwrap();
assert_eq!(exact_count(&follower), before);
assert_eq!(before, 10);
}
#[test]
fn open_without_config_derives_from_segments() {
let dir = tempfile::Builder::new()
.prefix("edge-ro-noconfig")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=10);
leader.flush().unwrap();
let expected = leader_exact_count(&leader);
drop(leader);
fs_err::remove_file(dir.path().join("edge_config.json")).unwrap();
let follower = open_follower(dir.path());
assert_eq!(exact_count(&follower), expected);
assert_eq!(expected, 10);
}
#[test]
fn provided_config_overrides_tunables_at_open() {
let dir = tempfile::Builder::new()
.prefix("edge-ro-provided-config")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=10);
leader.flush().unwrap();
let provided = EdgeConfig {
on_disk_payload: None,
vectors: HashMap::new(),
sparse_vectors: HashMap::new(),
hnsw_config: None,
quantization_config: None,
optimizers: None,
wal_options: None,
max_search_threads: Some(2),
search_pool_core: None,
};
let follower = ReadOnlyEdgeShard::<MmapFile>::open_with_enumerator(
MmapFs,
dir.path(),
LocalSegmentEnumerator::new(dir.path()),
Some(provided),
None,
)
.unwrap();
let config = follower.config_snapshot();
assert!(config.vectors.contains_key(VECTOR_NAME));
assert_eq!(config.on_disk_payload, Some(false));
assert_eq!(config.max_search_threads, Some(2));
assert_eq!(exact_count(&follower), 10);
upsert(&leader, 11..=15);
leader.flush().unwrap();
follower.refresh().unwrap();
let config = follower.config_snapshot();
assert!(config.vectors.contains_key(VECTOR_NAME));
assert_eq!(config.on_disk_payload, Some(false));
assert_eq!(config.max_search_threads, None);
assert_eq!(exact_count(&follower), 15);
}
struct ExcludingEnumerator {
segments_path: PathBuf,
exclude: Uuid,
}
impl SegmentEnumerator for ExcludingEnumerator {
fn list_segments(&self) -> OperationResult<HashMap<Uuid, PathBuf>> {
let mut segments = scan_segment_dirs(&self.segments_path)?;
segments.remove(&self.exclude);
Ok(segments)
}
}
#[test]
fn follower_uses_injected_enumerator() {
let dir = tempfile::Builder::new()
.prefix("edge-ro-enumerator")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=100);
leader.flush().unwrap();
let segments_path = dir.path().join(SEGMENTS_PATH);
let all_segments = scan_segment_dirs(&segments_path).unwrap();
assert!(!all_segments.is_empty());
let hidden = *all_segments.keys().next().unwrap();
let baseline = open_follower(dir.path());
assert_eq!(baseline.segments_count(), all_segments.len());
let follower = ReadOnlyEdgeShard::<MmapFile>::open_with_enumerator(
MmapFs,
dir.path(),
ExcludingEnumerator {
segments_path,
exclude: hidden,
},
None,
None,
)
.unwrap();
assert_eq!(follower.segments_count(), all_segments.len() - 1);
follower.refresh().unwrap();
assert_eq!(follower.segments_count(), all_segments.len() - 1);
}
#[test]
fn manifest_enumerator_requires_manifest() {
let dir = tempfile::Builder::new()
.prefix("edge-ro-manifest")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=100);
leader.flush().unwrap();
let segments_path = dir.path().join(SEGMENTS_PATH);
let on_disk: HashSet<Uuid> = scan_segment_dirs(&segments_path)
.unwrap()
.into_keys()
.collect();
assert!(!on_disk.is_empty());
let enumerator = ManifestSegmentEnumerator::new(MmapFs, dir.path());
let manifest_path = segment_manifest_path(dir.path());
let _ = fs_err::remove_file(&manifest_path);
assert!(enumerator.list_segments().is_err());
let active = *on_disk.iter().next().unwrap();
fs_err::write(&manifest_path, format!(r#"{{"{active}":"active"}}"#)).unwrap();
let from_manifest: HashSet<Uuid> = enumerator.list_segments().unwrap().into_keys().collect();
assert_eq!(from_manifest, HashSet::from([active]));
}
#[test]
fn leader_writes_manifest_and_follower_loads_it() {
let mut flags = FeatureFlags::default();
flags.write_segment_manifest = true;
init_feature_flags(flags);
if !crate::common::flags::feature_flags().write_segment_manifest {
return;
}
let dir = tempfile::Builder::new()
.prefix("edge-ro-manifest-write")
.tempdir()
.unwrap();
let leader = EdgeShard::new(dir.path(), test_config()).unwrap();
upsert(&leader, 1..=100);
leader.flush().unwrap();
let manifest_path = segment_manifest_path(dir.path());
assert!(
manifest_path.exists(),
"leader should write the segment manifest"
);
let on_disk: HashSet<Uuid> = scan_segment_dirs(&dir.path().join(SEGMENTS_PATH))
.unwrap()
.into_keys()
.collect();
let from_manifest: HashSet<Uuid> = ManifestSegmentEnumerator::new(MmapFs, dir.path())
.list_segments()
.unwrap()
.into_keys()
.collect();
assert_eq!(from_manifest, on_disk);
let follower = ReadOnlyEdgeShard::<MmapFile>::open_mmap(dir.path()).unwrap();
follower.refresh().unwrap();
assert_eq!(exact_count(&follower), 100);
}