1use std::cmp::Ordering;
22use std::collections::HashMap;
23use std::fmt::{Display, Formatter};
24use std::hash::Hash;
25use std::sync::Arc;
26
27use _serde::TableMetadataEnum;
28use chrono::{DateTime, Utc};
29use serde::{Deserialize, Serialize};
30use serde_repr::{Deserialize_repr, Serialize_repr};
31use uuid::Uuid;
32
33use super::snapshot::SnapshotReference;
34pub use super::table_metadata_builder::{TableMetadataBuildResult, TableMetadataBuilder};
35use super::{
36 DEFAULT_PARTITION_SPEC_ID, PartitionSpecRef, PartitionStatisticsFile, Schema, SchemaId,
37 SchemaRef, SnapshotRef, SnapshotRetention, SortOrder, SortOrderRef, StatisticsFile, StructType,
38 TableProperties, Transform,
39};
40use crate::catalog::{METADATA_FOLDER_NAME, MetadataLocation};
41use crate::compression::CompressionCodec;
42use crate::error::{Result, invalid_data, timestamp_ms_to_utc};
43use crate::io::FileIO;
44use crate::partitioning::compute_unified_partition_type;
45use crate::spec::EncryptedKey;
46use crate::{Error, ErrorKind};
47
48static MAIN_BRANCH: &str = "main";
49pub(crate) static ONE_MINUTE_MS: i64 = 60_000;
50
51pub(crate) static EMPTY_SNAPSHOT_ID: i64 = -1;
55pub(crate) static INITIAL_SEQUENCE_NUMBER: i64 = 0;
56
57pub const INITIAL_ROW_ID: u64 = 0;
59pub const MIN_FORMAT_VERSION_ROW_LINEAGE: FormatVersion = FormatVersion::V3;
61pub type TableMetadataRef = Arc<TableMetadata>;
63
64#[derive(Debug, PartialEq, Deserialize, Eq, Clone)]
65#[serde(try_from = "TableMetadataEnum")]
66pub struct TableMetadata {
71 pub(crate) format_version: FormatVersion,
73 pub(crate) table_uuid: Uuid,
75 pub(crate) location: String,
77 pub(crate) last_sequence_number: i64,
79 pub(crate) last_updated_ms: i64,
81 pub(crate) last_column_id: i32,
83 pub(crate) schemas: HashMap<i32, SchemaRef>,
85 pub(crate) current_schema_id: i32,
87 pub(crate) partition_specs: HashMap<i32, PartitionSpecRef>,
89 pub(crate) default_spec: PartitionSpecRef,
91 pub(crate) default_partition_type: StructType,
93 pub(crate) last_partition_id: i32,
95 pub(crate) properties: HashMap<String, String>,
99 pub(crate) current_snapshot_id: Option<i64>,
102 pub(crate) snapshots: HashMap<i64, SnapshotRef>,
107 pub(crate) snapshot_log: Vec<SnapshotLog>,
114
115 pub(crate) metadata_log: Vec<MetadataLog>,
122
123 pub(crate) sort_orders: HashMap<i64, SortOrderRef>,
125 pub(crate) default_sort_order_id: i64,
129 pub(crate) refs: HashMap<String, SnapshotReference>,
134 pub(crate) statistics: HashMap<i64, StatisticsFile>,
136 pub(crate) partition_statistics: HashMap<i64, PartitionStatisticsFile>,
138 pub(crate) encryption_keys: HashMap<String, EncryptedKey>,
140 pub(crate) next_row_id: u64,
142}
143
144impl TableMetadata {
145 #[must_use]
152 pub fn into_builder(self, current_file_location: Option<String>) -> TableMetadataBuilder {
153 TableMetadataBuilder::new_from_metadata(self, current_file_location)
154 }
155
156 #[inline]
158 pub(crate) fn partition_name_exists(&self, name: &str) -> bool {
159 self.partition_specs
160 .values()
161 .any(|spec| spec.fields().iter().any(|pf| pf.name == name))
162 }
163
164 #[inline]
166 pub(crate) fn name_exists_in_any_schema(&self, name: &str) -> bool {
167 self.schemas
168 .values()
169 .any(|schema| schema.field_by_name(name).is_some())
170 }
171
172 #[inline]
174 pub fn format_version(&self) -> FormatVersion {
175 self.format_version
176 }
177
178 #[inline]
180 pub fn uuid(&self) -> Uuid {
181 self.table_uuid
182 }
183
184 #[inline]
186 pub fn location(&self) -> &str {
187 self.location.as_str()
188 }
189
190 #[inline]
192 pub fn last_sequence_number(&self) -> i64 {
193 self.last_sequence_number
194 }
195
196 #[inline]
201 pub fn next_sequence_number(&self) -> i64 {
202 match self.format_version {
203 FormatVersion::V1 => INITIAL_SEQUENCE_NUMBER,
204 _ => self.last_sequence_number + 1,
205 }
206 }
207
208 #[inline]
210 pub fn last_column_id(&self) -> i32 {
211 self.last_column_id
212 }
213
214 #[inline]
216 pub fn last_partition_id(&self) -> i32 {
217 self.last_partition_id
218 }
219
220 #[inline]
222 pub fn last_updated_timestamp(&self) -> Result<DateTime<Utc>> {
223 timestamp_ms_to_utc(self.last_updated_ms)
224 }
225
226 #[inline]
228 pub fn last_updated_ms(&self) -> i64 {
229 self.last_updated_ms
230 }
231
232 #[inline]
234 pub fn schemas_iter(&self) -> impl ExactSizeIterator<Item = &SchemaRef> {
235 self.schemas.values()
236 }
237
238 #[inline]
240 pub fn schema_by_id(&self, schema_id: SchemaId) -> Option<&SchemaRef> {
241 self.schemas.get(&schema_id)
242 }
243
244 #[inline]
246 pub fn current_schema(&self) -> &SchemaRef {
247 self.schema_by_id(self.current_schema_id)
248 .expect("Current schema id set, but not found in table metadata")
249 }
250
251 #[inline]
253 pub fn current_schema_id(&self) -> SchemaId {
254 self.current_schema_id
255 }
256
257 #[inline]
259 pub fn partition_specs_iter(&self) -> impl ExactSizeIterator<Item = &PartitionSpecRef> {
260 self.partition_specs.values()
261 }
262
263 #[inline]
265 pub fn partition_spec_by_id(&self, spec_id: i32) -> Option<&PartitionSpecRef> {
266 self.partition_specs.get(&spec_id)
267 }
268
269 #[inline]
271 pub fn default_partition_spec(&self) -> &PartitionSpecRef {
272 &self.default_spec
273 }
274
275 #[inline]
277 pub fn default_partition_type(&self) -> &StructType {
278 &self.default_partition_type
279 }
280
281 pub fn unified_partition_type(&self, schema: &Schema) -> Result<StructType> {
288 compute_unified_partition_type(
289 self.partition_specs_iter().map(|spec| spec.as_ref()),
290 schema,
291 )
292 }
293
294 #[inline]
295 pub fn default_partition_spec_id(&self) -> i32 {
297 self.default_spec.spec_id()
298 }
299
300 #[inline]
302 pub fn snapshots(&self) -> impl ExactSizeIterator<Item = &SnapshotRef> {
303 self.snapshots.values()
304 }
305
306 #[inline]
308 pub fn snapshot_by_id(&self, snapshot_id: i64) -> Option<&SnapshotRef> {
309 self.snapshots.get(&snapshot_id)
310 }
311
312 #[inline]
314 pub fn history(&self) -> &[SnapshotLog] {
315 &self.snapshot_log
316 }
317
318 #[inline]
320 pub fn metadata_log(&self) -> &[MetadataLog] {
321 &self.metadata_log
322 }
323
324 #[inline]
326 pub fn current_snapshot(&self) -> Option<&SnapshotRef> {
327 self.current_snapshot_id.map(|s| {
328 self.snapshot_by_id(s)
329 .expect("Current snapshot id has been set, but doesn't exist in metadata")
330 })
331 }
332
333 #[inline]
335 pub fn current_snapshot_id(&self) -> Option<i64> {
336 self.current_snapshot_id
337 }
338
339 #[inline]
342 pub fn snapshot_for_ref(&self, ref_name: &str) -> Option<&SnapshotRef> {
343 self.refs.get(ref_name).map(|r| {
344 self.snapshot_by_id(r.snapshot_id)
345 .unwrap_or_else(|| panic!("Snapshot id of ref {ref_name} doesn't exist"))
346 })
347 }
348
349 #[inline]
351 pub fn sort_orders_iter(&self) -> impl ExactSizeIterator<Item = &SortOrderRef> {
352 self.sort_orders.values()
353 }
354
355 #[inline]
357 pub fn sort_order_by_id(&self, sort_order_id: i64) -> Option<&SortOrderRef> {
358 self.sort_orders.get(&sort_order_id)
359 }
360
361 #[inline]
363 pub fn default_sort_order(&self) -> &SortOrderRef {
364 self.sort_orders
365 .get(&self.default_sort_order_id)
366 .expect("Default order id has been set, but not found in table metadata!")
367 }
368
369 #[inline]
371 pub fn default_sort_order_id(&self) -> i64 {
372 self.default_sort_order_id
373 }
374
375 #[inline]
377 pub fn properties(&self) -> &HashMap<String, String> {
378 &self.properties
379 }
380
381 pub fn metadata_location(&self) -> Result<String> {
386 Ok(self
387 .table_properties()
388 .write_metadata_path()?
389 .unwrap_or_else(|| format!("{}/{}", self.location(), METADATA_FOLDER_NAME)))
390 }
391
392 pub fn metadata_compression_codec(&self) -> Result<CompressionCodec> {
401 self.table_properties().metadata_compression_codec()
402 }
403
404 #[inline]
406 pub fn table_properties(&self) -> TableProperties<'_> {
407 TableProperties::new(&self.properties)
408 }
409
410 #[inline]
412 pub fn statistics_iter(&self) -> impl ExactSizeIterator<Item = &StatisticsFile> {
413 self.statistics.values()
414 }
415
416 #[inline]
418 pub fn partition_statistics_iter(
419 &self,
420 ) -> impl ExactSizeIterator<Item = &PartitionStatisticsFile> {
421 self.partition_statistics.values()
422 }
423
424 #[inline]
426 pub fn statistics_for_snapshot(&self, snapshot_id: i64) -> Option<&StatisticsFile> {
427 self.statistics.get(&snapshot_id)
428 }
429
430 #[inline]
432 pub fn partition_statistics_for_snapshot(
433 &self,
434 snapshot_id: i64,
435 ) -> Option<&PartitionStatisticsFile> {
436 self.partition_statistics.get(&snapshot_id)
437 }
438
439 fn construct_refs(&mut self) {
440 if let Some(current_snapshot_id) = self.current_snapshot_id
441 && !self.refs.contains_key(MAIN_BRANCH)
442 {
443 self.refs
444 .insert(MAIN_BRANCH.to_string(), SnapshotReference {
445 snapshot_id: current_snapshot_id,
446 retention: SnapshotRetention::Branch {
447 min_snapshots_to_keep: None,
448 max_snapshot_age_ms: None,
449 max_ref_age_ms: None,
450 },
451 });
452 }
453 }
454
455 #[inline]
457 pub fn encryption_keys_iter(&self) -> impl ExactSizeIterator<Item = &EncryptedKey> {
458 self.encryption_keys.values()
459 }
460
461 #[inline]
463 pub fn encryption_key(&self, key_id: &str) -> Option<&EncryptedKey> {
464 self.encryption_keys.get(key_id)
465 }
466
467 #[inline]
469 pub fn next_row_id(&self) -> u64 {
470 self.next_row_id
471 }
472
473 pub async fn read_from(
475 file_io: &FileIO,
476 metadata_location: impl AsRef<str>,
477 ) -> Result<TableMetadata> {
478 let metadata_location = metadata_location.as_ref();
479 let input_file = file_io.new_input(metadata_location)?;
480 let metadata_content = input_file.read().await?;
481
482 let metadata = if metadata_content.len() > 2
484 && metadata_content[0] == 0x1F
485 && metadata_content[1] == 0x8B
486 {
487 let decompressed_data = CompressionCodec::gzip_default()
488 .decompress(metadata_content.to_vec())
489 .map_err(|e| {
490 invalid_data!("Trying to read compressed metadata file")
491 .with_context("file_path", metadata_location)
492 .with_source(e)
493 })?;
494 serde_json::from_slice(&decompressed_data)?
495 } else {
496 serde_json::from_slice(&metadata_content)?
497 };
498
499 Ok(metadata)
500 }
501
502 pub async fn write_to(
504 &self,
505 file_io: &FileIO,
506 metadata_location: &MetadataLocation,
507 ) -> Result<()> {
508 let json_data = serde_json::to_vec(self)?;
509
510 let codec = self.table_properties().metadata_compression_codec()?;
512
513 if codec != metadata_location.compression_codec() {
514 return Err(invalid_data!(
515 "Compression codec mismatch: metadata_location has {:?}, but table properties specify {:?}",
516 metadata_location.compression_codec(),
517 codec
518 ));
519 }
520
521 let data_to_write = match codec {
523 CompressionCodec::Gzip(_) => codec.compress(json_data)?,
524 CompressionCodec::None => json_data,
525 _ => {
526 return Err(invalid_data!(
527 "Unsupported metadata compression codec: {codec:?}"
528 ));
529 }
530 };
531
532 file_io
533 .new_output(metadata_location.to_string())?
534 .write(data_to_write.into())
535 .await
536 }
537
538 pub(super) fn try_normalize(&mut self) -> Result<&mut Self> {
546 self.validate_current_schema()?;
547 self.normalize_current_snapshot()?;
548 self.construct_refs();
549 self.validate_refs()?;
550 self.validate_chronological_snapshot_logs()?;
551 self.validate_chronological_metadata_logs()?;
552 self.location = self.location.trim_end_matches('/').to_string();
554 self.validate_snapshot_sequence_number()?;
555 self.validate_schema_format_compatibility()?;
556 self.try_normalize_partition_spec()?;
557 self.try_normalize_sort_order()?;
558 Ok(self)
559 }
560
561 fn try_normalize_partition_spec(&mut self) -> Result<()> {
563 for field in self.default_spec.fields() {
564 if field.transform != Transform::Void
567 && self.current_schema().field_by_id(field.source_id).is_none()
568 {
569 return Err(invalid_data!(
570 "Default partition spec {} references missing source field {} in current schema {}",
571 self.default_spec.spec_id(),
572 field.source_id,
573 self.current_schema_id
574 ));
575 }
576 }
577
578 if self
579 .partition_spec_by_id(self.default_spec.spec_id())
580 .is_none()
581 {
582 self.partition_specs.insert(
583 self.default_spec.spec_id(),
584 Arc::new(Arc::unwrap_or_clone(self.default_spec.clone())),
585 );
586 }
587
588 Ok(())
589 }
590
591 fn try_normalize_sort_order(&mut self) -> Result<()> {
593 if let Some(sort_order) = self.sort_order_by_id(SortOrder::UNSORTED_ORDER_ID)
595 && !sort_order.fields.is_empty()
596 {
597 return Err(Error::new(
598 ErrorKind::Unexpected,
599 format!(
600 "Sort order ID {} is reserved for unsorted order",
601 SortOrder::UNSORTED_ORDER_ID
602 ),
603 ));
604 }
605
606 if self.sort_order_by_id(self.default_sort_order_id).is_some() {
607 return Ok(());
608 }
609
610 if self.default_sort_order_id != SortOrder::UNSORTED_ORDER_ID {
611 return Err(invalid_data!(
612 "No sort order exists with the default sort order id {}.",
613 self.default_sort_order_id
614 ));
615 }
616
617 let sort_order = SortOrder::unsorted_order();
618 self.sort_orders
619 .insert(SortOrder::UNSORTED_ORDER_ID, Arc::new(sort_order));
620 Ok(())
621 }
622
623 fn validate_current_schema(&self) -> Result<()> {
625 if self.schema_by_id(self.current_schema_id).is_none() {
626 return Err(invalid_data!(
627 "No schema exists with the current schema id {}.",
628 self.current_schema_id
629 ));
630 }
631 Ok(())
632 }
633
634 fn normalize_current_snapshot(&mut self) -> Result<()> {
636 if let Some(current_snapshot_id) = self.current_snapshot_id {
637 if current_snapshot_id == EMPTY_SNAPSHOT_ID {
638 self.current_snapshot_id = None;
639 } else if self.snapshot_by_id(current_snapshot_id).is_none() {
640 return Err(invalid_data!(
641 "Snapshot for current snapshot id {current_snapshot_id} does not exist in the existing snapshots list"
642 ));
643 }
644 }
645 Ok(())
646 }
647
648 fn validate_refs(&self) -> Result<()> {
650 for (name, snapshot_ref) in self.refs.iter() {
651 if self.snapshot_by_id(snapshot_ref.snapshot_id).is_none() {
652 return Err(invalid_data!(
653 "Snapshot for reference {name} does not exist in the existing snapshots list"
654 ));
655 }
656 }
657
658 let main_ref = self.refs.get(MAIN_BRANCH);
659 if self.current_snapshot_id.is_some() {
660 if let Some(main_ref) = main_ref
661 && main_ref.snapshot_id != self.current_snapshot_id.unwrap_or_default()
662 {
663 return Err(invalid_data!(
664 "Current snapshot id does not match main branch ({:?} != {:?})",
665 self.current_snapshot_id.unwrap_or_default(),
666 main_ref.snapshot_id
667 ));
668 }
669 } else if main_ref.is_some() {
670 return Err(invalid_data!(
671 "Current snapshot is not set, but main branch exists"
672 ));
673 }
674
675 Ok(())
676 }
677
678 fn validate_snapshot_sequence_number(&self) -> Result<()> {
680 if self.format_version < FormatVersion::V2 && self.last_sequence_number != 0 {
681 return Err(invalid_data!(
682 "Last sequence number must be 0 in v1. Found {}",
683 self.last_sequence_number
684 ));
685 }
686
687 if self.format_version >= FormatVersion::V2
688 && let Some(snapshot) = self
689 .snapshots
690 .values()
691 .find(|snapshot| snapshot.sequence_number() > self.last_sequence_number)
692 {
693 return Err(invalid_data!(
694 "Invalid snapshot with id {} and sequence number {} greater than last sequence number {}",
695 snapshot.snapshot_id(),
696 snapshot.sequence_number(),
697 self.last_sequence_number
698 ));
699 }
700
701 Ok(())
702 }
703
704 fn validate_chronological_snapshot_logs(&self) -> Result<()> {
706 for window in self.snapshot_log.windows(2) {
707 let (prev, curr) = (&window[0], &window[1]);
708 if curr.timestamp_ms - prev.timestamp_ms < -ONE_MINUTE_MS {
711 return Err(invalid_data!("Expected sorted snapshot log entries"));
712 }
713 }
714
715 if let Some(last) = self.snapshot_log.last() {
716 if self.last_updated_ms - last.timestamp_ms < -ONE_MINUTE_MS {
719 return Err(invalid_data!(
720 "Invalid update timestamp {}: before last snapshot log entry at {}",
721 self.last_updated_ms,
722 last.timestamp_ms
723 ));
724 }
725 }
726 Ok(())
727 }
728
729 fn validate_chronological_metadata_logs(&self) -> Result<()> {
730 for window in self.metadata_log.windows(2) {
731 let (prev, curr) = (&window[0], &window[1]);
732 if curr.timestamp_ms - prev.timestamp_ms < -ONE_MINUTE_MS {
735 return Err(invalid_data!("Expected sorted metadata log entries"));
736 }
737 }
738
739 if let Some(last) = self.metadata_log.last() {
740 if self.last_updated_ms - last.timestamp_ms < -ONE_MINUTE_MS {
743 return Err(invalid_data!(
744 "Invalid update timestamp {}: before last metadata log entry at {}",
745 self.last_updated_ms,
746 last.timestamp_ms
747 ));
748 }
749 }
750
751 Ok(())
752 }
753
754 fn validate_schema_format_compatibility(&self) -> Result<()> {
757 self.current_schema()
758 .check_format_compatibility(self.format_version)
759 }
760}
761
762pub(super) mod _serde {
763 use std::borrow::BorrowMut;
764 use std::collections::HashMap;
769 use std::sync::Arc;
774
775 use serde::{Deserialize, Serialize};
776 use uuid::Uuid;
777
778 use super::{
779 DEFAULT_PARTITION_SPEC_ID, EMPTY_SNAPSHOT_ID, FormatVersion, MAIN_BRANCH, MetadataLog,
780 SnapshotLog, TableMetadata,
781 };
782 use crate::error::invalid_data;
783 use crate::spec::schema::_serde::{SchemaV1, SchemaV2};
784 use crate::spec::snapshot::_serde::{SnapshotV1, SnapshotV2, SnapshotV3};
785 use crate::spec::{
786 EncryptedKey, INITIAL_ROW_ID, PartitionField, PartitionSpec, PartitionSpecRef,
787 PartitionStatisticsFile, Schema, SchemaRef, Snapshot, SnapshotReference, SnapshotRetention,
788 SortOrder, StatisticsFile,
789 };
790 use crate::{Error, ErrorKind};
791
792 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
793 #[serde(untagged)]
794 pub(super) enum TableMetadataEnum {
795 V3(TableMetadataV3),
796 V2(TableMetadataV2),
797 V1(TableMetadataV1),
798 }
799
800 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
801 #[serde(rename_all = "kebab-case")]
802 pub(super) struct TableMetadataV3 {
804 pub format_version: VersionNumber<3>,
805 #[serde(flatten)]
806 pub shared: TableMetadataV2V3Shared,
807 pub next_row_id: u64,
808 #[serde(skip_serializing_if = "Option::is_none")]
809 pub encryption_keys: Option<Vec<EncryptedKey>>,
810 #[serde(skip_serializing_if = "Option::is_none")]
811 pub snapshots: Option<Vec<SnapshotV3>>,
812 }
813
814 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
815 #[serde(rename_all = "kebab-case")]
816 pub(super) struct TableMetadataV2V3Shared {
818 pub table_uuid: Uuid,
819 pub location: String,
820 pub last_sequence_number: i64,
821 pub last_updated_ms: i64,
822 pub last_column_id: i32,
823 pub schemas: Vec<SchemaV2>,
824 pub current_schema_id: i32,
825 pub partition_specs: Vec<PartitionSpec>,
826 pub default_spec_id: i32,
827 pub last_partition_id: i32,
828 #[serde(skip_serializing_if = "Option::is_none")]
829 pub properties: Option<HashMap<String, String>>,
830 #[serde(skip_serializing_if = "Option::is_none")]
831 pub current_snapshot_id: Option<i64>,
832 #[serde(skip_serializing_if = "Option::is_none")]
833 pub snapshot_log: Option<Vec<SnapshotLog>>,
834 #[serde(skip_serializing_if = "Option::is_none")]
835 pub metadata_log: Option<Vec<MetadataLog>>,
836 pub sort_orders: Vec<SortOrder>,
837 pub default_sort_order_id: i64,
838 #[serde(skip_serializing_if = "Option::is_none")]
839 pub refs: Option<HashMap<String, SnapshotReference>>,
840 #[serde(default, skip_serializing_if = "Vec::is_empty")]
841 pub statistics: Vec<StatisticsFile>,
842 #[serde(default, skip_serializing_if = "Vec::is_empty")]
843 pub partition_statistics: Vec<PartitionStatisticsFile>,
844 }
845
846 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
847 #[serde(rename_all = "kebab-case")]
848 pub(super) struct TableMetadataV2 {
850 pub format_version: VersionNumber<2>,
851 #[serde(flatten)]
852 pub shared: TableMetadataV2V3Shared,
853 #[serde(skip_serializing_if = "Option::is_none")]
854 pub snapshots: Option<Vec<SnapshotV2>>,
855 }
856
857 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
858 #[serde(rename_all = "kebab-case")]
859 pub(super) struct TableMetadataV1 {
861 pub format_version: VersionNumber<1>,
862 #[serde(skip_serializing_if = "Option::is_none")]
863 pub table_uuid: Option<Uuid>,
864 pub location: String,
865 pub last_updated_ms: i64,
866 pub last_column_id: i32,
867 pub schema: Option<SchemaV1>,
869 #[serde(skip_serializing_if = "Option::is_none")]
870 pub schemas: Option<Vec<SchemaV1>>,
871 #[serde(skip_serializing_if = "Option::is_none")]
872 pub current_schema_id: Option<i32>,
873 pub partition_spec: Option<Vec<PartitionField>>,
875 #[serde(skip_serializing_if = "Option::is_none")]
876 pub partition_specs: Option<Vec<PartitionSpec>>,
877 #[serde(skip_serializing_if = "Option::is_none")]
878 pub default_spec_id: Option<i32>,
879 #[serde(skip_serializing_if = "Option::is_none")]
880 pub last_partition_id: Option<i32>,
881 #[serde(skip_serializing_if = "Option::is_none")]
882 pub properties: Option<HashMap<String, String>>,
883 #[serde(skip_serializing_if = "Option::is_none")]
884 pub current_snapshot_id: Option<i64>,
885 #[serde(skip_serializing_if = "Option::is_none")]
886 pub snapshots: Option<Vec<SnapshotV1>>,
887 #[serde(skip_serializing_if = "Option::is_none")]
888 pub snapshot_log: Option<Vec<SnapshotLog>>,
889 #[serde(skip_serializing_if = "Option::is_none")]
890 pub metadata_log: Option<Vec<MetadataLog>>,
891 pub sort_orders: Option<Vec<SortOrder>>,
892 pub default_sort_order_id: Option<i64>,
893 #[serde(default, skip_serializing_if = "Vec::is_empty")]
894 pub statistics: Vec<StatisticsFile>,
895 #[serde(default, skip_serializing_if = "Vec::is_empty")]
896 pub partition_statistics: Vec<PartitionStatisticsFile>,
897 }
898
899 #[derive(Debug, PartialEq, Eq)]
901 pub(crate) struct VersionNumber<const V: u8>;
902
903 impl Serialize for TableMetadata {
904 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
905 where S: serde::Serializer {
906 let table_metadata_enum: TableMetadataEnum =
908 self.clone().try_into().map_err(serde::ser::Error::custom)?;
909
910 table_metadata_enum.serialize(serializer)
911 }
912 }
913
914 impl<const V: u8> Serialize for VersionNumber<V> {
915 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
916 where S: serde::Serializer {
917 serializer.serialize_u8(V)
918 }
919 }
920
921 impl<'de, const V: u8> Deserialize<'de> for VersionNumber<V> {
922 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
923 where D: serde::Deserializer<'de> {
924 let value = u8::deserialize(deserializer)?;
925 if value == V {
926 Ok(VersionNumber::<V>)
927 } else {
928 Err(serde::de::Error::custom("Invalid Version"))
929 }
930 }
931 }
932
933 impl TryFrom<TableMetadataEnum> for TableMetadata {
934 type Error = Error;
935 fn try_from(value: TableMetadataEnum) -> Result<Self, Error> {
936 match value {
937 TableMetadataEnum::V3(value) => value.try_into(),
938 TableMetadataEnum::V2(value) => value.try_into(),
939 TableMetadataEnum::V1(value) => value.try_into(),
940 }
941 }
942 }
943
944 impl TryFrom<TableMetadata> for TableMetadataEnum {
945 type Error = Error;
946 fn try_from(value: TableMetadata) -> Result<Self, Error> {
947 Ok(match value.format_version {
948 FormatVersion::V3 => TableMetadataEnum::V3(value.try_into()?),
949 FormatVersion::V2 => TableMetadataEnum::V2(value.into()),
950 FormatVersion::V1 => TableMetadataEnum::V1(value.try_into()?),
951 })
952 }
953 }
954
955 impl TryFrom<TableMetadataV3> for TableMetadata {
956 type Error = Error;
957 fn try_from(value: TableMetadataV3) -> Result<Self, Error> {
958 let TableMetadataV3 {
959 format_version: _,
960 shared: value,
961 next_row_id,
962 encryption_keys,
963 snapshots,
964 } = value;
965 let current_snapshot_id = if value.current_snapshot_id == Some(EMPTY_SNAPSHOT_ID) {
966 None
967 } else {
968 value.current_snapshot_id
969 };
970 let schemas = HashMap::from_iter(
971 value
972 .schemas
973 .into_iter()
974 .map(|schema| Ok((schema.schema_id, Arc::new(schema.try_into()?))))
975 .collect::<Result<Vec<_>, Error>>()?,
976 );
977
978 let current_schema: &SchemaRef =
979 schemas.get(&value.current_schema_id).ok_or_else(|| {
980 invalid_data!(
981 "No schema exists with the current schema id {}.",
982 value.current_schema_id
983 )
984 })?;
985 let partition_specs = HashMap::from_iter(
986 value
987 .partition_specs
988 .into_iter()
989 .map(|x| (x.spec_id(), Arc::new(x))),
990 );
991 let default_spec_id = value.default_spec_id;
992 let default_spec: PartitionSpecRef = partition_specs
993 .get(&value.default_spec_id)
994 .map(|spec| (**spec).clone())
995 .or_else(|| {
996 (DEFAULT_PARTITION_SPEC_ID == default_spec_id)
997 .then(PartitionSpec::unpartition_spec)
998 })
999 .ok_or_else(|| invalid_data!("Default partition spec {default_spec_id} not found"))?
1000 .into();
1001 let default_partition_type = default_spec.partition_type(current_schema)?;
1002
1003 let mut metadata = TableMetadata {
1004 format_version: FormatVersion::V3,
1005 table_uuid: value.table_uuid,
1006 location: value.location,
1007 last_sequence_number: value.last_sequence_number,
1008 last_updated_ms: value.last_updated_ms,
1009 last_column_id: value.last_column_id,
1010 current_schema_id: value.current_schema_id,
1011 schemas,
1012 partition_specs,
1013 default_partition_type,
1014 default_spec,
1015 last_partition_id: value.last_partition_id,
1016 properties: value.properties.unwrap_or_default(),
1017 current_snapshot_id,
1018 snapshots: snapshots
1019 .map(|snapshots| {
1020 HashMap::from_iter(
1021 snapshots
1022 .into_iter()
1023 .map(|x| (x.snapshot_id, Arc::new(x.into()))),
1024 )
1025 })
1026 .unwrap_or_default(),
1027 snapshot_log: value.snapshot_log.unwrap_or_default(),
1028 metadata_log: value.metadata_log.unwrap_or_default(),
1029 sort_orders: HashMap::from_iter(
1030 value
1031 .sort_orders
1032 .into_iter()
1033 .map(|x| (x.order_id, Arc::new(x))),
1034 ),
1035 default_sort_order_id: value.default_sort_order_id,
1036 refs: value.refs.unwrap_or_else(|| {
1037 if let Some(snapshot_id) = current_snapshot_id {
1038 HashMap::from_iter(vec![(MAIN_BRANCH.to_string(), SnapshotReference {
1039 snapshot_id,
1040 retention: SnapshotRetention::Branch {
1041 min_snapshots_to_keep: None,
1042 max_snapshot_age_ms: None,
1043 max_ref_age_ms: None,
1044 },
1045 })])
1046 } else {
1047 HashMap::new()
1048 }
1049 }),
1050 statistics: index_statistics(value.statistics),
1051 partition_statistics: index_partition_statistics(value.partition_statistics),
1052 encryption_keys: encryption_keys
1053 .map(|keys| {
1054 HashMap::from_iter(keys.into_iter().map(|key| (key.key_id.clone(), key)))
1055 })
1056 .unwrap_or_default(),
1057 next_row_id,
1058 };
1059
1060 metadata.borrow_mut().try_normalize()?;
1061 Ok(metadata)
1062 }
1063 }
1064
1065 impl TryFrom<TableMetadataV2> for TableMetadata {
1066 type Error = Error;
1067 fn try_from(value: TableMetadataV2) -> Result<Self, Error> {
1068 let snapshots = value.snapshots;
1069 let value = value.shared;
1070 let current_snapshot_id = if value.current_snapshot_id == Some(EMPTY_SNAPSHOT_ID) {
1071 None
1072 } else {
1073 value.current_snapshot_id
1074 };
1075 let schemas = HashMap::from_iter(
1076 value
1077 .schemas
1078 .into_iter()
1079 .map(|schema| Ok((schema.schema_id, Arc::new(schema.try_into()?))))
1080 .collect::<Result<Vec<_>, Error>>()?,
1081 );
1082
1083 let current_schema: &SchemaRef =
1084 schemas.get(&value.current_schema_id).ok_or_else(|| {
1085 invalid_data!(
1086 "No schema exists with the current schema id {}.",
1087 value.current_schema_id
1088 )
1089 })?;
1090 let partition_specs = HashMap::from_iter(
1091 value
1092 .partition_specs
1093 .into_iter()
1094 .map(|x| (x.spec_id(), Arc::new(x))),
1095 );
1096 let default_spec_id = value.default_spec_id;
1097 let default_spec: PartitionSpecRef = partition_specs
1098 .get(&value.default_spec_id)
1099 .map(|spec| (**spec).clone())
1100 .or_else(|| {
1101 (DEFAULT_PARTITION_SPEC_ID == default_spec_id)
1102 .then(PartitionSpec::unpartition_spec)
1103 })
1104 .ok_or_else(|| invalid_data!("Default partition spec {default_spec_id} not found"))?
1105 .into();
1106 let default_partition_type = default_spec.partition_type(current_schema)?;
1107
1108 let mut metadata = TableMetadata {
1109 format_version: FormatVersion::V2,
1110 table_uuid: value.table_uuid,
1111 location: value.location,
1112 last_sequence_number: value.last_sequence_number,
1113 last_updated_ms: value.last_updated_ms,
1114 last_column_id: value.last_column_id,
1115 current_schema_id: value.current_schema_id,
1116 schemas,
1117 partition_specs,
1118 default_partition_type,
1119 default_spec,
1120 last_partition_id: value.last_partition_id,
1121 properties: value.properties.unwrap_or_default(),
1122 current_snapshot_id,
1123 snapshots: snapshots
1124 .map(|snapshots| {
1125 HashMap::from_iter(
1126 snapshots
1127 .into_iter()
1128 .map(|x| (x.snapshot_id, Arc::new(x.into()))),
1129 )
1130 })
1131 .unwrap_or_default(),
1132 snapshot_log: value.snapshot_log.unwrap_or_default(),
1133 metadata_log: value.metadata_log.unwrap_or_default(),
1134 sort_orders: HashMap::from_iter(
1135 value
1136 .sort_orders
1137 .into_iter()
1138 .map(|x| (x.order_id, Arc::new(x))),
1139 ),
1140 default_sort_order_id: value.default_sort_order_id,
1141 refs: value.refs.unwrap_or_else(|| {
1142 if let Some(snapshot_id) = current_snapshot_id {
1143 HashMap::from_iter(vec![(MAIN_BRANCH.to_string(), SnapshotReference {
1144 snapshot_id,
1145 retention: SnapshotRetention::Branch {
1146 min_snapshots_to_keep: None,
1147 max_snapshot_age_ms: None,
1148 max_ref_age_ms: None,
1149 },
1150 })])
1151 } else {
1152 HashMap::new()
1153 }
1154 }),
1155 statistics: index_statistics(value.statistics),
1156 partition_statistics: index_partition_statistics(value.partition_statistics),
1157 encryption_keys: HashMap::new(),
1158 next_row_id: INITIAL_ROW_ID,
1159 };
1160
1161 metadata.borrow_mut().try_normalize()?;
1162 Ok(metadata)
1163 }
1164 }
1165
1166 impl TryFrom<TableMetadataV1> for TableMetadata {
1167 type Error = Error;
1168 fn try_from(value: TableMetadataV1) -> Result<Self, Error> {
1169 let current_snapshot_id = if value.current_snapshot_id == Some(EMPTY_SNAPSHOT_ID) {
1170 None
1171 } else {
1172 value.current_snapshot_id
1173 };
1174
1175 let (schemas, current_schema_id, current_schema) =
1176 if let (Some(schemas_vec), Some(schema_id)) =
1177 (&value.schemas, value.current_schema_id)
1178 {
1179 let schema_map = HashMap::from_iter(
1181 schemas_vec
1182 .clone()
1183 .into_iter()
1184 .map(|schema| {
1185 let schema: Schema = schema.try_into()?;
1186 Ok((schema.schema_id(), Arc::new(schema)))
1187 })
1188 .collect::<Result<Vec<_>, Error>>()?,
1189 );
1190
1191 let schema = schema_map
1192 .get(&schema_id)
1193 .ok_or_else(|| {
1194 invalid_data!(
1195 "No schema exists with the current schema id {schema_id}."
1196 )
1197 })?
1198 .clone();
1199 (schema_map, schema_id, schema)
1200 } else if let Some(schema) = value.schema {
1201 let schema: Schema = schema.try_into()?;
1203 let schema_id = schema.schema_id();
1204 let schema_arc = Arc::new(schema);
1205 let schema_map = HashMap::from_iter(vec![(schema_id, schema_arc.clone())]);
1206 (schema_map, schema_id, schema_arc)
1207 } else {
1208 return Err(invalid_data!(
1210 "No valid schema configuration found in table metadata"
1211 ));
1212 };
1213
1214 let partition_specs = if let Some(specs_vec) = value.partition_specs {
1216 specs_vec
1218 .into_iter()
1219 .map(|x| (x.spec_id(), Arc::new(x)))
1220 .collect::<HashMap<_, _>>()
1221 } else if let Some(partition_spec) = value.partition_spec {
1222 let spec = PartitionSpec::builder(current_schema.clone())
1224 .with_spec_id(DEFAULT_PARTITION_SPEC_ID)
1225 .add_unbound_fields(partition_spec.into_iter().map(|f| f.into_unbound()))?
1226 .build()?;
1227
1228 HashMap::from_iter(vec![(DEFAULT_PARTITION_SPEC_ID, Arc::new(spec))])
1229 } else {
1230 let spec = PartitionSpec::builder(current_schema.clone())
1232 .with_spec_id(DEFAULT_PARTITION_SPEC_ID)
1233 .build()?;
1234
1235 HashMap::from_iter(vec![(DEFAULT_PARTITION_SPEC_ID, Arc::new(spec))])
1236 };
1237
1238 let default_spec_id = value
1240 .default_spec_id
1241 .unwrap_or_else(|| partition_specs.keys().copied().max().unwrap_or_default());
1242
1243 let default_spec: PartitionSpecRef = partition_specs
1245 .get(&default_spec_id)
1246 .map(|x| Arc::unwrap_or_clone(x.clone()))
1247 .ok_or_else(|| invalid_data!("Default partition spec {default_spec_id} not found"))?
1248 .into();
1249 let default_partition_type = default_spec.partition_type(¤t_schema)?;
1250
1251 let mut metadata = TableMetadata {
1252 format_version: FormatVersion::V1,
1253 table_uuid: value.table_uuid.unwrap_or_default(),
1254 location: value.location,
1255 last_sequence_number: 0,
1256 last_updated_ms: value.last_updated_ms,
1257 last_column_id: value.last_column_id,
1258 current_schema_id,
1259 default_spec,
1260 default_partition_type,
1261 last_partition_id: value
1262 .last_partition_id
1263 .unwrap_or_else(|| partition_specs.keys().copied().max().unwrap_or_default()),
1264 partition_specs,
1265 schemas,
1266 properties: value.properties.unwrap_or_default(),
1267 current_snapshot_id,
1268 snapshots: value
1269 .snapshots
1270 .map(|snapshots| {
1271 Ok::<_, Error>(HashMap::from_iter(
1272 snapshots
1273 .into_iter()
1274 .map(|x| Ok((x.snapshot_id, Arc::new(x.try_into()?))))
1275 .collect::<Result<Vec<_>, Error>>()?,
1276 ))
1277 })
1278 .transpose()?
1279 .unwrap_or_default(),
1280 snapshot_log: value.snapshot_log.unwrap_or_default(),
1281 metadata_log: value.metadata_log.unwrap_or_default(),
1282 sort_orders: match value.sort_orders {
1283 Some(sort_orders) => HashMap::from_iter(
1284 sort_orders.into_iter().map(|x| (x.order_id, Arc::new(x))),
1285 ),
1286 None => HashMap::new(),
1287 },
1288 default_sort_order_id: value
1289 .default_sort_order_id
1290 .unwrap_or(SortOrder::UNSORTED_ORDER_ID),
1291 refs: if let Some(snapshot_id) = current_snapshot_id {
1292 HashMap::from_iter(vec![(MAIN_BRANCH.to_string(), SnapshotReference {
1293 snapshot_id,
1294 retention: SnapshotRetention::Branch {
1295 min_snapshots_to_keep: None,
1296 max_snapshot_age_ms: None,
1297 max_ref_age_ms: None,
1298 },
1299 })])
1300 } else {
1301 HashMap::new()
1302 },
1303 statistics: index_statistics(value.statistics),
1304 partition_statistics: index_partition_statistics(value.partition_statistics),
1305 encryption_keys: HashMap::new(),
1306 next_row_id: INITIAL_ROW_ID, };
1308
1309 metadata.borrow_mut().try_normalize()?;
1310 Ok(metadata)
1311 }
1312 }
1313
1314 impl TryFrom<TableMetadata> for TableMetadataV3 {
1315 type Error = Error;
1316
1317 fn try_from(mut v: TableMetadata) -> Result<Self, Self::Error> {
1318 let next_row_id = v.next_row_id;
1319 let encryption_keys = std::mem::take(&mut v.encryption_keys);
1320 let snapshots = std::mem::take(&mut v.snapshots);
1321 let shared = v.into();
1322
1323 Ok(TableMetadataV3 {
1324 format_version: VersionNumber::<3>,
1325 shared,
1326 next_row_id,
1327 encryption_keys: if encryption_keys.is_empty() {
1328 None
1329 } else {
1330 Some(encryption_keys.into_values().collect())
1331 },
1332 snapshots: if snapshots.is_empty() {
1333 None
1334 } else {
1335 Some(
1336 snapshots
1337 .into_values()
1338 .map(|s| SnapshotV3::try_from(Arc::unwrap_or_clone(s)))
1339 .collect::<Result<_, _>>()?,
1340 )
1341 },
1342 })
1343 }
1344 }
1345
1346 impl From<TableMetadata> for TableMetadataV2 {
1347 fn from(mut v: TableMetadata) -> Self {
1348 let snapshots = std::mem::take(&mut v.snapshots);
1349 let shared = v.into();
1350
1351 TableMetadataV2 {
1352 format_version: VersionNumber::<2>,
1353 shared,
1354 snapshots: if snapshots.is_empty() {
1355 None
1356 } else {
1357 Some(
1358 snapshots
1359 .into_values()
1360 .map(|s| SnapshotV2::from(Arc::unwrap_or_clone(s)))
1361 .collect(),
1362 )
1363 },
1364 }
1365 }
1366 }
1367
1368 impl From<TableMetadata> for TableMetadataV2V3Shared {
1369 fn from(v: TableMetadata) -> Self {
1370 TableMetadataV2V3Shared {
1371 table_uuid: v.table_uuid,
1372 location: v.location,
1373 last_sequence_number: v.last_sequence_number,
1374 last_updated_ms: v.last_updated_ms,
1375 last_column_id: v.last_column_id,
1376 schemas: v
1377 .schemas
1378 .into_values()
1379 .map(|x| {
1380 Arc::try_unwrap(x)
1381 .unwrap_or_else(|schema| schema.as_ref().clone())
1382 .into()
1383 })
1384 .collect(),
1385 current_schema_id: v.current_schema_id,
1386 partition_specs: v
1387 .partition_specs
1388 .into_values()
1389 .map(|x| Arc::try_unwrap(x).unwrap_or_else(|s| s.as_ref().clone()))
1390 .collect(),
1391 default_spec_id: v.default_spec.spec_id(),
1392 last_partition_id: v.last_partition_id,
1393 properties: if v.properties.is_empty() {
1394 None
1395 } else {
1396 Some(v.properties)
1397 },
1398 current_snapshot_id: v.current_snapshot_id,
1399 snapshot_log: if v.snapshot_log.is_empty() {
1400 None
1401 } else {
1402 Some(v.snapshot_log)
1403 },
1404 metadata_log: if v.metadata_log.is_empty() {
1405 None
1406 } else {
1407 Some(v.metadata_log)
1408 },
1409 sort_orders: v
1410 .sort_orders
1411 .into_values()
1412 .map(|x| Arc::try_unwrap(x).unwrap_or_else(|s| s.as_ref().clone()))
1413 .collect(),
1414 default_sort_order_id: v.default_sort_order_id,
1415 refs: Some(v.refs),
1416 statistics: v.statistics.into_values().collect(),
1417 partition_statistics: v.partition_statistics.into_values().collect(),
1418 }
1419 }
1420 }
1421
1422 impl TryFrom<TableMetadata> for TableMetadataV1 {
1423 type Error = Error;
1424 fn try_from(v: TableMetadata) -> Result<Self, Error> {
1425 Ok(TableMetadataV1 {
1426 format_version: VersionNumber::<1>,
1427 table_uuid: Some(v.table_uuid),
1428 location: v.location,
1429 last_updated_ms: v.last_updated_ms,
1430 last_column_id: v.last_column_id,
1431 schema: Some(
1432 v.schemas
1433 .get(&v.current_schema_id)
1434 .ok_or(Error::new(
1435 ErrorKind::Unexpected,
1436 "current_schema_id not found in schemas",
1437 ))?
1438 .as_ref()
1439 .clone()
1440 .into(),
1441 ),
1442 schemas: Some(
1443 v.schemas
1444 .into_values()
1445 .map(|x| {
1446 Arc::try_unwrap(x)
1447 .unwrap_or_else(|schema| schema.as_ref().clone())
1448 .into()
1449 })
1450 .collect(),
1451 ),
1452 current_schema_id: Some(v.current_schema_id),
1453 partition_spec: Some(v.default_spec.fields().to_vec()),
1454 partition_specs: Some(
1455 v.partition_specs
1456 .into_values()
1457 .map(|x| Arc::try_unwrap(x).unwrap_or_else(|s| s.as_ref().clone()))
1458 .collect(),
1459 ),
1460 default_spec_id: Some(v.default_spec.spec_id()),
1461 last_partition_id: Some(v.last_partition_id),
1462 properties: if v.properties.is_empty() {
1463 None
1464 } else {
1465 Some(v.properties)
1466 },
1467 current_snapshot_id: v.current_snapshot_id,
1468 snapshots: if v.snapshots.is_empty() {
1469 None
1470 } else {
1471 Some(
1472 v.snapshots
1473 .into_values()
1474 .map(|x| Snapshot::clone(&x).into())
1475 .collect(),
1476 )
1477 },
1478 snapshot_log: if v.snapshot_log.is_empty() {
1479 None
1480 } else {
1481 Some(v.snapshot_log)
1482 },
1483 metadata_log: if v.metadata_log.is_empty() {
1484 None
1485 } else {
1486 Some(v.metadata_log)
1487 },
1488 sort_orders: Some(
1489 v.sort_orders
1490 .into_values()
1491 .map(|s| Arc::try_unwrap(s).unwrap_or_else(|s| s.as_ref().clone()))
1492 .collect(),
1493 ),
1494 default_sort_order_id: Some(v.default_sort_order_id),
1495 statistics: v.statistics.into_values().collect(),
1496 partition_statistics: v.partition_statistics.into_values().collect(),
1497 })
1498 }
1499 }
1500
1501 fn index_statistics(statistics: Vec<StatisticsFile>) -> HashMap<i64, StatisticsFile> {
1502 statistics
1503 .into_iter()
1504 .rev()
1505 .map(|s| (s.snapshot_id, s))
1506 .collect()
1507 }
1508
1509 fn index_partition_statistics(
1510 statistics: Vec<PartitionStatisticsFile>,
1511 ) -> HashMap<i64, PartitionStatisticsFile> {
1512 statistics
1513 .into_iter()
1514 .rev()
1515 .map(|s| (s.snapshot_id, s))
1516 .collect()
1517 }
1518}
1519
1520#[derive(Debug, Serialize_repr, Deserialize_repr, PartialEq, Eq, Clone, Copy, Hash)]
1521#[repr(u8)]
1522pub enum FormatVersion {
1524 V1 = 1u8,
1526 V2 = 2u8,
1528 V3 = 3u8,
1530}
1531
1532impl PartialOrd for FormatVersion {
1533 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
1534 Some(self.cmp(other))
1535 }
1536}
1537
1538impl Ord for FormatVersion {
1539 fn cmp(&self, other: &Self) -> Ordering {
1540 (*self as u8).cmp(&(*other as u8))
1541 }
1542}
1543
1544impl Display for FormatVersion {
1545 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
1546 match self {
1547 FormatVersion::V1 => write!(f, "v1"),
1548 FormatVersion::V2 => write!(f, "v2"),
1549 FormatVersion::V3 => write!(f, "v3"),
1550 }
1551 }
1552}
1553
1554#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
1555#[serde(rename_all = "kebab-case")]
1556pub struct MetadataLog {
1558 pub metadata_file: String,
1560 pub timestamp_ms: i64,
1562}
1563
1564#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
1565#[serde(rename_all = "kebab-case")]
1566pub struct SnapshotLog {
1568 pub snapshot_id: i64,
1570 pub timestamp_ms: i64,
1572}
1573
1574impl SnapshotLog {
1575 pub fn timestamp(self) -> Result<DateTime<Utc>> {
1577 timestamp_ms_to_utc(self.timestamp_ms)
1578 }
1579
1580 #[inline]
1582 pub fn timestamp_ms(&self) -> i64 {
1583 self.timestamp_ms
1584 }
1585}
1586
1587#[cfg(test)]
1588mod tests {
1589 use std::collections::HashMap;
1590 use std::fs;
1591 use std::sync::Arc;
1592
1593 use anyhow::Result;
1594 use base64::Engine as _;
1595 use pretty_assertions::assert_eq;
1596 use tempfile::TempDir;
1597 use uuid::Uuid;
1598
1599 use super::{FormatVersion, MetadataLog, SnapshotLog, TableMetadataBuilder};
1600 use crate::catalog::MetadataLocation;
1601 use crate::compression::CompressionCodec;
1602 use crate::io::FileIO;
1603 use crate::spec::table_metadata::TableMetadata;
1604 use crate::spec::{
1605 BlobMetadata, EncryptedKey, INITIAL_ROW_ID, Literal, NestedField, NullOrder, Operation,
1606 PartitionSpec, PartitionStatisticsFile, PrimitiveLiteral, PrimitiveType, Schema, Snapshot,
1607 SnapshotReference, SnapshotRetention, SortDirection, SortField, SortOrder, StatisticsFile,
1608 Summary, TableProperties, Transform, Type, UnboundPartitionField, UnboundPartitionSpec,
1609 };
1610 use crate::{ErrorKind, TableCreation};
1611
1612 fn check_table_metadata_serde(json: &str, expected_type: TableMetadata) {
1613 let desered_type: TableMetadata = serde_json::from_str(json).unwrap();
1614 assert_eq!(desered_type, expected_type);
1615
1616 let sered_json = serde_json::to_string(&expected_type).unwrap();
1617 let parsed_json_value = serde_json::from_str::<TableMetadata>(&sered_json).unwrap();
1618
1619 assert_eq!(parsed_json_value, desered_type);
1620 }
1621
1622 fn get_test_table_metadata(file_name: &str) -> TableMetadata {
1623 let path = format!("testdata/table_metadata/{file_name}");
1624 let metadata: String = fs::read_to_string(path).unwrap();
1625
1626 serde_json::from_str(&metadata).unwrap()
1627 }
1628
1629 fn get_test_table_metadata_at(file_name: &str, location: &str) -> TableMetadata {
1632 TableMetadataBuilder::new_from_metadata(get_test_table_metadata(file_name), None)
1633 .set_location(location.to_string())
1634 .build()
1635 .unwrap()
1636 .metadata
1637 }
1638
1639 #[test]
1640 fn test_table_data_v2() {
1641 let data = r#"
1642 {
1643 "format-version" : 2,
1644 "table-uuid": "fb072c92-a02b-11e9-ae9c-1bb7bc9eca94",
1645 "location": "s3://b/wh/data.db/table",
1646 "last-sequence-number" : 1,
1647 "last-updated-ms": 1515100955770,
1648 "last-column-id": 1,
1649 "schemas": [
1650 {
1651 "schema-id" : 1,
1652 "type" : "struct",
1653 "fields" :[
1654 {
1655 "id": 1,
1656 "name": "struct_name",
1657 "required": true,
1658 "type": "fixed[1]"
1659 },
1660 {
1661 "id": 4,
1662 "name": "ts",
1663 "required": true,
1664 "type": "timestamp"
1665 }
1666 ]
1667 }
1668 ],
1669 "current-schema-id" : 1,
1670 "partition-specs": [
1671 {
1672 "spec-id": 0,
1673 "fields": [
1674 {
1675 "source-id": 4,
1676 "field-id": 1000,
1677 "name": "ts_day",
1678 "transform": "day"
1679 }
1680 ]
1681 }
1682 ],
1683 "default-spec-id": 0,
1684 "last-partition-id": 1000,
1685 "properties": {
1686 "commit.retry.num-retries": "1"
1687 },
1688 "metadata-log": [
1689 {
1690 "metadata-file": "s3://bucket/.../v1.json",
1691 "timestamp-ms": 1515100
1692 }
1693 ],
1694 "refs": {},
1695 "sort-orders": [
1696 {
1697 "order-id": 0,
1698 "fields": []
1699 }
1700 ],
1701 "default-sort-order-id": 0
1702 }
1703 "#;
1704
1705 let schema = Schema::builder()
1706 .with_schema_id(1)
1707 .with_fields(vec![
1708 Arc::new(NestedField::required(
1709 1,
1710 "struct_name",
1711 Type::Primitive(PrimitiveType::Fixed(1)),
1712 )),
1713 Arc::new(NestedField::required(
1714 4,
1715 "ts",
1716 Type::Primitive(PrimitiveType::Timestamp),
1717 )),
1718 ])
1719 .build()
1720 .unwrap();
1721
1722 let partition_spec = PartitionSpec::builder(schema.clone())
1723 .with_spec_id(0)
1724 .add_unbound_field(
1725 UnboundPartitionField::builder()
1726 .source_ids(vec![4])
1727 .field_id(1000)
1728 .name("ts_day".to_string())
1729 .transform(Transform::Day)
1730 .build()
1731 .unwrap(),
1732 )
1733 .unwrap()
1734 .build()
1735 .unwrap();
1736
1737 let default_partition_type = partition_spec.partition_type(&schema).unwrap();
1738 let expected = TableMetadata {
1739 format_version: FormatVersion::V2,
1740 table_uuid: Uuid::parse_str("fb072c92-a02b-11e9-ae9c-1bb7bc9eca94").unwrap(),
1741 location: "s3://b/wh/data.db/table".to_string(),
1742 last_updated_ms: 1515100955770,
1743 last_column_id: 1,
1744 schemas: HashMap::from_iter(vec![(1, Arc::new(schema))]),
1745 current_schema_id: 1,
1746 partition_specs: HashMap::from_iter(vec![(0, partition_spec.clone().into())]),
1747 default_partition_type,
1748 default_spec: partition_spec.into(),
1749 last_partition_id: 1000,
1750 default_sort_order_id: 0,
1751 sort_orders: HashMap::from_iter(vec![(0, SortOrder::unsorted_order().into())]),
1752 snapshots: HashMap::default(),
1753 current_snapshot_id: None,
1754 last_sequence_number: 1,
1755 properties: HashMap::from_iter(vec![(
1756 "commit.retry.num-retries".to_string(),
1757 "1".to_string(),
1758 )]),
1759 snapshot_log: Vec::new(),
1760 metadata_log: vec![MetadataLog {
1761 metadata_file: "s3://bucket/.../v1.json".to_string(),
1762 timestamp_ms: 1515100,
1763 }],
1764 refs: HashMap::new(),
1765 statistics: HashMap::new(),
1766 partition_statistics: HashMap::new(),
1767 encryption_keys: HashMap::new(),
1768 next_row_id: INITIAL_ROW_ID,
1769 };
1770
1771 let expected_json_value = serde_json::to_value(&expected).unwrap();
1772 check_table_metadata_serde(data, expected);
1773
1774 let json_value = serde_json::from_str::<serde_json::Value>(data).unwrap();
1775 assert_eq!(json_value, expected_json_value);
1776 }
1777
1778 #[test]
1779 fn test_table_data_v3() {
1780 let data = r#"
1781 {
1782 "format-version" : 3,
1783 "table-uuid": "fb072c92-a02b-11e9-ae9c-1bb7bc9eca94",
1784 "location": "s3://b/wh/data.db/table",
1785 "last-sequence-number" : 1,
1786 "last-updated-ms": 1515100955770,
1787 "last-column-id": 1,
1788 "next-row-id": 5,
1789 "schemas": [
1790 {
1791 "schema-id" : 1,
1792 "type" : "struct",
1793 "fields" :[
1794 {
1795 "id": 4,
1796 "name": "ts",
1797 "required": true,
1798 "type": "timestamp"
1799 }
1800 ]
1801 }
1802 ],
1803 "current-schema-id" : 1,
1804 "partition-specs": [
1805 {
1806 "spec-id": 0,
1807 "fields": [
1808 {
1809 "source-id": 4,
1810 "field-id": 1000,
1811 "name": "ts_day",
1812 "transform": "day"
1813 }
1814 ]
1815 }
1816 ],
1817 "default-spec-id": 0,
1818 "last-partition-id": 1000,
1819 "properties": {
1820 "commit.retry.num-retries": "1"
1821 },
1822 "metadata-log": [
1823 {
1824 "metadata-file": "s3://bucket/.../v1.json",
1825 "timestamp-ms": 1515100
1826 }
1827 ],
1828 "refs": {},
1829 "snapshots" : [ {
1830 "snapshot-id" : 1,
1831 "timestamp-ms" : 1662532818843,
1832 "sequence-number" : 0,
1833 "first-row-id" : 0,
1834 "added-rows" : 4,
1835 "key-id" : "key1",
1836 "summary" : {
1837 "operation" : "append"
1838 },
1839 "manifest-list" : "/home/iceberg/warehouse/nyc/taxis/metadata/snap-638933773299822130-1-7e6760f0-4f6c-4b23-b907-0a5a174e3863.avro",
1840 "schema-id" : 0
1841 }
1842 ],
1843 "encryption-keys": [
1844 {
1845 "key-id": "key1",
1846 "encrypted-by-id": "KMS",
1847 "encrypted-key-metadata": "c29tZS1lbmNyeXB0aW9uLWtleQ==",
1848 "properties": {
1849 "p1": "v1"
1850 }
1851 }
1852 ],
1853 "sort-orders": [
1854 {
1855 "order-id": 0,
1856 "fields": []
1857 }
1858 ],
1859 "default-sort-order-id": 0
1860 }
1861 "#;
1862
1863 let schema = Schema::builder()
1864 .with_schema_id(1)
1865 .with_fields(vec![Arc::new(NestedField::required(
1866 4,
1867 "ts",
1868 Type::Primitive(PrimitiveType::Timestamp),
1869 ))])
1870 .build()
1871 .unwrap();
1872
1873 let partition_spec = PartitionSpec::builder(schema.clone())
1874 .with_spec_id(0)
1875 .add_unbound_field(
1876 UnboundPartitionField::builder()
1877 .source_ids(vec![4])
1878 .field_id(1000)
1879 .name("ts_day".to_string())
1880 .transform(Transform::Day)
1881 .build()
1882 .unwrap(),
1883 )
1884 .unwrap()
1885 .build()
1886 .unwrap();
1887
1888 let snapshot = Snapshot::builder()
1889 .with_snapshot_id(1)
1890 .with_timestamp_ms(1662532818843)
1891 .with_sequence_number(0)
1892 .with_row_range(0, 4)
1893 .with_encryption_key_id(Some("key1".to_string()))
1894 .with_summary(Summary {
1895 operation: Operation::Append,
1896 additional_properties: HashMap::new(),
1897 })
1898 .with_manifest_list("/home/iceberg/warehouse/nyc/taxis/metadata/snap-638933773299822130-1-7e6760f0-4f6c-4b23-b907-0a5a174e3863.avro".to_string())
1899 .with_schema_id(0)
1900 .build();
1901
1902 let encryption_key = EncryptedKey::builder()
1903 .key_id("key1".to_string())
1904 .encrypted_by_id("KMS".to_string())
1905 .encrypted_key_metadata(
1906 base64::prelude::BASE64_STANDARD
1907 .decode("c29tZS1lbmNyeXB0aW9uLWtleQ==")
1908 .unwrap(),
1909 )
1910 .properties(HashMap::from_iter(vec![(
1911 "p1".to_string(),
1912 "v1".to_string(),
1913 )]))
1914 .build();
1915
1916 let default_partition_type = partition_spec.partition_type(&schema).unwrap();
1917 let expected = TableMetadata {
1918 format_version: FormatVersion::V3,
1919 table_uuid: Uuid::parse_str("fb072c92-a02b-11e9-ae9c-1bb7bc9eca94").unwrap(),
1920 location: "s3://b/wh/data.db/table".to_string(),
1921 last_updated_ms: 1515100955770,
1922 last_column_id: 1,
1923 schemas: HashMap::from_iter(vec![(1, Arc::new(schema))]),
1924 current_schema_id: 1,
1925 partition_specs: HashMap::from_iter(vec![(0, partition_spec.clone().into())]),
1926 default_partition_type,
1927 default_spec: partition_spec.into(),
1928 last_partition_id: 1000,
1929 default_sort_order_id: 0,
1930 sort_orders: HashMap::from_iter(vec![(0, SortOrder::unsorted_order().into())]),
1931 snapshots: HashMap::from_iter(vec![(1, snapshot.into())]),
1932 current_snapshot_id: None,
1933 last_sequence_number: 1,
1934 properties: HashMap::from_iter(vec![(
1935 "commit.retry.num-retries".to_string(),
1936 "1".to_string(),
1937 )]),
1938 snapshot_log: Vec::new(),
1939 metadata_log: vec![MetadataLog {
1940 metadata_file: "s3://bucket/.../v1.json".to_string(),
1941 timestamp_ms: 1515100,
1942 }],
1943 refs: HashMap::new(),
1944 statistics: HashMap::new(),
1945 partition_statistics: HashMap::new(),
1946 encryption_keys: HashMap::from_iter(vec![("key1".to_string(), encryption_key)]),
1947 next_row_id: 5,
1948 };
1949
1950 let expected_json_value = serde_json::to_value(&expected).unwrap();
1951 check_table_metadata_serde(data, expected);
1952
1953 let json_value = serde_json::from_str::<serde_json::Value>(data).unwrap();
1954 assert_eq!(json_value, expected_json_value);
1955 }
1956
1957 #[test]
1958 fn test_table_data_v1() {
1959 let data = r#"
1960 {
1961 "format-version" : 1,
1962 "table-uuid" : "df838b92-0b32-465d-a44e-d39936e538b7",
1963 "location" : "/home/iceberg/warehouse/nyc/taxis",
1964 "last-updated-ms" : 1662532818843,
1965 "last-column-id" : 5,
1966 "schema" : {
1967 "type" : "struct",
1968 "schema-id" : 0,
1969 "fields" : [ {
1970 "id" : 1,
1971 "name" : "vendor_id",
1972 "required" : false,
1973 "type" : "long"
1974 }, {
1975 "id" : 2,
1976 "name" : "trip_id",
1977 "required" : false,
1978 "type" : "long"
1979 }, {
1980 "id" : 3,
1981 "name" : "trip_distance",
1982 "required" : false,
1983 "type" : "float"
1984 }, {
1985 "id" : 4,
1986 "name" : "fare_amount",
1987 "required" : false,
1988 "type" : "double"
1989 }, {
1990 "id" : 5,
1991 "name" : "store_and_fwd_flag",
1992 "required" : false,
1993 "type" : "string"
1994 } ]
1995 },
1996 "partition-spec" : [ {
1997 "name" : "vendor_id",
1998 "transform" : "identity",
1999 "source-id" : 1,
2000 "field-id" : 1000
2001 } ],
2002 "last-partition-id" : 1000,
2003 "default-sort-order-id" : 0,
2004 "sort-orders" : [ {
2005 "order-id" : 0,
2006 "fields" : [ ]
2007 } ],
2008 "properties" : {
2009 "owner" : "root"
2010 },
2011 "current-snapshot-id" : 638933773299822130,
2012 "refs" : {
2013 "main" : {
2014 "snapshot-id" : 638933773299822130,
2015 "type" : "branch"
2016 }
2017 },
2018 "snapshots" : [ {
2019 "snapshot-id" : 638933773299822130,
2020 "timestamp-ms" : 1662532818843,
2021 "sequence-number" : 0,
2022 "summary" : {
2023 "operation" : "append",
2024 "spark.app.id" : "local-1662532784305",
2025 "added-data-files" : "4",
2026 "added-records" : "4",
2027 "added-files-size" : "6001"
2028 },
2029 "manifest-list" : "/home/iceberg/warehouse/nyc/taxis/metadata/snap-638933773299822130-1-7e6760f0-4f6c-4b23-b907-0a5a174e3863.avro",
2030 "schema-id" : 0
2031 } ],
2032 "snapshot-log" : [ {
2033 "timestamp-ms" : 1662532818843,
2034 "snapshot-id" : 638933773299822130
2035 } ],
2036 "metadata-log" : [ {
2037 "timestamp-ms" : 1662532805245,
2038 "metadata-file" : "/home/iceberg/warehouse/nyc/taxis/metadata/00000-8a62c37d-4573-4021-952a-c0baef7d21d0.metadata.json"
2039 } ]
2040 }
2041 "#;
2042
2043 let schema = Schema::builder()
2044 .with_fields(vec![
2045 Arc::new(NestedField::optional(
2046 1,
2047 "vendor_id",
2048 Type::Primitive(PrimitiveType::Long),
2049 )),
2050 Arc::new(NestedField::optional(
2051 2,
2052 "trip_id",
2053 Type::Primitive(PrimitiveType::Long),
2054 )),
2055 Arc::new(NestedField::optional(
2056 3,
2057 "trip_distance",
2058 Type::Primitive(PrimitiveType::Float),
2059 )),
2060 Arc::new(NestedField::optional(
2061 4,
2062 "fare_amount",
2063 Type::Primitive(PrimitiveType::Double),
2064 )),
2065 Arc::new(NestedField::optional(
2066 5,
2067 "store_and_fwd_flag",
2068 Type::Primitive(PrimitiveType::String),
2069 )),
2070 ])
2071 .build()
2072 .unwrap();
2073
2074 let schema = Arc::new(schema);
2075 let partition_spec = PartitionSpec::builder(schema.clone())
2076 .with_spec_id(0)
2077 .add_partition_field("vendor_id", "vendor_id", Transform::Identity)
2078 .unwrap()
2079 .build()
2080 .unwrap();
2081
2082 let sort_order = SortOrder::builder()
2083 .with_order_id(0)
2084 .build_unbound()
2085 .unwrap();
2086
2087 let snapshot = Snapshot::builder()
2088 .with_snapshot_id(638933773299822130)
2089 .with_timestamp_ms(1662532818843)
2090 .with_sequence_number(0)
2091 .with_schema_id(0)
2092 .with_manifest_list("/home/iceberg/warehouse/nyc/taxis/metadata/snap-638933773299822130-1-7e6760f0-4f6c-4b23-b907-0a5a174e3863.avro")
2093 .with_summary(Summary { operation: Operation::Append, additional_properties: HashMap::from_iter(vec![("spark.app.id".to_string(), "local-1662532784305".to_string()), ("added-data-files".to_string(), "4".to_string()), ("added-records".to_string(), "4".to_string()), ("added-files-size".to_string(), "6001".to_string())]) })
2094 .build();
2095
2096 let default_partition_type = partition_spec.partition_type(&schema).unwrap();
2097 let expected = TableMetadata {
2098 format_version: FormatVersion::V1,
2099 table_uuid: Uuid::parse_str("df838b92-0b32-465d-a44e-d39936e538b7").unwrap(),
2100 location: "/home/iceberg/warehouse/nyc/taxis".to_string(),
2101 last_updated_ms: 1662532818843,
2102 last_column_id: 5,
2103 schemas: HashMap::from_iter(vec![(0, schema)]),
2104 current_schema_id: 0,
2105 partition_specs: HashMap::from_iter(vec![(0, partition_spec.clone().into())]),
2106 default_partition_type,
2107 default_spec: Arc::new(partition_spec),
2108 last_partition_id: 1000,
2109 default_sort_order_id: 0,
2110 sort_orders: HashMap::from_iter(vec![(0, sort_order.into())]),
2111 snapshots: HashMap::from_iter(vec![(638933773299822130, Arc::new(snapshot))]),
2112 current_snapshot_id: Some(638933773299822130),
2113 last_sequence_number: 0,
2114 properties: HashMap::from_iter(vec![("owner".to_string(), "root".to_string())]),
2115 snapshot_log: vec![SnapshotLog {
2116 snapshot_id: 638933773299822130,
2117 timestamp_ms: 1662532818843,
2118 }],
2119 metadata_log: vec![MetadataLog { metadata_file: "/home/iceberg/warehouse/nyc/taxis/metadata/00000-8a62c37d-4573-4021-952a-c0baef7d21d0.metadata.json".to_string(), timestamp_ms: 1662532805245 }],
2120 refs: HashMap::from_iter(vec![("main".to_string(), SnapshotReference { snapshot_id: 638933773299822130, retention: SnapshotRetention::Branch { min_snapshots_to_keep: None, max_snapshot_age_ms: None, max_ref_age_ms: None } })]),
2121 statistics: HashMap::new(),
2122 partition_statistics: HashMap::new(),
2123 encryption_keys: HashMap::new(),
2124 next_row_id: INITIAL_ROW_ID,
2125 };
2126
2127 check_table_metadata_serde(data, expected);
2128 }
2129
2130 #[test]
2131 fn test_table_data_v2_no_snapshots() {
2132 let data = r#"
2133 {
2134 "format-version" : 2,
2135 "table-uuid": "fb072c92-a02b-11e9-ae9c-1bb7bc9eca94",
2136 "location": "s3://b/wh/data.db/table",
2137 "last-sequence-number" : 1,
2138 "last-updated-ms": 1515100955770,
2139 "last-column-id": 1,
2140 "schemas": [
2141 {
2142 "schema-id" : 1,
2143 "type" : "struct",
2144 "fields" :[
2145 {
2146 "id": 1,
2147 "name": "struct_name",
2148 "required": true,
2149 "type": "fixed[1]"
2150 }
2151 ]
2152 }
2153 ],
2154 "current-schema-id" : 1,
2155 "partition-specs": [
2156 {
2157 "spec-id": 0,
2158 "fields": []
2159 }
2160 ],
2161 "refs": {},
2162 "default-spec-id": 0,
2163 "last-partition-id": 1000,
2164 "metadata-log": [
2165 {
2166 "metadata-file": "s3://bucket/.../v1.json",
2167 "timestamp-ms": 1515100
2168 }
2169 ],
2170 "sort-orders": [
2171 {
2172 "order-id": 0,
2173 "fields": []
2174 }
2175 ],
2176 "default-sort-order-id": 0
2177 }
2178 "#;
2179
2180 let schema = Schema::builder()
2181 .with_schema_id(1)
2182 .with_fields(vec![Arc::new(NestedField::required(
2183 1,
2184 "struct_name",
2185 Type::Primitive(PrimitiveType::Fixed(1)),
2186 ))])
2187 .build()
2188 .unwrap();
2189
2190 let partition_spec = PartitionSpec::builder(schema.clone())
2191 .with_spec_id(0)
2192 .build()
2193 .unwrap();
2194
2195 let default_partition_type = partition_spec.partition_type(&schema).unwrap();
2196 let expected = TableMetadata {
2197 format_version: FormatVersion::V2,
2198 table_uuid: Uuid::parse_str("fb072c92-a02b-11e9-ae9c-1bb7bc9eca94").unwrap(),
2199 location: "s3://b/wh/data.db/table".to_string(),
2200 last_updated_ms: 1515100955770,
2201 last_column_id: 1,
2202 schemas: HashMap::from_iter(vec![(1, Arc::new(schema))]),
2203 current_schema_id: 1,
2204 partition_specs: HashMap::from_iter(vec![(0, partition_spec.clone().into())]),
2205 default_partition_type,
2206 default_spec: partition_spec.into(),
2207 last_partition_id: 1000,
2208 default_sort_order_id: 0,
2209 sort_orders: HashMap::from_iter(vec![(0, SortOrder::unsorted_order().into())]),
2210 snapshots: HashMap::default(),
2211 current_snapshot_id: None,
2212 last_sequence_number: 1,
2213 properties: HashMap::new(),
2214 snapshot_log: Vec::new(),
2215 metadata_log: vec![MetadataLog {
2216 metadata_file: "s3://bucket/.../v1.json".to_string(),
2217 timestamp_ms: 1515100,
2218 }],
2219 refs: HashMap::new(),
2220 statistics: HashMap::new(),
2221 partition_statistics: HashMap::new(),
2222 encryption_keys: HashMap::new(),
2223 next_row_id: INITIAL_ROW_ID,
2224 };
2225
2226 let expected_json_value = serde_json::to_value(&expected).unwrap();
2227 check_table_metadata_serde(data, expected);
2228
2229 let json_value = serde_json::from_str::<serde_json::Value>(data).unwrap();
2230 assert_eq!(json_value, expected_json_value);
2231 }
2232
2233 #[test]
2234 fn test_current_snapshot_id_must_match_main_branch() {
2235 let data = r#"
2236 {
2237 "format-version" : 2,
2238 "table-uuid": "fb072c92-a02b-11e9-ae9c-1bb7bc9eca94",
2239 "location": "s3://b/wh/data.db/table",
2240 "last-sequence-number" : 1,
2241 "last-updated-ms": 1515100955770,
2242 "last-column-id": 1,
2243 "schemas": [
2244 {
2245 "schema-id" : 1,
2246 "type" : "struct",
2247 "fields" :[
2248 {
2249 "id": 1,
2250 "name": "struct_name",
2251 "required": true,
2252 "type": "fixed[1]"
2253 },
2254 {
2255 "id": 4,
2256 "name": "ts",
2257 "required": true,
2258 "type": "timestamp"
2259 }
2260 ]
2261 }
2262 ],
2263 "current-schema-id" : 1,
2264 "partition-specs": [
2265 {
2266 "spec-id": 0,
2267 "fields": [
2268 {
2269 "source-id": 4,
2270 "field-id": 1000,
2271 "name": "ts_day",
2272 "transform": "day"
2273 }
2274 ]
2275 }
2276 ],
2277 "default-spec-id": 0,
2278 "last-partition-id": 1000,
2279 "properties": {
2280 "commit.retry.num-retries": "1"
2281 },
2282 "metadata-log": [
2283 {
2284 "metadata-file": "s3://bucket/.../v1.json",
2285 "timestamp-ms": 1515100
2286 }
2287 ],
2288 "sort-orders": [
2289 {
2290 "order-id": 0,
2291 "fields": []
2292 }
2293 ],
2294 "default-sort-order-id": 0,
2295 "current-snapshot-id" : 1,
2296 "refs" : {
2297 "main" : {
2298 "snapshot-id" : 2,
2299 "type" : "branch"
2300 }
2301 },
2302 "snapshots" : [ {
2303 "snapshot-id" : 1,
2304 "timestamp-ms" : 1662532818843,
2305 "sequence-number" : 0,
2306 "summary" : {
2307 "operation" : "append",
2308 "spark.app.id" : "local-1662532784305",
2309 "added-data-files" : "4",
2310 "added-records" : "4",
2311 "added-files-size" : "6001"
2312 },
2313 "manifest-list" : "/home/iceberg/warehouse/nyc/taxis/metadata/snap-638933773299822130-1-7e6760f0-4f6c-4b23-b907-0a5a174e3863.avro",
2314 "schema-id" : 0
2315 },
2316 {
2317 "snapshot-id" : 2,
2318 "timestamp-ms" : 1662532818844,
2319 "sequence-number" : 0,
2320 "summary" : {
2321 "operation" : "append",
2322 "spark.app.id" : "local-1662532784305",
2323 "added-data-files" : "4",
2324 "added-records" : "4",
2325 "added-files-size" : "6001"
2326 },
2327 "manifest-list" : "/home/iceberg/warehouse/nyc/taxis/metadata/snap-638933773299822130-1-7e6760f0-4f6c-4b23-b907-0a5a174e3863.avro",
2328 "schema-id" : 0
2329 } ]
2330 }
2331 "#;
2332
2333 let err = serde_json::from_str::<TableMetadata>(data).unwrap_err();
2334 assert!(
2335 err.to_string()
2336 .contains("Current snapshot id does not match main branch")
2337 );
2338 }
2339
2340 #[test]
2341 fn test_main_without_current() {
2342 let data = r#"
2343 {
2344 "format-version" : 2,
2345 "table-uuid": "fb072c92-a02b-11e9-ae9c-1bb7bc9eca94",
2346 "location": "s3://b/wh/data.db/table",
2347 "last-sequence-number" : 1,
2348 "last-updated-ms": 1515100955770,
2349 "last-column-id": 1,
2350 "schemas": [
2351 {
2352 "schema-id" : 1,
2353 "type" : "struct",
2354 "fields" :[
2355 {
2356 "id": 1,
2357 "name": "struct_name",
2358 "required": true,
2359 "type": "fixed[1]"
2360 },
2361 {
2362 "id": 4,
2363 "name": "ts",
2364 "required": true,
2365 "type": "timestamp"
2366 }
2367 ]
2368 }
2369 ],
2370 "current-schema-id" : 1,
2371 "partition-specs": [
2372 {
2373 "spec-id": 0,
2374 "fields": [
2375 {
2376 "source-id": 4,
2377 "field-id": 1000,
2378 "name": "ts_day",
2379 "transform": "day"
2380 }
2381 ]
2382 }
2383 ],
2384 "default-spec-id": 0,
2385 "last-partition-id": 1000,
2386 "properties": {
2387 "commit.retry.num-retries": "1"
2388 },
2389 "metadata-log": [
2390 {
2391 "metadata-file": "s3://bucket/.../v1.json",
2392 "timestamp-ms": 1515100
2393 }
2394 ],
2395 "sort-orders": [
2396 {
2397 "order-id": 0,
2398 "fields": []
2399 }
2400 ],
2401 "default-sort-order-id": 0,
2402 "refs" : {
2403 "main" : {
2404 "snapshot-id" : 1,
2405 "type" : "branch"
2406 }
2407 },
2408 "snapshots" : [ {
2409 "snapshot-id" : 1,
2410 "timestamp-ms" : 1662532818843,
2411 "sequence-number" : 0,
2412 "summary" : {
2413 "operation" : "append",
2414 "spark.app.id" : "local-1662532784305",
2415 "added-data-files" : "4",
2416 "added-records" : "4",
2417 "added-files-size" : "6001"
2418 },
2419 "manifest-list" : "/home/iceberg/warehouse/nyc/taxis/metadata/snap-638933773299822130-1-7e6760f0-4f6c-4b23-b907-0a5a174e3863.avro",
2420 "schema-id" : 0
2421 } ]
2422 }
2423 "#;
2424
2425 let err = serde_json::from_str::<TableMetadata>(data).unwrap_err();
2426 assert!(
2427 err.to_string()
2428 .contains("Current snapshot is not set, but main branch exists")
2429 );
2430 }
2431
2432 #[test]
2433 fn test_branch_snapshot_missing() {
2434 let data = r#"
2435 {
2436 "format-version" : 2,
2437 "table-uuid": "fb072c92-a02b-11e9-ae9c-1bb7bc9eca94",
2438 "location": "s3://b/wh/data.db/table",
2439 "last-sequence-number" : 1,
2440 "last-updated-ms": 1515100955770,
2441 "last-column-id": 1,
2442 "schemas": [
2443 {
2444 "schema-id" : 1,
2445 "type" : "struct",
2446 "fields" :[
2447 {
2448 "id": 1,
2449 "name": "struct_name",
2450 "required": true,
2451 "type": "fixed[1]"
2452 },
2453 {
2454 "id": 4,
2455 "name": "ts",
2456 "required": true,
2457 "type": "timestamp"
2458 }
2459 ]
2460 }
2461 ],
2462 "current-schema-id" : 1,
2463 "partition-specs": [
2464 {
2465 "spec-id": 0,
2466 "fields": [
2467 {
2468 "source-id": 4,
2469 "field-id": 1000,
2470 "name": "ts_day",
2471 "transform": "day"
2472 }
2473 ]
2474 }
2475 ],
2476 "default-spec-id": 0,
2477 "last-partition-id": 1000,
2478 "properties": {
2479 "commit.retry.num-retries": "1"
2480 },
2481 "metadata-log": [
2482 {
2483 "metadata-file": "s3://bucket/.../v1.json",
2484 "timestamp-ms": 1515100
2485 }
2486 ],
2487 "sort-orders": [
2488 {
2489 "order-id": 0,
2490 "fields": []
2491 }
2492 ],
2493 "default-sort-order-id": 0,
2494 "refs" : {
2495 "main" : {
2496 "snapshot-id" : 1,
2497 "type" : "branch"
2498 },
2499 "foo" : {
2500 "snapshot-id" : 2,
2501 "type" : "branch"
2502 }
2503 },
2504 "snapshots" : [ {
2505 "snapshot-id" : 1,
2506 "timestamp-ms" : 1662532818843,
2507 "sequence-number" : 0,
2508 "summary" : {
2509 "operation" : "append",
2510 "spark.app.id" : "local-1662532784305",
2511 "added-data-files" : "4",
2512 "added-records" : "4",
2513 "added-files-size" : "6001"
2514 },
2515 "manifest-list" : "/home/iceberg/warehouse/nyc/taxis/metadata/snap-638933773299822130-1-7e6760f0-4f6c-4b23-b907-0a5a174e3863.avro",
2516 "schema-id" : 0
2517 } ]
2518 }
2519 "#;
2520
2521 let err = serde_json::from_str::<TableMetadata>(data).unwrap_err();
2522 assert!(
2523 err.to_string().contains(
2524 "Snapshot for reference foo does not exist in the existing snapshots list"
2525 )
2526 );
2527 }
2528
2529 #[test]
2530 fn test_v2_wrong_max_snapshot_sequence_number() {
2531 let data = r#"
2532 {
2533 "format-version": 2,
2534 "table-uuid": "9c12d441-03fe-4693-9a96-a0705ddf69c1",
2535 "location": "s3://bucket/test/location",
2536 "last-sequence-number": 1,
2537 "last-updated-ms": 1602638573590,
2538 "last-column-id": 3,
2539 "current-schema-id": 0,
2540 "schemas": [
2541 {
2542 "type": "struct",
2543 "schema-id": 0,
2544 "fields": [
2545 {
2546 "id": 1,
2547 "name": "x",
2548 "required": true,
2549 "type": "long"
2550 }
2551 ]
2552 }
2553 ],
2554 "default-spec-id": 0,
2555 "partition-specs": [
2556 {
2557 "spec-id": 0,
2558 "fields": []
2559 }
2560 ],
2561 "last-partition-id": 1000,
2562 "default-sort-order-id": 0,
2563 "sort-orders": [
2564 {
2565 "order-id": 0,
2566 "fields": []
2567 }
2568 ],
2569 "properties": {},
2570 "current-snapshot-id": 3055729675574597004,
2571 "snapshots": [
2572 {
2573 "snapshot-id": 3055729675574597004,
2574 "timestamp-ms": 1555100955770,
2575 "sequence-number": 4,
2576 "summary": {
2577 "operation": "append"
2578 },
2579 "manifest-list": "s3://a/b/2.avro",
2580 "schema-id": 0
2581 }
2582 ],
2583 "statistics": [],
2584 "snapshot-log": [],
2585 "metadata-log": []
2586 }
2587 "#;
2588
2589 let err = serde_json::from_str::<TableMetadata>(data).unwrap_err();
2590 assert!(err.to_string().contains(
2591 "Invalid snapshot with id 3055729675574597004 and sequence number 4 greater than last sequence number 1"
2592 ));
2593
2594 let data = data.replace(
2596 r#""last-sequence-number": 1,"#,
2597 r#""last-sequence-number": 4,"#,
2598 );
2599 let metadata = serde_json::from_str::<TableMetadata>(data.as_str()).unwrap();
2600 assert_eq!(metadata.last_sequence_number, 4);
2601
2602 let data = data.replace(
2604 r#""last-sequence-number": 4,"#,
2605 r#""last-sequence-number": 5,"#,
2606 );
2607 let metadata = serde_json::from_str::<TableMetadata>(data.as_str()).unwrap();
2608 assert_eq!(metadata.last_sequence_number, 5);
2609 }
2610
2611 #[test]
2612 fn test_statistic_files() {
2613 let data = r#"
2614 {
2615 "format-version": 2,
2616 "table-uuid": "9c12d441-03fe-4693-9a96-a0705ddf69c1",
2617 "location": "s3://bucket/test/location",
2618 "last-sequence-number": 34,
2619 "last-updated-ms": 1602638573590,
2620 "last-column-id": 3,
2621 "current-schema-id": 0,
2622 "schemas": [
2623 {
2624 "type": "struct",
2625 "schema-id": 0,
2626 "fields": [
2627 {
2628 "id": 1,
2629 "name": "x",
2630 "required": true,
2631 "type": "long"
2632 }
2633 ]
2634 }
2635 ],
2636 "default-spec-id": 0,
2637 "partition-specs": [
2638 {
2639 "spec-id": 0,
2640 "fields": []
2641 }
2642 ],
2643 "last-partition-id": 1000,
2644 "default-sort-order-id": 0,
2645 "sort-orders": [
2646 {
2647 "order-id": 0,
2648 "fields": []
2649 }
2650 ],
2651 "properties": {},
2652 "current-snapshot-id": 3055729675574597004,
2653 "snapshots": [
2654 {
2655 "snapshot-id": 3055729675574597004,
2656 "timestamp-ms": 1555100955770,
2657 "sequence-number": 1,
2658 "summary": {
2659 "operation": "append"
2660 },
2661 "manifest-list": "s3://a/b/2.avro",
2662 "schema-id": 0
2663 }
2664 ],
2665 "statistics": [
2666 {
2667 "snapshot-id": 3055729675574597004,
2668 "statistics-path": "s3://a/b/stats.puffin",
2669 "file-size-in-bytes": 413,
2670 "file-footer-size-in-bytes": 42,
2671 "blob-metadata": [
2672 {
2673 "type": "ndv",
2674 "snapshot-id": 3055729675574597004,
2675 "sequence-number": 1,
2676 "fields": [
2677 1
2678 ]
2679 }
2680 ]
2681 }
2682 ],
2683 "snapshot-log": [],
2684 "metadata-log": []
2685 }
2686 "#;
2687
2688 let schema = Schema::builder()
2689 .with_schema_id(0)
2690 .with_fields(vec![Arc::new(NestedField::required(
2691 1,
2692 "x",
2693 Type::Primitive(PrimitiveType::Long),
2694 ))])
2695 .build()
2696 .unwrap();
2697 let partition_spec = PartitionSpec::builder(schema.clone())
2698 .with_spec_id(0)
2699 .build()
2700 .unwrap();
2701 let snapshot = Snapshot::builder()
2702 .with_snapshot_id(3055729675574597004)
2703 .with_timestamp_ms(1555100955770)
2704 .with_sequence_number(1)
2705 .with_manifest_list("s3://a/b/2.avro")
2706 .with_schema_id(0)
2707 .with_summary(Summary {
2708 operation: Operation::Append,
2709 additional_properties: HashMap::new(),
2710 })
2711 .build();
2712
2713 let default_partition_type = partition_spec.partition_type(&schema).unwrap();
2714 let expected = TableMetadata {
2715 format_version: FormatVersion::V2,
2716 table_uuid: Uuid::parse_str("9c12d441-03fe-4693-9a96-a0705ddf69c1").unwrap(),
2717 location: "s3://bucket/test/location".to_string(),
2718 last_updated_ms: 1602638573590,
2719 last_column_id: 3,
2720 schemas: HashMap::from_iter(vec![(0, Arc::new(schema))]),
2721 current_schema_id: 0,
2722 partition_specs: HashMap::from_iter(vec![(0, partition_spec.clone().into())]),
2723 default_partition_type,
2724 default_spec: Arc::new(partition_spec),
2725 last_partition_id: 1000,
2726 default_sort_order_id: 0,
2727 sort_orders: HashMap::from_iter(vec![(0, SortOrder::unsorted_order().into())]),
2728 snapshots: HashMap::from_iter(vec![(3055729675574597004, Arc::new(snapshot))]),
2729 current_snapshot_id: Some(3055729675574597004),
2730 last_sequence_number: 34,
2731 properties: HashMap::new(),
2732 snapshot_log: Vec::new(),
2733 metadata_log: Vec::new(),
2734 statistics: HashMap::from_iter(vec![(3055729675574597004, StatisticsFile {
2735 snapshot_id: 3055729675574597004,
2736 statistics_path: "s3://a/b/stats.puffin".to_string(),
2737 file_size_in_bytes: 413,
2738 file_footer_size_in_bytes: 42,
2739 key_metadata: None,
2740 blob_metadata: vec![BlobMetadata {
2741 snapshot_id: 3055729675574597004,
2742 sequence_number: 1,
2743 fields: vec![1],
2744 r#type: "ndv".to_string(),
2745 properties: HashMap::new(),
2746 }],
2747 })]),
2748 partition_statistics: HashMap::new(),
2749 refs: HashMap::from_iter(vec![("main".to_string(), SnapshotReference {
2750 snapshot_id: 3055729675574597004,
2751 retention: SnapshotRetention::Branch {
2752 min_snapshots_to_keep: None,
2753 max_snapshot_age_ms: None,
2754 max_ref_age_ms: None,
2755 },
2756 })]),
2757 encryption_keys: HashMap::new(),
2758 next_row_id: INITIAL_ROW_ID,
2759 };
2760
2761 check_table_metadata_serde(data, expected);
2762 }
2763
2764 #[test]
2765 fn test_partition_statistics_file() {
2766 let data = r#"
2767 {
2768 "format-version": 2,
2769 "table-uuid": "9c12d441-03fe-4693-9a96-a0705ddf69c1",
2770 "location": "s3://bucket/test/location",
2771 "last-sequence-number": 34,
2772 "last-updated-ms": 1602638573590,
2773 "last-column-id": 3,
2774 "current-schema-id": 0,
2775 "schemas": [
2776 {
2777 "type": "struct",
2778 "schema-id": 0,
2779 "fields": [
2780 {
2781 "id": 1,
2782 "name": "x",
2783 "required": true,
2784 "type": "long"
2785 }
2786 ]
2787 }
2788 ],
2789 "default-spec-id": 0,
2790 "partition-specs": [
2791 {
2792 "spec-id": 0,
2793 "fields": []
2794 }
2795 ],
2796 "last-partition-id": 1000,
2797 "default-sort-order-id": 0,
2798 "sort-orders": [
2799 {
2800 "order-id": 0,
2801 "fields": []
2802 }
2803 ],
2804 "properties": {},
2805 "current-snapshot-id": 3055729675574597004,
2806 "snapshots": [
2807 {
2808 "snapshot-id": 3055729675574597004,
2809 "timestamp-ms": 1555100955770,
2810 "sequence-number": 1,
2811 "summary": {
2812 "operation": "append"
2813 },
2814 "manifest-list": "s3://a/b/2.avro",
2815 "schema-id": 0
2816 }
2817 ],
2818 "partition-statistics": [
2819 {
2820 "snapshot-id": 3055729675574597004,
2821 "statistics-path": "s3://a/b/partition-stats.parquet",
2822 "file-size-in-bytes": 43
2823 }
2824 ],
2825 "snapshot-log": [],
2826 "metadata-log": []
2827 }
2828 "#;
2829
2830 let schema = Schema::builder()
2831 .with_schema_id(0)
2832 .with_fields(vec![Arc::new(NestedField::required(
2833 1,
2834 "x",
2835 Type::Primitive(PrimitiveType::Long),
2836 ))])
2837 .build()
2838 .unwrap();
2839 let partition_spec = PartitionSpec::builder(schema.clone())
2840 .with_spec_id(0)
2841 .build()
2842 .unwrap();
2843 let snapshot = Snapshot::builder()
2844 .with_snapshot_id(3055729675574597004)
2845 .with_timestamp_ms(1555100955770)
2846 .with_sequence_number(1)
2847 .with_manifest_list("s3://a/b/2.avro")
2848 .with_schema_id(0)
2849 .with_summary(Summary {
2850 operation: Operation::Append,
2851 additional_properties: HashMap::new(),
2852 })
2853 .build();
2854
2855 let default_partition_type = partition_spec.partition_type(&schema).unwrap();
2856 let expected = TableMetadata {
2857 format_version: FormatVersion::V2,
2858 table_uuid: Uuid::parse_str("9c12d441-03fe-4693-9a96-a0705ddf69c1").unwrap(),
2859 location: "s3://bucket/test/location".to_string(),
2860 last_updated_ms: 1602638573590,
2861 last_column_id: 3,
2862 schemas: HashMap::from_iter(vec![(0, Arc::new(schema))]),
2863 current_schema_id: 0,
2864 partition_specs: HashMap::from_iter(vec![(0, partition_spec.clone().into())]),
2865 default_spec: Arc::new(partition_spec),
2866 default_partition_type,
2867 last_partition_id: 1000,
2868 default_sort_order_id: 0,
2869 sort_orders: HashMap::from_iter(vec![(0, SortOrder::unsorted_order().into())]),
2870 snapshots: HashMap::from_iter(vec![(3055729675574597004, Arc::new(snapshot))]),
2871 current_snapshot_id: Some(3055729675574597004),
2872 last_sequence_number: 34,
2873 properties: HashMap::new(),
2874 snapshot_log: Vec::new(),
2875 metadata_log: Vec::new(),
2876 statistics: HashMap::new(),
2877 partition_statistics: HashMap::from_iter(vec![(
2878 3055729675574597004,
2879 PartitionStatisticsFile {
2880 snapshot_id: 3055729675574597004,
2881 statistics_path: "s3://a/b/partition-stats.parquet".to_string(),
2882 file_size_in_bytes: 43,
2883 },
2884 )]),
2885 refs: HashMap::from_iter(vec![("main".to_string(), SnapshotReference {
2886 snapshot_id: 3055729675574597004,
2887 retention: SnapshotRetention::Branch {
2888 min_snapshots_to_keep: None,
2889 max_snapshot_age_ms: None,
2890 max_ref_age_ms: None,
2891 },
2892 })]),
2893 encryption_keys: HashMap::new(),
2894 next_row_id: INITIAL_ROW_ID,
2895 };
2896
2897 check_table_metadata_serde(data, expected);
2898 }
2899
2900 #[test]
2901 fn test_invalid_table_uuid() -> Result<()> {
2902 let data = r#"
2903 {
2904 "format-version" : 2,
2905 "table-uuid": "xxxx"
2906 }
2907 "#;
2908 assert!(serde_json::from_str::<TableMetadata>(data).is_err());
2909 Ok(())
2910 }
2911
2912 #[test]
2913 fn test_deserialize_table_data_v2_invalid_format_version() -> Result<()> {
2914 let data = r#"
2915 {
2916 "format-version" : 1
2917 }
2918 "#;
2919 assert!(serde_json::from_str::<TableMetadata>(data).is_err());
2920 Ok(())
2921 }
2922
2923 #[test]
2924 fn test_table_metadata_v3_valid_minimal() {
2925 let metadata_str =
2926 fs::read_to_string("testdata/table_metadata/TableMetadataV3ValidMinimal.json").unwrap();
2927
2928 let table_metadata = serde_json::from_str::<TableMetadata>(&metadata_str).unwrap();
2929 assert_eq!(table_metadata.format_version, FormatVersion::V3);
2930
2931 let schema = Schema::builder()
2932 .with_schema_id(0)
2933 .with_fields(vec![
2934 Arc::new(
2935 NestedField::required(1, "x", Type::Primitive(PrimitiveType::Long))
2936 .with_initial_default(Literal::Primitive(PrimitiveLiteral::Long(1)))
2937 .with_write_default(Literal::Primitive(PrimitiveLiteral::Long(1))),
2938 ),
2939 Arc::new(
2940 NestedField::required(2, "y", Type::Primitive(PrimitiveType::Long))
2941 .with_doc("comment"),
2942 ),
2943 Arc::new(NestedField::required(
2944 3,
2945 "z",
2946 Type::Primitive(PrimitiveType::Long),
2947 )),
2948 ])
2949 .build()
2950 .unwrap();
2951
2952 let partition_spec = PartitionSpec::builder(schema.clone())
2953 .with_spec_id(0)
2954 .add_unbound_field(
2955 UnboundPartitionField::builder()
2956 .source_ids(vec![1])
2957 .field_id(1000)
2958 .name("x".to_string())
2959 .transform(Transform::Identity)
2960 .build()
2961 .unwrap(),
2962 )
2963 .unwrap()
2964 .build()
2965 .unwrap();
2966
2967 let sort_order = SortOrder::builder()
2968 .with_order_id(3)
2969 .with_sort_field(SortField {
2970 source_id: 2,
2971 transform: Transform::Identity,
2972 direction: SortDirection::Ascending,
2973 null_order: NullOrder::First,
2974 })
2975 .with_sort_field(SortField {
2976 source_id: 3,
2977 transform: Transform::Bucket(4),
2978 direction: SortDirection::Descending,
2979 null_order: NullOrder::Last,
2980 })
2981 .build_unbound()
2982 .unwrap();
2983
2984 let default_partition_type = partition_spec.partition_type(&schema).unwrap();
2985 let expected = TableMetadata {
2986 format_version: FormatVersion::V3,
2987 table_uuid: Uuid::parse_str("9c12d441-03fe-4693-9a96-a0705ddf69c1").unwrap(),
2988 location: "s3://bucket/test/location".to_string(),
2989 last_updated_ms: 1602638573590,
2990 last_column_id: 3,
2991 schemas: HashMap::from_iter(vec![(0, Arc::new(schema))]),
2992 current_schema_id: 0,
2993 partition_specs: HashMap::from_iter(vec![(0, partition_spec.clone().into())]),
2994 default_spec: Arc::new(partition_spec),
2995 default_partition_type,
2996 last_partition_id: 1000,
2997 default_sort_order_id: 3,
2998 sort_orders: HashMap::from_iter(vec![(3, sort_order.into())]),
2999 snapshots: HashMap::default(),
3000 current_snapshot_id: None,
3001 last_sequence_number: 34,
3002 properties: HashMap::new(),
3003 snapshot_log: Vec::new(),
3004 metadata_log: Vec::new(),
3005 refs: HashMap::new(),
3006 statistics: HashMap::new(),
3007 partition_statistics: HashMap::new(),
3008 encryption_keys: HashMap::new(),
3009 next_row_id: 0, };
3011
3012 check_table_metadata_serde(&metadata_str, expected);
3013 }
3014
3015 #[test]
3016 fn test_table_metadata_v2_file_valid() {
3017 let metadata =
3018 fs::read_to_string("testdata/table_metadata/TableMetadataV2Valid.json").unwrap();
3019
3020 let schema1 = Schema::builder()
3021 .with_schema_id(0)
3022 .with_fields(vec![Arc::new(NestedField::required(
3023 1,
3024 "x",
3025 Type::Primitive(PrimitiveType::Long),
3026 ))])
3027 .build()
3028 .unwrap();
3029
3030 let schema2 = Schema::builder()
3031 .with_schema_id(1)
3032 .with_fields(vec![
3033 Arc::new(NestedField::required(
3034 1,
3035 "x",
3036 Type::Primitive(PrimitiveType::Long),
3037 )),
3038 Arc::new(
3039 NestedField::required(2, "y", Type::Primitive(PrimitiveType::Long))
3040 .with_doc("comment"),
3041 ),
3042 Arc::new(NestedField::required(
3043 3,
3044 "z",
3045 Type::Primitive(PrimitiveType::Long),
3046 )),
3047 ])
3048 .with_identifier_field_ids(vec![1, 2])
3049 .build()
3050 .unwrap();
3051
3052 let partition_spec = PartitionSpec::builder(schema2.clone())
3053 .with_spec_id(0)
3054 .add_unbound_field(
3055 UnboundPartitionField::builder()
3056 .source_ids(vec![1])
3057 .field_id(1000)
3058 .name("x".to_string())
3059 .transform(Transform::Identity)
3060 .build()
3061 .unwrap(),
3062 )
3063 .unwrap()
3064 .build()
3065 .unwrap();
3066
3067 let sort_order = SortOrder::builder()
3068 .with_order_id(3)
3069 .with_sort_field(SortField {
3070 source_id: 2,
3071 transform: Transform::Identity,
3072 direction: SortDirection::Ascending,
3073 null_order: NullOrder::First,
3074 })
3075 .with_sort_field(SortField {
3076 source_id: 3,
3077 transform: Transform::Bucket(4),
3078 direction: SortDirection::Descending,
3079 null_order: NullOrder::Last,
3080 })
3081 .build_unbound()
3082 .unwrap();
3083
3084 let snapshot1 = Snapshot::builder()
3085 .with_snapshot_id(3051729675574597004)
3086 .with_timestamp_ms(1515100955770)
3087 .with_sequence_number(0)
3088 .with_manifest_list("s3://a/b/1.avro")
3089 .with_summary(Summary {
3090 operation: Operation::Append,
3091 additional_properties: HashMap::new(),
3092 })
3093 .build();
3094
3095 let snapshot2 = Snapshot::builder()
3096 .with_snapshot_id(3055729675574597004)
3097 .with_parent_snapshot_id(Some(3051729675574597004))
3098 .with_timestamp_ms(1555100955770)
3099 .with_sequence_number(1)
3100 .with_schema_id(1)
3101 .with_manifest_list("s3://a/b/2.avro")
3102 .with_summary(Summary {
3103 operation: Operation::Append,
3104 additional_properties: HashMap::new(),
3105 })
3106 .build();
3107
3108 let default_partition_type = partition_spec.partition_type(&schema2).unwrap();
3109 let expected = TableMetadata {
3110 format_version: FormatVersion::V2,
3111 table_uuid: Uuid::parse_str("9c12d441-03fe-4693-9a96-a0705ddf69c1").unwrap(),
3112 location: "s3://bucket/test/location".to_string(),
3113 last_updated_ms: 1602638573590,
3114 last_column_id: 3,
3115 schemas: HashMap::from_iter(vec![(0, Arc::new(schema1)), (1, Arc::new(schema2))]),
3116 current_schema_id: 1,
3117 partition_specs: HashMap::from_iter(vec![(0, partition_spec.clone().into())]),
3118 default_spec: Arc::new(partition_spec),
3119 default_partition_type,
3120 last_partition_id: 1000,
3121 default_sort_order_id: 3,
3122 sort_orders: HashMap::from_iter(vec![(3, sort_order.into())]),
3123 snapshots: HashMap::from_iter(vec![
3124 (3051729675574597004, Arc::new(snapshot1)),
3125 (3055729675574597004, Arc::new(snapshot2)),
3126 ]),
3127 current_snapshot_id: Some(3055729675574597004),
3128 last_sequence_number: 34,
3129 properties: HashMap::new(),
3130 snapshot_log: vec![
3131 SnapshotLog {
3132 snapshot_id: 3051729675574597004,
3133 timestamp_ms: 1515100955770,
3134 },
3135 SnapshotLog {
3136 snapshot_id: 3055729675574597004,
3137 timestamp_ms: 1555100955770,
3138 },
3139 ],
3140 metadata_log: Vec::new(),
3141 refs: HashMap::from_iter(vec![("main".to_string(), SnapshotReference {
3142 snapshot_id: 3055729675574597004,
3143 retention: SnapshotRetention::Branch {
3144 min_snapshots_to_keep: None,
3145 max_snapshot_age_ms: None,
3146 max_ref_age_ms: None,
3147 },
3148 })]),
3149 statistics: HashMap::new(),
3150 partition_statistics: HashMap::new(),
3151 encryption_keys: HashMap::new(),
3152 next_row_id: INITIAL_ROW_ID,
3153 };
3154
3155 check_table_metadata_serde(&metadata, expected);
3156 }
3157
3158 #[test]
3159 fn test_table_metadata_v2_file_valid_minimal() {
3160 let metadata =
3161 fs::read_to_string("testdata/table_metadata/TableMetadataV2ValidMinimal.json").unwrap();
3162
3163 let schema = Schema::builder()
3164 .with_schema_id(0)
3165 .with_fields(vec![
3166 Arc::new(NestedField::required(
3167 1,
3168 "x",
3169 Type::Primitive(PrimitiveType::Long),
3170 )),
3171 Arc::new(
3172 NestedField::required(2, "y", Type::Primitive(PrimitiveType::Long))
3173 .with_doc("comment"),
3174 ),
3175 Arc::new(NestedField::required(
3176 3,
3177 "z",
3178 Type::Primitive(PrimitiveType::Long),
3179 )),
3180 ])
3181 .build()
3182 .unwrap();
3183
3184 let partition_spec = PartitionSpec::builder(schema.clone())
3185 .with_spec_id(0)
3186 .add_unbound_field(
3187 UnboundPartitionField::builder()
3188 .source_ids(vec![1])
3189 .field_id(1000)
3190 .name("x".to_string())
3191 .transform(Transform::Identity)
3192 .build()
3193 .unwrap(),
3194 )
3195 .unwrap()
3196 .build()
3197 .unwrap();
3198
3199 let sort_order = SortOrder::builder()
3200 .with_order_id(3)
3201 .with_sort_field(SortField {
3202 source_id: 2,
3203 transform: Transform::Identity,
3204 direction: SortDirection::Ascending,
3205 null_order: NullOrder::First,
3206 })
3207 .with_sort_field(SortField {
3208 source_id: 3,
3209 transform: Transform::Bucket(4),
3210 direction: SortDirection::Descending,
3211 null_order: NullOrder::Last,
3212 })
3213 .build_unbound()
3214 .unwrap();
3215
3216 let default_partition_type = partition_spec.partition_type(&schema).unwrap();
3217 let expected = TableMetadata {
3218 format_version: FormatVersion::V2,
3219 table_uuid: Uuid::parse_str("9c12d441-03fe-4693-9a96-a0705ddf69c1").unwrap(),
3220 location: "s3://bucket/test/location".to_string(),
3221 last_updated_ms: 1602638573590,
3222 last_column_id: 3,
3223 schemas: HashMap::from_iter(vec![(0, Arc::new(schema))]),
3224 current_schema_id: 0,
3225 partition_specs: HashMap::from_iter(vec![(0, partition_spec.clone().into())]),
3226 default_partition_type,
3227 default_spec: Arc::new(partition_spec),
3228 last_partition_id: 1000,
3229 default_sort_order_id: 3,
3230 sort_orders: HashMap::from_iter(vec![(3, sort_order.into())]),
3231 snapshots: HashMap::default(),
3232 current_snapshot_id: None,
3233 last_sequence_number: 34,
3234 properties: HashMap::new(),
3235 snapshot_log: vec![],
3236 metadata_log: Vec::new(),
3237 refs: HashMap::new(),
3238 statistics: HashMap::new(),
3239 partition_statistics: HashMap::new(),
3240 encryption_keys: HashMap::new(),
3241 next_row_id: INITIAL_ROW_ID,
3242 };
3243
3244 check_table_metadata_serde(&metadata, expected);
3245 }
3246
3247 #[test]
3248 fn test_table_metadata_v1_file_valid() {
3249 let metadata =
3250 fs::read_to_string("testdata/table_metadata/TableMetadataV1Valid.json").unwrap();
3251
3252 let schema = Schema::builder()
3253 .with_schema_id(0)
3254 .with_fields(vec![
3255 Arc::new(NestedField::required(
3256 1,
3257 "x",
3258 Type::Primitive(PrimitiveType::Long),
3259 )),
3260 Arc::new(
3261 NestedField::required(2, "y", Type::Primitive(PrimitiveType::Long))
3262 .with_doc("comment"),
3263 ),
3264 Arc::new(NestedField::required(
3265 3,
3266 "z",
3267 Type::Primitive(PrimitiveType::Long),
3268 )),
3269 ])
3270 .build()
3271 .unwrap();
3272
3273 let partition_spec = PartitionSpec::builder(schema.clone())
3274 .with_spec_id(0)
3275 .add_unbound_field(
3276 UnboundPartitionField::builder()
3277 .source_ids(vec![1])
3278 .field_id(1000)
3279 .name("x".to_string())
3280 .transform(Transform::Identity)
3281 .build()
3282 .unwrap(),
3283 )
3284 .unwrap()
3285 .build()
3286 .unwrap();
3287
3288 let default_partition_type = partition_spec.partition_type(&schema).unwrap();
3289 let expected = TableMetadata {
3290 format_version: FormatVersion::V1,
3291 table_uuid: Uuid::parse_str("d20125c8-7284-442c-9aea-15fee620737c").unwrap(),
3292 location: "s3://bucket/test/location".to_string(),
3293 last_updated_ms: 1602638573874,
3294 last_column_id: 3,
3295 schemas: HashMap::from_iter(vec![(0, Arc::new(schema))]),
3296 current_schema_id: 0,
3297 partition_specs: HashMap::from_iter(vec![(0, partition_spec.clone().into())]),
3298 default_spec: Arc::new(partition_spec),
3299 default_partition_type,
3300 last_partition_id: 0,
3301 default_sort_order_id: 0,
3302 sort_orders: HashMap::from_iter(vec![(0, SortOrder::unsorted_order().into())]),
3304 snapshots: HashMap::new(),
3305 current_snapshot_id: None,
3306 last_sequence_number: 0,
3307 properties: HashMap::new(),
3308 snapshot_log: vec![],
3309 metadata_log: Vec::new(),
3310 refs: HashMap::new(),
3311 statistics: HashMap::new(),
3312 partition_statistics: HashMap::new(),
3313 encryption_keys: HashMap::new(),
3314 next_row_id: INITIAL_ROW_ID,
3315 };
3316
3317 check_table_metadata_serde(&metadata, expected);
3318 }
3319
3320 #[test]
3321 fn test_empty_snapshot_id_is_normalized_to_none() {
3322 let metadata =
3323 fs::read_to_string("testdata/table_metadata/TableMetadataV1Valid.json").unwrap();
3324 let deserialized: TableMetadata = serde_json::from_str(&metadata).unwrap();
3325 assert_eq!(
3326 deserialized.current_snapshot_id(),
3327 None,
3328 "current_snapshot_id of -1 should be deserialized as None"
3329 );
3330 }
3331
3332 #[test]
3333 fn test_table_metadata_v1_compat() {
3334 let metadata =
3335 fs::read_to_string("testdata/table_metadata/TableMetadataV1Compat.json").unwrap();
3336
3337 let desered_type: TableMetadata = serde_json::from_str(&metadata)
3339 .expect("Failed to deserialize TableMetadataV1Compat.json");
3340
3341 assert_eq!(desered_type.format_version(), FormatVersion::V1);
3343 assert_eq!(
3344 desered_type.uuid(),
3345 Uuid::parse_str("3276010d-7b1d-488c-98d8-9025fc4fde6b").unwrap()
3346 );
3347 assert_eq!(
3348 desered_type.location(),
3349 "s3://bucket/warehouse/iceberg/glue.db/table_name"
3350 );
3351 assert_eq!(desered_type.last_updated_ms(), 1727773114005);
3352 assert_eq!(desered_type.current_schema_id(), 0);
3353 }
3354
3355 #[test]
3356 fn test_table_metadata_v1_schemas_without_current_id() {
3357 let metadata = fs::read_to_string(
3358 "testdata/table_metadata/TableMetadataV1SchemasWithoutCurrentId.json",
3359 )
3360 .unwrap();
3361
3362 let desered_type: TableMetadata = serde_json::from_str(&metadata)
3364 .expect("Failed to deserialize TableMetadataV1SchemasWithoutCurrentId.json");
3365
3366 assert_eq!(desered_type.format_version(), FormatVersion::V1);
3368 assert_eq!(
3369 desered_type.uuid(),
3370 Uuid::parse_str("d20125c8-7284-442c-9aea-15fee620737c").unwrap()
3371 );
3372
3373 let schema = desered_type.current_schema();
3375 assert_eq!(schema.as_struct().fields().len(), 3);
3376 assert_eq!(schema.as_struct().fields()[0].name, "x");
3377 assert_eq!(schema.as_struct().fields()[1].name, "y");
3378 assert_eq!(schema.as_struct().fields()[2].name, "z");
3379 }
3380
3381 #[test]
3382 fn test_table_metadata_v1_no_valid_schema() {
3383 let metadata =
3384 fs::read_to_string("testdata/table_metadata/TableMetadataV1NoValidSchema.json")
3385 .unwrap();
3386
3387 let desered: Result<TableMetadata, serde_json::Error> = serde_json::from_str(&metadata);
3389
3390 assert!(desered.is_err());
3391 let error_message = desered.unwrap_err().to_string();
3392 assert!(
3393 error_message.contains("No valid schema configuration found"),
3394 "Expected error about no valid schema configuration, got: {error_message}"
3395 );
3396 }
3397
3398 #[test]
3399 fn test_table_metadata_v1_partition_specs_without_default_id() {
3400 let metadata = fs::read_to_string(
3401 "testdata/table_metadata/TableMetadataV1PartitionSpecsWithoutDefaultId.json",
3402 )
3403 .unwrap();
3404
3405 let desered_type: TableMetadata = serde_json::from_str(&metadata)
3407 .expect("Failed to deserialize TableMetadataV1PartitionSpecsWithoutDefaultId.json");
3408
3409 assert_eq!(desered_type.format_version(), FormatVersion::V1);
3411 assert_eq!(
3412 desered_type.uuid(),
3413 Uuid::parse_str("d20125c8-7284-442c-9aea-15fee620737c").unwrap()
3414 );
3415
3416 assert_eq!(desered_type.default_partition_spec_id(), 2); assert_eq!(desered_type.partition_specs.len(), 2);
3419
3420 let default_spec = &desered_type.default_spec;
3422 assert_eq!(default_spec.spec_id(), 2);
3423 assert_eq!(default_spec.fields().len(), 1);
3424 assert_eq!(default_spec.fields()[0].name, "y");
3425 assert_eq!(default_spec.fields()[0].transform, Transform::Identity);
3426 assert_eq!(default_spec.fields()[0].source_id, 2);
3427 }
3428
3429 #[test]
3430 fn test_table_metadata_v2_schema_not_found() {
3431 let metadata =
3432 fs::read_to_string("testdata/table_metadata/TableMetadataV2CurrentSchemaNotFound.json")
3433 .unwrap();
3434
3435 let desered: Result<TableMetadata, serde_json::Error> = serde_json::from_str(&metadata);
3436
3437 assert_eq!(
3438 desered.unwrap_err().to_string(),
3439 "DataInvalid => No schema exists with the current schema id 2."
3440 )
3441 }
3442
3443 #[test]
3444 fn test_table_metadata_v2_missing_sort_order() {
3445 let metadata =
3446 fs::read_to_string("testdata/table_metadata/TableMetadataV2MissingSortOrder.json")
3447 .unwrap();
3448
3449 let desered: Result<TableMetadata, serde_json::Error> = serde_json::from_str(&metadata);
3450
3451 assert_eq!(
3452 desered.unwrap_err().to_string(),
3453 "data did not match any variant of untagged enum TableMetadataEnum"
3454 )
3455 }
3456
3457 #[test]
3458 fn test_table_metadata_v2_missing_partition_specs() {
3459 let metadata =
3460 fs::read_to_string("testdata/table_metadata/TableMetadataV2MissingPartitionSpecs.json")
3461 .unwrap();
3462
3463 let desered: Result<TableMetadata, serde_json::Error> = serde_json::from_str(&metadata);
3464
3465 assert_eq!(
3466 desered.unwrap_err().to_string(),
3467 "data did not match any variant of untagged enum TableMetadataEnum"
3468 )
3469 }
3470
3471 #[test]
3472 fn test_table_metadata_v2_missing_last_partition_id() {
3473 let metadata = fs::read_to_string(
3474 "testdata/table_metadata/TableMetadataV2MissingLastPartitionId.json",
3475 )
3476 .unwrap();
3477
3478 let desered: Result<TableMetadata, serde_json::Error> = serde_json::from_str(&metadata);
3479
3480 assert_eq!(
3481 desered.unwrap_err().to_string(),
3482 "data did not match any variant of untagged enum TableMetadataEnum"
3483 )
3484 }
3485
3486 #[test]
3487 fn test_table_metadata_v2_missing_schemas() {
3488 let metadata =
3489 fs::read_to_string("testdata/table_metadata/TableMetadataV2MissingSchemas.json")
3490 .unwrap();
3491
3492 let desered: Result<TableMetadata, serde_json::Error> = serde_json::from_str(&metadata);
3493
3494 assert_eq!(
3495 desered.unwrap_err().to_string(),
3496 "data did not match any variant of untagged enum TableMetadataEnum"
3497 )
3498 }
3499
3500 #[test]
3501 fn test_table_metadata_v2_unsupported_version() {
3502 let metadata =
3503 fs::read_to_string("testdata/table_metadata/TableMetadataUnsupportedVersion.json")
3504 .unwrap();
3505
3506 let desered: Result<TableMetadata, serde_json::Error> = serde_json::from_str(&metadata);
3507
3508 assert_eq!(
3509 desered.unwrap_err().to_string(),
3510 "data did not match any variant of untagged enum TableMetadataEnum"
3511 )
3512 }
3513
3514 #[test]
3515 fn test_order_of_format_version() {
3516 assert!(FormatVersion::V1 < FormatVersion::V2);
3517 assert_eq!(FormatVersion::V1, FormatVersion::V1);
3518 assert_eq!(FormatVersion::V2, FormatVersion::V2);
3519 }
3520
3521 #[test]
3522 fn test_deserialize_default_partition_spec_with_dropped_source() {
3523 for version in [1, 2, 3] {
3524 for transform in ["identity", "truncate[4]", "bucket[16]"] {
3525 let mut metadata = serde_json::json!({
3526 "format-version": version,
3527 "table-uuid": "9c12d441-03fe-4693-9a96-a0705ddf69c1",
3528 "location": "s3://bucket/table",
3529 "last-sequence-number": 0,
3530 "last-updated-ms": 1602638573590_i64,
3531 "last-column-id": 2,
3532 "current-schema-id": 1,
3533 "schemas": [
3534 {
3535 "schema-id": 0,
3536 "type": "struct",
3537 "fields": [
3538 {"id": 1, "name": "x", "required": false, "type": "long"},
3539 {"id": 2, "name": "y", "required": false, "type": "long"}
3540 ]
3541 },
3542 {
3543 "schema-id": 1,
3544 "type": "struct",
3545 "fields": [
3546 {"id": 2, "name": "y", "required": false, "type": "long"}
3547 ]
3548 }
3549 ],
3550 "default-spec-id": 0,
3551 "partition-specs": [{"spec-id": 0, "fields": [{
3552 "source-id": 1, "field-id": 1000, "name": "x_partition",
3553 "transform": transform
3554 }]}],
3555 "last-partition-id": 1000,
3556 "default-sort-order-id": 0,
3557 "sort-orders": [{"order-id": 0, "fields": []}],
3558 "next-row-id": 0
3559 });
3560
3561 let err = serde_json::from_value::<TableMetadata>(metadata.clone()).unwrap_err();
3564 assert!(err.to_string().contains(
3565 "Default partition spec 0 references missing source field 1 in current schema 1"
3566 ));
3567
3568 metadata["partition-specs"]
3570 .as_array_mut()
3571 .unwrap()
3572 .push(serde_json::json!({"spec-id": 1, "fields": []}));
3573 metadata["default-spec-id"] = serde_json::json!(1);
3574 serde_json::from_value::<TableMetadata>(metadata.clone()).unwrap();
3575
3576 metadata["default-spec-id"] = serde_json::json!(0);
3578 metadata["partition-specs"][0]["fields"][0]["transform"] =
3579 serde_json::json!("void");
3580 serde_json::from_value::<TableMetadata>(metadata).unwrap();
3581 }
3582 }
3583 }
3584
3585 #[test]
3586 fn test_default_partition_spec() {
3587 let default_spec_id = 1234;
3588 let mut table_meta_data = get_test_table_metadata("TableMetadataV2Valid.json");
3589 let partition_spec = PartitionSpec::unpartition_spec();
3590 table_meta_data.default_spec = partition_spec.clone().into();
3591 table_meta_data
3592 .partition_specs
3593 .insert(default_spec_id, Arc::new(partition_spec));
3594
3595 assert_eq!(
3596 (*table_meta_data.default_partition_spec().clone()).clone(),
3597 (*table_meta_data
3598 .partition_spec_by_id(default_spec_id)
3599 .unwrap()
3600 .clone())
3601 .clone()
3602 );
3603 }
3604 #[test]
3605 fn test_default_sort_order() {
3606 let default_sort_order_id = 1234;
3607 let mut table_meta_data = get_test_table_metadata("TableMetadataV2Valid.json");
3608 table_meta_data.default_sort_order_id = default_sort_order_id;
3609 table_meta_data
3610 .sort_orders
3611 .insert(default_sort_order_id, Arc::new(SortOrder::default()));
3612
3613 assert_eq!(
3614 table_meta_data.default_sort_order(),
3615 table_meta_data
3616 .sort_orders
3617 .get(&default_sort_order_id)
3618 .unwrap()
3619 )
3620 }
3621
3622 #[test]
3623 fn test_table_metadata_builder_from_table_creation() {
3624 let table_creation = TableCreation::builder()
3625 .location("s3://db/table".to_string())
3626 .name("table".to_string())
3627 .properties(HashMap::new())
3628 .schema(Schema::builder().build().unwrap())
3629 .build();
3630 let table_metadata = TableMetadataBuilder::from_table_creation(table_creation)
3631 .unwrap()
3632 .build()
3633 .unwrap()
3634 .metadata;
3635 assert_eq!(table_metadata.location, "s3://db/table");
3636 assert_eq!(table_metadata.schemas.len(), 1);
3637 assert_eq!(
3638 table_metadata
3639 .schemas
3640 .get(&0)
3641 .unwrap()
3642 .as_struct()
3643 .fields()
3644 .len(),
3645 0
3646 );
3647 assert_eq!(table_metadata.properties.len(), 0);
3648 assert_eq!(
3649 table_metadata.partition_specs,
3650 HashMap::from([(
3651 0,
3652 Arc::new(
3653 PartitionSpec::builder(table_metadata.schemas.get(&0).unwrap().clone())
3654 .with_spec_id(0)
3655 .build()
3656 .unwrap()
3657 )
3658 )])
3659 );
3660 assert_eq!(
3661 table_metadata.sort_orders,
3662 HashMap::from([(
3663 0,
3664 Arc::new(SortOrder {
3665 order_id: 0,
3666 fields: vec![]
3667 })
3668 )])
3669 );
3670 }
3671
3672 #[tokio::test]
3673 async fn test_table_metadata_read_write() {
3674 let temp_dir = TempDir::new().unwrap();
3676 let temp_path = temp_dir.path().to_str().unwrap();
3677
3678 let file_io = FileIO::new_with_fs();
3680
3681 let original_metadata: TableMetadata =
3683 get_test_table_metadata_at("TableMetadataV2Valid.json", temp_path);
3684
3685 let metadata_location =
3687 MetadataLocation::try_new_with_metadata(&original_metadata).unwrap();
3688 let metadata_location_str = metadata_location.to_string();
3689
3690 original_metadata
3692 .write_to(&file_io, &metadata_location)
3693 .await
3694 .unwrap();
3695
3696 assert!(fs::metadata(&metadata_location_str).is_ok());
3698
3699 let read_metadata = TableMetadata::read_from(&file_io, &metadata_location_str)
3701 .await
3702 .unwrap();
3703
3704 assert_eq!(read_metadata, original_metadata);
3706 }
3707
3708 #[tokio::test]
3709 async fn test_table_metadata_read_compressed() {
3710 let temp_dir = TempDir::new().unwrap();
3711 let metadata_location = temp_dir.path().join("v1.gz.metadata.json");
3712
3713 let original_metadata: TableMetadata = get_test_table_metadata("TableMetadataV2Valid.json");
3714 let json = serde_json::to_string(&original_metadata).unwrap();
3715
3716 let compressed = CompressionCodec::gzip_default()
3717 .compress(json.into_bytes())
3718 .expect("failed to compress metadata");
3719 fs::write(&metadata_location, &compressed).expect("failed to write metadata");
3720
3721 let file_io = FileIO::new_with_fs();
3723 let metadata_location = metadata_location.to_str().unwrap();
3724 let read_metadata = TableMetadata::read_from(&file_io, metadata_location)
3725 .await
3726 .unwrap();
3727
3728 assert_eq!(read_metadata, original_metadata);
3730 }
3731
3732 #[tokio::test]
3733 async fn test_table_metadata_read_nonexistent_file() {
3734 let file_io = FileIO::new_with_fs();
3736
3737 let result = TableMetadata::read_from(&file_io, "/nonexistent/path/metadata.json").await;
3739
3740 assert!(result.is_err());
3742 }
3743
3744 #[tokio::test]
3745 async fn test_table_metadata_write_with_gzip_compression() {
3746 let temp_dir = TempDir::new().unwrap();
3747 let temp_path = temp_dir.path().to_str().unwrap();
3748 let file_io = FileIO::new_with_fs();
3749
3750 let original_metadata: TableMetadata =
3752 get_test_table_metadata_at("TableMetadataV2Valid.json", temp_path);
3753
3754 let mut props = original_metadata.properties.clone();
3756 props.insert(
3757 TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
3758 "GziP".to_string(),
3759 );
3760 let compressed_metadata =
3762 TableMetadataBuilder::new_from_metadata(original_metadata.clone(), None)
3763 .assign_uuid(original_metadata.table_uuid)
3764 .set_properties(props.clone())
3765 .unwrap()
3766 .build()
3767 .unwrap()
3768 .metadata;
3769
3770 let metadata_location =
3772 MetadataLocation::try_new_with_metadata(&compressed_metadata).unwrap();
3773 let metadata_location_str = metadata_location.to_string();
3774
3775 assert!(metadata_location_str.contains(".gz.metadata.json"));
3777
3778 compressed_metadata
3780 .write_to(&file_io, &metadata_location)
3781 .await
3782 .unwrap();
3783
3784 assert!(std::path::Path::new(&metadata_location_str).exists());
3786
3787 let raw_content = fs::read(&metadata_location_str).unwrap();
3789 assert!(raw_content.len() > 2);
3790 assert_eq!(raw_content[0], 0x1F); assert_eq!(raw_content[1], 0x8B); let read_metadata = TableMetadata::read_from(&file_io, &metadata_location_str)
3795 .await
3796 .unwrap();
3797
3798 assert_eq!(read_metadata, compressed_metadata);
3800 }
3801
3802 #[test]
3803 fn test_partition_name_exists() {
3804 let schema = Schema::builder()
3805 .with_fields(vec![
3806 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
3807 NestedField::required(2, "partition_col", Type::Primitive(PrimitiveType::Int))
3808 .into(),
3809 ])
3810 .build()
3811 .unwrap();
3812
3813 let spec1 = PartitionSpec::builder(schema.clone())
3814 .with_spec_id(1)
3815 .add_partition_field("data", "data_partition", Transform::Identity)
3816 .unwrap()
3817 .build()
3818 .unwrap();
3819
3820 let spec2 = PartitionSpec::builder(schema.clone())
3821 .with_spec_id(2)
3822 .add_partition_field("partition_col", "partition_bucket", Transform::Bucket(16))
3823 .unwrap()
3824 .build()
3825 .unwrap();
3826
3827 let metadata = TableMetadataBuilder::new(
3829 schema,
3830 spec1.clone().into_unbound(),
3831 SortOrder::unsorted_order(),
3832 "s3://test/location".to_string(),
3833 FormatVersion::V2,
3834 HashMap::new(),
3835 )
3836 .unwrap()
3837 .add_partition_spec(spec2.into_unbound())
3838 .unwrap()
3839 .build()
3840 .unwrap()
3841 .metadata;
3842
3843 assert!(metadata.partition_name_exists("data_partition"));
3844 assert!(metadata.partition_name_exists("partition_bucket"));
3845
3846 assert!(!metadata.partition_name_exists("nonexistent_field"));
3847 assert!(!metadata.partition_name_exists("data")); assert!(!metadata.partition_name_exists(""));
3849 }
3850
3851 #[test]
3852 fn test_partition_name_exists_empty_specs() {
3853 let schema = Schema::builder()
3855 .with_fields(vec![
3856 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
3857 ])
3858 .build()
3859 .unwrap();
3860
3861 let metadata = TableMetadataBuilder::new(
3862 schema,
3863 PartitionSpec::unpartition_spec().into_unbound(),
3864 SortOrder::unsorted_order(),
3865 "s3://test/location".to_string(),
3866 FormatVersion::V2,
3867 HashMap::new(),
3868 )
3869 .unwrap()
3870 .build()
3871 .unwrap()
3872 .metadata;
3873
3874 assert!(!metadata.partition_name_exists("any_field"));
3875 assert!(!metadata.partition_name_exists("data"));
3876 }
3877
3878 #[test]
3879 fn test_name_exists_in_any_schema() {
3880 let schema1 = Schema::builder()
3882 .with_schema_id(1)
3883 .with_fields(vec![
3884 NestedField::required(1, "field1", Type::Primitive(PrimitiveType::String)).into(),
3885 NestedField::required(2, "field2", Type::Primitive(PrimitiveType::Int)).into(),
3886 ])
3887 .build()
3888 .unwrap();
3889
3890 let schema2 = Schema::builder()
3891 .with_schema_id(2)
3892 .with_fields(vec![
3893 NestedField::required(1, "field1", Type::Primitive(PrimitiveType::String)).into(),
3894 NestedField::required(3, "field3", Type::Primitive(PrimitiveType::Long)).into(),
3895 ])
3896 .build()
3897 .unwrap();
3898
3899 let metadata = TableMetadataBuilder::new(
3900 schema1,
3901 PartitionSpec::unpartition_spec().into_unbound(),
3902 SortOrder::unsorted_order(),
3903 "s3://test/location".to_string(),
3904 FormatVersion::V2,
3905 HashMap::new(),
3906 )
3907 .unwrap()
3908 .add_current_schema(schema2)
3909 .unwrap()
3910 .build()
3911 .unwrap()
3912 .metadata;
3913
3914 assert!(metadata.name_exists_in_any_schema("field1")); assert!(metadata.name_exists_in_any_schema("field2")); assert!(metadata.name_exists_in_any_schema("field3")); assert!(!metadata.name_exists_in_any_schema("nonexistent_field"));
3919 assert!(!metadata.name_exists_in_any_schema("field4"));
3920 assert!(!metadata.name_exists_in_any_schema(""));
3921 }
3922
3923 #[test]
3924 fn test_name_exists_in_any_schema_empty_schemas() {
3925 let schema = Schema::builder().with_fields(vec![]).build().unwrap();
3926
3927 let metadata = TableMetadataBuilder::new(
3928 schema,
3929 PartitionSpec::unpartition_spec().into_unbound(),
3930 SortOrder::unsorted_order(),
3931 "s3://test/location".to_string(),
3932 FormatVersion::V2,
3933 HashMap::new(),
3934 )
3935 .unwrap()
3936 .build()
3937 .unwrap()
3938 .metadata;
3939
3940 assert!(!metadata.name_exists_in_any_schema("any_field"));
3941 }
3942
3943 #[test]
3944 fn test_helper_methods_multi_version_scenario() {
3945 let initial_schema = Schema::builder()
3947 .with_fields(vec![
3948 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
3949 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
3950 NestedField::required(
3951 3,
3952 "deprecated_field",
3953 Type::Primitive(PrimitiveType::String),
3954 )
3955 .into(),
3956 ])
3957 .build()
3958 .unwrap();
3959
3960 let metadata = TableMetadataBuilder::new(
3961 initial_schema,
3962 PartitionSpec::unpartition_spec().into_unbound(),
3963 SortOrder::unsorted_order(),
3964 "s3://test/location".to_string(),
3965 FormatVersion::V2,
3966 HashMap::new(),
3967 )
3968 .unwrap();
3969
3970 let evolved_schema = Schema::builder()
3971 .with_fields(vec![
3972 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
3973 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
3974 NestedField::required(
3975 3,
3976 "deprecated_field",
3977 Type::Primitive(PrimitiveType::String),
3978 )
3979 .into(),
3980 NestedField::required(4, "new_field", Type::Primitive(PrimitiveType::Double))
3981 .into(),
3982 ])
3983 .build()
3984 .unwrap();
3985
3986 let _final_schema = Schema::builder()
3988 .with_fields(vec![
3989 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
3990 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
3991 NestedField::required(4, "new_field", Type::Primitive(PrimitiveType::Double))
3992 .into(),
3993 NestedField::required(5, "latest_field", Type::Primitive(PrimitiveType::Boolean))
3994 .into(),
3995 ])
3996 .build()
3997 .unwrap();
3998
3999 let final_metadata = metadata
4000 .add_current_schema(evolved_schema)
4001 .unwrap()
4002 .build()
4003 .unwrap()
4004 .metadata;
4005
4006 assert!(!final_metadata.partition_name_exists("nonexistent_partition")); assert!(final_metadata.name_exists_in_any_schema("id")); assert!(final_metadata.name_exists_in_any_schema("name")); assert!(final_metadata.name_exists_in_any_schema("deprecated_field")); assert!(final_metadata.name_exists_in_any_schema("new_field")); assert!(!final_metadata.name_exists_in_any_schema("never_existed"));
4013 }
4014
4015 #[test]
4016 fn test_invalid_sort_order_id_zero_with_fields() {
4017 let metadata = r#"
4018 {
4019 "format-version": 2,
4020 "table-uuid": "9c12d441-03fe-4693-9a96-a0705ddf69c1",
4021 "location": "s3://bucket/test/location",
4022 "last-sequence-number": 111,
4023 "last-updated-ms": 1600000000000,
4024 "last-column-id": 3,
4025 "current-schema-id": 1,
4026 "schemas": [
4027 {
4028 "type": "struct",
4029 "schema-id": 1,
4030 "fields": [
4031 {"id": 1, "name": "x", "required": true, "type": "long"},
4032 {"id": 2, "name": "y", "required": true, "type": "long"}
4033 ]
4034 }
4035 ],
4036 "default-spec-id": 0,
4037 "partition-specs": [{"spec-id": 0, "fields": []}],
4038 "last-partition-id": 999,
4039 "default-sort-order-id": 0,
4040 "sort-orders": [
4041 {
4042 "order-id": 0,
4043 "fields": [
4044 {
4045 "transform": "identity",
4046 "source-id": 1,
4047 "direction": "asc",
4048 "null-order": "nulls-first"
4049 }
4050 ]
4051 }
4052 ],
4053 "properties": {},
4054 "current-snapshot-id": -1,
4055 "snapshots": []
4056 }
4057 "#;
4058
4059 let result: Result<TableMetadata, serde_json::Error> = serde_json::from_str(metadata);
4060
4061 assert!(
4063 result.is_err(),
4064 "Parsing should fail for sort order ID 0 with fields"
4065 );
4066 }
4067
4068 #[test]
4069 fn test_table_properties_with_defaults() {
4070 use crate::spec::TableProperties;
4071
4072 let schema = Schema::builder()
4073 .with_fields(vec![
4074 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
4075 ])
4076 .build()
4077 .unwrap();
4078
4079 let metadata = TableMetadataBuilder::new(
4080 schema,
4081 PartitionSpec::unpartition_spec().into_unbound(),
4082 SortOrder::unsorted_order(),
4083 "s3://test/location".to_string(),
4084 FormatVersion::V2,
4085 HashMap::new(),
4086 )
4087 .unwrap()
4088 .build()
4089 .unwrap()
4090 .metadata;
4091
4092 let props = metadata.table_properties();
4093
4094 assert_eq!(
4095 props.commit_num_retries().unwrap(),
4096 TableProperties::PROPERTY_COMMIT_NUM_RETRIES_DEFAULT
4097 );
4098 assert_eq!(
4099 props.write_target_file_size_bytes().unwrap(),
4100 TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES_DEFAULT
4101 );
4102 }
4103
4104 #[test]
4105 fn test_table_properties_with_custom_values() {
4106 use crate::spec::TableProperties;
4107
4108 let schema = Schema::builder()
4109 .with_fields(vec![
4110 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
4111 ])
4112 .build()
4113 .unwrap();
4114
4115 let properties = HashMap::from([
4116 (
4117 TableProperties::PROPERTY_COMMIT_NUM_RETRIES.to_string(),
4118 "10".to_string(),
4119 ),
4120 (
4121 TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES.to_string(),
4122 "1024".to_string(),
4123 ),
4124 ]);
4125
4126 let metadata = TableMetadataBuilder::new(
4127 schema,
4128 PartitionSpec::unpartition_spec().into_unbound(),
4129 SortOrder::unsorted_order(),
4130 "s3://test/location".to_string(),
4131 FormatVersion::V2,
4132 properties,
4133 )
4134 .unwrap()
4135 .build()
4136 .unwrap()
4137 .metadata;
4138
4139 let props = metadata.table_properties();
4140
4141 assert_eq!(props.commit_num_retries().unwrap(), 10);
4142 assert_eq!(props.write_target_file_size_bytes().unwrap(), 1024);
4143 }
4144
4145 #[test]
4146 fn test_deserialize_metadata_defers_invalid_table_property_errors() {
4147 let invalid_retries = "not_a_number";
4148 let invalid_codec = "unknown";
4149 let target_file_size = "1024";
4150
4151 for file_name in [
4152 "TableMetadataV1Valid.json",
4153 "TableMetadataV2ValidMinimal.json",
4154 "TableMetadataV3ValidMinimal.json",
4155 ] {
4156 let path = format!("testdata/table_metadata/{file_name}");
4157 let mut json: serde_json::Value =
4158 serde_json::from_str(&fs::read_to_string(path).unwrap()).unwrap();
4159 json["properties"] = serde_json::json!({
4160 (TableProperties::PROPERTY_COMMIT_NUM_RETRIES): invalid_retries,
4161 (TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC): invalid_codec,
4162 (TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES): target_file_size,
4163 });
4164
4165 let metadata: TableMetadata = serde_json::from_value(json).unwrap();
4166 assert_eq!(
4167 metadata
4168 .properties()
4169 .get(TableProperties::PROPERTY_COMMIT_NUM_RETRIES)
4170 .map(String::as_str),
4171 Some(invalid_retries)
4172 );
4173
4174 let table_properties = metadata.table_properties();
4175 let error = table_properties.commit_num_retries().unwrap_err();
4176 assert!(
4177 error
4178 .message()
4179 .contains(TableProperties::PROPERTY_COMMIT_NUM_RETRIES)
4180 );
4181 assert_eq!(
4182 table_properties.write_target_file_size_bytes().unwrap(),
4183 1024
4184 );
4185 let error = table_properties.metadata_compression_codec().unwrap_err();
4186 assert!(
4187 format!("{error}").contains(TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC)
4188 );
4189
4190 let serialized = serde_json::to_value(metadata).unwrap();
4191 assert_eq!(
4192 serialized["properties"][TableProperties::PROPERTY_COMMIT_NUM_RETRIES],
4193 invalid_retries
4194 );
4195 assert_eq!(
4196 serialized["properties"][TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC],
4197 invalid_codec
4198 );
4199 }
4200 }
4201
4202 #[test]
4203 fn test_table_properties_with_invalid_value() {
4204 let schema = Schema::builder()
4205 .with_fields(vec![
4206 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
4207 ])
4208 .build()
4209 .unwrap();
4210
4211 let properties = HashMap::from([
4212 (
4213 TableProperties::PROPERTY_COMMIT_NUM_RETRIES.to_string(),
4214 "not_a_number".to_string(),
4215 ),
4216 (
4217 TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES.to_string(),
4218 "1024".to_string(),
4219 ),
4220 ]);
4221
4222 let metadata = TableMetadataBuilder::new(
4223 schema,
4224 PartitionSpec::unpartition_spec().into_unbound(),
4225 SortOrder::unsorted_order(),
4226 "s3://test/location".to_string(),
4227 FormatVersion::V2,
4228 properties,
4229 )
4230 .unwrap()
4231 .build()
4232 .unwrap()
4233 .metadata;
4234
4235 let table_properties = metadata.table_properties();
4236 let err = table_properties.commit_num_retries().unwrap_err();
4237 assert_eq!(err.kind(), ErrorKind::DataInvalid);
4238 assert!(
4239 err.message()
4240 .contains(TableProperties::PROPERTY_COMMIT_NUM_RETRIES)
4241 );
4242 assert_eq!(
4243 table_properties.write_target_file_size_bytes().unwrap(),
4244 1024
4245 );
4246 }
4247
4248 #[test]
4249 fn test_v2_to_v3_upgrade_preserves_existing_snapshots_without_row_lineage() {
4250 let schema = Schema::builder()
4252 .with_fields(vec![
4253 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
4254 ])
4255 .build()
4256 .unwrap();
4257
4258 let v2_metadata = TableMetadataBuilder::new(
4259 schema,
4260 PartitionSpec::unpartition_spec().into_unbound(),
4261 SortOrder::unsorted_order(),
4262 "s3://bucket/test/location".to_string(),
4263 FormatVersion::V2,
4264 HashMap::new(),
4265 )
4266 .unwrap()
4267 .build()
4268 .unwrap()
4269 .metadata;
4270
4271 let snapshot = Snapshot::builder()
4273 .with_snapshot_id(1)
4274 .with_timestamp_ms(v2_metadata.last_updated_ms + 1)
4275 .with_sequence_number(1)
4276 .with_schema_id(0)
4277 .with_manifest_list("s3://bucket/test/metadata/snap-1.avro")
4278 .with_summary(Summary {
4279 operation: Operation::Append,
4280 additional_properties: HashMap::from([(
4281 "added-data-files".to_string(),
4282 "1".to_string(),
4283 )]),
4284 })
4285 .build();
4286
4287 let v2_with_snapshot = v2_metadata
4288 .into_builder(Some("s3://bucket/test/metadata/v00001.json".to_string()))
4289 .add_snapshot(snapshot)
4290 .unwrap()
4291 .set_ref("main", SnapshotReference {
4292 snapshot_id: 1,
4293 retention: SnapshotRetention::Branch {
4294 min_snapshots_to_keep: None,
4295 max_snapshot_age_ms: None,
4296 max_ref_age_ms: None,
4297 },
4298 })
4299 .unwrap()
4300 .build()
4301 .unwrap()
4302 .metadata;
4303
4304 let v2_json = serde_json::to_string(&v2_with_snapshot);
4306 assert!(v2_json.is_ok(), "v2 serialization should work");
4307
4308 let v3_metadata = v2_with_snapshot
4310 .into_builder(Some("s3://bucket/test/metadata/v00002.json".to_string()))
4311 .upgrade_format_version(FormatVersion::V3)
4312 .unwrap()
4313 .build()
4314 .unwrap()
4315 .metadata;
4316
4317 assert_eq!(v3_metadata.format_version, FormatVersion::V3);
4318 assert_eq!(v3_metadata.next_row_id, INITIAL_ROW_ID);
4319 assert_eq!(v3_metadata.snapshots.len(), 1);
4320
4321 let snapshot = v3_metadata.snapshots.values().next().unwrap();
4323 assert!(
4324 snapshot.row_range().is_none(),
4325 "Snapshot should have no row_range after upgrade"
4326 );
4327
4328 let v3_json = serde_json::to_string(&v3_metadata);
4330 assert!(
4331 v3_json.is_ok(),
4332 "v3 serialization should work for upgraded tables"
4333 );
4334
4335 let deserialized: TableMetadata = serde_json::from_str(&v3_json.unwrap()).unwrap();
4337 assert_eq!(deserialized.format_version, FormatVersion::V3);
4338 assert_eq!(deserialized.snapshots.len(), 1);
4339
4340 let deserialized_snapshot = deserialized.snapshots.values().next().unwrap();
4342 assert!(
4343 deserialized_snapshot.row_range().is_none(),
4344 "Deserialized snapshot should have no row_range"
4345 );
4346 }
4347
4348 #[test]
4349 fn test_v3_snapshot_with_row_lineage_serialization() {
4350 let schema = Schema::builder()
4352 .with_fields(vec![
4353 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
4354 ])
4355 .build()
4356 .unwrap();
4357
4358 let v3_metadata = TableMetadataBuilder::new(
4359 schema,
4360 PartitionSpec::unpartition_spec().into_unbound(),
4361 SortOrder::unsorted_order(),
4362 "s3://bucket/test/location".to_string(),
4363 FormatVersion::V3,
4364 HashMap::new(),
4365 )
4366 .unwrap()
4367 .build()
4368 .unwrap()
4369 .metadata;
4370
4371 let snapshot = Snapshot::builder()
4373 .with_snapshot_id(1)
4374 .with_timestamp_ms(v3_metadata.last_updated_ms + 1)
4375 .with_sequence_number(1)
4376 .with_schema_id(0)
4377 .with_manifest_list("s3://bucket/test/metadata/snap-1.avro")
4378 .with_summary(Summary {
4379 operation: Operation::Append,
4380 additional_properties: HashMap::from([(
4381 "added-data-files".to_string(),
4382 "1".to_string(),
4383 )]),
4384 })
4385 .with_row_range(100, 50) .build();
4387
4388 let v3_with_snapshot = v3_metadata
4389 .into_builder(Some("s3://bucket/test/metadata/v00001.json".to_string()))
4390 .add_snapshot(snapshot)
4391 .unwrap()
4392 .set_ref("main", SnapshotReference {
4393 snapshot_id: 1,
4394 retention: SnapshotRetention::Branch {
4395 min_snapshots_to_keep: None,
4396 max_snapshot_age_ms: None,
4397 max_ref_age_ms: None,
4398 },
4399 })
4400 .unwrap()
4401 .build()
4402 .unwrap()
4403 .metadata;
4404
4405 let snapshot = v3_with_snapshot.snapshots.values().next().unwrap();
4407 assert!(
4408 snapshot.row_range().is_some(),
4409 "Snapshot should have row_range"
4410 );
4411 let (first_row_id, added_rows) = snapshot.row_range().unwrap();
4412 assert_eq!(first_row_id, 100);
4413 assert_eq!(added_rows, 50);
4414
4415 let v3_json = serde_json::to_string(&v3_with_snapshot);
4417 assert!(
4418 v3_json.is_ok(),
4419 "v3 serialization should work for snapshots with row lineage"
4420 );
4421
4422 let deserialized: TableMetadata = serde_json::from_str(&v3_json.unwrap()).unwrap();
4424 assert_eq!(deserialized.format_version, FormatVersion::V3);
4425 assert_eq!(deserialized.snapshots.len(), 1);
4426
4427 let deserialized_snapshot = deserialized.snapshots.values().next().unwrap();
4429 assert!(
4430 deserialized_snapshot.row_range().is_some(),
4431 "Deserialized snapshot should have row_range"
4432 );
4433 let (deserialized_first_row_id, deserialized_added_rows) =
4434 deserialized_snapshot.row_range().unwrap();
4435 assert_eq!(deserialized_first_row_id, 100);
4436 assert_eq!(deserialized_added_rows, 50);
4437 }
4438
4439 #[test]
4440 fn test_metadata_location_default() {
4441 let metadata = get_test_table_metadata("TableMetadataV2Valid.json");
4443 assert_eq!(metadata.location(), "s3://bucket/test/location");
4444 assert_eq!(
4445 metadata.metadata_location().unwrap(),
4446 "s3://bucket/test/location/metadata"
4447 );
4448 }
4449
4450 #[test]
4451 fn test_metadata_location_honors_write_metadata_path() {
4452 let metadata = get_test_table_metadata("TableMetadataV2Valid.json")
4453 .into_builder(None)
4454 .set_properties(HashMap::from([(
4455 TableProperties::PROPERTY_WRITE_METADATA_PATH.to_string(),
4456 "s3://other-bucket/custom-meta".to_string(),
4457 )]))
4458 .unwrap()
4459 .build()
4460 .unwrap()
4461 .metadata;
4462 assert_eq!(
4463 metadata.metadata_location().unwrap(),
4464 "s3://other-bucket/custom-meta"
4465 );
4466 }
4467
4468 #[test]
4469 fn test_metadata_location_trims_trailing_slash() {
4470 let metadata = get_test_table_metadata("TableMetadataV2Valid.json")
4472 .into_builder(None)
4473 .set_properties(HashMap::from([(
4474 TableProperties::PROPERTY_WRITE_METADATA_PATH.to_string(),
4475 "s3://other-bucket/custom-meta/".to_string(),
4476 )]))
4477 .unwrap()
4478 .build()
4479 .unwrap()
4480 .metadata;
4481 assert_eq!(
4482 metadata.metadata_location().unwrap(),
4483 "s3://other-bucket/custom-meta"
4484 );
4485 }
4486
4487 #[test]
4488 fn test_unified_partition_type_spans_all_specs() {
4489 let schema = Schema::builder()
4490 .with_fields(vec![
4491 NestedField::required(1, "x", Type::Primitive(PrimitiveType::Long)).into(),
4492 NestedField::required(2, "y", Type::Primitive(PrimitiveType::Long)).into(),
4493 NestedField::required(3, "z", Type::Primitive(PrimitiveType::Long)).into(),
4494 ])
4495 .build()
4496 .unwrap();
4497
4498 let metadata = TableMetadataBuilder::new(
4499 schema.clone(),
4500 UnboundPartitionSpec::builder()
4501 .with_spec_id(0)
4502 .add_partition_field(
4503 UnboundPartitionField::builder()
4504 .source_ids(vec![2])
4505 .name("y")
4506 .transform(Transform::Identity)
4507 .build()
4508 .unwrap(),
4509 )
4510 .unwrap()
4511 .build(),
4512 SortOrder::unsorted_order(),
4513 "s3://bucket/table".to_string(),
4514 FormatVersion::V2,
4515 HashMap::new(),
4516 )
4517 .unwrap()
4518 .build()
4519 .unwrap()
4520 .metadata
4521 .into_builder(None)
4522 .add_partition_spec(
4523 UnboundPartitionSpec::builder()
4524 .add_partition_field(
4525 UnboundPartitionField::builder()
4526 .source_ids(vec![3])
4527 .name("z")
4528 .transform(Transform::Identity)
4529 .build()
4530 .unwrap(),
4531 )
4532 .unwrap()
4533 .build(),
4534 )
4535 .unwrap()
4536 .build()
4537 .unwrap()
4538 .metadata;
4539
4540 assert_eq!(metadata.default_partition_type().fields().len(), 1);
4542
4543 let unified = metadata.unified_partition_type(&schema).unwrap();
4544 let names: Vec<&str> = unified
4545 .fields()
4546 .iter()
4547 .map(|field| field.name.as_str())
4548 .collect();
4549 assert_eq!(names, vec!["y", "z"]);
4550 }
4551}