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