From 9178edddf6f68dbf00aeb17332821b8018a9c522 Mon Sep 17 00:00:00 2001 From: Vignesh C Date: Thu, 1 Oct 2026 10:58:19 +0530 Subject: [PATCH v2 2/2] Concurrent UPDATE/DELETE with the following CREATE/ALTER PUBLICATION commands can cause a race: CREATE/ALTER PUBLICATION ... FOR ALL TABLES CREATE/ALTER PUBLICATION ... ADD/SET TABLES IN SCHEMA ALTER PUBLICATION ... SET (publish = ...) The UPDATE/DELETE can pass the replica identity check using the old publication definition, while the publication is changed to include the table for replication concurrently before the update/delete completes. Logical decoding uses the definition visible when the change is committed. This can cause the publisher and decoder to use different definitions, leading to errors on the subscriber. Fix this by locking the scope of the publication change instead of locking all affected tables: Lock the schema OID for TABLES IN SCHEMA. Lock the pg_publication catalog OID for FOR ALL TABLES. Lock the pg_publication catalog OID when SET (publish = ...) enables pubupdate or pubdelete. CheckCmdReplicaIdentity() takes a matching RowExclusiveLock on these objects before checking the publication definition. This is needed only for tables without a local replica identity. RowExclusiveLock is self-compatible, so concurrent UPDATE/DELETE statements do not block each other. Only the publication DDL, which takes ShareRowExclusiveLock, waits for them to finish. Concurrent UPDATE/DELETE with the following CREATE/ALTER PUBLICATION commands can cause a race: CREATE/ALTER PUBLICATION ... FOR ALL TABLES CREATE/ALTER PUBLICATION ... ADD/SET TABLES IN SCHEMA ALTER PUBLICATION ... SET (publish = ...) UPDATE/DELETE can pass the replica identity check using the old publication definition, while the publication is concurrently changed to include the table before the operation completes. Logical decoding uses the publication definition visible at commit. This can cause the publisher and decoder to use different definitions, leading to errors on the subscriber. Fix this by locking the scope of the publication change: Lock the schema OID for TABLES IN SCHEMA. Lock the pg_publication catalog OID for FOR ALL TABLES. Lock the pg_publication catalog OID when SET (publish = ...) enables pubupdate or pubdelete. CheckCmdReplicaIdentity() takes a matching RowExclusiveLock on these objects before checking the publication definition. This is needed only for tables without a local replica identity. RowExclusiveLock is self-compatible, so concurrent UPDATE/DELETE statements do not block each other. Only the publication DDL, which takes ShareRowExclusiveLock, waits for them to finish. --- .../expected/invalidation_distribution.out | 8 +- .../specs/invalidation_distribution.spec | 10 +- src/backend/commands/publicationcmds.c | 50 ++++- src/backend/executor/execReplication.c | 34 ++- src/test/subscription/t/100_bugs.pl | 210 ++++++++++++++++++ 5 files changed, 297 insertions(+), 15 deletions(-) diff --git a/contrib/test_decoding/expected/invalidation_distribution.out b/contrib/test_decoding/expected/invalidation_distribution.out index ae53b1e61de..4eb239fc710 100644 --- a/contrib/test_decoding/expected/invalidation_distribution.out +++ b/contrib/test_decoding/expected/invalidation_distribution.out @@ -4,8 +4,9 @@ starting permutation: s1_insert_tbl1 s1_begin s1_insert_tbl1 s2_alter_pub_add_tb step s1_insert_tbl1: INSERT INTO tbl1 (val1, val2) VALUES (1, 1); step s1_begin: BEGIN; step s1_insert_tbl1: INSERT INTO tbl1 (val1, val2) VALUES (1, 1); -step s2_alter_pub_add_tbl: ALTER PUBLICATION pub ADD TABLE tbl1; +step s2_alter_pub_add_tbl: ALTER PUBLICATION pub ADD TABLE tbl1; step s1_commit: COMMIT; +step s2_alter_pub_add_tbl: <... completed> step s1_insert_tbl1: INSERT INTO tbl1 (val1, val2) VALUES (1, 1); step s2_get_binary_changes: SELECT count(data) FROM pg_logical_slot_get_binary_changes('isolation_slot', NULL, NULL, 'proto_version', '4', 'publication_names', 'pub') WHERE get_byte(data, 0) = 73; count @@ -24,14 +25,15 @@ step s1_begin: BEGIN; step s1_insert_tbl1: INSERT INTO tbl1 (val1, val2) VALUES (1, 1); step s3_begin: BEGIN; step s3_insert_tbl1: INSERT INTO tbl1 (val1, val2) VALUES (2, 2); -step s2_alter_pub_add_tbl: ALTER PUBLICATION pub ADD TABLE tbl1; +step s2_alter_pub_add_tbl: ALTER PUBLICATION pub ADD TABLE tbl1; step s1_insert_tbl1: INSERT INTO tbl1 (val1, val2) VALUES (1, 1); step s1_commit: COMMIT; step s3_commit: COMMIT; +step s2_alter_pub_add_tbl: <... completed> step s2_get_binary_changes: SELECT count(data) FROM pg_logical_slot_get_binary_changes('isolation_slot', NULL, NULL, 'proto_version', '4', 'publication_names', 'pub') WHERE get_byte(data, 0) = 73; count ----- - 1 + 0 (1 row) ?column? diff --git a/contrib/test_decoding/specs/invalidation_distribution.spec b/contrib/test_decoding/specs/invalidation_distribution.spec index 67d41969ac1..4677d07de21 100644 --- a/contrib/test_decoding/specs/invalidation_distribution.spec +++ b/contrib/test_decoding/specs/invalidation_distribution.spec @@ -37,7 +37,11 @@ step "s3_commit" { COMMIT; } # Expect to get one insert change. LOGICAL_REP_MSG_INSERT = 'I' permutation "s1_insert_tbl1" "s1_begin" "s1_insert_tbl1" "s2_alter_pub_add_tbl" "s1_commit" "s1_insert_tbl1" "s2_get_binary_changes" -# Expect to get one insert change with LOGICAL_REP_MSG_INSERT = 'I' from -# the second "s1_insert_tbl1" executed after adding the table tbl1 to the -# publication in "s2_alter_pub_add_tbl". +# ALTER PUBLICATION ... ADD TABLE takes a lock that conflicts with the +# RowExclusiveLock held by s1's and s3's open transactions, so it has to +# wait for both to commit. That makes its own commit happen strictly +# after s1_commit and s3_commit, so every insert in this permutation -- +# including the second "s1_insert_tbl1", despite appearing after +# "s2_alter_pub_add_tbl" in this list -- actually commits before tbl1 is +# added to the publication. Expect to get no insert changes. permutation "s1_begin" "s1_insert_tbl1" "s3_begin" "s3_insert_tbl1" "s2_alter_pub_add_tbl" "s1_insert_tbl1" "s1_commit" "s3_commit" "s2_get_binary_changes" diff --git a/src/backend/commands/publicationcmds.c b/src/backend/commands/publicationcmds.c index 2a8fd82c569..481b5b19d8b 100644 --- a/src/backend/commands/publicationcmds.c +++ b/src/backend/commands/publicationcmds.c @@ -940,6 +940,17 @@ CreatePublication(ParseState *pstate, CreatePublicationStmt *stmt) if (stmt->for_all_tables) { + /* + * Prevent concurrent data-modifying statements from relying on the + * publication not yet covering the table. CheckCmdReplicaIdentity() + * takes RowExclusiveLock on the same object, so this + * ShareRowExclusiveLock conflicts with it. + * + * We cannot lock all existing tables because FOR ALL TABLES also covers + * tables created after this point. + */ + LockRelationOid(PublicationRelationId, ShareRowExclusiveLock); + /* Process EXCEPT table list */ if (exceptrelations != NIL) { @@ -1126,6 +1137,18 @@ AlterPublicationOptions(ParseState *pstate, AlterPublicationStmt *stmt, if (publish_given) { + /* + * Enabling pubupdate or pubdelete can race with a concurrent UPDATE or + * DELETE on a table without a local replica identity. The statement + * maypass the replica identity check before this change takes effect. + * + * Take the same lock as CheckCmdReplicaIdentity() to prevent this + * race. + */ + if ((pubactions.pubupdate && !pubform->pubupdate) || + (pubactions.pubdelete && !pubform->pubdelete)) + LockRelationOid(PublicationRelationId, ShareRowExclusiveLock); + values[Anum_pg_publication_pubinsert - 1] = BoolGetDatum(pubactions.pubinsert); replaces[Anum_pg_publication_pubinsert - 1] = true; @@ -1610,6 +1633,17 @@ AlterPublicationAllFlags(AlterPublicationStmt *stmt, Relation rel, /* Update FOR ALL TABLES flag if changed */ if (stmt->for_all_tables != pubform->puballtables) { + /* + * Only widening (false -> true) needs to block concurrent + * data-modifying statements, since they may have already passed the + * replica identity check using the old publication definition. + * + * Lock the publication catalog once rather than locking each affected + * table, since the change can affect an unbounded set of tables. + */ + 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; @@ -2017,8 +2051,18 @@ CloseTableList(List *rels) } /* - * Lock the schemas specified in the schema list in AccessShareLock mode in - * order to prevent concurrent schema deletion. + * Lock the schemas specified in the schema list in ShareRowExclusiveLock mode + * in order to prevent concurrent schema deletion. + * + * ShareRowExclusiveLock, rather than AccessShareLock, is required because a + * data-modifying statement on a table in this schema may currently be + * relying on the schema not yet being part of any publication that requires + * a replica identity. Enumerating and locking every table in the schema + * instead does not scale to schemas with many tables, so we lock the schema + * itself: this conflicts with the RowExclusiveLock such a statement takes + * on the same object (see CheckCmdReplicaIdentity()), forcing this DDL to + * wait for it to finish and forcing the statement, if it starts afterward, + * to see this schema's post-DDL publication membership. */ static void LockSchemaList(List *schemalist) @@ -2031,7 +2075,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..ff191bba24b 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,31 @@ 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; + + /* + * A table without a local replica identity can race with publication DDL + * that can make the table require one, such as FOR ALL TABLES or TABLES + * IN SCHEMA. The publication definition seen by this statement may be + * stale if the DDL runs concurrently. + * + * Take locks on the table's schema and the publication catalog before + * checking the publication definition. These locks conflict with the + * corresponding publication DDL and ensure that any pending catalog + * invalidations are processed before building the descriptor. + * + * Tables with a local replica identity are not affected by this race, + * so no additional locking is needed for them. + */ + 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 +1131,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; /* diff --git a/src/test/subscription/t/100_bugs.pl b/src/test/subscription/t/100_bugs.pl index 808d592cb6e..b5e2d12b257 100644 --- a/src/test/subscription/t/100_bugs.pl +++ b/src/test/subscription/t/100_bugs.pl @@ -768,6 +768,69 @@ DROP SUBSCRIPTION sub_drop_refresh; qq{DROP PUBLICATION pub_drop_refresh, pub_seq_drop_refresh;}); } +# ============================================================================= +# ALTER PUBLICATION ... SET (publish = ...) can widen which actions a +# publication requires a replica identity for, on tables that are already +# members and do not change membership at all. This is the same race as FOR +# ALL TABLES / TABLES IN SCHEMA widening membership itself, and is closed the +# same way: CheckCmdReplicaIdentity() already takes RowExclusiveLock on the +# publication catalog unconditionally for any relation without a local +# replica identity, so AlterPublicationOptions() only needs to take the +# conflicting ShareRowExclusiveLock on the same object when pubupdate or +# pubdelete is turned on. +# ============================================================================= + +$node_publisher->safe_psql( + 'postgres', qq{ + CREATE TABLE tab_pubrace_options (id int, val int); + INSERT INTO tab_pubrace_options VALUES (1, 1); + CREATE PUBLICATION pub_pubrace_options FOR TABLE tab_pubrace_options + WITH (publish = 'insert'); +}); + +# Hold an UPDATE open on a table with no replica identity, currently +# published for insert only, so the UPDATE is allowed under the old +# definition. +my $options_dml = $node_publisher->background_psql('postgres'); +$options_dml->query_safe( + "BEGIN;UPDATE tab_pubrace_options SET val = 2 WHERE id = 1;"); + +# Issue the DDL without waiting for it: it must not be able to commit while +# the UPDATE is in progress. +my $options_ddl = $node_publisher->background_psql('postgres'); +$options_ddl->query_until( + qr/issued/, q{ + \echo issued + ALTER PUBLICATION pub_pubrace_options SET (publish = 'insert, update'); +}); + +ok( $node_publisher->poll_query_until( + 'postgres', qq{ + SELECT EXISTS (SELECT 1 FROM pg_locks + WHERE relation = 'pg_publication'::regclass + AND mode = 'ShareRowExclusiveLock' + AND NOT granted)}), + 'ALTER PUBLICATION SET (publish = ...) waits for a concurrent data-modifying statement' +); + +$options_dml->query_safe('COMMIT'); +$options_dml->quit; +$options_ddl->query_until(qr/finished/, "\\echo finished\n"); +$options_ddl->quit; + +# The publication now publishes updates for a table with no replica +# identity: further UPDATEs must be rejected. +($ret, $stdout, $stderr) = $node_publisher->psql('postgres', + 'UPDATE tab_pubrace_options SET val = 3 WHERE id = 1'); +ok( $stderr =~ + qr/cannot update table "tab_pubrace_options" because it does not have a replica identity and publishes updates/, + 'UPDATE is correctly rejected once SET (publish = ...) covers updates for the table' +); + +# Clean up +$node_publisher->safe_psql('postgres', + 'DROP PUBLICATION pub_pubrace_options; DROP TABLE tab_pubrace_options;'); + $node_publisher->stop('fast'); $node_subscriber->stop('fast'); @@ -881,6 +944,153 @@ $node_subscriber->safe_psql('postgres', 'DROP SUBSCRIPTION sub_pubrace'); $node_publisher->safe_psql('postgres', 'DROP PUBLICATION pub_pubrace_sync, pub_pubrace_filtered'); +# ============================================================================= +# FOR ALL TABLES and TABLES IN SCHEMA publication DDL cannot name every +# affected table up front, so it cannot take the per-table +# ShareRowExclusiveLock used above. Instead it locks a single object scoped to +# what it affects (the schema, or the publication catalog itself when even the +# schema isn't known, as with FOR ALL TABLES), and CheckCmdReplicaIdentity() +# takes a matching lock, on the same object, before trusting the publication +# descriptor for a table with no local replica identity. This interlock does +# not require any subscriber; it is verified directly against pg_locks. +# ============================================================================= + +# --- TABLES IN SCHEMA --- + +$node_publisher->safe_psql( + 'postgres', qq{ + CREATE SCHEMA sch_pubrace; + CREATE TABLE sch_pubrace.tab_pubrace_schema (id int, val int); + INSERT INTO sch_pubrace.tab_pubrace_schema VALUES (1, 1); + CREATE PUBLICATION pub_pubrace_schema; +}); + +# Hold an UPDATE open on a table with no replica identity, not yet covered by +# any publication. +my $schema_dml = $node_publisher->background_psql('postgres'); +$schema_dml->query_safe( + "BEGIN;UPDATE sch_pubrace.tab_pubrace_schema SET val = 2 WHERE id = 1;" +); + +# Issue the DDL without waiting for it: it must not be able to commit while +# the UPDATE is in progress. +my $schema_ddl = $node_publisher->background_psql('postgres'); +$schema_ddl->query_until( + qr/issued/, q{ + \echo issued + ALTER PUBLICATION pub_pubrace_schema ADD TABLES IN SCHEMA sch_pubrace; +}); + +ok( $node_publisher->poll_query_until( + 'postgres', qq{ + SELECT EXISTS (SELECT 1 FROM pg_locks + WHERE locktype = 'object' + AND classid = 'pg_namespace'::regclass + AND objid = 'sch_pubrace'::regnamespace + AND mode = 'ShareRowExclusiveLock' + AND NOT granted)}), + 'ALTER PUBLICATION ADD TABLES IN SCHEMA waits for a concurrent data-modifying statement' +); + +$schema_dml->query_safe('COMMIT'); +$schema_dml->quit; +$schema_ddl->query_until(qr/finished/, "\\echo finished\n"); +$schema_ddl->quit; + +# The table is now covered by the publication and has no replica identity: +# further UPDATEs must be rejected, proving the DDL's effect was not missed. +($ret, $stdout, $stderr) = $node_publisher->psql('postgres', + 'UPDATE sch_pubrace.tab_pubrace_schema SET val = 3 WHERE id = 1'); +ok( $stderr =~ + qr/cannot update table "tab_pubrace_schema" because it does not have a replica identity and publishes updates/, + 'UPDATE is correctly rejected once TABLES IN SCHEMA covers the table'); + +# Clean up +$node_publisher->safe_psql('postgres', + 'DROP PUBLICATION pub_pubrace_schema; DROP SCHEMA sch_pubrace CASCADE;'); + +# --- FOR ALL TABLES, via CREATE PUBLICATION --- + +$node_publisher->safe_psql( + 'postgres', qq{ + CREATE TABLE tab_pubrace_all (id int, val int); + INSERT INTO tab_pubrace_all VALUES (1, 1); +}); + +my $all_dml = $node_publisher->background_psql('postgres'); +$all_dml->query_safe( + "BEGIN;UPDATE tab_pubrace_all SET val = 2 WHERE id = 1;"); + +my $all_ddl = $node_publisher->background_psql('postgres'); +$all_ddl->query_until( + qr/issued/, q{ + \echo issued + CREATE PUBLICATION pub_pubrace_all FOR ALL TABLES; +}); + +ok( $node_publisher->poll_query_until( + 'postgres', qq{ + SELECT EXISTS (SELECT 1 FROM pg_locks + WHERE relation = 'pg_publication'::regclass + AND mode = 'ShareRowExclusiveLock' + AND NOT granted)}), + 'CREATE PUBLICATION FOR ALL TABLES waits for a concurrent data-modifying statement' +); + +$all_dml->query_safe('COMMIT'); +$all_dml->quit; +$all_ddl->query_until(qr/finished/, "\\echo finished\n"); +$all_ddl->quit; + +($ret, $stdout, $stderr) = $node_publisher->psql('postgres', + 'UPDATE tab_pubrace_all SET val = 3 WHERE id = 1'); +ok( $stderr =~ + qr/cannot update table "tab_pubrace_all" because it does not have a replica identity and publishes updates/, + 'UPDATE is correctly rejected once FOR ALL TABLES covers the table'); + +$node_publisher->safe_psql('postgres', 'DROP PUBLICATION pub_pubrace_all'); + +# --- FOR ALL TABLES, via ALTER PUBLICATION ... SET ALL TABLES --- + +$node_publisher->safe_psql('postgres', + 'CREATE PUBLICATION pub_pubrace_setall'); + +my $setall_dml = $node_publisher->background_psql('postgres'); +$setall_dml->query_safe( + "BEGIN;UPDATE tab_pubrace_all SET val = 4 WHERE id = 1;"); + +my $setall_ddl = $node_publisher->background_psql('postgres'); +$setall_ddl->query_until( + qr/issued/, q{ + \echo issued + ALTER PUBLICATION pub_pubrace_setall SET ALL TABLES; +}); + +ok( $node_publisher->poll_query_until( + 'postgres', qq{ + SELECT EXISTS (SELECT 1 FROM pg_locks + WHERE relation = 'pg_publication'::regclass + AND mode = 'ShareRowExclusiveLock' + AND NOT granted)}), + 'ALTER PUBLICATION SET ALL TABLES waits for a concurrent data-modifying statement' +); + +$setall_dml->query_safe('COMMIT'); +$setall_dml->quit; +$setall_ddl->query_until(qr/finished/, "\\echo finished\n"); +$setall_ddl->quit; + +($ret, $stdout, $stderr) = $node_publisher->psql('postgres', + 'UPDATE tab_pubrace_all SET val = 5 WHERE id = 1'); +ok( $stderr =~ + qr/cannot update table "tab_pubrace_all" because it does not have a replica identity and publishes updates/, + 'UPDATE is correctly rejected once ALTER PUBLICATION SET ALL TABLES covers the table' +); + +# Clean up +$node_publisher->safe_psql('postgres', + 'DROP PUBLICATION pub_pubrace_setall; DROP TABLE tab_pubrace_all;'); + $node_publisher->stop('fast'); $node_subscriber->stop('fast'); -- 2.55.0