qdrant_edge/segment/id_tracker/
mod.rs1pub mod compressed;
2pub mod id_tracker_base;
3pub mod immutable_id_tracker;
4pub mod in_memory_id_tracker;
5mod memory_reporter;
6pub mod mutable_id_tracker;
7pub mod point_mappings;
8
9use crate::common::types::PointOffsetType;
10pub use id_tracker_base::*;
11use itertools::Itertools as _;
12
13use crate::segment::types::{ExtendedPointId, PointIdType};
14
15#[derive(Debug, Clone, Copy)]
18pub struct MergedPointId {
19 pub external_id: ExtendedPointId,
22 pub tracker_index: usize,
25 pub internal_id: PointOffsetType,
27 pub version: u64,
29 pub is_deleted: bool,
37}
38
39pub fn for_each_unique_point<'a>(
46 id_trackers: impl Iterator<Item = &'a (impl IdTracker + ?Sized + 'a)>,
47 mut f: impl FnMut(MergedPointId),
48) {
49 let mut iter = id_trackers
50 .enumerate()
51 .map(|(segment_index, id_tracker)| {
52 id_tracker.point_mappings().iter_from(None).filter_map(
53 move |(external_id, internal_id)| {
54 let version = id_tracker.internal_version(internal_id);
55 let is_deleted = id_tracker.is_deleted_point(internal_id);
56 version.map(|version| MergedPointId {
58 external_id,
59 tracker_index: segment_index,
60 internal_id,
61 version,
62 is_deleted,
63 })
64 },
65 )
66 })
67 .kmerge_by(|a, b| a.external_id < b.external_id);
68
69 let Some(mut best_item) = iter.next() else {
70 return;
71 };
72
73 for item in iter {
74 if best_item.external_id == item.external_id {
75 if best_item.version < item.version {
76 best_item = item;
77 }
78 } else {
79 if !best_item.is_deleted {
80 f(best_item);
81 }
82 best_item = item;
83 }
84 }
85 if !best_item.is_deleted {
86 f(best_item);
87 }
88}
89
90impl From<&ExtendedPointId> for PointIdType {
91 fn from(point_id: &ExtendedPointId) -> Self {
92 match point_id {
93 ExtendedPointId::NumId(idx) => PointIdType::NumId(*idx),
94 ExtendedPointId::Uuid(uuid) => PointIdType::Uuid(*uuid),
95 }
96 }
97}
98
99#[cfg(test)]
100mod tests {
101 use std::collections::{HashMap, hash_map};
102
103 use in_memory_id_tracker::InMemoryIdTracker;
104 use rand::SeedableRng as _;
105 use rand::rngs::StdRng;
106 use rstest::rstest;
107
108 use super::*;
109
110 #[rstest]
111 fn test_for_each_unique_point(#[values(0, 1, 5)] tracker_count: usize) {
112 let mut rand = StdRng::seed_from_u64(42);
113
114 let id_trackers = (0..tracker_count)
115 .map(|_| InMemoryIdTracker::random(&mut rand, 1000, 500, 10))
116 .collect_vec();
117
118 let mut collisions = 0;
119
120 let mut expected = HashMap::<ExtendedPointId, MergedPointId>::new();
122 for (tracker_index, id_tracker) in id_trackers.iter().enumerate() {
123 for (external_id, internal_id) in id_tracker.point_mappings().iter_from(None) {
124 let version = id_tracker.internal_version(internal_id).unwrap();
125 let is_deleted = id_tracker.is_deleted_point(internal_id);
126 let merged_point_id = MergedPointId {
127 external_id,
128 tracker_index,
129 internal_id,
130 version,
131 is_deleted,
132 };
133 match expected.entry(external_id) {
134 hash_map::Entry::Occupied(mut entry) => {
135 collisions += 1;
136 if entry.get().version < version {
137 entry.insert(merged_point_id);
138 }
139 }
140 hash_map::Entry::Vacant(entry) => {
141 entry.insert(merged_point_id);
142 }
143 }
144 }
145 }
146
147 if tracker_count > 1 {
148 assert!(collisions > 500);
151 } else {
152 assert_eq!(collisions, 0);
154 }
155 if tracker_count == 0 {
156 assert!(expected.is_empty());
157 }
158
159 expected.retain(|_, v| !v.is_deleted);
162
163 for_each_unique_point(id_trackers.iter(), |merged_point_id| {
164 let v = expected.remove(&merged_point_id.external_id).unwrap();
165 assert_eq!(merged_point_id.tracker_index, v.tracker_index);
166 assert_eq!(merged_point_id.internal_id, v.internal_id);
167 assert_eq!(merged_point_id.version, v.version);
168 assert!(!merged_point_id.is_deleted);
169 });
170
171 assert!(expected.is_empty());
172 }
173}