Skip to main content

lance_table/format/
fragment.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4use std::collections::HashMap;
5use std::num::NonZero;
6use std::sync::Arc;
7
8use lance_core::Error;
9use lance_core::deepsize::DeepSizeOf;
10use lance_file::format::{MAJOR_VERSION, MINOR_VERSION};
11use lance_file::version::LanceFileVersion;
12use lance_io::utils::CachedFileSize;
13use object_store::path::Path;
14use serde::{Deserialize, Deserializer, Serialize, Serializer};
15
16use super::overlay::{DataOverlayFile, sort_overlays_newest_last};
17use crate::format::pb;
18
19use crate::rowids::version::{
20    RowDatasetVersionMeta, created_at_version_meta_to_pb, last_updated_at_version_meta_to_pb,
21};
22use lance_core::datatypes::Schema;
23use lance_core::error::Result;
24
25/// Lance Data File
26///
27/// A data file is one piece of file storing data.
28#[derive(Debug, Clone, PartialEq, Eq, DeepSizeOf)]
29pub struct DataFile {
30    /// Relative path of the data file to dataset root.
31    pub path: String,
32    /// The ids of fields in this file.
33    ///
34    /// When identical across many fragments (common case), multiple `DataFile`
35    /// instances share a single heap allocation via `Arc`, significantly
36    /// reducing manifest memory for large tables.
37    pub fields: Arc<[i32]>,
38    /// The offsets of the fields listed in `fields`, empty in v1 files
39    ///
40    /// Note that -1 is a possibility and it indices that the field has
41    /// no top-level column in the file.
42    ///
43    /// Columns that lack a field id may still exist as extra entries in
44    /// `column_indices`; such columns are ignored by field-id–based projection.
45    /// For example, some fields, such as blob fields, occupy multiple
46    /// columns in the file but only have a single field id.
47    pub column_indices: Arc<[i32]>,
48    /// The major version of the file format used to write this file.
49    pub file_major_version: u32,
50    /// The minor version of the file format used to write this file.
51    pub file_minor_version: u32,
52
53    /// The size of the file in bytes, if known.
54    pub file_size_bytes: CachedFileSize,
55
56    /// The base path of the datafile, when the datafile is outside the dataset.
57    pub base_id: Option<u32>,
58}
59
60// Custom Serialize: convert Arc<[i32]> to slice for transparent JSON output
61impl Serialize for DataFile {
62    fn serialize<S: Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
63        use serde::ser::SerializeStruct;
64        let mut s = serializer.serialize_struct("DataFile", 7)?;
65        s.serialize_field("path", &self.path)?;
66        s.serialize_field("fields", self.fields.as_ref())?;
67        s.serialize_field("column_indices", self.column_indices.as_ref())?;
68        s.serialize_field("file_major_version", &self.file_major_version)?;
69        s.serialize_field("file_minor_version", &self.file_minor_version)?;
70        s.serialize_field("file_size_bytes", &self.file_size_bytes)?;
71        s.serialize_field("base_id", &self.base_id)?;
72        s.end()
73    }
74}
75
76// Custom Deserialize: read Vec<i32> and convert to Arc<[i32]>
77impl<'de> Deserialize<'de> for DataFile {
78    fn deserialize<D: Deserializer<'de>>(deserializer: D) -> std::result::Result<Self, D::Error> {
79        #[derive(Deserialize)]
80        struct DataFileHelper {
81            path: String,
82            fields: Vec<i32>,
83            #[serde(default)]
84            column_indices: Vec<i32>,
85            #[serde(default)]
86            file_major_version: u32,
87            #[serde(default)]
88            file_minor_version: u32,
89            file_size_bytes: CachedFileSize,
90            base_id: Option<u32>,
91        }
92
93        let helper = DataFileHelper::deserialize(deserializer)?;
94        Ok(Self {
95            path: helper.path,
96            fields: Arc::from(helper.fields),
97            column_indices: Arc::from(helper.column_indices),
98            file_major_version: helper.file_major_version,
99            file_minor_version: helper.file_minor_version,
100            file_size_bytes: helper.file_size_bytes,
101            base_id: helper.base_id,
102        })
103    }
104}
105
106impl DataFile {
107    pub fn new(
108        path: impl Into<String>,
109        fields: Vec<i32>,
110        column_indices: Vec<i32>,
111        file_major_version: u32,
112        file_minor_version: u32,
113        file_size_bytes: Option<NonZero<u64>>,
114        base_id: Option<u32>,
115    ) -> Self {
116        Self {
117            path: path.into(),
118            fields: Arc::from(fields),
119            column_indices: Arc::from(column_indices),
120            file_major_version,
121            file_minor_version,
122            file_size_bytes: file_size_bytes.into(),
123            base_id,
124        }
125    }
126
127    /// Create a new `DataFile` with the expectation that fields and column_indices will be set later
128    pub fn new_unstarted(
129        path: impl Into<String>,
130        file_major_version: u32,
131        file_minor_version: u32,
132    ) -> Self {
133        Self {
134            path: path.into(),
135            fields: Arc::from([]),
136            column_indices: Arc::from([]),
137            file_major_version,
138            file_minor_version,
139            file_size_bytes: Default::default(),
140            base_id: None,
141        }
142    }
143
144    pub fn new_legacy_from_fields(
145        path: impl Into<String>,
146        fields: Vec<i32>,
147        base_id: Option<u32>,
148    ) -> Self {
149        Self::new(
150            path,
151            fields,
152            vec![],
153            MAJOR_VERSION as u32,
154            MINOR_VERSION as u32,
155            None,
156            base_id,
157        )
158    }
159
160    pub fn new_legacy(
161        path: impl Into<String>,
162        schema: &Schema,
163        file_size_bytes: Option<NonZero<u64>>,
164        base_id: Option<u32>,
165    ) -> Self {
166        let mut field_ids = schema.field_ids();
167        field_ids.sort();
168        Self::new(
169            path,
170            field_ids,
171            vec![],
172            MAJOR_VERSION as u32,
173            MINOR_VERSION as u32,
174            file_size_bytes,
175            base_id,
176        )
177    }
178
179    pub fn schema(&self, full_schema: &Schema) -> Schema {
180        full_schema.project_by_ids(&self.fields, false)
181    }
182
183    pub fn is_legacy_file(&self) -> bool {
184        self.file_major_version == 0 && self.file_minor_version < 3
185    }
186
187    pub fn validate(&self, base_path: &Path) -> Result<()> {
188        if self.is_legacy_file() {
189            if !self.fields.windows(2).all(|w| w[0] < w[1]) {
190                return Err(Error::corrupt_file(
191                    base_path.clone().join(self.path.clone()),
192                    "contained unsorted or duplicate field ids",
193                ));
194            }
195        } else if self.column_indices.len() < self.fields.len() {
196            // Every recorded field id must have a column index, but not every column needs
197            // to be associated with a field id (extra columns are allowed).
198            return Err(Error::corrupt_file(
199                base_path.clone().join(self.path.clone()),
200                "contained fewer column_indices than fields",
201            ));
202        }
203        Ok(())
204    }
205}
206
207impl From<&DataFile> for pb::DataFile {
208    fn from(df: &DataFile) -> Self {
209        Self {
210            path: df.path.clone(),
211            fields: df.fields.to_vec(),
212            column_indices: df.column_indices.to_vec(),
213            file_major_version: df.file_major_version,
214            file_minor_version: df.file_minor_version,
215            file_size_bytes: df.file_size_bytes.get().map_or(0, |v| v.get()),
216            base_id: df.base_id,
217        }
218    }
219}
220
221impl TryFrom<pb::DataFile> for DataFile {
222    type Error = Error;
223
224    fn try_from(proto: pb::DataFile) -> Result<Self> {
225        Ok(Self {
226            path: proto.path,
227            fields: Arc::from(proto.fields),
228            column_indices: Arc::from(proto.column_indices),
229            file_major_version: proto.file_major_version,
230            file_minor_version: proto.file_minor_version,
231            file_size_bytes: CachedFileSize::new(proto.file_size_bytes),
232            base_id: proto.base_id,
233        })
234    }
235}
236
237/// Interns repeated data so that fragments with identical content share a
238/// single heap allocation via `Arc`.
239///
240/// At 20M fragments the deduplication typically saves multiple GB of heap
241/// because every fragment in a homogeneous table carries the same field list,
242/// and post-compaction fragments share identical version metadata bytes.
243///
244/// Uses a `Vec`-based linear scan when the cache is small (<=16 entries)
245/// and upgrades to `HashMap` for larger caches. In the common homogeneous
246/// case (1-3 unique values), linear scan avoids per-fragment hashing overhead.
247#[derive(Default)]
248pub struct DataFileFieldInterner {
249    fields: InternCache<i32>,
250    column_indices: InternCache<i32>,
251    inline_bytes: InternCache<u8>,
252}
253
254/// A cache that uses linear scan for small sizes and HashMap for large.
255/// The threshold is chosen so that scan + compare is cheaper than hash for
256/// typical payload sizes (20-200 bytes).
257enum InternCache<T: Eq + std::hash::Hash + Clone> {
258    Small(Vec<Arc<[T]>>),
259    Large(HashMap<Arc<[T]>, ()>),
260}
261
262const INTERN_CACHE_UPGRADE_THRESHOLD: usize = 16;
263
264impl<T: Eq + std::hash::Hash + Clone> Default for InternCache<T> {
265    fn default() -> Self {
266        Self::Small(Vec::new())
267    }
268}
269
270impl<T: Eq + std::hash::Hash + Clone> InternCache<T> {
271    fn intern(&mut self, v: Vec<T>) -> Arc<[T]> {
272        match self {
273            Self::Small(entries) => {
274                for existing in entries.iter() {
275                    if existing.as_ref() == v.as_slice() {
276                        return existing.clone();
277                    }
278                }
279                let arc: Arc<[T]> = Arc::from(v);
280                entries.push(arc.clone());
281                if entries.len() > INTERN_CACHE_UPGRADE_THRESHOLD {
282                    let mut map = HashMap::with_capacity(entries.len());
283                    for e in entries.drain(..) {
284                        map.insert(e, ());
285                    }
286                    *self = Self::Large(map);
287                }
288                arc
289            }
290            Self::Large(map) => {
291                if let Some((existing, _)) = map.get_key_value(v.as_slice()) {
292                    existing.clone()
293                } else {
294                    let arc: Arc<[T]> = Arc::from(v);
295                    map.insert(arc.clone(), ());
296                    arc
297                }
298            }
299        }
300    }
301}
302
303impl DataFileFieldInterner {
304    /// Intern a `RowDatasetVersionMeta`, deduplicating inline byte payloads.
305    /// Accepts the protobuf oneof value directly to avoid an intermediate
306    /// `Arc<[u8]>` allocation that would need to be `.to_vec()`'d for the key lookup.
307    fn intern_last_updated_version_meta(
308        cache: &mut InternCache<u8>,
309        pb: pb::data_fragment::LastUpdatedAtVersionSequence,
310    ) -> Result<RowDatasetVersionMeta> {
311        match pb {
312            pb::data_fragment::LastUpdatedAtVersionSequence::InlineLastUpdatedAtVersions(data) => {
313                Ok(RowDatasetVersionMeta::Inline(cache.intern(data)))
314            }
315            pb::data_fragment::LastUpdatedAtVersionSequence::ExternalLastUpdatedAtVersions(
316                file,
317            ) => Ok(RowDatasetVersionMeta::External(ExternalFile {
318                path: file.path,
319                offset: file.offset,
320                size: file.size,
321            })),
322        }
323    }
324
325    /// Intern a `RowDatasetVersionMeta`, deduplicating inline byte payloads.
326    fn intern_created_version_meta(
327        cache: &mut InternCache<u8>,
328        pb: pb::data_fragment::CreatedAtVersionSequence,
329    ) -> Result<RowDatasetVersionMeta> {
330        match pb {
331            pb::data_fragment::CreatedAtVersionSequence::InlineCreatedAtVersions(data) => {
332                Ok(RowDatasetVersionMeta::Inline(cache.intern(data)))
333            }
334            pb::data_fragment::CreatedAtVersionSequence::ExternalCreatedAtVersions(file) => {
335                Ok(RowDatasetVersionMeta::External(ExternalFile {
336                    path: file.path,
337                    offset: file.offset,
338                    size: file.size,
339                }))
340            }
341        }
342    }
343
344    /// Convert a protobuf `DataFile`, interning `fields` and `column_indices`.
345    pub fn intern_data_file(&mut self, proto: pb::DataFile) -> Result<DataFile> {
346        Ok(DataFile {
347            path: proto.path,
348            fields: self.fields.intern(proto.fields),
349            column_indices: self.column_indices.intern(proto.column_indices),
350            file_major_version: proto.file_major_version,
351            file_minor_version: proto.file_minor_version,
352            file_size_bytes: CachedFileSize::new(proto.file_size_bytes),
353            base_id: proto.base_id,
354        })
355    }
356
357    /// Convert a protobuf `DataFragment`, interning fields and version metadata.
358    pub fn intern_fragment(&mut self, p: pb::DataFragment) -> Result<Fragment> {
359        let physical_rows = if p.physical_rows > 0 {
360            Some(p.physical_rows as usize)
361        } else {
362            None
363        };
364        let last_updated_at_version_meta = p
365            .last_updated_at_version_sequence
366            .map(|pb| Self::intern_last_updated_version_meta(&mut self.inline_bytes, pb))
367            .transpose()?;
368        let created_at_version_meta = p
369            .created_at_version_sequence
370            .map(|pb| Self::intern_created_version_meta(&mut self.inline_bytes, pb))
371            .transpose()?;
372        Ok(Fragment {
373            id: p.id,
374            files: p
375                .files
376                .into_iter()
377                .map(|f| self.intern_data_file(f))
378                .collect::<Result<_>>()?,
379            overlays: {
380                let mut overlays = p
381                    .overlays
382                    .into_iter()
383                    .map(DataOverlayFile::try_from)
384                    .collect::<Result<Vec<_>>>()?;
385                sort_overlays_newest_last(&mut overlays);
386                overlays
387            },
388            deletion_file: p.deletion_file.map(DeletionFile::try_from).transpose()?,
389            row_id_meta: p.row_id_sequence.map(RowIdMeta::try_from).transpose()?,
390            physical_rows,
391            last_updated_at_version_meta,
392            created_at_version_meta,
393        })
394    }
395}
396
397#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
398#[serde(rename_all = "lowercase")]
399pub enum DeletionFileType {
400    Array,
401    Bitmap,
402}
403
404impl DeletionFileType {
405    // TODO: pub(crate)
406    pub fn suffix(&self) -> &str {
407        match self {
408            Self::Array => "arrow",
409            Self::Bitmap => "bin",
410        }
411    }
412}
413
414#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
415pub struct DeletionFile {
416    pub read_version: u64,
417    pub id: u64,
418    pub file_type: DeletionFileType,
419    /// Number of deleted rows in this file. If None, this is unknown.
420    pub num_deleted_rows: Option<usize>,
421    pub base_id: Option<u32>,
422}
423
424impl TryFrom<pb::DeletionFile> for DeletionFile {
425    type Error = Error;
426
427    fn try_from(value: pb::DeletionFile) -> Result<Self> {
428        let file_type = match value.file_type {
429            0 => DeletionFileType::Array,
430            1 => DeletionFileType::Bitmap,
431            _ => {
432                return Err(Error::not_supported_source(
433                    "Unknown deletion file type".into(),
434                ));
435            }
436        };
437        let num_deleted_rows = if value.num_deleted_rows == 0 {
438            None
439        } else {
440            Some(value.num_deleted_rows as usize)
441        };
442        Ok(Self {
443            read_version: value.read_version,
444            id: value.id,
445            file_type,
446            num_deleted_rows,
447            base_id: value.base_id,
448        })
449    }
450}
451
452/// A reference to a part of a file.
453#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
454pub struct ExternalFile {
455    pub path: String,
456    pub offset: u64,
457    pub size: u64,
458}
459
460/// Metadata about location of the row id sequence.
461#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
462pub enum RowIdMeta {
463    Inline(Vec<u8>),
464    External(ExternalFile),
465}
466
467impl TryFrom<pb::data_fragment::RowIdSequence> for RowIdMeta {
468    type Error = Error;
469
470    fn try_from(value: pb::data_fragment::RowIdSequence) -> Result<Self> {
471        match value {
472            pb::data_fragment::RowIdSequence::InlineRowIds(data) => Ok(Self::Inline(data)),
473            pb::data_fragment::RowIdSequence::ExternalRowIds(file) => {
474                Ok(Self::External(ExternalFile {
475                    path: file.path.clone(),
476                    offset: file.offset,
477                    size: file.size,
478                }))
479            }
480        }
481    }
482}
483
484/// Data fragment.
485///
486/// A fragment is a set of files which represent the different columns of the same rows.
487/// If column exists in the schema, but the related file does not exist, treat this column as `nulls`.
488#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
489pub struct Fragment {
490    /// Fragment ID
491    pub id: u64,
492
493    /// Files within the fragment.
494    pub files: Vec<DataFile>,
495
496    /// Overlay files supplying new values for a subset of cells without
497    /// rewriting the base data files. Order is significant: a later entry is
498    /// newer than an earlier one. See [`DataOverlayFile`] for resolution rules.
499    #[serde(default, skip_serializing_if = "Vec::is_empty")]
500    pub overlays: Vec<DataOverlayFile>,
501
502    /// Optional file with deleted local row offsets.
503    #[serde(skip_serializing_if = "Option::is_none")]
504    pub deletion_file: Option<DeletionFile>,
505
506    /// RowIndex
507    #[serde(skip_serializing_if = "Option::is_none")]
508    pub row_id_meta: Option<RowIdMeta>,
509
510    /// Original number of rows in the fragment. If this is None, then it is
511    /// unknown. This is only optional for legacy reasons. All new tables should
512    /// have this set.
513    pub physical_rows: Option<usize>,
514
515    /// Last updated at version metadata
516    #[serde(skip_serializing_if = "Option::is_none")]
517    pub last_updated_at_version_meta: Option<RowDatasetVersionMeta>,
518
519    /// Created at version metadata
520    #[serde(skip_serializing_if = "Option::is_none")]
521    pub created_at_version_meta: Option<RowDatasetVersionMeta>,
522}
523
524impl Fragment {
525    pub fn new(id: u64) -> Self {
526        Self {
527            id,
528            files: vec![],
529            overlays: vec![],
530            deletion_file: None,
531            row_id_meta: None,
532            physical_rows: None,
533            last_updated_at_version_meta: None,
534            created_at_version_meta: None,
535        }
536    }
537
538    pub fn num_rows(&self) -> Option<usize> {
539        match (self.physical_rows, &self.deletion_file) {
540            // Known fragment length, no deletion file.
541            (Some(len), None) => Some(len),
542            // Known fragment length, but don't know deletion file size.
543            (
544                Some(len),
545                Some(DeletionFile {
546                    num_deleted_rows: Some(num_deleted_rows),
547                    ..
548                }),
549            ) => Some(len - num_deleted_rows),
550            _ => None,
551        }
552    }
553
554    pub fn from_json(json: &str) -> Result<Self> {
555        let fragment: Self = serde_json::from_str(json)?;
556        Ok(fragment)
557    }
558
559    /// Create a `Fragment` with one DataFile
560    pub fn with_file_legacy(
561        id: u64,
562        path: &str,
563        schema: &Schema,
564        physical_rows: Option<usize>,
565    ) -> Self {
566        Self {
567            id,
568            files: vec![DataFile::new_legacy(path, schema, None, None)],
569            overlays: vec![],
570            deletion_file: None,
571            physical_rows,
572            row_id_meta: None,
573            last_updated_at_version_meta: None,
574            created_at_version_meta: None,
575        }
576    }
577
578    pub fn with_file(
579        mut self,
580        path: impl Into<String>,
581        field_ids: Vec<i32>,
582        column_indices: Vec<i32>,
583        version: &LanceFileVersion,
584        file_size_bytes: Option<NonZero<u64>>,
585    ) -> Self {
586        let (major, minor) = version.to_numbers();
587        let data_file = DataFile::new(
588            path,
589            field_ids,
590            column_indices,
591            major,
592            minor,
593            file_size_bytes,
594            None,
595        );
596        self.files.push(data_file);
597        self
598    }
599
600    pub fn with_physical_rows(mut self, physical_rows: usize) -> Self {
601        self.physical_rows = Some(physical_rows);
602        self
603    }
604
605    pub fn add_file(
606        &mut self,
607        path: impl Into<String>,
608        field_ids: Vec<i32>,
609        column_indices: Vec<i32>,
610        version: &LanceFileVersion,
611        file_size_bytes: Option<NonZero<u64>>,
612    ) {
613        let (major, minor) = version.to_numbers();
614        self.files.push(DataFile::new(
615            path,
616            field_ids,
617            column_indices,
618            major,
619            minor,
620            file_size_bytes,
621            None,
622        ));
623    }
624
625    /// Add a new [`DataFile`] to this fragment.
626    pub fn add_file_legacy(&mut self, path: &str, schema: &Schema) {
627        self.files
628            .push(DataFile::new_legacy(path, schema, None, None));
629    }
630
631    // True if this fragment is made up of legacy v1 files, false otherwise
632    pub fn has_legacy_files(&self) -> bool {
633        // If any file in a fragment is legacy then all files in the fragment must be
634        self.files[0].is_legacy_file()
635    }
636
637    // Helper method to infer the Lance version from a set of fragments
638    //
639    // Returns None if there are no data files
640    // Returns an error if the data files have different versions
641    pub fn try_infer_version(fragments: &[Self]) -> Result<Option<LanceFileVersion>> {
642        // Otherwise we need to check the actual file versions
643        // Determine version from first file
644        let Some(sample_file) = fragments
645            .iter()
646            .find(|f| !f.files.is_empty())
647            .map(|f| &f.files[0])
648        else {
649            return Ok(None);
650        };
651        let file_version = LanceFileVersion::try_from_major_minor(
652            sample_file.file_major_version,
653            sample_file.file_minor_version,
654        )?;
655        // Ensure all files match
656        for frag in fragments {
657            for file in &frag.files {
658                let this_file_version = LanceFileVersion::try_from_major_minor(
659                    file.file_major_version,
660                    file.file_minor_version,
661                )?;
662                if file_version != this_file_version {
663                    return Err(Error::invalid_input(format!(
664                        "All data files must have the same version.  Detected both {} and {}",
665                        file_version, this_file_version
666                    )));
667                }
668            }
669        }
670        Ok(Some(file_version))
671    }
672}
673
674impl TryFrom<pb::DataFragment> for Fragment {
675    type Error = Error;
676
677    fn try_from(p: pb::DataFragment) -> Result<Self> {
678        let physical_rows = if p.physical_rows > 0 {
679            Some(p.physical_rows as usize)
680        } else {
681            None
682        };
683        Ok(Self {
684            id: p.id,
685            files: p
686                .files
687                .into_iter()
688                .map(DataFile::try_from)
689                .collect::<Result<_>>()?,
690            overlays: {
691                let mut overlays = p
692                    .overlays
693                    .into_iter()
694                    .map(DataOverlayFile::try_from)
695                    .collect::<Result<Vec<_>>>()?;
696                sort_overlays_newest_last(&mut overlays);
697                overlays
698            },
699            deletion_file: p.deletion_file.map(DeletionFile::try_from).transpose()?,
700            row_id_meta: p.row_id_sequence.map(RowIdMeta::try_from).transpose()?,
701            physical_rows,
702            last_updated_at_version_meta: p
703                .last_updated_at_version_sequence
704                .map(RowDatasetVersionMeta::try_from)
705                .transpose()?,
706            created_at_version_meta: p
707                .created_at_version_sequence
708                .map(RowDatasetVersionMeta::try_from)
709                .transpose()?,
710        })
711    }
712}
713
714impl From<&Fragment> for pb::DataFragment {
715    fn from(f: &Fragment) -> Self {
716        let deletion_file = f.deletion_file.as_ref().map(|f| {
717            let file_type = match f.file_type {
718                DeletionFileType::Array => pb::deletion_file::DeletionFileType::ArrowArray,
719                DeletionFileType::Bitmap => pb::deletion_file::DeletionFileType::Bitmap,
720            };
721            pb::DeletionFile {
722                read_version: f.read_version,
723                id: f.id,
724                file_type: file_type.into(),
725                num_deleted_rows: f.num_deleted_rows.unwrap_or_default() as u64,
726                base_id: f.base_id,
727            }
728        });
729
730        let row_id_sequence = f.row_id_meta.as_ref().map(|m| match m {
731            RowIdMeta::Inline(data) => pb::data_fragment::RowIdSequence::InlineRowIds(data.clone()),
732            RowIdMeta::External(file) => {
733                pb::data_fragment::RowIdSequence::ExternalRowIds(pb::ExternalFile {
734                    path: file.path.clone(),
735                    offset: file.offset,
736                    size: file.size,
737                })
738            }
739        });
740        let last_updated_at_version_sequence =
741            last_updated_at_version_meta_to_pb(&f.last_updated_at_version_meta);
742        let created_at_version_sequence = created_at_version_meta_to_pb(&f.created_at_version_meta);
743        Self {
744            id: f.id,
745            files: f.files.iter().map(pb::DataFile::from).collect(),
746            overlays: f.overlays.iter().map(pb::DataOverlayFile::from).collect(),
747            deletion_file,
748            row_id_sequence,
749            physical_rows: f.physical_rows.unwrap_or_default() as u64,
750            last_updated_at_version_sequence,
751            created_at_version_sequence,
752        }
753    }
754}
755
756#[cfg(test)]
757mod tests {
758    use super::*;
759    use crate::format::overlay::OverlayCoverage;
760    use arrow_schema::{
761        DataType, Field as ArrowField, Fields as ArrowFields, Schema as ArrowSchema,
762    };
763    use object_store::path::Path;
764    use roaring::RoaringBitmap;
765    use serde_json::{Value, json};
766
767    #[test]
768    fn test_data_overlay_roundtrip() {
769        // A fragment carrying a dense overlay round-trips through protobuf and
770        // back, and the parsed coverage bitmap is recovered per field.
771        let mut bitmap = RoaringBitmap::new();
772        bitmap.insert(1);
773        bitmap.insert(3);
774
775        let overlay = DataOverlayFile {
776            data_file: DataFile::new_legacy_from_fields("overlay-0.lance", vec![3], None),
777            coverage: OverlayCoverage::dense(bitmap.clone()),
778            committed_version: 7,
779        };
780        let mut fragment = Fragment::new(0);
781        fragment.files = vec![DataFile::new_legacy_from_fields(
782            "base.lance",
783            vec![1, 3],
784            None,
785        )];
786        fragment.overlays = vec![overlay];
787
788        let proto = pb::DataFragment::from(&fragment);
789        assert_eq!(proto.overlays.len(), 1);
790        let round_tripped = Fragment::try_from(proto).unwrap();
791        assert_eq!(round_tripped, fragment);
792
793        // Dense coverage applies to every field.
794        let recovered = round_tripped.overlays[0].coverage_for_field(0).unwrap();
795        assert_eq!(*recovered, bitmap);
796        assert_eq!(
797            *round_tripped.overlays[0].coverage_for_field(5).unwrap(),
798            bitmap
799        );
800    }
801
802    #[test]
803    fn test_data_overlay_sparse_per_field_coverage() {
804        // A sparse overlay carries one bitmap per field, recovered by position.
805        let name_coverage = RoaringBitmap::from_iter([2u32, 3]);
806        let embedding_coverage = RoaringBitmap::from_iter([1u32]);
807        let overlay = DataOverlayFile {
808            data_file: DataFile::new_legacy_from_fields("overlay-1.lance", vec![2, 4], None),
809            coverage: OverlayCoverage::sparse(vec![
810                name_coverage.clone(),
811                embedding_coverage.clone(),
812            ]),
813            committed_version: 3,
814        };
815        let mut fragment = Fragment::new(1);
816        fragment.overlays = vec![overlay];
817
818        let round_tripped = Fragment::try_from(pb::DataFragment::from(&fragment)).unwrap();
819        assert_eq!(
820            *round_tripped.overlays[0].coverage_for_field(0).unwrap(),
821            name_coverage
822        );
823        assert_eq!(
824            *round_tripped.overlays[0].coverage_for_field(1).unwrap(),
825            embedding_coverage
826        );
827    }
828
829    #[test]
830    fn test_overlays_sorted_newest_last_on_load() {
831        // Overlays load stable-sorted by committed_version (newest last), with
832        // list position preserved as the tiebreak for equal versions.
833        let mk = |version: u64, field: i32| DataOverlayFile {
834            data_file: DataFile::new_legacy_from_fields("o.lance", vec![field], None),
835            coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0u32])),
836            committed_version: version,
837        };
838        let mut fragment = Fragment::new(0);
839        // Written out of order: v5, v2, v2 (second), v3.
840        fragment.overlays = vec![mk(5, 1), mk(2, 2), mk(2, 3), mk(3, 4)];
841
842        let loaded = Fragment::try_from(pb::DataFragment::from(&fragment)).unwrap();
843        let versions: Vec<u64> = loaded
844            .overlays
845            .iter()
846            .map(|o| o.committed_version)
847            .collect();
848        assert_eq!(versions, vec![2, 2, 3, 5]);
849        // Stable: the two v2 overlays keep their original relative order (field 2
850        // before field 3).
851        assert_eq!(
852            loaded.overlays[0].data_file.fields.as_ref(),
853            [2i32].as_slice()
854        );
855        assert_eq!(
856            loaded.overlays[1].data_file.fields.as_ref(),
857            [3i32].as_slice()
858        );
859    }
860
861    #[test]
862    fn test_new_fragment() {
863        let path = "foobar.lance";
864
865        let arrow_schema = ArrowSchema::new(vec![
866            ArrowField::new(
867                "s",
868                DataType::Struct(ArrowFields::from(vec![
869                    ArrowField::new("si", DataType::Int32, false),
870                    ArrowField::new("sb", DataType::Binary, true),
871                ])),
872                true,
873            ),
874            ArrowField::new("bool", DataType::Boolean, true),
875        ]);
876        let schema = Schema::try_from(&arrow_schema).unwrap();
877        let fragment = Fragment::with_file_legacy(123, path, &schema, Some(10));
878
879        assert_eq!(123, fragment.id);
880        assert_eq!(
881            fragment.files,
882            vec![DataFile::new_legacy_from_fields(
883                path.to_string(),
884                vec![0, 1, 2, 3],
885                None,
886            )]
887        )
888    }
889
890    #[test]
891    fn test_roundtrip_fragment() {
892        let mut fragment = Fragment::new(123);
893        let schema = ArrowSchema::new(vec![ArrowField::new("x", DataType::Float16, true)]);
894        fragment.add_file_legacy("foobar.lance", &Schema::try_from(&schema).unwrap());
895        fragment.deletion_file = Some(DeletionFile {
896            read_version: 123,
897            id: 456,
898            file_type: DeletionFileType::Array,
899            num_deleted_rows: Some(10),
900            base_id: None,
901        });
902
903        let proto = pb::DataFragment::from(&fragment);
904        let fragment2 = Fragment::try_from(proto).unwrap();
905        assert_eq!(fragment, fragment2);
906
907        fragment.deletion_file = None;
908        let proto = pb::DataFragment::from(&fragment);
909        let fragment2 = Fragment::try_from(proto).unwrap();
910        assert_eq!(fragment, fragment2);
911    }
912
913    #[test]
914    fn test_to_json() {
915        let mut fragment = Fragment::new(123);
916        let schema = ArrowSchema::new(vec![ArrowField::new("x", DataType::Float16, true)]);
917        fragment.add_file_legacy("foobar.lance", &Schema::try_from(&schema).unwrap());
918        fragment.deletion_file = Some(DeletionFile {
919            read_version: 123,
920            id: 456,
921            file_type: DeletionFileType::Array,
922            num_deleted_rows: Some(10),
923            base_id: None,
924        });
925
926        let json = serde_json::to_string(&fragment).unwrap();
927
928        let value: Value = serde_json::from_str(&json).unwrap();
929        assert_eq!(
930            value,
931            json!({
932                "id": 123,
933                "files":[
934                    {"path": "foobar.lance", "fields": [0], "column_indices": [], 
935                     "file_major_version": MAJOR_VERSION, "file_minor_version": MINOR_VERSION,
936                     "file_size_bytes": null, "base_id": null }
937                ],
938                "deletion_file": {"read_version": 123, "id": 456, "file_type": "array",
939                                  "num_deleted_rows": 10, "base_id": null},
940                "physical_rows": None::<usize>}),
941        );
942
943        let frag2 = Fragment::from_json(&json).unwrap();
944        assert_eq!(fragment, frag2);
945    }
946
947    #[test]
948    fn data_file_validate_allows_extra_columns() {
949        let data_file = DataFile {
950            path: "foo.lance".to_string(),
951            fields: Arc::from([1, 2]),
952            // One extra column without a field id mapping
953            column_indices: Arc::from([0, 1, 2]),
954            file_major_version: MAJOR_VERSION as u32,
955            file_minor_version: MINOR_VERSION as u32,
956            file_size_bytes: Default::default(),
957            base_id: None,
958        };
959
960        let base_path = Path::from("base");
961        data_file
962            .validate(&base_path)
963            .expect("validation should allow extra columns without field ids");
964    }
965}