From 3af51541a02963563baac208a7bd768ac241a785 Mon Sep 17 00:00:00 2001 From: "lancedb-gatefixer[bot]" <313497061+lancedb-gatefixer[bot]@users.noreply.github.com> Date: Thu, 6 Aug 2026 16:46:49 +0800 Subject: [PATCH] test(rust): cover fixed-size binary merge insert regression (#3854) ## Summary - add a LanceDB regression for `merge_insert` with a non-nullable `FixedSizeBinary` column - exercise matched updates, unmatched inserts, and source-missing deletes - assert the exact merge statistics and final row count ## Root cause The Arrow `take` kernel previously ignored nulls in the index array for `FixedSizeBinary`. DataFusion uses that kernel while constructing outer-join results, so the join behind `when_not_matched_by_source_delete` could place invalid values into non-nullable columns. The current Arrow dependency contains the upstream fix; this test locks the corrected behavior at the LanceDB API boundary. ## Validation - `cargo fmt --all -- --check` - `cargo test --quiet --features remote --tests` - `cargo check --quiet --features remote --tests --examples` - `cargo clippy --quiet --features remote --tests --examples` Fixes #2869 Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com> --- rust/lancedb/src/table/merge.rs | 71 ++++++++++++++++++++++++++++++++- 1 file changed, 70 insertions(+), 1 deletion(-) diff --git a/rust/lancedb/src/table/merge.rs b/rust/lancedb/src/table/merge.rs index 13a633c67..82a1d1473 100644 --- a/rust/lancedb/src/table/merge.rs +++ b/rust/lancedb/src/table/merge.rs @@ -315,7 +315,10 @@ pub(crate) async fn execute_merge_insert( #[cfg(test)] mod tests { - use arrow_array::{Int32Array, RecordBatch, RecordBatchIterator, RecordBatchReader}; + use arrow_array::builder::FixedSizeBinaryBuilder; + use arrow_array::{ + Int32Array, RecordBatch, RecordBatchIterator, RecordBatchReader, StringArray, UInt64Array, + }; use arrow_schema::{DataType, Field, Schema}; use std::sync::Arc; @@ -337,6 +340,42 @@ mod tests { Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema)) } + fn fixed_size_binary_merge_batch( + id_range: std::ops::Range, + price: u64, + ) -> Box { + let ids = id_range.collect::>(); + let mut id_builder = FixedSizeBinaryBuilder::new(16); + for id in &ids { + let mut bytes = [0; 16]; + bytes[..8].copy_from_slice(&id.to_le_bytes()); + id_builder.append_value(bytes).unwrap(); + } + + let schema = Arc::new(Schema::new(vec![ + Field::new("id", DataType::FixedSizeBinary(16), false), + Field::new("id_as_int", DataType::UInt64, false), + Field::new("name", DataType::Utf8, false), + Field::new("market", DataType::Utf8, false), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(id_builder.finish()), + Arc::new(UInt64Array::from_iter_values(ids.iter().copied())), + Arc::new(StringArray::from_iter_values( + ids.iter().map(|id| format!("name{id}")), + )), + Arc::new(StringArray::from_iter_values(std::iter::repeat_n( + format!("market_{price}"), + ids.len(), + ))), + ], + ) + .unwrap(); + Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema)) + } + #[tokio::test] async fn test_merge_insert() { let conn = connect("memory://").execute().await.unwrap(); @@ -388,6 +427,36 @@ mod tests { ); } + #[tokio::test] + async fn test_merge_insert_fixed_size_binary_non_nullable() { + // Regression test for #2869: an unrelated FixedSizeBinary column used to corrupt the + // outer join that implements when_not_matched_by_source_delete. + let conn = connect("memory://").execute().await.unwrap(); + let table = conn + .create_table( + "fixed_size_binary_merge", + fixed_size_binary_merge_batch(0..256, 100), + ) + .execute() + .await + .unwrap(); + + let mut merge_insert = table.merge_insert(&["id_as_int"]); + merge_insert + .when_matched_update_all(None) + .when_not_matched_insert_all() + .when_not_matched_by_source_delete(None); + let result = merge_insert + .execute(fixed_size_binary_merge_batch(100..356, 200)) + .await + .unwrap(); + + assert_eq!(result.num_updated_rows, 156); + assert_eq!(result.num_inserted_rows, 100); + assert_eq!(result.num_deleted_rows, 100); + assert_eq!(table.count_rows(None).await.unwrap(), 256); + } + #[tokio::test] async fn test_merge_insert_use_index() { let conn = connect("memory://").execute().await.unwrap();