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::format::{MAJOR_VERSION, MINOR_VERSION};
11use lance_file::version::LanceFileVersion;
12use lance_io::utils::CachedFileSize;
13use object_store::path::Path;
14use serde::{Deserialize, Deserializer, Serialize, Serializer};
15
16use super::overlay::{DataOverlayFile, sort_overlays_newest_last};
17use crate::format::pb;
18
19use crate::rowids::version::{
20 RowDatasetVersionMeta, created_at_version_meta_to_pb, last_updated_at_version_meta_to_pb,
21};
22use lance_core::datatypes::Schema;
23use lance_core::error::Result;
24
25#[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(
108 path: impl Into<String>,
109 fields: Vec<i32>,
110 column_indices: Vec<i32>,
111 file_major_version: u32,
112 file_minor_version: u32,
113 file_size_bytes: Option<NonZero<u64>>,
114 base_id: Option<u32>,
115 ) -> Self {
116 Self {
117 path: path.into(),
118 fields: Arc::from(fields),
119 column_indices: Arc::from(column_indices),
120 file_major_version,
121 file_minor_version,
122 file_size_bytes: file_size_bytes.into(),
123 base_id,
124 }
125 }
126
127 pub fn new_unstarted(
129 path: impl Into<String>,
130 file_major_version: u32,
131 file_minor_version: u32,
132 ) -> Self {
133 Self {
134 path: path.into(),
135 fields: Arc::from([]),
136 column_indices: Arc::from([]),
137 file_major_version,
138 file_minor_version,
139 file_size_bytes: Default::default(),
140 base_id: None,
141 }
142 }
143
144 pub fn new_legacy_from_fields(
145 path: impl Into<String>,
146 fields: Vec<i32>,
147 base_id: Option<u32>,
148 ) -> Self {
149 Self::new(
150 path,
151 fields,
152 vec![],
153 MAJOR_VERSION as u32,
154 MINOR_VERSION as u32,
155 None,
156 base_id,
157 )
158 }
159
160 pub fn new_legacy(
161 path: impl Into<String>,
162 schema: &Schema,
163 file_size_bytes: Option<NonZero<u64>>,
164 base_id: Option<u32>,
165 ) -> Self {
166 let mut field_ids = schema.field_ids();
167 field_ids.sort();
168 Self::new(
169 path,
170 field_ids,
171 vec![],
172 MAJOR_VERSION as u32,
173 MINOR_VERSION as u32,
174 file_size_bytes,
175 base_id,
176 )
177 }
178
179 pub fn schema(&self, full_schema: &Schema) -> Schema {
180 full_schema.project_by_ids(&self.fields, false)
181 }
182
183 pub fn is_legacy_file(&self) -> bool {
184 self.file_major_version == 0 && self.file_minor_version < 3
185 }
186
187 pub fn validate(&self, base_path: &Path) -> Result<()> {
188 if self.is_legacy_file() {
189 if !self.fields.windows(2).all(|w| w[0] < w[1]) {
190 return Err(Error::corrupt_file(
191 base_path.clone().join(self.path.clone()),
192 "contained unsorted or duplicate field ids",
193 ));
194 }
195 } else if self.column_indices.len() < self.fields.len() {
196 return Err(Error::corrupt_file(
199 base_path.clone().join(self.path.clone()),
200 "contained fewer column_indices than fields",
201 ));
202 }
203 Ok(())
204 }
205}
206
207impl From<&DataFile> for pb::DataFile {
208 fn from(df: &DataFile) -> Self {
209 Self {
210 path: df.path.clone(),
211 fields: df.fields.to_vec(),
212 column_indices: df.column_indices.to_vec(),
213 file_major_version: df.file_major_version,
214 file_minor_version: df.file_minor_version,
215 file_size_bytes: df.file_size_bytes.get().map_or(0, |v| v.get()),
216 base_id: df.base_id,
217 }
218 }
219}
220
221impl TryFrom<pb::DataFile> for DataFile {
222 type Error = Error;
223
224 fn try_from(proto: pb::DataFile) -> Result<Self> {
225 Ok(Self {
226 path: proto.path,
227 fields: Arc::from(proto.fields),
228 column_indices: Arc::from(proto.column_indices),
229 file_major_version: proto.file_major_version,
230 file_minor_version: proto.file_minor_version,
231 file_size_bytes: CachedFileSize::new(proto.file_size_bytes),
232 base_id: proto.base_id,
233 })
234 }
235}
236
237#[derive(Default)]
248pub struct DataFileFieldInterner {
249 fields: InternCache<i32>,
250 column_indices: InternCache<i32>,
251 inline_bytes: InternCache<u8>,
252}
253
254enum InternCache<T: Eq + std::hash::Hash + Clone> {
258 Small(Vec<Arc<[T]>>),
259 Large(HashMap<Arc<[T]>, ()>),
260}
261
262const INTERN_CACHE_UPGRADE_THRESHOLD: usize = 16;
263
264impl<T: Eq + std::hash::Hash + Clone> Default for InternCache<T> {
265 fn default() -> Self {
266 Self::Small(Vec::new())
267 }
268}
269
270impl<T: Eq + std::hash::Hash + Clone> InternCache<T> {
271 fn intern(&mut self, v: Vec<T>) -> Arc<[T]> {
272 match self {
273 Self::Small(entries) => {
274 for existing in entries.iter() {
275 if existing.as_ref() == v.as_slice() {
276 return existing.clone();
277 }
278 }
279 let arc: Arc<[T]> = Arc::from(v);
280 entries.push(arc.clone());
281 if entries.len() > INTERN_CACHE_UPGRADE_THRESHOLD {
282 let mut map = HashMap::with_capacity(entries.len());
283 for e in entries.drain(..) {
284 map.insert(e, ());
285 }
286 *self = Self::Large(map);
287 }
288 arc
289 }
290 Self::Large(map) => {
291 if let Some((existing, _)) = map.get_key_value(v.as_slice()) {
292 existing.clone()
293 } else {
294 let arc: Arc<[T]> = Arc::from(v);
295 map.insert(arc.clone(), ());
296 arc
297 }
298 }
299 }
300 }
301}
302
303impl DataFileFieldInterner {
304 fn intern_last_updated_version_meta(
308 cache: &mut InternCache<u8>,
309 pb: pb::data_fragment::LastUpdatedAtVersionSequence,
310 ) -> Result<RowDatasetVersionMeta> {
311 match pb {
312 pb::data_fragment::LastUpdatedAtVersionSequence::InlineLastUpdatedAtVersions(data) => {
313 Ok(RowDatasetVersionMeta::Inline(cache.intern(data)))
314 }
315 pb::data_fragment::LastUpdatedAtVersionSequence::ExternalLastUpdatedAtVersions(
316 file,
317 ) => Ok(RowDatasetVersionMeta::External(ExternalFile {
318 path: file.path,
319 offset: file.offset,
320 size: file.size,
321 })),
322 }
323 }
324
325 fn intern_created_version_meta(
327 cache: &mut InternCache<u8>,
328 pb: pb::data_fragment::CreatedAtVersionSequence,
329 ) -> Result<RowDatasetVersionMeta> {
330 match pb {
331 pb::data_fragment::CreatedAtVersionSequence::InlineCreatedAtVersions(data) => {
332 Ok(RowDatasetVersionMeta::Inline(cache.intern(data)))
333 }
334 pb::data_fragment::CreatedAtVersionSequence::ExternalCreatedAtVersions(file) => {
335 Ok(RowDatasetVersionMeta::External(ExternalFile {
336 path: file.path,
337 offset: file.offset,
338 size: file.size,
339 }))
340 }
341 }
342 }
343
344 pub fn intern_data_file(&mut self, proto: pb::DataFile) -> Result<DataFile> {
346 Ok(DataFile {
347 path: proto.path,
348 fields: self.fields.intern(proto.fields),
349 column_indices: self.column_indices.intern(proto.column_indices),
350 file_major_version: proto.file_major_version,
351 file_minor_version: proto.file_minor_version,
352 file_size_bytes: CachedFileSize::new(proto.file_size_bytes),
353 base_id: proto.base_id,
354 })
355 }
356
357 pub fn intern_fragment(&mut self, p: pb::DataFragment) -> Result<Fragment> {
359 let physical_rows = if p.physical_rows > 0 {
360 Some(p.physical_rows as usize)
361 } else {
362 None
363 };
364 let last_updated_at_version_meta = p
365 .last_updated_at_version_sequence
366 .map(|pb| Self::intern_last_updated_version_meta(&mut self.inline_bytes, pb))
367 .transpose()?;
368 let created_at_version_meta = p
369 .created_at_version_sequence
370 .map(|pb| Self::intern_created_version_meta(&mut self.inline_bytes, pb))
371 .transpose()?;
372 Ok(Fragment {
373 id: p.id,
374 files: p
375 .files
376 .into_iter()
377 .map(|f| self.intern_data_file(f))
378 .collect::<Result<_>>()?,
379 overlays: {
380 let mut overlays = p
381 .overlays
382 .into_iter()
383 .map(DataOverlayFile::try_from)
384 .collect::<Result<Vec<_>>>()?;
385 sort_overlays_newest_last(&mut overlays);
386 overlays
387 },
388 deletion_file: p.deletion_file.map(DeletionFile::try_from).transpose()?,
389 row_id_meta: p.row_id_sequence.map(RowIdMeta::try_from).transpose()?,
390 physical_rows,
391 last_updated_at_version_meta,
392 created_at_version_meta,
393 })
394 }
395}
396
397#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
398#[serde(rename_all = "lowercase")]
399pub enum DeletionFileType {
400 Array,
401 Bitmap,
402}
403
404impl DeletionFileType {
405 pub fn suffix(&self) -> &str {
407 match self {
408 Self::Array => "arrow",
409 Self::Bitmap => "bin",
410 }
411 }
412}
413
414#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
415pub struct DeletionFile {
416 pub read_version: u64,
417 pub id: u64,
418 pub file_type: DeletionFileType,
419 pub num_deleted_rows: Option<usize>,
421 pub base_id: Option<u32>,
422}
423
424impl TryFrom<pb::DeletionFile> for DeletionFile {
425 type Error = Error;
426
427 fn try_from(value: pb::DeletionFile) -> Result<Self> {
428 let file_type = match value.file_type {
429 0 => DeletionFileType::Array,
430 1 => DeletionFileType::Bitmap,
431 _ => {
432 return Err(Error::not_supported_source(
433 "Unknown deletion file type".into(),
434 ));
435 }
436 };
437 let num_deleted_rows = if value.num_deleted_rows == 0 {
438 None
439 } else {
440 Some(value.num_deleted_rows as usize)
441 };
442 Ok(Self {
443 read_version: value.read_version,
444 id: value.id,
445 file_type,
446 num_deleted_rows,
447 base_id: value.base_id,
448 })
449 }
450}
451
452#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
454pub struct ExternalFile {
455 pub path: String,
456 pub offset: u64,
457 pub size: u64,
458}
459
460#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
462pub enum RowIdMeta {
463 Inline(Vec<u8>),
464 External(ExternalFile),
465}
466
467impl TryFrom<pb::data_fragment::RowIdSequence> for RowIdMeta {
468 type Error = Error;
469
470 fn try_from(value: pb::data_fragment::RowIdSequence) -> Result<Self> {
471 match value {
472 pb::data_fragment::RowIdSequence::InlineRowIds(data) => Ok(Self::Inline(data)),
473 pb::data_fragment::RowIdSequence::ExternalRowIds(file) => {
474 Ok(Self::External(ExternalFile {
475 path: file.path.clone(),
476 offset: file.offset,
477 size: file.size,
478 }))
479 }
480 }
481 }
482}
483
484#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
489pub struct Fragment {
490 pub id: u64,
492
493 pub files: Vec<DataFile>,
495
496 #[serde(default, skip_serializing_if = "Vec::is_empty")]
500 pub overlays: Vec<DataOverlayFile>,
501
502 #[serde(skip_serializing_if = "Option::is_none")]
504 pub deletion_file: Option<DeletionFile>,
505
506 #[serde(skip_serializing_if = "Option::is_none")]
508 pub row_id_meta: Option<RowIdMeta>,
509
510 pub physical_rows: Option<usize>,
514
515 #[serde(skip_serializing_if = "Option::is_none")]
517 pub last_updated_at_version_meta: Option<RowDatasetVersionMeta>,
518
519 #[serde(skip_serializing_if = "Option::is_none")]
521 pub created_at_version_meta: Option<RowDatasetVersionMeta>,
522}
523
524impl Fragment {
525 pub fn new(id: u64) -> Self {
526 Self {
527 id,
528 files: vec![],
529 overlays: vec![],
530 deletion_file: None,
531 row_id_meta: None,
532 physical_rows: None,
533 last_updated_at_version_meta: None,
534 created_at_version_meta: None,
535 }
536 }
537
538 pub fn num_rows(&self) -> Option<usize> {
539 match (self.physical_rows, &self.deletion_file) {
540 (Some(len), None) => Some(len),
542 (
544 Some(len),
545 Some(DeletionFile {
546 num_deleted_rows: Some(num_deleted_rows),
547 ..
548 }),
549 ) => Some(len - num_deleted_rows),
550 _ => None,
551 }
552 }
553
554 pub fn from_json(json: &str) -> Result<Self> {
555 let fragment: Self = serde_json::from_str(json)?;
556 Ok(fragment)
557 }
558
559 pub fn with_file_legacy(
561 id: u64,
562 path: &str,
563 schema: &Schema,
564 physical_rows: Option<usize>,
565 ) -> Self {
566 Self {
567 id,
568 files: vec![DataFile::new_legacy(path, schema, None, None)],
569 overlays: vec![],
570 deletion_file: None,
571 physical_rows,
572 row_id_meta: None,
573 last_updated_at_version_meta: None,
574 created_at_version_meta: None,
575 }
576 }
577
578 pub fn with_file(
579 mut self,
580 path: impl Into<String>,
581 field_ids: Vec<i32>,
582 column_indices: Vec<i32>,
583 version: &LanceFileVersion,
584 file_size_bytes: Option<NonZero<u64>>,
585 ) -> Self {
586 let (major, minor) = version.to_numbers();
587 let data_file = DataFile::new(
588 path,
589 field_ids,
590 column_indices,
591 major,
592 minor,
593 file_size_bytes,
594 None,
595 );
596 self.files.push(data_file);
597 self
598 }
599
600 pub fn with_physical_rows(mut self, physical_rows: usize) -> Self {
601 self.physical_rows = Some(physical_rows);
602 self
603 }
604
605 pub fn add_file(
606 &mut self,
607 path: impl Into<String>,
608 field_ids: Vec<i32>,
609 column_indices: Vec<i32>,
610 version: &LanceFileVersion,
611 file_size_bytes: Option<NonZero<u64>>,
612 ) {
613 let (major, minor) = version.to_numbers();
614 self.files.push(DataFile::new(
615 path,
616 field_ids,
617 column_indices,
618 major,
619 minor,
620 file_size_bytes,
621 None,
622 ));
623 }
624
625 pub fn add_file_legacy(&mut self, path: &str, schema: &Schema) {
627 self.files
628 .push(DataFile::new_legacy(path, schema, None, None));
629 }
630
631 pub fn has_legacy_files(&self) -> bool {
633 self.files[0].is_legacy_file()
635 }
636
637 pub fn try_infer_version(fragments: &[Self]) -> Result<Option<LanceFileVersion>> {
642 let Some(sample_file) = fragments
645 .iter()
646 .find(|f| !f.files.is_empty())
647 .map(|f| &f.files[0])
648 else {
649 return Ok(None);
650 };
651 let file_version = LanceFileVersion::try_from_major_minor(
652 sample_file.file_major_version,
653 sample_file.file_minor_version,
654 )?;
655 for frag in fragments {
657 for file in &frag.files {
658 let this_file_version = LanceFileVersion::try_from_major_minor(
659 file.file_major_version,
660 file.file_minor_version,
661 )?;
662 if file_version != this_file_version {
663 return Err(Error::invalid_input(format!(
664 "All data files must have the same version. Detected both {} and {}",
665 file_version, this_file_version
666 )));
667 }
668 }
669 }
670 Ok(Some(file_version))
671 }
672}
673
674impl TryFrom<pb::DataFragment> for Fragment {
675 type Error = Error;
676
677 fn try_from(p: pb::DataFragment) -> Result<Self> {
678 let physical_rows = if p.physical_rows > 0 {
679 Some(p.physical_rows as usize)
680 } else {
681 None
682 };
683 Ok(Self {
684 id: p.id,
685 files: p
686 .files
687 .into_iter()
688 .map(DataFile::try_from)
689 .collect::<Result<_>>()?,
690 overlays: {
691 let mut overlays = p
692 .overlays
693 .into_iter()
694 .map(DataOverlayFile::try_from)
695 .collect::<Result<Vec<_>>>()?;
696 sort_overlays_newest_last(&mut overlays);
697 overlays
698 },
699 deletion_file: p.deletion_file.map(DeletionFile::try_from).transpose()?,
700 row_id_meta: p.row_id_sequence.map(RowIdMeta::try_from).transpose()?,
701 physical_rows,
702 last_updated_at_version_meta: p
703 .last_updated_at_version_sequence
704 .map(RowDatasetVersionMeta::try_from)
705 .transpose()?,
706 created_at_version_meta: p
707 .created_at_version_sequence
708 .map(RowDatasetVersionMeta::try_from)
709 .transpose()?,
710 })
711 }
712}
713
714impl From<&Fragment> for pb::DataFragment {
715 fn from(f: &Fragment) -> Self {
716 let deletion_file = f.deletion_file.as_ref().map(|f| {
717 let file_type = match f.file_type {
718 DeletionFileType::Array => pb::deletion_file::DeletionFileType::ArrowArray,
719 DeletionFileType::Bitmap => pb::deletion_file::DeletionFileType::Bitmap,
720 };
721 pb::DeletionFile {
722 read_version: f.read_version,
723 id: f.id,
724 file_type: file_type.into(),
725 num_deleted_rows: f.num_deleted_rows.unwrap_or_default() as u64,
726 base_id: f.base_id,
727 }
728 });
729
730 let row_id_sequence = f.row_id_meta.as_ref().map(|m| match m {
731 RowIdMeta::Inline(data) => pb::data_fragment::RowIdSequence::InlineRowIds(data.clone()),
732 RowIdMeta::External(file) => {
733 pb::data_fragment::RowIdSequence::ExternalRowIds(pb::ExternalFile {
734 path: file.path.clone(),
735 offset: file.offset,
736 size: file.size,
737 })
738 }
739 });
740 let last_updated_at_version_sequence =
741 last_updated_at_version_meta_to_pb(&f.last_updated_at_version_meta);
742 let created_at_version_sequence = created_at_version_meta_to_pb(&f.created_at_version_meta);
743 Self {
744 id: f.id,
745 files: f.files.iter().map(pb::DataFile::from).collect(),
746 overlays: f.overlays.iter().map(pb::DataOverlayFile::from).collect(),
747 deletion_file,
748 row_id_sequence,
749 physical_rows: f.physical_rows.unwrap_or_default() as u64,
750 last_updated_at_version_sequence,
751 created_at_version_sequence,
752 }
753 }
754}
755
756#[cfg(test)]
757mod tests {
758 use super::*;
759 use crate::format::overlay::OverlayCoverage;
760 use arrow_schema::{
761 DataType, Field as ArrowField, Fields as ArrowFields, Schema as ArrowSchema,
762 };
763 use object_store::path::Path;
764 use roaring::RoaringBitmap;
765 use serde_json::{Value, json};
766
767 #[test]
768 fn test_data_overlay_roundtrip() {
769 let mut bitmap = RoaringBitmap::new();
772 bitmap.insert(1);
773 bitmap.insert(3);
774
775 let overlay = DataOverlayFile {
776 data_file: DataFile::new_legacy_from_fields("overlay-0.lance", vec![3], None),
777 coverage: OverlayCoverage::dense(bitmap.clone()),
778 committed_version: 7,
779 };
780 let mut fragment = Fragment::new(0);
781 fragment.files = vec![DataFile::new_legacy_from_fields(
782 "base.lance",
783 vec![1, 3],
784 None,
785 )];
786 fragment.overlays = vec![overlay];
787
788 let proto = pb::DataFragment::from(&fragment);
789 assert_eq!(proto.overlays.len(), 1);
790 let round_tripped = Fragment::try_from(proto).unwrap();
791 assert_eq!(round_tripped, fragment);
792
793 let recovered = round_tripped.overlays[0].coverage_for_field(0).unwrap();
795 assert_eq!(*recovered, bitmap);
796 assert_eq!(
797 *round_tripped.overlays[0].coverage_for_field(5).unwrap(),
798 bitmap
799 );
800 }
801
802 #[test]
803 fn test_data_overlay_sparse_per_field_coverage() {
804 let name_coverage = RoaringBitmap::from_iter([2u32, 3]);
806 let embedding_coverage = RoaringBitmap::from_iter([1u32]);
807 let overlay = DataOverlayFile {
808 data_file: DataFile::new_legacy_from_fields("overlay-1.lance", vec![2, 4], None),
809 coverage: OverlayCoverage::sparse(vec![
810 name_coverage.clone(),
811 embedding_coverage.clone(),
812 ]),
813 committed_version: 3,
814 };
815 let mut fragment = Fragment::new(1);
816 fragment.overlays = vec![overlay];
817
818 let round_tripped = Fragment::try_from(pb::DataFragment::from(&fragment)).unwrap();
819 assert_eq!(
820 *round_tripped.overlays[0].coverage_for_field(0).unwrap(),
821 name_coverage
822 );
823 assert_eq!(
824 *round_tripped.overlays[0].coverage_for_field(1).unwrap(),
825 embedding_coverage
826 );
827 }
828
829 #[test]
830 fn test_overlays_sorted_newest_last_on_load() {
831 let mk = |version: u64, field: i32| DataOverlayFile {
834 data_file: DataFile::new_legacy_from_fields("o.lance", vec![field], None),
835 coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0u32])),
836 committed_version: version,
837 };
838 let mut fragment = Fragment::new(0);
839 fragment.overlays = vec![mk(5, 1), mk(2, 2), mk(2, 3), mk(3, 4)];
841
842 let loaded = Fragment::try_from(pb::DataFragment::from(&fragment)).unwrap();
843 let versions: Vec<u64> = loaded
844 .overlays
845 .iter()
846 .map(|o| o.committed_version)
847 .collect();
848 assert_eq!(versions, vec![2, 2, 3, 5]);
849 assert_eq!(
852 loaded.overlays[0].data_file.fields.as_ref(),
853 [2i32].as_slice()
854 );
855 assert_eq!(
856 loaded.overlays[1].data_file.fields.as_ref(),
857 [3i32].as_slice()
858 );
859 }
860
861 #[test]
862 fn test_new_fragment() {
863 let path = "foobar.lance";
864
865 let arrow_schema = ArrowSchema::new(vec![
866 ArrowField::new(
867 "s",
868 DataType::Struct(ArrowFields::from(vec![
869 ArrowField::new("si", DataType::Int32, false),
870 ArrowField::new("sb", DataType::Binary, true),
871 ])),
872 true,
873 ),
874 ArrowField::new("bool", DataType::Boolean, true),
875 ]);
876 let schema = Schema::try_from(&arrow_schema).unwrap();
877 let fragment = Fragment::with_file_legacy(123, path, &schema, Some(10));
878
879 assert_eq!(123, fragment.id);
880 assert_eq!(
881 fragment.files,
882 vec![DataFile::new_legacy_from_fields(
883 path.to_string(),
884 vec![0, 1, 2, 3],
885 None,
886 )]
887 )
888 }
889
890 #[test]
891 fn test_roundtrip_fragment() {
892 let mut fragment = Fragment::new(123);
893 let schema = ArrowSchema::new(vec![ArrowField::new("x", DataType::Float16, true)]);
894 fragment.add_file_legacy("foobar.lance", &Schema::try_from(&schema).unwrap());
895 fragment.deletion_file = Some(DeletionFile {
896 read_version: 123,
897 id: 456,
898 file_type: DeletionFileType::Array,
899 num_deleted_rows: Some(10),
900 base_id: None,
901 });
902
903 let proto = pb::DataFragment::from(&fragment);
904 let fragment2 = Fragment::try_from(proto).unwrap();
905 assert_eq!(fragment, fragment2);
906
907 fragment.deletion_file = None;
908 let proto = pb::DataFragment::from(&fragment);
909 let fragment2 = Fragment::try_from(proto).unwrap();
910 assert_eq!(fragment, fragment2);
911 }
912
913 #[test]
914 fn test_to_json() {
915 let mut fragment = Fragment::new(123);
916 let schema = ArrowSchema::new(vec![ArrowField::new("x", DataType::Float16, true)]);
917 fragment.add_file_legacy("foobar.lance", &Schema::try_from(&schema).unwrap());
918 fragment.deletion_file = Some(DeletionFile {
919 read_version: 123,
920 id: 456,
921 file_type: DeletionFileType::Array,
922 num_deleted_rows: Some(10),
923 base_id: None,
924 });
925
926 let json = serde_json::to_string(&fragment).unwrap();
927
928 let value: Value = serde_json::from_str(&json).unwrap();
929 assert_eq!(
930 value,
931 json!({
932 "id": 123,
933 "files":[
934 {"path": "foobar.lance", "fields": [0], "column_indices": [],
935 "file_major_version": MAJOR_VERSION, "file_minor_version": MINOR_VERSION,
936 "file_size_bytes": null, "base_id": null }
937 ],
938 "deletion_file": {"read_version": 123, "id": 456, "file_type": "array",
939 "num_deleted_rows": 10, "base_id": null},
940 "physical_rows": None::<usize>}),
941 );
942
943 let frag2 = Fragment::from_json(&json).unwrap();
944 assert_eq!(fragment, frag2);
945 }
946
947 #[test]
948 fn data_file_validate_allows_extra_columns() {
949 let data_file = DataFile {
950 path: "foo.lance".to_string(),
951 fields: Arc::from([1, 2]),
952 column_indices: Arc::from([0, 1, 2]),
954 file_major_version: MAJOR_VERSION as u32,
955 file_minor_version: MINOR_VERSION as u32,
956 file_size_bytes: Default::default(),
957 base_id: None,
958 };
959
960 let base_path = Path::from("base");
961 data_file
962 .validate(&base_path)
963 .expect("validation should allow extra columns without field ids");
964 }
965}