1use std::collections::HashMap;
16use std::ops::Range;
17
18use common_time::Timestamp;
19use store_api::metadata::{RegionMetadata, RegionMetadataRef};
20
21use crate::compaction::run::primary_key_ranges_overlap;
22use crate::sst::file::{FileHandle, RegionFileId};
23use crate::sst::primary_key::PrimaryKeyRangeMapper;
24
25pub(super) struct FileOverlapIndex<'a> {
30 metadata: &'a RegionMetadata,
32 primary_key_mapper: PrimaryKeyRangeMapper,
33 files: Vec<FileHandle>,
36 by_start: Vec<usize>,
39 positions: HashMap<RegionFileId, usize>,
42 max_ends: Vec<Option<Timestamp>>,
47 leaf_base: usize,
50}
51
52struct OverlapQuery<'a> {
53 input: &'a FileHandle,
54 min_end: Timestamp,
55 candidates_end: usize,
56}
57
58impl<'a> FileOverlapIndex<'a> {
59 pub(super) fn new(files: Vec<FileHandle>, metadata: &'a RegionMetadataRef) -> Self {
60 let mut by_start: Vec<_> = (0..files.len()).collect();
61 by_start.sort_unstable_by_key(|&i| (files[i].time_range().0, i));
62 let leaf_base = files.len().next_power_of_two();
64 let mut max_ends = vec![None; 2 * leaf_base];
65 let mut positions = HashMap::with_capacity(files.len());
66 for (offset, &i) in by_start.iter().enumerate() {
67 max_ends[leaf_base + offset] = Some(files[i].time_range().1);
68 positions.insert(files[i].file_id(), leaf_base + offset);
69 }
70 for node in (1..leaf_base).rev() {
71 max_ends[node] = max_ends[2 * node].max(max_ends[2 * node + 1]);
72 }
73 Self {
74 metadata,
75 primary_key_mapper: PrimaryKeyRangeMapper::new(metadata.clone()),
76 files,
77 by_start,
78 positions,
79 max_ends,
80 leaf_base,
81 }
82 }
83
84 pub(super) fn remove(&mut self, file_id: RegionFileId) {
85 let Some(mut node) = self.positions.remove(&file_id) else {
86 return;
87 };
88 self.max_ends[node] = None;
89 while node > 1 {
90 node /= 2;
91 self.max_ends[node] = self.max_ends[2 * node].max(self.max_ends[2 * node + 1]);
92 }
93 }
94
95 pub(super) fn drain_overlaps(&mut self, input: &FileHandle) -> Vec<FileHandle> {
99 let mut matches = Vec::new();
100 while let Some(i) = self.find_overlap(input) {
101 self.remove(self.files[i].file_id());
102 matches.push(i);
103 }
104 matches.sort_unstable();
105 matches.into_iter().map(|i| self.files[i].clone()).collect()
106 }
107
108 fn find_overlap(&self, input: &FileHandle) -> Option<usize> {
109 let (start, end) = input.time_range();
110 let query = OverlapQuery {
111 input,
112 min_end: start,
113 candidates_end: self
114 .by_start
115 .partition_point(|&i| self.files[i].time_range().0 <= end),
116 };
117 self.find_in_subtree(1, 0..self.leaf_base, &query)
118 }
119
120 fn find_in_subtree(
121 &self,
122 node: usize,
123 leaves: Range<usize>,
124 query: &OverlapQuery<'_>,
125 ) -> Option<usize> {
126 if leaves.start >= query.candidates_end
127 || self.max_ends[node].is_none_or(|end| end < query.min_end)
128 {
129 return None;
130 }
131 if leaves.len() == 1 {
132 let i = self.by_start[leaves.start];
133 return files_may_overlap(
134 query.input,
135 &self.files[i],
136 self.metadata,
137 &self.primary_key_mapper,
138 )
139 .then_some(i);
140 }
141 let mid = leaves.start + leaves.len() / 2;
142 self.find_in_subtree(2 * node, leaves.start..mid, query)
143 .or_else(|| self.find_in_subtree(2 * node + 1, mid..leaves.end, query))
144 }
145}
146
147fn files_may_overlap(
150 lhs: &FileHandle,
151 rhs: &FileHandle,
152 metadata: &RegionMetadata,
153 mapper: &PrimaryKeyRangeMapper,
154) -> bool {
155 let (lhs_start, lhs_end) = lhs.time_range();
156 let (rhs_start, rhs_end) = rhs.time_range();
157 if lhs_start.max(rhs_start) > lhs_end.min(rhs_end) {
158 return false;
159 }
160
161 if metadata.schema_version != 0
167 || lhs.region_id() != metadata.region_id
168 || rhs.region_id() != metadata.region_id
169 {
170 return true;
171 }
172 match (lhs.primary_key_range(mapper), rhs.primary_key_range(mapper)) {
173 (Some(lhs), Some(rhs)) if lhs.0 <= lhs.1 && rhs.0 <= rhs.1 => {
174 primary_key_ranges_overlap(&lhs, &rhs)
175 }
176 _ => true,
177 }
178}
179
180#[cfg(test)]
181mod tests {
182 use std::collections::HashSet;
183 use std::sync::Arc;
184
185 use rand::{Rng, SeedableRng};
186 use store_api::storage::FileId;
187
188 use super::*;
189 use crate::compaction::test_util::{pk_range, primary_key_metadata_for_test};
190 use crate::sst::file::FileMeta;
191 use crate::test_util::memtable_util::metadata_for_test;
192 use crate::test_util::new_noop_file_purger;
193
194 fn file(start: i64, end: i64, pk: Option<(&str, &str)>) -> FileHandle {
195 FileHandle::new_with_primary_key_range(
196 FileMeta {
197 file_id: FileId::random(),
198 time_range: (
199 Timestamp::new_millisecond(start),
200 Timestamp::new_millisecond(end),
201 ),
202 ..Default::default()
203 },
204 new_noop_file_purger(),
205 pk.and_then(|(start, end)| pk_range(start.as_bytes(), end.as_bytes())),
206 )
207 }
208
209 #[rstest::rstest]
210 #[case(10, Some(("b", "c")), 0, false, true)]
211 #[case(11, Some(("b", "c")), 0, false, false)]
212 #[case(10, Some(("c", "d")), 0, false, false)]
213 #[case(10, None, 0, false, true)]
214 #[case(10, Some(("z", "a")), 0, false, true)]
215 #[case(10, Some(("c", "d")), 1, false, true)]
216 #[case(10, Some(("c", "d")), 0, true, true)]
217 #[case(11, None, 1, true, false)]
218 fn test_conservative_overlap(
219 #[case] start: i64,
220 #[case] pk: Option<(&str, &str)>,
221 #[case] schema_version: u64,
222 #[case] foreign: bool,
223 #[case] expected: bool,
224 ) {
225 let mut metadata = (*primary_key_metadata_for_test()).clone();
226 metadata.region_id = 0.into();
227 metadata.schema_version = schema_version;
228 let metadata = Arc::new(metadata);
229 let ranges = PrimaryKeyRangeMapper::new(metadata.clone());
230 let lhs = file(0, 10, Some(("a", "b")));
231 let mut rhs = file(start, start + 10, pk);
232 if foreign {
233 let mut meta = rhs.meta_ref().clone();
234 meta.region_id = 1.into();
235 rhs = FileHandle::new_with_primary_key_range(
236 meta,
237 new_noop_file_purger(),
238 rhs.raw_primary_key_range(),
239 );
240 }
241 assert_eq!(expected, files_may_overlap(&lhs, &rhs, &metadata, &ranges));
242 assert_eq!(expected, files_may_overlap(&rhs, &lhs, &metadata, &ranges));
243 }
244
245 #[test]
246 fn test_index_removal_distinguishes_region_owners() {
247 let metadata = metadata_for_test();
248 let lhs = file(0, 10, None);
249 let mut meta = lhs.meta_ref().clone();
250 meta.region_id = 1.into();
251 let rhs = FileHandle::new(meta, new_noop_file_purger());
252 let mut index = FileOverlapIndex::new(vec![lhs.clone(), rhs.clone()], &metadata);
253 index.remove(lhs.file_id());
254 let overlaps = index.drain_overlaps(&lhs);
255 assert_eq!(1, overlaps.len());
256 assert_eq!(rhs.file_id(), overlaps[0].file_id());
257 assert!(index.drain_overlaps(&lhs).is_empty());
258 }
259
260 #[test]
261 fn test_index_matches_linear_scan_after_removals() {
262 let metadata = primary_key_metadata_for_test();
263 let ranges = PrimaryKeyRangeMapper::new(metadata.clone());
264 let mut rng = rand::rngs::StdRng::seed_from_u64(9146);
265 for count in [0, 1, 7, 32, 127] {
266 let files: Vec<_> = (0..count)
267 .map(|_| {
268 let start = rng.random_range(-50..50);
269 let end = start + rng.random_range(0..40);
270 let pk = match rng.random_range(0..4) {
271 0 => None,
272 1 => Some(("a", "b")),
273 2 => Some(("b", "c")),
274 _ => Some(("z", "z")),
275 };
276 file(start, end, pk)
277 })
278 .collect();
279 let mut index = FileOverlapIndex::new(files.clone(), &metadata);
280 let mut removed = HashSet::new();
281 for f in files.iter().step_by(3) {
282 index.remove(f.file_id());
283 removed.insert(f.file_id());
284 }
285 for _ in 0..32 {
286 let start = rng.random_range(-60..60);
287 let query = file(start, start + rng.random_range(0..25), Some(("b", "c")));
288 let expected: Vec<_> = files
289 .iter()
290 .filter(|f| {
291 !removed.contains(&f.file_id())
292 && files_may_overlap(&query, f, &metadata, &ranges)
293 })
294 .map(FileHandle::file_id)
295 .collect();
296 let actual: Vec<_> = index
297 .drain_overlaps(&query)
298 .iter()
299 .map(FileHandle::file_id)
300 .collect();
301 assert_eq!(expected, actual, "snapshot size {count}");
302 removed.extend(expected);
303 }
304 }
305 }
306}