diff --git a/rust/lance-table/src/transaction/action.rs b/rust/lance-table/src/transaction/action.rs index 836d5b6a92d..95b1136a44e 100644 --- a/rust/lance-table/src/transaction/action.rs +++ b/rust/lance-table/src/transaction/action.rs @@ -68,6 +68,9 @@ macro_rules! for_each_action { SetDeletionFile, AlterField, DropField, + AddIndexSegment, + RemoveIndexSegment, + AdjustIndexCoverage, ReserveFragmentIds, ReserveRowIds, ResetTable, @@ -80,6 +83,8 @@ mod add_base; mod add_data_file; mod add_field; mod add_fragment; +mod add_index_segment; +mod adjust_index_coverage; mod alter_field; mod apply; mod config_update; @@ -87,6 +92,7 @@ mod drop_field; mod footprint; mod proto; mod remove_fragment; +mod remove_index_segment; mod reserve_fragment_ids; mod reserve_row_ids; mod reset_table; @@ -100,11 +106,14 @@ pub use add_base::AddBase; pub use add_data_file::AddDataFile; pub use add_field::AddField; pub use add_fragment::AddFragment; +pub use add_index_segment::AddIndexSegment; +pub use adjust_index_coverage::AdjustIndexCoverage; pub use alter_field::AlterField; pub use config_update::{ConfigUpdate, FieldMetadataUpdate}; pub use drop_field::DropField; pub use footprint::{ConfigMap, Coordinate, Footprint}; pub use remove_fragment::RemoveFragment; +pub use remove_index_segment::RemoveIndexSegment; pub use reserve_fragment_ids::ReserveFragmentIds; pub use reserve_row_ids::ReserveRowIds; pub use reset_table::ResetTable; diff --git a/rust/lance-table/src/transaction/action/add_index_segment.rs b/rust/lance-table/src/transaction/action/add_index_segment.rs new file mode 100644 index 00000000000..2db1b7c0ede --- /dev/null +++ b/rust/lance-table/src/transaction/action/add_index_segment.rs @@ -0,0 +1,914 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Add an index segment. + +use super::apply::ApplyState; +use super::footprint::IndexIdentity; +use super::proto::{data_change_from_wire, data_change_to_wire, non_empty, required}; +use super::{Footprint, Ref}; +use crate::format::{IndexFile, IndexMetadata, pb}; +use lance_core::deepsize::{Context, DeepSizeOf}; +use lance_core::{Error, Result}; +use roaring::RoaringBitmap; +use std::sync::Arc; +use uuid::Uuid; + +/// Add an index segment. +/// +/// The format has no first-class "index" apart from its segments: a logical +/// index is the set of segments sharing a `name`, and a brand-new index is +/// written as its first segment. So this one action covers both creating an +/// index and extending an existing one -- the difference is only whether any +/// segment already carries the name. +/// +/// The segment's `uuid` is chosen by the writer rather than minted from a +/// counter, so unlike a fragment or a field it needs no [`Ref`]: two writers +/// cannot pick the same one, and replaying the action onto another version +/// leaves it untouched. +#[derive(Debug, Clone, PartialEq)] +pub struct AddIndexSegment { + /// Identifies the segment, and names the directory its files live in. + pub uuid: Uuid, + /// The logical index this segment belongs to. + pub name: String, + /// The fields the segment is keyed on -- the columns it can be searched + /// by. Empty only for a system index keyed on no column at all, such as + /// `mem_wal` or `frag_reuse`. + /// + /// Keys only, unlike [`IndexMetadata::fields`] under the legacy contract, + /// where it means keyed columns followed by carried ones. `apply` derives + /// that representation. + pub fields: Vec, + /// The fields whose values the segment carries but is not necessarily + /// keyed on, independent of [`Self::fields`] rather than a subset of it. + /// Empty for a segment that carries no extra column. + /// + /// A field in both lists is one the segment is keyed on *and* carries. + /// That needs `FLAG_INDEPENDENT_COVERING_FIELDS`, which no release + /// implements, so `apply` rejects it for now. + pub covering_fields: Vec, + /// Index-type-specific metadata, opaque to the transaction layer. + pub index_details: Option>, + pub index_version: i32, + /// The fragments this segment covers, or `None` when the coverage was not + /// recorded -- what the system indices carry. An empty list is the + /// different statement that the segment covers nothing. + pub covered_fragments: Option>, + /// The segment's files and their sizes, empty when the writer did not + /// record them. + pub files: Vec, + /// The base path the files live under, for a segment imported from another + /// dataset. `None` means the dataset's own index directory. + pub base: Option, + /// When the segment was built, or `None` when the writer did not record + /// it. Carried rather than derived, because it describes a build that + /// replaying this action does not redo. + pub created_at: Option>, + /// The dataset version whose data this segment reflects, or `None` for the + /// version the operation reads -- what a freshly built segment reflects. + /// + /// A segment merged from older ones reflects only as much as its oldest + /// input, so this is not always the read version and is not derivable from + /// it. The overlay version gate reads it: an overlay committed at or before + /// it counts as already folded into the index. + pub dataset_version: Option, + pub data_change: bool, +} + +impl AddIndexSegment { + pub(super) fn apply(&self, state: &mut ApplyState) -> Result<()> { + let fields = self + .fields + .iter() + .map(|field| state.resolve_field(*field)) + .collect::>>()?; + let fragment_bitmap = self + .covered_fragments + .as_ref() + .map(|fragments| self.resolve_coverage(fragments, state)) + .transpose()?; + let base_id = self.base.map(|base| state.resolve_base(base)).transpose()?; + + let covering_fields = self + .covering_fields + .iter() + .map(|field| state.resolve_field(*field)) + .collect::>>()?; + let (fields, covering_fields) = Self::lower_covering(&self.name, fields, covering_fields)?; + + let metadata = IndexMetadata { + uuid: self.uuid, + name: self.name.clone(), + fields, + covering_fields, + dataset_version: self.reflected_version(state.read_version())?, + fragment_bitmap, + index_details: self.index_details.clone(), + index_version: self.index_version, + created_at: self.created_at, + base_id, + files: (!self.files.is_empty()).then(|| self.files.clone()), + }; + // Belt and braces: `lower_covering` builds the legacy form by + // construction, so a failure here means the two have drifted apart. + metadata.validate_covering_fields()?; + + state.add_index_segment(metadata) + } + + /// Render the action's independent key/covering declaration as the + /// `IndexMetadata` form a manifest can carry today. + /// + /// The action names keys and carried columns separately (the contract in + /// lance-format/lance#9159), while a manifest without + /// `FLAG_INDEPENDENT_COVERING_FIELDS` carries the legacy form: one + /// `fields` list of keyed columns followed by carried ones, with + /// `covering_fields` naming that trailing subset. + /// + /// A disjoint declaration converts exactly. An overlapping one -- a column + /// both keyed and carried -- has no legacy representation and is rejected + /// rather than silently flattened, which would drop the fact that the + /// column is carried. Both this lowering and the rejection come out when a + /// release implements the flag; the action's own shape does not change. + fn lower_covering( + name: &str, + keys: Vec, + covering: Vec, + ) -> Result<(Vec, Vec)> { + if covering.is_empty() { + return Ok((keys, covering)); + } + + let overlapping: Vec = covering + .iter() + .copied() + .filter(|field| keys.contains(field)) + .collect(); + if !overlapping.is_empty() { + return Err(Error::not_supported(format!( + "index '{name}' declares fields {overlapping:?} as both keyed and \ + covering; that requires FLAG_INDEPENDENT_COVERING_FIELDS, which no \ + release implements yet (keys {keys:?}, covering {covering:?})" + ))); + } + + let mut lowered = keys; + lowered.extend_from_slice(&covering); + Ok((lowered, covering)) + } + + /// The version this segment reflects, defaulting to the one the operation + /// reads. + fn reflected_version(&self, read_version: u64) -> Result { + let Some(version) = self.dataset_version else { + return Ok(read_version); + }; + if version > read_version { + return Err(Error::invalid_input(format!( + "AddIndexSegment for index '{}' reflects dataset version {version}, which is \ + newer than the version {read_version} this operation reads; a segment cannot \ + reflect data it could not have seen", + self.name + ))); + } + Ok(version) + } + + fn resolve_coverage(&self, fragments: &[Ref], state: &ApplyState) -> Result { + fragments + .iter() + .map(|fragment| { + let id = state.resolve_fragment(*fragment)?; + u32::try_from(id).map_err(|_| { + Error::invalid_input(format!( + "AddIndexSegment for index '{}' covers fragment {id}, which is beyond the \ + largest fragment id an index can record ({})", + self.name, + u32::MAX + )) + }) + }) + .collect() + } + + /// Left to the writer. An index is derived state, so building one normally + /// changes nothing a reader sees -- but the MemWAL index holds real + /// unflushed rows, so the answer is not the same for every segment. + pub(super) fn is_data_change(&self) -> bool { + self.data_change + } + + /// A claim on the logical index -- this is what the index is, and these are + /// the fragments the new segment describes -- plus a requirement on the + /// data the segment was built from. + /// + /// An index is derived state whose segments the query path unions, so two + /// writers may extend the same index at once. What they may not do is + /// describe the same fragment twice, which would double-count rows, or + /// disagree about what the index is. + /// + /// A segment that has fallen behind is pruned rather than being wrong, but + /// only the commit that does the changing can prune: a rewrite rebinds the + /// field and drops it from every segment the manifest carries. A segment + /// arriving *after* the rewrite prunes nothing, and the format gives a + /// reader no way to catch it -- an in-place column rewrite "cannot be + /// detected just by examining metadata", so a reader trusts + /// `fragment_bitmap` as it finds it. Requiring the data keeps that trust + /// earned. Deletions and overlays are excluded on purpose: both leave a + /// trace a reader can act on, so both may still run alongside a build. + pub(super) fn footprint(&self, footprint: &mut Footprint) { + footprint.build_index( + self.name.clone(), + IndexIdentity { + fields: self.fields.clone(), + covering_fields: self.covering_fields.clone(), + details: self.index_details.clone(), + index_version: self.index_version, + }, + self.covered_fragments.clone(), + ); + + // Keyed and carried columns alike: a carried column can be rewritten + // while the keyed one is untouched, and the segment would then answer + // from an obsolete carried value. + let dependencies: Vec = self + .fields + .iter() + .chain(self.covering_fields.iter()) + .copied() + .collect(); + // The segment describes values in the type the schema named when it + // was built. A cast landing first would leave it describing values as + // something they are not, and a drop would leave it over a field the + // manifest no longer has; both write the definition. Landing second, + // either prunes or discards the segment as it applies, so this is only + // required, not written. + for field in dependencies.iter().copied() { + footprint.require_field_definition(field); + } + let Some(covered) = &self.covered_fragments else { + // Unstated reach. The claim already collides with every other claim + // on the index; there is no fragment list to require. + return; + }; + for fragment in covered.iter().filter_map(|fragment| fragment.committed()) { + footprint.require_field_data(fragment, dependencies.iter().copied()); + } + } +} + +impl DeepSizeOf for AddIndexSegment { + fn deep_size_of_children(&self, context: &mut Context) -> usize { + // `prost_types::Any` does not implement DeepSizeOf, but it is only a + // type url and a byte string, so both are measured directly. + let index_details = self.index_details.as_ref().map_or(0, |details| { + std::mem::size_of::() + + details.type_url.capacity() + + details.value.capacity() + }); + self.uuid.as_bytes().deep_size_of_children(context) + + self.name.deep_size_of_children(context) + + self.fields.deep_size_of_children(context) + + self.covering_fields.deep_size_of_children(context) + + index_details + + self.covered_fragments.deep_size_of_children(context) + + self.files.deep_size_of_children(context) + } +} + +impl From<&AddIndexSegment> for pb::AddIndexSegment { + fn from(value: &AddIndexSegment) -> Self { + Self { + uuid: Some((&value.uuid).into()), + name: value.name.clone(), + fields: value.fields.iter().map(|field| (*field).into()).collect(), + covering_fields: value + .covering_fields + .iter() + .map(|field| (*field).into()) + .collect(), + index_details: value + .index_details + .as_ref() + .map(|details| details.as_ref().clone()), + index_version: Some(value.index_version), + covered_fragments: value.covered_fragments.as_ref().map(|fragments| { + pb::FragmentCoverage { + fragments: fragments.iter().map(|id| (*id).into()).collect(), + } + }), + files: value + .files + .iter() + .map(|file| pb::IndexFile { + path: file.path.clone(), + size_bytes: file.size_bytes, + }) + .collect(), + data_change: data_change_to_wire(value.data_change), + base: value.base.map(pb::Ref::from), + created_at: value + .created_at + .map(|created_at| created_at.timestamp_millis() as u64), + dataset_version: value.dataset_version, + } + } +} + +impl TryFrom for AddIndexSegment { + type Error = Error; + + fn try_from(message: pb::AddIndexSegment) -> Result { + let created_at = message + .created_at + .map(|millis| { + chrono::DateTime::from_timestamp_millis(millis as i64).ok_or_else(|| { + Error::invalid_input(format!( + "AddIndexSegment.created_at is {millis}ms since the epoch, which is not a \ + representable timestamp" + )) + }) + }) + .transpose()?; + + Ok(Self { + uuid: Uuid::try_from(&required(message.uuid, "AddIndexSegment.uuid")?)?, + name: non_empty(message.name, "AddIndexSegment.name")?, + fields: message + .fields + .into_iter() + .map(Ref::try_from) + .collect::>>()?, + covering_fields: message + .covering_fields + .into_iter() + .map(Ref::try_from) + .collect::>>()?, + index_details: message.index_details.map(Arc::new), + index_version: message.index_version.unwrap_or_default(), + covered_fragments: message + .covered_fragments + .map(|coverage| { + coverage + .fragments + .into_iter() + .map(Ref::try_from) + .collect::>>() + }) + .transpose()?, + files: message + .files + .into_iter() + .map(|file| IndexFile { + path: file.path, + size_bytes: file.size_bytes, + }) + .collect(), + base: message.base.map(Ref::try_from).transpose()?, + created_at, + dataset_version: message.dataset_version, + data_change: data_change_from_wire(message.data_change), + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::transaction::action::test_support::{ + added_field, apply_with_indices, backed_manifest, footprint, + }; + use crate::transaction::action::{ + Action, AddField, AddFragment, AlterField, CompositeOperation, DropField, Footprint, + RemoveFragment, TombstoneFieldData, UserAction, + }; + use crate::transaction::test_support::{default_build_config, sample_index_metadata}; + use crate::transaction::{Operation, Transaction}; + use rstest::rstest; + + fn segment(name: &str, fields: Vec) -> AddIndexSegment { + AddIndexSegment { + uuid: Uuid::new_v4(), + name: name.into(), + fields, + 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, + } + } + + #[test] + fn test_add_index_segment_records_coverage_and_the_version_it_describes() { + let manifest = backed_manifest(); + let action = AddIndexSegment { + files: vec![IndexFile { + path: "index.idx".into(), + size_bytes: 1024, + }], + ..segment("by_a", vec![Ref::Committed(0)]) + }; + + let (next, indices) = apply_with_indices( + &manifest, + vec![Action::AddIndexSegment(action.clone())], + Vec::new(), + ) + .unwrap(); + + assert_eq!(indices.len(), 1); + let index = &indices[0]; + assert_eq!(index.uuid, action.uuid); + assert_eq!(index.name, "by_a"); + assert_eq!(index.fields, vec![0]); + assert_eq!(index.fragment_bitmap, Some([0].into_iter().collect())); + assert_eq!(index.files, Some(action.files)); + // A segment that does not say what it reflects reflects the version the + // operation read, so replaying it elsewhere restamps it. + assert_eq!(index.dataset_version, manifest.version); + assert!(index.dataset_version < next.version); + } + + /// A segment that does not say what it reflects reflects the version its + /// writer read -- not whatever version the set ends up landing on. The two + /// are the same only when the commit is uncontended; after a lost race the + /// manifest has moved on while the segment has not. + #[test] + fn test_a_relocated_segment_reflects_the_version_it_was_built_against() { + // Stand in for a retry: the set was built against `read_version` and is + // being applied to a manifest two commits further along. + let mut manifest = backed_manifest(); + manifest.version += 2; + let read_version = manifest.version - 2; + + let transaction = Transaction::new( + read_version, + Operation::CompositeOperation(CompositeOperation::new(vec![UserAction::new( + "step", + vec![Action::AddIndexSegment(segment( + "by_a", + vec![Ref::Committed(0)], + ))], + )])), + None, + ); + let (_, indices) = transaction + .build_manifest( + Some(&manifest), + Vec::new(), + "tx.txn", + &default_build_config(), + ) + .unwrap(); + + assert_eq!(indices[0].dataset_version, read_version); + } + + /// The bound on an explicit `dataset_version` is the version the set read, + /// for the same reason: a segment cannot reflect data its writer could not + /// have seen, however far the manifest has moved since. + #[test] + fn test_a_segment_cannot_reflect_a_version_committed_after_the_one_it_read() { + let mut manifest = backed_manifest(); + manifest.version += 2; + let read_version = manifest.version - 2; + + let transaction = Transaction::new( + read_version, + Operation::CompositeOperation(CompositeOperation::new(vec![UserAction::new( + "step", + vec![Action::AddIndexSegment(AddIndexSegment { + dataset_version: Some(read_version + 1), + ..segment("by_a", vec![Ref::Committed(0)]) + })], + )])), + None, + ); + let error = transaction + .build_manifest( + Some(&manifest), + Vec::new(), + "tx.txn", + &default_build_config(), + ) + .unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!(error.to_string().contains("could not have seen"), "{error}"); + } + + #[test] + fn test_a_merged_segment_keeps_the_older_version_it_reflects() { + let manifest = backed_manifest(); + let (_, indices) = apply_with_indices( + &manifest, + vec![Action::AddIndexSegment(AddIndexSegment { + dataset_version: Some(manifest.version - 1), + ..segment("by_a", vec![Ref::Committed(0)]) + })], + Vec::new(), + ) + .unwrap(); + + assert_eq!(indices[0].dataset_version, manifest.version - 1); + } + + #[test] + fn test_a_segment_reflecting_a_future_version_is_rejected() { + let manifest = backed_manifest(); + let error = apply_with_indices( + &manifest, + vec![Action::AddIndexSegment(AddIndexSegment { + dataset_version: Some(manifest.version + 1), + ..segment("by_a", vec![Ref::Committed(0)]) + })], + Vec::new(), + ) + .unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains("could not have seen"), + "unexpected error: {error}" + ); + } + + /// The action names keys and carried columns independently, but a manifest + /// can only carry the legacy form -- one `fields` list with the carried + /// columns as its tail -- until a release implements + /// `FLAG_INDEPENDENT_COVERING_FIELDS`. Apply converts between the two. + #[test] + fn test_a_disjoint_covering_declaration_lowers_to_the_legacy_form() { + let (next, indices) = apply_with_indices( + &backed_manifest(), + vec![ + Action::AddField(AddField { + local: 0, + parent: None, + def: added_field("carried"), + }), + Action::AddIndexSegment(AddIndexSegment { + covering_fields: vec![Ref::Local(0)], + ..segment("by_a", vec![Ref::Committed(0)]) + }), + ], + Vec::new(), + ) + .unwrap(); + + let carried = next.schema.field("carried").unwrap().id; + assert_eq!(indices[0].fields, vec![0, carried]); + assert_eq!(indices[0].covering_fields, vec![carried]); + } + + /// A column both keyed and carried has no legacy representation: dropping + /// it from `covering_fields` would leave the manifest claiming only that + /// the index is keyed on the column, losing the fact that it also serves + /// the column's values. Reject rather than publish the weaker claim. + #[test] + fn test_a_covering_field_that_is_also_a_key_is_rejected_for_now() { + let error = apply_with_indices( + &backed_manifest(), + vec![Action::AddIndexSegment(AddIndexSegment { + covering_fields: vec![Ref::Committed(0)], + ..segment("by_a", vec![Ref::Committed(0)]) + })], + Vec::new(), + ) + .unwrap_err(); + + assert!(matches!(error, Error::NotSupported { .. }), "{error:?}"); + assert!( + error + .to_string() + .contains("FLAG_INDEPENDENT_COVERING_FIELDS"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_a_second_segment_extends_the_same_index() { + let existing = sample_index_metadata("by_a"); + let added = segment("by_a", vec![Ref::Committed(0)]); + + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![Action::AddIndexSegment(added.clone())], + vec![existing.clone()], + ) + .unwrap(); + + // A logical index is the set of segments sharing a name, so adding one + // leaves the other in place rather than replacing it. + assert_eq!(indices.len(), 2); + let uuids = indices.iter().map(|index| index.uuid).collect::>(); + assert!(uuids.contains(&existing.uuid)); + assert!(uuids.contains(&added.uuid)); + } + + #[test] + fn test_re_adding_a_segment_is_rejected() { + let existing = sample_index_metadata("by_a"); + let error = apply_with_indices( + &backed_manifest(), + vec![Action::AddIndexSegment(AddIndexSegment { + uuid: existing.uuid, + ..segment("by_a", vec![Ref::Committed(0)]) + })], + vec![existing], + ) + .unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains("added once"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_unrecorded_coverage_is_not_the_same_as_covering_nothing() { + let unknown = AddIndexSegment { + covered_fragments: None, + ..segment("system", vec![]) + }; + let empty = AddIndexSegment { + covered_fragments: Some(Vec::new()), + ..segment("empty", vec![]) + }; + + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![ + Action::AddIndexSegment(unknown), + Action::AddIndexSegment(empty), + ], + Vec::new(), + ) + .unwrap(); + + let bitmap_of = |name: &str| { + indices + .iter() + .find(|index| index.name == name) + .unwrap() + .fragment_bitmap + .clone() + }; + assert_eq!(bitmap_of("system"), None); + assert_eq!(bitmap_of("empty"), Some(RoaringBitmap::new())); + } + + #[test] + fn test_a_segment_can_cover_a_fragment_minted_in_the_same_operation() { + // Indexing what an operation just wrote is the point of composing the + // two, and the fragment has no committed id until this apply runs. + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![ + Action::AddFragment(AddFragment { + id: Ref::Local(0), + physical_rows: 10, + row_id_meta: None, + last_updated_at_version_meta: None, + created_at_version_meta: None, + data_change: true, + }), + Action::AddIndexSegment(AddIndexSegment { + covered_fragments: Some(vec![Ref::Committed(0), Ref::Local(0)]), + ..segment("by_a", vec![Ref::Committed(0)]) + }), + ], + Vec::new(), + ) + .unwrap(); + + assert_eq!( + indices[0].fragment_bitmap, + Some([0, 1].into_iter().collect()) + ); + } + + #[test] + fn test_a_segment_can_index_a_field_minted_in_the_same_operation() { + let (next, indices) = apply_with_indices( + &backed_manifest(), + vec![ + Action::AddField(AddField { + local: 0, + parent: None, + def: added_field("added"), + }), + Action::AddIndexSegment(segment("by_added", vec![Ref::Local(0)])), + ], + Vec::new(), + ) + .unwrap(); + + let field = next.schema.field("added").unwrap(); + assert_eq!(indices[0].fields, vec![field.id]); + } + + #[test] + fn test_a_segment_over_a_dropped_field_does_not_survive_the_commit() { + // The assembly prunes indices whose fields left the schema, so an + // operation that drops a field and indexes it cannot smuggle one in. + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![ + Action::DropField(DropField { + field: Ref::Committed(0), + }), + Action::AddIndexSegment(segment("by_a", vec![Ref::Committed(0)])), + ], + Vec::new(), + ) + .unwrap(); + + assert!(indices.is_empty()); + } + + /// Two builds over one column do not collide: a requirement is a claim + /// that nothing moved, not a claim on the column. + #[test] + fn test_two_indices_over_one_column_and_fragment_do_not_conflict() { + let ours = index_footprint(covering("by_a", Some(vec![0]))); + let theirs = index_footprint(covering("by_b", Some(vec![0]))); + + assert!(!ours.conflicts_with(&theirs)); + assert!(!theirs.conflicts_with(&ours)); + } + + /// The direction is the whole point, so it gets its own test. + /// + /// A segment requires the data it describes; a rewrite of that data writes + /// it. Arriving after the rewrite, the segment would publish coverage of + /// values it never saw and a reader has no way to detect that, so it must + /// be rejected. Arriving before, the rewrite prunes the segment's coverage + /// as it applies, so rejecting it would be a conflict over nothing. + #[test] + fn test_a_build_loses_to_a_committed_rewrite_but_not_the_reverse() { + let build = index_footprint(covering("by_a", Some(vec![0]))); + let rewrite = footprint(vec![Action::TombstoneFieldData(TombstoneFieldData { + fragment: Ref::Committed(0), + field_ids: vec![Ref::Committed(0)], + data_change: true, + })]); + + assert!( + build.conflicts_with(&rewrite), + "a segment cannot describe values a committed rewrite replaced" + ); + assert!( + !rewrite.conflicts_with(&build), + "the rewrite prunes the committed segment's coverage as it applies" + ); + } + + /// A rewrite of a column the segment merely carries invalidates it just as + /// a keyed one does -- the segment would answer from an obsolete carried + /// value. `covering_fields` is independent of `fields`, so the requirement + /// is over the union. + #[test] + fn test_a_build_requires_the_columns_it_carries_too() { + let mut segment = covering("by_a", Some(vec![0])); + segment.covering_fields = vec![Ref::Committed(1)]; + let build = index_footprint(segment); + + let rewrite = footprint(vec![Action::TombstoneFieldData(TombstoneFieldData { + fragment: Ref::Committed(0), + field_ids: vec![Ref::Committed(1)], + data_change: true, + })]); + + assert!(build.conflicts_with(&rewrite)); + } + + /// A cast landing first leaves the segment describing values in a type the + /// schema no longer names; a cast landing second rebinds the field in every + /// fragment and prunes the segment's coverage as it applies. Same shape as + /// the rewrite case above, one level up: the definition rather than the + /// data. + #[test] + fn test_a_build_loses_to_a_committed_cast_but_not_the_reverse() { + let build = index_footprint(covering("by_a", Some(vec![0]))); + let cast = footprint(vec![Action::AlterField(AlterField { + field: Ref::Committed(0), + name: None, + logical_type: Some("int64".into()), + nullable: None, + })]); + + assert!(build.conflicts_with(&cast)); + assert!(!cast.conflicts_with(&build)); + + // A segment of unstated reach depends on the definition just the same. + assert!(index_footprint(covering("by_a", None)).conflicts_with(&cast)); + } + + /// A segment requires the data of the fragments it covers, and that data + /// is gone if the fragment is: a compaction landing first would leave the + /// segment describing rows the dataset no longer has, in a fragment the + /// replacement does not cover. Landing second, the compaction is the one + /// answerable for the index, and the footprint has nothing to say. + #[test] + fn test_a_build_loses_to_a_committed_removal_of_a_covered_fragment() { + let build = index_footprint(covering("by_a", Some(vec![0]))); + let removal = |fragment: u64| { + footprint(vec![Action::RemoveFragment(RemoveFragment { + fragment: Ref::Committed(fragment), + data_change: false, + })]) + }; + + assert!(build.conflicts_with(&removal(0))); + assert!(!removal(0).conflicts_with(&build)); + assert!(!build.conflicts_with(&removal(1))); + } + + fn index_footprint(action: AddIndexSegment) -> Footprint { + footprint(vec![Action::AddIndexSegment(action)]) + } + + fn covering(name: &str, fragments: Option>) -> AddIndexSegment { + AddIndexSegment { + covered_fragments: fragments + .map(|ids| ids.into_iter().map(Ref::Committed).collect::>()), + ..segment(name, vec![Ref::Committed(0)]) + } + } + + #[rstest] + #[case::disjoint_coverage_of_one_index( + covering("by_a", Some(vec![0, 1])), + covering("by_a", Some(vec![2, 3])), + false, + )] + #[case::overlapping_coverage_of_one_index( + covering("by_a", Some(vec![0, 1])), + covering("by_a", Some(vec![1, 2])), + true, + )] + #[case::different_indices( + covering("by_a", Some(vec![0, 1])), + covering("by_b", Some(vec![0, 1])), + false, + )] + #[case::unstated_coverage_reaches_everywhere( + covering("by_a", None), + covering("by_a", Some(vec![7])), + true, + )] + #[case::two_segments_of_unstated_coverage(covering("by_a", None), covering("by_a", None), true)] + #[case::disagreeing_about_which_fields_the_index_is_over( + AddIndexSegment { fields: vec![Ref::Committed(1)], ..covering("by_a", Some(vec![0])) }, + covering("by_a", Some(vec![2])), + true, + )] + #[case::disagreeing_about_the_index_config( + AddIndexSegment { + index_details: Some(Arc::new(prost_types::Any { + type_url: "type.googleapis.com/lance.index.pb.VectorIndexDetails".into(), + value: vec![1, 2, 3], + })), + ..covering("by_a", Some(vec![0])) + }, + covering("by_a", Some(vec![2])), + true, + )] + #[case::disagreeing_about_the_index_version( + AddIndexSegment { index_version: 2, ..covering("by_a", Some(vec![0])) }, + covering("by_a", Some(vec![2])), + true, + )] + fn test_two_writers_extending_one_index( + #[case] ours: AddIndexSegment, + #[case] theirs: AddIndexSegment, + #[case] expected: bool, + ) { + let ours = index_footprint(ours); + let theirs = index_footprint(theirs); + assert_eq!(ours.conflicts_with(&theirs), expected); + assert_eq!(theirs.conflicts_with(&ours), expected); + } + + #[test] + fn test_two_writers_indexing_what_they_each_just_wrote_do_not_conflict() { + // Neither segment names a committed fragment, and the fragments they do + // name get distinct ids once both commits are ordered. + let local = AddIndexSegment { + covered_fragments: Some(vec![Ref::Local(0)]), + ..segment("by_a", vec![Ref::Committed(0)]) + }; + + let ours = index_footprint(local.clone()); + let theirs = index_footprint(local); + assert!(!ours.conflicts_with(&theirs)); + } +} diff --git a/rust/lance-table/src/transaction/action/adjust_index_coverage.rs b/rust/lance-table/src/transaction/action/adjust_index_coverage.rs new file mode 100644 index 00000000000..fd126ed4107 --- /dev/null +++ b/rust/lance-table/src/transaction/action/adjust_index_coverage.rs @@ -0,0 +1,409 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Adjust the fragment coverage of an index segment. + +use super::apply::ApplyState; +use super::proto::{non_empty, required}; +use super::{Coordinate, Footprint, Ref}; +use crate::format::pb; +use lance_core::deepsize::{Context, DeepSizeOf}; +use lance_core::{Error, Result}; +use uuid::Uuid; + +/// Adjust which fragments an existing index segment covers, without rewriting +/// the segment. +/// +/// This is how a segment picks up fragments a compaction produced, or lets go +/// of ones it no longer describes, in cases where the index files themselves are +/// still good. Rewriting a segment's contents is a +/// [`RemoveIndexSegment`](super::RemoveIndexSegment) plus an +/// [`AddIndexSegment`](super::AddIndexSegment) instead. +/// +/// Additions are applied before removals, so a fragment named on both sides ends +/// up outside the coverage. +/// +/// # When a fragment may be added +/// +/// Coverage says which fragments a segment's existing index files describe, so a +/// fragment may only be added when they already describe its rows. In practice +/// that means one case: a rewrite -- compaction, or an in-place column rewrite -- +/// moved rows the segment already covered into a new fragment, and the segment +/// reaches them through the fragment-reuse remapping. The matching removal of the +/// old fragment ids belongs in the same action. +/// +/// Adding a fragment of genuinely new rows is a writer error, even though nothing +/// here can detect it: the segment has no entries for those rows, so a query +/// served by the index would silently miss them. New rows are covered by building +/// a segment over them ([`AddIndexSegment`](super::AddIndexSegment)), not by +/// widening an existing one. +#[derive(Debug, Clone, PartialEq)] +pub struct AdjustIndexCoverage { + pub uuid: Uuid, + /// The logical index the segment belongs to. Conflict detection compares + /// coverage across the segments of one index, so it needs the index and not + /// just the segment. + pub name: String, + /// Fragments to bring into the coverage. A [`Ref::Local`] names one this + /// operation minted, which is how a segment follows its rows into the + /// fragment a rewrite in the same operation moved them to. + pub add_fragments: Vec, + /// Committed fragment ids to drop from the coverage. Unlike the additions + /// these take no [`Ref`]: a fragment minted in this same operation was not + /// in the coverage to begin with. + pub remove_fragments: Vec, +} + +impl AdjustIndexCoverage { + pub(super) fn apply(&self, state: &mut ApplyState) -> Result<()> { + let added = self + .add_fragments + .iter() + .map(|fragment| self.coverage_id(state.resolve_fragment(*fragment)?)) + .collect::>>()?; + let removed = self + .remove_fragments + .iter() + .map(|fragment| self.coverage_id(*fragment)) + .collect::>>()?; + + let segment = state.index_segment_mut(self.uuid, &self.name)?; + let Some(bitmap) = segment.fragment_bitmap.as_mut() else { + return Err(Error::invalid_input(format!( + "index segment {} records no fragment coverage, so there is nothing to adjust; \ + rewrite the segment to give it coverage", + self.uuid + ))); + }; + for fragment in added { + bitmap.insert(fragment); + } + for fragment in removed { + bitmap.remove(fragment); + } + Ok(()) + } + + fn coverage_id(&self, fragment: u64) -> Result { + u32::try_from(fragment).map_err(|_| { + Error::invalid_input(format!( + "AdjustIndexCoverage for segment {} names fragment {fragment}, which is beyond \ + the largest fragment id an index can record ({})", + self.uuid, + u32::MAX + )) + }) + } + + /// Coverage says which fragments an index describes, never what any of them + /// hold, so adjusting it cannot change a row a reader sees. + pub(super) fn is_data_change(&self) -> bool { + false + } + + /// The segment, by uuid, plus a claim on the fragments this brings under + /// the index. The claim is what stops one writer widening a segment onto a + /// fragment another writer is covering with a segment of its own. Those + /// fragments are also required to still be there: coverage of a fragment a + /// concurrent set removed describes rows the dataset no longer has. + /// + /// Removals claim nothing: coverage the index gives up cannot collide with + /// coverage another writer takes on. + /// + /// What this cannot require is the fragments' *data*. The action names no + /// fields -- the segment already has them -- so the footprint cannot say + /// which columns the widened coverage depends on, and a concurrent rewrite + /// of an indexed column in one of these fragments is not caught here. The + /// caller widening coverage over committed fragments is answerable for + /// having indexed them as they stand; the safe shape is to widen in the + /// same operation that writes the fragments, where they are minted and no + /// concurrent writer can touch them. + pub(super) fn footprint(&self, footprint: &mut Footprint) { + footprint.write(Coordinate::IndexSegment(self.uuid)); + footprint.extend_index_coverage(self.name.clone(), self.add_fragments.iter().copied()); + for fragment in self + .add_fragments + .iter() + .filter_map(|fragment| fragment.committed()) + { + footprint.require_fragment(fragment); + } + } +} + +impl DeepSizeOf for AdjustIndexCoverage { + fn deep_size_of_children(&self, context: &mut Context) -> usize { + self.uuid.as_bytes().deep_size_of_children(context) + + self.name.deep_size_of_children(context) + + self.add_fragments.deep_size_of_children(context) + + self.remove_fragments.deep_size_of_children(context) + } +} + +impl From<&AdjustIndexCoverage> for pb::AdjustIndexCoverage { + fn from(value: &AdjustIndexCoverage) -> Self { + Self { + uuid: Some((&value.uuid).into()), + name: value.name.clone(), + add_fragments: value + .add_fragments + .iter() + .map(|fragment| (*fragment).into()) + .collect(), + remove_fragments: value.remove_fragments.clone(), + } + } +} + +impl TryFrom for AdjustIndexCoverage { + type Error = Error; + + fn try_from(message: pb::AdjustIndexCoverage) -> Result { + Ok(Self { + uuid: Uuid::try_from(&required(message.uuid, "AdjustIndexCoverage.uuid")?)?, + name: non_empty(message.name, "AdjustIndexCoverage.name")?, + add_fragments: message + .add_fragments + .into_iter() + .map(Ref::try_from) + .collect::>>()?, + remove_fragments: message.remove_fragments, + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::format::IndexMetadata; + use crate::transaction::action::test_support::{ + apply_with_indices, backed_manifest, footprint, + }; + use crate::transaction::action::{Action, AddFragment, AddIndexSegment, RemoveFragment}; + use crate::transaction::test_support::sample_index_metadata; + + fn covering(name: &str, fragments: impl IntoIterator) -> IndexMetadata { + IndexMetadata { + fragment_bitmap: Some(fragments.into_iter().collect()), + ..sample_index_metadata(name) + } + } + + fn adjust(uuid: Uuid, name: &str, add: Vec, remove: Vec) -> Action { + Action::AdjustIndexCoverage(AdjustIndexCoverage { + uuid, + name: name.into(), + add_fragments: add, + remove_fragments: remove, + }) + } + + #[test] + fn test_adjust_index_coverage_adds_and_removes() { + let segment = covering("by_a", [0, 1]); + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![adjust( + segment.uuid, + "by_a", + vec![Ref::Committed(2)], + vec![1], + )], + vec![segment], + ) + .unwrap(); + + assert_eq!( + indices[0].fragment_bitmap, + Some([0, 2].into_iter().collect()) + ); + } + + #[test] + fn test_coverage_can_take_in_a_fragment_minted_in_the_same_operation() { + // Appending and extending an index's reach over what was appended is + // one operation, and the fragment has no committed id until apply. + let segment = covering("by_a", [0]); + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![ + Action::AddFragment(AddFragment { + id: Ref::Local(0), + physical_rows: 10, + row_id_meta: None, + last_updated_at_version_meta: None, + created_at_version_meta: None, + data_change: true, + }), + adjust(segment.uuid, "by_a", vec![Ref::Local(0)], vec![]), + ], + vec![segment], + ) + .unwrap(); + + assert_eq!( + indices[0].fragment_bitmap, + Some([0, 1].into_iter().collect()) + ); + } + + #[test] + fn test_a_fragment_named_on_both_sides_ends_up_uncovered() { + let segment = covering("by_a", [0]); + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![adjust( + segment.uuid, + "by_a", + vec![Ref::Committed(5)], + vec![5], + )], + vec![segment], + ) + .unwrap(); + + assert_eq!(indices[0].fragment_bitmap, Some([0].into_iter().collect())); + } + + #[test] + fn test_adjusting_a_segment_that_is_not_there_is_rejected() { + let error = apply_with_indices( + &backed_manifest(), + vec![adjust( + Uuid::from_u128(1), + "by_a", + vec![Ref::Committed(0)], + vec![], + )], + vec![covering("by_a", [0])], + ) + .unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains("is not part of the dataset"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_adjusting_a_segment_with_no_recorded_coverage_is_rejected() { + // Coverage of "unknown" is not an empty set to add to: turning it into + // a concrete set would narrow what the segment serves, silently. + let segment = IndexMetadata { + fragment_bitmap: None, + ..sample_index_metadata("system") + }; + let error = apply_with_indices( + &backed_manifest(), + vec![adjust( + segment.uuid, + "system", + vec![Ref::Committed(0)], + vec![], + )], + vec![segment], + ) + .unwrap_err(); + + assert!( + error.to_string().contains("records no fragment coverage"), + "unexpected error: {error}" + ); + } + + fn build(name: &str, fragment: u64) -> Action { + Action::AddIndexSegment(AddIndexSegment { + uuid: Uuid::new_v4(), + name: name.into(), + fields: vec![Ref::Committed(0)], + covering_fields: Vec::new(), + index_details: None, + index_version: 1, + covered_fragments: Some(vec![Ref::Committed(fragment)]), + files: Vec::new(), + base: None, + created_at: None, + dataset_version: None, + data_change: false, + }) + } + + #[test] + fn test_adjusting_a_segment_of_another_index_is_rejected() { + let segment = covering("by_a", [0]); + let error = apply_with_indices( + &backed_manifest(), + vec![adjust( + segment.uuid, + "by_b", + vec![Ref::Committed(2)], + vec![], + )], + vec![segment], + ) + .unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains("belongs to index 'by_a'"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_two_writers_adjusting_the_same_segment_conflict() { + let uuid = Uuid::from_u128(1); + let ours = footprint(vec![adjust(uuid, "by_a", vec![Ref::Committed(1)], vec![])]); + let same = footprint(vec![adjust(uuid, "by_a", vec![], vec![2])]); + let other = footprint(vec![adjust(Uuid::from_u128(2), "by_a", vec![], vec![2])]); + + assert!(ours.conflicts_with(&same)); + assert!(!ours.conflicts_with(&other)); + } + + #[test] + fn test_widening_onto_a_fragment_a_concurrent_segment_covers_conflicts() { + let widen = footprint(vec![adjust( + Uuid::from_u128(1), + "by_a", + vec![Ref::Committed(7)], + vec![], + )]); + + assert!(widen.conflicts_with(&footprint(vec![build("by_a", 7)]))); + assert!(!widen.conflicts_with(&footprint(vec![build("by_a", 8)]))); + assert!(!widen.conflicts_with(&footprint(vec![build("by_b", 7)]))); + } + + #[test] + fn test_widening_onto_a_fragment_a_concurrent_set_removes_conflicts() { + let widen = footprint(vec![adjust( + Uuid::from_u128(1), + "by_a", + vec![Ref::Committed(7)], + vec![], + )]); + let remove = |fragment: u64| { + footprint(vec![Action::RemoveFragment(RemoveFragment { + fragment: Ref::Committed(fragment), + data_change: false, + })]) + }; + + // Either order: coverage of a fragment that is gone describes rows the + // dataset no longer has. + assert!(widen.conflicts_with(&remove(7))); + assert!(remove(7).conflicts_with(&widen)); + assert!(!widen.conflicts_with(&remove(8))); + } + + #[test] + fn test_giving_up_coverage_claims_nothing() { + // Coverage an index drops cannot collide with coverage another writer + // takes on, so a removal-only adjustment leaves the index free. + let narrow = footprint(vec![adjust(Uuid::from_u128(1), "by_a", vec![], vec![7])]); + + assert!(!narrow.conflicts_with(&footprint(vec![build("by_a", 7)]))); + } +} diff --git a/rust/lance-table/src/transaction/action/apply.rs b/rust/lance-table/src/transaction/action/apply.rs index c6f40b94aeb..cdbd0cd477c 100644 --- a/rust/lance-table/src/transaction/action/apply.rs +++ b/rust/lance-table/src/transaction/action/apply.rs @@ -27,6 +27,7 @@ use crate::transaction::{LogicalIndexSegments, ReadVersionState, Transaction}; use lance_core::datatypes::Schema; use lance_core::{Error, Result}; use std::collections::{HashMap, HashSet}; +use uuid::Uuid; /// The field id written into a data file's field list once the file no longer /// backs that field. A file whose every slot is tombstoned is dropped. @@ -70,16 +71,15 @@ impl Transaction { .any(|index| index.name == MEM_WAL_INDEX_NAME) .then(|| Self::logical_index_segments(¤t_indices)); - let mut state = ApplyState::new(current_manifest); + let mut state = + ApplyState::new(current_manifest, current_indices, config, self.read_version); for action in composite_operation.iter_actions() { action.apply(&mut state)?; } state.into_manifest( self, - current_indices, transaction_file_path, - config, mem_wal_segments_before.as_ref(), read_version_state, ) @@ -91,12 +91,18 @@ impl Transaction { pub(super) struct ApplyState<'a> { /// The read version this delta applies to. current_manifest: &'a Manifest, + /// How the manifest this delta produces is to be assembled. + build_config: &'a ManifestBuildConfig, schema: Schema, fragments: Vec, /// Base paths minted by this operation. Kept apart from the manifest's own /// base paths, which the manifest assembly inherits from the read version. new_bases: Vec, existing_base_paths: HashMap, + /// The index segments, as the actions have left them so far. Index + /// actions edit this list; the assembly then prunes whatever the data + /// actions invalidated. + indices: Vec, /// The manifest's string maps. Unlike the schema and the fragment list, /// these are inherited wholesale by the manifest assembly, so an edit has /// to be written back over the assembled manifest. @@ -130,17 +136,26 @@ pub(super) struct ApplyState<'a> { /// a field no longer describes that fragment's contents. rebound_fields: HashMap>, - /// Whether the table was reset, which discards every index outright rather - /// than pruning fragments out of them. - reset: bool, + /// The version the transaction was built against. Not the current + /// manifest's version: on a retry the set is replayed onto something newer, + /// and what an action saw is still the version it read. + read_version: u64, } impl<'a> ApplyState<'a> { - fn new(manifest: &'a Manifest) -> Self { + fn new( + manifest: &'a Manifest, + indices: Vec, + build_config: &'a ManifestBuildConfig, + read_version: u64, + ) -> Self { Self { current_manifest: manifest, + read_version, + build_config, schema: manifest.schema.clone(), fragments: manifest.fragments.as_ref().clone(), + indices, new_bases: Vec::new(), existing_base_paths: manifest.base_paths.clone(), config: manifest.config.clone(), @@ -160,7 +175,6 @@ impl<'a> ApplyState<'a> { reserved_fragment_ids: None, reserved_row_ids: 0, rebound_fields: HashMap::new(), - reset: false, } } @@ -168,13 +182,12 @@ impl<'a> ApplyState<'a> { fn into_manifest( mut self, transaction: &Transaction, - current_indices: Vec, transaction_file_path: &str, - config: &ManifestBuildConfig, mem_wal_segments_before: Option<&LogicalIndexSegments>, read_version_state: Option>, ) -> Result<(Manifest, Vec)> { let current_manifest = self.current_manifest; + let config = self.build_config; let new_version = current_manifest.version + 1; let mut next_row_id = current_manifest @@ -186,19 +199,15 @@ impl<'a> ApplyState<'a> { let ApplyState { schema, mut fragments, + mut indices, new_bases, rebound_fields, reserved_fragment_ids, - reset, config: dataset_config, table_metadata, .. } = self; - let mut indices = current_indices; - if reset { - indices.clear(); - } prune_rebound_fields_from_indices(&mut indices, &rebound_fields); Transaction::retain_relevant_indices(&mut indices, &schema, &fragments); @@ -260,6 +269,75 @@ impl<'a> ApplyState<'a> { Ok((manifest, indices)) } + /// The version the transaction read, which is the newest data an index + /// segment added by this operation can have been built from. + /// + /// Deliberately not the current manifest's version. A set that loses a race + /// is replayed against whatever won, and a segment built before that still + /// reflects only what its writer could see -- stamping it with the newer + /// version would claim coverage of rows it never read. + pub(super) fn read_version(&self) -> u64 { + self.read_version + } + + /// 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) + /// followed by an add. + pub(super) fn add_index_segment(&mut self, segment: IndexMetadata) -> Result<()> { + if self.indices.iter().any(|index| index.uuid == segment.uuid) { + return Err(Error::invalid_input(format!( + "index segment {} is already part of the dataset; a segment is added once", + segment.uuid + ))); + } + self.indices.push(segment); + Ok(()) + } + + /// Drop an index segment by uuid. A removal that names a segment the + /// dataset does not have is a mistake rather than a no-op: it means the + /// operation was planned against a different set of segments. + pub(super) fn remove_index_segment(&mut self, uuid: Uuid, name: &str) -> Result<()> { + let Some(position) = self.index_segment_position(uuid, name)? else { + return Err(Error::invalid_input(format!( + "index segment {uuid} is not part of the dataset, so it cannot be removed" + ))); + }; + self.indices.remove(position); + Ok(()) + } + + /// 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. + fn index_segment_position(&self, uuid: Uuid, name: &str) -> Result> { + let Some(position) = self.indices.iter().position(|index| index.uuid == uuid) else { + return Ok(None); + }; + let found = &self.indices[position].name; + if found != name { + return Err(Error::invalid_input(format!( + "index segment {uuid} belongs to index '{found}', but the action names it as part \ + of index '{name}'" + ))); + } + Ok(Some(position)) + } + + pub(super) fn index_segment_mut( + &mut self, + uuid: Uuid, + name: &str, + ) -> Result<&mut IndexMetadata> { + let Some(position) = self.index_segment_position(uuid, name)? else { + return Err(Error::invalid_input(format!( + "index segment {uuid} is not part of the dataset, so it cannot be adjusted" + ))); + }; + Ok(&mut self.indices[position]) + } + pub(super) fn schema(&self) -> &Schema { &self.schema } @@ -350,9 +428,9 @@ impl<'a> ApplyState<'a> { self.schema.fields.clear(); self.schema.metadata.clear(); self.fragments.clear(); + self.indices.clear(); self.new_fragments.clear(); self.rebound_fields.clear(); - self.reset = true; } /// The base paths this apply can see: the read version's, plus the ones @@ -459,6 +537,19 @@ impl<'a> ApplyState<'a> { } } + pub(super) fn resolve_base(&self, reference: Ref) -> Result { + match reference { + Ref::Committed(id) => u32::try_from(id).map_err(|_| { + Error::invalid_input(format!("base id {id} in an action is out of range")) + }), + Ref::Local(token) => self + .base_tokens + .get(&token) + .copied() + .ok_or_else(|| unbound_token_err("base", token)), + } + } + pub(super) fn resolve_field(&self, reference: Ref) -> Result { match reference { Ref::Committed(id) => i32::try_from(id).map_err(|_| { diff --git a/rust/lance-table/src/transaction/action/footprint.rs b/rust/lance-table/src/transaction/action/footprint.rs index eae41e4d601..1f7e41bd6a3 100644 --- a/rust/lance-table/src/transaction/action/footprint.rs +++ b/rust/lance-table/src/transaction/action/footprint.rs @@ -19,6 +19,8 @@ use super::{CompositeOperation, Ref}; use crate::transaction::UpdateMap; use std::collections::{HashMap, HashSet}; +use std::sync::Arc; +use uuid::Uuid; /// One thing an action set writes. /// @@ -54,6 +56,8 @@ pub enum Coordinate { /// once, on one set of fields, so two writers declaring it on fields of /// their own collide even though they name different fields' metadata. UnenforcedKey(UnenforcedKey), + /// One index segment, by uuid. + IndexSegment(Uuid), } /// One of the string maps a manifest carries. @@ -88,7 +92,8 @@ impl Coordinate { | Self::BaseName(_) | Self::BaseLocation(_) | Self::ConfigEntry { .. } - | Self::UnenforcedKey(_) => None, + | Self::UnenforcedKey(_) + | Self::IndexSegment(_) => None, } } @@ -112,7 +117,9 @@ impl Coordinate { enum Mode { /// Reads the coordinate and needs it to still hold what it read. Data /// written for a committed field requires the field's definition, since - /// the values are encoded in the type it names. Unlike a write, two sets + /// the values are encoded in the type it names; an index segment requires + /// the data it describes, since the format gives a reader no way to notice + /// that those values were replaced underneath it. Unlike a write, two sets /// may require the same coordinate -- two readers of one column do not /// collide. Requires, @@ -201,6 +208,65 @@ 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, + /// 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, +} + +/// What an action set writes into one logical index. +/// +/// An index is the set of segments sharing a name, and the query path unions +/// them, so two writers may extend one index at the same time -- what they may +/// not do is describe the same rows twice, or disagree about what the index is. +#[derive(Debug, Clone, PartialEq)] +struct IndexClaim { + name: String, + /// What the set says the index is. `None` for an action that edits a + /// segment without restating the index's definition. + identity: Option, + /// The committed fragments this set brings under the index, or `None` when + /// the reach is not stated -- what the system indices carry, and what makes + /// a claim collide with every other claim on the same index. + /// + /// Fragments this operation mints are left out. They have no id in the read + /// version, so a concurrent writer cannot be covering one. + coverage: Option>, +} + +/// What an index is, for the purpose of deciding whether two writers are +/// building the same one. +/// +/// `details` is compared as the opaque blob it is. Two segments of one index +/// built by the same writer serialize identical config, so equality is the +/// right test until index config is lifted out of the per-segment details. +#[derive(Debug, Clone, PartialEq)] +pub(super) struct IndexIdentity { + pub fields: Vec, + /// Which of `fields` are merely carried: two writers that disagree about + /// this disagree about what the index answers for, not just what it holds. + pub covering_fields: Vec, + pub details: Option>, + pub index_version: i32, +} + +impl IndexClaim { + fn conflicts_with(&self, other: &Self) -> bool { + if self.name != other.name { + return false; + } + if let (Some(ours), Some(theirs)) = (&self.identity, &other.identity) + && ours != theirs + { + return true; + } + match (&self.coverage, &other.coverage) { + (Some(ours), Some(theirs)) => !ours.is_disjoint(theirs), + // An unstated reach could be any fragment, including one the other + // side is claiming. + _ => true, + } + } } impl Footprint { @@ -212,12 +278,14 @@ impl Footprint { /// That is where the one asymmetry lives: a requirement of this set is /// tested against what `committed` wrote, never the reverse, because /// `committed` serialized first and nothing arriving later can break what - /// it required. The distinction is not academic. Data written for a field - /// requires the field's definition; a cast of that field writes it. - /// Checking both directions would also reject the cast that arrives - /// *after* the data -- an order [`AlterField`](super::AlterField) already - /// handles, by rebinding the field in every fragment the manifest has by - /// then. + /// it required. The distinction is not academic. An index segment requires + /// the data it describes; a rewrite of that data writes it. Checking both + /// directions would also reject the rewrite that arrives *after* a segment + /// -- 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. 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 @@ -228,7 +296,17 @@ impl Footprint { if self.claims_conflict_with(committed) { return true; } - self.anchors_removed_by(committed) || committed.anchors_removed_by(self) + 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. + self.index_claims.iter().any(|ours| { + committed + .index_claims + .iter() + .any(|theirs| ours.conflicts_with(theirs)) + }) } /// Whether any coordinate touched by both sets, directly or through a @@ -337,6 +415,21 @@ impl Footprint { } } + /// Record that this set reads `fields` in `fragment` and needs them to + /// still hold what it read. + /// + /// A fragment this operation mints records nothing: no concurrent writer + /// can have replaced data that did not exist when they planned. + pub(super) fn require_field_data( + &mut self, + fragment: u64, + fields: impl IntoIterator, + ) { + for field in fields.into_iter().filter_map(committed_field) { + self.require(Coordinate::FieldData { fragment, field }); + } + } + /// A field's entry in the schema. pub(super) fn write_field_definition(&mut self, field: Ref) { if let Some(field) = committed_field(field) { @@ -358,6 +451,35 @@ impl Footprint { self.required_fragments.insert(fragment); } + /// Record that this set adds a segment to `name`, defining the index as + /// `identity` and describing `coverage`. + pub(super) fn build_index( + &mut self, + name: String, + identity: IndexIdentity, + coverage: Option>, + ) { + self.index_claims.push(IndexClaim { + name, + identity: Some(identity), + coverage: coverage.map(committed_fragments), + }); + } + + /// Record that this set brings `fragments` under `name` by widening a + /// segment that is already there, without restating what the index is. + pub(super) fn extend_index_coverage( + &mut self, + name: String, + fragments: impl IntoIterator, + ) { + self.index_claims.push(IndexClaim { + name, + identity: None, + coverage: Some(committed_fragments(fragments)), + }); + } + /// 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) { @@ -405,6 +527,13 @@ fn committed_field(reference: Ref) -> Option { i32::try_from(reference.committed()?).ok() } +fn committed_fragments(fragments: impl IntoIterator) -> HashSet { + fragments + .into_iter() + .filter_map(|fragment| fragment.committed()) + .collect() +} + impl From<&CompositeOperation> for Footprint { fn from(composite_operation: &CompositeOperation) -> Self { let mut footprint = Self::default(); diff --git a/rust/lance-table/src/transaction/action/proto.rs b/rust/lance-table/src/transaction/action/proto.rs index a3059657aa2..7a5506e3268 100644 --- a/rust/lance-table/src/transaction/action/proto.rs +++ b/rust/lance-table/src/transaction/action/proto.rs @@ -30,6 +30,17 @@ pub(super) fn required(value: Option, what: &str) -> Result { value.ok_or_else(|| Error::invalid_input(format!("{what} is required but was not set"))) } +/// A required string field. Protobuf gives an unset string the empty value, so +/// the two cases cannot be told apart and neither is usable. +pub(super) fn non_empty(value: String, what: &str) -> Result { + if value.is_empty() { + return Err(Error::invalid_input(format!( + "{what} is required but was empty" + ))); + } + Ok(value) +} + impl From for pb::Ref { fn from(value: Ref) -> Self { let kind = match value { @@ -142,18 +153,24 @@ for_each_action!(define_action_proto); #[cfg(test)] mod tests { use super::*; - use crate::format::{BasePath, DataFile, DeletionFile, DeletionFileType, RowIdMeta, pb}; + 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, AlterField, ConfigUpdate, DropField, - FieldMetadataUpdate, RemoveFragment, ReserveFragmentIds, ReserveRowIds, ResetTable, - SetDeletionFile, TombstoneFieldData, + AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, 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 rstest::rstest; use std::sync::Arc; + use uuid::Uuid; fn sample_data_file() -> DataFile { DataFile::new_unstarted("data/1.lance", ConcreteFileVersion::V2_0) @@ -215,6 +232,37 @@ mod tests { Action::DropField(DropField { field: Ref::Committed(3), }), + Action::AddIndexSegment(AddIndexSegment { + uuid: Uuid::from_u128(7), + name: "by_a".into(), + fields: vec![Ref::Committed(1), Ref::Local(3)], + covering_fields: vec![Ref::Local(3)], + index_details: Some(Arc::new(prost_types::Any { + type_url: "type.googleapis.com/lance.table.MemWalIndexDetails".into(), + value: vec![1, 2, 3], + })), + index_version: 2, + covered_fragments: Some(vec![Ref::Committed(4), Ref::Local(0)]), + files: vec![IndexFile { + path: "index.idx".into(), + size_bytes: 512, + }], + base: Some(Ref::Local(1)), + created_at: DateTime::from_timestamp_millis(1_700_000_000_000), + dataset_version: Some(3), + data_change: false, + }), + Action::RemoveIndexSegment(RemoveIndexSegment { + uuid: Uuid::from_u128(8), + name: "by_a".into(), + data_change: true, + }), + Action::AdjustIndexCoverage(AdjustIndexCoverage { + uuid: Uuid::from_u128(9), + name: "by_b".into(), + add_fragments: vec![Ref::Committed(1), Ref::Local(0)], + remove_fragments: vec![2, 3], + }), Action::ReserveFragmentIds(ReserveFragmentIds { count: 4 }), Action::ReserveRowIds(ReserveRowIds { count: 40 }), Action::ResetTable(ResetTable), @@ -293,6 +341,44 @@ mod tests { ); } + #[rstest] + #[case::addition(pb::Action { + action: Some(pb::action::Action::AddIndexSegment(pb::AddIndexSegment { + uuid: Some((&Uuid::from_u128(1)).into()), + name: String::new(), + ..Default::default() + })), + }, "AddIndexSegment.name")] + #[case::removal(pb::Action { + action: Some(pb::action::Action::RemoveIndexSegment(pb::RemoveIndexSegment { + uuid: Some((&Uuid::from_u128(1)).into()), + name: String::new(), + data_change: None, + })), + }, "RemoveIndexSegment.name")] + #[case::coverage(pb::Action { + action: Some(pb::action::Action::AdjustIndexCoverage(pb::AdjustIndexCoverage { + uuid: Some((&Uuid::from_u128(1)).into()), + name: String::new(), + add_fragments: Vec::new(), + remove_fragments: Vec::new(), + })), + }, "AdjustIndexCoverage.name")] + fn test_an_index_action_without_a_name_is_rejected( + #[case] message: pb::Action, + #[case] expected: &str, + ) { + // Protobuf cannot tell an unset string from an empty one, and an index + // action that does not say which index it edits cannot be checked for + // conflicts. + let error = Action::try_from(message).unwrap_err(); + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains(expected), + "unexpected message: {error}" + ); + } + #[test] fn test_empty_ref_is_rejected() { let error = Ref::try_from(pb::Ref { kind: None }).unwrap_err(); diff --git a/rust/lance-table/src/transaction/action/remove_index_segment.rs b/rust/lance-table/src/transaction/action/remove_index_segment.rs new file mode 100644 index 00000000000..a5753a01d8b --- /dev/null +++ b/rust/lance-table/src/transaction/action/remove_index_segment.rs @@ -0,0 +1,239 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Remove an index segment. + +use super::apply::ApplyState; +use super::proto::{data_change_from_wire, data_change_to_wire, non_empty, required}; +use super::{Coordinate, Footprint}; +use crate::format::pb; +use lance_core::deepsize::{Context, DeepSizeOf}; +use lance_core::{Error, Result}; +use uuid::Uuid; + +/// Remove an index segment. +/// +/// Dropping a whole logical index is one of these per segment carrying its +/// name, since the format knows only segments. Removing the last segment of an +/// index is what makes the index disappear. +#[derive(Debug, Clone, PartialEq)] +pub struct RemoveIndexSegment { + pub uuid: Uuid, + /// The logical index the segment belongs to. Not needed to find the + /// segment, which the uuid names outright; carried so apply can reject a + /// removal that was planned against a different set of segments. + pub name: String, + pub data_change: bool, +} + +impl RemoveIndexSegment { + pub(super) fn apply(&self, state: &mut ApplyState) -> Result<()> { + state.remove_index_segment(self.uuid, &self.name) + } + + /// Left to the writer, for the same reason as + /// [`AddIndexSegment`](super::AddIndexSegment): an index is normally derived + /// state, but the MemWAL index holds rows a reader can see. + pub(super) fn is_data_change(&self) -> bool { + self.data_change + } + + /// The segment, by uuid. Only another action naming the same segment + /// collides: a concurrent writer extending the same logical index is adding + /// a segment of its own, which this removal leaves alone, and one that + /// rebuilds what this segment covered leaves the index no worse off. + pub(super) fn footprint(&self, footprint: &mut Footprint) { + footprint.write(Coordinate::IndexSegment(self.uuid)); + } +} + +impl DeepSizeOf for RemoveIndexSegment { + fn deep_size_of_children(&self, context: &mut Context) -> usize { + self.uuid.as_bytes().deep_size_of_children(context) + + self.name.deep_size_of_children(context) + } +} + +impl From<&RemoveIndexSegment> for pb::RemoveIndexSegment { + fn from(value: &RemoveIndexSegment) -> Self { + Self { + uuid: Some((&value.uuid).into()), + name: value.name.clone(), + data_change: data_change_to_wire(value.data_change), + } + } +} + +impl TryFrom for RemoveIndexSegment { + type Error = Error; + + fn try_from(message: pb::RemoveIndexSegment) -> Result { + Ok(Self { + uuid: Uuid::try_from(&required(message.uuid, "RemoveIndexSegment.uuid")?)?, + name: non_empty(message.name, "RemoveIndexSegment.name")?, + data_change: data_change_from_wire(message.data_change), + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::transaction::action::test_support::{apply_with_indices, backed_manifest}; + use crate::transaction::action::{ + Action, AddIndexSegment, CompositeOperation, Footprint, Ref, UserAction, + }; + use crate::transaction::test_support::sample_index_metadata; + + fn remove(uuid: Uuid, name: &str) -> Action { + Action::RemoveIndexSegment(RemoveIndexSegment { + uuid, + name: name.into(), + data_change: false, + }) + } + + #[test] + fn test_remove_index_segment_drops_only_the_segment_it_names() { + let dropped = sample_index_metadata("by_a"); + let kept = sample_index_metadata("by_b"); + let dropped_uuid = dropped.uuid; + + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![remove(dropped_uuid, "by_a")], + vec![dropped, kept.clone()], + ) + .unwrap(); + + assert_eq!(indices.len(), 1); + assert_eq!(indices[0].uuid, kept.uuid); + } + + #[test] + fn test_dropping_an_index_removes_each_of_its_segments() { + let first = sample_index_metadata("by_a"); + let second = sample_index_metadata("by_a"); + + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![remove(first.uuid, "by_a"), remove(second.uuid, "by_a")], + vec![first, second], + ) + .unwrap(); + + assert!(indices.is_empty()); + } + + #[test] + fn test_removing_a_segment_that_is_not_there_is_rejected() { + let error = apply_with_indices( + &backed_manifest(), + vec![remove(Uuid::from_u128(1), "by_a")], + vec![sample_index_metadata("by_a")], + ) + .unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains("is not part of the dataset"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_a_segment_can_be_replaced_within_one_operation() { + let old = sample_index_metadata("by_a"); + let new_uuid = Uuid::from_u128(2); + let (_, indices) = apply_with_indices( + &backed_manifest(), + vec![ + remove(old.uuid, "by_a"), + Action::AddIndexSegment(AddIndexSegment { + uuid: new_uuid, + 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, + }), + ], + vec![old], + ) + .unwrap(); + + assert_eq!(indices.len(), 1); + assert_eq!(indices[0].uuid, new_uuid); + } + + #[test] + fn test_removing_a_segment_of_another_index_is_rejected() { + let segment = sample_index_metadata("by_a"); + let uuid = segment.uuid; + let error = apply_with_indices( + &backed_manifest(), + vec![remove(uuid, "by_b")], + vec![segment], + ) + .unwrap_err(); + + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error.to_string().contains("belongs to index 'by_a'"), + "unexpected error: {error}" + ); + } + + #[test] + fn test_removing_a_segment_leaves_a_concurrent_build_of_its_index_alone() { + // Dropping a stale segment while someone else builds a fresh one is the + // normal shape of index maintenance; the index ends up with the new + // segment and without the old one either way round. + let footprint = |actions| { + Footprint::from(&CompositeOperation::new(vec![UserAction::new( + "step", actions, + )])) + }; + let ours = footprint(vec![remove(Uuid::from_u128(1), "by_a")]); + let theirs = footprint(vec![Action::AddIndexSegment(AddIndexSegment { + uuid: Uuid::from_u128(2), + 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, + })]); + + assert!(!ours.conflicts_with(&theirs)); + assert!(!theirs.conflicts_with(&ours)); + } + + #[test] + fn test_two_writers_removing_the_same_segment_conflict() { + let footprint = |actions| { + Footprint::from(&CompositeOperation::new(vec![UserAction::new( + "step", actions, + )])) + }; + + let uuid = Uuid::from_u128(1); + let ours = footprint(vec![remove(uuid, "by_a")]); + let same = footprint(vec![remove(uuid, "by_a")]); + let other = footprint(vec![remove(Uuid::from_u128(2), "by_a")]); + + assert!(ours.conflicts_with(&same)); + assert!(!ours.conflicts_with(&other)); + } +} diff --git a/rust/lance/src/dataset/tests/dataset_transactions.rs b/rust/lance/src/dataset/tests/dataset_transactions.rs index 7dd56ebbf51..2e677c6db5b 100644 --- a/rust/lance/src/dataset/tests/dataset_transactions.rs +++ b/rust/lance/src/dataset/tests/dataset_transactions.rs @@ -1741,6 +1741,7 @@ mod composite { use std::sync::Arc; use crate::dataset::{CommitBuilder, InsertBuilder, WriteParams}; + use crate::index::load_all_indices; use crate::{Dataset, Error, Result}; use arrow_array::{Int32Array, RecordBatch}; use arrow_schema::{DataType, Field, Schema}; @@ -1751,9 +1752,10 @@ mod composite { }; use lance_table::rowids::{RowIdSequence, write_row_ids}; use lance_table::transaction::action::{ - Action, AddBase, AddDataFile, AddField, AddFragment, AlterField, CompositeOperation, - ConfigUpdate, DropField, Ref, RemoveFragment, ReserveFragmentIds, ReserveRowIds, - ResetTable, SetDeletionFile, TombstoneFieldData, UserAction, + Action, AddBase, AddDataFile, AddField, AddFragment, AddIndexSegment, AdjustIndexCoverage, + AlterField, CompositeOperation, ConfigUpdate, DropField, Ref, RemoveFragment, + RemoveIndexSegment, ReserveFragmentIds, ReserveRowIds, ResetTable, SetDeletionFile, + TombstoneFieldData, UserAction, }; use lance_table::transaction::{Operation, Transaction, UpdateMap, UpdateMapEntry}; use roaring::RoaringBitmap; @@ -2540,4 +2542,313 @@ mod composite { "unexpected error: {error}" ); } + + /// An index segment naming `covered`, with no files of its own -- these + /// tests inspect the manifest's index section rather than opening the index. + fn index_segment(name: &str, covered: Vec) -> AddIndexSegment { + AddIndexSegment { + uuid: Uuid::new_v4(), + name: name.into(), + fields: vec![Ref::Committed(0)], + covering_fields: Vec::new(), + index_details: None, + index_version: 1, + covered_fragments: Some(covered), + files: Vec::new(), + base: None, + created_at: None, + dataset_version: None, + data_change: false, + } + } + + async fn index_coverage(dataset: &Dataset, name: &str) -> Vec { + let indices = load_all_indices(dataset).await.unwrap(); + let segment = indices + .iter() + .find(|index| index.name == name) + .expect("index segment should be committed"); + segment + .fragment_bitmap + .as_ref() + .expect("coverage should be recorded") + .iter() + .collect() + } + + #[tokio::test] + async fn test_one_commit_appends_and_indexes_what_it_appended() { + let dataset = test_dataset(false).await; + let existing = dataset.fragments()[0].id; + let file = existing_data_file(&dataset, 0); + + 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, + }), + // The index covers the fragment this same commit minted, which + // has no id until the commit lands. + Action::AddIndexSegment(index_segment( + "by_a", + vec![Ref::Committed(existing), Ref::Local(0)], + )), + ], + ) + .await; + + let appended = dataset.fragments().last().unwrap().id; + assert_eq!( + index_coverage(&dataset, "by_a").await, + vec![existing as u32, appended as u32] + ); + } + + #[tokio::test] + async fn test_one_commit_replaces_an_index_segment() { + let dataset = test_dataset(false).await; + let old = index_segment("by_a", vec![Ref::Committed(0)]); + let old_uuid = old.uuid; + let dataset = commit(dataset, vec![Action::AddIndexSegment(old)]).await; + + let new = index_segment("by_a", vec![Ref::Committed(0), Ref::Committed(1)]); + let new_uuid = new.uuid; + let dataset = commit( + dataset, + vec![ + Action::RemoveIndexSegment(RemoveIndexSegment { + uuid: old_uuid, + name: "by_a".into(), + data_change: false, + }), + Action::AddIndexSegment(new), + ], + ) + .await; + + let uuids = load_all_indices(&dataset) + .await + .unwrap() + .iter() + .map(|index| index.uuid) + .collect::>(); + assert_eq!(uuids, vec![new_uuid]); + } + + /// Commit `actions` against `dataset` as it was, whatever has landed since. + /// This is what a concurrent writer does: it planned at that version and + /// arrives at the commit after someone else got there first. + async fn commit_from_stale(dataset: Arc, actions: Vec) -> Result { + let read_version = dataset.version().version; + CommitBuilder::new(dataset) + .with_experimental_composite_operations(true) + .execute(composite_txn(read_version, actions)) + .await + } + + #[tokio::test] + async fn test_two_writers_extend_one_index_over_disjoint_fragments() { + let dataset = test_dataset(false).await; + let stale = Arc::new(dataset.clone()); + + // The first writer wins the race; the second arrives against a version + // it has not seen. + commit( + dataset, + vec![Action::AddIndexSegment(index_segment( + "by_a", + vec![Ref::Committed(0)], + ))], + ) + .await; + + let dataset = commit_from_stale( + stale, + vec![Action::AddIndexSegment(index_segment( + "by_a", + vec![Ref::Committed(1)], + ))], + ) + .await + .expect("segments over disjoint fragments both belong to the index"); + + let mut coverage = load_all_indices(&dataset) + .await + .unwrap() + .iter() + .filter(|index| index.name == "by_a") + .map(|index| index.fragment_bitmap.as_ref().unwrap().iter().collect()) + .collect::>>(); + coverage.sort(); + assert_eq!(coverage, vec![vec![0], vec![1]]); + } + + #[tokio::test] + async fn test_two_writers_indexing_the_same_fragment_conflict() { + let dataset = test_dataset(false).await; + let stale = Arc::new(dataset.clone()); + + commit( + dataset, + vec![Action::AddIndexSegment(index_segment( + "by_a", + vec![Ref::Committed(0)], + ))], + ) + .await; + + let error = commit_from_stale( + stale, + vec![Action::AddIndexSegment(index_segment( + "by_a", + vec![Ref::Committed(0)], + ))], + ) + .await + .expect_err("two segments describing one fragment would double-count its rows"); + + assert!( + matches!(error, Error::RetryableCommitConflict { .. }), + "{error:?}" + ); + } + + /// Replace field 0's data in fragment 0, the way an in-place column + /// rewrite does. + fn rewrite_field_zero(file: DataFile) -> Vec { + vec![ + Action::TombstoneFieldData(TombstoneFieldData { + fragment: Ref::Committed(0), + field_ids: vec![Ref::Committed(0)], + data_change: true, + }), + Action::AddDataFile(AddDataFile { + fragment: Ref::Committed(0), + file, + field_ids: vec![Ref::Committed(0)], + data_change: true, + }), + ] + } + + /// The coverage of `name`, or `None` when no segment carries the name. + async fn index_coverage_opt(dataset: &Dataset, name: &str) -> Option> { + let indices = load_all_indices(dataset).await.unwrap(); + let segment = indices.iter().find(|index| index.name == name)?; + Some( + segment + .fragment_bitmap + .as_ref() + .expect("coverage should be recorded") + .iter() + .collect(), + ) + } + + /// Control for the test below: when the rewrite is the *later* commit, the + /// coverage is repaired. `TombstoneFieldData::apply` rebinds the field, + /// which prunes the fragment from every index the manifest carries. + #[tokio::test] + async fn test_a_rewrite_prunes_an_index_segment_committed_before_it() { + let dataset = test_dataset(false).await; + let file = existing_data_file(&dataset, 0); + + let dataset = commit( + dataset, + vec![Action::AddIndexSegment(index_segment( + "by_a", + vec![Ref::Committed(0)], + ))], + ) + .await; + assert_eq!(index_coverage(&dataset, "by_a").await, vec![0]); + + let dataset = commit(dataset, rewrite_field_zero(file)).await; + + assert_eq!( + index_coverage_opt(&dataset, "by_a") + .await + .unwrap_or_default(), + Vec::::new(), + "the rewrite should have dropped fragment 0 from the coverage" + ); + } + + /// The same two commits in the other order must reach the same place. + /// + /// The format spec's third invalidation case -- a fragment has had one of + /// the index's columns updated in place -- "cannot be detected just by + /// examining metadata", so a reader trusts `fragment_bitmap` and has no + /// recourse. Keeping it honest is a write-path obligation. + /// + /// Apply-time pruning only discharges that obligation when the rewrite is + /// the later commit, as the control above shows. A segment built before the + /// rewrite but committed after it carries no rebind of its own, so nothing + /// prunes it -- and its footprint records only an index claim, which is + /// never compared against the `FieldData` coordinates the rewrite writes. + #[tokio::test] + async fn test_an_index_segment_staged_before_a_rewrite_does_not_cover_it() { + let dataset = test_dataset(false).await; + let stale = Arc::new(dataset.clone()); + let file = existing_data_file(&dataset, 0); + + // The rewrite wins the race. + commit(dataset, rewrite_field_zero(file)).await; + + // The loser built its segment from the values the rewrite replaced. + let committed = commit_from_stale( + stale, + vec![Action::AddIndexSegment(index_segment( + "by_a", + vec![Ref::Committed(0)], + ))], + ) + .await; + + // Rejected, because the segment requires the data it describes and the + // rewrite wrote it. Dropping fragment 0 from the arriving segment's + // coverage would also be sound and would waste less work -- it is what + // the legacy path does for `Operation::Update` -- but an action set is + // never rewritten to move to a newer version, so there is nowhere to + // put that here. + let error = committed.expect_err("the segment describes replaced values"); + assert!( + matches!(error, Error::RetryableCommitConflict { .. }), + "{error:?}" + ); + } + + #[tokio::test] + async fn test_a_commit_extends_an_index_segments_coverage() { + let dataset = test_dataset(false).await; + let segment = index_segment("by_a", vec![Ref::Committed(0)]); + let uuid = segment.uuid; + let dataset = commit(dataset, vec![Action::AddIndexSegment(segment)]).await; + assert_eq!(index_coverage(&dataset, "by_a").await, vec![0]); + + let dataset = commit( + dataset, + vec![Action::AdjustIndexCoverage(AdjustIndexCoverage { + uuid, + name: "by_a".into(), + add_fragments: vec![Ref::Committed(1)], + remove_fragments: vec![0], + })], + ) + .await; + + assert_eq!(index_coverage(&dataset, "by_a").await, vec![1]); + } }