1use 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#[derive(Debug, Clone, PartialEq, Eq, DeepSizeOf)]
29pub struct DataFile {
30 pub path: String,
32 pub fields: Arc<[i32]>,
38 pub column_indices: Arc<[i32]>,
48 pub file_major_version: u32,
50 pub file_minor_version: u32,
52
53 pub file_size_bytes: CachedFileSize,
55
56 pub base_id: Option<u32>,
58}
59
60impl 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
76impl<'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(
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 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 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 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 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#[derive(Default)]
254pub struct DataFileFieldInterner {
255 fields: InternCache<i32>,
256 column_indices: InternCache<i32>,
257 inline_bytes: InternCache<u8>,
258}
259
260enum 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 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 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 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 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 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 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
463pub struct Fragment {
464 pub id: u64,
466
467 pub files: Vec<DataFile>,
469
470 #[serde(default, skip_serializing_if = "Vec::is_empty")]
474 pub overlays: Vec<DataOverlayFile>,
475
476 #[serde(skip_serializing_if = "Option::is_none")]
478 pub deletion_file: Option<DeletionFile>,
479
480 #[serde(skip_serializing_if = "Option::is_none")]
482 pub row_id_meta: Option<RowIdMeta>,
483
484 pub physical_rows: Option<usize>,
488
489 #[serde(skip_serializing_if = "Option::is_none")]
491 pub last_updated_at_version_meta: Option<RowDatasetVersionMeta>,
492
493 #[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 (Some(len), None) => Some(len),
516 (
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 pub fn referenced_lance_files(&self) -> impl Iterator<Item = &DataFile> + '_ {
536 let Self {
539 id: _,
540 files,
541 overlays,
542 deletion_file: _,
543 row_id_meta: _,
544 physical_rows: _,
545 last_updated_at_version_meta: _,
546 created_at_version_meta: _,
547 } = self;
548 files
549 .iter()
550 .chain(overlays.iter().map(|overlay| &overlay.data_file))
551 }
552
553 pub fn referenced_lance_files_mut(&mut self) -> impl Iterator<Item = &mut DataFile> + '_ {
557 let Self {
560 id: _,
561 files,
562 overlays,
563 deletion_file: _,
564 row_id_meta: _,
565 physical_rows: _,
566 last_updated_at_version_meta: _,
567 created_at_version_meta: _,
568 } = self;
569 files
570 .iter_mut()
571 .chain(overlays.iter_mut().map(|overlay| &mut overlay.data_file))
572 }
573
574 pub fn from_json(json: &str) -> Result<Self> {
575 let fragment: Self = serde_json::from_str(json)?;
576 Ok(fragment)
577 }
578
579 pub fn with_file_legacy(
581 id: u64,
582 path: &str,
583 schema: &Schema,
584 physical_rows: Option<usize>,
585 ) -> Self {
586 Self {
587 id,
588 files: vec![DataFile::new_legacy(path, schema, None, None)],
589 overlays: vec![],
590 deletion_file: None,
591 physical_rows,
592 row_id_meta: None,
593 last_updated_at_version_meta: None,
594 created_at_version_meta: None,
595 }
596 }
597
598 pub fn with_file(
599 mut self,
600 path: impl Into<String>,
601 field_ids: Vec<i32>,
602 column_indices: Vec<i32>,
603 version: ConcreteFileVersion,
604 file_size_bytes: Option<NonZero<u64>>,
605 ) -> Self {
606 let data_file = DataFile::new(
607 path,
608 field_ids,
609 column_indices,
610 version,
611 file_size_bytes,
612 None,
613 );
614 self.files.push(data_file);
615 self
616 }
617
618 pub fn with_physical_rows(mut self, physical_rows: usize) -> Self {
619 self.physical_rows = Some(physical_rows);
620 self
621 }
622
623 pub fn add_file(
624 &mut self,
625 path: impl Into<String>,
626 field_ids: Vec<i32>,
627 column_indices: Vec<i32>,
628 version: ConcreteFileVersion,
629 file_size_bytes: Option<NonZero<u64>>,
630 ) {
631 self.files.push(DataFile::new(
632 path,
633 field_ids,
634 column_indices,
635 version,
636 file_size_bytes,
637 None,
638 ));
639 }
640
641 pub fn add_file_legacy(&mut self, path: &str, schema: &Schema) {
643 self.files
644 .push(DataFile::new_legacy(path, schema, None, None));
645 }
646
647 pub fn try_infer_version(fragments: &[Self]) -> Result<Option<ConcreteFileVersion>> {
652 let Some(sample_file) = fragments
655 .iter()
656 .flat_map(Self::referenced_lance_files)
657 .next()
658 else {
659 return Ok(None);
660 };
661 let file_version = sample_file.file_version()?;
662 for frag in fragments {
664 for file in frag.referenced_lance_files() {
665 let this_file_version = file.file_version()?;
666 if file_version != this_file_version {
667 return Err(Error::invalid_input(format!(
668 "All data files must have the same version. Detected both {} and {}",
669 file_version, this_file_version
670 )));
671 }
672 }
673 }
674 Ok(Some(file_version))
675 }
676}
677
678impl TryFrom<pb::DataFragment> for Fragment {
679 type Error = Error;
680
681 fn try_from(p: pb::DataFragment) -> Result<Self> {
682 let physical_rows = if p.physical_rows > 0 {
683 Some(p.physical_rows as usize)
684 } else {
685 None
686 };
687 Ok(Self {
688 id: p.id,
689 files: p
690 .files
691 .into_iter()
692 .map(DataFile::try_from)
693 .collect::<Result<_>>()?,
694 overlays: {
695 let mut overlays = p
696 .overlays
697 .into_iter()
698 .map(DataOverlayFile::try_from)
699 .collect::<Result<Vec<_>>>()?;
700 sort_overlays_newest_last(&mut overlays);
701 overlays
702 },
703 deletion_file: p.deletion_file.map(DeletionFile::try_from).transpose()?,
704 row_id_meta: p.row_id_sequence.map(RowIdMeta::try_from).transpose()?,
705 physical_rows,
706 last_updated_at_version_meta: p
707 .last_updated_at_version_sequence
708 .map(RowDatasetVersionMeta::try_from)
709 .transpose()?,
710 created_at_version_meta: p
711 .created_at_version_sequence
712 .map(RowDatasetVersionMeta::try_from)
713 .transpose()?,
714 })
715 }
716}
717
718impl From<&Fragment> for pb::DataFragment {
719 fn from(f: &Fragment) -> Self {
720 let deletion_file = f.deletion_file.as_ref().map(|f| {
721 let file_type = match f.file_type {
722 DeletionFileType::Array => pb::deletion_file::DeletionFileType::ArrowArray,
723 DeletionFileType::Bitmap => pb::deletion_file::DeletionFileType::Bitmap,
724 };
725 pb::DeletionFile {
726 read_version: f.read_version,
727 id: f.id,
728 file_type: file_type.into(),
729 num_deleted_rows: f.num_deleted_rows.unwrap_or_default() as u64,
730 base_id: f.base_id,
731 }
732 });
733
734 let row_id_sequence = f.row_id_meta.as_ref().map(|m| match m {
735 RowIdMeta::Inline(data) => {
736 pb::data_fragment::RowIdSequence::InlineRowIds(data.to_vec())
737 }
738 RowIdMeta::External(file) => {
739 pb::data_fragment::RowIdSequence::ExternalRowIds(pb::ExternalFile {
740 path: file.path.clone(),
741 offset: file.offset,
742 size: file.size,
743 })
744 }
745 });
746 let last_updated_at_version_sequence =
747 last_updated_at_version_meta_to_pb(&f.last_updated_at_version_meta);
748 let created_at_version_sequence = created_at_version_meta_to_pb(&f.created_at_version_meta);
749 Self {
750 id: f.id,
751 files: f.files.iter().map(pb::DataFile::from).collect(),
752 overlays: f.overlays.iter().map(pb::DataOverlayFile::from).collect(),
753 deletion_file,
754 row_id_sequence,
755 physical_rows: f.physical_rows.unwrap_or_default() as u64,
756 last_updated_at_version_sequence,
757 created_at_version_sequence,
758 }
759 }
760}
761
762#[cfg(test)]
763mod tests {
764 use super::*;
765 use crate::format::overlay::OverlayCoverage;
766 use arrow_schema::{
767 DataType, Field as ArrowField, Fields as ArrowFields, Schema as ArrowSchema,
768 };
769 use lance_file::format::{MAJOR_VERSION, MINOR_VERSION};
770 use object_store::path::Path;
771 use roaring::RoaringBitmap;
772 use serde_json::{Value, json};
773
774 #[test]
775 fn test_data_overlay_roundtrip() {
776 let mut bitmap = RoaringBitmap::new();
779 bitmap.insert(1);
780 bitmap.insert(3);
781
782 let overlay = DataOverlayFile {
783 data_file: DataFile::new_legacy_from_fields("overlay-0.lance", vec![3], None),
784 coverage: OverlayCoverage::dense(bitmap.clone()),
785 committed_version: 7,
786 };
787 let mut fragment = Fragment::new(0);
788 fragment.files = vec![DataFile::new_legacy_from_fields(
789 "base.lance",
790 vec![1, 3],
791 None,
792 )];
793 fragment.overlays = vec![overlay];
794
795 let proto = pb::DataFragment::from(&fragment);
796 assert_eq!(proto.overlays.len(), 1);
797 let round_tripped = Fragment::try_from(proto).unwrap();
798 assert_eq!(round_tripped, fragment);
799
800 let recovered = round_tripped.overlays[0].coverage_for_field(0).unwrap();
802 assert_eq!(*recovered, bitmap);
803 assert_eq!(
804 *round_tripped.overlays[0].coverage_for_field(5).unwrap(),
805 bitmap
806 );
807 }
808
809 #[test]
810 fn test_data_overlay_sparse_per_field_coverage() {
811 let name_coverage = RoaringBitmap::from_iter([2u32, 3]);
813 let embedding_coverage = RoaringBitmap::from_iter([1u32]);
814 let overlay = DataOverlayFile {
815 data_file: DataFile::new_legacy_from_fields("overlay-1.lance", vec![2, 4], None),
816 coverage: OverlayCoverage::sparse(vec![
817 name_coverage.clone(),
818 embedding_coverage.clone(),
819 ]),
820 committed_version: 3,
821 };
822 let mut fragment = Fragment::new(1);
823 fragment.overlays = vec![overlay];
824
825 let round_tripped = Fragment::try_from(pb::DataFragment::from(&fragment)).unwrap();
826 assert_eq!(
827 *round_tripped.overlays[0].coverage_for_field(0).unwrap(),
828 name_coverage
829 );
830 assert_eq!(
831 *round_tripped.overlays[0].coverage_for_field(1).unwrap(),
832 embedding_coverage
833 );
834 }
835
836 #[test]
837 fn test_overlays_sorted_newest_last_on_load() {
838 let mk = |version: u64, field: i32| DataOverlayFile {
841 data_file: DataFile::new_legacy_from_fields("o.lance", vec![field], None),
842 coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0u32])),
843 committed_version: version,
844 };
845 let mut fragment = Fragment::new(0);
846 fragment.overlays = vec![mk(5, 1), mk(2, 2), mk(2, 3), mk(3, 4)];
848
849 let loaded = Fragment::try_from(pb::DataFragment::from(&fragment)).unwrap();
850 let versions: Vec<u64> = loaded
851 .overlays
852 .iter()
853 .map(|o| o.committed_version)
854 .collect();
855 assert_eq!(versions, vec![2, 2, 3, 5]);
856 assert_eq!(
859 loaded.overlays[0].data_file.fields.as_ref(),
860 [2i32].as_slice()
861 );
862 assert_eq!(
863 loaded.overlays[1].data_file.fields.as_ref(),
864 [3i32].as_slice()
865 );
866 }
867
868 #[test]
869 fn test_new_fragment() {
870 let path = "foobar.lance";
871
872 let arrow_schema = ArrowSchema::new(vec![
873 ArrowField::new(
874 "s",
875 DataType::Struct(ArrowFields::from(vec![
876 ArrowField::new("si", DataType::Int32, false),
877 ArrowField::new("sb", DataType::Binary, true),
878 ])),
879 true,
880 ),
881 ArrowField::new("bool", DataType::Boolean, true),
882 ]);
883 let schema = Schema::try_from(&arrow_schema).unwrap();
884 let fragment = Fragment::with_file_legacy(123, path, &schema, Some(10));
885
886 assert_eq!(123, fragment.id);
887 assert_eq!(
888 fragment.files,
889 vec![DataFile::new_legacy_from_fields(
890 path.to_string(),
891 vec![0, 1, 2, 3],
892 None,
893 )]
894 )
895 }
896
897 #[test]
898 fn test_roundtrip_fragment() {
899 let mut fragment = Fragment::new(123);
900 let schema = ArrowSchema::new(vec![ArrowField::new("x", DataType::Float16, true)]);
901 fragment.add_file_legacy("foobar.lance", &Schema::try_from(&schema).unwrap());
902 fragment.deletion_file = Some(DeletionFile {
903 read_version: 123,
904 id: 456,
905 file_type: DeletionFileType::Array,
906 num_deleted_rows: Some(10),
907 base_id: None,
908 });
909
910 let proto = pb::DataFragment::from(&fragment);
911 let fragment2 = Fragment::try_from(proto).unwrap();
912 assert_eq!(fragment, fragment2);
913
914 fragment.deletion_file = None;
915 let proto = pb::DataFragment::from(&fragment);
916 let fragment2 = Fragment::try_from(proto).unwrap();
917 assert_eq!(fragment, fragment2);
918 }
919
920 #[test]
921 fn infer_exact_file_version_and_reject_mixed_fragments() {
922 assert_eq!(Fragment::try_infer_version(&[]).unwrap(), None);
923
924 let v2_0 = Fragment::new(0).with_file(
925 "v2_0.lance",
926 vec![0],
927 vec![0],
928 ConcreteFileVersion::V2_0,
929 None,
930 );
931 assert_eq!(
932 (
933 v2_0.files[0].file_major_version,
934 v2_0.files[0].file_minor_version
935 ),
936 (2, 0)
937 );
938 let unstarted = DataFile::new_unstarted("unstarted.lance", ConcreteFileVersion::V2_0);
939 assert_eq!(
940 (unstarted.file_major_version, unstarted.file_minor_version),
941 (2, 0)
942 );
943 let v2_0_second = Fragment::new(1).with_file(
944 "v2_0_second.lance",
945 vec![0],
946 vec![0],
947 ConcreteFileVersion::V2_0,
948 None,
949 );
950 assert_eq!(
951 Fragment::try_infer_version(&[v2_0.clone(), v2_0_second]).unwrap(),
952 Some(ConcreteFileVersion::V2_0)
953 );
954
955 let mut mixed_overlay = v2_0.clone();
956 mixed_overlay.overlays.push(DataOverlayFile {
957 data_file: DataFile::new(
958 "overlay-v2_1.lance",
959 vec![0],
960 vec![0],
961 ConcreteFileVersion::V2_1,
962 None,
963 None,
964 ),
965 coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0])),
966 committed_version: 1,
967 });
968 let error = Fragment::try_infer_version(&[mixed_overlay]).unwrap_err();
969 assert!(error.to_string().contains("2.0"));
970 assert!(error.to_string().contains("2.1"));
971
972 let v2_1 = Fragment::new(2).with_file(
973 "v2_1.lance",
974 vec![0],
975 vec![0],
976 ConcreteFileVersion::V2_1,
977 None,
978 );
979 let error = Fragment::try_infer_version(&[v2_0, v2_1]).unwrap_err();
980 assert!(matches!(error, Error::InvalidInput { .. }));
981 let message = error.to_string();
982 assert!(message.contains("All data files must have the same version"));
983 assert!(message.contains("2.0"));
984 assert!(message.contains("2.1"));
985 }
986
987 #[test]
988 fn test_to_json() {
989 let mut fragment = Fragment::new(123);
990 let schema = ArrowSchema::new(vec![ArrowField::new("x", DataType::Float16, true)]);
991 fragment.add_file_legacy("foobar.lance", &Schema::try_from(&schema).unwrap());
992 fragment.deletion_file = Some(DeletionFile {
993 read_version: 123,
994 id: 456,
995 file_type: DeletionFileType::Array,
996 num_deleted_rows: Some(10),
997 base_id: None,
998 });
999
1000 let json = serde_json::to_string(&fragment).unwrap();
1001
1002 let value: Value = serde_json::from_str(&json).unwrap();
1003 assert_eq!(
1004 value,
1005 json!({
1006 "id": 123,
1007 "files":[
1008 {"path": "foobar.lance", "fields": [0], "column_indices": [],
1009 "file_major_version": MAJOR_VERSION, "file_minor_version": MINOR_VERSION,
1010 "file_size_bytes": null, "base_id": null }
1011 ],
1012 "deletion_file": {"read_version": 123, "id": 456, "file_type": "array",
1013 "num_deleted_rows": 10, "base_id": null},
1014 "physical_rows": None::<usize>}),
1015 );
1016
1017 let frag2 = Fragment::from_json(&json).unwrap();
1018 assert_eq!(fragment, frag2);
1019 }
1020
1021 #[test]
1022 fn data_file_validate_allows_extra_columns() {
1023 let data_file = DataFile {
1024 path: "foo.lance".to_string(),
1025 fields: Arc::from([1, 2]),
1026 column_indices: Arc::from([0, 1, 2]),
1028 file_major_version: MAJOR_VERSION as u32,
1029 file_minor_version: MINOR_VERSION as u32,
1030 file_size_bytes: Default::default(),
1031 base_id: None,
1032 };
1033
1034 let base_path = Path::from("base");
1035 data_file
1036 .validate(&base_path)
1037 .expect("validation should allow extra columns without field ids");
1038 }
1039}