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
284 changes: 284 additions & 0 deletions crates/omnigraph/tests/lance_surface_guards.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5867,3 +5867,287 @@ async fn packed_struct_refuses_null_values_and_lone_child_projection_lance_11()
"{error}"
);
}

// --- Group commit: composing inserts staged against older pins -------------
//
// Group commit (docs/rfcs/2026-09-28-group-commit.md, Composition) publishes
// several same-table inserts as one detached commit. Each entry staged its
// keyed merge-insert against the pin it captured, and other inserts may have
// landed on the table since. Composition is sound because Lance assigns what
// a fresh fragment needs at commit, from the manifest the commit is built
// on: `fragments_with_ids` fills only a zero fragment id, `assign_row_ids`
// fills only a missing or partial row-id sequence, and the row-version
// metadata is derived while the manifest is built. This fence pins that the
// raw staged fragments are unassigned, that the components' key filters
// share one Bloom configuration, so their union is a bitwise OR, and that
// one detached commit of the union from a later pin yields exactly the union
// of rows with fresh, unique fragment and row ids and one version stamp. A
// Lance release that assigned ids at staging would make composition reuse
// them: a design review for that RFC, not a test to weaken.

fn pk_rows(dataset: &Dataset, ids: &[&str]) -> RecordBatch {
let schema = Arc::new(Schema::from(dataset.schema()));
let values: Vec<i32> = (0..ids.len() as i32).collect();
let notes: Vec<Option<&str>> = ids.iter().map(|_| Some("group")).collect();
RecordBatch::try_new(
schema,
vec![
Arc::new(StringArray::from(ids.to_vec())),
Arc::new(Int32Array::from(values)),
Arc::new(StringArray::from(notes)),
],
)
.unwrap()
}

/// A keyed upsert of `ids`, staged uncommitted against `base`: the shape the
/// engine stages for an `insert` into a `@key` type.
async fn stage_keyed_upsert(base: &Dataset, ids: &[&str]) -> UncommittedMergeInsert {
stage_pk_merge(
Arc::new(base.clone()),
pk_rows(base, ids),
"id",
WhenMatched::UpdateAll,
WhenNotMatched::InsertAll,
None,
)
.await
}

async fn commit_detached_from(base: &Dataset, transaction: Transaction) -> Dataset {
CommitBuilder::new(Arc::new(base.clone()))
.with_detached(true)
.with_skip_auto_cleanup(true)
.execute(transaction)
.await
.unwrap()
}

/// `(key, stable row id, created-at version, last-updated version)` per row.
async fn keyed_row_identities(ds: &Dataset) -> Vec<(String, u64, u64, u64)> {
let mut scanner = ds.scan();
scanner.with_row_id();
scanner
.project(&["id", ROW_CREATED_AT_VERSION, ROW_LAST_UPDATED_AT_VERSION])
.unwrap();
let batches: Vec<RecordBatch> = scanner
.try_into_stream()
.await
.unwrap()
.try_collect()
.await
.unwrap();
let mut out = Vec::new();
for batch in batches {
let keys = batch["id"].as_string::<i32>();
let row_ids = batch[ROW_ID].as_primitive::<arrow_array::types::UInt64Type>();
let created =
batch[ROW_CREATED_AT_VERSION].as_primitive::<arrow_array::types::UInt64Type>();
let updated =
batch[ROW_LAST_UPDATED_AT_VERSION].as_primitive::<arrow_array::types::UInt64Type>();
for row in 0..batch.num_rows() {
out.push((
keys.value(row).to_string(),
row_ids.value(row),
created.value(row),
updated.value(row),
));
}
}
out.sort();
out
}

#[tokio::test]
async fn group_commit_composes_inserts_staged_against_older_pins_onto_a_later_pin() {
let dir = tempfile::tempdir().unwrap();
let uri = dir.path().join("t.lance");
let uri = uri.to_str().unwrap();
let p0 = fresh_pk_dataset(uri).await;

// Two entries capture P0 and stage keyed upserts of new keys.
let a = stage_keyed_upsert(&p0, &["b1", "b2"]).await;
let b = stage_keyed_upsert(&p0, &["c1"]).await;
// Another insert lands meanwhile: P1 is its detached commit from P0.
let landed = stage_keyed_upsert(&p0, &["d1"]).await;
let p1 = commit_detached_from(&p0, landed.transaction).await;
// A third entry captures P1.
let c = stage_keyed_upsert(&p1, &["e1"]).await;

let mut union = Vec::new();
let mut filters = Vec::new();
for staged in [&a, &b, &c] {
let Operation::Update {
removed_fragment_ids,
updated_fragments,
new_fragments,
..
} = &staged.transaction.operation
else {
panic!("a keyed merge-insert stages Operation::Update");
};
assert!(
removed_fragment_ids.is_empty() && updated_fragments.is_empty(),
"no key matched, so the staged transaction only appends"
);
for fragment in new_fragments {
assert_eq!(fragment.id, 0, "a raw staged fragment carries no id");
// Merge-insert records the captured ids of the rows it rewrote,
// rechunked per fragment and allowed to be incomplete; the commit
// fills the rest from its own base. An insert rewrote nothing, so
// the sequence is empty and every row id comes from the commit.
let carried = match &fragment.row_id_meta {
None => 0,
Some(lance_table::format::RowIdMeta::Inline(data)) => {
lance_table::rowids::read_row_ids(data).unwrap().len()
}
Some(other) => panic!("unexpected row-id metadata {other:?}"),
};
assert_eq!(
carried, 0,
"a raw staged insert fragment carries no row ids of its own"
);
}
union.extend(new_fragments.iter().cloned());
filters.push(
staged_inserted_rows_filter(staged)
.expect("a primary-key merge-insert carries a key filter")
.clone(),
);
}
let composed_fragments = union.len();

// The union of Bloom filters under one configuration is their bitwise OR.
let FilterType::Bloom {
bitmap: mut union_bits,
num_bits,
number_of_items,
probability,
} = filters[0].filter.clone()
else {
panic!("pinned Lance emits a Bloom key filter")
};
for filter in &filters[1..] {
assert_eq!(filter.field_ids, filters[0].field_ids, "one key column");
let FilterType::Bloom {
bitmap,
num_bits: other_bits,
number_of_items: other_items,
probability: other_probability,
} = &filter.filter
else {
panic!("pinned Lance emits a Bloom key filter")
};
assert_eq!(
(*other_bits, *other_items),
(num_bits, number_of_items),
"one Bloom configuration"
);
assert!((other_probability - probability).abs() < f64::EPSILON);
for (byte, other) in union_bits.iter_mut().zip(bitmap) {
*byte |= other;
}
}
let union_filter = KeyExistenceFilter {
field_ids: filters[0].field_ids.clone(),
filter: FilterType::Bloom {
bitmap: union_bits,
num_bits,
number_of_items,
probability,
},
};
for filter in &filters {
assert!(union_filter.intersects(filter).unwrap().0);
}

// Compose onto P1: the first component's operation, carrying the union.
let mut operation = a.transaction.operation.clone();
let Operation::Update {
new_fragments,
inserted_rows_filter,
..
} = &mut operation
else {
unreachable!()
};
*new_fragments = union;
*inserted_rows_filter = Some(union_filter);
let composed =
commit_detached_from(&p1, Transaction::new(p1.version().version, operation, None)).await;
assert!(lance_table::format::is_detached_version(
composed.version().version
));

// Exactly the union of rows, and P1's rows carried unchanged.
let rows = keyed_row_identities(&composed).await;
let keys: Vec<&str> = rows.iter().map(|row| row.0.as_str()).collect();
assert_eq!(keys, ["alice", "b1", "b2", "bob", "c1", "d1", "e1"]);
let p1_rows = keyed_row_identities(&p1).await;
for row in &p1_rows {
assert!(rows.contains(row), "P1's row changed identity: {row:?}");
}

// Fresh, unique stable row ids from P1's counter.
let row_ids: HashSet<u64> = rows.iter().map(|row| row.1).collect();
assert_eq!(row_ids.len(), rows.len(), "stable row ids are unique");
let p1_max_row_id = p1_rows.iter().map(|row| row.1).max().unwrap();
for row in rows.iter().filter(|row| !p1_rows.contains(row)) {
assert!(row.1 > p1_max_row_id, "reused a row id: {row:?}");
}

// Fresh, unique fragment ids above P1's.
let p1_fragment_ids: HashSet<u64> = p1
.get_fragments()
.iter()
.map(|fragment| fragment.id() as u64)
.collect();
let fragment_ids: Vec<u64> = composed
.get_fragments()
.iter()
.map(|fragment| fragment.id() as u64)
.collect();
assert_eq!(
fragment_ids.iter().collect::<HashSet<_>>().len(),
fragment_ids.len(),
"fragment ids are unique"
);
assert_eq!(
fragment_ids.len(),
p1_fragment_ids.len() + composed_fragments,
"every composed fragment is added"
);
let p1_max_fragment_id = *p1_fragment_ids.iter().max().unwrap();
for id in fragment_ids
.iter()
.filter(|id| !p1_fragment_ids.contains(id))
{
assert!(*id > p1_max_fragment_id, "reused fragment id {id}");
}

// One commit, one stamp, distinct from the version that landed first.
let stamp = |key: &str| {
let row = rows.iter().find(|row| row.0 == key).unwrap();
(row.2, row.3)
};
for key in ["b2", "c1", "e1"] {
assert_eq!(
stamp(key),
stamp("b1"),
"{key} is stamped by the composed commit"
);
}
assert_ne!(
stamp("b1"),
stamp("d1"),
"the landed insert keeps its own stamp"
);

// The composed version reopens by id.
let reopened = DatasetBuilder::from_uri(uri)
.with_version(composed.version().version)
.load()
.await
.unwrap();
assert_eq!(keyed_row_identities(&reopened).await, rows);
}
3 changes: 2 additions & 1 deletion docs/rfcs/0067-detached-table-commits.md
Original file line number Diff line number Diff line change
Expand Up @@ -1277,7 +1277,8 @@ claims).
per-chunk promotion cost near the 1,024-chunk bound.
- Group commit at the branch gate is a separate proposal; here it is N
detached commits, one publication, N promotions. Decided in that proposal,
when it opens.
when it opens. Proposed in [RFC: Group commit](2026-09-28-group-commit.md),
which publishes a batch as one graph commit rather than as N commit rows.
- How a blocked table is unblocked. A foreign linear commit occupying a
pin's target never resolves on its own, and every later pin on that table
waits behind it. The candidates are a `repair` arm that restores the
Expand Down
Loading
Loading