From 99a68db78c149ceb3342be358c6b3178d18f0a2c 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:40:01 +0800 Subject: [PATCH] test(rust): cover concurrent appends during compaction (#3878) ## Summary - add a LanceDB core regression for compaction overlapping appends through separate table handles - verify concurrent commits preserve fragment ID order on an indexed table - run the follow-up compaction that exposed the original row-ID ordering failure and verify all rows remain ## Root cause Older Lance versions could reserve fragment IDs for compaction, allow concurrent appends to commit later IDs, and then commit the reserved compaction fragments at the end of the manifest. A later compaction could consequently receive row IDs out of order. Current Lance sorts fragments at the transaction boundary; this adds the missing LanceDB-level regression coverage for the Node-visible concurrency contract. ## Validation - `cargo fmt --all` - focused regression passed once with output and 20 repeated runs - `cargo test --quiet --features remote -p lancedb table::optimize::tests` (14 passed) - `cargo check --quiet --features remote --tests --examples` - `cargo clippy --quiet --features remote --tests --examples` - `cargo test --quiet --features remote --tests` (867 passed, 1 ignored) Fixes #1498 Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com> --- rust/lancedb/src/table/optimize.rs | 90 ++++++++++++++++++++++++++++++ 1 file changed, 90 insertions(+) diff --git a/rust/lancedb/src/table/optimize.rs b/rust/lancedb/src/table/optimize.rs index e29445b2d..3fdad0533 100644 --- a/rust/lancedb/src/table/optimize.rs +++ b/rust/lancedb/src/table/optimize.rs @@ -304,6 +304,96 @@ mod tests { assert_eq!(all_values, expected); } + #[tokio::test] + async fn test_compact_with_concurrent_add() { + const NUM_FRAGMENTS: usize = 5; + const ROWS_PER_FRAGMENT: i32 = 300; + + let tmpdir = tempfile::tempdir().unwrap(); + let conn = connect(tmpdir.path().to_str().unwrap()) + .execute() + .await + .unwrap(); + let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)])); + let batch = RecordBatch::try_new( + schema, + vec![Arc::new(Int32Array::from_iter_values(0..ROWS_PER_FRAGMENT))], + ) + .unwrap(); + + let table = conn + .create_table("test_concurrent_compact", batch.clone()) + .execute() + .await + .unwrap(); + table + .create_index(&["id"], Index::BTree(BTreeIndexBuilder::default())) + .execute() + .await + .unwrap(); + for _ in 0..NUM_FRAGMENTS { + table.add(batch.clone()).execute().await.unwrap(); + } + + // Use separate handles so the two writes actually overlap, as they can + // when different Node connections operate on the same S3 table. + let compact_table = conn + .open_table("test_concurrent_compact") + .execute() + .await + .unwrap(); + let append_table = conn + .open_table("test_concurrent_compact") + .execute() + .await + .unwrap(); + let compact_task = tokio::spawn(async move { + compact_table + .optimize(OptimizeAction::Compact { + options: CompactionOptions { + target_rows_per_fragment: 1_000, + ..Default::default() + }, + remap_options: None, + }) + .await + }); + tokio::task::yield_now().await; + for _ in 0..NUM_FRAGMENTS { + append_table.add(batch.clone()).execute().await.unwrap(); + } + compact_task.await.unwrap().unwrap(); + + let table = conn + .open_table("test_concurrent_compact") + .execute() + .await + .unwrap(); + let dataset = table.dataset().unwrap().get().await.unwrap(); + let fragment_ids = dataset + .get_fragments() + .iter() + .map(|fragment| fragment.id()) + .collect::>(); + assert!(fragment_ids.windows(2).all(|ids| ids[0] < ids[1])); + + // A second compaction exposed the original out-of-order row-id bug. + table + .optimize(OptimizeAction::Compact { + options: CompactionOptions { + target_rows_per_fragment: 1_000, + ..Default::default() + }, + remap_options: None, + }) + .await + .unwrap(); + assert_eq!( + table.count_rows(None).await.unwrap(), + ROWS_PER_FRAGMENT as usize * (NUM_FRAGMENTS * 2 + 1) + ); + } + #[tokio::test] async fn test_optimize_prune_versions() { let conn = connect("memory://").execute().await.unwrap();