1use std::ops::{Range, RangeInclusive};
15mod bitmap;
19mod encoded_array;
20mod index;
21pub mod segment;
22mod serde;
23pub mod version;
24
25use lance_core::deepsize::DeepSizeOf;
26pub use index::FragmentRowIdIndex;
28pub use index::RowIdIndex;
29use lance_core::{Error, Result};
30use lance_io::ReadBatchParams;
31use lance_select::{RowAddrMask, RowAddrTreeMap, RowSetOps};
32pub use serde::{read_row_ids, write_row_ids};
33
34use crate::utils::LanceIteratorExtension;
35use segment::U64Segment;
36use tracing::instrument;
37
38#[derive(Debug, Clone, DeepSizeOf, PartialEq, Eq, Default)]
51pub struct RowIdSequence(Vec<U64Segment>);
52
53impl std::fmt::Display for RowIdSequence {
54 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
55 let mut iter = self.iter();
56 let mut first_10 = Vec::new();
57 let mut last_10 = Vec::new();
58 for row_id in iter.by_ref() {
59 first_10.push(row_id);
60 if first_10.len() > 10 {
61 break;
62 }
63 }
64
65 while let Some(row_id) = iter.next_back() {
66 last_10.push(row_id);
67 if last_10.len() > 10 {
68 break;
69 }
70 }
71 last_10.reverse();
72
73 let theres_more = iter.next().is_some();
74
75 write!(f, "[")?;
76 for row_id in first_10 {
77 write!(f, "{}", row_id)?;
78 }
79 if theres_more {
80 write!(f, ", ...")?;
81 }
82 for row_id in last_10 {
83 write!(f, ", {}", row_id)?;
84 }
85 write!(f, "]")
86 }
87}
88
89impl From<Range<u64>> for RowIdSequence {
90 fn from(range: Range<u64>) -> Self {
91 Self(vec![U64Segment::Range(range)])
92 }
93}
94
95impl From<&[u64]> for RowIdSequence {
96 fn from(row_ids: &[u64]) -> Self {
97 Self(vec![U64Segment::from_slice(row_ids)])
98 }
99}
100
101impl RowIdSequence {
102 pub fn new() -> Self {
103 Self::default()
104 }
105
106 pub fn iter(&self) -> impl DoubleEndedIterator<Item = u64> + '_ {
107 self.0.iter().flat_map(|segment| segment.iter())
108 }
109
110 pub fn len(&self) -> u64 {
111 self.0.iter().map(|segment| segment.len() as u64).sum()
112 }
113
114 pub fn is_empty(&self) -> bool {
115 self.0.is_empty()
116 }
117
118 pub fn row_id_range(&self) -> Option<RangeInclusive<u64>> {
126 let min = self
127 .0
128 .iter()
129 .filter_map(|s| s.range())
130 .map(|r| *r.start())
131 .min()?;
132 let max = self
133 .0
134 .iter()
135 .filter_map(|s| s.range())
136 .map(|r| *r.end())
137 .max()?;
138 Some(min..=max)
139 }
140
141 pub fn extend(&mut self, other: Self) {
143 if let (Some(U64Segment::Range(range1)), Some(U64Segment::Range(range2))) =
147 (self.0.last(), other.0.first())
148 && range1.end == range2.start
149 {
150 let new_range = U64Segment::Range(range1.start..range2.end);
151 self.0.pop();
152 self.0.push(new_range);
153 self.0.extend(other.0.into_iter().skip(1));
154 return;
155 }
156 self.0.extend(other.0);
158 }
159
160 pub fn delete(&mut self, row_ids: impl IntoIterator<Item = u64>) {
162 let (row_ids, offsets) = self.find_ids(row_ids);
164
165 let capacity = self.0.capacity();
166 let old_segments = std::mem::replace(&mut self.0, Vec::with_capacity(capacity));
167 let mut remaining_segments = old_segments.as_slice();
168
169 for (segment_idx, range) in offsets {
170 let segments_handled = old_segments.len() - remaining_segments.len();
171 let segments_to_add = segment_idx - segments_handled;
172 self.0
173 .extend_from_slice(&remaining_segments[..segments_to_add]);
174 remaining_segments = &remaining_segments[segments_to_add..];
175
176 let segment;
177 (segment, remaining_segments) = remaining_segments.split_first().unwrap();
178
179 let segment_ids = &row_ids[range];
180 self.0.push(segment.delete(segment_ids));
181 }
182
183 self.0.extend_from_slice(remaining_segments);
185 }
186
187 pub fn mask(&mut self, positions: impl IntoIterator<Item = u32>) -> Result<()> {
189 let mut local_positions = Vec::new();
190 let mut positions_iter = positions.into_iter();
191 let mut curr_position = positions_iter.next();
192 let mut offset = 0;
193 let mut cutoff = 0;
194
195 for segment in &mut self.0 {
196 cutoff += segment.len() as u32;
198 while let Some(position) = curr_position {
199 if position >= cutoff {
200 break;
201 }
202 local_positions.push(position - offset);
203 curr_position = positions_iter.next();
204 }
205
206 if !local_positions.is_empty() {
207 segment.mask(&local_positions);
208 local_positions.clear();
209 }
210 offset = cutoff;
211 }
212
213 self.0.retain(|segment| !segment.is_empty());
214
215 Ok(())
216 }
217
218 fn find_ids(
224 &self,
225 row_ids: impl IntoIterator<Item = u64>,
226 ) -> (Vec<u64>, Vec<(usize, Range<usize>)>) {
227 let mut segment_iter = self.0.iter().enumerate().cycle();
231
232 let mut segment_matches = vec![Vec::new(); self.0.len()];
233
234 row_ids.into_iter().for_each(|row_id| {
235 let mut i = 0;
236 while i < self.0.len() {
238 let (segment_idx, segment) = segment_iter.next().unwrap();
239 if segment.range().is_some_and(|range| range.contains(&row_id))
240 && let Some(offset) = segment.position(row_id)
241 {
242 segment_matches.get_mut(segment_idx).unwrap().push(offset);
243 }
245 i += 1;
246 }
247 });
248 for matches in &mut segment_matches {
249 matches.sort_unstable();
250 }
251
252 let mut offset = 0;
253 let segment_ranges = segment_matches
254 .iter()
255 .enumerate()
256 .filter(|(_, matches)| !matches.is_empty())
257 .map(|(segment_idx, matches)| {
258 let range = offset..offset + matches.len();
259 offset += matches.len();
260 (segment_idx, range)
261 })
262 .collect();
263 let row_ids = segment_matches
264 .into_iter()
265 .enumerate()
266 .flat_map(|(segment_idx, offset)| {
267 offset
268 .into_iter()
269 .map(move |offset| self.0[segment_idx].get(offset).unwrap())
270 })
271 .collect();
272
273 (row_ids, segment_ranges)
274 }
275
276 pub fn slice(&self, offset: usize, len: usize) -> RowIdSeqSlice<'_> {
277 if len == 0 {
278 return RowIdSeqSlice {
279 segments: &[],
280 offset_start: 0,
281 offset_last: 0,
282 };
283 }
284
285 let mut offset_start = offset;
287 let mut segment_offset = 0;
288 for segment in &self.0 {
289 let segment_len = segment.len();
290 if offset_start < segment_len {
291 break;
292 }
293 offset_start -= segment_len;
294 segment_offset += 1;
295 }
296
297 let mut offset_last = offset_start + len;
299 let mut segment_offset_last = segment_offset;
300 for segment in &self.0[segment_offset..] {
301 let segment_len = segment.len();
302 if offset_last <= segment_len {
303 break;
304 }
305 offset_last -= segment_len;
306 segment_offset_last += 1;
307 }
308
309 RowIdSeqSlice {
310 segments: &self.0[segment_offset..=segment_offset_last],
311 offset_start,
312 offset_last,
313 }
314 }
315
316 pub fn get(&self, index: usize) -> Option<u64> {
320 let mut offset = 0;
321 for segment in &self.0 {
322 let segment_len = segment.len();
323 if index < offset + segment_len {
324 return segment.get(index - offset);
325 }
326 offset += segment_len;
327 }
328 None
329 }
330
331 pub fn select<'a>(
339 &'a self,
340 selection: impl Iterator<Item = usize> + 'a,
341 ) -> impl Iterator<Item = u64> + 'a {
342 let mut seg_iter = self.0.iter();
343 let mut cur_seg = seg_iter.next();
344 let mut rows_passed = 0;
345 let mut cur_seg_len = cur_seg.map(|seg| seg.len()).unwrap_or(0);
346 let mut last_index = 0;
347 selection.filter_map(move |index| {
348 if index < last_index {
349 panic!("Selection is not sorted");
350 }
351 last_index = index;
352
353 cur_seg?;
354
355 while (index - rows_passed) >= cur_seg_len {
356 rows_passed += cur_seg_len;
357 cur_seg = seg_iter.next();
358 cur_seg_len = cur_seg?.len();
359 }
360
361 Some(cur_seg.unwrap().get(index - rows_passed).unwrap())
362 })
363 }
364
365 #[instrument(level = "debug", skip_all)]
380 pub fn mask_to_offset_ranges(&self, mask: &RowAddrMask) -> Vec<Range<u64>> {
381 let mut offset = 0;
382 let mut ranges = Vec::new();
383 for segment in &self.0 {
384 match segment {
385 U64Segment::Range(range) => {
386 let mut ids = RowAddrTreeMap::from(range.clone());
387 ids.mask(mask);
388 let mut cur: Option<Range<u64>> = None;
392 for (fragment, run) in unsafe { ids.iter_runs() } {
393 let frag = u64::from(fragment);
394 let run_start = (frag << 32) | u64::from(*run.start());
395 let run_end_excl = (frag << 32) | (u64::from(*run.end()) + 1);
396 let start = run_start - range.start + offset;
397 let end = run_end_excl - range.start + offset;
398 match cur.as_mut() {
399 Some(c) if c.end == start => c.end = end,
400 Some(c) => {
401 ranges.push(std::mem::replace(c, start..end));
402 }
403 None => cur = Some(start..end),
404 }
405 }
406 if let Some(c) = cur {
407 ranges.push(c);
408 }
409 offset += range.end - range.start;
410 }
411 U64Segment::RangeWithHoles { range, holes } => {
412 let offset_start = offset;
413 let mut ids = RowAddrTreeMap::from(range.clone());
414 offset += range.end - range.start;
415 for hole in holes.iter() {
416 if ids.remove(hole) {
417 offset -= 1;
418 }
419 }
420 ids.mask(mask);
421
422 let mut sorted_holes = holes.clone().into_iter().collect::<Vec<_>>();
424 sorted_holes.sort_unstable();
425 let mut next_holes_iter = sorted_holes.into_iter().peekable();
426 let mut holes_passed = 0;
427 ranges.extend(GroupingIterator::new(unsafe { ids.into_addr_iter() }.map(
428 |addr| {
429 while let Some(next_hole) = next_holes_iter.peek() {
430 if *next_hole < addr {
431 next_holes_iter.next();
432 holes_passed += 1;
433 } else {
434 break;
435 }
436 }
437 addr - range.start + offset_start - holes_passed
438 },
439 )));
440 }
441 U64Segment::RangeWithBitmap { range, bitmap } => {
442 let mut ids = RowAddrTreeMap::from(range.clone());
443 let offset_start = offset;
444 offset += range.end - range.start;
445 for (i, val) in range.clone().enumerate() {
446 if !bitmap.get(i) && ids.remove(val) {
447 offset -= 1;
448 }
449 }
450 ids.mask(mask);
451 let mut bitmap_iter = bitmap.iter();
452 let mut bitmap_iter_pos = 0;
453 let mut holes_passed = 0;
454 ranges.extend(GroupingIterator::new(unsafe { ids.into_addr_iter() }.map(
455 |addr| {
456 let position_in_range = addr - range.start;
457 while bitmap_iter_pos < position_in_range {
458 if !bitmap_iter.next().unwrap() {
459 holes_passed += 1;
460 }
461 bitmap_iter_pos += 1;
462 }
463 offset_start + position_in_range - holes_passed
464 },
465 )));
466 }
467 U64Segment::SortedArray(array) | U64Segment::Array(array) => {
468 ranges.extend(GroupingIterator::new(array.iter().enumerate().filter_map(
470 |(off, id)| {
471 if mask.selected(id) {
472 Some(off as u64 + offset)
473 } else {
474 None
475 }
476 },
477 )));
478 offset += array.len() as u64;
479 }
480 }
481 }
482 ranges
483 }
484}
485
486struct GroupingIterator<I: Iterator<Item = u64>> {
491 iter: I,
492 cur_range: Option<Range<u64>>,
493}
494
495impl<I: Iterator<Item = u64>> GroupingIterator<I> {
496 fn new(iter: I) -> Self {
497 Self {
498 iter,
499 cur_range: None,
500 }
501 }
502}
503
504impl<I: Iterator<Item = u64>> Iterator for GroupingIterator<I> {
505 type Item = Range<u64>;
506
507 fn next(&mut self) -> Option<Self::Item> {
508 for id in self.iter.by_ref() {
509 if let Some(range) = self.cur_range.as_mut() {
510 if range.end == id {
511 range.end = id + 1;
512 } else {
513 let ret = Some(range.clone());
514 self.cur_range = Some(id..id + 1);
515 return ret;
516 }
517 } else {
518 self.cur_range = Some(id..id + 1);
519 }
520 }
521 self.cur_range.take()
522 }
523}
524
525impl From<&RowIdSequence> for RowAddrTreeMap {
526 fn from(row_ids: &RowIdSequence) -> Self {
527 let mut tree_map = Self::new();
528 for segment in &row_ids.0 {
529 let mut seg = Self::new();
530 match segment {
531 U64Segment::Range(range) => {
532 seg.insert_range(range.clone());
533 }
534 U64Segment::RangeWithBitmap { range, bitmap } => {
535 seg.insert_range(range.clone());
536 for (i, val) in range.clone().enumerate() {
537 if !bitmap.get(i) {
538 seg.remove(val);
539 }
540 }
541 }
542 U64Segment::RangeWithHoles { range, holes } => {
543 seg.insert_range(range.clone());
544 for hole in holes.iter() {
545 seg.remove(hole);
546 }
547 }
548 U64Segment::SortedArray(array) | U64Segment::Array(array) => {
549 for val in array.iter() {
550 seg.insert(val);
551 }
552 }
553 }
554 tree_map |= seg;
555 }
556 tree_map
557 }
558}
559
560#[derive(Debug)]
561pub struct RowIdSeqSlice<'a> {
562 segments: &'a [U64Segment],
564 offset_start: usize,
566 offset_last: usize,
568}
569
570impl RowIdSeqSlice<'_> {
571 pub fn iter(&self) -> impl Iterator<Item = u64> + '_ {
572 let mut known_size = self.segments.iter().map(|segment| segment.len()).sum();
573 known_size -= self.offset_start;
574 known_size -= self.segments.last().map(|s| s.len()).unwrap_or_default() - self.offset_last;
575
576 let end = self.segments.len().saturating_sub(1);
577 self.segments
578 .iter()
579 .enumerate()
580 .flat_map(move |(i, segment)| {
581 match i {
582 0 if self.segments.len() == 1 => {
583 let len = self.offset_last - self.offset_start;
584 Box::new(segment.iter().skip(self.offset_start).take(len))
587 as Box<dyn Iterator<Item = u64>>
588 }
589 0 => Box::new(segment.iter().skip(self.offset_start)),
590 i if i == end => Box::new(segment.iter().take(self.offset_last)),
591 _ => Box::new(segment.iter()),
592 }
593 })
594 .exact_size(known_size)
595 }
596}
597
598pub fn rechunk_sequences(
612 sequences: impl IntoIterator<Item = RowIdSequence>,
613 chunk_sizes: impl IntoIterator<Item = u64>,
614 allow_incomplete: bool,
615) -> Result<Vec<RowIdSequence>> {
616 let chunk_sizes_vec: Vec<u64> = chunk_sizes.into_iter().collect();
618 let total_chunks = chunk_sizes_vec.len();
619 let mut chunked_sequences = Vec::with_capacity(total_chunks);
620 let mut segment_iter = sequences
621 .into_iter()
622 .flat_map(|sequence| sequence.0.into_iter())
623 .peekable();
624
625 let too_few_segments_error = |chunk_index: usize, expected_chunk_size: u64, remaining: u64| {
626 Error::invalid_input(format!(
627 "Got too few segments for chunk {}. Expected chunk size: {}, remaining needed: {}",
628 chunk_index, expected_chunk_size, remaining
629 ))
630 };
631
632 let too_many_segments_error = |processed_chunks: usize, total_chunk_sizes: usize| {
633 Error::invalid_input(format!(
634 "Got too many segments for the provided chunk lengths. Processed {} chunks out of {} expected",
635 processed_chunks, total_chunk_sizes
636 ))
637 };
638
639 let mut segment_offset = 0_u64;
640
641 for (chunk_index, chunk_size) in chunk_sizes_vec.iter().enumerate() {
642 let chunk_size = *chunk_size;
643 let mut sequence = RowIdSequence(Vec::new());
644 let mut remaining = chunk_size;
645
646 while remaining > 0 {
647 let remaining_in_segment = segment_iter
648 .peek()
649 .map_or(0, |segment| segment.len() as u64 - segment_offset);
650
651 if remaining_in_segment == 0 {
653 if segment_iter.next().is_some() {
654 segment_offset = 0;
655 continue;
656 } else {
657 if allow_incomplete {
659 break;
660 } else {
661 return Err(too_few_segments_error(chunk_index, chunk_size, remaining));
662 }
663 }
664 }
665
666 match remaining_in_segment.cmp(&remaining) {
668 std::cmp::Ordering::Greater => {
669 let segment = segment_iter
671 .peek()
672 .ok_or_else(|| too_few_segments_error(chunk_index, chunk_size, remaining))?
673 .slice(segment_offset as usize, remaining as usize);
674 sequence.extend(RowIdSequence(vec![segment]));
675 segment_offset += remaining;
676 remaining = 0;
677 }
678 std::cmp::Ordering::Equal | std::cmp::Ordering::Less => {
679 let segment = segment_iter
683 .next()
684 .ok_or_else(|| too_few_segments_error(chunk_index, chunk_size, remaining))?
685 .slice(segment_offset as usize, remaining_in_segment as usize);
686 sequence.extend(RowIdSequence(vec![segment]));
687 segment_offset = 0;
688 remaining -= remaining_in_segment;
689 }
690 }
691 }
692
693 chunked_sequences.push(sequence);
694 }
695
696 if segment_iter.peek().is_some() {
697 return Err(too_many_segments_error(
698 chunked_sequences.len(),
699 total_chunks,
700 ));
701 }
702
703 Ok(chunked_sequences)
704}
705
706pub fn select_row_ids<'a>(
708 sequence: &'a RowIdSequence,
709 offsets: &'a ReadBatchParams,
710) -> Result<Vec<u64>> {
711 let out_of_bounds_err = |offset: u32| {
712 Error::invalid_input(format!(
713 "Index out of bounds: {} for sequence of length {}",
714 offset,
715 sequence.len()
716 ))
717 };
718
719 match offsets {
720 ReadBatchParams::Indices(indices) => indices
722 .values()
723 .iter()
724 .map(|index| {
725 sequence
726 .get(*index as usize)
727 .ok_or_else(|| out_of_bounds_err(*index))
728 })
729 .collect(),
730 ReadBatchParams::Range(range) => {
731 if range.end > sequence.len() as usize {
732 return Err(out_of_bounds_err(range.end as u32));
733 }
734 let sequence = sequence.slice(range.start, range.end - range.start);
735 Ok(sequence.iter().collect())
736 }
737 ReadBatchParams::Ranges(ranges) => {
738 let num_rows = ranges
739 .iter()
740 .map(|r| (r.end - r.start) as usize)
741 .sum::<usize>();
742 let mut result = Vec::with_capacity(num_rows);
743 for range in ranges.as_ref() {
744 if range.end > sequence.len() {
745 return Err(out_of_bounds_err(range.end as u32));
746 }
747 let sequence =
748 sequence.slice(range.start as usize, (range.end - range.start) as usize);
749 result.extend(sequence.iter());
750 }
751 Ok(result)
752 }
753
754 ReadBatchParams::RangeFull => Ok(sequence.iter().collect()),
755 ReadBatchParams::RangeTo(to) => {
756 if to.end > sequence.len() as usize {
757 return Err(out_of_bounds_err(to.end as u32));
758 }
759 let len = to.end;
760 let sequence = sequence.slice(0, len);
761 Ok(sequence.iter().collect())
762 }
763 ReadBatchParams::RangeFrom(from) => {
764 let sequence = sequence.slice(from.start, sequence.len() as usize - from.start);
765 Ok(sequence.iter().collect())
766 }
767 }
768}
769
770#[cfg(test)]
771mod test {
772 use super::*;
773
774 use pretty_assertions::assert_eq;
775 use test::bitmap::Bitmap;
776
777 #[test]
778 fn test_row_id_sequence_from_range() {
779 let sequence = RowIdSequence::from(0..10);
780 assert_eq!(sequence.len(), 10);
781 assert_eq!(sequence.is_empty(), false);
782
783 let iter = sequence.iter();
784 assert_eq!(iter.collect::<Vec<_>>(), (0..10).collect::<Vec<_>>());
785 }
786
787 #[test]
788 fn test_row_id_sequence_extend() {
789 let mut sequence = RowIdSequence::from(0..10);
790 sequence.extend(RowIdSequence::from(10..20));
791 assert_eq!(sequence.0, vec![U64Segment::Range(0..20)]);
792
793 let mut sequence = RowIdSequence::from(0..10);
794 sequence.extend(RowIdSequence::from(20..30));
795 assert_eq!(
796 sequence.0,
797 vec![U64Segment::Range(0..10), U64Segment::Range(20..30)]
798 );
799 }
800
801 #[test]
802 fn test_row_id_sequence_delete() {
803 let mut sequence = RowIdSequence::from(0..10);
804 sequence.delete(vec![1, 3, 5, 7, 9]);
805 let mut expected_bitmap = Bitmap::new_empty(9);
806 for i in [0, 2, 4, 6, 8] {
807 expected_bitmap.set(i as usize);
808 }
809 assert_eq!(
810 sequence.0,
811 vec![U64Segment::RangeWithBitmap {
812 range: 0..9,
813 bitmap: expected_bitmap
814 },]
815 );
816
817 let mut sequence = RowIdSequence::from(0..10);
818 sequence.extend(RowIdSequence::from(12..20));
819 sequence.delete(vec![0, 9, 10, 11, 12, 13]);
820 assert_eq!(
821 sequence.0,
822 vec![U64Segment::Range(1..9), U64Segment::Range(14..20),]
823 );
824
825 let mut sequence = RowIdSequence::from(0..10);
826 sequence.delete(vec![0, 1, 2, 3, 4, 5, 6, 7, 8, 9]);
827 assert_eq!(sequence.0, vec![U64Segment::Range(0..0)]);
828 }
829
830 #[test]
831 fn test_row_id_slice() {
832 let sequence = RowIdSequence(vec![
835 U64Segment::Range(30..35), U64Segment::RangeWithHoles {
837 range: 50..60,
839 holes: vec![53, 54].into(),
840 },
841 U64Segment::SortedArray(vec![7, 9].into()), U64Segment::RangeWithBitmap {
843 range: 0..5,
844 bitmap: [true, false, true, false, true].as_slice().into(),
845 },
846 U64Segment::Array(vec![35, 39].into()),
847 U64Segment::Range(40..50),
848 ]);
849
850 for offset in 0..sequence.len() as usize {
852 for len in 0..sequence.len() as usize {
853 if offset + len > sequence.len() as usize {
854 continue;
855 }
856 let slice = sequence.slice(offset, len);
857
858 let actual = slice.iter().collect::<Vec<_>>();
859 let expected = sequence.iter().skip(offset).take(len).collect::<Vec<_>>();
860 assert_eq!(
861 actual, expected,
862 "Failed for offset {} and len {}",
863 offset, len
864 );
865
866 let (claimed_size, claimed_max) = slice.iter().size_hint();
867 assert_eq!(claimed_max, Some(claimed_size)); assert_eq!(claimed_size, actual.len()); }
870 }
871 }
872
873 #[test]
874 fn test_row_id_slice_empty() {
875 let sequence = RowIdSequence::from(0..10);
876 let slice = sequence.slice(10, 0);
877 assert_eq!(slice.iter().collect::<Vec<_>>(), Vec::<u64>::new());
878 }
879
880 #[test]
881 fn test_row_id_sequence_rechunk() {
882 fn assert_rechunked(
883 input: Vec<RowIdSequence>,
884 chunk_sizes: Vec<u64>,
885 expected: Vec<RowIdSequence>,
886 ) {
887 let chunked = rechunk_sequences(input, chunk_sizes, false).unwrap();
888 assert_eq!(chunked, expected);
889 }
890
891 let many_segments = vec![
893 RowIdSequence(vec![U64Segment::Range(0..5), U64Segment::Range(35..40)]),
894 RowIdSequence::from(10..18),
895 RowIdSequence::from(18..28),
896 RowIdSequence::from(28..30),
897 ];
898 let fewer_segments = vec![
899 RowIdSequence(vec![U64Segment::Range(0..5), U64Segment::Range(35..40)]),
900 RowIdSequence::from(10..30),
901 ];
902 assert_rechunked(
903 many_segments.clone(),
904 fewer_segments.iter().map(|seq| seq.len()).collect(),
905 fewer_segments.clone(),
906 );
907
908 assert_rechunked(
910 fewer_segments,
911 many_segments.iter().map(|seq| seq.len()).collect(),
912 many_segments.clone(),
913 );
914
915 assert_rechunked(
917 many_segments.clone(),
918 many_segments.iter().map(|seq| seq.len()).collect(),
919 many_segments.clone(),
920 );
921
922 let result = rechunk_sequences(many_segments.clone(), vec![100], false);
924 assert!(result.is_err());
925
926 let result = rechunk_sequences(many_segments, vec![5], false);
928 assert!(result.is_err());
929 }
930
931 #[test]
932 fn test_select_row_ids() {
933 let offsets = [
935 ReadBatchParams::Indices(vec![1, 3, 9, 5, 7, 6].into()),
936 ReadBatchParams::Range(2..8),
937 ReadBatchParams::RangeFull,
938 ReadBatchParams::RangeTo(..5),
939 ReadBatchParams::RangeFrom(5..),
940 ReadBatchParams::Ranges(vec![2..3, 5..10].into()),
941 ];
942
943 let sequences = [
946 RowIdSequence(vec![
947 U64Segment::Range(0..5),
948 U64Segment::RangeWithHoles {
949 range: 50..60,
950 holes: vec![53, 54].into(),
951 },
952 U64Segment::SortedArray(vec![7, 9].into()),
953 ]),
954 RowIdSequence(vec![
955 U64Segment::RangeWithBitmap {
956 range: 0..5,
957 bitmap: [true, false, true, false, true].as_slice().into(),
958 },
959 U64Segment::Array(vec![30, 20, 10].into()),
960 U64Segment::Range(40..50),
961 ]),
962 ];
963
964 for params in offsets {
965 for sequence in &sequences {
966 let row_ids = select_row_ids(sequence, ¶ms).unwrap();
967 let flat_sequence = sequence.iter().collect::<Vec<_>>();
968
969 let selection: Vec<usize> = match ¶ms {
971 ReadBatchParams::RangeFull => (0..flat_sequence.len()).collect(),
972 ReadBatchParams::RangeTo(to) => (0..to.end).collect(),
973 ReadBatchParams::RangeFrom(from) => (from.start..flat_sequence.len()).collect(),
974 ReadBatchParams::Range(range) => range.clone().collect(),
975 ReadBatchParams::Ranges(ranges) => ranges
976 .iter()
977 .flat_map(|r| r.start as usize..r.end as usize)
978 .collect(),
979 ReadBatchParams::Indices(indices) => {
980 indices.values().iter().map(|i| *i as usize).collect()
981 }
982 };
983
984 let expected = selection
985 .into_iter()
986 .map(|i| flat_sequence[i])
987 .collect::<Vec<_>>();
988 assert_eq!(
989 row_ids, expected,
990 "Failed for params {:?} on the sequence {:?}",
991 ¶ms, sequence
992 );
993 }
994 }
995 }
996
997 #[test]
998 fn test_select_row_ids_out_of_bounds() {
999 let offsets = [
1000 ReadBatchParams::Indices(vec![1, 1000, 4].into()),
1001 ReadBatchParams::Range(2..1000),
1002 ReadBatchParams::RangeTo(..1000),
1003 ];
1004
1005 let sequence = RowIdSequence::from(0..10);
1006
1007 for params in offsets {
1008 let result = select_row_ids(&sequence, ¶ms);
1009 assert!(result.is_err());
1010 assert!(matches!(result.unwrap_err(), Error::InvalidInput { .. }));
1011 }
1012 }
1013
1014 #[test]
1015 fn test_row_id_sequence_to_treemap() {
1016 let sequence = RowIdSequence(vec![
1017 U64Segment::Range(0..5),
1018 U64Segment::RangeWithHoles {
1019 range: 50..60,
1020 holes: vec![53, 54].into(),
1021 },
1022 U64Segment::SortedArray(vec![7, 9].into()),
1023 U64Segment::RangeWithBitmap {
1024 range: 10..15,
1025 bitmap: [true, false, true, false, true].as_slice().into(),
1026 },
1027 U64Segment::Array(vec![35, 39].into()),
1028 U64Segment::Range(40..50),
1029 ]);
1030
1031 let tree_map = RowAddrTreeMap::from(&sequence);
1032 let expected = vec![
1033 0, 1, 2, 3, 4, 7, 9, 10, 12, 14, 35, 39, 40, 41, 42, 43, 44, 45, 46, 47, 48, 49, 50,
1034 51, 52, 55, 56, 57, 58, 59,
1035 ]
1036 .into_iter()
1037 .collect::<RowAddrTreeMap>();
1038 assert_eq!(tree_map, expected);
1039 }
1040
1041 #[test]
1042 fn test_row_id_sequence_to_treemap_overlapping_segments() {
1043 let sequence = RowIdSequence(vec![
1047 U64Segment::RangeWithBitmap {
1048 range: 0..6,
1049 bitmap: [true, false, true, false, true, false].as_slice().into(),
1050 },
1051 U64Segment::RangeWithBitmap {
1052 range: 0..6,
1053 bitmap: [false, true, false, true, false, true].as_slice().into(),
1054 },
1055 ]);
1056
1057 let expected = sequence.iter().collect::<RowAddrTreeMap>();
1058 assert_eq!(expected, (0..6).collect::<RowAddrTreeMap>());
1059 assert_eq!(RowAddrTreeMap::from(&sequence), expected);
1060 }
1061
1062 #[test]
1063 fn test_row_addr_mask() {
1064 let sequence = RowIdSequence(vec![
1070 U64Segment::Range(0..5),
1071 U64Segment::RangeWithHoles {
1072 range: 50..60,
1073 holes: vec![53, 54].into(),
1074 },
1075 U64Segment::SortedArray(vec![7, 9].into()),
1076 U64Segment::RangeWithBitmap {
1077 range: 10..15,
1078 bitmap: [true, false, true, false, true].as_slice().into(),
1079 },
1080 U64Segment::Array(vec![35, 39].into()),
1081 ]);
1082
1083 let values_to_remove = [4, 55, 7, 12, 39];
1085 let positions_to_remove = sequence
1086 .iter()
1087 .enumerate()
1088 .filter_map(|(i, val)| {
1089 if values_to_remove.contains(&val) {
1090 Some(i as u32)
1091 } else {
1092 None
1093 }
1094 })
1095 .collect::<Vec<_>>();
1096 let mut sequence = sequence;
1097 sequence.mask(positions_to_remove).unwrap();
1098 let expected = RowIdSequence(vec![
1099 U64Segment::Range(0..4),
1100 U64Segment::RangeWithBitmap {
1101 range: 50..60,
1102 bitmap: [
1103 true, true, true, false, false, false, true, true, true, true,
1104 ]
1105 .as_slice()
1106 .into(),
1107 },
1108 U64Segment::Range(9..10),
1109 U64Segment::RangeWithBitmap {
1110 range: 10..15,
1111 bitmap: [true, false, false, false, true].as_slice().into(),
1112 },
1113 U64Segment::Array(vec![35].into()),
1114 ]);
1115 assert_eq!(sequence, expected);
1116 }
1117
1118 #[test]
1119 fn test_row_addr_mask_everything() {
1120 let mut sequence = RowIdSequence(vec![
1121 U64Segment::Range(0..5),
1122 U64Segment::SortedArray(vec![7, 9].into()),
1123 ]);
1124 sequence.mask(0..sequence.len() as u32).unwrap();
1125 let expected = RowIdSequence(vec![]);
1126 assert_eq!(sequence, expected);
1127 }
1128
1129 #[test]
1130 fn test_selection() {
1131 let sequence = RowIdSequence(vec![
1132 U64Segment::Range(0..5),
1133 U64Segment::Range(10..15),
1134 U64Segment::Range(20..25),
1135 ]);
1136 let selection = sequence.select(vec![2, 4, 13, 14, 57].into_iter());
1137 assert_eq!(selection.collect::<Vec<_>>(), vec![2, 4, 23, 24]);
1138 }
1139
1140 #[test]
1141 #[should_panic]
1142 fn test_selection_unsorted() {
1143 let sequence = RowIdSequence(vec![
1144 U64Segment::Range(0..5),
1145 U64Segment::Range(10..15),
1146 U64Segment::Range(20..25),
1147 ]);
1148 let _ = sequence
1149 .select(vec![2, 4, 3].into_iter())
1150 .collect::<Vec<_>>();
1151 }
1152
1153 #[test]
1154 fn test_mask_to_offset_ranges() {
1155 let sequence = RowIdSequence(vec![U64Segment::Range(0..10)]);
1157 let mask = RowAddrMask::from_allowed(RowAddrTreeMap::from_iter(&[0, 2, 4, 6, 8]));
1158 let ranges = sequence.mask_to_offset_ranges(&mask);
1159 assert_eq!(ranges, vec![0..1, 2..3, 4..5, 6..7, 8..9]);
1160
1161 let sequence = RowIdSequence(vec![U64Segment::Range(40..60)]);
1162 let mask = RowAddrMask::from_allowed(RowAddrTreeMap::from_iter(&[54]));
1163 let ranges = sequence.mask_to_offset_ranges(&mask);
1164 assert_eq!(ranges, vec![14..15]);
1165
1166 let sequence = RowIdSequence(vec![U64Segment::Range(40..60)]);
1167 let mask = RowAddrMask::from_block(RowAddrTreeMap::from_iter(&[54]));
1168 let ranges = sequence.mask_to_offset_ranges(&mask);
1169 assert_eq!(ranges, vec![0..14, 15..20]);
1170
1171 let sequence = RowIdSequence(vec![U64Segment::RangeWithHoles {
1174 range: 0..10,
1175 holes: vec![2, 6].into(),
1176 }]);
1177 let mask = RowAddrMask::from_allowed(RowAddrTreeMap::from_iter(&[0, 2, 4, 6, 8]));
1178 let ranges = sequence.mask_to_offset_ranges(&mask);
1179 assert_eq!(ranges, vec![0..1, 3..4, 6..7]);
1180
1181 let sequence = RowIdSequence(vec![U64Segment::RangeWithHoles {
1182 range: 40..60,
1183 holes: vec![47, 43].into(),
1184 }]);
1185 let mask = RowAddrMask::from_allowed(RowAddrTreeMap::from_iter(&[44]));
1186 let ranges = sequence.mask_to_offset_ranges(&mask);
1187 assert_eq!(ranges, vec![3..4]);
1188
1189 let sequence = RowIdSequence(vec![U64Segment::RangeWithHoles {
1190 range: 40..60,
1191 holes: vec![47, 43].into(),
1192 }]);
1193 let mask = RowAddrMask::from_block(RowAddrTreeMap::from_iter(&[44]));
1194 let ranges = sequence.mask_to_offset_ranges(&mask);
1195 assert_eq!(ranges, vec![0..3, 4..18]);
1196
1197 let sequence = RowIdSequence(vec![U64Segment::RangeWithBitmap {
1200 range: 0..10,
1201 bitmap: [
1202 true, true, false, false, true, true, true, true, false, false,
1203 ]
1204 .as_slice()
1205 .into(),
1206 }]);
1207 let mask = RowAddrMask::from_allowed(RowAddrTreeMap::from_iter(&[0, 2, 4, 6, 8]));
1208 let ranges = sequence.mask_to_offset_ranges(&mask);
1209 assert_eq!(ranges, vec![0..1, 2..3, 4..5]);
1210
1211 let sequence = RowIdSequence(vec![U64Segment::RangeWithBitmap {
1212 range: 40..45,
1213 bitmap: [true, true, false, false, true].as_slice().into(),
1214 }]);
1215 let mask = RowAddrMask::from_allowed(RowAddrTreeMap::from_iter(&[44]));
1216 let ranges = sequence.mask_to_offset_ranges(&mask);
1217 assert_eq!(ranges, vec![2..3]);
1218
1219 let sequence = RowIdSequence(vec![U64Segment::RangeWithBitmap {
1220 range: 40..45,
1221 bitmap: [true, true, false, false, true].as_slice().into(),
1222 }]);
1223 let mask = RowAddrMask::from_block(RowAddrTreeMap::from_iter(&[44]));
1224 let ranges = sequence.mask_to_offset_ranges(&mask);
1225 assert_eq!(ranges, vec![0..2]);
1226
1227 let sequence = RowIdSequence(vec![U64Segment::SortedArray(vec![0, 2, 4, 6, 8].into())]);
1229 let mask = RowAddrMask::from_allowed(RowAddrTreeMap::from_iter(&[0, 6, 8]));
1230 let ranges = sequence.mask_to_offset_ranges(&mask);
1231 assert_eq!(ranges, vec![0..1, 3..5]);
1232
1233 let sequence = RowIdSequence(vec![U64Segment::Array(vec![8, 2, 6, 0, 4].into())]);
1234 let mask = RowAddrMask::from_allowed(RowAddrTreeMap::from_iter(&[0, 6, 8]));
1235 let ranges = sequence.mask_to_offset_ranges(&mask);
1236 assert_eq!(ranges, vec![0..1, 2..4]);
1237
1238 let sequence = RowIdSequence(vec![
1243 U64Segment::Range(0..5),
1244 U64Segment::RangeWithHoles {
1245 range: 100..105,
1246 holes: vec![103].into(),
1247 },
1248 U64Segment::SortedArray(vec![44, 46, 78].into()),
1249 ]);
1250 let mask = RowAddrMask::from_allowed(RowAddrTreeMap::from_iter(&[0, 2, 46, 100, 104]));
1251 let ranges = sequence.mask_to_offset_ranges(&mask);
1252 assert_eq!(ranges, vec![0..1, 2..3, 5..6, 8..9, 10..11]);
1253
1254 let sequence = RowIdSequence(vec![U64Segment::Range(0..10)]);
1256 let mask = RowAddrMask::default();
1257 let ranges = sequence.mask_to_offset_ranges(&mask);
1258 assert_eq!(ranges, vec![0..10]);
1259
1260 let sequence = RowIdSequence(vec![U64Segment::Range(0..10)]);
1262 let mask = RowAddrMask::allow_nothing();
1263 let ranges = sequence.mask_to_offset_ranges(&mask);
1264 assert_eq!(ranges, vec![]);
1265 }
1266
1267 #[test]
1268 fn test_row_id_sequence_rechunk_with_empty_segments() {
1269 let input_sequences = vec![
1271 RowIdSequence::from(0..2), RowIdSequence::from(20..23), ];
1274 let chunk_sizes = vec![2, 3]; let result = rechunk_sequences(input_sequences, chunk_sizes, false).unwrap();
1277 assert_eq!(result.len(), 2);
1278 assert_eq!(result[0].len(), 2);
1279 assert_eq!(result[1].len(), 3);
1280
1281 let first_chunk: Vec<u64> = result[0].iter().collect();
1282 let second_chunk: Vec<u64> = result[1].iter().collect();
1283 assert_eq!(first_chunk, vec![0, 1]);
1284 assert_eq!(second_chunk, vec![20, 21, 22]);
1285
1286 let input_sequences = vec![
1288 RowIdSequence::from(0..2), RowIdSequence::from(20..21), RowIdSequence::from(30..32), ];
1292 let chunk_sizes = vec![5]; let result = rechunk_sequences(input_sequences, chunk_sizes, false).unwrap();
1295 assert_eq!(result.len(), 1);
1296 assert_eq!(result[0].len(), 5);
1297
1298 let elements: Vec<u64> = result[0].iter().collect();
1299 assert_eq!(elements, vec![0, 1, 20, 30, 31]);
1300
1301 let input_sequences = vec![
1303 RowIdSequence::from(0..2), RowIdSequence::from(10..10), RowIdSequence::from(20..22), ];
1307 let chunk_sizes = vec![3, 1];
1308 let result = rechunk_sequences(input_sequences, chunk_sizes, false).unwrap();
1309
1310 assert_eq!(result.len(), 2);
1311 assert_eq!(result[0].len(), 3);
1312 assert_eq!(result[1].len(), 1);
1313
1314 let first_chunk_elements: Vec<u64> = result[0].iter().collect();
1315 let second_chunk_elements: Vec<u64> = result[1].iter().collect();
1316 assert_eq!(first_chunk_elements, vec![0, 1, 20]);
1317 assert_eq!(second_chunk_elements, vec![21]);
1318
1319 let input_sequences = vec![
1321 RowIdSequence::from(0..1), RowIdSequence::from(10..10), RowIdSequence::from(20..20), RowIdSequence::from(30..32), ];
1326 let chunk_sizes = vec![3];
1327 let result = rechunk_sequences(input_sequences, chunk_sizes, false).unwrap();
1328
1329 assert_eq!(result.len(), 1);
1330 assert_eq!(result[0].len(), 3);
1331
1332 let elements: Vec<u64> = result[0].iter().collect();
1333 assert_eq!(elements, vec![0, 30, 31]);
1334
1335 let input_sequences = vec![
1337 RowIdSequence::from(0..3), RowIdSequence::from(10..10), RowIdSequence::from(20..22), ];
1341 let chunk_sizes = vec![3, 2];
1342 let result = rechunk_sequences(input_sequences, chunk_sizes, false).unwrap();
1343
1344 assert_eq!(result.len(), 2);
1345 assert_eq!(result[0].len(), 3);
1346 assert_eq!(result[1].len(), 2);
1347
1348 let first_chunk_elements: Vec<u64> = result[0].iter().collect();
1349 let second_chunk_elements: Vec<u64> = result[1].iter().collect();
1350 assert_eq!(first_chunk_elements, vec![0, 1, 2]);
1351 assert_eq!(second_chunk_elements, vec![20, 21]);
1352
1353 let input_sequences = vec![
1355 RowIdSequence::from(0..2), RowIdSequence::from(10..10), ];
1358 let chunk_sizes = vec![5]; let result = rechunk_sequences(input_sequences, chunk_sizes, true).unwrap();
1360
1361 assert_eq!(result.len(), 1);
1362 assert_eq!(result[0].len(), 2);
1363
1364 let elements: Vec<u64> = result[0].iter().collect();
1365 assert_eq!(elements, vec![0, 1]);
1366 }
1367
1368 #[test]
1369 fn test_row_id_range_empty() {
1370 let seq = RowIdSequence::from(0u64..0);
1371 assert_eq!(seq.row_id_range(), None);
1372 }
1373
1374 #[test]
1375 fn test_row_id_range_single_contiguous() {
1376 let seq = RowIdSequence::from(10u64..20);
1377 assert_eq!(seq.row_id_range(), Some(10..=19));
1378 }
1379
1380 #[test]
1381 fn test_row_id_range_unsorted_array() {
1382 let seq = RowIdSequence::from([50u64, 10, 30].as_slice());
1384 let r = seq.row_id_range().unwrap();
1385 assert!(*r.start() <= 10);
1386 assert!(*r.end() >= 50);
1387 }
1388
1389 #[test]
1390 fn test_row_id_range_multi_segment() {
1391 let mut seq = RowIdSequence::from(0u64..5);
1393 seq.extend(RowIdSequence::from(100u64..105));
1394 let r = seq.row_id_range().unwrap();
1395 assert_eq!(*r.start(), 0);
1396 assert_eq!(*r.end(), 104);
1397 }
1398}