From e6f30f76f60c747eace13fc8edc15f34e8abfa75 Mon Sep 17 00:00:00 2001 From: Will Jones Date: Tue, 18 Aug 2026 20:40:55 -0700 Subject: [PATCH 01/10] feat(transaction): add the AddOverlays action Appends overlay files to a fragment, supplying new values for a subset of its (row offset, field) cells without rewriting its base data files. Each overlay's `committed_version` is stamped with the version the commit produces, so a retry against a newer manifest re-stamps rather than backdates. Overlays are appended, never replaced, so the action writes no coordinate of its own -- two concurrent overlays over the same cells both land and the newer version wins. It does record that the fragment must still be there, which is a new kind of entry in the footprint: a dependency rather than a write. --- rust/lance-table/src/transaction/action.rs | 3 + .../src/transaction/action/add_overlays.rs | 242 ++++++++++++++++++ .../src/transaction/action/apply.rs | 8 + .../src/transaction/action/footprint.rs | 9 +- .../src/transaction/action/proto.rs | 20 +- 5 files changed, 276 insertions(+), 6 deletions(-) create mode 100644 rust/lance-table/src/transaction/action/add_overlays.rs diff --git a/rust/lance-table/src/transaction/action.rs b/rust/lance-table/src/transaction/action.rs index 95b1136a44e..835d59df7d0 100644 --- a/rust/lance-table/src/transaction/action.rs +++ b/rust/lance-table/src/transaction/action.rs @@ -66,6 +66,7 @@ macro_rules! for_each_action { TombstoneFieldData, RemoveFragment, SetDeletionFile, + AddOverlays, AlterField, DropField, AddIndexSegment, @@ -84,6 +85,7 @@ mod add_data_file; mod add_field; mod add_fragment; mod add_index_segment; +mod add_overlays; mod adjust_index_coverage; mod alter_field; mod apply; @@ -107,6 +109,7 @@ pub use add_data_file::AddDataFile; pub use add_field::AddField; pub use add_fragment::AddFragment; pub use add_index_segment::AddIndexSegment; +pub use add_overlays::AddOverlays; pub use adjust_index_coverage::AdjustIndexCoverage; pub use alter_field::AlterField; pub use config_update::{ConfigUpdate, FieldMetadataUpdate}; diff --git a/rust/lance-table/src/transaction/action/add_overlays.rs b/rust/lance-table/src/transaction/action/add_overlays.rs new file mode 100644 index 00000000000..930e39cf85a --- /dev/null +++ b/rust/lance-table/src/transaction/action/add_overlays.rs @@ -0,0 +1,242 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Append overlay files to a fragment. + +use super::apply::ApplyState; +use super::proto::{data_change_from_wire, data_change_to_wire, required}; +use super::{Footprint, Ref}; +use crate::format::overlay::DataOverlayFile; +use crate::format::pb; +use lance_core::deepsize::DeepSizeOf; +use lance_core::{Error, Result}; + +/// Append overlay files to one fragment, supplying new values for a subset of +/// its `(row offset, field)` cells without rewriting its base data files. +/// +/// Overlays are appended rather than replaced, so overlays a concurrent writer +/// added survive. Within the fragment they are ordered newest-last, which is +/// what appending at the version this commit produces gives. +#[derive(Debug, Clone, PartialEq, DeepSizeOf)] +pub struct AddOverlays { + pub fragment: Ref, + /// The overlays to append, oldest first. Each one's `committed_version` is + /// ignored and stamped with the version this commit produces, so replaying + /// the action onto a newer version re-stamps rather than backdates. + pub overlays: Vec, + pub data_change: bool, +} + +impl AddOverlays { + pub(super) fn apply(&self, state: &mut ApplyState) -> Result<()> { + let fragment_id = state.resolve_fragment(self.fragment)?; + let committed_version = state.new_version(); + let fragment = state.fragment_mut(fragment_id, "AddOverlays")?; + fragment + .overlays + .extend(self.overlays.iter().cloned().map(|mut overlay| { + overlay.committed_version = committed_version; + overlay + })); + Ok(()) + } + + /// An overlay supplies new cell values, so by default it changes what a + /// reader sees. The writer can still mark a restatement of existing values + /// as no change. + pub(super) fn is_data_change(&self) -> bool { + self.data_change + } + + /// Only that the fragment must still be there. An overlay writes no + /// coordinate of its own: two concurrent overlays over the same cells both + /// land, and the newer `committed_version` decides which value wins. + pub(super) fn footprint(&self, footprint: &mut Footprint) { + if let Some(fragment) = self.fragment.committed() { + footprint.require_fragment(fragment); + } + } +} + +impl From<&AddOverlays> for pb::AddOverlays { + fn from(value: &AddOverlays) -> Self { + Self { + fragment: Some(value.fragment.into()), + overlays: value + .overlays + .iter() + .map(pb::DataOverlayFile::from) + .collect(), + data_change: data_change_to_wire(value.data_change), + } + } +} + +impl TryFrom for AddOverlays { + type Error = Error; + + fn try_from(message: pb::AddOverlays) -> Result { + Ok(Self { + fragment: required(message.fragment, "AddOverlays.fragment")?.try_into()?, + overlays: message + .overlays + .into_iter() + .map(DataOverlayFile::try_from) + .collect::>>()?, + data_change: data_change_from_wire(message.data_change), + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::format::DataFile; + use crate::format::overlay::OverlayCoverage; + use crate::transaction::action::test_support::{apply, backed_manifest}; + use crate::transaction::action::{ + Action, AddFragment, CompositeOperation, RemoveFragment, UserAction, + }; + use roaring::RoaringBitmap; + use std::sync::Arc; + + fn overlay(path: &str, offsets: &[u32]) -> DataOverlayFile { + DataOverlayFile { + data_file: DataFile::new(path, vec![0], vec![0], 2, 0, None, None), + coverage: OverlayCoverage::Shared(Arc::new( + offsets.iter().copied().collect::(), + )), + // Whatever the writer left here is overwritten at apply. + committed_version: 0, + } + } + + fn add(fragment: Ref, overlays: Vec) -> Action { + Action::AddOverlays(AddOverlays { + fragment, + overlays, + data_change: true, + }) + } + + fn footprint(actions: Vec) -> Footprint { + Footprint::from(&CompositeOperation::new(vec![UserAction::new( + "step", actions, + )])) + } + + #[test] + fn test_overlays_are_stamped_with_the_version_the_commit_produces() { + let manifest = backed_manifest(); + let expected = manifest.version + 1; + + let out = apply( + &manifest, + vec![add(Ref::Committed(0), vec![overlay("a.lance", &[1, 2])])], + ) + .unwrap(); + + let overlays = &out.fragments[0].overlays; + assert_eq!(overlays.len(), 1); + assert_eq!(overlays[0].committed_version, expected); + } + + #[test] + fn test_overlays_are_appended_to_the_ones_already_there() { + let mut manifest = backed_manifest(); + let mut fragment = manifest.fragments[0].clone(); + fragment.overlays.push(overlay("old.lance", &[0])); + manifest.fragments = Arc::new(vec![fragment]); + + let out = apply( + &manifest, + vec![add(Ref::Committed(0), vec![overlay("new.lance", &[1])])], + ) + .unwrap(); + + let paths = out.fragments[0] + .overlays + .iter() + .map(|overlay| overlay.data_file.path.as_str()) + .collect::>(); + assert_eq!(paths, vec!["old.lance", "new.lance"]); + } + + #[test] + fn test_several_overlays_keep_the_order_they_were_given_in() { + let out = apply( + &backed_manifest(), + vec![add( + Ref::Committed(0), + vec![overlay("first.lance", &[0]), overlay("second.lance", &[0])], + )], + ) + .unwrap(); + + let paths = out.fragments[0] + .overlays + .iter() + .map(|overlay| overlay.data_file.path.as_str()) + .collect::>(); + assert_eq!(paths, vec!["first.lance", "second.lance"]); + } + + #[test] + fn test_an_overlay_can_target_a_fragment_this_operation_minted() { + let out = apply( + &backed_manifest(), + vec![ + Action::AddFragment(AddFragment { + id: Ref::Local(0), + physical_rows: 4, + row_id_meta: None, + last_updated_at_version_meta: None, + created_at_version_meta: None, + data_change: true, + }), + add(Ref::Local(0), vec![overlay("a.lance", &[0])]), + ], + ) + .unwrap(); + + let minted = out.fragments.iter().find(|f| f.id == 1).unwrap(); + assert_eq!(minted.overlays.len(), 1); + } + + #[test] + fn test_overlaying_a_fragment_that_is_not_there_is_rejected() { + let error = apply( + &backed_manifest(), + vec![add(Ref::Committed(42), vec![overlay("a.lance", &[0])])], + ) + .unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error + .to_string() + .contains("AddOverlays targets fragment 42, which does not exist"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_two_writers_overlaying_the_same_fragment_do_not_conflict() { + let ours = footprint(vec![add(Ref::Committed(0), vec![overlay("a.lance", &[0])])]); + let theirs = footprint(vec![add(Ref::Committed(0), vec![overlay("b.lance", &[0])])]); + + assert!(!ours.conflicts_with(&theirs)); + } + + #[test] + fn test_overlaying_a_fragment_a_concurrent_writer_removes_conflicts() { + let ours = footprint(vec![add(Ref::Committed(0), vec![overlay("a.lance", &[0])])]); + let theirs = footprint(vec![Action::RemoveFragment(RemoveFragment { + fragment: Ref::Committed(0), + data_change: true, + })]); + + assert!(ours.conflicts_with(&theirs)); + assert!(theirs.conflicts_with(&ours)); + } +} diff --git a/rust/lance-table/src/transaction/action/apply.rs b/rust/lance-table/src/transaction/action/apply.rs index cdbd0cd477c..be771cf797b 100644 --- a/rust/lance-table/src/transaction/action/apply.rs +++ b/rust/lance-table/src/transaction/action/apply.rs @@ -280,6 +280,14 @@ impl<'a> ApplyState<'a> { self.read_version } + /// The version this delta produces, for the actions that stamp it into what + /// they write. Distinct from [`Self::read_version`]: a retry against a newer + /// manifest re-runs the apply, so anything stamped with this is re-stamped + /// rather than carried over. + pub(super) fn new_version(&self) -> u64 { + self.current_manifest.version + 1 + } + /// Add an index segment. A segment's uuid identifies it, so re-adding one /// that is already there is a mistake rather than a replacement -- swapping /// a segment out is a [`RemoveIndexSegment`](super::RemoveIndexSegment) diff --git a/rust/lance-table/src/transaction/action/footprint.rs b/rust/lance-table/src/transaction/action/footprint.rs index 1f7e41bd6a3..1ed32ac62ff 100644 --- a/rust/lance-table/src/transaction/action/footprint.rs +++ b/rust/lance-table/src/transaction/action/footprint.rs @@ -194,9 +194,12 @@ pub struct Footprint { regions: HashMap, /// Fragments this set needs to still be there, without writing anything a /// concurrent set could name inside them. Data for a field this set mints, - /// written into a committed fragment, is the case this exists for: the - /// field id is invisible to a concurrent writer, so the cells are not a - /// coordinate, but they are gone if the fragment is. + /// written into a committed fragment, is one case: the field id is + /// invisible to a concurrent writer, so the cells are not a coordinate, but + /// they are gone if the fragment is. An overlay is the other: two + /// concurrent overlays over the same cells both land and the newer one + /// wins, so they must not collide with each other -- but neither survives a + /// concurrent writer dropping the fragment out from under them. /// /// Not a `Requires` claim on the region, because the check is symmetric: a /// removal that lands second destroys the cells just as surely as one that diff --git a/rust/lance-table/src/transaction/action/proto.rs b/rust/lance-table/src/transaction/action/proto.rs index 7a5506e3268..4481e8fdc38 100644 --- a/rust/lance-table/src/transaction/action/proto.rs +++ b/rust/lance-table/src/transaction/action/proto.rs @@ -153,21 +153,24 @@ for_each_action!(define_action_proto); #[cfg(test)] mod tests { use super::*; + use crate::format::overlay::{DataOverlayFile, OverlayCoverage}; use crate::format::{ BasePath, DataFile, DeletionFile, DeletionFileType, IndexFile, RowIdMeta, pb, }; use crate::rowids::version::RowDatasetVersionMeta; use crate::transaction::UpdateMap; use crate::transaction::action::{ - AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, AdjustIndexCoverage, - AlterField, ConfigUpdate, DropField, FieldMetadataUpdate, RemoveFragment, - RemoveIndexSegment, ReserveFragmentIds, ReserveRowIds, ResetTable, SetDeletionFile, + AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, AddOverlays, + AdjustIndexCoverage, AlterField, ConfigUpdate, DropField, FieldMetadataUpdate, + RemoveFragment, RemoveIndexSegment, ReserveFragmentIds, ReserveRowIds, ResetTable, + SetDeletionFile, TombstoneFieldData, }; use arrow_schema::{DataType, Field as ArrowField}; use chrono::DateTime; use lance_core::datatypes::Field; use lance_file::version::ConcreteFileVersion; + use roaring::RoaringBitmap; use rstest::rstest; use std::sync::Arc; use uuid::Uuid; @@ -223,6 +226,17 @@ mod tests { }), data_change: true, }), + Action::AddOverlays(AddOverlays { + fragment: Ref::Committed(6), + overlays: vec![DataOverlayFile { + data_file: sample_data_file(), + coverage: OverlayCoverage::PerField(vec![Arc::new( + [1u32, 4].into_iter().collect::(), + )]), + committed_version: 11, + }], + data_change: true, + }), Action::AlterField(AlterField { field: Ref::Committed(2), name: Some("renamed".into()), From d1d8cb614226232dae80f736de56376120a9c562 Mon Sep 17 00:00:00 2001 From: Will Jones Date: Tue, 18 Aug 2026 20:44:09 -0700 Subject: [PATCH 02/10] feat(transaction): add the RefreshRowVersionMetadata action Restamps the per-row `last_updated_at_version` sequence of fragments whose columns were rewritten in place, which is what a legacy Merge does implicitly. Nothing else in an operation restates when those rows last changed, because rewriting columns in place leaves the rows where they are. `created_at_version` is left alone: the rows are the same rows, and a row this operation mints gets both stamps from the AddFragment that minted it. Naming a fragment on a dataset without stable row ids is rejected rather than fabricating sequences that have nowhere to live. --- rust/lance-table/src/transaction/action.rs | 3 + .../src/transaction/action/apply.rs | 6 + .../src/transaction/action/footprint.rs | 6 +- .../src/transaction/action/proto.rs | 17 +- .../action/refresh_row_version_metadata.rs | 198 ++++++++++++++++++ rust/lance-table/src/transaction/proto.rs | 9 +- 6 files changed, 226 insertions(+), 13 deletions(-) create mode 100644 rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs diff --git a/rust/lance-table/src/transaction/action.rs b/rust/lance-table/src/transaction/action.rs index 835d59df7d0..7ebad35eaaf 100644 --- a/rust/lance-table/src/transaction/action.rs +++ b/rust/lance-table/src/transaction/action.rs @@ -67,6 +67,7 @@ macro_rules! for_each_action { RemoveFragment, SetDeletionFile, AddOverlays, + RefreshRowVersionMetadata, AlterField, DropField, AddIndexSegment, @@ -93,6 +94,7 @@ mod config_update; mod drop_field; mod footprint; mod proto; +mod refresh_row_version_metadata; mod remove_fragment; mod remove_index_segment; mod reserve_fragment_ids; @@ -115,6 +117,7 @@ pub use alter_field::AlterField; pub use config_update::{ConfigUpdate, FieldMetadataUpdate}; pub use drop_field::DropField; pub use footprint::{ConfigMap, Coordinate, Footprint}; +pub use refresh_row_version_metadata::RefreshRowVersionMetadata; pub use remove_fragment::RemoveFragment; pub use remove_index_segment::RemoveIndexSegment; pub use reserve_fragment_ids::ReserveFragmentIds; diff --git a/rust/lance-table/src/transaction/action/apply.rs b/rust/lance-table/src/transaction/action/apply.rs index be771cf797b..6863a289841 100644 --- a/rust/lance-table/src/transaction/action/apply.rs +++ b/rust/lance-table/src/transaction/action/apply.rs @@ -346,6 +346,12 @@ impl<'a> ApplyState<'a> { Ok(&mut self.indices[position]) } + /// Whether the dataset carries stable row ids, and with them the per-row + /// version sequences. + pub(super) fn uses_stable_row_ids(&self) -> bool { + self.current_manifest.uses_stable_row_ids() + } + pub(super) fn schema(&self) -> &Schema { &self.schema } diff --git a/rust/lance-table/src/transaction/action/footprint.rs b/rust/lance-table/src/transaction/action/footprint.rs index 1ed32ac62ff..4470f55ed87 100644 --- a/rust/lance-table/src/transaction/action/footprint.rs +++ b/rust/lance-table/src/transaction/action/footprint.rs @@ -33,6 +33,8 @@ pub enum Coordinate { FragmentExistence(u64), /// A committed fragment's deletion file. FragmentDeletions(u64), + /// A committed fragment's per-row version sequences. + FragmentRowVersions(u64), /// The data backing one field within one committed fragment. FieldData { fragment: u64, field: i32 }, /// A field's definition in the schema. @@ -85,7 +87,9 @@ impl Coordinate { /// The fragment this coordinate lives in, if it is fragment-scoped. fn fragment(&self) -> Option { match self { - Self::FragmentExistence(id) | Self::FragmentDeletions(id) => Some(*id), + Self::FragmentExistence(id) + | Self::FragmentDeletions(id) + | Self::FragmentRowVersions(id) => Some(*id), Self::FieldData { fragment, .. } => Some(*fragment), Self::FieldDefinition(_) | Self::FieldName(_) diff --git a/rust/lance-table/src/transaction/action/proto.rs b/rust/lance-table/src/transaction/action/proto.rs index 4481e8fdc38..59fb8aab5ff 100644 --- a/rust/lance-table/src/transaction/action/proto.rs +++ b/rust/lance-table/src/transaction/action/proto.rs @@ -162,9 +162,8 @@ mod tests { use crate::transaction::action::{ AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, AddOverlays, AdjustIndexCoverage, AlterField, ConfigUpdate, DropField, FieldMetadataUpdate, - RemoveFragment, RemoveIndexSegment, ReserveFragmentIds, ReserveRowIds, ResetTable, - SetDeletionFile, - TombstoneFieldData, + RefreshRowVersionMetadata, RemoveFragment, RemoveIndexSegment, ReserveFragmentIds, + ReserveRowIds, ResetTable, SetDeletionFile, TombstoneFieldData, }; use arrow_schema::{DataType, Field as ArrowField}; use chrono::DateTime; @@ -237,6 +236,9 @@ mod tests { }], data_change: true, }), + Action::RefreshRowVersionMetadata(RefreshRowVersionMetadata { + fragment_ids: vec![4, 6], + }), Action::AlterField(AlterField { field: Ref::Committed(2), name: Some("renamed".into()), @@ -329,11 +331,10 @@ mod tests { #[test] fn test_unimplemented_action_is_rejected() { let message = pb::Action { - action: Some(pb::action::Action::RefreshRowVersionMetadata( - pb::RefreshRowVersionMetadata { - fragment_ids: vec![1], - }, - )), + action: Some(pb::action::Action::AssertUniqueKeys(pb::AssertUniqueKeys { + key_fields: vec![], + filter: None, + })), }; let error = Action::try_from(message).unwrap_err(); assert!( diff --git a/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs b/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs new file mode 100644 index 00000000000..96fef3914fa --- /dev/null +++ b/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs @@ -0,0 +1,198 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Restamp the row-version metadata of fragments rewritten in place. + +use super::apply::ApplyState; +use super::{Coordinate, Footprint}; +use crate::format::pb; +use crate::rowids::version::refresh_row_latest_update_meta_for_full_frag_rewrite_cols; +use lance_core::deepsize::DeepSizeOf; +use lance_core::{Error, Result}; + +/// Record that every row of these fragments was updated by this operation. +/// +/// Under stable row ids a fragment carries a `last_updated_at_version` sequence +/// per row. Rewriting a fragment's columns in place leaves the rows where they +/// are, so nothing else in the operation restates when they last changed; this +/// action does, mirroring the refresh a legacy `Merge` performs implicitly. +/// +/// `created_at_version` is deliberately untouched: the rows are the same rows, +/// and a row minted by this operation gets both stamps from the +/// [`AddFragment`](super::AddFragment) that minted it. +/// +/// Fragments are named by committed id. A fragment this operation minted has +/// nothing to restamp. +#[derive(Debug, Clone, PartialEq, DeepSizeOf)] +pub struct RefreshRowVersionMetadata { + pub fragment_ids: Vec, +} + +impl RefreshRowVersionMetadata { + pub(super) fn apply(&self, state: &mut ApplyState) -> Result<()> { + if !self.fragment_ids.is_empty() && !state.uses_stable_row_ids() { + return Err(Error::invalid_input( + "RefreshRowVersionMetadata restamps the per-row version sequences, which only \ + exist on a dataset using stable row ids", + )); + } + let new_version = state.new_version(); + for fragment_id in &self.fragment_ids { + let fragment = state.fragment_mut(*fragment_id, "RefreshRowVersionMetadata")?; + refresh_row_latest_update_meta_for_full_frag_rewrite_cols(fragment, new_version)?; + } + Ok(()) + } + + /// The rows themselves are restated elsewhere in the operation -- by the + /// data files it writes -- so this is bookkeeping about that change, not a + /// change of its own. + pub(super) fn is_data_change(&self) -> bool { + false + } + + pub(super) fn footprint(&self, footprint: &mut Footprint) { + for fragment_id in &self.fragment_ids { + footprint.add(Coordinate::FragmentRowVersions(*fragment_id)); + } + } +} + +impl From<&RefreshRowVersionMetadata> for pb::RefreshRowVersionMetadata { + fn from(value: &RefreshRowVersionMetadata) -> Self { + Self { + fragment_ids: value.fragment_ids.clone(), + } + } +} + +impl TryFrom for RefreshRowVersionMetadata { + type Error = Error; + + fn try_from(message: pb::RefreshRowVersionMetadata) -> Result { + Ok(Self { + fragment_ids: message.fragment_ids, + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::format::{DataFile, Fragment, RowIdMeta}; + use crate::rowids::{RowIdSequence, write_row_ids}; + use crate::transaction::action::test_support::{apply, backed_manifest}; + use crate::transaction::action::{Action, CompositeOperation, UserAction}; + use crate::transaction::test_support::make_stable_row_id_manifest; + use std::sync::Arc; + + fn refresh(fragment_ids: Vec) -> Action { + Action::RefreshRowVersionMetadata(RefreshRowVersionMetadata { fragment_ids }) + } + + /// A stable-row-id manifest whose fragment 0 has `rows` rows, all last + /// updated at version 1. + fn manifest_with_rows(rows: usize) -> crate::format::Manifest { + let row_ids = RowIdSequence::from((0..rows as u64).collect::>().as_slice()); + let mut fragment = Fragment { + id: 0, + files: vec![DataFile::new( + "data.lance", + vec![0], + vec![0], + 2, + 0, + None, + None, + )], + overlays: vec![], + deletion_file: None, + row_id_meta: Some(RowIdMeta::Inline(write_row_ids(&row_ids))), + physical_rows: Some(rows), + last_updated_at_version_meta: None, + created_at_version_meta: None, + }; + refresh_row_latest_update_meta_for_full_frag_rewrite_cols(&mut fragment, 1).unwrap(); + make_stable_row_id_manifest(vec![fragment]) + } + + fn last_updated_versions(fragment: &Fragment) -> Vec { + let meta = fragment + .last_updated_at_version_meta + .as_ref() + .expect("the fragment should carry a last-updated sequence"); + let sequence = meta.load_sequence().unwrap(); + (0..fragment.physical_rows.unwrap()) + .map(|offset| sequence.version_at(offset).unwrap()) + .collect() + } + + #[test] + fn test_refresh_stamps_every_row_with_the_version_the_commit_produces() { + let manifest = manifest_with_rows(3); + let expected = manifest.version + 1; + + let out = apply(&manifest, vec![refresh(vec![0])]).unwrap(); + + assert_eq!(last_updated_versions(&out.fragments[0]), vec![expected; 3]); + } + + #[test] + fn test_refresh_leaves_the_created_at_sequence_alone() { + let mut manifest = manifest_with_rows(3); + let mut fragment = manifest.fragments[0].clone(); + fragment.created_at_version_meta = fragment.last_updated_at_version_meta.clone(); + let created_at = fragment.created_at_version_meta.clone(); + manifest.fragments = Arc::new(vec![fragment]); + + let out = apply(&manifest, vec![refresh(vec![0])]).unwrap(); + + assert_eq!(out.fragments[0].created_at_version_meta, created_at); + } + + #[test] + fn test_refreshing_no_fragments_is_allowed_on_any_dataset() { + // An empty list is what a writer emits when nothing was rewritten in + // place, so it must not depend on stable row ids. + apply(&backed_manifest(), vec![refresh(vec![])]).unwrap(); + } + + #[test] + fn test_refreshing_without_stable_row_ids_is_rejected() { + let error = apply(&backed_manifest(), vec![refresh(vec![0])]).unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains("stable row ids"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_refreshing_a_fragment_that_is_not_there_is_rejected() { + let error = apply(&manifest_with_rows(3), vec![refresh(vec![7])]).unwrap_err(); + + assert!( + error + .to_string() + .contains("RefreshRowVersionMetadata targets fragment 7, which does not exist"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_two_writers_refreshing_the_same_fragment_conflict() { + let footprint = |actions| { + Footprint::from(&CompositeOperation::new(vec![UserAction::new( + "step", actions, + )])) + }; + + let ours = footprint(vec![refresh(vec![0])]); + let same = footprint(vec![refresh(vec![0, 1])]); + let other = footprint(vec![refresh(vec![1])]); + + assert!(ours.conflicts_with(&same)); + assert!(!ours.conflicts_with(&other)); + } +} diff --git a/rust/lance-table/src/transaction/proto.rs b/rust/lance-table/src/transaction/proto.rs index 0ce47f74334..139ed30ba0b 100644 --- a/rust/lance-table/src/transaction/proto.rs +++ b/rust/lance-table/src/transaction/proto.rs @@ -919,11 +919,12 @@ mod tests { operation: Some(pb::transaction::Operation::CompositeOperation( pb::CompositeOperation { actions: vec![pb::UserAction { - description: "refresh row versions".to_string(), + description: "assert unique keys".to_string(), actions: vec![pb::Action { - action: Some(pb::action::Action::RefreshRowVersionMetadata( - pb::RefreshRowVersionMetadata { - fragment_ids: vec![1], + action: Some(pb::action::Action::AssertUniqueKeys( + pb::AssertUniqueKeys { + key_fields: vec![], + filter: None, }, )), }], From 3bf9bfeb049bb1f231d9e45480ecb6eea447f83a Mon Sep 17 00:00:00 2001 From: Will Jones Date: Tue, 18 Aug 2026 20:46:30 -0700 Subject: [PATCH 03/10] feat(transaction): add the UpdateCompactedSsTables action Records which MemWAL SSTables have been compacted into the base table, in the MemWAL system index. Per shard the highest generation wins, so replaying an older commit over a newer one cannot walk the progress backwards. The rows were already readable through the WAL, so this is bookkeeping about where they live rather than a change to them. The drafted `update_compacted_sstables` oneof field is renamed to `update_compacted_ss_tables` so the generated variant name matches the message name, which is what the action vocabulary keys the wire encoding off. The tag is unchanged. --- rust/lance-table/src/transaction/action.rs | 3 + .../src/transaction/action/apply.rs | 15 +- .../src/transaction/action/footprint.rs | 10 + .../src/transaction/action/proto.rs | 9 +- .../action/refresh_row_version_metadata.rs | 2 +- .../action/update_compacted_sstables.rs | 227 ++++++++++++++++++ 6 files changed, 263 insertions(+), 3 deletions(-) create mode 100644 rust/lance-table/src/transaction/action/update_compacted_sstables.rs diff --git a/rust/lance-table/src/transaction/action.rs b/rust/lance-table/src/transaction/action.rs index 7ebad35eaaf..6477f1e720b 100644 --- a/rust/lance-table/src/transaction/action.rs +++ b/rust/lance-table/src/transaction/action.rs @@ -68,6 +68,7 @@ macro_rules! for_each_action { SetDeletionFile, AddOverlays, RefreshRowVersionMetadata, + UpdateCompactedSsTables, AlterField, DropField, AddIndexSegment, @@ -102,6 +103,7 @@ mod reserve_row_ids; mod reset_table; mod set_deletion_file; mod tombstone_field_data; +mod update_compacted_sstables; #[cfg(test)] mod test_support; @@ -125,6 +127,7 @@ pub use reserve_row_ids::ReserveRowIds; pub use reset_table::ResetTable; pub use set_deletion_file::SetDeletionFile; pub use tombstone_field_data::TombstoneFieldData; +pub use update_compacted_sstables::UpdateCompactedSsTables; use apply::ApplyState; use lance_core::Result; diff --git a/rust/lance-table/src/transaction/action/apply.rs b/rust/lance-table/src/transaction/action/apply.rs index 6863a289841..0dc846f2d48 100644 --- a/rust/lance-table/src/transaction/action/apply.rs +++ b/rust/lance-table/src/transaction/action/apply.rs @@ -22,7 +22,9 @@ use super::{CompositeOperation, Ref}; use crate::format::{BasePath, Fragment, IndexMetadata, Manifest, ManifestBuildConfig, RowIdMeta}; use crate::rowids::read_row_ids; use crate::rowids::version::build_version_meta; -use crate::system_index::mem_wal::MEM_WAL_INDEX_NAME; +use crate::system_index::mem_wal::{ + CompactedSsTable, MEM_WAL_INDEX_NAME, update_mem_wal_index_compacted_sstables, +}; use crate::transaction::{LogicalIndexSegments, ReadVersionState, Transaction}; use lance_core::datatypes::Schema; use lance_core::{Error, Result}; @@ -316,6 +318,17 @@ impl<'a> ApplyState<'a> { Ok(()) } + /// Record MemWAL SSTable compaction progress against the table's existing + /// MemWAL index, which has to already be there -- progress against an index + /// no one built is rejected rather than inventing the metadata. + pub(super) fn update_compacted_sstables( + &mut self, + compacted_sstables: Vec, + ) -> Result<()> { + let new_version = self.new_version(); + update_mem_wal_index_compacted_sstables(&mut self.indices, new_version, compacted_sstables) + } + /// Where the segment `uuid` lives, checking on the way that it really /// belongs to `name`. A mismatch means the operation was planned against a /// different set of segments, the same way a missing segment does. diff --git a/rust/lance-table/src/transaction/action/footprint.rs b/rust/lance-table/src/transaction/action/footprint.rs index 4470f55ed87..0aa9b00d250 100644 --- a/rust/lance-table/src/transaction/action/footprint.rs +++ b/rust/lance-table/src/transaction/action/footprint.rs @@ -487,6 +487,16 @@ impl Footprint { }); } + /// Record that this set rewrites `name` in a way whose reach it cannot + /// state, colliding with any concurrent write to the same index. + pub(super) fn rewrite_index(&mut self, name: String) { + self.index_claims.push(IndexClaim { + name, + identity: None, + coverage: None, + }); + } + /// Record that this set removes `fragment` outright: its existence, and /// with it every coordinate inside it. pub(super) fn remove_fragment(&mut self, fragment: u64) { diff --git a/rust/lance-table/src/transaction/action/proto.rs b/rust/lance-table/src/transaction/action/proto.rs index 59fb8aab5ff..7ee04dcfb1c 100644 --- a/rust/lance-table/src/transaction/action/proto.rs +++ b/rust/lance-table/src/transaction/action/proto.rs @@ -158,12 +158,13 @@ mod tests { BasePath, DataFile, DeletionFile, DeletionFileType, IndexFile, RowIdMeta, pb, }; use crate::rowids::version::RowDatasetVersionMeta; + use crate::system_index::mem_wal::CompactedSsTable; use crate::transaction::UpdateMap; use crate::transaction::action::{ AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, AddOverlays, AdjustIndexCoverage, AlterField, ConfigUpdate, DropField, FieldMetadataUpdate, RefreshRowVersionMetadata, RemoveFragment, RemoveIndexSegment, ReserveFragmentIds, - ReserveRowIds, ResetTable, SetDeletionFile, TombstoneFieldData, + ReserveRowIds, ResetTable, SetDeletionFile, TombstoneFieldData, UpdateCompactedSsTables, }; use arrow_schema::{DataType, Field as ArrowField}; use chrono::DateTime; @@ -279,6 +280,12 @@ mod tests { add_fragments: vec![Ref::Committed(1), Ref::Local(0)], remove_fragments: vec![2, 3], }), + Action::UpdateCompactedSsTables(UpdateCompactedSsTables { + compacted_sstables: vec![ + CompactedSsTable::new(Uuid::from_u128(10), 2), + CompactedSsTable::new(Uuid::from_u128(11), 5), + ], + }), Action::ReserveFragmentIds(ReserveFragmentIds { count: 4 }), Action::ReserveRowIds(ReserveRowIds { count: 40 }), Action::ResetTable(ResetTable), diff --git a/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs b/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs index 96fef3914fa..16ba603e1d5 100644 --- a/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs +++ b/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs @@ -53,7 +53,7 @@ impl RefreshRowVersionMetadata { pub(super) fn footprint(&self, footprint: &mut Footprint) { for fragment_id in &self.fragment_ids { - footprint.add(Coordinate::FragmentRowVersions(*fragment_id)); + footprint.write(Coordinate::FragmentRowVersions(*fragment_id)); } } } diff --git a/rust/lance-table/src/transaction/action/update_compacted_sstables.rs b/rust/lance-table/src/transaction/action/update_compacted_sstables.rs new file mode 100644 index 00000000000..b9a37b26590 --- /dev/null +++ b/rust/lance-table/src/transaction/action/update_compacted_sstables.rs @@ -0,0 +1,227 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Record which MemWAL SSTables have been compacted into the base table. + +use super::Footprint; +use super::apply::ApplyState; +use crate::format::pb; +use crate::system_index::mem_wal::{CompactedSsTable, MEM_WAL_INDEX_NAME}; +use lance_core::deepsize::DeepSizeOf; +use lance_core::{Error, Result}; + +/// Mark MemWAL SSTables as compacted into the base table. +/// +/// The rows were already readable through the WAL, so this records where they +/// are rather than changing them. Per shard the highest generation wins, so +/// replaying an older commit over a newer one cannot walk the progress +/// backwards. +/// +/// This is the one action that edits the MemWAL system index rather than the +/// data, which is why it exists at all: the index is a segment like any other, +/// but its contents are compaction bookkeeping that no +/// [`AddIndexSegment`](super::AddIndexSegment) could express as a delta. +#[derive(Debug, Clone, PartialEq, DeepSizeOf)] +pub struct UpdateCompactedSsTables { + pub compacted_sstables: Vec, +} + +impl UpdateCompactedSsTables { + pub(super) fn apply(&self, state: &mut ApplyState) -> Result<()> { + if self.compacted_sstables.is_empty() { + return Err(Error::invalid_input( + "UpdateCompactedSsTables names no SSTable, so there is no progress to record", + )); + } + state.update_compacted_sstables(self.compacted_sstables.clone()) + } + + pub(super) fn is_data_change(&self) -> bool { + false + } + + /// A rewrite of the whole MemWAL index. Unlike a segment joining an + /// ordinary index, this restates the index entry outright, so a concurrent + /// writer touching that index at all would lose what this recorded -- + /// which is what an unstated reach claims. + pub(super) fn footprint(&self, footprint: &mut Footprint) { + footprint.rewrite_index(MEM_WAL_INDEX_NAME.to_string()); + } +} + +impl From<&UpdateCompactedSsTables> for pb::UpdateCompactedSsTables { + fn from(value: &UpdateCompactedSsTables) -> Self { + Self { + compacted_sstables: value + .compacted_sstables + .iter() + .map(pb::CompactedSsTable::from) + .collect(), + } + } +} + +impl TryFrom for UpdateCompactedSsTables { + type Error = Error; + + fn try_from(message: pb::UpdateCompactedSsTables) -> Result { + Ok(Self { + compacted_sstables: message + .compacted_sstables + .into_iter() + .map(CompactedSsTable::try_from) + .collect::>>()?, + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::system_index::mem_wal::load_mem_wal_index_details; + use crate::transaction::action::test_support::{apply_with_indices, backed_manifest}; + use crate::transaction::action::{ + Action, AddIndexSegment, CompositeOperation, Ref, UserAction, + }; + use crate::transaction::test_support::sample_index_metadata; + use uuid::Uuid; + + fn update(sstables: Vec<(u128, u64)>) -> Action { + Action::UpdateCompactedSsTables(UpdateCompactedSsTables { + compacted_sstables: sstables + .into_iter() + .map(|(shard, generation)| { + CompactedSsTable::new(Uuid::from_u128(shard), generation) + }) + .collect(), + }) + } + + /// The compaction progress the MemWAL index records, as (shard, generation) + /// pairs sorted by shard. + fn progress(indices: &[crate::format::IndexMetadata]) -> Vec<(u128, u64)> { + let mem_wal = indices + .iter() + .find(|index| index.name == MEM_WAL_INDEX_NAME) + .expect("the MemWAL index should be there"); + let details = load_mem_wal_index_details(mem_wal.clone()).unwrap(); + let mut progress = details + .compacted_sstables + .iter() + .map(|sstable| (sstable.shard_id.as_u128(), sstable.generation)) + .collect::>(); + progress.sort(); + progress + } + + #[test] + fn test_recording_progress_creates_the_mem_wal_index_when_there_is_none() { + let (_, indices) = + apply_with_indices(&backed_manifest(), vec![update(vec![(1, 7)])], Vec::new()).unwrap(); + + assert_eq!(progress(&indices), vec![(1, 7)]); + } + + #[test] + fn test_a_later_generation_supersedes_the_one_recorded_for_that_shard() { + let manifest = backed_manifest(); + let (_, indices) = + apply_with_indices(&manifest, vec![update(vec![(1, 3)])], Vec::new()).unwrap(); + let (_, indices) = + apply_with_indices(&manifest, vec![update(vec![(1, 9)])], indices).unwrap(); + + assert_eq!(progress(&indices), vec![(1, 9)]); + } + + #[test] + fn test_an_earlier_generation_does_not_walk_a_shard_backwards() { + let manifest = backed_manifest(); + let (_, indices) = + apply_with_indices(&manifest, vec![update(vec![(1, 9)])], Vec::new()).unwrap(); + let (_, indices) = + apply_with_indices(&manifest, vec![update(vec![(1, 3)])], indices).unwrap(); + + assert_eq!(progress(&indices), vec![(1, 9)]); + } + + #[test] + fn test_shards_are_tracked_independently() { + let manifest = backed_manifest(); + let (_, indices) = + apply_with_indices(&manifest, vec![update(vec![(1, 3)])], Vec::new()).unwrap(); + let (_, indices) = + apply_with_indices(&manifest, vec![update(vec![(2, 5)])], indices).unwrap(); + + assert_eq!(progress(&indices), vec![(1, 3), (2, 5)]); + } + + #[test] + fn test_recording_no_sstables_is_rejected() { + let error = + apply_with_indices(&backed_manifest(), vec![update(vec![])], Vec::new()).unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains("names no SSTable"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_recording_progress_leaves_other_indices_alone() { + let kept = sample_index_metadata("by_a"); + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![update(vec![(1, 7)])], + vec![kept.clone()], + ) + .unwrap(); + + assert!(indices.iter().any(|index| index.uuid == kept.uuid)); + } + + #[test] + fn test_two_writers_recording_progress_conflict() { + let footprint = |actions| { + Footprint::from(&CompositeOperation::new(vec![UserAction::new( + "step", actions, + )])) + }; + + let ours = footprint(vec![update(vec![(1, 3)])]); + let theirs = footprint(vec![update(vec![(2, 5)])]); + + assert!(ours.conflicts_with(&theirs)); + } + + #[test] + fn test_recording_progress_conflicts_with_rebuilding_the_memwal_index() { + // Recording progress rewrites the whole MemWAL entry, so it collides + // with a concurrent writer replacing that entry -- unlike a regular + // index, where two writers may hold disjoint segments. + let footprint = |actions| { + Footprint::from(&CompositeOperation::new(vec![UserAction::new( + "step", actions, + )])) + }; + + let progress = footprint(vec![update(vec![(1, 3)])]); + let rebuild = footprint(vec![Action::AddIndexSegment(AddIndexSegment { + uuid: Uuid::from_u128(20), + name: MEM_WAL_INDEX_NAME.into(), + fields: vec![Ref::Committed(0)], + covering_fields: Vec::new(), + index_details: None, + index_version: 1, + covered_fragments: None, + files: Vec::new(), + base: None, + created_at: None, + dataset_version: None, + data_change: true, + })]); + + assert!(progress.conflicts_with(&rebuild)); + assert!(rebuild.conflicts_with(&progress)); + } +} From 8f8b8fa0e3b46d234663cd447c72fc668bf071f4 Mon Sep 17 00:00:00 2001 From: Will Jones Date: Tue, 18 Aug 2026 20:50:07 -0700 Subject: [PATCH 04/10] feat(transaction): add the AssertUniqueKeys action A precondition rather than a delta: the keys this operation inserts must not collide with keys a concurrent commit inserted. The key columns are an unenforced primary key, so nothing in the manifest records which keys exist -- the filter of inserted key hashes has to travel with the operation because it cannot be recovered from any post-image. This is the first thing two footprints compare that is not a coordinate, so the footprint grows a row-insertion marker (set by an AddFragment that is a data change) and the assertions themselves. Two sets are compatible when both say which keys they insert, over the same columns, and the filters provably do not intersect; an unqualified insert, different key columns, or filters built with incomparable parameters all leave the assertion unverifiable, which counts as a conflict. With this the implemented vocabulary covers the whole draft, so the "drafted but not implemented" rejection has nothing left to reject and is replaced by one for an action written by a newer Lance -- which protobuf decodes as no variant at all. --- rust/lance-table/src/transaction/action.rs | 18 +- .../src/transaction/action/add_fragment.rs | 19 +- .../transaction/action/assert_unique_keys.rs | 240 ++++++++++++++++++ .../src/transaction/action/footprint.rs | 74 +++++- .../src/transaction/action/proto.rs | 61 +++-- rust/lance-table/src/transaction/proto.rs | 24 +- 6 files changed, 375 insertions(+), 61 deletions(-) create mode 100644 rust/lance-table/src/transaction/action/assert_unique_keys.rs diff --git a/rust/lance-table/src/transaction/action.rs b/rust/lance-table/src/transaction/action.rs index 6477f1e720b..8e72ab220d9 100644 --- a/rust/lance-table/src/transaction/action.rs +++ b/rust/lance-table/src/transaction/action.rs @@ -11,8 +11,9 @@ //! //! The wire format and the reasoning behind it live in //! `protos/transaction/actions.proto`; the two definitions must stay in step. -//! Only the subset of the specified vocabulary that is implemented appears here -- -//! an action this build does not know is rejected on load rather than skipped. +//! Every action the specification defines is implemented here, so the two +//! vocabularies now coincide -- an action a newer Lance writes is rejected on +//! load rather than skipped. //! //! Each action lives in its own module and owns everything about itself: its //! definition, how it is applied, which coordinates it writes, and its wire @@ -78,6 +79,7 @@ macro_rules! for_each_action { ReserveRowIds, ResetTable, ConfigUpdate, + AssertUniqueKeys, } }; } @@ -91,6 +93,7 @@ mod add_overlays; mod adjust_index_coverage; mod alter_field; mod apply; +mod assert_unique_keys; mod config_update; mod drop_field; mod footprint; @@ -116,6 +119,7 @@ pub use add_index_segment::AddIndexSegment; pub use add_overlays::AddOverlays; pub use adjust_index_coverage::AdjustIndexCoverage; pub use alter_field::AlterField; +pub use assert_unique_keys::AssertUniqueKeys; pub use config_update::{ConfigUpdate, FieldMetadataUpdate}; pub use drop_field::DropField; pub use footprint::{ConfigMap, Coordinate, Footprint}; @@ -275,12 +279,12 @@ impl UserAction { macro_rules! define_action { ($($variant:ident,)*) => { - /// A single granular change to the manifest. + /// A single granular change to the manifest, or an assertion about the + /// version it lands on. /// - /// The specified vocabulary is larger than this; the variants here are the - /// ones this build implements end to end. Each one is defined, applied, - /// and encoded in the module named after it, and appears here only - /// because it is listed in `for_each_action!`. + /// Each variant is defined, applied, and encoded in the module named + /// after it, and appears here only because it is listed in + /// `for_each_action!`. #[derive(Debug, Clone, PartialEq, DeepSizeOf)] pub enum Action { $($variant($variant),)* diff --git a/rust/lance-table/src/transaction/action/add_fragment.rs b/rust/lance-table/src/transaction/action/add_fragment.rs index d896325a5b8..dd80e561709 100644 --- a/rust/lance-table/src/transaction/action/add_fragment.rs +++ b/rust/lance-table/src/transaction/action/add_fragment.rs @@ -72,17 +72,26 @@ impl AddFragment { self.data_change } - /// Nothing for a local token: the fragment does not exist in the read + /// No coordinate for a local token: the fragment does not exist in the read /// version, so no concurrent writer can be naming it. /// - /// A committed id does write one coordinate. A reservation is meant to be - /// one writer's alone, but nothing in the format enforces that, so two - /// operations handed the same range would otherwise each add a fragment at - /// the same id and the second commit would silently win. + /// A committed id does write one. A reservation is meant to be one writer's + /// alone, but nothing in the format enforces that, so two operations handed + /// the same range would otherwise each add a fragment at the same id and + /// the second commit would silently win. + /// + /// Either form records that rows arrive, which is the one thing about an + /// added fragment a concurrent + /// [`AssertUniqueKeys`](super::AssertUniqueKeys) has to know. A fragment + /// that is not a data change holds rows that were already in the dataset -- + /// a compaction rewrite -- and brings in no new key. pub(super) fn footprint(&self, footprint: &mut Footprint) { if let Some(id) = self.id.committed() { footprint.write(Coordinate::FragmentExistence(id)); } + if self.data_change { + footprint.insert_rows(); + } } } diff --git a/rust/lance-table/src/transaction/action/assert_unique_keys.rs b/rust/lance-table/src/transaction/action/assert_unique_keys.rs new file mode 100644 index 00000000000..6a0e765fd4b --- /dev/null +++ b/rust/lance-table/src/transaction/action/assert_unique_keys.rs @@ -0,0 +1,240 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Assert that no concurrent commit inserted a colliding key. + +use super::apply::ApplyState; +use super::proto::required; +use super::{Footprint, Ref}; +use crate::format::key_existence::KeyExistenceFilter; +use crate::format::pb; +use lance_core::deepsize::DeepSizeOf; +use lance_core::{Error, Result}; + +/// A precondition rather than a delta: the keys this operation inserts must not +/// collide with keys a concurrent commit inserted. +/// +/// This is the home for merge-insert's strict primary-key conflict detection. +/// The key columns are an unenforced primary key, so nothing in the manifest +/// records which keys exist -- the filter of inserted key hashes travels with +/// the operation because it cannot be recovered from any post-image. +/// +/// Two operations that both carry one are compatible when they agree on the key +/// columns and their filters do not intersect. An operation carrying one is not +/// compatible with a concurrent operation that inserts rows without saying which +/// keys they carry, because there is nothing to compare against. +#[derive(Debug, Clone, PartialEq, DeepSizeOf)] +pub struct AssertUniqueKeys { + /// The key columns, in order. This is the authoritative list; the field ids + /// the filter carries are an artifact of the shared filter type and are left + /// empty here. + pub key_fields: Vec, + pub filter: KeyExistenceFilter, +} + +impl AssertUniqueKeys { + /// Nothing to apply -- the assertion is checked when two operations are + /// compared, not when one is folded into a manifest. The key columns are + /// still resolved and looked up, so an assertion naming a field that is not + /// there fails at the commit that carries it rather than silently guarding + /// nothing. + pub(super) fn apply(&self, state: &mut ApplyState) -> Result<()> { + if self.key_fields.is_empty() { + return Err(Error::invalid_input( + "AssertUniqueKeys names no key column, so there is nothing for it to assert", + )); + } + for key_field in &self.key_fields { + let field_id = state.resolve_field(*key_field)?; + if state.schema().field_by_id(field_id).is_none() { + return Err(Error::invalid_input(format!( + "AssertUniqueKeys names key field {field_id}, which is not in the schema" + ))); + } + } + Ok(()) + } + + pub(super) fn is_data_change(&self) -> bool { + false + } + + pub(super) fn footprint(&self, footprint: &mut Footprint) { + footprint.assert_unique_keys(self.key_fields.clone(), self.filter.clone()); + } +} + +impl From<&AssertUniqueKeys> for pb::AssertUniqueKeys { + fn from(value: &AssertUniqueKeys) -> Self { + let mut filter = pb::KeyExistenceFilter::from(&value.filter); + // `key_fields` is the authoritative list; the filter's own copy is an + // artifact of the type it shares with the legacy Update operation. + filter.field_ids.clear(); + Self { + key_fields: value.key_fields.iter().map(|&f| f.into()).collect(), + filter: Some(filter), + } + } +} + +impl TryFrom for AssertUniqueKeys { + type Error = Error; + + fn try_from(message: pb::AssertUniqueKeys) -> Result { + let filter = required(message.filter, "AssertUniqueKeys.filter")?; + Ok(Self { + key_fields: message + .key_fields + .into_iter() + .map(Ref::try_from) + .collect::>>()?, + filter: KeyExistenceFilter::try_from(&filter)?, + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::format::key_existence::FilterType; + use crate::transaction::action::test_support::{apply, backed_manifest}; + use crate::transaction::action::{ + Action, AddFragment, CompositeOperation, RemoveFragment, UserAction, + }; + + fn assertion(key_fields: Vec, hashes: &[u64]) -> Action { + Action::AssertUniqueKeys(AssertUniqueKeys { + key_fields, + filter: KeyExistenceFilter { + field_ids: Vec::new(), + filter: FilterType::ExactSet(hashes.iter().copied().collect()), + }, + }) + } + + fn append(local: u32) -> Action { + Action::AddFragment(AddFragment { + id: Ref::Local(local), + physical_rows: 5, + row_id_meta: None, + last_updated_at_version_meta: None, + created_at_version_meta: None, + data_change: true, + }) + } + + fn compact(local: u32) -> Action { + let Action::AddFragment(mut fragment) = append(local) else { + unreachable!() + }; + fragment.data_change = false; + Action::AddFragment(fragment) + } + + fn footprint(actions: Vec) -> Footprint { + Footprint::from(&CompositeOperation::new(vec![UserAction::new( + "step", actions, + )])) + } + + #[test] + fn test_an_assertion_changes_nothing() { + let manifest = backed_manifest(); + let out = apply( + &manifest, + vec![assertion(vec![Ref::Committed(0)], &[1, 2, 3])], + ) + .unwrap(); + + assert_eq!(out.fragments, manifest.fragments); + assert_eq!(out.schema, manifest.schema); + } + + #[test] + fn test_asserting_over_no_key_column_is_rejected() { + let error = apply(&backed_manifest(), vec![assertion(vec![], &[1])]).unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains("names no key column"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_asserting_over_a_field_that_is_not_there_is_rejected() { + let error = apply( + &backed_manifest(), + vec![assertion(vec![Ref::Committed(9)], &[1])], + ) + .unwrap_err(); + + assert!( + error.to_string().contains("not in the schema"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_inserts_with_disjoint_keys_do_not_conflict() { + let ours = footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[1, 2])]); + let theirs = footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[3, 4])]); + + assert!(!ours.conflicts_with(&theirs)); + assert!(!theirs.conflicts_with(&ours)); + } + + #[test] + fn test_inserts_sharing_a_key_conflict() { + let ours = footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[1, 2])]); + let theirs = footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[2, 3])]); + + assert!(ours.conflicts_with(&theirs)); + assert!(theirs.conflicts_with(&ours)); + } + + #[test] + fn test_assertions_over_different_key_columns_conflict() { + // Two filters over different columns hash different values, so neither + // says anything about the other's keys. + let ours = footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[1])]); + let theirs = footprint(vec![append(0), assertion(vec![Ref::Committed(1)], &[2])]); + + assert!(ours.conflicts_with(&theirs)); + } + + #[test] + fn test_an_assertion_conflicts_with_an_unqualified_insert() { + let ours = footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[1])]); + let plain_append = footprint(vec![append(0)]); + + assert!(ours.conflicts_with(&plain_append)); + assert!(plain_append.conflicts_with(&ours)); + } + + #[test] + fn test_two_plain_appends_still_do_not_conflict() { + assert!(!footprint(vec![append(0)]).conflicts_with(&footprint(vec![append(0)]))); + } + + #[test] + fn test_an_assertion_ignores_a_concurrent_change_that_inserts_no_rows() { + let ours = footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[1])]); + let removal = footprint(vec![Action::RemoveFragment(RemoveFragment { + fragment: Ref::Committed(3), + data_change: true, + })]); + + assert!(!ours.conflicts_with(&removal)); + } + + #[test] + fn test_a_compaction_is_not_an_insert() { + // Compaction rewrites rows that are already there, so it cannot have + // introduced a key the assertion would have to rule out. + let ours = footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[1])]); + let compaction = footprint(vec![compact(0)]); + + assert!(!ours.conflicts_with(&compaction)); + } +} diff --git a/rust/lance-table/src/transaction/action/footprint.rs b/rust/lance-table/src/transaction/action/footprint.rs index 0aa9b00d250..a1b874120d7 100644 --- a/rust/lance-table/src/transaction/action/footprint.rs +++ b/rust/lance-table/src/transaction/action/footprint.rs @@ -17,6 +17,7 @@ //! module. This module holds the coordinate space and the comparison. use super::{CompositeOperation, Ref}; +use crate::format::key_existence::KeyExistenceFilter; use crate::transaction::UpdateMap; use std::collections::{HashMap, HashSet}; use std::sync::Arc; @@ -215,12 +216,27 @@ pub struct Footprint { /// append's rows would vanish or survive the reset depending on which /// commit landed first -- so it is a flag rather than a region. exclusive: bool, + /// Whether this set brings rows into the dataset that were not there + /// before. Rows are not a coordinate -- a row a concurrent writer inserts + /// has no id anyone could name -- so what the key assertions below compare + /// is this flag plus the filters. + inserts_rows: bool, + /// The unique-key preconditions this set carries, one per + /// [`AssertUniqueKeys`](super::AssertUniqueKeys). + key_assertions: Vec, /// What this set writes into logical indices. A segment's uuid is a /// coordinate, but the thing two writers can collide over is the index the /// segment joins, which is not a set of ids -- see [`IndexClaim`]. index_claims: Vec, } +/// A claim about which keys an action set inserts, and over which columns. +#[derive(Debug, Clone, PartialEq)] +struct KeyAssertion { + key_fields: Vec, + filter: KeyExistenceFilter, +} + /// What an action set writes into one logical index. /// /// An index is the set of segments sharing a name, and the query path unions @@ -306,8 +322,18 @@ impl Footprint { if self.anchors_removed_by(committed) || committed.anchors_removed_by(self) { return true; } - // Index claims are not coordinates and are compared on their own terms, - // symmetrically. + // The value predicates below are not coordinates and are compared on + // their own terms, symmetrically. + // + // Symmetric, unlike the coordinate requirements above. Dropping the + // second direction would let a plain append land a duplicate of a key + // a committed merge-insert asserted was absent -- correct if the + // assertion is a precondition on that commit, wrong if it is an + // invariant on the table. Until that is settled, keep the stricter + // reading. + if self.key_assertion_violated_by(committed) || committed.key_assertion_violated_by(self) { + return true; + } self.index_claims.iter().any(|ours| { committed .index_claims @@ -359,6 +385,39 @@ impl Footprint { .any(|fragment| other.regions.contains_key(&Region::Fragment(*fragment))) } + /// Whether `other` may have inserted a key this set asserts is not there. + /// + /// A set that asserts nothing has nothing to violate, and a set that + /// inserts no rows cannot have inserted a key. Otherwise the two are only + /// compatible if `other` says which keys it inserted, over the same columns, + /// and the two filters provably do not intersect. Anything less -- an + /// unqualified insert, different key columns, filters built with + /// incomparable parameters -- leaves the assertion unverifiable, which + /// counts as a conflict. + fn key_assertion_violated_by(&self, other: &Self) -> bool { + if self.key_assertions.is_empty() || !other.inserts_rows { + return false; + } + if other.key_assertions.is_empty() { + return true; + } + for ours in &self.key_assertions { + for theirs in &other.key_assertions { + if ours.key_fields != theirs.key_fields { + return true; + } + match ours.filter.intersects(&theirs.filter) { + Ok((false, _)) => {} + // Either the keys really do overlap, or the two filters were + // built with parameters that cannot be compared. Neither + // clears the assertion. + Ok((true, _)) | Err(_) => return true, + } + } + } + false + } + /// Record a claim on `coordinate`, keeping the strongest mode if the set /// already touches it another way. See [`pair_conflicts`] for why the /// strongest mode is the right summary. @@ -497,6 +556,17 @@ impl Footprint { }); } + /// Note that this set brings in rows that were not in the dataset before. + pub(super) fn insert_rows(&mut self) { + self.inserts_rows = true; + } + + /// Record a precondition that no concurrent commit inserted a colliding key. + pub(super) fn assert_unique_keys(&mut self, key_fields: Vec, filter: KeyExistenceFilter) { + self.key_assertions + .push(KeyAssertion { key_fields, filter }); + } + /// Record that this set removes `fragment` outright: its existence, and /// with it every coordinate inside it. pub(super) fn remove_fragment(&mut self, fragment: u64) { diff --git a/rust/lance-table/src/transaction/action/proto.rs b/rust/lance-table/src/transaction/action/proto.rs index 7ee04dcfb1c..4649e54687c 100644 --- a/rust/lance-table/src/transaction/action/proto.rs +++ b/rust/lance-table/src/transaction/action/proto.rs @@ -7,10 +7,13 @@ //! [`Ref`], [`CompositeOperation`], and [`UserAction`] wrappers, the dispatch over //! the `oneof`, and the helpers the per-action conversions share. //! -//! Reading is fail-closed: an action this build does not implement is an error, +//! Reading is fail-closed: an action this build does not recognize is an error, //! never a silently skipped element. The commit path collects concurrent //! transactions with `try_collect`, so a transaction carrying an unknown action -//! must abort the commit rather than be treated as a no-op. +//! must abort the commit rather than be treated as a no-op. Every drafted action +//! is implemented, so an unrecognized one can only come from a newer Lance -- +//! which protobuf decodes as no variant at all, since it drops the field it does +//! not know. use super::{Action, CompositeOperation, Ref, UserAction}; use crate::format::pb; @@ -132,15 +135,13 @@ macro_rules! define_action_proto { $(Some(pb::action::Action::$variant(action)) => { Ok(Self::$variant(action.try_into()?)) })* - // The specified vocabulary is larger than what is implemented. - // Reject rather than skip: silently dropping an action would - // apply a partial transaction. - Some(other) => Err(Error::not_supported(format!( - "the action-based transaction uses action {other:?}, which is specified \ - but not implemented by this version of Lance", - ))), - None => Err(Error::invalid_input( - "an Action in a user operation was empty", + // An action written by a newer Lance decodes to no known + // variant, because protobuf drops the field it does not + // know. Reject rather than skip: silently dropping an + // action would apply a partial transaction. + None => Err(Error::not_supported( + "an Action in a user operation carried no change this version of Lance \ + understands; it was either empty or written by a newer version", )), } } @@ -153,6 +154,7 @@ for_each_action!(define_action_proto); #[cfg(test)] mod tests { use super::*; + use crate::format::key_existence::{FilterType, KeyExistenceFilter}; use crate::format::overlay::{DataOverlayFile, OverlayCoverage}; use crate::format::{ BasePath, DataFile, DeletionFile, DeletionFileType, IndexFile, RowIdMeta, pb, @@ -162,9 +164,10 @@ mod tests { use crate::transaction::UpdateMap; use crate::transaction::action::{ AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, AddOverlays, - AdjustIndexCoverage, AlterField, ConfigUpdate, DropField, FieldMetadataUpdate, - RefreshRowVersionMetadata, RemoveFragment, RemoveIndexSegment, ReserveFragmentIds, - ReserveRowIds, ResetTable, SetDeletionFile, TombstoneFieldData, UpdateCompactedSsTables, + AdjustIndexCoverage, AlterField, AssertUniqueKeys, ConfigUpdate, DropField, + FieldMetadataUpdate, RefreshRowVersionMetadata, RemoveFragment, RemoveIndexSegment, + ReserveFragmentIds, ReserveRowIds, ResetTable, SetDeletionFile, TombstoneFieldData, + UpdateCompactedSsTables, }; use arrow_schema::{DataType, Field as ArrowField}; use chrono::DateTime; @@ -286,6 +289,13 @@ mod tests { CompactedSsTable::new(Uuid::from_u128(11), 5), ], }), + Action::AssertUniqueKeys(AssertUniqueKeys { + key_fields: vec![Ref::Committed(1), Ref::Local(3)], + filter: KeyExistenceFilter { + field_ids: Vec::new(), + filter: FilterType::ExactSet([7u64, 9].into_iter().collect()), + }, + }), Action::ReserveFragmentIds(ReserveFragmentIds { count: 4 }), Action::ReserveRowIds(ReserveRowIds { count: 40 }), Action::ResetTable(ResetTable), @@ -336,29 +346,16 @@ mod tests { } #[test] - fn test_unimplemented_action_is_rejected() { - let message = pb::Action { - action: Some(pb::action::Action::AssertUniqueKeys(pb::AssertUniqueKeys { - key_fields: vec![], - filter: None, - })), - }; - let error = Action::try_from(message).unwrap_err(); + fn test_an_unrecognized_action_is_rejected() { + // An empty action and one a newer Lance wrote decode the same way: + // protobuf drops the field this build does not know. + let error = Action::try_from(pb::Action { action: None }).unwrap_err(); assert!( matches!(error, Error::NotSupported { .. }), "expected NotSupported, got {error:?}" ); assert!( - error.to_string().contains("not implemented"), - "unexpected message: {error}" - ); - } - - #[test] - fn test_empty_action_is_rejected() { - let error = Action::try_from(pb::Action { action: None }).unwrap_err(); - assert!( - error.to_string().contains("was empty"), + error.to_string().contains("written by a newer version"), "unexpected message: {error}" ); } diff --git a/rust/lance-table/src/transaction/proto.rs b/rust/lance-table/src/transaction/proto.rs index 139ed30ba0b..b96099c4559 100644 --- a/rust/lance-table/src/transaction/proto.rs +++ b/rust/lance-table/src/transaction/proto.rs @@ -907,27 +907,21 @@ mod tests { } #[test] - fn test_unimplemented_action_fails_closed_on_load() { - // The vocabulary is larger than what is implemented. Loading a - // transaction that uses an unimplemented action must fail rather than - // parse leniently: load_and_sort_new_transactions collects concurrent - // transactions with try_collect, so this aborts an in-flight commit - // instead of letting it proceed against a change it cannot see. + fn test_unrecognized_action_fails_closed_on_load() { + // An action a newer Lance wrote decodes to no known variant, because + // protobuf drops the field this build does not know. Loading it must + // fail rather than parse leniently: load_and_sort_new_transactions + // collects concurrent transactions with try_collect, so this aborts an + // in-flight commit instead of letting it proceed against a change it + // cannot see. let message = pb::Transaction { read_version: 1, uuid: Uuid::new_v4().to_string(), operation: Some(pb::transaction::Operation::CompositeOperation( pb::CompositeOperation { actions: vec![pb::UserAction { - description: "assert unique keys".to_string(), - actions: vec![pb::Action { - action: Some(pb::action::Action::AssertUniqueKeys( - pb::AssertUniqueKeys { - key_fields: vec![], - filter: None, - }, - )), - }], + description: "something from the future".to_string(), + actions: vec![pb::Action { action: None }], }], }, )), From 9b963ef7fc900ba12cf93002f3207dc922b8a00e Mon Sep 17 00:00:00 2001 From: Will Jones Date: Tue, 18 Aug 2026 20:57:57 -0700 Subject: [PATCH 05/10] test(lance): commit the remaining actions against a real dataset Four commits through the real commit path: appending a fragment and overlaying an existing one in the same version, restamping row versions for a fragment whose column was rewritten in place, recording MemWAL compaction progress, and carrying a key assertion alongside the insert it guards. --- .../src/dataset/tests/dataset_transactions.rs | 170 +++++++++++++++++- 1 file changed, 166 insertions(+), 4 deletions(-) diff --git a/rust/lance/src/dataset/tests/dataset_transactions.rs b/rust/lance/src/dataset/tests/dataset_transactions.rs index 2e677c6db5b..b0cde14377c 100644 --- a/rust/lance/src/dataset/tests/dataset_transactions.rs +++ b/rust/lance/src/dataset/tests/dataset_transactions.rs @@ -1747,15 +1747,21 @@ mod composite { use arrow_schema::{DataType, Field, Schema}; use lance_core::datatypes::Field as LanceField; use lance_table::feature_flags::FLAG_COVERED_INDEX_METADATA; + use lance_table::format::key_existence::{FilterType, KeyExistenceFilter}; + use lance_table::format::overlay::{DataOverlayFile, OverlayCoverage}; use lance_table::format::{ BasePath, DataFile, DeletionFile, DeletionFileType, IndexMetadata, RowIdMeta, }; use lance_table::rowids::{RowIdSequence, write_row_ids}; + use lance_table::system_index::mem_wal::{ + CompactedSsTable, MEM_WAL_INDEX_NAME, load_mem_wal_index_details, + }; use lance_table::transaction::action::{ - Action, AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, AdjustIndexCoverage, - AlterField, CompositeOperation, ConfigUpdate, DropField, Ref, RemoveFragment, - RemoveIndexSegment, ReserveFragmentIds, ReserveRowIds, ResetTable, SetDeletionFile, - TombstoneFieldData, UserAction, + Action, AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, AddOverlays, + AdjustIndexCoverage, AlterField, AssertUniqueKeys, CompositeOperation, ConfigUpdate, + DropField, Ref, RefreshRowVersionMetadata, RemoveFragment, RemoveIndexSegment, + ReserveFragmentIds, ReserveRowIds, ResetTable, SetDeletionFile, TombstoneFieldData, + UpdateCompactedSsTables, UserAction, }; use lance_table::transaction::{Operation, Transaction, UpdateMap, UpdateMapEntry}; use roaring::RoaringBitmap; @@ -2851,4 +2857,160 @@ mod composite { assert_eq!(index_coverage(&dataset, "by_a").await, vec![1]); } + + #[tokio::test] + async fn test_one_commit_appends_and_overlays_what_was_already_there() { + let dataset = test_dataset(false).await; + let overlaid = dataset.fragments()[0].id; + let file = existing_data_file(&dataset, 0); + let overlay_file = existing_data_file(&dataset, 1); + let expected_version = dataset.version().version + 1; + + let dataset = commit( + dataset, + vec![ + Action::AddFragment(AddFragment { + id: Ref::Local(0), + physical_rows: 5, + row_id_meta: None, + last_updated_at_version_meta: None, + created_at_version_meta: None, + data_change: true, + }), + Action::AddDataFile(AddDataFile { + fragment: Ref::Local(0), + file, + field_ids: vec![Ref::Committed(0)], + data_change: true, + }), + Action::AddOverlays(AddOverlays { + fragment: Ref::Committed(overlaid), + overlays: vec![DataOverlayFile { + data_file: overlay_file, + coverage: OverlayCoverage::dense([0u32, 2].into_iter().collect()), + // Left for the commit to stamp. + committed_version: 0, + }], + data_change: true, + }), + ], + ) + .await; + + assert_eq!(dataset.fragments().len(), 3); + let fragment = dataset + .fragments() + .iter() + .find(|fragment| fragment.id == overlaid) + .unwrap(); + assert_eq!(fragment.overlays.len(), 1); + assert_eq!(fragment.overlays[0].committed_version, expected_version); + } + + #[tokio::test] + async fn test_a_commit_restamps_row_versions_for_a_fragment_it_rewrote() { + let dataset = test_dataset(true).await; + let rewritten = dataset.fragments()[0].id; + let file = existing_data_file(&dataset, 1); + let expected_version = dataset.version().version + 1; + + let dataset = commit( + dataset, + vec![ + // Rewriting a column in place leaves the rows where they are, so + // nothing else in the commit says when they last changed. + Action::AddDataFile(AddDataFile { + fragment: Ref::Committed(rewritten), + file, + field_ids: vec![Ref::Committed(0)], + data_change: true, + }), + Action::RefreshRowVersionMetadata(RefreshRowVersionMetadata { + fragment_ids: vec![rewritten], + }), + ], + ) + .await; + + let fragment = dataset + .fragments() + .iter() + .find(|fragment| fragment.id == rewritten) + .unwrap(); + let sequence = fragment + .last_updated_at_version_meta + .as_ref() + .expect("the rewritten fragment should carry a last-updated sequence") + .load_sequence() + .unwrap(); + let versions = (0..fragment.physical_rows.unwrap()) + .map(|offset| sequence.version_at(offset).unwrap()) + .collect::>(); + assert_eq!(versions, vec![expected_version; versions.len()]); + } + + #[tokio::test] + async fn test_a_commit_records_mem_wal_compaction_progress() { + let dataset = test_dataset(false).await; + let shard = Uuid::new_v4(); + + let dataset = commit( + dataset, + vec![Action::UpdateCompactedSsTables(UpdateCompactedSsTables { + compacted_sstables: vec![CompactedSsTable::new(shard, 3)], + })], + ) + .await; + + let indices = load_all_indices(&dataset).await.unwrap(); + let mem_wal = indices + .iter() + .find(|index| index.name == MEM_WAL_INDEX_NAME) + .expect("the MemWAL index should have been created"); + let details = load_mem_wal_index_details((*mem_wal).clone()).unwrap(); + assert_eq!(details.compacted_sstables.len(), 1); + assert_eq!(details.compacted_sstables[0].shard_id, shard); + assert_eq!(details.compacted_sstables[0].generation, 3); + + // The data is untouched: recording where rows live is not a change to + // them. + assert_eq!(dataset.fragments().len(), 2); + } + + #[tokio::test] + async fn test_a_commit_carrying_a_key_assertion_changes_nothing_itself() { + let dataset = test_dataset(false).await; + let file = existing_data_file(&dataset, 0); + let before = dataset.fragments().len(); + + let dataset = commit( + dataset, + vec![ + Action::AddFragment(AddFragment { + id: Ref::Local(0), + physical_rows: 5, + row_id_meta: None, + last_updated_at_version_meta: None, + created_at_version_meta: None, + data_change: true, + }), + Action::AddDataFile(AddDataFile { + fragment: Ref::Local(0), + file, + field_ids: vec![Ref::Committed(0)], + data_change: true, + }), + Action::AssertUniqueKeys(AssertUniqueKeys { + key_fields: vec![Ref::Committed(0)], + filter: KeyExistenceFilter { + field_ids: Vec::new(), + filter: FilterType::ExactSet([11u64, 12].into_iter().collect()), + }, + }), + ], + ) + .await; + + assert_eq!(dataset.fragments().len(), before + 1); + } } From a28762d2b3858f63bdd99fa52288743f565f2690 Mon Sep 17 00:00:00 2001 From: Will Jones Date: Wed, 19 Aug 2026 13:06:26 -0700 Subject: [PATCH 06/10] fix(transaction): adapt the remaining actions to upstream API changes `update_mem_wal_index_compacted_sstables` stopped creating the index and stopped tolerating a stale generation, so `UpdateCompactedSsTables` now inherits both rejections. Its docs and tests say so, and the tests seed the index the way a real table would have it. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/transaction/action/add_overlays.rs | 10 ++- .../action/refresh_row_version_metadata.rs | 6 +- .../action/update_compacted_sstables.rs | 72 +++++++++++++------ .../src/dataset/tests/dataset_transactions.rs | 23 +++++- 4 files changed, 84 insertions(+), 27 deletions(-) diff --git a/rust/lance-table/src/transaction/action/add_overlays.rs b/rust/lance-table/src/transaction/action/add_overlays.rs index 930e39cf85a..b735d01ada4 100644 --- a/rust/lance-table/src/transaction/action/add_overlays.rs +++ b/rust/lance-table/src/transaction/action/add_overlays.rs @@ -97,12 +97,20 @@ mod tests { use crate::transaction::action::{ Action, AddFragment, CompositeOperation, RemoveFragment, UserAction, }; + use lance_file::version::ConcreteFileVersion; use roaring::RoaringBitmap; use std::sync::Arc; fn overlay(path: &str, offsets: &[u32]) -> DataOverlayFile { DataOverlayFile { - data_file: DataFile::new(path, vec![0], vec![0], 2, 0, None, None), + data_file: DataFile::new( + path, + vec![0], + vec![0], + ConcreteFileVersion::V2_0, + None, + None, + ), coverage: OverlayCoverage::Shared(Arc::new( offsets.iter().copied().collect::(), )), diff --git a/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs b/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs index 16ba603e1d5..7ab0fda5906 100644 --- a/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs +++ b/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs @@ -84,6 +84,7 @@ mod tests { use crate::transaction::action::test_support::{apply, backed_manifest}; use crate::transaction::action::{Action, CompositeOperation, UserAction}; use crate::transaction::test_support::make_stable_row_id_manifest; + use lance_file::version::ConcreteFileVersion; use std::sync::Arc; fn refresh(fragment_ids: Vec) -> Action { @@ -100,14 +101,13 @@ mod tests { "data.lance", vec![0], vec![0], - 2, - 0, + ConcreteFileVersion::V2_0, None, None, )], overlays: vec![], deletion_file: None, - row_id_meta: Some(RowIdMeta::Inline(write_row_ids(&row_ids))), + row_id_meta: Some(RowIdMeta::Inline(write_row_ids(&row_ids).into())), physical_rows: Some(rows), last_updated_at_version_meta: None, created_at_version_meta: None, diff --git a/rust/lance-table/src/transaction/action/update_compacted_sstables.rs b/rust/lance-table/src/transaction/action/update_compacted_sstables.rs index b9a37b26590..bfc21cd18c4 100644 --- a/rust/lance-table/src/transaction/action/update_compacted_sstables.rs +++ b/rust/lance-table/src/transaction/action/update_compacted_sstables.rs @@ -13,9 +13,9 @@ use lance_core::{Error, Result}; /// Mark MemWAL SSTables as compacted into the base table. /// /// The rows were already readable through the WAL, so this records where they -/// are rather than changing them. Per shard the highest generation wins, so -/// replaying an older commit over a newer one cannot walk the progress -/// backwards. +/// are rather than changing them. A shard's generation may only move forward, +/// and the table must already carry a MemWAL index: progress against a shard +/// nothing corroborates is rejected rather than invented. /// /// This is the one action that edits the MemWAL system index rather than the /// data, which is why it exists at all: the index is a segment like any other, @@ -78,7 +78,9 @@ impl TryFrom for UpdateCompactedSsTables { #[cfg(test)] mod tests { use super::*; - use crate::system_index::mem_wal::load_mem_wal_index_details; + use crate::system_index::mem_wal::{ + MemWalIndexDetails, load_mem_wal_index_details, new_mem_wal_index_meta, + }; use crate::transaction::action::test_support::{apply_with_indices, backed_manifest}; use crate::transaction::action::{ Action, AddIndexSegment, CompositeOperation, Ref, UserAction, @@ -114,19 +116,33 @@ mod tests { progress } + /// A MemWAL index recording no compaction progress yet, which recording + /// progress requires the table to already have. + fn empty_mem_wal_index() -> crate::format::IndexMetadata { + new_mem_wal_index_meta(1, MemWalIndexDetails::default()).unwrap() + } + #[test] - fn test_recording_progress_creates_the_mem_wal_index_when_there_is_none() { - let (_, indices) = - apply_with_indices(&backed_manifest(), vec![update(vec![(1, 7)])], Vec::new()).unwrap(); + fn test_recording_progress_without_a_mem_wal_index_is_rejected() { + let error = apply_with_indices(&backed_manifest(), vec![update(vec![(1, 7)])], Vec::new()) + .unwrap_err(); - assert_eq!(progress(&indices), vec![(1, 7)]); + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains("does not exist on this table"), + "unexpected error: {error}" + ); } #[test] fn test_a_later_generation_supersedes_the_one_recorded_for_that_shard() { let manifest = backed_manifest(); - let (_, indices) = - apply_with_indices(&manifest, vec![update(vec![(1, 3)])], Vec::new()).unwrap(); + let (_, indices) = apply_with_indices( + &manifest, + vec![update(vec![(1, 3)])], + vec![empty_mem_wal_index()], + ) + .unwrap(); let (_, indices) = apply_with_indices(&manifest, vec![update(vec![(1, 9)])], indices).unwrap(); @@ -134,21 +150,31 @@ mod tests { } #[test] - fn test_an_earlier_generation_does_not_walk_a_shard_backwards() { + fn test_an_earlier_generation_is_rejected_rather_than_walking_a_shard_backwards() { let manifest = backed_manifest(); - let (_, indices) = - apply_with_indices(&manifest, vec![update(vec![(1, 9)])], Vec::new()).unwrap(); - let (_, indices) = - apply_with_indices(&manifest, vec![update(vec![(1, 3)])], indices).unwrap(); + let (_, indices) = apply_with_indices( + &manifest, + vec![update(vec![(1, 9)])], + vec![empty_mem_wal_index()], + ) + .unwrap(); + let error = apply_with_indices(&manifest, vec![update(vec![(1, 3)])], indices).unwrap_err(); - assert_eq!(progress(&indices), vec![(1, 9)]); + assert!( + error.to_string().contains("Stale SSTable compaction"), + "unexpected error: {error}" + ); } #[test] fn test_shards_are_tracked_independently() { let manifest = backed_manifest(); - let (_, indices) = - apply_with_indices(&manifest, vec![update(vec![(1, 3)])], Vec::new()).unwrap(); + let (_, indices) = apply_with_indices( + &manifest, + vec![update(vec![(1, 3)])], + vec![empty_mem_wal_index()], + ) + .unwrap(); let (_, indices) = apply_with_indices(&manifest, vec![update(vec![(2, 5)])], indices).unwrap(); @@ -157,8 +183,12 @@ mod tests { #[test] fn test_recording_no_sstables_is_rejected() { - let error = - apply_with_indices(&backed_manifest(), vec![update(vec![])], Vec::new()).unwrap_err(); + let error = apply_with_indices( + &backed_manifest(), + vec![update(vec![])], + vec![empty_mem_wal_index()], + ) + .unwrap_err(); assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); assert!( @@ -173,7 +203,7 @@ mod tests { let (_, indices) = apply_with_indices( &backed_manifest(), vec![update(vec![(1, 7)])], - vec![kept.clone()], + vec![kept.clone(), empty_mem_wal_index()], ) .unwrap(); diff --git a/rust/lance/src/dataset/tests/dataset_transactions.rs b/rust/lance/src/dataset/tests/dataset_transactions.rs index b0cde14377c..384d817f4bf 100644 --- a/rust/lance/src/dataset/tests/dataset_transactions.rs +++ b/rust/lance/src/dataset/tests/dataset_transactions.rs @@ -1754,7 +1754,8 @@ mod composite { }; use lance_table::rowids::{RowIdSequence, write_row_ids}; use lance_table::system_index::mem_wal::{ - CompactedSsTable, MEM_WAL_INDEX_NAME, load_mem_wal_index_details, + CompactedSsTable, MEM_WAL_INDEX_NAME, MemWalIndexDetails, load_mem_wal_index_details, + new_mem_wal_index_meta, }; use lance_table::transaction::action::{ Action, AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, AddOverlays, @@ -2954,6 +2955,24 @@ mod composite { let dataset = test_dataset(false).await; let shard = Uuid::new_v4(); + // Progress is only recordable against a table that already carries the + // index, so put one there the way the MemWAL writer would. + let read_version = dataset.version().version; + let dataset = CommitBuilder::new(Arc::new(dataset)) + .execute(Transaction::new( + read_version, + Operation::CreateIndex { + new_indices: vec![ + new_mem_wal_index_meta(read_version, MemWalIndexDetails::default()) + .unwrap(), + ], + removed_indices: Vec::new(), + }, + None, + )) + .await + .unwrap(); + let dataset = commit( dataset, vec![Action::UpdateCompactedSsTables(UpdateCompactedSsTables { @@ -2966,7 +2985,7 @@ mod composite { let mem_wal = indices .iter() .find(|index| index.name == MEM_WAL_INDEX_NAME) - .expect("the MemWAL index should have been created"); + .expect("the MemWAL index should still be there"); let details = load_mem_wal_index_details((*mem_wal).clone()).unwrap(); assert_eq!(details.compacted_sstables.len(), 1); assert_eq!(details.compacted_sstables[0].shard_id, shard); From 5898a64e68fbf6791f288c27203fd123f88a247e Mon Sep 17 00:00:00 2001 From: Will Jones Date: Wed, 19 Aug 2026 13:57:14 -0700 Subject: [PATCH 07/10] docs(transaction): correct why key uniqueness is not a coordinate The comment claimed rows are not a coordinate because a row a concurrent writer inserts has no id anyone could name. Both halves are wrong: rows do have ids, and a new data file can carry rows that already existed. The actual reason is that two writers inserting the same user-supplied key write it into fragments of their own, so their coordinates stay disjoint however badly the keys collide. Also records that the flag over-approximates -- the rows a merge insert updates arrive in a new fragment too -- and why `Update` has no translation yet. Co-Authored-By: Claude Opus 5 (1M context) --- .../lance-table/src/transaction/action/footprint.rs | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/rust/lance-table/src/transaction/action/footprint.rs b/rust/lance-table/src/transaction/action/footprint.rs index a1b874120d7..42cb3b741e8 100644 --- a/rust/lance-table/src/transaction/action/footprint.rs +++ b/rust/lance-table/src/transaction/action/footprint.rs @@ -217,9 +217,16 @@ pub struct Footprint { /// commit landed first -- so it is a flag rather than a region. exclusive: bool, /// Whether this set brings rows into the dataset that were not there - /// before. Rows are not a coordinate -- a row a concurrent writer inserts - /// has no id anyone could name -- so what the key assertions below compare - /// is this flag plus the filters. + /// before. Coordinates cannot answer what the key assertions ask: two + /// writers inserting the same key write it into fragments of their own, so + /// their coordinates stay disjoint however badly the keys collide. Key + /// uniqueness is a claim about values, not about structure, so it is + /// checked by comparing this flag against the filters below. + /// + /// Set by any [`AddFragment`](super::AddFragment) that is a data change, + /// which over-approximates: the rows a merge insert updates arrive in a new + /// fragment too, and those carry no new key. The cost is a conflict between + /// two writers who only ever touched keys that were already there. inserts_rows: bool, /// The unique-key preconditions this set carries, one per /// [`AssertUniqueKeys`](super::AssertUniqueKeys). From e0941316aa377d9d5f0cfe49feb8d96e0a73a87f Mon Sep 17 00:00:00 2001 From: Will Jones Date: Mon, 21 Sep 2026 10:39:45 -0700 Subject: [PATCH 08/10] docs(transaction): call Transaction V2 experimental, not a draft Completes the rename on this layer, where the vocabulary is finished: "every drafted action is implemented" becomes "every specified action". Co-Authored-By: Claude Opus 5 (1M context) --- rust/lance-table/src/transaction/action/proto.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rust/lance-table/src/transaction/action/proto.rs b/rust/lance-table/src/transaction/action/proto.rs index 4649e54687c..d5c5b00fb0a 100644 --- a/rust/lance-table/src/transaction/action/proto.rs +++ b/rust/lance-table/src/transaction/action/proto.rs @@ -10,7 +10,7 @@ //! Reading is fail-closed: an action this build does not recognize is an error, //! never a silently skipped element. The commit path collects concurrent //! transactions with `try_collect`, so a transaction carrying an unknown action -//! must abort the commit rather than be treated as a no-op. Every drafted action +//! must abort the commit rather than be treated as a no-op. Every specified action //! is implemented, so an unrecognized one can only come from a newer Lance -- //! which protobuf decodes as no variant at all, since it drops the field it does //! not know. From 6d18f6ee1838d5b2114f3a9ed53d429e06f652e2 Mon Sep 17 00:00:00 2001 From: Will Jones Date: Tue, 22 Sep 2026 16:36:32 -0700 Subject: [PATCH 09/10] fix(transaction): give overlays and key assertions the coordinates they were missing Overlays. An overlay recorded only that its fragment must still exist, which let a full rewrite of the same column land on either side of it. Landing second, the rewrite tombstones the overlay as it applies and replaces every cell from a snapshot that never saw it, so the overlay's values vanished without a trace -- the case the legacy Update-vs- DataOverlay check exists for. The footprint now distinguishes a partial write from a full one: two partial writes to one coordinate commute (both overlays land, the newer wins), a partial and a full write do not. An overlay is a partial write of each overlaid field's data in the fragment, and requires each field's definition for the same reason a data file does, so an overlay landing after a cast or a drop of its field is a conflict too. `required_fragments` keeps only its add-column case; the overlay's fragment now follows from its partial writes. Key assertions. A key could arrive without a new row: a column rewrite or an overlay on the key column puts any value there, and the writer asserts nothing about which. Such a write was invisible to the assertion check, which only looked at sets that insert rows. Any write, whole or partial, into a key column now violates an assertion over it. And a set carrying assertions over two key column sets no longer conflicts with every other asserting set: each assertion is matched against the other side's assertions over the same columns, and only an insert that says nothing about those columns is unverifiable. Also says, on RefreshRowVersionMetadata, why its fragment-wide coordinate is over-strict and why that is acceptable. Co-Authored-By: Claude Fable 5.1 --- .../src/transaction/action/add_overlays.rs | 85 +++++++++++-- .../transaction/action/assert_unique_keys.rs | 59 ++++++++- .../src/transaction/action/footprint.rs | 113 ++++++++++++------ .../action/refresh_row_version_metadata.rs | 5 + 4 files changed, 220 insertions(+), 42 deletions(-) diff --git a/rust/lance-table/src/transaction/action/add_overlays.rs b/rust/lance-table/src/transaction/action/add_overlays.rs index b735d01ada4..95edd8869d0 100644 --- a/rust/lance-table/src/transaction/action/add_overlays.rs +++ b/rust/lance-table/src/transaction/action/add_overlays.rs @@ -5,7 +5,7 @@ use super::apply::ApplyState; use super::proto::{data_change_from_wire, data_change_to_wire, required}; -use super::{Footprint, Ref}; +use super::{Coordinate, Footprint, Ref}; use crate::format::overlay::DataOverlayFile; use crate::format::pb; use lance_core::deepsize::DeepSizeOf; @@ -48,12 +48,37 @@ impl AddOverlays { self.data_change } - /// Only that the fragment must still be there. An overlay writes no - /// coordinate of its own: two concurrent overlays over the same cells both - /// land, and the newer `committed_version` decides which value wins. + /// A partial write of each overlaid field's data in the fragment, and a + /// requirement on each of those fields' definitions. + /// + /// Partial, because two concurrent overlays over the same cells both land + /// and the newer `committed_version` decides which value wins. A full + /// rewrite of the same column is another matter: it replaces every cell + /// from a snapshot that never saw this overlay, and + /// [`TombstoneFieldData`](super::TombstoneFieldData) tombstones the overlay + /// as it applies, so the overlay's values would be lost without a trace. + /// That pair conflicts in either order, as the legacy Update-vs-DataOverlay + /// check already does. + /// + /// The definition is required for the same reason a data file requires it: + /// the overlay's values are encoded in the type the schema names, and a + /// concurrent cast or drop landing first would leave them described as + /// something they are not, or over a field the manifest no longer has. + /// + /// A fragment this operation mints records nothing; no concurrent writer + /// can name a cell inside it. pub(super) fn footprint(&self, footprint: &mut Footprint) { - if let Some(fragment) = self.fragment.committed() { - footprint.require_fragment(fragment); + let Some(fragment) = self.fragment.committed() else { + return; + }; + for field in self + .overlays + .iter() + .flat_map(|overlay| overlay.data_file.fields.iter().copied()) + .filter(|field| *field >= 0) + { + footprint.write_part(Coordinate::FieldData { fragment, field }); + footprint.require(Coordinate::FieldDefinition(field)); } } } @@ -95,7 +120,8 @@ mod tests { use crate::format::overlay::OverlayCoverage; use crate::transaction::action::test_support::{apply, backed_manifest}; use crate::transaction::action::{ - Action, AddFragment, CompositeOperation, RemoveFragment, UserAction, + Action, AddFragment, AlterField, CompositeOperation, DropField, RemoveFragment, + TombstoneFieldData, UserAction, }; use lance_file::version::ConcreteFileVersion; use roaring::RoaringBitmap; @@ -236,6 +262,51 @@ mod tests { assert!(!ours.conflicts_with(&theirs)); } + /// A full rewrite of the overlaid column replaces every cell from a + /// snapshot that never saw the overlay, and tombstones the overlay as it + /// applies; the overlay's values would be gone without a trace. Either + /// order loses, so the pair is symmetric. + #[test] + fn test_an_overlay_and_a_rewrite_of_the_same_column_conflict() { + let overlaid = footprint(vec![add(Ref::Committed(0), vec![overlay("a.lance", &[0])])]); + let rewrite = |field: i32| { + footprint(vec![Action::TombstoneFieldData(TombstoneFieldData { + fragment: Ref::Committed(0), + field_ids: vec![Ref::Committed(field as u64)], + data_change: true, + })]) + }; + + assert!(overlaid.conflicts_with(&rewrite(0))); + assert!(rewrite(0).conflicts_with(&overlaid)); + // Another column of the same fragment is not touched by the overlay. + assert!(!overlaid.conflicts_with(&rewrite(1))); + assert!(!rewrite(1).conflicts_with(&overlaid)); + } + + /// An overlay landing after its field was cast or dropped would carry + /// values in a type the schema no longer names, or for a field that is + /// gone. Landing first, the cast rebinds the field everywhere and the drop + /// tombstones the overlay, so those orders are fine. + #[test] + fn test_an_overlay_needs_its_field_defined_as_it_was() { + let overlaid = footprint(vec![add(Ref::Committed(0), vec![overlay("a.lance", &[0])])]); + let cast = footprint(vec![Action::AlterField(AlterField { + field: Ref::Committed(0), + name: None, + logical_type: Some("int64".into()), + nullable: None, + })]); + let dropped = footprint(vec![Action::DropField(DropField { + field: Ref::Committed(0), + })]); + + assert!(overlaid.conflicts_with(&cast)); + assert!(!cast.conflicts_with(&overlaid)); + assert!(overlaid.conflicts_with(&dropped)); + assert!(!dropped.conflicts_with(&overlaid)); + } + #[test] fn test_overlaying_a_fragment_a_concurrent_writer_removes_conflicts() { let ours = footprint(vec![add(Ref::Committed(0), vec![overlay("a.lance", &[0])])]); diff --git a/rust/lance-table/src/transaction/action/assert_unique_keys.rs b/rust/lance-table/src/transaction/action/assert_unique_keys.rs index 6a0e765fd4b..290758584c6 100644 --- a/rust/lance-table/src/transaction/action/assert_unique_keys.rs +++ b/rust/lance-table/src/transaction/action/assert_unique_keys.rs @@ -23,6 +23,13 @@ use lance_core::{Error, Result}; /// columns and their filters do not intersect. An operation carrying one is not /// compatible with a concurrent operation that inserts rows without saying which /// keys they carry, because there is nothing to compare against. +/// +/// An assertion speaks for *every* row the operation inserts, over its key +/// columns. An operation may carry several, one per key column set, but each +/// must still cover all of its inserted rows: two merge-inserts over different +/// keys squashed into one operation cannot each assert only their own rows, +/// because the other's rows carry values in those columns too. Conflict +/// detection matches assertions by key column set on that assumption. #[derive(Debug, Clone, PartialEq, DeepSizeOf)] pub struct AssertUniqueKeys { /// The key columns, in order. This is the authoritative list; the field ids @@ -99,7 +106,7 @@ mod tests { use crate::format::key_existence::FilterType; use crate::transaction::action::test_support::{apply, backed_manifest}; use crate::transaction::action::{ - Action, AddFragment, CompositeOperation, RemoveFragment, UserAction, + Action, AddFragment, CompositeOperation, RemoveFragment, TombstoneFieldData, UserAction, }; fn assertion(key_fields: Vec, hashes: &[u64]) -> Action { @@ -203,6 +210,56 @@ mod tests { assert!(ours.conflicts_with(&theirs)); } + /// Each assertion speaks for every row the set inserts, over its own key + /// columns. A set carrying two of them -- a table with two unique keys, + /// each asserted over all the inserted rows -- is compared column set by + /// column set, not every assertion against every other. + #[test] + fn test_assertions_are_matched_by_key_column() { + let ours = footprint(vec![ + append(0), + assertion(vec![Ref::Committed(0)], &[1]), + assertion(vec![Ref::Committed(1)], &[5]), + ]); + let disjoint = footprint(vec![ + append(0), + assertion(vec![Ref::Committed(0)], &[2]), + assertion(vec![Ref::Committed(1)], &[6]), + ]); + let overlapping_on_the_second = footprint(vec![ + append(0), + assertion(vec![Ref::Committed(0)], &[2]), + assertion(vec![Ref::Committed(1)], &[5]), + ]); + let silent_on_the_second = + footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[2])]); + + assert!(!ours.conflicts_with(&disjoint)); + assert!(ours.conflicts_with(&overlapping_on_the_second)); + // An insert that says nothing about one of our key columns could have + // put anything in it. + assert!(ours.conflicts_with(&silent_on_the_second)); + } + + /// A key does not have to arrive in a new row. Rewriting a key column in + /// place -- a column rewrite, an overlay -- can put any value in it, and the + /// writer asserts nothing about which, so the assertion cannot be checked. + #[test] + fn test_an_assertion_conflicts_with_an_in_place_write_to_its_key_column() { + let ours = footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[1])]); + let rewrite = |field: u64| { + footprint(vec![Action::TombstoneFieldData(TombstoneFieldData { + fragment: Ref::Committed(0), + field_ids: vec![Ref::Committed(field)], + data_change: true, + })]) + }; + + assert!(ours.conflicts_with(&rewrite(0))); + assert!(rewrite(0).conflicts_with(&ours)); + assert!(!ours.conflicts_with(&rewrite(1))); + } + #[test] fn test_an_assertion_conflicts_with_an_unqualified_insert() { let ours = footprint(vec![append(0), assertion(vec![Ref::Committed(0)], &[1])]); diff --git a/rust/lance-table/src/transaction/action/footprint.rs b/rust/lance-table/src/transaction/action/footprint.rs index 42cb3b741e8..4d8fb2be372 100644 --- a/rust/lance-table/src/transaction/action/footprint.rs +++ b/rust/lance-table/src/transaction/action/footprint.rs @@ -128,6 +128,11 @@ enum Mode { /// may require the same coordinate -- two readers of one column do not /// collide. Requires, + /// Writes some of the coordinate and leaves the rest as it was. An overlay + /// is the case: new values for a subset of one field's cells in one + /// fragment. Two of these commute -- both land, the newer wins where they + /// overlap -- which is what sets it apart from a full write. + WritesPart, /// Replaces the coordinate. Two of these on one coordinate never commute. Writes, } @@ -137,8 +142,8 @@ enum Mode { /// containment: a region write is a write of every coordinate the region holds. #[derive(Debug, Clone, PartialEq, Eq, Hash)] enum Region { - /// A committed fragment: its existence, its deletions, and the data of - /// every field in it. + /// A committed fragment: its existence, deletions, row versions, and the + /// data of every field in it. Fragment(u64), /// One of the manifest's string maps, every key included. Map(ConfigMap), @@ -157,22 +162,30 @@ impl Region { /// `committed` on the same coordinate, where the `committed` set landed after /// the `committing` set read. /// -/// | committing \ committed | Writes | Requires | -/// |------------------------|----------|----------| -/// | Writes | conflict | ok | -/// | Requires | conflict | ok | +/// | committing \ committed | Writes | WritesPart | Requires | +/// |------------------------|----------|------------|----------| +/// | Writes | conflict | conflict | ok | +/// | WritesPart | conflict | ok | ok | +/// | Requires | conflict | ok | ok | /// /// The `Requires` column is all "ok": what the committed set required held when /// it committed, and it serialized first, so nothing arriving later can /// retroactively break it. The `Requires` row is the open question, because -/// the committing set read before the other landed. +/// the committing set read before the other landed. It collides with a full +/// write and not a partial one: a segment built over a column an overlay then +/// landed on is stamped with the version it read, the overlay carries the later +/// version it committed at, and the read path masks the overlaid cells out of +/// any segment older than the overlay. /// /// The table is monotone along [`Mode`]'s order in every row and column, which /// is what lets one coordinate touched two ways by one set be summarized by /// its strongest mode. fn pair_conflicts(committing: Mode, committed: Mode) -> bool { match (committing, committed) { - (Mode::Writes | Mode::Requires, Mode::Writes) => true, + (Mode::Writes, Mode::Writes | Mode::WritesPart) => true, + (Mode::WritesPart | Mode::Requires, Mode::Writes) => true, + (Mode::WritesPart, Mode::WritesPart) => false, + (Mode::Requires, Mode::WritesPart) => false, (_, Mode::Requires) => false, } } @@ -199,16 +212,13 @@ pub struct Footprint { regions: HashMap, /// Fragments this set needs to still be there, without writing anything a /// concurrent set could name inside them. Data for a field this set mints, - /// written into a committed fragment, is one case: the field id is - /// invisible to a concurrent writer, so the cells are not a coordinate, but - /// they are gone if the fragment is. An overlay is the other: two - /// concurrent overlays over the same cells both land and the newer one - /// wins, so they must not collide with each other -- but neither survives a - /// concurrent writer dropping the fragment out from under them. + /// written into a committed fragment, is the case this exists for: the + /// field id is invisible to a concurrent writer, so the cells are not a + /// coordinate, but they are gone if the fragment is. /// /// Not a `Requires` claim on the region, because the check is symmetric: a /// removal that lands second destroys the cells just as surely as one that - /// lands first. And not a write, because it must not collide with a + /// lands first. And not a `WritesPart`, because it must not collide with a /// concurrent write to some other coordinate in the same fragment. required_fragments: HashSet, /// Whether this set rewrites the table wholesale. Such a set collides with @@ -314,8 +324,8 @@ impl Footprint { /// -- an order the system already handles, by pruning the stale fragment /// out of the segment's coverage as the rewrite applies. /// - /// Index claims are not coordinates and are compared on their own terms, - /// symmetrically. + /// The value predicates -- key assertions, index claims -- are not + /// coordinates and are compared on their own terms, symmetrically. pub fn conflicts_with(&self, committed: &Self) -> bool { // A wholesale rewrite leaves nothing for a concurrent set to land on -- // not even an append, whose rows the reset would discard or resurrect @@ -329,9 +339,6 @@ impl Footprint { if self.anchors_removed_by(committed) || committed.anchors_removed_by(self) { return true; } - // The value predicates below are not coordinates and are compared on - // their own terms, symmetrically. - // // Symmetric, unlike the coordinate requirements above. Dropping the // second direction would let a plain append land a duplicate of a key // a committed merge-insert asserted was absent -- correct if the @@ -392,27 +399,45 @@ impl Footprint { .any(|fragment| other.regions.contains_key(&Region::Fragment(*fragment))) } - /// Whether `other` may have inserted a key this set asserts is not there. + /// Whether `other` may have introduced a key this set asserts is not there. /// - /// A set that asserts nothing has nothing to violate, and a set that - /// inserts no rows cannot have inserted a key. Otherwise the two are only - /// compatible if `other` says which keys it inserted, over the same columns, - /// and the two filters provably do not intersect. Anything less -- an - /// unqualified insert, different key columns, filters built with - /// incomparable parameters -- leaves the assertion unverifiable, which - /// counts as a conflict. + /// A set that asserts nothing has nothing to violate. A key can arrive two + /// ways. Writing into a key column of rows that are already there -- a + /// column rewrite, an overlay -- can put any value in it and says nothing + /// about which, so that is a violation outright. Inserting rows is the + /// other, and is only compatible if `other` says which keys it inserted + /// over the same columns, and the two filters provably do not intersect. + /// Anything less -- an unqualified insert, no assertion over these columns, + /// filters built with incomparable parameters -- leaves the assertion + /// unverifiable, which counts as a conflict. fn key_assertion_violated_by(&self, other: &Self) -> bool { - if self.key_assertions.is_empty() || !other.inserts_rows { + if self.key_assertions.is_empty() { return false; } - if other.key_assertions.is_empty() { + if self + .key_assertions + .iter() + .any(|ours| other.writes_into_any_of(&ours.key_fields)) + { return true; } + if !other.inserts_rows { + return false; + } for ours in &self.key_assertions { - for theirs in &other.key_assertions { - if ours.key_fields != theirs.key_fields { - return true; - } + // Each assertion speaks for every row its set inserts, over its own + // key columns (see `AssertUniqueKeys`). An assertion of theirs + // over other columns hashes other values, so it says nothing about + // ours either way; one over the same columns says everything. + let mut over_same_columns = other + .key_assertions + .iter() + .filter(|theirs| theirs.key_fields == ours.key_fields) + .peekable(); + if over_same_columns.peek().is_none() { + return true; + } + for theirs in over_same_columns { match ours.filter.intersects(&theirs.filter) { Ok((false, _)) => {} // Either the keys really do overlap, or the two filters were @@ -425,6 +450,20 @@ impl Footprint { false } + /// Whether this set writes, wholly or in part, the data of any of `fields` + /// in some committed fragment. + fn writes_into_any_of(&self, fields: &[Ref]) -> bool { + self.claims + .iter() + .filter(|(_, mode)| **mode >= Mode::WritesPart) + .any(|(coordinate, _)| match coordinate { + Coordinate::FieldData { field, .. } => { + fields.contains(&Ref::Committed(*field as u64)) + } + _ => false, + }) + } + /// Record a claim on `coordinate`, keeping the strongest mode if the set /// already touches it another way. See [`pair_conflicts`] for why the /// strongest mode is the right summary. @@ -440,6 +479,12 @@ impl Footprint { self.claim(coordinate, Mode::Writes); } + /// Record that this set writes some of `coordinate`, leaving the rest as + /// it was. See [`Mode::WritesPart`]. + pub(super) fn write_part(&mut self, coordinate: Coordinate) { + self.claim(coordinate, Mode::WritesPart); + } + /// Record that this set reads `coordinate` and needs it to still hold what /// it read. See [`Mode::Requires`]. pub(super) fn require(&mut self, coordinate: Coordinate) { diff --git a/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs b/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs index 7ab0fda5906..51428c5eefe 100644 --- a/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs +++ b/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs @@ -51,6 +51,11 @@ impl RefreshRowVersionMetadata { false } + /// Each fragment's version sequence, whole. The sequence is one value per + /// fragment, restamped for every row at once, so two sets refreshing the + /// same fragment collide even when the column rewrites that prompted them + /// touched different fields. Over-strict, and cheap: the alternative is a + /// sequence that records which columns each row's version speaks for. pub(super) fn footprint(&self, footprint: &mut Footprint) { for fragment_id in &self.fragment_ids { footprint.write(Coordinate::FragmentRowVersions(*fragment_id)); From af1ddba78838dde072da60023731c6750540e5cc Mon Sep 17 00:00:00 2001 From: Will Jones Date: Wed, 23 Sep 2026 09:14:35 -0700 Subject: [PATCH 10/10] test(transaction): pin the overlay and key-assertion pair answers State the answers nothing asserted yet: an index segment may be built over a column a committed overlay landed on, because the segment is stamped with the version it read and the read path masks any overlay committed after that; and a set that both overlays a cell and rewrites the whole column presents the full write to a concurrent set, so another overlay of that column no longer commutes with it. The pair tests here build their sides with the shared `footprint` fixture. Co-Authored-By: Claude Fable 5.1 --- .../src/transaction/action/add_overlays.rs | 67 ++++++++++++++++--- .../transaction/action/assert_unique_keys.rs | 12 +--- 2 files changed, 60 insertions(+), 19 deletions(-) diff --git a/rust/lance-table/src/transaction/action/add_overlays.rs b/rust/lance-table/src/transaction/action/add_overlays.rs index 95edd8869d0..f2ac57a7504 100644 --- a/rust/lance-table/src/transaction/action/add_overlays.rs +++ b/rust/lance-table/src/transaction/action/add_overlays.rs @@ -118,10 +118,10 @@ mod tests { use super::*; use crate::format::DataFile; use crate::format::overlay::OverlayCoverage; - use crate::transaction::action::test_support::{apply, backed_manifest}; + use crate::transaction::action::test_support::{apply, backed_manifest, footprint}; use crate::transaction::action::{ - Action, AddFragment, AlterField, CompositeOperation, DropField, RemoveFragment, - TombstoneFieldData, UserAction, + Action, AddFragment, AddIndexSegment, AlterField, DropField, RemoveFragment, + TombstoneFieldData, }; use lance_file::version::ConcreteFileVersion; use roaring::RoaringBitmap; @@ -153,12 +153,6 @@ mod tests { }) } - fn footprint(actions: Vec) -> Footprint { - Footprint::from(&CompositeOperation::new(vec![UserAction::new( - "step", actions, - )])) - } - #[test] fn test_overlays_are_stamped_with_the_version_the_commit_produces() { let manifest = backed_manifest(); @@ -307,6 +301,61 @@ mod tests { assert!(!dropped.conflicts_with(&overlaid)); } + /// An index built over a column an overlay landed on first is accepted: + /// the segment is stamped with the version it read, the overlay carries + /// the later version it committed at, and the read path masks the + /// overlaid cells out of any segment older than the overlay + /// (`collect_overlay_stale_frags`). A requirement therefore does not + /// collide with a committed partial write, only with a full one. + #[test] + fn test_an_index_may_be_built_over_a_committed_overlay() { + let build = footprint(vec![Action::AddIndexSegment(AddIndexSegment { + uuid: uuid::Uuid::from_u128(1), + name: "by_a".into(), + fields: vec![Ref::Committed(0)], + covering_fields: Vec::new(), + index_details: None, + index_version: 1, + covered_fragments: Some(vec![Ref::Committed(0)]), + files: Vec::new(), + base: None, + created_at: None, + dataset_version: None, + data_change: false, + })]); + let overlaid = footprint(vec![add(Ref::Committed(0), vec![overlay("a.lance", &[0])])]); + + assert!(!build.conflicts_with(&overlaid)); + assert!(!overlaid.conflicts_with(&build)); + } + + /// One set may overlay a cell and then rewrite the whole column in the + /// same operation (an overlay being materialized). The full write is what + /// a concurrent set sees: another overlay of the column no longer + /// commutes with it, while a rewrite of a different column still does. + #[test] + fn test_a_full_write_in_the_same_set_dominates_a_partial_one() { + let materialize = footprint(vec![ + add(Ref::Committed(0), vec![overlay("a.lance", &[0])]), + Action::TombstoneFieldData(TombstoneFieldData { + fragment: Ref::Committed(0), + field_ids: vec![Ref::Committed(0)], + data_change: true, + }), + ]); + let other_overlay = footprint(vec![add(Ref::Committed(0), vec![overlay("b.lance", &[0])])]); + let other_column = footprint(vec![Action::TombstoneFieldData(TombstoneFieldData { + fragment: Ref::Committed(0), + field_ids: vec![Ref::Committed(1)], + data_change: true, + })]); + + assert!(materialize.conflicts_with(&other_overlay)); + assert!(other_overlay.conflicts_with(&materialize)); + assert!(!materialize.conflicts_with(&other_column)); + assert!(!other_column.conflicts_with(&materialize)); + } + #[test] fn test_overlaying_a_fragment_a_concurrent_writer_removes_conflicts() { let ours = footprint(vec![add(Ref::Committed(0), vec![overlay("a.lance", &[0])])]); diff --git a/rust/lance-table/src/transaction/action/assert_unique_keys.rs b/rust/lance-table/src/transaction/action/assert_unique_keys.rs index 290758584c6..3f643776c8e 100644 --- a/rust/lance-table/src/transaction/action/assert_unique_keys.rs +++ b/rust/lance-table/src/transaction/action/assert_unique_keys.rs @@ -104,10 +104,8 @@ impl TryFrom for AssertUniqueKeys { mod tests { use super::*; use crate::format::key_existence::FilterType; - use crate::transaction::action::test_support::{apply, backed_manifest}; - use crate::transaction::action::{ - Action, AddFragment, CompositeOperation, RemoveFragment, TombstoneFieldData, UserAction, - }; + use crate::transaction::action::test_support::{apply, backed_manifest, footprint}; + use crate::transaction::action::{Action, AddFragment, RemoveFragment, TombstoneFieldData}; fn assertion(key_fields: Vec, hashes: &[u64]) -> Action { Action::AssertUniqueKeys(AssertUniqueKeys { @@ -138,12 +136,6 @@ mod tests { Action::AddFragment(fragment) } - fn footprint(actions: Vec) -> Footprint { - Footprint::from(&CompositeOperation::new(vec![UserAction::new( - "step", actions, - )])) - } - #[test] fn test_an_assertion_changes_nothing() { let manifest = backed_manifest();