diff --git a/crates/omnigraph/tests/lance_surface_guards.rs b/crates/omnigraph/tests/lance_surface_guards.rs index ba3971f19..7132be10f 100644 --- a/crates/omnigraph/tests/lance_surface_guards.rs +++ b/crates/omnigraph/tests/lance_surface_guards.rs @@ -7171,3 +7171,287 @@ async fn packed_struct_refuses_null_values_and_lone_child_projection_lance_11() 0 ); } + +// --- Group commit: composing inserts staged against older pins ------------- +// +// Group commit (docs/rfcs/2026-09-28-group-commit.md, Composition) publishes +// several same-table inserts as one detached commit. Each entry staged its +// keyed merge-insert against the pin it captured, and other inserts may have +// landed on the table since. Composition is sound because Lance assigns what +// a fresh fragment needs at commit, from the manifest the commit is built +// on: `fragments_with_ids` fills only a zero fragment id, `assign_row_ids` +// fills only a missing or partial row-id sequence, and the row-version +// metadata is derived while the manifest is built. This fence pins that the +// raw staged fragments are unassigned, that the components' key filters +// share one Bloom configuration, so their union is a bitwise OR, and that +// one detached commit of the union from a later pin yields exactly the union +// of rows with fresh, unique fragment and row ids and one version stamp. A +// Lance release that assigned ids at staging would make composition reuse +// them: a design review for that RFC, not a test to weaken. + +fn pk_rows(dataset: &Dataset, ids: &[&str]) -> RecordBatch { + let schema = Arc::new(Schema::from(dataset.schema())); + let values: Vec = (0..ids.len() as i32).collect(); + let notes: Vec> = ids.iter().map(|_| Some("group")).collect(); + RecordBatch::try_new( + schema, + vec![ + Arc::new(StringArray::from(ids.to_vec())), + Arc::new(Int32Array::from(values)), + Arc::new(StringArray::from(notes)), + ], + ) + .unwrap() +} + +/// A keyed upsert of `ids`, staged uncommitted against `base`: the shape the +/// engine stages for an `insert` into a `@key` type. +async fn stage_keyed_upsert(base: &Dataset, ids: &[&str]) -> UncommittedMergeInsert { + stage_pk_merge( + Arc::new(base.clone()), + pk_rows(base, ids), + "id", + WhenMatched::UpdateAll, + WhenNotMatched::InsertAll, + None, + ) + .await +} + +async fn commit_detached_from(base: &Dataset, transaction: Transaction) -> Dataset { + CommitBuilder::new(Arc::new(base.clone())) + .with_detached(true) + .with_skip_auto_cleanup(true) + .execute(transaction) + .await + .unwrap() +} + +/// `(key, stable row id, created-at version, last-updated version)` per row. +async fn keyed_row_identities(ds: &Dataset) -> Vec<(String, u64, u64, u64)> { + let mut scanner = ds.scan(); + scanner.with_row_id(); + scanner + .project(&["id", ROW_CREATED_AT_VERSION, ROW_LAST_UPDATED_AT_VERSION]) + .unwrap(); + let batches: Vec = scanner + .try_into_stream() + .await + .unwrap() + .try_collect() + .await + .unwrap(); + let mut out = Vec::new(); + for batch in batches { + let keys = batch["id"].as_string::(); + let row_ids = batch[ROW_ID].as_primitive::(); + let created = + batch[ROW_CREATED_AT_VERSION].as_primitive::(); + let updated = + batch[ROW_LAST_UPDATED_AT_VERSION].as_primitive::(); + for row in 0..batch.num_rows() { + out.push(( + keys.value(row).to_string(), + row_ids.value(row), + created.value(row), + updated.value(row), + )); + } + } + out.sort(); + out +} + +#[tokio::test] +async fn group_commit_composes_inserts_staged_against_older_pins_onto_a_later_pin() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().join("t.lance"); + let uri = uri.to_str().unwrap(); + let p0 = fresh_pk_dataset(uri).await; + + // Two entries capture P0 and stage keyed upserts of new keys. + let a = stage_keyed_upsert(&p0, &["b1", "b2"]).await; + let b = stage_keyed_upsert(&p0, &["c1"]).await; + // Another insert lands meanwhile: P1 is its detached commit from P0. + let landed = stage_keyed_upsert(&p0, &["d1"]).await; + let p1 = commit_detached_from(&p0, landed.transaction).await; + // A third entry captures P1. + let c = stage_keyed_upsert(&p1, &["e1"]).await; + + let mut union = Vec::new(); + let mut filters = Vec::new(); + for staged in [&a, &b, &c] { + let Operation::Update { + removed_fragment_ids, + updated_fragments, + new_fragments, + .. + } = &staged.transaction.operation + else { + panic!("a keyed merge-insert stages Operation::Update"); + }; + assert!( + removed_fragment_ids.is_empty() && updated_fragments.is_empty(), + "no key matched, so the staged transaction only appends" + ); + for fragment in new_fragments { + assert_eq!(fragment.id, 0, "a raw staged fragment carries no id"); + // Merge-insert records the captured ids of the rows it rewrote, + // rechunked per fragment and allowed to be incomplete; the commit + // fills the rest from its own base. An insert rewrote nothing, so + // the sequence is empty and every row id comes from the commit. + let carried = match &fragment.row_id_meta { + None => 0, + Some(lance_table::format::RowIdMeta::Inline(data)) => { + lance_table::rowids::read_row_ids(data).unwrap().len() + } + Some(other) => panic!("unexpected row-id metadata {other:?}"), + }; + assert_eq!( + carried, 0, + "a raw staged insert fragment carries no row ids of its own" + ); + } + union.extend(new_fragments.iter().cloned()); + filters.push( + staged_inserted_rows_filter(staged) + .expect("a primary-key merge-insert carries a key filter") + .clone(), + ); + } + let composed_fragments = union.len(); + + // The union of Bloom filters under one configuration is their bitwise OR. + let FilterType::Bloom { + bitmap: mut union_bits, + num_bits, + number_of_items, + probability, + } = filters[0].filter.clone() + else { + panic!("pinned Lance emits a Bloom key filter") + }; + for filter in &filters[1..] { + assert_eq!(filter.field_ids, filters[0].field_ids, "one key column"); + let FilterType::Bloom { + bitmap, + num_bits: other_bits, + number_of_items: other_items, + probability: other_probability, + } = &filter.filter + else { + panic!("pinned Lance emits a Bloom key filter") + }; + assert_eq!( + (*other_bits, *other_items), + (num_bits, number_of_items), + "one Bloom configuration" + ); + assert!((other_probability - probability).abs() < f64::EPSILON); + for (byte, other) in union_bits.iter_mut().zip(bitmap) { + *byte |= other; + } + } + let union_filter = KeyExistenceFilter { + field_ids: filters[0].field_ids.clone(), + filter: FilterType::Bloom { + bitmap: union_bits, + num_bits, + number_of_items, + probability, + }, + }; + for filter in &filters { + assert!(union_filter.intersects(filter).unwrap().0); + } + + // Compose onto P1: the first component's operation, carrying the union. + let mut operation = a.transaction.operation.clone(); + let Operation::Update { + new_fragments, + inserted_rows_filter, + .. + } = &mut operation + else { + unreachable!() + }; + *new_fragments = union; + *inserted_rows_filter = Some(union_filter); + let composed = + commit_detached_from(&p1, Transaction::new(p1.version().version, operation, None)).await; + assert!(lance_table::format::is_detached_version( + composed.version().version + )); + + // Exactly the union of rows, and P1's rows carried unchanged. + let rows = keyed_row_identities(&composed).await; + let keys: Vec<&str> = rows.iter().map(|row| row.0.as_str()).collect(); + assert_eq!(keys, ["alice", "b1", "b2", "bob", "c1", "d1", "e1"]); + let p1_rows = keyed_row_identities(&p1).await; + for row in &p1_rows { + assert!(rows.contains(row), "P1's row changed identity: {row:?}"); + } + + // Fresh, unique stable row ids from P1's counter. + let row_ids: HashSet = rows.iter().map(|row| row.1).collect(); + assert_eq!(row_ids.len(), rows.len(), "stable row ids are unique"); + let p1_max_row_id = p1_rows.iter().map(|row| row.1).max().unwrap(); + for row in rows.iter().filter(|row| !p1_rows.contains(row)) { + assert!(row.1 > p1_max_row_id, "reused a row id: {row:?}"); + } + + // Fresh, unique fragment ids above P1's. + let p1_fragment_ids: HashSet = p1 + .get_fragments() + .iter() + .map(|fragment| fragment.id() as u64) + .collect(); + let fragment_ids: Vec = composed + .get_fragments() + .iter() + .map(|fragment| fragment.id() as u64) + .collect(); + assert_eq!( + fragment_ids.iter().collect::>().len(), + fragment_ids.len(), + "fragment ids are unique" + ); + assert_eq!( + fragment_ids.len(), + p1_fragment_ids.len() + composed_fragments, + "every composed fragment is added" + ); + let p1_max_fragment_id = *p1_fragment_ids.iter().max().unwrap(); + for id in fragment_ids + .iter() + .filter(|id| !p1_fragment_ids.contains(id)) + { + assert!(*id > p1_max_fragment_id, "reused fragment id {id}"); + } + + // One commit, one stamp, distinct from the version that landed first. + let stamp = |key: &str| { + let row = rows.iter().find(|row| row.0 == key).unwrap(); + (row.2, row.3) + }; + for key in ["b2", "c1", "e1"] { + assert_eq!( + stamp(key), + stamp("b1"), + "{key} is stamped by the composed commit" + ); + } + assert_ne!( + stamp("b1"), + stamp("d1"), + "the landed insert keeps its own stamp" + ); + + // The composed version reopens by id. + let reopened = DatasetBuilder::from_uri(uri) + .with_version(composed.version().version) + .load() + .await + .unwrap(); + assert_eq!(keyed_row_identities(&reopened).await, rows); +} diff --git a/docs/rfcs/0067-detached-table-commits.md b/docs/rfcs/0067-detached-table-commits.md index 8db0d8b3e..2f1b5aa96 100644 --- a/docs/rfcs/0067-detached-table-commits.md +++ b/docs/rfcs/0067-detached-table-commits.md @@ -1277,7 +1277,8 @@ claims). per-chunk promotion cost near the 1,024-chunk bound. - Group commit at the branch gate is a separate proposal; here it is N detached commits, one publication, N promotions. Decided in that proposal, - when it opens. + when it opens. Proposed in [RFC: Group commit](2026-09-28-group-commit.md), + which publishes a batch as one graph commit rather than as N commit rows. - How a blocked table is unblocked. A foreign linear commit occupying a pin's target never resolves on its own, and every later pin on that table waits behind it. The candidates are a `repair` arm that restores the diff --git a/docs/rfcs/2026-09-28-group-commit.md b/docs/rfcs/2026-09-28-group-commit.md new file mode 100644 index 000000000..6ad7c0adf --- /dev/null +++ b/docs/rfcs/2026-09-28-group-commit.md @@ -0,0 +1,819 @@ +--- +rfc: "2026-09-28-group-commit" +title: "Group commit" +track: maintainer +status: draft +implementation: not-started +authors: + - ragnorc +created: 2026-09-28 +updated: 2026-10-03 +discussion: "https://github.com/ModernRelay/omnigraph/pull/785" +supersedes: [] +superseded_by: [] +blocked_on: + - "Engine half of the composition guard: a composed commit carries the markers its components carry, and the change feed's one-transaction proof and merge's insert-absence certificate accept it. The Lance half is pinned by `lance_surface_guards.rs::group_commit_composes_inserts_staged_against_older_pins_onto_a_later_pin`." + - "Correctness: the admission matrix under Evidence and tests, failpoints around the composed detached commits and the compare-and-swap including the in-doubt read-back, and the DST concurrent universe with the publisher on (every acknowledged write visible exactly once, one actor per commit)." + - "Instrument: `concurrent-writes` with eight writers on one branch, local and on RustFS at +30 ms per round trip, above one writer's rate on the same build (31.7 and 0.87 commits/s on `b14c22c5` with PR #783), with one writer within noise of today and per-writer commit counts within a factor of two of each other (#784)." +--- + +# RFC: Group commit + +**Depends on:** [RFC 0067](0067-detached-table-commits.md), whose +Throughput section sketches this as step 3 of its path; +[RFC: Detached-only tables](2026-09-21-detached-only-tables.md), which leaves +no work after a publication and defines the staging witness the collector +judges; [RFC: Shared schema gate](2026-09-18-shared-schema-gate.md), whose +decision log holds the measurement that motivates this. +**Surveyed:** OmniGraph `main` at `b14c22c5` and PR #783 at `8863d607`: the +publication path in `omnigraph-catalog` (`publisher.rs`, `commit.rs`, +`state.rs`, `commit_graph.rs`, `retention.rs`), the engine's write path +(`exec/staging.rs`, `exec/mutation.rs`, `validate.rs`, `loader/mod.rs`, +`table_store.rs`, `db/omnigraph.rs`, `db/write_queue.rs`), the change feed +(`changes/`), merge discovery (`exec/merge.rs`) and the collector +(`db/omnigraph/collector.rs`); the complete Lance +[distributed write](https://lance.org/guide/distributed_write/) guide and +[transaction specification](https://lance.org/format/table/transaction/), +read 2026-09-28. Rechecked 2026-09-30 against `main` at `c6f24757`, after: +- #783 merged; +- #813, which reuses a write's captured authority at revalidation; +- #797 and #794, the `.gqt` concurrent block and request measurement. + +## Summary + +Several writes on one branch publish with one `__manifest` compare-and-swap +as one graph commit. A publisher per branch in each process takes over the +part of a write that runs under the branch gate today. Writers prepare as +they do now and submit what they staged. While one publication is in flight, +new submissions queue. When it returns, the publisher admits the queued +writes that cannot invalidate each other or anything published since they +were captured, commits each table's admitted effects as one detached commit, +and publishes them through the unchanged single-commit publisher. Every +admitted write is acknowledged with that commit. + +Two decisions keep this small: + +- **A batch is one graph commit.** Its effect is the union of its writes. + The writes were in flight together, so no reader could observe an order + between them, and no intermediate state was ever a branch head. Every + surface that maps a commit to its `__manifest` version keeps working: time + travel, diffs, the change feed, merge bases, retention roots and the + lost-acknowledgement read-back. RFC 0067's sketch, N commit rows in one + publication, breaks every one of them (Alternatives). +- **Same-table inserts compose by key.** Fragments that independent writers + wrote against one table commit as one transaction, the shape Lance + documents for distributed writes. This covers strict inserts, and upserts + that inserted every row because none of their keys existed yet. That is + the case the measurement needs: every writer in the `concurrent-writes` + benchmark inserts new keys into one `@key` table, which the engine stages + as upserts, so admission by disjoint tables alone would batch nothing. + +A batch holds the writes of one actor, because a commit row names one actor. +A write that names an expected graph head always publishes alone. No storage +format, wire shape or API field changes. + +## Motivation + +### The same-branch ceiling + +The measurement is recorded in the shared schema gate RFC's decision log +(entries dated 2026-09-28, added by PR #783). The instrument was `concurrent-writes`, closed-loop, with +insert-only mutations and disjoint keys, on the local filesystem and on +RustFS behind toxiproxy adding about 30 ms per round trip. + +| | Local | +30 ms | +|---|---|---| +| One writer on one branch | 31.7 commits/s | 0.87 commits/s | +| Eight writers on one branch, PR #783 | 20–25 | 0.40 | +| Branch-gate hold for one commit | 7.3 ms | 670 ms | +| of which revalidation / detached commit / publication | 0.6 / 1.1 / 5.6 ms | 298 / 85 / 284 ms | +| Failed revalidations per publication, eight writers | 6.6 | 2.9 | + +Revalidation fails whenever the branch head has moved, so every publication +invalidates every other prepared write on the branch. Eight writers on one +branch therefore cannot exceed one writer. At +30 ms each failed attempt +holds the gate for about 300 ms before giving up, about 19 s of a 30 s +window, which is more than the successful writes used. Failing faster from +the handle's in-memory view was prototyped and measured. It raised the ++30 ms rate from 0.40 to 0.53 commits/s, gave nothing locally, and left the +ceiling where it was. The local one-branch cell drifts by about 20% between +sessions. Interleaved reruns put `main`, PR #783 and a diagnostic build all +at 20–25 commits/s, so only same-session comparisons count. Issue #784 is the same mechanism seen from the +writers' side: at +30 ms one writer makes all twelve commits of each run. + +### What a batch buys + +The work under the branch gate is fixed per publication: + +- revalidation's round trips; +- one detached commit per touched table; +- the catalog's copy-on-write rewrite at the next `__manifest` version, with + its compare-and-swap (`commit::overwrite` since #781). + +With B writes per publication that cost is paid once, and the detached +commits of a batch's tables run in parallel. No prepared write is +invalidated by a write it does not conflict with. + +The model: at +30 ms a batch costs about the 670 ms one write costs today, +so eight writes per batch would publish several times the current rate. +This is a model, not a claim. The instrument in `blocked_on` decides, and +RFC 0039's rule that claims need open-loop driving applies to anything +stated later. + +The history-dependent costs RFC 0067 and RFC 0068 measure grow with +publications, not with writes. A publication adds one registration row per +touched table and two lineage rows. Under load, B same-table inserts +produce one registration row instead of B, so `__manifest` grows B times +more slowly. + +### Why disjoint tables are not enough + +RFC 0067 admits an entry when the tables it read and wrote are disjoint from +the tables other entries wrote. All eight benchmark writers insert `Chunk` +rows, so that rule admits one entry per batch. Insert-heavy ingestion into +a few types is the common production shape too. + +`Chunk` declares `slug: String @key`. An `insert` into a keyed node type is +staged as an upsert (`PendingMode::Upsert`, `exec/mutation.rs`), so that a +later update in the same query coalesces with it. Composition therefore has +to cover upserts, not only strict inserts. An upsert whose keys all turned +out absent has made the same read a strict insert makes, key absence at its +base, and its staged transaction only appends. + +Lance itself tracks conflicts between concurrent inserts at key grain: a +merge-insert `Update` carries `inserted_rows`, a key-existence filter "used +for conflict detection" (transaction specification, Update). That filter is +a Bloom filter, so the exact ids cannot be recovered from it. The engine +knows an insert's ids while staging, but it does not keep them. A doc +comment on `StagedWrite` (`table_store.rs`) still describes such a field; +the field itself is gone. Keeping the ids is a prerequisite (Rollout). + +## User and operational behavior + +- **Commit grain.** Concurrent writes by one actor on one branch may publish + as one graph commit. + - Each write's receipt names that commit. + - `commit list` shows one commit for all of them. + - That commit's change-feed block, diff and time-travel snapshot contain + every one of its writes. + - Without contention every write is its own commit, as today. + - `graph_manifest_version` in `CommitOutput` stays the commit's own + version. +- **Conditional writes.** A write that names an expected graph head + (`Omnigraph-If-Graph-Commit`, or `expected_head` in the SDK) is always its + own commit, and that commit's parent is the named head. +- **Errors.** A write the publisher refuses at admission gets the typed + read-set conflict (`ReadSetChanged`), as it does today when revalidation + fails. The existing retry contract is unchanged: + - an insert-only mutation re-prepares, up to 32 times; + - an Append or Merge load re-prepares; + - an update, a delete and an Overwrite load return the conflict to the + caller (`exec/mutation.rs`, `loader/mod.rs`). + + A publication whose outcome is in doubt gives every write in the batch the + in-doubt error, naming the batch's commit id. +- **Completion evidence** ([Owned server operations](2026-09-30-owned-server-operations.md)). + Each entry's outcome carries the evidence the same failure carries on + today's path, so the server's classification is unchanged: + - An entry refused at admission has made no detached commit, since + commits follow admission. It gets the typed read-set conflict, which + the server treats as settled, as it does today's refusal by + revalidation before `commit_all`'s first commit. + - A batch that loses the compare-and-swap leaves its detached commits + unreachable and gives its entries the typed conflict, not in doubt. + - An in-doubt batch makes every entry in doubt. The server closes + admission once, as it does for one in-doubt write today. + - An admitted entry whose batch outcome the engine did not observe (its + leader's future panicked or was dropped) gets `Uncertain`, so the + server closes admission as it does for a panicking write today. An + entry not yet admitted has made no effect and is refused. +- **Fewer conflicts.** Today any publication on the branch refuses every + write prepared before it. Under this RFC, a write is refused only when a + publication since its capture touched its footprint. An update that today + fails because an unrelated table changed now succeeds. +- **Fairness.** Entries are admitted in arrival order. A refused entry that + re-prepares keeps its original place, so a writer that just published + cannot overtake the writers it made stale (#784). This holds for the kinds + that re-prepare automatically, which include the benchmark's writes. +- **Cancellation.** + - An HTTP write is an owned operation: its caller's disconnect does not + drop the engine future, so it never removes an entry. Dropping applies + to an embedded caller's future, and at shutdown the server waits for + registered operations, so for the batches their entries are in, under + its one deadline. + - A caller that goes away before its entry is admitted is dropped from the + queue. + - Once admitted, the publication runs on its leader's task. An HTTP + leader is an owned operation and is never dropped. If an embedded + caller drops a leader's future, the batch's entries are refused and + re-prepare when the leader had made no effect yet, and get `Uncertain` + otherwise; the next waiter is promoted either way. + - A write whose caller disconnected after admission may land. That caller + is in the same position as a lost acknowledgement today. +- **Unchanged.** + - Merge, schema apply, index builds, Optimize, cleanup and branch controls + publish as today, and the branch gate serializes them with batches. + - Reads. + - Writers in different processes still race only at the compare-and-swap. +- **Documentation.** `docs/dev/writes.md` gets the write path. The user + guide's commit pages say that a commit may carry several writes. The + release notes record the commit-grain change. + +## Design + +### Entries + +A mutation or load prepares exactly as today, and nothing is committed +before submission: + +1. capture the authority; +2. execute; +3. validate at the pinned base; +4. stage: fragments are written and each table's Lance transaction is built. + +Instead of taking the gates in `commit_all`, it submits an **entry** and +waits for its outcome. An entry holds: + +- the captured authority. That is the branch incarnation, the schema + identity, the caller's expected head if any, and the two heads the write + transaction already distinguishes (`WriteTxn`, `db/omnigraph.rs`): + - the **materialized head** H0, the branch's `graph_head` row, which the + publisher's `ExactGraphHead` compares; + - the **effective head**, which a fresh named branch inherits from its + fork point before its first publication, and which the caller's + expected head is compared against. +- for each written table, the staged write, the pin it was staged from, and + its class (below); +- the **footprint**: every table the write read or wrote, which is the union + of: + - its written tables; + - its **execution reads**: every table execution opened through + `ensure_path`. That includes a predicate scan that matched nothing and + therefore staged nothing. These are the tables + `MutationStaging.expected_versions` records today (`exec/staging.rs`). + - its **validation reads**, which today are derivable from the changeset + and the catalog but recorded nowhere: + - edge-endpoint existence reads the two endpoint node tables; + - a node delete's cascade and referential checks read every edge table + incident to its type; + - an overwrite load's referential checks read edge tables outside the + load. + + Recording them is the first rollout step. +- the actor and an arrival ticket. + +Each written table is in one of two classes: + +- **Append-only**, when all of these hold: + - the write on that table is a strict insert, or an upsert; + - its staged transaction on that table adds fragments and removes or + updates none, so an upsert matched no existing key; + - it carries the exact ids of its inserted rows, within the per-entry cap; + - the table has no non-key `@unique` group, whether it is a node table or + an edge table (`validate.rs` checks both against committed rows); + - an edge table has no bounded `@card`. + + Validation of such a write reads only key absence on its own table (the + strict-insert probe, or the upsert's join that matched nothing) and, for + an edge, the existence of its endpoints. +- **Exclusive** covers everything else: + - updates, deletes, upserts that matched a key, and overwrites; + - cascades; + - non-key uniqueness; + - bounded cardinality, whose check reads a fresh live version + (`validate.rs`); + - oversized id sets. + + An entry with any exclusive table is an exclusive entry. + +### The publisher + +Each process keeps one publisher per graph root and branch: a queue and a +leader flag in the process-global write-queue manager that already owns the +gates (`db/write_queue.rs`, keyed by root). It spawns no task. The writers +form a write group: an entry that arrives while no batch is publishing +leads, publishes the batch that holds its own entry on its own task, then +promotes the next waiting entry to lead the next batch, or clears the flag +when none waits. A writer never publishes others' writes at the expense of +its own next one, and every publication runs inside the owned operation of +the writer leading it. + +A batch runs: + +1. **Gates.** Take the schema gate's shared side and the branch gate, which + is what `commit_all` takes today. The per-table gates it also takes add + no exclusion inside the branch gate (the shared gate's decision log), so + the publisher does not take them. +2. **Revalidate once** with the checks `revalidate_write_txn` runs per + write today: the branch authority and the schema contract, which since + #834 is a row of `__manifest` read with it. Since #813, a write to a branch other + than the handle's bound branch reuses its captured authority when a + probe finds the `__manifest` version unchanged; the publisher can reuse + its own last published authority the same way. This yields the current materialized head Hc, the + current effective head, and the current pins. On a fresh named branch + before its first publication the materialized head is absent while the + effective head is the fork point. Below, Hc always means the + materialized head. +3. **Log check.** If Hc is not the head the publisher's last batch + produced, another publication landed; clear the effect log. +4. **Admit** queued entries in ticket order, up to caps on entries and on + composed fragments. The fragment cap keeps a composed transaction inline + in its manifest. + - An entry of another actor stays queued and leads the next batch. + - An entry with an expected head is admitted only under all of these, + and it then closes the batch: + - it is the batch's first entry; + - the current effective head equals its expected head; + - Hc equals the materialized head it captured. + + This is exactly the pair of checks `revalidate_write_txn` makes today. +5. **Commit** each table's effect detached from the table's pin at Hc, + stamped with the staging witness (branch incarnation, Hc). Hc is the + exact materialized head the publication's precondition names. Tables + commit in parallel. +6. **Publish** with `publish_with_precondition`, passing: + - one `DatasetUpdate` per table; + - the expected versions at Hc for every table in the batch's footprints; + - one `LineageIntent` carrying the batch's actor; + - `ExactGraphHead(Hc)`. + + The publisher API is unchanged. +7. **Complete** every entry with the outcome. A submitting handle applies + the published snapshot to its own coordinator before its receipt + returns, so a read on that handle sees its write. No handle's + coordinator is held across the compare-and-swap. That is expected to + remove the same-handle read wait the shared gate RFC lists as + unresolved; the implementation has to show it. + +Two failure outcomes: + +- **Confirmed refusal** (`ReadSetChanged`): delete the batch's detached + manifests, which is the existing rule for a refused publication, refresh, + and judge the entries again from step 3. Entries captured before the + foreign commit are refused. +- **In doubt:** every entry gets the in-doubt outcome. The read-back is + unchanged, because it looks for exactly one commit. + +### Admission + +Take an entry k captured at H0 and a batch forming at Hc. Let *writers* +be every commit in (H0, Hc] together with every entry admitted earlier in +this batch. k is admitted when all three hold: + +1. **Authority.** Its branch incarnation and schema identity equal the + batch's. +2. **Horizon.** Every commit in (H0, Hc] is in the effect log, so this + publisher published it, and there are at most K of them. +3. **Conflicts.** Against k's footprint: + - If k is an exclusive entry, no writer wrote any table in its + footprint. + - If k is append-only: + - no writer wrote any table in its footprint exclusively; + - for every table that k and a writer both appended to, their id sets + are disjoint. + +The second case follows from what the entry read. An append-only entry read +three things: +- the absence of its own ids; +- the existence of its endpoints; +- any tables its execution opened. For a write that only inserts, those are + its written tables. + +Appends by other writers with other ids preserve all three; an exclusive +write to any footprint table may not. + +An exclusive entry read rows or absences of arbitrary shape, such as a +predicate scan (including one that matched nothing), a cascade, or +referential emptiness. So any write to its footprint invalidates it. Leaving +out a table that execution scanned but did not write would admit write skew. +For example, T1 updates A and scans B for a value only T2 would set, while +T2 does the reverse. Admitted together, both find nothing in their second +scan, which is a result no serial order gives. Execution reads are +therefore part of the footprint. + +The batch is conflict-serializable in admission order: + +- an entry's reads are unaffected by every earlier writer; +- a later entry's writes come after its reads in that order; +- a write that lands after a read is not a conflict in that order. + +Constraints validated at an entry's base therefore hold at Hc. A constraint +the rule cannot carry forward makes the write exclusive, and then its whole +footprint must be unchanged. + +### Composition + +Each table in a batch has either one exclusive entry or any number of +append-only entries; the admission rule excludes a mix. + +- **Exclusive:** the entry's transaction commits as staged. Its base is + still the table's pin, because nothing wrote the table since H0. +- **Append-only:** one transaction whose new fragments are the union of the + entries' fragments, read from the table's pin at Hc. That pin may be later + than an entry's base, but only appends with disjoint ids came in between. + Lance assigns what the union needs at commit, and never earlier: + - fragment ids: "not assigned at transaction creation time; they are + assigned during manifest construction" (Append); + - stable row ids: "only the commit knows which values are free" + (distributed write); + - row-version metadata: "leave them as `None`. Lance derives both while + building the manifest" (distributed write). + + The union must be taken from each staged write's raw transaction. There, + every new fragment has fragment id 0, and its row-id sequence is empty for + an insert: merge-insert records only the ids of the rows it rewrote, + chunked per fragment and allowed to be incomplete. Lance then assigns both + at commit from Hc; `assign_row_ids` fills an incomplete sequence from the + commit's counter. + + `StagedWrite::new_fragments()` is not usable here. It is the + read-your-writes copy, with fragment and row ids already provisionally + assigned from the entry's own base, and Lance keeps nonzero fragment ids + and complete row-id sequences as given. + + The composed transaction keeps the insertion-only shape the components + have, a merge-insert `Update` in `RewriteRows` mode with new fragments + only, and its markers: + - the key-existence filter is the bitwise OR of the components' Bloom + filters, which share one configuration and one key column, so the OR is + exactly their union and no key is re-encoded; + - the field lists stay those of the components, which are identical for + one table under one schema; + - a marker is kept only when every component carries it. + `omnigraph.no_by_source_delete` is stamped at the keyed merge-insert + chokepoint (`table_store.rs`). + `omnigraph.insert_absence` is certified for strict inserts and for + pure-insert upserts. When every component carries it, it is certified + again for the composed base: the admission proof establishes absence at + Hc with no further read, and a guard pins that merge's certificate + reader accepts it. The change feed's fast path accepts either marker + (`transaction_is_row_set_preserving`). + + Copying the properties onto an arbitrary `Append` would not do, because + merge's certificate reader checks the shape. + +The commit goes through the sealed table adapter, as a new entry the +durable-call guard registers. The change feed's fast path still applies. It +requires the pinned transaction's read version to be the previous pin +(`changes/candidate_scan.rs`), and a composed commit's read version is the +table's pin at Hc. + +### The effect log + +The effect log is in memory, per publisher. For each of the publisher's +last K batches it records: + +- the commit id; +- the head that batch produced; +- for each table the batch wrote, either "exclusive" or the ids it appended. + +An append whose ids exceed the cap is logged as exclusive. The log is +bounded by K times the cap. + +It caches facts about immutable, published commits, which invariant 12 +permits. It is never commit authority; `ExactGraphHead(Hc)` is. If the +process restarts, the log is lost, and entries captured before the restart +fail the horizon rule and re-prepare. + +### Why the collector's staging rule does not change + +The collector judges an unpublished staging dead once the branch head has +moved past the head it recorded (`judge_staging` in the collector). Under +this RFC, every detached commit is written by the publisher after +admission and stamped with Hc. A batch that loses the compare-and-swap +leaves stagings the existing rule proves dead once the head moves past Hc. + +RFC 0067's sketch had writers commit detached effects before submitting. +Those effects would carry the entry's own H0, and an admitted entry whose +H0 is older than Hc would be judged dead while it was about to publish. +Committing after admission keeps the rule. It also keeps the shared gate's +finding that detached commits made ahead of admission are mostly garbage +under contention. + +An entry's fragments are unreferenced from preparation until its batch's +detached commit. The collector deletes a file that no manifest references +only when the file is older than `UNVERIFIED_THRESHOLD_DAYS`, which is seven +days. Every write already relies on that window between preparation and +its detached commit. The queue and the horizon bound an entry's wait in +publications, not in wall-clock time: a suspended process can outlive the +window today, and it still can. This RFC inherits that assumption; it +neither adds to it nor removes it. A refused entry's fragments are the +orphans every failed revalidation leaves today, and the same age rule +reclaims them. + +## Invariants + +- **2 (one publication door).** Unchanged: one publication per batch. +- **3 (one coherent view).** + - Each entry holds one captured view, and the batch revalidates the + complete authority once before any of its effects. + - An entry admitted at Hc later than its H0 publishes against state that + differs from its capture only by commits the admission rule proves did + not touch its footprint. That is backward validation, the + optimistic-concurrency shape databases use. + - The invariant's wording, "revalidates that complete authority before + effects" and "never combines fresh and stale facts", needs a precise + amendment in the implementation change: an entry's footprint is proven + unchanged, not its whole view. +- **4 (publish once).** Every write publishes once, in exactly one commit. +- **5 (recovery).** Unchanged. Detached effects are committed only after + admission, and are garbage under the unchanged rule otherwise. +- **7 (derived acceleration).** A composed append leaves an unindexed tail + like any append. +- **8 (loud integrity).** Every constraint an entry validated holds at + publication by the admission rule. Those it cannot carry forward make the + write exclusive. +- **10 (policy).** Enforced per write at its `_as` entry point before + submission. Unchanged. +- **11 (bounded).** + - The queue is bounded by entries and bytes, and the effect log by K and + the cap. For HTTP writes the server's admission already bounds the + operations that can wait in it; the queue's own bound covers embedded + callers. + - A refusal goes through the existing bounded re-prepare loop and is + counted as a re-prepare. +- **12 (one source of truth).** The effect log caches immutable facts. The + queue holds requests, not state derivable from the manifest. +- **Deny-list.** + - The publisher is process-local and is not presented as fencing; the + compare-and-swap remains the fence between processes. + - It adds no job queue for manifest-derived work. + +## Compatibility and reversibility + +- **Storage:** none. Commit, head and registration rows are unchanged, and + an older binary reads a batch commit as an ordinary commit. +- **Wire:** none. +- **Behavior:** commit grain under concurrent writes by one actor (User and + operational behavior). +- **Reversal:** two settings together reproduce today's behavior, and they + are the rollback: + - a batch cap of one entry; + - a horizon of zero, which admits only entries captured at Hc. + + A cap of one alone does not: an entry captured before an unrelated + publication would still be admitted, where today it is refused. + +## Alternatives + +- **N commit rows in one publication**, as RFC 0067 sketched. Both code + surveys for this draft found it breaks every mapping from a commit to its + `__manifest` version: + - Snapshot reads and diffs of a non-last member fail their head check + (`graph_coordinator.rs`, `feed.rs`), and so does the change feed. + - The feed's one-transaction proof needs one detached commit per interval. + - Merge bases and the collector's merge-base roots go through + `pinned_graph_commit` (`retention.rs`), which aborts the whole cleanup + plan on a non-head member. + - Parent resolution orders by `(graph_manifest_version, created_at, id)`, + so it can silently fork lineage. + - `commit::overwrite` keeps duplicate head rows within one publication. + - The lost-acknowledgement read-back recognizes only the head commit. + - `ExactGraphHead` holds one expectation. + - Schema apply records its parent before publication. + + Per-commit pins would need registration rows keyed by commit, a format + change. Its interior commits would name states no reader ever observed, + and with same-table composition they would have no table versions at all. +- **Admission by disjoint tables only.** No gain on the measured workload. +- **Chaining same-table writes**, staging each on the previous write's + detached version the way merge chunks chain. This serializes preparation + per table, and #784 puts one preparation at about 2.4 s at +30 ms under + contention. Refusing one link refuses its successors. Composition needs no + re-preparation. +- **A timed batching window.** Adds latency at low load. Queuing while a + publication is in flight batches without a timer. +- **A dedicated publisher task.** A task the engine spawns is a child + producer no single owned operation owns ([Server runtime and online + deployment](2026-09-29-server-runtime-and-online-deployment.md), Operation + ownership). Its requests also fall outside every session's `--measure` row + and every `.gqt` `order:`. The spike measured nothing it would add over + the write group (decision log, 2026-10-03). +- **Failing fast before revalidation.** Measured, and bounded by the + single-writer ceiling (shared gate decision log). +- **Batches across actors.** Needs per-write attribution in the commit row, a + format change. Its natural home is beside #513's per-write idempotency key. +- **RFC 0068's commit record.** Lowers the per-publication cost. It is + orthogonal and composes with this: one record per batch. + +## Evidence and tests + +- **`lance_surface_guards.rs`: the composition guard, Lance half.** + - Test: `group_commit_composes_inserts_staged_against_older_pins_onto_a_later_pin`, + passing on Lance 11.0.0 since 2026-09-30. + - Setup: two keyed upserts of new keys staged against P0 and one staged + against P1, where P1 is a third insert that landed detached from P0. + The three are committed as one detached transaction from P1. + - Pinned: + - the raw staged fragments carry fragment id 0 and an empty row-id + sequence; + - the components' key filters share one Bloom configuration, and + their OR covers each; + - the result is exactly the union of rows, and P1's rows keep their + row ids and stamps; + - new row ids and fragment ids are unique and above P1's; + - the composed rows share one version stamp, distinct from the + insert that landed first; + - the composed version reopens by id. +- **The composition guard, engine half.** Owned by `changes.rs` and the + merge owners. They assert that a composed commit carries its components' + markers, and that the feed's fast path and the merge chain's + insert-absence certificate accept it. +- **The admission matrix.** Each cell asserts the commit count, the rows, + and which entry re-prepared or failed. + + The cells that differ only in outcome are `.gqt` cases with a + `--- concurrent` block (#797), the testing guide's default for behavior + visible in outcomes: + - two to four sessions on one handle; + - an `order:` that holds each writer after its preparation until every + writer has prepared. It parks at a request made in the writer's own + future and under no lock another session needs. Which request + qualifies is not settled: fragment writes may run on Lance's pool, + which no order can name. Finding one is part of rollout step 2's + evidence. + - one `ok` or `error:` line per session; + - the rows checked in the step after the block. + + Commit counts and which entry re-prepared are not expressible there, so + those assertions stay in `writes.rs`. The cells are: + - Appends with disjoint ids make one commit. This covers strict inserts, + and `insert` into a `@key` type, which is staged as an upsert. + - The same id: one commit, and the other entry re-prepares. A strict + insert then fails with its key conflict; an upsert becomes an update. + - An append against an exclusive write on the same table. + - Edges with a non-key `@unique` group are exclusive: two edge inserts + with the same unique value but different generated ids never share a + batch. + - Write skew through execution reads: T1 updates A and scans B for a value + that matches nothing, and T2 does the reverse. They never share a batch, + and the second fails with the conflict. + - An edge insert against a delete on its endpoint type. + - Bounded `@card` is exclusive. + - An expected head publishes alone, including the first conditional write + on a fresh named branch, whose materialized head is absent + (`writes.rs` already guards that case for single writes). + - Two actors make two commits. + - A publication by merge or by another handle between capture and batch. + - Horizon zero with a cap of one reproduces today's refusals. +- **`failpoints.rs`.** For each window, assert what did not move: + - after the composed detached commits and before publication; + - a refused compare-and-swap, whose detached manifests are deleted; + - in doubt, where every entry is in doubt and the read-back resolves the + batch; + - a waiter cancelled before and after admission. +- **DST: the concurrent universe with the publisher on.** + - The writers already share the process-global manager, so batching + happens. + - Every publication runs on a writer's task, so the seam scheduler's + existing writer actors cover it. + - Oracle: every acknowledged write is visible exactly once, and each + commit has one actor. + - The collector oracle runs with batches in flight. +- **`write_cost.rs`.** A batch of one costs what a write costs today. +- **`concurrent-writes`.** One and eight writers on one branch, local and at + +30 ms, with the per-writer commit distribution. + +## Rollout + +Each step leaves `main` shippable. + +1. On the mutation and load staging, record: + - the footprint, which is `expected_versions` plus the validation reads; + - each staged table's class; + - the exact ids of inserted rows on the staged write. + + Remove `StagedWrite`'s dangling doc comment for the id field that no + longer exists. No behavior change. +2. Add the publisher with batches of one. It replaces the gate section of + `commit_all`, behavior is identical, and the cost owner and DST stay + green. +3. Batch exclusive entries by footprint. +4. Compose append-only entries, behind the surface guard. +5. Run the instrument; update the documentation and add the release-note + fragment. + +## Unresolved questions + +- **K and the caps.** Chosen from the instrument by the implementer. +- **Loads in batches.** A small load composes under the same rules, and a + large one exceeds the id cap and is exclusive. Recommended: admit loads. + Decided in step 3. +- **Merge, index builds and Optimize as exclusive entries**, instead of + taking the branch gate between batches. Later, once batches are measured. +- **Admitting across a foreign commit** by reading that commit's per-table + change sets, instead of refusing every entry captured before it. The + change sets exist (`omnigraph.deleted_ids` and the fragments of each + detached transaction). Later. +- **Where per-write attribution belongs.** It is needed for batches across + actors and for #513's per-write idempotency key: in this format or in RFC + 0068's record. Decided when either is proposed. +- **The next lever after this one.** Once publication is amortized, a + writer's cycle at +30 ms is dominated by preparation's own round trips. + Since #834 one insert makes 13 requests: 6 writes, the key check, the + table open, and 4 `__manifest` reads. #813 removed the branch probe's + reopen for writes to a branch other than the handle's bound branch. A + capture served from the publisher's known head is the next step, filed as + #837. Not part of this RFC. + +## Decision log + +- 2026-09-28 — Drafted from the shared schema gate's step-2 measurement and + three code surveys of `main` plus PR #783 (lineage and time travel, the + change feed and per-commit changes, write footprints and the publisher). + - Chose one graph commit per batch over RFC 0067's N commit rows, because + the surveys found every commit-to-version mapping assumes it. + - Chose key-grain composition of same-table inserts, because the measured + workload writes one table. + - Chose detached commits after admission, because the collector's staging + rule would otherwise sweep admitted entries. +- 2026-09-28 — An independent review (Codex, `gpt-6-astra`) of the first + draft found the admission rule unsound in two cases and five further + gaps. All seven were checked against the code and corrected. + - Edges with a non-key `@unique` group were admissible as append-only, + which lets two edges with one unique value publish together. Non-key + uniqueness is now exclusive for edge tables too. + - The footprint counted only validation reads. A predicate scan that + matched nothing stages nothing, so two updates could each read the + table the other wrote and be admitted together (write skew). The + footprint now includes every table execution opened, which is what + `expected_versions` already records. + - Hc conflated the materialized head, which the publisher compares, with + the effective head, which a caller's expected head is compared against. + They differ on a fresh named branch before its first publication. The + two are now distinct. + - A batch cap of one did not restore today's behavior, because + cross-publication admission stayed on. The rollback is now a cap of one + with a horizon of zero. + - The fairness claim assumed every refused entry re-prepares. Only + insert-only mutations and Append or Merge loads do. The text now + states the existing retry contract. + - `insert` into a keyed type is staged as an upsert, so the motivating + benchmark would not have composed. The append-only class now includes + upserts whose staged transaction only appends. + - The exact inserted ids are not carried on `StagedWrite`; only a stale + doc comment remains. Carrying them is now a rollout prerequisite. + + The review confirmed the rest: one commit per batch leaves the + commit-to-version consumers intact, Lance 11 can compose the raw + unassigned fragments at a later pin, and the collector's staging rule + holds with the materialized head as the witness. From the same review, + the RFC now also states: + - composition must use the raw transaction's fragments; + - the composed commit keeps a marker only when every component carries + it; + - the seven-day orphan window is an inherited wall-clock assumption, not + one this RFC bounds. +- 2026-09-30 — Rechecked against `main` after #783, #813, #797 and #794 + merged. + - #783's dependency left `blocked_on`. + - The outcome cells of the admission matrix moved to `.gqt` concurrent + blocks. The commit-count assertions stay in Rust. + - The GQT concurrent block cannot name a request made on a spawned task. + That turned publisher-task attribution into an unresolved question with + a recommendation. + - #813's reuse of captured authority is noted where revalidation and the + next lever are described. + +- 2026-09-30 — The Lance half of the composition guard is checked in and + passes. + - Its first run corrected a fact the review had stated and this RFC + repeated: that the raw staged fragments carry no row-id metadata. They + do carry a row-id sequence, the captured ids of rows the merge-insert + rewrote, and for an insert it is empty. Lance's `assign_row_ids` fills + an incomplete sequence at commit, so composition holds; the guard now + decodes the sequence and pins that it is empty. + - The key-existence filter's union is the bitwise OR of the components' + Bloom filters, pinned by the same guard, in place of rebuilding the + filter from the ids. + +- 2026-10-01 — Rechecked against `main` after #823, #824 and #833. + - #824 made every HTTP write a server-owned operation and classifies a + failed write by completion evidence. Errors now state each entry's + evidence, Cancellation states that a disconnect never removes an HTTP + entry, and invariant 11 notes that server admission bounds the HTTP + side of the queue. + - #823's merge receipts and #833's changelog fragments change nothing + here; step 5's release note becomes a fragment. + - #834 moved the schema contract into `__manifest` and removed the + schema-apply sentinel, so revalidation no longer reads it and an insert + makes 13 requests instead of 20. The next-lever paragraph now cites + #837. + +- 2026-10-03 — The publisher is a write group, not a spawned task. + - The spike posted on the PR ran a leader-follower write group and reached + the instrument gate: one writer unchanged, eight writers on one branch + at 9.23 against 0.87 commits/s at +30 ms, per-writer counts within 1.3×. + Its first version kept the leader leading while entries waited, which + published other writers' writes instead of the leader's next one; 2 + writers ran below baseline and one run starved a writer. A writer now + leads only the batch that holds its own entry, then promotes the next. + - The alternative rejected earlier, cancellation abandoning followers, + narrows under the accepted server runtime RFC: an HTTP write is an owned + operation that a disconnect never drops. An embedded leader that is + dropped refuses its batch before any effect and reports `Uncertain` + after one, and promotes the next waiter. + - A spawned publisher would be a child producer no single owned operation + owns, so it moves to Alternatives, and the unresolved question of where + publication requests are attributed closes: they run on the leader's + own task. diff --git a/docs/rfcs/README.md b/docs/rfcs/README.md index a5c0ace52..4b02fd8fd 100644 --- a/docs/rfcs/README.md +++ b/docs/rfcs/README.md @@ -233,6 +233,7 @@ then dated RFCs by date. | [2026-09-21](2026-09-21-detached-only-tables.md) | Detached-only tables | maintainer | accepted | in-progress | | [2026-09-24](2026-09-24-shared-expression-model.md) | Shared expression model | maintainer | draft | in-progress | | [2026-09-26](2026-09-26-self-contained-server-testing.md) | Self-contained server testing with GQT and DST | maintainer | draft | not-started | +| [2026-09-28](2026-09-28-group-commit.md) | Group commit | maintainer | draft | not-started | | [2026-09-29](2026-09-29-server-runtime-and-online-deployment.md) | Server runtime and online deployment | maintainer | accepted | in-progress | | [2026-09-30](2026-09-30-typed-edge-alternation.md) | Typed edge alternation and bounded wildcard traversal | maintainer | accepted | complete | | [2026-09-30](2026-09-30-v012-http-admission.md) | v0.12 HTTP admission | maintainer | accepted | complete |