Skip to content
Open
Changes from 2 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
103 changes: 102 additions & 1 deletion crates/catalog/rest/src/catalog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ use std::sync::{Arc, OnceLock};
use async_trait::async_trait;
use iceberg::encryption::kms::{KeyManagementClient, KmsClientFactory};
use iceberg::io::{FileIO, FileIOBuilder, StorageFactory};
use iceberg::spec::TableProperties;
use iceberg::table::Table;
use iceberg::{
Catalog, CatalogBuilder, Error, ErrorKind, Namespace, NamespaceIdent, Result, Runtime,
Expand Down Expand Up @@ -1133,6 +1134,22 @@ impl SessionCatalog for RestSessionCatalog {

let table_ident = TableIdent::new(namespace.clone(), creation.name.clone());

let mut properties = creation.properties;

if properties.contains_key(TableProperties::PROPERTY_FORMAT_VERSION) {
return Err(Error::new(
ErrorKind::DataInvalid,
format!(
"Table properties should not contain reserved properties, but got: [{}]. Set `TableCreation::format_version` instead",
TableProperties::PROPERTY_FORMAT_VERSION
),
));
}
Comment on lines +1139 to +1147

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a break technically but I think it's a legitimate break that we can document in the upgrade docs for 0.11.0

We could also allow the property IF it matches the format_version provided?

properties.insert(
TableProperties::PROPERTY_FORMAT_VERSION.to_string(),
(creation.format_version as u8).to_string(),
);

let request = HttpRequest::build(
client
.http_client
Expand All @@ -1144,7 +1161,7 @@ impl SessionCatalog for RestSessionCatalog {
partition_spec: creation.partition_spec,
write_order: creation.sort_order,
stage_create: Some(false),
properties: creation.properties,
properties,
}),
)?;

Expand Down Expand Up @@ -3928,6 +3945,90 @@ mod tests {
create_table_mock.assert_async().await;
}

fn single_column_schema() -> Schema {
Schema::builder()
.with_fields(vec![
NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
])
.build()
.unwrap()
}

#[tokio::test]
async fn test_create_table_sends_the_requested_format_version_as_a_property() {
for (format_version, expected) in [
(None, r#""format-version":"2""#),
(Some(FormatVersion::V1), r#""format-version":"1""#),
(Some(FormatVersion::V3), r#""format-version":"3""#),
] {
let mut server = Server::new_async().await;
let config_mock = create_config_mock(&mut server).await;
let create_table_mock = server
.mock("POST", "/v1/namespaces/ns1/tables")
.match_body(mockito::Matcher::Regex(expected.to_string()))
.with_status(200)
.with_body_from_file(format!(
"{}/testdata/{}",
env!("CARGO_MANIFEST_DIR"),
"create_table_response.json"
))
.create_async()
.await;

let mut creation = TableCreation::builder()
.name("test1".to_string())
.schema(single_column_schema())
.build();
if let Some(format_version) = format_version {
creation.format_version = format_version;
}

let catalog = session_catalog(RestCatalogConfig::builder().uri(server.url()).build());
catalog
.create_table(
&SessionContext::empty(),
&NamespaceIdent::from_strs(["ns1"]).unwrap(),
creation,
)
.await
.unwrap();

config_mock.assert_async().await;
create_table_mock.assert_async().await;
}
}

#[tokio::test]
async fn test_create_table_rejects_a_format_version_property() {
let mut server = Server::new_async().await;
let config_mock = create_config_mock(&mut server).await;

let catalog = session_catalog(RestCatalogConfig::builder().uri(server.url()).build());
let error = catalog
.create_table(
&SessionContext::empty(),
&NamespaceIdent::from_strs(["ns1"]).unwrap(),
TableCreation::builder()
.name("test1".to_string())
.schema(single_column_schema())
.properties(HashMap::from([(
"format-version".to_string(),
"3".to_string(),
)]))
.build(),
)
.await
.unwrap_err();

assert_eq!(ErrorKind::DataInvalid, error.kind());
assert!(
error.message().contains("format-version"),
"unexpected message: {}",
error.message()
);
config_mock.assert_async().await;
}

#[tokio::test]
async fn test_create_table_409() {
let mut server = Server::new_async().await;
Expand Down
Loading