diff --git a/crates/integrations/datafusion/src/procedures.rs b/crates/integrations/datafusion/src/procedures.rs index 369c81d9e..f3d3eec5f 100644 --- a/crates/integrations/datafusion/src/procedures.rs +++ b/crates/integrations/datafusion/src/procedures.rs @@ -66,7 +66,7 @@ use paimon::catalog::{Catalog, Identifier, RESTCatalog}; use paimon::lumina::LUMINA_IDENTIFIER; use paimon::spec::Snapshot; use paimon::table::{ - normalize_global_index_type_for_drop, SnapshotManager, Table, TagManager, + normalize_global_index_type_for_drop, BranchManager, SnapshotManager, Table, TagManager, SUPPORTED_GLOBAL_INDEX_TYPES_FOR_DROP, }; use paimon::vindex::is_vindex_index_type; @@ -170,6 +170,7 @@ fn declared_parameters(proc_name: &str) -> Option<&'static [&'static str]> { "rollback_to" => &["table", "snapshot_id", "tag"], "rollback_to_timestamp" => &["table", "timestamp"], "create_tag_from_timestamp" => &["table", "tag", "timestamp"], + "delete_branch" => &["table", "branch"], "create_global_index" => &["table", "index_column", "index_type", "options"], // `partitions`/`dry_run` are declared but not yet implemented; they still reach // their own "not supported yet" error rather than being reported as unknown. @@ -283,6 +284,7 @@ pub async fn execute_call( "create_tag_from_timestamp" => { proc_create_tag_from_timestamp(ctx, catalog, catalog_name, &args).await } + "delete_branch" => proc_delete_branch(ctx, catalog, catalog_name, &args).await, "create_global_index" => proc_create_global_index(ctx, catalog, catalog_name, &args).await, "drop_global_index" => proc_drop_global_index(ctx, catalog, catalog_name, &args).await, "create_lumina_index" => proc_create_lumina_index(ctx, catalog, catalog_name, &args).await, @@ -453,6 +455,43 @@ async fn proc_create_tag( ok_result(ctx) } +async fn proc_delete_branch( + ctx: &SessionContext, + catalog: &Arc, + catalog_name: &str, + args: &HashMap, +) -> DFResult { + let table = get_table(catalog, catalog_name, args).await?; + let branch_str = require_arg(args, "branch")?; + + let bm = BranchManager::new(table.file_io().clone(), table.location().to_string()); + let options = table.schema().options(); + for branch_name in branch_str.split(',') { + let branch_name = branch_name.trim(); + if branch_name.is_empty() { + continue; + } + // Validate the logical name before any existence check or deletion: a + // separator-bearing name like `prod/schema` would otherwise slip past the + // configured-branch guard (it is not literally `prod`) and recursively + // delete the inner directory of a protected branch. + BranchManager::validate_branch_name(branch_name).map_err(to_datafusion_error)?; + BranchManager::ensure_branch_deletable(options, branch_name) + .map_err(to_datafusion_error)?; + if !bm + .branch_exists(branch_name) + .await + .map_err(to_datafusion_error)? + { + continue; + } + bm.drop_branch(branch_name) + .await + .map_err(to_datafusion_error)?; + } + ok_result(ctx) +} + async fn proc_delete_tag( ctx: &SessionContext, catalog: &Arc, diff --git a/crates/integrations/datafusion/tests/procedures.rs b/crates/integrations/datafusion/tests/procedures.rs index 73e4f143e..ef01bc8cd 100644 --- a/crates/integrations/datafusion/tests/procedures.rs +++ b/crates/integrations/datafusion/tests/procedures.rs @@ -17,7 +17,13 @@ mod common; -use common::{assert_sql_error, collect_id_name, exec, row_count, setup_sql_context}; +use common::{ + assert_sql_error, collect_id_name, create_sql_context, create_test_env, exec, row_count, + setup_sql_context, +}; +use paimon::catalog::Identifier; +use paimon::table::BranchManager; +use paimon::Catalog; async fn setup_table_with_snapshots() -> (tempfile::TempDir, paimon_datafusion::SQLContext) { let (tmp, sql_context) = setup_sql_context().await; @@ -85,6 +91,194 @@ async fn test_create_tag_with_snapshot_id() { assert_eq!(count, 1); } +#[tokio::test] +async fn test_delete_branch() { + let (_tmp, catalog) = create_test_env(); + let sql_context = create_sql_context(catalog.clone()).await; + exec(&sql_context, "CREATE SCHEMA paimon.test_db").await; + exec( + &sql_context, + "CREATE TABLE paimon.test_db.t1 (id INT, name VARCHAR(100), PRIMARY KEY (id))", + ) + .await; + exec( + &sql_context, + "INSERT INTO paimon.test_db.t1 VALUES (1, 'alice')", + ) + .await; + + // Seed a branch through the core manager (create_branch is a separate PR). + let table = catalog + .get_table(&Identifier::new("test_db", "t1")) + .await + .unwrap(); + let bm = BranchManager::new(table.file_io().clone(), table.location().to_string()); + bm.create_branch("b1").await.unwrap(); + assert!(bm.branch_exists("b1").await.unwrap(), "seed branch exists"); + + exec( + &sql_context, + "CALL sys.delete_branch(table => 'test_db.t1', branch => 'b1')", + ) + .await; + + assert!( + !bm.branch_exists("b1").await.unwrap(), + "delete_branch should remove the branch" + ); + let count = row_count( + &sql_context, + "SELECT * FROM paimon.test_db.`t1$branches` WHERE branch_name = 'b1'", + ) + .await; + assert_eq!(count, 0); + + // Deleting a branch that no longer exists is a no-op, not an error. + exec( + &sql_context, + "CALL sys.delete_branch(table => 'test_db.t1', branch => 'b1')", + ) + .await; +} + +#[tokio::test] +async fn test_delete_branch_preserves_scan_configured_branch() { + // A branch named by `scan.primary-branch` / `scan.fallback-branch` is on a + // reader's path, so `delete_branch` must refuse it (mirrors Java + // `AbstractFileStoreTable.deleteBranch`) while still allowing unrelated + // branches to be dropped. + for (table_name, option_key) in [ + ("tp", "scan.primary-branch"), + ("tf", "scan.fallback-branch"), + ] { + let (_tmp, catalog) = create_test_env(); + let sql_context = create_sql_context(catalog.clone()).await; + exec(&sql_context, "CREATE SCHEMA paimon.test_db").await; + exec( + &sql_context, + &format!( + "CREATE TABLE paimon.test_db.{table_name} (id INT, name VARCHAR(100), PRIMARY KEY (id)) WITH ('{option_key}' = 'prod')" + ), + ) + .await; + exec( + &sql_context, + &format!("INSERT INTO paimon.test_db.{table_name} VALUES (1, 'alice')"), + ) + .await; + + let table = catalog + .get_table(&Identifier::new("test_db", table_name)) + .await + .unwrap(); + let bm = BranchManager::new(table.file_io().clone(), table.location().to_string()); + bm.create_branch("prod").await.unwrap(); + bm.create_branch("tmp").await.unwrap(); + + // The configured branch cannot be deleted, and the error names the option. + assert_sql_error( + &sql_context, + &format!("CALL sys.delete_branch(table => 'test_db.{table_name}', branch => 'prod')"), + option_key, + ) + .await; + assert!( + bm.branch_exists("prod").await.unwrap(), + "{option_key}: protected branch must remain" + ); + + // An unrelated branch is still deletable, so the guard is not over-broad. + exec( + &sql_context, + &format!("CALL sys.delete_branch(table => 'test_db.{table_name}', branch => 'tmp')"), + ) + .await; + assert!( + !bm.branch_exists("tmp").await.unwrap(), + "{option_key}: unrelated branch should be deletable" + ); + } +} + +#[tokio::test] +async fn test_delete_branch_batch_stops_at_protected_branch() { + // A comma-separated request must still catch a protected branch even when + // it is not listed first. Matches Java `Table.deleteBranches`, which loops + // `deleteBranch` per name; the protected branch is never dropped. + let (_tmp, catalog) = create_test_env(); + let sql_context = create_sql_context(catalog.clone()).await; + exec(&sql_context, "CREATE SCHEMA paimon.test_db").await; + exec( + &sql_context, + "CREATE TABLE paimon.test_db.t1 (id INT, name VARCHAR(100), PRIMARY KEY (id)) WITH ('scan.primary-branch' = 'prod')", + ) + .await; + exec( + &sql_context, + "INSERT INTO paimon.test_db.t1 VALUES (1, 'alice')", + ) + .await; + + let table = catalog + .get_table(&Identifier::new("test_db", "t1")) + .await + .unwrap(); + let bm = BranchManager::new(table.file_io().clone(), table.location().to_string()); + bm.create_branch("prod").await.unwrap(); + bm.create_branch("keep").await.unwrap(); + + // 'prod' is protected and appears second; the call must fail and leave it. + assert_sql_error( + &sql_context, + "CALL sys.delete_branch(table => 'test_db.t1', branch => 'keep,prod')", + "scan.primary-branch", + ) + .await; + assert!( + bm.branch_exists("prod").await.unwrap(), + "protected branch must survive the batch" + ); +} + +#[tokio::test] +async fn test_delete_branch_rejects_path_separated_name() { + // A separator-bearing name slips past the configured-branch guard (it is not + // literally `prod`), so without name validation `branch_exists` finds the real + // `branch-prod/schema` directory and `drop_branch` deletes it, destroying the + // protected branch's schema. The name must be rejected before any deletion. + let (_tmp, catalog) = create_test_env(); + let sql_context = create_sql_context(catalog.clone()).await; + exec(&sql_context, "CREATE SCHEMA paimon.test_db").await; + exec( + &sql_context, + "CREATE TABLE paimon.test_db.t1 (id INT, name VARCHAR(100), PRIMARY KEY (id)) WITH ('scan.primary-branch' = 'prod')", + ) + .await; + exec( + &sql_context, + "INSERT INTO paimon.test_db.t1 VALUES (1, 'alice')", + ) + .await; + + let table = catalog + .get_table(&Identifier::new("test_db", "t1")) + .await + .unwrap(); + let bm = BranchManager::new(table.file_io().clone(), table.location().to_string()); + bm.create_branch("prod").await.unwrap(); + + assert_sql_error( + &sql_context, + "CALL sys.delete_branch(table => 'test_db.t1', branch => 'prod/schema')", + "path separator", + ) + .await; + + // The protected branch and its schema survive and stay openable. + assert!(bm.branch_exists("prod").await.unwrap(), "prod must remain"); + table.copy_with_branch("prod").await.unwrap(); +} + #[tokio::test] async fn test_create_lumina_index_requires_index_column() { let (_tmp, sql_context) = setup_table_with_snapshots().await; diff --git a/crates/paimon/src/table/branch_manager.rs b/crates/paimon/src/table/branch_manager.rs index bd94cdb15..54c8e7499 100644 --- a/crates/paimon/src/table/branch_manager.rs +++ b/crates/paimon/src/table/branch_manager.rs @@ -24,6 +24,12 @@ use crate::catalog::DEFAULT_MAIN_BRANCH; use crate::io::FileIO; use crate::table::{SchemaManager, SnapshotManager, TagManager}; use opendal::raw::get_basename; +use std::collections::HashMap; + +/// Table option a reader falls back to when the primary branch has no data. +const SCAN_FALLBACK_BRANCH: &str = "scan.fallback-branch"; +/// Table option naming the branch a reader treats as primary. +const SCAN_PRIMARY_BRANCH: &str = "scan.primary-branch"; const BRANCH_DIR: &str = "branch"; const BRANCH_PREFIX: &str = "branch-"; @@ -64,11 +70,12 @@ impl BranchManager { /// Validate branch name format. /// - /// Rules: - /// - Cannot be "main" - /// - Cannot be blank or whitespace only - /// - Cannot be a pure numeric string - fn validate_branch_name(branch_name: &str) -> crate::Result<()> { + /// Enforces the catalog/table reader contract (`copy_with_branch`, + /// `$branch_...` resolution): rejects blank, `.`/`..`, path separators and + /// control characters, plus the manager's own main/pure-numeric rules. A + /// name accepted here is always openable by a reader. + pub fn validate_branch_name(branch_name: &str) -> crate::Result<()> { + crate::catalog::validate_branch_name(branch_name)?; if branch_name == DEFAULT_MAIN_BRANCH { return Err(crate::Error::DataInvalid { message: format!( @@ -78,12 +85,6 @@ impl BranchManager { source: None, }); } - if branch_name.trim().is_empty() { - return Err(crate::Error::DataInvalid { - message: format!("Branch name '{}' is blank.", branch_name), - source: None, - }); - } if branch_name.chars().all(|c| c.is_ascii_digit()) { return Err(crate::Error::DataInvalid { message: format!( @@ -203,6 +204,33 @@ impl BranchManager { Ok(()) } + /// Reject deleting a branch a reader is configured to consult. + /// + /// A branch named by `scan.primary-branch` or `scan.fallback-branch` is part + /// of a table's read path; dropping it would break reads that fall back to + /// it. Mirrors Java `AbstractFileStoreTable.deleteBranch`, which refuses the + /// deletion and asks the caller to unset the option first. Callers pass the + /// table options because these keys live on the schema, not the manager. + pub fn ensure_branch_deletable( + options: &HashMap, + branch_name: &str, + ) -> crate::Result<()> { + for key in [SCAN_PRIMARY_BRANCH, SCAN_FALLBACK_BRANCH] { + if options + .get(key) + .is_some_and(|configured| configured == branch_name) + { + return Err(crate::Error::DataInvalid { + message: format!( + "Cannot delete branch '{branch_name}' because it is configured as '{key}'. Unset '{key}' first." + ), + source: None, + }); + } + } + Ok(()) + } + /// Rename an existing branch. pub async fn rename_branch(&self, from: &str, to: &str) -> crate::Result<()> { if from == DEFAULT_MAIN_BRANCH { @@ -392,7 +420,19 @@ mod tests { let result = BranchManager::validate_branch_name(""); assert!(result.is_err()); let msg = format!("{}", result.unwrap_err()); - assert!(msg.contains("blank")); + assert!(msg.contains("empty"), "got: {msg}"); + } + + #[tokio::test] + async fn test_validate_branch_name_rejects_reader_unopenable_names() { + // `.`/`..`, path separators and control chars must be rejected to match + // the catalog/table reader contract. + for name in [".", "..", "a/b", "a\\b", "a\u{001C}b"] { + assert!( + BranchManager::validate_branch_name(name).is_err(), + "{name:?} should be rejected" + ); + } } #[tokio::test] @@ -501,6 +541,30 @@ mod tests { assert!(msg.contains("doesn't exist")); } + #[test] + fn test_ensure_branch_deletable_rejects_scan_branches() { + for key in ["scan.primary-branch", "scan.fallback-branch"] { + let mut options = HashMap::new(); + options.insert(key.to_string(), "prod".to_string()); + + // The branch the option points at is part of a read path. + let err = BranchManager::ensure_branch_deletable(&options, "prod").unwrap_err(); + let msg = format!("{err}"); + assert!(msg.contains("prod"), "{msg}"); + assert!(msg.contains(key), "{msg}"); + + // A different branch is unaffected even while the option is set, + // so the guard keys on the configured value, not its mere presence. + BranchManager::ensure_branch_deletable(&options, "feature").unwrap(); + } + } + + #[test] + fn test_ensure_branch_deletable_without_scan_options() { + let options = HashMap::new(); + BranchManager::ensure_branch_deletable(&options, "any").unwrap(); + } + #[tokio::test] async fn test_rename_branch() { let file_io = test_file_io();