diff --git a/crates/tracedecay-rusqlite-runtime/src/ledger/commit.rs b/crates/tracedecay-rusqlite-runtime/src/ledger/commit.rs index 49f0fa5711..59173006ea 100644 --- a/crates/tracedecay-rusqlite-runtime/src/ledger/commit.rs +++ b/crates/tracedecay-rusqlite-runtime/src/ledger/commit.rs @@ -75,11 +75,12 @@ fn record_with_bookkeeping( }; let receipt_json = encode_json(&receipt, "original_receipt_json")?; checkpoint::persist(transaction, &submission, &checkpoint, &receipt)?; - // The persisted checkpoint is the authority for which records are now - // unreachable, so pruning reads it after persist. Every commit makes one - // bounded cleanup pass; the record inserted below is current and cannot be - // selected by that pass. - prune::prune_superseded(transaction, &metadata.shard_id)?; + // The persisted checkpoint is the validated authority for which of this + // incarnation's records are now unreachable, so the prune runs after + // persist and carries that checkpoint. Every commit makes one bounded pass + // so a backlog converges; the record inserted below sits at the persisted + // epoch and is never eligible. + prune::prune_superseded(transaction, &submission, &checkpoint)?; idempotency::insert(transaction, &submission, &receipt, &receipt_json)?; match bookkeeping { RuntimeBookkeeping::None => {} diff --git a/crates/tracedecay-rusqlite-runtime/src/ledger/prune.rs b/crates/tracedecay-rusqlite-runtime/src/ledger/prune.rs index 220edd0eeb..49c0dc22e0 100644 --- a/crates/tracedecay-rusqlite-runtime/src/ledger/prune.rs +++ b/crates/tracedecay-rusqlite-runtime/src/ledger/prune.rs @@ -8,7 +8,7 @@ //! //! `checkpoint::next` rejects any submission whose authority epoch is below the //! epoch persisted for its incarnation, so once the checkpoint for an -//! incarnation advances to epoch `E`, every idempotency record for that +//! incarnation stands at epoch `E`, every idempotency record for that //! incarnation below `E` is permanently unreachable: no submission carrying it //! can reach the ledger insert, and the writer rolls its savepoint back on the //! resulting `StaleAuthority` error. Deleting those rows therefore cannot admit @@ -21,38 +21,52 @@ //! authority is a provable end to the duplicate window. use rusqlite::params; -use tracedecay_store::StoreShardIdV1; use super::{ LedgerError, - sqlite::{LedgerTransaction, encode_json}, + checkpoint::NextCheckpoint, + sqlite::{LedgerTransaction, Submission, sqlite_u64}, }; -/// Maximum superseded records removed by one foreground commit. -const MAX_PRUNED_ROWS_PER_COMMIT: i64 = 256; +/// Maximum superseded records one foreground commit may remove. +/// +/// The cleanup shares the user mutation's savepoint and the process's sole +/// SQLite writer transaction, so it must never scale with the size of the +/// backlog it discovers. A legacy ledger can hold hundreds of thousands of +/// superseded rows; deleting them in one statement would monopolise admission, +/// and an interruption or `SQLITE_FULL` would roll back the checkpoint advance +/// with it, so every retry would re-attempt the same delete and the new epoch +/// would never commit. One bounded batch keeps the epoch advance committable no +/// matter how much retention work remains. +pub(super) const MAX_PRUNED_ROWS_PER_COMMIT: i64 = 256; -/// Deletes a bounded set of rows whose `(incarnation, authority_epoch)` can no -/// longer be carried by any admissible submission. +/// Deletes at most one bounded batch of rows whose `(incarnation, +/// authority_epoch)` can no longer be carried by any admissible submission. +/// +/// The candidate set is restricted to the single incarnation whose checkpoint +/// this commit just decoded, validated, and persisted, at that checkpoint's +/// validated epoch. No other incarnation's checkpoint row is consulted: reading +/// a neighbour's raw scalar `authority_epoch` would trust a value that nothing +/// has validated, so one inconsistent row - a scalar corrupted to `999` while +/// its watermark and receipt still encode `7` - would silently retire that +/// incarnation's live receipts and re-admit the duplicate writes they exist to +/// stop. A neighbour's superseded records are retired by that incarnation's own +/// commits, under its own validated checkpoint. /// -/// The join restricts candidates to incarnations that already have a -/// checkpoint, so an unknown authority position never authorises a delete. -/// Ordering by the table's composite primary key makes each bounded pass -/// deterministic and guarantees that repeated commits converge on the -/// remaining backlog. +/// The candidate scan matches the table's primary-key prefix and is therefore +/// already in key order, which makes each bounded pass deterministic and lets +/// repeated commits converge on the remaining backlog. const DELETE_SUPERSEDED: &str = r#" DELETE FROM td_runtime_writer_idempotency_v1 WHERE (shard_json, incarnation, authority_epoch, idempotency_key) IN ( SELECT candidate.shard_json, candidate.incarnation, candidate.authority_epoch, candidate.idempotency_key FROM td_runtime_writer_idempotency_v1 AS candidate - JOIN td_runtime_writer_checkpoint_v1 AS checkpoint - ON checkpoint.shard_json = candidate.shard_json - AND checkpoint.incarnation = candidate.incarnation WHERE candidate.shard_json = ?1 - AND candidate.authority_epoch < checkpoint.authority_epoch - ORDER BY candidate.shard_json, candidate.incarnation, - candidate.authority_epoch, candidate.idempotency_key - LIMIT ?2 + AND candidate.incarnation = ?2 + AND candidate.authority_epoch < ?3 + ORDER BY candidate.authority_epoch, candidate.idempotency_key + LIMIT ?4 ) "#; @@ -61,15 +75,26 @@ WHERE (shard_json, incarnation, authority_epoch, idempotency_key) IN ( /// /// Runs in the caller's transaction, like every other ledger operation, so the /// deletion shares the commit boundary of the mutation that advances cleanup. -/// Every new commit performs one bounded pass so legacy or partially drained -/// backlogs converge without monopolising the writer transaction. +/// Every commit makes one bounded pass, so a backlog left by an earlier pass - +/// or already present in an upgraded database - converges over subsequent +/// commits instead of being tied to the single transition commit that first +/// discovered it. pub(super) fn prune_superseded( transaction: &impl LedgerTransaction, - shard_id: &StoreShardIdV1, + submission: &Submission<'_>, + checkpoint: &NextCheckpoint, ) -> Result { - let shard_json = encode_json(shard_id, "shard_json")?; + let persisted_epoch = sqlite_u64( + checkpoint.watermark.authority_epoch.get(), + "authority epoch", + )?; Ok(transaction.execute( DELETE_SUPERSEDED, - params![&shard_json, MAX_PRUNED_ROWS_PER_COMMIT], + params![ + &submission.binding_key.shard_json, + submission.binding_key.incarnation_sql, + persisted_epoch, + MAX_PRUNED_ROWS_PER_COMMIT, + ], )?) } diff --git a/crates/tracedecay-rusqlite-runtime/src/ledger/tests.rs b/crates/tracedecay-rusqlite-runtime/src/ledger/tests.rs index d3f3e29491..f69816d551 100644 --- a/crates/tracedecay-rusqlite-runtime/src/ledger/tests.rs +++ b/crates/tracedecay-rusqlite-runtime/src/ledger/tests.rs @@ -247,15 +247,93 @@ fn advancing_authority_prunes_only_the_superseded_records() { ); } -/// Authority rotation may discover an arbitrarily large legacy ledger. The -/// foreground commit that discovers it must make bounded progress instead of -/// turning the whole backlog into one writer transaction. +/// The safety property that bounds the whole retention rule: a submission whose +/// idempotency record was pruned must still not be able to commit a second +/// time. It fails closed on stale authority instead of being admitted as new. +#[test] +fn a_pruned_record_fails_closed_rather_than_admitting_a_duplicate() { + let mut connection = Connection::open_in_memory().unwrap(); + let transaction = connection.transaction().unwrap(); + initialize_schema(&transaction).unwrap(); + let original = at_authority("operation.original", "key.duplicate", 'a', 1, 7); + let receipt = commit(&transaction, &original); + + let advanced = at_authority("operation.advance", "key.advance", 'a', 1, 8); + commit(&transaction, &advanced); + assert_eq!( + idempotency_rows(&transaction), + vec![(1, 8)], + "the epoch-7 record backing the original receipt has been pruned" + ); + + // The exact same submission arrives again under its now-revoked authority. + let duplicate = record_commit(&transaction, &original, &scope(&original), None); + assert!( + matches!( + duplicate, + Err(LedgerError::StaleAuthority { + persisted, + requested, + }) if persisted == advanced.authority_epoch + && requested == original.authority_epoch + ), + "a resubmission under superseded authority must be refused, got {duplicate:?}" + ); + assert_eq!( + idempotency_rows(&transaction), + vec![(1, 8)], + "the refused duplicate wrote no ledger record" + ); + assert_eq!( + current_watermark(&transaction, &binding(&advanced)) + .unwrap() + .unwrap() + .commit_sequence + .0, + receipt.commit_sequence.0 + 1, + "the refused duplicate did not advance the commit sequence" + ); +} + +/// Without an authority advance the ledger must retain everything: a replay has +/// to keep returning the original receipt. +#[test] +fn records_are_retained_while_their_authority_still_stands() { + let mut connection = Connection::open_in_memory().unwrap(); + let transaction = connection.transaction().unwrap(); + initialize_schema(&transaction).unwrap(); + let original = at_authority("operation.stable", "key.stable", 'a', 1, 7); + let receipt = commit(&transaction, &original); + for index in 0..3 { + let next = at_authority(&format!("operation.more.{index}"), "key.more", 'a', 1, 7); + let _ = record_commit(&transaction, &next, &scope(&next), None).unwrap(); + } + assert_eq!( + idempotency_rows(&transaction).len(), + 2, + "the distinct keys are retained at the standing epoch" + ); + assert!( + matches!( + record_commit(&transaction, &original, &scope(&original), None).unwrap(), + LedgerDisposition::Replay(found) if found == receipt + ), + "a replay under standing authority still returns the original receipt" + ); +} + +/// Authority rotation can discover an arbitrarily large legacy backlog. The +/// foreground commit that discovers it must remove at most one bounded batch, +/// so the epoch advance commits on its own terms instead of dragging the whole +/// backlog into the user mutation's writer transaction. #[test] -fn one_authority_advance_prunes_at_most_256_superseded_records() { +fn one_foreground_commit_prunes_at_most_one_bounded_batch() { + let batch = usize::try_from(prune::MAX_PRUNED_ROWS_PER_COMMIT).unwrap(); + let seeded = batch + 2; let mut connection = Connection::open_in_memory().unwrap(); let transaction = connection.transaction().unwrap(); initialize_schema(&transaction).unwrap(); - for index in 0..258 { + for index in 0..seeded { let metadata = at_authority( &format!("operation.bounded.{index}"), &format!("key.bounded.{index}"), @@ -278,25 +356,27 @@ fn one_authority_advance_prunes_at_most_256_superseded_records() { let rows = idempotency_rows(&transaction); assert_eq!( rows.iter().filter(|(_, epoch)| *epoch == 7).count(), - 2, - "one foreground commit may prune at most 256 legacy records" + seeded - batch, + "one foreground commit removes at most one bounded batch of superseded \ + records, never the whole backlog" ); assert_eq!( rows.iter().filter(|(_, epoch)| *epoch == 8).count(), 1, - "the authority-advancing commit remains recorded" + "the authority-advancing commit still records its own receipt" ); } -/// Cleanup cannot be tied only to the transition commit: that bounded pass can -/// leave a backlog, and an upgraded database can already have its current -/// checkpoint. Later commits at the standing epoch must keep draining it. +/// A bounded pass leaves a backlog, so cleanup cannot be tied to the transition +/// commit alone. Later commits at the standing epoch must keep draining it, or +/// the retention rule never converges. #[test] -fn same_epoch_commits_converge_a_bounded_superseded_backlog() { +fn later_commits_drain_the_remaining_superseded_backlog() { + let batch = usize::try_from(prune::MAX_PRUNED_ROWS_PER_COMMIT).unwrap(); let mut connection = Connection::open_in_memory().unwrap(); let transaction = connection.transaction().unwrap(); initialize_schema(&transaction).unwrap(); - for index in 0..257 { + for index in 0..batch + 1 { let metadata = at_authority( &format!("operation.converge.{index}"), &format!("key.converge.{index}"), @@ -321,7 +401,7 @@ fn same_epoch_commits_converge_a_bounded_superseded_backlog() { .filter(|(_, epoch)| *epoch == 7) .count(), 1, - "the bounded transition pass leaves one legacy record" + "the bounded transition pass leaves the remainder of the backlog" ); let follow_up = at_authority( @@ -337,7 +417,7 @@ fn same_epoch_commits_converge_a_bounded_superseded_backlog() { assert_eq!( rows.iter().filter(|(_, epoch)| *epoch == 7).count(), 0, - "a same-epoch commit continues draining the superseded backlog" + "a later commit at the standing epoch continues draining the backlog" ); assert_eq!( rows.iter().filter(|(_, epoch)| *epoch == 8).count(), @@ -346,77 +426,46 @@ fn same_epoch_commits_converge_a_bounded_superseded_backlog() { ); } -/// The safety property that bounds the whole retention rule: a submission whose -/// idempotency record was pruned must still not be able to commit a second -/// time. It fails closed on stale authority instead of being admitted as new. +/// Only the incarnation whose checkpoint this commit decoded and validated may +/// have its records retired. A neighbouring incarnation whose checkpoint row is +/// inconsistent must never have its receipts deleted on the strength of that +/// row's raw scalar: losing a receipt silently re-admits a duplicate write. #[test] -fn a_pruned_record_fails_closed_rather_than_admitting_a_duplicate() { +fn a_corrupt_neighbouring_checkpoint_cannot_retire_its_receipts() { let mut connection = Connection::open_in_memory().unwrap(); let transaction = connection.transaction().unwrap(); initialize_schema(&transaction).unwrap(); - let original = at_authority("operation.original", "key.duplicate", 'a', 1, 7); - let receipt = commit(&transaction, &original); - - let advanced = at_authority("operation.advance", "key.advance", 'a', 1, 8); - commit(&transaction, &advanced); - assert_eq!( - idempotency_rows(&transaction), - vec![(1, 8)], - "the epoch-7 record backing the original receipt has been pruned" - ); + let neighbour = at_authority("operation.neighbour", "key.neighbour", 'a', 2, 7); + let neighbour_binding = binding(&neighbour); + commit(&transaction, &neighbour); + let seed = at_authority("operation.seed", "key.seed", 'a', 1, 7); + commit(&transaction, &seed); - // The exact same submission arrives again under its now-revoked authority. - let duplicate = record_commit(&transaction, &original, &scope(&original), None); + // Incarnation 2's checkpoint scalar now claims epoch 999 while its + // watermark and receipt still encode 7. Loading that checkpoint fails + // closed, so nothing may act on the raw scalar either. + transaction + .execute( + "UPDATE td_runtime_writer_checkpoint_v1 SET authority_epoch = 999 + WHERE incarnation = 2", + [], + ) + .unwrap(); assert!( matches!( - duplicate, - Err(LedgerError::StaleAuthority { - persisted, - requested, - }) if persisted == advanced.authority_epoch - && requested == original.authority_epoch + current_watermark(&transaction, &neighbour_binding), + Err(LedgerError::Corrupt { .. }) ), - "a resubmission under superseded authority must be refused, got {duplicate:?}" - ); - assert_eq!( - idempotency_rows(&transaction), - vec![(1, 8)], - "the refused duplicate wrote no ledger record" - ); - assert_eq!( - current_watermark(&transaction, &binding(&advanced)) - .unwrap() - .unwrap() - .commit_sequence - .0, - receipt.commit_sequence.0 + 1, - "the refused duplicate did not advance the commit sequence" + "the neighbouring checkpoint is corrupt and fails closed when loaded" ); -} -/// Without an authority advance the ledger must retain everything: a replay has -/// to keep returning the original receipt. -#[test] -fn records_are_retained_while_their_authority_still_stands() { - let mut connection = Connection::open_in_memory().unwrap(); - let transaction = connection.transaction().unwrap(); - initialize_schema(&transaction).unwrap(); - let original = at_authority("operation.stable", "key.stable", 'a', 1, 7); - let receipt = commit(&transaction, &original); - for index in 0..3 { - let next = at_authority(&format!("operation.more.{index}"), "key.more", 'a', 1, 7); - let _ = record_commit(&transaction, &next, &scope(&next), None).unwrap(); - } - assert_eq!( - idempotency_rows(&transaction).len(), - 2, - "the distinct keys are retained at the standing epoch" - ); + // A completely unrelated, fully validated transition on incarnation 1. + let advanced = at_authority("operation.unrelated", "key.unrelated", 'a', 1, 8); + commit(&transaction, &advanced); + assert!( - matches!( - record_commit(&transaction, &original, &scope(&original), None).unwrap(), - LedgerDisposition::Replay(found) if found == receipt - ), - "a replay under standing authority still returns the original receipt" + idempotency_rows(&transaction).contains(&(2, 7)), + "an unvalidated neighbouring checkpoint must not authorise deleting \ + that incarnation's receipts" ); }