From a2fcac33cee3ac893c350b27aaa1798003df0332 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Thu, 24 Sep 2026 19:36:31 +0800 Subject: [PATCH 1/4] feat(datafusion): support the delete_branch procedure Java exposes `sys.delete_branch(table, branch)` to drop one or more comma-separated branches. The core BranchManager::drop_branch primitive already existed; only the DataFusion procedure was missing. Register delete_branch: declare its table/branch parameters (so a misspelled argument is rejected, like Java's binding) and dispatch to BranchManager, dropping each named branch and skipping one that does not exist -- the same forgiving, comma-splitting shape as delete_tag. --- .../integrations/datafusion/src/procedures.rs | 33 ++++++++++- .../datafusion/tests/procedures.rs | 58 ++++++++++++++++++- 2 files changed, 89 insertions(+), 2 deletions(-) diff --git a/crates/integrations/datafusion/src/procedures.rs b/crates/integrations/datafusion/src/procedures.rs index 369c81d9e..1c2c66c8e 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,35 @@ 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()); + for branch_name in branch_str.split(',') { + let branch_name = branch_name.trim(); + if branch_name.is_empty() { + continue; + } + 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..72e639275 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,56 @@ 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_create_lumina_index_requires_index_column() { let (_tmp, sql_context) = setup_table_with_snapshots().await; From 32bd9154e3682712ebcf2511d9ea684746aa345d Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Sat, 26 Sep 2026 09:44:23 +0800 Subject: [PATCH 2/4] fix(datafusion): reject deleting a branch used for scan reads Review follow-up (#941): `sys.delete_branch` called `BranchManager::drop_branch` directly, so it would drop a branch named by `scan.primary-branch` or `scan.fallback-branch` and silently break the table's read path. Paimon Java `AbstractFileStoreTable.deleteBranch` refuses this and asks the caller to unset the option first. Add a `BranchManager::ensure_branch_deletable` guard reading those two options and call it before dropping each branch. --- .../integrations/datafusion/src/procedures.rs | 3 + crates/paimon/src/table/branch_manager.rs | 57 +++++++++++++++++++ 2 files changed, 60 insertions(+) diff --git a/crates/integrations/datafusion/src/procedures.rs b/crates/integrations/datafusion/src/procedures.rs index 1c2c66c8e..57f0f2b67 100644 --- a/crates/integrations/datafusion/src/procedures.rs +++ b/crates/integrations/datafusion/src/procedures.rs @@ -465,11 +465,14 @@ async fn proc_delete_branch( 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; } + BranchManager::ensure_branch_deletable(options, branch_name) + .map_err(to_datafusion_error)?; if !bm .branch_exists(branch_name) .await diff --git a/crates/paimon/src/table/branch_manager.rs b/crates/paimon/src/table/branch_manager.rs index bd94cdb15..ebb88857c 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-"; @@ -203,6 +209,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 { @@ -501,6 +534,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(); From b5228f8a0a064fa7b814cb80bdf7103c7c2ddef5 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Wed, 30 Sep 2026 20:46:27 +0800 Subject: [PATCH 3/4] test(datafusion): cover delete_branch preserving scan-configured branches Add the SQL regressions requested in review: at the CALL boundary, `sys.delete_branch` must refuse a branch named by `scan.primary-branch` or `scan.fallback-branch` and leave it in place, while still deleting unrelated branches; and a comma-separated request must catch a protected branch even when it is not listed first. --- .../datafusion/tests/procedures.rs | 99 +++++++++++++++++++ 1 file changed, 99 insertions(+) diff --git a/crates/integrations/datafusion/tests/procedures.rs b/crates/integrations/datafusion/tests/procedures.rs index 72e639275..a379ab10b 100644 --- a/crates/integrations/datafusion/tests/procedures.rs +++ b/crates/integrations/datafusion/tests/procedures.rs @@ -141,6 +141,105 @@ async fn test_delete_branch() { .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_create_lumina_index_requires_index_column() { let (_tmp, sql_context) = setup_table_with_snapshots().await; From 40d1ed92ed9dd511aa9138b63e4da46081c14b01 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Fri, 2 Oct 2026 08:15:49 +0800 Subject: [PATCH 4/4] fix(datafusion): validate the delete_branch name before any deletion `proc_delete_branch` compared each name literally against the configured `scan.primary-branch`/`scan.fallback-branch`, so a separator-bearing name such as `prod/schema` passed the guard (it is not literally `prod`); `branch_exists` then matched the real `branch-prod/schema` directory and `drop_branch` deleted it, destroying the protected branch's schema. Validate each logical name first via the shared `BranchManager::validate_branch_name`, which now delegates to the catalog/table reader contract (rejects blank, `.`/`..`, path separators and control characters) and keeps the manager's main/numeric rules. The configured-branch guard then runs on a valid name. --- .../integrations/datafusion/src/procedures.rs | 5 +++ .../datafusion/tests/procedures.rs | 39 +++++++++++++++++++ crates/paimon/src/table/branch_manager.rs | 31 +++++++++------ 3 files changed, 63 insertions(+), 12 deletions(-) diff --git a/crates/integrations/datafusion/src/procedures.rs b/crates/integrations/datafusion/src/procedures.rs index 57f0f2b67..f3d3eec5f 100644 --- a/crates/integrations/datafusion/src/procedures.rs +++ b/crates/integrations/datafusion/src/procedures.rs @@ -471,6 +471,11 @@ async fn proc_delete_branch( 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 diff --git a/crates/integrations/datafusion/tests/procedures.rs b/crates/integrations/datafusion/tests/procedures.rs index a379ab10b..ef01bc8cd 100644 --- a/crates/integrations/datafusion/tests/procedures.rs +++ b/crates/integrations/datafusion/tests/procedures.rs @@ -240,6 +240,45 @@ async fn test_delete_branch_batch_stops_at_protected_branch() { ); } +#[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 ebb88857c..54c8e7499 100644 --- a/crates/paimon/src/table/branch_manager.rs +++ b/crates/paimon/src/table/branch_manager.rs @@ -70,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!( @@ -84,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!( @@ -425,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]