From a691b11a68ac650678a30ab08b54ab486bfa9e88 Mon Sep 17 00:00:00 2001 From: Jack Ye Date: Sat, 11 Jul 2026 02:09:14 -0700 Subject: [PATCH 1/3] feat: add CompositeOperation for multi-statement transactions Adds a Composite operation that applies an ordered list of existing operations as one atomic commit. Sub-operations reuse the existing build_manifest apply logic (a sequential fold with the version pinned to the single committed version) and the existing TransactionRebase conflict rules by delegation on both sides of a conflict check. Composite transactions are not inlined into the manifest so released readers can still open post-composite versions; the eager inline transaction decode in load_manifest is also made tolerant so future operation types no longer prevent opening a dataset. Restore, Clone, nested composites, and empty composites are rejected. --- docs/src/format/table/transaction.md | 34 ++ protos/transaction.proto | 20 + python/python/lance/dataset.py | 58 +- python/python/tests/test_dataset.py | 52 ++ python/src/transaction.rs | 12 + rust/lance/src/dataset.rs | 50 +- rust/lance/src/dataset/transaction.rs | 532 +++++++++++++++++- rust/lance/src/dataset/write/commit.rs | 245 +++++++- rust/lance/src/io/commit/conflict_resolver.rs | 357 ++++++++++++ 9 files changed, 1341 insertions(+), 19 deletions(-) diff --git a/docs/src/format/table/transaction.md b/docs/src/format/table/transaction.md index 78dd5301fb8..8fc52d23f53 100644 --- a/docs/src/format/table/transaction.md +++ b/docs/src/format/table/transaction.md @@ -498,6 +498,40 @@ An UpdateBases operation only modifies the base paths. As a result, it only conf UpdateBases operations and even then only conflicts if the two operations have base paths with the same id, name, or path. +### Composite + +An ordered list of operations applied sequentially as one atomic commit. Each sub-operation is +applied to the manifest produced by the previous one, and the final manifest is committed as a +single new version. This enables multi-statement transactions — for example, an append and an +index update over the appended fragments (with fragment IDs reserved via ReserveFragments ahead +of time) that become visible atomically. + +Each sub-transaction carries only its `operation`; the outer transaction's `read_version`, `uuid`, +`tag`, and `transaction_properties` govern the commit, and the corresponding inner fields are left +unset by writers and ignored by readers. The sub-operations must not contain Restore, Clone, or a +nested Composite, and the list must not be empty. + +
+CompositeOperation protobuf message + +```protobuf +%%% proto.message.CompositeOperation %%% +``` + +
+ +#### Composite Compatibility + +A Composite operation conflicts with a concurrent transaction if and only if at least one of its +sub-operations conflicts with it, applying each sub-operation's own compatibility rules. When the +concurrent transaction is itself a Composite, it is decomposed the same way, so two Composite +operations conflict if any pair of their sub-operations conflicts. Rebases are likewise delegated: +each sub-operation is rebased under its own rules and the results are reassembled in order. + +Writers do not inline Composite transactions into the manifest's transaction section: released +readers eagerly decode that section when opening a dataset and would fail on an unknown operation. +The transaction remains available in the external transaction file. + ## Conflict Resolution ### Terminology diff --git a/protos/transaction.proto b/protos/transaction.proto index e72e95025a4..934c7583674 100644 --- a/protos/transaction.proto +++ b/protos/transaction.proto @@ -329,6 +329,23 @@ message Transaction { repeated BasePath new_bases = 1; } + // An ordered list of operations applied sequentially as one atomic commit. + // + // Each sub-transaction's `operation` is applied to the manifest produced by + // the previous one; the resulting manifest is committed as a single new + // version. Only the `operation` field of each sub-transaction is meaningful: + // the outer transaction's `read_version`, `uuid`, `tag` and + // `transaction_properties` govern the commit, and the corresponding inner + // fields are left unset by writers and ignored by readers. Inner fields are + // reserved for future per-statement metadata (e.g. merging independently + // created transactions while preserving their identity). + // + // Sub-transactions must not contain `Restore`, `Clone`, or a nested + // `CompositeOperation`, and the list must not be empty. + message CompositeOperation { + repeated Transaction transactions = 1; + } + // The operation of this transaction. oneof operation { Append append = 100; @@ -346,6 +363,9 @@ message Transaction { UpdateMemWalState update_mem_wal_state = 112; Clone clone = 113; UpdateBases update_bases = 114; + // Field 115 is intentionally skipped: it is used by an in-flight proposal + // (https://github.com/lance-format/lance/pull/7674). + CompositeOperation composite = 116; } // Fields 200/202 (`blob_append` / `blob_overwrite`) previously represented blob dataset ops. diff --git a/python/python/lance/dataset.py b/python/python/lance/dataset.py index 21402c3d276..df5b4cb748b 100644 --- a/python/python/lance/dataset.py +++ b/python/python/lance/dataset.py @@ -4580,10 +4580,15 @@ def commit( and not isinstance( operation, (LanceOperation.Overwrite, LanceOperation.Restore) ) + and not ( + isinstance(operation, LanceOperation.Composite) + and len(operation.operations) > 0 + and isinstance(operation.operations[0], LanceOperation.Overwrite) + ) ): raise ValueError( "read_version is required for all operations except " - "Overwrite and Restore" + "Overwrite, Restore, and Composite starting with Overwrite" ) # Storage options provider is automatically created in Rust when @@ -6073,6 +6078,57 @@ class UpdateConfig(BaseOperation): schema_metadata_updates: Optional[LanceOperation.UpdateMap] = None field_metadata_updates: Optional[Dict[int, LanceOperation.UpdateMap]] = None + @dataclass + class Composite(BaseOperation): + """ + An ordered list of operations applied sequentially as one atomic commit. + + Each operation is applied to the result of the previous one, and the + final state is committed as a single new version. This allows + multi-statement transactions such as an append and an index update + that become visible atomically. + + The operations must not contain ``Restore`` or a nested ``Composite``. + + Attributes + ---------- + operations: list[LanceOperation.BaseOperation] + The operations to apply in order. + + Examples + -------- + + >>> import lance + >>> import pyarrow as pa + >>> tab1 = pa.table({"a": [1, 2], "b": ["a", "b"]}) + >>> dataset = lance.write_dataset(tab1, "example") + >>> tab2 = pa.table({"a": [3, 4], "b": ["c", "d"]}) + >>> fragment = lance.fragment.LanceFragment.create("example", tab2) + >>> operation = lance.LanceOperation.Composite( + ... [ + ... lance.LanceOperation.Append([fragment]), + ... lance.LanceOperation.Delete([], [1], "a >= 3"), + ... ] + ... ) + >>> dataset = lance.LanceDataset.commit("example", operation, + ... read_version=dataset.version) + >>> dataset.to_table().to_pandas() + a b + 0 1 a + 1 2 b + """ + + operations: Iterable[LanceOperation.BaseOperation] + + def __post_init__(self): + if not isinstance(self.operations, list) or not all( + isinstance(op, LanceOperation.BaseOperation) for op in self.operations + ): + raise TypeError( + "operations must be list[LanceOperation.BaseOperation], " + f"got {type(self.operations)}" + ) + @dataclass class ColumnOrdering: diff --git a/python/python/tests/test_dataset.py b/python/python/tests/test_dataset.py index 500ee2e17bf..c2a8942d32e 100644 --- a/python/python/tests/test_dataset.py +++ b/python/python/tests/test_dataset.py @@ -1838,6 +1838,58 @@ def test_append_with_commit(tmp_path: Path): assert tbl == expected +def test_composite_with_commit(tmp_path: Path): + table = pa.Table.from_pydict({"a": range(100), "b": range(100)}) + base_dir = tmp_path / "test" + dataset = lance.write_dataset(table, base_dir) + + fragment = lance.fragment.LanceFragment.create(base_dir, table) + composite = lance.LanceOperation.Composite( + [ + lance.LanceOperation.Append([fragment]), + lance.LanceOperation.UpdateConfig( + config_updates=lance.LanceOperation.UpdateMap( + updates={"composite.key": "value"}, replace=False + ) + ), + ] + ) + + with pytest.raises(ValueError): + # Must specify read version + lance.LanceDataset.commit(dataset, composite) + + dataset = lance.LanceDataset.commit(dataset, composite, read_version=1) + + # Both effects are visible in a single new version. + assert dataset.version == 2 + assert dataset.count_rows() == 200 + assert dataset.config()["composite.key"] == "value" + + # The committed transaction reads back as a composite. + transaction = dataset.read_transaction(2) + assert isinstance(transaction.operation, lance.LanceOperation.Composite) + assert len(transaction.operation.operations) == 2 + assert isinstance(transaction.operation.operations[0], lance.LanceOperation.Append) + + +def test_composite_validation(): + with pytest.raises(TypeError): + lance.LanceOperation.Composite(operations="not-a-list") + + data = pa.Table.from_pydict({"a": range(10), "b": range(10)}) + dataset = lance.write_dataset(data, "memory://composite-validation") + empty = lance.LanceOperation.Composite(operations=[]) + with pytest.raises(Exception, match="at least one"): + lance.LanceDataset.commit(dataset, empty, read_version=1) + + nested = lance.LanceOperation.Composite( + operations=[lance.LanceOperation.Composite(operations=[])] + ) + with pytest.raises(Exception, match="cannot contain"): + lance.LanceDataset.commit(dataset, nested, read_version=1) + + def test_commit_batch_append(): data1 = pa.Table.from_pydict({"a": range(100), "b": range(100)}) dataset = lance.write_dataset(data1, "memory://test") diff --git a/python/src/transaction.rs b/python/src/transaction.rs index 1b659395099..8089f772c8d 100644 --- a/python/src/transaction.rs +++ b/python/src/transaction.rs @@ -390,6 +390,11 @@ impl FromPyObject<'_, '_> for PyLance { }; Ok(Self(op)) } + "Composite" => { + let operations = extract_vec(&ob.getattr("operations")?)?; + let op = Operation::Composite { operations }; + Ok(Self(op)) + } unsupported => Err(PyValueError::new_err(format!( "Unsupported operation: {unsupported}", ))), @@ -595,6 +600,13 @@ impl<'py> IntoPyObject<'py> for PyLance<&Operation> { base_op.call0() } } + Operation::Composite { operations } => { + let operations = export_vec(py, operations.as_slice())?; + let cls = namespace + .getattr("Composite") + .expect("Failed to get Composite class"); + cls.call1((operations,)) + } _ => todo!(), } } diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index 72b13e903b5..449ccf6683e 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -738,16 +738,32 @@ impl Dataset { let message_len = LittleEndian::read_u32(&last_block[offset_in_block..offset_in_block + 4]) as usize; let message_data = &last_block[offset_in_block + 4..offset_in_block + 4 + message_len]; - let transaction: Transaction = - lance_table::format::pb::Transaction::decode(message_data)?.try_into()?; - - let metadata_cache = session.metadata_cache.for_dataset(uri); - let metadata_key = TransactionKey { - version: manifest_location.version, - }; - metadata_cache - .insert_with_key(&metadata_key, Arc::new(transaction)) - .await; + // This is an opportunistic cache warm-up: a transaction written by + // a newer version may contain an operation this version cannot + // decode, and that must not prevent opening the dataset. Paths + // that need the transaction contents surface the error themselves. + match lance_table::format::pb::Transaction::decode(message_data) + .map_err(Error::from) + .and_then(Transaction::try_from) + { + Ok(transaction) => { + let metadata_cache = session.metadata_cache.for_dataset(uri); + let metadata_key = TransactionKey { + version: manifest_location.version, + }; + metadata_cache + .insert_with_key(&metadata_key, Arc::new(transaction)) + .await; + } + Err(err) => { + log::warn!( + "Failed to decode the inline transaction of version {}; \ + it may have been written by a newer version of Lance: {}", + manifest_location.version, + err + ); + } + } } if manifest.should_use_legacy_format() { @@ -3674,6 +3690,20 @@ pub(crate) async fn write_manifest_file( naming_scheme: ManifestNamingScheme, mut transaction: Option<&Transaction>, ) -> std::result::Result { + // Composite transactions are not inlined into the manifest: released + // readers eagerly decode the inline transaction section when opening a + // dataset and fail on unknown operation types, which would make this + // version unopenable for them. Without the inline section they open + // normally; only paths that need the transaction contents (conflict + // resolution, history) read the external transaction file. + // TODO: inline composites once readers with tolerant decoding are + // widespread. + if matches!( + transaction.map(|t| &t.operation), + Some(Operation::Composite { .. }) + ) { + transaction = None; + } if config.auto_set_feature_flags { // build_manifest may have already set FLAG_STABLE_ROW_IDS on the manifest. // Preserve it here so this second apply_feature_flags call does not clear it diff --git a/rust/lance/src/dataset/transaction.rs b/rust/lance/src/dataset/transaction.rs index 1b929f95e3c..e3e0de47171 100644 --- a/rust/lance/src/dataset/transaction.rs +++ b/rust/lance/src/dataset/transaction.rs @@ -453,6 +453,17 @@ pub enum Operation { /// The new base paths to add to the manifest. new_bases: Vec, }, + + /// An ordered list of operations applied sequentially as one atomic commit. + /// + /// Each operation is applied to the manifest produced by the previous one, + /// and the result is committed as a single new version. Conflict resolution + /// delegates to each sub-operation's existing rules. + /// + /// Sub-operations must not contain [`Operation::Restore`], + /// [`Operation::Clone`], or a nested composite, and the list must not be + /// empty. + Composite { operations: Vec }, } #[derive(Debug, Clone, PartialEq, DeepSizeOf)] @@ -502,6 +513,16 @@ impl std::fmt::Display for Operation { Self::Clone { .. } => write!(f, "Clone"), Self::UpdateMemWalState { .. } => write!(f, "UpdateMemWalState"), Self::UpdateBases { .. } => write!(f, "UpdateBases"), + Self::Composite { operations } => { + write!(f, "Composite[")?; + for (i, op) in operations.iter().enumerate() { + if i > 0 { + write!(f, ", ")?; + } + write!(f, "{}", op.name())?; + } + write!(f, "]") + } } } } @@ -1345,6 +1366,10 @@ impl PartialEq for Operation { (Self::Clone { .. }, Self::UpdateBases { .. }) => { std::mem::discriminant(self) == std::mem::discriminant(other) } + // Sub-operation order is semantic (operations apply sequentially), + // so composites compare element-wise in order. + (Self::Composite { operations: a }, Self::Composite { operations: b }) => a == b, + (Self::Composite { .. }, _) | (_, Self::Composite { .. }) => false, } } } @@ -1410,6 +1435,10 @@ impl Operation { } }) .collect(), + Self::Composite { operations } => operations + .iter() + .flat_map(|op| op.get_upsert_config_keys()) + .collect(), _ => Vec::::new(), } } @@ -1431,6 +1460,10 @@ impl Operation { } }) .collect(), + Self::Composite { operations } => operations + .iter() + .flat_map(|op| op.get_delete_config_keys()) + .collect(), _ => Vec::::new(), } } @@ -1469,6 +1502,12 @@ impl Operation { } false } + (Self::Composite { operations }, _) => { + operations.iter().any(|op| op.modifies_same_metadata(other)) + } + (_, Self::Composite { operations }) => operations + .iter() + .any(|other_op| self.modifies_same_metadata(other_op)), _ => false, } } @@ -1524,8 +1563,36 @@ impl Operation { Self::UpdateMemWalState { .. } => "UpdateMemWalState", Self::Clone { .. } => "Clone", Self::UpdateBases { .. } => "UpdateBases", + Self::Composite { .. } => "Composite", } } + + /// Validate the structural rules for a composite operation's + /// sub-operations: the list must be non-empty and must not contain + /// `Restore`, `Clone`, or a nested `Composite`. `Restore` and `Clone` are + /// applied outside the normal manifest-building path in the commit driver, + /// so they cannot participate in a sequential fold; nesting adds nothing + /// over a flat list. + pub(crate) fn validate_composite(operations: &[Self]) -> Result<()> { + if operations.is_empty() { + return Err(Error::invalid_input( + "Composite operation must contain at least one sub-operation", + )); + } + for (position, op) in operations.iter().enumerate() { + match op { + Self::Restore { .. } | Self::Clone { .. } | Self::Composite { .. } => { + return Err(Error::invalid_input(format!( + "Composite operation cannot contain {} (at position {})", + op.name(), + position + ))); + } + _ => {} + } + } + Ok(()) + } } /// Helper function to apply UpdateMap changes to a HashMap @@ -1766,6 +1833,95 @@ impl Transaction { current_indices: Vec, transaction_file_path: &str, config: &ManifestWriteConfig, + ) -> Result<(Manifest, Vec)> { + let new_version = current_manifest.map_or(1, |m| m.version + 1); + if let Operation::Composite { operations } = &self.operation { + return self.build_manifest_composite( + operations, + current_manifest, + current_indices, + transaction_file_path, + config, + new_version, + ); + } + self.build_manifest_impl( + current_manifest, + current_indices, + transaction_file_path, + config, + new_version, + ) + } + + /// Apply a composite operation by folding its sub-operations: each + /// sub-operation is applied to the manifest produced by the previous one. + /// + /// All sub-operations are applied with the same `new_version` (the single + /// version this commit will create) so that row-level version metadata is + /// stamped correctly, and each intermediate manifest's version is pinned + /// back to `new_version` because [`Manifest::new_from_previous`] increments + /// the version on every application. + fn build_manifest_composite( + &self, + operations: &[Operation], + current_manifest: Option<&Manifest>, + current_indices: Vec, + transaction_file_path: &str, + config: &ManifestWriteConfig, + new_version: u64, + ) -> Result<(Manifest, Vec)> { + Operation::validate_composite(operations)?; + let mut state: Option<(Manifest, Vec)> = None; + let mut initial_indices = Some(current_indices); + for operation in operations { + let sub_transaction = Self { + read_version: self.read_version, + uuid: self.uuid.clone(), + operation: operation.clone(), + tag: None, + transaction_properties: None, + }; + let (mut manifest, indices) = match state.take() { + None => { + // Validate each sub-operation against the manifest it will + // actually apply to; entry-point validation only sees the + // base manifest and cannot account for schema changes made + // by earlier sub-operations. + validate_operation(current_manifest, operation)?; + sub_transaction.build_manifest_impl( + current_manifest, + initial_indices.take().expect("first fold step"), + transaction_file_path, + config, + new_version, + )? + } + Some((prev_manifest, prev_indices)) => { + validate_operation(Some(&prev_manifest), operation)?; + sub_transaction.build_manifest_impl( + Some(&prev_manifest), + prev_indices, + transaction_file_path, + config, + new_version, + )? + } + }; + manifest.version = new_version; + state = Some((manifest, indices)); + } + // validate_composite rejects empty operation lists, so state is set. + Ok(state.expect("non-empty composite")) + } + + fn build_manifest_impl( + &self, + current_manifest: Option<&Manifest>, + current_indices: Vec, + transaction_file_path: &str, + config: &ManifestWriteConfig, + new_version: u64, ) -> Result<(Manifest, Vec)> { if config.use_stable_row_ids && current_manifest @@ -1861,14 +2017,18 @@ impl Transaction { )) }); - let new_version = current_manifest.map_or(1, |m| m.version + 1); - match &self.operation { Operation::Clone { .. } => { return Err(Error::internal( "Clone operation should not enter build_manifest.".to_string(), )); } + Operation::Composite { .. } => { + return Err(Error::internal( + "Composite operation should be applied via build_manifest, not build_manifest_impl." + .to_string(), + )); + } Operation::Append { fragments } => { final_fragments.extend(maybe_existing_fragments?.clone()); let mut new_fragments = @@ -3345,6 +3505,16 @@ impl TryFrom for Transaction { })) => Operation::UpdateBases { new_bases: new_bases.into_iter().map(BasePath::from).collect(), }, + Some(pb::transaction::Operation::Composite(pb::transaction::CompositeOperation { + transactions, + })) => { + let operations = transactions + .into_iter() + .map(|sub| Self::try_from(sub).map(|t| t.operation)) + .collect::>>()?; + Operation::validate_composite(&operations)?; + Operation::Composite { operations } + } None => { return Err(Error::internal( "Transaction message did not contain an operation".to_string(), @@ -3632,6 +3802,25 @@ impl From<&Transaction> for pb::Transaction { .collect::>(), }) } + Operation::Composite { operations } => { + // Only the sub-transactions' `operation` field is meaningful; + // the outer transaction governs read_version/uuid/tag/properties. + let transactions = operations + .iter() + .map(|op| { + Self::from(&Transaction { + read_version: 0, + uuid: String::new(), + operation: op.clone(), + tag: None, + transaction_properties: None, + }) + }) + .collect(); + pb::transaction::Operation::Composite(pb::transaction::CompositeOperation { + transactions, + }) + } }; let transaction_properties = value @@ -3705,6 +3894,18 @@ pub fn validate_operation(manifest: Option<&Manifest>, operation: &Operation) -> return Ok(()); } (None, Operation::Clone { .. }) => return Ok(()), + (None, Operation::Composite { operations }) => { + Operation::validate_composite(operations)?; + // A composite can create a dataset only by starting with Overwrite. + // Sub-operations are validated against their intermediate manifests + // during the build_manifest fold. + if !matches!(operations.first(), Some(Operation::Overwrite { .. })) { + return Err(Error::invalid_input( + "A composite operation on a non-existent dataset must start with Overwrite", + )); + } + return Ok(()); + } (Some(manifest), _) => manifest, (None, _) => { return Err(Error::invalid_input(format!( @@ -3745,6 +3946,9 @@ pub fn validate_operation(manifest: Option<&Manifest>, operation: &Operation) -> schema_fragments_valid(Some(manifest), &manifest.schema, updated_fragments)?; schema_fragments_valid(Some(manifest), &manifest.schema, new_fragments) } + // Structural rules only; each sub-operation is validated against its + // intermediate manifest during the build_manifest fold. + Operation::Composite { operations } => Operation::validate_composite(operations), _ => Ok(()), } } @@ -4136,6 +4340,330 @@ mod tests { ); } + fn sample_update_config(key: &str, value: &str) -> Operation { + Operation::UpdateConfig { + config_updates: Some(UpdateMap { + update_entries: vec![(key, value).into()], + replace: false, + }), + table_metadata_updates: None, + schema_metadata_updates: None, + field_metadata_updates: HashMap::new(), + } + } + + #[test] + fn test_composite_pb_roundtrip() { + let operations = vec![ + Operation::Append { + fragments: vec![Fragment::new(0)], + }, + sample_update_config("k", "v"), + ]; + let transaction = Transaction::new_from_version( + 3, + Operation::Composite { + operations: operations.clone(), + }, + ); + + let pb_transaction = pb::Transaction::from(&transaction); + let Some(pb::transaction::Operation::Composite(ref composite)) = pb_transaction.operation + else { + panic!("expected a composite operation"); + }; + assert_eq!(composite.transactions.len(), 2); + for inner in &composite.transactions { + // Only the operation of each sub-transaction is meaningful. + assert_eq!(inner.read_version, 0); + assert!(inner.uuid.is_empty()); + } + + let roundtripped = Transaction::try_from(pb_transaction).unwrap(); + assert_eq!(roundtripped.read_version, 3); + assert_eq!(roundtripped.operation, Operation::Composite { operations }); + } + + #[test] + fn test_composite_validation() { + assert!(Operation::validate_composite(&[]).is_err()); + for op in [ + Operation::Restore { version: 1 }, + Operation::Clone { + is_shallow: true, + ref_name: None, + ref_version: 1, + ref_path: "path".to_string(), + branch_name: None, + }, + Operation::Composite { + operations: vec![Operation::Append { fragments: vec![] }], + }, + ] { + let err = Operation::validate_composite(std::slice::from_ref(&op)).unwrap_err(); + assert!( + err.to_string().contains(op.name()), + "error should name the offending operation: {}", + err + ); + } + Operation::validate_composite(&[Operation::Append { fragments: vec![] }]).unwrap(); + } + + #[test] + fn test_composite_pb_decode_rejects_invalid() { + let nested = pb::Transaction { + operation: Some(pb::transaction::Operation::Composite( + pb::transaction::CompositeOperation { + transactions: vec![pb::Transaction { + operation: Some(pb::transaction::Operation::Composite( + pb::transaction::CompositeOperation { + transactions: vec![pb::Transaction { + operation: Some(pb::transaction::Operation::Append( + pb::transaction::Append { fragments: vec![] }, + )), + ..Default::default() + }], + }, + )), + ..Default::default() + }], + }, + )), + ..Default::default() + }; + assert!(Transaction::try_from(nested).is_err()); + + let empty = pb::Transaction { + operation: Some(pb::transaction::Operation::Composite( + pb::transaction::CompositeOperation { + transactions: vec![], + }, + )), + ..Default::default() + }; + assert!(Transaction::try_from(empty).is_err()); + } + + #[test] + fn test_composite_build_manifest_folds_sequentially() { + let manifest = sample_manifest(); + assert_eq!(manifest.version, 1); + + let transaction = Transaction::new_from_version( + 1, + Operation::Composite { + operations: vec![ + Operation::Append { + fragments: vec![Fragment::new(0)], + }, + Operation::Append { + fragments: vec![Fragment::new(0)], + }, + Operation::Delete { + updated_fragments: vec![], + deleted_fragment_ids: vec![0], + predicate: "true".to_string(), + }, + ], + }, + ); + + let (new_manifest, _) = transaction + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .unwrap(); + + // The whole composite creates exactly one new version. + assert_eq!(new_manifest.version, 2); + // The two appends received distinct sequential fragment ids and the + // delete removed the pre-existing fragment 0, including a fragment + // interleaving added and pre-existing ids. + let ids = new_manifest + .fragments + .iter() + .map(|f| f.id) + .collect::>(); + assert_eq!(ids, vec![1, 2]); + assert_eq!(new_manifest.max_fragment_id(), Some(2)); + } + + #[test] + fn test_composite_build_manifest_config_and_index() { + let manifest = sample_manifest(); + let index = sample_index_metadata("vector_idx"); + let transaction = Transaction::new_from_version( + 1, + Operation::Composite { + operations: vec![ + sample_update_config("dataset.key", "value"), + Operation::CreateIndex { + new_indices: vec![index.clone()], + removed_indices: vec![], + }, + ], + }, + ); + + let (new_manifest, final_indices) = transaction + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .unwrap(); + + assert_eq!(new_manifest.version, 2); + assert_eq!( + new_manifest.config.get("dataset.key"), + Some(&"value".to_string()) + ); + assert_eq!(final_indices.len(), 1); + assert_eq!(final_indices[0].uuid, index.uuid); + } + + #[test] + fn test_composite_reserved_fragment_ids_append_and_index() { + let manifest = sample_manifest(); + // Reserve fragment ids in a prior transaction so data and index files + // can be written against known ids before the composite commit. + let (manifest, _) = + Transaction::new_from_version(1, Operation::ReserveFragments { num_fragments: 2 }) + .build_manifest( + Some(&manifest), + vec![], + "txn-reserve", + &ManifestWriteConfig::default(), + ) + .unwrap(); + assert_eq!(manifest.max_fragment_id, Some(2)); + + let mut index = sample_index_metadata("vector_idx"); + index.fragment_bitmap = Some([1, 2].into_iter().collect()); + let transaction = Transaction::new_from_version( + 2, + Operation::Composite { + operations: vec![ + Operation::Append { + fragments: vec![Fragment::new(1), Fragment::new(2)], + }, + Operation::CreateIndex { + new_indices: vec![index], + removed_indices: vec![], + }, + ], + }, + ); + + let (new_manifest, final_indices) = transaction + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .unwrap(); + + // Pre-reserved fragment ids pass through the append unchanged, so the + // index's fragment bitmap references the final ids. + let ids = new_manifest + .fragments + .iter() + .map(|f| f.id) + .collect::>(); + assert_eq!(ids, vec![0, 1, 2]); + assert_eq!(final_indices.len(), 1); + assert_eq!( + final_indices[0] + .fragment_bitmap + .as_ref() + .unwrap() + .iter() + .collect::>(), + vec![1, 2] + ); + } + + #[test] + fn test_composite_build_manifest_rejects_invalid() { + let manifest = sample_manifest(); + for operations in [ + vec![], + vec![Operation::Restore { version: 1 }], + vec![Operation::Composite { + operations: vec![Operation::Append { fragments: vec![] }], + }], + ] { + let transaction = Transaction::new_from_version(1, Operation::Composite { operations }); + assert!( + transaction + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .is_err() + ); + } + } + + #[test] + fn test_composite_stable_row_ids_version_meta_uses_final_version() { + // A create-mode composite: Overwrite followed by Append with stable + // row ids. Both sub-operations must stamp row version metadata with + // the single committed version (1), not per-step intermediate versions. + let schema = ArrowSchema::new(vec![ArrowField::new("id", DataType::Int32, false)]); + let lance_schema = LanceSchema::try_from(&schema).unwrap(); + let mut frag1 = Fragment::new(0); + frag1.physical_rows = Some(10); + let mut frag2 = Fragment::new(0); + frag2.physical_rows = Some(5); + + let config = ManifestWriteConfig { + use_stable_row_ids: true, + ..Default::default() + }; + let transaction = Transaction::new_from_version( + 0, + Operation::Composite { + operations: vec![ + Operation::Overwrite { + fragments: vec![frag1], + schema: lance_schema, + config_upsert_values: None, + initial_bases: None, + }, + Operation::Append { + fragments: vec![frag2], + }, + ], + }, + ); + + let (manifest, _) = transaction + .build_manifest(None, vec![], "txn", &config) + .unwrap(); + + assert_eq!(manifest.version, 1); + assert_eq!(manifest.fragments.len(), 2); + let fragment_ids = manifest.fragments.iter().map(|f| f.id).collect::>(); + for fragment_id in fragment_ids { + let versions = created_at_versions(&manifest, fragment_id); + assert!( + versions.iter().all(|v| *v == 1), + "fragment {} has created-at versions {:?}, expected all 1", + fragment_id, + versions + ); + } + } + #[test] fn test_remove_tombstoned_data_files() { // Create a fragment with mixed data files: some normal, some fully tombstoned diff --git a/rust/lance/src/dataset/write/commit.rs b/rust/lance/src/dataset/write/commit.rs index f8d09d3c55e..e64020cbd59 100644 --- a/rust/lance/src/dataset/write/commit.rs +++ b/rust/lance/src/dataset/write/commit.rs @@ -304,12 +304,16 @@ impl<'a> CommitBuilder<'a> { } }; - if dest.dataset().is_none() - && !matches!( - transaction.operation, - Operation::Overwrite { .. } | Operation::Clone { .. } - ) - { + let can_create_dataset = match &transaction.operation { + Operation::Overwrite { .. } | Operation::Clone { .. } => true, + // A composite can create a dataset by starting with Overwrite; + // validate_operation enforces this below. + Operation::Composite { operations } => { + matches!(operations.first(), Some(Operation::Overwrite { .. })) + } + _ => false, + }; + if dest.dataset().is_none() && !can_create_dataset { return Err(Error::dataset_not_found( base_path.to_string(), "The dataset must already exist unless the operation is Overwrite".into(), @@ -1002,6 +1006,235 @@ mod tests { assert_eq!(transaction.read_version, 1); } + fn sample_update_config_op(key: &str, value: &str) -> Operation { + use crate::dataset::transaction::UpdateMap; + Operation::UpdateConfig { + config_updates: Some(UpdateMap { + update_entries: vec![(key, value).into()], + replace: false, + }), + table_metadata_updates: None, + schema_metadata_updates: None, + field_metadata_updates: HashMap::new(), + } + } + + fn composite_transaction(read_version: u64, operations: Vec) -> Transaction { + Transaction { + uuid: uuid::Uuid::new_v4().hyphenated().to_string(), + operation: Operation::Composite { operations }, + read_version, + tag: None, + transaction_properties: None, + } + } + + async fn sample_dataset(uri: &str) -> (Arc, RecordBatch) { + let schema = Arc::new(ArrowSchema::new(vec![ArrowField::new( + "i", + DataType::Int32, + false, + )])); + let batch = RecordBatch::try_new( + schema, + vec![Arc::new(Int32Array::from_iter_values(0..10_i32))], + ) + .unwrap(); + let dataset = InsertBuilder::new(uri) + .execute(vec![batch.clone()]) + .await + .unwrap(); + (Arc::new(dataset), batch) + } + + async fn uncommitted_append(dataset: &Arc, batch: &RecordBatch) -> Operation { + use crate::dataset::WriteMode; + let transaction = InsertBuilder::new(dataset.clone()) + .with_params(&WriteParams { + mode: WriteMode::Append, + ..Default::default() + }) + .execute_uncommitted(vec![batch.clone()]) + .await + .unwrap(); + transaction.operation + } + + #[tokio::test] + async fn test_commit_composite() { + let (dataset, batch) = sample_dataset("memory://composite").await; + let append = uncommitted_append(&dataset, &batch).await; + + let composite = composite_transaction( + dataset.manifest.version, + vec![append, sample_update_config_op("composite.key", "value")], + ); + let new_ds = CommitBuilder::new(dataset.clone()) + .execute(composite) + .await + .unwrap(); + + // Both effects land in a single new version. + assert_eq!(new_ds.manifest.version, 2); + assert_eq!(new_ds.count_rows(None).await.unwrap(), 20); + assert_eq!( + new_ds.manifest.config.get("composite.key"), + Some(&"value".to_string()) + ); + + // The transaction is not inlined into the manifest so that released + // readers, which decode the inline section eagerly, can still open + // this version. It is still readable from the transaction file. + assert!(new_ds.manifest.transaction_section.is_none()); + let read_back = new_ds.read_transaction().await.unwrap().unwrap(); + assert!(matches!( + read_back.operation, + Operation::Composite { ref operations } if operations.len() == 2 + )); + } + + #[tokio::test] + async fn test_commit_composite_rebases_on_concurrent_append() { + let (dataset, batch) = sample_dataset("memory://composite-rebase").await; + let composite_append = uncommitted_append(&dataset, &batch).await; + + // A concurrent plain append lands first, at version 2. + let concurrent_append = uncommitted_append(&dataset, &batch).await; + CommitBuilder::new(dataset.clone()) + .execute(Transaction { + uuid: uuid::Uuid::new_v4().hyphenated().to_string(), + operation: concurrent_append, + read_version: 1, + tag: None, + transaction_properties: None, + }) + .await + .unwrap(); + + // The composite, with a stale read version, rebases and succeeds. + let composite = composite_transaction( + 1, + vec![ + composite_append, + sample_update_config_op("composite.key", "value"), + ], + ); + let new_ds = CommitBuilder::new(dataset.clone()) + .execute(composite) + .await + .unwrap(); + + assert_eq!(new_ds.manifest.version, 3); + assert_eq!(new_ds.count_rows(None).await.unwrap(), 30); + assert_eq!( + new_ds.manifest.config.get("composite.key"), + Some(&"value".to_string()) + ); + } + + #[tokio::test] + async fn test_commit_composite_incompatible_on_config_conflict() { + let (dataset, _batch) = sample_dataset("memory://composite-conflict").await; + + // A plain config update lands first, at version 2. + CommitBuilder::new(dataset.clone()) + .execute(Transaction { + uuid: uuid::Uuid::new_v4().hyphenated().to_string(), + operation: sample_update_config_op("composite.key", "first"), + read_version: 1, + tag: None, + transaction_properties: None, + }) + .await + .unwrap(); + + // A composite upserting the same key from the same read version must + // surface the sub-operation's incompatibility. + let composite = + composite_transaction(1, vec![sample_update_config_op("composite.key", "second")]); + let res = CommitBuilder::new(dataset.clone()).execute(composite).await; + assert!( + matches!(res, Err(Error::IncompatibleTransaction { .. })), + "got {res:?}" + ); + } + + #[tokio::test] + async fn test_commit_plain_rebases_on_concurrent_composite() { + let (dataset, batch) = sample_dataset("memory://composite-other").await; + let plain_append = uncommitted_append(&dataset, &batch).await; + + // A composite lands first, at version 2. Composites are not inlined + // into the manifest, so the plain transaction's conflict resolution + // must read it from the external transaction file and decompose it. + let composite_append = uncommitted_append(&dataset, &batch).await; + CommitBuilder::new(dataset.clone()) + .execute(composite_transaction( + 1, + vec![ + composite_append, + sample_update_config_op("composite.key", "value"), + ], + )) + .await + .unwrap(); + + let new_ds = CommitBuilder::new(dataset.clone()) + .execute(Transaction { + uuid: uuid::Uuid::new_v4().hyphenated().to_string(), + operation: plain_append, + read_version: 1, + tag: None, + transaction_properties: None, + }) + .await + .unwrap(); + + assert_eq!(new_ds.manifest.version, 3); + assert_eq!(new_ds.count_rows(None).await.unwrap(), 30); + } + + #[tokio::test] + async fn test_commit_composite_creates_dataset() { + let schema = ArrowSchema::new(vec![ArrowField::new("i", DataType::Int32, false)]); + let lance_schema = lance_core::datatypes::Schema::try_from(&schema).unwrap(); + + // A composite starting with Overwrite can create a dataset. + let composite = composite_transaction( + 0, + vec![ + Operation::Overwrite { + fragments: vec![], + schema: lance_schema, + config_upsert_values: None, + initial_bases: None, + }, + sample_update_config_op("created.by", "composite"), + ], + ); + let dataset = CommitBuilder::new("memory://composite-create") + .execute(composite) + .await + .unwrap(); + assert_eq!(dataset.manifest.version, 1); + assert_eq!(dataset.count_rows(None).await.unwrap(), 0); + assert_eq!( + dataset.manifest.config.get("created.by"), + Some(&"composite".to_string()) + ); + + // A composite not starting with Overwrite cannot. + let composite = + composite_transaction(0, vec![sample_update_config_op("created.by", "composite")]); + let res = CommitBuilder::new("memory://composite-create-2") + .execute(composite) + .await; + assert!( + matches!(res, Err(Error::DatasetNotFound { .. })), + "got {res:?}" + ); + } + /// On non-lexically-ordered stores (e.g. S3 Express) a commit should use the /// version hint (a few HEAD probes, O(k)) instead of a full O(n) listing. #[tokio::test] diff --git a/rust/lance/src/io/commit/conflict_resolver.rs b/rust/lance/src/io/commit/conflict_resolver.rs index d95821dd130..facb4505789 100644 --- a/rust/lance/src/io/commit/conflict_resolver.rs +++ b/rust/lance/src/io/commit/conflict_resolver.rs @@ -35,6 +35,10 @@ pub struct TransactionRebase<'a> { /// Merged generations from conflicting UpdateMemWalState transactions. /// Used when rebasing CreateIndex of MemWalIndex. conflicting_mem_wal_merged_gens: Vec, + /// One rebase per sub-operation when the transaction is a composite; + /// empty otherwise. Composite conflict checks and rebases delegate to + /// these so each sub-operation keeps its existing rules. + sub_rebases: Vec, } impl<'a> TransactionRebase<'a> { @@ -42,6 +46,39 @@ impl<'a> TransactionRebase<'a> { dataset: &Dataset, transaction: Transaction, affected_rows: Option<&'a RowAddrTreeMap>, + ) -> Result { + if let Operation::Composite { operations } = &transaction.operation { + // A single affected-rows map is ambiguous across multiple + // Delete/Update sub-operations, so sub-rebases fall back to + // fragment-level conflict detection. + let mut sub_rebases = Vec::with_capacity(operations.len()); + for operation in operations { + let sub_transaction = Transaction { + read_version: transaction.read_version, + uuid: transaction.uuid.clone(), + operation: operation.clone(), + tag: None, + transaction_properties: None, + }; + sub_rebases.push(Self::try_new_single(dataset, sub_transaction, None).await?); + } + return Ok(Self { + transaction, + initial_fragments: HashMap::new(), + modified_fragment_ids: HashSet::new(), + affected_rows: None, + conflicting_frag_reuse_indices: Vec::new(), + conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases, + }); + } + Self::try_new_single(dataset, transaction, affected_rows).await + } + + async fn try_new_single( + dataset: &Dataset, + transaction: Transaction, + affected_rows: Option<&'a RowAddrTreeMap>, ) -> Result { match &transaction.operation { // These operations add new fragments or don't modify any. @@ -61,6 +98,7 @@ impl<'a> TransactionRebase<'a> { modified_fragment_ids: HashSet::new(), conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }), Operation::Delete { updated_fragments, @@ -89,6 +127,7 @@ impl<'a> TransactionRebase<'a> { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }); } @@ -102,6 +141,7 @@ impl<'a> TransactionRebase<'a> { modified_fragment_ids, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }) } Operation::Rewrite { groups, .. } => { @@ -120,6 +160,7 @@ impl<'a> TransactionRebase<'a> { modified_fragment_ids, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }) } Operation::DataReplacement { replacements } => { @@ -135,6 +176,7 @@ impl<'a> TransactionRebase<'a> { modified_fragment_ids, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }) } Operation::Merge { fragments, .. } => { @@ -149,8 +191,12 @@ impl<'a> TransactionRebase<'a> { modified_fragment_ids, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }) } + Operation::Composite { .. } => Err(Error::internal( + "Composite sub-operations cannot be nested composites".to_string(), + )), } } @@ -204,6 +250,29 @@ impl<'a> TransactionRebase<'a> { /// Will return an error if the transaction is not valid. Otherwise, it will /// return Ok(()). pub fn check_txn(&mut self, other_transaction: &Transaction, other_version: u64) -> Result<()> { + // A committed composite conflicts iff any of its sub-operations + // conflicts: decompose it and check against each one. + if let Operation::Composite { operations } = &other_transaction.operation { + for operation in operations { + let sub_other = Transaction { + read_version: other_transaction.read_version, + uuid: other_transaction.uuid.clone(), + operation: operation.clone(), + tag: None, + transaction_properties: None, + }; + self.check_txn(&sub_other, other_version)?; + } + return Ok(()); + } + // A composite in flight delegates to its sub-operations' rebases so + // each keeps its existing conflict rules. + if matches!(self.transaction.operation, Operation::Composite { .. }) { + for sub_rebase in self.sub_rebases.iter_mut() { + sub_rebase.check_txn(other_transaction, other_version)?; + } + return Ok(()); + } let op = &self.transaction.operation; match op { Operation::Delete { .. } => self.check_delete_txn(other_transaction, other_version), @@ -235,6 +304,8 @@ impl<'a> TransactionRebase<'a> { Operation::UpdateBases { .. } => { self.check_add_bases_txn(other_transaction, other_version) } + // Handled by the delegation above. + Operation::Composite { .. } => unreachable!("composite check_txn delegates above"), } } @@ -245,6 +316,10 @@ impl<'a> TransactionRebase<'a> { ) -> Result<()> { if let Operation::Delete { .. } = &self.transaction.operation { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), Operation::CreateIndex { .. } | Operation::ReserveFragments { .. } | Operation::Clone { .. } @@ -393,6 +468,10 @@ impl<'a> TransactionRebase<'a> { } match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), Operation::CreateIndex { .. } | Operation::ReserveFragments { .. } | Operation::Project { .. } @@ -512,6 +591,10 @@ impl<'a> TransactionRebase<'a> { } = &self.transaction.operation { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), Operation::Append { .. } | Operation::Clone { .. } | Operation::UpdateBases { .. } => Ok(()), @@ -665,6 +748,10 @@ impl<'a> TransactionRebase<'a> { } = &self.transaction.operation { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), // Rewrite is only compatible with operations that don't touch // existing fragments or update fragments we don't touch. Operation::Append { .. } @@ -841,6 +928,10 @@ impl<'a> TransactionRebase<'a> { other_version: u64, ) -> Result<()> { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), Operation::Overwrite { .. } => { if self .transaction @@ -889,6 +980,10 @@ impl<'a> TransactionRebase<'a> { other_version: u64, ) -> Result<()> { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), // Append is not compatible with any operation that completely // overwrites the schema. Operation::Overwrite { .. } @@ -918,6 +1013,10 @@ impl<'a> TransactionRebase<'a> { ) -> Result<()> { if let Operation::DataReplacement { replacements } = &self.transaction.operation { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), Operation::Append { .. } | Operation::Clone { .. } | Operation::UpdateConfig { .. } @@ -1061,6 +1160,10 @@ impl<'a> TransactionRebase<'a> { other_version: u64, ) -> Result<()> { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), Operation::CreateIndex { .. } | Operation::ReserveFragments { .. } | Operation::Clone { .. } @@ -1090,6 +1193,10 @@ impl<'a> TransactionRebase<'a> { other_version: u64, ) -> Result<()> { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), Operation::Append { .. } | Operation::Delete { .. } | Operation::Overwrite { .. } @@ -1116,6 +1223,10 @@ impl<'a> TransactionRebase<'a> { other_version: u64, ) -> Result<()> { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), Operation::Overwrite { .. } | Operation::Restore { .. } => { Err(self.incompatible_conflict_err(other_transaction, other_version)) } @@ -1141,6 +1252,10 @@ impl<'a> TransactionRebase<'a> { other_version: u64, ) -> Result<()> { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), // Project is compatible with anything that doesn't change the schema Operation::Append { .. } | Operation::Update { .. } @@ -1176,6 +1291,10 @@ impl<'a> TransactionRebase<'a> { } = &self.transaction.operation { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), Operation::Overwrite { .. } => { // Updates to schema metadata or field metadata conflict with any kind // of overwrite. @@ -1235,6 +1354,10 @@ impl<'a> TransactionRebase<'a> { } = &self.transaction.operation { match &other_transaction.operation { + Operation::Composite { .. } => Err(Error::internal( + "composite transactions are decomposed before per-operation conflict checks" + .to_string(), + )), Operation::UpdateMemWalState { merged_generations: other_merged_generations, } => { @@ -1368,6 +1491,20 @@ impl<'a> TransactionRebase<'a> { /// Writes pub async fn finish(self, dataset: &Dataset) -> Result { + if matches!(self.transaction.operation, Operation::Composite { .. }) { + let mut operations = Vec::with_capacity(self.sub_rebases.len()); + for sub_rebase in self.sub_rebases { + let finished = sub_rebase.finish_single(dataset).await?; + operations.push(finished.operation); + } + let mut transaction = self.transaction; + transaction.operation = Operation::Composite { operations }; + return Ok(transaction); + } + self.finish_single(dataset).await + } + + async fn finish_single(self, dataset: &Dataset) -> Result { match &self.transaction.operation { Operation::Delete { .. } | Operation::Update { .. } => { self.finish_delete_update(dataset).await @@ -1385,6 +1522,9 @@ impl<'a> TransactionRebase<'a> { | Operation::UpdateConfig { .. } | Operation::UpdateMemWalState { .. } | Operation::UpdateBases { .. } => Ok(self.transaction), + Operation::Composite { .. } => Err(Error::internal( + "Composite sub-operations cannot be nested composites".to_string(), + )), } } @@ -2739,6 +2879,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; for (other, expected_conflict) in other_transactions.iter().zip(expected_conflicts) { @@ -2781,6 +2922,209 @@ mod tests { } } + #[test] + fn test_composite_conflicts() { + use io::commit::conflict_resolver::tests::{ConflictResult::*, modified_fragment_ids}; + + let fragment0 = Fragment::new(0); + let fragment1 = Fragment::new(1); + + let sub_operations = vec![ + Operation::Append { + fragments: vec![fragment1], + }, + Operation::Delete { + updated_fragments: vec![fragment0.clone()], + deleted_fragment_ids: vec![], + predicate: "x > 2".to_string(), + }, + create_update_config_for_test( + Some(HashMap::from_iter(vec![( + "lance.test".to_string(), + "value".to_string(), + )])), + None, + None, + None, + ), + ]; + let composite_transaction = Transaction::new( + 0, + Operation::Composite { + operations: sub_operations.clone(), + }, + None, + ); + + // A plain in-flight transaction checked against a committed composite: + // the composite decomposes and each sub-operation's rules apply. + let plain_cases = [ + ( + // Does not overlap any sub-operation. + Operation::Append { + fragments: vec![Fragment::new(3)], + }, + Compatible, + ), + ( + // Overlaps the composite's Delete of fragment 0. + Operation::Delete { + updated_fragments: vec![fragment0.clone()], + deleted_fragment_ids: vec![], + predicate: "x > 2".to_string(), + }, + Retryable, + ), + ( + // Upserts the same config key as the composite's UpdateConfig. + create_update_config_for_test( + Some(HashMap::from_iter(vec![( + "lance.test".to_string(), + "other-value".to_string(), + )])), + None, + None, + None, + ), + NotCompatible, + ), + ]; + for (operation, expected) in &plain_cases { + let mut rebase = TransactionRebase { + transaction: Transaction::new(0, operation.clone(), None), + initial_fragments: HashMap::new(), + modified_fragment_ids: modified_fragment_ids(operation).collect::>(), + affected_rows: None, + conflicting_frag_reuse_indices: Vec::new(), + conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), + }; + let result = rebase.check_txn(&composite_transaction, 1); + assert_conflict_result( + result, + expected, + operation, + &composite_transaction.operation, + ); + } + + // An in-flight composite delegates to its sub-operations' rebases. + let make_composite_rebase = || TransactionRebase { + transaction: composite_transaction.clone(), + initial_fragments: HashMap::new(), + modified_fragment_ids: HashSet::new(), + affected_rows: None, + conflicting_frag_reuse_indices: Vec::new(), + conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: sub_operations + .iter() + .map(|op| TransactionRebase { + transaction: Transaction::new(0, op.clone(), None), + initial_fragments: HashMap::new(), + modified_fragment_ids: modified_fragment_ids(op).collect::>(), + affected_rows: None, + conflicting_frag_reuse_indices: Vec::new(), + conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), + }) + .collect(), + }; + let composite_cases = [ + ( + Operation::Append { + fragments: vec![Fragment::new(3)], + }, + Compatible, + ), + ( + // Overlaps the composite's Delete sub-operation. + Operation::Delete { + updated_fragments: vec![fragment0.clone()], + deleted_fragment_ids: vec![], + predicate: "x > 2".to_string(), + }, + Retryable, + ), + ( + // The composite's Append sub-operation is incompatible with a + // concurrent Overwrite. + Operation::Overwrite { + fragments: vec![fragment0], + schema: lance_core::datatypes::Schema::default(), + config_upsert_values: None, + initial_bases: None, + }, + NotCompatible, + ), + ( + // Upserts the same config key as the composite's UpdateConfig. + create_update_config_for_test( + Some(HashMap::from_iter(vec![( + "lance.test".to_string(), + "other-value".to_string(), + )])), + None, + None, + None, + ), + NotCompatible, + ), + ]; + for (operation, expected) in &composite_cases { + let mut rebase = make_composite_rebase(); + let other = Transaction::new(0, operation.clone(), None); + let result = rebase.check_txn(&other, 1); + assert_conflict_result( + result, + expected, + &composite_transaction.operation, + operation, + ); + } + + // Composite vs composite: decomposed on both sides; the overlapping + // Delete sub-operations conflict. + let mut rebase = make_composite_rebase(); + let result = rebase.check_txn(&composite_transaction, 1); + assert_conflict_result( + result, + &Retryable, + &composite_transaction.operation, + &composite_transaction.operation, + ); + } + + fn assert_conflict_result( + result: Result<()>, + expected: &ConflictResult, + operation: &Operation, + other: &Operation, + ) { + match expected { + ConflictResult::Compatible => assert!( + result.is_ok(), + "Transaction {:?} should be compatible with {:?}, but was {:?}", + operation, + other, + result + ), + ConflictResult::NotCompatible => assert!( + matches!(result, Err(Error::IncompatibleTransaction { .. })), + "Transaction {:?} should be incompatible with {:?}, but was {:?}", + operation, + other, + result + ), + ConflictResult::Retryable => assert!( + matches!(result, Err(Error::RetryableCommitConflict { .. })), + "Transaction {:?} should be retryable with {:?}, but was {:?}", + operation, + other, + result + ), + } + } + #[test] fn test_create_index_conflicts_only_on_same_name() { let index0 = IndexMetadata { @@ -2816,6 +3160,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; let same_name = Transaction::new( @@ -2870,6 +3215,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; let different_name_result = rebase.check_txn(&different_name, 1); assert!( @@ -3255,6 +3601,9 @@ mod tests { Operation::DataReplacement { replacements } => { Box::new(replacements.iter().map(|r| r.0)) } + Operation::Composite { operations } => { + Box::new(operations.iter().flat_map(modified_fragment_ids)) + } } } @@ -3507,6 +3856,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; let result = rebase.check_txn(&txn2, 1); @@ -3570,6 +3920,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; let result = rebase.check_txn(&committed_txn, 1); @@ -3608,6 +3959,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; let result = rebase.check_txn(&committed_txn, 1); @@ -3647,6 +3999,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; let result = rebase.check_txn(&committed_txn, 1); @@ -3686,6 +4039,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; let result = rebase.check_txn(&committed_txn, 1); @@ -3736,6 +4090,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; let result = rebase.check_txn(&committed_txn, 1); @@ -3761,6 +4116,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; let result_higher = rebase_higher.check_txn(&committed_txn, 1); @@ -3807,6 +4163,7 @@ mod tests { affected_rows: None, conflicting_frag_reuse_indices: Vec::new(), conflicting_mem_wal_merged_gens: Vec::new(), + sub_rebases: Vec::new(), }; // CreateIndex of MemWalIndex should be compatible with UpdateMemWalState From 75c78f84d073319781ba9cd4b2c6bbbec45b5d18 Mon Sep 17 00:00:00 2001 From: Jack Ye Date: Sat, 11 Jul 2026 02:55:31 -0700 Subject: [PATCH 2/3] fix: address composite operation review findings - Reject sub-operations that reference fragment ids assigned to fragments added earlier in the same composite: those ids shift when the commit is rebuilt on conflict retry and would silently retarget another writer's fragment. Explicitly reserved ids stay allowed. - Apply the outer transaction's tag to the final manifest. - Keep the inline transaction section when disable_transaction_file is set so a composite commit always leaves a transaction record. - Stamp row-version metadata with the true previous version for non-first Update sub-operations on stable-row-id datasets. - Advance the outer read_version after sub-operation rebases. - Recognize overwrite-leading composites in Dataset::commit and the CommitBuilder storage-format check. - Surface unsupported operations as errors instead of panics across the JNI boundary; warn when affected_rows is ignored for composites. --- java/lance-jni/src/transaction.rs | 16 +- python/python/lance/dataset.py | 28 ++- rust/lance/src/dataset.rs | 29 ++- rust/lance/src/dataset/transaction.rs | 210 +++++++++++++++++- rust/lance/src/dataset/write/commit.rs | 19 +- rust/lance/src/io/commit/conflict_resolver.rs | 12 + 6 files changed, 279 insertions(+), 35 deletions(-) diff --git a/java/lance-jni/src/transaction.rs b/java/lance-jni/src/transaction.rs index 4f899f56ff2..6afb29e6510 100644 --- a/java/lance-jni/src/transaction.rs +++ b/java/lance-jni/src/transaction.rs @@ -575,7 +575,12 @@ fn convert_to_java_operation_inner<'local>( "(I)V", &[JValue::Int(num_fragments as i32)], )?), - _ => unimplemented!(), + // A panic here unwinds across the JNI boundary and aborts the JVM, + // so unmapped operations must surface as errors instead. + other => Err(Error::unsupported_error(format!( + "operation {} is not supported by the Java binding yet", + other.name() + ))), } } @@ -1303,7 +1308,14 @@ fn convert_to_rust_operation( removed_indices, }); } - _ => unimplemented!(), + // A panic here unwinds across the JNI boundary and aborts the JVM, + // so unmapped operations must surface as errors instead. + other => { + return Err(Error::unsupported_error(format!( + "operation {} is not supported by the Java binding yet", + other + ))); + } }; Ok(op) } diff --git a/python/python/lance/dataset.py b/python/python/lance/dataset.py index df5b4cb748b..c1ddc7ef0d8 100644 --- a/python/python/lance/dataset.py +++ b/python/python/lance/dataset.py @@ -6088,7 +6088,12 @@ class Composite(BaseOperation): multi-statement transactions such as an append and an index update that become visible atomically. - The operations must not contain ``Restore`` or a nested ``Composite``. + The operations must not contain ``Restore``, ``Clone``, or a nested + ``Composite``. Sub-operations must not reference fragment IDs that + are assigned to fragments added earlier in the same composite — + such IDs are not stable across commit retries. To reference new + fragments (e.g. an index over appended data), reserve fragment IDs + first and set them on the fragments explicitly. Attributes ---------- @@ -6107,27 +6112,32 @@ class Composite(BaseOperation): >>> operation = lance.LanceOperation.Composite( ... [ ... lance.LanceOperation.Append([fragment]), - ... lance.LanceOperation.Delete([], [1], "a >= 3"), + ... lance.LanceOperation.Delete([], [0], "a < 3"), ... ] ... ) >>> dataset = lance.LanceDataset.commit("example", operation, ... read_version=dataset.version) >>> dataset.to_table().to_pandas() a b - 0 1 a - 1 2 b + 0 3 c + 1 4 d """ operations: Iterable[LanceOperation.BaseOperation] def __post_init__(self): - if not isinstance(self.operations, list) or not all( - isinstance(op, LanceOperation.BaseOperation) for op in self.operations - ): + operations = list(self.operations) + invalid = [ + op + for op in operations + if not isinstance(op, LanceOperation.BaseOperation) + ] + if invalid: raise TypeError( - "operations must be list[LanceOperation.BaseOperation], " - f"got {type(self.operations)}" + "operations must all be LanceOperation.BaseOperation, " + f"got {[type(op) for op in invalid]}" ) + self.operations = operations @dataclass diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index 449ccf6683e..8a65202e652 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -1384,11 +1384,15 @@ impl Dataset { detached: bool, ) -> Result { let read_version = read_version.map_or_else( - || match operation { - Operation::Overwrite { .. } | Operation::Restore { .. } => Ok(0), - _ => Err(Error::invalid_input( - "read_version must be specified for this operation", - )), + || { + if operation.overwrites_dataset() || matches!(operation, Operation::Restore { .. }) + { + Ok(0) + } else { + Err(Error::invalid_input( + "read_version must be specified for this operation", + )) + } }, Ok, )?; @@ -3695,13 +3699,18 @@ pub(crate) async fn write_manifest_file( // dataset and fail on unknown operation types, which would make this // version unopenable for them. Without the inline section they open // normally; only paths that need the transaction contents (conflict - // resolution, history) read the external transaction file. + // resolution, history) read the external transaction file. When the + // external file is disabled the inline section is the only transaction + // record, so it must be kept — dropping both would break conflict + // resolution for every later writer. // TODO: inline composites once readers with tolerant decoding are // widespread. - if matches!( - transaction.map(|t| &t.operation), - Some(Operation::Composite { .. }) - ) { + if !config.disable_transaction_file + && matches!( + transaction.map(|t| &t.operation), + Some(Operation::Composite { .. }) + ) + { transaction = None; } if config.auto_set_feature_flags { diff --git a/rust/lance/src/dataset/transaction.rs b/rust/lance/src/dataset/transaction.rs index e3e0de47171..b08d35516e4 100644 --- a/rust/lance/src/dataset/transaction.rs +++ b/rust/lance/src/dataset/transaction.rs @@ -1593,6 +1593,72 @@ impl Operation { } Ok(()) } + + /// Whether this operation replaces the dataset wholesale: a plain + /// Overwrite, or a composite whose first sub-operation is an Overwrite. + pub(crate) fn overwrites_dataset(&self) -> bool { + match self { + Self::Overwrite { .. } => true, + Self::Composite { operations } => { + matches!(operations.first(), Some(Self::Overwrite { .. })) + } + _ => false, + } + } + + /// Fragment ids this operation expects to already exist in the manifest + /// it applies to, as opposed to ids of fragments it adds. + fn referenced_fragment_ids(&self) -> HashSet { + match self { + Self::Delete { + updated_fragments, + deleted_fragment_ids, + .. + } => updated_fragments + .iter() + .map(|f| f.id) + .chain(deleted_fragment_ids.iter().copied()) + .collect(), + Self::Update { + updated_fragments, + removed_fragment_ids, + .. + } => updated_fragments + .iter() + .map(|f| f.id) + .chain(removed_fragment_ids.iter().copied()) + .collect(), + Self::Rewrite { groups, .. } => groups + .iter() + .flat_map(|g| g.old_fragments.iter().map(|f| f.id)) + .collect(), + Self::DataReplacement { replacements } => replacements.iter().map(|r| r.0).collect(), + Self::Merge { fragments, .. } => fragments.iter().map(|f| f.id).collect(), + Self::CreateIndex { new_indices, .. } => new_indices + .iter() + .filter_map(|index| index.fragment_bitmap.as_ref()) + .flat_map(|bitmap| bitmap.iter().map(u64::from)) + .collect(), + _ => HashSet::new(), + } + } + + /// Ids explicitly pre-assigned (non-zero, e.g. via ReserveFragments) on + /// fragments this operation adds. Unlike ids assigned during + /// `build_manifest`, these are stable across commit retries. + fn explicit_new_fragment_ids(&self) -> HashSet { + let added: Box> = match self { + Self::Append { fragments } | Self::Overwrite { fragments, .. } => { + Box::new(fragments.iter()) + } + Self::Update { new_fragments, .. } => Box::new(new_fragments.iter()), + Self::Rewrite { groups, .. } => { + Box::new(groups.iter().flat_map(|g| g.new_fragments.iter())) + } + _ => Box::new(std::iter::empty()), + }; + added.filter(|f| f.id != 0).map(|f| f.id).collect() + } } /// Helper function to apply UpdateMap changes to a HashMap @@ -1874,7 +1940,33 @@ impl Transaction { Operation::validate_composite(operations)?; let mut state: Option<(Manifest, Vec)> = None; let mut initial_indices = Some(current_indices); - for operation in operations { + // Ids assigned during this fold to fragments that arrived with id 0. + // Such ids change whenever the commit is rebuilt against a newer + // manifest (conflict retry), so later sub-operations must not + // reference them: the reference would silently retarget another + // writer's fragment. Explicitly pre-assigned ids (ReserveFragments) + // are stable and stay referenceable. + let mut unstable_created_ids: HashSet = HashSet::new(); + let mut previous_ids: HashSet = current_manifest + .map(|m| m.fragments.iter().map(|f| f.id).collect()) + .unwrap_or_default(); + for (position, operation) in operations.iter().enumerate() { + let forward_references = operation + .referenced_fragment_ids() + .intersection(&unstable_created_ids) + .copied() + .collect::>(); + if !forward_references.is_empty() { + return Err(Error::invalid_input(format!( + "Composite sub-operation {} ({}) references fragment ids {:?} that were \ + assigned at commit time to fragments added earlier in this composite; \ + these ids are not stable across commit retries. Reserve fragment ids \ + with ReserveFragments and set them on the new fragments instead", + position, + operation.name(), + forward_references, + ))); + } let sub_transaction = Self { read_version: self.read_version, uuid: self.uuid.clone(), @@ -1887,7 +1979,10 @@ impl Transaction { // Validate each sub-operation against the manifest it will // actually apply to; entry-point validation only sees the // base manifest and cannot account for schema changes made - // by earlier sub-operations. + // by earlier sub-operations. Note this validation also + // reruns on every commit retry against the caught-up + // manifest, which is stricter than the entry-point-only + // validation of plain operations. validate_operation(current_manifest, operation)?; sub_transaction.build_manifest_impl( current_manifest, @@ -1909,10 +2004,21 @@ impl Transaction { } }; manifest.version = new_version; + let current_ids: HashSet = manifest.fragments.iter().map(|f| f.id).collect(); + let explicit_ids = operation.explicit_new_fragment_ids(); + unstable_created_ids.extend( + current_ids + .difference(&previous_ids) + .filter(|id| !explicit_ids.contains(id)), + ); + previous_ids = current_ids; state = Some((manifest, indices)); } // validate_composite rejects empty operation lists, so state is set. - Ok(state.expect("non-empty composite")) + let (mut manifest, indices) = state.expect("non-empty composite"); + // The outer transaction governs the tag; sub-transactions carry none. + manifest.tag.clone_from(&self.tag); + Ok((manifest, indices)) } fn build_manifest_impl( @@ -2106,7 +2212,12 @@ impl Transaction { && let Some(UpdatedFragmentOffsets(off_map)) = updated_fragment_offsets && !off_map.is_empty() { - let prev_version = current_manifest.map(|m| m.version).unwrap_or(0); + // Derived from new_version rather than current_manifest: + // inside a composite fold the intermediate manifest's + // version is already pinned to new_version, which would + // fabricate "updated at this commit" stamps for uncovered + // rows. new_version is always current version + 1. + let prev_version = new_version - 1; for fragment in final_fragments.iter_mut() { let Some(bitmap) = off_map.get(&fragment.id) else { continue; @@ -4589,6 +4700,97 @@ mod tests { ); } + #[test] + fn test_composite_rejects_references_to_assigned_ids() { + let manifest = sample_manifest(); + + // Deleting the id the append will receive: the id is assigned at + // commit time and shifts under concurrent-append rebase, so this + // must be rejected. + let delete_guessed = Transaction::new_from_version( + 1, + Operation::Composite { + operations: vec![ + Operation::Append { + fragments: vec![Fragment::new(0)], + }, + Operation::Delete { + updated_fragments: vec![], + deleted_fragment_ids: vec![1], + predicate: "true".to_string(), + }, + ], + }, + ); + let err = delete_guessed + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .unwrap_err(); + assert!( + err.to_string().contains("ReserveFragments"), + "error should point at reserving ids: {}", + err + ); + + // Same for an index bitmap over an unreserved append. + let mut index = sample_index_metadata("vector_idx"); + index.fragment_bitmap = Some([1].into_iter().collect()); + let index_guessed = Transaction::new_from_version( + 1, + Operation::Composite { + operations: vec![ + Operation::Append { + fragments: vec![Fragment::new(0)], + }, + Operation::CreateIndex { + new_indices: vec![index], + removed_indices: vec![], + }, + ], + }, + ); + assert!( + index_guessed + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .is_err() + ); + + // Referencing a pre-existing fragment is fine even when the composite + // also appends (covered by test_composite_build_manifest_folds_sequentially), + // and explicitly pre-assigned ids stay referenceable (covered by + // test_composite_reserved_fragment_ids_append_and_index). + } + + #[test] + fn test_composite_build_manifest_applies_outer_tag() { + let manifest = sample_manifest(); + let transaction = Transaction::new( + 1, + Operation::Composite { + operations: vec![sample_update_config("k", "v")], + }, + Some("v1-tag".to_string()), + ); + let (new_manifest, _) = transaction + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .unwrap(); + assert_eq!(new_manifest.tag.as_deref(), Some("v1-tag")); + } + #[test] fn test_composite_build_manifest_rejects_invalid() { let manifest = sample_manifest(); diff --git a/rust/lance/src/dataset/write/commit.rs b/rust/lance/src/dataset/write/commit.rs index e64020cbd59..1ac09641e8e 100644 --- a/rust/lance/src/dataset/write/commit.rs +++ b/rust/lance/src/dataset/write/commit.rs @@ -174,6 +174,10 @@ impl<'a> CommitBuilder<'a> { /// Provide the set of row addresses that were deleted or updated. This is /// used to perform fast conflict resolution. + /// + /// Ignored for composite transactions: a single map is ambiguous across + /// multiple sub-operations, so they fall back to fragment-level conflict + /// detection. pub fn with_affected_rows(mut self, affected_rows: RowAddrTreeMap) -> Self { self.affected_rows = Some(affected_rows); self @@ -304,15 +308,10 @@ impl<'a> CommitBuilder<'a> { } }; - let can_create_dataset = match &transaction.operation { - Operation::Overwrite { .. } | Operation::Clone { .. } => true, - // A composite can create a dataset by starting with Overwrite; - // validate_operation enforces this below. - Operation::Composite { operations } => { - matches!(operations.first(), Some(Operation::Overwrite { .. })) - } - _ => false, - }; + // A composite can create a dataset by starting with Overwrite; + // validate_operation enforces this below. + let can_create_dataset = transaction.operation.overwrites_dataset() + || matches!(transaction.operation, Operation::Clone { .. }); if dest.dataset().is_none() && !can_create_dataset { return Err(Error::dataset_not_found( base_path.to_string(), @@ -355,7 +354,7 @@ impl<'a> CommitBuilder<'a> { { let passed_storage_format = DataStorageFormat::new(storage_format); if ds.manifest.data_storage_format != passed_storage_format - && !matches!(transaction.operation, Operation::Overwrite { .. }) + && !transaction.operation.overwrites_dataset() { return Err(Error::invalid_input_source(format!( "Storage format mismatch. Existing dataset uses {:?}, but new data uses {:?}", diff --git a/rust/lance/src/io/commit/conflict_resolver.rs b/rust/lance/src/io/commit/conflict_resolver.rs index facb4505789..ed2f233bf26 100644 --- a/rust/lance/src/io/commit/conflict_resolver.rs +++ b/rust/lance/src/io/commit/conflict_resolver.rs @@ -51,6 +51,12 @@ impl<'a> TransactionRebase<'a> { // A single affected-rows map is ambiguous across multiple // Delete/Update sub-operations, so sub-rebases fall back to // fragment-level conflict detection. + if affected_rows.is_some() { + log::warn!( + "affected_rows is ignored for composite transactions; sub-operations \ + use fragment-level conflict detection" + ); + } let mut sub_rebases = Vec::with_capacity(operations.len()); for operation in operations { let sub_transaction = Transaction { @@ -1493,11 +1499,17 @@ impl<'a> TransactionRebase<'a> { pub async fn finish(self, dataset: &Dataset) -> Result { if matches!(self.transaction.operation, Operation::Composite { .. }) { let mut operations = Vec::with_capacity(self.sub_rebases.len()); + let mut read_version = self.transaction.read_version; for sub_rebase in self.sub_rebases { let finished = sub_rebase.finish_single(dataset).await?; + // finish_delete_update advances the read version of a sub- + // transaction whose deletion files were rebased against the + // current dataset; the outer transaction must reflect that. + read_version = read_version.max(finished.read_version); operations.push(finished.operation); } let mut transaction = self.transaction; + transaction.read_version = read_version; transaction.operation = Operation::Composite { operations }; return Ok(transaction); } From b307de934df3da4bdb4303fbf37b9e6893a199bd Mon Sep 17 00:00:00 2001 From: Jack Ye Date: Sat, 11 Jul 2026 08:23:29 -0700 Subject: [PATCH 3/3] feat: inline composite transactions like any other operation Composites now write the inline manifest transaction section normally. The inability of released readers to decode an unknown inline operation is an independent robustness problem: the opportunistic decode on dataset open now tolerates undecodable transactions (log and skip) instead of failing the open, which protects against any future operation type. --- docs/src/format/table/transaction.md | 4 - rust/lance/src/dataset.rs | 77 ++++++++----------- .../src/dataset/tests/dataset_transactions.rs | 29 +++++++ rust/lance/src/dataset/write/commit.rs | 7 +- 4 files changed, 66 insertions(+), 51 deletions(-) diff --git a/docs/src/format/table/transaction.md b/docs/src/format/table/transaction.md index 8fc52d23f53..2b5c2055a60 100644 --- a/docs/src/format/table/transaction.md +++ b/docs/src/format/table/transaction.md @@ -528,10 +528,6 @@ concurrent transaction is itself a Composite, it is decomposed the same way, so operations conflict if any pair of their sub-operations conflicts. Rebases are likewise delegated: each sub-operation is rebased under its own rules and the results are reassembled in order. -Writers do not inline Composite transactions into the manifest's transaction section: released -readers eagerly decode that section when opening a dataset and would fail on an unknown operation. -The transaction remains available in the external transaction file. - ## Conflict Resolution ### Terminology diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index 8a65202e652..c98b580467b 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -738,31 +738,16 @@ impl Dataset { let message_len = LittleEndian::read_u32(&last_block[offset_in_block..offset_in_block + 4]) as usize; let message_data = &last_block[offset_in_block + 4..offset_in_block + 4 + message_len]; - // This is an opportunistic cache warm-up: a transaction written by - // a newer version may contain an operation this version cannot - // decode, and that must not prevent opening the dataset. Paths - // that need the transaction contents surface the error themselves. - match lance_table::format::pb::Transaction::decode(message_data) - .map_err(Error::from) - .and_then(Transaction::try_from) + if let Some(transaction) = + decode_inline_transaction(message_data, manifest_location.version) { - Ok(transaction) => { - let metadata_cache = session.metadata_cache.for_dataset(uri); - let metadata_key = TransactionKey { - version: manifest_location.version, - }; - metadata_cache - .insert_with_key(&metadata_key, Arc::new(transaction)) - .await; - } - Err(err) => { - log::warn!( - "Failed to decode the inline transaction of version {}; \ - it may have been written by a newer version of Lance: {}", - manifest_location.version, - err - ); - } + let metadata_cache = session.metadata_cache.for_dataset(uri); + let metadata_key = TransactionKey { + version: manifest_location.version, + }; + metadata_cache + .insert_with_key(&metadata_key, Arc::new(transaction)) + .await; } } @@ -3682,6 +3667,31 @@ impl ManifestWriteConfig { } } +/// Decode an inline transaction section for opportunistic caching. +/// +/// Returns `None` instead of failing when the transaction cannot be decoded: +/// the section may have been written by a newer version of Lance with an +/// operation type this version does not know, and that must not prevent +/// opening the dataset. Paths that need the transaction contents surface the +/// error at their call sites instead. +fn decode_inline_transaction(message_data: &[u8], version: u64) -> Option { + match lance_table::format::pb::Transaction::decode(message_data) + .map_err(Error::from) + .and_then(Transaction::try_from) + { + Ok(transaction) => Some(transaction), + Err(err) => { + log::warn!( + "Failed to decode the inline transaction of version {}; \ + it may have been written by a newer version of Lance: {}", + version, + err + ); + None + } + } +} + /// Commit a manifest file and create a copy at the latest manifest path. #[allow(clippy::too_many_arguments)] pub(crate) async fn write_manifest_file( @@ -3694,25 +3704,6 @@ pub(crate) async fn write_manifest_file( naming_scheme: ManifestNamingScheme, mut transaction: Option<&Transaction>, ) -> std::result::Result { - // Composite transactions are not inlined into the manifest: released - // readers eagerly decode the inline transaction section when opening a - // dataset and fail on unknown operation types, which would make this - // version unopenable for them. Without the inline section they open - // normally; only paths that need the transaction contents (conflict - // resolution, history) read the external transaction file. When the - // external file is disabled the inline section is the only transaction - // record, so it must be kept — dropping both would break conflict - // resolution for every later writer. - // TODO: inline composites once readers with tolerant decoding are - // widespread. - if !config.disable_transaction_file - && matches!( - transaction.map(|t| &t.operation), - Some(Operation::Composite { .. }) - ) - { - transaction = None; - } if config.auto_set_feature_flags { // build_manifest may have already set FLAG_STABLE_ROW_IDS on the manifest. // Preserve it here so this second apply_feature_flags call does not clear it diff --git a/rust/lance/src/dataset/tests/dataset_transactions.rs b/rust/lance/src/dataset/tests/dataset_transactions.rs index 3e2a4caa3b3..548c6985cb3 100644 --- a/rust/lance/src/dataset/tests/dataset_transactions.rs +++ b/rust/lance/src/dataset/tests/dataset_transactions.rs @@ -208,6 +208,35 @@ async fn test_session_store_registry() { assert_eq!(registry.active_stores().len(), 0); } +#[test] +fn test_decode_inline_transaction_tolerates_unknown_operations() { + use crate::dataset::decode_inline_transaction; + use lance_table::format::pb; + use prost::Message; + + // A transaction written by a newer version of Lance may carry an operation + // this version cannot decode; prost surfaces it as a missing oneof. This + // must not fail (it would prevent opening the dataset), only skip caching. + let unknown_operation = pb::Transaction { + read_version: 1, + uuid: "test".to_string(), + ..Default::default() + }; + assert!(decode_inline_transaction(&unknown_operation.encode_to_vec(), 42).is_none()); + + // Corrupt bytes are likewise tolerated. + assert!(decode_inline_transaction(&[0xff, 0xff, 0xff], 42).is_none()); + + // A decodable transaction is returned. + let known = pb::Transaction::from(&Transaction::new( + 1, + Operation::Append { fragments: vec![] }, + None, + )); + let decoded = decode_inline_transaction(&known.encode_to_vec(), 42).unwrap(); + assert!(matches!(decoded.operation, Operation::Append { .. })); +} + #[tokio::test] async fn test_migrate_v2_manifest_paths() { let test_uri = TempStrDir::default(); diff --git a/rust/lance/src/dataset/write/commit.rs b/rust/lance/src/dataset/write/commit.rs index 1ac09641e8e..fd30d17c716 100644 --- a/rust/lance/src/dataset/write/commit.rs +++ b/rust/lance/src/dataset/write/commit.rs @@ -1081,10 +1081,9 @@ mod tests { Some(&"value".to_string()) ); - // The transaction is not inlined into the manifest so that released - // readers, which decode the inline section eagerly, can still open - // this version. It is still readable from the transaction file. - assert!(new_ds.manifest.transaction_section.is_none()); + // Composites are inlined into the manifest like any other operation + // and read back as a composite. + assert!(new_ds.manifest.transaction_section.is_some()); let read_back = new_ds.read_transaction().await.unwrap().unwrap(); assert!(matches!( read_back.operation,