From 453b0d90fc2f20740e583062096fd34e5c0746df Mon Sep 17 00:00:00 2001 From: Khanna Date: Fri, 22 Nov 2024 12:03:13 +0530 Subject: [PATCH v3] Add support for two-phase commit in pg_createsubscriber This patch introduces the '--enable-two-phase' option to the 'pg_createsubscriber' utility, allowing users to enable two-phase commit for subscriptions during their creation. By default, two-phase commit is disabled if the option is not provided. When two-phase commit is enabled, prepared transactions are sent to the subscriber at the time of 'PREPARE TRANSACTION', and they are processed as two-phase transactions on the subscriber as well. If disabled, prepared transactions are sent only when committed and are processed immediately by the subscriber. Documentation has been updated to reflect the new option, and test cases have been added to validate various scenarios, including proper validation of the '--enable-two-phase' option and its combinations with other options. --- doc/src/sgml/ref/pg_createsubscriber.sgml | 11 +++++++++ src/bin/pg_basebackup/pg_createsubscriber.c | 21 ++++++++++++---- .../t/040_pg_createsubscriber.pl | 24 ++++++++++++++++--- 3 files changed, 48 insertions(+), 8 deletions(-) diff --git a/doc/src/sgml/ref/pg_createsubscriber.sgml b/doc/src/sgml/ref/pg_createsubscriber.sgml index df1a92b4da..d2bbef3567 100644 --- a/doc/src/sgml/ref/pg_createsubscriber.sgml +++ b/doc/src/sgml/ref/pg_createsubscriber.sgml @@ -161,6 +161,17 @@ PostgreSQL documentation + + + + + + Enables two_phase + commit for the subscription. The default is false. + + + + diff --git a/src/bin/pg_basebackup/pg_createsubscriber.c b/src/bin/pg_basebackup/pg_createsubscriber.c index e96370a9ec..01a203983d 100644 --- a/src/bin/pg_basebackup/pg_createsubscriber.c +++ b/src/bin/pg_basebackup/pg_createsubscriber.c @@ -38,6 +38,7 @@ struct CreateSubscriberOptions char *socket_dir; /* directory for Unix-domain socket, if any */ char *sub_port; /* subscriber port number */ const char *sub_username; /* subscriber username */ + bool two_phase; /* two-phase option */ SimpleStringList database_names; /* list of database names */ SimpleStringList pub_names; /* list of publication names */ SimpleStringList sub_names; /* list of subscription names */ @@ -53,6 +54,7 @@ struct LogicalRepInfo char *pubname; /* publication name */ char *subname; /* subscription name */ char *replslotname; /* replication slot name */ + bool two_phase; /* two-phase enabled for the subscription */ bool made_replslot; /* replication slot was created */ bool made_publication; /* publication was created */ @@ -227,6 +229,7 @@ usage(void) printf(_(" -P, --publisher-server=CONNSTR publisher connection string\n")); printf(_(" -s, --socketdir=DIR socket directory to use (default current dir.)\n")); printf(_(" -t, --recovery-timeout=SECS seconds to wait for recovery to end\n")); + printf(_(" -T, --enable-two-phase enable two-phase commit for the subscription\n")); printf(_(" -U, --subscriber-username=NAME user name for subscriber connection\n")); printf(_(" -v, --verbose output verbose messages\n")); printf(_(" --config-file=FILENAME use specified main server configuration\n" @@ -456,6 +459,7 @@ store_pub_sub_info(const struct CreateSubscriberOptions *opt, conninfo = concat_conninfo_dbname(pub_base_conninfo, cell->val); dbinfo[i].pubconninfo = conninfo; dbinfo[i].dbname = cell->val; + dbinfo[i].two_phase = opt->two_phase; if (num_pubs > 0) dbinfo[i].pubname = pubcell->val; else @@ -479,9 +483,10 @@ store_pub_sub_info(const struct CreateSubscriberOptions *opt, dbinfo[i].pubname ? dbinfo[i].pubname : "(auto)", dbinfo[i].replslotname ? dbinfo[i].replslotname : "(auto)", dbinfo[i].pubconninfo); - pg_log_debug("subscriber(%d): subscription: %s ; connection string: %s", i, + pg_log_debug("subscriber(%d): subscription: %s ; connection string: %s, two_phase: %s", i, dbinfo[i].subname ? dbinfo[i].subname : "(auto)", - dbinfo[i].subconninfo); + dbinfo[i].subconninfo, + dbinfo[i].two_phase ? "true" : "false"); if (num_pubs > 0) pubcell = pubcell->next; @@ -1699,8 +1704,9 @@ create_subscription(PGconn *conn, const struct LogicalRepInfo *dbinfo) appendPQExpBuffer(str, "CREATE SUBSCRIPTION %s CONNECTION %s PUBLICATION %s " "WITH (create_slot = false, enabled = false, " - "slot_name = %s, copy_data = false)", - subname_esc, pubconninfo_esc, pubname_esc, replslotname_esc); + "slot_name = %s, copy_data = false, two_phase = %s)", + subname_esc, pubconninfo_esc, pubname_esc, replslotname_esc, + dbinfo->two_phase ? "true" : "false"); pg_free(pubname_esc); pg_free(subname_esc); @@ -1872,6 +1878,7 @@ main(int argc, char **argv) {"publisher-server", required_argument, NULL, 'P'}, {"socketdir", required_argument, NULL, 's'}, {"recovery-timeout", required_argument, NULL, 't'}, + {"enable-two-phase", no_argument, NULL, 'T'}, {"subscriber-username", required_argument, NULL, 'U'}, {"verbose", no_argument, NULL, 'v'}, {"version", no_argument, NULL, 'V'}, @@ -1927,6 +1934,7 @@ main(int argc, char **argv) opt.socket_dir = NULL; opt.sub_port = DEFAULT_SUB_PORT; opt.sub_username = NULL; + opt.two_phase = false; opt.database_names = (SimpleStringList) { 0 @@ -1949,7 +1957,7 @@ main(int argc, char **argv) get_restricted_token(); - while ((c = getopt_long(argc, argv, "d:D:np:P:s:t:U:v", + while ((c = getopt_long(argc, argv, "d:D:np:P:s:t:T:U:v", long_options, &option_index)) != -1) { switch (c) @@ -1986,6 +1994,9 @@ main(int argc, char **argv) case 't': opt.recovery_timeout = atoi(optarg); break; + case 'T': + opt.two_phase = true; + break; case 'U': opt.sub_username = pg_strdup(optarg); break; diff --git a/src/bin/pg_basebackup/t/040_pg_createsubscriber.pl b/src/bin/pg_basebackup/t/040_pg_createsubscriber.pl index 0a900edb65..4db196052d 100644 --- a/src/bin/pg_basebackup/t/040_pg_createsubscriber.pl +++ b/src/bin/pg_basebackup/t/040_pg_createsubscriber.pl @@ -249,14 +249,16 @@ command_fails( $db2 ], 'primary contains unmet conditions on node P'); -# Restore default settings here but only apply it after testing standby. Some -# standby settings should not be a lower setting than on the primary. +# Restore default settings here (except for max_prepared_transactions as this is +# required for --enable-two-phase) but only apply it after testing standby. +# Some standby settings should not be a lower setting than on the primary. $node_p->append_conf( 'postgresql.conf', q{ wal_level = logical max_replication_slots = 10 max_wal_senders = 10 max_worker_processes = 8 +max_prepared_transactions = 10 }); # Check some unmet conditions on node S @@ -283,6 +285,7 @@ $node_s->append_conf( max_replication_slots = 10 max_logical_replication_workers = 4 max_worker_processes = 8 +max_prepared_transactions = 10 }); # Restore default settings on both servers $node_p->restart; @@ -357,6 +360,7 @@ command_ok( 'run pg_createsubscriber without --databases'); # Run pg_createsubscriber on node S +# In passing, also test the --enable-two-phase option command_ok( [ 'pg_createsubscriber', '--verbose', @@ -371,7 +375,7 @@ command_ok( 'replslot1', '--replication-slot', 'replslot2', '--database', $db1, '--database', - $db2 + $db2, '--enable-two-phase' ], 'run pg_createsubscriber on node S'); @@ -387,9 +391,23 @@ is($result, qq(0), $node_p->safe_psql($db1, "INSERT INTO tbl1 VALUES('third row')"); $node_p->safe_psql($db2, "INSERT INTO tbl2 VALUES('row 1')"); +# Prepare a transaction on the publisher +$node_p->safe_psql( + $db1, qq[ + BEGIN; + INSERT INTO tbl1 SELECT generate_series(1, 10); + PREPARE TRANSACTION 'test_prepare'; +]); + # Start subscriber $node_s->start; +# Verify that the prepared transaction is replicated to the subscriber +my $count_prepared_s = + $node_s->safe_psql($db1, "SELECT count(*) FROM pg_prepared_xacts;"); + +is($count_prepared_s, qq(1), 'Prepared transaction replicated to subscriber'); + # Confirm the pre-existing subscription has been removed $result = $node_s->safe_psql( 'postgres', qq( -- 2.41.0.windows.3