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
40 changes: 39 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"],
"create_branch" => &["table", "branch", "tag", "ignore_if_exists"],
"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
}
"create_branch" => proc_create_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,42 @@ async fn proc_create_tag(
ok_result(ctx)
}

async fn proc_create_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_name = require_arg(args, "branch")?;
let tag = args.get("tag").map(String::as_str);
let ignore_if_exists = args
.get("ignore_if_exists")
.map(|s| s.eq_ignore_ascii_case("true"))
.unwrap_or(false);

let bm = BranchManager::new(table.file_io().clone(), table.location().to_string());
if ignore_if_exists
&& bm
.branch_exists(branch_name)
.await
.map_err(to_datafusion_error)?
{
return ok_result(ctx);
}
match tag {
Some(tag_name) => bm
.create_branch_from_tag(branch_name, tag_name)
.await
.map_err(to_datafusion_error)?,
None => bm
.create_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
68 changes: 68 additions & 0 deletions crates/integrations/datafusion/tests/procedures.rs
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,74 @@ async fn test_create_tag_with_snapshot_id() {
assert_eq!(count, 1);
}

#[tokio::test]
async fn test_create_branch() {
let (_tmp, sql_context) = setup_table_with_snapshots().await;

exec(
&sql_context,
"CALL sys.create_branch(table => 'test_db.t1', branch => 'b1')",
)
.await;

// The branch now shows up in the $branches system table.
let count = row_count(
&sql_context,
"SELECT * FROM paimon.test_db.`t1$branches` WHERE branch_name = 'b1'",
)
.await;
assert_eq!(count, 1);

// Recreating the same branch without ignore_if_exists is an error...
assert_sql_error(
&sql_context,
"CALL sys.create_branch(table => 'test_db.t1', branch => 'b1')",
"already exists",
)
.await;

// ...while ignore_if_exists makes it a no-op that leaves the one branch.
exec(
&sql_context,
"CALL sys.create_branch(table => 'test_db.t1', branch => 'b1', ignore_if_exists => true)",
)
.await;
let count = row_count(
&sql_context,
"SELECT * FROM paimon.test_db.`t1$branches` WHERE branch_name = 'b1'",
)
.await;
assert_eq!(count, 1);
}

#[tokio::test]
async fn test_create_branch_rejects_path_separator() {
let (tmp, sql_context) = setup_table_with_snapshots().await;
exec(
&sql_context,
"CALL sys.create_branch(table => 'test_db.t1', branch => 'b1')",
)
.await;

// `b1/hidden` would be created under `branch-b1/` and never be listed as a branch.
assert_sql_error(
&sql_context,
"CALL sys.create_branch(table => 'test_db.t1', branch => 'b1/hidden')",
"path separator",
)
.await;

// The name is rejected before any schema file is copied into the nested path.
let b1_dir = tmp.path().join("test_db.db/t1/branch/branch-b1");
assert!(b1_dir.is_dir(), "{} is missing", b1_dir.display());
assert!(
!b1_dir.join("hidden").exists(),
"the nested branch was created"
);
let count = row_count(&sql_context, "SELECT * FROM paimon.test_db.`t1$branches`").await;
assert_eq!(count, 1);
}

#[tokio::test]
async fn test_create_lumina_index_requires_index_column() {
let (_tmp, sql_context) = setup_table_with_snapshots().await;
Expand Down
22 changes: 22 additions & 0 deletions crates/paimon/src/table/branch_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,17 @@ impl BranchManager {
source: None,
});
}
// A path separator would place the branch directory at a nested path
// (`branch-b1/hidden`), so the branch is created but never listed back by
// `$branches`, leaving it silently orphaned. Reject it up front.
if branch_name.contains('/') || branch_name.contains('\\') {
return Err(crate::Error::DataInvalid {
message: format!(
"Branch name '{branch_name}' must not contain a path separator ('/' or '\\')."
),
source: None,
});
}
Ok(())
}

Expand Down Expand Up @@ -409,6 +420,17 @@ mod tests {
assert!(BranchManager::validate_branch_name("branch-1").is_ok());
}

#[tokio::test]
async fn test_validate_branch_name_rejects_path_separator() {
// A '/'-bearing name creates an unlistable, orphaned branch directory.
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_create_branch() {
let file_io = test_file_io();
Expand Down
Loading