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.