From 4704e330b49c01ae43fb478034d91cfdb66fdee3 Mon Sep 17 00:00:00 2001 From: Hasyimi Bahrudin Date: Mon, 5 Oct 2026 22:47:00 +0800 Subject: [PATCH] fix: TRUNCATE keeps rows written earlier in the same batch MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `expand_truncates` turned a TRUNCATE into deletes of the keys in the FileIndex — the rows already materialized — but rows written earlier in the same fold step aren't there yet, so they survived it. The TRUNCATE now also drops the events before it in the step; earlier steps of the unit are in the FileIndex already. Also: the delete half of a split key change carried the update's unchanged TOAST columns, so TOAST resolution ran for a delete. Replayed after a crash between claim and ack, that delete comes after the first copy moved the row, and resolution failed for good ("TOAST resolution failed: PK … is not in any indexed data file"). The split now leaves the delete none, and resolution skips deletes. Co-Authored-By: Claude Opus 5.5 --- crates/pg2iceberg-iceberg/src/materialize.rs | 3 ++- crates/pg2iceberg-logical/src/materializer.rs | 7 ++++-- crates/pg2iceberg-logical/src/pipeline.rs | 2 ++ crates/pg2iceberg-tests/tests/dst.rs | 25 +++++++++++++++++++ 4 files changed, 34 insertions(+), 3 deletions(-) diff --git a/crates/pg2iceberg-iceberg/src/materialize.rs b/crates/pg2iceberg-iceberg/src/materialize.rs index 76632cf..10b10ac 100644 --- a/crates/pg2iceberg-iceberg/src/materialize.rs +++ b/crates/pg2iceberg-iceberg/src/materialize.rs @@ -63,7 +63,8 @@ pub fn resolve_unchanged_cols( prior_rows_by_path: &BTreeMap>, ) -> Result<()> { for r in rows { - if r.unchanged_cols.is_empty() { + // A delete writes only its key: nothing to resolve. + if r.unchanged_cols.is_empty() || r.op == Op::Delete { continue; } let key = r diff --git a/crates/pg2iceberg-logical/src/materializer.rs b/crates/pg2iceberg-logical/src/materializer.rs index e3e53d5..9a1fb1b 100644 --- a/crates/pg2iceberg-logical/src/materializer.rs +++ b/crates/pg2iceberg-logical/src/materializer.rs @@ -1951,7 +1951,9 @@ fn now_micros() -> i64 { /// Expand `Op::Truncate` events into per-PK `Op::Delete` events /// against the current FileIndex. Mirrors PG's TRUNCATE semantics: /// every row known to Iceberg right now is wiped, so the materializer -/// emits an equality-delete for each. +/// emits an equality-delete for each — and the events before it in this +/// step, rows the FileIndex doesn't hold yet, are dropped. (Earlier +/// steps of the unit are in the FileIndex already.) /// /// Subsequent post-truncate Insert/Update events for the same PK /// will overwrite the synthetic Delete during the fold (last-write- @@ -1974,6 +1976,7 @@ fn expand_truncates( out.push(evt); continue; } + out.clear(); // Decode each PK key back into a PK-only Row, over the PK // columns in the order every key was built from (`pk_cols`). // Keys that don't fit (shouldn't happen — FileIndex stores what @@ -2019,7 +2022,7 @@ fn collect_toast_paths( ) -> BTreeSet { let mut paths = BTreeSet::new(); for r in rows { - if r.unchanged_cols.is_empty() { + if r.unchanged_cols.is_empty() || r.op == Op::Delete { continue; } let key = r diff --git a/crates/pg2iceberg-logical/src/pipeline.rs b/crates/pg2iceberg-logical/src/pipeline.rs index 6ae8d97..8ab4f67 100644 --- a/crates/pg2iceberg-logical/src/pipeline.rs +++ b/crates/pg2iceberg-logical/src/pipeline.rs @@ -356,6 +356,8 @@ impl Pipeline { let mut del = evt.clone(); del.op = Op::Delete; del.after = None; + // A delete has no TOASTed values to resolve. + del.unchanged_cols.clear(); self.sink.record_change(del)?; let mut upd = evt; diff --git a/crates/pg2iceberg-tests/tests/dst.rs b/crates/pg2iceberg-tests/tests/dst.rs index 13763c4..7c8ed0d 100644 --- a/crates/pg2iceberg-tests/tests/dst.rs +++ b/crates/pg2iceberg-tests/tests/dst.rs @@ -2520,6 +2520,31 @@ fn truncate_removes_rows_not_yet_materialized() { check_invariants(&mut h).unwrap(); } +/// A key change with a TOASTed column unchanged, replayed after a crash +/// between its claim and the slot ack: the replayed copy's delete of the +/// old key comes after the first copy's moved the row. A delete has no +/// TOASTed value to resolve, so it must not look for one. +#[test] +fn replayed_key_change_with_toast_deletes_without_resolving() { + let mut h = DstHarness::boot(); + h.run_step(&Step::Insert { id: 1, qty: 0 }); + h.run_step(&Step::DriveFlush); + h.run_step(&Step::MaterializerCycle); + h.run_step(&Step::ChangePk { + from: 1, + to: 5, + toast: true, + }); + // Its Begin, change and Commit, claimed before any keepalive: the + // restart resends the transaction committing at the claimed LSN. + h.run_step(&Step::DrivePartial { n: 3 }); + block_on(h.pipeline.flush()).unwrap(); + h.run_step(&Step::CrashMidStream); + h.run_step(&Step::DriveFlush); + h.run_step(&Step::MaterializerCycle); + check_invariants(&mut h).unwrap(); +} + /// A row updated twice in one batch with its TOASTed column unchanged /// both times keeps that column's value: the second update must not /// take the first one's unchanged placeholder for a value.