Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 35 additions & 1 deletion crates/integrations/datafusion/src/procedures.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -453,6 +455,38 @@ async fn proc_create_tag(
ok_result(ctx)
}

async fn proc_delete_branch(
ctx: &SessionContext,
catalog: &Arc<dyn Catalog>,
catalog_name: &str,
args: &HashMap<String, String>,
) -> DFResult<DataFrame> {
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;
}
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<dyn Catalog>,
Expand Down
157 changes: 156 additions & 1 deletion crates/integrations/datafusion/tests/procedures.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -85,6 +91,155 @@ 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_create_lumina_index_requires_index_column() {
let (_tmp, sql_context) = setup_table_with_snapshots().await;
Expand Down
57 changes: 57 additions & 0 deletions crates/paimon/src/table/branch_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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-";
Expand Down Expand Up @@ -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<String, String>,
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 {
Expand Down Expand Up @@ -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();
Expand Down
Loading