Skip to main content

iceberg/catalog/
mod.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//! Catalog API for Apache Iceberg
19
20pub mod memory;
21mod metadata_location;
22mod session;
23pub(crate) mod utils;
24
25use std::collections::HashMap;
26use std::fmt::{Debug, Display};
27use std::future::Future;
28use std::mem::take;
29use std::ops::Deref;
30use std::str::FromStr;
31use std::sync::Arc;
32
33use _serde::{deserialize_snapshot, serialize_snapshot};
34use async_trait::async_trait;
35pub use memory::MemoryCatalog;
36pub use metadata_location::*;
37#[cfg(test)]
38use mockall::automock;
39use serde_derive::{Deserialize, Serialize};
40pub use session::*;
41use typed_builder::TypedBuilder;
42use uuid::Uuid;
43
44use crate::encryption::kms::KmsClientFactory;
45use crate::error::invalid_data;
46use crate::io::StorageFactory;
47use crate::runtime::Runtime;
48use crate::spec::{
49    EncryptedKey, FormatVersion, PartitionStatisticsFile, Schema, SchemaId, Snapshot,
50    SnapshotReference, SortOrder, StatisticsFile, TableMetadata, TableMetadataBuilder,
51    UnboundPartitionSpec, ViewFormatVersion, ViewRepresentations, ViewVersion,
52};
53use crate::table::Table;
54use crate::{Error, ErrorKind, Result};
55
56/// The catalog API for Iceberg Rust.
57#[async_trait]
58#[cfg_attr(test, automock)]
59pub trait Catalog: Debug + Sync + Send {
60    /// List namespaces inside the catalog.
61    async fn list_namespaces(&self, parent: Option<&NamespaceIdent>)
62    -> Result<Vec<NamespaceIdent>>;
63
64    /// Create a new namespace inside the catalog.
65    async fn create_namespace(
66        &self,
67        namespace: &NamespaceIdent,
68        properties: HashMap<String, String>,
69    ) -> Result<Namespace>;
70
71    /// Get a namespace information from the catalog.
72    async fn get_namespace(&self, namespace: &NamespaceIdent) -> Result<Namespace>;
73
74    /// Check if namespace exists in catalog.
75    async fn namespace_exists(&self, namespace: &NamespaceIdent) -> Result<bool>;
76
77    /// Update a namespace inside the catalog.
78    ///
79    /// # Behavior
80    ///
81    /// The properties must be the full set of namespace.
82    async fn update_namespace(
83        &self,
84        namespace: &NamespaceIdent,
85        properties: HashMap<String, String>,
86    ) -> Result<()>;
87
88    /// Drop a namespace from the catalog, or returns error if it doesn't exist.
89    async fn drop_namespace(&self, namespace: &NamespaceIdent) -> Result<()>;
90
91    /// List tables from namespace.
92    async fn list_tables(&self, namespace: &NamespaceIdent) -> Result<Vec<TableIdent>>;
93
94    /// Create a new table inside the namespace.
95    async fn create_table(
96        &self,
97        namespace: &NamespaceIdent,
98        creation: TableCreation,
99    ) -> Result<Table>;
100
101    /// Load table from the catalog.
102    async fn load_table(&self, table: &TableIdent) -> Result<Table>;
103
104    /// Drop a table from the catalog, or returns error if it doesn't exist.
105    async fn drop_table(&self, table: &TableIdent) -> Result<()>;
106
107    /// Drop a table from the catalog and delete the underlying table data.
108    ///
109    /// Implementations should load the table metadata, drop the table
110    /// from the catalog, then delete all associated data and metadata files.
111    /// The [`drop_table_data`](utils::drop_table_data) utility function can
112    /// be used for the file cleanup step.
113    async fn purge_table(&self, table: &TableIdent) -> Result<()>;
114
115    /// Check if a table exists in the catalog.
116    async fn table_exists(&self, table: &TableIdent) -> Result<bool>;
117
118    /// Rename a table in the catalog.
119    async fn rename_table(&self, src: &TableIdent, dest: &TableIdent) -> Result<()>;
120
121    /// Register an existing table to the catalog.
122    async fn register_table(&self, table: &TableIdent, metadata_location: String) -> Result<Table>;
123
124    /// Update a table to the catalog.
125    async fn update_table(&self, commit: TableCommit) -> Result<Table>;
126}
127
128/// Common interface for all catalog builders.
129pub trait CatalogBuilder: Default + Debug + Send + Sync {
130    /// The catalog type that this builder creates.
131    type C: Catalog;
132
133    /// Set a custom StorageFactory to use for storage operations.
134    ///
135    /// When a StorageFactory is provided, the catalog will use it to build FileIO
136    /// instances for all storage operations instead of using the default factory.
137    ///
138    /// # Arguments
139    ///
140    /// * `storage_factory` - The StorageFactory to use for creating storage instances
141    ///
142    /// # Example
143    ///
144    /// ```rust,ignore
145    /// use iceberg::CatalogBuilder;
146    /// use iceberg::io::StorageFactory;
147    /// use iceberg_storage_opendal::OpenDalStorageFactory;
148    /// use std::sync::Arc;
149    ///
150    /// let catalog = MyCatalogBuilder::default()
151    ///     .with_storage_factory(Arc::new(OpenDalStorageFactory::S3 {
152    ///         customized_credential_load: None,
153    ///     }))
154    ///     .load("my_catalog", props)
155    ///     .await?;
156    /// ```
157    fn with_storage_factory(self, storage_factory: Arc<dyn StorageFactory>) -> Self;
158
159    /// Set a [`KmsClientFactory`] to enable table encryption.
160    ///
161    /// When provided, the catalog calls the factory once during
162    /// [`load`](Self::load) with the catalog properties to create a shared
163    /// [`KeyManagementClient`](crate::encryption::KeyManagementClient).
164    /// That client is then passed to each table's `TableBuilder` so tables
165    /// with `encryption.key-id` set can construct an `EncryptionManager`.
166    ///
167    /// # Example
168    ///
169    /// ```rust,ignore
170    /// use iceberg::CatalogBuilder;
171    /// use iceberg::encryption::kms::KmsClientFactory;
172    /// use std::sync::Arc;
173    ///
174    /// let catalog = MyCatalogBuilder::default()
175    ///     .with_kms_client_factory(Arc::new(MyKmsClientFactory))
176    ///     .load("my_catalog", props)
177    ///     .await?;
178    /// ```
179    fn with_kms_client_factory(self, kms_client_factory: Arc<dyn KmsClientFactory>) -> Self;
180
181    /// Set a custom tokio Runtime to use for spawning async tasks.
182    ///
183    /// When a Runtime is provided, the catalog will propagate it to all tables
184    /// it creates. Tasks such as scan planning and delete file processing
185    /// will be spawned on this runtime.
186    fn with_runtime(self, runtime: Runtime) -> Self;
187
188    /// Create a new catalog instance.
189    fn load(
190        self,
191        name: impl Into<String>,
192        props: HashMap<String, String>,
193    ) -> impl Future<Output = Result<Self::C>> + Send;
194}
195
196/// NamespaceIdent represents the identifier of a namespace in the catalog.
197///
198/// The namespace identifier is a list of strings, where each string is a
199/// component of the namespace. It's the catalog implementer's responsibility to
200/// handle the namespace identifier correctly.
201#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
202pub struct NamespaceIdent(Vec<String>);
203
204impl NamespaceIdent {
205    /// Create a new namespace identifier with only one level.
206    pub fn new(name: String) -> Self {
207        Self(vec![name])
208    }
209
210    /// Create a multi-level namespace identifier from vector.
211    pub fn from_vec(names: Vec<String>) -> Result<Self> {
212        if names.is_empty() {
213            return Err(invalid_data!("Namespace identifier can't be empty!"));
214        }
215        Ok(Self(names))
216    }
217
218    /// Try to create namespace identifier from an iterator of string.
219    pub fn from_strs(iter: impl IntoIterator<Item = impl ToString>) -> Result<Self> {
220        Self::from_vec(iter.into_iter().map(|s| s.to_string()).collect())
221    }
222
223    /// Returns a string for used in url.
224    pub fn to_url_string(&self) -> String {
225        self.as_ref().join("\u{001f}")
226    }
227
228    /// Returns inner strings.
229    pub fn inner(self) -> Vec<String> {
230        self.0
231    }
232
233    /// Get the parent of this namespace.
234    /// Returns None if this namespace only has a single element and thus has no parent.
235    pub fn parent(&self) -> Option<Self> {
236        self.0.split_last().and_then(|(_, parent)| {
237            if parent.is_empty() {
238                None
239            } else {
240                Some(Self(parent.to_vec()))
241            }
242        })
243    }
244}
245
246impl AsRef<Vec<String>> for NamespaceIdent {
247    fn as_ref(&self) -> &Vec<String> {
248        &self.0
249    }
250}
251
252impl Deref for NamespaceIdent {
253    type Target = [String];
254
255    fn deref(&self) -> &Self::Target {
256        &self.0
257    }
258}
259
260/// Namespace represents a namespace in the catalog.
261#[derive(Debug, Clone, PartialEq, Eq)]
262pub struct Namespace {
263    name: NamespaceIdent,
264    properties: HashMap<String, String>,
265}
266
267impl Namespace {
268    /// Create a new namespace.
269    pub fn new(name: NamespaceIdent) -> Self {
270        Self::with_properties(name, HashMap::default())
271    }
272
273    /// Create a new namespace with properties.
274    pub fn with_properties(name: NamespaceIdent, properties: HashMap<String, String>) -> Self {
275        Self { name, properties }
276    }
277
278    /// Get the name of the namespace.
279    pub fn name(&self) -> &NamespaceIdent {
280        &self.name
281    }
282
283    /// Get the properties of the namespace.
284    pub fn properties(&self) -> &HashMap<String, String> {
285        &self.properties
286    }
287}
288
289impl Display for NamespaceIdent {
290    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
291        write!(f, "{}", self.0.join("."))
292    }
293}
294
295/// TableIdent represents the identifier of a table in the catalog.
296#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
297pub struct TableIdent {
298    /// Namespace of the table.
299    pub namespace: NamespaceIdent,
300    /// Table name.
301    pub name: String,
302}
303
304impl TableIdent {
305    /// Create a new table identifier.
306    pub fn new(namespace: NamespaceIdent, name: String) -> Self {
307        Self { namespace, name }
308    }
309
310    /// Get the namespace of the table.
311    pub fn namespace(&self) -> &NamespaceIdent {
312        &self.namespace
313    }
314
315    /// Get the name of the table.
316    pub fn name(&self) -> &str {
317        &self.name
318    }
319
320    /// Try to create table identifier from an iterator of string.
321    pub fn from_strs(iter: impl IntoIterator<Item = impl ToString>) -> Result<Self> {
322        let mut vec: Vec<String> = iter.into_iter().map(|s| s.to_string()).collect();
323        let table_name = vec
324            .pop()
325            .ok_or_else(|| invalid_data!("Table identifier can't be empty!"))?;
326        let namespace_ident = NamespaceIdent::from_vec(vec)?;
327
328        Ok(Self {
329            namespace: namespace_ident,
330            name: table_name,
331        })
332    }
333}
334
335impl Display for TableIdent {
336    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
337        write!(f, "{}.{}", self.namespace, self.name)
338    }
339}
340
341/// TableCreation represents the creation of a table in the catalog.
342#[derive(Debug, TypedBuilder)]
343pub struct TableCreation {
344    /// The name of the table.
345    pub name: String,
346    /// The location of the table.
347    #[builder(default, setter(strip_option(fallback = location_opt)))]
348    pub location: Option<String>,
349    /// The schema of the table.
350    pub schema: Schema,
351    /// The partition spec of the table, could be None.
352    #[builder(default, setter(strip_option(fallback = partition_spec_opt), into))]
353    pub partition_spec: Option<UnboundPartitionSpec>,
354    /// The sort order of the table.
355    #[builder(default, setter(strip_option(fallback = sort_order_opt)))]
356    pub sort_order: Option<SortOrder>,
357    /// The properties of the table.
358    #[builder(default, setter(transform = |props: impl IntoIterator<Item=(String, String)>| {
359        props.into_iter().collect()
360    }))]
361    pub properties: HashMap<String, String>,
362    /// Format version of the table. Defaults to V2.
363    #[builder(default = FormatVersion::V2)]
364    pub format_version: FormatVersion,
365}
366
367/// TableCommit represents the commit of a table in the catalog.
368///
369/// The builder is marked as private since it's dangerous and error-prone to construct
370/// [`TableCommit`] directly.
371/// Users are supposed to use [`crate::transaction::Transaction`] to update table.
372#[derive(Debug, TypedBuilder)]
373#[builder(build_method(vis = "pub(crate)"))]
374pub struct TableCommit {
375    /// The table ident.
376    ident: TableIdent,
377    /// The requirements of the table.
378    ///
379    /// Commit will fail if the requirements are not met.
380    requirements: Vec<TableRequirement>,
381    /// The updates of the table.
382    updates: Vec<TableUpdate>,
383}
384
385impl TableCommit {
386    /// Return the table identifier.
387    pub fn identifier(&self) -> &TableIdent {
388        &self.ident
389    }
390
391    /// Take all requirements.
392    pub fn take_requirements(&mut self) -> Vec<TableRequirement> {
393        take(&mut self.requirements)
394    }
395
396    /// Take all updates.
397    pub fn take_updates(&mut self) -> Vec<TableUpdate> {
398        take(&mut self.updates)
399    }
400
401    /// Applies this [`TableCommit`] to the given [`Table`] as part of a catalog update.
402    /// Typically used by [`Catalog::update_table`] to validate requirements and apply metadata updates.
403    ///
404    /// Returns a new [`Table`] with updated metadata,
405    /// or an error if validation or application fails.
406    pub fn apply(self, table: Table) -> Result<Table> {
407        // check requirements
408        for requirement in self.requirements {
409            requirement.check(Some(table.metadata()))?;
410        }
411
412        // get current metadata location
413        let current_metadata_location = table.metadata_location_result()?;
414
415        // apply updates to metadata builder
416        let mut metadata_builder = table
417            .metadata()
418            .clone()
419            .into_builder(Some(current_metadata_location.to_string()));
420        for update in self.updates {
421            metadata_builder = update.apply(metadata_builder)?;
422        }
423
424        // Build the new metadata
425        let new_metadata = metadata_builder.build()?.metadata;
426
427        let new_metadata_location = MetadataLocation::from_str(current_metadata_location)?
428            .with_next_version()
429            .try_with_new_metadata(&new_metadata)?
430            .to_string();
431
432        Ok(table
433            .with_metadata(Arc::new(new_metadata))
434            .with_metadata_location(new_metadata_location))
435    }
436}
437
438/// TableRequirement represents a requirement for a table in the catalog.
439#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
440#[serde(tag = "type")]
441pub enum TableRequirement {
442    /// The table must not already exist; used for create transactions
443    #[serde(rename = "assert-create")]
444    NotExist,
445    /// The table UUID must match the requirement.
446    #[serde(rename = "assert-table-uuid")]
447    UuidMatch {
448        /// Uuid of original table.
449        uuid: Uuid,
450    },
451    /// The table branch or tag identified by the requirement's `reference` must
452    /// reference the requirement's `snapshot-id`.
453    #[serde(rename = "assert-ref-snapshot-id")]
454    RefSnapshotIdMatch {
455        /// The reference of the table to assert.
456        r#ref: String,
457        /// The snapshot id of the table to assert.
458        /// If the id is `None`, the ref must not already exist.
459        #[serde(rename = "snapshot-id")]
460        snapshot_id: Option<i64>,
461    },
462    /// The table's last assigned column id must match the requirement.
463    #[serde(rename = "assert-last-assigned-field-id")]
464    LastAssignedFieldIdMatch {
465        /// The last assigned field id of the table to assert.
466        #[serde(rename = "last-assigned-field-id")]
467        last_assigned_field_id: i32,
468    },
469    /// The table's current schema id must match the requirement.
470    #[serde(rename = "assert-current-schema-id")]
471    CurrentSchemaIdMatch {
472        /// Current schema id of the table to assert.
473        #[serde(rename = "current-schema-id")]
474        current_schema_id: SchemaId,
475    },
476    /// The table's last assigned partition id must match the
477    /// requirement.
478    #[serde(rename = "assert-last-assigned-partition-id")]
479    LastAssignedPartitionIdMatch {
480        /// Last assigned partition id of the table to assert.
481        #[serde(rename = "last-assigned-partition-id")]
482        last_assigned_partition_id: i32,
483    },
484    /// The table's default spec id must match the requirement.
485    #[serde(rename = "assert-default-spec-id")]
486    DefaultSpecIdMatch {
487        /// Default spec id of the table to assert.
488        #[serde(rename = "default-spec-id")]
489        default_spec_id: i32,
490    },
491    /// The table's default sort order id must match the requirement.
492    #[serde(rename = "assert-default-sort-order-id")]
493    DefaultSortOrderIdMatch {
494        /// Default sort order id of the table to assert.
495        #[serde(rename = "default-sort-order-id")]
496        default_sort_order_id: i64,
497    },
498}
499
500/// TableUpdate represents an update to a table in the catalog.
501#[derive(Debug, Serialize, Deserialize, PartialEq, Clone)]
502#[serde(tag = "action", rename_all = "kebab-case")]
503#[allow(clippy::large_enum_variant)]
504pub enum TableUpdate {
505    /// Upgrade table's format version
506    #[serde(rename_all = "kebab-case")]
507    UpgradeFormatVersion {
508        /// Target format upgrade to.
509        format_version: FormatVersion,
510    },
511    /// Assign a new UUID to the table
512    #[serde(rename_all = "kebab-case")]
513    AssignUuid {
514        /// The new UUID to assign.
515        uuid: Uuid,
516    },
517    /// Add a new schema to the table
518    #[serde(rename_all = "kebab-case")]
519    AddSchema {
520        /// The schema to add.
521        schema: Schema,
522    },
523    /// Set table's current schema
524    #[serde(rename_all = "kebab-case")]
525    SetCurrentSchema {
526        /// Schema ID to set as current, or -1 to set last added schema
527        schema_id: i32,
528    },
529    /// Add a new partition spec to the table
530    AddSpec {
531        /// The partition spec to add.
532        spec: UnboundPartitionSpec,
533    },
534    /// Set table's default spec
535    #[serde(rename_all = "kebab-case")]
536    SetDefaultSpec {
537        /// Partition spec ID to set as the default, or -1 to set last added spec
538        spec_id: i32,
539    },
540    /// Add sort order to table.
541    #[serde(rename_all = "kebab-case")]
542    AddSortOrder {
543        /// Sort order to add.
544        sort_order: SortOrder,
545    },
546    /// Set table's default sort order
547    #[serde(rename_all = "kebab-case")]
548    SetDefaultSortOrder {
549        /// Sort order ID to set as the default, or -1 to set last added sort order
550        sort_order_id: i64,
551    },
552    /// Add snapshot to table.
553    #[serde(rename_all = "kebab-case")]
554    AddSnapshot {
555        /// Snapshot to add.
556        #[serde(
557            deserialize_with = "deserialize_snapshot",
558            serialize_with = "serialize_snapshot"
559        )]
560        snapshot: Snapshot,
561    },
562    /// Set table's snapshot ref.
563    #[serde(rename_all = "kebab-case")]
564    SetSnapshotRef {
565        /// Name of snapshot reference to set.
566        ref_name: String,
567        /// Snapshot reference to set.
568        #[serde(flatten)]
569        reference: SnapshotReference,
570    },
571    /// Remove table's snapshots
572    #[serde(rename_all = "kebab-case")]
573    RemoveSnapshots {
574        /// Snapshot ids to remove.
575        snapshot_ids: Vec<i64>,
576    },
577    /// Remove snapshot reference
578    #[serde(rename_all = "kebab-case")]
579    RemoveSnapshotRef {
580        /// Name of snapshot reference to remove.
581        ref_name: String,
582    },
583    /// Update table's location
584    SetLocation {
585        /// New location for table.
586        location: String,
587    },
588    /// Update table's properties
589    SetProperties {
590        /// Properties to update for table.
591        updates: HashMap<String, String>,
592    },
593    /// Remove table's properties
594    RemoveProperties {
595        /// Properties to remove
596        removals: Vec<String>,
597    },
598    /// Remove partition specs
599    #[serde(rename_all = "kebab-case")]
600    RemovePartitionSpecs {
601        /// Partition spec ids to remove.
602        spec_ids: Vec<i32>,
603    },
604    /// Set statistics for a snapshot
605    #[serde(with = "_serde_set_statistics")]
606    SetStatistics {
607        /// File containing the statistics
608        statistics: StatisticsFile,
609    },
610    /// Remove statistics for a snapshot
611    #[serde(rename_all = "kebab-case")]
612    RemoveStatistics {
613        /// Snapshot id to remove statistics for.
614        snapshot_id: i64,
615    },
616    /// Set partition statistics for a snapshot
617    #[serde(rename_all = "kebab-case")]
618    SetPartitionStatistics {
619        /// File containing the partition statistics
620        partition_statistics: PartitionStatisticsFile,
621    },
622    /// Remove partition statistics for a snapshot
623    #[serde(rename_all = "kebab-case")]
624    RemovePartitionStatistics {
625        /// Snapshot id to remove partition statistics for.
626        snapshot_id: i64,
627    },
628    /// Remove schemas
629    #[serde(rename_all = "kebab-case")]
630    RemoveSchemas {
631        /// Schema IDs to remove.
632        schema_ids: Vec<i32>,
633    },
634    /// Add an encryption key
635    #[serde(rename_all = "kebab-case")]
636    AddEncryptionKey {
637        /// The encryption key to add.
638        encryption_key: EncryptedKey,
639    },
640    /// Remove an encryption key
641    #[serde(rename_all = "kebab-case")]
642    RemoveEncryptionKey {
643        /// The id of the encryption key to remove.
644        key_id: String,
645    },
646}
647
648impl TableUpdate {
649    /// Applies the update to the table metadata builder.
650    pub fn apply(self, builder: TableMetadataBuilder) -> Result<TableMetadataBuilder> {
651        match self {
652            TableUpdate::AssignUuid { uuid } => Ok(builder.assign_uuid(uuid)),
653            TableUpdate::AddSchema { schema, .. } => Ok(builder.add_schema(schema)?),
654            TableUpdate::SetCurrentSchema { schema_id } => builder.set_current_schema(schema_id),
655            TableUpdate::AddSpec { spec } => builder.add_partition_spec(spec),
656            TableUpdate::SetDefaultSpec { spec_id } => builder.set_default_partition_spec(spec_id),
657            TableUpdate::AddSortOrder { sort_order } => builder.add_sort_order(sort_order),
658            TableUpdate::SetDefaultSortOrder { sort_order_id } => {
659                builder.set_default_sort_order(sort_order_id)
660            }
661            TableUpdate::AddSnapshot { snapshot } => builder.add_snapshot(snapshot),
662            TableUpdate::SetSnapshotRef {
663                ref_name,
664                reference,
665            } => builder.set_ref(&ref_name, reference),
666            TableUpdate::RemoveSnapshots { snapshot_ids } => {
667                Ok(builder.remove_snapshots(&snapshot_ids))
668            }
669            TableUpdate::RemoveSnapshotRef { ref_name } => Ok(builder.remove_ref(&ref_name)),
670            TableUpdate::SetLocation { location } => Ok(builder.set_location(location)),
671            TableUpdate::SetProperties { updates } => builder.set_properties(updates),
672            TableUpdate::RemoveProperties { removals } => builder.remove_properties(&removals),
673            TableUpdate::UpgradeFormatVersion { format_version } => {
674                builder.upgrade_format_version(format_version)
675            }
676            TableUpdate::RemovePartitionSpecs { spec_ids } => {
677                builder.remove_partition_specs(&spec_ids)
678            }
679            TableUpdate::SetStatistics { statistics } => Ok(builder.set_statistics(statistics)),
680            TableUpdate::RemoveStatistics { snapshot_id } => {
681                Ok(builder.remove_statistics(snapshot_id))
682            }
683            TableUpdate::SetPartitionStatistics {
684                partition_statistics,
685            } => Ok(builder.set_partition_statistics(partition_statistics)),
686            TableUpdate::RemovePartitionStatistics { snapshot_id } => {
687                Ok(builder.remove_partition_statistics(snapshot_id))
688            }
689            TableUpdate::RemoveSchemas { schema_ids } => builder.remove_schemas(&schema_ids),
690            TableUpdate::AddEncryptionKey { encryption_key } => {
691                Ok(builder.add_encryption_key(encryption_key))
692            }
693            TableUpdate::RemoveEncryptionKey { key_id } => {
694                Ok(builder.remove_encryption_key(&key_id))
695            }
696        }
697    }
698}
699
700impl TableRequirement {
701    /// Check that the requirement is met by the table metadata.
702    /// If the requirement is not met, an appropriate error is returned.
703    ///
704    /// Provide metadata as `None` if the table does not exist.
705    pub fn check(&self, metadata: Option<&TableMetadata>) -> Result<()> {
706        if let Some(metadata) = metadata {
707            match self {
708                TableRequirement::NotExist => {
709                    return Err(Error::new(
710                        ErrorKind::CatalogCommitConflicts,
711                        format!(
712                            "Requirement failed: Table with id {} already exists",
713                            metadata.uuid()
714                        ),
715                    )
716                    .with_retryable(true));
717                }
718                TableRequirement::UuidMatch { uuid } => {
719                    if &metadata.uuid() != uuid {
720                        return Err(Error::new(
721                            ErrorKind::CatalogCommitConflicts,
722                            "Requirement failed: Table UUID does not match",
723                        )
724                        .with_context("expected", *uuid)
725                        .with_context("found", metadata.uuid())
726                        .with_retryable(true));
727                    }
728                }
729                TableRequirement::CurrentSchemaIdMatch { current_schema_id } => {
730                    // ToDo: Harmonize the types of current_schema_id
731                    if metadata.current_schema_id != *current_schema_id {
732                        return Err(Error::new(
733                            ErrorKind::CatalogCommitConflicts,
734                            "Requirement failed: Current schema id does not match",
735                        )
736                        .with_context("expected", current_schema_id.to_string())
737                        .with_context("found", metadata.current_schema_id.to_string())
738                        .with_retryable(true));
739                    }
740                }
741                TableRequirement::DefaultSortOrderIdMatch {
742                    default_sort_order_id,
743                } => {
744                    if metadata.default_sort_order().order_id != *default_sort_order_id {
745                        return Err(Error::new(
746                            ErrorKind::CatalogCommitConflicts,
747                            "Requirement failed: Default sort order id does not match",
748                        )
749                        .with_context("expected", default_sort_order_id.to_string())
750                        .with_context("found", metadata.default_sort_order().order_id.to_string())
751                        .with_retryable(true));
752                    }
753                }
754                TableRequirement::RefSnapshotIdMatch { r#ref, snapshot_id } => {
755                    let snapshot_ref = metadata.snapshot_for_ref(r#ref);
756                    if let Some(snapshot_id) = snapshot_id {
757                        let snapshot_ref = snapshot_ref.ok_or(
758                            Error::new(
759                                ErrorKind::CatalogCommitConflicts,
760                                format!("Requirement failed: Branch or tag `{ref}` not found"),
761                            )
762                            .with_retryable(true),
763                        )?;
764                        if snapshot_ref.snapshot_id() != *snapshot_id {
765                            return Err(Error::new(
766                                ErrorKind::CatalogCommitConflicts,
767                                format!(
768                                    "Requirement failed: Branch or tag `{ref}`'s snapshot has changed"
769                                ),
770                            )
771                            .with_context("expected", snapshot_id.to_string())
772                            .with_context("found", snapshot_ref.snapshot_id().to_string())
773                            .with_retryable(true));
774                        }
775                    } else if snapshot_ref.is_some() {
776                        // a null snapshot ID means the ref should not exist already
777                        return Err(Error::new(
778                            ErrorKind::CatalogCommitConflicts,
779                            format!("Requirement failed: Branch or tag `{ref}` already exists"),
780                        )
781                        .with_retryable(true));
782                    }
783                }
784                TableRequirement::DefaultSpecIdMatch { default_spec_id } => {
785                    // ToDo: Harmonize the types of default_spec_id
786                    if metadata.default_partition_spec_id() != *default_spec_id {
787                        return Err(Error::new(
788                            ErrorKind::CatalogCommitConflicts,
789                            "Requirement failed: Default partition spec id does not match",
790                        )
791                        .with_context("expected", default_spec_id.to_string())
792                        .with_context("found", metadata.default_partition_spec_id().to_string())
793                        .with_retryable(true));
794                    }
795                }
796                TableRequirement::LastAssignedPartitionIdMatch {
797                    last_assigned_partition_id,
798                } => {
799                    if metadata.last_partition_id != *last_assigned_partition_id {
800                        return Err(Error::new(
801                            ErrorKind::CatalogCommitConflicts,
802                            "Requirement failed: Last assigned partition id does not match",
803                        )
804                        .with_context("expected", last_assigned_partition_id.to_string())
805                        .with_context("found", metadata.last_partition_id.to_string())
806                        .with_retryable(true));
807                    }
808                }
809                TableRequirement::LastAssignedFieldIdMatch {
810                    last_assigned_field_id,
811                } => {
812                    if &metadata.last_column_id != last_assigned_field_id {
813                        return Err(Error::new(
814                            ErrorKind::CatalogCommitConflicts,
815                            "Requirement failed: Last assigned field id does not match",
816                        )
817                        .with_context("expected", last_assigned_field_id.to_string())
818                        .with_context("found", metadata.last_column_id.to_string())
819                        .with_retryable(true));
820                    }
821                }
822            };
823        } else {
824            match self {
825                TableRequirement::NotExist => {}
826                _ => {
827                    return Err(Error::new(
828                        ErrorKind::TableNotFound,
829                        "Requirement failed: Table does not exist",
830                    ));
831                }
832            }
833        }
834
835        Ok(())
836    }
837}
838
839pub(super) mod _serde {
840    use serde::{Deserialize as _, Deserializer, Serialize as _};
841
842    use super::*;
843    use crate::spec::{SchemaId, Summary};
844
845    pub(super) fn deserialize_snapshot<'de, D>(
846        deserializer: D,
847    ) -> std::result::Result<Snapshot, D::Error>
848    where D: Deserializer<'de> {
849        let buf = CatalogSnapshot::deserialize(deserializer)?;
850        Ok(buf.into())
851    }
852
853    pub(super) fn serialize_snapshot<S>(
854        snapshot: &Snapshot,
855        serializer: S,
856    ) -> std::result::Result<S::Ok, S::Error>
857    where
858        S: serde::Serializer,
859    {
860        let buf: CatalogSnapshot = snapshot.clone().into();
861        buf.serialize(serializer)
862    }
863
864    #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
865    #[serde(rename_all = "kebab-case")]
866    /// Defines the structure of a v2 snapshot for the catalog.
867    /// Main difference to SnapshotV2 is that sequence-number is optional
868    /// in the rest catalog spec to allow for backwards compatibility with v1.
869    struct CatalogSnapshot {
870        snapshot_id: i64,
871        #[serde(skip_serializing_if = "Option::is_none")]
872        parent_snapshot_id: Option<i64>,
873        #[serde(default)]
874        sequence_number: i64,
875        timestamp_ms: i64,
876        manifest_list: String,
877        summary: Summary,
878        #[serde(skip_serializing_if = "Option::is_none")]
879        schema_id: Option<SchemaId>,
880        #[serde(skip_serializing_if = "Option::is_none")]
881        first_row_id: Option<u64>,
882        #[serde(skip_serializing_if = "Option::is_none")]
883        added_rows: Option<u64>,
884        #[serde(skip_serializing_if = "Option::is_none")]
885        key_id: Option<String>,
886    }
887
888    impl From<CatalogSnapshot> for Snapshot {
889        fn from(snapshot: CatalogSnapshot) -> Self {
890            let CatalogSnapshot {
891                snapshot_id,
892                parent_snapshot_id,
893                sequence_number,
894                timestamp_ms,
895                manifest_list,
896                schema_id,
897                summary,
898                first_row_id,
899                added_rows,
900                key_id,
901            } = snapshot;
902            let builder = Snapshot::builder()
903                .with_snapshot_id(snapshot_id)
904                .with_parent_snapshot_id(parent_snapshot_id)
905                .with_sequence_number(sequence_number)
906                .with_timestamp_ms(timestamp_ms)
907                .with_manifest_list(manifest_list)
908                .with_summary(summary)
909                .with_encryption_key_id(key_id);
910            let row_range = first_row_id.zip(added_rows);
911            match (schema_id, row_range) {
912                (None, None) => builder.build(),
913                (Some(schema_id), None) => builder.with_schema_id(schema_id).build(),
914                (None, Some((first_row_id, last_row_id))) => {
915                    builder.with_row_range(first_row_id, last_row_id).build()
916                }
917                (Some(schema_id), Some((first_row_id, last_row_id))) => builder
918                    .with_schema_id(schema_id)
919                    .with_row_range(first_row_id, last_row_id)
920                    .build(),
921            }
922        }
923    }
924
925    impl From<Snapshot> for CatalogSnapshot {
926        fn from(snapshot: Snapshot) -> Self {
927            let first_row_id = snapshot.first_row_id();
928            let added_rows = snapshot.added_rows_count();
929            let Snapshot {
930                snapshot_id,
931                parent_snapshot_id,
932                sequence_number,
933                timestamp_ms,
934                manifest_list,
935                summary,
936                schema_id,
937                row_range: _,
938                encryption_key_id: key_id,
939            } = snapshot;
940            CatalogSnapshot {
941                snapshot_id,
942                parent_snapshot_id,
943                sequence_number,
944                timestamp_ms,
945                manifest_list,
946                summary,
947                schema_id,
948                first_row_id,
949                added_rows,
950                key_id,
951            }
952        }
953    }
954}
955
956/// ViewCreation represents the creation of a view in the catalog.
957#[derive(Debug, TypedBuilder)]
958pub struct ViewCreation {
959    /// The name of the view.
960    pub name: String,
961    /// The view's base location; used to create metadata file locations
962    pub location: String,
963    /// Representations for the view.
964    pub representations: ViewRepresentations,
965    /// The schema of the view.
966    pub schema: Schema,
967    /// The properties of the view.
968    #[builder(default)]
969    pub properties: HashMap<String, String>,
970    /// The default namespace to use when a reference in the SELECT is a single identifier
971    pub default_namespace: NamespaceIdent,
972    /// Default catalog to use when a reference in the SELECT does not contain a catalog
973    #[builder(default)]
974    pub default_catalog: Option<String>,
975    /// A string to string map of summary metadata about the version
976    /// Typical keys are "engine-name" and "engine-version"
977    #[builder(default)]
978    pub summary: HashMap<String, String>,
979}
980
981/// ViewUpdate represents an update to a view in the catalog.
982#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
983#[serde(tag = "action", rename_all = "kebab-case")]
984#[allow(clippy::large_enum_variant)]
985pub enum ViewUpdate {
986    /// Assign a new UUID to the view
987    #[serde(rename_all = "kebab-case")]
988    AssignUuid {
989        /// The new UUID to assign.
990        uuid: Uuid,
991    },
992    /// Upgrade view's format version
993    #[serde(rename_all = "kebab-case")]
994    UpgradeFormatVersion {
995        /// Target format upgrade to.
996        format_version: ViewFormatVersion,
997    },
998    /// Add a new schema to the view
999    #[serde(rename_all = "kebab-case")]
1000    AddSchema {
1001        /// The schema to add.
1002        schema: Schema,
1003        /// The last column id of the view.
1004        last_column_id: Option<i32>,
1005    },
1006    /// Set view's current schema
1007    #[serde(rename_all = "kebab-case")]
1008    SetLocation {
1009        /// New location for view.
1010        location: String,
1011    },
1012    /// Set view's properties
1013    ///
1014    /// Matching keys are updated, and non-matching keys are left unchanged.
1015    #[serde(rename_all = "kebab-case")]
1016    SetProperties {
1017        /// Properties to update for view.
1018        updates: HashMap<String, String>,
1019    },
1020    /// Remove view's properties
1021    #[serde(rename_all = "kebab-case")]
1022    RemoveProperties {
1023        /// Properties to remove
1024        removals: Vec<String>,
1025    },
1026    /// Add a new version to the view
1027    #[serde(rename_all = "kebab-case")]
1028    AddViewVersion {
1029        /// The view version to add.
1030        view_version: ViewVersion,
1031    },
1032    /// Set view's current version
1033    #[serde(rename_all = "kebab-case")]
1034    SetCurrentViewVersion {
1035        /// View version id to set as current, or -1 to set last added version
1036        view_version_id: i32,
1037    },
1038}
1039
1040mod _serde_set_statistics {
1041    // The rest spec requires an additional field `snapshot-id`
1042    // that is redundant with the `snapshot_id` field in the statistics file.
1043    use serde::{Deserialize, Deserializer, Serialize, Serializer};
1044
1045    use super::*;
1046
1047    #[derive(Debug, Serialize, Deserialize)]
1048    #[serde(rename_all = "kebab-case")]
1049    struct SetStatistics {
1050        snapshot_id: Option<i64>,
1051        statistics: StatisticsFile,
1052    }
1053
1054    pub fn serialize<S>(
1055        value: &StatisticsFile,
1056        serializer: S,
1057    ) -> std::result::Result<S::Ok, S::Error>
1058    where
1059        S: Serializer,
1060    {
1061        SetStatistics {
1062            snapshot_id: Some(value.snapshot_id),
1063            statistics: value.clone(),
1064        }
1065        .serialize(serializer)
1066    }
1067
1068    pub fn deserialize<'de, D>(deserializer: D) -> std::result::Result<StatisticsFile, D::Error>
1069    where D: Deserializer<'de> {
1070        let SetStatistics {
1071            snapshot_id,
1072            statistics,
1073        } = SetStatistics::deserialize(deserializer)?;
1074        if let Some(snapshot_id) = snapshot_id
1075            && snapshot_id != statistics.snapshot_id
1076        {
1077            return Err(serde::de::Error::custom(format!(
1078                "Snapshot id to set {snapshot_id} does not match the statistics file snapshot id {}",
1079                statistics.snapshot_id
1080            )));
1081        }
1082
1083        Ok(statistics)
1084    }
1085}
1086
1087#[cfg(test)]
1088mod tests {
1089    use std::collections::HashMap;
1090    use std::fmt::Debug;
1091    use std::fs::File;
1092    use std::io::BufReader;
1093
1094    use base64::Engine as _;
1095    use serde::Serialize;
1096    use serde::de::DeserializeOwned;
1097    use uuid::uuid;
1098
1099    use super::ViewUpdate;
1100    use crate::io::FileIO;
1101    use crate::spec::{
1102        BlobMetadata, EncryptedKey, FormatVersion, MAIN_BRANCH, NestedField, NullOrder, Operation,
1103        PartitionStatisticsFile, PrimitiveType, Schema, Snapshot, SnapshotReference,
1104        SnapshotRetention, SortDirection, SortField, SortOrder, SqlViewRepresentation,
1105        StatisticsFile, Summary, TableMetadata, TableMetadataBuilder, Transform, Type,
1106        UnboundPartitionField, UnboundPartitionSpec, ViewFormatVersion, ViewRepresentation,
1107        ViewRepresentations, ViewVersion,
1108    };
1109    use crate::table::Table;
1110    use crate::test_utils::test_runtime;
1111    use crate::{
1112        NamespaceIdent, TableCommit, TableCreation, TableIdent, TableRequirement, TableUpdate,
1113    };
1114
1115    #[test]
1116    fn test_parent_namespace() {
1117        let ns1 = NamespaceIdent::from_strs(vec!["ns1"]).unwrap();
1118        let ns2 = NamespaceIdent::from_strs(vec!["ns1", "ns2"]).unwrap();
1119        let ns3 = NamespaceIdent::from_strs(vec!["ns1", "ns2", "ns3"]).unwrap();
1120
1121        assert_eq!(ns1.parent(), None);
1122        assert_eq!(ns2.parent(), Some(ns1.clone()));
1123        assert_eq!(ns3.parent(), Some(ns2.clone()));
1124    }
1125
1126    #[test]
1127    fn test_create_table_id() {
1128        let table_id = TableIdent {
1129            namespace: NamespaceIdent::from_strs(vec!["ns1"]).unwrap(),
1130            name: "t1".to_string(),
1131        };
1132
1133        assert_eq!(table_id, TableIdent::from_strs(vec!["ns1", "t1"]).unwrap());
1134    }
1135
1136    #[test]
1137    fn test_table_creation_iterator_properties() {
1138        let builder = TableCreation::builder()
1139            .name("table".to_string())
1140            .schema(Schema::builder().build().unwrap());
1141
1142        fn s(k: &str, v: &str) -> (String, String) {
1143            (k.to_string(), v.to_string())
1144        }
1145
1146        let table_creation = builder
1147            .properties([s("key", "value"), s("foo", "bar")])
1148            .build();
1149
1150        assert_eq!(
1151            HashMap::from([s("key", "value"), s("foo", "bar")]),
1152            table_creation.properties
1153        );
1154    }
1155
1156    fn test_serde_json<T: Serialize + DeserializeOwned + PartialEq + Debug>(
1157        json: impl ToString,
1158        expected: T,
1159    ) {
1160        let json_str = json.to_string();
1161        let actual: T = serde_json::from_str(&json_str).expect("Failed to parse from json");
1162        assert_eq!(actual, expected, "Parsed value is not equal to expected");
1163
1164        let restored: T = serde_json::from_str(
1165            &serde_json::to_string(&actual).expect("Failed to serialize to json"),
1166        )
1167        .expect("Failed to parse from serialized json");
1168
1169        assert_eq!(
1170            restored, expected,
1171            "Parsed restored value is not equal to expected"
1172        );
1173    }
1174
1175    fn metadata() -> TableMetadata {
1176        let tbl_creation = TableCreation::builder()
1177            .name("table".to_string())
1178            .location("/path/to/table".to_string())
1179            .schema(Schema::builder().build().unwrap())
1180            .build();
1181
1182        TableMetadataBuilder::from_table_creation(tbl_creation)
1183            .unwrap()
1184            .assign_uuid(uuid::Uuid::nil())
1185            .build()
1186            .unwrap()
1187            .metadata
1188    }
1189
1190    #[test]
1191    fn test_check_requirement_not_exist() {
1192        let metadata = metadata();
1193        let requirement = TableRequirement::NotExist;
1194
1195        assert!(requirement.check(Some(&metadata)).is_err());
1196        assert!(requirement.check(None).is_ok());
1197    }
1198
1199    #[test]
1200    fn test_check_table_uuid() {
1201        let metadata = metadata();
1202
1203        let requirement = TableRequirement::UuidMatch {
1204            uuid: uuid::Uuid::now_v7(),
1205        };
1206        assert!(requirement.check(Some(&metadata)).is_err());
1207
1208        let requirement = TableRequirement::UuidMatch {
1209            uuid: uuid::Uuid::nil(),
1210        };
1211        assert!(requirement.check(Some(&metadata)).is_ok());
1212    }
1213
1214    #[test]
1215    fn test_check_ref_snapshot_id() {
1216        let metadata = metadata();
1217
1218        // Ref does not exist but should
1219        let requirement = TableRequirement::RefSnapshotIdMatch {
1220            r#ref: "my_branch".to_string(),
1221            snapshot_id: Some(1),
1222        };
1223        assert!(requirement.check(Some(&metadata)).is_err());
1224
1225        // Ref does not exist and should not
1226        let requirement = TableRequirement::RefSnapshotIdMatch {
1227            r#ref: "my_branch".to_string(),
1228            snapshot_id: None,
1229        };
1230        assert!(requirement.check(Some(&metadata)).is_ok());
1231
1232        // Add snapshot
1233        let snapshot = Snapshot::builder()
1234            .with_snapshot_id(3051729675574597004)
1235            .with_sequence_number(10)
1236            .with_timestamp_ms(9992191116217)
1237            .with_manifest_list("s3://b/wh/.../s1.avro".to_string())
1238            .with_schema_id(0)
1239            .with_summary(Summary {
1240                operation: Operation::Append,
1241                additional_properties: HashMap::new(),
1242            })
1243            .build();
1244
1245        let builder = metadata.into_builder(None);
1246        let builder = TableUpdate::AddSnapshot {
1247            snapshot: snapshot.clone(),
1248        }
1249        .apply(builder)
1250        .unwrap();
1251        let metadata = TableUpdate::SetSnapshotRef {
1252            ref_name: MAIN_BRANCH.to_string(),
1253            reference: SnapshotReference {
1254                snapshot_id: snapshot.snapshot_id(),
1255                retention: SnapshotRetention::Branch {
1256                    min_snapshots_to_keep: Some(10),
1257                    max_snapshot_age_ms: None,
1258                    max_ref_age_ms: None,
1259                },
1260            },
1261        }
1262        .apply(builder)
1263        .unwrap()
1264        .build()
1265        .unwrap()
1266        .metadata;
1267
1268        // Ref exists and should match
1269        let requirement = TableRequirement::RefSnapshotIdMatch {
1270            r#ref: "main".to_string(),
1271            snapshot_id: Some(3051729675574597004),
1272        };
1273        assert!(requirement.check(Some(&metadata)).is_ok());
1274
1275        // Ref exists but does not match
1276        let requirement = TableRequirement::RefSnapshotIdMatch {
1277            r#ref: "main".to_string(),
1278            snapshot_id: Some(1),
1279        };
1280        assert!(requirement.check(Some(&metadata)).is_err());
1281    }
1282
1283    #[test]
1284    fn test_check_last_assigned_field_id() {
1285        let metadata = metadata();
1286
1287        let requirement = TableRequirement::LastAssignedFieldIdMatch {
1288            last_assigned_field_id: 1,
1289        };
1290        assert!(requirement.check(Some(&metadata)).is_err());
1291
1292        let requirement = TableRequirement::LastAssignedFieldIdMatch {
1293            last_assigned_field_id: 0,
1294        };
1295        assert!(requirement.check(Some(&metadata)).is_ok());
1296    }
1297
1298    #[test]
1299    fn test_check_current_schema_id() {
1300        let metadata = metadata();
1301
1302        let requirement = TableRequirement::CurrentSchemaIdMatch {
1303            current_schema_id: 1,
1304        };
1305        assert!(requirement.check(Some(&metadata)).is_err());
1306
1307        let requirement = TableRequirement::CurrentSchemaIdMatch {
1308            current_schema_id: 0,
1309        };
1310        assert!(requirement.check(Some(&metadata)).is_ok());
1311    }
1312
1313    #[test]
1314    fn test_check_last_assigned_partition_id() {
1315        let metadata = metadata();
1316        let requirement = TableRequirement::LastAssignedPartitionIdMatch {
1317            last_assigned_partition_id: 0,
1318        };
1319        assert!(requirement.check(Some(&metadata)).is_err());
1320
1321        let requirement = TableRequirement::LastAssignedPartitionIdMatch {
1322            last_assigned_partition_id: 999,
1323        };
1324        assert!(requirement.check(Some(&metadata)).is_ok());
1325    }
1326
1327    #[test]
1328    fn test_check_default_spec_id() {
1329        let metadata = metadata();
1330
1331        let requirement = TableRequirement::DefaultSpecIdMatch { default_spec_id: 1 };
1332        assert!(requirement.check(Some(&metadata)).is_err());
1333
1334        let requirement = TableRequirement::DefaultSpecIdMatch { default_spec_id: 0 };
1335        assert!(requirement.check(Some(&metadata)).is_ok());
1336    }
1337
1338    #[test]
1339    fn test_check_default_sort_order_id() {
1340        let metadata = metadata();
1341
1342        let requirement = TableRequirement::DefaultSortOrderIdMatch {
1343            default_sort_order_id: 1,
1344        };
1345        assert!(requirement.check(Some(&metadata)).is_err());
1346
1347        let requirement = TableRequirement::DefaultSortOrderIdMatch {
1348            default_sort_order_id: 0,
1349        };
1350        assert!(requirement.check(Some(&metadata)).is_ok());
1351    }
1352
1353    #[test]
1354    fn test_table_uuid() {
1355        test_serde_json(
1356            r#"
1357{
1358    "type": "assert-table-uuid",
1359    "uuid": "2cc52516-5e73-41f2-b139-545d41a4e151"
1360}
1361        "#,
1362            TableRequirement::UuidMatch {
1363                uuid: uuid!("2cc52516-5e73-41f2-b139-545d41a4e151"),
1364            },
1365        );
1366    }
1367
1368    #[test]
1369    fn test_assert_table_not_exists() {
1370        test_serde_json(
1371            r#"
1372{
1373    "type": "assert-create"
1374}
1375        "#,
1376            TableRequirement::NotExist,
1377        );
1378    }
1379
1380    #[test]
1381    fn test_assert_ref_snapshot_id() {
1382        test_serde_json(
1383            r#"
1384{
1385    "type": "assert-ref-snapshot-id",
1386    "ref": "snapshot-name",
1387    "snapshot-id": null
1388}
1389        "#,
1390            TableRequirement::RefSnapshotIdMatch {
1391                r#ref: "snapshot-name".to_string(),
1392                snapshot_id: None,
1393            },
1394        );
1395
1396        test_serde_json(
1397            r#"
1398{
1399    "type": "assert-ref-snapshot-id",
1400    "ref": "snapshot-name",
1401    "snapshot-id": 1
1402}
1403        "#,
1404            TableRequirement::RefSnapshotIdMatch {
1405                r#ref: "snapshot-name".to_string(),
1406                snapshot_id: Some(1),
1407            },
1408        );
1409    }
1410
1411    #[test]
1412    fn test_assert_last_assigned_field_id() {
1413        test_serde_json(
1414            r#"
1415{
1416    "type": "assert-last-assigned-field-id",
1417    "last-assigned-field-id": 12
1418}
1419        "#,
1420            TableRequirement::LastAssignedFieldIdMatch {
1421                last_assigned_field_id: 12,
1422            },
1423        );
1424    }
1425
1426    #[test]
1427    fn test_assert_current_schema_id() {
1428        test_serde_json(
1429            r#"
1430{
1431    "type": "assert-current-schema-id",
1432    "current-schema-id": 4
1433}
1434        "#,
1435            TableRequirement::CurrentSchemaIdMatch {
1436                current_schema_id: 4,
1437            },
1438        );
1439    }
1440
1441    #[test]
1442    fn test_assert_last_assigned_partition_id() {
1443        test_serde_json(
1444            r#"
1445{
1446    "type": "assert-last-assigned-partition-id",
1447    "last-assigned-partition-id": 1004
1448}
1449        "#,
1450            TableRequirement::LastAssignedPartitionIdMatch {
1451                last_assigned_partition_id: 1004,
1452            },
1453        );
1454    }
1455
1456    #[test]
1457    fn test_assert_default_spec_id() {
1458        test_serde_json(
1459            r#"
1460{
1461    "type": "assert-default-spec-id",
1462    "default-spec-id": 5
1463}
1464        "#,
1465            TableRequirement::DefaultSpecIdMatch { default_spec_id: 5 },
1466        );
1467    }
1468
1469    #[test]
1470    fn test_assert_default_sort_order() {
1471        let json = r#"
1472{
1473    "type": "assert-default-sort-order-id",
1474    "default-sort-order-id": 10
1475}
1476        "#;
1477
1478        let update = TableRequirement::DefaultSortOrderIdMatch {
1479            default_sort_order_id: 10,
1480        };
1481
1482        test_serde_json(json, update);
1483    }
1484
1485    #[test]
1486    fn test_parse_assert_invalid() {
1487        assert!(
1488            serde_json::from_str::<TableRequirement>(
1489                r#"
1490{
1491    "default-sort-order-id": 10
1492}
1493"#
1494            )
1495            .is_err(),
1496            "Table requirements should not be parsed without type."
1497        );
1498    }
1499
1500    #[test]
1501    fn test_assign_uuid() {
1502        test_serde_json(
1503            r#"
1504{
1505    "action": "assign-uuid",
1506    "uuid": "2cc52516-5e73-41f2-b139-545d41a4e151"
1507}
1508        "#,
1509            TableUpdate::AssignUuid {
1510                uuid: uuid!("2cc52516-5e73-41f2-b139-545d41a4e151"),
1511            },
1512        );
1513    }
1514
1515    #[test]
1516    fn test_upgrade_format_version() {
1517        test_serde_json(
1518            r#"
1519{
1520    "action": "upgrade-format-version",
1521    "format-version": 2
1522}
1523        "#,
1524            TableUpdate::UpgradeFormatVersion {
1525                format_version: FormatVersion::V2,
1526            },
1527        );
1528    }
1529
1530    #[test]
1531    fn test_add_schema() {
1532        let test_schema = Schema::builder()
1533            .with_schema_id(1)
1534            .with_identifier_field_ids(vec![2])
1535            .with_fields(vec![
1536                NestedField::optional(1, "foo", Type::Primitive(PrimitiveType::String)).into(),
1537                NestedField::required(2, "bar", Type::Primitive(PrimitiveType::Int)).into(),
1538                NestedField::optional(3, "baz", Type::Primitive(PrimitiveType::Boolean)).into(),
1539            ])
1540            .build()
1541            .unwrap();
1542        test_serde_json(
1543            r#"
1544{
1545    "action": "add-schema",
1546    "schema": {
1547        "type": "struct",
1548        "schema-id": 1,
1549        "fields": [
1550            {
1551                "id": 1,
1552                "name": "foo",
1553                "required": false,
1554                "type": "string"
1555            },
1556            {
1557                "id": 2,
1558                "name": "bar",
1559                "required": true,
1560                "type": "int"
1561            },
1562            {
1563                "id": 3,
1564                "name": "baz",
1565                "required": false,
1566                "type": "boolean"
1567            }
1568        ],
1569        "identifier-field-ids": [
1570            2
1571        ]
1572    },
1573    "last-column-id": 3
1574}
1575        "#,
1576            TableUpdate::AddSchema {
1577                schema: test_schema.clone(),
1578            },
1579        );
1580
1581        test_serde_json(
1582            r#"
1583{
1584    "action": "add-schema",
1585    "schema": {
1586        "type": "struct",
1587        "schema-id": 1,
1588        "fields": [
1589            {
1590                "id": 1,
1591                "name": "foo",
1592                "required": false,
1593                "type": "string"
1594            },
1595            {
1596                "id": 2,
1597                "name": "bar",
1598                "required": true,
1599                "type": "int"
1600            },
1601            {
1602                "id": 3,
1603                "name": "baz",
1604                "required": false,
1605                "type": "boolean"
1606            }
1607        ],
1608        "identifier-field-ids": [
1609            2
1610        ]
1611    }
1612}
1613        "#,
1614            TableUpdate::AddSchema {
1615                schema: test_schema.clone(),
1616            },
1617        );
1618    }
1619
1620    #[test]
1621    fn test_set_current_schema() {
1622        test_serde_json(
1623            r#"
1624{
1625   "action": "set-current-schema",
1626   "schema-id": 23
1627}
1628        "#,
1629            TableUpdate::SetCurrentSchema { schema_id: 23 },
1630        );
1631    }
1632
1633    #[test]
1634    fn test_add_spec() {
1635        test_serde_json(
1636            r#"
1637{
1638    "action": "add-spec",
1639    "spec": {
1640        "fields": [
1641            {
1642                "source-id": 4,
1643                "name": "ts_day",
1644                "transform": "day"
1645            },
1646            {
1647                "source-id": 1,
1648                "name": "id_bucket",
1649                "transform": "bucket[16]"
1650            },
1651            {
1652                "source-id": 2,
1653                "name": "id_truncate",
1654                "transform": "truncate[4]"
1655            }
1656        ]
1657    }
1658}
1659        "#,
1660            TableUpdate::AddSpec {
1661                spec: UnboundPartitionSpec::builder()
1662                    .add_partition_field(
1663                        UnboundPartitionField::builder()
1664                            .source_ids(vec![4])
1665                            .name("ts_day")
1666                            .transform(Transform::Day)
1667                            .build()
1668                            .unwrap(),
1669                    )
1670                    .unwrap()
1671                    .add_partition_field(
1672                        UnboundPartitionField::builder()
1673                            .source_ids(vec![1])
1674                            .name("id_bucket")
1675                            .transform(Transform::Bucket(16))
1676                            .build()
1677                            .unwrap(),
1678                    )
1679                    .unwrap()
1680                    .add_partition_field(
1681                        UnboundPartitionField::builder()
1682                            .source_ids(vec![2])
1683                            .name("id_truncate")
1684                            .transform(Transform::Truncate(4))
1685                            .build()
1686                            .unwrap(),
1687                    )
1688                    .unwrap()
1689                    .build(),
1690            },
1691        );
1692    }
1693
1694    #[test]
1695    fn test_set_default_spec() {
1696        test_serde_json(
1697            r#"
1698{
1699    "action": "set-default-spec",
1700    "spec-id": 1
1701}
1702        "#,
1703            TableUpdate::SetDefaultSpec { spec_id: 1 },
1704        )
1705    }
1706
1707    #[test]
1708    fn test_add_sort_order() {
1709        let json = r#"
1710{
1711    "action": "add-sort-order",
1712    "sort-order": {
1713        "order-id": 1,
1714        "fields": [
1715            {
1716                "transform": "identity",
1717                "source-id": 2,
1718                "direction": "asc",
1719                "null-order": "nulls-first"
1720            },
1721            {
1722                "transform": "bucket[4]",
1723                "source-id": 3,
1724                "direction": "desc",
1725                "null-order": "nulls-last"
1726            }
1727        ]
1728    }
1729}
1730        "#;
1731
1732        let update = TableUpdate::AddSortOrder {
1733            sort_order: SortOrder::builder()
1734                .with_order_id(1)
1735                .with_sort_field(
1736                    SortField::builder()
1737                        .source_id(2)
1738                        .direction(SortDirection::Ascending)
1739                        .null_order(NullOrder::First)
1740                        .transform(Transform::Identity)
1741                        .build(),
1742                )
1743                .with_sort_field(
1744                    SortField::builder()
1745                        .source_id(3)
1746                        .direction(SortDirection::Descending)
1747                        .null_order(NullOrder::Last)
1748                        .transform(Transform::Bucket(4))
1749                        .build(),
1750                )
1751                .build_unbound()
1752                .unwrap(),
1753        };
1754
1755        test_serde_json(json, update);
1756    }
1757
1758    #[test]
1759    fn test_set_default_order() {
1760        let json = r#"
1761{
1762    "action": "set-default-sort-order",
1763    "sort-order-id": 2
1764}
1765        "#;
1766        let update = TableUpdate::SetDefaultSortOrder { sort_order_id: 2 };
1767
1768        test_serde_json(json, update);
1769    }
1770
1771    #[test]
1772    fn test_add_snapshot() {
1773        let json = r#"
1774{
1775    "action": "add-snapshot",
1776    "snapshot": {
1777        "snapshot-id": 3055729675574597000,
1778        "parent-snapshot-id": 3051729675574597000,
1779        "timestamp-ms": 1555100955770,
1780        "sequence-number": 1,
1781        "summary": {
1782            "operation": "append"
1783        },
1784        "manifest-list": "s3://a/b/2.avro",
1785        "schema-id": 1
1786    }
1787}
1788        "#;
1789
1790        let update = TableUpdate::AddSnapshot {
1791            snapshot: Snapshot::builder()
1792                .with_snapshot_id(3055729675574597000)
1793                .with_parent_snapshot_id(Some(3051729675574597000))
1794                .with_timestamp_ms(1555100955770)
1795                .with_sequence_number(1)
1796                .with_manifest_list("s3://a/b/2.avro")
1797                .with_schema_id(1)
1798                .with_summary(Summary {
1799                    operation: Operation::Append,
1800                    additional_properties: HashMap::default(),
1801                })
1802                .build(),
1803        };
1804
1805        test_serde_json(json, update);
1806    }
1807
1808    #[test]
1809    fn test_add_snapshot_v1() {
1810        let json = r#"
1811{
1812    "action": "add-snapshot",
1813    "snapshot": {
1814        "snapshot-id": 3055729675574597000,
1815        "parent-snapshot-id": 3051729675574597000,
1816        "timestamp-ms": 1555100955770,
1817        "summary": {
1818            "operation": "append"
1819        },
1820        "manifest-list": "s3://a/b/2.avro"
1821    }
1822}
1823    "#;
1824
1825        let update = TableUpdate::AddSnapshot {
1826            snapshot: Snapshot::builder()
1827                .with_snapshot_id(3055729675574597000)
1828                .with_parent_snapshot_id(Some(3051729675574597000))
1829                .with_timestamp_ms(1555100955770)
1830                .with_sequence_number(0)
1831                .with_manifest_list("s3://a/b/2.avro")
1832                .with_summary(Summary {
1833                    operation: Operation::Append,
1834                    additional_properties: HashMap::default(),
1835                })
1836                .build(),
1837        };
1838
1839        let actual: TableUpdate = serde_json::from_str(json).expect("Failed to parse from json");
1840        assert_eq!(actual, update, "Parsed value is not equal to expected");
1841    }
1842
1843    #[test]
1844    fn test_add_snapshot_v3() {
1845        let json = serde_json::json!(
1846        {
1847            "action": "add-snapshot",
1848            "snapshot": {
1849                "snapshot-id": 3055729675574597000i64,
1850                "parent-snapshot-id": 3051729675574597000i64,
1851                "timestamp-ms": 1555100955770i64,
1852                "first-row-id":0,
1853                "added-rows":2,
1854                "key-id":"key123",
1855                "summary": {
1856                    "operation": "append"
1857                },
1858                "manifest-list": "s3://a/b/2.avro"
1859            }
1860        });
1861
1862        let update = TableUpdate::AddSnapshot {
1863            snapshot: Snapshot::builder()
1864                .with_snapshot_id(3055729675574597000)
1865                .with_parent_snapshot_id(Some(3051729675574597000))
1866                .with_timestamp_ms(1555100955770)
1867                .with_sequence_number(0)
1868                .with_manifest_list("s3://a/b/2.avro")
1869                .with_row_range(0, 2)
1870                .with_encryption_key_id(Some("key123".to_string()))
1871                .with_summary(Summary {
1872                    operation: Operation::Append,
1873                    additional_properties: HashMap::default(),
1874                })
1875                .build(),
1876        };
1877
1878        let actual: TableUpdate = serde_json::from_value(json).expect("Failed to parse from json");
1879        assert_eq!(actual, update, "Parsed value is not equal to expected");
1880        let restored: TableUpdate = serde_json::from_str(
1881            &serde_json::to_string(&actual).expect("Failed to serialize to json"),
1882        )
1883        .expect("Failed to parse from serialized json");
1884        assert_eq!(restored, update);
1885    }
1886
1887    #[test]
1888    fn test_remove_snapshots() {
1889        let json = r#"
1890{
1891    "action": "remove-snapshots",
1892    "snapshot-ids": [
1893        1,
1894        2
1895    ]
1896}
1897        "#;
1898
1899        let update = TableUpdate::RemoveSnapshots {
1900            snapshot_ids: vec![1, 2],
1901        };
1902        test_serde_json(json, update);
1903    }
1904
1905    #[test]
1906    fn test_remove_snapshot_ref() {
1907        let json = r#"
1908{
1909    "action": "remove-snapshot-ref",
1910    "ref-name": "snapshot-ref"
1911}
1912        "#;
1913
1914        let update = TableUpdate::RemoveSnapshotRef {
1915            ref_name: "snapshot-ref".to_string(),
1916        };
1917        test_serde_json(json, update);
1918    }
1919
1920    #[test]
1921    fn test_set_snapshot_ref_tag() {
1922        let json = r#"
1923{
1924    "action": "set-snapshot-ref",
1925    "type": "tag",
1926    "ref-name": "hank",
1927    "snapshot-id": 1,
1928    "max-ref-age-ms": 1
1929}
1930        "#;
1931
1932        let update = TableUpdate::SetSnapshotRef {
1933            ref_name: "hank".to_string(),
1934            reference: SnapshotReference {
1935                snapshot_id: 1,
1936                retention: SnapshotRetention::Tag {
1937                    max_ref_age_ms: Some(1),
1938                },
1939            },
1940        };
1941
1942        test_serde_json(json, update);
1943    }
1944
1945    #[test]
1946    fn test_set_snapshot_ref_branch() {
1947        let json = r#"
1948{
1949    "action": "set-snapshot-ref",
1950    "type": "branch",
1951    "ref-name": "hank",
1952    "snapshot-id": 1,
1953    "min-snapshots-to-keep": 2,
1954    "max-snapshot-age-ms": 3,
1955    "max-ref-age-ms": 4
1956}
1957        "#;
1958
1959        let update = TableUpdate::SetSnapshotRef {
1960            ref_name: "hank".to_string(),
1961            reference: SnapshotReference {
1962                snapshot_id: 1,
1963                retention: SnapshotRetention::Branch {
1964                    min_snapshots_to_keep: Some(2),
1965                    max_snapshot_age_ms: Some(3),
1966                    max_ref_age_ms: Some(4),
1967                },
1968            },
1969        };
1970
1971        test_serde_json(json, update);
1972    }
1973
1974    #[test]
1975    fn test_set_properties() {
1976        let json = r#"
1977{
1978    "action": "set-properties",
1979    "updates": {
1980        "prop1": "v1",
1981        "prop2": "v2"
1982    }
1983}
1984        "#;
1985
1986        let update = TableUpdate::SetProperties {
1987            updates: vec![
1988                ("prop1".to_string(), "v1".to_string()),
1989                ("prop2".to_string(), "v2".to_string()),
1990            ]
1991            .into_iter()
1992            .collect(),
1993        };
1994
1995        test_serde_json(json, update);
1996    }
1997
1998    #[test]
1999    fn test_remove_properties() {
2000        let json = r#"
2001{
2002    "action": "remove-properties",
2003    "removals": [
2004        "prop1",
2005        "prop2"
2006    ]
2007}
2008        "#;
2009
2010        let update = TableUpdate::RemoveProperties {
2011            removals: vec!["prop1".to_string(), "prop2".to_string()],
2012        };
2013
2014        test_serde_json(json, update);
2015    }
2016
2017    #[test]
2018    fn test_set_location() {
2019        let json = r#"
2020{
2021    "action": "set-location",
2022    "location": "s3://bucket/warehouse/tbl_location"
2023}
2024    "#;
2025
2026        let update = TableUpdate::SetLocation {
2027            location: "s3://bucket/warehouse/tbl_location".to_string(),
2028        };
2029
2030        test_serde_json(json, update);
2031    }
2032
2033    #[test]
2034    fn test_table_update_apply() {
2035        let table_creation = TableCreation::builder()
2036            .location("s3://db/table".to_string())
2037            .name("table".to_string())
2038            .properties(HashMap::new())
2039            .schema(Schema::builder().build().unwrap())
2040            .build();
2041        let table_metadata = TableMetadataBuilder::from_table_creation(table_creation)
2042            .unwrap()
2043            .build()
2044            .unwrap()
2045            .metadata;
2046        let table_metadata_builder = TableMetadataBuilder::new_from_metadata(
2047            table_metadata,
2048            Some("s3://db/table/metadata/metadata1.gz.json".to_string()),
2049        );
2050
2051        let uuid = uuid::Uuid::new_v4();
2052        let update = TableUpdate::AssignUuid { uuid };
2053        let updated_metadata = update
2054            .apply(table_metadata_builder)
2055            .unwrap()
2056            .build()
2057            .unwrap()
2058            .metadata;
2059        assert_eq!(updated_metadata.uuid(), uuid);
2060    }
2061
2062    #[test]
2063    fn test_view_assign_uuid() {
2064        test_serde_json(
2065            r#"
2066{
2067    "action": "assign-uuid",
2068    "uuid": "2cc52516-5e73-41f2-b139-545d41a4e151"
2069}
2070        "#,
2071            ViewUpdate::AssignUuid {
2072                uuid: uuid!("2cc52516-5e73-41f2-b139-545d41a4e151"),
2073            },
2074        );
2075    }
2076
2077    #[test]
2078    fn test_view_upgrade_format_version() {
2079        test_serde_json(
2080            r#"
2081{
2082    "action": "upgrade-format-version",
2083    "format-version": 1
2084}
2085        "#,
2086            ViewUpdate::UpgradeFormatVersion {
2087                format_version: ViewFormatVersion::V1,
2088            },
2089        );
2090    }
2091
2092    #[test]
2093    fn test_view_add_schema() {
2094        let test_schema = Schema::builder()
2095            .with_schema_id(1)
2096            .with_identifier_field_ids(vec![2])
2097            .with_fields(vec![
2098                NestedField::optional(1, "foo", Type::Primitive(PrimitiveType::String)).into(),
2099                NestedField::required(2, "bar", Type::Primitive(PrimitiveType::Int)).into(),
2100                NestedField::optional(3, "baz", Type::Primitive(PrimitiveType::Boolean)).into(),
2101            ])
2102            .build()
2103            .unwrap();
2104        test_serde_json(
2105            r#"
2106{
2107    "action": "add-schema",
2108    "schema": {
2109        "type": "struct",
2110        "schema-id": 1,
2111        "fields": [
2112            {
2113                "id": 1,
2114                "name": "foo",
2115                "required": false,
2116                "type": "string"
2117            },
2118            {
2119                "id": 2,
2120                "name": "bar",
2121                "required": true,
2122                "type": "int"
2123            },
2124            {
2125                "id": 3,
2126                "name": "baz",
2127                "required": false,
2128                "type": "boolean"
2129            }
2130        ],
2131        "identifier-field-ids": [
2132            2
2133        ]
2134    },
2135    "last-column-id": 3
2136}
2137        "#,
2138            ViewUpdate::AddSchema {
2139                schema: test_schema.clone(),
2140                last_column_id: Some(3),
2141            },
2142        );
2143    }
2144
2145    #[test]
2146    fn test_view_set_location() {
2147        test_serde_json(
2148            r#"
2149{
2150    "action": "set-location",
2151    "location": "s3://db/view"
2152}
2153        "#,
2154            ViewUpdate::SetLocation {
2155                location: "s3://db/view".to_string(),
2156            },
2157        );
2158    }
2159
2160    #[test]
2161    fn test_view_set_properties() {
2162        test_serde_json(
2163            r#"
2164{
2165    "action": "set-properties",
2166    "updates": {
2167        "prop1": "v1",
2168        "prop2": "v2"
2169    }
2170}
2171        "#,
2172            ViewUpdate::SetProperties {
2173                updates: vec![
2174                    ("prop1".to_string(), "v1".to_string()),
2175                    ("prop2".to_string(), "v2".to_string()),
2176                ]
2177                .into_iter()
2178                .collect(),
2179            },
2180        );
2181    }
2182
2183    #[test]
2184    fn test_view_remove_properties() {
2185        test_serde_json(
2186            r#"
2187{
2188    "action": "remove-properties",
2189    "removals": [
2190        "prop1",
2191        "prop2"
2192    ]
2193}
2194        "#,
2195            ViewUpdate::RemoveProperties {
2196                removals: vec!["prop1".to_string(), "prop2".to_string()],
2197            },
2198        );
2199    }
2200
2201    #[test]
2202    fn test_view_add_view_version() {
2203        test_serde_json(
2204            r#"
2205{
2206    "action": "add-view-version",
2207    "view-version": {
2208            "version-id" : 1,
2209            "timestamp-ms" : 1573518431292,
2210            "schema-id" : 1,
2211            "default-catalog" : "prod",
2212            "default-namespace" : [ "default" ],
2213            "summary" : {
2214              "engine-name" : "Spark"
2215            },
2216            "representations" : [ {
2217              "type" : "sql",
2218              "sql" : "SELECT\n    COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2",
2219              "dialect" : "spark"
2220            } ]
2221    }
2222}
2223        "#,
2224            ViewUpdate::AddViewVersion {
2225                view_version: ViewVersion::builder()
2226                    .with_version_id(1)
2227                    .with_timestamp_ms(1573518431292)
2228                    .with_schema_id(1)
2229                    .with_default_catalog(Some("prod".to_string()))
2230                    .with_default_namespace(NamespaceIdent::from_strs(vec!["default"]).unwrap())
2231                    .with_summary(
2232                        vec![("engine-name".to_string(), "Spark".to_string())]
2233                            .into_iter()
2234                            .collect(),
2235                    )
2236                    .with_representations(ViewRepresentations(vec![ViewRepresentation::Sql(SqlViewRepresentation {
2237                        sql: "SELECT\n    COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2".to_string(),
2238                        dialect: "spark".to_string(),
2239                    })]))
2240                    .build(),
2241            },
2242        );
2243    }
2244
2245    #[test]
2246    fn test_view_set_current_view_version() {
2247        test_serde_json(
2248            r#"
2249{
2250    "action": "set-current-view-version",
2251    "view-version-id": 1
2252}
2253        "#,
2254            ViewUpdate::SetCurrentViewVersion { view_version_id: 1 },
2255        );
2256    }
2257
2258    #[test]
2259    fn test_remove_partition_specs_update() {
2260        test_serde_json(
2261            r#"
2262{
2263    "action": "remove-partition-specs",
2264    "spec-ids": [1, 2]
2265}
2266        "#,
2267            TableUpdate::RemovePartitionSpecs {
2268                spec_ids: vec![1, 2],
2269            },
2270        );
2271    }
2272
2273    #[test]
2274    fn test_set_statistics_file() {
2275        test_serde_json(
2276            r#"
2277        {
2278                "action": "set-statistics",
2279                "snapshot-id": 1940541653261589030,
2280                "statistics": {
2281                        "snapshot-id": 1940541653261589030,
2282                        "statistics-path": "s3://bucket/warehouse/stats.puffin",
2283                        "file-size-in-bytes": 124,
2284                        "file-footer-size-in-bytes": 27,
2285                        "blob-metadata": [
2286                                {
2287                                        "type": "boring-type",
2288                                        "snapshot-id": 1940541653261589030,
2289                                        "sequence-number": 2,
2290                                        "fields": [
2291                                                1
2292                                        ],
2293                                        "properties": {
2294                                                "prop-key": "prop-value"
2295                                        }
2296                                }
2297                        ]
2298                }
2299        }
2300        "#,
2301            TableUpdate::SetStatistics {
2302                statistics: StatisticsFile {
2303                    snapshot_id: 1940541653261589030,
2304                    statistics_path: "s3://bucket/warehouse/stats.puffin".to_string(),
2305                    file_size_in_bytes: 124,
2306                    file_footer_size_in_bytes: 27,
2307                    key_metadata: None,
2308                    blob_metadata: vec![BlobMetadata {
2309                        r#type: "boring-type".to_string(),
2310                        snapshot_id: 1940541653261589030,
2311                        sequence_number: 2,
2312                        fields: vec![1],
2313                        properties: vec![("prop-key".to_string(), "prop-value".to_string())]
2314                            .into_iter()
2315                            .collect(),
2316                    }],
2317                },
2318            },
2319        );
2320    }
2321
2322    #[test]
2323    fn test_remove_statistics_file() {
2324        test_serde_json(
2325            r#"
2326        {
2327                "action": "remove-statistics",
2328                "snapshot-id": 1940541653261589030
2329        }
2330        "#,
2331            TableUpdate::RemoveStatistics {
2332                snapshot_id: 1940541653261589030,
2333            },
2334        );
2335    }
2336
2337    #[test]
2338    fn test_set_partition_statistics_file() {
2339        test_serde_json(
2340            r#"
2341            {
2342                "action": "set-partition-statistics",
2343                "partition-statistics": {
2344                    "snapshot-id": 1940541653261589030,
2345                    "statistics-path": "s3://bucket/warehouse/stats1.parquet",
2346                    "file-size-in-bytes": 43
2347                }
2348            }
2349            "#,
2350            TableUpdate::SetPartitionStatistics {
2351                partition_statistics: PartitionStatisticsFile {
2352                    snapshot_id: 1940541653261589030,
2353                    statistics_path: "s3://bucket/warehouse/stats1.parquet".to_string(),
2354                    file_size_in_bytes: 43,
2355                },
2356            },
2357        )
2358    }
2359
2360    #[test]
2361    fn test_remove_partition_statistics_file() {
2362        test_serde_json(
2363            r#"
2364            {
2365                "action": "remove-partition-statistics",
2366                "snapshot-id": 1940541653261589030
2367            }
2368            "#,
2369            TableUpdate::RemovePartitionStatistics {
2370                snapshot_id: 1940541653261589030,
2371            },
2372        )
2373    }
2374
2375    #[test]
2376    fn test_remove_schema_update() {
2377        test_serde_json(
2378            r#"
2379                {
2380                    "action": "remove-schemas",
2381                    "schema-ids": [1, 2]
2382                }        
2383            "#,
2384            TableUpdate::RemoveSchemas {
2385                schema_ids: vec![1, 2],
2386            },
2387        );
2388    }
2389
2390    #[test]
2391    fn test_add_encryption_key() {
2392        let key_bytes = "key".as_bytes();
2393        let encoded_key = base64::engine::general_purpose::STANDARD.encode(key_bytes);
2394        test_serde_json(
2395            format!(
2396                r#"
2397                {{
2398                    "action": "add-encryption-key",
2399                    "encryption-key": {{
2400                        "key-id": "a",
2401                        "encrypted-key-metadata": "{encoded_key}",
2402                        "encrypted-by-id": "b"
2403                    }}
2404                }}        
2405            "#
2406            ),
2407            TableUpdate::AddEncryptionKey {
2408                encryption_key: EncryptedKey::builder()
2409                    .key_id("a")
2410                    .encrypted_key_metadata(key_bytes.to_vec())
2411                    .encrypted_by_id("b")
2412                    .build(),
2413            },
2414        );
2415    }
2416
2417    #[test]
2418    fn test_remove_encryption_key() {
2419        test_serde_json(
2420            r#"
2421                {
2422                    "action": "remove-encryption-key",
2423                    "key-id": "a"
2424                }        
2425            "#,
2426            TableUpdate::RemoveEncryptionKey {
2427                key_id: "a".to_string(),
2428            },
2429        );
2430    }
2431
2432    #[test]
2433    fn test_table_commit() {
2434        let table = {
2435            let file = File::open(format!(
2436                "{}/testdata/table_metadata/{}",
2437                env!("CARGO_MANIFEST_DIR"),
2438                "TableMetadataV2Valid.json"
2439            ))
2440            .unwrap();
2441            let reader = BufReader::new(file);
2442            let resp = serde_json::from_reader::<_, TableMetadata>(reader).unwrap();
2443
2444            Table::builder()
2445                .metadata(resp)
2446                .metadata_location("s3://bucket/test/location/metadata/00000-8a62c37d-4573-4021-952a-c0baef7d21d0.metadata.json")
2447                .identifier(TableIdent::from_strs(["ns1", "test1"]).unwrap())
2448                .file_io(FileIO::new_with_memory())
2449                .runtime(test_runtime())
2450                .build()
2451                .unwrap()
2452        };
2453
2454        let updates = vec![
2455            TableUpdate::SetLocation {
2456                location: "s3://bucket/test/new_location/data".to_string(),
2457            },
2458            TableUpdate::SetProperties {
2459                updates: vec![
2460                    ("prop1".to_string(), "v1".to_string()),
2461                    ("prop2".to_string(), "v2".to_string()),
2462                ]
2463                .into_iter()
2464                .collect(),
2465            },
2466        ];
2467
2468        let requirements = vec![TableRequirement::UuidMatch {
2469            uuid: table.metadata().table_uuid,
2470        }];
2471
2472        let table_commit = TableCommit::builder()
2473            .ident(table.identifier().to_owned())
2474            .updates(updates)
2475            .requirements(requirements)
2476            .build();
2477
2478        let updated_table = table_commit.apply(table).unwrap();
2479
2480        assert_eq!(
2481            updated_table.metadata().properties.get("prop1").unwrap(),
2482            "v1"
2483        );
2484        assert_eq!(
2485            updated_table.metadata().properties.get("prop2").unwrap(),
2486            "v2"
2487        );
2488
2489        // metadata version should be bumped, and the metadata directory should be
2490        // re-derived from the updated table location
2491        assert!(
2492            updated_table
2493                .metadata_location()
2494                .unwrap()
2495                .starts_with("s3://bucket/test/new_location/data/metadata/00001-")
2496        );
2497
2498        assert_eq!(
2499            updated_table.metadata().location,
2500            "s3://bucket/test/new_location/data",
2501        );
2502    }
2503}