Skip to main content

lance_table/system_index/
frag_reuse.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4use std::{collections::HashMap, io::Cursor, sync::Arc};
5
6use arrow_array::cast::AsArray;
7use arrow_array::types::UInt64Type;
8use arrow_array::{Array, ArrayRef, PrimitiveArray, RecordBatch, UInt64Array};
9use lance_core::deepsize::{Context, DeepSizeOf};
10use lance_core::utils::row_addr_remap::{GroupInputWithLayout, RowAddrRemap};
11use lance_core::{Error, Result};
12use lance_select::RowAddrTreeMap;
13use roaring::{RoaringBitmap, RoaringTreemap};
14use serde::{Deserialize, Serialize};
15use uuid::Uuid;
16
17use crate::format::pb::fragment_reuse_index_details::InlineContent;
18use crate::format::{ExternalFile, Fragment, pb};
19
20pub const FRAG_REUSE_INDEX_NAME: &str = "__lance_frag_reuse";
21pub const FRAG_REUSE_DETAILS_FILE_NAME: &str = "details.binpb";
22
23#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
24pub struct FragDigest {
25    pub id: u64,
26    pub physical_rows: usize,
27    pub num_deleted_rows: usize,
28}
29
30impl From<&FragDigest> for pb::fragment_reuse_index_details::FragmentDigest {
31    fn from(digest: &FragDigest) -> Self {
32        Self {
33            id: digest.id,
34            physical_rows: digest.physical_rows as u64,
35            num_deleted_rows: digest.num_deleted_rows as u64,
36        }
37    }
38}
39
40impl From<&Fragment> for FragDigest {
41    fn from(fragment: &Fragment) -> Self {
42        Self {
43            id: fragment.id,
44            physical_rows: fragment
45                .physical_rows
46                .expect("Fragment doesn't have physical rows recorded"),
47            num_deleted_rows: fragment
48                .deletion_file
49                .as_ref()
50                .and_then(|d| d.num_deleted_rows)
51                .unwrap_or(0),
52        }
53    }
54}
55
56impl TryFrom<pb::fragment_reuse_index_details::FragmentDigest> for FragDigest {
57    type Error = Error;
58
59    fn try_from(digest: pb::fragment_reuse_index_details::FragmentDigest) -> Result<Self> {
60        Ok(Self {
61            id: digest.id,
62            physical_rows: digest.physical_rows as usize,
63            num_deleted_rows: digest.num_deleted_rows as usize,
64        })
65    }
66}
67
68#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
69pub struct FragReuseGroup {
70    pub changed_row_addrs: Vec<u8>,
71    pub old_frags: Vec<FragDigest>,
72    pub new_frags: Vec<FragDigest>,
73}
74
75impl From<&FragReuseGroup> for pb::fragment_reuse_index_details::Group {
76    fn from(group: &FragReuseGroup) -> Self {
77        Self {
78            changed_row_addrs: group.changed_row_addrs.clone(),
79            old_fragments: group.old_frags.iter().map(|f| f.into()).collect(),
80            new_fragments: group.new_frags.iter().map(|f| f.into()).collect(),
81        }
82    }
83}
84
85impl TryFrom<pb::fragment_reuse_index_details::Group> for FragReuseGroup {
86    type Error = Error;
87
88    fn try_from(group: pb::fragment_reuse_index_details::Group) -> Result<Self> {
89        Ok(Self {
90            changed_row_addrs: group.changed_row_addrs,
91            old_frags: group
92                .old_fragments
93                .into_iter()
94                .map(FragDigest::try_from)
95                .collect::<Result<_>>()?,
96            new_frags: group
97                .new_fragments
98                .into_iter()
99                .map(FragDigest::try_from)
100                .collect::<Result<_>>()?,
101        })
102    }
103}
104
105#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
106pub struct FragReuseVersion {
107    pub dataset_version: u64,
108    pub groups: Vec<FragReuseGroup>,
109}
110
111impl From<&FragReuseVersion> for pb::fragment_reuse_index_details::Version {
112    fn from(version: &FragReuseVersion) -> Self {
113        Self {
114            dataset_version: version.dataset_version,
115            groups: version.groups.iter().map(|g| g.into()).collect(),
116        }
117    }
118}
119
120impl TryFrom<pb::fragment_reuse_index_details::Version> for FragReuseVersion {
121    type Error = Error;
122
123    fn try_from(version: pb::fragment_reuse_index_details::Version) -> Result<Self> {
124        Ok(Self {
125            dataset_version: version.dataset_version,
126            groups: version
127                .groups
128                .into_iter()
129                .map(FragReuseGroup::try_from)
130                .collect::<Result<_>>()?,
131        })
132    }
133}
134
135impl FragReuseVersion {
136    pub fn old_frag_ids(&self) -> Vec<u64> {
137        self.groups
138            .iter()
139            .flat_map(|g| g.old_frags.iter().map(|f| f.id))
140            .collect::<Vec<_>>()
141    }
142
143    pub fn new_frag_ids(&self) -> Vec<u64> {
144        self.groups
145            .iter()
146            .flat_map(|g| g.new_frags.iter().map(|f| f.id))
147            .collect::<Vec<_>>()
148    }
149
150    pub fn new_frag_bitmap(&self) -> RoaringBitmap {
151        RoaringBitmap::from_iter(self.new_frag_ids().iter().map(|&id| id as u32))
152    }
153}
154
155#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
156pub enum FragReuseIndexDetailsContentType {
157    Inline(FragReuseIndexDetails),
158    External(ExternalFile),
159}
160
161#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
162pub struct FragReuseIndexDetails {
163    pub versions: Vec<FragReuseVersion>,
164}
165
166impl From<&FragReuseIndexDetails> for InlineContent {
167    fn from(details: &FragReuseIndexDetails) -> Self {
168        let mut versions: Vec<pb::fragment_reuse_index_details::Version> =
169            details.versions.iter().map(|m| m.into()).collect();
170        // sort from oldest to latest version
171        versions.sort_by_key(|v| v.dataset_version);
172        Self { versions }
173    }
174}
175
176impl TryFrom<InlineContent> for FragReuseIndexDetails {
177    type Error = Error;
178
179    fn try_from(content: InlineContent) -> Result<Self> {
180        Ok(Self {
181            versions: content
182                .versions
183                .into_iter()
184                .map(|m| m.try_into())
185                .collect::<Result<Vec<_>>>()?,
186        })
187    }
188}
189
190impl FragReuseIndexDetails {
191    pub fn new_frag_bitmap(&self) -> RoaringBitmap {
192        RoaringBitmap::from_iter(
193            self.versions
194                .iter()
195                .flat_map(|v| v.new_frag_ids().into_iter().map(|id| id as u32)),
196        )
197    }
198}
199
200/// An index that stores materialized row ID maps.
201///
202/// This type is retained for API and serde compatibility. Dataset loading uses
203/// [`CompactFragReuseIndex`] so persisted FRI details are not expanded into a
204/// hash-map entry for every affected row.
205#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
206pub struct FragReuseIndex {
207    pub uuid: Uuid,
208    pub row_id_maps: Vec<HashMap<u64, Option<u64>>>,
209    pub details: FragReuseIndexDetails,
210}
211
212impl DeepSizeOf for FragReuseIndex {
213    fn deep_size_of_children(&self, cx: &mut Context) -> usize {
214        self.row_id_maps.deep_size_of_children(cx) + self.details.deep_size_of_children(cx)
215    }
216}
217
218impl FragReuseIndex {
219    pub fn new(
220        uuid: Uuid,
221        row_id_maps: Vec<HashMap<u64, Option<u64>>>,
222        details: FragReuseIndexDetails,
223    ) -> Self {
224        Self {
225            uuid,
226            row_id_maps,
227            details,
228        }
229    }
230
231    pub fn is_empty(&self) -> bool {
232        self.row_id_maps.iter().all(HashMap::is_empty)
233    }
234
235    pub fn remap_row_id(&self, row_id: u64) -> Option<u64> {
236        let mut mapped = Some(row_id);
237        for row_id_map in &self.row_id_maps {
238            if let Some(current) = mapped {
239                mapped = row_id_map.get(&current).copied().unwrap_or(mapped);
240            }
241        }
242        mapped
243    }
244
245    pub fn remap_row_ids_in_place(&self, row_ids: &mut [Option<u64>]) {
246        for row_id_map in &self.row_id_maps {
247            for row_id in row_ids.iter_mut() {
248                if let Some(current) = *row_id
249                    && let Some(mapped) = row_id_map.get(&current)
250                {
251                    *row_id = *mapped;
252                }
253            }
254        }
255    }
256
257    pub fn remap_row_addrs_tree_map(&self, row_addrs: &RowAddrTreeMap) -> RowAddrTreeMap {
258        RowAddrTreeMap::from_iter(
259            row_addrs
260                .row_addrs()
261                .unwrap()
262                .filter_map(|addr| self.remap_row_id(u64::from(addr))),
263        )
264    }
265
266    pub fn remap_row_ids_roaring_tree_map(&self, row_ids: &RoaringTreemap) -> RoaringTreemap {
267        RoaringTreemap::from_iter(row_ids.iter().filter_map(|addr| self.remap_row_id(addr)))
268    }
269
270    pub fn remap_row_ids_record_batch(
271        &self,
272        batch: RecordBatch,
273        row_id_idx: usize,
274    ) -> Result<RecordBatch> {
275        remap_row_ids_record_batch(batch, row_id_idx, |row_ids| {
276            self.remap_row_ids_in_place(row_ids)
277        })
278    }
279
280    pub fn remap_row_ids_array(&self, array: ArrayRef) -> PrimitiveArray<UInt64Type> {
281        remap_row_ids_array(array, |row_ids| self.remap_row_ids_in_place(row_ids))
282    }
283
284    pub fn remap_fragment_bitmap(&self, fragment_bitmap: &mut RoaringBitmap) -> Result<()> {
285        remap_fragment_bitmap(&self.details, fragment_bitmap)
286    }
287}
288
289/// A compact row-address remap chain for deferred compactions.
290///
291/// Each FRI version retains rewritten-row bitmaps and fragment layouts. Queries
292/// use bitmap rank plus the ordered new-fragment ranges instead of storing one
293/// hash-map entry per affected row.
294#[derive(Debug, Clone, PartialEq, Eq)]
295pub struct CompactFragReuseIndex {
296    pub uuid: Uuid,
297    row_addr_remap: RowAddrRemap,
298    pub details: FragReuseIndexDetails,
299}
300
301impl DeepSizeOf for CompactFragReuseIndex {
302    fn deep_size_of_children(&self, cx: &mut Context) -> usize {
303        self.row_addr_remap.deep_size_of_children(cx) + self.details.deep_size_of_children(cx)
304    }
305}
306
307impl CompactFragReuseIndex {
308    #[doc(hidden)]
309    pub fn from_row_id_maps(
310        uuid: Uuid,
311        row_id_maps: Vec<HashMap<u64, Option<u64>>>,
312        details: FragReuseIndexDetails,
313    ) -> Self {
314        Self {
315            uuid,
316            row_addr_remap: RowAddrRemap::chained(
317                row_id_maps.into_iter().map(RowAddrRemap::direct),
318            ),
319            details,
320        }
321    }
322
323    /// Build a queryable index directly from serialized FRI details without
324    /// expanding each affected row into a hash map.
325    pub fn try_new(uuid: Uuid, details: FragReuseIndexDetails) -> Result<Self> {
326        let mut version_remaps = Vec::with_capacity(details.versions.len());
327        for (version_idx, version) in details.versions.iter().enumerate() {
328            let mut groups = Vec::with_capacity(version.groups.len());
329            for (group_idx, group) in version.groups.iter().enumerate() {
330                let changed_row_addrs = RoaringTreemap::deserialize_from(Cursor::new(
331                    &group.changed_row_addrs,
332                ))
333                .map_err(|error| {
334                    Error::index(format!(
335                        "failed to deserialize changed row addresses for FRI version {version_idx}, group {group_idx}: {error}"
336                    ))
337                })?;
338                let old_frags = group
339                    .old_frags
340                    .iter()
341                    .map(|frag| fragment_layout(frag, "old", version_idx, group_idx))
342                    .collect::<Result<Vec<_>>>()?;
343                let new_frags = group
344                    .new_frags
345                    .iter()
346                    .map(|frag| fragment_layout(frag, "new", version_idx, group_idx))
347                    .collect::<Result<Vec<_>>>()?;
348                groups.push(GroupInputWithLayout {
349                    rewritten_old_row_addrs: changed_row_addrs,
350                    old_frags,
351                    new_frags,
352                });
353            }
354            let remap = RowAddrRemap::compact_with_layout(groups).map_err(|error| {
355                Error::index(format!(
356                    "failed to build compact remap for FRI version {version_idx}: {error}"
357                ))
358            })?;
359            version_remaps.push(remap);
360        }
361
362        Ok(Self {
363            uuid,
364            row_addr_remap: RowAddrRemap::chained(version_remaps),
365            details,
366        })
367    }
368
369    /// The ordered remap chain used by index and transaction remapping paths.
370    pub fn row_addr_remap(&self) -> &RowAddrRemap {
371        &self.row_addr_remap
372    }
373
374    /// Returns whether the index contains no row-address remapping.
375    pub fn is_empty(&self) -> bool {
376        self.row_addr_remap.is_empty()
377    }
378
379    pub fn remap_row_id(&self, row_id: u64) -> Option<u64> {
380        self.row_addr_remap.get(row_id).unwrap_or(Some(row_id))
381    }
382
383    /// Apply all FRI versions to row addresses in place. `None` values remain
384    /// deleted and missing mappings pass through unchanged.
385    pub fn remap_row_ids_in_place(&self, row_ids: &mut [Option<u64>]) {
386        self.row_addr_remap.remap_in_place(row_ids);
387    }
388
389    pub fn remap_row_addrs_tree_map(&self, row_addrs: &RowAddrTreeMap) -> RowAddrTreeMap {
390        RowAddrTreeMap::from_iter(row_addrs.row_addrs().unwrap().filter_map(|addr| {
391            let addr_as_u64 = u64::from(addr);
392            self.remap_row_id(addr_as_u64)
393        }))
394    }
395
396    pub fn remap_row_ids_roaring_tree_map(&self, row_ids: &RoaringTreemap) -> RoaringTreemap {
397        RoaringTreemap::from_iter(row_ids.iter().filter_map(|addr| self.remap_row_id(addr)))
398    }
399
400    /// Remap a record batch that contains a row_id column at index `row_id_idx`
401    /// Currently this assumes there are only 2 columns in the schema,
402    /// which is the case for all indexes.
403    /// For example, for btree, the schema is (value, row_id).
404    /// For vector index storage, the schema is (row_id, vector).
405    pub fn remap_row_ids_record_batch(
406        &self,
407        batch: RecordBatch,
408        row_id_idx: usize,
409    ) -> Result<RecordBatch> {
410        remap_row_ids_record_batch(batch, row_id_idx, |row_ids| {
411            self.remap_row_ids_in_place(row_ids)
412        })
413    }
414
415    pub fn remap_row_ids_array(&self, array: ArrayRef) -> PrimitiveArray<UInt64Type> {
416        remap_row_ids_array(array, |row_ids| self.remap_row_ids_in_place(row_ids))
417    }
418
419    pub fn remap_fragment_bitmap(&self, fragment_bitmap: &mut RoaringBitmap) -> Result<()> {
420        remap_fragment_bitmap(&self.details, fragment_bitmap)
421    }
422}
423
424fn remap_row_ids_record_batch(
425    batch: RecordBatch,
426    row_id_idx: usize,
427    remap: impl FnOnce(&mut [Option<u64>]),
428) -> Result<RecordBatch> {
429    assert_eq!(batch.schema().fields().len(), 2);
430    let other_column_idx = 1 - row_id_idx;
431    let row_ids = batch.column(row_id_idx).as_primitive::<UInt64Type>();
432    let mut remapped_row_ids = row_ids
433        .values()
434        .iter()
435        .copied()
436        .map(Some)
437        .collect::<Vec<_>>();
438    remap(&mut remapped_row_ids);
439    let (val_indices, new_row_ids): (Vec<u64>, Vec<u64>) = remapped_row_ids
440        .iter()
441        .enumerate()
442        .filter_map(|(idx, new_id)| new_id.map(|new_id| (idx as u64, new_id)))
443        .unzip();
444    let new_val_indices = UInt64Array::from_iter_values(val_indices);
445    let new_vals = arrow::compute::take(batch.column(other_column_idx), &new_val_indices, None)?;
446
447    let mut batch_data: Vec<(usize, ArrayRef)> = vec![
448        (
449            row_id_idx,
450            Arc::new(UInt64Array::from_iter_values(new_row_ids)) as ArrayRef,
451        ),
452        (other_column_idx, Arc::new(new_vals)),
453    ];
454    batch_data.sort_by_key(|(i, _)| *i);
455    Ok(RecordBatch::try_new(
456        batch.schema(),
457        batch_data.into_iter().map(|(_, item)| item).collect(),
458    )?)
459}
460
461fn remap_row_ids_array(
462    array: ArrayRef,
463    remap: impl FnOnce(&mut [Option<u64>]),
464) -> PrimitiveArray<UInt64Type> {
465    let primitive_array = array
466        .as_any()
467        .downcast_ref::<PrimitiveArray<UInt64Type>>()
468        .expect("expected row IDs to be uint64 array");
469    let mut remapped = (0..primitive_array.len())
470        .map(|i| {
471            if primitive_array.is_null(i) {
472                None
473            } else {
474                Some(primitive_array.value(i))
475            }
476        })
477        .collect::<Vec<_>>();
478    remap(&mut remapped);
479    PrimitiveArray::from(remapped)
480}
481
482fn remap_fragment_bitmap(
483    details: &FragReuseIndexDetails,
484    fragment_bitmap: &mut RoaringBitmap,
485) -> Result<()> {
486    for version in details.versions.iter() {
487        for group in version.groups.iter() {
488            let mut removed = 0;
489            for old_frag in group.old_frags.iter() {
490                if fragment_bitmap.remove(old_frag.id as u32) {
491                    removed += 1;
492                }
493            }
494
495            if removed > 0 {
496                if removed != group.old_frags.len() {
497                    // Straddle: the index covered only part of this rewrite
498                    // group. Caused by the bug fixed in
499                    // <https://github.com/lance-format/lance/pull/6610>.
500                    // We've already removed the indexed old_frags from the
501                    // bitmap above; deliberately do NOT insert new_frags,
502                    // since the merged fragment also contains rows that
503                    // were never indexed. Affected rows fall through to
504                    // flat scan until the next optimize_indices. The fix
505                    // is persisted on the next write via build_manifest.
506                    tracing::warn!(
507                        "Healing straddling fragment-reuse rewrite group in index bitmap: \
508                             group {:?} was only partially indexed ({} of {} old fragments). \
509                             Affected rows will use flat scan until the next optimize_indices.",
510                        group.old_frags,
511                        removed,
512                        group.old_frags.len(),
513                    );
514                    continue;
515                }
516
517                for new_frag in group.new_frags.iter() {
518                    fragment_bitmap.insert(new_frag.id as u32);
519                }
520            }
521        }
522    }
523    Ok(())
524}
525
526fn fragment_layout(
527    frag: &FragDigest,
528    role: &str,
529    version_idx: usize,
530    group_idx: usize,
531) -> Result<(u32, u32)> {
532    let fragment_id = u32::try_from(frag.id).map_err(|_| {
533        Error::index(format!(
534            "FRI version {version_idx}, group {group_idx} has {role} fragment id {} outside the row-address range",
535            frag.id
536        ))
537    })?;
538    let physical_rows = u32::try_from(frag.physical_rows).map_err(|_| {
539        Error::index(format!(
540            "FRI version {version_idx}, group {group_idx} has {role} fragment {fragment_id} with physical_rows={} outside the row-address range",
541            frag.physical_rows
542        ))
543    })?;
544    Ok((fragment_id, physical_rows))
545}
546
547#[cfg(test)]
548mod tests {
549
550    use super::*;
551    use rstest::rstest;
552
553    fn addr(fragment_id: u32, offset: u32) -> u64 {
554        u64::from(lance_core::utils::address::RowAddress::new_from_parts(
555            fragment_id,
556            offset,
557        ))
558    }
559
560    fn serialize_changed(addrs: impl IntoIterator<Item = u64>) -> Vec<u8> {
561        let changed = RoaringTreemap::from_iter(addrs);
562        let mut bytes = Vec::with_capacity(changed.serialized_size());
563        changed.serialize_into(&mut bytes).unwrap();
564        bytes
565    }
566
567    fn digest(id: u64, physical_rows: usize) -> FragDigest {
568        FragDigest {
569            id,
570            physical_rows,
571            num_deleted_rows: 0,
572        }
573    }
574
575    #[test]
576    fn test_compact_fri_tristate_one_to_many_and_chain() {
577        let details = FragReuseIndexDetails {
578            versions: vec![
579                FragReuseVersion {
580                    dataset_version: 1,
581                    groups: vec![
582                        // One old fragment is split into two new fragments.
583                        FragReuseGroup {
584                            changed_row_addrs: serialize_changed([
585                                addr(1, 0),
586                                addr(1, 2),
587                                addr(1, 3),
588                            ]),
589                            old_frags: vec![digest(1, 4)],
590                            new_frags: vec![digest(10, 1), digest(11, 2)],
591                        },
592                        // A separate rewrite group deletes an entire fragment.
593                        FragReuseGroup {
594                            changed_row_addrs: serialize_changed([]),
595                            old_frags: vec![digest(3, 2)],
596                            new_frags: vec![],
597                        },
598                    ],
599                },
600                FragReuseVersion {
601                    dataset_version: 2,
602                    groups: vec![FragReuseGroup {
603                        changed_row_addrs: serialize_changed([addr(10, 0), addr(11, 1)]),
604                        old_frags: vec![digest(10, 1), digest(11, 2)],
605                        new_frags: vec![digest(20, 2)],
606                    }],
607                },
608            ],
609        };
610        let details = FragReuseIndexDetails::try_from(InlineContent::from(&details)).unwrap();
611        let fri = CompactFragReuseIndex::try_new(Uuid::new_v4(), details).unwrap();
612
613        // Surviving rows follow both versions in oldest-to-newest order.
614        assert_eq!(fri.remap_row_id(addr(1, 0)), Some(addr(20, 0)));
615        assert_eq!(fri.remap_row_id(addr(1, 3)), Some(addr(20, 1)));
616        // Deletes can happen in either the first or a later version.
617        assert_eq!(fri.remap_row_id(addr(1, 1)), None);
618        assert_eq!(fri.remap_row_id(addr(1, 2)), None);
619        assert_eq!(fri.remap_row_id(addr(3, 0)), None);
620        // Uncovered fragments and out-of-range offsets retain the existing
621        // missing-map pass-through semantics.
622        assert_eq!(fri.remap_row_id(addr(2, 0)), Some(addr(2, 0)));
623        assert_eq!(fri.remap_row_id(addr(1, 4)), Some(addr(1, 4)));
624
625        let mut batch = vec![
626            Some(addr(1, 0)),
627            Some(addr(1, 1)),
628            Some(addr(1, 2)),
629            Some(addr(1, 3)),
630            Some(addr(2, 0)),
631            None,
632        ];
633        fri.remap_row_ids_in_place(&mut batch);
634        assert_eq!(
635            batch,
636            vec![
637                Some(addr(20, 0)),
638                None,
639                None,
640                Some(addr(20, 1)),
641                Some(addr(2, 0)),
642                None,
643            ]
644        );
645    }
646
647    #[test]
648    fn test_compact_fri_rejects_invalid_changed_row_bitmap() {
649        let details = FragReuseIndexDetails {
650            versions: vec![FragReuseVersion {
651                dataset_version: 1,
652                groups: vec![FragReuseGroup {
653                    changed_row_addrs: vec![1, 2, 3],
654                    old_frags: vec![digest(1, 1)],
655                    new_frags: vec![digest(2, 1)],
656                }],
657            }],
658        };
659        let error = CompactFragReuseIndex::try_new(Uuid::new_v4(), details).unwrap_err();
660        assert!(matches!(error, Error::Index { .. }));
661        assert!(
662            error
663                .to_string()
664                .contains("failed to deserialize changed row addresses for FRI version 0, group 0"),
665            "{error}"
666        );
667    }
668
669    #[rstest]
670    #[case::unknown_fragment(
671        vec![addr(2, 0)],
672        vec![digest(1, 1)],
673        "from fragments [2] not in its old fragments"
674    )]
675    #[case::offset_out_of_range(
676        vec![addr(1, 1)],
677        vec![digest(1, 1)],
678        "row offset outside old fragment 1 with physical_rows=1"
679    )]
680    #[case::duplicate_old_fragment(
681        vec![addr(1, 0)],
682        vec![digest(1, 1), digest(1, 1)],
683        "old fragment 1 more than once"
684    )]
685    fn test_compact_fri_preserves_layout_validation(
686        #[case] changed_addrs: Vec<u64>,
687        #[case] old_frags: Vec<FragDigest>,
688        #[case] expected_message: &str,
689    ) {
690        let details = FragReuseIndexDetails {
691            versions: vec![FragReuseVersion {
692                dataset_version: 1,
693                groups: vec![FragReuseGroup {
694                    changed_row_addrs: serialize_changed(changed_addrs),
695                    old_frags,
696                    new_frags: vec![digest(10, 1)],
697                }],
698            }],
699        };
700
701        let error = CompactFragReuseIndex::try_new(Uuid::new_v4(), details).unwrap_err();
702        assert!(matches!(error, Error::Index { .. }));
703        let message = error.to_string();
704        assert!(message.contains("FRI version 0"), "{message}");
705        assert!(message.contains("rewrite group 0"), "{message}");
706        assert!(message.contains(expected_message), "{message}");
707    }
708
709    #[tokio::test]
710    async fn test_serialize_deserialize_index_details() {
711        // Create sample FragReuseVersions with different dataset versions
712        let version1 = FragReuseVersion {
713            dataset_version: 2,
714            groups: vec![FragReuseGroup {
715                changed_row_addrs: vec![1, 2, 3],
716                old_frags: vec![FragDigest {
717                    id: 1,
718                    physical_rows: 1,
719                    num_deleted_rows: 0,
720                }],
721                new_frags: vec![
722                    FragDigest {
723                        id: 2,
724                        physical_rows: 1,
725                        num_deleted_rows: 0,
726                    },
727                    FragDigest {
728                        id: 3,
729                        physical_rows: 1,
730                        num_deleted_rows: 0,
731                    },
732                ],
733            }],
734        };
735
736        let version2 = FragReuseVersion {
737            dataset_version: 1,
738            groups: vec![FragReuseGroup {
739                changed_row_addrs: vec![4, 5, 6],
740                old_frags: vec![FragDigest {
741                    id: 2,
742                    physical_rows: 1,
743                    num_deleted_rows: 0,
744                }],
745                new_frags: vec![
746                    FragDigest {
747                        id: 4,
748                        physical_rows: 1,
749                        num_deleted_rows: 0,
750                    },
751                    FragDigest {
752                        id: 5,
753                        physical_rows: 1,
754                        num_deleted_rows: 0,
755                    },
756                ],
757            }],
758        };
759
760        // Create FragReuseIndexDetails with versions in reverse order
761        let details = FragReuseIndexDetails {
762            versions: vec![version1, version2],
763        };
764
765        // Convert to protobuf format
766        let inline_content: InlineContent = (&details).into();
767
768        // Convert back to FragReuseIndexDetails
769        let roundtrip_details = FragReuseIndexDetails::try_from(inline_content).unwrap();
770
771        // Verify the roundtrip
772        assert_eq!(roundtrip_details.versions.len(), 2);
773
774        // Verify versions are sorted by dataset_version (oldest to latest)
775        assert_eq!(roundtrip_details.versions[0].dataset_version, 1);
776        assert_eq!(
777            roundtrip_details.versions[0].groups[0].changed_row_addrs,
778            vec![4, 5, 6]
779        );
780        assert_eq!(
781            roundtrip_details.versions[0].groups[0].new_frags,
782            vec![
783                FragDigest {
784                    id: 4,
785                    physical_rows: 1,
786                    num_deleted_rows: 0,
787                },
788                FragDigest {
789                    id: 5,
790                    physical_rows: 1,
791                    num_deleted_rows: 0,
792                }
793            ]
794        );
795        assert_eq!(
796            roundtrip_details.versions[0].groups[0].old_frags,
797            vec![FragDigest {
798                id: 2,
799                physical_rows: 1,
800                num_deleted_rows: 0,
801            }]
802        );
803
804        assert_eq!(roundtrip_details.versions[1].dataset_version, 2);
805        assert_eq!(
806            roundtrip_details.versions[1].groups[0].changed_row_addrs,
807            vec![1, 2, 3]
808        );
809        assert_eq!(
810            roundtrip_details.versions[1].groups[0].new_frags,
811            vec![
812                FragDigest {
813                    id: 2,
814                    physical_rows: 1,
815                    num_deleted_rows: 0,
816                },
817                FragDigest {
818                    id: 3,
819                    physical_rows: 1,
820                    num_deleted_rows: 0,
821                }
822            ]
823        );
824        assert_eq!(
825            roundtrip_details.versions[1].groups[0].old_frags,
826            vec![FragDigest {
827                id: 1,
828                physical_rows: 1,
829                num_deleted_rows: 0,
830            }]
831        );
832    }
833}