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