From c7d68e7606be1ae37bce9ad245cd30524e64d54c Mon Sep 17 00:00:00 2001 From: thesyncim Date: Wed, 26 Aug 2026 00:04:54 +0100 Subject: [PATCH] fix(cdc): keep single deletes on primary key --- internal/cdc/applier.go | 28 +--------------------------- internal/cdc/pipeline_test.go | 6 ++---- 2 files changed, 3 insertions(+), 31 deletions(-) diff --git a/internal/cdc/applier.go b/internal/cdc/applier.go index 15181e0..76d7a0c 100644 --- a/internal/cdc/applier.go +++ b/internal/cdc/applier.go @@ -4917,7 +4917,6 @@ func appendPrimaryKeyDeletePredicate( primary []targetColumn, tuple Tuple, ) error { - positions := make([]int, len(primary)) for i, column := range primary { datum := tuple[column.sourceIndex] if datum.Kind == DatumUnchangedToast { @@ -4928,37 +4927,12 @@ func appendPrimaryKeyDeletePredicate( return err } *params = append(*params, param) - positions[i] = len(*params) if i != 0 { sql.WriteString(" AND ") } sql.WriteString(column.quoted) - fmt.Fprintf(sql, " = $%d", positions[i]) + fmt.Fprintf(sql, " = $%d", len(*params)) } - if len(primary) < 2 { - return nil - } - writeBound := func(operator string) { - sql.WriteString(" AND ROW(") - for i, column := range primary { - if i != 0 { - sql.WriteByte(',') - } - sql.WriteString(column.quoted) - } - sql.WriteByte(')') - sql.WriteString(operator) - sql.WriteString("ROW(") - for i, position := range positions { - if i != 0 { - sql.WriteByte(',') - } - fmt.Fprintf(sql, "$%d", position) - } - sql.WriteByte(')') - } - writeBound(">=") - writeBound("<=") return nil } diff --git a/internal/cdc/pipeline_test.go b/internal/cdc/pipeline_test.go index bf8f51f..1ad74c3 100644 --- a/internal/cdc/pipeline_test.go +++ b/internal/cdc/pipeline_test.go @@ -454,7 +454,7 @@ func TestPrimaryKeyDeleteUsesCatalogIndexOrder(t *testing.T) { } } -func TestSinglePrimaryKeyDeleteAddsExactCatalogOrderedBounds(t *testing.T) { +func TestSinglePrimaryKeyDeleteUsesExactCatalogOrderedEqualities(t *testing.T) { t.Parallel() relation := &targetRelation{ quoted: `"shard_schema"."read_state"`, @@ -482,9 +482,7 @@ func TestSinglePrimaryKeyDeleteAddsExactCatalogOrderedBounds(t *testing.T) { if err := appendPrimaryKeyDeletePredicate(&sql, ¶ms, relation, primary, tuple); err != nil { t.Fatal(err) } - want := `"app_pk" = $1 AND "user_id" = $2 AND "channel_cid" = $3` + - ` AND ROW("app_pk","user_id","channel_cid")>=ROW($1,$2,$3)` + - ` AND ROW("app_pk","user_id","channel_cid")<=ROW($1,$2,$3)` + want := `"app_pk" = $1 AND "user_id" = $2 AND "channel_cid" = $3` if got := sql.String(); got != want { t.Fatalf("single delete predicate = %q, want %q", got, want) }