Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
b5c808e
feat(write): route postpone batch writes to fixed buckets
XiaoHongbo-Hope Aug 3, 2026
a18eede
fix(write): make fixed-bucket postpone writers one-shot
XiaoHongbo-Hope Aug 3, 2026
dd300d5
fix(write): align postpone fixed-bucket planning
XiaoHongbo-Hope Aug 3, 2026
2bd3e68
fix(write): isolate postpone batch planning
XiaoHongbo-Hope Aug 5, 2026
7c63459
fix(write): share postpone bucket plans
XiaoHongbo-Hope Aug 5, 2026
90a1b72
fix(write): enable fixed buckets for batch writes
XiaoHongbo-Hope Aug 5, 2026
dd8008f
fix(write): make postpone fixed buckets explicit
XiaoHongbo-Hope Aug 6, 2026
a3d8266
fix(spec): remove unused fixed-bucket toggle
XiaoHongbo-Hope Aug 6, 2026
6031680
docs(write): trim postpone fixed-bucket comments
XiaoHongbo-Hope Aug 8, 2026
024ea43
refactor(write): simplify postpone fixed-bucket planning
XiaoHongbo-Hope Aug 8, 2026
2ec4d3c
refactor(write): reduce fixed-bucket boilerplate
XiaoHongbo-Hope Aug 8, 2026
25aa45e
refactor(write): consolidate fixed-bucket implementation
XiaoHongbo-Hope Aug 8, 2026
291e6a9
test(write): consolidate fixed-bucket coverage
XiaoHongbo-Hope Aug 8, 2026
afbb593
fix(write): enforce fixed-bucket writer ownership
XiaoHongbo-Hope Aug 10, 2026
bc852ec
refactor(write): separate fixed-bucket planning
XiaoHongbo-Hope Aug 10, 2026
65bd2d0
refactor(c): align fixed-bucket commit wiring
XiaoHongbo-Hope Aug 10, 2026
0cb0b20
refactor(write): require explicit postpone bucket plan
XiaoHongbo-Hope Aug 10, 2026
cd206dd
refactor(write): split postpone fixed-bucket modules
XiaoHongbo-Hope Aug 10, 2026
0b96b9f
refactor(c): model fixed-bucket builders explicitly
XiaoHongbo-Hope Aug 10, 2026
7668df7
refactor(write): keep postpone router internal
XiaoHongbo-Hope Aug 10, 2026
aa1215c
refactor(write): remove redundant constructor layer
XiaoHongbo-Hope Aug 10, 2026
5112cfc
refactor(commit): isolate postpone ownership validation
XiaoHongbo-Hope Aug 10, 2026
22bb670
test(write): colocate fixed-bucket coverage
XiaoHongbo-Hope Aug 10, 2026
df30d8f
test(write): colocate fixed-bucket builder coverage
XiaoHongbo-Hope Aug 10, 2026
fa0ef4d
fix(commit): enforce fixed-bucket ownership globally
XiaoHongbo-Hope Aug 10, 2026
516231a
test(write): cover fixed-bucket delete rows
XiaoHongbo-Hope Aug 10, 2026
633efff
fix(commit): reject concurrent fixed-bucket writers
XiaoHongbo-Hope Aug 10, 2026
dfc5092
fix(write): reject cross-partition postpone writes
XiaoHongbo-Hope Aug 10, 2026
f8c42fc
fix(write): harden postpone fixed-bucket commits
XiaoHongbo-Hope Aug 10, 2026
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
347 changes: 346 additions & 1 deletion bindings/c/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ use arrow_array::{Array, Int32Array, RecordBatch, StringArray, StructArray};
use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as ArrowSchema};
use paimon::catalog::Identifier;
use paimon::io::FileIOBuilder;
use paimon::spec::{DataType, IntType, Schema, TableSchema, VarCharType};
use paimon::spec::{CommitKind, DataType, IntType, Schema, TableSchema, VarCharType};
use paimon::table::{SnapshotManager, Table};

use crate::error::*;
Expand Down Expand Up @@ -73,6 +73,19 @@ fn not_null_table_schema() -> TableSchema {
TableSchema::new(0, &schema)
}

fn partitioned_postpone_table_schema() -> TableSchema {
let schema = Schema::builder()
.column("pt", DataType::VarChar(VarCharType::string_type()))
.column("id", DataType::Int(IntType::new()))
.column("name", DataType::VarChar(VarCharType::string_type()))
.primary_key(["pt", "id"])
.partition_keys(["pt"])
.option("bucket", "-2")
.build()
.unwrap();
TableSchema::new(0, &schema)
}

unsafe fn wrap_table(table: Table) -> *mut paimon_table {
let inner = Box::into_raw(Box::new(table)) as *mut c_void;
Box::into_raw(Box::new(paimon_table { inner }))
Expand Down Expand Up @@ -104,6 +117,38 @@ fn make_batch(ids: Vec<i32>, names: Vec<&str>) -> RecordBatch {
.unwrap()
}

fn make_partitioned_write_batch(pts: Vec<&str>, ids: Vec<i32>, names: Vec<&str>) -> RecordBatch {
let schema = Arc::new(ArrowSchema::new(vec![
ArrowField::new("pt", ArrowDataType::Utf8, false),
ArrowField::new("id", ArrowDataType::Int32, false),
ArrowField::new("name", ArrowDataType::Utf8, true),
]));
RecordBatch::try_new(
schema,
vec![
Arc::new(StringArray::from(pts)),
Arc::new(Int32Array::from(ids)),
Arc::new(StringArray::from(names)),
],
)
.unwrap()
}

fn make_postpone_bucket_plan_batch(partitions: Vec<&str>, counts: Vec<i32>) -> RecordBatch {
let schema = Arc::new(ArrowSchema::new(vec![
ArrowField::new("pt", ArrowDataType::Utf8, false),
ArrowField::new("total_buckets", ArrowDataType::Int32, false),
]));
RecordBatch::try_new(
schema,
vec![
Arc::new(StringArray::from(partitions)),
Arc::new(Int32Array::from(counts)),
],
)
.unwrap()
}

fn make_type_mismatch_batch(ids: Vec<&str>, names: Vec<&str>) -> RecordBatch {
let schema = Arc::new(ArrowSchema::new(vec![
ArrowField::new("id", ArrowDataType::Utf8, false),
Expand Down Expand Up @@ -1271,6 +1316,101 @@ fn test_commit_rejects_messages_from_different_builder_identity() {
}
}

#[test]
fn test_commit_rejects_mismatched_write_kind_and_overwrite_mode() {
let path = "memory:/test_commit_write_context";
let file_io = memory_file_io();
setup_table_dirs(&file_io, path);
let table = Table::new(
file_io.clone(),
Identifier::new("default", "test"),
path.to_string(),
partitioned_postpone_table_schema(),
None,
);
let handle = unsafe { wrap_table(table) };
let commit_user = CString::new("write-context-job").unwrap();

unsafe {
let standard_wb =
paimon_table_new_write_builder_with_commit_user(handle, commit_user.as_ptr())
.write_builder;
let standard_tw = paimon_write_builder_new_write(standard_wb).write;
let (array, schema) = export_batch_to_ffi(make_partitioned_write_batch(
vec!["p"],
vec![1],
vec!["standard"],
));
assert!(paimon_table_write_write_arrow_batch(
standard_tw,
(&**array) as *const FFI_ArrowArray as *mut c_void,
(&**schema) as *const FFI_ArrowSchema as *mut c_void,
)
.is_null());
let standard_messages = paimon_table_write_prepare_commit(standard_tw).messages;

let fixed_wb = paimon_table_new_postpone_fixed_bucket_write_builder_with_commit_user(
handle,
commit_user.as_ptr(),
)
.write_builder;
assert!(paimon_write_builder_with_overwrite(fixed_wb).is_null());
let (array, schema) =
export_batch_to_ffi(make_postpone_bucket_plan_batch(vec!["p"], vec![1]));
assert!(paimon_write_builder_with_postpone_bucket_plan(
fixed_wb,
(&**array) as *const FFI_ArrowArray as *mut c_void,
(&**schema) as *const FFI_ArrowSchema as *mut c_void,
)
.is_null());
let fixed_tw = paimon_write_builder_new_write(fixed_wb).write;
let (array, schema) = export_batch_to_ffi(make_partitioned_write_batch(
vec!["p"],
vec![1],
vec!["fixed"],
));
assert!(paimon_table_write_write_arrow_batch(
fixed_tw,
(&**array) as *const FFI_ArrowArray as *mut c_void,
(&**schema) as *const FFI_ArrowSchema as *mut c_void,
)
.is_null());
let fixed_messages = paimon_table_write_prepare_commit(fixed_tw).messages;

let error = paimon_commit_messages_merge(standard_messages, fixed_messages);
assert!(!error.is_null());
assert!(error_message(error).contains("write kind and overwrite mode"));
paimon_error_free(error);

let fixed_commit = paimon_write_builder_new_commit(fixed_wb).commit;
let error = paimon_table_commit_commit(fixed_commit, standard_messages);
assert!(!error.is_null());
assert!(error_message(error).contains("different write kind or overwrite mode"));
paimon_error_free(error);

let standard_commit = paimon_write_builder_new_commit(standard_wb).commit;
let error = paimon_table_commit_overwrite(standard_commit, standard_messages);
assert!(!error.is_null());
assert!(error_message(error).contains("append messages cannot be committed"));
paimon_error_free(error);

assert!(crate::runtime()
.block_on(SnapshotManager::new(file_io, path.to_string()).get_latest_snapshot())
.unwrap()
.is_none());

paimon_table_commit_free(standard_commit);
paimon_table_commit_free(fixed_commit);
paimon_commit_messages_free(fixed_messages);
paimon_commit_messages_free(standard_messages);
paimon_table_write_free(fixed_tw);
paimon_write_builder_free(fixed_wb);
paimon_table_write_free(standard_tw);
paimon_write_builder_free(standard_wb);
unwrap_table(handle);
}
}

#[test]
fn test_commit_messages_live_until_explicit_free() {
const CHILD_ENV: &str = "PAIMON_C_MESSAGES_LIFETIME_CHILD";
Expand Down Expand Up @@ -1442,6 +1582,198 @@ fn test_commit_messages_merge_preserves_all_writer_files() {
}
}

#[test]
fn test_distributed_postpone_writers_share_bucket_plan() {
let path = "memory:/test_distributed_postpone_bucket_plan";
let file_io = memory_file_io();
setup_table_dirs(&file_io, path);
let table = Table::new(
file_io,
Identifier::new("default", "test"),
path.to_string(),
partitioned_postpone_table_schema(),
None,
);
let handle = unsafe { wrap_table(table) };
let commit_user = CString::new("distributed-postpone-job").unwrap();

unsafe {
let normal = paimon_table_new_write_builder(handle);
assert!(normal.error.is_null());
assert!(matches!(
(&*((*normal.write_builder).inner as *const WriteBuilderState)).kind,
WriteBuilderKind::Standard
));
paimon_write_builder_free(normal.write_builder);

let fixed = paimon_table_new_postpone_fixed_bucket_write_builder(handle);
assert!(fixed.error.is_null());
assert!(matches!(
(&*((*fixed.write_builder).inner as *const WriteBuilderState)).kind,
WriteBuilderKind::PostponeFixed { .. }
));
let write = paimon_write_builder_new_write(fixed.write_builder);
assert!(write.write.is_null());
assert!(error_message(write.error).contains("bucket plan is required"));
paimon_error_free(write.error);
paimon_write_builder_free(fixed.write_builder);

let wb1 = paimon_table_new_postpone_fixed_bucket_write_builder_with_commit_user(
handle,
commit_user.as_ptr(),
)
.write_builder;
let wb2 = paimon_table_new_postpone_fixed_bucket_write_builder_with_commit_user(
handle,
commit_user.as_ptr(),
)
.write_builder;

for wb in [wb1, wb2] {
assert!(paimon_write_builder_with_overwrite(wb).is_null());
let (array, schema) = export_batch_to_ffi(make_postpone_bucket_plan_batch(
vec!["p1", "p2"],
vec![3, 3],
));
let error = paimon_write_builder_with_postpone_bucket_plan(
wb,
(&**array) as *const FFI_ArrowArray as *mut c_void,
(&**schema) as *const FFI_ArrowSchema as *mut c_void,
);
assert!(error.is_null());
}

let tw1 = paimon_write_builder_new_write(wb1).write;
let tw2 = paimon_write_builder_new_write(wb2).write;
for (tw, partitions, ids, names) in [
(tw1, vec!["p1"], vec![1], vec!["a"]),
(
tw2,
vec!["p2", "p2", "p2", "p2"],
vec![2, 3, 4, 5],
vec!["b", "c", "d", "e"],
),
] {
let (array, schema) =
export_batch_to_ffi(make_partitioned_write_batch(partitions, ids, names));
let error = paimon_table_write_write_arrow_batch(
tw,
(&**array) as *const FFI_ArrowArray as *mut c_void,
(&**schema) as *const FFI_ArrowSchema as *mut c_void,
);
assert!(error.is_null());
}

let messages1 = paimon_table_write_prepare_commit(tw1).messages;
let messages2 = paimon_table_write_prepare_commit(tw2).messages;
for messages in [messages1, messages2] {
let state = &*((*messages).inner as *const CommitMessagesState);
assert!(!state.messages.is_empty());
assert!(state
.messages
.iter()
.all(|message| message.total_buckets == Some(3)));
}
let error = paimon_commit_messages_merge(messages1, messages2);
assert!(error.is_null());
let commit = paimon_write_builder_new_commit(wb1).commit;
assert!(matches!(
(&*((*commit).inner as *const TableCommitState)).commit,
TableCommitKind::PostponeFixed(_)
));
let error = paimon_table_commit_commit(commit, messages1);
assert!(error.is_null());
let snapshot = crate::runtime()
.block_on(
SnapshotManager::new(table_ref(handle).file_io().clone(), path.to_string())
.get_latest_snapshot(),
)
.unwrap()
.unwrap();
assert_eq!(snapshot.commit_kind(), &CommitKind::OVERWRITE);

paimon_table_commit_free(commit);
paimon_commit_messages_free(messages2);
paimon_commit_messages_free(messages1);
paimon_table_write_free(tw2);
paimon_table_write_free(tw1);
paimon_write_builder_free(wb2);
paimon_write_builder_free(wb1);
unwrap_table(handle);
}
}

#[test]
fn test_distributed_postpone_writers_reject_overlapping_bucket_ownership() {
let path = "memory:/test_distributed_postpone_overlapping_ownership";
let file_io = memory_file_io();
setup_table_dirs(&file_io, path);
let table = Table::new(
file_io,
Identifier::new("default", "test"),
path.to_string(),
partitioned_postpone_table_schema(),
None,
);
let handle = unsafe { wrap_table(table) };
let commit_user = CString::new("overlapping-postpone-job").unwrap();

unsafe {
let wb1 = paimon_table_new_postpone_fixed_bucket_write_builder_with_commit_user(
handle,
commit_user.as_ptr(),
)
.write_builder;
let wb2 = paimon_table_new_postpone_fixed_bucket_write_builder_with_commit_user(
handle,
commit_user.as_ptr(),
)
.write_builder;
for wb in [wb1, wb2] {
let (array, schema) =
export_batch_to_ffi(make_postpone_bucket_plan_batch(vec!["p"], vec![1]));
assert!(paimon_write_builder_with_postpone_bucket_plan(
wb,
(&**array) as *const FFI_ArrowArray as *mut c_void,
(&**schema) as *const FFI_ArrowSchema as *mut c_void,
)
.is_null());
}

let tw1 = paimon_write_builder_new_write(wb1).write;
let tw2 = paimon_write_builder_new_write(wb2).write;
for (tw, name) in [(tw1, "first"), (tw2, "second")] {
let (array, schema) =
export_batch_to_ffi(make_partitioned_write_batch(vec!["p"], vec![1], vec![name]));
assert!(paimon_table_write_write_arrow_batch(
tw,
(&**array) as *const FFI_ArrowArray as *mut c_void,
(&**schema) as *const FFI_ArrowSchema as *mut c_void,
)
.is_null());
}

let messages1 = paimon_table_write_prepare_commit(tw1).messages;
let messages2 = paimon_table_write_prepare_commit(tw2).messages;
let error = paimon_commit_messages_merge(messages1, messages2);
assert!(error.is_null());
let commit = paimon_write_builder_new_commit(wb1).commit;
let error = paimon_table_commit_commit(commit, messages1);
assert!(!error.is_null());
assert!(error_message(error).contains("writer ownership conflict for bucket 0"));
paimon_error_free(error);

paimon_table_commit_free(commit);
paimon_commit_messages_free(messages2);
paimon_commit_messages_free(messages1);
paimon_table_write_free(tw2);
paimon_table_write_free(tw1);
paimon_write_builder_free(wb2);
paimon_write_builder_free(wb1);
unwrap_table(handle);
}
}

#[test]
fn test_write_multiple_batches() {
let path = "memory:/test_write_multi_batch";
Expand Down Expand Up @@ -1721,6 +2053,11 @@ fn test_null_pointer_handling() {
assert!(result.write_builder.is_null());
paimon_error_free(result.error);

let result = paimon_table_new_postpone_fixed_bucket_write_builder(ptr::null());
assert!(!result.error.is_null());
assert!(result.write_builder.is_null());
paimon_error_free(result.error);

let result = paimon_write_builder_new_write(ptr::null());
assert!(!result.error.is_null());
assert!(result.write.is_null());
Expand All @@ -1731,6 +2068,14 @@ fn test_null_pointer_handling() {
assert!(result.commit.is_null());
paimon_error_free(result.error);

let err = paimon_write_builder_with_postpone_bucket_plan(
ptr::null_mut(),
ptr::null_mut(),
ptr::null_mut(),
);
assert!(!err.is_null());
paimon_error_free(err);

let err =
paimon_table_write_write_arrow_batch(ptr::null_mut(), ptr::null_mut(), ptr::null_mut());
assert!(!err.is_null());
Expand Down
Loading
Loading