diff --git a/.github/labeler-area.yml b/.github/labeler-area.yml index 998e752fba6..2045aab45c5 100644 --- a/.github/labeler-area.yml +++ b/.github/labeler-area.yml @@ -44,7 +44,7 @@ A-format: - "protos/index_old.proto" - "protos/rowids.proto" - "protos/table.proto" - - "protos/transaction.proto" + - "protos/transaction/*.proto" - "docs/src/format/**" # Drives the format-spec vote gate (.github/workflows/format-vote-gate.yml): @@ -62,7 +62,7 @@ format-change: - "protos/index_old.proto" - "protos/rowids.proto" - "protos/table.proto" - - "protos/transaction.proto" + - "protos/transaction/*.proto" - "docs/src/format/**" # Lockfiles are intentionally not excluded: a pure dependency bump gets both diff --git a/ci/test_labeler_area.py b/ci/test_labeler_area.py index ae4b74bd806..48b5aef1664 100644 --- a/ci/test_labeler_area.py +++ b/ci/test_labeler_area.py @@ -24,11 +24,19 @@ def test_format_labels_use_the_same_paths(): def test_only_persisted_protos_are_format_changes(): + """Every persisted proto is gated, and nothing else is. + + The config entries are globs and protos live in subdirectories, so both + sides are resolved to concrete files before being compared. + """ detected_proto_paths = { - path for path in paths_for("format-change") if path.startswith("protos/") + matched.relative_to(ROOT).as_posix() + for glob in paths_for("format-change") + if glob.startswith("protos/") + for matched in ROOT.glob(glob) } all_proto_paths = { - path.relative_to(ROOT).as_posix() for path in (ROOT / "protos").glob("*.proto") + path.relative_to(ROOT).as_posix() for path in (ROOT / "protos").rglob("*.proto") } assert detected_proto_paths == all_proto_paths - EXECUTION_PROTO_PATHS diff --git a/docs/src/format/table/transaction.md b/docs/src/format/table/transaction.md index c88c3170988..f7de5e2d54a 100644 --- a/docs/src/format/table/transaction.md +++ b/docs/src/format/table/transaction.md @@ -47,7 +47,7 @@ For detailed conflict detection and resolution mechanisms, see the [Conflict Res ## Transaction Types -The authoritative specification for transaction types is defined in [`protos/transaction.proto`](https://github.com/lancedb/lance/blob/main/protos/transaction.proto). +The authoritative specification for transaction types is defined in [`protos/transaction/`](https://github.com/lancedb/lance/tree/main/protos/transaction). Each transaction contains a `read_version` field indicating the table version from which the transaction was built, a `uuid` field uniquely identifying the transaction, and an `operation` field specifying one of the following transaction types: @@ -556,6 +556,97 @@ An UpdateBases operation only modifies the base paths. As a result, it only conf UpdateBases operations and even then only conflicts if the two operations have base paths with the same id, name, or path. +## Action-Based Transactions (Transaction V2) + +!!! warning "Experimental" + + Transaction V2 is an experimental format feature under the process in + [Voting](../../community/voting.md). Its protobuf field numbers are not a + stable contract, breaking changes may be made without a separate vote, and + the feature is removed from the specification and the codebase if its + stabilization vote does not pass. Discussion: + [#5960](https://github.com/lance-format/lance/discussions/5960). + +!!! danger "Writing Transaction V2 breaks compatibility with readers before v12.0.0" + + A table version committed as a `CompositeOperation` **cannot be opened** by + Lance before v12.0.0 — not merely read as history. Those releases decode the + manifest's inline transaction section while opening the table, and an + operation they do not recognize decodes to no operation at all, which fails + the open. The fix + ([#7740](https://github.com/lance-format/lance/pull/7740)) shipped in + v12.0.0 and is not expected to be backported to earlier release lines; see + [#9454](https://github.com/lance-format/lance/issues/9454). + + Treat this as durable, not transitional. Writing a `CompositeOperation` is + opt-in for exactly this reason, and a table that has one anywhere in its + history raises the minimum reader version for that table to v12.0.0. + +Every operation above is a single named verb. A transaction may instead carry a +`CompositeOperation`: an ordered list of granular *actions* that apply +atomically as one manifest change. This lets one commit express a change no +single named operation covers, such as appending a fragment and updating an +index together, and makes the transaction a true diff of the manifest rather +than a description of intent. + +
+CompositeOperation protobuf message + +```protobuf +%%% proto.message.CompositeOperation %%% +``` + +
+ +A `CompositeOperation` holds `UserAction`s, each a human-recognizable step with +a description, which in turn hold the granular `Action`s applied to the +manifest. The two levels exist so the transaction history stays readable when +several operations are squashed into one: the descriptions survive, while the +action lists are flattened on apply. + +Actions are deltas rather than post-images, and an action that allocates a new +identifier (a field, fragment, or base id) names it with a local placeholder +that resolves against the target's counters at apply time. Both properties are +what let a composite operation be replayed against a newer version, or onto a +different branch, without rewriting it. + +The full action vocabulary and the reasoning behind its shape are defined in +`protos/transaction/actions.proto`. + +One vocabulary detail is a format contract rather than an implementation +choice. `AddIndexSegment.fields` names only the columns the segment is keyed +on, and `covering_fields` is an independent declaration of the columns whose +values it carries, so the segment'''s dependency set is the union of the two. +This is the contract an index action targets, not the legacy one where +`IndexMetadata.fields` means keyed columns followed by carried columns. A +declaration whose two lists overlap can only be represented by a manifest +declaring `FLAG_INDEPENDENT_COVERING_FIELDS`; until a release implements that +flag, apply rejects the overlapping form and lowers a disjoint one to the +legacy subset representation. + +### Compatibility + +A `CompositeOperation` produces an ordinary manifest, so a reader that scans a +table whose latest version was written this way needs no knowledge of the +feature. Readers are affected in two places: + +- A reader that decodes the transaction itself — reading the transaction + history, or checking a concurrent commit for conflicts — sees an operation it + does not recognize. Implementations must reject it rather than treat it as a + no-op, so that a concurrent V2 commit aborts an in-flight commit instead of + being silently skipped. +- Transactions are usually inlined into the manifest. A reader that fails when + an inline transaction does not decode cannot open such a table at all. + Implementations must tolerate an undecodable inline transaction and continue + opening the table, because the transaction contents are not needed to read + data. In the Rust implementation this has been true only since v12.0.0, which + is what makes writing Transaction V2 a compatibility break against earlier + releases rather than a graceful degradation. + +Conflict resolution between two `CompositeOperation`s is computed from the +actions themselves. Between a `CompositeOperation` and any named operation it +fails closed: the commit is rejected as a conflict rather than compared. + ## Conflict Resolution ### Terminology diff --git a/protos/transaction/actions.proto b/protos/transaction/actions.proto new file mode 100644 index 00000000000..bb7b5ff4b87 --- /dev/null +++ b/protos/transaction/actions.proto @@ -0,0 +1,543 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +syntax = "proto3"; + +import "file.proto"; +import "table.proto"; +import "transaction/common.proto"; +import "google/protobuf/any.proto"; + +package lance.table; + +/* + * Action-based transactions (Transaction V2) — EXPERIMENTAL. + * + * A `CompositeOperation` replaces the legacy `Operation` (see + * transaction.proto) with an ordered list of `Action`s committing + * atomically as one manifest change. This is an experimental feature under + * the process in docs/src/community/voting.md: the field numbers are not a + * stable contract, may change without a separate vote, and are removed if + * the stabilization vote does not pass. Support here is read-side + * fail-closed — rejected on load, with no write path. + * + * COMPATIBILITY WARNING. A table version carrying one of these cannot be + * opened at all by Lance before v12.0.0, and the fix is not backported + * (issue #9454). Writing V2 is opt-in for that reason. + * + * Why this shape: + * + * 1. Actions are deltas, not post-images: each records *the change* to the + * manifest, not the resulting state. That is what lets actions compose + * into one commit and replay onto another branch with no special cases. + * + * 2. `Ref = Committed | Local`. A minting action cannot know its + * counter-allocated id (field, fragment, base) until commit, and gets a + * different one on replay onto another target, so it names the id with + * a `Local` token resolved against the target's counters at apply. A + * writer that must bake an id into a file it writes in the same commit + * instead takes `Committed` ids from a range an earlier commit reserved + * (`ReserveFragmentIds`, `ReserveRowIds`). + * + * 3. Schema changes are per field (`AddField` / `DropField` / `AlterField`), + * never a wholesale replacement, so disjoint edits commute. `AlterField` + * preserves the field id and makes each mutable facet an independent + * optional, so concurrent alters of different facets commute too. + * + * 4. Assertions carry only preconditions no post-image can reconstruct, e.g. + * the keys a merge-insert inserted (`AssertUniqueKeys`). + * + * Off the wire and recomputed instead: conflict footprints, a pure function + * of each action, and large row-level deltas recoverable from the + * read-version state. + */ + +/* + * A reference to a counter-allocated identifier (field id, fragment id, or base + * id) that may not be committed yet. + * + * `committed` is an already-assigned id. `local` is a placeholder token minted + * by an `Add*` action earlier in the same `CompositeOperation`; it resolves to a + * freshly-allocated committed id at apply, and re-resolves against the target's + * counters on merge/rebase. `local` tokens are scoped to a single + * `CompositeOperation` and must be distinct within it (validated by the writer once + * the write path exists). + */ +message Ref { + oneof kind { + uint64 committed = 1; + uint32 local = 2; + } +} + +/* + * A user-facing, composable transaction: an ordered list of user actions that + * commit atomically as a single manifest change. + */ +message CompositeOperation { + // The ordered list of user actions applied by this operation. + repeated UserAction actions = 1; +} + +/* + * A single user-recognizable step within a CompositeOperation (e.g. "append batch", + * "rebuild index"). + * + * The description keeps the transaction history human-readable. When a range of + * transactions is squashed, each original CompositeOperation collapses into one + * UserAction so the readable sequence survives; the action lists are flattened + * when applied to the manifest. + */ +message UserAction { + // Human-readable description of this step. + string description = 1; + // The granular manifest changes this step expands to. + repeated Action actions = 2; +} + +/* + * A single granular change (a delta) or assertion. Legacy operations decompose + * into an ordered list of these. + */ +message Action { + oneof action { + // -- Minting actions (allocate a new counter-based id via a Local token) -- + AddFragment add_fragment = 1; + AddField add_field = 3; + AddBase add_base = 4; + // -- Fragment / data-file deltas -- + AddDataFile add_data_file = 2; + TombstoneFieldData tombstone_field_data = 5; + RemoveFragment remove_fragment = 6; + // -- Reference-stable changes (committed coordinates; post-image at rest) -- + SetDeletionFile set_deletion_file = 7; + ConfigUpdate config_update = 8; + AddOverlays add_overlays = 9; + RefreshRowVersionMetadata refresh_row_version_metadata = 10; + /* Named `update_compacted_ss_tables` rather than `..._sstables` so the + * generated oneof variant matches the message name. + */ + UpdateCompactedSsTables update_compacted_ss_tables = 11; + // -- Schema (field-level; no wholesale SetSchema) -- + DropField drop_field = 12; + AlterField alter_field = 13; + // -- Index segments -- + AddIndexSegment add_index_segment = 14; + RemoveIndexSegment remove_index_segment = 15; + AdjustIndexCoverage adjust_index_coverage = 16; + // -- Wholesale-within-operation -- + ResetTable reset_table = 17; + ReserveFragmentIds reserve_fragment_ids = 18; + ReserveRowIds reserve_row_ids = 20; + // -- Assertions (preconditions, not deltas) -- + AssertUniqueKeys assert_unique_keys = 19; + } +} + +/* + * Add a new, empty fragment. + * + * Its data files arrive via AddDataFile actions referencing this fragment by + * the same `id`. A newly added fragment has no deletion vector; deletions are + * applied by a later operation via SetDeletionFile. + */ +message AddFragment { + /* + * The id to add the fragment at. + * + * Local is the usual form: a placeholder token resolved to a freshly + * allocated id at apply. Committed names an id out of a range an earlier + * commit reserved with ReserveFragmentIds, for a writer that must know the + * fragment's id before committing (see rationale 2); apply rejects an id + * that is not reserved or is already in use. + * + * AddField and AddBase have no Committed form, because no reservation + * mechanism exists for field or base ids. + */ + Ref id = 1; + // Number of physical rows (including rows later tombstoned). + uint64 physical_rows = 2; + /* + * Stable-row-id and row-version sequences, carried exactly as on + * DataFragment. Absent on datasets without stable row ids. + */ + oneof row_id_sequence { + bytes inline_row_ids = 3; + ExternalFile external_row_ids = 4; + } + oneof last_updated_at_version_sequence { + bytes inline_last_updated_at_versions = 5; + ExternalFile external_last_updated_at_versions = 6; + } + oneof created_at_version_sequence { + bytes inline_created_at_versions = 7; + ExternalFile external_created_at_versions = 8; + } + /* + * false => rearrangement only (rows unchanged, e.g. compaction); CDC and + * streaming consumers may skip it. Absent is treated as true (real change). + */ + optional bool data_change = 9; +} + +/* + * Add a data file to a fragment. + */ +message AddDataFile { + // The fragment to add the file to (Committed, or a same-op Local fragment). + Ref fragment = 1; + /* + * The data file. Its `fields` are left unset (-1) and stamped in at apply + * once `field_ids` resolve; `field_ids` below is the authority for the + * column -> field mapping. + */ + DataFile file = 2; + /* + * One entry per column in `file`, in order. Committed for existing fields, + * Local for fields minted by an AddField in the same operation. + */ + repeated Ref field_ids = 3; + // Data-change marker; see AddFragment.data_change. + optional bool data_change = 4; +} + +/* + * Mint a new schema field. + * + * A struct/list column that introduces several fields is expressed as several + * ordered AddField actions (parent before children, children referencing the + * parent via a Local ref). + */ +message AddField { + // Placeholder token for the new field id, resolved at apply. + uint32 local = 1; + /* + * The parent field: Committed to reparent under an existing field, Local for + * a sibling minted in the same operation. Absent => top-level column. + */ + optional Ref parent = 2; + // Field definition; its `id` and `parent_id` are ignored (see `local`, `parent`). + lance.file.Field def = 3; +} + +/* + * Mint a new base path. + */ +message AddBase { + // Placeholder token for the base id, resolved at apply. + uint32 local = 1; + // Base path; its `id` is ignored and stamped in at apply. + BasePath base = 2; +} + +/* + * Tombstone the data-file binding of one or more committed fields. + * + * Each field's slot in whatever data file currently backs it is set to -2, and + * any file left with no live fields is pruned at apply. This is how a column's + * data is dropped or superseded (data files have no id of their own; a live + * field is backed by exactly one file). A column re-encode is + * TombstoneFieldData(X) followed by AddDataFile(new file, field X). + */ +message TombstoneFieldData { + Ref fragment = 1; + // The fields whose current backing is tombstoned. + repeated Ref field_ids = 2; + // Data-change marker; see AddFragment.data_change. + optional bool data_change = 3; +} + +/* + * Remove a fragment entirely (all rows deleted, or replaced by compaction). + */ +message RemoveFragment { + Ref fragment = 1; + // Data-change marker; see AddFragment.data_change. + optional bool data_change = 2; +} + +/* + * Set (replace) a fragment's deletion file. + * + * Reference-stable: the fragment id is committed and physical row offsets are + * stable, so the post-image deletion file is safe to store. The newly-deleted + * rows (the delta used for rebase and conflict) are derived by diffing against + * the read-version deletion file and are not serialized here. + */ +message SetDeletionFile { + /* The fragment, by committed id: unlike the sibling fragment actions this + * takes no Ref, because a fragment minted in the same operation has no rows + * to delete yet (see AddFragment). + */ + uint64 fragment = 1; + DeletionFile deletion_file = 2; + // Data-change marker; see AddFragment.data_change. + optional bool data_change = 3; +} + +/* + * Apply config / metadata updates. + * + * Reference-stable key-merges over stable coordinates (config keys, field ids). + * Mirrors the new-style fields of the legacy UpdateConfig operation. + */ +message ConfigUpdate { + optional UpdateMap config = 1; + optional UpdateMap table_metadata = 2; + optional UpdateMap schema_metadata = 3; + /* + * Per-field metadata updates, keyed by Ref so a field minted in the same + * operation (Local) can receive metadata. + */ + repeated FieldMetadata field_metadata = 4; + + message FieldMetadata { + Ref field = 1; + UpdateMap updates = 2; + } +} + +/* + * Append overlay files to fragments (see DataOverlayFile in table.proto). + * + * Reference-stable: overlays target committed fragments/offsets/fields and are + * appended, never replacing existing overlays. + */ +message AddOverlays { + Ref fragment = 1; + repeated DataOverlayFile overlays = 2; + // Data-change marker; see AddFragment.data_change. + optional bool data_change = 3; +} + +/* + * Refresh row-version (created-at / last-updated-at) metadata for fragments + * whose columns were merged in place, mirroring the implicit refresh Merge + * performs on stable-row-id datasets. + */ +message RefreshRowVersionMetadata { + repeated uint64 fragment_ids = 1; +} + +/* + * Mark MemWAL SSTables as compacted into the base table. Covers + * UpdateMemWalState and the compaction bookkeeping of Update. + */ +message UpdateCompactedSsTables { + repeated CompactedSsTable compacted_sstables = 1; +} + +/* + * Remove a field from the schema (and, at apply, prune any data file left with + * no live fields). References an existing committed field id. + */ +message DropField { + Ref field = 1; +} + +/* + * Alter facets of an existing field in place, preserving its identity (id). + * + * Each facet is optional: present means "change this facet", absent means + * "leave unchanged". Keying conflict on (field id, facet) lets independent + * facet changes to the same field commute (e.g. a cast concurrent with a + * nullability change). A cast additionally needs TombstoneFieldData plus a new + * AddDataFile to rewrite the data. New facets are added as optional fields. + */ +message AlterField { + Ref field = 1; + optional string name = 2; + // New Arrow logical type (see Field.logical_type). The cast. + optional string logical_type = 3; + optional bool nullable = 4; +} + +/* + * Add an index segment. + * + * A brand-new index is expressed as its first segment; the set of segments + * sharing `name` constitutes one logical index. + */ +message AddIndexSegment { + UUID uuid = 1; + string name = 2; + /* + * The field ids this segment is keyed on -- the columns it can be searched + * by (Committed, or Local for same-op minted fields). Empty only for a + * system index keyed on no column at all, such as mem_wal or frag_reuse. + * + * Keys only. `IndexMetadata.fields` in a manifest written under the legacy + * contract instead means keyed columns followed by carried ones; this + * action does not inherit that, because it has no historical encoding to + * stay compatible with. Apply derives the manifest's representation. + */ + repeated Ref fields = 3; + /* + * Field ids whose values the segment carries but is not necessarily keyed + * on. Independent of `fields` rather than a subset of it, and the segment's + * dependency set is the union of the two. Empty for a segment that carries + * no extra column. + * + * An id may appear in both lists, meaning a column the segment is keyed on + * *and* carries values for. That form requires the manifest to declare + * FLAG_INDEPENDENT_COVERING_FIELDS, which no release implements yet, so + * apply rejects an overlapping declaration until it does. A disjoint + * declaration lowers to the legacy subset form and commits today. + */ + repeated Ref covering_fields = 12; + google.protobuf.Any index_details = 4; + optional int32 index_version = 5; + /* + * Fragments covered by this segment (resolved to a fragment_bitmap at apply, + * and remapped under fragment relocation). + * + * Absent means the coverage is unknown, which is what the system indices + * (MemWAL, fragment reuse) carry and what pre-bitmap segments read back as. + * That is a different statement from a present-but-empty list, which says the + * segment covers no fragment at all -- the query path serves a segment of + * unknown coverage and skips one that covers nothing. + */ + optional FragmentCoverage covered_fragments = 6; + /* + * Index files with sizes, when produced by the writer (e.g. a compaction that + * rewrites the segment). Empty when unavailable. + */ + repeated IndexFile files = 7; + /* Data-change marker; see AddFragment.data_change. false marks a segment + * rebuild/compaction that does not reflect any change to the indexed data. + */ + optional bool data_change = 8; + /* + * The base path this segment's files live under, for a segment imported from + * another dataset. Committed, or Local for a base minted by an AddBase in the + * same operation. Absent => the dataset's own index directory. + */ + optional Ref base = 9; + /* + * When this segment was built, in milliseconds since the Unix epoch. Absent + * for a segment whose build time was not recorded. + */ + optional uint64 created_at = 10; + /* + * The dataset version whose data this segment reflects. + * + * This is a correctness gate, not provenance: an overlay committed at or + * before it is treated as already folded into the index. A segment merged + * from several older ones reflects only as much as its oldest input, so this + * is genuinely below the version the operation reads and cannot be derived + * from it. Absent means the version the operation reads, which is what a + * freshly built segment reflects. It may never exceed that version. + */ + optional uint64 dataset_version = 11; +} + +/* + * A list of fragments, wrapped so that "no coverage recorded" is distinguishable + * from "covers nothing" (a bare `repeated` field cannot tell them apart). + */ +message FragmentCoverage { + // Committed, or Local for a fragment minted in the same operation. + repeated Ref fragments = 1; +} + +/* + * Remove an index segment by uuid. + */ +message RemoveIndexSegment { + UUID uuid = 1; + /* Data-change marker; see AddFragment.data_change. false marks removal of a + * segment superseded by compaction, not a change to the indexed data. + */ + optional bool data_change = 2; + /* + * The logical index the segment belongs to. Carried so apply can check the + * removal was planned against the segment it actually names; a mismatch means + * the operation was built against a different set of segments. + */ + string name = 3; +} + +/* + * Adjust the fragment coverage of an existing index segment without rewriting + * it. + * + * NOTE: coverage representation is an acknowledged open design area; this shape + * is provisional pending how much coverage change is derivable from fragment + * relocation versus stated explicitly here. + */ +message AdjustIndexCoverage { + UUID uuid = 1; + // Fragments to add to the segment's coverage (Committed or same-op Local). + repeated Ref add_fragments = 2; + // Committed fragment ids to drop from the segment's coverage. + repeated uint64 remove_fragments = 3; + /* + * The logical index the segment belongs to. Conflict detection compares + * coverage across the segments of one index, so it needs the index this + * adjustment widens, not just the segment. Apply also checks the named + * segment really carries this name. + */ + string name = 4; +} + +/* + * Reset the table to an empty state, in preparation for a fresh schema and data + * written by later actions in the same operation (the decomposition of a full + * Overwrite / CREATE OR REPLACE). + * + * Drops the entire schema, schema metadata, all fragments, and all indices; + * preserves table config, table metadata, and base paths (change those via + * ConfigUpdate / AddBase in the same operation). Its write footprint spans the + * whole schema, all fragments, and all indices, so it conflicts with any + * concurrent change to them (DROP TABLE-style exclusivity). + */ +message ResetTable {} + +/* + * Reserve a contiguous range of fragment ids from the counter, for a later + * (possibly distributed) writer to populate. Fragments written against a + * reserved range reference those ids as Committed. + */ +message ReserveFragmentIds { + uint32 count = 1; +} + +/* + * Reserve a contiguous range of stable row ids from the counter, for a later + * writer to populate. + * + * The counterpart of ReserveFragmentIds for the row id space, and needed + * alongside it: reserving a fragment id fixes the high half of a row address, + * but on a dataset with stable row ids an index records row ids instead, and + * those come off their own counter. + * + * The reserved range is the `count` ids ending at the committed manifest's + * next_row_id, so the reserving writer learns which ids it got by reading that + * back. A fragment written against the range carries those ids in its row id + * sequence, which apply then leaves as it finds them. + * + * An error on a dataset that does not use stable row ids, which has no row id + * counter to reserve from. + */ +message ReserveRowIds { + uint64 count = 1; +} + +/* + * An assertion (precondition), not a delta: the keys this operation inserts + * must not collide with keys inserted by a concurrent commit. + * + * Carries a bloom / exact-set filter of the inserted key hashes (not derivable + * from any post-image); conflict = the two filters intersect. Home for + * merge-insert's strict primary-key conflict detection. + */ +message AssertUniqueKeys { + /* + * Field ids of the key columns (unenforced primary key). Committed, or Local + * for same-op minted fields. This is the authoritative key-column list; + * `filter.field_ids` (an artifact of the shared KeyExistenceFilter type) is + * ignored here and left empty. + */ + repeated Ref key_fields = 1; + KeyExistenceFilter filter = 2; +} diff --git a/protos/transaction/common.proto b/protos/transaction/common.proto new file mode 100644 index 00000000000..ad290bb485a --- /dev/null +++ b/protos/transaction/common.proto @@ -0,0 +1,97 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +syntax = "proto3"; + +package lance.table; + +/* + * Shared transaction building blocks, referenced by both the legacy + * operations in transaction.proto and the action-based transactions in + * actions.proto. + * + * These types are transaction *machinery* (metadata merges, conflict-detection + * filters), not persisted table state, so they live here rather than in + * table.proto. Kept in a dedicated file so both transaction models can import + * them without a circular dependency (transaction.proto imports actions.proto, + * so actions.proto must not import transaction.proto). + */ + +/* + * An entry for a map update. + * + * If `value` is not set, the key is removed from the map. + */ +message UpdateMapEntry { + // The key of the map entry to update. + string key = 1; + // The value to set for the key. Absent removes the key. + optional string value = 2; +} + +/* + * A set of updates to apply to a string map (table config, table metadata, + * schema metadata, or per-field metadata). + */ +message UpdateMap { + repeated UpdateMapEntry update_entries = 1; + /* + * If true, the map is replaced entirely with these entries. If false, the + * entries are merged into the existing map. + */ + bool replace = 2; +} + +/* + * Exact set of key hashes for conflict detection. + * + * Used when the number of inserted rows is small. + */ +message ExactKeySetFilter { + // 64-bit hashes of the inserted row keys. + repeated uint64 key_hashes = 1; +} + +/* + * Bloom filter for key existence tests. + * + * Used when the number of inserted rows is large. + */ +message BloomFilter { + // Bitset backing the bloom filter (SBBF format). + bytes bitmap = 1; + // Number of bits in the bitmap. + uint32 num_bits = 2; + /* + * Number of items the filter was sized for. Used for intersection validation: + * filters with different sizes cannot be compared. Default: 8192. + */ + uint64 number_of_items = 3; + /* + * False positive probability the filter was sized for. Used for intersection + * validation: filters with different parameters cannot be compared. + * Default: 0.00057. + */ + double probability = 4; +} + +/* + * A filter for checking key existence in the set of rows inserted by a + * merge-insert operation. + * + * Only created when the merge-insert's ON columns match the schema's unenforced + * primary key; its presence indicates strict primary-key conflict detection. + * Conflict detection intersects two filters, so the underlying representation + * is either an exact set (small row counts) or a Bloom filter (large counts). + */ +message KeyExistenceFilter { + // Field ids of the columns participating in the key (the unenforced primary key). + repeated int32 field_ids = 1; + // The underlying data structure storing the key hashes. + oneof data { + // Exact set of key hashes (small number of rows). + ExactKeySetFilter exact = 2; + // Bloom filter (large number of rows). + BloomFilter bloom = 3; + } +} diff --git a/protos/transaction.proto b/protos/transaction/transaction.proto similarity index 85% rename from protos/transaction.proto rename to protos/transaction/transaction.proto index 6e01831b4af..7f020a827b6 100644 --- a/protos/transaction.proto +++ b/protos/transaction/transaction.proto @@ -5,6 +5,8 @@ syntax = "proto3"; import "file.proto"; import "table.proto"; +import "transaction/actions.proto"; +import "transaction/common.proto"; import "google/protobuf/any.proto"; package lance.table; @@ -220,51 +222,6 @@ message Transaction { optional string branch_name = 5; } - /* Exact set of key hashes for conflict detection. - * Used when the number of inserted rows is small. - */ - message ExactKeySetFilter { - // 64-bit hashes of the inserted row keys. - repeated uint64 key_hashes = 1; - } - - /* Bloom filter for key existence tests. - * Used when the number of rows is large. - */ - message BloomFilter { - // Bitset backing the bloom filter (SBBF format). - bytes bitmap = 1; - // Number of bits in the bitmap. - uint32 num_bits = 2; - /* Number of items the filter was sized for. - * Used for intersection validation (filters with different sizes cannot be compared). - * Default: 8192 - */ - uint64 number_of_items = 3; - /* False positive probability the filter was sized for. - * Used for intersection validation (filters with different parameters cannot be compared). - * Default: 0.00057 - */ - double probability = 4; - } - - /* A filter for checking key existence in set of rows inserted by a merge insert operation. - * Only created when the merge insert's ON columns match the schema's unenforced primary key. - * The presence of this filter indicates strict primary key conflict detection should be used. - * Can use either an exact set (for small row counts) or a Bloom filter (for large row counts). - */ - message KeyExistenceFilter { - // Field IDs of columns participating in the key (must match unenforced primary key). - repeated int32 field_ids = 1; - // The underlying data structure storing the key hashes. - oneof data { - // Exact set of key hashes (used for small number of rows). - ExactKeySetFilter exact = 2; - // Bloom filter (used for large number of rows). - BloomFilter bloom = 3; - } - } - // Serialized as sorted distinct local physical row offsets within the fragment (0-based). message UInt32List { repeated uint32 values = 1; @@ -321,22 +278,6 @@ message Transaction { REWRITE_COLUMNS = 1; } - // An entry for a map update. If value is not set, the key will be removed from the map. - message UpdateMapEntry { - // The key of the map entry to update. - string key = 1; - // The value to set for the key. - optional string value = 2; - } - - message UpdateMap { - repeated UpdateMapEntry update_entries = 1; - /* If true, the map will be replaced entirely with the new entries. - * If false, the new entries will be merged with the existing map. - */ - bool replace = 2; - } - /* An operation that updates the table config, table metadata, schema metadata, * or field metadata. */ @@ -425,6 +366,18 @@ message Transaction { Clone clone = 113; UpdateBases update_bases = 114; DataOverlay data_overlay = 115; + /* Action-based transaction (Transaction V2). See actions.proto. + * + * EXPERIMENTAL: currently rejected on load; no write path. + * + * A table version carrying this arm cannot be opened at all by Lance + * before v12.0.0. Those releases decode the inline transaction while + * opening the table, and an arm they do not know decodes to no operation, + * which fails the open (issue #9454). The fix landed in v12.0.0 and is not + * backported, so this is a durable compatibility break rather than a + * transitional one. + */ + CompositeOperation composite_operation = 116; } // Fields 200/202 (`blob_append` / `blob_overwrite`) previously represented blob dataset ops. diff --git a/rust/lance-table/build.rs b/rust/lance-table/build.rs index 03216636b30..103109adfdb 100644 --- a/rust/lance-table/build.rs +++ b/rust/lance-table/build.rs @@ -19,7 +19,9 @@ fn main() -> Result<()> { prost_build.compile_protos( &[ "./protos/table.proto", - "./protos/transaction.proto", + "./protos/transaction/common.proto", + "./protos/transaction/actions.proto", + "./protos/transaction/transaction.proto", "./protos/rowids.proto", ], &["./protos"], diff --git a/rust/lance-table/src/format/key_existence.rs b/rust/lance-table/src/format/key_existence.rs index 210ef5f3836..da88deb62ed 100644 --- a/rust/lance-table/src/format/key_existence.rs +++ b/rust/lance-table/src/format/key_existence.rs @@ -131,18 +131,16 @@ impl KeyExistenceFilterBuilder { } } -impl From<&KeyExistenceFilterBuilder> for pb::transaction::KeyExistenceFilter { +impl From<&KeyExistenceFilterBuilder> for pb::KeyExistenceFilter { fn from(builder: &KeyExistenceFilterBuilder) -> Self { Self { field_ids: builder.field_ids.clone(), - data: Some(pb::transaction::key_existence_filter::Data::Bloom( - pb::transaction::BloomFilter { - bitmap: builder.sbbf.to_bytes(), - num_bits: (builder.sbbf.size_bytes() as u32) * 8, - number_of_items: BLOOM_FILTER_DEFAULT_NUMBER_OF_ITEMS, - probability: BLOOM_FILTER_DEFAULT_PROBABILITY, - }, - )), + data: Some(pb::key_existence_filter::Data::Bloom(pb::BloomFilter { + bitmap: builder.sbbf.to_bytes(), + num_bits: (builder.sbbf.size_bytes() as u32) * 8, + number_of_items: BLOOM_FILTER_DEFAULT_NUMBER_OF_ITEMS, + probability: BLOOM_FILTER_DEFAULT_PROBABILITY, + })), } } } @@ -212,13 +210,13 @@ impl KeyExistenceFilter { } } -impl From<&KeyExistenceFilter> for pb::transaction::KeyExistenceFilter { +impl From<&KeyExistenceFilter> for pb::KeyExistenceFilter { fn from(filter: &KeyExistenceFilter) -> Self { match &filter.filter { FilterType::ExactSet(hashes) => Self { field_ids: filter.field_ids.clone(), - data: Some(pb::transaction::key_existence_filter::Data::Exact( - pb::transaction::ExactKeySetFilter { + data: Some(pb::key_existence_filter::Data::Exact( + pb::ExactKeySetFilter { key_hashes: hashes.iter().copied().collect(), }, )), @@ -230,28 +228,26 @@ impl From<&KeyExistenceFilter> for pb::transaction::KeyExistenceFilter { probability, } => Self { field_ids: filter.field_ids.clone(), - data: Some(pb::transaction::key_existence_filter::Data::Bloom( - pb::transaction::BloomFilter { - bitmap: bitmap.clone(), - num_bits: *num_bits, - number_of_items: *number_of_items, - probability: *probability, - }, - )), + data: Some(pb::key_existence_filter::Data::Bloom(pb::BloomFilter { + bitmap: bitmap.clone(), + num_bits: *num_bits, + number_of_items: *number_of_items, + probability: *probability, + })), }, } } } -impl TryFrom<&pb::transaction::KeyExistenceFilter> for KeyExistenceFilter { +impl TryFrom<&pb::KeyExistenceFilter> for KeyExistenceFilter { type Error = lance_core::Error; - fn try_from(message: &pb::transaction::KeyExistenceFilter) -> Result { + fn try_from(message: &pb::KeyExistenceFilter) -> Result { let filter = match message.data.as_ref() { - Some(pb::transaction::key_existence_filter::Data::Exact(exact)) => { + Some(pb::key_existence_filter::Data::Exact(exact)) => { FilterType::ExactSet(exact.key_hashes.iter().copied().collect()) } - Some(pb::transaction::key_existence_filter::Data::Bloom(b)) => { + Some(pb::key_existence_filter::Data::Bloom(b)) => { // Use defaults for backwards compatibility let number_of_items = if b.number_of_items == 0 { BLOOM_FILTER_DEFAULT_NUMBER_OF_ITEMS diff --git a/rust/lance-table/src/transaction/proto.rs b/rust/lance-table/src/transaction/proto.rs index c51c8f16719..ef6ae7bc5e9 100644 --- a/rust/lance-table/src/transaction/proto.rs +++ b/rust/lance-table/src/transaction/proto.rs @@ -401,6 +401,19 @@ impl TryFrom for Transaction { .map(DataOverlayGroup::try_from) .collect::>>()?, }, + Some(pb::transaction::Operation::CompositeOperation(_)) => { + // Action-based transactions (Transaction V2) are an experimental + // wire format (OSS-1530). This version of Lance recognizes the message + // but has no support for it: reject on load, fail-closed. Because + // load_and_sort_new_transactions collects transactions with + // try_collect, a concurrent V2 commit in the conflict window + // aborts the whole commit rather than being silently skipped. + // Do NOT make this parsing lenient. + return Err(Error::not_supported( + "action-based transactions (Transaction V2) are not supported \ + by this version of Lance; please upgrade", + )); + } None => { return Err(Error::internal( "Transaction message did not contain an operation".to_string(), @@ -652,20 +665,12 @@ impl From<&Transaction> for pb::Transaction { schema_metadata_updates, field_metadata_updates, } => pb::transaction::Operation::UpdateConfig(pb::transaction::UpdateConfig { - config_updates: config_updates - .as_ref() - .map(pb::transaction::UpdateMap::from), - table_metadata_updates: table_metadata_updates - .as_ref() - .map(pb::transaction::UpdateMap::from), - schema_metadata_updates: schema_metadata_updates - .as_ref() - .map(pb::transaction::UpdateMap::from), + config_updates: config_updates.as_ref().map(pb::UpdateMap::from), + table_metadata_updates: table_metadata_updates.as_ref().map(pb::UpdateMap::from), + schema_metadata_updates: schema_metadata_updates.as_ref().map(pb::UpdateMap::from), field_metadata_updates: field_metadata_updates .iter() - .map(|(field_id, update_map)| { - (*field_id, pb::transaction::UpdateMap::from(update_map)) - }) + .map(|(field_id, update_map)| (*field_id, pb::UpdateMap::from(update_map))) .collect(), // Leave old fields empty - we only write new-style fields upsert_values: Default::default(), @@ -764,13 +769,13 @@ impl From<&RewriteGroup> for pb::transaction::rewrite::RewriteGroup { } } -impl From<&UpdateMap> for pb::transaction::UpdateMap { +impl From<&UpdateMap> for pb::UpdateMap { fn from(update_map: &UpdateMap) -> Self { Self { update_entries: update_map .update_entries .iter() - .map(|entry| pb::transaction::UpdateMapEntry { + .map(|entry| pb::UpdateMapEntry { key: entry.key.clone(), value: entry.value.clone(), }) @@ -780,8 +785,8 @@ impl From<&UpdateMap> for pb::transaction::UpdateMap { } } -impl From<&pb::transaction::UpdateMap> for UpdateMap { - fn from(pb_update_map: &pb::transaction::UpdateMap) -> Self { +impl From<&pb::UpdateMap> for UpdateMap { + fn from(pb_update_map: &pb::UpdateMap) -> Self { Self { update_entries: pb_update_map .update_entries @@ -854,4 +859,40 @@ mod tests { other => panic!("expected DataOverlay, got {other:?}"), } } + + #[test] + fn test_composite_operation_rejected_on_load() { + // Action-based transactions (Transaction V2) are an experimental wire + // format that this version of Lance does not support. Loading one must fail closed + // (never be silently skipped or leniently parsed), so that a concurrent + // V2 commit in the conflict window aborts an in-flight commit. + 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: "append batch".to_string(), + actions: vec![pb::Action { + action: Some(pb::action::Action::AddFragment(pb::AddFragment { + id: Some(pb::Ref { + kind: Some(pb::r#ref::Kind::Local(0)), + }), + physical_rows: 1, + data_change: Some(true), + ..Default::default() + })), + }], + }], + }, + )), + ..Default::default() + }; + + let err = Transaction::try_from(message).unwrap_err(); + assert!( + matches!(err, Error::NotSupported { .. }), + "expected NotSupported, got: {err:?}" + ); + } }