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, 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#[derive(Debug, Clone, PartialEq, Eq, DeepSizeOf)]
28pub struct DataFile {
29 pub path: String,
31 pub fields: Arc<[i32]>,
37 pub column_indices: Arc<[i32]>,
47 pub file_major_version: u32,
49 pub file_minor_version: u32,
51
52 pub file_size_bytes: CachedFileSize,
54
55 pub base_id: Option<u32>,
57}
58
59impl 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
75impl<'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 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 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 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 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#[derive(Default)]
244pub struct DataFileFieldInterner {
245 fields: InternCache<i32>,
246 column_indices: InternCache<i32>,
247 inline_bytes: InternCache<u8>,
248}
249
250enum 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 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 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 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 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 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 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#[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#[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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
485pub struct Fragment {
486 pub id: u64,
488
489 pub files: Vec<DataFile>,
491
492 #[serde(default, skip_serializing_if = "Vec::is_empty")]
496 pub overlays: Vec<DataOverlayFile>,
497
498 #[serde(skip_serializing_if = "Option::is_none")]
500 pub deletion_file: Option<DeletionFile>,
501
502 #[serde(skip_serializing_if = "Option::is_none")]
504 pub row_id_meta: Option<RowIdMeta>,
505
506 pub physical_rows: Option<usize>,
510
511 #[serde(skip_serializing_if = "Option::is_none")]
513 pub last_updated_at_version_meta: Option<RowDatasetVersionMeta>,
514
515 #[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 (Some(len), None) => Some(len),
538 (
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 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 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 pub fn has_legacy_files(&self) -> bool {
625 self.files[0].is_legacy_file()
627 }
628
629 pub fn try_infer_version(fragments: &[Self]) -> Result<Option<ConcreteFileVersion>> {
634 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 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 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 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 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 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 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 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 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}