1use async_trait::async_trait;
5use chrono::prelude::*;
6use lance_core::deepsize::DeepSizeOf;
7use lance_file::datatypes::{Fields, FieldsWithMeta};
8use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
9use lance_file::versions::v1::{
10 encoding::populate_schema_dictionaries, reader::FileReader as V1FileReader,
11};
12use lance_io::traits::{ProtoStruct, Reader};
13use object_store::path::Path;
14use prost::Message;
15use prost_types::Timestamp;
16use std::collections::{BTreeMap, HashMap};
17use std::ops::Range;
18use std::sync::Arc;
19
20use super::Fragment;
21use crate::feature_flags::{FLAG_STABLE_ROW_IDS, has_deprecated_v2_feature_flag};
22use crate::format::fragment::DataFileFieldInterner;
23use crate::format::pb;
24use lance_core::cache::LanceCache;
25use lance_core::datatypes::Schema;
26use lance_core::{Error, Result};
27use lance_io::object_store::{ObjectStore, ObjectStoreRegistry};
28use lance_io::utils::read_struct;
29
30#[derive(Debug, Clone, PartialEq, DeepSizeOf)]
37pub struct Manifest {
38 pub schema: Schema,
40
41 pub version: u64,
43
44 pub branch: Option<String>,
46
47 pub writer_version: Option<WriterVersion>,
49
50 pub fragments: Arc<Vec<Fragment>>,
55
56 pub version_aux_data: usize,
58
59 pub index_section: Option<usize>,
61
62 pub timestamp_nanos: u128,
64
65 pub tag: Option<String>,
67
68 pub reader_feature_flags: u64,
70
71 pub writer_feature_flags: u64,
73
74 pub max_fragment_id: Option<u32>,
77
78 pub transaction_file: Option<String>,
80
81 pub transaction_section: Option<usize>,
83
84 fragment_offsets: Vec<usize>,
87
88 pub next_row_id: u64,
90
91 pub data_storage_format: DataStorageFormat,
93
94 pub config: HashMap<String, String>,
96
97 pub table_metadata: HashMap<String, String>,
103
104 pub base_paths: HashMap<u32, BasePath>,
106}
107
108pub const DETACHED_VERSION_MASK: u64 = 0x8000_0000_0000_0000;
110
111pub fn is_detached_version(version: u64) -> bool {
112 version & DETACHED_VERSION_MASK != 0
113}
114
115fn compute_fragment_offsets(fragments: &[Fragment]) -> Vec<usize> {
116 fragments
117 .iter()
118 .map(|f| f.num_rows().unwrap_or_default())
119 .chain([0]) .scan(0_usize, |offset, len| {
121 let start = *offset;
122 *offset += len;
123 Some(start)
124 })
125 .collect()
126}
127
128#[derive(Default)]
129pub struct ManifestSummary {
130 pub total_fragments: u64,
131 pub total_data_files: u64,
132 pub total_files_size: u64,
133 pub total_deletion_files: u64,
134 pub total_data_file_rows: u64,
135 pub total_deletion_file_rows: u64,
136 pub total_rows: u64,
137}
138
139impl From<ManifestSummary> for BTreeMap<String, String> {
140 fn from(summary: ManifestSummary) -> Self {
141 let mut stats_map = Self::new();
142 stats_map.insert(
143 "total_fragments".to_string(),
144 summary.total_fragments.to_string(),
145 );
146 stats_map.insert(
147 "total_data_files".to_string(),
148 summary.total_data_files.to_string(),
149 );
150 stats_map.insert(
151 "total_files_size".to_string(),
152 summary.total_files_size.to_string(),
153 );
154 stats_map.insert(
155 "total_deletion_files".to_string(),
156 summary.total_deletion_files.to_string(),
157 );
158 stats_map.insert(
159 "total_data_file_rows".to_string(),
160 summary.total_data_file_rows.to_string(),
161 );
162 stats_map.insert(
163 "total_deletion_file_rows".to_string(),
164 summary.total_deletion_file_rows.to_string(),
165 );
166 stats_map.insert("total_rows".to_string(), summary.total_rows.to_string());
167 stats_map
168 }
169}
170
171impl Manifest {
172 pub fn new(
173 schema: Schema,
174 fragments: Arc<Vec<Fragment>>,
175 data_storage_format: DataStorageFormat,
176 base_paths: HashMap<u32, BasePath>,
177 ) -> Self {
178 let fragment_offsets = compute_fragment_offsets(&fragments);
179
180 Self {
181 schema,
182 version: 1,
183 branch: None,
184 writer_version: Some(WriterVersion::default()),
185 fragments,
186 version_aux_data: 0,
187 index_section: None,
188 timestamp_nanos: 0,
189 tag: None,
190 reader_feature_flags: 0,
191 writer_feature_flags: 0,
192 max_fragment_id: None,
193 transaction_file: None,
194 transaction_section: None,
195 fragment_offsets,
196 next_row_id: 0,
197 data_storage_format,
198 config: HashMap::new(),
199 table_metadata: HashMap::new(),
200 base_paths,
201 }
202 }
203
204 pub fn new_from_previous(
205 previous: &Self,
206 schema: Schema,
207 fragments: Arc<Vec<Fragment>>,
208 ) -> Self {
209 let fragment_offsets = compute_fragment_offsets(&fragments);
210
211 Self {
212 schema,
213 version: previous.version + 1,
214 branch: previous.branch.clone(),
215 writer_version: Some(WriterVersion::default()),
216 fragments,
217 version_aux_data: 0,
218 index_section: None, timestamp_nanos: 0, tag: None,
221 reader_feature_flags: 0, writer_feature_flags: 0, max_fragment_id: previous.max_fragment_id,
224 transaction_file: None,
225 transaction_section: None,
226 fragment_offsets,
227 next_row_id: previous.next_row_id,
228 data_storage_format: previous.data_storage_format.clone(),
229 config: previous.config.clone(),
230 table_metadata: previous.table_metadata.clone(),
231 base_paths: previous.base_paths.clone(),
232 }
233 }
234
235 pub fn shallow_clone(
240 &self,
241 ref_name: Option<String>,
242 ref_path: String,
243 ref_base_id: u32,
244 branch_name: Option<String>,
245 transaction_file: String,
246 ) -> Self {
247 let cloned_fragments = self
248 .fragments
249 .as_ref()
250 .iter()
251 .map(|fragment| {
252 let mut cloned_fragment = fragment.clone();
253 for file in &mut cloned_fragment.files {
254 if file.base_id.is_none() {
255 file.base_id = Some(ref_base_id);
256 }
257 }
258
259 if let Some(deletion) = &mut cloned_fragment.deletion_file
260 && deletion.base_id.is_none()
261 {
262 deletion.base_id = Some(ref_base_id);
263 }
264 cloned_fragment
265 })
266 .collect::<Vec<_>>();
267
268 Self {
269 schema: self.schema.clone(),
270 version: self.version,
271 branch: branch_name,
272 writer_version: self.writer_version.clone(),
273 fragments: Arc::new(cloned_fragments),
274 version_aux_data: self.version_aux_data,
275 index_section: None, timestamp_nanos: self.timestamp_nanos,
277 tag: None,
278 reader_feature_flags: 0, writer_feature_flags: 0, max_fragment_id: self.max_fragment_id,
281 transaction_file: Some(transaction_file),
282 transaction_section: None,
283 fragment_offsets: self.fragment_offsets.clone(),
284 next_row_id: self.next_row_id,
285 data_storage_format: self.data_storage_format.clone(),
286 config: self.config.clone(),
287 base_paths: {
288 let mut base_paths = self.base_paths.clone();
289 let base_path = BasePath::new(ref_base_id, ref_path, ref_name, true);
290 base_paths.insert(ref_base_id, base_path);
291 base_paths
292 },
293 table_metadata: self.table_metadata.clone(),
294 }
295 }
296
297 pub fn timestamp(&self) -> DateTime<Utc> {
299 let nanos = self.timestamp_nanos % 1_000_000_000;
300 let seconds = ((self.timestamp_nanos - nanos) / 1_000_000_000) as i64;
301 Utc.from_utc_datetime(
302 &DateTime::from_timestamp(seconds, nanos as u32)
303 .unwrap_or_default()
304 .naive_utc(),
305 )
306 }
307
308 pub fn set_timestamp(&mut self, nanos: u128) {
310 self.timestamp_nanos = nanos;
311 }
312
313 pub fn config_mut(&mut self) -> &mut HashMap<String, String> {
315 &mut self.config
316 }
317
318 pub fn table_metadata_mut(&mut self) -> &mut HashMap<String, String> {
320 &mut self.table_metadata
321 }
322
323 pub fn schema_metadata_mut(&mut self) -> &mut HashMap<String, String> {
325 &mut self.schema.metadata
326 }
327
328 pub fn field_metadata_mut(&mut self, field_id: i32) -> Option<&mut HashMap<String, String>> {
332 self.schema
333 .field_by_id_mut(field_id)
334 .map(|field| &mut field.metadata)
335 }
336
337 #[deprecated(note = "Use config_mut() for direct access to config HashMap")]
339 pub fn update_config(&mut self, upsert_values: impl IntoIterator<Item = (String, String)>) {
340 self.config.extend(upsert_values);
341 }
342
343 #[deprecated(note = "Use config_mut() for direct access to config HashMap")]
345 pub fn delete_config_keys(&mut self, delete_keys: &[&str]) {
346 self.config
347 .retain(|key, _| !delete_keys.contains(&key.as_str()));
348 }
349
350 #[deprecated(note = "Use schema_metadata_mut() for direct access to schema metadata HashMap")]
352 pub fn replace_schema_metadata(&mut self, new_metadata: HashMap<String, String>) {
353 self.schema.metadata = new_metadata;
354 }
355
356 #[deprecated(
360 note = "Use field_metadata_mut(field_id) for direct access to field metadata HashMap"
361 )]
362 pub fn replace_field_metadata(
363 &mut self,
364 field_id: i32,
365 new_metadata: HashMap<String, String>,
366 ) -> Result<()> {
367 if let Some(field) = self.schema.field_by_id_mut(field_id) {
368 field.metadata = new_metadata;
369 Ok(())
370 } else {
371 Err(Error::invalid_input(format!(
372 "Field with id {} does not exist for replace_field_metadata",
373 field_id
374 )))
375 }
376 }
377
378 pub fn update_max_fragment_id(&mut self) {
380 if self.fragments.is_empty() {
382 return;
383 }
384
385 let max_fragment_id = self
386 .fragments
387 .iter()
388 .map(|f| f.id)
389 .max()
390 .unwrap() .try_into()
392 .unwrap();
393
394 match self.max_fragment_id {
395 None => {
396 self.max_fragment_id = Some(max_fragment_id);
398 }
399 Some(current_max) => {
400 if max_fragment_id > current_max {
403 self.max_fragment_id = Some(max_fragment_id);
404 }
405 }
406 }
407 }
408
409 pub fn max_fragment_id(&self) -> Option<u64> {
414 if let Some(max_id) = self.max_fragment_id {
415 Some(max_id.into())
417 } else {
418 self.fragments.iter().map(|f| f.id).max()
420 }
421 }
422
423 pub fn max_field_id(&self) -> i32 {
428 let schema_max_id = self.schema.max_field_id().unwrap_or(-1);
429 let fragment_max_id = self
430 .fragments
431 .iter()
432 .flat_map(|f| f.files.iter().flat_map(|file| file.fields.iter()))
433 .max()
434 .copied();
435 let fragment_max_id = fragment_max_id.unwrap_or(-1);
436 schema_max_id.max(fragment_max_id)
437 }
438
439 pub fn fragments_since(&self, since: &Self) -> Result<Vec<Fragment>> {
442 if since.version >= self.version {
443 return Err(Error::invalid_input(format!(
444 "fragments_since: given version {} is newer than manifest version {}",
445 since.version, self.version
446 )));
447 }
448 let start = since.max_fragment_id();
449 Ok(self
450 .fragments
451 .iter()
452 .filter(|&f| start.map(|s| f.id > s).unwrap_or(true))
453 .cloned()
454 .collect())
455 }
456
457 pub fn fragments_by_offset_range(&self, range: Range<usize>) -> Vec<(usize, &Fragment)> {
473 let start = range.start;
474 let end = range.end;
475 let idx = self
476 .fragment_offsets
477 .binary_search(&start)
478 .unwrap_or_else(|idx| idx - 1);
479
480 let mut fragments = vec![];
481 for i in idx..self.fragments.len() {
482 if self.fragment_offsets[i] >= end
483 || self.fragment_offsets[i] + self.fragments[i].num_rows().unwrap_or_default()
484 <= start
485 {
486 break;
487 }
488 fragments.push((self.fragment_offsets[i], &self.fragments[i]));
489 }
490
491 fragments
492 }
493
494 pub fn uses_stable_row_ids(&self) -> bool {
496 self.reader_feature_flags & FLAG_STABLE_ROW_IDS != 0
497 }
498
499 pub fn serialized(&self) -> Vec<u8> {
502 let pb_manifest: pb::Manifest = self.into();
503 pb_manifest.encode_to_vec()
504 }
505
506 pub fn should_use_legacy_format(&self) -> bool {
507 self.data_storage_format.version == ConcreteFileVersion::V1
508 }
509
510 pub fn summary(&self) -> ManifestSummary {
521 let mut summary =
523 self.fragments
524 .iter()
525 .fold(ManifestSummary::default(), |mut summary, f| {
526 summary.total_data_files += f.files.len() as u64;
528 if let Some(num_rows) = f.num_rows() {
530 summary.total_rows += num_rows as u64;
531 }
532 for data_file in &f.files {
534 if let Some(size_bytes) = data_file.file_size_bytes.get() {
535 summary.total_files_size += size_bytes.get();
536 }
537 }
538 if f.deletion_file.is_some() {
540 summary.total_deletion_files += 1;
541 }
542 if let Some(deletion_file) = &f.deletion_file
544 && let Some(num_deleted) = deletion_file.num_deleted_rows
545 {
546 summary.total_deletion_file_rows += num_deleted as u64;
547 }
548 summary
549 });
550 summary.total_fragments = self.fragments.len() as u64;
551 summary.total_data_file_rows = summary.total_rows + summary.total_deletion_file_rows;
552
553 summary
554 }
555}
556
557pub async fn populate_manifest_schema_dictionaries(
576 manifest: &mut Manifest,
577 reader: &dyn Reader,
578) -> Result<()> {
579 match manifest.data_storage_format.version {
580 ConcreteFileVersion::V1 => {
581 populate_schema_dictionaries(&mut manifest.schema, reader).await?;
582 }
583 ConcreteFileVersion::V2_0
584 | ConcreteFileVersion::V2_1
585 | ConcreteFileVersion::V2_2
586 | ConcreteFileVersion::V2_3 => {}
587 }
588 Ok(())
589}
590
591#[derive(Debug, Clone, PartialEq)]
592pub struct BasePath {
593 pub id: u32,
594 pub name: Option<String>,
595 pub is_dataset_root: bool,
596 pub path: String,
598}
599
600impl BasePath {
601 pub fn new(id: u32, path: String, name: Option<String>, is_dataset_root: bool) -> Self {
610 Self {
611 id,
612 name,
613 is_dataset_root,
614 path,
615 }
616 }
617
618 pub fn extract_path(&self, registry: Arc<ObjectStoreRegistry>) -> Result<Path> {
622 ObjectStore::extract_path_from_uri(registry, &self.path)
623 }
624}
625
626impl DeepSizeOf for BasePath {
627 fn deep_size_of_children(&self, context: &mut lance_core::deepsize::Context) -> usize {
628 self.name.deep_size_of_children(context)
629 + self.path.deep_size_of_children(context) * 2
630 + size_of::<bool>()
631 }
632}
633
634#[derive(Debug, Clone, PartialEq, DeepSizeOf)]
635pub struct WriterVersion {
636 pub library: String,
637 pub version: String,
638 pub prerelease: Option<String>,
639 pub build_metadata: Option<String>,
640}
641
642#[derive(Debug, Clone, PartialEq, DeepSizeOf)]
643pub struct DataStorageFormat {
644 pub file_format: String,
645 pub version: ConcreteFileVersion,
646}
647
648const LANCE_FORMAT_NAME: &str = "lance";
649
650impl DataStorageFormat {
651 pub fn new(version: ConcreteFileVersion) -> Self {
652 Self {
653 file_format: LANCE_FORMAT_NAME.to_string(),
654 version,
655 }
656 }
657
658 pub fn lance_file_format(&self) -> ConcreteFileVersion {
660 self.version
661 }
662
663 pub fn lance_file_version(&self) -> Result<LanceFileVersion> {
665 Ok(self.version.into())
666 }
667}
668
669impl Default for DataStorageFormat {
670 fn default() -> Self {
671 Self::new(ConcreteFileVersion::from(LanceFileVersion::Stable))
672 }
673}
674
675impl TryFrom<pb::manifest::DataStorageFormat> for DataStorageFormat {
676 type Error = Error;
677
678 fn try_from(pb: pb::manifest::DataStorageFormat) -> Result<Self> {
679 Ok(Self {
680 file_format: pb.file_format,
681 version: ConcreteFileVersion::from_manifest_string(&pb.version)?,
682 })
683 }
684}
685
686#[derive(Debug, Clone, Copy, PartialEq, Eq)]
687pub enum VersionPart {
688 Major,
689 Minor,
690 Patch,
691}
692
693fn bump_version(version: &mut semver::Version, part: VersionPart) {
694 match part {
695 VersionPart::Major => {
696 version.major += 1;
697 version.minor = 0;
698 version.patch = 0;
699 }
700 VersionPart::Minor => {
701 version.minor += 1;
702 version.patch = 0;
703 }
704 VersionPart::Patch => {
705 version.patch += 1;
706 }
707 }
708}
709
710impl WriterVersion {
711 fn split_version(full_version: &str) -> Option<(String, Option<String>, Option<String>)> {
721 let mut parsed = semver::Version::parse(full_version).ok()?;
722
723 let prerelease = if parsed.pre.is_empty() {
724 None
725 } else {
726 Some(parsed.pre.to_string())
727 };
728
729 let build_metadata = if parsed.build.is_empty() {
730 None
731 } else {
732 Some(parsed.build.to_string())
733 };
734
735 parsed.pre = semver::Prerelease::EMPTY;
737 parsed.build = semver::BuildMetadata::EMPTY;
738 Some((parsed.to_string(), prerelease, build_metadata))
739 }
740
741 #[deprecated(note = "Use `lance_lib_version()` instead")]
744 pub fn semver(&self) -> Option<(u32, u32, u32, Option<&str>)> {
745 let (version_part, tag) = if let Some(dash_idx) = self.version.find('-') {
747 (
748 &self.version[..dash_idx],
749 Some(&self.version[dash_idx + 1..]),
750 )
751 } else {
752 (self.version.as_str(), None)
753 };
754
755 let mut parts = version_part.split('.');
756 let major = parts.next().unwrap_or("0").parse().ok()?;
757 let minor = parts.next().unwrap_or("0").parse().ok()?;
758 let patch = parts.next().unwrap_or("0").parse().ok()?;
759
760 Some((major, minor, patch, tag))
761 }
762
763 pub fn lance_lib_version(&self) -> Option<semver::Version> {
771 if self.library != "lance" {
772 return None;
773 }
774
775 let mut version = semver::Version::parse(&self.version).ok()?;
776
777 if let Some(ref prerelease) = self.prerelease {
778 version.pre = semver::Prerelease::new(prerelease).ok()?;
779 }
780
781 if let Some(ref build_metadata) = self.build_metadata {
782 version.build = semver::BuildMetadata::new(build_metadata).ok()?;
783 }
784
785 Some(version)
786 }
787
788 #[deprecated(
789 note = "Use `lance_lib_version()` instead, which safely checks the library field and returns Option"
790 )]
791 #[allow(deprecated)]
792 pub fn semver_or_panic(&self) -> (u32, u32, u32, Option<&str>) {
793 self.semver()
794 .unwrap_or_else(|| panic!("Invalid writer version: {}", self.version))
795 }
796
797 #[deprecated(note = "Use `lance_lib_version()` and its `older_than` method instead.")]
803 pub fn older_than(&self, major: u32, minor: u32, patch: u32) -> bool {
804 let version = self
805 .lance_lib_version()
806 .expect("Not lance library or invalid version");
807 let other = semver::Version {
808 major: major.into(),
809 minor: minor.into(),
810 patch: patch.into(),
811 pre: semver::Prerelease::EMPTY,
812 build: semver::BuildMetadata::EMPTY,
813 };
814 version < other
815 }
816
817 #[deprecated(note = "This is meant for testing and will be made private in future version.")]
818 pub fn bump(&self, part: VersionPart, keep_tag: bool) -> Self {
819 let mut version = self.lance_lib_version().expect("Should be lance version");
820 bump_version(&mut version, part);
821 if !keep_tag {
822 version.pre = semver::Prerelease::EMPTY;
823 }
824 let (clean_version, prerelease, build_metadata) = Self::split_version(&version.to_string())
825 .expect("Bumped version should be valid semver");
826 Self {
827 library: self.library.clone(),
828 version: clean_version,
829 prerelease,
830 build_metadata,
831 }
832 }
833}
834
835impl Default for WriterVersion {
836 #[cfg(not(test))]
837 fn default() -> Self {
838 let full_version = env!("CARGO_PKG_VERSION");
839 let (version, prerelease, build_metadata) =
840 Self::split_version(full_version).expect("CARGO_PKG_VERSION should be valid semver");
841 Self {
842 library: "lance".to_string(),
843 version,
844 prerelease,
845 build_metadata,
846 }
847 }
848
849 #[cfg(test)]
851 #[allow(deprecated)]
852 fn default() -> Self {
853 let full_version = env!("CARGO_PKG_VERSION");
854 let (version, prerelease, build_metadata) =
855 Self::split_version(full_version).expect("CARGO_PKG_VERSION should be valid semver");
856 Self {
857 library: "lance".to_string(),
858 version,
859 prerelease,
860 build_metadata,
861 }
862 .bump(VersionPart::Patch, true)
863 }
864}
865
866impl ProtoStruct for Manifest {
867 type Proto = pb::Manifest;
868}
869
870impl From<pb::BasePath> for BasePath {
871 fn from(p: pb::BasePath) -> Self {
872 Self::new(p.id, p.path, p.name, p.is_dataset_root)
873 }
874}
875
876impl From<BasePath> for pb::BasePath {
877 fn from(p: BasePath) -> Self {
878 Self {
879 id: p.id,
880 name: p.name,
881 is_dataset_root: p.is_dataset_root,
882 path: p.path,
883 }
884 }
885}
886
887impl TryFrom<pb::Manifest> for Manifest {
888 type Error = Error;
889
890 fn try_from(p: pb::Manifest) -> Result<Self> {
891 let timestamp_nanos = p.timestamp.map(|ts| {
892 let sec = ts.seconds as u128 * 1e9 as u128;
893 let nanos = ts.nanos as u128;
894 sec + nanos
895 });
896 let writer_version = match p.writer_version {
898 Some(pb::manifest::WriterVersion {
899 library,
900 version,
901 prerelease,
902 build_metadata,
903 }) => Some(WriterVersion {
904 library,
905 version,
906 prerelease,
907 build_metadata,
908 }),
909 _ => None,
910 };
911 let mut interner = DataFileFieldInterner::default();
912 let fragments = Arc::new(
913 p.fragments
914 .into_iter()
915 .map(|f| interner.intern_fragment(f))
916 .collect::<Result<Vec<_>>>()?,
917 );
918 let fragment_offsets = compute_fragment_offsets(fragments.as_slice());
919 let fields_with_meta = FieldsWithMeta {
920 fields: Fields(p.fields),
921 metadata: p.schema_metadata,
922 };
923
924 if FLAG_STABLE_ROW_IDS & p.reader_feature_flags != 0
925 && !fragments.iter().all(|frag| frag.row_id_meta.is_some())
926 {
927 return Err(Error::internal("All fragments must have row ids"));
928 }
929
930 let data_storage_format = match p.data_format {
931 None => {
932 if let Some(inferred_version) = Fragment::try_infer_version(fragments.as_ref())? {
933 DataStorageFormat::new(inferred_version)
935 } else {
936 if has_deprecated_v2_feature_flag(p.writer_feature_flags) {
938 DataStorageFormat::new(ConcreteFileVersion::from(LanceFileVersion::Stable))
939 } else {
940 DataStorageFormat::new(ConcreteFileVersion::V1)
941 }
942 }
943 }
944 Some(format) => DataStorageFormat::try_from(format)?,
945 };
946
947 let schema = Schema::try_from(fields_with_meta)?;
948
949 Ok(Self {
950 schema,
951 version: p.version,
952 branch: p.branch,
953 writer_version,
954 version_aux_data: p.version_aux_data as usize,
955 index_section: p.index_section.map(|i| i as usize),
956 timestamp_nanos: timestamp_nanos.unwrap_or(0),
957 tag: if p.tag.is_empty() { None } else { Some(p.tag) },
958 reader_feature_flags: p.reader_feature_flags,
959 writer_feature_flags: p.writer_feature_flags,
960 max_fragment_id: p.max_fragment_id,
961 fragments,
962 transaction_file: if p.transaction_file.is_empty() {
963 None
964 } else {
965 Some(p.transaction_file)
966 },
967 transaction_section: p.transaction_section.map(|i| i as usize),
968 fragment_offsets,
969 next_row_id: p.next_row_id,
970 data_storage_format,
971 config: p.config,
972 table_metadata: p.table_metadata,
973 base_paths: p
974 .base_paths
975 .iter()
976 .map(|item| (item.id, item.clone().into()))
977 .collect(),
978 })
979 }
980}
981
982impl From<&Manifest> for pb::Manifest {
983 fn from(m: &Manifest) -> Self {
984 let timestamp_nanos = if m.timestamp_nanos == 0 {
985 None
986 } else {
987 let nanos = m.timestamp_nanos % 1e9 as u128;
988 let seconds = ((m.timestamp_nanos - nanos) / 1e9 as u128) as i64;
989 Some(Timestamp {
990 seconds,
991 nanos: nanos as i32,
992 })
993 };
994 let fields_with_meta: FieldsWithMeta = (&m.schema).into();
995 Self {
996 fields: fields_with_meta.fields.0,
997 schema_metadata: m
998 .schema
999 .metadata
1000 .iter()
1001 .map(|(k, v)| (k.clone(), v.as_bytes().to_vec()))
1002 .collect(),
1003 version: m.version,
1004 branch: m.branch.clone(),
1005 writer_version: m
1006 .writer_version
1007 .as_ref()
1008 .map(|wv| pb::manifest::WriterVersion {
1009 library: wv.library.clone(),
1010 version: wv.version.clone(),
1011 prerelease: wv.prerelease.clone(),
1012 build_metadata: wv.build_metadata.clone(),
1013 }),
1014 fragments: m.fragments.iter().map(pb::DataFragment::from).collect(),
1015 table_metadata: m.table_metadata.clone(),
1016 version_aux_data: m.version_aux_data as u64,
1017 index_section: m.index_section.map(|i| i as u64),
1018 timestamp: timestamp_nanos,
1019 tag: m.tag.clone().unwrap_or_default(),
1020 reader_feature_flags: m.reader_feature_flags,
1021 writer_feature_flags: m.writer_feature_flags,
1022 max_fragment_id: m.max_fragment_id,
1023 transaction_file: m.transaction_file.clone().unwrap_or_default(),
1024 next_row_id: m.next_row_id,
1025 data_format: Some(pb::manifest::DataStorageFormat {
1026 file_format: m.data_storage_format.file_format.clone(),
1027 version: m
1028 .data_storage_format
1029 .version
1030 .to_manifest_string()
1031 .to_string(),
1032 }),
1033 config: m.config.clone(),
1034 base_paths: m
1035 .base_paths
1036 .values()
1037 .map(|base_path| pb::BasePath {
1038 id: base_path.id,
1039 name: base_path.name.clone(),
1040 is_dataset_root: base_path.is_dataset_root,
1041 path: base_path.path.clone(),
1042 })
1043 .collect(),
1044 transaction_section: m.transaction_section.map(|i| i as u64),
1045 }
1046 }
1047}
1048
1049#[async_trait]
1050pub trait SelfDescribingFileReader {
1051 async fn try_new_self_described(
1059 object_store: &ObjectStore,
1060 path: &Path,
1061 cache: Option<&LanceCache>,
1062 ) -> Result<Self>
1063 where
1064 Self: Sized,
1065 {
1066 let reader = object_store.open(path).await?;
1067 Self::try_new_self_described_from_reader(reader.into(), cache).await
1068 }
1069
1070 async fn try_new_self_described_from_reader(
1071 reader: Arc<dyn Reader>,
1072 cache: Option<&LanceCache>,
1073 ) -> Result<Self>
1074 where
1075 Self: Sized;
1076}
1077
1078#[async_trait]
1079impl SelfDescribingFileReader for V1FileReader {
1080 async fn try_new_self_described_from_reader(
1081 reader: Arc<dyn Reader>,
1082 cache: Option<&LanceCache>,
1083 ) -> Result<Self> {
1084 let metadata = Self::read_metadata(reader.as_ref(), cache).await?;
1085 let manifest_position = metadata.manifest_position.ok_or(Error::internal(format!(
1086 "Attempt to open file at {} as self-describing but it did not contain a manifest",
1087 reader.path(),
1088 )))?;
1089 let mut manifest: Manifest = read_struct(reader.as_ref(), manifest_position).await?;
1090 populate_manifest_schema_dictionaries(&mut manifest, reader.as_ref()).await?;
1091 let schema = manifest.schema;
1092 let max_field_id = schema.max_field_id().unwrap_or_default();
1093 Self::try_new_from_reader(
1094 reader.path(),
1095 reader.clone(),
1096 Some(metadata),
1097 schema,
1098 0,
1099 0,
1100 max_field_id,
1101 cache,
1102 )
1103 .await
1104 }
1105}
1106
1107#[cfg(test)]
1108mod tests {
1109 use crate::feature_flags::FLAG_USE_V2_FORMAT_DEPRECATED;
1110 use crate::format::{DataFile, DeletionFile, DeletionFileType};
1111 use std::num::NonZero;
1112
1113 use super::*;
1114
1115 use arrow_schema::{Field as ArrowField, Schema as ArrowSchema};
1116 use lance_core::datatypes::Field;
1117
1118 #[test]
1119 fn old_empty_manifest_recovers_v1_or_current_stable() {
1120 let old_manifest = pb::Manifest {
1121 data_format: None,
1122 ..Default::default()
1123 };
1124 let recovered_v1 = Manifest::try_from(old_manifest.clone()).unwrap();
1125 assert_eq!(
1126 recovered_v1.data_storage_format.lance_file_format(),
1127 ConcreteFileVersion::V1
1128 );
1129
1130 let recovered_stable = Manifest::try_from(pb::Manifest {
1131 writer_feature_flags: FLAG_USE_V2_FORMAT_DEPRECATED,
1132 ..old_manifest
1133 })
1134 .unwrap();
1135 assert_eq!(
1136 recovered_stable.data_storage_format.lance_file_format(),
1137 ConcreteFileVersion::from(LanceFileVersion::Stable)
1138 );
1139 }
1140
1141 #[test]
1142 fn manifest_persistence_rejects_selectors_and_public_aliases() {
1143 for version in ["stable", "next", "legacy", "0.3"] {
1144 let manifest = pb::Manifest {
1145 data_format: Some(pb::manifest::DataStorageFormat {
1146 file_format: LANCE_FORMAT_NAME.to_string(),
1147 version: version.to_string(),
1148 }),
1149 ..Default::default()
1150 };
1151 assert!(Manifest::try_from(manifest).is_err(), "accepted {version}");
1152 }
1153 }
1154
1155 #[test]
1156 fn manifest_codec_writes_canonical_exact_string() {
1157 let manifest = Manifest::new(
1158 Schema::default(),
1159 Arc::new(Vec::new()),
1160 DataStorageFormat::new(ConcreteFileVersion::V2_0),
1161 HashMap::new(),
1162 );
1163 let encoded = pb::Manifest::from(&manifest);
1164 assert_eq!(encoded.data_format.unwrap().version, "2.0");
1165 }
1166
1167 #[test]
1168 fn missing_format_infers_exact_version_and_rejects_mixed_files() {
1169 let v2_0 = Fragment::new(0).with_file(
1170 "v2_0.lance",
1171 vec![0],
1172 vec![0],
1173 ConcreteFileVersion::V2_0,
1174 None,
1175 );
1176 let manifest = Manifest::new(
1177 Schema::default(),
1178 Arc::new(vec![v2_0.clone()]),
1179 DataStorageFormat::new(ConcreteFileVersion::V1),
1180 HashMap::new(),
1181 );
1182 let mut encoded = pb::Manifest::from(&manifest);
1183 encoded.data_format = None;
1184 let recovered = Manifest::try_from(encoded).unwrap();
1185 assert_eq!(
1186 recovered.data_storage_format.lance_file_format(),
1187 ConcreteFileVersion::V2_0
1188 );
1189
1190 let v2_1 = Fragment::new(1).with_file(
1191 "v2_1.lance",
1192 vec![0],
1193 vec![0],
1194 ConcreteFileVersion::V2_1,
1195 None,
1196 );
1197 let mixed_manifest = Manifest::new(
1198 Schema::default(),
1199 Arc::new(vec![v2_0, v2_1]),
1200 DataStorageFormat::new(ConcreteFileVersion::V2_0),
1201 HashMap::new(),
1202 );
1203 let mut encoded = pb::Manifest::from(&mixed_manifest);
1204 encoded.data_format = None;
1205 let error = Manifest::try_from(encoded).unwrap_err();
1206 assert!(
1207 error
1208 .to_string()
1209 .contains("All data files must have the same version")
1210 );
1211 }
1212
1213 #[test]
1214 fn test_writer_version() {
1215 let wv = WriterVersion::default();
1216 assert_eq!(wv.library, "lance");
1217
1218 let cargo_version = env!("CARGO_PKG_VERSION");
1220 let expected_tag = if cargo_version.contains('-') {
1221 Some(cargo_version.split('-').nth(1).unwrap())
1222 } else {
1223 None
1224 };
1225
1226 let version_parts: Vec<&str> = wv.version.split('.').collect();
1228 assert_eq!(
1229 version_parts.len(),
1230 3,
1231 "Version should be major.minor.patch"
1232 );
1233 assert!(
1234 !wv.version.contains('-'),
1235 "Version field should not contain prerelease"
1236 );
1237
1238 assert_eq!(wv.prerelease.as_deref(), expected_tag);
1240 assert_eq!(wv.build_metadata, None);
1242
1243 let version = wv.lance_lib_version().unwrap();
1245 assert_eq!(
1246 version.major,
1247 env!("CARGO_PKG_VERSION_MAJOR").parse::<u64>().unwrap()
1248 );
1249 assert_eq!(
1250 version.minor,
1251 env!("CARGO_PKG_VERSION_MINOR").parse::<u64>().unwrap()
1252 );
1253 assert_eq!(
1254 version.patch,
1255 env!("CARGO_PKG_VERSION_PATCH").parse::<u64>().unwrap() + 1
1257 );
1258 assert_eq!(version.pre.as_str(), expected_tag.unwrap_or(""));
1259
1260 for part in &[VersionPart::Major, VersionPart::Minor, VersionPart::Patch] {
1261 let mut bumped_version = version.clone();
1262 bump_version(&mut bumped_version, *part);
1263 assert!(version < bumped_version);
1264 }
1265 }
1266
1267 #[test]
1268 fn test_writer_version_split() {
1269 let (version, prerelease, build_metadata) =
1271 WriterVersion::split_version("2.0.0-rc.1").unwrap();
1272 assert_eq!(version, "2.0.0");
1273 assert_eq!(prerelease, Some("rc.1".to_string()));
1274 assert_eq!(build_metadata, None);
1275
1276 let (version, prerelease, build_metadata) = WriterVersion::split_version("2.0.0").unwrap();
1278 assert_eq!(version, "2.0.0");
1279 assert_eq!(prerelease, None);
1280 assert_eq!(build_metadata, None);
1281
1282 let (version, prerelease, build_metadata) =
1284 WriterVersion::split_version("2.0.0-rc.1+build.123").unwrap();
1285 assert_eq!(version, "2.0.0");
1286 assert_eq!(prerelease, Some("rc.1".to_string()));
1287 assert_eq!(build_metadata, Some("build.123".to_string()));
1288
1289 let (version, prerelease, build_metadata) =
1291 WriterVersion::split_version("2.0.0+build.123").unwrap();
1292 assert_eq!(version, "2.0.0");
1293 assert_eq!(prerelease, None);
1294 assert_eq!(build_metadata, Some("build.123".to_string()));
1295
1296 assert!(WriterVersion::split_version("not-a-version").is_none());
1298 }
1299
1300 #[test]
1301 fn test_writer_version_comparison_with_prerelease() {
1302 let v1 = WriterVersion {
1303 library: "lance".to_string(),
1304 version: "2.0.0".to_string(),
1305 prerelease: Some("rc.1".to_string()),
1306 build_metadata: None,
1307 };
1308
1309 let v2 = WriterVersion {
1310 library: "lance".to_string(),
1311 version: "2.0.0".to_string(),
1312 prerelease: None,
1313 build_metadata: None,
1314 };
1315
1316 let semver1 = v1.lance_lib_version().unwrap();
1317 let semver2 = v2.lance_lib_version().unwrap();
1318
1319 assert!(semver1 < semver2);
1321 }
1322
1323 #[test]
1324 fn test_writer_version_with_build_metadata() {
1325 let v = WriterVersion {
1326 library: "lance".to_string(),
1327 version: "2.0.0".to_string(),
1328 prerelease: Some("rc.1".to_string()),
1329 build_metadata: Some("build.123".to_string()),
1330 };
1331
1332 let semver = v.lance_lib_version().unwrap();
1333 assert_eq!(semver.to_string(), "2.0.0-rc.1+build.123");
1334 assert_eq!(semver.major, 2);
1335 assert_eq!(semver.minor, 0);
1336 assert_eq!(semver.patch, 0);
1337 assert_eq!(semver.pre.as_str(), "rc.1");
1338 assert_eq!(semver.build.as_str(), "build.123");
1339 }
1340
1341 #[test]
1342 fn test_writer_version_non_semver() {
1343 let v = WriterVersion {
1345 library: "lance".to_string(),
1346 version: "custom-build-v1".to_string(),
1347 prerelease: None,
1348 build_metadata: None,
1349 };
1350
1351 assert!(v.lance_lib_version().is_none());
1353
1354 assert_eq!(v.library, "lance");
1356 assert_eq!(v.version, "custom-build-v1");
1357 }
1358
1359 #[test]
1360 #[allow(deprecated)]
1361 fn test_older_than_with_prerelease() {
1362 let v_rc = WriterVersion {
1364 library: "lance".to_string(),
1365 version: "2.0.0".to_string(),
1366 prerelease: Some("rc.1".to_string()),
1367 build_metadata: None,
1368 };
1369
1370 assert!(v_rc.older_than(2, 0, 0));
1372
1373 assert!(v_rc.older_than(2, 0, 1));
1375
1376 assert!(!v_rc.older_than(1, 9, 9));
1378
1379 let v_release = WriterVersion {
1380 library: "lance".to_string(),
1381 version: "2.0.0".to_string(),
1382 prerelease: None,
1383 build_metadata: None,
1384 };
1385
1386 assert!(!v_release.older_than(2, 0, 0));
1388
1389 assert!(v_release.older_than(2, 0, 1));
1391 }
1392
1393 #[test]
1394 fn test_fragments_by_offset_range() {
1395 let arrow_schema = ArrowSchema::new(vec![ArrowField::new(
1396 "a",
1397 arrow_schema::DataType::Int64,
1398 false,
1399 )]);
1400 let schema = Schema::try_from(&arrow_schema).unwrap();
1401 let fragments = vec![
1402 Fragment::with_file_legacy(0, "path1", &schema, Some(10)),
1403 Fragment::with_file_legacy(1, "path2", &schema, Some(15)),
1404 Fragment::with_file_legacy(2, "path3", &schema, Some(20)),
1405 ];
1406 let manifest = Manifest::new(
1407 schema,
1408 Arc::new(fragments),
1409 DataStorageFormat::default(),
1410 HashMap::new(),
1411 );
1412
1413 let actual = manifest.fragments_by_offset_range(0..10);
1414 assert_eq!(actual.len(), 1);
1415 assert_eq!(actual[0].0, 0);
1416 assert_eq!(actual[0].1.id, 0);
1417
1418 let actual = manifest.fragments_by_offset_range(5..15);
1419 assert_eq!(actual.len(), 2);
1420 assert_eq!(actual[0].0, 0);
1421 assert_eq!(actual[0].1.id, 0);
1422 assert_eq!(actual[1].0, 10);
1423 assert_eq!(actual[1].1.id, 1);
1424
1425 let actual = manifest.fragments_by_offset_range(15..50);
1426 assert_eq!(actual.len(), 2);
1427 assert_eq!(actual[0].0, 10);
1428 assert_eq!(actual[0].1.id, 1);
1429 assert_eq!(actual[1].0, 25);
1430 assert_eq!(actual[1].1.id, 2);
1431
1432 let actual = manifest.fragments_by_offset_range(45..100);
1434 assert!(actual.is_empty());
1435
1436 assert!(manifest.fragments_by_offset_range(200..400).is_empty());
1437 }
1438
1439 #[test]
1440 fn test_max_field_id() {
1441 let mut field0 =
1443 Field::try_from(ArrowField::new("a", arrow_schema::DataType::Int64, false)).unwrap();
1444 field0.set_id(-1, &mut 0);
1445 let mut field2 =
1446 Field::try_from(ArrowField::new("b", arrow_schema::DataType::Int64, false)).unwrap();
1447 field2.set_id(-1, &mut 2);
1448
1449 let schema = Schema {
1450 fields: vec![field0, field2],
1451 metadata: Default::default(),
1452 };
1453 let fragments = vec![
1454 Fragment {
1455 id: 0,
1456 files: vec![DataFile::new_legacy_from_fields(
1457 "path1",
1458 vec![0, 1, 2],
1459 None,
1460 )],
1461 overlays: vec![],
1462 deletion_file: None,
1463 row_id_meta: None,
1464 physical_rows: None,
1465 created_at_version_meta: None,
1466 last_updated_at_version_meta: None,
1467 },
1468 Fragment {
1469 id: 1,
1470 files: vec![
1471 DataFile::new_legacy_from_fields("path2", vec![0, 1, 43], None),
1472 DataFile::new_legacy_from_fields("path3", vec![2], None),
1473 ],
1474 overlays: vec![],
1475 deletion_file: None,
1476 row_id_meta: None,
1477 physical_rows: None,
1478 created_at_version_meta: None,
1479 last_updated_at_version_meta: None,
1480 },
1481 ];
1482
1483 let manifest = Manifest::new(
1484 schema,
1485 Arc::new(fragments),
1486 DataStorageFormat::default(),
1487 HashMap::new(),
1488 );
1489
1490 assert_eq!(manifest.max_field_id(), 43);
1491 }
1492
1493 #[test]
1494 fn test_config() {
1495 let arrow_schema = ArrowSchema::new(vec![ArrowField::new(
1496 "a",
1497 arrow_schema::DataType::Int64,
1498 false,
1499 )]);
1500 let schema = Schema::try_from(&arrow_schema).unwrap();
1501 let fragments = vec![
1502 Fragment::with_file_legacy(0, "path1", &schema, Some(10)),
1503 Fragment::with_file_legacy(1, "path2", &schema, Some(15)),
1504 Fragment::with_file_legacy(2, "path3", &schema, Some(20)),
1505 ];
1506 let mut manifest = Manifest::new(
1507 schema,
1508 Arc::new(fragments),
1509 DataStorageFormat::default(),
1510 HashMap::new(),
1511 );
1512
1513 let mut config = manifest.config.clone();
1514 config.insert("lance.test".to_string(), "value".to_string());
1515 config.insert("other-key".to_string(), "other-value".to_string());
1516
1517 manifest.config_mut().extend(config.clone());
1518 assert_eq!(manifest.config, config.clone());
1519
1520 config.remove("other-key");
1521 manifest.config_mut().remove("other-key");
1522 assert_eq!(manifest.config, config);
1523 }
1524
1525 #[test]
1526 fn test_manifest_summary() {
1527 let arrow_schema = ArrowSchema::new(vec![
1529 ArrowField::new("id", arrow_schema::DataType::Int64, false),
1530 ArrowField::new("name", arrow_schema::DataType::Utf8, true),
1531 ]);
1532 let schema = Schema::try_from(&arrow_schema).unwrap();
1533
1534 let empty_manifest = Manifest::new(
1535 schema.clone(),
1536 Arc::new(vec![]),
1537 DataStorageFormat::default(),
1538 HashMap::new(),
1539 );
1540
1541 let empty_summary = empty_manifest.summary();
1542 assert_eq!(empty_summary.total_rows, 0);
1543 assert_eq!(empty_summary.total_files_size, 0);
1544 assert_eq!(empty_summary.total_fragments, 0);
1545 assert_eq!(empty_summary.total_data_files, 0);
1546 assert_eq!(empty_summary.total_deletion_file_rows, 0);
1547 assert_eq!(empty_summary.total_data_file_rows, 0);
1548 assert_eq!(empty_summary.total_deletion_files, 0);
1549
1550 let empty_fragments = vec![
1552 Fragment::with_file_legacy(0, "empty_file1.lance", &schema, Some(0)),
1553 Fragment::with_file_legacy(1, "empty_file2.lance", &schema, Some(0)),
1554 ];
1555
1556 let empty_files_manifest = Manifest::new(
1557 schema.clone(),
1558 Arc::new(empty_fragments),
1559 DataStorageFormat::default(),
1560 HashMap::new(),
1561 );
1562
1563 let empty_files_summary = empty_files_manifest.summary();
1564 assert_eq!(empty_files_summary.total_rows, 0);
1565 assert_eq!(empty_files_summary.total_files_size, 0);
1566 assert_eq!(empty_files_summary.total_fragments, 2);
1567 assert_eq!(empty_files_summary.total_data_files, 2);
1568 assert_eq!(empty_files_summary.total_deletion_file_rows, 0);
1569 assert_eq!(empty_files_summary.total_data_file_rows, 0);
1570 assert_eq!(empty_files_summary.total_deletion_files, 0);
1571
1572 let real_fragments = vec![
1574 Fragment::with_file_legacy(0, "data_file1.lance", &schema, Some(100)),
1575 Fragment::with_file_legacy(1, "data_file2.lance", &schema, Some(250)),
1576 Fragment::with_file_legacy(2, "data_file3.lance", &schema, Some(75)),
1577 ];
1578
1579 let real_data_manifest = Manifest::new(
1580 schema.clone(),
1581 Arc::new(real_fragments),
1582 DataStorageFormat::default(),
1583 HashMap::new(),
1584 );
1585
1586 let real_data_summary = real_data_manifest.summary();
1587 assert_eq!(real_data_summary.total_rows, 425); assert_eq!(real_data_summary.total_files_size, 0); assert_eq!(real_data_summary.total_fragments, 3);
1590 assert_eq!(real_data_summary.total_data_files, 3);
1591 assert_eq!(real_data_summary.total_deletion_file_rows, 0);
1592 assert_eq!(real_data_summary.total_data_file_rows, 425);
1593 assert_eq!(real_data_summary.total_deletion_files, 0);
1594
1595 let mut fragment_with_deletion = Fragment::new(0)
1597 .with_file(
1598 "data_with_deletion.lance",
1599 vec![0, 1],
1600 vec![0, 1],
1601 ConcreteFileVersion::from(LanceFileVersion::Stable),
1602 NonZero::new(1000),
1603 )
1604 .with_physical_rows(50);
1605 fragment_with_deletion.deletion_file = Some(DeletionFile {
1606 read_version: 123,
1607 id: 456,
1608 file_type: DeletionFileType::Array,
1609 num_deleted_rows: Some(10),
1610 base_id: None,
1611 });
1612
1613 let manifest_with_deletion = Manifest::new(
1614 schema,
1615 Arc::new(vec![fragment_with_deletion]),
1616 DataStorageFormat::default(),
1617 HashMap::new(),
1618 );
1619
1620 let deletion_summary = manifest_with_deletion.summary();
1621 assert_eq!(deletion_summary.total_rows, 40); assert_eq!(deletion_summary.total_files_size, 1000);
1623 assert_eq!(deletion_summary.total_fragments, 1);
1624 assert_eq!(deletion_summary.total_data_files, 1);
1625 assert_eq!(deletion_summary.total_deletion_file_rows, 10);
1626 assert_eq!(deletion_summary.total_data_file_rows, 50);
1627 assert_eq!(deletion_summary.total_deletion_files, 1);
1628
1629 let stats_map: BTreeMap<String, String> = deletion_summary.into();
1631 assert_eq!(stats_map.len(), 7)
1632 }
1633}