diff --git a/rust/lance-table/src/transaction/action.rs b/rust/lance-table/src/transaction/action.rs index 95b1136a44e..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 @@ -66,6 +67,9 @@ macro_rules! for_each_action { TombstoneFieldData, RemoveFragment, SetDeletionFile, + AddOverlays, + RefreshRowVersionMetadata, + UpdateCompactedSsTables, AlterField, DropField, AddIndexSegment, @@ -75,6 +79,7 @@ macro_rules! for_each_action { ReserveRowIds, ResetTable, ConfigUpdate, + AssertUniqueKeys, } }; } @@ -84,13 +89,16 @@ 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; +mod assert_unique_keys; 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; @@ -98,6 +106,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; @@ -107,11 +116,14 @@ 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 assert_unique_keys::AssertUniqueKeys; 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; @@ -119,6 +131,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; @@ -266,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/add_overlays.rs b/rust/lance-table/src/transaction/action/add_overlays.rs new file mode 100644 index 00000000000..f2ac57a7504 --- /dev/null +++ b/rust/lance-table/src/transaction/action/add_overlays.rs @@ -0,0 +1,370 @@ +// 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::{Coordinate, 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 + } + + /// 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) { + 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)); + } + } +} + +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, footprint}; + use crate::transaction::action::{ + Action, AddFragment, AddIndexSegment, AlterField, DropField, RemoveFragment, + TombstoneFieldData, + }; + 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], + ConcreteFileVersion::V2_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, + }) + } + + #[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)); + } + + /// 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)); + } + + /// 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])])]); + 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..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}; @@ -280,6 +282,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) @@ -308,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. @@ -338,6 +359,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/assert_unique_keys.rs b/rust/lance-table/src/transaction/action/assert_unique_keys.rs new file mode 100644 index 00000000000..3f643776c8e --- /dev/null +++ b/rust/lance-table/src/transaction/action/assert_unique_keys.rs @@ -0,0 +1,289 @@ +// 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. +/// +/// 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 + /// 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, footprint}; + use crate::transaction::action::{Action, AddFragment, RemoveFragment, TombstoneFieldData}; + + 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) + } + + #[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)); + } + + /// 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])]); + 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 1f7e41bd6a3..4d8fb2be372 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; @@ -33,6 +34,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 +88,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(_) @@ -123,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, } @@ -132,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), @@ -152,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, } } @@ -200,7 +218,7 @@ pub struct Footprint { /// /// 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 @@ -208,12 +226,34 @@ 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. 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). + 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 @@ -284,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 @@ -299,8 +339,15 @@ 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. + // 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 @@ -352,6 +399,71 @@ impl Footprint { .any(|fragment| other.regions.contains_key(&Region::Fragment(*fragment))) } + /// Whether `other` may have introduced a key this set asserts is not there. + /// + /// 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() { + return false; + } + 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 { + // 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 + // built with parameters that cannot be compared. Neither + // clears the assertion. + Ok((true, _)) | Err(_) => return true, + } + } + } + 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. @@ -367,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) { @@ -480,6 +598,27 @@ 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, + }); + } + + /// 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 7a5506e3268..d5c5b00fb0a 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 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. 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,21 +154,26 @@ 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, }; 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, AdjustIndexCoverage, - AlterField, ConfigUpdate, DropField, FieldMetadataUpdate, RemoveFragment, - RemoveIndexSegment, ReserveFragmentIds, ReserveRowIds, ResetTable, SetDeletionFile, - TombstoneFieldData, + AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, AddOverlays, + 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; 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 +229,20 @@ 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::RefreshRowVersionMetadata(RefreshRowVersionMetadata { + fragment_ids: vec![4, 6], + }), Action::AlterField(AlterField { field: Ref::Committed(2), name: Some("renamed".into()), @@ -263,6 +283,19 @@ 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::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), @@ -313,30 +346,16 @@ 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], - }, - )), - }; - 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/action/refresh_row_version_metadata.rs b/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs new file mode 100644 index 00000000000..51428c5eefe --- /dev/null +++ b/rust/lance-table/src/transaction/action/refresh_row_version_metadata.rs @@ -0,0 +1,203 @@ +// 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 + } + + /// 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)); + } + } +} + +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 lance_file::version::ConcreteFileVersion; + 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], + ConcreteFileVersion::V2_0, + None, + None, + )], + overlays: vec![], + deletion_file: None, + 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, + }; + 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/action/update_compacted_sstables.rs b/rust/lance-table/src/transaction/action/update_compacted_sstables.rs new file mode 100644 index 00000000000..bfc21cd18c4 --- /dev/null +++ b/rust/lance-table/src/transaction/action/update_compacted_sstables.rs @@ -0,0 +1,257 @@ +// 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. 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, +/// 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::{ + 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, + }; + 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 + } + + /// 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_without_a_mem_wal_index_is_rejected() { + let error = apply_with_indices(&backed_manifest(), vec![update(vec![(1, 7)])], Vec::new()) + .unwrap_err(); + + 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![empty_mem_wal_index()], + ) + .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_is_rejected_rather_than_walking_a_shard_backwards() { + let manifest = backed_manifest(); + 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!( + 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![empty_mem_wal_index()], + ) + .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![empty_mem_wal_index()], + ) + .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(), empty_mem_wal_index()], + ) + .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)); + } +} diff --git a/rust/lance-table/src/transaction/proto.rs b/rust/lance-table/src/transaction/proto.rs index 0ce47f74334..b96099c4559 100644 --- a/rust/lance-table/src/transaction/proto.rs +++ b/rust/lance-table/src/transaction/proto.rs @@ -907,26 +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: "refresh row versions".to_string(), - actions: vec![pb::Action { - action: Some(pb::action::Action::RefreshRowVersionMetadata( - pb::RefreshRowVersionMetadata { - fragment_ids: vec![1], - }, - )), - }], + description: "something from the future".to_string(), + actions: vec![pb::Action { action: None }], }], }, )), diff --git a/rust/lance/src/dataset/tests/dataset_transactions.rs b/rust/lance/src/dataset/tests/dataset_transactions.rs index 2e677c6db5b..384d817f4bf 100644 --- a/rust/lance/src/dataset/tests/dataset_transactions.rs +++ b/rust/lance/src/dataset/tests/dataset_transactions.rs @@ -1747,15 +1747,22 @@ 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, MemWalIndexDetails, load_mem_wal_index_details, + new_mem_wal_index_meta, + }; 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 +2858,178 @@ 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(); + + // 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 { + 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 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); + 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); + } }