1use crate::lsm_tree::merge::ItemOp::{Discard, Keep, Replace};
6use crate::lsm_tree::merge::{MergeLayerIterator, MergeResult};
7use crate::lsm_tree::types::{Item, LayerIterator};
8use crate::object_store::Extent;
9use crate::object_store::allocator::{AllocatorKey, AllocatorValue};
10use anyhow::Error;
11use std::collections::HashSet;
12
13pub fn merge(
14 left: &MergeLayerIterator<'_, AllocatorKey, AllocatorValue>,
15 right: &MergeLayerIterator<'_, AllocatorKey, AllocatorValue>,
16) -> MergeResult<AllocatorKey, AllocatorValue> {
17 if left.key().device_range.end < right.key().device_range.start {
27 return MergeResult::EmitLeft;
28 }
29
30 if left.key().device_range.end == right.key().device_range.start {
35 if *left.value() == *right.value() {
37 return MergeResult::Other {
38 emit: None,
39 left: Discard,
40 right: Replace(
41 Item::new(
42 AllocatorKey {
43 device_range: Extent(
44 left.key().device_range.start..right.key().device_range.end,
45 ),
46 },
47 left.value().clone(),
48 )
49 .boxed(),
50 ),
51 };
52 } else {
53 return MergeResult::EmitLeft;
54 }
55 }
56 if left.key().device_range.start == right.key().device_range.start {
57 if left.key().device_range.end < right.key().device_range.end {
62 if left.layer_index < right.layer_index {
64 return MergeResult::Other {
65 emit: None,
66 left: Keep,
67 right: if left.key().device_range.end == right.key().device_range.end {
68 Discard
69 } else {
70 Replace(
71 Item::new(
72 AllocatorKey {
73 device_range: Extent(
74 left.key().device_range.end..right.key().device_range.end,
75 ),
76 },
77 right.value().clone(),
78 )
79 .boxed(),
80 )
81 },
82 };
83 } else {
84 return MergeResult::Other { emit: None, left: Discard, right: Keep };
86 }
87
88 } else {
93 if right.layer_index < left.layer_index {
95 return MergeResult::Other {
96 emit: None,
97 left: if right.key().device_range.end == left.key().device_range.end {
98 Discard
99 } else {
100 Replace(
101 Item::new(
102 AllocatorKey {
103 device_range: Extent(
104 right.key().device_range.end..left.key().device_range.end,
105 ),
106 },
107 left.value().clone(),
108 )
109 .boxed(),
110 )
111 },
112 right: Keep,
113 };
114 } else {
115 return MergeResult::Other { emit: None, left: Keep, right: Discard };
117 }
118 }
119 }
120 debug_assert!(left.key().device_range.end >= right.key().device_range.start);
125 MergeResult::Other {
126 emit: Some(
127 Item::new(
128 AllocatorKey {
129 device_range: Extent(
130 left.key().device_range.start..right.key().device_range.start,
131 ),
132 },
133 left.value().clone(),
134 )
135 .boxed(),
136 ),
137 left: Replace(
138 Item::new(
139 AllocatorKey {
140 device_range: Extent(
141 right.key().device_range.start..left.key().device_range.end,
142 ),
143 },
144 left.value().clone(),
145 )
146 .boxed(),
147 ),
148 right: Keep,
149 }
150}
151
152pub fn filter_tombstones<'a>(
153 iter: impl LayerIterator<AllocatorKey, AllocatorValue> + 'a,
154) -> impl Future<Output = Result<impl LayerIterator<AllocatorKey, AllocatorValue> + 'a, Error>> {
155 iter.filter(|i| *i.value != AllocatorValue::None)
156}
157
158pub fn filter_marked_for_deletion<'a>(
159 iter: impl LayerIterator<AllocatorKey, AllocatorValue> + 'a,
160 marked_for_deletion: HashSet<u64>,
161) -> impl Future<Output = Result<impl LayerIterator<AllocatorKey, AllocatorValue> + 'a, Error>> {
162 iter.filter(move |i| {
163 if let AllocatorValue::Abs { owner_object_id, .. } = i.value {
164 !marked_for_deletion.contains(owner_object_id)
165 } else {
166 true
167 }
168 })
169}
170
171#[cfg(test)]
172mod tests {
173 use crate::lsm_tree::types::{Item, ItemRef, LayerIterator};
174 use crate::lsm_tree::{LSMTree, Query};
175 use crate::object_store::allocator::merge::{filter_tombstones, merge};
176 use crate::object_store::allocator::{AllocatorKey, AllocatorValue};
177 use std::ops::Range;
178
179 async fn test_merge(
181 left: (Range<u64>, AllocatorValue),
182 right: (Range<u64>, AllocatorValue),
183 expected: &[(Range<u64>, AllocatorValue)],
184 ) {
185 let tree = LSMTree::new(merge, None);
186 tree.insert(Item::new(AllocatorKey { device_range: right.0.into() }, right.1))
187 .expect("insert error");
188 tree.seal();
189 tree.insert(Item::new(AllocatorKey { device_range: left.0.into() }, left.1))
190 .expect("insert error");
191 let layer_set = tree.layer_set();
192 let mut merger = layer_set.merger();
193 let mut iter = filter_tombstones(merger.query(Query::FullScan).await.expect("seek failed"))
194 .await
195 .expect("filter failed");
196 for e in expected {
197 let ItemRef { key, value, .. } = iter.get().expect("get failed");
198 assert_eq!((key, value), (&AllocatorKey { device_range: e.0.clone().into() }, &e.1));
199 iter.advance().await.expect("advance failed");
200 }
201 assert!(iter.get().is_none());
202 }
203
204 #[fuchsia::test]
205 async fn test_no_overlap() {
206 test_merge(
207 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
208 (200..300, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
209 &[
210 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
211 (200..300, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
212 ],
213 )
214 .await;
215 }
216
217 #[fuchsia::test]
218 async fn test_touching() {
219 test_merge(
220 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
221 (100..200, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
222 &[(0..200, AllocatorValue::Abs { count: 1, owner_object_id: 1 })],
223 )
224 .await;
225 }
226
227 #[fuchsia::test]
228 async fn test_identical() {
229 test_merge(
230 (0..100, AllocatorValue::Abs { count: 2, owner_object_id: 1 }),
231 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
232 &[(0..100, AllocatorValue::Abs { count: 2, owner_object_id: 1 })],
233 )
234 .await;
235 test_merge(
236 (0..100, AllocatorValue::None),
237 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
238 &[],
239 )
240 .await;
241 }
242
243 #[fuchsia::test]
244 async fn test_left_smaller_than_right_with_same_start() {
245 test_merge(
246 (0..100, AllocatorValue::Abs { count: 2, owner_object_id: 1 }),
247 (0..200, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
248 &[
249 (0..100, AllocatorValue::Abs { count: 2, owner_object_id: 1 }),
250 (100..200, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
251 ],
252 )
253 .await;
254 test_merge(
255 (0..100, AllocatorValue::None),
256 (0..200, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
257 &[(100..200, AllocatorValue::Abs { count: 1, owner_object_id: 1 })],
258 )
259 .await;
260 }
261
262 #[fuchsia::test]
263 async fn test_left_starts_before_right_with_overlap() {
264 test_merge(
265 (0..200, AllocatorValue::Abs { count: 2, owner_object_id: 1 }),
266 (100..150, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
267 &[
268 (0..100, AllocatorValue::Abs { count: 2, owner_object_id: 1 }),
269 (100..200, AllocatorValue::Abs { count: 2, owner_object_id: 1 }),
270 ],
271 )
272 .await;
273 }
274
275 #[fuchsia::test]
276 async fn test_different_object_id() {
277 test_merge(
279 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
280 (200..300, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
281 &[
282 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
283 (200..300, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
284 ],
285 )
286 .await;
287 test_merge(
289 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
290 (100..200, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
291 &[
292 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
293 (100..200, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
294 ],
295 )
296 .await;
297 test_merge(
299 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
300 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
301 &[(0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 })],
302 )
303 .await;
304 test_merge(
306 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
307 (0..200, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
308 &[
309 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
310 (100..200, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
311 ],
312 )
313 .await;
314 test_merge(
316 (0..200, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
317 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
318 &[(0..200, AllocatorValue::Abs { count: 1, owner_object_id: 1 })],
319 )
320 .await;
321 test_merge(
323 (0..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
324 (50..150, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
325 &[
326 (0..50, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
327 (50..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
328 (100..150, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
329 ],
330 )
331 .await;
332 }
333
334 #[fuchsia::test]
335 async fn test_tombstones() {
336 let key = AllocatorKey { device_range: (0..100 * 512).into() };
344 let lower_bound = AllocatorKey::lower_bound_for_merge_into(&key);
345 let tree = LSMTree::new(merge, None);
346 tree.merge_into(
347 Item::new(key.clone(), AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
348 &lower_bound,
349 );
350 tree.seal();
351 tree.merge_into(
352 Item::new(key.clone(), AllocatorValue::Abs { count: 2, owner_object_id: 1 }),
353 &lower_bound,
354 );
355 tree.seal();
356 tree.merge_into(Item::new(key.clone(), AllocatorValue::None), &lower_bound);
357 tree.merge_into(
358 Item::new(key.clone(), AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
359 &lower_bound,
360 );
361 tree.seal();
362 tree.merge_into(Item::new(key.clone(), AllocatorValue::None), &lower_bound);
363 tree.merge_into(
364 Item::new(key.clone(), AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
365 &lower_bound,
366 );
367 let layer_set = tree.layer_set();
368 let mut merger = layer_set.merger();
369 let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
370 let ItemRef { key: k, value, .. } = iter.get().expect("get failed");
371 assert_eq!((k, value), (&key, &AllocatorValue::Abs { count: 1, owner_object_id: 1 }));
372 iter.advance().await.expect("advance failed");
373 assert!(iter.get().is_none());
374 }
375
376 #[fuchsia::test]
377 async fn test_merge_adjacent_in_mutable_layer() {
378 {
380 let tree = LSMTree::new(merge, None);
381
382 let key1 = AllocatorKey { device_range: (4096..135168).into() };
383 let val1 = AllocatorValue::Abs { count: 1, owner_object_id: 3 };
384 tree.merge_into(
385 Item::new(key1.clone(), val1.clone()),
386 &key1.lower_bound_for_merge_into(),
387 );
388
389 let key2 = AllocatorKey { device_range: (135168..139264).into() };
390 let val2 = AllocatorValue::Abs { count: 1, owner_object_id: 3 };
391 tree.merge_into(
392 Item::new(key2.clone(), val2.clone()),
393 &key2.lower_bound_for_merge_into(),
394 );
395
396 let layer_set = tree.layer_set();
398 let mut merger = layer_set.merger();
399 let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
400
401 let ItemRef { key, value, .. } = iter.get().expect("get failed");
402 assert_eq!(
403 (key, value),
404 (
405 &AllocatorKey { device_range: (4096..139264).into() },
406 &AllocatorValue::Abs { count: 1, owner_object_id: 3 }
407 )
408 );
409 iter.advance().await.expect("advance failed");
410 assert!(iter.get().is_none());
411 }
412
413 {
415 let tree = LSMTree::new(merge, None);
416
417 let key2 = AllocatorKey { device_range: (135168..139264).into() };
418 let val2 = AllocatorValue::Abs { count: 1, owner_object_id: 3 };
419 tree.merge_into(
420 Item::new(key2.clone(), val2.clone()),
421 &key2.lower_bound_for_merge_into(),
422 );
423
424 let key1 = AllocatorKey { device_range: (4096..135168).into() };
425 let val1 = AllocatorValue::Abs { count: 1, owner_object_id: 3 };
426 tree.merge_into(
427 Item::new(key1.clone(), val1.clone()),
428 &key1.lower_bound_for_merge_into(),
429 );
430
431 let layer_set = tree.layer_set();
433 let mut merger = layer_set.merger();
434 let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
435
436 let ItemRef { key, value, .. } = iter.get().expect("get failed");
437 assert_eq!(
438 (key, value),
439 (
440 &AllocatorKey { device_range: (4096..139264).into() },
441 &AllocatorValue::Abs { count: 1, owner_object_id: 3 }
442 )
443 );
444 iter.advance().await.expect("advance failed");
445 assert!(iter.get().is_none());
446 }
447 }
448
449 #[fuchsia::test]
450 async fn test_overlapping_boundaries() {
451 let base = (50..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 });
452
453 let test_val = AllocatorValue::Abs { count: 1, owner_object_id: 2 };
454
455 test_merge(
460 base.clone(),
461 (49..100, test_val.clone()),
462 &[
463 (49..50, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
464 (50..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
465 ],
466 )
467 .await;
468 test_merge(
470 (49..100, test_val.clone()),
471 base.clone(),
472 &[
473 (49..50, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
474 (50..100, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
475 ],
476 )
477 .await;
478
479 test_merge(
482 base.clone(),
483 (51..100, test_val.clone()),
484 &[
485 (50..51, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
486 (51..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
487 ],
488 )
489 .await;
490 test_merge(
492 (51..100, test_val.clone()),
493 base.clone(),
494 &[
495 (50..51, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
496 (51..100, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
497 ],
498 )
499 .await;
500
501 test_merge(
505 base.clone(),
506 (50..99, test_val.clone()),
507 &[(50..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 })],
508 )
509 .await;
510 test_merge(
512 (50..99, test_val.clone()),
513 base.clone(),
514 &[
515 (50..99, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
516 (99..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
517 ],
518 )
519 .await;
520
521 test_merge(
524 base.clone(),
525 (50..101, test_val.clone()),
526 &[
527 (50..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
528 (100..101, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
529 ],
530 )
531 .await;
532 test_merge(
534 (50..101, test_val.clone()),
535 base.clone(),
536 &[(50..101, AllocatorValue::Abs { count: 1, owner_object_id: 2 })],
537 )
538 .await;
539
540 test_merge(
544 base.clone(),
545 (49..99, test_val.clone()),
546 &[
547 (49..50, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
548 (50..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
549 ],
550 )
551 .await;
552 test_merge(
554 (49..99, test_val.clone()),
555 base.clone(),
556 &[
557 (49..50, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
558 (50..99, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
559 (99..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
560 ],
561 )
562 .await;
563
564 test_merge(
567 base.clone(),
568 (51..101, test_val.clone()),
569 &[
570 (50..51, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
571 (51..100, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
572 (100..101, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
573 ],
574 )
575 .await;
576 test_merge(
578 (51..101, test_val.clone()),
579 base.clone(),
580 &[
581 (50..51, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
582 (51..101, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
583 ],
584 )
585 .await;
586 }
587
588 #[fuchsia::test]
589 async fn test_length_1_ranges() {
590 test_merge(
592 (10..11, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
593 (11..12, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
594 &[
595 (10..11, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
596 (11..12, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
597 ],
598 )
599 .await;
600
601 test_merge(
603 (10..11, AllocatorValue::Abs { count: 1, owner_object_id: 1 }),
604 (10..11, AllocatorValue::Abs { count: 1, owner_object_id: 2 }),
605 &[(10..11, AllocatorValue::Abs { count: 1, owner_object_id: 1 })],
606 )
607 .await;
608 }
609}