From a9f4a8f7e47e45252084544b9ebd8d58b0f7d753 Mon Sep 17 00:00:00 2001 From: Zhijie Hou Date: Fri, 2 Oct 2026 16:18:43 +0800 Subject: [PATCH vPOC 1/2] Lock tables against writers when altering publication table membership Publication DDL that names a table took only ShareUpdateExclusiveLock, which does not conflict with the RowExclusiveLock held by concurrent data-modifying statements. This allowed the publication definition to change while a statement was in progress. A data-modifying statement uses the publication definition visible to it to decide whether the operation is allowed and what tuple data must be written to WAL. Logical decoding, however, uses the publication definition visible when the change is committed. If the publication definition changes in between, the publisher and the decoder can use different definitions. For example, an UPDATE on a table without a replica identity could be allowed based on the old definition, but be decoded after a row filter referencing a non-replica-identity column was added, producing a change that the subscriber cannot apply and that makes the apply worker fail repeatedly with "publisher did not send replica identity column expected by the logical replication target relation". Fix this by taking ShareRowExclusiveLock when opening the tables a publication DDL operates on. This conflicts with RowExclusiveLock, preventing writers from running concurrently with the DDL, while still allowing readers to proceed. Naming a partitioned table publishes its partitions implicitly, so lock its whole partition tree in the same mode. The same applies to the partitions of a partitioned table whose exclusion from a FOR ALL TABLES publication is dropped. For the same reason, ALTER PUBLICATION ... SET (publish = ...) now locks the explicitly listed member tables, so that an in-progress writer cannot be checked against the old publication actions. Publications defined by schemas or FOR ALL TABLES have no explicit member list and are not handled here. Reported-by: vignesh C Discussion: https://postgr.es/m/CALDaNm2-UtYX+pFQMEn44bpgZwYX9F9kV2Q4+HnC0TWp-Nx3ig@mail.gmail.com --- .../expected/invalidation_distribution.out | 8 +- .../expected/publication-ddl-update.out | 65 +++++++++++ contrib/test_decoding/meson.build | 1 + .../specs/invalidation_distribution.spec | 14 ++- .../specs/publication-ddl-update.spec | 106 ++++++++++++++++++ src/backend/commands/publicationcmds.c | 46 +++++++- 6 files changed, 227 insertions(+), 13 deletions(-) create mode 100644 contrib/test_decoding/expected/publication-ddl-update.out create mode 100644 contrib/test_decoding/specs/publication-ddl-update.spec 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/expected/publication-ddl-update.out b/contrib/test_decoding/expected/publication-ddl-update.out new file mode 100644 index 00000000000..17303f9c15c --- /dev/null +++ b/contrib/test_decoding/expected/publication-ddl-update.out @@ -0,0 +1,65 @@ +Parsed test spec with 7 sessions + +starting permutation: s1_begin s1_update s2_add_table s1_commit s2_get_binary_changes +step s1_begin: BEGIN; +step s1_update: UPDATE tab1 SET val = 2 WHERE id = 1; +step s2_add_table: ALTER PUBLICATION pub ADD TABLE tab1 WHERE (val = 2 OR val = 1); +step s1_commit: COMMIT; +step s2_add_table: <... completed> +s2: WARNING: skipped loading publication "pub" +DETAIL: The publication does not exist at this point in the WAL. +HINT: Create the publication if it does not exist. +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) = 85; +count +----- + 0 +(1 row) + +?column? +-------- +stop +(1 row) + + +starting permutation: s3_begin s3_update s4_set_publish s3_commit s4_check +step s3_begin: BEGIN; +step s3_update: UPDATE tab2 SET val = 2 WHERE id = 1; +step s4_set_publish: ALTER PUBLICATION pub_ins SET (publish = 'insert, update'); +step s3_commit: COMMIT; +step s4_set_publish: <... completed> +step s4_check: UPDATE tab2 SET val = 3 WHERE id = 1; +ERROR: cannot update table "tab2" because it does not have a replica identity and publishes updates +?column? +-------- +stop +(1 row) + + +starting permutation: s5_begin s5_update_leaf s7_add_parent s5_commit s5_check_leaf +step s5_begin: BEGIN; +step s5_update_leaf: UPDATE part1_1 SET val = 2 WHERE id = 1; +step s7_add_parent: ALTER PUBLICATION pub ADD TABLE part1; +step s5_commit: COMMIT; +step s7_add_parent: <... completed> +step s5_check_leaf: UPDATE part1_1 SET val = 3 WHERE id = 1; +ERROR: cannot update table "part1_1" because it does not have a replica identity and publishes updates +?column? +-------- +stop +(1 row) + + +starting permutation: s7_create_exc s6_begin s6_update_exc s7_unexcept s6_commit s6_check_exc +step s7_create_exc: CREATE PUBLICATION pub_exc FOR ALL TABLES EXCEPT (TABLE part_exc); +step s6_begin: BEGIN; +step s6_update_exc: UPDATE part_exc1 SET val = 2 WHERE id = 1; +step s7_unexcept: ALTER PUBLICATION pub_exc SET ALL TABLES EXCEPT (TABLE t_other); +step s6_commit: COMMIT; +step s7_unexcept: <... completed> +step s6_check_exc: UPDATE part_exc1 SET val = 3 WHERE id = 1; +ERROR: cannot update table "part_exc1" because it does not have a replica identity and publishes updates +?column? +-------- +stop +(1 row) + diff --git a/contrib/test_decoding/meson.build b/contrib/test_decoding/meson.build index ac655853d26..ba49d4a69a9 100644 --- a/contrib/test_decoding/meson.build +++ b/contrib/test_decoding/meson.build @@ -65,6 +65,7 @@ tests += { 'slot_creation_error', 'skip_snapshot_restore', 'invalidation_distribution', + 'publication-ddl-update', 'parallel_session_origin', ], 'regress_args': [ diff --git a/contrib/test_decoding/specs/invalidation_distribution.spec b/contrib/test_decoding/specs/invalidation_distribution.spec index 67d41969ac1..323d595c34e 100644 --- a/contrib/test_decoding/specs/invalidation_distribution.spec +++ b/contrib/test_decoding/specs/invalidation_distribution.spec @@ -34,10 +34,16 @@ step "s3_begin" { BEGIN; } step "s3_insert_tbl1" { INSERT INTO tbl1 (val1, val2) VALUES (2, 2); } step "s3_commit" { COMMIT; } -# Expect to get one insert change. LOGICAL_REP_MSG_INSERT = 'I' +# ALTER PUBLICATION ... ADD TABLE takes a lock that conflicts with the +# RowExclusiveLock held by s1's open transaction, so it has to wait for +# s1_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/contrib/test_decoding/specs/publication-ddl-update.spec b/contrib/test_decoding/specs/publication-ddl-update.spec new file mode 100644 index 00000000000..5feaa67884f --- /dev/null +++ b/contrib/test_decoding/specs/publication-ddl-update.spec @@ -0,0 +1,106 @@ +# Test that publication DDL naming a table serializes with a concurrent +# data-modifying statement on it: the DDL waits for the writer, so the +# writer's change is decoded while the table is not yet published and is not +# sent at all. Without that lock, the DDL could commit while the statement +# is in progress, and the change would be decoded under the new publication +# definition and sent without replica identity data, which the subscriber +# cannot apply. + +setup +{ + SELECT 'init' FROM pg_create_logical_replication_slot('isolation_slot', 'pgoutput'); + CREATE TABLE tab1 (id int, val int); + INSERT INTO tab1 VALUES (1, 1); + CREATE PUBLICATION pub; + CREATE TABLE tab2 (id int, val int); + INSERT INTO tab2 VALUES (1, 1); + CREATE PUBLICATION pub_ins FOR TABLE tab2 WITH (publish = 'insert'); + CREATE TABLE part1 (id int, val int) PARTITION BY RANGE (id); + CREATE TABLE part1_1 PARTITION OF part1 FOR VALUES FROM (0) TO (1000); + INSERT INTO part1_1 VALUES (1, 1); + CREATE TABLE part_exc (id int, val int) PARTITION BY RANGE (id); + CREATE TABLE part_exc1 PARTITION OF part_exc FOR VALUES FROM (0) TO (1000); + INSERT INTO part_exc1 VALUES (1, 1); + CREATE TABLE t_other (id int, val int); + INSERT INTO t_other VALUES (1, 1); +} + +teardown +{ + DROP TABLE tab1; + DROP TABLE tab2; + DROP TABLE part1 CASCADE; + DROP TABLE part_exc CASCADE; + DROP TABLE t_other; + DROP PUBLICATION pub; + DROP PUBLICATION pub_ins; + DROP PUBLICATION IF EXISTS pub_exc; + SELECT 'stop' FROM pg_drop_replication_slot('isolation_slot'); +} + +session "s1" +setup { SET synchronous_commit=on; } + +step "s1_begin" { BEGIN; } +step "s1_update" { UPDATE tab1 SET val = 2 WHERE id = 1; } +step "s1_commit" { COMMIT; } + +session "s2" +setup { SET synchronous_commit=on; } + +step "s2_add_table" { ALTER PUBLICATION pub ADD TABLE tab1 WHERE (val = 2 OR val = 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) = 85; } + +session "s3" +setup { SET synchronous_commit=on; } + +step "s3_begin" { BEGIN; } +step "s3_update" { UPDATE tab2 SET val = 2 WHERE id = 1; } +step "s3_commit" { COMMIT; } + +session "s4" + +step "s4_set_publish" { ALTER PUBLICATION pub_ins SET (publish = 'insert, update'); } +step "s4_check" { UPDATE tab2 SET val = 3 WHERE id = 1; } + +session "s5" + +step "s5_begin" { BEGIN; } +step "s5_update_leaf" { UPDATE part1_1 SET val = 2 WHERE id = 1; } +step "s5_commit" { COMMIT; } +step "s5_check_leaf" { UPDATE part1_1 SET val = 3 WHERE id = 1; } + +session "s6" + +step "s6_begin" { BEGIN; } +step "s6_update_exc" { UPDATE part_exc1 SET val = 2 WHERE id = 1; } +step "s6_commit" { COMMIT; } +step "s6_check_exc" { UPDATE part_exc1 SET val = 3 WHERE id = 1; } + +session "s7" + +step "s7_add_parent" { ALTER PUBLICATION pub ADD TABLE part1; } +step "s7_create_exc" { CREATE PUBLICATION pub_exc FOR ALL TABLES EXCEPT (TABLE part_exc); } +step "s7_unexcept" { ALTER PUBLICATION pub_exc SET ALL TABLES EXCEPT (TABLE t_other); } + +# LOGICAL_REP_MSG_UPDATE = 'U' = 85. ALTER PUBLICATION ... ADD TABLE waits +# for the in-progress UPDATE, so the UPDATE commits before the table is +# published and is not sent. Without the lock, the DDL would commit first +# and the change would be sent without replica identity data. +permutation "s1_begin" "s1_update" "s2_add_table" "s1_commit" "s2_get_binary_changes" + +# ALTER PUBLICATION ... SET (publish = ...) enabling UPDATE waits for an +# in-progress writer on a member table; afterwards, writers on it are +# rejected since it has no replica identity. +permutation "s3_begin" "s3_update" "s4_set_publish" "s3_commit" "s4_check" + +# A partitioned table's partitions are implicitly published with it, so +# adding the parent waits for an in-progress writer on a partition without a +# replica identity; afterwards, writers on it are rejected. +permutation "s5_begin" "s5_update_leaf" "s7_add_parent" "s5_commit" "s5_check_leaf" + +# Removing an exclusion of a partitioned table from a FOR ALL TABLES +# publication publishes its partitions, so the DDL waits for an in-progress +# writer on a partition without a replica identity; afterwards, writers on +# it are rejected. +permutation "s7_create_exc" "s6_begin" "s6_update_exc" "s7_unexcept" "s6_commit" "s6_check_exc" diff --git a/src/backend/commands/publicationcmds.c b/src/backend/commands/publicationcmds.c index 301050104dd..fdb83eafcff 100644 --- a/src/backend/commands/publicationcmds.c +++ b/src/backend/commands/publicationcmds.c @@ -1191,6 +1191,16 @@ AlterPublicationOptions(ParseState *pstate, AlterPublicationStmt *stmt, lfirst_oid(lc)); } + /* + * Changing the actions being published can change the replica identity + * requirements for member tables, and make statements checked against + * the old definition produce changes the subscriber cannot apply. Lock + * the explicitly listed member tables against concurrent writers, like + * OpenTableList() does for membership changes. + */ + foreach_oid(relid, relids) + LockRelationOid(relid, ShareRowExclusiveLock); + schemarelids = GetAllSchemaPublicationRelations(pubform->oid, PUBLICATION_PART_ALL); relids = list_concat_unique_oid(relids, schemarelids); @@ -1389,7 +1399,9 @@ AlterPublicationTables(AlterPublicationStmt *stmt, HeapTuple tup, oldrel->columns = NIL; oldrel->except = false; oldrel->relation = table_open(oldrelid, - ShareUpdateExclusiveLock); + ShareRowExclusiveLock); + (void) find_all_inheritors(oldrelid, ShareRowExclusiveLock, + NULL); delrels = lappend(delrels, oldrel); } } @@ -1836,8 +1848,22 @@ RemovePublicationSchemaById(Oid psoid) /* * Open relations specified by a PublicationTable list. - * The returned tables are locked in ShareUpdateExclusiveLock mode in order to - * add them to a publication. + * + * The returned tables are locked in ShareRowExclusiveLock mode in order to + * add them to a publication. The lock must conflict with the + * RowExclusiveLock held by data-modifying statements. Adding a table to a + * publication, or changing its row filter or column list, can change both + * whether an UPDATE or DELETE on the table is allowed and what tuple data + * the statement writes to WAL. Those decisions are made from the + * publication definition visible when the statement performs its replica + * identity checks and cannot be revisited afterwards, whereas logical + * decoding uses the publication definition visible at commit time. If this + * DDL were allowed to commit while a modification is in progress, WAL + * logging and logical decoding could use different publication definitions, + * resulting in a change that the subscriber cannot apply. Conflicting with + * RowExclusiveLock ensures that writers are blocked while the DDL runs and + * that a writer starting afterwards sees the new publication definition, + * while readers remain unaffected. */ static List * OpenTableList(List *tables) @@ -1862,7 +1888,7 @@ OpenTableList(List *tables) /* Allow query cancel in case this takes a long time */ CHECK_FOR_INTERRUPTS(); - rel = table_openrv(t->relation, ShareUpdateExclusiveLock); + rel = table_openrv(t->relation, ShareRowExclusiveLock); myrelid = RelationGetRelid(rel); /* @@ -1888,7 +1914,7 @@ OpenTableList(List *tables) errmsg("conflicting or redundant column lists for table \"%s\"", RelationGetRelationName(rel)))); - table_close(rel, ShareUpdateExclusiveLock); + table_close(rel, ShareRowExclusiveLock); continue; } @@ -1917,7 +1943,7 @@ OpenTableList(List *tables) List *children; ListCell *child; - children = find_all_inheritors(myrelid, ShareUpdateExclusiveLock, + children = find_all_inheritors(myrelid, ShareRowExclusiveLock, NULL); foreach(child, children) @@ -1980,6 +2006,13 @@ OpenTableList(List *tables) relids_with_collist = lappend_oid(relids_with_collist, childrelid); } } + + /* + * A partitioned table's partitions are implicitly published with it, + * so lock them against concurrent writers, like for the table itself. + */ + else if (rel->rd_rel->relkind == RELKIND_PARTITIONED_TABLE) + (void) find_all_inheritors(myrelid, ShareRowExclusiveLock, NULL); } return rels; @@ -2033,6 +2066,7 @@ LockSchemaList(List *schemalist) } } + /* * Add listed tables to the publication. */ -- 2.34.1