Skip to main content

iceberg/spec/
table_metadata.rs

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