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