diff --git a/crates/integrations/datafusion/src/procedures.rs b/crates/integrations/datafusion/src/procedures.rs index 369c81d9e..1eb4ceb00 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"], + "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. @@ -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, @@ -453,6 +455,42 @@ async fn proc_create_tag( ok_result(ctx) } +async fn proc_create_branch( + ctx: &SessionContext, + catalog: &Arc, + catalog_name: &str, + args: &HashMap, +) -> DFResult { + 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, diff --git a/crates/integrations/datafusion/tests/procedures.rs b/crates/integrations/datafusion/tests/procedures.rs index 73e4f143e..302234b99 100644 --- a/crates/integrations/datafusion/tests/procedures.rs +++ b/crates/integrations/datafusion/tests/procedures.rs @@ -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; diff --git a/crates/paimon/src/table/branch_manager.rs b/crates/paimon/src/table/branch_manager.rs index bd94cdb15..bccfe73a1 100644 --- a/crates/paimon/src/table/branch_manager.rs +++ b/crates/paimon/src/table/branch_manager.rs @@ -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 would place 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!( @@ -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!( @@ -176,12 +176,19 @@ impl BranchManager { let tag_dst = tag_manager.with_branch(branch_name).tag_path(tag_name); self.file_io.copy_file(&tag_src, &tag_dst).await?; - // Copy snapshot file to branch + // Copy snapshot file to branch. A tag can outlive its main snapshot JSON + // (expiration keeps tagged data but may delete `snapshot/snapshot-`), + // so when the live file is gone, materialize the snapshot already resolved + // from the tag instead of failing. Mirrors Java + // `FileSystemBranchManager.createBranch`. let snap_src = snapshot_manager.snapshot_path(snapshot.id()); - let snap_dst = snapshot_manager - .with_branch(branch_name) - .snapshot_path(snapshot.id()); - self.file_io.copy_file(&snap_src, &snap_dst).await?; + let branch_snapshot_manager = snapshot_manager.with_branch(branch_name); + let snap_dst = branch_snapshot_manager.snapshot_path(snapshot.id()); + if self.file_io.exists(&snap_src).await? { + self.file_io.copy_file(&snap_src, &snap_dst).await?; + } else { + branch_snapshot_manager.commit_snapshot(&snapshot).await?; + } // Copy schemas to branch self.copy_schemas_to_branch(branch_name, snapshot.schema_id()) @@ -392,7 +399,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] @@ -409,6 +428,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(); @@ -476,6 +506,41 @@ mod tests { assert_eq!(schemas.len(), 1); } + #[tokio::test] + async fn test_create_branch_from_tag_materializes_missing_live_snapshot() { + // A tag can outlive its main snapshot JSON: expiration deletes + // `snapshot/snapshot-` but keeps the tag, schema and data. Creating a + // branch from such a tag must materialize the tag's snapshot, not fail + // copying a missing file. + let file_io = test_file_io(); + let table_path = "memory:/test_create_branch_tag_retained".to_string(); + let schema_manager = SchemaManager::new(file_io.clone(), table_path.clone()); + let snapshot_manager = SnapshotManager::new(file_io.clone(), table_path.clone()); + let tag_manager = TagManager::new(file_io.clone(), table_path.clone()); + + write_schema(&file_io, &schema_manager, &test_schema()).await; + let snap = test_snapshot(1); + write_snapshot(&file_io, &snapshot_manager, &snap).await; + write_tag(&file_io, &tag_manager, "v1", &snap).await; + + // Expire the live snapshot JSON while keeping the tag. + snapshot_manager.delete_snapshot(1).await.unwrap(); + assert!( + !file_io + .exists(&snapshot_manager.snapshot_path(1)) + .await + .unwrap(), + "live snapshot JSON should be gone" + ); + + let bm = BranchManager::new(file_io.clone(), table_path.clone()); + bm.create_branch_from_tag("retained", "v1").await.unwrap(); + + // The branch snapshot was materialized from the tag and is readable. + let branch_snap_manager = snapshot_manager.with_branch("retained"); + assert_eq!(branch_snap_manager.get_snapshot(1).await.unwrap().id(), 1); + } + #[tokio::test] async fn test_drop_branch() { let file_io = test_file_io();