From 94d3df62187297caa438f6c5abe3eee16041a1ce Mon Sep 17 00:00:00 2001 From: Zhijie Hou Date: Sun, 4 Oct 2026 09:17:14 +0800 Subject: [PATCH vPOC 2/2] Synchronize publication scope widening with concurrent writers A data-modifying statement checks the replica identity requirement and decides what tuple data to write to WAL using the publication definition visible to the statement, while logical decoding uses the definition in effect when the change is committed. For publications whose member tables are not named explicitly (FOR ALL TABLES, TABLES IN SCHEMA), and for ALTER PUBLICATION ... SET (publish = ...) enabling UPDATE or DELETE, the set of affected tables is not known up front, so the per-table locking used for explicitly named tables does not apply. A concurrent UPDATE or DELETE on a table without a replica identity could thus still pass the replica identity check using the old definition, producing a change the subscriber cannot apply. Close the race mostly on the DDL side, keeping it out of the DML hot path: * When publication DDL widens the set of changes some tables publish, lock the affected plain tables that have no usable replica identity (no usable replica identity index, and not REPLICA IDENTITY FULL) in ShareRowExclusiveLock mode, blocking concurrent writers on them. Tables with a usable replica identity need no locking, because their changes carry the identity data regardless of the publication definition. Nothing is locked when the publication does not publish UPDATEs or DELETEs, since only those actions require a replica identity. * The DDL then holds ShareLock on the pg_publication catalog until commit. DDL that removes a table's replica identity (ALTER TABLE ... REPLICA IDENTITY NOTHING or DEFAULT without a usable primary key, dropping the replica identity index) takes a conflicting RowExclusiveLock on the same catalog, so the replica identity state surveyed by the publication DDL stays valid until it commits. ShareLock is the weakest mode conflicting with that RowExclusiveLock, and it does not serialize concurrent publication DDLs against each other. * RelationBuildPublicationDesc() takes RowExclusiveLock on the pg_publication catalog and holds it until end of transaction when it builds the publication descriptor of a table that has no usable replica identity and does not publish UPDATEs and DELETEs, so that the DDL cannot commit in between. This covers tables that did not exist or were not visible when the DDL surveyed the existing tables; acquiring the lock also processes pending invalidation messages, so the descriptor is rebuilt from up-to-date catalogs if it was invalidated meanwhile. No lock is taken when the table publishes UPDATEs and DELETEs, since any UPDATE or DELETE on it is rejected anyway. To avoid deadlocks, all three places acquire locks in table-first, catalog-second order: DML and replica identity removal necessarily hold a lock on their table at that point, so the publication DDL cannot take the catalog lock first. Since each locked table consumes an entry in the shared lock table, the locked tables are counted along the way and a clear error is raised if they exceed the lock table size (minus one per-process share reserved for the transaction's own locks), instead of failing later with "out of shared memory". Discussion: https://postgr.es/m/CALDaNm2-UtYX+pFQMEn44bpgZwYX9F9kV2Q4+HnC0TWp-Nx3ig@mail.gmail.com --- src/backend/catalog/pg_publication.c | 34 ++++ src/backend/commands/publicationcmds.c | 174 +++++++++++++++++- src/backend/utils/cache/relcache.c | 8 + src/include/catalog/pg_publication.h | 1 + .../expected/publication-ddl-race.out | 92 +++++++++ src/test/isolation/isolation_schedule | 1 + .../isolation/specs/publication-ddl-race.spec | 148 +++++++++++++++ 7 files changed, 457 insertions(+), 1 deletion(-) create mode 100644 src/test/isolation/expected/publication-ddl-race.out create mode 100644 src/test/isolation/specs/publication-ddl-race.spec diff --git a/src/backend/catalog/pg_publication.c b/src/backend/catalog/pg_publication.c index f11ea8839ca..719f360fe35 100644 --- a/src/backend/catalog/pg_publication.c +++ b/src/backend/catalog/pg_publication.c @@ -1301,6 +1301,40 @@ GetAllSchemaPublicationRelations(Oid pubid, PublicationPartOpt pub_partopt) return result; } +/* + * Get the list of all plain tables in the database that could be published + * by a FOR ALL TABLES publication. + */ +List * +GetAllPublishableTables(void) +{ + Relation classRel; + TableScanDesc scan; + HeapTuple tuple; + List *result = NIL; + + classRel = table_open(RelationRelationId, AccessShareLock); + + scan = table_beginscan_catalog(classRel, 0, NULL); + while ((tuple = heap_getnext(scan, ForwardScanDirection)) != NULL) + { + Form_pg_class relForm = (Form_pg_class) GETSTRUCT(tuple); + + /* Only plain tables receive the data-modifying statements. */ + if (relForm->relkind != RELKIND_RELATION) + continue; + + if (!is_publishable_class(relForm->oid, relForm)) + continue; + + result = lappend_oid(result, relForm->oid); + } + + table_endscan(scan); + table_close(classRel, AccessShareLock); + return result; +} + /* * Get publication using oid * diff --git a/src/backend/commands/publicationcmds.c b/src/backend/commands/publicationcmds.c index fdb83eafcff..fc82ec15ca9 100644 --- a/src/backend/commands/publicationcmds.c +++ b/src/backend/commands/publicationcmds.c @@ -15,7 +15,9 @@ #include "postgres.h" #include "access/htup_details.h" +#include "access/relation.h" #include "access/table.h" +#include "access/twophase.h" #include "access/xact.h" #include "catalog/catalog.h" #include "catalog/indexing.h" @@ -39,6 +41,7 @@ #include "parser/parse_relation.h" #include "rewrite/rewriteHandler.h" #include "storage/lmgr.h" +#include "storage/lock.h" #include "utils/acl.h" #include "utils/builtins.h" #include "utils/inval.h" @@ -64,6 +67,8 @@ typedef struct rf_context static List *OpenTableList(List *tables); static void CloseTableList(List *rels); static void LockSchemaList(List *schemalist); +static void LockTablesWithoutReplicaIdentity(Form_pg_publication pubform, + List *schemas); static void PublicationAddTables(Oid pubid, List *rels, bool if_not_exists, AlterPublicationStmt *stmt); static void PublicationDropTables(Oid pubid, List *rels, bool missing_ok); @@ -925,7 +930,6 @@ CreatePublication(ParseState *pstate, CreatePublicationStmt *stmt) /* Insert tuple into catalog. */ CatalogTupleInsert(rel, tup); - heap_freetuple(tup); recordDependencyOnOwner(PublicationRelationId, puboid, GetUserId()); @@ -938,6 +942,14 @@ CreatePublication(ParseState *pstate, CreatePublicationStmt *stmt) ObjectsInPublicationToOids(stmt->pubobjects, pstate, &relations, &exceptrelations, &schemaidlist); + /* + * Prevent concurrent writers on tables that have no usable replica + * identity from racing with this DDL; see + * LockTablesWithoutReplicaIdentity(). + */ + LockTablesWithoutReplicaIdentity((Form_pg_publication) GETSTRUCT(tup), + schemaidlist); + if (stmt->for_all_tables) { /* Process EXCEPT table list */ @@ -1161,6 +1173,15 @@ AlterPublicationOptions(ParseState *pstate, AlterPublicationStmt *stmt, pubform = (Form_pg_publication) GETSTRUCT(tup); + /* + * Enabling publication of UPDATEs or DELETEs makes member tables that + * have no usable replica identity require one; prevent concurrent writers + * on such tables. see LockTablesWithoutReplicaIdentity(). + */ + if (publish_given) + LockTablesWithoutReplicaIdentity(pubform, + GetPublicationSchemas(pubform->oid)); + /* Invalidate the relcache. */ if (pubform->puballtables) { @@ -1440,6 +1461,17 @@ AlterPublicationSchemas(AlterPublicationStmt *stmt, if (!schemaidlist && stmt->action != AP_SetObjects) return; + /* + * Prevent concurrent writers on the schemas' tables. See + * LockTablesWithoutReplicaIdentity(). + * + * Note that the catalog lock taken there must not be taken when no schema + * is being added: acquiring it before the table locks taken by later + * parts of the DDL would invert the lock ordering. + */ + if (stmt->action != AP_DropObjects) + LockTablesWithoutReplicaIdentity(pubform, schemaidlist); + /* * Schema lock is held until the publication is altered to prevent * concurrent schema deletion. @@ -1644,6 +1676,13 @@ AlterPublicationAllFlags(AlterPublicationStmt *stmt, Relation rel, CatalogTupleUpdate(rel, &tup->t_self, tup); CommandCounterIncrement(); + /* + * Prevent concurrent writers on the publication's tables. See + * LockTablesWithoutReplicaIdentity(). + */ + LockTablesWithoutReplicaIdentity((Form_pg_publication) GETSTRUCT(tup), + NIL); + /* For ALL TABLES, we must invalidate all relcache entries */ if (replaces[Anum_pg_publication_puballtables - 1]) CacheInvalidateRelcacheAll(); @@ -2066,6 +2105,139 @@ LockSchemaList(List *schemalist) } } +/* + * Prevent data-modifying statements from racing with a publication DDL that + * widens the set of changes tables publish (adding TABLES IN SCHEMA or ALL + * TABLES to a publication, enabling publication of UPDATE or DELETE), by + * locking the covered plain tables that have no usable replica identity in + * ShareRowExclusiveLock mode. + * + * An UPDATE or DELETE checks the replica identity requirement and decides + * what tuple data to write to WAL using the publication definition visible + * to the statement, while logical decoding uses the definition in effect at + * commit time. If the DDL could commit while such a statement is in + * progress on a table without a replica identity, the change could be + * decoded and sent without the replica identity data the subscriber needs + * to apply it. The following lock protocol prevents that: + * + * - The DDL locks the covered tables without a usable replica identity in + * ShareRowExclusiveLock mode, which conflicts with the RowExclusiveLock + * held by writers, so it waits for in-progress writers, and writers + * starting afterwards see the new publication definition. Tables with a + * usable replica identity need not be locked: their changes carry the + * identity data regardless of the publication definition. + * + * - The DDL then holds ShareLock on the pg_publication catalog until + * commit, and RelationBuildPublicationDesc() takes a conflicting + * RowExclusiveLock on it, held until end of transaction, when building + * the publication descriptor of a table that has no usable replica + * identity. The descriptor is (re)built exactly when the table may be + * new to the backend, e.g. a table created after the DDL surveyed the + * existing tables, so the DDL cannot commit while such a statement is in + * progress. Note that this serializes the DDL with writers on tables + * that are out of the DDL's scope too, but only with those touching + * tables without a usable replica identity. + * + * No lock is needed to serialize with DDL that removes a table's replica + * identity: such removal takes AccessExclusiveLock on the table, so a writer + * on the table either ends before the removal commits (its WAL was written + * while the identity existed), or starts after it and then takes the catalog + * lock when rebuiling publication descriptor as described above, blocking the + * publication DDL's commit. + * + * To avoid deadlocks, the catalog lock is taken only after the table locks: + * a data-modifying statement necessarily holds a lock on its table when + * building the publication descriptor, so the DDL cannot take the catalog + * lock first. Note that a transaction holding the catalog lock from an + * earlier descriptor build can still deadlock with this DDL when it later + * requests a lock on a table locked here; such cycles are resolved by the + * deadlock detector. + * + * Since each locked table consumes an entry in the shared lock table, the + * tables locked are counted along the way, and a clear error is raised if + * they exceed half of the lock table size (cf. NLOCKENTS() in lock.c). This + * is expected to be rare, since a table without a replica identity cannot + * be updated or deleted through the publication anyway. + */ +static void +LockTablesWithoutReplicaIdentity(Form_pg_publication pubform, List *schemas) +{ + List *relids = NIL; + uint64 maxlocks; + int num_without_ri = 0; + + Assert(!(pubform->puballtables && schemas)); + + if (!pubform->pubupdate && !pubform->pubdelete) + return; + + /* + * Nothing is covered by the publication, now or in the future; no lock is + * needed. Note that an empty schema still needs the catalog lock below, + * to cover tables created in it while this DDL is in progress. + */ + if (!pubform->puballtables && schemas == NIL) + return; + + if (pubform->puballtables) + relids = GetAllPublishableTables(); + + foreach_oid(schemaid, schemas) + relids = list_concat(relids, + GetSchemaPublicationRelations(schemaid, + PUBLICATION_PART_LEAF)); + + maxlocks = (uint64) (max_locks_per_xact * (MaxBackends + max_prepared_xacts - 1)) / 2; + + /* + * Lock the tables in OID order, so that concurrent publication DDLs + * cannot deadlock with each other by locking them in different orders. + */ + list_sort(relids, list_oid_cmp); + list_deduplicate_oid(relids); + + foreach_oid(relid, relids) + { + Relation rel; + + /* Allow query cancel in case this takes a long time */ + CHECK_FOR_INTERRUPTS(); + + rel = try_relation_open(relid, ShareRowExclusiveLock); + if (rel == NULL) + continue; /* concurrently dropped */ + + /* + * Skip if not a publishable table, is a sequence, or has a usable + * replica identity. + */ + if (!is_publishable_relation(rel) || + rel->rd_rel->relkind == RELKIND_SEQUENCE || + OidIsValid(RelationGetReplicaIndex(rel)) || + rel->rd_rel->relreplident == REPLICA_IDENTITY_FULL) + { + table_close(rel, ShareRowExclusiveLock); + continue; + } + + if (++num_without_ri > maxlocks) + ereport(ERROR, + (errcode(ERRCODE_CONFIGURATION_LIMIT_EXCEEDED), + errmsg("too many tables without a replica identity in the publication"), + errhint("Set a replica identity on the tables that need to be updated or deleted, or exclude them from the publication."))); + + /* Hold the lock until end of transaction */ + table_close(rel, NoLock); + } + + /* + * Serialize with concurrent publication descriptor builds. ShareLock is + * the weakest mode that conflicts with the RowExclusiveLock those take on + * this catalog, and, unlike ShareRowExclusiveLock, it does not needlessly + * serialize concurrent publication DDLs against each other. + */ + LockRelationOid(PublicationRelationId, ShareLock); +} /* * Add listed tables to the publication. diff --git a/src/backend/utils/cache/relcache.c b/src/backend/utils/cache/relcache.c index d8f04a05309..da99c3c62bc 100644 --- a/src/backend/utils/cache/relcache.c +++ b/src/backend/utils/cache/relcache.c @@ -5866,6 +5866,14 @@ RelationBuildPublicationDesc(Relation relation, PublicationDesc *pubdesc) return; } + /* + * Prevent concurrent publication DDL to publish the table lacking replica + * identity (see LockTablesWithoutReplicaIdentity). + */ + if (!OidIsValid(RelationGetReplicaIndex(relation)) && + relation->rd_rel->relreplident != REPLICA_IDENTITY_FULL) + LockRelationOid(PublicationRelationId, RowExclusiveLock); + memset(pubdesc, 0, sizeof(PublicationDesc)); pubdesc->rf_valid_for_update = true; pubdesc->rf_valid_for_delete = true; diff --git a/src/include/catalog/pg_publication.h b/src/include/catalog/pg_publication.h index 5d1e6c54a85..4e2d3fd1036 100644 --- a/src/include/catalog/pg_publication.h +++ b/src/include/catalog/pg_publication.h @@ -188,6 +188,7 @@ extern List *GetSchemaPublicationRelations(Oid schemaid, PublicationPartOpt pub_partopt); extern List *GetAllSchemaPublicationRelations(Oid pubid, PublicationPartOpt pub_partopt); +extern List *GetAllPublishableTables(void); extern List *GetPubPartitionOptionRelations(List *result, PublicationPartOpt pub_partopt, Oid relid); diff --git a/src/test/isolation/expected/publication-ddl-race.out b/src/test/isolation/expected/publication-ddl-race.out new file mode 100644 index 00000000000..690157f4cdf --- /dev/null +++ b/src/test/isolation/expected/publication-ddl-race.out @@ -0,0 +1,92 @@ +Parsed test spec with 3 sessions + +starting permutation: s1_begin s1_update_nori s2_add_schema s3_update_ri s1_commit s3_update_nori +step s1_begin: BEGIN; +step s1_update_nori: UPDATE sch1.t_nori SET val = 2 WHERE id = 1; +step s2_add_schema: ALTER PUBLICATION pub1 ADD TABLES IN SCHEMA sch1; +step s3_update_ri: UPDATE sch1.t_ri SET val = 2 WHERE id = 1; +step s1_commit: COMMIT; +step s2_add_schema: <... completed> +step s3_update_nori: UPDATE sch1.t_nori SET val = 3 WHERE id = 1; +ERROR: cannot update table "t_nori" because it does not have a replica identity and publishes updates + +starting permutation: s1_begin s1_update_nori s2_add_schema_ins s3_update_ins s1_commit +step s1_begin: BEGIN; +step s1_update_nori: UPDATE sch1.t_nori SET val = 2 WHERE id = 1; +step s2_add_schema_ins: ALTER PUBLICATION pub_ins ADD TABLES IN SCHEMA sch2; +step s3_update_ins: UPDATE sch2.t_ins SET val = 2 WHERE id = 1; +step s1_commit: COMMIT; + +starting permutation: s1_begin s1_update_w s2_set_publish_w s1_commit s3_update_w +step s1_begin: BEGIN; +step s1_update_w: UPDATE sch3.t_w SET val = 2 WHERE id = 1; +step s2_set_publish_w: ALTER PUBLICATION pub_w SET (publish = 'insert, update'); +step s1_commit: COMMIT; +step s2_set_publish_w: <... completed> +step s3_update_w: UPDATE sch3.t_w SET val = 3 WHERE id = 1; +ERROR: cannot update table "t_w" because it does not have a replica identity and publishes updates + +starting permutation: s1_begin s1_update_nori s2_create_all s1_commit s3_update_nori s2_drop_all +step s1_begin: BEGIN; +step s1_update_nori: UPDATE sch1.t_nori SET val = 2 WHERE id = 1; +step s2_create_all: CREATE PUBLICATION pub_all FOR ALL TABLES; +step s1_commit: COMMIT; +step s2_create_all: <... completed> +step s3_update_nori: UPDATE sch1.t_nori SET val = 3 WHERE id = 1; +ERROR: cannot update table "t_nori" because it does not have a replica identity and publishes updates +step s2_drop_all: DROP PUBLICATION pub_all; + +starting permutation: s1_begin s1_update_nori s2_create_schema s1_commit s3_update_nori +step s1_begin: BEGIN; +step s1_update_nori: UPDATE sch1.t_nori SET val = 2 WHERE id = 1; +step s2_create_schema: CREATE PUBLICATION pub_create FOR TABLES IN SCHEMA sch1; +step s1_commit: COMMIT; +step s2_create_schema: <... completed> +step s3_update_nori: UPDATE sch1.t_nori SET val = 3 WHERE id = 1; +ERROR: cannot update table "t_nori" because it does not have a replica identity and publishes updates + +starting permutation: s2_begin s2_add_empty s3_create_in_emp s3_begin s3_update_emp s2_commit s3_commit +step s2_begin: BEGIN; +step s2_add_empty: ALTER PUBLICATION pub_emp ADD TABLES IN SCHEMA sch_emp; +step s3_create_in_emp: CREATE TABLE sch_emp.t_new (id int, val int); +step s3_begin: BEGIN; +step s3_update_emp: UPDATE sch_emp.t_new SET val = 1 WHERE id = 1; +step s2_commit: COMMIT; +step s3_update_emp: <... completed> +ERROR: cannot update table "t_new" because it does not have a replica identity and publishes updates +step s3_commit: COMMIT; + +starting permutation: s1_begin s1_update_nori s2_begin s2_setall s1_commit s3_create_new s3_begin s3_update_new s2_commit s3_commit +step s1_begin: BEGIN; +step s1_update_nori: UPDATE sch1.t_nori SET val = 2 WHERE id = 1; +step s2_begin: BEGIN; +step s2_setall: ALTER PUBLICATION pub_setall SET ALL TABLES; +step s1_commit: COMMIT; +step s2_setall: <... completed> +step s3_create_new: CREATE TABLE t_new (id int, val int); +step s3_begin: BEGIN; +step s3_update_new: UPDATE t_new SET val = 1 WHERE id = 1; +step s2_commit: COMMIT; +step s3_update_new: <... completed> +ERROR: cannot update table "t_new" because it does not have a replica identity and publishes updates +step s3_commit: COMMIT; + +starting permutation: s2_begin s2_add_schema_flip s3_flip s3_begin s3_update_ri2 s2_commit s3_commit +step s2_begin: BEGIN; +step s2_add_schema_flip: ALTER PUBLICATION pub_flip ADD TABLES IN SCHEMA sch1; +step s3_flip: ALTER TABLE sch1.t_ri REPLICA IDENTITY NOTHING; +step s3_begin: BEGIN; +step s3_update_ri2: UPDATE sch1.t_ri SET val = 3 WHERE id = 1; +step s2_commit: COMMIT; +step s3_update_ri2: <... completed> +ERROR: cannot update table "t_ri" because it does not have a replica identity and publishes updates +step s3_commit: COMMIT; + +starting permutation: s1_begin s1_update_attach s2_attach s1_commit s3_update_attach +step s1_begin: BEGIN; +step s1_update_attach: UPDATE t_attach SET val = 2 WHERE id = 1; +step s2_attach: ALTER TABLE schp.part ATTACH PARTITION t_attach FOR VALUES FROM (1000) TO (2000); +step s1_commit: COMMIT; +step s2_attach: <... completed> +step s3_update_attach: UPDATE t_attach SET val = 3 WHERE id = 1; +ERROR: cannot update table "t_attach" because it does not have a replica identity and publishes updates diff --git a/src/test/isolation/isolation_schedule b/src/test/isolation/isolation_schedule index f1676a961f9..50f68de6462 100644 --- a/src/test/isolation/isolation_schedule +++ b/src/test/isolation/isolation_schedule @@ -131,3 +131,4 @@ test: ddl-dependency-locking test: tablespace-dependency-locking test: pub-concurrent-drop test: drop-owned-grant +test: publication-ddl-race diff --git a/src/test/isolation/specs/publication-ddl-race.spec b/src/test/isolation/specs/publication-ddl-race.spec new file mode 100644 index 00000000000..7e3eb75ab5f --- /dev/null +++ b/src/test/isolation/specs/publication-ddl-race.spec @@ -0,0 +1,148 @@ +# Test that publication DDL that widens the set of changes tables publish +# (adding TABLES IN SCHEMA or ALL TABLES, enabling publication of UPDATE or +# DELETE) serializes with concurrent data-modifying statements: the DDL must +# wait for in-progress writers, and writers must not pass the replica identity +# check with the old publication definition while the DDL commits underneath +# them. +# +# The interlocks are verified via the lock waits they produce. The +# correctness of the outcome is verified by the replica identity error that a +# writer on a table without a replica identity gets once the DDL has +# committed. + +setup +{ + CREATE SCHEMA sch1; + CREATE TABLE sch1.t_nori (id int, val int); + INSERT INTO sch1.t_nori VALUES (1, 1); + CREATE TABLE sch1.t_ri (id int PRIMARY KEY, val int); + INSERT INTO sch1.t_ri VALUES (1, 1); + CREATE TABLE t_outside (id int, val int); + INSERT INTO t_outside VALUES (1, 1); + CREATE SCHEMA sch2; + CREATE TABLE sch2.t_ins (id int, val int); + INSERT INTO sch2.t_ins VALUES (1, 1); + CREATE SCHEMA sch3; + CREATE TABLE sch3.t_w (id int, val int); + INSERT INTO sch3.t_w VALUES (1, 1); + CREATE SCHEMA schp; + CREATE SCHEMA schleaf; + CREATE TABLE schp.part (id int, val int) PARTITION BY RANGE (id); + CREATE TABLE schleaf.leaf1 PARTITION OF schp.part + FOR VALUES FROM (0) TO (1000); + INSERT INTO schleaf.leaf1 VALUES (1, 1); + CREATE TABLE t_attach (id int, val int); + INSERT INTO t_attach VALUES (1500, 1); + CREATE PUBLICATION pub1; + CREATE PUBLICATION pub_emp; + CREATE SCHEMA sch_emp; + CREATE PUBLICATION pub_ins WITH (publish = 'insert'); + CREATE PUBLICATION pub_w FOR TABLES IN SCHEMA sch3 WITH (publish = 'insert'); + CREATE PUBLICATION pub_setall; + CREATE PUBLICATION pub_flip; + CREATE PUBLICATION pub_part FOR TABLE schp.part; +} + +teardown +{ + DROP TABLE IF EXISTS t_attach; + DROP SCHEMA sch1 CASCADE; + DROP SCHEMA sch2 CASCADE; + DROP SCHEMA sch3 CASCADE; + DROP SCHEMA schp CASCADE; + DROP SCHEMA schleaf CASCADE; + DROP TABLE t_outside; + DROP PUBLICATION pub1; + DROP PUBLICATION pub_emp; + DROP SCHEMA sch_emp CASCADE; + DROP PUBLICATION pub_ins; + DROP PUBLICATION pub_w; + DROP PUBLICATION pub_setall; + DROP PUBLICATION pub_flip; + DROP PUBLICATION pub_part; +} + +session "s1" + +step "s1_begin" { BEGIN; } +step "s1_update_nori" { UPDATE sch1.t_nori SET val = 2 WHERE id = 1; } +step "s1_update_w" { UPDATE sch3.t_w SET val = 2 WHERE id = 1; } +step "s1_update_attach" { UPDATE t_attach SET val = 2 WHERE id = 1; } +step "s1_commit" { COMMIT; } + +session "s2" + +step "s2_add_schema" { ALTER PUBLICATION pub1 ADD TABLES IN SCHEMA sch1; } +step "s2_create_schema" { CREATE PUBLICATION pub_create FOR TABLES IN SCHEMA sch1; } +step "s2_add_empty" { ALTER PUBLICATION pub_emp ADD TABLES IN SCHEMA sch_emp; } +step "s2_add_schema_ins" { ALTER PUBLICATION pub_ins ADD TABLES IN SCHEMA sch2; } +step "s2_set_publish_w" { ALTER PUBLICATION pub_w SET (publish = 'insert, update'); } +step "s2_create_all" { CREATE PUBLICATION pub_all FOR ALL TABLES; } +step "s2_drop_all" { DROP PUBLICATION pub_all; } +step "s2_begin" { BEGIN; } +step "s2_setall" { ALTER PUBLICATION pub_setall SET ALL TABLES; } +step "s2_commit" { COMMIT; } +step "s2_add_schema_flip" { ALTER PUBLICATION pub_flip ADD TABLES IN SCHEMA sch1; } +step "s2_attach" { ALTER TABLE schp.part ATTACH PARTITION t_attach FOR VALUES FROM (1000) TO (2000); } + +session "s3" + +step "s3_update_ri" { UPDATE sch1.t_ri SET val = 2 WHERE id = 1; } +step "s3_update_nori" { UPDATE sch1.t_nori SET val = 3 WHERE id = 1; } +step "s3_update_ins" { UPDATE sch2.t_ins SET val = 2 WHERE id = 1; } +step "s3_update_w" { UPDATE sch3.t_w SET val = 3 WHERE id = 1; } +step "s3_create_new" { CREATE TABLE t_new (id int, val int); } +step "s3_create_in_emp" { CREATE TABLE sch_emp.t_new (id int, val int); } +step "s3_begin" { BEGIN; } +step "s3_update_new" { UPDATE t_new SET val = 1 WHERE id = 1; } +step "s3_update_emp" { UPDATE sch_emp.t_new SET val = 1 WHERE id = 1; } +step "s3_commit" { COMMIT; } +step "s3_flip" { ALTER TABLE sch1.t_ri REPLICA IDENTITY NOTHING; } +step "s3_update_ri2" { UPDATE sch1.t_ri SET val = 3 WHERE id = 1; } +step "s3_update_attach" { UPDATE t_attach SET val = 3 WHERE id = 1; } + +# ALTER PUBLICATION ... ADD TABLES IN SCHEMA waits for an in-progress writer +# on a table in the schema without a replica identity, but a writer on a +# table with a replica identity is not blocked. Afterwards, writers on the +# table are correctly rejected. +permutation "s1_begin" "s1_update_nori" "s2_add_schema" "s3_update_ri" "s1_commit" "s3_update_nori" + +# An insert-only publication does not need a replica identity on its tables, +# so adding a schema to it does not block writers. +permutation "s1_begin" "s1_update_nori" "s2_add_schema_ins" "s3_update_ins" "s1_commit" + +# ALTER PUBLICATION ... SET (publish = ...) enabling UPDATE on a TABLES IN +# SCHEMA publication waits for an in-progress writer on a schema member +# without a replica identity. +permutation "s1_begin" "s1_update_w" "s2_set_publish_w" "s1_commit" "s3_update_w" + +# CREATE PUBLICATION ... FOR ALL TABLES waits for an in-progress writer on a +# table without a replica identity. +permutation "s1_begin" "s1_update_nori" "s2_create_all" "s1_commit" "s3_update_nori" "s2_drop_all" + +# CREATE PUBLICATION ... FOR TABLES IN SCHEMA also waits for an in-progress +# writer on a table in the schema without a replica identity. +permutation "s1_begin" "s1_update_nori" "s2_create_schema" "s1_commit" "s3_update_nori" + +# Adding an empty schema still takes the publication catalog lock, so a table +# created in it while the DDL is in progress is covered: a writer on it waits +# for the DDL's catalog lock, then is correctly rejected. +permutation "s2_begin" "s2_add_empty" "s3_create_in_emp" "s3_begin" "s3_update_emp" "s2_commit" "s3_commit" + +# ALTER PUBLICATION ... SET ALL TABLES waits for an in-progress writer; a +# table created after it surveyed the existing tables is covered by the +# writer building its publication descriptor taking a conflicting lock on the +# publication catalog, which the DDL holds until commit. +permutation "s1_begin" "s1_update_nori" "s2_begin" "s2_setall" "s1_commit" "s3_create_new" "s3_begin" "s3_update_new" "s2_commit" "s3_commit" + +# A table whose replica identity is removed after the DDL surveyed it: a +# writer on it builds its publication descriptor, and waits for the DDL's +# lock on the publication catalog. Once the DDL has committed, the writer is +# correctly rejected. +permutation "s2_begin" "s2_add_schema_flip" "s3_flip" "s3_begin" "s3_update_ri2" "s2_commit" "s3_commit" + +# ATTACH PARTITION takes AccessExclusiveLock on the new partition, so it waits +# for an in-progress writer on it; afterwards the partition is covered by the +# parent's publication and writers are correctly rejected. +permutation "s1_begin" "s1_update_attach" "s2_attach" "s1_commit" "s3_update_attach" + -- 2.34.1