Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion crates/pg2iceberg-iceberg/src/materialize.rs
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,8 @@ pub fn resolve_unchanged_cols(
prior_rows_by_path: &BTreeMap<String, Vec<Row>>,
) -> 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
Expand Down
7 changes: 5 additions & 2 deletions crates/pg2iceberg-logical/src/materializer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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-
Expand All @@ -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
Expand Down Expand Up @@ -2019,7 +2022,7 @@ fn collect_toast_paths(
) -> BTreeSet<String> {
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
Expand Down
2 changes: 2 additions & 0 deletions crates/pg2iceberg-logical/src/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -356,6 +356,8 @@ impl<C: Coordinator + ?Sized> Pipeline<C> {
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;
Expand Down
25 changes: 25 additions & 0 deletions crates/pg2iceberg-tests/tests/dst.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading