From a6a0e292829a3ec46e4dd3634f6d9a1b1b495f90 Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Wed, 15 Jul 2026 09:32:40 +0800 Subject: [PATCH] fix(table-catalog): support Spark REST commits (#4788) * fix(table-catalog): support Spark REST commits * chore(deps): update s3s SigV4 revision --------- Co-authored-by: Henry Guo --- Cargo.lock | 3 +- Cargo.toml | 2 +- deny.toml | 7 +- rustfs/Cargo.toml | 2 +- rustfs/src/admin/handlers/table_catalog.rs | 236 ++++++++++++++---- rustfs/src/server/http.rs | 24 +- rustfs/src/server/layer.rs | 220 ++++++++++++++++ scripts/table-catalog/engine_compatibility.py | 1 + .../test_engine_compatibility.py | 1 + 9 files changed, 431 insertions(+), 65 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index d8c9fdb57..6286ed958 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -10296,8 +10296,7 @@ checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" [[package]] name = "s3s" version = "0.14.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "abe1bd31748cb69848c2cf4028cecdadf5f8299ef3b02a83e21b20e7bda3aba9" +source = "git+https://github.com/s3s-project/s3s.git?rev=ce69c3f10824535c7c24b2f71cdb2aaa4dffb5e0#ce69c3f10824535c7c24b2f71cdb2aaa4dffb5e0" dependencies = [ "arc-swap", "arrayvec", diff --git a/Cargo.toml b/Cargo.toml index d09d6af44..762967b65 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -283,7 +283,7 @@ redis = { version = "1.3.0", features = ["connection-manager", "tokio-rustls-com rustix = { version = "1.1.4", features = ["fs"] } rust-embed = { version = "8.12.0" } rustc-hash = { version = "2.1.3" } -s3s = { version = "0.14.1", features = ["minio"] } +s3s = { git = "https://github.com/s3s-project/s3s.git", rev = "ce69c3f10824535c7c24b2f71cdb2aaa4dffb5e0", features = ["minio"] } serial_test = "3.5.0" shadow-rs = { version = "2.0.0", default-features = false } siphasher = "1.0.3" diff --git a/deny.toml b/deny.toml index 8699f3d2a..f9fab889b 100644 --- a/deny.toml +++ b/deny.toml @@ -36,9 +36,10 @@ unknown-registry = "deny" unknown-git = "deny" allow-registry = ["https://github.com/rust-lang/crates.io-index"] allow-git = [ - # Custom S3 server library with minio compatibility patches not yet upstreamed. - # Pinned to a specific commit in workspace Cargo.toml. - # "https://github.com/rustfs/s3s", + # Official s3s repository. Temporarily pinned for generic REST SigV4 + # payload-checksum compatibility while upstream PR 631 is reviewed. + # owner: marshawcoco review: 2026-10 + "https://github.com/s3s-project/s3s.git", "https://github.com/apache/datafusion.git", # hyper pinned (via [patch.crates-io]) to a rev carrying the HTTP/1 # flush-before-shutdown fix (hyperium/hyper#4018, commit 72046cc7) that is diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 0d710d10e..81c2d7c16 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -139,6 +139,7 @@ bytes = { workspace = true } chrono = { workspace = true, features = ["serde"] } flatbuffers.workspace = true rmp-serde.workspace = true +quick-xml.workspace = true rustfs-signer.workspace = true serde.workspace = true serde_json.workspace = true @@ -202,7 +203,6 @@ libmimalloc-sys = { version = "0.1.49", features = ["extended"] } [dev-dependencies] uuid = { workspace = true, features = ["v4"] } serial_test = { workspace = true } -quick-xml = { workspace = true } tempfile = { workspace = true } aws-config = { workspace = true } anyhow = { workspace = true } diff --git a/rustfs/src/admin/handlers/table_catalog.rs b/rustfs/src/admin/handlers/table_catalog.rs index 6ea5ae241..ba749c760 100644 --- a/rustfs/src/admin/handlers/table_catalog.rs +++ b/rustfs/src/admin/handlers/table_catalog.rs @@ -36,7 +36,7 @@ use rustfs_policy::{ action::{Action, AdminAction}, }, }; -use s3s::{Body, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error}; +use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error}; use serde::{Deserialize, Serialize, de::DeserializeOwned}; use std::collections::{BTreeMap, BTreeSet, HashMap}; use std::time::{Duration as StdDuration, Instant}; @@ -50,6 +50,16 @@ const DEFAULT_TABLE_CATALOG_CREDENTIAL_TTL_SECONDS: i64 = 15 * 60; const MIN_TABLE_CATALOG_CREDENTIAL_TTL_SECONDS: i64 = 60; const MAX_TABLE_CATALOG_CREDENTIAL_TTL_SECONDS: i64 = 60 * 60; const WAREHOUSE_PROPERTY: &str = "warehouse"; +const PREFIX_PROPERTY: &str = "prefix"; +const ICEBERG_ERROR_ALREADY_EXISTS: &str = "AlreadyExistsException"; +const ICEBERG_ERROR_BAD_REQUEST: &str = "BadRequestException"; +const ICEBERG_ERROR_COMMIT_FAILED: &str = "CommitFailedException"; +const ICEBERG_ERROR_NAMESPACE_NOT_EMPTY: &str = "NamespaceNotEmptyException"; +const ICEBERG_ERROR_NO_SUCH_NAMESPACE: &str = "NoSuchNamespaceException"; +const ICEBERG_ERROR_NO_SUCH_RESOURCE: &str = "NoSuchResourceException"; +const ICEBERG_ERROR_NO_SUCH_TABLE: &str = "NoSuchTableException"; +const ICEBERG_ERROR_NO_SUCH_VIEW: &str = "NoSuchViewException"; +const ICEBERG_ERROR_REST: &str = "RESTException"; const CATALOG_ENDPOINT_PREFIX_CONFIG_KEY: &str = "rustfs.catalog-endpoint-prefix"; const CATALOG_COMPAT_ENDPOINT_PREFIX_CONFIG_KEY: &str = "rustfs.catalog-compat-endpoint-prefix"; const CATALOG_BACKING_CONFIG_KEY: &str = "rustfs.catalog-backing"; @@ -199,8 +209,8 @@ static ROLLBACK_TABLE_CATALOG_HANDLER: RollbackTableCatalogHandler = RollbackTab #[derive(Debug, Serialize)] struct CatalogConfigResponse { - defaults: BTreeMap<&'static str, &'static str>, - overrides: BTreeMap<&'static str, &'static str>, + defaults: BTreeMap, + overrides: BTreeMap, endpoints: Vec<&'static str>, admin_discovery: CatalogAdminDiscovery, } @@ -1031,20 +1041,30 @@ fn register_table_catalog_prefix_routes(r: &mut S3Router, prefix Ok(()) } -fn catalog_config_response() -> S3Result { +fn catalog_config_response(warehouse: Option<&str>) -> S3Result { let usecase = default_admin_usecase(); let backing_mode = crate::table_catalog::TableCatalogBackingMode::from_env().map_err(catalog_store_error)?; let mut overrides = BTreeMap::new(); if backing_mode != crate::table_catalog::TableCatalogBackingMode::ObjectBacked { - overrides.insert(CATALOG_BACKING_CONFIG_KEY, backing_mode.as_str()); + overrides.insert(CATALOG_BACKING_CONFIG_KEY.to_string(), backing_mode.as_str().to_string()); + } + let mut defaults = BTreeMap::from([ + (WAREHOUSE_PROPERTY.to_string(), DEFAULT_WAREHOUSE_ID.to_string()), + (CATALOG_ENDPOINT_PREFIX_CONFIG_KEY.to_string(), TABLE_CATALOG_PREFIX.to_string()), + ( + CATALOG_COMPAT_ENDPOINT_PREFIX_CONFIG_KEY.to_string(), + TABLE_CATALOG_COMPAT_PREFIX.to_string(), + ), + ( + CATALOG_BACKING_CONFIG_KEY.to_string(), + crate::table_catalog::TABLE_CATALOG_BACKING_OBJECT.to_string(), + ), + ]); + if let Some(warehouse) = warehouse { + defaults.insert(PREFIX_PROPERTY.to_string(), warehouse.to_string()); } Ok(CatalogConfigResponse { - defaults: BTreeMap::from([ - (WAREHOUSE_PROPERTY, DEFAULT_WAREHOUSE_ID), - (CATALOG_ENDPOINT_PREFIX_CONFIG_KEY, TABLE_CATALOG_PREFIX), - (CATALOG_COMPAT_ENDPOINT_PREFIX_CONFIG_KEY, TABLE_CATALOG_COMPAT_PREFIX), - (CATALOG_BACKING_CONFIG_KEY, crate::table_catalog::TABLE_CATALOG_BACKING_OBJECT), - ]), + defaults, overrides, endpoints: TABLE_CATALOG_ENDPOINTS.to_vec(), admin_discovery: CatalogAdminDiscovery { @@ -1272,6 +1292,26 @@ fn warehouse_from_params(params: &Params<'_, '_>) -> S3Result { Ok(warehouse.to_string()) } +fn warehouse_from_config_query(uri: &http::Uri) -> S3Result> { + let Some(query) = uri.query() else { + return Ok(None); + }; + let mut warehouse = None; + for (key, value) in url::form_urlencoded::parse(query.as_bytes()) { + if key != WAREHOUSE_PROPERTY { + continue; + } + if warehouse.is_some() { + return Err(s3_error!(InvalidRequest, "warehouse query parameter must not be repeated")); + } + if value.is_empty() { + return Err(s3_error!(InvalidRequest, "warehouse query parameter must not be empty")); + } + warehouse = Some(value.into_owned()); + } + Ok(warehouse) +} + fn namespace_from_params(params: &Params<'_, '_>) -> S3Result { let namespace = params.get("namespace").unwrap_or(""); crate::table_catalog::Namespace::parse(namespace).map_err(|err| s3_error!(InvalidRequest, "invalid namespace: {}", err)) @@ -3332,23 +3372,51 @@ fn namespace_entry_from_create_request( }) } -fn catalog_store_error(err: crate::table_catalog::TableCatalogStoreError) -> s3s::S3Error { +fn iceberg_rest_error(error_type: &str, status: StatusCode, message: impl Into) -> S3Error { + let mut err = S3Error::with_message(S3ErrorCode::Custom(error_type.into()), message.into()); + err.set_status_code(status); + err +} + +fn catalog_store_error(err: crate::table_catalog::TableCatalogStoreError) -> S3Error { match err { crate::table_catalog::TableCatalogStoreError::NotFound(message) => { - s3_error!(InvalidRequest, "{message}") + iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_RESOURCE, StatusCode::NOT_FOUND, message) } crate::table_catalog::TableCatalogStoreError::Conflict(message) => { - s3_error!(PreconditionFailed, "{message}") + iceberg_rest_error(ICEBERG_ERROR_COMMIT_FAILED, StatusCode::CONFLICT, message) } crate::table_catalog::TableCatalogStoreError::Invalid(message) => { - s3_error!(InvalidRequest, "{message}") + iceberg_rest_error(ICEBERG_ERROR_BAD_REQUEST, StatusCode::BAD_REQUEST, message) } crate::table_catalog::TableCatalogStoreError::Internal(message) => { - s3_error!(InternalError, "{message}") + iceberg_rest_error(ICEBERG_ERROR_REST, StatusCode::INTERNAL_SERVER_ERROR, message) } } } +fn catalog_store_conflict_error(err: crate::table_catalog::TableCatalogStoreError, conflict_type: &'static str) -> S3Error { + match err { + crate::table_catalog::TableCatalogStoreError::Conflict(message) => { + iceberg_rest_error(conflict_type, StatusCode::CONFLICT, message) + } + err => catalog_store_error(err), + } +} + +fn catalog_store_already_exists_error(err: crate::table_catalog::TableCatalogStoreError) -> S3Error { + catalog_store_conflict_error(err, ICEBERG_ERROR_ALREADY_EXISTS) +} + +fn catalog_store_namespace_drop_error(err: crate::table_catalog::TableCatalogStoreError) -> S3Error { + match err { + crate::table_catalog::TableCatalogStoreError::NotFound(message) => { + iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_NAMESPACE, StatusCode::NOT_FOUND, message) + } + err => catalog_store_conflict_error(err, ICEBERG_ERROR_NAMESPACE_NOT_EMPTY), + } +} + async fn create_namespace_response( store: &S, bucket: &str, @@ -3360,7 +3428,10 @@ where { let entry = namespace_entry_from_create_request(bucket, request)?; ensure_table_bucket_entry(store, bucket, table_bucket_enabled).await?; - store.create_namespace(entry.clone()).await.map_err(catalog_store_error)?; + store + .create_namespace(entry.clone()) + .await + .map_err(catalog_store_already_exists_error)?; namespace_response_from_entry(entry) } @@ -3385,7 +3456,11 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "namespace not found")); + return Err(iceberg_rest_error( + ICEBERG_ERROR_NO_SUCH_NAMESPACE, + StatusCode::NOT_FOUND, + "namespace not found", + )); }; namespace_response_from_entry(entry) } @@ -3406,7 +3481,10 @@ async fn drop_namespace_in_store(store: &S, bucket: &str, namespace: &str) -> where S: crate::table_catalog::TableCatalogStore + ?Sized, { - store.drop_namespace(bucket, namespace).await.map_err(catalog_store_error) + store + .drop_namespace(bucket, namespace) + .await + .map_err(catalog_store_namespace_drop_error) } async fn register_table_response( @@ -3425,7 +3503,10 @@ where let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?; validate_metadata_table_location_in_bucket(bucket, &metadata)?; adopt_registered_metadata_identity(&mut entry, &metadata)?; - store.register_table(entry.clone()).await.map_err(catalog_store_error)?; + store + .register_table(entry.clone()) + .await + .map_err(catalog_store_already_exists_error)?; Ok(load_table_response_from_entry(entry, metadata)) } @@ -3452,8 +3533,11 @@ where crate::table_catalog::TableCatalogPutPrecondition::IfAbsent, ) .await - .map_err(catalog_store_error)?; - store.create_table(entry.clone()).await.map_err(catalog_store_error)?; + .map_err(catalog_store_already_exists_error)?; + store + .create_table(entry.clone()) + .await + .map_err(catalog_store_already_exists_error)?; Ok(load_table_response_from_entry(entry, metadata)) } @@ -3480,8 +3564,11 @@ where crate::table_catalog::TableCatalogPutPrecondition::IfAbsent, ) .await - .map_err(catalog_store_error)?; - store.create_view(entry.clone()).await.map_err(catalog_store_error)?; + .map_err(catalog_store_already_exists_error)?; + store + .create_view(entry.clone()) + .await + .map_err(catalog_store_already_exists_error)?; Ok(load_view_response_from_entry(entry, metadata)) } @@ -3535,7 +3622,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "table not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_TABLE, StatusCode::NOT_FOUND, "table not found")); }; let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?; Ok(load_table_response_from_entry(entry, metadata)) @@ -3571,7 +3658,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "view not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_VIEW, StatusCode::NOT_FOUND, "view not found")); }; let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?; Ok(load_view_response_from_entry(entry, metadata)) @@ -3610,7 +3697,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "view not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_VIEW, StatusCode::NOT_FOUND, "view not found")); }; let current_metadata = read_table_metadata_json(metadata_backend, bucket, ¤t.metadata_location).await?; validate_view_commit_requirements(¤t_metadata, &request.requirements)?; @@ -3700,7 +3787,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "table not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_TABLE, StatusCode::NOT_FOUND, "table not found")); }; load_credentials_response_from_entry(&entry, issuer, principal).await } @@ -3719,7 +3806,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "table not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_TABLE, StatusCode::NOT_FOUND, "table not found")); }; Ok(table_metadata_location_response_from_entry(entry)) } @@ -3740,7 +3827,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "table not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_TABLE, StatusCode::NOT_FOUND, "table not found")); }; let table_name = crate::table_catalog::IdentifierSegment::parse(table.to_string()) .map_err(|err| s3_error!(InvalidRequest, "invalid table name: {}", err))?; @@ -3790,7 +3877,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "table not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_TABLE, StatusCode::NOT_FOUND, "table not found")); }; let table_name = crate::table_catalog::IdentifierSegment::parse(table.to_string()) .map_err(|err| s3_error!(InvalidRequest, "invalid table name: {}", err))?; @@ -3822,7 +3909,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "table not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_TABLE, StatusCode::NOT_FOUND, "table not found")); }; let table_name = crate::table_catalog::IdentifierSegment::parse(table.to_string()) .map_err(|err| s3_error!(InvalidRequest, "invalid table name: {}", err))?; @@ -4018,7 +4105,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "table not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_TABLE, StatusCode::NOT_FOUND, "table not found")); }; if report.table_id != current.table_id || report.current_metadata_location != current.metadata_location { return Err(s3_error!(PreconditionFailed, "snapshot expiration plan is stale")); @@ -4099,7 +4186,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "table not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_TABLE, StatusCode::NOT_FOUND, "table not found")); }; let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?; let current_snapshot_id = metadata.get("current-snapshot-id").and_then(serde_json::Value::as_i64); @@ -4213,7 +4300,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "table not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_TABLE, StatusCode::NOT_FOUND, "table not found")); }; let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?; let reference = metadata @@ -4662,7 +4749,7 @@ where .await .map_err(catalog_store_error)? else { - return Err(s3_error!(InvalidRequest, "table not found")); + return Err(iceberg_rest_error(ICEBERG_ERROR_NO_SUCH_TABLE, StatusCode::NOT_FOUND, "table not found")); }; let table_name = crate::table_catalog::IdentifierSegment::parse(table.to_string()) .map_err(|err| s3_error!(InvalidRequest, "invalid table name: {}", err))?; @@ -4701,7 +4788,8 @@ pub struct GetCatalogConfigHandler {} impl Operation for GetCatalogConfigHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { authorize_table_catalog_request(&req, AdminAction::GetTableCatalogAction).await?; - build_json_response(StatusCode::OK, &catalog_config_response()?) + let warehouse = warehouse_from_config_query(&req.uri)?; + build_json_response(StatusCode::OK, &catalog_config_response(warehouse.as_deref())?) } } @@ -5587,19 +5675,27 @@ mod tests { #[test] #[serial_test::serial] fn catalog_config_response_lists_standard_rest_endpoints() { - let response = temp_env::with_var_unset(crate::table_catalog::ENV_TABLE_CATALOG_BACKING, catalog_config_response) - .expect("catalog config should build"); + let response = + temp_env::with_var_unset(crate::table_catalog::ENV_TABLE_CATALOG_BACKING, || catalog_config_response(None)) + .expect("catalog config should build"); - assert_eq!(response.defaults.get(WAREHOUSE_PROPERTY), Some(&DEFAULT_WAREHOUSE_ID)); - assert_eq!(response.defaults.get(CATALOG_ENDPOINT_PREFIX_CONFIG_KEY), Some(&TABLE_CATALOG_PREFIX)); + assert_eq!(response.defaults.get(WAREHOUSE_PROPERTY).map(String::as_str), Some(DEFAULT_WAREHOUSE_ID)); assert_eq!( - response.defaults.get(CATALOG_COMPAT_ENDPOINT_PREFIX_CONFIG_KEY), - Some(&TABLE_CATALOG_COMPAT_PREFIX) + response.defaults.get(CATALOG_ENDPOINT_PREFIX_CONFIG_KEY).map(String::as_str), + Some(TABLE_CATALOG_PREFIX) ); assert_eq!( - response.defaults.get(CATALOG_BACKING_CONFIG_KEY), - Some(&crate::table_catalog::TABLE_CATALOG_BACKING_OBJECT) + response + .defaults + .get(CATALOG_COMPAT_ENDPOINT_PREFIX_CONFIG_KEY) + .map(String::as_str), + Some(TABLE_CATALOG_COMPAT_PREFIX) ); + assert_eq!( + response.defaults.get(CATALOG_BACKING_CONFIG_KEY).map(String::as_str), + Some(crate::table_catalog::TABLE_CATALOG_BACKING_OBJECT) + ); + assert!(!response.defaults.contains_key(PREFIX_PROPERTY)); assert!(response.overrides.is_empty()); assert_eq!(response.admin_discovery.runtime_capabilities, "/rustfs/admin/v4/runtime/capabilities"); assert_eq!(response.admin_discovery.cluster_snapshot, "/rustfs/admin/v4/cluster/snapshot"); @@ -5710,16 +5806,60 @@ mod tests { let response = temp_env::with_var( crate::table_catalog::ENV_TABLE_CATALOG_BACKING, Some(crate::table_catalog::TABLE_CATALOG_BACKING_DURABLE_STRONG), - catalog_config_response, + || catalog_config_response(None), ) .expect("catalog config should build"); assert_eq!( - response.overrides.get(CATALOG_BACKING_CONFIG_KEY), - Some(&crate::table_catalog::TABLE_CATALOG_BACKING_DURABLE_STRONG) + response.overrides.get(CATALOG_BACKING_CONFIG_KEY).map(String::as_str), + Some(crate::table_catalog::TABLE_CATALOG_BACKING_DURABLE_STRONG) ); } + #[test] + #[serial_test::serial] + fn catalog_config_response_uses_requested_warehouse_as_standard_prefix() { + let response = temp_env::with_var_unset(crate::table_catalog::ENV_TABLE_CATALOG_BACKING, || { + catalog_config_response(Some("analytics")) + }) + .expect("catalog config should build"); + + assert_eq!(response.defaults.get(PREFIX_PROPERTY).map(String::as_str), Some("analytics")); + } + + #[test] + fn warehouse_config_query_rejects_empty_and_repeated_values() { + let uri = "/iceberg/v1/config?warehouse=analytics".parse().expect("URI"); + assert_eq!(warehouse_from_config_query(&uri).expect("warehouse query"), Some("analytics".to_string())); + + let uri = "/iceberg/v1/config?warehouse=".parse().expect("URI"); + assert!(warehouse_from_config_query(&uri).is_err()); + + let uri = "/iceberg/v1/config?warehouse=one&warehouse=two".parse().expect("URI"); + assert!(warehouse_from_config_query(&uri).is_err()); + } + + #[test] + fn catalog_conflicts_use_operation_specific_iceberg_errors() { + let already_exists = catalog_store_already_exists_error(crate::table_catalog::TableCatalogStoreError::Conflict( + "table already exists".to_string(), + )); + assert_eq!(already_exists.code(), &S3ErrorCode::Custom(ICEBERG_ERROR_ALREADY_EXISTS.into())); + assert_eq!(already_exists.status_code(), Some(StatusCode::CONFLICT)); + + let namespace_not_empty = catalog_store_namespace_drop_error(crate::table_catalog::TableCatalogStoreError::Conflict( + "namespace is not empty".to_string(), + )); + assert_eq!(namespace_not_empty.code(), &S3ErrorCode::Custom(ICEBERG_ERROR_NAMESPACE_NOT_EMPTY.into())); + assert_eq!(namespace_not_empty.status_code(), Some(StatusCode::CONFLICT)); + + let namespace_not_found = catalog_store_namespace_drop_error(crate::table_catalog::TableCatalogStoreError::NotFound( + "namespace not found".to_string(), + )); + assert_eq!(namespace_not_found.code(), &S3ErrorCode::Custom(ICEBERG_ERROR_NO_SUCH_NAMESPACE.into())); + assert_eq!(namespace_not_found.status_code(), Some(StatusCode::NOT_FOUND)); + } + #[test] fn table_catalog_admin_operation_result_labels_are_stable() { let success: Result<(), ()> = Ok(()); diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index 84871073f..9a7ff13d9 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -23,8 +23,9 @@ use crate::server::{ hybrid::hybrid, layer::{ BodylessStatusFixLayer, ConditionalCorsLayer, DoubleSlashListBucketsCompatLayer, EmptyBodyContentLengthCompatLayer, - HeadRequestBodyFixLayer, ObjectAttributesEtagFixLayer, PublicHealthEndpointLayer, RedirectLayer, RequestContextLayer, - RequestLoggingLayer, S3ErrorMessageCompatLayer, VirtualHostStyleHintLayer, redact_sensitive_uri_query, + HeadRequestBodyFixLayer, IcebergRestErrorCompatLayer, ObjectAttributesEtagFixLayer, PublicHealthEndpointLayer, + RedirectLayer, RequestContextLayer, RequestLoggingLayer, S3ErrorMessageCompatLayer, VirtualHostStyleHintLayer, + redact_sensitive_uri_query, }, tls_material::{ TlsAcceptorHolder, TlsHandshakeFailureKind, build_acceptor_from_loaded, load_tls_material, spawn_reload_loop, @@ -1190,14 +1191,15 @@ fn process_connection( // 13. CompressionLayer — response compression (whitelist, path-aware) // 14. PathCategoryInjectionLayer — injects path category for compression predicate // 15. S3ErrorMessageCompatLayer — missing S3 error message compatibility - // 16. ObjectAttributesEtagFixLayer — ETag fix for GetObjectAttributes - // 17. ConditionalCorsLayer — S3 API CORS - // 18. RedirectLayer — console redirect (conditional) - // 19. BodylessStatusFixLayer — clears body for 1xx/204/205/304 responses - // 20. HeadRequestBodyFixLayer — strips actual body bytes from HEAD responses - // 21. PublicHealthEndpointLayer — handles public health before s3s host parsing - // 22. VirtualHostStyleHintLayer — actionable error for unroutable virtual-hosted-style (conditional) - // 23. DoubleSlashListBucketsCompatLayer — rewrites `GET //` to `GET /` for ListBuckets (MinIO browser compat) + // 16. IcebergRestErrorCompatLayer — Iceberg REST JSON error compatibility + // 17. ObjectAttributesEtagFixLayer — ETag fix for GetObjectAttributes + // 18. ConditionalCorsLayer — S3 API CORS + // 19. RedirectLayer — console redirect (conditional) + // 20. BodylessStatusFixLayer — clears body for 1xx/204/205/304 responses + // 21. HeadRequestBodyFixLayer — strips actual body bytes from HEAD responses + // 22. PublicHealthEndpointLayer — handles public health before s3s host parsing + // 23. VirtualHostStyleHintLayer — actionable error for unroutable virtual-hosted-style (conditional) + // 24. DoubleSlashListBucketsCompatLayer — rewrites `GET //` to `GET /` for ListBuckets (MinIO browser compat) // ───────────────────────────────────────────────────────────── // Batch 1 intentionally keeps the external and internode stacks behaviorally // identical while giving each path family a named construction boundary. @@ -1368,6 +1370,7 @@ fn process_connection( .layer(CompressionLayer::new().compress_when(PathAwareHttpCompressionPredicate::new(compression_config.clone()))) .layer(PathCategoryInjectionLayer) .layer(S3ErrorMessageCompatLayer) + .layer(IcebergRestErrorCompatLayer) .layer(ObjectAttributesEtagFixLayer) .layer(ConditionalCorsLayer::new()) .option_layer(if is_console { Some(RedirectLayer) } else { None }) @@ -1527,6 +1530,7 @@ fn process_connection( .layer(CompressionLayer::new().compress_when(PathAwareHttpCompressionPredicate::new(compression_config.clone()))) .layer(PathCategoryInjectionLayer) .layer(S3ErrorMessageCompatLayer) + .layer(IcebergRestErrorCompatLayer) .layer(ObjectAttributesEtagFixLayer) .layer(ConditionalCorsLayer::new()) .option_layer(if is_console { Some(RedirectLayer) } else { None }) diff --git a/rustfs/src/server/layer.rs b/rustfs/src/server/layer.rs index 017a18e54..0c1d0aaec 100644 --- a/rustfs/src/server/layer.rs +++ b/rustfs/src/server/layer.rs @@ -37,6 +37,7 @@ use rustfs_trusted_proxies::ClientInfo; use rustfs_utils::get_env_opt_str; use rustfs_utils::http::headers::AMZ_REQUEST_ID; use s3s::S3ErrorCode; +use serde::{Deserialize, Serialize}; use std::future::Future; use std::pin::Pin; use std::sync::Arc; @@ -628,6 +629,165 @@ where } } +#[derive(Clone)] +pub struct IcebergRestErrorCompatLayer; + +impl Layer for IcebergRestErrorCompatLayer { + type Service = IcebergRestErrorCompatService; + + fn layer(&self, inner: S) -> Self::Service { + IcebergRestErrorCompatService { inner } + } +} + +#[derive(Clone)] +pub struct IcebergRestErrorCompatService { + inner: S, +} + +impl Service> for IcebergRestErrorCompatService +where + S: Service, Response = Response>> + Clone + Send + 'static, + S::Future: Send + 'static, + S::Error: Send + 'static, + RestBody: Body + From + Send + 'static, + RestBody::Error: Into + Send + 'static, + GrpcBody: Send + 'static, +{ + type Response = Response>; + type Error = S::Error; + type Future = Pin> + Send>>; + + fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll> { + self.inner.poll_ready(cx) + } + + fn call(&mut self, req: HttpRequest) -> Self::Future { + let catalog_path = + (req.method() != Method::HEAD && is_table_catalog_path(req.uri().path())).then(|| req.uri().path().to_string()); + let mut inner = self.inner.clone(); + + Box::pin(async move { + let response = inner.call(req).await?; + let (parts, body) = response.into_parts(); + let should_convert = catalog_path.is_some() && !parts.status.is_success() && is_xml_response(&parts.headers); + + let response = match body { + HybridBody::Rest { rest_body } if should_convert => { + let (rest_body, converted_status) = convert_iceberg_error_in_xml( + rest_body, + parts.status, + catalog_path.as_deref().expect("catalog path was checked"), + ) + .await + .map_err(Into::into)?; + let mut parts = parts; + if let Some(status) = converted_status { + parts.status = status; + parts.headers.remove(http::header::CONTENT_LENGTH); + parts + .headers + .insert(http::header::CONTENT_TYPE, HeaderValue::from_static("application/json")); + } + Response::from_parts(parts, HybridBody::Rest { rest_body }) + } + HybridBody::Rest { rest_body } => Response::from_parts(parts, HybridBody::Rest { rest_body }), + HybridBody::Grpc { grpc_body } => Response::from_parts(parts, HybridBody::Grpc { grpc_body }), + }; + + Ok(response) + }) + } +} + +#[derive(Debug, Deserialize)] +struct SerializedS3Error { + #[serde(rename = "Code")] + code: String, + #[serde(default, rename = "Message")] + message: Option, +} + +#[derive(Debug, Serialize)] +struct IcebergRestErrorEnvelope { + error: IcebergRestError, +} + +#[derive(Debug, Serialize)] +struct IcebergRestError { + message: String, + #[serde(rename = "type")] + error_type: String, + code: u16, +} + +async fn convert_iceberg_error_in_xml( + body: RestBody, + status: StatusCode, + path: &str, +) -> Result<(RestBody, Option), RestBody::Error> +where + RestBody: Body + From, +{ + let bytes = BodyExt::collect(body).await?.to_bytes(); + let Some(parsed) = std::str::from_utf8(&bytes) + .ok() + .and_then(|xml| quick_xml::de::from_str::(xml).ok()) + else { + return Ok((RestBody::from(bytes), None)); + }; + let status = iceberg_rest_status(status, &parsed.code); + let envelope = IcebergRestErrorEnvelope { + error: IcebergRestError { + message: parsed.message.unwrap_or_else(|| parsed.code.clone()), + error_type: iceberg_rest_error_type(&parsed.code, status, path), + code: status.as_u16(), + }, + }; + let Ok(json) = serde_json::to_vec(&envelope) else { + return Ok((RestBody::from(bytes), None)); + }; + Ok((RestBody::from(Bytes::from(json)), Some(status))) +} + +fn iceberg_rest_status(status: StatusCode, error_code: &str) -> StatusCode { + if status == StatusCode::PRECONDITION_FAILED || error_code == "PreconditionFailed" { + StatusCode::CONFLICT + } else { + status + } +} + +fn iceberg_rest_error_type(error_code: &str, status: StatusCode, path: &str) -> String { + if error_code.ends_with("Exception") { + return error_code.to_string(); + } + match error_code { + "AccessDenied" | "SignatureDoesNotMatch" => "ForbiddenException".to_string(), + "InvalidAccessKeyId" => "NotAuthorizedException".to_string(), + "InvalidArgument" | "InvalidRequest" => "BadRequestException".to_string(), + "PreconditionFailed" => "CommitFailedException".to_string(), + _ => match status { + StatusCode::UNAUTHORIZED => "NotAuthorizedException".to_string(), + StatusCode::FORBIDDEN => "ForbiddenException".to_string(), + StatusCode::NOT_FOUND => iceberg_not_found_error_type(path).to_string(), + StatusCode::CONFLICT => "CommitFailedException".to_string(), + status if status.is_server_error() => "RESTException".to_string(), + _ => "BadRequestException".to_string(), + }, + } +} + +fn iceberg_not_found_error_type(path: &str) -> &'static str { + if path.contains("/views/") { + "NoSuchViewException" + } else if path.contains("/tables/") { + "NoSuchTableException" + } else { + "NoSuchNamespaceException" + } +} + #[derive(Clone)] pub struct ObjectAttributesEtagFixLayer; @@ -2724,6 +2884,66 @@ mod tests { assert_eq!(bytes, input); } + #[tokio::test] + async fn iceberg_rest_error_conversion_returns_standard_json_envelope() { + let body = Full::from(Bytes::from_static( + b"NoSuchTableExceptiontable not found", + )); + + let (body, status) = + convert_iceberg_error_in_xml(body, StatusCode::NOT_FOUND, "/iceberg/v1/warehouse/namespaces/ns/tables/events") + .await + .expect("convert Iceberg error"); + let value: serde_json::Value = + serde_json::from_slice(&BodyExt::collect(body).await.expect("collect converted body").to_bytes()) + .expect("JSON error envelope"); + + assert_eq!(status, Some(StatusCode::NOT_FOUND)); + assert_eq!(value["error"]["code"], 404); + assert_eq!(value["error"]["type"], "NoSuchTableException"); + assert_eq!(value["error"]["message"], "table not found"); + } + + #[tokio::test] + async fn iceberg_rest_error_conversion_maps_precondition_to_commit_conflict() { + let body = Full::from(Bytes::from_static( + b"PreconditionFailedversion token changed", + )); + + let (body, status) = convert_iceberg_error_in_xml( + body, + StatusCode::PRECONDITION_FAILED, + "/_iceberg/v1/warehouse/namespaces/ns/tables/events", + ) + .await + .expect("convert Iceberg conflict"); + let value: serde_json::Value = + serde_json::from_slice(&BodyExt::collect(body).await.expect("collect converted body").to_bytes()) + .expect("JSON error envelope"); + + assert_eq!(status, Some(StatusCode::CONFLICT)); + assert_eq!(value["error"]["code"], 409); + assert_eq!(value["error"]["type"], "CommitFailedException"); + } + + #[tokio::test] + async fn iceberg_rest_error_conversion_preserves_forbidden_status() { + let body = Full::from(Bytes::from_static( + b"SignatureDoesNotMatchsignature mismatch", + )); + + let (body, status) = convert_iceberg_error_in_xml(body, StatusCode::FORBIDDEN, "/iceberg/v1/config") + .await + .expect("convert Iceberg auth error"); + let value: serde_json::Value = + serde_json::from_slice(&BodyExt::collect(body).await.expect("collect converted body").to_bytes()) + .expect("JSON error envelope"); + + assert_eq!(status, Some(StatusCode::FORBIDDEN)); + assert_eq!(value["error"]["code"], 403); + assert_eq!(value["error"]["type"], "ForbiddenException"); + } + #[test] fn test_insert_missing_signature_error_message() { let (fixed, changed) = diff --git a/scripts/table-catalog/engine_compatibility.py b/scripts/table-catalog/engine_compatibility.py index b048653c3..9f7a241fd 100644 --- a/scripts/table-catalog/engine_compatibility.py +++ b/scripts/table-catalog/engine_compatibility.py @@ -257,6 +257,7 @@ def spark_catalog_config( (f"{prefix}.type", "rest"), (f"{prefix}.uri", f"{endpoint}{rest_path}"), (f"{prefix}.warehouse", warehouse), + (f"{prefix}.prefix", warehouse), (f"{prefix}.io-impl", "org.apache.iceberg.aws.s3.S3FileIO"), (f"{prefix}.s3.endpoint", endpoint), (f"{prefix}.s3.path-style-access", "true"), diff --git a/scripts/table-catalog/test_engine_compatibility.py b/scripts/table-catalog/test_engine_compatibility.py index 2a8d665e6..2636d1a21 100644 --- a/scripts/table-catalog/test_engine_compatibility.py +++ b/scripts/table-catalog/test_engine_compatibility.py @@ -57,6 +57,7 @@ class EngineCompatibilityTest(unittest.TestCase): self.assertEqual(config["spark.sql.catalog.rustfs.type"], "rest") self.assertEqual(config["spark.sql.catalog.rustfs.uri"], "http://127.0.0.1:9000/iceberg") self.assertEqual(config["spark.sql.catalog.rustfs.warehouse"], "rustfs-s3table-smoke") + self.assertEqual(config["spark.sql.catalog.rustfs.prefix"], "rustfs-s3table-smoke") self.assertEqual(config["spark.sql.catalog.rustfs.io-impl"], "org.apache.iceberg.aws.s3.S3FileIO") self.assertEqual(config["spark.sql.catalog.rustfs.s3.endpoint"], "http://127.0.0.1:9000") self.assertEqual(config["spark.sql.catalog.rustfs.rest.signing-name"], "s3")