1use std::collections::{HashMap, HashSet};
19use std::sync::Arc;
20
21use uuid::Uuid;
22
23use super::{
24 DEFAULT_PARTITION_SPEC_ID, DEFAULT_SCHEMA_ID, FormatVersion, MAIN_BRANCH, MetadataLog,
25 ONE_MINUTE_MS, PartitionSpec, PartitionSpecBuilder, PartitionStatisticsFile, Schema, SchemaRef,
26 Snapshot, SnapshotLog, SnapshotReference, SnapshotRetention, SortOrder, SortOrderRef,
27 StatisticsFile, StructType, TableMetadata, TableProperties, UNPARTITIONED_LAST_ASSIGNED_ID,
28 UnboundPartitionSpec,
29};
30use crate::error::{Error, ErrorKind, Result, invalid_data};
31use crate::spec::{EncryptedKey, INITIAL_ROW_ID, MIN_FORMAT_VERSION_ROW_LINEAGE};
32use crate::{TableCreation, TableUpdate};
33
34pub(crate) const FIRST_FIELD_ID: i32 = 1;
35
36#[derive(Debug, Clone)]
49pub struct TableMetadataBuilder {
50 metadata: TableMetadata,
51 changes: Vec<TableUpdate>,
52 last_added_schema_id: Option<i32>,
53 last_added_spec_id: Option<i32>,
54 last_added_order_id: Option<i64>,
55 previous_history_entry: Option<MetadataLog>,
57 last_updated_ms: Option<i64>,
58}
59
60#[derive(Debug, Clone, PartialEq)]
61pub struct TableMetadataBuildResult {
63 pub metadata: TableMetadata,
65 pub changes: Vec<TableUpdate>,
67 pub expired_metadata_logs: Vec<MetadataLog>,
69}
70
71impl TableMetadataBuilder {
72 pub const LAST_ADDED: i32 = -1;
74
75 pub fn new(
80 schema: Schema,
81 spec: impl Into<UnboundPartitionSpec>,
82 sort_order: SortOrder,
83 location: String,
84 format_version: FormatVersion,
85 properties: HashMap<String, String>,
86 ) -> Result<Self> {
87 let (fresh_schema, fresh_spec, fresh_sort_order) =
89 Self::reassign_ids(schema, spec.into(), sort_order)?;
90 let schema_id = fresh_schema.schema_id();
91
92 let builder = Self {
93 metadata: TableMetadata {
94 format_version,
95 table_uuid: Uuid::now_v7(),
96 location: "".to_string(), last_sequence_number: 0,
98 last_updated_ms: 0, last_column_id: -1, current_schema_id: -1, schemas: HashMap::new(),
102 partition_specs: HashMap::new(),
103 default_spec: Arc::new(
104 PartitionSpec::unpartition_spec().with_spec_id(-1),
110 ), default_partition_type: StructType::new(vec![]),
112 last_partition_id: UNPARTITIONED_LAST_ASSIGNED_ID,
113 properties: HashMap::new(),
114 current_snapshot_id: None,
115 snapshots: HashMap::new(),
116 snapshot_log: vec![],
117 sort_orders: HashMap::new(),
118 metadata_log: vec![],
119 default_sort_order_id: -1, refs: HashMap::default(),
121 statistics: HashMap::new(),
122 partition_statistics: HashMap::new(),
123 encryption_keys: HashMap::new(),
124 next_row_id: INITIAL_ROW_ID,
125 },
126 last_updated_ms: None,
127 changes: vec![],
128 last_added_schema_id: Some(schema_id),
129 last_added_spec_id: None,
130 last_added_order_id: None,
131 previous_history_entry: None,
132 };
133
134 builder
135 .set_location(location)
136 .add_current_schema(fresh_schema)?
137 .add_default_partition_spec(fresh_spec.into_unbound())?
138 .add_default_sort_order(fresh_sort_order)?
139 .set_properties(properties)
140 }
141
142 #[must_use]
148 pub fn new_from_metadata(
149 previous: TableMetadata,
150 current_file_location: Option<String>,
151 ) -> Self {
152 Self {
153 previous_history_entry: current_file_location.map(|l| MetadataLog {
154 metadata_file: l,
155 timestamp_ms: previous.last_updated_ms,
156 }),
157 metadata: previous,
158 changes: Vec::default(),
159 last_added_schema_id: None,
160 last_added_spec_id: None,
161 last_added_order_id: None,
162 last_updated_ms: None,
163 }
164 }
165
166 pub fn from_table_creation(table_creation: TableCreation) -> Result<Self> {
168 let TableCreation {
169 name: _,
170 location,
171 schema,
172 partition_spec,
173 sort_order,
174 properties,
175 format_version,
176 } = table_creation;
177
178 let location =
179 location.ok_or_else(|| invalid_data!("Can't create table without location"))?;
180 let partition_spec = partition_spec.unwrap_or(UnboundPartitionSpec {
181 spec_id: None,
182 fields: vec![],
183 });
184
185 Self::new(
186 schema,
187 partition_spec,
188 sort_order.unwrap_or(SortOrder::unsorted_order()),
189 location,
190 format_version,
191 properties,
192 )
193 }
194
195 pub fn assign_uuid(mut self, uuid: Uuid) -> Self {
197 if self.metadata.table_uuid != uuid {
198 self.metadata.table_uuid = uuid;
199 self.changes.push(TableUpdate::AssignUuid { uuid });
200 }
201
202 self
203 }
204
205 pub fn upgrade_format_version(mut self, format_version: FormatVersion) -> Result<Self> {
210 if format_version < self.metadata.format_version {
211 return Err(invalid_data!(
212 "Cannot downgrade FormatVersion from {} to {}",
213 self.metadata.format_version,
214 format_version
215 ));
216 }
217
218 if format_version != self.metadata.format_version {
219 match format_version {
220 FormatVersion::V1 => {
221 }
223 FormatVersion::V2 => {
224 self.metadata.format_version = format_version;
225 self.changes
226 .push(TableUpdate::UpgradeFormatVersion { format_version });
227 }
228 FormatVersion::V3 => {
229 self.metadata.format_version = format_version;
230 self.metadata.next_row_id = INITIAL_ROW_ID;
232 self.changes
233 .push(TableUpdate::UpgradeFormatVersion { format_version });
234 }
235 }
236 }
237
238 Ok(self)
239 }
240
241 pub fn set_properties(mut self, properties: HashMap<String, String>) -> Result<Self> {
250 let reserved_properties = properties
252 .keys()
253 .filter(|key| TableProperties::RESERVED_PROPERTIES.contains(&key.as_str()))
254 .map(ToString::to_string)
255 .collect::<Vec<_>>();
256
257 if !reserved_properties.is_empty() {
258 return Err(invalid_data!(
259 "Table properties should not contain reserved properties, but got: [{}]",
260 reserved_properties.join(", ")
261 ));
262 }
263
264 if properties.is_empty() {
265 return Ok(self);
266 }
267
268 self.metadata.properties.extend(properties.clone());
269 self.changes.push(TableUpdate::SetProperties {
270 updates: properties,
271 });
272
273 Ok(self)
274 }
275
276 pub fn remove_properties(mut self, properties: &[String]) -> Result<Self> {
282 let properties = properties.iter().cloned().collect::<HashSet<_>>();
284
285 let reserved_properties = properties
287 .iter()
288 .filter(|key| TableProperties::RESERVED_PROPERTIES.contains(&key.as_str()))
289 .map(ToString::to_string)
290 .collect::<Vec<_>>();
291
292 if !reserved_properties.is_empty() {
293 return Err(invalid_data!(
294 "Table properties to remove contain reserved properties: [{}]",
295 reserved_properties.join(", ")
296 ));
297 }
298
299 for property in &properties {
300 self.metadata.properties.remove(property);
301 }
302
303 if !properties.is_empty() {
304 self.changes.push(TableUpdate::RemoveProperties {
305 removals: properties.into_iter().collect(),
306 });
307 }
308
309 Ok(self)
310 }
311
312 pub fn set_location(mut self, location: String) -> Self {
314 let location = location.trim_end_matches('/').to_string();
315 if self.metadata.location != location {
316 self.changes.push(TableUpdate::SetLocation {
317 location: location.clone(),
318 });
319 self.metadata.location = location;
320 }
321
322 self
323 }
324
325 pub fn add_snapshot(mut self, snapshot: Snapshot) -> Result<Self> {
334 if self
335 .metadata
336 .snapshots
337 .contains_key(&snapshot.snapshot_id())
338 {
339 return Err(invalid_data!(
340 "Snapshot already exists for: '{}'",
341 snapshot.snapshot_id()
342 ));
343 }
344
345 if self.metadata.format_version != FormatVersion::V1
346 && snapshot.sequence_number() <= self.metadata.last_sequence_number
347 && snapshot.parent_snapshot_id().is_some()
348 {
349 return Err(invalid_data!(
350 "Cannot add snapshot with sequence number {} older than last sequence number {}",
351 snapshot.sequence_number(),
352 self.metadata.last_sequence_number
353 ));
354 }
355
356 if let Some(last) = self.metadata.snapshot_log.last() {
357 if snapshot.timestamp_ms() - last.timestamp_ms < -ONE_MINUTE_MS {
360 return Err(invalid_data!(
361 "Invalid snapshot timestamp {}: before last snapshot timestamp {}",
362 snapshot.timestamp_ms(),
363 last.timestamp_ms
364 ));
365 }
366 }
367
368 let max_last_updated = self
369 .last_updated_ms
370 .unwrap_or_default()
371 .max(self.metadata.last_updated_ms);
372 if snapshot.timestamp_ms() - max_last_updated < -ONE_MINUTE_MS {
373 return Err(invalid_data!(
374 "Invalid snapshot timestamp {}: before last updated timestamp {}",
375 snapshot.timestamp_ms(),
376 max_last_updated
377 ));
378 }
379
380 let mut added_rows = None;
381 if self.metadata.format_version >= MIN_FORMAT_VERSION_ROW_LINEAGE {
382 if let Some((first_row_id, added_rows_count)) = snapshot.row_range() {
383 if first_row_id < self.metadata.next_row_id {
384 return Err(invalid_data!(
385 "Cannot add a snapshot, first-row-id is behind table next-row-id: {first_row_id} < {}",
386 self.metadata.next_row_id
387 ));
388 }
389
390 added_rows = Some(added_rows_count);
391 } else {
392 return Err(invalid_data!(
393 "Cannot add a snapshot: first-row-id is null. first-row-id must be set for format version >= {MIN_FORMAT_VERSION_ROW_LINEAGE}",
394 ));
395 }
396 }
397
398 if let Some(added_rows) = added_rows {
399 self.metadata.next_row_id = self
400 .metadata
401 .next_row_id
402 .checked_add(added_rows)
403 .ok_or_else(|| {
404 invalid_data!(
405 "Cannot add snapshot: next-row-id overflowed when adding added-rows"
406 )
407 })?;
408 }
409
410 self.changes.push(TableUpdate::AddSnapshot {
412 snapshot: snapshot.clone(),
413 });
414
415 self.last_updated_ms = Some(snapshot.timestamp_ms());
416 self.metadata.last_sequence_number = snapshot.sequence_number();
417 self.metadata
418 .snapshots
419 .insert(snapshot.snapshot_id(), snapshot.into());
420
421 Ok(self)
422 }
423
424 pub fn set_branch_snapshot(self, snapshot: Snapshot, branch: &str) -> Result<Self> {
430 let reference = self.metadata.refs.get(branch).cloned();
431
432 let reference = if let Some(mut reference) = reference {
433 if !reference.is_branch() {
434 return Err(invalid_data!(
435 "Cannot append snapshot to non-branch reference '{branch}'",
436 ));
437 }
438
439 reference.snapshot_id = snapshot.snapshot_id();
440 reference
441 } else {
442 SnapshotReference {
443 snapshot_id: snapshot.snapshot_id(),
444 retention: SnapshotRetention::Branch {
445 min_snapshots_to_keep: None,
446 max_snapshot_age_ms: None,
447 max_ref_age_ms: None,
448 },
449 }
450 };
451
452 self.add_snapshot(snapshot)?.set_ref(branch, reference)
453 }
454
455 pub fn remove_snapshots(mut self, snapshot_ids: &[i64]) -> Self {
459 let mut removed_snapshots = Vec::with_capacity(snapshot_ids.len());
460
461 self.metadata.snapshots.retain(|k, _| {
462 if snapshot_ids.contains(k) {
463 removed_snapshots.push(*k);
464 false
465 } else {
466 true
467 }
468 });
469
470 if !removed_snapshots.is_empty() {
471 self.changes.push(TableUpdate::RemoveSnapshots {
472 snapshot_ids: removed_snapshots,
473 });
474 }
475
476 self.metadata
478 .refs
479 .retain(|_, v| self.metadata.snapshots.contains_key(&v.snapshot_id));
480
481 self
482 }
483
484 pub fn set_ref(mut self, ref_name: &str, reference: SnapshotReference) -> Result<Self> {
489 if self
490 .metadata
491 .refs
492 .get(ref_name)
493 .is_some_and(|snap_ref| snap_ref.eq(&reference))
494 {
495 return Ok(self);
496 }
497
498 let Some(snapshot) = self.metadata.snapshots.get(&reference.snapshot_id) else {
499 return Err(invalid_data!(
500 "Cannot set '{ref_name}' to unknown snapshot: '{}'",
501 reference.snapshot_id
502 ));
503 };
504
505 let is_added_snapshot = self.changes.iter().any(|update| {
507 matches!(update, TableUpdate::AddSnapshot { snapshot: snap } if snap.snapshot_id() == snapshot.snapshot_id())
508 });
509 if is_added_snapshot {
510 self.last_updated_ms = Some(snapshot.timestamp_ms());
511 }
512
513 if ref_name == MAIN_BRANCH {
515 self.metadata.current_snapshot_id = Some(snapshot.snapshot_id());
516 let timestamp_ms = if let Some(last_updated_ms) = self.last_updated_ms {
517 last_updated_ms
518 } else {
519 let last_updated_ms = chrono::Utc::now().timestamp_millis();
520 self.last_updated_ms = Some(last_updated_ms);
521 last_updated_ms
522 };
523
524 self.metadata.snapshot_log.push(SnapshotLog {
525 snapshot_id: snapshot.snapshot_id(),
526 timestamp_ms,
527 });
528 }
529
530 self.changes.push(TableUpdate::SetSnapshotRef {
531 ref_name: ref_name.to_string(),
532 reference: reference.clone(),
533 });
534 self.metadata.refs.insert(ref_name.to_string(), reference);
535
536 Ok(self)
537 }
538
539 pub fn remove_ref(mut self, ref_name: &str) -> Self {
543 if ref_name == MAIN_BRANCH {
544 self.metadata.current_snapshot_id = None;
545 }
546
547 if self.metadata.refs.remove(ref_name).is_some() || ref_name == MAIN_BRANCH {
548 self.changes.push(TableUpdate::RemoveSnapshotRef {
549 ref_name: ref_name.to_string(),
550 });
551 }
552
553 self
554 }
555
556 pub fn set_statistics(mut self, statistics: StatisticsFile) -> Self {
558 self.metadata
559 .statistics
560 .insert(statistics.snapshot_id, statistics.clone());
561 self.changes.push(TableUpdate::SetStatistics {
562 statistics: statistics.clone(),
563 });
564 self
565 }
566
567 pub fn remove_statistics(mut self, snapshot_id: i64) -> Self {
569 let previous = self.metadata.statistics.remove(&snapshot_id);
570 if previous.is_some() {
571 self.changes
572 .push(TableUpdate::RemoveStatistics { snapshot_id });
573 }
574 self
575 }
576
577 pub fn set_partition_statistics(
579 mut self,
580 partition_statistics_file: PartitionStatisticsFile,
581 ) -> Self {
582 self.metadata.partition_statistics.insert(
583 partition_statistics_file.snapshot_id,
584 partition_statistics_file.clone(),
585 );
586 self.changes.push(TableUpdate::SetPartitionStatistics {
587 partition_statistics: partition_statistics_file,
588 });
589 self
590 }
591
592 pub fn remove_partition_statistics(mut self, snapshot_id: i64) -> Self {
594 let previous = self.metadata.partition_statistics.remove(&snapshot_id);
595 if previous.is_some() {
596 self.changes
597 .push(TableUpdate::RemovePartitionStatistics { snapshot_id });
598 }
599 self
600 }
601
602 pub fn add_schema(mut self, schema: Schema) -> Result<Self> {
609 self.validate_schema_field_names(&schema)?;
611
612 let new_schema_id = self.reuse_or_create_new_schema_id(&schema);
613 let schema_found = self.metadata.schemas.contains_key(&new_schema_id);
614
615 if schema_found {
616 if self.last_added_schema_id != Some(new_schema_id) {
617 self.changes.push(TableUpdate::AddSchema {
618 schema: schema.clone(),
619 });
620 self.last_added_schema_id = Some(new_schema_id);
621 }
622
623 return Ok(self);
624 }
625
626 self.metadata.last_column_id =
629 std::cmp::max(self.metadata.last_column_id, schema.highest_field_id());
630
631 let schema = match new_schema_id == schema.schema_id() {
633 true => schema,
634 false => schema.with_schema_id(new_schema_id),
635 };
636
637 self.metadata
638 .schemas
639 .insert(new_schema_id, schema.clone().into());
640
641 self.changes.push(TableUpdate::AddSchema { schema });
642
643 self.last_added_schema_id = Some(new_schema_id);
644
645 Ok(self)
646 }
647
648 pub fn set_current_schema(mut self, mut schema_id: i32) -> Result<Self> {
656 if schema_id == Self::LAST_ADDED {
657 schema_id = self.last_added_schema_id.ok_or_else(|| {
658 invalid_data!(
659 "Cannot set current schema to last added schema: no schema has been added."
660 )
661 })?;
662 };
663 let schema_id = schema_id; if schema_id == self.metadata.current_schema_id {
666 return Ok(self);
667 }
668
669 let _schema = self.metadata.schemas.get(&schema_id).ok_or_else(|| {
670 invalid_data!("Cannot set current schema to unknown schema with id: '{schema_id}'")
671 })?;
672
673 self.metadata.current_schema_id = schema_id;
679
680 if self.last_added_schema_id == Some(schema_id) {
681 self.changes.push(TableUpdate::SetCurrentSchema {
682 schema_id: Self::LAST_ADDED,
683 });
684 } else {
685 self.changes
686 .push(TableUpdate::SetCurrentSchema { schema_id });
687 }
688
689 Ok(self)
690 }
691
692 pub fn add_current_schema(self, schema: Schema) -> Result<Self> {
694 self.add_schema(schema)?
695 .set_current_schema(Self::LAST_ADDED)
696 }
697
698 fn validate_schema_field_names(&self, schema: &Schema) -> Result<()> {
707 if self.metadata.schemas.is_empty() {
708 return Ok(());
709 }
710
711 for field_name in schema.field_id_to_name_map().values() {
712 let has_partition_conflict = self.metadata.partition_name_exists(field_name);
713 let is_new_field = !self.metadata.name_exists_in_any_schema(field_name);
714
715 if has_partition_conflict && is_new_field {
716 return Err(invalid_data!(
717 "Cannot add schema field '{field_name}' because it conflicts with existing partition field name. \
718 Schema evolution cannot introduce field names that match existing partition field names."
719 ));
720 }
721 }
722
723 Ok(())
724 }
725
726 fn validate_partition_field_names(&self, unbound_spec: &UnboundPartitionSpec) -> Result<()> {
736 if self.metadata.schemas.is_empty() {
737 return Ok(());
738 }
739
740 let current_schema = self.get_current_schema()?;
741 for partition_field in unbound_spec.fields() {
742 let exists_in_any_schema = self
743 .metadata
744 .name_exists_in_any_schema(partition_field.name());
745
746 if !exists_in_any_schema {
748 continue;
749 }
750
751 if let Some(schema_field) = current_schema.field_by_name(partition_field.name()) {
753 let is_identity_transform =
754 partition_field.transform() == crate::spec::Transform::Identity;
755 let has_matching_source_id = schema_field.id == partition_field.source_id()?;
756
757 if !is_identity_transform {
758 return Err(invalid_data!(
759 "Cannot create partition with name '{}' that conflicts with schema field and is not an identity transform.",
760 partition_field.name()
761 ));
762 }
763
764 if !has_matching_source_id {
765 return Err(invalid_data!(
766 "Cannot create identity partition sourced from different field in schema. \
767 Field name '{}' has id `{}` in schema but partition source id is `{}`",
768 partition_field.name(),
769 schema_field.id,
770 partition_field.source_id()?
771 ));
772 }
773 }
774 }
775
776 Ok(())
777 }
778
779 pub fn add_partition_spec(mut self, unbound_spec: UnboundPartitionSpec) -> Result<Self> {
790 let schema = self.get_current_schema()?.clone();
791
792 self.validate_partition_field_names(&unbound_spec)?;
794
795 let unbound_spec = self.reuse_partition_field_ids(unbound_spec)?;
797
798 let spec = PartitionSpecBuilder::new_from_unbound(unbound_spec.clone(), schema)?
799 .with_last_assigned_field_id(self.metadata.last_partition_id)
800 .build()?;
801
802 let new_spec_id = self.reuse_or_create_new_spec_id(&spec);
803 let spec_found = self.metadata.partition_specs.contains_key(&new_spec_id);
804 let spec = spec.with_spec_id(new_spec_id);
805 let unbound_spec = unbound_spec.with_spec_id(new_spec_id);
806
807 if spec_found {
808 if self.last_added_spec_id != Some(new_spec_id) {
809 self.changes
810 .push(TableUpdate::AddSpec { spec: unbound_spec });
811 self.last_added_spec_id = Some(new_spec_id);
812 }
813
814 return Ok(self);
815 }
816
817 if self.metadata.format_version <= FormatVersion::V1 && !spec.has_sequential_ids() {
818 return Err(invalid_data!(
819 "Cannot add partition spec with non-sequential field ids to format version 1 table"
820 ));
821 }
822
823 let highest_field_id = spec
824 .highest_field_id()
825 .unwrap_or(UNPARTITIONED_LAST_ASSIGNED_ID);
826 self.metadata
827 .partition_specs
828 .insert(new_spec_id, Arc::new(spec));
829 self.changes
830 .push(TableUpdate::AddSpec { spec: unbound_spec });
831
832 self.last_added_spec_id = Some(new_spec_id);
833 self.metadata.last_partition_id =
834 std::cmp::max(self.metadata.last_partition_id, highest_field_id);
835
836 Ok(self)
837 }
838
839 fn reuse_partition_field_ids(
844 &self,
845 unbound_spec: UnboundPartitionSpec,
846 ) -> Result<UnboundPartitionSpec> {
847 let equivalent_field_ids: HashMap<_, _> = self
849 .metadata
850 .partition_specs
851 .values()
852 .flat_map(|spec| spec.fields())
853 .map(|field| ((vec![field.source_id], field.transform), field.field_id))
854 .collect();
855
856 let fields = unbound_spec
858 .fields
859 .into_iter()
860 .map(|field| {
861 if field.field_id().is_none()
862 && let Some(&existing_field_id) =
863 equivalent_field_ids.get(&(field.source_ids().to_vec(), field.transform()))
864 {
865 field.with_field_id(existing_field_id)
866 } else {
867 field
868 }
869 })
870 .collect();
871
872 Ok(UnboundPartitionSpec {
873 spec_id: unbound_spec.spec_id,
874 fields,
875 })
876 }
877
878 pub fn set_default_partition_spec(mut self, mut spec_id: i32) -> Result<Self> {
884 if spec_id == Self::LAST_ADDED {
885 spec_id = self.last_added_spec_id.ok_or_else(|| {
886 invalid_data!(
887 "Cannot set default partition spec to last added spec: no spec has been added."
888 )
889 })?;
890 }
891
892 if self.metadata.default_spec.spec_id() == spec_id {
893 return Ok(self);
894 }
895
896 if !self.metadata.partition_specs.contains_key(&spec_id) {
897 return Err(invalid_data!(
898 "Cannot set default partition spec to unknown spec with id: '{spec_id}'",
899 ));
900 }
901
902 let schemaless_spec = self
903 .metadata
904 .partition_specs
905 .get(&spec_id)
906 .ok_or_else(|| {
907 invalid_data!(
908 "Cannot set default partition spec to unknown spec with id: '{spec_id}'",
909 )
910 })?
911 .clone();
912 let spec = Arc::unwrap_or_clone(schemaless_spec);
913 let spec_type = spec.partition_type(self.get_current_schema()?)?;
914 self.metadata.default_spec = Arc::new(spec);
915 self.metadata.default_partition_type = spec_type;
916
917 if self.last_added_spec_id == Some(spec_id) {
918 self.changes.push(TableUpdate::SetDefaultSpec {
919 spec_id: Self::LAST_ADDED,
920 });
921 } else {
922 self.changes.push(TableUpdate::SetDefaultSpec { spec_id });
923 }
924
925 Ok(self)
926 }
927
928 pub fn add_default_partition_spec(self, unbound_spec: UnboundPartitionSpec) -> Result<Self> {
930 self.add_partition_spec(unbound_spec)?
931 .set_default_partition_spec(Self::LAST_ADDED)
932 }
933
934 pub fn remove_partition_specs(mut self, spec_ids: &[i32]) -> Result<Self> {
941 if spec_ids.contains(&self.metadata.default_spec.spec_id()) {
942 return Err(invalid_data!("Cannot remove default partition spec"));
943 }
944
945 let mut removed_specs = Vec::with_capacity(spec_ids.len());
946 spec_ids.iter().for_each(|id| {
947 if self.metadata.partition_specs.remove(id).is_some() {
948 removed_specs.push(*id);
949 }
950 });
951
952 if !removed_specs.is_empty() {
953 self.changes.push(TableUpdate::RemovePartitionSpecs {
954 spec_ids: removed_specs,
955 });
956 }
957
958 Ok(self)
959 }
960
961 pub fn add_sort_order(mut self, sort_order: SortOrder) -> Result<Self> {
972 let new_order_id = self.reuse_or_create_new_sort_id(&sort_order);
973 let sort_order_found = self.metadata.sort_orders.contains_key(&new_order_id);
974
975 if sort_order_found {
976 if self.last_added_order_id != Some(new_order_id) {
977 self.changes.push(TableUpdate::AddSortOrder {
978 sort_order: sort_order.clone().with_order_id(new_order_id),
979 });
980 self.last_added_order_id = Some(new_order_id);
981 }
982
983 return Ok(self);
984 }
985
986 let schema = self.get_current_schema()?.clone().as_ref().clone();
987 let sort_order = SortOrder::builder()
988 .with_order_id(new_order_id)
989 .with_fields(sort_order.fields)
990 .build(&schema)
991 .map_err(|e| {
992 invalid_data!("Sort order to add is incompatible with current schema: {e}")
993 .with_source(e)
994 })?;
995
996 self.last_added_order_id = Some(new_order_id);
997 self.metadata
998 .sort_orders
999 .insert(new_order_id, sort_order.clone().into());
1000 self.changes.push(TableUpdate::AddSortOrder { sort_order });
1001
1002 Ok(self)
1003 }
1004
1005 pub fn set_default_sort_order(mut self, mut sort_order_id: i64) -> Result<Self> {
1011 if sort_order_id == Self::LAST_ADDED as i64 {
1012 sort_order_id = self.last_added_order_id.ok_or_else(|| {
1013 invalid_data!(
1014 "Cannot set default sort order to last added order: no order has been added."
1015 )
1016 })?;
1017 }
1018
1019 if self.metadata.default_sort_order_id == sort_order_id {
1020 return Ok(self);
1021 }
1022
1023 if !self.metadata.sort_orders.contains_key(&sort_order_id) {
1024 return Err(invalid_data!(
1025 "Cannot set default sort order to unknown order with id: '{sort_order_id}'"
1026 ));
1027 }
1028
1029 self.metadata.default_sort_order_id = sort_order_id;
1030
1031 if self.last_added_order_id == Some(sort_order_id) {
1032 self.changes.push(TableUpdate::SetDefaultSortOrder {
1033 sort_order_id: Self::LAST_ADDED as i64,
1034 });
1035 } else {
1036 self.changes
1037 .push(TableUpdate::SetDefaultSortOrder { sort_order_id });
1038 }
1039
1040 Ok(self)
1041 }
1042
1043 fn add_default_sort_order(self, sort_order: SortOrder) -> Result<Self> {
1045 self.add_sort_order(sort_order)?
1046 .set_default_sort_order(Self::LAST_ADDED as i64)
1047 }
1048
1049 pub fn add_encryption_key(mut self, key: EncryptedKey) -> Self {
1051 let key_id = key.key_id().to_string();
1052 if self.metadata.encryption_keys.contains_key(&key_id) {
1053 return self;
1055 }
1056
1057 self.metadata.encryption_keys.insert(key_id, key.clone());
1058 self.changes.push(TableUpdate::AddEncryptionKey {
1059 encryption_key: key,
1060 });
1061 self
1062 }
1063
1064 pub fn remove_encryption_key(mut self, key_id: &str) -> Self {
1066 if self.metadata.encryption_keys.remove(key_id).is_some() {
1067 self.changes.push(TableUpdate::RemoveEncryptionKey {
1068 key_id: key_id.to_string(),
1069 });
1070 }
1071 self
1072 }
1073
1074 pub fn build(mut self) -> Result<TableMetadataBuildResult> {
1076 self.metadata.last_updated_ms = self
1077 .last_updated_ms
1078 .unwrap_or_else(|| chrono::Utc::now().timestamp_millis());
1079
1080 let schema = self.get_current_schema()?.clone();
1084 let sort_order = Arc::unwrap_or_clone(self.get_default_sort_order()?);
1085
1086 self.metadata.default_spec = Arc::new(
1087 Arc::unwrap_or_clone(self.metadata.default_spec)
1088 .into_unbound()
1089 .bind(schema.clone())?,
1090 );
1091 self.metadata.default_partition_type =
1092 self.metadata.default_spec.partition_type(&schema)?;
1093 SortOrder::builder()
1094 .with_fields(sort_order.fields)
1095 .build(&schema)?;
1096
1097 self.update_snapshot_log()?;
1098 self.metadata.try_normalize()?;
1099
1100 if let Some(hist_entry) = self.previous_history_entry.take() {
1101 self.metadata.metadata_log.push(hist_entry);
1102 }
1103 let expired_metadata_logs = self.expire_metadata_log();
1104
1105 Ok(TableMetadataBuildResult {
1106 metadata: self.metadata,
1107 changes: self.changes,
1108 expired_metadata_logs,
1109 })
1110 }
1111
1112 fn expire_metadata_log(&mut self) -> Vec<MetadataLog> {
1113 let max_size = self
1114 .metadata
1115 .properties
1116 .get(TableProperties::PROPERTY_METADATA_PREVIOUS_VERSIONS_MAX)
1117 .and_then(|v| v.parse::<usize>().ok())
1118 .unwrap_or(TableProperties::PROPERTY_METADATA_PREVIOUS_VERSIONS_MAX_DEFAULT)
1119 .max(1);
1120
1121 if self.metadata.metadata_log.len() > max_size {
1122 self.metadata
1123 .metadata_log
1124 .drain(0..self.metadata.metadata_log.len() - max_size)
1125 .collect()
1126 } else {
1127 Vec::new()
1128 }
1129 }
1130
1131 fn update_snapshot_log(&mut self) -> Result<()> {
1132 let intermediate_snapshots = self.get_intermediate_snapshots();
1133 let has_removed_snapshots = self
1134 .changes
1135 .iter()
1136 .any(|update| matches!(update, TableUpdate::RemoveSnapshots { .. }));
1137
1138 if intermediate_snapshots.is_empty() && !has_removed_snapshots {
1139 return Ok(());
1140 }
1141
1142 let mut new_snapshot_log = Vec::new();
1143 for log_entry in &self.metadata.snapshot_log {
1144 let snapshot_id = log_entry.snapshot_id;
1145 if self.metadata.snapshots.contains_key(&snapshot_id) {
1146 if !intermediate_snapshots.contains(&snapshot_id) {
1147 new_snapshot_log.push(log_entry.clone());
1148 }
1149 } else if has_removed_snapshots {
1150 new_snapshot_log.clear();
1156 }
1157 }
1158
1159 if let Some(current_snapshot_id) = self.metadata.current_snapshot_id {
1160 let last_id = new_snapshot_log.last().map(|entry| entry.snapshot_id);
1161 if last_id != Some(current_snapshot_id) {
1162 return Err(invalid_data!(
1163 "Cannot set invalid snapshot log: latest entry is not the current snapshot"
1164 ));
1165 }
1166 };
1167
1168 self.metadata.snapshot_log = new_snapshot_log;
1169 Ok(())
1170 }
1171
1172 fn get_intermediate_snapshots(&self) -> HashSet<i64> {
1182 let added_snapshot_ids = self
1183 .changes
1184 .iter()
1185 .filter_map(|update| match update {
1186 TableUpdate::AddSnapshot { snapshot } => Some(snapshot.snapshot_id()),
1187 _ => None,
1188 })
1189 .collect::<HashSet<_>>();
1190
1191 self.changes
1192 .iter()
1193 .filter_map(|update| match update {
1194 TableUpdate::SetSnapshotRef {
1195 ref_name,
1196 reference,
1197 } => {
1198 if added_snapshot_ids.contains(&reference.snapshot_id)
1199 && ref_name == MAIN_BRANCH
1200 && reference.snapshot_id
1201 != self
1202 .metadata
1203 .current_snapshot_id
1204 .unwrap_or(i64::from(Self::LAST_ADDED))
1205 {
1206 Some(reference.snapshot_id)
1207 } else {
1208 None
1209 }
1210 }
1211 _ => None,
1212 })
1213 .collect()
1214 }
1215
1216 fn reassign_ids(
1217 schema: Schema,
1218 spec: UnboundPartitionSpec,
1219 sort_order: SortOrder,
1220 ) -> Result<(Schema, PartitionSpec, SortOrder)> {
1221 let previous_id_to_name = schema.field_id_to_name_map().clone();
1223 let fresh_schema = schema
1224 .into_builder()
1225 .with_schema_id(DEFAULT_SCHEMA_ID)
1226 .with_reassigned_field_ids(FIRST_FIELD_ID)
1227 .build()?;
1228
1229 let mut fresh_spec = PartitionSpecBuilder::new(fresh_schema.clone());
1231 for field in spec.fields() {
1232 let source_id = field.source_id()?;
1233 let source_field_name = previous_id_to_name.get(&source_id).ok_or_else(|| {
1234 invalid_data!(
1235 "Cannot find source column with id {} for partition column {} in schema.",
1236 source_id,
1237 field.name()
1238 )
1239 })?;
1240 fresh_spec = fresh_spec.add_partition_field(
1241 source_field_name,
1242 field.name(),
1243 field.transform(),
1244 )?;
1245 }
1246 let fresh_spec = fresh_spec.build()?;
1247
1248 let mut fresh_order = SortOrder::builder();
1250 for mut field in sort_order.fields {
1251 let source_field_name = previous_id_to_name.get(&field.source_id).ok_or_else(|| {
1252 invalid_data!(
1253 "Cannot find source column with id {} for sort column in schema.",
1254 field.source_id
1255 )
1256 })?;
1257 let new_field_id = fresh_schema
1258 .field_by_name(source_field_name)
1259 .ok_or_else(|| {
1260 Error::new(
1261 ErrorKind::Unexpected,
1262 format!(
1263 "Cannot find source column with name {source_field_name} for sort column in re-assigned schema."
1264 ),
1265 )
1266 })?.id;
1267 field.source_id = new_field_id;
1268 fresh_order.with_sort_field(field);
1269 }
1270 let fresh_sort_order = fresh_order.build(&fresh_schema)?;
1271
1272 Ok((fresh_schema, fresh_spec, fresh_sort_order))
1273 }
1274
1275 fn reuse_or_create_new_schema_id(&self, new_schema: &Schema) -> i32 {
1276 self.metadata
1277 .schemas
1278 .iter()
1279 .find_map(|(id, schema)| new_schema.is_same_schema(schema).then_some(*id))
1280 .unwrap_or_else(|| self.get_highest_schema_id() + 1)
1281 }
1282
1283 fn get_highest_schema_id(&self) -> i32 {
1284 *self
1285 .metadata
1286 .schemas
1287 .keys()
1288 .max()
1289 .unwrap_or(&self.metadata.current_schema_id)
1290 }
1291
1292 fn get_current_schema(&self) -> Result<&SchemaRef> {
1293 self.metadata
1294 .schemas
1295 .get(&self.metadata.current_schema_id)
1296 .ok_or_else(|| {
1297 invalid_data!(
1298 "Current schema with id '{}' not found in table metadata.",
1299 self.metadata.current_schema_id
1300 )
1301 })
1302 }
1303
1304 fn get_default_sort_order(&self) -> Result<SortOrderRef> {
1305 self.metadata
1306 .sort_orders
1307 .get(&self.metadata.default_sort_order_id)
1308 .cloned()
1309 .ok_or_else(|| {
1310 invalid_data!(
1311 "Default sort order with id '{}' not found in table metadata.",
1312 self.metadata.default_sort_order_id
1313 )
1314 })
1315 }
1316
1317 fn reuse_or_create_new_spec_id(&self, new_spec: &PartitionSpec) -> i32 {
1319 self.metadata
1320 .partition_specs
1321 .iter()
1322 .find_map(|(id, old_spec)| new_spec.is_compatible_with(old_spec).then_some(*id))
1323 .unwrap_or_else(|| {
1324 self.get_highest_spec_id()
1325 .map(|id| id + 1)
1326 .unwrap_or(DEFAULT_PARTITION_SPEC_ID)
1327 })
1328 }
1329
1330 fn get_highest_spec_id(&self) -> Option<i32> {
1331 self.metadata.partition_specs.keys().max().copied()
1332 }
1333
1334 fn reuse_or_create_new_sort_id(&self, new_sort_order: &SortOrder) -> i64 {
1336 if new_sort_order.is_unsorted() {
1337 return SortOrder::unsorted_order().order_id;
1338 }
1339
1340 self.metadata
1341 .sort_orders
1342 .iter()
1343 .find_map(|(id, sort_order)| {
1344 sort_order.fields.eq(&new_sort_order.fields).then_some(*id)
1345 })
1346 .unwrap_or_else(|| {
1347 self.highest_sort_order_id()
1348 .unwrap_or(SortOrder::unsorted_order().order_id)
1349 + 1
1350 })
1351 }
1352
1353 fn highest_sort_order_id(&self) -> Option<i64> {
1354 self.metadata.sort_orders.keys().max().copied()
1355 }
1356
1357 pub fn remove_schemas(mut self, schema_id_to_remove: &[i32]) -> Result<Self> {
1360 if schema_id_to_remove.contains(&self.metadata.current_schema_id) {
1361 return Err(invalid_data!("Cannot remove current schema"));
1362 }
1363
1364 if schema_id_to_remove.is_empty() {
1365 return Ok(self);
1366 }
1367
1368 let mut removed_schemas = Vec::with_capacity(schema_id_to_remove.len());
1369 self.metadata.schemas.retain(|id, _schema| {
1370 if schema_id_to_remove.contains(id) {
1371 removed_schemas.push(*id);
1372 false
1373 } else {
1374 true
1375 }
1376 });
1377
1378 self.changes.push(TableUpdate::RemoveSchemas {
1379 schema_ids: removed_schemas,
1380 });
1381
1382 Ok(self)
1383 }
1384}
1385
1386impl From<TableMetadataBuildResult> for TableMetadata {
1387 fn from(result: TableMetadataBuildResult) -> Self {
1388 result.metadata
1389 }
1390}
1391
1392#[cfg(test)]
1393mod tests {
1394 use std::fs::File;
1395 use std::io::BufReader;
1396 use std::thread::sleep;
1397
1398 use super::*;
1399 use crate::TableIdent;
1400 use crate::io::FileIO;
1401 use crate::spec::{
1402 BlobMetadata, NestedField, NullOrder, Operation, PartitionSpec, PrimitiveType, Schema,
1403 SnapshotRetention, SortDirection, SortField, StructType, Summary, TableProperties,
1404 Transform, Type, UnboundPartitionField,
1405 };
1406 use crate::table::Table;
1407 use crate::test_utils::test_runtime;
1408
1409 const TEST_LOCATION: &str = "s3://bucket/test/location";
1410 const LAST_ASSIGNED_COLUMN_ID: i32 = 3;
1411
1412 fn schema() -> Schema {
1413 Schema::builder()
1414 .with_fields(vec![
1415 NestedField::required(1, "x", Type::Primitive(PrimitiveType::Long)).into(),
1416 NestedField::required(2, "y", Type::Primitive(PrimitiveType::Long)).into(),
1417 NestedField::required(3, "z", Type::Primitive(PrimitiveType::Long)).into(),
1418 ])
1419 .build()
1420 .unwrap()
1421 }
1422
1423 fn sort_order() -> SortOrder {
1424 let schema = schema();
1425 SortOrder::builder()
1426 .with_order_id(1)
1427 .with_sort_field(SortField {
1428 source_id: 3,
1429 transform: Transform::Bucket(4),
1430 direction: SortDirection::Descending,
1431 null_order: NullOrder::First,
1432 })
1433 .build(&schema)
1434 .unwrap()
1435 }
1436
1437 fn partition_spec() -> UnboundPartitionSpec {
1438 UnboundPartitionSpec::builder()
1439 .with_spec_id(0)
1440 .add_partition_field(
1441 UnboundPartitionField::builder()
1442 .source_ids(vec![2])
1443 .name("y")
1444 .transform(Transform::Identity)
1445 .build()
1446 .unwrap(),
1447 )
1448 .unwrap()
1449 .build()
1450 }
1451
1452 fn builder_without_changes(format_version: FormatVersion) -> TableMetadataBuilder {
1453 TableMetadataBuilder::new(
1454 schema(),
1455 partition_spec(),
1456 sort_order(),
1457 TEST_LOCATION.to_string(),
1458 format_version,
1459 HashMap::new(),
1460 )
1461 .unwrap()
1462 .build()
1463 .unwrap()
1464 .metadata
1465 .into_builder(Some(
1466 "s3://bucket/test/location/metadata/metadata1.json".to_string(),
1467 ))
1468 }
1469
1470 #[test]
1471 fn test_minimal_build() {
1472 let metadata = TableMetadataBuilder::new(
1473 schema(),
1474 partition_spec(),
1475 sort_order(),
1476 TEST_LOCATION.to_string(),
1477 FormatVersion::V1,
1478 HashMap::new(),
1479 )
1480 .unwrap()
1481 .build()
1482 .unwrap()
1483 .metadata;
1484
1485 assert_eq!(metadata.format_version, FormatVersion::V1);
1486 assert_eq!(metadata.location, TEST_LOCATION);
1487 assert_eq!(metadata.current_schema_id, 0);
1488 assert_eq!(metadata.default_spec.spec_id(), 0);
1489 assert_eq!(metadata.default_sort_order_id, 1);
1490 assert_eq!(metadata.last_partition_id, 1000);
1491 assert_eq!(metadata.last_column_id, 3);
1492 assert_eq!(metadata.snapshots.len(), 0);
1493 assert_eq!(metadata.current_snapshot_id, None);
1494 assert_eq!(metadata.refs.len(), 0);
1495 assert_eq!(metadata.properties.len(), 0);
1496 assert_eq!(metadata.metadata_log.len(), 0);
1497 assert_eq!(metadata.last_sequence_number, 0);
1498 assert_eq!(metadata.last_column_id, LAST_ASSIGNED_COLUMN_ID);
1499
1500 let _ = serde_json::to_string(&metadata).unwrap();
1502
1503 let metadata = metadata
1505 .into_builder(Some(
1506 "s3://bucket/test/location/metadata/metadata1.json".to_string(),
1507 ))
1508 .upgrade_format_version(FormatVersion::V2)
1509 .unwrap()
1510 .build()
1511 .unwrap()
1512 .metadata;
1513
1514 assert_eq!(metadata.format_version, FormatVersion::V2);
1515 let _ = serde_json::to_string(&metadata).unwrap();
1516 }
1517
1518 #[test]
1519 fn test_build_unpartitioned_unsorted() {
1520 let schema = Schema::builder().build().unwrap();
1521 let metadata = TableMetadataBuilder::new(
1522 schema.clone(),
1523 PartitionSpec::unpartition_spec(),
1524 SortOrder::unsorted_order(),
1525 TEST_LOCATION.to_string(),
1526 FormatVersion::V2,
1527 HashMap::new(),
1528 )
1529 .unwrap()
1530 .build()
1531 .unwrap()
1532 .metadata;
1533
1534 assert_eq!(metadata.format_version, FormatVersion::V2);
1535 assert_eq!(metadata.location, TEST_LOCATION);
1536 assert_eq!(metadata.current_schema_id, 0);
1537 assert_eq!(metadata.default_spec.spec_id(), 0);
1538 assert_eq!(metadata.default_sort_order_id, 0);
1539 assert_eq!(metadata.last_partition_id, UNPARTITIONED_LAST_ASSIGNED_ID);
1540 assert_eq!(metadata.last_column_id, 0);
1541 assert_eq!(metadata.snapshots.len(), 0);
1542 assert_eq!(metadata.current_snapshot_id, None);
1543 assert_eq!(metadata.refs.len(), 0);
1544 assert_eq!(metadata.properties.len(), 0);
1545 assert_eq!(metadata.metadata_log.len(), 0);
1546 assert_eq!(metadata.last_sequence_number, 0);
1547 }
1548
1549 #[test]
1550 fn test_reassigns_ids() {
1551 let schema = Schema::builder()
1552 .with_schema_id(10)
1553 .with_fields(vec![
1554 NestedField::required(11, "a", Type::Primitive(PrimitiveType::Long)).into(),
1555 NestedField::required(12, "b", Type::Primitive(PrimitiveType::Long)).into(),
1556 NestedField::required(
1557 13,
1558 "struct",
1559 Type::Struct(StructType::new(vec![
1560 NestedField::required(14, "nested", Type::Primitive(PrimitiveType::Long))
1561 .into(),
1562 ])),
1563 )
1564 .into(),
1565 NestedField::required(15, "c", Type::Primitive(PrimitiveType::Long)).into(),
1566 ])
1567 .build()
1568 .unwrap();
1569 let spec = PartitionSpec::builder(schema.clone())
1570 .with_spec_id(20)
1571 .add_partition_field("a", "a", Transform::Identity)
1572 .unwrap()
1573 .add_partition_field("struct.nested", "nested_partition", Transform::Identity)
1574 .unwrap()
1575 .build()
1576 .unwrap();
1577 let sort_order = SortOrder::builder()
1578 .with_fields(vec![SortField {
1579 source_id: 11,
1580 transform: Transform::Identity,
1581 direction: SortDirection::Ascending,
1582 null_order: NullOrder::First,
1583 }])
1584 .with_order_id(10)
1585 .build(&schema)
1586 .unwrap();
1587
1588 let (fresh_schema, fresh_spec, fresh_sort_order) =
1589 TableMetadataBuilder::reassign_ids(schema, spec.into_unbound(), sort_order).unwrap();
1590
1591 let expected_schema = Schema::builder()
1592 .with_fields(vec![
1593 NestedField::required(1, "a", Type::Primitive(PrimitiveType::Long)).into(),
1594 NestedField::required(2, "b", Type::Primitive(PrimitiveType::Long)).into(),
1595 NestedField::required(
1596 3,
1597 "struct",
1598 Type::Struct(StructType::new(vec![
1599 NestedField::required(5, "nested", Type::Primitive(PrimitiveType::Long))
1600 .into(),
1601 ])),
1602 )
1603 .into(),
1604 NestedField::required(4, "c", Type::Primitive(PrimitiveType::Long)).into(),
1605 ])
1606 .build()
1607 .unwrap();
1608
1609 let expected_spec = PartitionSpec::builder(expected_schema.clone())
1610 .with_spec_id(0)
1611 .add_partition_field("a", "a", Transform::Identity)
1612 .unwrap()
1613 .add_partition_field("struct.nested", "nested_partition", Transform::Identity)
1614 .unwrap()
1615 .build()
1616 .unwrap();
1617
1618 let expected_sort_order = SortOrder::builder()
1619 .with_fields(vec![SortField {
1620 source_id: 1,
1621 transform: Transform::Identity,
1622 direction: SortDirection::Ascending,
1623 null_order: NullOrder::First,
1624 }])
1625 .with_order_id(1)
1626 .build(&expected_schema)
1627 .unwrap();
1628
1629 assert_eq!(fresh_schema, expected_schema);
1630 assert_eq!(fresh_spec, expected_spec);
1631 assert_eq!(fresh_sort_order, expected_sort_order);
1632 }
1633
1634 #[test]
1635 fn test_ids_are_reassigned_for_new_metadata() {
1636 let schema = schema().into_builder().with_schema_id(10).build().unwrap();
1637
1638 let metadata = TableMetadataBuilder::new(
1639 schema,
1640 partition_spec(),
1641 sort_order(),
1642 TEST_LOCATION.to_string(),
1643 FormatVersion::V1,
1644 HashMap::new(),
1645 )
1646 .unwrap()
1647 .build()
1648 .unwrap()
1649 .metadata;
1650
1651 assert_eq!(metadata.current_schema_id, 0);
1652 assert_eq!(metadata.current_schema().schema_id(), 0);
1653 }
1654
1655 #[test]
1656 fn test_new_metadata_changes() {
1657 let changes = TableMetadataBuilder::new(
1658 schema(),
1659 partition_spec(),
1660 sort_order(),
1661 TEST_LOCATION.to_string(),
1662 FormatVersion::V1,
1663 HashMap::from_iter(vec![("property 1".to_string(), "value 1".to_string())]),
1664 )
1665 .unwrap()
1666 .build()
1667 .unwrap()
1668 .changes;
1669
1670 pretty_assertions::assert_eq!(changes, vec![
1671 TableUpdate::SetLocation {
1672 location: TEST_LOCATION.to_string()
1673 },
1674 TableUpdate::AddSchema { schema: schema() },
1675 TableUpdate::SetCurrentSchema { schema_id: -1 },
1676 TableUpdate::AddSpec {
1677 spec: PartitionSpec::builder(schema())
1680 .with_spec_id(0)
1681 .add_unbound_field(
1682 UnboundPartitionField::builder()
1683 .source_ids(vec![2])
1684 .field_id(1000)
1685 .name("y".to_string())
1686 .transform(Transform::Identity)
1687 .build()
1688 .unwrap()
1689 )
1690 .unwrap()
1691 .build()
1692 .unwrap()
1693 .into_unbound(),
1694 },
1695 TableUpdate::SetDefaultSpec { spec_id: -1 },
1696 TableUpdate::AddSortOrder {
1697 sort_order: sort_order(),
1698 },
1699 TableUpdate::SetDefaultSortOrder { sort_order_id: -1 },
1700 TableUpdate::SetProperties {
1701 updates: HashMap::from_iter(vec![(
1702 "property 1".to_string(),
1703 "value 1".to_string()
1704 )]),
1705 }
1706 ]);
1707 }
1708
1709 #[test]
1710 fn test_new_metadata_changes_unpartitioned_unsorted() {
1711 let schema = Schema::builder().build().unwrap();
1712 let changes = TableMetadataBuilder::new(
1713 schema.clone(),
1714 PartitionSpec::unpartition_spec().into_unbound(),
1715 SortOrder::unsorted_order(),
1716 TEST_LOCATION.to_string(),
1717 FormatVersion::V1,
1718 HashMap::new(),
1719 )
1720 .unwrap()
1721 .build()
1722 .unwrap()
1723 .changes;
1724
1725 pretty_assertions::assert_eq!(changes, vec![
1726 TableUpdate::SetLocation {
1727 location: TEST_LOCATION.to_string()
1728 },
1729 TableUpdate::AddSchema {
1730 schema: Schema::builder().build().unwrap(),
1731 },
1732 TableUpdate::SetCurrentSchema { schema_id: -1 },
1733 TableUpdate::AddSpec {
1734 spec: PartitionSpec::builder(schema)
1737 .with_spec_id(0)
1738 .build()
1739 .unwrap()
1740 .into_unbound(),
1741 },
1742 TableUpdate::SetDefaultSpec { spec_id: -1 },
1743 TableUpdate::AddSortOrder {
1744 sort_order: SortOrder::unsorted_order(),
1745 },
1746 TableUpdate::SetDefaultSortOrder { sort_order_id: -1 },
1747 ]);
1748 }
1749
1750 #[test]
1751 fn test_add_partition_spec() {
1752 let builder = builder_without_changes(FormatVersion::V2);
1753
1754 let added_spec = UnboundPartitionSpec::builder()
1755 .with_spec_id(10)
1756 .add_partition_fields(vec![
1757 UnboundPartitionField::builder()
1759 .source_ids(vec![2])
1760 .field_id(1000)
1761 .name("y".to_string())
1762 .transform(Transform::Identity)
1763 .build()
1764 .unwrap(),
1765 UnboundPartitionField::builder()
1767 .source_ids(vec![3])
1768 .name("z".to_string())
1769 .transform(Transform::Identity)
1770 .build()
1771 .unwrap(),
1772 ])
1773 .unwrap()
1774 .build();
1775
1776 let build_result = builder
1777 .add_partition_spec(added_spec.clone())
1778 .unwrap()
1779 .build()
1780 .unwrap();
1781
1782 let expected_change = added_spec.with_spec_id(1);
1784 let expected_spec = PartitionSpec::builder(schema())
1785 .with_spec_id(1)
1786 .add_unbound_field(
1787 UnboundPartitionField::builder()
1788 .source_ids(vec![2])
1789 .field_id(1000)
1790 .name("y".to_string())
1791 .transform(Transform::Identity)
1792 .build()
1793 .unwrap(),
1794 )
1795 .unwrap()
1796 .add_unbound_field(
1797 UnboundPartitionField::builder()
1798 .source_ids(vec![3])
1799 .field_id(1001)
1800 .name("z".to_string())
1801 .transform(Transform::Identity)
1802 .build()
1803 .unwrap(),
1804 )
1805 .unwrap()
1806 .build()
1807 .unwrap();
1808
1809 assert_eq!(build_result.changes.len(), 1);
1810 assert_eq!(
1811 build_result.metadata.partition_spec_by_id(1),
1812 Some(&Arc::new(expected_spec))
1813 );
1814 assert_eq!(build_result.metadata.default_spec.spec_id(), 0);
1815 assert_eq!(build_result.metadata.last_partition_id, 1001);
1816 pretty_assertions::assert_eq!(build_result.changes[0], TableUpdate::AddSpec {
1817 spec: expected_change
1818 });
1819
1820 let build_result = build_result
1822 .metadata
1823 .into_builder(Some(
1824 "s3://bucket/test/location/metadata/metadata1.json".to_string(),
1825 ))
1826 .remove_partition_specs(&[1])
1827 .unwrap()
1828 .build()
1829 .unwrap();
1830
1831 assert_eq!(build_result.changes.len(), 1);
1832 assert_eq!(build_result.metadata.partition_specs.len(), 1);
1833 assert!(build_result.metadata.partition_spec_by_id(1).is_none());
1834 }
1835
1836 #[test]
1837 fn test_set_default_partition_spec() {
1838 let builder = builder_without_changes(FormatVersion::V2);
1839 let schema = builder.get_current_schema().unwrap().clone();
1840 let added_spec = UnboundPartitionSpec::builder()
1841 .with_spec_id(10)
1842 .add_partition_field(
1843 UnboundPartitionField::builder()
1844 .source_ids(vec![1])
1845 .name("y_bucket[2]")
1846 .transform(Transform::Bucket(2))
1847 .build()
1848 .unwrap(),
1849 )
1850 .unwrap()
1851 .build();
1852
1853 let build_result = builder
1854 .add_partition_spec(added_spec.clone())
1855 .unwrap()
1856 .set_default_partition_spec(-1)
1857 .unwrap()
1858 .build()
1859 .unwrap();
1860
1861 let expected_spec = PartitionSpec::builder(schema)
1862 .with_spec_id(1)
1863 .add_unbound_field(
1864 UnboundPartitionField::builder()
1865 .source_ids(vec![1])
1866 .field_id(1001)
1867 .name("y_bucket[2]".to_string())
1868 .transform(Transform::Bucket(2))
1869 .build()
1870 .unwrap(),
1871 )
1872 .unwrap()
1873 .build()
1874 .unwrap();
1875
1876 assert_eq!(build_result.changes.len(), 2);
1877 assert_eq!(build_result.metadata.default_spec, Arc::new(expected_spec));
1878 assert_eq!(build_result.changes, vec![
1879 TableUpdate::AddSpec {
1880 spec: added_spec.with_spec_id(1)
1882 },
1883 TableUpdate::SetDefaultSpec { spec_id: -1 }
1884 ]);
1885 }
1886
1887 #[test]
1888 fn test_set_existing_default_partition_spec() {
1889 let builder = builder_without_changes(FormatVersion::V2);
1890 let unbound_spec = UnboundPartitionSpec::builder().with_spec_id(1).build();
1892 let build_result = builder
1893 .add_partition_spec(unbound_spec.clone())
1894 .unwrap()
1895 .set_default_partition_spec(-1)
1896 .unwrap()
1897 .build()
1898 .unwrap();
1899
1900 assert_eq!(build_result.changes.len(), 2);
1901 assert_eq!(build_result.changes[0], TableUpdate::AddSpec {
1902 spec: unbound_spec.clone()
1903 });
1904 assert_eq!(build_result.changes[1], TableUpdate::SetDefaultSpec {
1905 spec_id: -1
1906 });
1907 assert_eq!(
1908 build_result.metadata.default_spec,
1909 Arc::new(
1910 unbound_spec
1911 .bind(build_result.metadata.current_schema().clone())
1912 .unwrap()
1913 )
1914 );
1915
1916 let build_result = build_result
1918 .metadata
1919 .into_builder(Some(
1920 "s3://bucket/test/location/metadata/metadata1.json".to_string(),
1921 ))
1922 .set_default_partition_spec(0)
1923 .unwrap()
1924 .build()
1925 .unwrap();
1926
1927 assert_eq!(build_result.changes.len(), 1);
1928 assert_eq!(build_result.changes[0], TableUpdate::SetDefaultSpec {
1929 spec_id: 0
1930 });
1931 assert_eq!(
1932 build_result.metadata.default_spec,
1933 Arc::new(
1934 partition_spec()
1935 .bind(build_result.metadata.current_schema().clone())
1936 .unwrap()
1937 )
1938 );
1939 }
1940
1941 #[test]
1942 fn test_add_sort_order() {
1943 let builder = builder_without_changes(FormatVersion::V2);
1944
1945 let added_sort_order = SortOrder::builder()
1946 .with_order_id(10)
1947 .with_fields(vec![SortField {
1948 source_id: 1,
1949 transform: Transform::Identity,
1950 direction: SortDirection::Ascending,
1951 null_order: NullOrder::First,
1952 }])
1953 .build(&schema())
1954 .unwrap();
1955
1956 let build_result = builder
1957 .add_sort_order(added_sort_order.clone())
1958 .unwrap()
1959 .build()
1960 .unwrap();
1961
1962 let expected_sort_order = added_sort_order.with_order_id(2);
1963
1964 assert_eq!(build_result.changes.len(), 1);
1965 assert_eq!(build_result.metadata.sort_orders.keys().max(), Some(&2));
1966 pretty_assertions::assert_eq!(
1967 build_result.metadata.sort_order_by_id(2),
1968 Some(&Arc::new(expected_sort_order.clone()))
1969 );
1970 pretty_assertions::assert_eq!(build_result.changes[0], TableUpdate::AddSortOrder {
1971 sort_order: expected_sort_order
1972 });
1973 }
1974
1975 #[test]
1976 fn test_add_compatible_schema() {
1977 let builder = builder_without_changes(FormatVersion::V2);
1978
1979 let added_schema = Schema::builder()
1980 .with_schema_id(1)
1981 .with_fields(vec![
1982 NestedField::required(1, "x", Type::Primitive(PrimitiveType::Long)).into(),
1983 NestedField::required(2, "y", Type::Primitive(PrimitiveType::Long)).into(),
1984 NestedField::required(3, "z", Type::Primitive(PrimitiveType::Long)).into(),
1985 NestedField::required(4, "a", Type::Primitive(PrimitiveType::Long)).into(),
1986 ])
1987 .build()
1988 .unwrap();
1989
1990 let build_result = builder
1991 .add_current_schema(added_schema.clone())
1992 .unwrap()
1993 .build()
1994 .unwrap();
1995
1996 assert_eq!(build_result.changes.len(), 2);
1997 assert_eq!(build_result.metadata.schemas.keys().max(), Some(&1));
1998 pretty_assertions::assert_eq!(
1999 build_result.metadata.schema_by_id(1),
2000 Some(&Arc::new(added_schema.clone()))
2001 );
2002 pretty_assertions::assert_eq!(build_result.changes[0], TableUpdate::AddSchema {
2003 schema: added_schema
2004 });
2005 assert_eq!(build_result.changes[1], TableUpdate::SetCurrentSchema {
2006 schema_id: -1
2007 });
2008 }
2009
2010 #[test]
2011 fn test_set_current_schema_change_is_minus_one_if_schema_was_added_in_this_change() {
2012 let builder = builder_without_changes(FormatVersion::V2);
2013
2014 let added_schema = Schema::builder()
2015 .with_schema_id(1)
2016 .with_fields(vec![
2017 NestedField::required(1, "x", Type::Primitive(PrimitiveType::Long)).into(),
2018 NestedField::required(2, "y", Type::Primitive(PrimitiveType::Long)).into(),
2019 NestedField::required(3, "z", Type::Primitive(PrimitiveType::Long)).into(),
2020 NestedField::required(4, "a", Type::Primitive(PrimitiveType::Long)).into(),
2021 ])
2022 .build()
2023 .unwrap();
2024
2025 let build_result = builder
2026 .add_schema(added_schema.clone())
2027 .unwrap()
2028 .set_current_schema(1)
2029 .unwrap()
2030 .build()
2031 .unwrap();
2032
2033 assert_eq!(build_result.changes.len(), 2);
2034 assert_eq!(build_result.changes[1], TableUpdate::SetCurrentSchema {
2035 schema_id: -1
2036 });
2037 }
2038
2039 #[test]
2040 fn test_no_metadata_log_for_create_table() {
2041 let build_result = TableMetadataBuilder::new(
2042 schema(),
2043 partition_spec(),
2044 sort_order(),
2045 TEST_LOCATION.to_string(),
2046 FormatVersion::V2,
2047 HashMap::new(),
2048 )
2049 .unwrap()
2050 .build()
2051 .unwrap();
2052
2053 assert_eq!(build_result.metadata.metadata_log.len(), 0);
2054 }
2055
2056 #[test]
2057 fn test_table_properties_view_reflects_metadata_updates() {
2058 let property = TableProperties::PROPERTY_COMMIT_NUM_RETRIES.to_string();
2059 let metadata = builder_without_changes(FormatVersion::V2)
2060 .set_properties(HashMap::from([(property.clone(), "7".to_string())]))
2061 .unwrap()
2062 .build()
2063 .unwrap()
2064 .metadata;
2065
2066 assert_eq!(metadata.table_properties().commit_num_retries().unwrap(), 7);
2067
2068 let metadata = metadata
2069 .into_builder(None)
2070 .remove_properties(&[property])
2071 .unwrap()
2072 .build()
2073 .unwrap()
2074 .metadata;
2075
2076 assert_eq!(
2077 metadata.table_properties().commit_num_retries().unwrap(),
2078 TableProperties::PROPERTY_COMMIT_NUM_RETRIES_DEFAULT
2079 );
2080 }
2081
2082 #[test]
2083 fn test_no_metadata_log_entry_for_no_previous_location() {
2084 let metadata = builder_without_changes(FormatVersion::V2)
2086 .build()
2087 .unwrap()
2088 .metadata;
2089 assert_eq!(metadata.metadata_log.len(), 1);
2090
2091 let build_result = metadata
2092 .into_builder(None)
2093 .set_properties(HashMap::from_iter(vec![(
2094 "foo".to_string(),
2095 "bar".to_string(),
2096 )]))
2097 .unwrap()
2098 .build()
2099 .unwrap();
2100
2101 assert_eq!(build_result.metadata.metadata_log.len(), 1);
2102 }
2103
2104 #[test]
2105 fn test_from_metadata_generates_metadata_log() {
2106 let metadata_path = "s3://bucket/test/location/metadata/metadata1.json";
2107 let builder = TableMetadataBuilder::new(
2108 schema(),
2109 partition_spec(),
2110 sort_order(),
2111 TEST_LOCATION.to_string(),
2112 FormatVersion::V2,
2113 HashMap::new(),
2114 )
2115 .unwrap()
2116 .build()
2117 .unwrap()
2118 .metadata
2119 .into_builder(Some(metadata_path.to_string()));
2120
2121 let builder = builder
2122 .add_default_sort_order(SortOrder::unsorted_order())
2123 .unwrap();
2124
2125 let build_result = builder.build().unwrap();
2126
2127 assert_eq!(build_result.metadata.metadata_log.len(), 1);
2128 assert_eq!(
2129 build_result.metadata.metadata_log[0].metadata_file,
2130 metadata_path
2131 );
2132 }
2133
2134 #[test]
2135 fn test_set_ref() {
2136 let builder = builder_without_changes(FormatVersion::V2);
2137
2138 let snapshot = Snapshot::builder()
2139 .with_snapshot_id(1)
2140 .with_timestamp_ms(builder.metadata.last_updated_ms + 1)
2141 .with_sequence_number(0)
2142 .with_schema_id(0)
2143 .with_manifest_list("/snap-1.avro")
2144 .with_summary(Summary {
2145 operation: Operation::Append,
2146 additional_properties: HashMap::from_iter(vec![
2147 (
2148 "spark.app.id".to_string(),
2149 "local-1662532784305".to_string(),
2150 ),
2151 ("added-data-files".to_string(), "4".to_string()),
2152 ("added-records".to_string(), "4".to_string()),
2153 ("added-files-size".to_string(), "6001".to_string()),
2154 ]),
2155 })
2156 .build();
2157
2158 let builder = builder.add_snapshot(snapshot.clone()).unwrap();
2159
2160 assert!(
2161 builder
2162 .clone()
2163 .set_ref(MAIN_BRANCH, SnapshotReference {
2164 snapshot_id: 10,
2165 retention: SnapshotRetention::Branch {
2166 min_snapshots_to_keep: Some(10),
2167 max_snapshot_age_ms: None,
2168 max_ref_age_ms: None,
2169 },
2170 })
2171 .unwrap_err()
2172 .to_string()
2173 .contains("Cannot set 'main' to unknown snapshot: '10'")
2174 );
2175
2176 let build_result = builder
2177 .set_ref(MAIN_BRANCH, SnapshotReference {
2178 snapshot_id: 1,
2179 retention: SnapshotRetention::Branch {
2180 min_snapshots_to_keep: Some(10),
2181 max_snapshot_age_ms: None,
2182 max_ref_age_ms: None,
2183 },
2184 })
2185 .unwrap()
2186 .build()
2187 .unwrap();
2188 assert_eq!(build_result.metadata.snapshots.len(), 1);
2189 assert_eq!(
2190 build_result.metadata.snapshot_by_id(1),
2191 Some(&Arc::new(snapshot.clone()))
2192 );
2193 assert_eq!(build_result.metadata.snapshot_log, vec![SnapshotLog {
2194 snapshot_id: 1,
2195 timestamp_ms: snapshot.timestamp_ms()
2196 }])
2197 }
2198
2199 #[test]
2200 fn test_snapshot_log_skips_intermediates() {
2201 let builder = builder_without_changes(FormatVersion::V2);
2202
2203 let snapshot_1 = Snapshot::builder()
2204 .with_snapshot_id(1)
2205 .with_timestamp_ms(builder.metadata.last_updated_ms + 1)
2206 .with_sequence_number(0)
2207 .with_schema_id(0)
2208 .with_manifest_list("/snap-1.avro")
2209 .with_summary(Summary {
2210 operation: Operation::Append,
2211 additional_properties: HashMap::from_iter(vec![
2212 (
2213 "spark.app.id".to_string(),
2214 "local-1662532784305".to_string(),
2215 ),
2216 ("added-data-files".to_string(), "4".to_string()),
2217 ("added-records".to_string(), "4".to_string()),
2218 ("added-files-size".to_string(), "6001".to_string()),
2219 ]),
2220 })
2221 .build();
2222
2223 let snapshot_2 = Snapshot::builder()
2224 .with_snapshot_id(2)
2225 .with_timestamp_ms(builder.metadata.last_updated_ms + 1)
2226 .with_sequence_number(0)
2227 .with_schema_id(0)
2228 .with_manifest_list("/snap-1.avro")
2229 .with_summary(Summary {
2230 operation: Operation::Append,
2231 additional_properties: HashMap::from_iter(vec![
2232 (
2233 "spark.app.id".to_string(),
2234 "local-1662532784305".to_string(),
2235 ),
2236 ("added-data-files".to_string(), "4".to_string()),
2237 ("added-records".to_string(), "4".to_string()),
2238 ("added-files-size".to_string(), "6001".to_string()),
2239 ]),
2240 })
2241 .build();
2242
2243 let result = builder
2244 .add_snapshot(snapshot_1)
2245 .unwrap()
2246 .set_ref(MAIN_BRANCH, SnapshotReference {
2247 snapshot_id: 1,
2248 retention: SnapshotRetention::Branch {
2249 min_snapshots_to_keep: Some(10),
2250 max_snapshot_age_ms: None,
2251 max_ref_age_ms: None,
2252 },
2253 })
2254 .unwrap()
2255 .set_branch_snapshot(snapshot_2.clone(), MAIN_BRANCH)
2256 .unwrap()
2257 .build()
2258 .unwrap();
2259
2260 assert_eq!(result.metadata.snapshot_log, vec![SnapshotLog {
2261 snapshot_id: 2,
2262 timestamp_ms: snapshot_2.timestamp_ms()
2263 }]);
2264 assert_eq!(result.metadata.current_snapshot().unwrap().snapshot_id(), 2);
2265 }
2266
2267 #[test]
2268 fn test_remove_main_ref_keeps_snapshot_log() {
2269 let builder = builder_without_changes(FormatVersion::V2);
2270
2271 let snapshot = Snapshot::builder()
2272 .with_snapshot_id(1)
2273 .with_timestamp_ms(builder.metadata.last_updated_ms + 1)
2274 .with_sequence_number(0)
2275 .with_schema_id(0)
2276 .with_manifest_list("/snap-1.avro")
2277 .with_summary(Summary {
2278 operation: Operation::Append,
2279 additional_properties: HashMap::from_iter(vec![
2280 (
2281 "spark.app.id".to_string(),
2282 "local-1662532784305".to_string(),
2283 ),
2284 ("added-data-files".to_string(), "4".to_string()),
2285 ("added-records".to_string(), "4".to_string()),
2286 ("added-files-size".to_string(), "6001".to_string()),
2287 ]),
2288 })
2289 .build();
2290
2291 let result = builder
2292 .add_snapshot(snapshot.clone())
2293 .unwrap()
2294 .set_ref(MAIN_BRANCH, SnapshotReference {
2295 snapshot_id: 1,
2296 retention: SnapshotRetention::Branch {
2297 min_snapshots_to_keep: Some(10),
2298 max_snapshot_age_ms: None,
2299 max_ref_age_ms: None,
2300 },
2301 })
2302 .unwrap()
2303 .build()
2304 .unwrap();
2305
2306 assert_eq!(result.metadata.snapshot_log.len(), 1);
2308 assert_eq!(result.metadata.snapshot_log[0].snapshot_id, 1);
2309 assert_eq!(result.metadata.current_snapshot_id, Some(1));
2310
2311 let result_after_remove = result
2313 .metadata
2314 .into_builder(Some(
2315 "s3://bucket/test/location/metadata/metadata2.json".to_string(),
2316 ))
2317 .remove_ref(MAIN_BRANCH)
2318 .build()
2319 .unwrap();
2320
2321 assert_eq!(result_after_remove.metadata.snapshot_log.len(), 1);
2323 assert_eq!(result_after_remove.metadata.snapshot_log[0].snapshot_id, 1);
2324 assert_eq!(result_after_remove.metadata.current_snapshot_id, None);
2325 assert_eq!(result_after_remove.changes.len(), 1);
2326 assert_eq!(
2327 result_after_remove.changes[0],
2328 TableUpdate::RemoveSnapshotRef {
2329 ref_name: MAIN_BRANCH.to_string()
2330 }
2331 );
2332 }
2333
2334 #[test]
2335 fn test_set_branch_snapshot_creates_branch_if_not_exists() {
2336 let builder = builder_without_changes(FormatVersion::V2);
2337
2338 let snapshot = Snapshot::builder()
2339 .with_snapshot_id(2)
2340 .with_timestamp_ms(builder.metadata.last_updated_ms + 1)
2341 .with_sequence_number(0)
2342 .with_schema_id(0)
2343 .with_manifest_list("/snap-1.avro")
2344 .with_summary(Summary {
2345 operation: Operation::Append,
2346 additional_properties: HashMap::new(),
2347 })
2348 .build();
2349
2350 let build_result = builder
2351 .set_branch_snapshot(snapshot.clone(), "new_branch")
2352 .unwrap()
2353 .build()
2354 .unwrap();
2355
2356 let reference = SnapshotReference {
2357 snapshot_id: 2,
2358 retention: SnapshotRetention::Branch {
2359 min_snapshots_to_keep: None,
2360 max_snapshot_age_ms: None,
2361 max_ref_age_ms: None,
2362 },
2363 };
2364
2365 assert_eq!(build_result.metadata.refs.len(), 1);
2366 assert_eq!(
2367 build_result.metadata.refs.get("new_branch"),
2368 Some(&reference)
2369 );
2370 assert_eq!(build_result.changes, vec![
2371 TableUpdate::AddSnapshot { snapshot },
2372 TableUpdate::SetSnapshotRef {
2373 ref_name: "new_branch".to_string(),
2374 reference
2375 }
2376 ]);
2377 }
2378
2379 #[test]
2380 fn test_cannot_add_duplicate_snapshot_id() {
2381 let builder = builder_without_changes(FormatVersion::V2);
2382
2383 let snapshot = Snapshot::builder()
2384 .with_snapshot_id(2)
2385 .with_timestamp_ms(builder.metadata.last_updated_ms + 1)
2386 .with_sequence_number(0)
2387 .with_schema_id(0)
2388 .with_manifest_list("/snap-1.avro")
2389 .with_summary(Summary {
2390 operation: Operation::Append,
2391 additional_properties: HashMap::from_iter(vec![
2392 (
2393 "spark.app.id".to_string(),
2394 "local-1662532784305".to_string(),
2395 ),
2396 ("added-data-files".to_string(), "4".to_string()),
2397 ("added-records".to_string(), "4".to_string()),
2398 ("added-files-size".to_string(), "6001".to_string()),
2399 ]),
2400 })
2401 .build();
2402
2403 let builder = builder.add_snapshot(snapshot.clone()).unwrap();
2404 builder.add_snapshot(snapshot).unwrap_err();
2405 }
2406
2407 #[test]
2408 fn test_add_incompatible_current_schema_fails() {
2409 let builder = builder_without_changes(FormatVersion::V2);
2410
2411 let added_schema = Schema::builder()
2412 .with_schema_id(1)
2413 .with_fields(vec![])
2414 .build()
2415 .unwrap();
2416
2417 let err = builder
2418 .add_current_schema(added_schema)
2419 .unwrap()
2420 .build()
2421 .unwrap_err();
2422
2423 assert!(
2424 err.to_string()
2425 .contains("Cannot find partition source field")
2426 );
2427 }
2428
2429 #[test]
2430 fn test_add_partition_spec_for_v1_requires_sequential_ids() {
2431 let builder = builder_without_changes(FormatVersion::V1);
2432
2433 let added_spec = UnboundPartitionSpec::builder()
2434 .with_spec_id(10)
2435 .add_partition_fields(vec![
2436 UnboundPartitionField::builder()
2437 .source_ids(vec![2])
2438 .field_id(1000)
2439 .name("y".to_string())
2440 .transform(Transform::Identity)
2441 .build()
2442 .unwrap(),
2443 UnboundPartitionField::builder()
2444 .source_ids(vec![3])
2445 .field_id(1002)
2446 .name("z".to_string())
2447 .transform(Transform::Identity)
2448 .build()
2449 .unwrap(),
2450 ])
2451 .unwrap()
2452 .build();
2453
2454 let err = builder.add_partition_spec(added_spec).unwrap_err();
2455 assert!(err.to_string().contains(
2456 "Cannot add partition spec with non-sequential field ids to format version 1 table"
2457 ));
2458 }
2459
2460 #[test]
2461 fn test_expire_metadata_log() {
2462 let builder = builder_without_changes(FormatVersion::V2);
2463 let metadata = builder
2464 .set_properties(HashMap::from_iter(vec![(
2465 TableProperties::PROPERTY_METADATA_PREVIOUS_VERSIONS_MAX.to_string(),
2466 "2".to_string(),
2467 )]))
2468 .unwrap()
2469 .build()
2470 .unwrap();
2471 assert_eq!(metadata.metadata.metadata_log.len(), 1);
2472 assert_eq!(metadata.expired_metadata_logs.len(), 0);
2473
2474 let metadata = metadata
2475 .metadata
2476 .into_builder(Some("path2".to_string()))
2477 .set_properties(HashMap::from_iter(vec![(
2478 "change_nr".to_string(),
2479 "1".to_string(),
2480 )]))
2481 .unwrap()
2482 .build()
2483 .unwrap();
2484
2485 assert_eq!(metadata.metadata.metadata_log.len(), 2);
2486 assert_eq!(metadata.expired_metadata_logs.len(), 0);
2487
2488 let metadata = metadata
2489 .metadata
2490 .into_builder(Some("path2".to_string()))
2491 .set_properties(HashMap::from_iter(vec![(
2492 "change_nr".to_string(),
2493 "2".to_string(),
2494 )]))
2495 .unwrap()
2496 .build()
2497 .unwrap();
2498 assert_eq!(metadata.metadata.metadata_log.len(), 2);
2499 assert_eq!(metadata.expired_metadata_logs.len(), 1);
2500 }
2501
2502 #[test]
2503 fn test_v2_sequence_number_cannot_decrease() {
2504 let builder = builder_without_changes(FormatVersion::V2);
2505
2506 let snapshot = Snapshot::builder()
2507 .with_snapshot_id(1)
2508 .with_timestamp_ms(builder.metadata.last_updated_ms + 1)
2509 .with_sequence_number(1)
2510 .with_schema_id(0)
2511 .with_manifest_list("/snap-1")
2512 .with_summary(Summary {
2513 operation: Operation::Append,
2514 additional_properties: HashMap::new(),
2515 })
2516 .build();
2517
2518 let builder = builder
2519 .add_snapshot(snapshot.clone())
2520 .unwrap()
2521 .set_ref(MAIN_BRANCH, SnapshotReference {
2522 snapshot_id: 1,
2523 retention: SnapshotRetention::Branch {
2524 min_snapshots_to_keep: Some(10),
2525 max_snapshot_age_ms: None,
2526 max_ref_age_ms: None,
2527 },
2528 })
2529 .unwrap();
2530
2531 let snapshot = Snapshot::builder()
2532 .with_snapshot_id(2)
2533 .with_timestamp_ms(builder.metadata.last_updated_ms + 1)
2534 .with_sequence_number(0)
2535 .with_schema_id(0)
2536 .with_manifest_list("/snap-0")
2537 .with_parent_snapshot_id(Some(1))
2538 .with_summary(Summary {
2539 operation: Operation::Append,
2540 additional_properties: HashMap::new(),
2541 })
2542 .build();
2543
2544 let err = builder
2545 .set_branch_snapshot(snapshot, MAIN_BRANCH)
2546 .unwrap_err();
2547 assert!(
2548 err.to_string()
2549 .contains("Cannot add snapshot with sequence number")
2550 );
2551 }
2552
2553 #[test]
2554 fn test_default_spec_cannot_be_removed() {
2555 let builder = builder_without_changes(FormatVersion::V2);
2556
2557 builder.remove_partition_specs(&[0]).unwrap_err();
2558 }
2559
2560 #[test]
2561 fn test_statistics() {
2562 let builder = builder_without_changes(FormatVersion::V2);
2563
2564 let statistics = StatisticsFile {
2565 snapshot_id: 3055729675574597004,
2566 statistics_path: "s3://a/b/stats.puffin".to_string(),
2567 file_size_in_bytes: 413,
2568 file_footer_size_in_bytes: 42,
2569 key_metadata: None,
2570 blob_metadata: vec![BlobMetadata {
2571 snapshot_id: 3055729675574597004,
2572 sequence_number: 1,
2573 fields: vec![1],
2574 r#type: "ndv".to_string(),
2575 properties: HashMap::new(),
2576 }],
2577 };
2578 let build_result = builder.set_statistics(statistics.clone()).build().unwrap();
2579
2580 assert_eq!(
2581 build_result.metadata.statistics,
2582 HashMap::from_iter(vec![(3055729675574597004, statistics.clone())])
2583 );
2584 assert_eq!(build_result.changes, vec![TableUpdate::SetStatistics {
2585 statistics: statistics.clone()
2586 }]);
2587
2588 let builder = build_result.metadata.into_builder(None);
2590 let build_result = builder
2591 .remove_statistics(statistics.snapshot_id)
2592 .build()
2593 .unwrap();
2594
2595 assert_eq!(build_result.metadata.statistics.len(), 0);
2596 assert_eq!(build_result.changes, vec![TableUpdate::RemoveStatistics {
2597 snapshot_id: statistics.snapshot_id
2598 }]);
2599
2600 let builder = build_result.metadata.into_builder(None);
2602 let build_result = builder
2603 .remove_statistics(statistics.snapshot_id)
2604 .build()
2605 .unwrap();
2606 assert_eq!(build_result.metadata.statistics.len(), 0);
2607 assert_eq!(build_result.changes.len(), 0);
2608 }
2609
2610 #[test]
2611 fn test_add_partition_statistics() {
2612 let builder = builder_without_changes(FormatVersion::V2);
2613
2614 let statistics = PartitionStatisticsFile {
2615 snapshot_id: 3055729675574597004,
2616 statistics_path: "s3://a/b/partition-stats.parquet".to_string(),
2617 file_size_in_bytes: 43,
2618 };
2619
2620 let build_result = builder
2621 .set_partition_statistics(statistics.clone())
2622 .build()
2623 .unwrap();
2624 assert_eq!(
2625 build_result.metadata.partition_statistics,
2626 HashMap::from_iter(vec![(3055729675574597004, statistics.clone())])
2627 );
2628 assert_eq!(build_result.changes, vec![
2629 TableUpdate::SetPartitionStatistics {
2630 partition_statistics: statistics.clone()
2631 }
2632 ]);
2633
2634 let builder = build_result.metadata.into_builder(None);
2636 let build_result = builder
2637 .remove_partition_statistics(statistics.snapshot_id)
2638 .build()
2639 .unwrap();
2640 assert_eq!(build_result.metadata.partition_statistics.len(), 0);
2641 assert_eq!(build_result.changes, vec![
2642 TableUpdate::RemovePartitionStatistics {
2643 snapshot_id: statistics.snapshot_id
2644 }
2645 ]);
2646
2647 let builder = build_result.metadata.into_builder(None);
2649 let build_result = builder
2650 .remove_partition_statistics(statistics.snapshot_id)
2651 .build()
2652 .unwrap();
2653 assert_eq!(build_result.metadata.partition_statistics.len(), 0);
2654 assert_eq!(build_result.changes.len(), 0);
2655 }
2656
2657 #[test]
2658 fn last_update_increased_for_property_only_update() {
2659 let builder = builder_without_changes(FormatVersion::V2);
2660
2661 let metadata = builder.build().unwrap().metadata;
2662 let last_updated_ms = metadata.last_updated_ms;
2663 sleep(std::time::Duration::from_millis(2));
2664
2665 let build_result = metadata
2666 .into_builder(Some(
2667 "s3://bucket/test/location/metadata/metadata1.json".to_string(),
2668 ))
2669 .set_properties(HashMap::from_iter(vec![(
2670 "foo".to_string(),
2671 "bar".to_string(),
2672 )]))
2673 .unwrap()
2674 .build()
2675 .unwrap();
2676
2677 assert!(
2678 build_result.metadata.last_updated_ms > last_updated_ms,
2679 "{} > {}",
2680 build_result.metadata.last_updated_ms,
2681 last_updated_ms
2682 );
2683 }
2684
2685 #[test]
2686 fn test_construct_default_main_branch() {
2687 let file = File::open(format!(
2689 "{}/testdata/table_metadata/{}",
2690 env!("CARGO_MANIFEST_DIR"),
2691 "TableMetadataV2Valid.json"
2692 ))
2693 .unwrap();
2694 let reader = BufReader::new(file);
2695 let resp = serde_json::from_reader::<_, TableMetadata>(reader).unwrap();
2696
2697 let table = Table::builder()
2698 .metadata(resp)
2699 .metadata_location("s3://bucket/test/location/metadata/v1.json")
2700 .identifier(TableIdent::from_strs(["ns1", "test1"]).unwrap())
2701 .file_io(FileIO::new_with_memory())
2702 .runtime(test_runtime())
2703 .build()
2704 .unwrap();
2705
2706 assert_eq!(
2707 table.metadata().refs.get(MAIN_BRANCH).unwrap().snapshot_id,
2708 table.metadata().current_snapshot_id().unwrap()
2709 );
2710 }
2711
2712 #[test]
2713 fn test_active_schema_cannot_be_removed() {
2714 let builder = builder_without_changes(FormatVersion::V2);
2715 builder.remove_schemas(&[0]).unwrap_err();
2716 }
2717
2718 #[test]
2719 fn test_remove_schemas() {
2720 let file = File::open(format!(
2721 "{}/testdata/table_metadata/{}",
2722 env!("CARGO_MANIFEST_DIR"),
2723 "TableMetadataV2Valid.json"
2724 ))
2725 .unwrap();
2726 let reader = BufReader::new(file);
2727 let resp = serde_json::from_reader::<_, TableMetadata>(reader).unwrap();
2728
2729 let table = Table::builder()
2730 .metadata(resp)
2731 .metadata_location("s3://bucket/test/location/metadata/v1.json")
2732 .identifier(TableIdent::from_strs(["ns1", "test1"]).unwrap())
2733 .file_io(FileIO::new_with_memory())
2734 .runtime(test_runtime())
2735 .build()
2736 .unwrap();
2737
2738 assert_eq!(2, table.metadata().schemas.len());
2739
2740 {
2741 let meta_data_builder = table.metadata().clone().into_builder(None);
2743 meta_data_builder.remove_schemas(&[1]).unwrap_err();
2744 }
2745
2746 let mut meta_data_builder = table.metadata().clone().into_builder(None);
2747 meta_data_builder = meta_data_builder.remove_schemas(&[0]).unwrap();
2748 let build_result = meta_data_builder.build().unwrap();
2749 assert_eq!(1, build_result.metadata.schemas.len());
2750 assert_eq!(1, build_result.metadata.current_schema_id);
2751 assert_eq!(1, build_result.metadata.current_schema().schema_id());
2752 assert_eq!(1, build_result.changes.len());
2753
2754 let remove_schema_ids =
2755 if let TableUpdate::RemoveSchemas { schema_ids } = &build_result.changes[0] {
2756 schema_ids
2757 } else {
2758 unreachable!("Expected RemoveSchema change")
2759 };
2760 assert_eq!(remove_schema_ids, &[0]);
2761 }
2762
2763 #[test]
2764 fn test_schema_evolution_now_correctly_validates_partition_field_name_conflicts() {
2765 let initial_schema = Schema::builder()
2766 .with_fields(vec![
2767 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
2768 ])
2769 .build()
2770 .unwrap();
2771
2772 let partition_spec_with_bucket = UnboundPartitionSpec::builder()
2773 .with_spec_id(0)
2774 .add_partition_field(
2775 UnboundPartitionField::builder()
2776 .source_ids(vec![1])
2777 .name("bucket_data")
2778 .transform(Transform::Bucket(16))
2779 .build()
2780 .unwrap(),
2781 )
2782 .unwrap()
2783 .build();
2784
2785 let metadata = TableMetadataBuilder::new(
2786 initial_schema,
2787 partition_spec_with_bucket,
2788 SortOrder::unsorted_order(),
2789 TEST_LOCATION.to_string(),
2790 FormatVersion::V2,
2791 HashMap::new(),
2792 )
2793 .unwrap()
2794 .build()
2795 .unwrap()
2796 .metadata;
2797
2798 let partition_field_names: Vec<String> = metadata
2799 .default_partition_spec()
2800 .fields()
2801 .iter()
2802 .map(|f| f.name.clone())
2803 .collect();
2804 assert!(partition_field_names.contains(&"bucket_data".to_string()));
2805
2806 let evolved_schema = Schema::builder()
2807 .with_fields(vec![
2808 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
2809 NestedField::required(2, "bucket_data", Type::Primitive(PrimitiveType::Int)).into(),
2811 ])
2812 .build()
2813 .unwrap();
2814
2815 let builder = metadata.into_builder(Some(
2816 "s3://bucket/test/location/metadata/metadata1.json".to_string(),
2817 ));
2818
2819 let result = builder.add_current_schema(evolved_schema);
2821
2822 assert!(result.is_err());
2823 let error = result.unwrap_err();
2824 let error_message = error.message();
2825 assert!(error_message.contains("Cannot add schema field 'bucket_data' because it conflicts with existing partition field name"));
2826 assert!(error_message.contains("Schema evolution cannot introduce field names that match existing partition field names"));
2827 }
2828
2829 #[test]
2830 fn test_schema_evolution_should_validate_on_schema_add_not_metadata_build() {
2831 let initial_schema = Schema::builder()
2832 .with_fields(vec![
2833 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
2834 ])
2835 .build()
2836 .unwrap();
2837
2838 let partition_spec = UnboundPartitionSpec::builder()
2839 .with_spec_id(0)
2840 .add_partition_field(
2841 UnboundPartitionField::builder()
2842 .source_ids(vec![1])
2843 .name("partition_col")
2844 .transform(Transform::Bucket(16))
2845 .build()
2846 .unwrap(),
2847 )
2848 .unwrap()
2849 .build();
2850
2851 let metadata = TableMetadataBuilder::new(
2852 initial_schema,
2853 partition_spec,
2854 SortOrder::unsorted_order(),
2855 TEST_LOCATION.to_string(),
2856 FormatVersion::V2,
2857 HashMap::new(),
2858 )
2859 .unwrap()
2860 .build()
2861 .unwrap()
2862 .metadata;
2863
2864 let non_conflicting_schema = Schema::builder()
2865 .with_fields(vec![
2866 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
2867 NestedField::required(2, "new_field", Type::Primitive(PrimitiveType::Int)).into(),
2868 ])
2869 .build()
2870 .unwrap();
2871
2872 let result = metadata
2874 .clone()
2875 .into_builder(Some("test_location".to_string()))
2876 .add_current_schema(non_conflicting_schema)
2877 .unwrap()
2878 .build();
2879
2880 assert!(result.is_ok());
2881 }
2882
2883 #[test]
2884 fn test_partition_spec_evolution_validates_schema_field_name_conflicts() {
2885 let initial_schema = Schema::builder()
2886 .with_fields(vec![
2887 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
2888 NestedField::required(2, "existing_field", Type::Primitive(PrimitiveType::Int))
2889 .into(),
2890 ])
2891 .build()
2892 .unwrap();
2893
2894 let partition_spec = UnboundPartitionSpec::builder()
2895 .with_spec_id(0)
2896 .add_partition_field(
2897 UnboundPartitionField::builder()
2898 .source_ids(vec![1])
2899 .name("data_bucket")
2900 .transform(Transform::Bucket(16))
2901 .build()
2902 .unwrap(),
2903 )
2904 .unwrap()
2905 .build();
2906
2907 let metadata = TableMetadataBuilder::new(
2908 initial_schema,
2909 partition_spec,
2910 SortOrder::unsorted_order(),
2911 TEST_LOCATION.to_string(),
2912 FormatVersion::V2,
2913 HashMap::new(),
2914 )
2915 .unwrap()
2916 .build()
2917 .unwrap()
2918 .metadata;
2919
2920 let builder = metadata.into_builder(Some(
2921 "s3://bucket/test/location/metadata/metadata1.json".to_string(),
2922 ));
2923
2924 let conflicting_partition_spec = UnboundPartitionSpec::builder()
2925 .with_spec_id(1)
2926 .add_partition_field(
2927 UnboundPartitionField::builder()
2928 .source_ids(vec![1])
2929 .name("existing_field")
2930 .transform(Transform::Bucket(8))
2931 .build()
2932 .unwrap(),
2933 )
2934 .unwrap()
2935 .build();
2936
2937 let result = builder.add_partition_spec(conflicting_partition_spec);
2938
2939 assert!(result.is_err());
2940 let error = result.unwrap_err();
2941 let error_message = error.message();
2942 assert!(error_message.contains(
2944 "Cannot create partition with name 'existing_field' that conflicts with schema field"
2945 ));
2946 assert!(error_message.contains("and is not an identity transform"));
2947 }
2948
2949 #[test]
2950 fn test_schema_evolution_validates_against_all_historical_schemas() {
2951 let initial_schema = Schema::builder()
2953 .with_fields(vec![
2954 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
2955 NestedField::required(2, "existing_field", Type::Primitive(PrimitiveType::Int))
2956 .into(),
2957 ])
2958 .build()
2959 .unwrap();
2960
2961 let partition_spec = UnboundPartitionSpec::builder()
2962 .with_spec_id(0)
2963 .add_partition_field(
2964 UnboundPartitionField::builder()
2965 .source_ids(vec![1])
2966 .name("bucket_data")
2967 .transform(Transform::Bucket(16))
2968 .build()
2969 .unwrap(),
2970 )
2971 .unwrap()
2972 .build();
2973
2974 let metadata = TableMetadataBuilder::new(
2975 initial_schema,
2976 partition_spec,
2977 SortOrder::unsorted_order(),
2978 TEST_LOCATION.to_string(),
2979 FormatVersion::V2,
2980 HashMap::new(),
2981 )
2982 .unwrap()
2983 .build()
2984 .unwrap()
2985 .metadata;
2986
2987 let second_schema = Schema::builder()
2989 .with_fields(vec![
2990 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
2991 NestedField::required(3, "new_field", Type::Primitive(PrimitiveType::String))
2992 .into(),
2993 ])
2994 .build()
2995 .unwrap();
2996
2997 let metadata = metadata
2998 .into_builder(Some("test_location".to_string()))
2999 .add_current_schema(second_schema)
3000 .unwrap()
3001 .build()
3002 .unwrap()
3003 .metadata;
3004
3005 let third_schema = Schema::builder()
3009 .with_fields(vec![
3010 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
3011 NestedField::required(3, "new_field", Type::Primitive(PrimitiveType::String))
3012 .into(),
3013 NestedField::required(4, "existing_field", Type::Primitive(PrimitiveType::Int))
3014 .into(),
3015 ])
3016 .build()
3017 .unwrap();
3018
3019 let builder = metadata
3020 .clone()
3021 .into_builder(Some("test_location".to_string()));
3022
3023 let result = builder.add_current_schema(third_schema);
3025 assert!(result.is_ok());
3026
3027 let conflicting_schema = Schema::builder()
3030 .with_fields(vec![
3031 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
3032 NestedField::required(3, "new_field", Type::Primitive(PrimitiveType::String))
3033 .into(),
3034 NestedField::required(4, "existing_field", Type::Primitive(PrimitiveType::Int))
3035 .into(),
3036 NestedField::required(5, "bucket_data", Type::Primitive(PrimitiveType::String))
3037 .into(), ])
3039 .build()
3040 .unwrap();
3041
3042 let builder2 = metadata.into_builder(Some("test_location".to_string()));
3043 let result2 = builder2.add_current_schema(conflicting_schema);
3044
3045 assert!(result2.is_err());
3048 let error = result2.unwrap_err();
3049 assert!(error.message().contains("Cannot add schema field 'bucket_data' because it conflicts with existing partition field name"));
3050 }
3051
3052 #[test]
3053 fn test_schema_evolution_allows_existing_partition_field_if_exists_in_historical_schema() {
3054 let initial_schema = Schema::builder()
3056 .with_fields(vec![
3057 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
3058 NestedField::required(2, "partition_data", Type::Primitive(PrimitiveType::Int))
3059 .into(),
3060 ])
3061 .build()
3062 .unwrap();
3063
3064 let partition_spec = UnboundPartitionSpec::builder()
3065 .with_spec_id(0)
3066 .add_partition_field(
3067 UnboundPartitionField::builder()
3068 .source_ids(vec![2])
3069 .name("partition_data")
3070 .transform(Transform::Identity)
3071 .build()
3072 .unwrap(),
3073 )
3074 .unwrap()
3075 .build();
3076
3077 let metadata = TableMetadataBuilder::new(
3078 initial_schema,
3079 partition_spec,
3080 SortOrder::unsorted_order(),
3081 TEST_LOCATION.to_string(),
3082 FormatVersion::V2,
3083 HashMap::new(),
3084 )
3085 .unwrap()
3086 .build()
3087 .unwrap()
3088 .metadata;
3089
3090 let evolved_schema = Schema::builder()
3092 .with_fields(vec![
3093 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
3094 NestedField::required(2, "partition_data", Type::Primitive(PrimitiveType::Int))
3095 .into(),
3096 NestedField::required(3, "new_field", Type::Primitive(PrimitiveType::String))
3097 .into(),
3098 ])
3099 .build()
3100 .unwrap();
3101
3102 let result = metadata
3104 .into_builder(Some("test_location".to_string()))
3105 .add_current_schema(evolved_schema);
3106
3107 assert!(result.is_ok());
3108 }
3109
3110 #[test]
3111 fn test_schema_evolution_prevents_new_field_conflicting_with_partition_field() {
3112 let initial_schema = Schema::builder()
3114 .with_fields(vec![
3115 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
3116 ])
3117 .build()
3118 .unwrap();
3119
3120 let partition_spec = UnboundPartitionSpec::builder()
3121 .with_spec_id(0)
3122 .add_partition_field(
3123 UnboundPartitionField::builder()
3124 .source_ids(vec![1])
3125 .name("bucket_data")
3126 .transform(Transform::Bucket(16))
3127 .build()
3128 .unwrap(),
3129 )
3130 .unwrap()
3131 .build();
3132
3133 let metadata = TableMetadataBuilder::new(
3134 initial_schema,
3135 partition_spec,
3136 SortOrder::unsorted_order(),
3137 TEST_LOCATION.to_string(),
3138 FormatVersion::V2,
3139 HashMap::new(),
3140 )
3141 .unwrap()
3142 .build()
3143 .unwrap()
3144 .metadata;
3145
3146 let conflicting_schema = Schema::builder()
3148 .with_fields(vec![
3149 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
3150 NestedField::required(2, "bucket_data", Type::Primitive(PrimitiveType::Int)).into(),
3152 ])
3153 .build()
3154 .unwrap();
3155
3156 let builder = metadata.into_builder(Some("test_location".to_string()));
3157 let result = builder.add_current_schema(conflicting_schema);
3158
3159 assert!(result.is_err());
3162 let error = result.unwrap_err();
3163 assert!(error.message().contains("Cannot add schema field 'bucket_data' because it conflicts with existing partition field name"));
3164 }
3165
3166 #[test]
3167 fn test_partition_spec_evolution_allows_non_conflicting_names() {
3168 let initial_schema = Schema::builder()
3169 .with_fields(vec![
3170 NestedField::required(1, "data", Type::Primitive(PrimitiveType::String)).into(),
3171 NestedField::required(2, "existing_field", Type::Primitive(PrimitiveType::Int))
3172 .into(),
3173 ])
3174 .build()
3175 .unwrap();
3176
3177 let partition_spec = UnboundPartitionSpec::builder()
3178 .with_spec_id(0)
3179 .add_partition_field(
3180 UnboundPartitionField::builder()
3181 .source_ids(vec![1])
3182 .name("data_bucket")
3183 .transform(Transform::Bucket(16))
3184 .build()
3185 .unwrap(),
3186 )
3187 .unwrap()
3188 .build();
3189
3190 let metadata = TableMetadataBuilder::new(
3191 initial_schema,
3192 partition_spec,
3193 SortOrder::unsorted_order(),
3194 TEST_LOCATION.to_string(),
3195 FormatVersion::V2,
3196 HashMap::new(),
3197 )
3198 .unwrap()
3199 .build()
3200 .unwrap()
3201 .metadata;
3202
3203 let builder = metadata.into_builder(Some(
3204 "s3://bucket/test/location/metadata/metadata1.json".to_string(),
3205 ));
3206
3207 let non_conflicting_partition_spec = UnboundPartitionSpec::builder()
3209 .with_spec_id(1)
3210 .add_partition_field(
3211 UnboundPartitionField::builder()
3212 .source_ids(vec![2])
3213 .name("new_partition_field")
3214 .transform(Transform::Bucket(8))
3215 .build()
3216 .unwrap(),
3217 )
3218 .unwrap()
3219 .build();
3220
3221 let result = builder.add_partition_spec(non_conflicting_partition_spec);
3222
3223 assert!(result.is_ok());
3224 }
3225
3226 #[test]
3227 fn test_row_lineage_addition() {
3228 let new_rows = 30;
3229 let base = builder_without_changes(FormatVersion::V3)
3230 .build()
3231 .unwrap()
3232 .metadata;
3233 let add_rows = Snapshot::builder()
3234 .with_snapshot_id(0)
3235 .with_timestamp_ms(base.last_updated_ms + 1)
3236 .with_sequence_number(0)
3237 .with_schema_id(0)
3238 .with_manifest_list("foo")
3239 .with_parent_snapshot_id(None)
3240 .with_summary(Summary {
3241 operation: Operation::Append,
3242 additional_properties: HashMap::new(),
3243 })
3244 .with_row_range(base.next_row_id(), new_rows)
3245 .build();
3246
3247 let first_addition = base
3248 .into_builder(None)
3249 .add_snapshot(add_rows.clone())
3250 .unwrap()
3251 .build()
3252 .unwrap()
3253 .metadata;
3254
3255 assert_eq!(first_addition.next_row_id(), new_rows);
3256
3257 let add_more_rows = Snapshot::builder()
3258 .with_snapshot_id(1)
3259 .with_timestamp_ms(first_addition.last_updated_ms + 1)
3260 .with_sequence_number(1)
3261 .with_schema_id(0)
3262 .with_manifest_list("foo")
3263 .with_parent_snapshot_id(Some(0))
3264 .with_summary(Summary {
3265 operation: Operation::Append,
3266 additional_properties: HashMap::new(),
3267 })
3268 .with_row_range(first_addition.next_row_id(), new_rows)
3269 .build();
3270
3271 let second_addition = first_addition
3272 .into_builder(None)
3273 .add_snapshot(add_more_rows)
3274 .unwrap()
3275 .build()
3276 .unwrap()
3277 .metadata;
3278 assert_eq!(second_addition.next_row_id(), new_rows * 2);
3279 }
3280
3281 #[test]
3282 fn test_row_lineage_invalid_snapshot() {
3283 let new_rows = 30;
3284 let base = builder_without_changes(FormatVersion::V3)
3285 .build()
3286 .unwrap()
3287 .metadata;
3288
3289 let add_rows = Snapshot::builder()
3291 .with_snapshot_id(0)
3292 .with_timestamp_ms(base.last_updated_ms + 1)
3293 .with_sequence_number(0)
3294 .with_schema_id(0)
3295 .with_manifest_list("foo")
3296 .with_parent_snapshot_id(None)
3297 .with_summary(Summary {
3298 operation: Operation::Append,
3299 additional_properties: HashMap::new(),
3300 })
3301 .with_row_range(base.next_row_id(), new_rows)
3302 .build();
3303
3304 let added = base
3305 .into_builder(None)
3306 .add_snapshot(add_rows)
3307 .unwrap()
3308 .build()
3309 .unwrap()
3310 .metadata;
3311
3312 let invalid_new_rows = Snapshot::builder()
3313 .with_snapshot_id(1)
3314 .with_timestamp_ms(added.last_updated_ms + 1)
3315 .with_sequence_number(1)
3316 .with_schema_id(0)
3317 .with_manifest_list("foo")
3318 .with_parent_snapshot_id(Some(0))
3319 .with_summary(Summary {
3320 operation: Operation::Append,
3321 additional_properties: HashMap::new(),
3322 })
3323 .with_row_range(added.next_row_id() - 1, 10)
3325 .build();
3326
3327 let err = added
3328 .into_builder(None)
3329 .add_snapshot(invalid_new_rows)
3330 .unwrap_err();
3331 assert!(
3332 err.to_string().contains(
3333 "Cannot add a snapshot, first-row-id is behind table next-row-id: 29 < 30"
3334 )
3335 );
3336 }
3337
3338 #[test]
3339 fn test_row_lineage_append_branch() {
3340 let branch = "some_branch";
3344
3345 let base = builder_without_changes(FormatVersion::V3)
3347 .build()
3348 .unwrap()
3349 .metadata;
3350
3351 assert_eq!(base.next_row_id(), 0);
3353
3354 let branch_snapshot_1 = Snapshot::builder()
3356 .with_snapshot_id(1)
3357 .with_timestamp_ms(base.last_updated_ms + 1)
3358 .with_sequence_number(0)
3359 .with_schema_id(0)
3360 .with_manifest_list("foo")
3361 .with_parent_snapshot_id(None)
3362 .with_summary(Summary {
3363 operation: Operation::Append,
3364 additional_properties: HashMap::new(),
3365 })
3366 .with_row_range(base.next_row_id(), 30)
3367 .build();
3368
3369 let table_after_branch_1 = base
3370 .into_builder(None)
3371 .set_branch_snapshot(branch_snapshot_1.clone(), branch)
3372 .unwrap()
3373 .build()
3374 .unwrap()
3375 .metadata;
3376
3377 assert!(table_after_branch_1.current_snapshot().is_none());
3379
3380 let branch_ref = table_after_branch_1.refs.get(branch).unwrap();
3382 let branch_snap_1 = table_after_branch_1
3383 .snapshots
3384 .get(&branch_ref.snapshot_id)
3385 .unwrap();
3386 assert_eq!(branch_snap_1.first_row_id(), Some(0));
3387
3388 assert_eq!(table_after_branch_1.next_row_id(), 30);
3390
3391 let main_snapshot = Snapshot::builder()
3393 .with_snapshot_id(2)
3394 .with_timestamp_ms(table_after_branch_1.last_updated_ms + 1)
3395 .with_sequence_number(1)
3396 .with_schema_id(0)
3397 .with_manifest_list("bar")
3398 .with_parent_snapshot_id(None)
3399 .with_summary(Summary {
3400 operation: Operation::Append,
3401 additional_properties: HashMap::new(),
3402 })
3403 .with_row_range(table_after_branch_1.next_row_id(), 28)
3404 .build();
3405
3406 let table_after_main = table_after_branch_1
3407 .into_builder(None)
3408 .add_snapshot(main_snapshot.clone())
3409 .unwrap()
3410 .set_ref(MAIN_BRANCH, SnapshotReference {
3411 snapshot_id: main_snapshot.snapshot_id(),
3412 retention: SnapshotRetention::Branch {
3413 min_snapshots_to_keep: None,
3414 max_snapshot_age_ms: None,
3415 max_ref_age_ms: None,
3416 },
3417 })
3418 .unwrap()
3419 .build()
3420 .unwrap()
3421 .metadata;
3422
3423 let current_snapshot = table_after_main.current_snapshot().unwrap();
3425 assert_eq!(current_snapshot.first_row_id(), Some(30));
3426
3427 assert_eq!(table_after_main.next_row_id(), 58);
3429
3430 let branch_snapshot_2 = Snapshot::builder()
3432 .with_snapshot_id(3)
3433 .with_timestamp_ms(table_after_main.last_updated_ms + 1)
3434 .with_sequence_number(2)
3435 .with_schema_id(0)
3436 .with_manifest_list("baz")
3437 .with_parent_snapshot_id(Some(branch_snapshot_1.snapshot_id()))
3438 .with_summary(Summary {
3439 operation: Operation::Append,
3440 additional_properties: HashMap::new(),
3441 })
3442 .with_row_range(table_after_main.next_row_id(), 21)
3443 .build();
3444
3445 let table_after_branch_2 = table_after_main
3446 .into_builder(None)
3447 .set_branch_snapshot(branch_snapshot_2.clone(), branch)
3448 .unwrap()
3449 .build()
3450 .unwrap()
3451 .metadata;
3452
3453 let branch_ref_2 = table_after_branch_2.refs.get(branch).unwrap();
3455 let branch_snap_2 = table_after_branch_2
3456 .snapshots
3457 .get(&branch_ref_2.snapshot_id)
3458 .unwrap();
3459 assert_eq!(branch_snap_2.first_row_id(), Some(58));
3460
3461 assert_eq!(table_after_branch_2.next_row_id(), 79);
3463 }
3464
3465 #[test]
3466 fn test_encryption_keys() {
3467 let builder = builder_without_changes(FormatVersion::V2);
3468
3469 let encryption_key_1 = EncryptedKey::builder()
3471 .key_id("key-1")
3472 .encrypted_key_metadata(vec![1, 2, 3, 4])
3473 .encrypted_by_id("encryption-service-1")
3474 .properties(HashMap::from_iter(vec![(
3475 "algorithm".to_string(),
3476 "AES-256".to_string(),
3477 )]))
3478 .build();
3479
3480 let encryption_key_2 = EncryptedKey::builder()
3481 .key_id("key-2")
3482 .encrypted_key_metadata(vec![5, 6, 7, 8])
3483 .encrypted_by_id("encryption-service-2")
3484 .properties(HashMap::new())
3485 .build();
3486
3487 let build_result = builder
3489 .add_encryption_key(encryption_key_1.clone())
3490 .build()
3491 .unwrap();
3492
3493 assert_eq!(build_result.changes.len(), 1);
3494 assert_eq!(build_result.metadata.encryption_keys.len(), 1);
3495 assert_eq!(
3496 build_result.metadata.encryption_key("key-1"),
3497 Some(&encryption_key_1)
3498 );
3499 assert_eq!(build_result.changes[0], TableUpdate::AddEncryptionKey {
3500 encryption_key: encryption_key_1.clone()
3501 });
3502
3503 let build_result = build_result
3505 .metadata
3506 .into_builder(Some(
3507 "s3://bucket/test/location/metadata/metadata1.json".to_string(),
3508 ))
3509 .add_encryption_key(encryption_key_2.clone())
3510 .build()
3511 .unwrap();
3512
3513 assert_eq!(build_result.changes.len(), 1);
3514 assert_eq!(build_result.metadata.encryption_keys.len(), 2);
3515 assert_eq!(
3516 build_result.metadata.encryption_key("key-1"),
3517 Some(&encryption_key_1)
3518 );
3519 assert_eq!(
3520 build_result.metadata.encryption_key("key-2"),
3521 Some(&encryption_key_2)
3522 );
3523 assert_eq!(build_result.changes[0], TableUpdate::AddEncryptionKey {
3524 encryption_key: encryption_key_2.clone()
3525 });
3526
3527 let build_result = build_result
3529 .metadata
3530 .into_builder(Some(
3531 "s3://bucket/test/location/metadata/metadata2.json".to_string(),
3532 ))
3533 .add_encryption_key(encryption_key_1.clone())
3534 .build()
3535 .unwrap();
3536
3537 assert_eq!(build_result.changes.len(), 0);
3538 assert_eq!(build_result.metadata.encryption_keys.len(), 2);
3539
3540 let build_result = build_result
3542 .metadata
3543 .into_builder(Some(
3544 "s3://bucket/test/location/metadata/metadata3.json".to_string(),
3545 ))
3546 .remove_encryption_key("key-1")
3547 .build()
3548 .unwrap();
3549
3550 assert_eq!(build_result.changes.len(), 1);
3551 assert_eq!(build_result.metadata.encryption_keys.len(), 1);
3552 assert_eq!(build_result.metadata.encryption_key("key-1"), None);
3553 assert_eq!(
3554 build_result.metadata.encryption_key("key-2"),
3555 Some(&encryption_key_2)
3556 );
3557 assert_eq!(build_result.changes[0], TableUpdate::RemoveEncryptionKey {
3558 key_id: "key-1".to_string()
3559 });
3560
3561 let build_result = build_result
3563 .metadata
3564 .into_builder(Some(
3565 "s3://bucket/test/location/metadata/metadata4.json".to_string(),
3566 ))
3567 .remove_encryption_key("non-existent-key")
3568 .build()
3569 .unwrap();
3570
3571 assert_eq!(build_result.changes.len(), 0);
3572 assert_eq!(build_result.metadata.encryption_keys.len(), 1);
3573
3574 let keys = build_result
3576 .metadata
3577 .encryption_keys_iter()
3578 .collect::<Vec<_>>();
3579 assert_eq!(keys.len(), 1);
3580 assert_eq!(keys[0], &encryption_key_2);
3581
3582 let build_result = build_result
3584 .metadata
3585 .into_builder(Some(
3586 "s3://bucket/test/location/metadata/metadata5.json".to_string(),
3587 ))
3588 .remove_encryption_key("key-2")
3589 .build()
3590 .unwrap();
3591
3592 assert_eq!(build_result.changes.len(), 1);
3593 assert_eq!(build_result.metadata.encryption_keys.len(), 0);
3594 assert_eq!(build_result.metadata.encryption_key("key-2"), None);
3595 assert_eq!(build_result.changes[0], TableUpdate::RemoveEncryptionKey {
3596 key_id: "key-2".to_string()
3597 });
3598
3599 let keys = build_result.metadata.encryption_keys_iter();
3601 assert_eq!(keys.len(), 0);
3602 }
3603
3604 #[test]
3605 fn test_partition_field_id_reuse_across_specs() {
3606 let schema = Schema::builder()
3607 .with_fields(vec![
3608 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
3609 NestedField::required(2, "data", Type::Primitive(PrimitiveType::String)).into(),
3610 NestedField::required(3, "timestamp", Type::Primitive(PrimitiveType::Timestamp))
3611 .into(),
3612 ])
3613 .build()
3614 .unwrap();
3615
3616 let initial_spec = UnboundPartitionSpec::builder()
3618 .add_partition_field(
3619 UnboundPartitionField::builder()
3620 .source_ids(vec![1])
3621 .name("id")
3622 .transform(Transform::Identity)
3623 .build()
3624 .unwrap(),
3625 )
3626 .unwrap()
3627 .build();
3628
3629 let mut metadata = TableMetadataBuilder::new(
3630 schema,
3631 initial_spec,
3632 SortOrder::unsorted_order(),
3633 "s3://bucket/table".to_string(),
3634 FormatVersion::V2,
3635 HashMap::new(),
3636 )
3637 .unwrap()
3638 .build()
3639 .unwrap()
3640 .metadata;
3641
3642 let spec1 = UnboundPartitionSpec::builder()
3644 .add_partition_field(
3645 UnboundPartitionField::builder()
3646 .source_ids(vec![2])
3647 .name("data_bucket")
3648 .transform(Transform::Bucket(10))
3649 .build()
3650 .unwrap(),
3651 )
3652 .unwrap()
3653 .build();
3654 let builder = metadata.into_builder(Some("s3://bucket/table/metadata/v1.json".to_string()));
3655 let result = builder.add_partition_spec(spec1).unwrap().build().unwrap();
3656 metadata = result.metadata;
3657
3658 let spec2 = UnboundPartitionSpec::builder()
3661 .add_partition_field(
3662 UnboundPartitionField::builder()
3663 .source_ids(vec![1])
3664 .name("id")
3665 .transform(Transform::Identity)
3666 .build()
3667 .unwrap(),
3668 ) .unwrap()
3670 .add_partition_field(
3671 UnboundPartitionField::builder()
3672 .source_ids(vec![2])
3673 .name("data_bucket")
3674 .transform(Transform::Bucket(10))
3675 .build()
3676 .unwrap(),
3677 ) .unwrap()
3679 .add_partition_field(
3680 UnboundPartitionField::builder()
3681 .source_ids(vec![3])
3682 .name("year")
3683 .transform(Transform::Year)
3684 .build()
3685 .unwrap(),
3686 ) .unwrap()
3688 .build();
3689 let builder = metadata.into_builder(Some("s3://bucket/table/metadata/v2.json".to_string()));
3690 let result = builder.add_partition_spec(spec2).unwrap().build().unwrap();
3691
3692 let spec2 = result.metadata.partition_spec_by_id(2).unwrap();
3694 let field_ids: Vec<i32> = spec2.fields().iter().map(|f| f.field_id).collect();
3695 assert_eq!(field_ids, vec![1000, 1001, 1002]); }
3697}