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
21 changes: 20 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"],
"rename_branch" => &["table", "from_branch", "to_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
}
"rename_branch" => proc_rename_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,23 @@ async fn proc_create_tag(
ok_result(ctx)
}

async fn proc_rename_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 from_branch = require_arg(args, "from_branch")?;
let to_branch = require_arg(args, "to_branch")?;

let bm = BranchManager::new(table.file_io().clone(), table.location().to_string());
bm.rename_branch(from_branch, to_branch)
.await
.map_err(to_datafusion_error)?;
ok_result(ctx)
}

async fn proc_delete_tag(
ctx: &SessionContext,
catalog: &Arc<dyn Catalog>,
Expand Down
160 changes: 159 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,158 @@ async fn test_create_tag_with_snapshot_id() {
assert_eq!(count, 1);
}

#[tokio::test]
async fn test_rename_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();

exec(
&sql_context,
"CALL sys.rename_branch(table => 'test_db.t1', from_branch => 'b1', to_branch => 'b2')",
)
.await;

assert!(!bm.branch_exists("b1").await.unwrap(), "old branch gone");
assert!(bm.branch_exists("b2").await.unwrap(), "new branch present");
let old = row_count(
&sql_context,
"SELECT * FROM paimon.test_db.`t1$branches` WHERE branch_name = 'b1'",
)
.await;
assert_eq!(old, 0);
let new = row_count(
&sql_context,
"SELECT * FROM paimon.test_db.`t1$branches` WHERE branch_name = 'b2'",
)
.await;
assert_eq!(new, 1);

// Renaming a branch that does not exist is an error.
assert_sql_error(
&sql_context,
"CALL sys.rename_branch(table => 'test_db.t1', from_branch => 'b1', to_branch => 'b3')",
"doesn't exist",
)
.await;
}

#[tokio::test]
async fn test_rename_branch_rejects_path_separator() {
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;

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();

// Renaming to `foo/bar` would move the branch under `branch-foo/` and hide it
// from `$branches`; it must be rejected and leave `b1` untouched.
assert_sql_error(
&sql_context,
"CALL sys.rename_branch(table => 'test_db.t1', from_branch => 'b1', to_branch => 'foo/bar')",
"path separator",
)
.await;

assert!(bm.branch_exists("b1").await.unwrap(), "b1 must remain");
assert!(
!bm.branch_exists("foo/bar").await.unwrap(),
"foo/bar must not exist"
);
let visible = row_count(
&sql_context,
"SELECT * FROM paimon.test_db.`t1$branches` WHERE branch_name = 'b1'",
)
.await;
assert_eq!(visible, 1);
}

#[tokio::test]
async fn test_rename_branch_rejects_unopenable_source_and_target() {
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;

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();

// A path-separated SOURCE would rename `branch-b1/schema` (an inner directory),
// not the branch, orphaning b1's metadata. It must be rejected before any move.
assert_sql_error(
&sql_context,
"CALL sys.rename_branch(table => 'test_db.t1', from_branch => 'b1/schema', to_branch => 'stolen')",
"path separator",
)
.await;
assert!(bm.branch_exists("b1").await.unwrap(), "b1 must remain");
assert!(
!bm.branch_exists("stolen").await.unwrap(),
"stolen must not exist"
);

// A `..` TARGET is a single directory segment but no reader can open it, so the
// rename must be rejected rather than moving b1 to an unopenable name.
assert_sql_error(
&sql_context,
"CALL sys.rename_branch(table => 'test_db.t1', from_branch => 'b1', to_branch => '..')",
"'.' or '..'",
)
.await;
assert!(
bm.branch_exists("b1").await.unwrap(),
"b1 must survive a rejected rename"
);
// b1 is still openable through the table reader.
table.copy_with_branch("b1").await.unwrap();
}

#[tokio::test]
async fn test_create_lumina_index_requires_index_column() {
let (_tmp, sql_context) = setup_table_with_snapshots().await;
Expand Down
42 changes: 35 additions & 7 deletions crates/paimon/src/table/branch_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,12 @@ impl BranchManager {
/// - Cannot be blank or whitespace only
/// - Cannot be a pure numeric string
fn validate_branch_name(branch_name: &str) -> crate::Result<()> {
// Enforce the same name contract the catalog/table reader applies
// (`copy_with_branch` and `$branch_...` resolution): reject blank,
// `.`/`..`, path separators and control characters. A name accepted here
// must be openable by a reader; otherwise create/rename/delete would move
// metadata under a directory the table API can never resolve.
crate::catalog::validate_branch_name(branch_name)?;
if branch_name == DEFAULT_MAIN_BRANCH {
return Err(crate::Error::DataInvalid {
message: format!(
Expand All @@ -78,12 +84,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!(
Expand Down Expand Up @@ -211,6 +211,10 @@ impl BranchManager {
source: None,
});
}
// Validate the source name before touching the filesystem: a malformed
// logical name like `b1/schema` would otherwise rename a branch's inner
// directory, not the branch, orphaning its metadata.
Self::validate_branch_name(from)?;
if !self.branch_exists(from).await? {
return Err(crate::Error::DataInvalid {
message: format!("Branch name '{}' doesn't exist.", from),
Expand Down Expand Up @@ -392,7 +396,7 @@ 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]
Expand All @@ -409,6 +413,30 @@ mod tests {
assert!(BranchManager::validate_branch_name("branch-1").is_ok());
}

#[tokio::test]
async fn test_validate_branch_name_rejects_path_separator() {
// A '/'-bearing rename target creates an unlistable, orphaned branch.
for name in ["b1/hidden", "a\\b"] {
let result = BranchManager::validate_branch_name(name);
assert!(result.is_err(), "'{name}' should be rejected");
let msg = format!("{}", result.unwrap_err());
assert!(msg.contains("path separator"), "got: {msg}");
}
}

#[tokio::test]
async fn test_validate_branch_name_rejects_reader_unopenable_names() {
// Names the catalog/table reader rejects (`.`, `..`, control chars) must
// also be rejected here, so a created/renamed branch is always openable
// via `copy_with_branch` / `$branch_...`.
for name in [".", "..", "a\u{0007}b", "a\u{001C}b"] {
assert!(
BranchManager::validate_branch_name(name).is_err(),
"{name:?} should be rejected"
);
}
}

#[tokio::test]
async fn test_create_branch() {
let file_io = test_file_io();
Expand Down
Loading