From 605e56bf48fb32a2245c3e7a2da39d70127673a5 Mon Sep 17 00:00:00 2001 From: Mikhail Nikalayeu Date: Wed, 23 Sep 2026 06:58:12 +0530 Subject: [PATCH v6_REL17] Fix tuple search during apply after concurrent index DDL. When the apply worker searches the local relation by index, it can take the first match as the row only if the index is the relation's replica identity or primary key. Otherwise every match has to be compared with the search slot, which holds a complete row only under REPLICA IDENTITY FULL. Only the index OID was saved, so the scan worked this out a second time from the catalogs, and the two answers can differ. Apply holds only RowExclusiveLock, which doesn't conflict with DROP INDEX CONCURRENTLY or REINDEX CONCURRENTLY, so either can demote the chosen index in between. Nothing matches, so the change is silently dropped as an update_missing conflict, and assert-enabled builds fail. Fix this by recording the answer as idxisreplident in the relation map entry and passing the entry down to FindReplTupleInLocalRel(), so it is settled once. Oversight in 89e46da5e5. Author: Mikhail Nikalayeu Reviewed-by: Amit Kapila Reviewed-by: vignesh C Reviewed-by: Zhijie Hou Discussion: https://postgr.es/m/CADzfLwUJovFcnknCC9wjZKECX9xecgnGzC2r2TMV8h4QDD_jwQ@mail.gmail.com Backpatch-through: 16, where it was introduced --- src/backend/executor/execReplication.c | 37 ++++-- src/backend/replication/logical/relation.c | 18 ++- src/backend/replication/logical/worker.c | 87 ++++++------ src/include/executor/executor.h | 5 + src/include/replication/logicalrelation.h | 4 + src/test/subscription/Makefile | 4 +- src/test/subscription/meson.build | 5 +- .../subscription/t/032_subscribe_use_index.pl | 125 ++++++++++++++++++ 8 files changed, 234 insertions(+), 51 deletions(-) diff --git a/src/backend/executor/execReplication.c b/src/backend/executor/execReplication.c index cb1202e4506..9f6da0035ca 100644 --- a/src/backend/executor/execReplication.c +++ b/src/backend/executor/execReplication.c @@ -171,12 +171,17 @@ build_replindex_scan_key(ScanKey skey, Relation rel, Relation idxrel, * * If a matching tuple is found, lock it with lockmode, fill the slot with its * contents, and return true. Return false otherwise. + * + * 'skipduplicates' specifies whether the first matching tuple can be used + * without comparing it against 'searchslot'. If false, all matching tuples are + * compared against 'searchslot', which must contain a complete row. */ bool -RelationFindReplTupleByIndex(Relation rel, Oid idxoid, - LockTupleMode lockmode, - TupleTableSlot *searchslot, - TupleTableSlot *outslot) +RelationFindReplTupleByIndexExt(Relation rel, Oid idxoid, + bool skipduplicates, + LockTupleMode lockmode, + TupleTableSlot *searchslot, + TupleTableSlot *outslot) { ScanKeyData skey[INDEX_MAX_KEYS]; int skey_attoff; @@ -186,13 +191,10 @@ RelationFindReplTupleByIndex(Relation rel, Oid idxoid, Relation idxrel; bool found; TypeCacheEntry **eq = NULL; - bool isIdxSafeToSkipDuplicates; /* Open the index. */ idxrel = index_open(idxoid, RowExclusiveLock); - isIdxSafeToSkipDuplicates = (GetRelationIdentityOrPK(rel) == idxoid); - InitDirtySnapshot(snap); /* Build scan key. */ @@ -213,7 +215,7 @@ retry: * Avoid expensive equality check if the index is primary key or * replica identity index. */ - if (!isIdxSafeToSkipDuplicates) + if (!skipduplicates) { if (eq == NULL) eq = palloc0(sizeof(*eq) * outslot->tts_tupleDescriptor->natts); @@ -298,6 +300,25 @@ retry: return found; } +/* + * ABI-compatible wrapper to emulate old version of the above function. + * Do not call this version in new code. + */ +bool +RelationFindReplTupleByIndex(Relation rel, Oid idxoid, + LockTupleMode lockmode, + TupleTableSlot *searchslot, + TupleTableSlot *outslot) +{ + bool skipduplicates; + + skipduplicates = (GetRelationIdentityOrPK(rel) == idxoid); + + return RelationFindReplTupleByIndexExt(rel, idxoid, + skipduplicates, + lockmode, searchslot, outslot); +} + /* * Compare the tuples in the slots by checking if they have equal values. */ diff --git a/src/backend/replication/logical/relation.c b/src/backend/replication/logical/relation.c index 8794574f42a..1909de8e570 100644 --- a/src/backend/replication/logical/relation.c +++ b/src/backend/replication/logical/relation.c @@ -55,7 +55,7 @@ typedef struct LogicalRepPartMapEntry } LogicalRepPartMapEntry; static Oid FindLogicalRepLocalIndex(Relation localrel, LogicalRepRelation *remoterel, - AttrMap *attrMap); + AttrMap *attrMap, bool *idxisreplident); /* * Relcache invalidation callback for our relation map cache. @@ -453,7 +453,8 @@ logicalrep_rel_open(LogicalRepRelId remoteid, LOCKMODE lockmode) * on the relation). */ entry->localindexoid = FindLogicalRepLocalIndex(entry->localrel, remoterel, - entry->attrmap); + entry->attrmap, + &entry->idxisreplident); entry->localrelvalid = true; } @@ -725,7 +726,8 @@ logicalrep_partition_open(LogicalRepRelMapEntry *root, * anything in the LogicalRepPartMapContext (hence CacheMemoryContext). */ entry->localindexoid = FindLogicalRepLocalIndex(partrel, remoterel, - entry->attrmap); + entry->attrmap, + &entry->idxisreplident); entry->localrelvalid = true; @@ -877,13 +879,18 @@ GetRelationIdentityOrPK(Relation rel) /* * Returns the index oid if we can use an index for subscriber. Otherwise, * returns InvalidOid. + * + * '*idxisreplident' is true if the returned index is the relation's replica + * identity or primary key, and false otherwise. */ static Oid FindLogicalRepLocalIndex(Relation localrel, LogicalRepRelation *remoterel, - AttrMap *attrMap) + AttrMap *attrMap, bool *idxisreplident) { Oid idxoid; + *idxisreplident = false; + /* * We never need index oid for partitioned tables, always rely on leaf * partition's index. @@ -896,7 +903,10 @@ FindLogicalRepLocalIndex(Relation localrel, LogicalRepRelation *remoterel, */ idxoid = GetRelationIdentityOrPK(localrel); if (OidIsValid(idxoid)) + { + *idxisreplident = true; return idxoid; + } if (remoterel->replident == REPLICA_IDENTITY_FULL) { diff --git a/src/backend/replication/logical/worker.c b/src/backend/replication/logical/worker.c index e8eaddbe149..b75cdd411a3 100644 --- a/src/backend/replication/logical/worker.c +++ b/src/backend/replication/logical/worker.c @@ -182,6 +182,7 @@ #include "utils/acl.h" #include "utils/dynahash.h" #include "utils/guc.h" +#include "utils/injection_point.h" #include "utils/inval.h" #include "utils/lsyscache.h" #include "utils/memutils.h" @@ -385,15 +386,13 @@ static void apply_handle_insert_internal(ApplyExecutionData *edata, static void apply_handle_update_internal(ApplyExecutionData *edata, ResultRelInfo *relinfo, TupleTableSlot *remoteslot, - LogicalRepTupleData *newtup, - Oid localindexoid); + LogicalRepTupleData *newtup); static void apply_handle_delete_internal(ApplyExecutionData *edata, ResultRelInfo *relinfo, TupleTableSlot *remoteslot, - Oid localindexoid); + LogicalRepRelMapEntry *relmapentry); static bool FindReplTupleInLocalRel(ApplyExecutionData *edata, Relation localrel, - LogicalRepRelation *remoterel, - Oid localidxoid, + LogicalRepRelMapEntry *relmapentry, TupleTableSlot *remoteslot, TupleTableSlot **localslot); static void apply_handle_tuple_routing(ApplyExecutionData *edata, @@ -2513,11 +2512,8 @@ check_relation_updatable(LogicalRepRelMapEntry *rel) if (rel->updatable) return; - /* - * We are in error mode so it's fine this is somewhat slow. It's better to - * give user correct error. - */ - if (OidIsValid(GetRelationIdentityOrPK(rel->localrel))) + /* Use the entry, so this matches what updatable was decided from. */ + if (rel->idxisreplident) { ereport(ERROR, (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), @@ -2644,7 +2640,7 @@ apply_handle_update(StringInfo s) remoteslot, &newtup, CMD_UPDATE); else apply_handle_update_internal(edata, edata->targetRelInfo, - remoteslot, &newtup, rel->localindexoid); + remoteslot, &newtup); finish_edata(edata); @@ -2668,8 +2664,7 @@ static void apply_handle_update_internal(ApplyExecutionData *edata, ResultRelInfo *relinfo, TupleTableSlot *remoteslot, - LogicalRepTupleData *newtup, - Oid localindexoid) + LogicalRepTupleData *newtup) { EState *estate = edata->estate; LogicalRepRelMapEntry *relmapentry = edata->targetRel; @@ -2680,11 +2675,12 @@ apply_handle_update_internal(ApplyExecutionData *edata, MemoryContext oldctx; EvalPlanQualInit(&epqstate, estate, NULL, NIL, -1, NIL); + + INJECTION_POINT("apply-update-before-open-indices"); + ExecOpenIndices(relinfo, false); - found = FindReplTupleInLocalRel(edata, localrel, - &relmapentry->remoterel, - localindexoid, + found = FindReplTupleInLocalRel(edata, localrel, relmapentry, remoteslot, &localslot); ExecClearTuple(remoteslot); @@ -2803,7 +2799,7 @@ apply_handle_delete(StringInfo s) ExecOpenIndices(relinfo, false); apply_handle_delete_internal(edata, relinfo, - remoteslot, rel->localindexoid); + remoteslot, rel); ExecCloseIndices(relinfo); } @@ -2829,11 +2825,10 @@ static void apply_handle_delete_internal(ApplyExecutionData *edata, ResultRelInfo *relinfo, TupleTableSlot *remoteslot, - Oid localindexoid) + LogicalRepRelMapEntry *relmapentry) { EState *estate = edata->estate; Relation localrel = relinfo->ri_RelationDesc; - LogicalRepRelation *remoterel = &edata->targetRel->remoterel; EPQState epqstate; TupleTableSlot *localslot; bool found; @@ -2845,7 +2840,7 @@ apply_handle_delete_internal(ApplyExecutionData *edata, !localrel->rd_rel->relhasindex || RelationGetIndexList(localrel) == NIL); - found = FindReplTupleInLocalRel(edata, localrel, remoterel, localindexoid, + found = FindReplTupleInLocalRel(edata, localrel, relmapentry, remoteslot, &localslot); /* If found delete it. */ @@ -2880,16 +2875,20 @@ apply_handle_delete_internal(ApplyExecutionData *edata, * the corresponding local relation using either replica identity index, * primary key, index or if needed, sequential scan. * + * 'relmapentry' is the relation map entry for 'localrel'; it tells which + * index to use, if any, and whether that index is the relation's replica + * identity or primary key. + * * Local tuple, if found, is returned in '*localslot'. */ static bool FindReplTupleInLocalRel(ApplyExecutionData *edata, Relation localrel, - LogicalRepRelation *remoterel, - Oid localidxoid, + LogicalRepRelMapEntry *relmapentry, TupleTableSlot *remoteslot, TupleTableSlot **localslot) { EState *estate = edata->estate; + Oid localidxoid = relmapentry->localindexoid; bool found; /* @@ -2901,23 +2900,41 @@ FindReplTupleInLocalRel(ApplyExecutionData *edata, Relation localrel, *localslot = table_slot_create(localrel, &estate->es_tupleTable); Assert(OidIsValid(localidxoid) || - (remoterel->replident == REPLICA_IDENTITY_FULL)); + (relmapentry->remoterel.replident == REPLICA_IDENTITY_FULL)); if (OidIsValid(localidxoid)) { #ifdef USE_ASSERT_CHECKING Relation idxrel = index_open(localidxoid, AccessShareLock); - /* Index must be PK, RI, or usable for REPLICA IDENTITY FULL tables */ - Assert(GetRelationIdentityOrPK(idxrel) == localidxoid || - IsIndexUsableForReplicaIdentityFull(BuildIndexInfo(idxrel), - edata->targetRel->attrmap)); + if (relmapentry->idxisreplident) + { + /* + * We cannot assert this is still the replica identity or primary + * key. DROP INDEX CONCURRENTLY and REINDEX CONCURRENTLY clear + * indisvalid and indisreplident without conflicting with our + * RowExclusiveLock, so GetRelationIdentityOrPK() may no longer + * return it. Unique and non-partial is what the scan actually + * relies on, and no DDL can take those away. + */ + Assert(idxrel->rd_index->indisunique); + Assert(heap_attisnull(idxrel->rd_indextuple, + Anum_pg_index_indpred, NULL)); + } + else + { + /* Otherwise every match is compared, so we need a whole row. */ + Assert(relmapentry->remoterel.replident == REPLICA_IDENTITY_FULL); + Assert(IsIndexUsableForReplicaIdentityFull(BuildIndexInfo(idxrel), + relmapentry->attrmap)); + } index_close(idxrel, AccessShareLock); #endif - found = RelationFindReplTupleByIndex(localrel, localidxoid, - LockTupleExclusive, - remoteslot, *localslot); + found = RelationFindReplTupleByIndexExt(localrel, localidxoid, + relmapentry->idxisreplident, + LockTupleExclusive, + remoteslot, *localslot); } else found = RelationFindReplTupleSeq(localrel, LockTupleExclusive, @@ -3017,8 +3034,7 @@ apply_handle_tuple_routing(ApplyExecutionData *edata, case CMD_DELETE: apply_handle_delete_internal(edata, partrelinfo, - remoteslot_part, - part_entry->localindexoid); + remoteslot_part, part_entry); break; case CMD_UPDATE: @@ -3036,9 +3052,7 @@ apply_handle_tuple_routing(ApplyExecutionData *edata, bool found; /* Get the matching local tuple from the partition. */ - found = FindReplTupleInLocalRel(edata, partrel, - &part_entry->remoterel, - part_entry->localindexoid, + found = FindReplTupleInLocalRel(edata, partrel, part_entry, remoteslot_part, &localslot); if (!found) { @@ -3133,8 +3147,7 @@ apply_handle_tuple_routing(ApplyExecutionData *edata, /* DELETE old tuple found in the old partition. */ apply_handle_delete_internal(edata, partrelinfo, - localslot, - part_entry->localindexoid); + localslot, part_entry); /* INSERT new tuple into the new partition. */ diff --git a/src/include/executor/executor.h b/src/include/executor/executor.h index 7b8b07bd034..d26c79d3b71 100644 --- a/src/include/executor/executor.h +++ b/src/include/executor/executor.h @@ -656,6 +656,11 @@ extern bool RelationFindReplTupleByIndex(Relation rel, Oid idxoid, LockTupleMode lockmode, TupleTableSlot *searchslot, TupleTableSlot *outslot); +extern bool RelationFindReplTupleByIndexExt(Relation rel, Oid idxoid, + bool skipduplicates, + LockTupleMode lockmode, + TupleTableSlot *searchslot, + TupleTableSlot *outslot); extern bool RelationFindReplTupleSeq(Relation rel, LockTupleMode lockmode, TupleTableSlot *searchslot, TupleTableSlot *outslot); diff --git a/src/include/replication/logicalrelation.h b/src/include/replication/logicalrelation.h index e687b40a566..b92927a3670 100644 --- a/src/include/replication/logicalrelation.h +++ b/src/include/replication/logicalrelation.h @@ -32,6 +32,10 @@ typedef struct LogicalRepRelMapEntry Relation localrel; /* relcache entry (NULL when closed) */ AttrMap *attrmap; /* map of local attributes to remote ones */ bool updatable; /* Can apply updates/deletes? */ + bool idxisreplident; /* is localindexoid the relation's replica + * identity or primary key, rather than an + * index usable for a REPLICA IDENTITY FULL + * remote relation? */ Oid localindexoid; /* which index to use, or InvalidOid if none */ /* Sync state. */ diff --git a/src/test/subscription/Makefile b/src/test/subscription/Makefile index ce1ca430095..bb9ded3115b 100644 --- a/src/test/subscription/Makefile +++ b/src/test/subscription/Makefile @@ -13,9 +13,11 @@ subdir = src/test/subscription top_builddir = ../../.. include $(top_builddir)/src/Makefile.global -EXTRA_INSTALL = contrib/hstore +EXTRA_INSTALL = contrib/hstore \ + src/test/modules/injection_points export with_icu +export enable_injection_points check: $(prove_check) diff --git a/src/test/subscription/meson.build b/src/test/subscription/meson.build index c591cd7d619..0b5abbecbb5 100644 --- a/src/test/subscription/meson.build +++ b/src/test/subscription/meson.build @@ -5,7 +5,10 @@ tests += { 'sd': meson.current_source_dir(), 'bd': meson.current_build_dir(), 'tap': { - 'env': {'with_icu': icu.found() ? 'yes' : 'no'}, + 'env': { + 'with_icu': icu.found() ? 'yes' : 'no', + 'enable_injection_points': get_option('injection_points') ? 'yes' : 'no', + }, 'tests': [ 't/001_rep_changes.pl', 't/002_types.pl', diff --git a/src/test/subscription/t/032_subscribe_use_index.pl b/src/test/subscription/t/032_subscribe_use_index.pl index ce812e927fa..e2e3913f63b 100644 --- a/src/test/subscription/t/032_subscribe_use_index.pl +++ b/src/test/subscription/t/032_subscribe_use_index.pl @@ -606,6 +606,131 @@ $node_subscriber->safe_psql('postgres', "DROP TABLE test_invalid"); # Testcase end: Subscription does not use an invalid index # ============================================================================= +# ============================================================================= +# Testcase start: Subscription keeps using an index that concurrent DDL has +# demoted from replica identity +# +# DROP INDEX CONCURRENTLY clears indisvalid and indisreplident and commits that +# before waiting for the lock the apply worker holds, so the worker can be left +# holding an index the catalogs no longer call the replica identity. The drop +# gets no further while apply holds the table, so the index is still complete +# and still maintained, and the change must be applied through it rather than +# dropped as a missing-tuple conflict. +# +# REINDEX CONCURRENTLY reaches the same state by swapping a new index in, but +# the apply path is the same one, so it is not tested separately. + +SKIP: +{ + skip 'Injection points not supported by this build', 5 + unless $ENV{enable_injection_points} eq 'yes'; + + $node_subscriber->safe_psql('postgres', + 'CREATE EXTENSION injection_points'); + + # create tables pub and sub, using a unique index as replica identity + $node_publisher->safe_psql( + 'postgres', q[ + CREATE TABLE test_dropri (x int NOT NULL, y int); + CREATE UNIQUE INDEX test_dropri_ri ON test_dropri (x); + ALTER TABLE test_dropri REPLICA IDENTITY USING INDEX test_dropri_ri; + INSERT INTO test_dropri SELECT i, i FROM generate_series(1,20) i; + CREATE PUBLICATION tap_pub_dropri FOR TABLE test_dropri; + ]); + $node_subscriber->safe_psql( + 'postgres', q[ + CREATE TABLE test_dropri (x int NOT NULL, y int); + CREATE UNIQUE INDEX test_dropri_ri ON test_dropri (x); + ALTER TABLE test_dropri REPLICA IDENTITY USING INDEX test_dropri_ri; + ]); + $node_subscriber->safe_psql('postgres', + "CREATE SUBSCRIPTION tap_sub_dropri CONNECTION '$publisher_connstr application_name=dropri' PUBLICATION tap_pub_dropri" + ); + + # wait for initial table synchronization to finish + $node_subscriber->wait_for_subscription_sync($node_publisher, 'dropri'); + + # Let the worker take the index and stop before opening the relation's + # indexes. The point is attached server-wide: the apply worker is not a + # session this test can attach anything in. + $node_subscriber->safe_psql('postgres', + "SELECT injection_points_attach('apply-update-before-open-indices', 'wait')" + ); + $node_publisher->safe_psql('postgres', + "UPDATE test_dropri SET y = 99 WHERE x = 7"); + $node_subscriber->wait_for_event( + 'logical replication apply worker', + 'apply-update-before-open-indices'); + + # This commits the loss of the replica identity, leaving relreplident set + # to 'i' with no index claiming to be that identity, then parks waiting + # for the apply worker's lock on the table. + my $log_offset = -s $node_subscriber->logfile; + my $drop = $node_subscriber->background_psql('postgres'); + $drop->query_until( + qr/starting_drop/, q[ + \echo starting_drop + DROP INDEX CONCURRENTLY test_dropri_ri; + ]); + $node_subscriber->poll_query_until('postgres', + "SELECT count(*) = 0 FROM pg_index" + . " WHERE indrelid = 'test_dropri'::regclass AND indisreplident") + or die "timed out waiting for the identity index to be invalidated"; + + # Detach before waking, so the worker cannot park on the point again. + $node_subscriber->safe_psql( + 'postgres', + "SELECT injection_points_detach('apply-update-before-open-indices'); + SELECT injection_points_wakeup('apply-update-before-open-indices');" + ); + + # The straddling change goes through, found by the demoted index. + $node_publisher->wait_for_catchup('dropri'); + $result = $node_subscriber->safe_psql('postgres', + "SELECT y FROM test_dropri WHERE x = 7"); + is($result, qq(99), 'change straddling the drop is applied'); + ok($drop->quit, 'DROP INDEX CONCURRENTLY completes'); + + # From the next change on, the relation map entry is rebuilt, finds no + # replica identity, and apply stops with the usual error. + $node_publisher->safe_psql('postgres', + "UPDATE test_dropri SET y = 123 WHERE x = 8"); + ok( $node_subscriber->poll_query_until( + 'postgres', q[ + SELECT apply_error_count > 0 FROM pg_stat_subscription_stats + WHERE subname = 'tap_sub_dropri']), + 'later changes wait for a replica identity'); + like( + slurp_file($node_subscriber->logfile, $log_offset), + qr/logical replication target relation "public\.test_dropri" has neither REPLICA IDENTITY index nor PRIMARY KEY/, + 'and say why'); + + # Give the relation a replica identity again and they resume. + $node_subscriber->safe_psql( + 'postgres', q[ + CREATE UNIQUE INDEX test_dropri_ri2 ON test_dropri (x); + ALTER TABLE test_dropri REPLICA IDENTITY USING INDEX test_dropri_ri2; + ]); + $node_publisher->wait_for_catchup('dropri'); + $result = $node_subscriber->safe_psql('postgres', + "SELECT y FROM test_dropri WHERE x = 8"); + is($result, qq(123), 'replication resumes'); + + # cleanup pub + $node_publisher->safe_psql('postgres', "DROP PUBLICATION tap_pub_dropri"); + $node_publisher->safe_psql('postgres', "DROP TABLE test_dropri"); + # cleanup sub + $node_subscriber->safe_psql('postgres', + "DROP SUBSCRIPTION tap_sub_dropri"); + $node_subscriber->safe_psql('postgres', "DROP TABLE test_dropri"); + $node_subscriber->safe_psql('postgres', + "DROP EXTENSION injection_points"); +} + +# Testcase end: Subscription keeps using an index that concurrent DDL has +# demoted from replica identity +# ============================================================================= + $node_subscriber->stop('fast'); $node_publisher->stop('fast'); -- 2.55.0