Skip to main content

qdrant_edge/segment/id_tracker/
mod.rs

1pub 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/// Calling [`for_each_unique_point`] will yield this struct for each unique
16/// point.
17#[derive(Debug, Clone, Copy)]
18pub struct MergedPointId {
19    /// Unique external ID. If the same external ID is present in multiple
20    /// trackers, the item with the highest version takes precedence.
21    pub external_id: ExtendedPointId,
22    /// An index within `id_trackers` iterator that points to the [`IdTracker`]
23    /// that contains this point.
24    pub tracker_index: usize,
25    /// The internal ID of the point within the [`IdTracker`] that contains it.
26    pub internal_id: PointOffsetType,
27    /// The version of the point within the [`IdTracker`] that contains it.
28    pub version: u64,
29    /// Whether the point is marked deleted in the source tracker. Used by
30    /// [`for_each_unique_point`] to skip yielding a point whose
31    /// highest-versioned copy is a tombstone — without this filter,
32    /// `segment_builder` would resurrect deletes that were applied to a source
33    /// segment between when it was first picked for optimization and when the
34    /// merge actually reads from it (e.g. via snapshot teardown's
35    /// `propagate_to_wrapped`).
36    pub is_deleted: bool,
37}
38
39/// Calls a closure for each unique point from multiple ID trackers.
40///
41/// Discards points that have no version (their flush was interrupted) and
42/// points whose highest-versioned copy across the input trackers is marked
43/// deleted (so deletes applied to a source between optimization start and
44/// merge are not silently dropped).
45pub 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                    // a point without a version had an interrupted flush sequence and should be discarded
57                    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        // Naive HashMap-based implementation of for_each_unique_point.
121        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            // Ensure generated id_trackers have a lot of common points, so we
149            // are testing the merge logic.
150            assert!(collisions > 500);
151        } else {
152            // No collisions expected for a single or no id_trackers.
153            assert_eq!(collisions, 0);
154        }
155        if tracker_count == 0 {
156            assert!(expected.is_empty());
157        }
158
159        // `for_each_unique_point` skips winners whose source has them marked
160        // deleted, so drop those from the expected set before comparing.
161        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}