diff --git a/.gitignore b/.gitignore index fe9d61dd8e..aebb6b4dea 100644 --- a/.gitignore +++ b/.gitignore @@ -21,3 +21,4 @@ venv .env .python-version .claude/ +/docs/plans/ diff --git a/lib/segment/benches/boolean_filtering.rs b/lib/segment/benches/boolean_filtering.rs index 2eb37c6b8c..5901f594ed 100644 --- a/lib/segment/benches/boolean_filtering.rs +++ b/lib/segment/benches/boolean_filtering.rs @@ -47,7 +47,7 @@ pub fn plain_boolean_query_points(c: &mut Criterion) { b.iter(|| { let filter = random_bool_filter(&mut rng); result_size += plain_index - .query_points(&filter, &hw_counter, &is_stopped, None) + .query_points(&filter, &hw_counter, &is_stopped) .unwrap() .len(); query_count += 1; @@ -77,7 +77,7 @@ pub fn struct_boolean_query_points(c: &mut Criterion) { b.iter(|| { let filter = random_bool_filter(&mut rng); result_size += struct_index - .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped)) .unwrap() .len(); query_count += 1; @@ -131,7 +131,7 @@ pub fn keyword_index_boolean_query_points(c: &mut Criterion) { b.iter(|| { let filter = random_bool_filter(&mut rng); result_size += index - .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped)) .unwrap() .len(); query_count += 1; diff --git a/lib/segment/benches/conditional_search.rs b/lib/segment/benches/conditional_search.rs index 0d55b6b1d1..fba63acbf7 100644 --- a/lib/segment/benches/conditional_search.rs +++ b/lib/segment/benches/conditional_search.rs @@ -38,7 +38,7 @@ fn conditional_plain_search_benchmark(c: &mut Criterion) { b.iter(|| { let filter = random_must_filter(&mut rng, 2); result_size += plain_index - .query_points(&filter, &hw_counter, &is_stopped, None) + .query_points(&filter, &hw_counter, &is_stopped) .unwrap() .len(); query_count += 1; @@ -56,7 +56,7 @@ fn conditional_plain_search_benchmark(c: &mut Criterion) { b.iter(|| { let filter = random_must_filter(&mut rng, 1); result_size += plain_index - .query_points(&filter, &hw_counter, &is_stopped, None) + .query_points(&filter, &hw_counter, &is_stopped) .unwrap() .len(); query_count += 1; @@ -156,7 +156,7 @@ fn conditional_struct_search_benchmark(c: &mut Criterion) { b.iter(|| { let filter = random_must_filter(&mut rng, 2); result_size += struct_index - .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped)) .unwrap() .len(); query_count += 1; diff --git a/lib/segment/benches/range_filtering.rs b/lib/segment/benches/range_filtering.rs index 008b9f60eb..5378acef58 100644 --- a/lib/segment/benches/range_filtering.rs +++ b/lib/segment/benches/range_filtering.rs @@ -111,7 +111,7 @@ fn range_filtering(c: &mut Criterion) { || random_range_filter(&mut rng, FLT_KEY), |filter| { result_size += index - .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped)) .unwrap() .len(); query_count += 1; @@ -125,7 +125,7 @@ fn range_filtering(c: &mut Criterion) { || random_range_filter(&mut rng, INT_KEY), |filter| { result_size += index - .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped)) .unwrap() .len(); query_count += 1; @@ -154,7 +154,7 @@ fn range_filtering(c: &mut Criterion) { || random_range_filter(&mut rng, FLT_KEY), |filter| { result_size += index - .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped)) .unwrap() .len(); query_count += 1; @@ -168,7 +168,7 @@ fn range_filtering(c: &mut Criterion) { || random_range_filter(&mut rng, INT_KEY), |filter| { result_size += index - .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped)) .unwrap() .len(); query_count += 1; diff --git a/lib/segment/benches/sparse_index_build.rs b/lib/segment/benches/sparse_index_build.rs index 27d3d5fac9..5db7978d9b 100644 --- a/lib/segment/benches/sparse_index_build.rs +++ b/lib/segment/benches/sparse_index_build.rs @@ -92,7 +92,6 @@ fn sparse_vector_index_build_benchmark(c: &mut Criterion) { path: index_dir.path(), stopped: &stopped, tick_progress: || (), - deferred_internal_id: None, }) .unwrap(); assert_eq!(sparse_vector_index.indexed_vector_count(), NUM_VECTORS); @@ -109,7 +108,6 @@ fn sparse_vector_index_build_benchmark(c: &mut Criterion) { path: index_dir.path(), stopped: &stopped, tick_progress: || (), - deferred_internal_id: None, }) .unwrap(); diff --git a/lib/segment/benches/sparse_index_search.rs b/lib/segment/benches/sparse_index_search.rs index 299c8bd810..10a60d943d 100644 --- a/lib/segment/benches/sparse_index_search.rs +++ b/lib/segment/benches/sparse_index_search.rs @@ -113,7 +113,6 @@ fn sparse_vector_index_search_benchmark_impl( path: mmap_index_dir.path(), stopped: &stopped, tick_progress: || pb.inc(1), - deferred_internal_id: None, }) .unwrap(); pb.finish_and_clear(); diff --git a/lib/segment/benches/vector_search.rs b/lib/segment/benches/vector_search.rs index 3ab588c347..5f72d32188 100644 --- a/lib/segment/benches/vector_search.rs +++ b/lib/segment/benches/vector_search.rs @@ -77,7 +77,7 @@ fn benchmark(c: id_tracker.deleted_point_bitslice(), 10, ) - .peek_top_all(&DEFAULT_STOPPED, None) + .peek_top_all(&DEFAULT_STOPPED) .expect("points scored") }, BatchSize::SmallInput, diff --git a/lib/segment/src/data_types/query_context.rs b/lib/segment/src/data_types/query_context.rs index abdaca987d..f5cbddcbfb 100644 --- a/lib/segment/src/data_types/query_context.rs +++ b/lib/segment/src/data_types/query_context.rs @@ -6,7 +6,7 @@ use common::bitvec::BitSlice; use common::counter::hardware_accumulator::HwMeasurementAcc; use common::counter::hardware_counter::HardwareCounterCell; use common::cow::SimpleCow; -use common::types::{PointOffsetType, ScoreType}; +use common::types::ScoreType; use sparse::common::types::{DimId, DimWeight}; use crate::data_types::tiny_map; @@ -142,11 +142,7 @@ impl<'a> SegmentQueryContext<'a> { self.query_context.available_point_count() } - pub fn get_vector_context( - &self, - vector_name: &VectorName, - deferred_internal_id: Option, - ) -> VectorQueryContext<'_> { + pub fn get_vector_context(&self, vector_name: &VectorName) -> VectorQueryContext<'_> { VectorQueryContext { search_optimized_threshold_kb: self.query_context.search_optimized_threshold_kb, is_stopped: Some(&self.query_context.is_stopped), @@ -159,7 +155,6 @@ impl<'a> SegmentQueryContext<'a> { .copied(), deleted_points: self.deleted_points, hardware_counter: self.hardware_counter.fork(), - deferred_internal_id, } } @@ -197,8 +192,6 @@ pub struct VectorQueryContext<'a> { deleted_points: Option<&'a BitSlice>, hardware_counter: HardwareCounterCell, - - deferred_internal_id: Option, } impl VectorQueryContext<'_> { @@ -249,10 +242,6 @@ impl VectorQueryContext<'_> { pub fn is_require_idf(&self) -> bool { self.idf.is_some() && self.indexed_vectors.is_some() } - - pub fn deferred_internal_id(&self) -> Option { - self.deferred_internal_id - } } #[cfg(feature = "testing")] @@ -265,7 +254,6 @@ impl Default for VectorQueryContext<'_> { indexed_vectors: None, deleted_points: None, hardware_counter: HardwareCounterCell::new(), - deferred_internal_id: None, } } } diff --git a/lib/segment/src/fixtures/sparse_fixtures.rs b/lib/segment/src/fixtures/sparse_fixtures.rs index 71205dd138..871fdb44d4 100644 --- a/lib/segment/src/fixtures/sparse_fixtures.rs +++ b/lib/segment/src/fixtures/sparse_fixtures.rs @@ -83,7 +83,6 @@ pub fn fixture_sparse_index_from_iter( path: index_dir, stopped: &stopped, tick_progress: || (), - deferred_internal_id: None, })?; assert_eq!( diff --git a/lib/segment/src/id_tracker/id_tracker_base/point_mappings_ref.rs b/lib/segment/src/id_tracker/id_tracker_base/point_mappings_ref.rs index 4b0a2df483..2a10c4cf42 100644 --- a/lib/segment/src/id_tracker/id_tracker_base/point_mappings_ref.rs +++ b/lib/segment/src/id_tracker/id_tracker_base/point_mappings_ref.rs @@ -1,6 +1,7 @@ use atomic_refcell::AtomicRef; use common::bitvec::{BitSlice, BitSliceExt as _}; -use common::types::PointOffsetType; +use common::types::{DeferredBehavior, PointOffsetType}; +use itertools::Either; use self_cell::self_cell; use super::tracker_enum::IdTrackerEnum; @@ -77,12 +78,10 @@ impl<'a> PointMappingsRefEnum<'a> { ) } - /// Iterate over all internal IDs. Optionally filter all deferred points. - pub fn iter_internal_visible( - &self, - deferred_internal_id: Option, - ) -> Box + '_> { - match deferred_internal_id { + /// Iterate over all internal IDs, filtering deferred points using the + /// mapping's own threshold. + pub fn iter_internal_visible(self) -> Box + 'a> { + match self.deferred_internal_id() { None => self.iter_internal(), Some(deferred_internal_id) => Box::new( self.iter_internal() @@ -91,13 +90,50 @@ impl<'a> PointMappingsRefEnum<'a> { } } - /// Iterate starting from a given ID. Optionally filter all deferred points. + /// Iterate over all internal IDs, with deferred filtering selected by + /// `deferred_behavior`: + /// - [`DeferredBehavior::Exclude`] applies the mapping's own threshold; + /// - [`DeferredBehavior::IncludeAll`] yields every point regardless of the + /// threshold. + pub fn iter_internal_with_behavior( + self, + deferred_behavior: DeferredBehavior, + ) -> Box + 'a> { + if deferred_behavior.include_all_points() { + self.iter_internal() + } else { + self.iter_internal_visible() + } + } + + /// Wrap an iterator of internal IDs so that points at or above the + /// mapping's deferred threshold are excluded. + /// + /// For [`DeferredBehavior::IncludeAll`] — or when the mapping has no + /// deferred threshold — the iterator is returned unchanged. Intended for + /// iterators sourced outside the mapping (e.g., field-index outputs) where + /// the threshold isn't applied implicitly. + pub fn filter_deferred( + self, + iter: I, + deferred_behavior: DeferredBehavior, + ) -> impl Iterator + where + I: Iterator, + { + match deferred_behavior.apply(self.deferred_internal_id()) { + None => Either::Left(iter), + Some(cutoff) => Either::Right(iter.filter(move |&id| id < cutoff)), + } + } + + /// Iterate starting from a given ID, filtering deferred points using the + /// mapping's own threshold. pub fn iter_from_visible( - &self, + self, external_id: Option, - deferred_internal_id: Option, - ) -> Box + '_> { - match deferred_internal_id { + ) -> Box + 'a> { + match self.deferred_internal_id() { None => self.iter_from(external_id), Some(deferred_internal_id) => Box::new( self.iter_from(external_id) @@ -106,11 +142,12 @@ impl<'a> PointMappingsRefEnum<'a> { } } + /// Iterate over internal IDs in random order, filtering deferred points + /// using the mapping's own threshold. pub fn iter_random_visible( - &self, - deferred_internal_id: Option, - ) -> Box + '_> { - match deferred_internal_id { + self, + ) -> Box + 'a> { + match self.deferred_internal_id() { None => self.iter_random(), Some(deferred_internal_id) => Box::new( self.iter_random() @@ -120,6 +157,21 @@ impl<'a> PointMappingsRefEnum<'a> { ), } } + + /// Deferred threshold attached to this mapping, if any. + /// + /// Compressed mappings (used by immutable / read-only-immutable id trackers) + /// never carry a deferred threshold. + /// + /// Kept private so callers go through the dispatch helpers + /// ([`Self::iter_internal_with_behavior`], [`Self::external_iter_cutoff`]) + /// instead of leaking the raw threshold. + fn deferred_internal_id(self) -> Option { + match self { + PointMappingsRefEnum::Plain(m) => m.deferred_internal_id(), + PointMappingsRefEnum::Compressed(_) => None, + } + } } self_cell! { diff --git a/lib/segment/src/id_tracker/id_tracker_base/read_only_tracker_enum.rs b/lib/segment/src/id_tracker/id_tracker_base/read_only_tracker_enum.rs index 05122751b9..0bcbe07060 100644 --- a/lib/segment/src/id_tracker/id_tracker_base/read_only_tracker_enum.rs +++ b/lib/segment/src/id_tracker/id_tracker_base/read_only_tracker_enum.rs @@ -100,4 +100,18 @@ impl IdTrackerRead for ReadOnlyIdTrackerEnum { ReadOnlyIdTrackerEnum::Immutable(id_tracker) => id_tracker.iter_internal_versions(), } } + + fn deferred_internal_id(&self) -> Option { + match self { + ReadOnlyIdTrackerEnum::Appendable(id_tracker) => id_tracker.deferred_internal_id(), + ReadOnlyIdTrackerEnum::Immutable(id_tracker) => id_tracker.deferred_internal_id(), + } + } + + fn deferred_deleted_count(&self) -> usize { + match self { + ReadOnlyIdTrackerEnum::Appendable(id_tracker) => id_tracker.deferred_deleted_count(), + ReadOnlyIdTrackerEnum::Immutable(id_tracker) => id_tracker.deferred_deleted_count(), + } + } } diff --git a/lib/segment/src/id_tracker/id_tracker_base/tracker_enum.rs b/lib/segment/src/id_tracker/id_tracker_base/tracker_enum.rs index 13472893d0..8bb0530832 100644 --- a/lib/segment/src/id_tracker/id_tracker_base/tracker_enum.rs +++ b/lib/segment/src/id_tracker/id_tracker_base/tracker_enum.rs @@ -111,6 +111,22 @@ impl IdTrackerRead for IdTrackerEnum { IdTrackerEnum::InMemoryIdTracker(id_tracker) => id_tracker.iter_internal_versions(), } } + + fn deferred_internal_id(&self) -> Option { + match self { + IdTrackerEnum::MutableIdTracker(id_tracker) => id_tracker.deferred_internal_id(), + IdTrackerEnum::ImmutableIdTracker(id_tracker) => id_tracker.deferred_internal_id(), + IdTrackerEnum::InMemoryIdTracker(id_tracker) => id_tracker.deferred_internal_id(), + } + } + + fn deferred_deleted_count(&self) -> usize { + match self { + IdTrackerEnum::MutableIdTracker(id_tracker) => id_tracker.deferred_deleted_count(), + IdTrackerEnum::ImmutableIdTracker(id_tracker) => id_tracker.deferred_deleted_count(), + IdTrackerEnum::InMemoryIdTracker(id_tracker) => id_tracker.deferred_deleted_count(), + } + } } impl IdTracker for IdTrackerEnum { diff --git a/lib/segment/src/id_tracker/id_tracker_base/trait_def.rs b/lib/segment/src/id_tracker/id_tracker_base/trait_def.rs index 197aeb1525..074d199f19 100644 --- a/lib/segment/src/id_tracker/id_tracker_base/trait_def.rs +++ b/lib/segment/src/id_tracker/id_tracker_base/trait_def.rs @@ -2,7 +2,7 @@ use std::fmt; use std::path::PathBuf; use common::bitvec::{BitSlice, BitSliceExt as _}; -use common::types::PointOffsetType; +use common::types::{DeferredBehavior, PointOffsetType}; use rand::rngs::StdRng; use rand::{RngExt, SeedableRng}; @@ -174,4 +174,53 @@ pub trait IdTrackerRead { fn iter_internal_versions( &self, ) -> Box + '_>; + + /// Internal-id threshold above which points are hidden from reads. + /// + /// Only appendable trackers can carry a non-`None` value. + fn deferred_internal_id(&self) -> Option { + None + } + + /// Number of soft-deleted points at or above the deferred threshold. + fn deferred_deleted_count(&self) -> usize { + 0 + } + + /// Translate external point ids into two parallel vectors of `(ids, + /// offsets)` in a single pass. + /// + /// Applies deferred filtering according to `deferred_behavior` inline + /// (no separate `point_is_deferred` lookup), missing points will be ignored + /// + /// The parallel-vector return shape lets downstream batched fetchers + /// consume `&offsets` directly. + /// + /// Centralising this here keeps the deferred threshold from leaking out + /// of the id tracker — callers go through this entry point instead of + /// reading `deferred_internal_id()` themselves. + fn resolve_external_ids( + &self, + point_ids: &[PointIdType], + deferred_behavior: DeferredBehavior, + ) -> (Vec, Vec) { + // Non-appendable trackers never carry a deferred threshold (it's set + // only via `MutableIdTracker::open`, guarded by `appendable_flag` in + // the segment constructor), so we don't need an appendable check here. + let deferred_cutoff = deferred_behavior.apply(self.deferred_internal_id()); + + let mut ids = Vec::with_capacity(point_ids.len()); + let mut offsets = Vec::with_capacity(point_ids.len()); + for &point_id in point_ids { + let Some(internal_id) = self.internal_id(point_id) else { + continue; + }; + if deferred_cutoff.is_some_and(|cutoff| internal_id >= cutoff) { + continue; + } + ids.push(point_id); + offsets.push(internal_id); + } + (ids, offsets) + } } diff --git a/lib/segment/src/id_tracker/in_memory_id_tracker.rs b/lib/segment/src/id_tracker/in_memory_id_tracker.rs index ff64da774e..70310e15c5 100644 --- a/lib/segment/src/id_tracker/in_memory_id_tracker.rs +++ b/lib/segment/src/id_tracker/in_memory_id_tracker.rs @@ -100,6 +100,14 @@ impl IdTrackerRead for InMemoryIdTracker { "in memory id tracker" } + fn deferred_internal_id(&self) -> Option { + self.mappings.deferred_internal_id() + } + + fn deferred_deleted_count(&self) -> usize { + self.mappings.deferred_deleted_count() + } + fn iter_internal_versions( &self, ) -> Box + '_> { diff --git a/lib/segment/src/id_tracker/mutable_id_tracker/mappings_storage.rs b/lib/segment/src/id_tracker/mutable_id_tracker/mappings_storage.rs index 0bc08e1d24..0ea2093231 100644 --- a/lib/segment/src/id_tracker/mutable_id_tracker/mappings_storage.rs +++ b/lib/segment/src/id_tracker/mutable_id_tracker/mappings_storage.rs @@ -128,12 +128,15 @@ fn write_mapping_changes( /// Returns loaded point mappings and the number of bytes read from the file. /// /// If the file ends with an incomplete entry, it is truncated from the file. -pub(super) fn load_mappings(mappings_path: &Path) -> OperationResult<(PointMappings, u64)> { +pub(super) fn load_mappings( + mappings_path: &Path, + deferred_internal_id: Option, +) -> OperationResult<(PointMappings, u64)> { let file = OneshotFile::open(mappings_path)?; let file_len = file.metadata()?.len(); let mut reader = BufReader::new(file); - let mappings = read_mappings(&mut reader)?; + let mappings = read_mappings(&mut reader, deferred_internal_id)?; let read_to = reader.stream_position()?; reader.into_inner().drop_cache()?; @@ -198,7 +201,10 @@ where /// Read point mappings from the given reader /// /// Returns loaded point mappings. -pub(super) fn read_mappings(reader: R) -> OperationResult +pub(super) fn read_mappings( + reader: R, + deferred_internal_id: Option, +) -> OperationResult where R: Read + Seek, { @@ -285,6 +291,7 @@ where internal_to_external, external_to_internal_num, external_to_internal_uuid, + deferred_internal_id, ); Ok(mappings) diff --git a/lib/segment/src/id_tracker/mutable_id_tracker/mod.rs b/lib/segment/src/id_tracker/mutable_id_tracker/mod.rs index 4105150e4a..0ffb6ffcf4 100644 --- a/lib/segment/src/id_tracker/mutable_id_tracker/mod.rs +++ b/lib/segment/src/id_tracker/mutable_id_tracker/mod.rs @@ -69,7 +69,10 @@ pub struct MutableIdTracker { } impl MutableIdTracker { - pub fn open(segment_path: impl Into) -> OperationResult { + pub fn open( + segment_path: impl Into, + deferred_internal_id: Option, + ) -> OperationResult { let segment_path = segment_path.into(); let (mappings_path, versions_path) = @@ -93,11 +96,18 @@ impl MutableIdTracker { } let (mappings, mappings_expected_len) = if has_mappings { - load_mappings(&mappings_path).map_err(|err| { + load_mappings(&mappings_path, deferred_internal_id).map_err(|err| { OperationError::service_error(format!("Failed to load ID tracker mappings: {err}")) })? } else { - (PointMappings::default(), 0) + let mappings = PointMappings::new( + Default::default(), + Default::default(), + Default::default(), + Default::default(), + deferred_internal_id, + ); + (mappings, 0) }; let internal_to_version = if has_versions { @@ -210,6 +220,14 @@ impl IdTrackerRead for MutableIdTracker { fn name(&self) -> &'static str { "mutable id tracker" } + + fn deferred_internal_id(&self) -> Option { + self.mappings.deferred_internal_id() + } + + fn deferred_deleted_count(&self) -> usize { + self.mappings.deferred_deleted_count() + } } impl IdTracker for MutableIdTracker { diff --git a/lib/segment/src/id_tracker/mutable_id_tracker/read_only/id_tracker_read.rs b/lib/segment/src/id_tracker/mutable_id_tracker/read_only/id_tracker_read.rs index 32ad81fb9a..313ef7a72d 100644 --- a/lib/segment/src/id_tracker/mutable_id_tracker/read_only/id_tracker_read.rs +++ b/lib/segment/src/id_tracker/mutable_id_tracker/read_only/id_tracker_read.rs @@ -46,6 +46,14 @@ impl IdTrackerRead for ReadOnlyAppendableIdTracker { "read-only appendable id tracker" } + fn deferred_internal_id(&self) -> Option { + self.mappings.deferred_internal_id() + } + + fn deferred_deleted_count(&self) -> usize { + self.mappings.deferred_deleted_count() + } + fn iter_internal_versions( &self, ) -> Box + '_> { diff --git a/lib/segment/src/id_tracker/mutable_id_tracker/tests.rs b/lib/segment/src/id_tracker/mutable_id_tracker/tests.rs index 15d582dac7..608fb9cd75 100644 --- a/lib/segment/src/id_tracker/mutable_id_tracker/tests.rs +++ b/lib/segment/src/id_tracker/mutable_id_tracker/tests.rs @@ -27,7 +27,7 @@ const DEFAULT_VERSION: SeqNumberType = 42; fn test_iterator() { let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap(); - let mut id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap(); + let mut id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap(); id_tracker.set_link(200.into(), 0).unwrap(); id_tracker.set_link(100.into(), 1).unwrap(); @@ -97,7 +97,7 @@ fn test_load_store() { (id_tracker.mappings, id_tracker.internal_to_version) }; - let mut loaded_id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap(); + let mut loaded_id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap(); assert_eq!( old_versions.len(), @@ -155,7 +155,7 @@ fn test_store_load_mutated() { (dropped_points, custom_version) }; - let id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap(); + let id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap(); for (index, point) in TEST_POINTS.iter().enumerate() { let internal_id = index as PointOffsetType; @@ -251,7 +251,7 @@ fn test_point_deletion_persists_reload() { }; // Point should still be gone - let id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap(); + let id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap(); assert_eq!(id_tracker.internal_id(point_to_delete), None); old_mappings @@ -297,28 +297,28 @@ fn test_point_mappings_de_serialization_single() { fn test_point_mappings_deserializing_special() { // Empty reader creates empty mappings let buf = Cursor::new(b""); - assert_eq!(read_mappings(buf).unwrap().total_point_count(), 0); + assert_eq!(read_mappings(buf, None).unwrap().total_point_count(), 0); // Corrupt if reading invalid type byte let buf = Cursor::new(b"\x00"); assert!( - read_mappings(buf) + read_mappings(buf, None) .unwrap_err() .to_string() .contains("Corrupted ID tracker mapping storage") ); let buf = Cursor::new(b"malformed!"); - assert!(read_mappings(buf).is_err()); + assert!(read_mappings(buf, None).is_err()); // Empty if change is not fully written let buf = Cursor::new(b"\x01\x01\x00\x00\x00\x00"); - assert_eq!(read_mappings(buf).unwrap().total_point_count(), 0); + assert_eq!(read_mappings(buf, None).unwrap().total_point_count(), 0); // Exactly one entry let buf = Cursor::new(b"\x01\x01\x00\x00\x00\x00\x00\x00\x00\x02\x00\x00\x00"); assert_eq!( - read_mappings(buf) + read_mappings(buf, None) .unwrap() .internal_id(&PointIdType::NumId(1)), Some(2) @@ -327,7 +327,7 @@ fn test_point_mappings_deserializing_special() { // Exactly one entry and a malformed second one let buf = Cursor::new(b"\x01\x01\x00\x00\x00\x00\x00\x00\x00\x02\x00\x00\x00\x00\x01\x00"); assert!( - read_mappings(buf) + read_mappings(buf, None) .unwrap_err() .to_string() .contains("Corrupted ID tracker mapping storage") @@ -335,7 +335,7 @@ fn test_point_mappings_deserializing_special() { // Exactly one entry and an incomplete second one let buf = Cursor::new(b"\x01\x01\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x03\x01\x00"); - let mappings = read_mappings(buf).unwrap(); + let mappings = read_mappings(buf, None).unwrap(); assert_eq!(mappings.total_point_count(), 1); assert_eq!(mappings.internal_id(&PointIdType::NumId(1)), Some(0)); } @@ -397,7 +397,7 @@ fn test_point_mappings_truncation() { .unwrap(); assert_eq!(fs::metadata(&mappings_path).unwrap().len(), 13); assert_eq!( - load_mappings(&mappings_path) + load_mappings(&mappings_path, None) .unwrap() .0 .internal_id(&PointIdType::NumId(1)), @@ -413,7 +413,7 @@ fn test_point_mappings_truncation() { .unwrap(); assert_eq!(fs::metadata(&mappings_path).unwrap().len(), 14); assert_eq!( - load_mappings(&mappings_path) + load_mappings(&mappings_path, None) .unwrap() .0 .internal_id(&PointIdType::NumId(1)), @@ -429,7 +429,7 @@ fn test_point_mappings_truncation() { .unwrap(); assert_eq!(fs::metadata(&mappings_path).unwrap().len(), 16); assert_eq!( - load_mappings(&mappings_path) + load_mappings(&mappings_path, None) .unwrap() .0 .internal_id(&PointIdType::NumId(1)), @@ -444,7 +444,7 @@ fn test_point_mappings_truncation() { ).unwrap(); assert_eq!(fs::metadata(&mappings_path).unwrap().len(), 28); assert_eq!( - load_mappings(&mappings_path) + load_mappings(&mappings_path, None) .unwrap() .0 .internal_id(&PointIdType::NumId(1)), @@ -468,7 +468,8 @@ fn make_in_memory_tracker_from_memory() -> InMemoryIdTracker { } fn make_mutable_tracker(path: &Path) -> MutableIdTracker { - let mut id_tracker = MutableIdTracker::open(path).expect("failed to open mutable ID tracker"); + let mut id_tracker = + MutableIdTracker::open(path, None).expect("failed to open mutable ID tracker"); for value in TEST_POINTS.iter() { let internal_id = id_tracker.total_point_count() as PointOffsetType; @@ -537,7 +538,7 @@ fn simple_id_tracker_vs_mutable_tracker_congruence() { let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap(); let db = open_db(segment_dir.path(), &[DB_VECTOR_CF]).unwrap(); - let mut mutable_id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap(); + let mut mutable_id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap(); let mut simple_id_tracker = SimpleIdTracker::open(db).unwrap(); // Insert 100 random points into id_tracker @@ -612,7 +613,7 @@ fn simple_id_tracker_vs_mutable_tracker_congruence() { mutable_id_tracker.mapping_flusher()().unwrap(); mutable_id_tracker.versions_flusher()().unwrap(); drop(mutable_id_tracker); - let mutable_id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap(); + let mutable_id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap(); check_trackers(&simple_id_tracker, &mutable_id_tracker); } diff --git a/lib/segment/src/id_tracker/point_mappings.rs b/lib/segment/src/id_tracker/point_mappings.rs index e894e965a6..fbd34dbf4d 100644 --- a/lib/segment/src/id_tracker/point_mappings.rs +++ b/lib/segment/src/id_tracker/point_mappings.rs @@ -35,6 +35,14 @@ pub struct PointMappings { // Having two separate maps allows us iterating only over one type at a time without having to filter. external_to_internal_num: BTreeMap, external_to_internal_uuid: BTreeMap, + + /// Points with internal id >= this value are hidden from reads. + /// Only set for appendable segments with deferred points. + deferred_internal_id: Option, + + /// Number of deleted deferred points. Maintained incrementally so we can + /// derive the visible deferred count without re-scanning the deleted bitslice. + deferred_deleted_count: usize, } impl PointMappings { @@ -43,12 +51,25 @@ impl PointMappings { internal_to_external: Vec, external_to_internal_num: BTreeMap, external_to_internal_uuid: BTreeMap, + deferred_internal_id: Option, ) -> Self { + let deferred_deleted_count = deferred_internal_id + .map(|deferred_from| { + let total = deleted.len(); + if total <= deferred_from as usize { + 0 + } else { + deleted[deferred_from as usize..total].count_ones() + } + }) + .unwrap_or(0); Self { deleted, internal_to_external, external_to_internal_num, external_to_internal_uuid, + deferred_internal_id, + deferred_deleted_count, } } @@ -109,7 +130,22 @@ impl PointMappings { } if let Some(internal_id) = &internal_id { + let was_already_deleted = *self + .deleted + .get(*internal_id as usize) + .as_deref() + .unwrap_or(&true); self.deleted.set(*internal_id as usize, true); + + // Count newly-deleted deferred points so we can report visible deferred totals + // without rescanning the deleted bitslice. + if !was_already_deleted + && self + .deferred_internal_id + .is_some_and(|deferred_from| *internal_id >= deferred_from) + { + self.deferred_deleted_count += 1; + } } internal_id @@ -268,6 +304,14 @@ impl PointMappings { self.internal_to_external.len() } + pub(crate) fn deferred_internal_id(&self) -> Option { + self.deferred_internal_id + } + + pub(crate) fn deferred_deleted_count(&self) -> usize { + self.deferred_deleted_count + } + /// Generate a random [`PointMappings`]. #[cfg(test)] pub fn random(rand: &mut StdRng, total_size: u32) -> Self { @@ -326,6 +370,8 @@ impl PointMappings { internal_to_external, external_to_internal_num, external_to_internal_uuid, + deferred_internal_id: None, + deferred_deleted_count: 0, } } @@ -349,6 +395,8 @@ impl PointMappings { internal_to_external, external_to_internal_num, external_to_internal_uuid, + deferred_internal_id: _, + deferred_deleted_count: _, } = self; let deleted_bytes = deleted.capacity().div_ceil(u8::BITS as usize); diff --git a/lib/segment/src/index/hnsw_index/hnsw/build.rs b/lib/segment/src/index/hnsw_index/hnsw/build.rs index 9c4c7e308c..9d216907c1 100644 --- a/lib/segment/src/index/hnsw_index/hnsw/build.rs +++ b/lib/segment/src/index/hnsw_index/hnsw/build.rs @@ -8,7 +8,7 @@ use common::cow::BoxCow; #[cfg(target_os = "linux")] use common::cpu::linux_low_thread_priority; use common::progress_tracker::ProgressTracker; -use common::types::PointOffsetType; +use common::types::{DeferredBehavior, PointOffsetType}; use fs_err as fs; use log::{debug, trace}; use rand::Rng; @@ -459,7 +459,6 @@ impl HNSWIndex { let points_to_index = condition_points( payload_block.condition, - id_tracker_ref.deref(), &payload_index_ref, &vector_storage_ref, stopped, @@ -601,7 +600,6 @@ impl HNSWIndex { /// Get list of points for indexing, associated with payload block filtering condition fn condition_points( condition: FieldCondition, - id_tracker: &IdTrackerEnum, payload_index: &StructPayloadIndex, vector_storage: &VectorStorageEnum, stopped: &AtomicBool, @@ -612,18 +610,14 @@ fn condition_points( let deleted_bitslice = vector_storage.deleted_vector_bitslice(); - let point_mappings = id_tracker.point_mappings(); - payload_index.with_view(|v| { let cardinality_estimation = v.estimate_cardinality(&filter, &disposed_hw_counter)?; Ok(v.iter_filtered_points( &filter, - id_tracker, - &point_mappings, &cardinality_estimation, &disposed_hw_counter, stopped, - None, + DeferredBehavior::IncludeAll, )? .filter(|&point_id| !deleted_bitslice.get_bit(point_id as usize).unwrap_or(false)) .collect()) diff --git a/lib/segment/src/index/hnsw_index/hnsw/search.rs b/lib/segment/src/index/hnsw_index/hnsw/search.rs index b3d8117189..e5dc6bc062 100644 --- a/lib/segment/src/index/hnsw_index/hnsw/search.rs +++ b/lib/segment/src/index/hnsw_index/hnsw/search.rs @@ -1,7 +1,7 @@ use common::bitvec::BitSlice; use common::counter::hardware_counter::HardwareCounterCell; use common::cow::BoxCow; -use common::types::{PointOffsetType, ScoredPointOffset}; +use common::types::{DeferredBehavior, PointOffsetType, ScoredPointOffset}; use super::HNSWIndex; use crate::common::operation_error::OperationResult; @@ -298,21 +298,17 @@ impl HNSWIndex { let hw_counter = &vector_query_context.hardware_counter(); let is_stopped = &vector_query_context.is_stopped(); - let id_tracker = self.id_tracker.borrow(); let payload_index = self.payload_index.borrow(); - let point_mappings = id_tracker.point_mappings(); // Assume query is already estimated to be small enough so we can iterate over all matched ids let filtered_points: Vec = payload_index.with_view(|v| { let query_cardinality = v.estimate_cardinality(filter, hw_counter)?; v.iter_filtered_points( filter, - &*id_tracker, - &point_mappings, &query_cardinality, hw_counter, is_stopped, // No deferred filtering here since it's HNSW index. - None, + DeferredBehavior::IncludeAll, ) .map(|it| it.collect()) })?; diff --git a/lib/segment/src/index/hnsw_index/point_scorer.rs b/lib/segment/src/index/hnsw_index/point_scorer.rs index 193bf96457..4427dcf69e 100644 --- a/lib/segment/src/index/hnsw_index/point_scorer.rs +++ b/lib/segment/src/index/hnsw_index/point_scorer.rs @@ -347,21 +347,34 @@ impl<'a> BatchFilteredSearcher<'a> { } } - pub fn peek_top_all( - self, - is_stopped: &AtomicBool, - deferred_internal_id: Option, - ) -> OperationResult>> { - let iter = self - .filters + /// Iterator over every internal point ID that isn't soft-deleted in this + /// searcher's `point_deleted` bitslice. + /// + /// Does not apply deferred-point filtering — wrap with + /// `PointMappingsRefEnum::filter_deferred` (or compose otherwise) before + /// passing to [`Self::peek_top_iter`] when deferred awareness is needed. + /// + /// The returned iterator borrows the underlying bitslice (lifetime `'a`), + /// independent of `&self`, so it can be composed and then passed into + /// `peek_top_iter(self, ...)` which consumes the searcher. + pub fn iter_not_deleted(&self) -> impl Iterator + 'a { + self.filters .point_deleted .iter_zeros() .map(|p| p as PointOffsetType) - .take_while(|&point_id| { - // Early exit if we hit the max point ID (e.g. a deferred point). - point_id < deferred_internal_id.unwrap_or(PointOffsetType::MAX) - }); + } + /// Score every non-deleted point without deferred filtering. + /// + /// Production paths compose `iter_not_deleted` with + /// `PointMappingsRefEnum::filter_deferred` and call + /// [`Self::peek_top_iter`] directly. + #[cfg(feature = "testing")] + pub fn peek_top_all( + self, + is_stopped: &AtomicBool, + ) -> OperationResult>> { + let iter = self.iter_not_deleted(); self.peek_top_iter(iter, is_stopped) } diff --git a/lib/segment/src/index/payload_index_base.rs b/lib/segment/src/index/payload_index_base.rs index 9bed35405c..9014c67f85 100644 --- a/lib/segment/src/index/payload_index_base.rs +++ b/lib/segment/src/index/payload_index_base.rs @@ -4,7 +4,7 @@ use std::sync::atomic::AtomicBool; use ahash::AHashMap; use common::counter::hardware_counter::HardwareCounterCell; -use common::types::{PointOffsetType, ScoreType}; +use common::types::{DeferredBehavior, PointOffsetType, ScoreType}; use serde_json::Value; use super::field_index::numeric_index::NumericFieldIndexRead; @@ -13,7 +13,6 @@ use super::query_optimization::rescore_formula::FormulaScorer; use super::query_optimization::rescore_formula::parsed_formula::ParsedFormula; use crate::common::Flusher; use crate::common::operation_error::OperationResult; -use crate::id_tracker::{IdTrackerRead, PointMappingsRefEnum}; use crate::index::field_index::{CardinalityEstimation, PayloadBlockCondition}; use crate::json_path::JsonPath; use crate::payload_storage::FilterContext; @@ -66,7 +65,6 @@ pub trait PayloadIndexRead { filter: &Filter, hw_counter: &HardwareCounterCell, is_stopped: &AtomicBool, - deferred_internal_id: Option, ) -> OperationResult>; /// Return number of points, indexed by this field @@ -106,19 +104,16 @@ pub trait PayloadIndexRead { /// Iterate point offsets that match the filter. /// - /// Generic over `I: IdTrackerRead` so callers pass their concrete tracker - /// without dynamic dispatch; the iterator return uses RPITIT so each impl - /// keeps its own zero-cost concrete chain. - #[allow(clippy::too_many_arguments)] - fn iter_filtered_points<'a, I: IdTrackerRead>( + /// The iterator return uses RPITIT so each impl keeps its own zero-cost + /// concrete chain. The id tracker is read from `&self`, so impls reach it + /// through their own field rather than receiving a separate parameter. + fn iter_filtered_points<'a>( &'a self, filter: &'a Filter, - id_tracker: &'a I, - point_mappings: &'a PointMappingsRefEnum<'a>, query_cardinality: &'a CardinalityEstimation, hw_counter: &'a HardwareCounterCell, is_stopped: &'a AtomicBool, - deferred_internal_id: Option, + deferred_behavior: DeferredBehavior, ) -> OperationResult + 'a>; /// Iterate conditions for payload blocks with minimum size of `threshold` diff --git a/lib/segment/src/index/plain_payload_index.rs b/lib/segment/src/index/plain_payload_index.rs index bc93aa3830..44fbcd18cc 100644 --- a/lib/segment/src/index/plain_payload_index.rs +++ b/lib/segment/src/index/plain_payload_index.rs @@ -7,7 +7,7 @@ use ahash::AHashMap; use atomic_refcell::AtomicRefCell; use common::counter::hardware_counter::HardwareCounterCell; use common::iterator_ext::IteratorExt; -use common::types::{PointOffsetType, ScoreType}; +use common::types::{DeferredBehavior, PointOffsetType, ScoreType}; use fs_err as fs; use schemars::_serde_json::Value; @@ -15,7 +15,7 @@ use super::field_index::FieldIndex; use super::payload_config::PayloadFieldSchemaWithIndexType; use crate::common::Flusher; use crate::common::operation_error::{OperationError, OperationResult}; -use crate::id_tracker::{IdTrackerEnum, IdTrackerRead, PointMappingsRefEnum}; +use crate::id_tracker::{IdTrackerEnum, IdTrackerRead}; use crate::index::field_index::facet_index::FacetIndexEnum; use crate::index::field_index::numeric_index::{NumericFieldIndex, NumericFieldIndexRead}; use crate::index::field_index::{CardinalityEstimation, FacetIndex, PayloadBlockCondition}; @@ -111,12 +111,11 @@ impl PayloadIndexRead for PlainPayloadIndex { filter: &Filter, hw_counter: &HardwareCounterCell, is_stopped: &AtomicBool, - deferred_internal_id: Option, ) -> OperationResult> { let filter_context = self.filter_context(filter, hw_counter)?; let id_tracker = self.id_tracker.borrow(); let point_mappings = id_tracker.point_mappings(); - let all_points_iter = point_mappings.iter_internal_visible(deferred_internal_id); + let all_points_iter = point_mappings.iter_internal_visible(); Ok(all_points_iter .stop_if(is_stopped) .filter(|id| filter_context.check(*id)) @@ -190,21 +189,26 @@ impl PayloadIndexRead for PlainPayloadIndex { )) } - fn iter_filtered_points<'a, I: IdTrackerRead>( + fn iter_filtered_points<'a>( &'a self, filter: &'a Filter, - _id_tracker: &'a I, - point_mappings: &'a PointMappingsRefEnum<'a>, _query_cardinality: &'a CardinalityEstimation, hw_counter: &'a HardwareCounterCell, is_stopped: &'a AtomicBool, - deferred_internal_id: Option, + deferred_behavior: DeferredBehavior, ) -> OperationResult + 'a> { let filter_context = self.filter_context(filter, hw_counter)?; - let all_points_iter = point_mappings.iter_internal_visible(deferred_internal_id); - Ok(all_points_iter + // `self.id_tracker` is an `Arc>`, so the mapping borrow is + // local; collect eagerly to detach the iterator from the borrow. + let matched: Vec = self + .id_tracker + .borrow() + .point_mappings() + .iter_internal_with_behavior(deferred_behavior) .stop_if(is_stopped) - .filter(move |id| filter_context.check(*id))) + .filter(|id| filter_context.check(*id)) + .collect(); + Ok(matched.into_iter()) } } diff --git a/lib/segment/src/index/plain_vector_index.rs b/lib/segment/src/index/plain_vector_index.rs index 34206ccc82..1eb9c8a770 100644 --- a/lib/segment/src/index/plain_vector_index.rs +++ b/lib/segment/src/index/plain_vector_index.rs @@ -4,7 +4,7 @@ use std::sync::Arc; use atomic_refcell::AtomicRefCell; use common::counter::hardware_counter::HardwareCounterCell; -use common::types::{PointOffsetType, ScoredPointOffset, TelemetryDetail}; +use common::types::{DeferredBehavior, PointOffsetType, ScoredPointOffset, TelemetryDetail}; use parking_lot::Mutex; use sparse::common::types::DimId; @@ -137,16 +137,20 @@ impl VectorIndexRead for PlainVectorIndex { query_context.hardware_counter(), )?; - let deferred_internal_id = query_context.deferred_internal_id(); - let mut search_results = match filter { Some(filter) => { - let filtered_ids_vec = self.payload_index.borrow().with_view(|v| { - v.query_points(filter, &hw_counter, &is_stopped, deferred_internal_id) - })?; + let filtered_ids_vec = self + .payload_index + .borrow() + .with_view(|v| v.query_points(filter, &hw_counter, &is_stopped))?; batch_searcher.peek_top_iter(filtered_ids_vec.iter().copied(), &is_stopped)? } - None => batch_searcher.peek_top_all(&is_stopped, deferred_internal_id)?, + None => { + let iter = id_tracker + .point_mappings() + .filter_deferred(batch_searcher.iter_not_deleted(), DeferredBehavior::Exclude); + batch_searcher.peek_top_iter(iter, &is_stopped)? + } }; for (search_result, query_vector) in search_results.iter_mut().zip(query_vectors) { diff --git a/lib/segment/src/index/sparse_index/sparse_vector_index.rs b/lib/segment/src/index/sparse_index/sparse_vector_index.rs index fceedd442f..7036afd656 100644 --- a/lib/segment/src/index/sparse_index/sparse_vector_index.rs +++ b/lib/segment/src/index/sparse_index/sparse_vector_index.rs @@ -8,7 +8,6 @@ use atomic_refcell::AtomicRefCell; use common::counter::hardware_counter::HardwareCounterCell; use common::generic_consts::Random; use common::storage_version::StorageVersion as _; -use common::types::PointOffsetType; use fs_err as fs; use sparse::common::scores_memory_pool::ScoresMemoryPool; use sparse::common::sparse_vector::SparseVector; @@ -37,7 +36,6 @@ pub struct SparseVectorIndex { searches_telemetry: SparseSearchesTelemetry, indices_tracker: IndicesTracker, scores_memory_pool: ScoresMemoryPool, - deferred_internal_id: Option, } /// Getters for internals, used for testing. @@ -72,7 +70,6 @@ pub struct SparseVectorIndexOpenArgs<'a, F: FnMut()> { pub path: &'a Path, pub stopped: &'a AtomicBool, pub tick_progress: F, - pub deferred_internal_id: Option, } impl SparseVectorIndex { @@ -86,7 +83,6 @@ impl SparseVectorIndex { path, stopped, tick_progress, - deferred_internal_id, } = args; let config_path = SparseIndexConfig::get_config_path(path); @@ -145,7 +141,6 @@ impl SparseVectorIndex { searches_telemetry, indices_tracker, scores_memory_pool, - deferred_internal_id, }) } diff --git a/lib/segment/src/index/sparse_index/sparse_vector_index/search.rs b/lib/segment/src/index/sparse_index/sparse_vector_index/search.rs index 973f01615d..e3d83f191c 100644 --- a/lib/segment/src/index/sparse_index/sparse_vector_index/search.rs +++ b/lib/segment/src/index/sparse_index/sparse_vector_index/search.rs @@ -1,5 +1,5 @@ use common::counter::hardware_counter::HardwareCounterCell; -use common::types::{PointOffsetType, ScoredPointOffset}; +use common::types::{DeferredBehavior, PointOffsetType, ScoredPointOffset}; use itertools::Itertools; use sparse::common::sparse_vector::SparseVector; use sparse::index::inverted_index::InvertedIndex; @@ -71,14 +71,10 @@ impl SparseVectorIndex { // `prefiltered_points` always contains visible points only so we don't need additional filtering here. Some(filtered_points) => filtered_points.iter().copied(), None => { - let filtered_points = self.payload_index.borrow().with_view(|v| { - v.query_points( - filter, - &hw_counter, - &is_stopped, - vector_query_context.deferred_internal_id(), - ) - })?; + let filtered_points = self + .payload_index + .borrow() + .with_view(|v| v.query_points(filter, &hw_counter, &is_stopped))?; *prefiltered_points = Some(filtered_points); prefiltered_points.as_ref().unwrap().iter().copied() } @@ -86,7 +82,10 @@ impl SparseVectorIndex { searcher.peek_top_iter(filtered_points, &is_stopped)? } None => { - searcher.peek_top_all(&is_stopped, vector_query_context.deferred_internal_id())? + let iter = id_tracker + .point_mappings() + .filter_deferred(searcher.iter_not_deleted(), DeferredBehavior::Exclude); + searcher.peek_top_iter(iter, &is_stopped)? } }; let res = results.pop().expect("single element results"); @@ -119,14 +118,10 @@ impl SparseVectorIndex { // so no additional filtering is required in that case. Some(filtered_points) => filtered_points.iter(), None => { - let filtered_points = self.payload_index.borrow().with_view(|v| { - v.query_points( - filter, - &hw_counter, - &is_stopped, - vector_query_context.deferred_internal_id(), - ) - })?; + let filtered_points = self + .payload_index + .borrow() + .with_view(|v| v.query_points(filter, &hw_counter, &is_stopped))?; *prefiltered_points = Some(filtered_points); prefiltered_points.as_ref().unwrap().iter() } diff --git a/lib/segment/src/index/sparse_index/sparse_vector_index/vector_index_impl.rs b/lib/segment/src/index/sparse_index/sparse_vector_index/vector_index_impl.rs index 56e96dbe81..66499161a8 100644 --- a/lib/segment/src/index/sparse_index/sparse_vector_index/vector_index_impl.rs +++ b/lib/segment/src/index/sparse_index/sparse_vector_index/vector_index_impl.rs @@ -14,6 +14,7 @@ use crate::common::operation_error::{OperationError, OperationResult, check_proc use crate::data_types::named_vectors::CowVector; use crate::data_types::query_context::VectorQueryContext; use crate::data_types::vectors::{QueryVector, VectorInternal, VectorRef}; +use crate::id_tracker::IdTrackerRead; use crate::index::sparse_index::indices_tracker::IndicesTracker; use crate::index::sparse_index::sparse_index_config::{SparseIndexConfig, SparseIndexType}; use crate::index::{VectorIndex, VectorIndexRead}; @@ -34,12 +35,6 @@ impl VectorIndexRead for SparseVectorIndex VectorIndex for SparseVectorIndex= deferred); if point_is_deferred { diff --git a/lib/segment/src/index/struct_payload_index/read_view/payload_index_read.rs b/lib/segment/src/index/struct_payload_index/read_view/payload_index_read.rs index 85fccbdfde..10ae5afdc6 100644 --- a/lib/segment/src/index/struct_payload_index/read_view/payload_index_read.rs +++ b/lib/segment/src/index/struct_payload_index/read_view/payload_index_read.rs @@ -6,11 +6,11 @@ use common::counter::hardware_counter::HardwareCounterCell; use common::counter::iterator_hw_measurement::HwMeasurementIteratorExt; use common::either_variant::EitherVariant; use common::iterator_ext::IteratorExt; -use common::types::{PointOffsetType, ScoreType}; +use common::types::{DeferredBehavior, PointOffsetType, ScoreType}; use super::StructPayloadIndexReadView; use crate::common::operation_error::OperationResult; -use crate::id_tracker::{IdTrackerRead, PointMappingsRefEnum}; +use crate::id_tracker::IdTrackerRead; use crate::index::PayloadIndexRead; use crate::index::field_index::numeric_index::NumericFieldIndexRead; use crate::index::field_index::{ @@ -68,20 +68,16 @@ where filter: &Filter, hw_counter: &HardwareCounterCell, is_stopped: &AtomicBool, - deferred_internal_id: Option, ) -> OperationResult> { // Assume query is already estimated to be small enough so we can iterate over all matched ids let query_cardinality = self.estimate_cardinality(filter, hw_counter)?; - let point_mappings = self.id_tracker.point_mappings(); let result = self .iter_filtered_points( filter, - self.id_tracker, - &point_mappings, &query_cardinality, hw_counter, is_stopped, - deferred_internal_id, + DeferredBehavior::Exclude, )? .collect(); Ok(result) @@ -143,18 +139,18 @@ where )) } - fn iter_filtered_points<'b, IT: IdTrackerRead>( + fn iter_filtered_points<'b>( &'b self, filter: &'b Filter, - id_tracker: &'b IT, - point_mappings: &'b PointMappingsRefEnum<'b>, query_cardinality: &'b CardinalityEstimation, hw_counter: &'b HardwareCounterCell, is_stopped: &'b AtomicBool, - deferred_internal_id: Option, + deferred_behavior: DeferredBehavior, ) -> OperationResult + 'b> { + let point_mappings = self.id_tracker.point_mappings(); + if query_cardinality.primary_clauses.is_empty() { - let full_scan_iterator = point_mappings.iter_internal_visible(deferred_internal_id); + let full_scan_iterator = point_mappings.iter_internal_with_behavior(deferred_behavior); let struct_filtered_context = self.struct_filtered_context(filter, hw_counter)?; // Worst case: query expected to return few matches, but index can't be used @@ -165,7 +161,7 @@ where Ok(EitherVariant::A(matched_points)) } else { // CPU-optimized strategy here: points are made unique before applying other filters. - let mut visited_list = self.visited_pool.get(id_tracker.total_point_count()); + let mut visited_list = self.visited_pool.get(self.id_tracker.total_point_count()); // If even one iterator is None, we should replace the whole thing with // an iterator over all ids. @@ -180,14 +176,12 @@ where .iter_conditions() .all(|condition| query_cardinality.is_primary(condition)); - let joined_primary_iterator = primary_iterators - .into_iter() - // Filter out deferred points. - // This iterator (and each primary iterator too) can yield items in non sorted order, depending on the type of index and primary condition. - .flatten() - .filter(move |&internal_id| { - internal_id < deferred_internal_id.unwrap_or(PointOffsetType::MAX) - }) + // Primary clause iterators come from field indexes and don't go through + // the mapping, so deferred filtering must be applied to them explicitly. + // Each primary iterator (and the flattened stream) can yield items in + // non-sorted order depending on the field-index type and primary condition. + let joined_primary_iterator = point_mappings + .filter_deferred(primary_iterators.into_iter().flatten(), deferred_behavior) .stop_if(is_stopped); return Ok(if all_conditions_are_primary { @@ -212,7 +206,7 @@ where // and applying full filter. let struct_filtered_context = self.struct_filtered_context(filter, hw_counter)?; - let id_tracker_iterator = point_mappings.iter_internal_visible(deferred_internal_id); + let id_tracker_iterator = point_mappings.iter_internal_with_behavior(deferred_behavior); let iter = id_tracker_iterator .stop_if(is_stopped) diff --git a/lib/segment/src/index/struct_payload_index/read_view/tests.rs b/lib/segment/src/index/struct_payload_index/read_view/tests.rs index cf20041cd7..1ccf6e263c 100644 --- a/lib/segment/src/index/struct_payload_index/read_view/tests.rs +++ b/lib/segment/src/index/struct_payload_index/read_view/tests.rs @@ -55,7 +55,7 @@ fn smoke_view_over_in_memory_backends() { // `query_points` over an empty filter on an empty tracker returns nothing. let empty_filter = Filter::default(); let result = view - .query_points(&empty_filter, &hw_counter, &is_stopped, None) + .query_points(&empty_filter, &hw_counter, &is_stopped) .expect("query_points"); assert!(result.is_empty(), "no points in tracker"); diff --git a/lib/segment/src/segment/as_view.rs b/lib/segment/src/segment/as_view.rs index 6932f34177..e71efd0605 100644 --- a/lib/segment/src/segment/as_view.rs +++ b/lib/segment/src/segment/as_view.rs @@ -20,7 +20,6 @@ impl Segment { payload_storage: payload_storage.deref(), vector_data: &self.vector_data, segment_config: &self.segment_config, - deferred_point_status: self.deferred_point_status.as_ref(), appendable_flag: self.appendable_flag, }; diff --git a/lib/segment/src/segment/memory.rs b/lib/segment/src/segment/memory.rs index 0764bb8f95..4868182f34 100644 --- a/lib/segment/src/segment/memory.rs +++ b/lib/segment/src/segment/memory.rs @@ -53,7 +53,6 @@ impl Segment { segment_type: _, segment_config, error_status: _, - deferred_point_status: _, } = self; let sparse_names = &segment_config.sparse_vector_data; diff --git a/lib/segment/src/segment/mod.rs b/lib/segment/src/segment/mod.rs index e5b925472a..4c0c1154b2 100644 --- a/lib/segment/src/segment/mod.rs +++ b/lib/segment/src/segment/mod.rs @@ -21,7 +21,6 @@ use std::sync::Arc; use atomic_refcell::AtomicRefCell; use common::is_alive_lock::IsAliveLock; use common::storage_version::StorageVersion; -use common::types::PointOffsetType; use parking_lot::Mutex; use uuid::Uuid; @@ -87,18 +86,6 @@ pub struct Segment { /// Last unhandled error /// If not None, all update operations will be aborted until original operation is performed properly pub error_status: Option, - pub(crate) deferred_point_status: Option, -} - -#[derive(Debug)] -pub struct DeferredPointStatus { - /// Points with internal id >= this value are hidden from reads. - /// Available for appendable segments only. - pub(crate) deferred_internal_id: PointOffsetType, - - /// Amount of deleted deferred points. Must kept track of properly to be able - /// to calculate the amount of available deferred and visible points. - pub(crate) deferred_deleted_count: usize, } pub struct VectorData { diff --git a/lib/segment/src/segment/read_view/deferred.rs b/lib/segment/src/segment/read_view/deferred.rs index 17beed1424..f629d57a72 100644 --- a/lib/segment/src/segment/read_view/deferred.rs +++ b/lib/segment/src/segment/read_view/deferred.rs @@ -15,11 +15,11 @@ where TVD: VectorDataRead, { pub(super) fn deferred_internal_id(&self) -> Option { - self.deferred_point_status.map(|s| s.deferred_internal_id) + self.id_tracker.deferred_internal_id() } - pub(super) fn deferred_deleted_count(&self) -> Option { - self.deferred_point_status.map(|s| s.deferred_deleted_count) + pub(super) fn deferred_deleted_count(&self) -> usize { + self.id_tracker.deferred_deleted_count() } pub fn deferred_point_count(&self) -> usize { @@ -28,7 +28,7 @@ where .id_tracker .total_point_count() .saturating_sub(internal_id as usize) - .saturating_sub(self.deferred_deleted_count().unwrap_or_default()), + .saturating_sub(self.deferred_deleted_count()), None => 0, } } diff --git a/lib/segment/src/segment/read_view/facet.rs b/lib/segment/src/segment/read_view/facet.rs index 09d49adc48..bbe28d858c 100644 --- a/lib/segment/src/segment/read_view/facet.rs +++ b/lib/segment/src/segment/read_view/facet.rs @@ -2,7 +2,7 @@ use std::collections::{BTreeSet, HashMap}; use std::sync::atomic::AtomicBool; use common::counter::hardware_counter::HardwareCounterCell; -use common::types::PointOffsetType; +use common::types::{DeferredBehavior, PointOffsetType}; use itertools::Itertools; use crate::common::operation_error::{OperationError, OperationResult, check_process_stopped}; @@ -67,17 +67,14 @@ where if use_iterative_approach { // Go over the filtered points and aggregate the values (read from other indexes). - let point_mappings = self.id_tracker.point_mappings(); let points = self .payload_index .iter_filtered_points( filter, - self.id_tracker, - &point_mappings, &filter_cardinality, hw_counter, is_stopped, - self.deferred_internal_id(), + DeferredBehavior::Exclude, )? .filter(|&point_id| !self.id_tracker.is_deleted_point(point_id)); facet_index.for_points_values(points, hw_counter, |_point_id, iter| { @@ -150,18 +147,15 @@ where let filter_cardinality = self .payload_index .estimate_cardinality(filter, hw_counter)?; - let point_mappings = self.id_tracker.point_mappings(); let points = self .payload_index .iter_filtered_points( filter, - self.id_tracker, - &point_mappings, &filter_cardinality, hw_counter, is_stopped, - self.deferred_internal_id(), + DeferredBehavior::Exclude, )? .filter(|&point_id| !self.id_tracker.is_deleted_point(point_id)); facet_index.for_points_values(points, hw_counter, |_point_id, iter| { diff --git a/lib/segment/src/segment/read_view/info.rs b/lib/segment/src/segment/read_view/info.rs index 5c872b9ce0..0ddf4d3230 100644 --- a/lib/segment/src/segment/read_view/info.rs +++ b/lib/segment/src/segment/read_view/info.rs @@ -83,7 +83,7 @@ where num_indexed_vectors: num_indexed_vectors_total, num_points: self.id_tracker.available_point_count(), num_deferred_points: Some(self.deferred_point_count()), - num_deleted_deferred_points: Some(self.deferred_deleted_count().unwrap_or_default()), + num_deleted_deferred_points: Some(self.deferred_deleted_count()), num_deleted_vectors: self.id_tracker.deleted_point_count(), vectors_size_bytes, // Considers vector storage, but not indices. payloads_size_bytes, // Considers payload storage, but not indices. diff --git a/lib/segment/src/segment/read_view/mod.rs b/lib/segment/src/segment/read_view/mod.rs index 8d669f22e4..4826954115 100644 --- a/lib/segment/src/segment/read_view/mod.rs +++ b/lib/segment/src/segment/read_view/mod.rs @@ -18,8 +18,8 @@ use crate::index::field_index::FieldIndex; use crate::index::struct_payload_index::StructPayloadIndexReadView; use crate::payload_storage::PayloadStorageRead; use crate::payload_storage::payload_storage_enum::PayloadStorageEnum; +use crate::segment::VectorData; use crate::segment::vector_data_read::VectorDataRead; -use crate::segment::{DeferredPointStatus, VectorData}; use crate::types::{PointIdType, SegmentConfig, SeqNumberType, VectorNameBuf}; use crate::vector_storage::VectorStorageEnum; @@ -41,7 +41,6 @@ where pub(crate) payload_storage: &'s TPayloadStorage, pub(crate) vector_data: &'s HashMap, pub(crate) segment_config: &'s SegmentConfig, - pub(crate) deferred_point_status: Option<&'s DeferredPointStatus>, pub(crate) appendable_flag: bool, } diff --git a/lib/segment/src/segment/read_view/order_by.rs b/lib/segment/src/segment/read_view/order_by.rs index 59fc66cbcd..dcbcf728a0 100644 --- a/lib/segment/src/segment/read_view/order_by.rs +++ b/lib/segment/src/segment/read_view/order_by.rs @@ -46,19 +46,14 @@ where let start_from = order_by.start_from(); - let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id()); - - let point_mappings = self.id_tracker.point_mappings(); let values_ids_iterator = self .payload_index .iter_filtered_points( condition, - self.id_tracker, - &point_mappings, &cardinality_estimation, hw_counter, is_stopped, - effective_deferred_id, + deferred_behavior, )? .flat_map(|internal_id| { // Repeat a point for as many values as it has. diff --git a/lib/segment/src/segment/read_view/sampling.rs b/lib/segment/src/segment/read_view/sampling.rs index dfd51c420d..5651898993 100644 --- a/lib/segment/src/segment/read_view/sampling.rs +++ b/lib/segment/src/segment/read_view/sampling.rs @@ -2,6 +2,7 @@ use std::sync::atomic::AtomicBool; use common::counter::hardware_counter::HardwareCounterCell; use common::iterator_ext::IteratorExt; +use common::types::DeferredBehavior; use rand::seq::{IteratorRandom, SliceRandom}; use crate::common::operation_error::OperationResult; @@ -22,7 +23,7 @@ where pub fn read_by_random_id(&self, limit: usize) -> Vec { self.id_tracker .point_mappings() - .iter_random_visible(self.deferred_internal_id()) + .iter_random_visible() .map(|x| x.0) .take(limit) .collect() @@ -35,7 +36,6 @@ where is_stopped: &AtomicBool, hw_counter: &HardwareCounterCell, ) -> OperationResult> { - let point_mappings = self.id_tracker.point_mappings(); let cardinality_estimation = self .payload_index .estimate_cardinality(condition, hw_counter)?; @@ -43,12 +43,10 @@ where .payload_index .iter_filtered_points( condition, - self.id_tracker, - &point_mappings, &cardinality_estimation, hw_counter, is_stopped, - self.deferred_internal_id(), + DeferredBehavior::Exclude, )? .filter_map(|internal_id| self.id_tracker.external_id(internal_id)); @@ -69,7 +67,7 @@ where Ok(self .id_tracker .point_mappings() - .iter_random_visible(self.deferred_internal_id()) + .iter_random_visible() .stop_if(is_stopped) .filter(move |(_, internal_id)| filter_context.check(*internal_id)) .map(|(external_id, _)| external_id) diff --git a/lib/segment/src/segment/read_view/scroll.rs b/lib/segment/src/segment/read_view/scroll.rs index 8740ae389f..254e0f5bab 100644 --- a/lib/segment/src/segment/read_view/scroll.rs +++ b/lib/segment/src/segment/read_view/scroll.rs @@ -59,12 +59,14 @@ where limit: Option, deferred_behavior: DeferredBehavior, ) -> Vec { - let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id()); + let point_mappings = self.id_tracker.point_mappings(); + let iter = if deferred_behavior.include_all_points() { + point_mappings.iter_from(offset) + } else { + point_mappings.iter_from_visible(offset) + }; - self.id_tracker - .point_mappings() - .iter_from_visible(offset, effective_deferred_id) - .map(|x| x.0) + iter.map(|x| x.0) .take(limit.unwrap_or(usize::MAX)) .collect() } @@ -78,23 +80,18 @@ where hw_counter: &HardwareCounterCell, deferred_behavior: DeferredBehavior, ) -> OperationResult> { - let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id()); - let cardinality_estimation = self .payload_index .estimate_cardinality(condition, hw_counter)?; - let point_mappings = self.id_tracker.point_mappings(); let ids_iterator = self .payload_index .iter_filtered_points( condition, - self.id_tracker, - &point_mappings, &cardinality_estimation, hw_counter, is_stopped, - effective_deferred_id, + deferred_behavior, )? .filter_map(|internal_id| { let external_id = self.id_tracker.external_id(internal_id)?; @@ -121,13 +118,15 @@ where hw_counter: &HardwareCounterCell, deferred_behavior: DeferredBehavior, ) -> OperationResult> { - let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id()); - let filter_context = self.payload_index.filter_context(condition, hw_counter)?; - Ok(self - .id_tracker - .point_mappings() - .iter_from_visible(offset, effective_deferred_id) + let point_mappings = self.id_tracker.point_mappings(); + let iter = if deferred_behavior.include_all_points() { + point_mappings.iter_from(offset) + } else { + point_mappings.iter_from_visible(offset) + }; + + Ok(iter .stop_if(is_stopped) .filter(move |(_, internal_id)| filter_context.check(*internal_id)) .map(|(external_id, _)| external_id) diff --git a/lib/segment/src/segment/read_view/search.rs b/lib/segment/src/segment/read_view/search.rs index de4e6b8b54..3acee3467d 100644 --- a/lib/segment/src/segment/read_view/search.rs +++ b/lib/segment/src/segment/read_view/search.rs @@ -2,13 +2,14 @@ use std::sync::atomic::AtomicBool; use ahash::AHashMap; use common::counter::hardware_counter::HardwareCounterCell; +use common::iterator_ext::IteratorExt; use common::types::{DeferredBehavior, ScoredPointOffset}; use crate::common::operation_error::{OperationError, OperationResult}; use crate::common::{check_query_vectors, check_stopped}; use crate::data_types::query_context::{QueryContext, QueryIdfStats, SegmentQueryContext}; use crate::data_types::segment_record::{NamedVectorsOwned, SegmentRecord}; -use crate::data_types::vectors::{QueryVector, VectorInternal, VectorStructInternal}; +use crate::data_types::vectors::{QueryVector, VectorStructInternal}; use crate::id_tracker::IdTrackerRead; use crate::index::{PayloadIndexRead, VectorIndexRead}; use crate::payload_storage::PayloadStorageRead; @@ -37,89 +38,87 @@ where is_stopped: &AtomicBool, deferred_behavior: DeferredBehavior, ) -> OperationResult> { - let mut records = AHashMap::with_capacity(point_ids.len()); + // Stage 1: resolve external → internal once, into two parallel vectors. + // The id tracker owns this: deferred filtering happens inline (no + // `point_is_deferred` lookup). The parallel-vector shape lets + // a future batched payload / vector fetcher consume `&offsets` + // straight without unzipping first. + let (resolved_ids, resolved_offsets) = self + .id_tracker + .resolve_external_ids(point_ids, deferred_behavior); + debug_assert_eq!(resolved_ids.len(), resolved_offsets.len()); - // Filter out deferred points. This is done in two stages to prevent cloning `point_ids` - // and iterating more than needed but still satisfy Rust's ownership constraints. - let behavior_allows_filtering = !deferred_behavior.include_all_points(); - let filter_deferred = self.has_deferred_points() && behavior_allows_filtering; - let filtered_point_ids = filter_deferred.then(|| { - point_ids - .iter() - .filter(|&&point_id| !self.point_is_deferred(point_id)) - .copied() - .collect::>() - }); + // Stage 2: pre-allocate one record per resolved point. The `vectors` + // slot is initialised here according to `with_vector`, so the + // `WithVector::Bool(false)` path needs no separate clearing pass. + let needs_vectors = match with_vector { + WithVector::Bool(true) | WithVector::Selector(_) => true, + WithVector::Bool(false) => false, + }; + let mut records: AHashMap = resolved_ids + .iter() + .map(|&id| { + let record = SegmentRecord { + id, + vectors: needs_vectors.then(NamedVectorsOwned::default), + payload: None, + }; + (id, record) + }) + .collect(); - // Stage two: select the correct slice and shadow `point_ids`. - let point_ids = filtered_point_ids.as_deref().unwrap_or(point_ids); - - let mut update_record_vector = - |vector_name: &VectorNameBuf, - point_id: PointIdType, - vector_internal: VectorInternal| { - let point_record = records - .entry(point_id) - .or_insert_with(|| SegmentRecord::empty(point_id)); - - point_record - .vectors - .get_or_insert_with(NamedVectorsOwned::default) - .push((vector_name.clone(), vector_internal)); + // Stage 3: vectors. The external id rides along as the read's user + // data — it comes back unchanged in the callback, so no extra + // `offset → id` lookup is needed. + if needs_vectors { + let mut process_vectors = |vector_name: &VectorNameBuf| -> OperationResult<()> { + let keys = resolved_ids + .iter() + .zip(&resolved_offsets) + .map(|(&id, &offset)| (id, offset)) + .stop_if(is_stopped); + self.vectors_by_offsets(vector_name, keys, hw_counter, |id, _offset, vec| { + if let Some(record) = records.get_mut(&id) { + record + .vectors + .as_mut() + .expect("needs_vectors path keeps vectors as Some") + .push((vector_name.clone(), vec)); + } + }) }; - match with_vector { - WithVector::Bool(true) => { - for vector_name in self.vector_data.keys() { - self.read_vectors( - vector_name, - point_ids, - hw_counter, - is_stopped, - |point_id, vec| { - update_record_vector(vector_name, point_id, vec); - }, - )?; + match with_vector { + WithVector::Bool(true) => { + for vector_name in self.vector_data.keys() { + process_vectors(vector_name)?; + } } - } - WithVector::Bool(false) => { - // Do not display empty `vectors: {}` if disabled. - for &point_id in point_ids { - let point_record = records - .entry(point_id) - .or_insert_with(|| SegmentRecord::empty(point_id)); - point_record.vectors = None; - } - } - WithVector::Selector(selector) => { - for vector_name in selector { - self.read_vectors( - vector_name, - point_ids, - hw_counter, - is_stopped, - |point_id, vec| { - update_record_vector(vector_name, point_id, vec); - }, - )?; + WithVector::Selector(names) => { + for vector_name in names { + process_vectors(vector_name)?; + } } + WithVector::Bool(false) => unreachable!("guarded by needs_vectors"), } } - for &point_id in point_ids { - let payload = if with_payload.enable { - if let Some(selector) = &with_payload.payload_selector { - Some(selector.process(self.payload(point_id, hw_counter)?)) - } else { - Some(self.payload(point_id, hw_counter)?) + // Stage 4: payload. Use the already-resolved offsets to skip another + // external→internal lookup per point. The per-iteration shape here + // mirrors what a future batched payload fetcher would consume — + // `&resolved_offsets` becomes its input directly. + if with_payload.enable { + for (&id, &offset) in resolved_ids.iter().zip(&resolved_offsets) { + check_stopped(is_stopped)?; + let payload = self.payload_by_offset(offset, hw_counter)?; + let payload = match &with_payload.payload_selector { + Some(selector) => selector.process(payload), + None => payload, + }; + if let Some(record) = records.get_mut(&id) { + record.payload = Some(payload); } - } else { - None - }; - let point_record = records - .entry(point_id) - .or_insert_with(|| SegmentRecord::empty(point_id)); - point_record.payload = payload; + } } Ok(records) @@ -226,8 +225,7 @@ where .vector_data .get(vector_name) .ok_or_else(|| OperationError::vector_name_not_exists(vector_name))?; - let vector_query_context = - query_context.get_vector_context(vector_name, self.deferred_internal_id()); + let vector_query_context = query_context.get_vector_context(vector_name); let internal_results = vector_data.vector_index().search( query_vectors, filter, diff --git a/lib/segment/src/segment/read_view/vectors.rs b/lib/segment/src/segment/read_view/vectors.rs index eb64e0ce1c..8f8f7b3971 100644 --- a/lib/segment/src/segment/read_view/vectors.rs +++ b/lib/segment/src/segment/read_view/vectors.rs @@ -1,8 +1,5 @@ -use std::sync::atomic::AtomicBool; - use common::counter::hardware_counter::HardwareCounterCell; use common::generic_consts::Random; -use common::iterator_ext::IteratorExt; use common::types::PointOffsetType; use crate::common::check_vector_name; @@ -35,9 +32,9 @@ where let mut result = None; self.vectors_by_offsets( vector_name, - std::iter::once(point_offset), + std::iter::once(((), point_offset)), hw_counter, - |_, vector_internal| { + |(), _, vector_internal| { result = Some(vector_internal); }, )?; @@ -45,12 +42,19 @@ where } /// Retrieve multiple vectors by internal ID. - pub fn vectors_by_offsets( + /// + /// Each input is tagged with caller-supplied user data `U` (any `Copy` + /// type — the external id, the input position, …). The data is threaded + /// back to the callback unchanged, so callers can map results into a + /// parallel input array without keeping a separate `offset → ...` lookup + /// table. Deleted points are filtered out lazily — entries with deleted + /// vectors or deleted points are simply not delivered to the callback. + pub fn vectors_by_offsets( &self, vector_name: &VectorName, - point_offsets: impl IntoIterator, + keys: impl IntoIterator, hw_counter: &HardwareCounterCell, - mut callback: impl FnMut(PointOffsetType, VectorInternal), + mut callback: impl FnMut(U, PointOffsetType, VectorInternal), ) -> OperationResult<()> { check_vector_name(vector_name, self.segment_config)?; let vector_data = self @@ -61,7 +65,9 @@ where let total_vectors = vector_storage.total_vector_count(); let id_tracker = self.id_tracker; - let non_deleted_offsets = point_offsets.into_iter().filter(|&point_offset| { + // The user-data tag rides alongside each offset, so the deletion + // filter stays lazy — no parallel `Vec<(orig_idx, offset)>` needed. + let live_keys = keys.into_iter().filter(|&(_, point_offset)| { if total_vectors <= point_offset as usize { debug_assert!( false, @@ -70,59 +76,23 @@ where ); return false; } - let is_vector_deleted = vector_storage.is_deleted_vector(point_offset); let is_point_deleted = id_tracker.is_deleted_point(point_offset); !is_vector_deleted && !is_point_deleted }); - vector_storage.read_vectors::(non_deleted_offsets, |point_offset, cow_vector| { - if vector_storage.is_on_disk() { - hw_counter - .vector_io_read() - .incr_delta(cow_vector.estimate_size_in_bytes()); - } - callback(point_offset, cow_vector.to_owned()); - }); - - Ok(()) - } - - /// Retrieve named vectors for a list of external point IDs, invoking the - /// callback for each (point_id, vector) pair. - pub fn read_vectors( - &self, - vector_name: &VectorName, - point_ids: &[PointIdType], - hw_counter: &HardwareCounterCell, - is_stopped: &AtomicBool, - mut callback: impl FnMut(PointIdType, VectorInternal), - ) -> OperationResult<()> { - let mut error = None; - let internal_ids = point_ids - .iter() - .copied() - .stop_if(is_stopped) - .filter_map(|point_id| match self.lookup_internal_id(point_id) { - Ok(point_offset) => Some(point_offset), - Err(err) => { - error = Some(err); - None - } - }); - self.vectors_by_offsets( - vector_name, - internal_ids, - hw_counter, - |point_offset, vector_internal| { - if let Some(point_id) = self.id_tracker.external_id(point_offset) { - callback(point_id, vector_internal); + vector_storage.read_vectors::( + live_keys, + |user_data, point_offset, cow_vector| { + if vector_storage.is_on_disk() { + hw_counter + .vector_io_read() + .incr_delta(cow_vector.estimate_size_in_bytes()); } + callback(user_data, point_offset, cow_vector.to_owned()); }, - )?; - if let Some(err) = error { - return Err(err); - } + ); + Ok(()) } diff --git a/lib/segment/src/segment/segment_ops.rs b/lib/segment/src/segment/segment_ops.rs index de81fd4e49..38080353ff 100644 --- a/lib/segment/src/segment/segment_ops.rs +++ b/lib/segment/src/segment/segment_ops.rs @@ -315,21 +315,11 @@ impl Segment { let mut id_tracker = self.id_tracker.borrow_mut(); - let is_point_already_deleted = id_tracker.is_deleted_point(internal_id); - + // `drop_internal` updates the id tracker's deferred-deleted counter when + // the point sits at or above the deferred threshold, with double-delete + // protection inside `PointMappings::drop`. id_tracker.drop_internal(internal_id)?; - let deferred_point_status = self.deferred_point_status.as_mut(); - - // Increase counter for deleted points. - if let Some(deferred_point_status) = deferred_point_status - && internal_id >= deferred_point_status.deferred_internal_id - // Don't count the deletion of the same point twice - && !is_point_already_deleted - { - deferred_point_status.deferred_deleted_count += 1; - } - // Before, we propagated point deletions to also delete its vectors. This turns // out to be problematic because this sometimes makes us lose vector data // because we cannot control the order of segment flushes. @@ -583,29 +573,6 @@ impl Segment { pub fn fix_id_tracker_inconsistencies(&mut self) -> OperationResult> { self.id_tracker.borrow_mut().fix_inconsistencies() } - - /// Calculates the amount of deleted deferred points by iterating over all points in the ID tracker. Therefore this operation - /// can be expensive and should only be run once at segment creation. - pub(crate) fn calculate_deleted_deferred_point_count(&self) -> usize { - let Some(deferred_from) = self.deferred_internal_id() else { - return 0; - }; - - let id_tracker = self.id_tracker.borrow(); - let total_points = id_tracker.total_point_count(); - - if total_points < deferred_from as usize { - return 0; - } - - id_tracker.deleted_point_bitslice()[deferred_from as usize..total_points].count_ones() - } - - pub(crate) fn deferred_internal_id(&self) -> Option { - self.deferred_point_status - .as_ref() - .map(|i| i.deferred_internal_id) - } } fn restore_snapshot_in_place(snapshot_path: &Path) -> OperationResult<()> { diff --git a/lib/segment/src/segment/tests/mod.rs b/lib/segment/src/segment/tests/mod.rs index f6c36ef827..c55ff8255b 100644 --- a/lib/segment/src/segment/tests/mod.rs +++ b/lib/segment/src/segment/tests/mod.rs @@ -7,7 +7,7 @@ use ahash::AHashSet; use common::counter::hardware_counter::HardwareCounterCell; use common::tar_ext; use common::tar_unpack::tar_unpack_file; -use common::types::DeferredBehavior; +use common::types::{DeferredBehavior, PointOffsetType}; use fs_err as fs; use fs_err::File; use ordered_float::OrderedFloat; @@ -32,7 +32,7 @@ use crate::entry::entry_point::{ NonAppendableSegmentEntry as _, ReadSegmentEntry as _, SegmentEntry as _, }; use crate::entry::{SnapshotEntry as _, StorageSegmentEntry as _}; -use crate::id_tracker::IdTracker; +use crate::id_tracker::{IdTracker, IdTrackerRead}; use crate::index::sparse_index::sparse_index_config::{SparseIndexConfig, SparseIndexType}; use crate::json_path::JsonPath; use crate::segment_constructor::simple_segment_constructor::{ @@ -910,7 +910,10 @@ fn create_deferred_segment( // Now we should have deferred points assert_eq!(segment.has_deferred_points(), n_deferred > 0); if n_deferred > 0 { - assert_eq!(segment.deferred_internal_id(), Some(n_vectors as u32)); + assert_eq!( + segment.id_tracker.borrow().deferred_internal_id(), + Some(n_vectors as u32) + ); } // Points 1 to n_vectors should NOT be deferred @@ -1007,7 +1010,7 @@ fn test_dense_deferred_points() { "Segment should still have deferred points after reopening" ); assert_eq!( - segment.deferred_internal_id(), + segment.id_tracker.borrow().deferred_internal_id(), Some(13), "Deferred internal ID should still be `DEFERRED_POINTS_ID` after reopening" ); @@ -1071,7 +1074,7 @@ fn test_deferred_point_estimation_with_filter() { // For consistency we also test that the same cardinality is estimated if no deferred points exist. if n_deferred == 0 { - assert_eq!(segment.deferred_internal_id(), None); + assert_eq!(segment.id_tracker.borrow().deferred_internal_id(), None); let estimation = segment .estimate_point_count(Some(&filter), &hw_counter) .unwrap(); @@ -1301,7 +1304,7 @@ fn test_deferred_point_facets() { ); let dir = Builder::new().prefix("segment_dir").tempdir().unwrap(); - let mut segment = create_deferred_segment(&dir, 5, N_POINTS, n_deferred); + let segment = create_deferred_segment(&dir, 5, N_POINTS, n_deferred); let request = FacetParams { key: key.clone(), @@ -1314,14 +1317,17 @@ fn test_deferred_point_facets() { .facet(&request, &AtomicBool::new(false), &hw_counter) .unwrap(); - let old_status = segment.deferred_point_status.take(); - if n_deferred > 0 { - assert!(old_status.is_some()); - } - let facet_res = segment + // Compare against the same point set without deferred mode by + // rebuilding the segment with `deferred_internal_id = None`. + let no_deferred_dir = Builder::new() + .prefix("segment_dir_no_deferred") + .tempdir() + .unwrap(); + let no_deferred_segment = + create_deferred_segment(&no_deferred_dir, 5, N_POINTS + n_deferred, 0); + let facet_res = no_deferred_segment .facet(&request, &AtomicBool::new(false), &hw_counter) .unwrap(); - segment.deferred_point_status = old_status; let expected_deferred = if filter.is_some() { n_deferred.div_ceil(3) @@ -1446,7 +1452,7 @@ fn assert_deferred_points_excluded( log::debug!(" => deferred points = {n_deferred}; filter-set ID = {filter_set_id}",); let dir = Builder::new().prefix("segment_dir").tempdir().unwrap(); - let mut segment = create_deferred_segment(&dir, 5, N_POINTS, n_deferred); + let segment = create_deferred_segment(&dir, 5, N_POINTS, n_deferred); // Search with deferred mode let search_res_deferred = operation(&segment, filter_set.filter.as_ref()); @@ -1460,29 +1466,34 @@ fn assert_deferred_points_excluded( } // Disable deferred points and search again. - if need_rebuilt_segment { - // Don't run this on windows because this test is already extremely slow. - // Recreating the segment here would double that time. - if cfg!(target_os = "windows") { - drop(segment); - dir.close().unwrap(); - continue; - } - - let dir = Builder::new().prefix("segment_dir_2").tempdir().unwrap(); - segment = create_deferred_segment(&dir, 5, N_POINTS + n_deferred, 0); - } else { - segment.deferred_point_status = None; + // On Windows segment creation is extremely IO-heavy; skip the rebuild path + // for tests where it noticeably slows the suite down. + if need_rebuilt_segment && cfg!(target_os = "windows") { + drop(segment); + dir.close().unwrap(); + continue; } - let search_res_normal = operation(&segment, filter_set.filter.as_ref()); + // Deferred state is owned by the id tracker and only set at segment + // construction time, so we rebuild a fresh non-deferred segment with the + // same total point count to compare against. + let no_deferred_dir = Builder::new() + .prefix("segment_dir_no_deferred") + .tempdir() + .unwrap(); + let no_deferred_segment = + create_deferred_segment(&no_deferred_dir, 5, N_POINTS + n_deferred, 0); + + let search_res_normal = operation(&no_deferred_segment, filter_set.filter.as_ref()); assert_eq!( search_res_normal.len(), filter_set.expected_visible + (filter_set.expected_deferred)(n_deferred) ); drop(segment); + drop(no_deferred_segment); dir.close().unwrap(); + no_deferred_dir.close().unwrap(); } } } @@ -1510,7 +1521,7 @@ fn test_deleted_deferred_point_count() { assert_eq!(segment.available_point_count_without_deferred(), N_POINTS); for d in 0..n_deferred { - let delete_id = segment.deferred_internal_id().unwrap() + d as u32; + let delete_id = segment.id_tracker.borrow().deferred_internal_id().unwrap() + d as u32; segment .delete_point_internal(delete_id, &hw_counter) .unwrap(); @@ -1521,7 +1532,7 @@ fn test_deleted_deferred_point_count() { n_deferred.checked_sub(deleted_count).unwrap() ); assert_eq!( - segment.calculate_deleted_deferred_point_count(), + segment.id_tracker.borrow().deferred_deleted_count(), deleted_count, ); @@ -1536,7 +1547,7 @@ fn test_deleted_deferred_point_count() { ); assert_eq!( - segment.calculate_deleted_deferred_point_count(), + segment.id_tracker.borrow().deferred_deleted_count(), deleted_count ); @@ -1545,7 +1556,10 @@ fn test_deleted_deferred_point_count() { // We delete all deferred points in the segment. assert_eq!(segment.deferred_point_count(), 0); - assert_eq!(segment.calculate_deleted_deferred_point_count(), n_deferred); + assert_eq!( + segment.id_tracker.borrow().deferred_deleted_count(), + n_deferred + ); assert_eq!(segment.available_point_count_without_deferred(), N_POINTS); } } diff --git a/lib/segment/src/segment/vector_name_ops.rs b/lib/segment/src/segment/vector_name_ops.rs index 2ca4cba0c8..bb6dcdb532 100644 --- a/lib/segment/src/segment/vector_name_ops.rs +++ b/lib/segment/src/segment/vector_name_ops.rs @@ -169,7 +169,6 @@ impl Segment { path: &vector_index_path, stopped: &stopped, tick_progress: || (), - deferred_internal_id: None, })?; // Register the new storage with the payload index so `has_vector` diff --git a/lib/segment/src/segment_constructor/segment_builder.rs b/lib/segment/src/segment_constructor/segment_builder.rs index defafa9ba4..76cac217ad 100644 --- a/lib/segment/src/segment_constructor/segment_builder.rs +++ b/lib/segment/src/segment_constructor/segment_builder.rs @@ -86,7 +86,9 @@ impl SegmentBuilder { let temp_dir = create_temp_dir(temp_dir)?; let id_tracker = if segment_config.is_appendable() { - IdTrackerEnum::MutableIdTracker(create_mutable_id_tracker(temp_dir.path())?) + // Deferred state is applied when the freshly built segment is reloaded + // via `load_segment`. The transient builder tracker doesn't need it. + IdTrackerEnum::MutableIdTracker(create_mutable_id_tracker(temp_dir.path(), None)?) } else { IdTrackerEnum::InMemoryIdTracker(InMemoryIdTracker::new()) }; @@ -658,8 +660,6 @@ impl SegmentBuilder { path: &vector_index_path, stopped, tick_progress: || (), - // We don't use the `index` returned here so we always set deferred to `None`. It's been loaded properly later. - deferred_internal_id: None, })?; if sparse_vector_config.storage_type.is_on_disk() { diff --git a/lib/segment/src/segment_constructor/segment_constructor_base.rs b/lib/segment/src/segment_constructor/segment_constructor_base.rs index 6cb6e9b80d..afae3586bb 100644 --- a/lib/segment/src/segment_constructor/segment_constructor_base.rs +++ b/lib/segment/src/segment_constructor/segment_constructor_base.rs @@ -39,9 +39,7 @@ use crate::index::sparse_index::sparse_vector_index::{ use crate::index::struct_payload_index::StructPayloadIndex; use crate::payload_storage::mmap_payload_storage::MmapPayloadStorage; use crate::payload_storage::payload_storage_enum::PayloadStorageEnum; -use crate::segment::{ - DeferredPointStatus, SEGMENT_STATE_FILE, Segment, SegmentVersion, VectorData, -}; +use crate::segment::{SEGMENT_STATE_FILE, Segment, SegmentVersion, VectorData}; use crate::types::{ Distance, HnswGlobalConfig, Indexes, PayloadStorageType, SegmentConfig, SegmentState, SegmentType, SeqNumberType, SparseVectorStorageType, VectorDataConfig, VectorName, @@ -223,8 +221,11 @@ pub(crate) fn create_payload_storage( Ok(payload_storage) } -pub(crate) fn create_mutable_id_tracker(segment_path: &Path) -> OperationResult { - MutableIdTracker::open(segment_path) +pub(crate) fn create_mutable_id_tracker( + segment_path: &Path, + deferred_internal_id: Option, +) -> OperationResult { + MutableIdTracker::open(segment_path, deferred_internal_id) } pub(crate) fn get_payload_index_path(segment_path: &Path) -> PathBuf { @@ -411,7 +412,8 @@ fn create_segment( let use_mutable_id_tracker = appendable_flag || !immutable_id_tracker::mappings_path(segment_path).is_file(); let started = Instant::now(); - let id_tracker = create_segment_id_tracker(use_mutable_id_tracker, segment_path)?; + let id_tracker = + create_segment_id_tracker(use_mutable_id_tracker, segment_path, deferred_internal_id)?; log_load_timing(segment_path, "id_tracker", started); let mut vector_storages = HashMap::new(); @@ -556,7 +558,6 @@ fn create_segment( path: &vector_index_path, stopped, tick_progress: || (), - deferred_internal_id, })?); log_load_timing( segment_path, @@ -582,7 +583,7 @@ fn create_segment( SegmentType::Plain }; - let mut segment = Segment { + Ok(Segment { uuid, initial_version, version, @@ -598,22 +599,13 @@ fn create_segment( payload_storage, segment_config: config.clone(), error_status: None, - deferred_point_status: None, - }; - - if let Some(deferred_internal_id) = deferred_internal_id { - segment.deferred_point_status = Some(DeferredPointStatus { - deferred_internal_id, - deferred_deleted_count: segment.calculate_deleted_deferred_point_count(), - }); - } - - Ok(segment) + }) } fn create_segment_id_tracker( mutable_id_tracker: bool, segment_path: &Path, + deferred_internal_id: Option, ) -> OperationResult>> { if !mutable_id_tracker { return Ok(sp(IdTrackerEnum::ImmutableIdTracker( @@ -622,7 +614,7 @@ fn create_segment_id_tracker( } Ok(sp(IdTrackerEnum::MutableIdTracker( - create_mutable_id_tracker(segment_path)?, + create_mutable_id_tracker(segment_path, deferred_internal_id)?, ))) } diff --git a/lib/segment/src/vector_storage/dense/dense_vector_storage.rs b/lib/segment/src/vector_storage/dense/dense_vector_storage.rs index 9bfa84248e..e79a525d50 100644 --- a/lib/segment/src/vector_storage/dense/dense_vector_storage.rs +++ b/lib/segment/src/vector_storage/dense/dense_vector_storage.rs @@ -242,21 +242,22 @@ where .expect("Vector not found") } - fn read_vectors( + fn read_vectors( &self, - keys: impl IntoIterator, - mut callback: impl FnMut(PointOffsetType, CowVector<'_>), + keys: impl IntoIterator, + mut callback: impl FnMut(U, PointOffsetType, CowVector<'_>), ) { - let point_offsets: Vec<_> = keys.into_iter().collect(); + // Split into parallel arrays in one pass: `for_each_in_batch` needs an + // offsets slice (it chunks it for batched reads), but we still want + // `user_data[idx]` available inside the callback. + let (user_data, point_offsets): (Vec, Vec) = keys.into_iter().unzip(); self.vectors .as_ref() .unwrap() .for_each_in_batch(&point_offsets, |idx, vector| { - let point_offset = point_offsets[idx]; let vector = CowVector::from(T::slice_to_float_cow(Cow::Borrowed(vector))); - - callback(point_offset, vector); + callback(user_data[idx], point_offsets[idx], vector); }); } @@ -495,7 +496,7 @@ mod tests { 2, ); let res = searcher - .peek_top_all(&DEFAULT_STOPPED, None) + .peek_top_all(&DEFAULT_STOPPED) .unwrap() .into_iter() .exactly_one() @@ -640,7 +641,7 @@ mod tests { 5, ); let closest = searcher - .peek_top_all(&DEFAULT_STOPPED, None) + .peek_top_all(&DEFAULT_STOPPED) .unwrap() .into_iter() .exactly_one() diff --git a/lib/segment/src/vector_storage/read_only/mod.rs b/lib/segment/src/vector_storage/read_only/mod.rs index a7121d69c8..94c122c195 100644 --- a/lib/segment/src/vector_storage/read_only/mod.rs +++ b/lib/segment/src/vector_storage/read_only/mod.rs @@ -105,22 +105,26 @@ impl VectorStorageRead for VectorStorageReadEnum { } } - fn read_vectors( + fn read_vectors( &self, - keys: impl IntoIterator, - callback: impl FnMut(PointOffsetType, CowVector<'_>), + keys: impl IntoIterator, + callback: impl FnMut(U, PointOffsetType, CowVector<'_>), ) { match self { - VectorStorageReadEnum::Dense(s) => s.read_vectors::

(keys, callback), - VectorStorageReadEnum::DenseByte(s) => s.read_vectors::

(keys, callback), - VectorStorageReadEnum::DenseHalf(s) => s.read_vectors::

(keys, callback), - VectorStorageReadEnum::DenseChunked(s) => s.read_vectors::

(keys, callback), - VectorStorageReadEnum::DenseChunkedByte(s) => s.read_vectors::

(keys, callback), - VectorStorageReadEnum::DenseChunkedHalf(s) => s.read_vectors::

(keys, callback), - VectorStorageReadEnum::MultiDenseChunked(s) => s.read_vectors::

(keys, callback), - VectorStorageReadEnum::MultiDenseChunkedByte(s) => s.read_vectors::

(keys, callback), - VectorStorageReadEnum::MultiDenseChunkedHalf(s) => s.read_vectors::

(keys, callback), - VectorStorageReadEnum::Sparse(s) => s.read_vectors::

(keys, callback), + VectorStorageReadEnum::Dense(s) => s.read_vectors::(keys, callback), + VectorStorageReadEnum::DenseByte(s) => s.read_vectors::(keys, callback), + VectorStorageReadEnum::DenseHalf(s) => s.read_vectors::(keys, callback), + VectorStorageReadEnum::DenseChunked(s) => s.read_vectors::(keys, callback), + VectorStorageReadEnum::DenseChunkedByte(s) => s.read_vectors::(keys, callback), + VectorStorageReadEnum::DenseChunkedHalf(s) => s.read_vectors::(keys, callback), + VectorStorageReadEnum::MultiDenseChunked(s) => s.read_vectors::(keys, callback), + VectorStorageReadEnum::MultiDenseChunkedByte(s) => { + s.read_vectors::(keys, callback) + } + VectorStorageReadEnum::MultiDenseChunkedHalf(s) => { + s.read_vectors::(keys, callback) + } + VectorStorageReadEnum::Sparse(s) => s.read_vectors::(keys, callback), } } diff --git a/lib/segment/src/vector_storage/tests/test_appendable_dense_vector_storage.rs b/lib/segment/src/vector_storage/tests/test_appendable_dense_vector_storage.rs index 29c78e5636..a7fec3f610 100644 --- a/lib/segment/src/vector_storage/tests/test_appendable_dense_vector_storage.rs +++ b/lib/segment/src/vector_storage/tests/test_appendable_dense_vector_storage.rs @@ -119,7 +119,7 @@ fn do_test_delete_points(storage: &mut VectorStorageEnum) { 5, ); let closest = searcher - .peek_top_all(&DEFAULT_STOPPED, None) + .peek_top_all(&DEFAULT_STOPPED) .unwrap() .into_iter() .exactly_one() diff --git a/lib/segment/src/vector_storage/tests/test_appendable_multi_dense_vector_storage.rs b/lib/segment/src/vector_storage/tests/test_appendable_multi_dense_vector_storage.rs index 5e5774cce9..98a4d5aac5 100644 --- a/lib/segment/src/vector_storage/tests/test_appendable_multi_dense_vector_storage.rs +++ b/lib/segment/src/vector_storage/tests/test_appendable_multi_dense_vector_storage.rs @@ -182,7 +182,7 @@ fn do_test_delete_points(vector_dim: usize, vec_count: usize, storage: &mut Vect 5, ); let closest = searcher - .peek_top_all(&DEFAULT_STOPPED, None) + .peek_top_all(&DEFAULT_STOPPED) .unwrap() .pop() .unwrap(); diff --git a/lib/segment/src/vector_storage/vector_storage_base.rs b/lib/segment/src/vector_storage/vector_storage_base.rs index 65c726012f..0502626720 100644 --- a/lib/segment/src/vector_storage/vector_storage_base.rs +++ b/lib/segment/src/vector_storage/vector_storage_base.rs @@ -93,15 +93,21 @@ pub trait VectorStorageRead { /// Get the vector by the given key with potential optimizations for sequential reads. fn get_vector(&self, key: PointOffsetType) -> CowVector<'_>; - /// Get multiple vectors by the given keys + /// Get multiple vectors by the given keys. /// Potentially optimized for internal parallel reads. - fn read_vectors( + /// + /// Each input is tagged with caller-supplied user data `U` (e.g., an + /// external point id, or the input position). The data is threaded + /// straight back to the callback alongside its offset and vector, so + /// callers can map results into a parallel input array without keeping a + /// separate `offset → ...` lookup table. + fn read_vectors( &self, - keys: impl IntoIterator, - mut callback: impl FnMut(PointOffsetType, CowVector<'_>), + keys: impl IntoIterator, + mut callback: impl FnMut(U, PointOffsetType, CowVector<'_>), ) { - for key in keys { - callback(key, self.get_vector::

(key)); + for (user_data, key) in keys { + callback(user_data, key, self.get_vector::

(key)); } } @@ -778,47 +784,53 @@ impl VectorStorageRead for VectorStorageEnum { } } - fn read_vectors( + fn read_vectors( &self, - keys: impl IntoIterator, - callback: impl FnMut(PointOffsetType, CowVector<'_>), + keys: impl IntoIterator, + callback: impl FnMut(U, PointOffsetType, CowVector<'_>), ) { match self { - VectorStorageEnum::DenseVolatile(v) => v.read_vectors::

(keys, callback), + VectorStorageEnum::DenseVolatile(v) => v.read_vectors::(keys, callback), #[cfg(test)] - VectorStorageEnum::DenseVolatileByte(v) => v.read_vectors::

(keys, callback), + VectorStorageEnum::DenseVolatileByte(v) => v.read_vectors::(keys, callback), #[cfg(test)] - VectorStorageEnum::DenseVolatileHalf(v) => v.read_vectors::

(keys, callback), - VectorStorageEnum::DenseMemmap(v) => v.read_vectors::

(keys, callback), - VectorStorageEnum::DenseMemmapByte(v) => v.read_vectors::

(keys, callback), - VectorStorageEnum::DenseMemmapHalf(v) => v.read_vectors::

(keys, callback), + VectorStorageEnum::DenseVolatileHalf(v) => v.read_vectors::(keys, callback), + VectorStorageEnum::DenseMemmap(v) => v.read_vectors::(keys, callback), + VectorStorageEnum::DenseMemmapByte(v) => v.read_vectors::(keys, callback), + VectorStorageEnum::DenseMemmapHalf(v) => v.read_vectors::(keys, callback), #[cfg(target_os = "linux")] - VectorStorageEnum::DenseUring(v) => v.read_vectors::

(keys, callback), + VectorStorageEnum::DenseUring(v) => v.read_vectors::(keys, callback), #[cfg(target_os = "linux")] - VectorStorageEnum::DenseUringByte(v) => v.read_vectors::

(keys, callback), + VectorStorageEnum::DenseUringByte(v) => v.read_vectors::(keys, callback), #[cfg(target_os = "linux")] - VectorStorageEnum::DenseUringHalf(v) => v.read_vectors::

(keys, callback), + VectorStorageEnum::DenseUringHalf(v) => v.read_vectors::(keys, callback), - VectorStorageEnum::DenseAppendableMemmap(v) => v.read_vectors::

(keys, callback), - VectorStorageEnum::DenseAppendableMemmapByte(v) => v.read_vectors::

(keys, callback), - VectorStorageEnum::DenseAppendableMemmapHalf(v) => v.read_vectors::

(keys, callback), - VectorStorageEnum::SparseVolatile(v) => v.read_vectors::

(keys, callback), - VectorStorageEnum::SparseMmap(v) => v.read_vectors::

(keys, callback), - VectorStorageEnum::MultiDenseVolatile(v) => v.read_vectors::

(keys, callback), + VectorStorageEnum::DenseAppendableMemmap(v) => v.read_vectors::(keys, callback), + VectorStorageEnum::DenseAppendableMemmapByte(v) => { + v.read_vectors::(keys, callback) + } + VectorStorageEnum::DenseAppendableMemmapHalf(v) => { + v.read_vectors::(keys, callback) + } + VectorStorageEnum::SparseVolatile(v) => v.read_vectors::(keys, callback), + VectorStorageEnum::SparseMmap(v) => v.read_vectors::(keys, callback), + VectorStorageEnum::MultiDenseVolatile(v) => v.read_vectors::(keys, callback), #[cfg(test)] - VectorStorageEnum::MultiDenseVolatileByte(v) => v.read_vectors::

(keys, callback), + VectorStorageEnum::MultiDenseVolatileByte(v) => v.read_vectors::(keys, callback), #[cfg(test)] - VectorStorageEnum::MultiDenseVolatileHalf(v) => v.read_vectors::

(keys, callback), - VectorStorageEnum::MultiDenseAppendableMemmap(v) => v.read_vectors::

(keys, callback), + VectorStorageEnum::MultiDenseVolatileHalf(v) => v.read_vectors::(keys, callback), + VectorStorageEnum::MultiDenseAppendableMemmap(v) => { + v.read_vectors::(keys, callback) + } VectorStorageEnum::MultiDenseAppendableMemmapByte(v) => { - v.read_vectors::

(keys, callback) + v.read_vectors::(keys, callback) } VectorStorageEnum::MultiDenseAppendableMemmapHalf(v) => { - v.read_vectors::

(keys, callback) + v.read_vectors::(keys, callback) } - VectorStorageEnum::EmptyDense(v) => v.read_vectors::

(keys, callback), - VectorStorageEnum::EmptySparse(v) => v.read_vectors::

(keys, callback), + VectorStorageEnum::EmptyDense(v) => v.read_vectors::(keys, callback), + VectorStorageEnum::EmptySparse(v) => v.read_vectors::(keys, callback), } } diff --git a/lib/segment/tests/integration/exact_search_test.rs b/lib/segment/tests/integration/exact_search_test.rs index 5fe275227d..98743156db 100644 --- a/lib/segment/tests/integration/exact_search_test.rs +++ b/lib/segment/tests/integration/exact_search_test.rs @@ -113,7 +113,7 @@ fn exact_search_test() { let px = payload_index_ptr.borrow(); let filter = Filter::new_must(Condition::Field(block.condition.clone())); let points = px - .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped)) .unwrap(); for point in points { coverage.insert(point, coverage.get(&point).unwrap_or(&0) + 1); diff --git a/lib/segment/tests/integration/filtrable_hnsw_test.rs b/lib/segment/tests/integration/filtrable_hnsw_test.rs index 45cbc5d1c3..512204ff6f 100644 --- a/lib/segment/tests/integration/filtrable_hnsw_test.rs +++ b/lib/segment/tests/integration/filtrable_hnsw_test.rs @@ -136,7 +136,7 @@ fn _test_filterable_hnsw( for block in &blocks { let filter = Filter::new_must(Condition::Field(block.condition.clone())); let points = px - .with_view(|v| v.query_points(&filter, &hw_counter, &stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &stopped)) .unwrap(); for point in points { coverage.insert(point, coverage.get(&point).unwrap_or(&0) + 1); diff --git a/lib/segment/tests/integration/nested_filtering_test.rs b/lib/segment/tests/integration/nested_filtering_test.rs index 23587b513f..61889a5916 100644 --- a/lib/segment/tests/integration/nested_filtering_test.rs +++ b/lib/segment/tests/integration/nested_filtering_test.rs @@ -146,7 +146,7 @@ fn test_filtering_context_consistency() { let nested_filter_0 = Filter::new_must(nested_condition_0); let (res0, check_res0) = index.with_view(|v| { let res0 = v - .query_points(&nested_filter_0, &hw_counter, &is_stopped, None) + .query_points(&nested_filter_0, &hw_counter, &is_stopped) .unwrap(); let filter_context = v.filter_context(&nested_filter_0, &hw_counter).unwrap(); let check_res0: Vec<_> = (0..NUM_POINTS as PointOffsetType) @@ -187,7 +187,7 @@ fn test_filtering_context_consistency() { let (res1, check_res1) = index.with_view(|v| { let res1 = v - .query_points(&nested_filter_1, &hw_counter, &is_stopped, None) + .query_points(&nested_filter_1, &hw_counter, &is_stopped) .unwrap(); let filter_context = v.filter_context(&nested_filter_1, &hw_counter).unwrap(); let check_res1: Vec<_> = (0..NUM_POINTS as PointOffsetType) @@ -225,7 +225,7 @@ fn test_filtering_context_consistency() { let (res2, check_res2) = index.with_view(|v| { let res2 = v - .query_points(&nested_filter_2, &hw_counter, &is_stopped, None) + .query_points(&nested_filter_2, &hw_counter, &is_stopped) .unwrap(); let filter_context = v.filter_context(&nested_filter_2, &hw_counter).unwrap(); let check_res2: Vec<_> = (0..NUM_POINTS as PointOffsetType) @@ -273,7 +273,7 @@ fn test_filtering_context_consistency() { let (res3, check_res3) = index.with_view(|v| { let res3 = v - .query_points(&nested_filter_3, &hw_counter, &is_stopped, None) + .query_points(&nested_filter_3, &hw_counter, &is_stopped) .unwrap(); let filter_context = v.filter_context(&nested_filter_3, &hw_counter).unwrap(); let check_res3: Vec<_> = (0..NUM_POINTS as PointOffsetType) diff --git a/lib/segment/tests/integration/payload_index_test.rs b/lib/segment/tests/integration/payload_index_test.rs index 354aaa213c..45583847be 100644 --- a/lib/segment/tests/integration/payload_index_test.rs +++ b/lib/segment/tests/integration/payload_index_test.rs @@ -636,7 +636,7 @@ fn test_is_empty_conditions(test_segments: &TestSegments) -> Result<()> { .plain_segment .payload_index .borrow() - .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped)) .unwrap(); let real_number = plain_result.len(); @@ -646,7 +646,7 @@ fn test_is_empty_conditions(test_segments: &TestSegments) -> Result<()> { .struct_segment .payload_index .borrow() - .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None)) + .with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped)) .unwrap() .into_iter() // null index does not track deleted points, so we need to filter them out here. In callsites, diff --git a/lib/segment/tests/integration/sparse_discover_test.rs b/lib/segment/tests/integration/sparse_discover_test.rs index 3538b6baa6..ab76b9d7f9 100644 --- a/lib/segment/tests/integration/sparse_discover_test.rs +++ b/lib/segment/tests/integration/sparse_discover_test.rs @@ -184,7 +184,6 @@ fn sparse_index_discover_test() { path: index_dir.path(), stopped: &stopped, tick_progress: || (), - deferred_internal_id: None, }) .unwrap(); @@ -219,7 +218,7 @@ fn sparse_index_discover_test() { let query_context = QueryContext::default(); let segment_query_context = query_context.get_segment_query_context(); - let vector_context = segment_query_context.get_vector_context(SPARSE_VECTOR_NAME, None); + let vector_context = segment_query_context.get_vector_context(SPARSE_VECTOR_NAME); let sparse_search_result = sparse_index .search(&[&sparse_query], None, top, None, &vector_context) @@ -302,7 +301,6 @@ fn sparse_index_hardware_measurement_test() { path: index_dir.path(), stopped: &stopped, tick_progress: || (), - deferred_internal_id: None, }) .unwrap(); @@ -312,7 +310,7 @@ fn sparse_index_hardware_measurement_test() { let query_context = QueryContext::default(); let segment_query_context = query_context.get_segment_query_context(); - let vector_context = segment_query_context.get_vector_context(SPARSE_VECTOR_NAME, None); + let vector_context = segment_query_context.get_vector_context(SPARSE_VECTOR_NAME); let cpu_usage = query_context.hardware_usage_accumulator().get_cpu(); assert_eq!(cpu_usage, 0); diff --git a/lib/segment/tests/integration/sparse_vector_index_search_tests.rs b/lib/segment/tests/integration/sparse_vector_index_search_tests.rs index fd629ae9ca..8a6212ded3 100644 --- a/lib/segment/tests/integration/sparse_vector_index_search_tests.rs +++ b/lib/segment/tests/integration/sparse_vector_index_search_tests.rs @@ -230,7 +230,6 @@ fn sparse_vector_index_consistent_with_storage() { path: mmap_index_dir.path(), stopped: &stopped, tick_progress: || (), - deferred_internal_id: None, }) .unwrap(); @@ -257,7 +256,6 @@ fn sparse_vector_index_consistent_with_storage() { path: mmap_index_dir.path(), stopped: &stopped, tick_progress: || (), - deferred_internal_id: None, }) .unwrap(); @@ -694,7 +692,6 @@ fn check_persistence( path: inverted_index_dir.path(), stopped: &stopped, tick_progress: || (), - deferred_internal_id: None, }) .unwrap() };