From cd96491b39a86d25e56344f333dbaa9e30c0af6b Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Tue, 16 Jun 2026 21:19:27 +0800 Subject: [PATCH] feat(table-catalog): add refs and views semantics (#3502) Co-authored-by: Henry Guo --- rustfs/src/admin/handlers/table_catalog.rs | 1162 ++++++++++++++++++- rustfs/src/admin/route_policy.rs | 53 +- rustfs/src/admin/route_registration_test.rs | 60 + rustfs/src/table_catalog.rs | 380 ++++++ scripts/table-catalog/README.md | 10 +- scripts/table-catalog/pyiceberg_smoke.py | 6 - 6 files changed, 1624 insertions(+), 47 deletions(-) diff --git a/rustfs/src/admin/handlers/table_catalog.rs b/rustfs/src/admin/handlers/table_catalog.rs index 05758fe53..58fcf40db 100644 --- a/rustfs/src/admin/handlers/table_catalog.rs +++ b/rustfs/src/admin/handlers/table_catalog.rs @@ -71,6 +71,7 @@ const S3_SECRET_ACCESS_KEY_CONFIG_KEY: &str = "s3.secret-access-key"; const S3_SESSION_TOKEN_CONFIG_KEY: &str = "s3.session-token"; const TABLE_CATALOG_NAMESPACE_RESOURCE_ROOT: &str = "namespaces"; const TABLE_CATALOG_TABLE_RESOURCE_ROOT: &str = "tables"; +const TABLE_CATALOG_VIEW_RESOURCE_ROOT: &str = "views"; const TABLE_CATALOG_ADMIN_OPERATION_SLOW_LOG_THRESHOLD: StdDuration = StdDuration::from_secs(2); const DEFAULT_TABLE_MAINTENANCE_WORKER_ID: &str = "rustfs-maintenance-worker"; const TABLE_CATALOG_ENDPOINTS: &[&str] = &[ @@ -105,9 +106,12 @@ const TABLE_CATALOG_ENDPOINTS: &[&str] = &[ "POST /{warehouse}/namespaces/{namespace}/tables/{table}", "DELETE /{warehouse}/namespaces/{namespace}/tables/{table}", "GET /{warehouse}/namespaces/{namespace}/views/{view}", + "HEAD /{warehouse}/namespaces/{namespace}/views/{view}", "POST /{warehouse}/namespaces/{namespace}/views/{view}", "DELETE /{warehouse}/namespaces/{namespace}/views/{view}", "GET /{warehouse}/namespaces/{namespace}/tables/{table}/refs", + "PUT /{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + "DELETE /{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", "POST /{warehouse}/namespaces/{namespace}/tables/{table}/maintenance/metadata", "GET /{warehouse}/namespaces/{namespace}/tables/{table}/metadata-location", "PUT /{warehouse}/namespaces/{namespace}/tables/{table}/metadata-location", @@ -143,9 +147,12 @@ static LOAD_CREDENTIALS_HANDLER: RestLoadCredentialsHandler = RestLoadCredential static COMMIT_TABLE_HANDLER: RestCommitTableHandler = RestCommitTableHandler {}; static DROP_TABLE_HANDLER: RestDropTableHandler = RestDropTableHandler {}; static LOAD_VIEW_HANDLER: RestLoadViewHandler = RestLoadViewHandler {}; +static VIEW_EXISTS_HANDLER: RestViewExistsHandler = RestViewExistsHandler {}; static REPLACE_VIEW_HANDLER: RestReplaceViewHandler = RestReplaceViewHandler {}; static DROP_VIEW_HANDLER: RestDropViewHandler = RestDropViewHandler {}; static LIST_TABLE_REFS_HANDLER: ListTableRefsHandler = ListTableRefsHandler {}; +static PUT_TABLE_REF_HANDLER: PutTableRefHandler = PutTableRefHandler {}; +static DELETE_TABLE_REF_HANDLER: DeleteTableRefHandler = DeleteTableRefHandler {}; static GET_TABLE_METADATA_LOCATION_HANDLER: GetTableMetadataLocationHandler = GetTableMetadataLocationHandler {}; static UPDATE_TABLE_METADATA_LOCATION_HANDLER: UpdateTableMetadataLocationHandler = UpdateTableMetadataLocationHandler {}; static TABLE_METADATA_MAINTENANCE_HANDLER: RestTableMetadataMaintenanceHandler = RestTableMetadataMaintenanceHandler {}; @@ -203,6 +210,19 @@ struct CreateTableRequest { properties: BTreeMap, } +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct CreateViewRequest { + name: String, + #[serde(default)] + location: Option, + schema: serde_json::Value, + #[serde(rename = "view-version")] + view_version: serde_json::Value, + #[serde(default)] + properties: BTreeMap, +} + #[derive(Debug, Deserialize)] #[serde(deny_unknown_fields)] struct RestCommitTableRequest { @@ -228,6 +248,61 @@ struct RestCommitTableRequest { writer: Option, } +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct RestCommitViewRequest { + #[serde(default, rename = "commit-id")] + commit_id: Option, + #[serde(default, rename = "expected-version-token")] + expected_version_token: Option, + #[serde(default, rename = "expected-metadata-location")] + expected_metadata_location: Option, + #[serde(default, rename = "new-metadata-location")] + new_metadata_location: Option, + #[serde(default)] + requirements: Vec, + #[serde(default)] + updates: Vec, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct PutTableRefRequest { + #[serde(rename = "snapshot-id")] + snapshot_id: i64, + #[serde(rename = "type")] + ref_type: String, + #[serde(default, rename = "expected-snapshot-id")] + expected_snapshot_id: Option, + #[serde(default, rename = "min-snapshots-to-keep")] + min_snapshots_to_keep: Option, + #[serde(default, rename = "max-snapshot-age-ms")] + max_snapshot_age_ms: Option, + #[serde(default, rename = "max-ref-age-ms")] + max_ref_age_ms: Option, + #[serde(default, rename = "commit-id")] + commit_id: Option, + #[serde(default, rename = "idempotency-key")] + idempotency_key: Option, + #[serde(default)] + writer: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct DeleteTableRefRequest { + #[serde(default, rename = "expected-snapshot-id")] + expected_snapshot_id: Option, + #[serde(default)] + force: bool, + #[serde(default, rename = "commit-id")] + commit_id: Option, + #[serde(default, rename = "idempotency-key")] + idempotency_key: Option, + #[serde(default)] + writer: Option, +} + #[derive(Debug, Deserialize)] #[serde(deny_unknown_fields)] struct TableMetadataMaintenanceRequest { @@ -302,15 +377,6 @@ struct RollbackTableRequest { idempotency_key: Option, } -#[derive(Debug, Serialize)] -#[serde(rename_all = "kebab-case")] -struct UnsupportedCatalogFeatureResponse { - feature: &'static str, - status: &'static str, - reason: &'static str, - replacement: &'static str, -} - #[derive(Debug, Serialize)] #[serde(rename_all = "kebab-case")] struct TableRefsResponse { @@ -390,6 +456,19 @@ struct RestListTablesResponse { identifiers: Vec, } +#[derive(Debug, Serialize)] +struct RestListViewsResponse { + identifiers: Vec, +} + +#[derive(Debug, Serialize)] +struct RestLoadViewResponse { + #[serde(rename = "metadata-location")] + metadata_location: String, + metadata: serde_json::Value, + config: BTreeMap, +} + #[derive(Debug, Serialize)] struct RestStorageCredential { prefix: String, @@ -671,6 +750,11 @@ fn register_table_catalog_prefix_routes(r: &mut S3Router, prefix format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}/views/{{view}}").as_str(), AdminOperation(&LOAD_VIEW_HANDLER), )?; + r.insert( + Method::HEAD, + format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}/views/{{view}}").as_str(), + AdminOperation(&VIEW_EXISTS_HANDLER), + )?; r.insert( Method::POST, format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}/views/{{view}}").as_str(), @@ -686,6 +770,16 @@ fn register_table_catalog_prefix_routes(r: &mut S3Router, prefix format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}/tables/{{table}}/refs").as_str(), AdminOperation(&LIST_TABLE_REFS_HANDLER), )?; + r.insert( + Method::PUT, + format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}/tables/{{table}}/refs/{{ref}}").as_str(), + AdminOperation(&PUT_TABLE_REF_HANDLER), + )?; + r.insert( + Method::DELETE, + format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}/tables/{{table}}/refs/{{ref}}").as_str(), + AdminOperation(&DELETE_TABLE_REF_HANDLER), + )?; r.insert( Method::GET, format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}/tables/{{table}}/metadata-location").as_str(), @@ -783,21 +877,6 @@ fn empty_response(status: StatusCode) -> S3Response<(StatusCode, Body)> { S3Response::new((status, Body::default())) } -fn unsupported_catalog_feature_response( - feature: &'static str, - replacement: &'static str, -) -> S3Result> { - build_json_response( - StatusCode::NOT_IMPLEMENTED, - &UnsupportedCatalogFeatureResponse { - feature, - status: "unsupported", - reason: "the catalog advertises this advanced Iceberg surface explicitly but does not implement it yet", - replacement, - }, - ) -} - fn duration_millis_u64(duration: StdDuration) -> u64 { u64::try_from(duration.as_millis()).unwrap_or(u64::MAX) } @@ -883,6 +962,7 @@ struct TableCatalogResource<'a> { warehouse: &'a str, namespace: Option, table: Option, + view: Option, } impl<'a> TableCatalogResource<'a> { @@ -891,6 +971,7 @@ impl<'a> TableCatalogResource<'a> { warehouse, namespace: None, table: None, + view: None, } } @@ -899,6 +980,7 @@ impl<'a> TableCatalogResource<'a> { warehouse, namespace: Some(namespace.storage_id()), table: None, + view: None, } } @@ -907,16 +989,29 @@ impl<'a> TableCatalogResource<'a> { warehouse, namespace: Some(namespace.storage_id()), table: Some(table.to_string()), + view: None, + } + } + + fn view(warehouse: &'a str, namespace: &crate::table_catalog::Namespace, view: &str) -> Self { + Self { + warehouse, + namespace: Some(namespace.storage_id()), + table: None, + view: Some(view.to_string()), } } fn object_path(&self) -> Option { - match (&self.namespace, &self.table) { - (Some(namespace), Some(table)) => Some(format!( + match (&self.namespace, &self.table, &self.view) { + (Some(namespace), Some(table), None) => Some(format!( "{TABLE_CATALOG_NAMESPACE_RESOURCE_ROOT}/{namespace}/{TABLE_CATALOG_TABLE_RESOURCE_ROOT}/{table}" )), - (Some(namespace), None) => Some(format!("{TABLE_CATALOG_NAMESPACE_RESOURCE_ROOT}/{namespace}")), - (None, _) => None, + (Some(namespace), None, Some(view)) => Some(format!( + "{TABLE_CATALOG_NAMESPACE_RESOURCE_ROOT}/{namespace}/{TABLE_CATALOG_VIEW_RESOURCE_ROOT}/{view}" + )), + (Some(namespace), None, None) => Some(format!("{TABLE_CATALOG_NAMESPACE_RESOURCE_ROOT}/{namespace}")), + _ => None, } } } @@ -993,6 +1088,13 @@ fn view_name_from_params(params: &Params<'_, '_>) -> S3Result { Ok(view.to_string()) } +fn ref_name_from_params(params: &Params<'_, '_>) -> S3Result { + let ref_name = params.get("ref").unwrap_or(""); + crate::table_catalog::IdentifierSegment::parse(ref_name.to_string()) + .map_err(|err| s3_error!(InvalidRequest, "invalid ref name: {}", err))?; + Ok(ref_name.to_string()) +} + fn job_id_from_params(params: &Params<'_, '_>) -> S3Result { let job = params.get("job").unwrap_or(""); if job.is_empty() { @@ -1151,6 +1253,21 @@ fn list_tables_response_from_entries(entries: Vec) -> S3Result { + let identifiers = entries + .into_iter() + .map(|entry| { + let namespace = crate::table_catalog::Namespace::parse(&entry.namespace) + .map_err(|err| s3_error!(InternalError, "persisted view entry namespace is invalid: {}", err))?; + Ok(RestTableIdentifier { + namespace: namespace_segments(&namespace), + name: entry.view, + }) + }) + .collect::>>()?; + Ok(RestListViewsResponse { identifiers }) +} + fn table_credential_vending_enabled() -> bool { std::env::var(ENV_TABLE_CATALOG_CREDENTIAL_VENDING) .ok() @@ -1304,6 +1421,21 @@ fn load_table_response_from_entry(entry: crate::table_catalog::TableEntry, metad } } +fn load_view_response_from_entry(entry: crate::table_catalog::ViewEntry, metadata: serde_json::Value) -> RestLoadViewResponse { + let mut config = BTreeMap::new(); + let warehouse_location = entry.warehouse_location.clone(); + config.insert("warehouse-location".to_string(), warehouse_location.clone()); + config.insert(CREDENTIAL_SCOPE_CONFIG_KEY.to_string(), CREDENTIAL_SCOPE_TABLE_PREFIX.to_string()); + config.insert(CREDENTIAL_SCOPE_PREFIX_CONFIG_KEY.to_string(), warehouse_location); + config.insert(CREDENTIAL_MODE_CONFIG_KEY.to_string(), CREDENTIAL_MODE_CLIENT_PROVIDED.to_string()); + + RestLoadViewResponse { + metadata_location: entry.metadata_location, + metadata, + config, + } +} + async fn load_credentials_response_from_entry( entry: &crate::table_catalog::TableEntry, issuer: &dyn TableCredentialIssuer, @@ -1414,6 +1546,11 @@ fn validate_metadata_table_location_in_bucket(bucket: &str, metadata: &serde_jso validate_table_location_in_bucket(bucket, location) } +fn validate_metadata_view_location_in_bucket(bucket: &str, metadata: &serde_json::Value) -> S3Result<()> { + let location = metadata_table_location(metadata)?; + validate_table_location_in_bucket(bucket, location) +} + fn validate_metadata_matches_current_metadata( current_metadata: &serde_json::Value, target_metadata: &serde_json::Value, @@ -1431,6 +1568,28 @@ fn validate_metadata_matches_current_metadata( Ok(()) } +fn metadata_view_uuid(metadata: &serde_json::Value) -> S3Result<&str> { + metadata + .get("view-uuid") + .and_then(serde_json::Value::as_str) + .filter(|uuid| !uuid.is_empty()) + .ok_or_else(|| s3_error!(InvalidRequest, "view metadata is missing view-uuid")) +} + +fn validate_metadata_matches_current_view_metadata( + current_metadata: &serde_json::Value, + target_metadata: &serde_json::Value, +) -> S3Result<()> { + let expected_view_uuid = metadata_view_uuid(current_metadata)?; + metadata_format_version(current_metadata)?; + let target_view_uuid = metadata_view_uuid(target_metadata)?; + metadata_format_version(target_metadata)?; + if target_view_uuid != expected_view_uuid { + return Err(s3_error!(InvalidRequest, "view metadata view-uuid does not match current view metadata")); + } + Ok(()) +} + fn adopt_registered_metadata_identity( entry: &mut crate::table_catalog::TableEntry, metadata: &serde_json::Value, @@ -1559,6 +1718,41 @@ fn table_entry_from_create_table_request( Ok((entry, metadata)) } +fn view_entry_from_create_view_request( + bucket: &str, + namespace: &crate::table_catalog::Namespace, + request: CreateViewRequest, +) -> S3Result<(crate::table_catalog::ViewEntry, serde_json::Value)> { + let view = crate::table_catalog::IdentifierSegment::parse(request.name) + .map_err(|err| s3_error!(InvalidRequest, "invalid view name: {}", err))?; + let view_id = Uuid::new_v4().to_string(); + let view_uuid = Uuid::new_v4().to_string(); + let warehouse_location = request.location.unwrap_or_else(|| format!("s3://{bucket}/views/{view_id}")); + validate_table_location_in_bucket(bucket, &warehouse_location)?; + let metadata_location = crate::table_catalog::default_view_metadata_file_path(namespace, &view, "00001.metadata.json"); + + let entry = crate::table_catalog::ViewEntry { + version: crate::table_catalog::TABLE_CATALOG_ENTRY_VERSION, + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + view: view.as_str().to_string(), + view_id, + view_uuid, + format: "ICEBERG_VIEW".to_string(), + format_version: 1, + warehouse_location, + metadata_location, + version_token: format!("token-{}", Uuid::new_v4()), + generation: 1, + state: crate::table_catalog::TableCatalogEntryState::Active, + properties: request.properties, + created_at: None, + updated_at: None, + }; + let metadata = initial_view_metadata_json(&entry, request.schema, request.view_version, entry.properties.clone())?; + Ok((entry, metadata)) +} + fn initial_table_metadata_json( entry: &crate::table_catalog::TableEntry, mut schema: serde_json::Value, @@ -1641,6 +1835,60 @@ fn initial_table_metadata_json( })) } +fn initial_view_metadata_json( + entry: &crate::table_catalog::ViewEntry, + mut schema: serde_json::Value, + mut view_version: serde_json::Value, + properties: BTreeMap, +) -> S3Result { + let schema_object = schema + .as_object_mut() + .ok_or_else(|| s3_error!(InvalidRequest, "schema must be a JSON object"))?; + schema_object + .entry("schema-id".to_string()) + .or_insert_with(|| serde_json::Value::from(0)); + let schema_id = schema_object + .get("schema-id") + .and_then(serde_json::Value::as_i64) + .ok_or_else(|| s3_error!(InvalidRequest, "schema-id must be an integer"))?; + + let view_version_object = view_version + .as_object_mut() + .ok_or_else(|| s3_error!(InvalidRequest, "view-version must be a JSON object"))?; + view_version_object + .entry("version-id".to_string()) + .or_insert_with(|| serde_json::Value::from(1)); + view_version_object + .entry("schema-id".to_string()) + .or_insert_with(|| serde_json::Value::from(schema_id)); + view_version_object + .entry("timestamp-ms".to_string()) + .or_insert_with(|| serde_json::Value::from(current_time_millis())); + let version_id = view_version_object + .get("version-id") + .and_then(serde_json::Value::as_i64) + .ok_or_else(|| s3_error!(InvalidRequest, "view-version version-id must be an integer"))?; + let timestamp_ms = view_version_object + .get("timestamp-ms") + .and_then(serde_json::Value::as_i64) + .unwrap_or_else(current_time_millis); + + Ok(serde_json::json!({ + "format-version": entry.format_version, + "view-uuid": entry.view_uuid, + "location": entry.warehouse_location, + "current-version-id": version_id, + "schemas": [schema], + "versions": [view_version], + "version-log": [{ + "timestamp-ms": timestamp_ms, + "version-id": version_id + }], + "metadata-log": [], + "properties": properties + })) +} + fn current_time_millis() -> i64 { let now = OffsetDateTime::now_utc(); now.unix_timestamp() @@ -1861,6 +2109,73 @@ fn apply_table_commit_updates( Ok(metadata) } +fn validate_view_commit_requirements(metadata: &serde_json::Value, requirements: &[serde_json::Value]) -> S3Result<()> { + for requirement in requirements { + let requirement_type = requirement + .get("type") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| s3_error!(InvalidRequest, "commit requirement type is required"))?; + match requirement_type { + "assert-view-uuid" => { + let expected = requirement + .get("uuid") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| s3_error!(InvalidRequest, "assert-view-uuid requires uuid"))?; + let actual = metadata + .get("view-uuid") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| s3_error!(InvalidRequest, "current view metadata is missing view-uuid"))?; + if actual != expected { + return Err(s3_error!(PreconditionFailed, "commit requirement failed: view uuid changed")); + } + } + "assert-current-view-version-id" => { + validate_i64_requirement_with_metadata_key( + metadata, + requirement, + "current-view-version-id", + "current-version-id", + "current view version id", + )?; + } + _ => return Err(s3_error!(NotImplemented, "unsupported view commit requirement: {requirement_type}")), + } + } + Ok(()) +} + +fn apply_view_commit_updates( + mut metadata: serde_json::Value, + updates: &[serde_json::Value], + previous_metadata_location: &str, +) -> S3Result { + if !metadata.is_object() { + return Err(s3_error!(InvalidRequest, "current view metadata must be a JSON object")); + } + + for update in updates { + let action = update + .get("action") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| s3_error!(InvalidRequest, "view update action is required"))?; + match action { + "assign-uuid" => apply_assign_uuid_update(&mut metadata, update)?, + "add-schema" => apply_add_schema_update(&mut metadata, update)?, + "set-current-schema" => apply_set_current_schema_update(&mut metadata, update)?, + "add-view-version" => apply_add_view_version_update(&mut metadata, update)?, + "set-current-view-version" => apply_set_current_view_version_update(&mut metadata, update)?, + "set-location" => apply_set_location_update(&mut metadata, update)?, + "set-properties" => apply_set_properties_update(&mut metadata, update)?, + "remove-properties" => apply_remove_properties_update(&mut metadata, update)?, + _ => return Err(s3_error!(NotImplemented, "unsupported view update: {action}")), + } + } + + append_previous_metadata_log(&mut metadata, previous_metadata_location)?; + metadata_object_mut(&mut metadata)?.insert("last-updated-ms".to_string(), serde_json::Value::from(current_time_millis())); + Ok(metadata) +} + fn apply_assign_uuid_update(metadata: &mut serde_json::Value, update: &serde_json::Value) -> S3Result<()> { let uuid = update .get("uuid") @@ -1876,6 +2191,48 @@ fn apply_assign_uuid_update(metadata: &mut serde_json::Value, update: &serde_jso Ok(()) } +fn apply_add_view_version_update(metadata: &mut serde_json::Value, update: &serde_json::Value) -> S3Result<()> { + let mut view_version = update + .get("view-version") + .cloned() + .ok_or_else(|| s3_error!(InvalidRequest, "add-view-version requires view-version"))?; + if !view_version.is_object() { + return Err(s3_error!(InvalidRequest, "view-version must be a JSON object")); + } + if view_version.get("version-id").is_none() { + let next_id = next_array_object_i64(metadata, "versions", "version-id")?; + view_version + .as_object_mut() + .ok_or_else(|| s3_error!(InvalidRequest, "view-version must be a JSON object"))? + .insert("version-id".to_string(), serde_json::Value::from(next_id)); + } + view_version + .as_object_mut() + .ok_or_else(|| s3_error!(InvalidRequest, "view-version must be a JSON object"))? + .entry("timestamp-ms".to_string()) + .or_insert_with(|| serde_json::Value::from(current_time_millis())); + ensure_array_field(metadata, "versions")?.push(view_version); + Ok(()) +} + +fn apply_set_current_view_version_update(metadata: &mut serde_json::Value, update: &serde_json::Value) -> S3Result<()> { + let requested_id = update + .get("view-version-id") + .and_then(serde_json::Value::as_i64) + .ok_or_else(|| s3_error!(InvalidRequest, "set-current-view-version requires view-version-id"))?; + let version_id = if requested_id == -1 { + last_array_object_i64(metadata, "versions", "version-id")? + } else { + requested_id + }; + metadata_object_mut(metadata)?.insert("current-version-id".to_string(), serde_json::Value::from(version_id)); + ensure_array_field(metadata, "version-log")?.push(serde_json::json!({ + "timestamp-ms": current_time_millis(), + "version-id": version_id + })); + Ok(()) +} + fn apply_upgrade_format_version_update(metadata: &mut serde_json::Value, update: &serde_json::Value) -> S3Result<()> { let version = update .get("format-version") @@ -2344,6 +2701,34 @@ where Ok(load_table_response_from_entry(entry, metadata)) } +async fn create_view_response( + store: &S, + metadata_backend: &impl crate::table_catalog::TableCatalogObjectBackend, + bucket: &str, + namespace: &crate::table_catalog::Namespace, + request: CreateViewRequest, + table_bucket_enabled: bool, +) -> S3Result +where + S: crate::table_catalog::TableCatalogStore + ?Sized, +{ + let (entry, metadata) = view_entry_from_create_view_request(bucket, namespace, request)?; + ensure_table_bucket_entry(store, bucket, table_bucket_enabled).await?; + let metadata_data = serde_json::to_vec(&metadata) + .map_err(|err| s3_error!(InternalError, "failed to serialize initial view metadata: {}", err))?; + metadata_backend + .put_object( + bucket, + &entry.metadata_location, + metadata_data, + crate::table_catalog::TableCatalogPutPrecondition::IfAbsent, + ) + .await + .map_err(catalog_store_error)?; + store.create_view(entry.clone()).await.map_err(catalog_store_error)?; + Ok(load_view_response_from_entry(entry, metadata)) +} + async fn read_table_metadata_json( metadata_backend: &impl crate::table_catalog::TableCatalogObjectBackend, bucket: &str, @@ -2400,6 +2785,132 @@ where Ok(load_table_response_from_entry(entry, metadata)) } +async fn list_views_response( + store: &S, + bucket: &str, + namespace: &crate::table_catalog::Namespace, +) -> S3Result +where + S: crate::table_catalog::TableCatalogStore + ?Sized, +{ + let entries = store + .list_views(bucket, &namespace.public_name()) + .await + .map_err(catalog_store_error)?; + list_views_response_from_entries(entries) +} + +async fn load_view_response( + store: &S, + metadata_backend: &impl crate::table_catalog::TableCatalogObjectBackend, + bucket: &str, + namespace: &crate::table_catalog::Namespace, + view: &str, +) -> S3Result +where + S: crate::table_catalog::TableCatalogStore + ?Sized, +{ + let Some(entry) = store + .load_view(bucket, &namespace.public_name(), view) + .await + .map_err(catalog_store_error)? + else { + return Err(s3_error!(InvalidRequest, "view not found")); + }; + let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?; + Ok(load_view_response_from_entry(entry, metadata)) +} + +async fn view_exists_status( + store: &S, + bucket: &str, + namespace: &crate::table_catalog::Namespace, + view: &str, +) -> S3Result +where + S: crate::table_catalog::TableCatalogStore + ?Sized, +{ + let exists = store + .load_view(bucket, &namespace.public_name(), view) + .await + .map_err(catalog_store_error)? + .is_some(); + Ok(exists_status(exists)) +} + +async fn replace_view_response( + store: &S, + metadata_backend: &impl crate::table_catalog::TableCatalogObjectBackend, + bucket: &str, + namespace: &crate::table_catalog::Namespace, + view: &str, + request: RestCommitViewRequest, +) -> S3Result +where + S: crate::table_catalog::TableCatalogStore + ?Sized, +{ + let Some(current) = store + .load_view(bucket, &namespace.public_name(), view) + .await + .map_err(catalog_store_error)? + else { + return Err(s3_error!(InvalidRequest, "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)?; + let view_name = crate::table_catalog::IdentifierSegment::parse(view.to_string()) + .map_err(|err| s3_error!(InvalidRequest, "invalid view name: {}", err))?; + let (next_metadata_location, next_metadata) = if let Some(new_metadata_location) = request.new_metadata_location { + if !crate::table_catalog::is_valid_view_metadata_location(namespace, &view_name, &new_metadata_location) { + return Err(s3_error!(InvalidRequest, "metadata location must be inside the view metadata directory")); + } + let target_metadata = read_table_metadata_json(metadata_backend, bucket, &new_metadata_location).await?; + validate_metadata_view_location_in_bucket(bucket, &target_metadata)?; + validate_metadata_matches_current_view_metadata(¤t_metadata, &target_metadata)?; + (new_metadata_location, target_metadata) + } else { + let next_metadata = apply_view_commit_updates(current_metadata.clone(), &request.updates, ¤t.metadata_location)?; + validate_metadata_view_location_in_bucket(bucket, &next_metadata)?; + validate_metadata_matches_current_view_metadata(¤t_metadata, &next_metadata)?; + let (_, metadata_file_token) = standard_commit_ids(request.commit_id); + let next_generation = current.generation.saturating_add(1); + let next_metadata_location = crate::table_catalog::default_view_metadata_file_path( + namespace, + &view_name, + &next_metadata_file_name(next_generation, &metadata_file_token), + ); + let next_metadata_data = serde_json::to_vec(&next_metadata) + .map_err(|err| s3_error!(InternalError, "failed to serialize view metadata update: {}", err))?; + metadata_backend + .put_object( + bucket, + &next_metadata_location, + next_metadata_data, + crate::table_catalog::TableCatalogPutPrecondition::IfAbsent, + ) + .await + .map_err(catalog_store_error)?; + (next_metadata_location, next_metadata) + }; + + let result = store + .replace_view(crate::table_catalog::ViewCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + view: view.to_string(), + expected_version_token: request + .expected_version_token + .unwrap_or_else(|| current.version_token.clone()), + expected_metadata_location: request + .expected_metadata_location + .unwrap_or_else(|| current.metadata_location.clone()), + new_metadata_location: next_metadata_location, + }) + .await + .map_err(catalog_store_error)?; + Ok(load_view_response_from_entry(result.view, next_metadata)) +} + async fn table_exists_status( store: &S, bucket: &str, @@ -2611,6 +3122,16 @@ where .map_err(catalog_store_error) } +async fn drop_view_in_store(store: &S, bucket: &str, namespace: &crate::table_catalog::Namespace, view: &str) -> S3Result<()> +where + S: crate::table_catalog::TableCatalogStore + ?Sized, +{ + store + .drop_view(bucket, &namespace.public_name(), view) + .await + .map_err(catalog_store_error) +} + async fn table_metadata_maintenance_response( store: &crate::table_catalog::ObjectTableCatalogStore, metadata_backend: &B, @@ -2846,6 +3367,135 @@ where }) } +async fn put_table_ref_response( + store: &S, + metadata_backend: &impl crate::table_catalog::TableCatalogObjectBackend, + bucket: &str, + namespace: &crate::table_catalog::Namespace, + table: &str, + ref_name: &str, + request: PutTableRefRequest, +) -> S3Result +where + S: crate::table_catalog::TableCatalogStore + ?Sized, +{ + if !matches!(request.ref_type.as_str(), "branch" | "tag") { + return Err(s3_error!(InvalidRequest, "snapshot ref type must be branch or tag")); + } + let mut update = serde_json::json!({ + "action": "set-snapshot-ref", + "ref-name": ref_name, + "type": request.ref_type, + "snapshot-id": request.snapshot_id + }); + if let Some(value) = request.min_snapshots_to_keep { + update["min-snapshots-to-keep"] = serde_json::Value::from(value); + } + if let Some(value) = request.max_snapshot_age_ms { + update["max-snapshot-age-ms"] = serde_json::Value::from(value); + } + if let Some(value) = request.max_ref_age_ms { + update["max-ref-age-ms"] = serde_json::Value::from(value); + } + let mut requirements = Vec::new(); + if let Some(expected_snapshot_id) = request.expected_snapshot_id { + requirements.push(serde_json::json!({ + "type": "assert-ref-snapshot-id", + "ref": ref_name, + "snapshot-id": expected_snapshot_id + })); + } + standard_commit_table_response( + store, + metadata_backend, + bucket, + namespace, + table, + RestCommitTableRequest { + _identifier: None, + commit_id: request.commit_id, + idempotency_key: request.idempotency_key, + operation: Some("set-snapshot-ref".to_string()), + expected_version_token: None, + expected_metadata_location: None, + new_metadata_location: None, + requirements, + updates: vec![update], + writer: request.writer.or_else(|| Some("rustfs-ref-api".to_string())), + }, + ) + .await +} + +async fn delete_table_ref_response( + store: &S, + metadata_backend: &impl crate::table_catalog::TableCatalogObjectBackend, + bucket: &str, + namespace: &crate::table_catalog::Namespace, + table: &str, + ref_name: &str, + request: DeleteTableRefRequest, +) -> S3Result +where + S: crate::table_catalog::TableCatalogStore + ?Sized, +{ + if ref_name == "main" { + return Err(s3_error!(InvalidRequest, "main snapshot ref cannot be deleted")); + } + let Some(entry) = store + .load_table(bucket, &namespace.public_name(), table) + .await + .map_err(catalog_store_error)? + else { + return Err(s3_error!(InvalidRequest, "table not found")); + }; + let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?; + let reference = metadata + .get("refs") + .and_then(serde_json::Value::as_object) + .and_then(|refs| refs.get(ref_name)); + if reference.is_some_and(snapshot_ref_has_explicit_retention) && !request.force { + return Err(s3_error!(InvalidRequest, "snapshot ref has retention policy; force is required")); + } + let mut requirements = Vec::new(); + if let Some(expected_snapshot_id) = request.expected_snapshot_id { + requirements.push(serde_json::json!({ + "type": "assert-ref-snapshot-id", + "ref": ref_name, + "snapshot-id": expected_snapshot_id + })); + } + standard_commit_table_response( + store, + metadata_backend, + bucket, + namespace, + table, + RestCommitTableRequest { + _identifier: None, + commit_id: request.commit_id, + idempotency_key: request.idempotency_key, + operation: Some("remove-snapshot-ref".to_string()), + expected_version_token: None, + expected_metadata_location: None, + new_metadata_location: None, + requirements, + updates: vec![serde_json::json!({ + "action": "remove-snapshot-ref", + "ref-name": ref_name + })], + writer: request.writer.or_else(|| Some("rustfs-ref-api".to_string())), + }, + ) + .await +} + +fn snapshot_ref_has_explicit_retention(reference: &serde_json::Value) -> bool { + reference.get("min-snapshots-to-keep").is_some() + || reference.get("max-snapshot-age-ms").is_some() + || reference.get("max-ref-age-ms").is_some() +} + fn external_catalog_bridge_response( bucket: &str, namespace: &crate::table_catalog::Namespace, @@ -3151,7 +3801,9 @@ impl Operation for RestListViewsHandler { let namespace = namespace_from_params(¶ms)?; let resource = TableCatalogResource::namespace(&warehouse, &namespace); authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataAction).await?; - unsupported_catalog_feature_response("iceberg-views", "use Iceberg tables; view catalog routes are reserved") + let store = table_catalog_store()?; + let response = list_views_response(&store, &warehouse, &namespace).await?; + build_json_response(StatusCode::OK, &response) } } @@ -3164,7 +3816,13 @@ impl Operation for RestCreateViewHandler { let namespace = namespace_from_params(¶ms)?; let resource = TableCatalogResource::namespace(&warehouse, &namespace); authorize_table_catalog_resource_request(&req, &resource, AdminAction::CreateTableAction).await?; - unsupported_catalog_feature_response("iceberg-views", "create a table or register an existing Iceberg table") + let request = read_json_body::(req.input).await?; + let metadata_backend = table_catalog_backend()?; + let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone()); + let table_bucket_enabled = table_bucket_enabled_from_metadata(&warehouse).await?; + let response = + create_view_response(&store, &metadata_backend, &warehouse, &namespace, request, table_bucket_enabled).await?; + build_json_response(StatusCode::OK, &response) } } @@ -3259,10 +3917,28 @@ impl Operation for RestLoadViewHandler { async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; let namespace = namespace_from_params(¶ms)?; - let _view = view_name_from_params(¶ms)?; - let resource = TableCatalogResource::namespace(&warehouse, &namespace); + let view = view_name_from_params(¶ms)?; + let resource = TableCatalogResource::view(&warehouse, &namespace, &view); authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataAction).await?; - unsupported_catalog_feature_response("iceberg-views", "use Iceberg tables; view load is reserved") + let metadata_backend = table_catalog_backend()?; + let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone()); + let response = load_view_response(&store, &metadata_backend, &warehouse, &namespace, &view).await?; + build_json_response(StatusCode::OK, &response) + } +} + +pub struct RestViewExistsHandler {} + +#[async_trait::async_trait] +impl Operation for RestViewExistsHandler { + async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { + let warehouse = warehouse_from_params(¶ms)?; + let namespace = namespace_from_params(¶ms)?; + let view = view_name_from_params(¶ms)?; + let resource = TableCatalogResource::view(&warehouse, &namespace, &view); + authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableAction).await?; + let store = table_catalog_store()?; + Ok(empty_response(view_exists_status(&store, &warehouse, &namespace, &view).await?)) } } @@ -3273,10 +3949,14 @@ impl Operation for RestReplaceViewHandler { async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; let namespace = namespace_from_params(¶ms)?; - let _view = view_name_from_params(¶ms)?; - let resource = TableCatalogResource::namespace(&warehouse, &namespace); + let view = view_name_from_params(¶ms)?; + let resource = TableCatalogResource::view(&warehouse, &namespace, &view); authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?; - unsupported_catalog_feature_response("iceberg-views", "commit table metadata updates; view replacement is reserved") + let request = read_json_body::(req.input).await?; + let metadata_backend = table_catalog_backend()?; + let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone()); + let response = replace_view_response(&store, &metadata_backend, &warehouse, &namespace, &view, request).await?; + build_json_response(StatusCode::OK, &response) } } @@ -3287,10 +3967,12 @@ impl Operation for RestDropViewHandler { async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; let namespace = namespace_from_params(¶ms)?; - let _view = view_name_from_params(¶ms)?; - let resource = TableCatalogResource::namespace(&warehouse, &namespace); + let view = view_name_from_params(¶ms)?; + let resource = TableCatalogResource::view(&warehouse, &namespace, &view); authorize_table_catalog_resource_request(&req, &resource, AdminAction::DeleteTableAction).await?; - unsupported_catalog_feature_response("iceberg-views", "drop table resources; view deletion is reserved") + let store = table_catalog_store()?; + drop_view_in_store(&store, &warehouse, &namespace, &view).await?; + Ok(empty_response(StatusCode::NO_CONTENT)) } } @@ -3311,6 +3993,46 @@ impl Operation for ListTableRefsHandler { } } +pub struct PutTableRefHandler {} + +#[async_trait::async_trait] +impl Operation for PutTableRefHandler { + async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { + let warehouse = warehouse_from_params(¶ms)?; + let namespace = namespace_from_params(¶ms)?; + let table = table_name_from_params(¶ms)?; + let ref_name = ref_name_from_params(¶ms)?; + let resource = TableCatalogResource::table(&warehouse, &namespace, &table); + authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?; + let request = read_json_body::(req.input).await?; + let metadata_backend = table_catalog_backend()?; + let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone()); + let response = + put_table_ref_response(&store, &metadata_backend, &warehouse, &namespace, &table, &ref_name, request).await?; + build_json_response(StatusCode::OK, &response) + } +} + +pub struct DeleteTableRefHandler {} + +#[async_trait::async_trait] +impl Operation for DeleteTableRefHandler { + async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { + let warehouse = warehouse_from_params(¶ms)?; + let namespace = namespace_from_params(¶ms)?; + let table = table_name_from_params(¶ms)?; + let ref_name = ref_name_from_params(¶ms)?; + let resource = TableCatalogResource::table(&warehouse, &namespace, &table); + authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?; + let request = read_json_body::(req.input).await?; + let metadata_backend = table_catalog_backend()?; + let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone()); + let response = + delete_table_ref_response(&store, &metadata_backend, &warehouse, &namespace, &table, &ref_name, request).await?; + build_json_response(StatusCode::OK, &response) + } +} + pub struct GetTableMetadataLocationHandler {} #[async_trait::async_trait] @@ -3659,6 +4381,11 @@ mod tests { ); assert!(response.endpoints.contains(&"GET /{warehouse}/namespaces/{namespace}/views")); assert!(response.endpoints.contains(&"POST /{warehouse}/namespaces/{namespace}/views")); + assert!( + response + .endpoints + .contains(&"HEAD /{warehouse}/namespaces/{namespace}/views/{view}") + ); assert!( response .endpoints @@ -3689,6 +4416,16 @@ mod tests { .endpoints .contains(&"GET /{warehouse}/namespaces/{namespace}/tables/{table}/refs") ); + assert!( + response + .endpoints + .contains(&"PUT /{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}") + ); + assert!( + response + .endpoints + .contains(&"DELETE /{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}") + ); assert!( response .endpoints @@ -5422,6 +6159,260 @@ mod tests { assert!(validate_metadata_table_location_in_bucket("warehouse", &updated).is_err()); } + #[test] + fn create_view_request_accepts_standard_iceberg_rest_shape() { + let request: CreateViewRequest = serde_json::from_value(serde_json::json!({ + "name": "recent_events", + "schema": { + "type": "struct", + "schema-id": 0, + "fields": [ + { + "id": 1, + "name": "id", + "required": true, + "type": "long" + } + ] + }, + "view-version": { + "version-id": 1, + "schema-id": 0, + "summary": { + "engine-name": "spark", + "engine-version": "3.5.0" + }, + "default-catalog": "warehouse", + "default-namespace": ["analytics"], + "representations": [ + { + "type": "sql", + "sql": "SELECT id FROM analytics.events WHERE ts >= current_date()", + "dialect": "spark" + } + ] + }, + "properties": { + "comment": "recent event ids" + } + })) + .expect("standard create view request should parse"); + + assert_eq!(request.name, "recent_events"); + assert_eq!(request.properties.get("comment").map(String::as_str), Some("recent event ids")); + } + + #[tokio::test] + async fn view_catalog_responses_persist_replace_and_drop_view_metadata() { + let store = TestTableCatalogStore::default(); + let metadata_backend = TestTableCatalogObjectBackend::default(); + let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse"); + ensure_table_bucket_entry(&store, "warehouse", true) + .await + .expect("table bucket entry should be seeded"); + create_namespace_response( + &store, + "warehouse", + CreateNamespaceRequest { + namespace: vec!["analytics".to_string()], + properties: BTreeMap::new(), + }, + true, + ) + .await + .expect("namespace should be created"); + + let create_request: CreateViewRequest = serde_json::from_value(serde_json::json!({ + "name": "recent_events", + "schema": { + "type": "struct", + "schema-id": 0, + "fields": [ + { + "id": 1, + "name": "id", + "required": true, + "type": "long" + } + ] + }, + "view-version": { + "version-id": 1, + "schema-id": 0, + "summary": { + "engine-name": "spark" + }, + "default-catalog": "warehouse", + "default-namespace": ["analytics"], + "representations": [ + { + "type": "sql", + "sql": "SELECT id FROM analytics.events", + "dialect": "spark" + } + ] + } + })) + .expect("standard create view request should parse"); + + let created = create_view_response(&store, &metadata_backend, "warehouse", &namespace, create_request, true) + .await + .expect("view should be created"); + assert_eq!(created.metadata["format-version"], 1); + assert_eq!(created.metadata["current-version-id"], 1); + assert_eq!(created.metadata["versions"][0]["representations"][0]["dialect"], "spark"); + assert!( + metadata_backend + .object_exists("warehouse", &created.metadata_location) + .await + .expect("view metadata object lookup should succeed") + ); + + let listed = list_views_response(&store, "warehouse", &namespace) + .await + .expect("views should list"); + assert_eq!(listed.identifiers.len(), 1); + assert_eq!(listed.identifiers[0].name, "recent_events"); + + let loaded = load_view_response(&store, &metadata_backend, "warehouse", &namespace, "recent_events") + .await + .expect("view should load"); + assert_eq!(loaded.metadata_location, created.metadata_location); + let replace_request: RestCommitViewRequest = serde_json::from_value(serde_json::json!({ + "updates": [ + { + "action": "add-view-version", + "view-version": { + "version-id": 2, + "schema-id": 0, + "summary": { + "engine-name": "spark" + }, + "default-catalog": "warehouse", + "default-namespace": ["analytics"], + "representations": [ + { + "type": "sql", + "sql": "SELECT id FROM analytics.events WHERE id > 10", + "dialect": "spark" + } + ] + } + }, + { + "action": "set-current-view-version", + "view-version-id": 2 + } + ] + })) + .expect("replace view request should parse"); + let replaced = + replace_view_response(&store, &metadata_backend, "warehouse", &namespace, "recent_events", replace_request) + .await + .expect("view should replace"); + assert_ne!(replaced.metadata_location, created.metadata_location); + assert_eq!(replaced.metadata["current-version-id"], 2); + assert_eq!( + replaced.metadata["version-log"] + .as_array() + .expect("version log should be an array") + .len(), + 2 + ); + + drop_view_in_store(&store, "warehouse", &namespace, "recent_events") + .await + .expect("view should drop"); + let listed = list_views_response(&store, "warehouse", &namespace) + .await + .expect("views should list after drop"); + assert!(listed.identifiers.is_empty()); + } + + #[tokio::test] + async fn table_ref_write_responses_commit_retention_refs_and_protect_deletes() { + let store = TestTableCatalogStore::default(); + let metadata_backend = TestTableCatalogObjectBackend::default(); + let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse"); + create_standard_events_table(&store, &metadata_backend, &namespace).await; + + let append_request: RestCommitTableRequest = serde_json::from_value(serde_json::json!({ + "updates": [ + { + "action": "add-snapshot", + "snapshot": { + "snapshot-id": 10, + "sequence-number": 1, + "timestamp-ms": 1234, + "manifest-list": "s3://warehouse/tables/table-id/metadata/snap-10.avro", + "summary": { + "operation": "append" + } + } + }, + { + "action": "set-snapshot-ref", + "ref-name": "main", + "snapshot-id": 10, + "type": "branch" + } + ] + })) + .expect("append request should parse"); + commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", append_request) + .await + .expect("append should commit"); + + let ref_request: PutTableRefRequest = serde_json::from_value(serde_json::json!({ + "snapshot-id": 10, + "type": "tag", + "max-ref-age-ms": 86400000, + "expected-snapshot-id": null + })) + .expect("ref put request should parse"); + put_table_ref_response(&store, &metadata_backend, "warehouse", &namespace, "events", "audit", ref_request) + .await + .expect("ref put should commit"); + + let refs = table_refs_response(&store, &metadata_backend, "warehouse", &namespace, "events") + .await + .expect("refs should load"); + assert_eq!(refs.refs["audit"]["type"], "tag"); + assert_eq!(refs.refs["audit"]["max-ref-age-ms"], 86400000); + + let delete_without_force: DeleteTableRefRequest = + serde_json::from_value(serde_json::json!({})).expect("ref delete request should parse"); + let error = delete_table_ref_response( + &store, + &metadata_backend, + "warehouse", + &namespace, + "events", + "audit", + delete_without_force, + ) + .await + .expect_err("retention refs should require force delete"); + assert_eq!(error.code(), &s3s::S3ErrorCode::InvalidRequest); + + let force_delete: DeleteTableRefRequest = + serde_json::from_value(serde_json::json!({ "force": true })).expect("ref force delete should parse"); + delete_table_ref_response(&store, &metadata_backend, "warehouse", &namespace, "events", "audit", force_delete) + .await + .expect("force delete should commit"); + let refs = table_refs_response(&store, &metadata_backend, "warehouse", &namespace, "events") + .await + .expect("refs should load after delete"); + assert!(!refs.refs.contains_key("audit")); + + let main_delete: DeleteTableRefRequest = + serde_json::from_value(serde_json::json!({ "force": true })).expect("main delete request should parse"); + let error = delete_table_ref_response(&store, &metadata_backend, "warehouse", &namespace, "events", "main", main_delete) + .await + .expect_err("main ref should remain protected"); + assert_eq!(error.code(), &s3s::S3ErrorCode::InvalidRequest); + } + #[test] fn load_table_response_includes_rest_metadata_payload() { let metadata = serde_json::json!({ @@ -5753,6 +6744,7 @@ mod tests { table_buckets: tokio::sync::Mutex>, namespaces: tokio::sync::Mutex>, tables: tokio::sync::Mutex>, + views: tokio::sync::Mutex>, commits: tokio::sync::Mutex>, fail_put_table_bucket: tokio::sync::Mutex, } @@ -6189,6 +7181,98 @@ mod tests { Ok(()) } + async fn create_view(&self, entry: crate::table_catalog::ViewEntry) -> crate::table_catalog::TableCatalogStoreResult<()> { + if self.get_table_bucket(&entry.table_bucket).await?.is_none() { + return Err(crate::table_catalog::TableCatalogStoreError::NotFound(format!( + "table bucket {}", + entry.table_bucket + ))); + } + if self.get_namespace(&entry.table_bucket, &entry.namespace).await?.is_none() { + return Err(crate::table_catalog::TableCatalogStoreError::NotFound(format!( + "namespace {}/{}", + entry.table_bucket, entry.namespace + ))); + } + self.views.lock().await.push(entry); + Ok(()) + } + + async fn list_views( + &self, + table_bucket: &str, + namespace: &str, + ) -> crate::table_catalog::TableCatalogStoreResult> { + Ok(self + .views + .lock() + .await + .iter() + .filter(|entry| entry.table_bucket == table_bucket && entry.namespace == namespace) + .cloned() + .collect()) + } + + async fn load_view( + &self, + table_bucket: &str, + namespace: &str, + view: &str, + ) -> crate::table_catalog::TableCatalogStoreResult> { + Ok(self + .views + .lock() + .await + .iter() + .find(|entry| entry.table_bucket == table_bucket && entry.namespace == namespace && entry.view == view) + .cloned()) + } + + async fn replace_view( + &self, + request: crate::table_catalog::ViewCommitRequest, + ) -> crate::table_catalog::TableCatalogStoreResult { + let mut views = self.views.lock().await; + let Some(index) = views.iter().position(|entry| { + entry.table_bucket == request.table_bucket && entry.namespace == request.namespace && entry.view == request.view + }) else { + return Err(crate::table_catalog::TableCatalogStoreError::NotFound(format!( + "view {}/{}/{}", + request.table_bucket, request.namespace, request.view + ))); + }; + let current = views[index].clone(); + if current.version_token != request.expected_version_token { + return Err(crate::table_catalog::TableCatalogStoreError::Conflict( + "current view version token does not match expected token".to_string(), + )); + } + if current.metadata_location != request.expected_metadata_location { + return Err(crate::table_catalog::TableCatalogStoreError::Conflict( + "current view metadata location does not match expected location".to_string(), + )); + } + let mut next = current; + next.metadata_location = request.new_metadata_location; + next.version_token = "token-view-committed".to_string(); + next.generation = next.generation.saturating_add(1); + views[index] = next.clone(); + Ok(crate::table_catalog::ViewCommitResult { view: next }) + } + + async fn drop_view( + &self, + table_bucket: &str, + namespace: &str, + view: &str, + ) -> crate::table_catalog::TableCatalogStoreResult<()> { + self.views + .lock() + .await + .retain(|entry| !(entry.table_bucket == table_bucket && entry.namespace == namespace && entry.view == view)); + Ok(()) + } + async fn get_commit_by_id( &self, _table_bucket: &str, diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index 6485a1c9e..2c66b504c 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -718,6 +718,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ GET_TABLE_METADATA, RouteRiskLevel::Sensitive, ), + admin( + HttpMethod::Head, + "/iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}", + GET_TABLE, + RouteRiskLevel::Sensitive, + ), admin( HttpMethod::Post, "/iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}", @@ -736,6 +742,18 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ GET_TABLE_METADATA, RouteRiskLevel::Sensitive, ), + admin( + HttpMethod::Put, + "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + COMMIT_TABLE, + RouteRiskLevel::High, + ), + admin( + HttpMethod::Delete, + "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + COMMIT_TABLE, + RouteRiskLevel::High, + ), admin( HttpMethod::Get, "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/metadata-location", @@ -929,6 +947,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ GET_TABLE_METADATA, RouteRiskLevel::Sensitive, ), + admin( + HttpMethod::Head, + "/_iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}", + GET_TABLE, + RouteRiskLevel::Sensitive, + ), admin( HttpMethod::Post, "/_iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}", @@ -947,6 +971,18 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ GET_TABLE_METADATA, RouteRiskLevel::Sensitive, ), + admin( + HttpMethod::Put, + "/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + COMMIT_TABLE, + RouteRiskLevel::High, + ), + admin( + HttpMethod::Delete, + "/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + COMMIT_TABLE, + RouteRiskLevel::High, + ), admin( HttpMethod::Get, "/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/metadata-location", @@ -1214,7 +1250,7 @@ mod tests { let table_specs = ADMIN_ROUTE_POLICY_SPECS .iter() .filter(|spec| spec.path().starts_with("/iceberg/v1") || spec.path().starts_with("/_iceberg/v1")); - assert_eq!(table_specs.count(), 72); + assert_eq!(table_specs.count(), 78); assert_action(HttpMethod::Put, "/iceberg/v1/buckets/{warehouse}", SET_TABLE_BUCKET); assert_action(HttpMethod::Get, "/_iceberg/v1/buckets/{warehouse}", GET_TABLE_BUCKET); assert_action(HttpMethod::Get, "/iceberg/v1/{warehouse}/namespaces", GET_TABLE_NAMESPACE); @@ -1244,6 +1280,11 @@ mod tests { "/iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}", GET_TABLE_METADATA, ); + assert_action( + HttpMethod::Head, + "/_iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}", + GET_TABLE, + ); assert_action( HttpMethod::Post, "/_iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}", @@ -1259,6 +1300,16 @@ mod tests { "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs", GET_TABLE_METADATA, ); + assert_action( + HttpMethod::Put, + "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + COMMIT_TABLE, + ); + assert_action( + HttpMethod::Delete, + "/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + COMMIT_TABLE, + ); assert_action( HttpMethod::Get, "/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}", diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index 14b355d31..1c385d2d2 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -352,6 +352,11 @@ fn expected_admin_route_matrix() -> Vec { "/{warehouse}/namespaces/{namespace}/views/{view}", "/analytics/namespaces/sales/views/recent_orders", ), + table_route_sample( + Method::HEAD, + "/{warehouse}/namespaces/{namespace}/views/{view}", + "/analytics/namespaces/sales/views/recent_orders", + ), table_route_sample( Method::POST, "/{warehouse}/namespaces/{namespace}/views/{view}", @@ -367,6 +372,16 @@ fn expected_admin_route_matrix() -> Vec { "/{warehouse}/namespaces/{namespace}/tables/{table}/refs", "/analytics/namespaces/sales/tables/orders/refs", ), + table_route_sample( + Method::PUT, + "/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + "/analytics/namespaces/sales/tables/orders/refs/audit", + ), + table_route_sample( + Method::DELETE, + "/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + "/analytics/namespaces/sales/tables/orders/refs/audit", + ), table_route_sample( Method::POST, "/{warehouse}/namespaces/{namespace}/tables/{table}/maintenance/metadata", @@ -500,6 +515,11 @@ fn expected_admin_route_matrix() -> Vec { "/{warehouse}/namespaces/{namespace}/views/{view}", "/analytics/namespaces/sales/views/recent_orders", ), + compat_table_route_sample( + Method::HEAD, + "/{warehouse}/namespaces/{namespace}/views/{view}", + "/analytics/namespaces/sales/views/recent_orders", + ), compat_table_route_sample( Method::POST, "/{warehouse}/namespaces/{namespace}/views/{view}", @@ -515,6 +535,16 @@ fn expected_admin_route_matrix() -> Vec { "/{warehouse}/namespaces/{namespace}/tables/{table}/refs", "/analytics/namespaces/sales/tables/orders/refs", ), + compat_table_route_sample( + Method::PUT, + "/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + "/analytics/namespaces/sales/tables/orders/refs/audit", + ), + compat_table_route_sample( + Method::DELETE, + "/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}", + "/analytics/namespaces/sales/tables/orders/refs/audit", + ), compat_table_route_sample( Method::POST, "/{warehouse}/namespaces/{namespace}/tables/{table}/maintenance/metadata", @@ -719,6 +749,11 @@ fn test_register_routes_cover_representative_admin_paths() { Method::GET, &table_catalog_path("/analytics/namespaces/sales/views/recent_orders"), ); + assert_route( + &router, + Method::HEAD, + &table_catalog_path("/analytics/namespaces/sales/views/recent_orders"), + ); assert_route( &router, Method::POST, @@ -734,6 +769,16 @@ fn test_register_routes_cover_representative_admin_paths() { Method::GET, &table_catalog_path("/analytics/namespaces/sales/tables/orders/refs"), ); + assert_route( + &router, + Method::PUT, + &table_catalog_path("/analytics/namespaces/sales/tables/orders/refs/audit"), + ); + assert_route( + &router, + Method::DELETE, + &table_catalog_path("/analytics/namespaces/sales/tables/orders/refs/audit"), + ); assert_route( &router, Method::POST, @@ -841,6 +886,11 @@ fn test_register_routes_cover_representative_admin_paths() { Method::GET, &compat_table_catalog_path("/analytics/namespaces/sales/views/recent_orders"), ); + assert_route( + &router, + Method::HEAD, + &compat_table_catalog_path("/analytics/namespaces/sales/views/recent_orders"), + ); assert_route( &router, Method::POST, @@ -856,6 +906,16 @@ fn test_register_routes_cover_representative_admin_paths() { Method::GET, &compat_table_catalog_path("/analytics/namespaces/sales/tables/orders/refs"), ); + assert_route( + &router, + Method::PUT, + &compat_table_catalog_path("/analytics/namespaces/sales/tables/orders/refs/audit"), + ); + assert_route( + &router, + Method::DELETE, + &compat_table_catalog_path("/analytics/namespaces/sales/tables/orders/refs/audit"), + ); assert_route( &router, Method::POST, diff --git a/rustfs/src/table_catalog.rs b/rustfs/src/table_catalog.rs index 39dd5e95b..580a8d465 100644 --- a/rustfs/src/table_catalog.rs +++ b/rustfs/src/table_catalog.rs @@ -67,6 +67,7 @@ pub const TABLE_RESERVED_PREFIX: &str = BUCKET_TABLE_RESERVED_PREFIX; const WAREHOUSE_ROOT: &str = "warehouses"; const NAMESPACE_ROOT: &str = "namespaces"; const TABLE_ROOT: &str = "tables"; +const VIEW_ROOT: &str = "views"; const NAMESPACE_MARKER_FILE: &str = "namespace.json"; const TABLE_MARKER_FILE: &str = "table.json"; const CURRENT_POINTER_FILE: &str = "current.json"; @@ -77,6 +78,7 @@ const DELETE_DIR: &str = "delete"; const TABLE_BUCKET_ENTRY_FILE: &str = "table-bucket.json"; const NAMESPACE_ENTRY_FILE: &str = "namespace-entry.json"; const TABLE_ENTRY_FILE: &str = "table-entry.json"; +const VIEW_ENTRY_FILE: &str = "view-entry.json"; const INTERNAL_CATALOG_ROOT: &str = BUCKET_TABLE_CATALOG_META_PREFIX; const TABLE_BUCKET_ROOT: &str = BUCKET_TABLE_CATALOG_TABLE_BUCKETS_PREFIX; const COMMIT_LOG_ROOT: &str = "commits"; @@ -221,6 +223,28 @@ pub(crate) struct TableEntry { pub updated_at: Option, } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct ViewEntry { + pub version: u16, + pub table_bucket: String, + pub namespace: String, + pub view: String, + pub view_id: String, + pub view_uuid: String, + pub format: String, + pub format_version: u16, + pub warehouse_location: String, + pub metadata_location: String, + pub version_token: String, + pub generation: u64, + pub state: TableCatalogEntryState, + #[serde(default)] + pub properties: BTreeMap, + pub created_at: Option, + pub updated_at: Option, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "SCREAMING_SNAKE_CASE")] pub(crate) enum CommitLogStatus { @@ -273,6 +297,23 @@ pub(crate) struct TableCommitResult { pub commit_log: CommitLogEntry, } +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct ViewCommitRequest { + pub table_bucket: String, + pub namespace: String, + pub view: String, + pub expected_version_token: String, + pub expected_metadata_location: String, + pub new_metadata_location: String, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct ViewCommitResult { + pub view: ViewEntry, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub(crate) struct TableMaintenanceConfig { @@ -984,6 +1025,16 @@ pub(crate) trait TableCatalogStore: Send + Sync { async fn drop_table(&self, table_bucket: &str, namespace: &str, table: &str) -> TableCatalogStoreResult<()>; + async fn create_view(&self, entry: ViewEntry) -> TableCatalogStoreResult<()>; + + async fn list_views(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult>; + + async fn load_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult>; + + async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult; + + async fn drop_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult<()>; + async fn get_commit_by_id( &self, table_bucket: &str, @@ -1104,6 +1155,19 @@ impl TableCatalogObjectPaths { ) } + pub fn view_entries_prefix(&self, table_bucket: &str, namespace: &Namespace) -> String { + format!("{}{}/{}/", self.namespace_entries_prefix(table_bucket), namespace.storage_id(), VIEW_ROOT) + } + + pub fn view_entry_path(&self, table_bucket: &str, namespace: &Namespace, view: &IdentifierSegment) -> String { + format!( + "{}{}/{}", + self.view_entries_prefix(table_bucket, namespace), + view.as_str(), + VIEW_ENTRY_FILE + ) + } + pub fn table_maintenance_config_path( &self, table_bucket: &str, @@ -1338,6 +1402,25 @@ where Ok(Some((entry, etag))) } + async fn read_view_with_etag_unlocked( + &self, + table_bucket: &str, + namespace: &Namespace, + view: &IdentifierSegment, + ) -> TableCatalogStoreResult> { + let view_path = self.paths.view_entry_path(table_bucket, namespace, view); + let Some((entry, etag)) = self + .read_entry_unlocked::(self.catalog_bucket(), &view_path) + .await? + else { + return Ok(None); + }; + let Some(etag) = etag else { + return Err(TableCatalogStoreError::Internal(format!("catalog view entry has no etag: {view_path}"))); + }; + Ok(Some((entry, etag))) + } + async fn write_table_entry( &self, entry: TableEntry, @@ -1358,6 +1441,22 @@ where .await } + async fn write_view_entry(&self, entry: ViewEntry, precondition: TableCatalogPutPrecondition) -> TableCatalogStoreResult<()> { + validate_catalog_entry_version("view", entry.version)?; + self.require_table_bucket(&entry.table_bucket).await?; + let namespace = parse_namespace_for_store(&entry.namespace)?; + let view = parse_table_for_store(&entry.view)?; + if self.get_namespace(&entry.table_bucket, &entry.namespace).await?.is_none() { + return Err(TableCatalogStoreError::NotFound(format!( + "namespace {}/{}", + entry.table_bucket, entry.namespace + ))); + } + let view_path = self.paths.view_entry_path(&entry.table_bucket, &namespace, &view); + self.write_entry(self.catalog_bucket(), &view_path, &entry, precondition) + .await + } + async fn read_commit_by_path(&self, object: &str) -> TableCatalogStoreResult> { self.read_entry::(self.catalog_bucket(), object) .await @@ -2912,6 +3011,13 @@ where namespace.public_name() ))); } + if !self.list_views(table_bucket, &namespace.public_name()).await?.is_empty() { + return Err(TableCatalogStoreError::Conflict(format!( + "namespace {}/{} is not empty", + table_bucket, + namespace.public_name() + ))); + } self.backend .delete_object(self.catalog_bucket(), &self.paths.namespace_entry_path(table_bucket, &namespace)) .await @@ -3234,6 +3340,115 @@ where self.backend.delete_object(self.catalog_bucket(), &object).await } + async fn create_view(&self, entry: ViewEntry) -> TableCatalogStoreResult<()> { + self.write_view_entry(entry, TableCatalogPutPrecondition::IfAbsent).await + } + + async fn list_views(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult> { + let namespace = parse_namespace_for_store(namespace)?; + let mut entries = Vec::new(); + for object in self + .backend + .list_objects(self.catalog_bucket(), &self.paths.view_entries_prefix(table_bucket, &namespace)) + .await? + { + if !object.ends_with(VIEW_ENTRY_FILE) { + continue; + } + if let Some((entry, _)) = self.read_entry::(self.catalog_bucket(), &object).await? { + entries.push(entry); + } + } + entries.sort_by(|left, right| left.view.cmp(&right.view)); + Ok(entries) + } + + async fn load_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult> { + let namespace = parse_namespace_for_store(namespace)?; + let view = parse_table_for_store(view)?; + self.read_entry::(self.catalog_bucket(), &self.paths.view_entry_path(table_bucket, &namespace, &view)) + .await + .map(|entry| entry.map(|(view, _)| view)) + } + + async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult { + let namespace = parse_namespace_for_store(&request.namespace)?; + let view = parse_table_for_store(&request.view)?; + let view_path = self.paths.view_entry_path(&request.table_bucket, &namespace, &view); + let _guard = self.backend.acquire_write_lock(self.catalog_bucket(), &view_path).await?; + let Some((current, current_etag)) = self + .read_view_with_etag_unlocked(&request.table_bucket, &namespace, &view) + .await? + else { + return Err(TableCatalogStoreError::NotFound(format!( + "view {}/{}/{}", + request.table_bucket, request.namespace, request.view + ))); + }; + if current.version_token != request.expected_version_token { + return Err(TableCatalogStoreError::Conflict( + "current view version token does not match expected token".to_string(), + )); + } + if current.metadata_location != request.expected_metadata_location { + return Err(TableCatalogStoreError::Conflict( + "current view metadata location does not match expected location".to_string(), + )); + } + if !is_valid_view_metadata_location(&namespace, &view, &request.new_metadata_location) { + return Err(TableCatalogStoreError::Invalid( + "new metadata location must be inside the view metadata directory".to_string(), + )); + } + let Some(new_metadata_object) = self + .backend + .read_object(&request.table_bucket, &request.new_metadata_location) + .await? + else { + return Err(TableCatalogStoreError::NotFound(format!( + "new view metadata object {}", + request.new_metadata_location + ))); + }; + let next_warehouse_location = + metadata_warehouse_location(&request.table_bucket, &request.new_metadata_location, &new_metadata_object)?; + + let mut next = current; + next.metadata_location = request.new_metadata_location; + if let Some(warehouse_location) = next_warehouse_location { + next.warehouse_location = warehouse_location; + } + next.version_token = format!("token-{}", Uuid::new_v4()); + next.generation = next.generation.saturating_add(1); + self.write_entry_unlocked( + self.catalog_bucket(), + &view_path, + &next, + TableCatalogPutPrecondition::IfMatch(current_etag), + ) + .await?; + Ok(ViewCommitResult { view: next }) + } + + async fn drop_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult<()> { + let namespace = parse_namespace_for_store(namespace)?; + let view = parse_table_for_store(view)?; + let object = self.paths.view_entry_path(table_bucket, &namespace, &view); + if self + .load_view(table_bucket, &namespace.public_name(), view.as_str()) + .await? + .is_none() + { + return Err(TableCatalogStoreError::NotFound(format!( + "view {}/{}/{}", + table_bucket, + namespace.public_name(), + view.as_str() + ))); + } + self.backend.delete_object(self.catalog_bucket(), &object).await + } + async fn get_commit_by_id( &self, table_bucket: &str, @@ -6163,6 +6378,22 @@ pub(crate) fn default_table_delete_dir_path(namespace: &Namespace, table: &Ident format!("{}{}/{}", default_table_root_prefix(namespace), table.as_str(), DELETE_DIR) } +pub(crate) fn default_view_root_prefix(namespace: &Namespace) -> String { + format!("{}{}/{}/", default_namespace_root_prefix(), namespace.storage_id(), VIEW_ROOT) +} + +pub(crate) fn default_view_metadata_dir_path(namespace: &Namespace, view: &IdentifierSegment) -> String { + format!("{}{}/{}", default_view_root_prefix(namespace), view.as_str(), METADATA_DIR) +} + +pub(crate) fn default_view_metadata_file_path( + namespace: &Namespace, + view: &IdentifierSegment, + metadata_file_name: &str, +) -> String { + format!("{}/{}", default_view_metadata_dir_path(namespace, view), metadata_file_name) +} + pub(crate) fn default_table_metadata_file_path( namespace: &Namespace, table: &IdentifierSegment, @@ -6229,6 +6460,17 @@ pub(crate) fn is_valid_table_metadata_location( .is_some_and(is_valid_table_metadata_file_name) } +pub(crate) fn is_valid_view_metadata_location(namespace: &Namespace, view: &IdentifierSegment, metadata_location: &str) -> bool { + if metadata_location.is_empty() { + return false; + } + + let metadata_prefix = format!("{}/", default_view_metadata_dir_path(namespace, view)); + metadata_location + .strip_prefix(&metadata_prefix) + .is_some_and(is_valid_table_metadata_file_name) +} + pub(crate) fn is_valid_table_metadata_file_name(metadata_file_name: &str) -> bool { if metadata_file_name.is_empty() || metadata_file_name.len() > TABLE_METADATA_FILE_NAME_MAX_LEN @@ -6554,6 +6796,50 @@ mod tests { Ok(()) } + async fn create_view(&self, _entry: ViewEntry) -> TableCatalogStoreResult<()> { + Ok(()) + } + + async fn list_views(&self, _table_bucket: &str, _namespace: &str) -> TableCatalogStoreResult> { + Ok(Vec::new()) + } + + async fn load_view( + &self, + _table_bucket: &str, + _namespace: &str, + _view: &str, + ) -> TableCatalogStoreResult> { + Ok(None) + } + + async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult { + Ok(ViewCommitResult { + view: ViewEntry { + version: TABLE_CATALOG_ENTRY_VERSION, + table_bucket: request.table_bucket, + namespace: request.namespace, + view: request.view, + view_id: "view-id".to_string(), + view_uuid: "view-uuid".to_string(), + format: "ICEBERG_VIEW".to_string(), + format_version: 1, + warehouse_location: "s3://analytics/views/view-id".to_string(), + metadata_location: request.new_metadata_location, + version_token: "token-v2".to_string(), + generation: 2, + state: TableCatalogEntryState::Active, + properties: BTreeMap::new(), + created_at: None, + updated_at: None, + }, + }) + } + + async fn drop_view(&self, _table_bucket: &str, _namespace: &str, _view: &str) -> TableCatalogStoreResult<()> { + Ok(()) + } + async fn get_commit_by_id( &self, _table_bucket: &str, @@ -6636,6 +6922,10 @@ mod tests { paths.table_entry_path(bucket, &namespace, &table), format!("{bucket_root}namespaces/analytics/daily_events/tables/events/table-entry.json") ); + assert_eq!( + paths.view_entry_path(bucket, &namespace, &table), + format!("{bucket_root}namespaces/analytics/daily_events/views/events/view-entry.json") + ); let commit_path = paths.commit_log_entry_path("table/../bucket", "table/../id", "commit/%2f\nid"); let idempotency_path = paths.commit_idempotency_entry_path("table/../bucket", "table/../id", "client/%2f\nrequest"); @@ -7045,6 +7335,27 @@ mod tests { } } + fn test_view_entry(bucket: &str, namespace: &Namespace, view: &IdentifierSegment, metadata_location: String) -> ViewEntry { + ViewEntry { + version: TABLE_CATALOG_ENTRY_VERSION, + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + view: view.as_str().to_string(), + view_id: "view-id".to_string(), + view_uuid: "view-uuid".to_string(), + format: "ICEBERG_VIEW".to_string(), + format_version: 1, + warehouse_location: format!("s3://{bucket}/views/view-id"), + metadata_location, + version_token: "token-v1".to_string(), + generation: 1, + state: TableCatalogEntryState::Active, + properties: BTreeMap::new(), + created_at: None, + updated_at: None, + } + } + async fn seed_table_for_metadata_maintenance( store: &ObjectTableCatalogStore, bucket: &str, @@ -7078,6 +7389,75 @@ mod tests { assert_eq!(object_buckets, BTreeSet::from([rustfs_ecstore::disk::RUSTFS_META_BUCKET])); } + #[tokio::test] + async fn object_table_catalog_store_persists_view_entries_and_blocks_non_empty_namespace_drop() { + let backend = TestCatalogObjectBackend::default(); + let store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let view = IdentifierSegment::parse("recent_orders").unwrap(); + let current_metadata = default_view_metadata_file_path(&namespace, &view, "00001.metadata.json"); + let next_metadata = default_view_metadata_file_path(&namespace, &view, "00002.metadata.json"); + + store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap(); + store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .unwrap(); + store + .create_view(test_view_entry(bucket, &namespace, &view, current_metadata.clone())) + .await + .unwrap(); + + assert_eq!(store.list_views(bucket, &namespace.public_name()).await.unwrap()[0].view, "recent_orders"); + assert!( + store + .load_view(bucket, &namespace.public_name(), view.as_str()) + .await + .unwrap() + .is_some() + ); + assert!(matches!( + store.drop_namespace(bucket, &namespace.public_name()).await, + Err(TableCatalogStoreError::Conflict(_)) + )); + + backend + .seed_object( + bucket, + &next_metadata, + serde_json::to_vec(&serde_json::json!({ + "format-version": 1, + "view-uuid": "view-uuid", + "location": format!("s3://{bucket}/views/view-id") + })) + .unwrap(), + ) + .await; + let result = store + .replace_view(ViewCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + view: view.as_str().to_string(), + expected_version_token: "token-v1".to_string(), + expected_metadata_location: current_metadata, + new_metadata_location: next_metadata.clone(), + }) + .await + .unwrap(); + + assert_eq!(result.view.metadata_location, next_metadata); + assert_eq!(result.view.generation, 2); + assert_ne!(result.view.version_token, "token-v1"); + + store + .drop_view(bucket, &namespace.public_name(), view.as_str()) + .await + .unwrap(); + assert!(store.list_views(bucket, &namespace.public_name()).await.unwrap().is_empty()); + store.drop_namespace(bucket, &namespace.public_name()).await.unwrap(); + } + #[tokio::test] async fn maintenance_dry_run_keeps_current_metadata() { let backend = TestCatalogObjectBackend::default(); diff --git a/scripts/table-catalog/README.md b/scripts/table-catalog/README.md index ea48d17ed..e95faf2cc 100644 --- a/scripts/table-catalog/README.md +++ b/scripts/table-catalog/README.md @@ -120,6 +120,15 @@ work items. They are intentionally conservative: only PyIceberg is automated by this script today; other engines are documented until a repeatable harness is added. +RustFS also exposes catalog-backed advanced Iceberg surfaces that are not part +of the PyIceberg append smoke path yet: + +- table refs can be listed, created or replaced, and deleted through catalog + commits; refs with explicit retention policy require a forced delete, and + `main` cannot be deleted +- Iceberg views support basic create, list, load, replace, existence check, and + drop routes with persisted view metadata and view-scoped authorization + ## Client Matrix | Client | Current status | Claim | @@ -154,7 +163,6 @@ current unsupported inventory is: - snapshot expiration dry-run planning and manual catalog commit: supported through metadata maintenance reports - automatic maintenance scheduling: external scheduler hook supported through the worker run endpoint; built-in periodic scheduling is not claimed - compaction rewrite: controlled run-once support for unpartitioned Parquet binpack through metadata maintenance; built-in periodic scheduling, sort compaction, delete-file rewrite, and row-level compaction are not claimed -- Iceberg views: stable unsupported routes are registered and return explicit unsupported JSON - external catalog bridges: metadata import/register is supported, but Polaris/Glue/DLF/Hive synchronization is unsupported - multi-table transactions: not a short-term production claim diff --git a/scripts/table-catalog/pyiceberg_smoke.py b/scripts/table-catalog/pyiceberg_smoke.py index cb82469c4..f8e043d7b 100755 --- a/scripts/table-catalog/pyiceberg_smoke.py +++ b/scripts/table-catalog/pyiceberg_smoke.py @@ -156,12 +156,6 @@ UNSUPPORTED_INVENTORY: list[dict[str, str]] = [ "roadmap_area": "snapshot-maintenance", "expected_behavior": "metadata maintenance can plan binpack candidates and commit a safe unpartitioned Parquet rewrite through the catalog; built-in periodic scheduling, sort compaction, delete-file rewrite, and row-level compaction are not claimed", }, - { - "capability": "iceberg-views", - "status": "stable-unsupported-routes", - "roadmap_area": "view-api", - "expected_behavior": "view routes are registered and should return a stable unsupported JSON response until implemented", - }, { "capability": "external-catalog-bridge", "status": "metadata-import-only",