Skip to main content

qdrant_edge/segment/id_tracker/
mod.rs

1pub mod compressed;
2pub mod disk_id_tracker;
3pub mod format_detection;
4pub mod id_tracker_base;
5pub mod immutable_id_tracker;
6pub mod in_memory_id_tracker;
7mod memory_reporter;
8pub mod mutable_id_tracker;
9pub mod point_mappings;
10
11use crate::common::types::PointOffsetType;
12pub use format_detection::IdTrackerFormat;
13pub use id_tracker_base::*;
14use itertools::Itertools as _;
15
16use crate::segment::types::{ExtendedPointId, PointIdType};
17
18/// Calling [`for_each_unique_point`] will yield this struct for each unique
19/// point.
20#[derive(Debug, Clone, Copy)]
21pub struct MergedPointId {
22    /// Unique external ID. If the same external ID is present in multiple
23    /// trackers, the item with the highest version takes precedence.
24    pub external_id: ExtendedPointId,
25    /// An index within `id_trackers` iterator that points to the [`IdTracker`]
26    /// that contains this point.
27    pub tracker_index: usize,
28    /// The internal ID of the point within the [`IdTracker`] that contains it.
29    pub internal_id: PointOffsetType,
30    /// The version of the point within the [`IdTracker`] that contains it.
31    pub version: u64,
32    /// Whether the point is marked deleted in the source tracker. Used by
33    /// [`for_each_unique_point`] to skip yielding a point whose
34    /// highest-versioned copy is a tombstone — without this filter,
35    /// `segment_builder` would resurrect deletes that were applied to a source
36    /// segment between when it was first picked for optimization and when the
37    /// merge actually reads from it (e.g. via snapshot teardown's
38    /// `propagate_to_wrapped`).
39    pub is_deleted: bool,
40}
41
42/// Calls a closure for each unique point from multiple ID trackers.
43///
44/// Discards points that have no version (their flush was interrupted) and
45/// points whose highest-versioned copy across the input trackers is marked
46/// deleted (so deletes applied to a source between optimization start and
47/// merge are not silently dropped).
48pub fn for_each_unique_point<'a>(
49    id_trackers: impl Iterator<Item = &'a (impl IdTracker + ?Sized + 'a)>,
50    mut f: impl FnMut(MergedPointId),
51) {
52    let mut iter = id_trackers
53        .enumerate()
54        .map(|(segment_index, id_tracker)| {
55            id_tracker.point_mappings().iter_from(None).filter_map(
56                move |(external_id, internal_id)| {
57                    let version = id_tracker.internal_version(internal_id);
58                    let is_deleted = id_tracker.is_deleted_point(internal_id);
59                    // a point without a version had an interrupted flush sequence and should be discarded
60                    version.map(|version| MergedPointId {
61                        external_id,
62                        tracker_index: segment_index,
63                        internal_id,
64                        version,
65                        is_deleted,
66                    })
67                },
68            )
69        })
70        .kmerge_by(|a, b| a.external_id < b.external_id);
71
72    let Some(mut best_item) = iter.next() else {
73        return;
74    };
75
76    for item in iter {
77        if best_item.external_id == item.external_id {
78            if best_item.version < item.version {
79                best_item = item;
80            }
81        } else {
82            if !best_item.is_deleted {
83                f(best_item);
84            }
85            best_item = item;
86        }
87    }
88    if !best_item.is_deleted {
89        f(best_item);
90    }
91}
92
93impl From<&ExtendedPointId> for PointIdType {
94    fn from(point_id: &ExtendedPointId) -> Self {
95        match point_id {
96            ExtendedPointId::NumId(idx) => PointIdType::NumId(*idx),
97            ExtendedPointId::Uuid(uuid) => PointIdType::Uuid(*uuid),
98        }
99    }
100}
101
102#[cfg(test)]
103mod tests {
104    use std::collections::{HashMap, hash_map};
105
106    use in_memory_id_tracker::InMemoryIdTracker;
107    use rand::SeedableRng as _;
108    use rand::rngs::StdRng;
109    use rstest::rstest;
110    use tempfile::Builder;
111
112    use super::*;
113    use crate::segment::id_tracker::mutable_id_tracker::MutableIdTracker;
114
115    /// Review finding #1 (optimizer merge consequence): `for_each_unique_point`
116    /// is the merge primitive `segment_builder::update_from` uses to pick the
117    /// winning copy per external id when rolling segments together. It walks
118    /// `iter_from(None)` and keeps the highest version. For a shadowed point
119    /// (active head below the cutoff + deferred head above it, the result of an
120    /// in-place mutation), `iter_from` resolves the collision to the *active*
121    /// (older) slot and never yields the deferred head — so the merge keeps the
122    /// pre-mutation copy and the mutation is silently dropped on optimization.
123    ///
124    /// On `dev` this worked: a deferred write tombstoned the active, so the
125    /// single map pointed at the deferred head and the merge pulled the latest.
126    #[test]
127    fn for_each_unique_point_keeps_deferred_head_for_shadowed_point() {
128        let dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
129        // Appendable tracker with deferred cutoff = 5.
130        let mut tracker = MutableIdTracker::open(dir.path(), Some(5)).unwrap();
131        let p7 = PointIdType::NumId(7);
132
133        // ext 7: active@2 is the pre-mutation copy (version 5); deferred@9 is
134        // the latest in-place mutation (version 8). The active is shadowed.
135        tracker.set_link(p7, 2).unwrap();
136        tracker.set_internal_version(2, 5).unwrap();
137        tracker.set_link(p7, 9).unwrap();
138        tracker.set_internal_version(9, 8).unwrap();
139
140        let trackers = [tracker];
141        let mut merged = Vec::new();
142        for_each_unique_point(trackers.iter(), |m| {
143            merged.push((m.external_id, m.internal_id, m.version));
144        });
145
146        // The merge must carry the latest (deferred) copy so the mutation
147        // survives optimization. It instead yields the stale active copy
148        // (internal 2, version 5), silently dropping the version-8 mutation.
149        assert_eq!(
150            merged,
151            vec![(p7, 9, 8)],
152            "for_each_unique_point dropped the deferred (latest) copy of a \
153             shadowed point, keeping the pre-mutation active version",
154        );
155    }
156
157    #[rstest]
158    fn test_for_each_unique_point(#[values(0, 1, 5)] tracker_count: usize) {
159        let mut rand = StdRng::seed_from_u64(42);
160
161        let id_trackers = (0..tracker_count)
162            .map(|_| InMemoryIdTracker::random(&mut rand, 1000, 500, 10))
163            .collect_vec();
164
165        let mut collisions = 0;
166
167        // Naive HashMap-based implementation of for_each_unique_point.
168        let mut expected = HashMap::<ExtendedPointId, MergedPointId>::new();
169        for (tracker_index, id_tracker) in id_trackers.iter().enumerate() {
170            for (external_id, internal_id) in id_tracker.point_mappings().iter_from(None) {
171                let version = id_tracker.internal_version(internal_id).unwrap();
172                let is_deleted = id_tracker.is_deleted_point(internal_id);
173                let merged_point_id = MergedPointId {
174                    external_id,
175                    tracker_index,
176                    internal_id,
177                    version,
178                    is_deleted,
179                };
180                match expected.entry(external_id) {
181                    hash_map::Entry::Occupied(mut entry) => {
182                        collisions += 1;
183                        if entry.get().version < version {
184                            entry.insert(merged_point_id);
185                        }
186                    }
187                    hash_map::Entry::Vacant(entry) => {
188                        entry.insert(merged_point_id);
189                    }
190                }
191            }
192        }
193
194        if tracker_count > 1 {
195            // Ensure generated id_trackers have a lot of common points, so we
196            // are testing the merge logic.
197            assert!(collisions > 500);
198        } else {
199            // No collisions expected for a single or no id_trackers.
200            assert_eq!(collisions, 0);
201        }
202        if tracker_count == 0 {
203            assert!(expected.is_empty());
204        }
205
206        // `for_each_unique_point` skips winners whose source has them marked
207        // deleted, so drop those from the expected set before comparing.
208        expected.retain(|_, v| !v.is_deleted);
209
210        for_each_unique_point(id_trackers.iter(), |merged_point_id| {
211            let v = expected.remove(&merged_point_id.external_id).unwrap();
212            assert_eq!(merged_point_id.tracker_index, v.tracker_index);
213            assert_eq!(merged_point_id.internal_id, v.internal_id);
214            assert_eq!(merged_point_id.version, v.version);
215            assert!(!merged_point_id.is_deleted);
216        });
217
218        assert!(expected.is_empty());
219    }
220}