Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
77 changes: 50 additions & 27 deletions crates/catalog/loader/tests/schema_update_suite.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ use std::collections::HashMap;

use common::{CatalogKind, cleanup_namespace_dyn, load_catalog};
use iceberg::spec::{NestedField, PrimitiveType, Schema, StructType, Type};
use iceberg::transaction::{AddColumn, ApplyTransactionAction, Transaction};
use iceberg::transaction::{AddColumn, DeleteColumn, Transaction};
use iceberg::{ErrorKind, NamespaceIdent, Result, TableCreation, TableIdent};
use iceberg_test_utils::normalize_test_name_with_parts;
use rstest::rstest;
Expand All @@ -52,6 +52,8 @@ fn base_schema() -> Schema {
#[case::memory_catalog(CatalogKind::Memory)]
#[tokio::test]
async fn test_catalog_schema_add_column(#[case] kind: CatalogKind) -> Result<()> {
use iceberg::transaction::ApplyTransactionAction;

let Some(harness) = load_catalog(kind).await else {
return Ok(());
};
Expand All @@ -78,11 +80,13 @@ async fn test_catalog_schema_add_column(#[case] kind: CatalogKind) -> Result<()>

let tx = Transaction::new(&table);
let tx = tx
.update_schema()
.add_column(AddColumn::optional(
"a",
Type::Primitive(PrimitiveType::Int),
))
.update_schema()?
.add(
AddColumn::builder()
.name("a")
.r#type(PrimitiveType::Int.into())
.build(),
)?
.apply(tx)?;
let updated = tx.commit(catalog.as_ref()).await?;

Expand All @@ -105,6 +109,8 @@ async fn test_catalog_schema_add_column(#[case] kind: CatalogKind) -> Result<()>
#[case::memory_catalog(CatalogKind::Memory)]
#[tokio::test]
async fn test_catalog_schema_add_nested_and_delete_column(#[case] kind: CatalogKind) -> Result<()> {
use iceberg::transaction::ApplyTransactionAction;

let Some(harness) = load_catalog(kind).await else {
return Ok(());
};
Expand Down Expand Up @@ -135,28 +141,33 @@ async fn test_catalog_schema_add_nested_and_delete_column(#[case] kind: CatalogK
// First transaction: add a nested struct column.
let tx = Transaction::new(&table);
let tx = tx
.update_schema()
.add_column(AddColumn::optional(
"info",
Type::Struct(StructType::new(vec![
NestedField::optional(0, "city", Type::Primitive(PrimitiveType::String)).into(),
])),
))
.update_schema()?
.add(
AddColumn::builder()
.name("info")
.r#type(
StructType::new(vec![
NestedField::optional(0, "city", PrimitiveType::String.into()).into(),
])
.into(),
)
.build(),
)?
.apply(tx)?;
let table = tx.commit(catalog.as_ref()).await?;

// Second transaction: add a sub-field to the nested struct and delete a top-level column.
let tx = Transaction::new(&table);
let tx = tx
.update_schema()
.add_column(
.update_schema()?
.add(
AddColumn::builder()
.name("zip")
.field_type(Type::Primitive(PrimitiveType::String))
.parent("info")
.r#type(PrimitiveType::String.into())
.parent(Some("info".into()))
.build(),
)
.delete_column("baz")
)?
.delete(DeleteColumn::new("baz"))?
.apply(tx)?;
let table = tx.commit(catalog.as_ref()).await?;

Expand All @@ -179,6 +190,8 @@ async fn test_catalog_schema_add_nested_and_delete_column(#[case] kind: CatalogK
#[case::memory_catalog(CatalogKind::Memory)]
#[tokio::test]
async fn test_catalog_schema_delete_invalid_column_errors(#[case] kind: CatalogKind) -> Result<()> {
use iceberg::transaction::ApplyTransactionAction;

let Some(harness) = load_catalog(kind).await else {
return Ok(());
};
Expand Down Expand Up @@ -208,14 +221,20 @@ async fn test_catalog_schema_delete_invalid_column_errors(#[case] kind: CatalogK

// Deleting an identifier field must fail.
let tx = Transaction::new(&table);
let tx = tx.update_schema().delete_column("bar").apply(tx)?;
let tx = tx
.update_schema()?
.delete(DeleteColumn::new("bar"))?
.apply(tx)?;
let err = tx.commit(catalog.as_ref()).await.unwrap_err();
assert_eq!(err.kind(), ErrorKind::PreconditionFailed);

// Deleting a nonexistent field must fail.
let tx = Transaction::new(&table);
let tx = tx.update_schema().delete_column("nonexistent").apply(tx)?;
let err = tx.commit(catalog.as_ref()).await.unwrap_err();
let err = tx
.update_schema()?
.delete(DeleteColumn::new("nonexistent"))
.err()
.unwrap();
assert_eq!(err.kind(), ErrorKind::PreconditionFailed);

Ok(())
Expand All @@ -233,6 +252,8 @@ async fn test_catalog_schema_delete_invalid_column_errors(#[case] kind: CatalogK
async fn test_catalog_schema_update_persisted_after_reload(
#[case] kind: CatalogKind,
) -> Result<()> {
use iceberg::transaction::ApplyTransactionAction;

let Some(harness) = load_catalog(kind).await else {
return Ok(());
};
Expand Down Expand Up @@ -263,11 +284,13 @@ async fn test_catalog_schema_update_persisted_after_reload(

let tx = Transaction::new(&table);
let tx = tx
.update_schema()
.add_column(AddColumn::optional(
"new_field",
Type::Primitive(PrimitiveType::Long),
))
.update_schema()?
.add(
AddColumn::builder()
.name("new_field")
.r#type(PrimitiveType::Long.into())
.build(),
)?
.apply(tx)?;
tx.commit(catalog.as_ref()).await?;

Expand Down
Loading
Loading