diff --git a/src/backend/commands/publicationcmds.c b/src/backend/commands/publicationcmds.c index 301050104dd..9015a5a5136 100644 --- a/src/backend/commands/publicationcmds.c +++ b/src/backend/commands/publicationcmds.c @@ -940,6 +940,14 @@ CreatePublication(ParseState *pstate, CreatePublicationStmt *stmt) if (stmt->for_all_tables) { + /* + * Block concurrent data-modifying statements that may rely on this + * table not being covered by a FOR ALL TABLES publication requiring a + * replica identity. This conflicts with the RowExclusiveLock taken by + * CheckCmdReplicaIdentity() on the same object. + */ + LockRelationOid(PublicationRelationId, ShareRowExclusiveLock); + /* Process EXCEPT table list */ if (exceptrelations != NIL) { @@ -1610,6 +1618,15 @@ AlterPublicationAllFlags(AlterPublicationStmt *stmt, Relation rel, /* Update FOR ALL TABLES flag if changed */ if (stmt->for_all_tables != pubform->puballtables) { + /* + * Block concurrent data-modifying statements when setting FOR ALL + * TABLES publication requiring a replica identity. This conflicts with + * the RowExclusiveLock taken by CheckCmdReplicaIdentity() on the same + * object. + */ + if (stmt->for_all_tables) + LockRelationOid(PublicationRelationId, ShareRowExclusiveLock); + values[Anum_pg_publication_puballtables - 1] = BoolGetDatum(stmt->for_all_tables); replaces[Anum_pg_publication_puballtables - 1] = true; @@ -2005,8 +2022,9 @@ CloseTableList(List *rels) } /* - * Lock the schemas specified in the schema list in AccessShareLock mode in - * order to prevent concurrent schema deletion. + * Lock the schemas in ShareRowExclusiveLock mode to prevent concurrent + * schema deletion and data-modifying statements from racing with the + * publication change. */ static void LockSchemaList(List *schemalist) @@ -2019,7 +2037,7 @@ LockSchemaList(List *schemalist) /* Allow query cancel in case this takes a long time */ CHECK_FOR_INTERRUPTS(); - LockDatabaseObject(NamespaceRelationId, schemaid, 0, AccessShareLock); + LockDatabaseObject(NamespaceRelationId, schemaid, 0, ShareRowExclusiveLock); /* * It is possible that by the time we acquire the lock on schema, diff --git a/src/backend/executor/execReplication.c b/src/backend/executor/execReplication.c index dd42acc13e2..fe7b229e908 100644 --- a/src/backend/executor/execReplication.c +++ b/src/backend/executor/execReplication.c @@ -24,6 +24,7 @@ #include "access/xact.h" #include "access/heapam.h" #include "catalog/pg_am_d.h" +#include "catalog/pg_namespace.h" #include "commands/trigger.h" #include "executor/executor.h" #include "executor/nodeModifyTable.h" @@ -1031,6 +1032,7 @@ void CheckCmdReplicaIdentity(Relation rel, CmdType cmd) { PublicationDesc pubdesc; + bool has_local_identity; /* * Skip checking the replica identity for partitioned tables, because the @@ -1043,6 +1045,30 @@ CheckCmdReplicaIdentity(Relation rel, CmdType cmd) if (cmd != CMD_UPDATE && cmd != CMD_DELETE) return; + /* If relation has replica identity we are always good. */ + has_local_identity = OidIsValid(RelationGetReplicaIndex(rel)) || + rel->rd_rel->relreplident == REPLICA_IDENTITY_FULL; + + /* + * Tables without a local replica identity can race with publication DDL + * for FOR ALL TABLES and TABLES IN SCHEMA. Take locks on the table's + * schema and the publication catalog before using the publication + * descriptor. + * + * The locks conflict with the corresponding publication DDL locks and also + * ensure that any pending catalog invalidations are processed before we + * build the descriptor. + * + * Tables with a local replica identity are not affected by this race, so + * they do not need these additional locks. + */ + if (!has_local_identity) + { + LockDatabaseObject(NamespaceRelationId, RelationGetNamespace(rel), 0, + RowExclusiveLock); + LockRelationOid(PublicationRelationId, RowExclusiveLock); + } + /* * It is only safe to execute UPDATE/DELETE if the relation does not * publish UPDATEs or DELETEs, or all the following conditions are @@ -1104,12 +1130,7 @@ CheckCmdReplicaIdentity(Relation rel, CmdType cmd) RelationGetRelationName(rel)), errdetail("Replica identity must not contain unpublished generated columns."))); - /* If relation has replica identity we are always good. */ - if (OidIsValid(RelationGetReplicaIndex(rel))) - return; - - /* REPLICA IDENTITY FULL is also good for UPDATE/DELETE. */ - if (rel->rd_rel->relreplident == REPLICA_IDENTITY_FULL) + if (has_local_identity) return; /*