1use std::ops::RangeInclusive;
5use std::sync::Arc;
6
7use super::bitmap::{Bitmap, count_ones};
8use super::{RowIdSequence, U64Segment};
9use lance_core::deepsize::DeepSizeOf;
10use lance_core::utils::address::RowAddress;
11use lance_core::utils::deletion::DeletionVector;
12use lance_core::{Error, Result};
13use rangemap::RangeInclusiveMap;
14
15const MAX_PROBE_DEPTH: u64 = 64;
19
20const MERGE_ROWS_BUDGET: u64 = 1 << 20;
23
24const RANK_BLOCK_BYTES: usize = 256;
28
29#[derive(Debug)]
45pub struct RowIdIndex {
46 fragments: Vec<FragmentEntry>,
48 end_tree: Vec<u64>,
51 merged: Option<MergedIndex>,
52}
53
54type MergedIndex = RangeInclusiveMap<u64, (U64Segment, U64Segment)>;
55
56pub struct FragmentRowIdIndex {
57 pub fragment_id: u32,
58 pub row_id_sequence: Arc<RowIdSequence>,
59 pub deletion_vector: Arc<DeletionVector>,
60}
61
62impl RowIdIndex {
63 pub fn new(fragment_indices: &[FragmentRowIdIndex]) -> Result<Self> {
65 let mut fragments: Vec<FragmentEntry> = fragment_indices
66 .iter()
67 .filter_map(FragmentEntry::new)
68 .collect();
69 fragments.sort_unstable_by_key(|entry| entry.start);
70
71 let mut index = Self {
72 end_tree: build_end_tree(&fragments),
73 fragments,
74 merged: None,
75 };
76 if !probing_beats_merging(&index.fragments) {
77 index.merged = Some(index.build_merged()?);
78 }
79 Ok(index)
80 }
81
82 fn build_merged(&self) -> Result<MergedIndex> {
83 let sources: Vec<FragmentRowIdIndex> = self
84 .fragments
85 .iter()
86 .map(|entry| FragmentRowIdIndex {
87 fragment_id: entry.fragment_id,
88 row_id_sequence: entry.sequence.clone(),
89 deletion_vector: entry.deletion_vector.clone(),
90 })
91 .collect();
92 let chunks = sources
93 .iter()
94 .flat_map(decompose_sequence)
95 .collect::<Vec<_>>();
96
97 let mut final_chunks = Vec::new();
98 for processed_chunk in prep_index_chunks(chunks) {
99 match processed_chunk {
100 RawIndexChunk::NonOverlapping(chunk) => {
101 final_chunks.push(chunk);
102 }
103 RawIndexChunk::Overlapping(_range, overlapping_chunks) => {
104 let merged_chunk = merge_overlapping_chunks(overlapping_chunks)?;
108 final_chunks.push(merged_chunk);
109 }
110 }
111 }
112
113 Ok(RangeInclusiveMap::from_iter(final_chunks))
114 }
115
116 pub fn get(&self, row_id: u64) -> Result<Option<RowAddress>> {
125 if let Some(merged) = &self.merged {
126 return Ok(merged_get(merged, row_id));
127 }
128 self.probe(row_id)
129 }
130
131 pub fn get_many(&self, row_ids: &[u64]) -> Result<Vec<Option<RowAddress>>> {
143 let n = row_ids.len();
144 let mut out = vec![None; n];
145 if n == 0 {
146 return Ok(out);
147 }
148
149 let mut sorted: Vec<(u64, usize)> = row_ids.iter().copied().zip(0..n).collect();
150 sorted.sort_unstable_by_key(|&(id, _)| id);
151
152 let Some(merged) = &self.merged else {
153 for (id, orig_idx) in sorted {
155 out[orig_idx] = self.probe(id)?;
156 }
157 return Ok(out);
158 };
159
160 let mut chunks = merged.iter().peekable();
161 for (id, orig_idx) in sorted {
162 while let Some((range, _)) = chunks.peek() {
164 if *range.end() < id {
165 chunks.next();
166 } else {
167 break;
168 }
169 }
170 let Some((range, (row_id_seg, addr_seg))) = chunks.peek() else {
171 break;
172 };
173 if id < *range.start() {
174 continue; }
176 if let Some(pos) = row_id_seg.position(id)
177 && let Some(addr) = addr_seg.get(pos)
178 {
179 out[orig_idx] = Some(RowAddress::from(addr));
180 }
181 }
182 Ok(out)
183 }
184
185 fn probe(&self, row_id: u64) -> Result<Option<RowAddress>> {
191 let fragments = self.fragments.len();
192 if fragments == 0 {
193 return Ok(None);
194 }
195 let upper = self
197 .fragments
198 .partition_point(|entry| entry.start <= row_id);
199 if upper == 0 {
200 return Ok(None);
201 }
202 let mut found: Option<RowAddress> = None;
203 let (mut node, mut lo, mut width) = (1usize, 0usize, self.end_tree.len() / 2);
208 loop {
209 if lo >= upper {
210 break;
212 }
213 if self.end_tree[node] >= row_id {
214 if width > 1 {
215 node *= 2;
216 width /= 2;
217 continue;
218 }
219 if lo < fragments
220 && let Some(candidate) = self.fragments[lo].resolve(row_id)
221 {
222 if found.is_some() {
223 return Err(Error::internal(format!(
224 "row id index corrupt: stable row id {row_id} is \
225 live in multiple fragments",
226 )));
227 }
228 found = Some(candidate);
229 }
230 }
231 while node & 1 == 1 {
232 if node == 1 {
233 return Ok(found);
234 }
235 node /= 2;
236 lo -= width;
237 width *= 2;
238 }
239 node += 1;
240 lo += width;
241 }
242 Ok(found)
243 }
244}
245
246fn merged_get(merged: &MergedIndex, row_id: u64) -> Option<RowAddress> {
247 let (row_id_segment, address_segment) = merged.get(&row_id)?;
248 let pos = row_id_segment.position(row_id)?;
249 let address = address_segment.get(pos)?;
250 Some(RowAddress::from(address))
251}
252
253#[derive(Debug)]
255struct SegmentEntry {
256 seq_idx: usize,
257 range: RangeInclusive<u64>,
258 start_offset: u32,
259 lookup: SegmentLookup,
260}
261
262#[derive(Debug)]
264enum SegmentLookup {
265 Native,
267 Positions(Vec<(u64, u32)>),
270 BitmapRank(Vec<u32>),
274}
275
276impl SegmentEntry {
277 fn position(&self, sequence: &RowIdSequence, row_id: u64) -> Option<usize> {
279 let segment = &sequence.0[self.seq_idx];
280 match &self.lookup {
281 SegmentLookup::Native => segment.position(row_id),
282 SegmentLookup::Positions(positions) => positions
283 .binary_search_by_key(&row_id, |(id, _)| *id)
284 .ok()
285 .map(|found| positions[found].1 as usize),
286 SegmentLookup::BitmapRank(rank) => {
287 let U64Segment::RangeWithBitmap { range, bitmap } = segment else {
288 return None;
289 };
290 if !range.contains(&row_id) {
291 return None;
292 }
293 let offset = (row_id - range.start) as usize;
294 if !bitmap.get(offset) {
295 return None;
296 }
297 let block_start = offset / (RANK_BLOCK_BYTES * 8) * (RANK_BLOCK_BYTES * 8);
298 let ones_before = rank[offset / (RANK_BLOCK_BYTES * 8)] as usize
299 + bitmap.slice(block_start, offset - block_start).count_ones();
300 Some(ones_before)
301 }
302 }
303 }
304}
305
306fn build_lookup(segment: &U64Segment) -> (SegmentLookup, usize) {
309 match segment {
310 U64Segment::Array(_) => {
311 let mut positions: Vec<(u64, u32)> = segment
312 .iter()
313 .enumerate()
314 .map(|(position, row_id)| (row_id, position as u32))
315 .collect();
316 positions.sort_unstable();
317 positions.dedup_by_key(|(row_id, _)| *row_id);
319 (SegmentLookup::Positions(positions), segment.len())
320 }
321 U64Segment::RangeWithBitmap { bitmap, .. } => {
322 let (rank, ones) = build_bitmap_rank(bitmap);
323 (SegmentLookup::BitmapRank(rank), ones)
324 }
325 _ => (SegmentLookup::Native, segment.len()),
326 }
327}
328
329fn build_bitmap_rank(bitmap: &Bitmap) -> (Vec<u32>, usize) {
331 let mut rank = Vec::with_capacity(bitmap.data.len().div_ceil(RANK_BLOCK_BYTES));
332 let mut ones = 0usize;
333 for block in bitmap.data.chunks(RANK_BLOCK_BYTES) {
334 rank.push(ones as u32);
335 ones += count_ones(block);
336 }
337 (rank, ones)
338}
339
340#[derive(Debug)]
341struct FragmentEntry {
342 fragment_id: u32,
343 sequence: Arc<RowIdSequence>,
344 deletion_vector: Arc<DeletionVector>,
345 segments: Vec<SegmentEntry>,
346 start: u64,
347 end: u64,
348 merge_rows: u64,
350}
351
352impl FragmentEntry {
353 fn new(source: &FragmentRowIdIndex) -> Option<Self> {
354 let mut segments: Vec<SegmentEntry> = Vec::new();
355 let mut start_offset: u32 = 0;
356 let mut merge_rows: u64 = 0;
357 let deleted = !source.deletion_vector.is_empty();
358 for (seq_idx, segment) in source.row_id_sequence.0.iter().enumerate() {
359 let (lookup, len) = build_lookup(segment);
360 if deleted || !matches!(segment, U64Segment::Range(_)) {
362 merge_rows += len as u64;
363 }
364 if len > 0
367 && let Some(range) = segment.range()
368 {
369 segments.push(SegmentEntry {
370 seq_idx,
371 range,
372 start_offset,
373 lookup,
374 });
375 }
376 start_offset += len as u32;
377 }
378 let start = segments.iter().map(|entry| *entry.range.start()).min()?;
379 let end = segments.iter().map(|entry| *entry.range.end()).max()?;
380 Some(Self {
381 fragment_id: source.fragment_id,
382 sequence: source.row_id_sequence.clone(),
383 deletion_vector: source.deletion_vector.clone(),
384 segments,
385 start,
386 end,
387 merge_rows,
388 })
389 }
390
391 fn resolve(&self, row_id: u64) -> Option<RowAddress> {
394 for entry in &self.segments {
395 if !entry.range.contains(&row_id) {
396 continue;
397 }
398 let Some(position) = entry.position(&self.sequence, row_id) else {
399 continue;
400 };
401 let row_offset = entry.start_offset + position as u32;
402 if self.deletion_vector.contains(row_offset) {
403 continue;
404 }
405 return Some(RowAddress::new_from_parts(self.fragment_id, row_offset));
406 }
407 None
408 }
409}
410
411fn probing_beats_merging(fragments: &[FragmentEntry]) -> bool {
418 let merge_rows: u64 = fragments.iter().map(|entry| entry.merge_rows).sum();
419 merge_rows > MERGE_ROWS_BUDGET && max_overlap_depth(fragments) <= MAX_PROBE_DEPTH
420}
421
422fn max_overlap_depth(fragments: &[FragmentEntry]) -> u64 {
424 let mut ends: Vec<u64> = fragments.iter().map(|entry| entry.end).collect();
425 ends.sort_unstable();
426 let mut closed = 0;
427 let mut depth: u64 = 0;
428 for (opened, entry) in fragments.iter().enumerate() {
429 while closed < ends.len() && ends[closed] < entry.start {
430 closed += 1;
431 }
432 depth = depth.max((opened + 1 - closed) as u64);
433 }
434 depth
435}
436
437fn build_end_tree(fragments: &[FragmentEntry]) -> Vec<u64> {
440 if fragments.is_empty() {
441 return Vec::new();
442 }
443 let leaves = fragments.len().next_power_of_two();
444 let mut tree = vec![0_u64; 2 * leaves];
445 for (slot, entry) in fragments.iter().enumerate() {
446 tree[leaves + slot] = entry.end;
447 }
448 for node in (1..leaves).rev() {
449 tree[node] = tree[2 * node].max(tree[2 * node + 1]);
450 }
451 tree
452}
453
454impl DeepSizeOf for RowIdIndex {
455 fn deep_size_of_children(&self, context: &mut lance_core::deepsize::Context) -> usize {
458 let fragment_bytes: usize = self
459 .fragments
460 .iter()
461 .map(|entry| {
462 entry.sequence.deep_size_of_children(context)
463 + entry.deletion_vector.deep_size_of_children(context)
464 + entry.segments.capacity() * std::mem::size_of::<SegmentEntry>()
465 + entry
466 .segments
467 .iter()
468 .map(|segment| match &segment.lookup {
469 SegmentLookup::Native => 0,
470 SegmentLookup::Positions(positions) => {
471 positions.capacity() * std::mem::size_of::<(u64, u32)>()
472 }
473 SegmentLookup::BitmapRank(rank) => {
474 rank.capacity() * std::mem::size_of::<u32>()
475 }
476 })
477 .sum::<usize>()
478 })
479 .sum();
480 let merged_bytes: usize = self
481 .merged
482 .as_ref()
483 .map(|merged| {
484 merged
485 .iter()
486 .map(|(_, (row_id_segment, address_segment))| {
487 (2 * std::mem::size_of::<u64>())
488 + std::mem::size_of::<(U64Segment, U64Segment)>()
489 + row_id_segment.deep_size_of_children(context)
490 + address_segment.deep_size_of_children(context)
491 })
492 .sum()
493 })
494 .unwrap_or(0);
495 fragment_bytes
496 + merged_bytes
497 + self.fragments.capacity() * std::mem::size_of::<FragmentEntry>()
498 + self.end_tree.capacity() * std::mem::size_of::<u64>()
499 }
500}
501
502fn decompose_sequence(
503 frag_index: &FragmentRowIdIndex,
504) -> Vec<(RangeInclusive<u64>, (U64Segment, U64Segment))> {
505 let mut start_address: u64 = RowAddress::first_row(frag_index.fragment_id).into();
506 let mut current_offset = 0u32;
507 let no_deletions = frag_index.deletion_vector.is_empty();
508
509 frag_index
510 .row_id_sequence
511 .0
512 .iter()
513 .filter_map(|segment| {
514 let segment_len = segment.len();
515
516 let result = if no_deletions {
517 decompose_segment_no_deletions(segment, start_address)
518 } else {
519 decompose_segment_with_deletions(
520 segment,
521 start_address,
522 current_offset,
523 &frag_index.deletion_vector,
524 )
525 };
526
527 current_offset += segment_len as u32;
528 start_address += segment_len as u64;
529
530 result
531 })
532 .collect()
533}
534
535fn build_chunk_from_pairs(mut pairs: Vec<(u64, u64)>) -> Option<IndexChunk> {
537 if pairs.is_empty() {
538 return None;
539 }
540 pairs.sort_unstable_by_key(|(row_id, _)| *row_id);
543 let (row_ids, addresses): (Vec<u64>, Vec<u64>) = pairs.into_iter().unzip();
544 let row_id_segment = U64Segment::from_iter(row_ids);
545 let address_segment = U64Segment::from_iter(addresses);
546 let coverage = row_id_segment.range()?;
547 Some((coverage, (row_id_segment, address_segment)))
548}
549
550fn decompose_segment_no_deletions(segment: &U64Segment, start_address: u64) -> Option<IndexChunk> {
552 match segment {
553 U64Segment::Range(range) if !range.is_empty() => {
554 let len = range.end - range.start;
555 let row_id_segment = U64Segment::Range(range.clone());
556 let address_segment = U64Segment::Range(start_address..start_address + len);
557 let coverage = range.start..=range.end - 1;
558 Some((coverage, (row_id_segment, address_segment)))
559 }
560 _ if segment.is_empty() => None,
561 _ => {
562 let pairs: Vec<(u64, u64)> = segment
564 .iter()
565 .enumerate()
566 .map(|(i, row_id)| (row_id, start_address + i as u64))
567 .collect();
568 build_chunk_from_pairs(pairs)
569 }
570 }
571}
572
573fn decompose_segment_with_deletions(
575 segment: &U64Segment,
576 start_address: u64,
577 current_offset: u32,
578 deletion_vector: &DeletionVector,
579) -> Option<IndexChunk> {
580 let pairs: Vec<(u64, u64)> = segment
581 .iter()
582 .enumerate()
583 .filter_map(|(i, row_id)| {
584 let row_offset = current_offset + i as u32;
585 if !deletion_vector.contains(row_offset) {
586 Some((row_id, start_address + i as u64))
587 } else {
588 None
589 }
590 })
591 .collect();
592 build_chunk_from_pairs(pairs)
593}
594
595type IndexChunk = (RangeInclusive<u64>, (U64Segment, U64Segment));
596
597#[derive(Debug)]
598enum RawIndexChunk {
599 NonOverlapping(IndexChunk),
600 Overlapping(RangeInclusive<u64>, Vec<IndexChunk>),
601}
602
603impl RawIndexChunk {
604 fn range_end(&self) -> u64 {
605 match self {
606 Self::NonOverlapping((range, _)) => *range.end(),
607 Self::Overlapping(range, _) => *range.end(),
608 }
609 }
610}
611
612fn prep_index_chunks(mut chunks: Vec<IndexChunk>) -> impl Iterator<Item = RawIndexChunk> {
617 chunks.sort_by_key(|(range, _)| u64::MAX - *range.start());
618
619 let mut output = Vec::new();
620
621 if let Some(first_chunk) = chunks.pop() {
623 output.push(RawIndexChunk::NonOverlapping(first_chunk));
624 } else {
625 return output.into_iter();
627 }
628
629 let mut current_range = 0..=0;
630 let mut current_overlap = Vec::new();
631 while let Some(chunk) = chunks.pop() {
632 debug_assert_eq!(
633 current_overlap
634 .iter()
635 .map(|(range, _): &IndexChunk| *range.start())
636 .min()
637 .unwrap_or_default(),
638 *current_range.start(),
639 );
640 debug_assert_eq!(
641 current_overlap
642 .iter()
643 .map(|(range, _): &IndexChunk| *range.end())
644 .max()
645 .unwrap_or_default(),
646 *current_range.end(),
647 );
648
649 if current_overlap.is_empty() {
650 let last_chunk_end = output.last().unwrap().range_end();
652 if *chunk.0.start() <= last_chunk_end {
653 match output.pop().unwrap() {
655 RawIndexChunk::NonOverlapping(chunk) => {
656 current_overlap.push(chunk);
657 }
658 _ => unreachable!(),
659 }
660 current_overlap.push(chunk);
661
662 let range_start = *current_overlap.first().unwrap().0.start();
663 let range_end = *current_overlap
664 .last()
665 .unwrap()
666 .0
667 .end()
668 .max(current_overlap.first().unwrap().0.end());
669 current_range = range_start..=range_end;
670 } else {
671 output.push(RawIndexChunk::NonOverlapping(chunk));
673 }
674 } else {
675 if chunk.0.start() <= current_range.end() {
677 let range_end = *chunk.0.end().max(current_range.end());
679 current_range = *current_range.start()..=range_end;
680
681 current_overlap.push(chunk);
682 } else {
683 output.push(RawIndexChunk::Overlapping(
685 std::mem::replace(&mut current_range, 0..=0),
686 std::mem::take(&mut current_overlap),
687 ));
688 output.push(RawIndexChunk::NonOverlapping(chunk));
689 }
690 }
691 }
692 debug_assert_eq!(
693 current_overlap
694 .iter()
695 .map(|(range, _): &IndexChunk| *range.start())
696 .min()
697 .unwrap_or_default(),
698 *current_range.start(),
699 );
700 debug_assert_eq!(
701 current_overlap
702 .iter()
703 .map(|(range, _): &IndexChunk| *range.end())
704 .max()
705 .unwrap_or_default(),
706 *current_range.end(),
707 );
708
709 if !current_overlap.is_empty() {
710 output.push(RawIndexChunk::Overlapping(
711 current_range.clone(),
712 current_overlap,
713 ));
714 }
715
716 output.into_iter()
717}
718
719fn merge_overlapping_chunks(overlapping_chunks: Vec<IndexChunk>) -> Result<IndexChunk> {
720 let total_capacity = overlapping_chunks
721 .iter()
722 .map(|(_, (row_ids, _))| row_ids.len())
723 .sum();
724 let mut values = Vec::with_capacity(total_capacity);
725 for (_, (row_ids, row_addrs)) in overlapping_chunks.iter() {
726 values.extend(row_ids.iter().zip(row_addrs.iter()));
727 }
728 values.sort_by_key(|(row_id, _)| *row_id);
729 if let Some(w) = values.windows(2).find(|w| w[0].0 == w[1].0) {
732 return Err(Error::internal(format!(
733 "row id index corrupt: stable row id {} is live in multiple fragments",
734 w[0].0
735 )));
736 }
737 let row_id_segment = U64Segment::from_iter(values.iter().map(|(row_id, _)| *row_id));
738 let address_segment = U64Segment::from_iter(values.iter().map(|(_, row_addr)| *row_addr));
739
740 let range = row_id_segment.range().unwrap();
741
742 Ok((range, (row_id_segment, address_segment)))
743}
744
745#[cfg(test)]
746impl RowIdIndex {
747 fn probing(fragment_indices: &[FragmentRowIdIndex]) -> Result<Self> {
749 let mut fragments: Vec<FragmentEntry> = fragment_indices
750 .iter()
751 .filter_map(FragmentEntry::new)
752 .collect();
753 fragments.sort_unstable_by_key(|entry| entry.start);
754 Ok(Self {
755 end_tree: build_end_tree(&fragments),
756 fragments,
757 merged: None,
758 })
759 }
760}
761
762#[cfg(test)]
763mod tests {
764 use super::*;
765 use proptest::{
766 prelude::{Just, Strategy, any},
767 prop_assert, prop_assert_eq,
768 };
769
770 fn sparse_sequence(len: u64) -> RowIdSequence {
772 RowIdSequence(vec![U64Segment::SortedArray(
773 (0..len).map(|value| value * 2).collect::<Vec<u64>>().into(),
774 )])
775 }
776
777 fn fragment(fragment_id: u32, sequence: RowIdSequence) -> FragmentRowIdIndex {
778 FragmentRowIdIndex {
779 fragment_id,
780 row_id_sequence: Arc::new(sequence),
781 deletion_vector: Arc::new(DeletionVector::default()),
782 }
783 }
784
785 #[test]
786 fn test_new_builds_the_merged_map_unless_probing_wins() {
787 let ranges = fragment(1, RowIdSequence(vec![U64Segment::Range(0..1_000_000)]));
790 assert!(RowIdIndex::new(&[ranges]).unwrap().merged.is_some());
791 let small = fragment(1, sparse_sequence(16));
792 assert!(RowIdIndex::new(&[small]).unwrap().merged.is_some());
793
794 let wide = fragment(1, sparse_sequence(MERGE_ROWS_BUDGET + 1));
796 let index = RowIdIndex::new(&[wide]).unwrap();
797 assert!(index.merged.is_none());
798 assert_eq!(
799 index.get(6).unwrap(),
800 Some(RowAddress::new_from_parts(1, 3))
801 );
802 }
803
804 #[test]
805 fn test_deep_overlap_merges_however_many_rows_it_reads() {
806 let fragments = MAX_PROBE_DEPTH + 1;
809 let rows_per_fragment = MERGE_ROWS_BUDGET / fragments + 1;
810 let deep: Vec<FragmentRowIdIndex> = (0..fragments as u32)
811 .map(|id| {
812 let ids: Vec<u64> = (0..rows_per_fragment)
813 .map(|value| value * fragments + id as u64)
814 .collect();
815 fragment(id, RowIdSequence(vec![U64Segment::SortedArray(ids.into())]))
816 })
817 .collect();
818
819 assert!(RowIdIndex::new(&deep).unwrap().merged.is_some());
820 }
821
822 #[test]
823 fn test_probe_resolves_a_row_id_the_merged_map_rejects() {
824 let sources = [
825 fragment(1, RowIdSequence::from(&[0, 2][..])),
826 fragment(2, RowIdSequence::from(&[1, 2][..])),
827 ];
828 assert!(RowIdIndex::new(&sources).is_err());
829
830 let index = RowIdIndex::probing(&sources[..1]).unwrap();
831 assert_eq!(
832 index.get(2).unwrap(),
833 Some(RowAddress::new_from_parts(1, 1))
834 );
835 }
836
837 #[test]
838 fn test_probe_errors_when_two_fragments_hold_an_id_live() {
839 let sources = [
840 fragment(1, RowIdSequence::from(&[0, 2][..])),
841 fragment(2, RowIdSequence::from(&[1, 2][..])),
842 ];
843 let index = RowIdIndex::probing(&sources).unwrap();
844 assert_eq!(
845 index.get(0).unwrap(),
846 Some(RowAddress::new_from_parts(1, 0))
847 );
848
849 let error = index.get(2).unwrap_err();
850 assert!(matches!(&error, Error::Internal { .. }));
851 assert!(
852 error
853 .to_string()
854 .contains("stable row id 2 is live in multiple fragments")
855 );
856
857 let error = index.get_many(&[0, 2]).unwrap_err();
858 assert!(matches!(&error, Error::Internal { .. }));
859 }
860
861 #[test]
862 fn test_probe_ranks_a_bitmap_wider_than_one_block() {
863 let span = (RANK_BLOCK_BYTES * 8 * 3 + 17) as u64;
866 let present: Vec<bool> = (0..span).map(|slot| slot % 3 != 1).collect();
867 let segment = U64Segment::RangeWithBitmap {
868 range: 1000..1000 + span,
869 bitmap: present.as_slice().into(),
870 };
871 let index = RowIdIndex::probing(&[fragment(7, RowIdSequence(vec![segment]))]).unwrap();
872 let mut position = 0;
873 for slot in 0..span {
874 let found = index.get(1000 + slot).unwrap();
875 if slot % 3 == 1 {
876 assert_eq!(found, None);
877 } else {
878 assert_eq!(found, Some(RowAddress::new_from_parts(7, position)));
879 position += 1;
880 }
881 }
882 assert_eq!(index.get(999).unwrap(), None);
883 assert_eq!(index.get(1000 + span).unwrap(), None);
884 }
885
886 #[test]
887 fn test_probe_finds_every_position_of_an_unsorted_array() {
888 let row_ids: Vec<u64> = (0..2048).map(|value| (value * 7919) % 2048).collect();
889 let index = RowIdIndex::probing(&[fragment(
890 3,
891 RowIdSequence(vec![U64Segment::Array(row_ids.clone().into())]),
892 )])
893 .unwrap();
894 for (offset, row_id) in row_ids.iter().enumerate() {
895 assert_eq!(
896 index.get(*row_id).unwrap(),
897 Some(RowAddress::new_from_parts(3, offset as u32))
898 );
899 }
900 assert!(index.merged.is_none());
901 }
902
903 #[test]
904 fn test_new_index() {
905 let fragment_indices = vec![
906 FragmentRowIdIndex {
907 fragment_id: 10,
908 row_id_sequence: Arc::new(RowIdSequence(vec![
909 U64Segment::Range(0..10),
910 U64Segment::RangeWithHoles {
911 range: 10..17,
912 holes: vec![12, 15].into(),
913 },
914 U64Segment::SortedArray(vec![20, 25, 30].into()),
915 ])),
916 deletion_vector: Arc::new(DeletionVector::default()),
917 },
918 FragmentRowIdIndex {
919 fragment_id: 20,
920 row_id_sequence: Arc::new(RowIdSequence(vec![
921 U64Segment::RangeWithBitmap {
922 range: 17..20,
923 bitmap: [true, false, true].as_slice().into(),
924 },
925 U64Segment::Array(vec![40, 50, 60].into()),
926 ])),
927 deletion_vector: Arc::new(DeletionVector::default()),
928 },
929 ];
930
931 let index = RowIdIndex::new(&fragment_indices).unwrap();
932
933 assert_eq!(
935 index.get(0).unwrap(),
936 Some(RowAddress::new_from_parts(10, 0))
937 );
938 assert_eq!(index.get(15).unwrap(), None);
939 assert_eq!(
940 index.get(16).unwrap(),
941 Some(RowAddress::new_from_parts(10, 14))
942 );
943 assert_eq!(
944 index.get(17).unwrap(),
945 Some(RowAddress::new_from_parts(20, 0))
946 );
947 assert_eq!(
948 index.get(25).unwrap(),
949 Some(RowAddress::new_from_parts(10, 16))
950 );
951 assert_eq!(
952 index.get(40).unwrap(),
953 Some(RowAddress::new_from_parts(20, 2))
954 );
955 assert_eq!(
956 index.get(60).unwrap(),
957 Some(RowAddress::new_from_parts(20, 4))
958 );
959 assert_eq!(index.get(61).unwrap(), None);
960 }
961
962 #[test]
963 fn test_new_index_overlap() {
964 let fragment_indices = vec![
965 FragmentRowIdIndex {
966 fragment_id: 23,
967 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::SortedArray(
968 vec![3, 6, 9].into(),
969 )])),
970 deletion_vector: Arc::new(DeletionVector::default()),
971 },
972 FragmentRowIdIndex {
973 fragment_id: 42,
974 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::SortedArray(
975 vec![2, 5, 8].into(),
976 )])),
977 deletion_vector: Arc::new(DeletionVector::default()),
978 },
979 FragmentRowIdIndex {
980 fragment_id: 10,
981 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::SortedArray(
982 vec![1, 4, 7].into(),
983 )])),
984 deletion_vector: Arc::new(DeletionVector::default()),
985 },
986 ];
987
988 let index = RowIdIndex::new(&fragment_indices).unwrap();
989
990 assert_eq!(
992 index.get(1).unwrap(),
993 Some(RowAddress::new_from_parts(10, 0))
994 );
995 assert_eq!(
996 index.get(2).unwrap(),
997 Some(RowAddress::new_from_parts(42, 0))
998 );
999 assert_eq!(
1000 index.get(3).unwrap(),
1001 Some(RowAddress::new_from_parts(23, 0))
1002 );
1003 assert_eq!(
1004 index.get(4).unwrap(),
1005 Some(RowAddress::new_from_parts(10, 1))
1006 );
1007 assert_eq!(
1008 index.get(5).unwrap(),
1009 Some(RowAddress::new_from_parts(42, 1))
1010 );
1011 assert_eq!(
1012 index.get(6).unwrap(),
1013 Some(RowAddress::new_from_parts(23, 1))
1014 );
1015 assert_eq!(
1016 index.get(7).unwrap(),
1017 Some(RowAddress::new_from_parts(10, 2))
1018 );
1019 assert_eq!(
1020 index.get(8).unwrap(),
1021 Some(RowAddress::new_from_parts(42, 2))
1022 );
1023 assert_eq!(
1024 index.get(9).unwrap(),
1025 Some(RowAddress::new_from_parts(23, 2))
1026 );
1027 }
1028
1029 #[test]
1030 fn test_new_index_unsorted_row_ids() {
1031 let fragment_indices = vec![
1033 FragmentRowIdIndex {
1034 fragment_id: 10,
1035 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Array(
1036 vec![9, 3, 6].into(), )])),
1038 deletion_vector: Arc::new(DeletionVector::default()),
1039 },
1040 FragmentRowIdIndex {
1041 fragment_id: 20,
1042 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Array(
1043 vec![8, 2, 5].into(), )])),
1045 deletion_vector: Arc::new(DeletionVector::default()),
1046 },
1047 FragmentRowIdIndex {
1048 fragment_id: 30,
1049 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Array(
1050 vec![7, 1, 4].into(), )])),
1052 deletion_vector: Arc::new(DeletionVector::default()),
1053 },
1054 ];
1055
1056 let index = RowIdIndex::new(&fragment_indices).unwrap();
1057
1058 assert_eq!(
1060 index.get(1).unwrap(),
1061 Some(RowAddress::new_from_parts(30, 1))
1062 );
1063 assert_eq!(
1064 index.get(2).unwrap(),
1065 Some(RowAddress::new_from_parts(20, 1))
1066 );
1067 assert_eq!(
1068 index.get(3).unwrap(),
1069 Some(RowAddress::new_from_parts(10, 1))
1070 );
1071 assert_eq!(
1072 index.get(4).unwrap(),
1073 Some(RowAddress::new_from_parts(30, 2))
1074 );
1075 assert_eq!(
1076 index.get(5).unwrap(),
1077 Some(RowAddress::new_from_parts(20, 2))
1078 );
1079 assert_eq!(
1080 index.get(6).unwrap(),
1081 Some(RowAddress::new_from_parts(10, 2))
1082 );
1083 assert_eq!(
1084 index.get(7).unwrap(),
1085 Some(RowAddress::new_from_parts(30, 0))
1086 );
1087 assert_eq!(
1088 index.get(8).unwrap(),
1089 Some(RowAddress::new_from_parts(20, 0))
1090 );
1091 assert_eq!(
1092 index.get(9).unwrap(),
1093 Some(RowAddress::new_from_parts(10, 0))
1094 );
1095
1096 assert_eq!(index.get(0).unwrap(), None);
1098 assert_eq!(index.get(10).unwrap(), None);
1099 }
1100
1101 #[test]
1102 fn test_new_index_partial_overlap() {
1103 let fragment_indices = vec![
1104 FragmentRowIdIndex {
1105 fragment_id: 0,
1106 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::RangeWithHoles {
1107 range: 0..100,
1108 holes: vec![50].into(),
1109 }])),
1110 deletion_vector: Arc::new(DeletionVector::default()),
1111 },
1112 FragmentRowIdIndex {
1113 fragment_id: 1,
1114 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Range(50..51)])),
1115 deletion_vector: Arc::new(DeletionVector::default()),
1116 },
1117 ];
1118
1119 let index = RowIdIndex::new(&fragment_indices).unwrap();
1120
1121 assert_eq!(
1123 index.get(0).unwrap(),
1124 Some(RowAddress::new_from_parts(0, 0))
1125 );
1126 assert_eq!(
1127 index.get(49).unwrap(),
1128 Some(RowAddress::new_from_parts(0, 49))
1129 );
1130 assert_eq!(
1131 index.get(50).unwrap(),
1132 Some(RowAddress::new_from_parts(1, 0))
1133 );
1134 assert_eq!(
1135 index.get(51).unwrap(),
1136 Some(RowAddress::new_from_parts(0, 50))
1137 );
1138 assert_eq!(
1139 index.get(99).unwrap(),
1140 Some(RowAddress::new_from_parts(0, 98))
1141 );
1142 }
1143
1144 #[test]
1145 fn test_overlapping_chunks_sparse_with_deletions() {
1146 let fragment_indices = vec![
1149 FragmentRowIdIndex {
1150 fragment_id: 10,
1151 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::SortedArray(
1152 vec![1, 3, 5, 7, 9].into(),
1153 )])),
1154 deletion_vector: Arc::new(DeletionVector::default()),
1155 },
1156 FragmentRowIdIndex {
1157 fragment_id: 20,
1158 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::SortedArray(
1159 vec![0, 2, 4, 6, 8].into(),
1160 )])),
1161 deletion_vector: Arc::new(DeletionVector::from_iter(vec![2])),
1163 },
1164 ];
1165
1166 let index = RowIdIndex::new(&fragment_indices).unwrap();
1167
1168 assert_eq!(
1169 index.get(0).unwrap(),
1170 Some(RowAddress::new_from_parts(20, 0))
1171 );
1172 assert_eq!(
1173 index.get(1).unwrap(),
1174 Some(RowAddress::new_from_parts(10, 0))
1175 );
1176 assert_eq!(
1177 index.get(2).unwrap(),
1178 Some(RowAddress::new_from_parts(20, 1))
1179 );
1180 assert_eq!(
1181 index.get(3).unwrap(),
1182 Some(RowAddress::new_from_parts(10, 1))
1183 );
1184 assert_eq!(index.get(4).unwrap(), None);
1185 assert_eq!(
1187 index.get(6).unwrap(),
1188 Some(RowAddress::new_from_parts(20, 3))
1189 );
1190 assert_eq!(
1191 index.get(8).unwrap(),
1192 Some(RowAddress::new_from_parts(20, 4))
1193 );
1194 assert_eq!(
1195 index.get(9).unwrap(),
1196 Some(RowAddress::new_from_parts(10, 4))
1197 );
1198 }
1199
1200 #[test]
1201 fn test_index_with_deletion_vector() {
1202 let deletion_vector = DeletionVector::from_iter(vec![2, 3]);
1203
1204 let fragment_indices = vec![FragmentRowIdIndex {
1205 fragment_id: 10,
1206 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Range(0..6)])),
1207 deletion_vector: Arc::new(deletion_vector),
1208 }];
1209
1210 let index = RowIdIndex::new(&fragment_indices).unwrap();
1211
1212 assert_eq!(
1213 index.get(0).unwrap(),
1214 Some(RowAddress::new_from_parts(10, 0))
1215 );
1216 assert_eq!(
1217 index.get(1).unwrap(),
1218 Some(RowAddress::new_from_parts(10, 1))
1219 );
1220 assert_eq!(
1221 index.get(4).unwrap(),
1222 Some(RowAddress::new_from_parts(10, 4))
1223 );
1224 assert_eq!(
1225 index.get(5).unwrap(),
1226 Some(RowAddress::new_from_parts(10, 5))
1227 );
1228
1229 assert_eq!(index.get(2).unwrap(), None);
1230 assert_eq!(index.get(3).unwrap(), None);
1231 }
1232
1233 #[test]
1234 fn test_empty_fragment_sequences() {
1235 let fragment_indices = vec![
1236 FragmentRowIdIndex {
1237 fragment_id: 10,
1238 row_id_sequence: Arc::new(RowIdSequence(vec![])),
1239 deletion_vector: Arc::new(DeletionVector::default()),
1240 },
1241 FragmentRowIdIndex {
1242 fragment_id: 20,
1243 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Range(5..8)])),
1244 deletion_vector: Arc::new(DeletionVector::default()),
1245 },
1246 ];
1247
1248 let index = RowIdIndex::new(&fragment_indices).unwrap();
1249
1250 assert_eq!(
1251 index.get(5).unwrap(),
1252 Some(RowAddress::new_from_parts(20, 0))
1253 );
1254 assert_eq!(
1255 index.get(7).unwrap(),
1256 Some(RowAddress::new_from_parts(20, 2))
1257 );
1258 assert_eq!(index.get(4).unwrap(), None);
1259 }
1260
1261 #[test]
1262 fn test_completely_empty_index() {
1263 let fragment_indices = vec![];
1264 let index = RowIdIndex::new(&fragment_indices).unwrap();
1265
1266 assert_eq!(index.get(0).unwrap(), None);
1267 assert_eq!(index.get(100).unwrap(), None);
1268 }
1269
1270 #[test]
1271 fn test_non_overlapping_ranges() {
1272 let fragment_indices = vec![
1273 FragmentRowIdIndex {
1274 fragment_id: 10,
1275 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Range(0..5)])),
1276 deletion_vector: Arc::new(DeletionVector::default()),
1277 },
1278 FragmentRowIdIndex {
1279 fragment_id: 20,
1280 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Range(5..10)])),
1281 deletion_vector: Arc::new(DeletionVector::default()),
1282 },
1283 FragmentRowIdIndex {
1284 fragment_id: 30,
1285 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Range(10..15)])),
1286 deletion_vector: Arc::new(DeletionVector::default()),
1287 },
1288 ];
1289
1290 let index = RowIdIndex::new(&fragment_indices).unwrap();
1291
1292 assert_eq!(
1293 index.get(0).unwrap(),
1294 Some(RowAddress::new_from_parts(10, 0))
1295 );
1296 assert_eq!(
1297 index.get(4).unwrap(),
1298 Some(RowAddress::new_from_parts(10, 4))
1299 );
1300 assert_eq!(
1301 index.get(5).unwrap(),
1302 Some(RowAddress::new_from_parts(20, 0))
1303 );
1304 assert_eq!(
1305 index.get(9).unwrap(),
1306 Some(RowAddress::new_from_parts(20, 4))
1307 );
1308 assert_eq!(
1309 index.get(10).unwrap(),
1310 Some(RowAddress::new_from_parts(30, 0))
1311 );
1312 assert_eq!(
1313 index.get(14).unwrap(),
1314 Some(RowAddress::new_from_parts(30, 4))
1315 );
1316 }
1317
1318 fn arbitrary_row_ids(
1319 num_fragments_range: std::ops::Range<usize>,
1320 frag_size_range: std::ops::Range<usize>,
1321 ) -> impl Strategy<Value = Vec<(u32, Arc<RowIdSequence>)>> {
1322 let fragment_sizes = proptest::collection::vec(frag_size_range, num_fragments_range);
1323 fragment_sizes.prop_flat_map(|fragment_sizes| {
1324 let num_rows = fragment_sizes.iter().sum::<usize>() as u64;
1325 let row_ids = 0..num_rows;
1326 let row_ids = row_ids.collect::<Vec<_>>();
1327 let row_ids_shuffled = proptest::strategy::Just(row_ids).prop_shuffle();
1328 row_ids_shuffled.prop_map(move |row_ids| {
1329 let mut sequences = Vec::with_capacity(fragment_sizes.len());
1330 let mut i = 0;
1331 for size in &fragment_sizes {
1332 let end = i + size;
1333 let sequence =
1334 RowIdSequence(vec![U64Segment::from_slice(row_ids[i..end].into())]);
1335 sequences.push((i as u32, Arc::new(sequence)));
1336 i = end;
1337 }
1338 sequences
1339 })
1340 })
1341 }
1342
1343 fn arbitrary_row_ids_with_deletions(
1344 num_fragments_range: std::ops::Range<usize>,
1345 frag_size_range: std::ops::Range<usize>,
1346 ) -> impl Strategy<Value = Vec<(u32, Arc<RowIdSequence>, Arc<DeletionVector>)>> {
1347 arbitrary_row_ids(num_fragments_range, frag_size_range)
1348 .prop_flat_map(|row_ids| {
1349 let num_rows = row_ids
1350 .iter()
1351 .map(|(_, sequence)| sequence.len() as usize)
1352 .sum::<usize>();
1353 (
1354 Just(row_ids),
1355 proptest::collection::vec(any::<bool>(), num_rows),
1356 )
1357 })
1358 .prop_map(|(row_ids, deleted_rows)| {
1359 let mut deleted_rows = deleted_rows.into_iter();
1360 row_ids
1361 .into_iter()
1362 .map(|(fragment_id, sequence)| {
1363 let mut deletion_bitmap = roaring::RoaringBitmap::new();
1364 for offset in 0..sequence.len() as u32 {
1365 if deleted_rows.next().unwrap() {
1366 deletion_bitmap.insert(offset);
1367 }
1368 }
1369 (
1370 fragment_id,
1371 sequence,
1372 Arc::new(DeletionVector::Bitmap(deletion_bitmap)),
1373 )
1374 })
1375 .collect()
1376 })
1377 }
1378
1379 #[test]
1380 fn test_large_range_segments_no_deletions() {
1381 let rows_per_fragment = 250_000u64;
1385 let num_fragments = 100u32;
1386 let mut offset = 0u64;
1387
1388 let fragment_indices: Vec<FragmentRowIdIndex> = (0..num_fragments)
1389 .map(|frag_id| {
1390 let start = offset;
1391 offset += rows_per_fragment;
1392 FragmentRowIdIndex {
1393 fragment_id: frag_id,
1394 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Range(
1395 start..start + rows_per_fragment,
1396 )])),
1397 deletion_vector: Arc::new(DeletionVector::default()),
1398 }
1399 })
1400 .collect();
1401
1402 let start = std::time::Instant::now();
1403 let index = RowIdIndex::new(&fragment_indices).unwrap();
1404 let elapsed = start.elapsed();
1405
1406 assert_eq!(
1408 index.get(0).unwrap(),
1409 Some(RowAddress::new_from_parts(0, 0))
1410 );
1411 assert_eq!(
1412 index.get(rows_per_fragment - 1).unwrap(),
1413 Some(RowAddress::new_from_parts(0, rows_per_fragment as u32 - 1))
1414 );
1415 assert_eq!(
1416 index.get(rows_per_fragment).unwrap(),
1417 Some(RowAddress::new_from_parts(1, 0))
1418 );
1419 let last_row = num_fragments as u64 * rows_per_fragment - 1;
1420 assert_eq!(
1421 index.get(last_row).unwrap(),
1422 Some(RowAddress::new_from_parts(
1423 num_fragments - 1,
1424 rows_per_fragment as u32 - 1
1425 ))
1426 );
1427 assert_eq!(index.get(last_row + 1).unwrap(), None);
1428
1429 assert!(
1432 elapsed.as_secs() < 1,
1433 "Index build took {:?} for {} fragments x {} rows = {} total rows. \
1434 This suggests the O(rows) -> O(fragments) optimization is not working.",
1435 elapsed,
1436 num_fragments,
1437 rows_per_fragment,
1438 num_fragments as u64 * rows_per_fragment,
1439 );
1440 }
1441
1442 #[test]
1443 fn test_large_range_segments_with_deletions() {
1444 let rows_per_fragment = 1_000u64;
1445 let num_fragments = 10u32;
1446 let mut offset = 0u64;
1447
1448 let fragment_indices: Vec<FragmentRowIdIndex> = (0..num_fragments)
1449 .map(|frag_id| {
1450 let start = offset;
1451 offset += rows_per_fragment;
1452
1453 let mut deleted = roaring::RoaringBitmap::new();
1455 for i in (0..rows_per_fragment as u32).step_by(3) {
1456 deleted.insert(i);
1457 }
1458
1459 FragmentRowIdIndex {
1460 fragment_id: frag_id,
1461 row_id_sequence: Arc::new(RowIdSequence(vec![U64Segment::Range(
1462 start..start + rows_per_fragment,
1463 )])),
1464 deletion_vector: Arc::new(DeletionVector::Bitmap(deleted)),
1465 }
1466 })
1467 .collect();
1468
1469 let index = RowIdIndex::new(&fragment_indices).unwrap();
1470
1471 assert_eq!(index.get(0).unwrap(), None);
1474 assert_eq!(index.get(3).unwrap(), None);
1476
1477 assert_eq!(
1480 index.get(1).unwrap(),
1481 Some(RowAddress::new_from_parts(0, 1))
1482 );
1483 assert_eq!(
1485 index.get(2).unwrap(),
1486 Some(RowAddress::new_from_parts(0, 2))
1487 );
1488 assert_eq!(
1490 index.get(4).unwrap(),
1491 Some(RowAddress::new_from_parts(0, 4))
1492 );
1493
1494 assert_eq!(index.get(rows_per_fragment).unwrap(), None);
1497 assert_eq!(
1499 index.get(rows_per_fragment + 1).unwrap(),
1500 Some(RowAddress::new_from_parts(1, 1))
1501 );
1502
1503 let last_row = num_fragments as u64 * rows_per_fragment - 1;
1506 assert_eq!(index.get(last_row).unwrap(), None);
1507 assert_eq!(
1509 index.get(last_row - 1).unwrap(),
1510 Some(RowAddress::new_from_parts(num_fragments - 1, 998))
1511 );
1512
1513 assert_eq!(index.get(last_row + 1).unwrap(), None);
1515 }
1516
1517 proptest::proptest! {
1518 #[test]
1519 fn test_new_index_robustness(
1520 row_ids in arbitrary_row_ids_with_deletions(0..5, 0..32)
1521 ) {
1522 let fragment_indices: Vec<FragmentRowIdIndex> = row_ids
1523 .iter()
1524 .map(|(frag_id, sequence, deletion_vector)| FragmentRowIdIndex {
1525 fragment_id: *frag_id,
1526 row_id_sequence: sequence.clone(),
1527 deletion_vector: deletion_vector.clone(),
1528 })
1529 .collect();
1530
1531 let merged = RowIdIndex::new(&fragment_indices).unwrap();
1532 let probing = RowIdIndex::probing(&fragment_indices).unwrap();
1533 for index in [&merged, &probing] {
1534 for (frag_id, sequence, deletion_vector) in row_ids.iter() {
1535 for (local_offset, row_id) in sequence.iter().enumerate() {
1536 let expected = if deletion_vector.contains(local_offset as u32) {
1537 None
1538 } else {
1539 Some(RowAddress::new_from_parts(*frag_id, local_offset as u32))
1540 };
1541 prop_assert_eq!(
1542 index.get(row_id).unwrap(),
1543 expected,
1544 "Row id {} in sequence {:?} not found in index {:?}",
1545 row_id,
1546 sequence,
1547 index
1548 );
1549 }
1550 }
1551 }
1552 }
1553
1554 #[test]
1555 fn test_new_index_moved_row_id(
1556 row_id in any::<u64>(),
1557 source_fragment in 0u32..1024,
1558 fragment_delta in 1u32..1024,
1559 ) {
1560 let target_fragment = source_fragment + fragment_delta;
1561 let fragment_indices = [
1562 FragmentRowIdIndex {
1563 fragment_id: source_fragment,
1564 row_id_sequence: Arc::new(RowIdSequence::from(&[row_id][..])),
1565 deletion_vector: Arc::new(DeletionVector::Bitmap(
1566 roaring::RoaringBitmap::from_iter([0]),
1567 )),
1568 },
1569 FragmentRowIdIndex {
1570 fragment_id: target_fragment,
1571 row_id_sequence: Arc::new(RowIdSequence::from(&[row_id][..])),
1572 deletion_vector: Arc::new(DeletionVector::default()),
1573 },
1574 ];
1575
1576 let index = RowIdIndex::new(&fragment_indices).unwrap();
1577 prop_assert_eq!(
1578 index.get(row_id).unwrap(),
1579 Some(RowAddress::new_from_parts(target_fragment, 0))
1580 );
1581 }
1582
1583 #[test]
1584 fn test_new_index_rejects_duplicate_live_row_id(
1585 row_id in any::<u64>(),
1586 first_fragment in 0u32..1024,
1587 fragment_delta in 1u32..1024,
1588 ) {
1589 let second_fragment = first_fragment + fragment_delta;
1590 let fragment_indices = [
1591 FragmentRowIdIndex {
1592 fragment_id: first_fragment,
1593 row_id_sequence: Arc::new(RowIdSequence::from(&[row_id][..])),
1594 deletion_vector: Arc::new(DeletionVector::default()),
1595 },
1596 FragmentRowIdIndex {
1597 fragment_id: second_fragment,
1598 row_id_sequence: Arc::new(RowIdSequence::from(&[row_id][..])),
1599 deletion_vector: Arc::new(DeletionVector::default()),
1600 },
1601 ];
1602
1603 let error = RowIdIndex::new(&fragment_indices).unwrap_err();
1604 let is_internal = matches!(&error, Error::Internal { .. });
1605 let expected_message =
1606 format!("stable row id {row_id} is live in multiple fragments");
1607 let error_message = error.to_string();
1608 prop_assert!(is_internal);
1609 prop_assert!(error_message.contains(&expected_message));
1610 }
1611 }
1612}