diff --git a/ci_scripts/perl/PostgreSQL/Test/TdeCluster.pm b/ci_scripts/perl/PostgreSQL/Test/TdeCluster.pm index 6ac0c09bc..5d0bdedb2 100644 --- a/ci_scripts/perl/PostgreSQL/Test/TdeCluster.pm +++ b/ci_scripts/perl/PostgreSQL/Test/TdeCluster.pm @@ -57,6 +57,8 @@ my %wal_skip = ( 'pg_basebackup without -E from server with encrypted WAL produces broken backups', 'src/bin/pg_combinebackup/t/006_db_file_copy.pl' => 'pg_basebackup without -E from server with encrypted WAL produces broken backups', + 'src/bin/pg_combinebackup/t/012_vm_consistency.pl' => + 'pg_basebackup without -E from server with encrypted WAL produces broken backups', 'src/bin/pg_rewind/t/001_basic.pl' => 'copies WAL directly to archive without using archive_command', 'src/bin/pg_verifybackup/t/009_extract.pl' => diff --git a/documentation/docs/release-notes/release-notes-v2.2.2.md b/documentation/docs/release-notes/release-notes-v2.2.2.md index 721b0b6e2..b8abf83a5 100644 --- a/documentation/docs/release-notes/release-notes-v2.2.2.md +++ b/documentation/docs/release-notes/release-notes-v2.2.2.md @@ -40,3 +40,4 @@ Changes introduced in `pg_tde` 2.2.2: - [PG-2492](https://perconadev.atlassian.net/browse/PG-2492) - Fixed crash when empty certificate parameters are passed to `pg_tde_add_global_key_provider_kmip()` or `pg_tde_add_database_key_provider_kmip()` - [PG-2608](https://perconadev.atlassian.net/browse/PG-2608) - Fixed a race condition in the Vault key provider that could occur when multiple processes accessed the same cURL handle after a fork. - [PG-2609](https://perconadev.atlassian.net/browse/PG-2609) - Fixed WAL archiving (`pg_tde_archive_decrypt` and `pg_tde_restore_encrypt`) when pg_wal dir is a symlink. +- Include bugfixes for frontend tools from PostgreSQL 16.15, 17.11 and 18.5 diff --git a/fetools/pg16/pg_basebackup/pg_recvlogical.c b/fetools/pg16/pg_basebackup/pg_recvlogical.c index f3c7937a1..6bf7021a4 100644 --- a/fetools/pg16/pg_basebackup/pg_recvlogical.c +++ b/fetools/pg16/pg_basebackup/pg_recvlogical.c @@ -230,8 +230,9 @@ StreamLogicalLog(void) /* Initiate the replication stream at specified location */ query = createPQExpBuffer(); - appendPQExpBuffer(query, "START_REPLICATION SLOT \"%s\" LOGICAL %X/%X", - replication_slot, LSN_FORMAT_ARGS(startpos)); + appendPQExpBufferStr(query, "START_REPLICATION SLOT "); + AppendQuotedIdentifier(query, replication_slot); + appendPQExpBuffer(query, " LOGICAL %X/%X", LSN_FORMAT_ARGS(startpos)); /* print options if there are any */ if (noptions) @@ -244,11 +245,14 @@ StreamLogicalLog(void) appendPQExpBufferStr(query, ", "); /* write option name */ - appendPQExpBuffer(query, "\"%s\"", options[(i * 2)]); + AppendQuotedIdentifier(query, options[i * 2]); /* write option value if specified */ - if (options[(i * 2) + 1] != NULL) - appendPQExpBuffer(query, " '%s'", options[(i * 2) + 1]); + if (options[i * 2 + 1] != NULL) + { + appendPQExpBufferChar(query, ' '); + AppendQuotedLiteral(query, options[i * 2 + 1]); + } } if (noptions) @@ -327,7 +331,7 @@ StreamLogicalLog(void) outfd = fileno(stdout); else outfd = open(outfile, O_CREAT | O_APPEND | O_WRONLY | PG_BINARY, - S_IRUSR | S_IWUSR); + pg_file_create_mode); if (outfd == -1) { pg_log_error("could not open log file \"%s\": %m", outfile); diff --git a/fetools/pg16/pg_basebackup/receivelog.c b/fetools/pg16/pg_basebackup/receivelog.c index 01f15e347..f98cb6801 100644 --- a/fetools/pg16/pg_basebackup/receivelog.c +++ b/fetools/pg16/pg_basebackup/receivelog.c @@ -457,8 +457,7 @@ CheckServerVersionForStreaming(PGconn *conn) bool ReceiveXlogStream(PGconn *conn, StreamCtl *stream) { - char query[128]; - char slotcmd[128]; + PQExpBuffer query; PGresult *res; XLogRecPtr stoppos; @@ -483,7 +482,6 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) if (stream->replication_slot != NULL) { reportFlushPosition = true; - sprintf(slotcmd, "SLOT \"%s\" ", stream->replication_slot); } else { @@ -491,7 +489,6 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) reportFlushPosition = true; else reportFlushPosition = false; - slotcmd[0] = 0; } if (stream->sysidentifier != NULL) @@ -540,8 +537,10 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) */ if (!existsTimeLineHistoryFile(stream)) { - snprintf(query, sizeof(query), "TIMELINE_HISTORY %u", stream->timeline); - res = PQexec(conn, query); + query = createPQExpBuffer(); + appendPQExpBuffer(query, "TIMELINE_HISTORY %u", stream->timeline); + res = PQexec(conn, query->data); + destroyPQExpBuffer(query); if (PQresultStatus(res) != PGRES_TUPLES_OK) { /* FIXME: we might send it ok, but get an error */ @@ -577,11 +576,18 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) return true; /* Initiate the replication stream at specified location */ - snprintf(query, sizeof(query), "START_REPLICATION %s%X/%X TIMELINE %u", - slotcmd, - LSN_FORMAT_ARGS(stream->startpos), - stream->timeline); - res = PQexec(conn, query); + query = createPQExpBuffer(); + appendPQExpBufferStr(query, "START_REPLICATION"); + if (stream->replication_slot != NULL) + { + appendPQExpBufferStr(query, " SLOT "); + AppendQuotedIdentifier(query, stream->replication_slot); + } + appendPQExpBuffer(query, " %X/%X TIMELINE %u", + LSN_FORMAT_ARGS(stream->startpos), + stream->timeline); + res = PQexec(conn, query->data); + destroyPQExpBuffer(query); if (PQresultStatus(res) != PGRES_COPY_BOTH) { pg_log_error("could not send replication command \"%s\": %s", diff --git a/fetools/pg16/pg_basebackup/streamutil.c b/fetools/pg16/pg_basebackup/streamutil.c index 15514599c..c8fa57fa4 100644 --- a/fetools/pg16/pg_basebackup/streamutil.c +++ b/fetools/pg16/pg_basebackup/streamutil.c @@ -492,7 +492,8 @@ GetSlotInformation(PGconn *conn, const char *slot_name, *restart_tli = tli_loc; query = createPQExpBuffer(); - appendPQExpBuffer(query, "READ_REPLICATION_SLOT %s", slot_name); + appendPQExpBufferStr(query, "READ_REPLICATION_SLOT "); + AppendQuotedIdentifier(query, slot_name); res = PQexec(conn, query->data); destroyPQExpBuffer(query); @@ -588,13 +589,17 @@ CreateReplicationSlot(PGconn *conn, const char *slot_name, const char *plugin, Assert(slot_name != NULL); /* Build base portion of query */ - appendPQExpBuffer(query, "CREATE_REPLICATION_SLOT \"%s\"", slot_name); + appendPQExpBufferStr(query, "CREATE_REPLICATION_SLOT "); + AppendQuotedIdentifier(query, slot_name); if (is_temporary) appendPQExpBufferStr(query, " TEMPORARY"); if (is_physical) appendPQExpBufferStr(query, " PHYSICAL"); else - appendPQExpBuffer(query, " LOGICAL \"%s\"", plugin); + { + appendPQExpBufferStr(query, " LOGICAL "); + AppendQuotedIdentifier(query, plugin); + } /* Add any requested options */ if (use_new_option_syntax) @@ -690,8 +695,8 @@ DropReplicationSlot(PGconn *conn, const char *slot_name) query = createPQExpBuffer(); /* Build query */ - appendPQExpBuffer(query, "DROP_REPLICATION_SLOT \"%s\"", - slot_name); + appendPQExpBufferStr(query, "DROP_REPLICATION_SLOT "); + AppendQuotedIdentifier(query, slot_name); res = PQexec(conn, query->data); if (PQresultStatus(res) != PGRES_COMMAND_OK) { @@ -719,6 +724,29 @@ DropReplicationSlot(PGconn *conn, const char *slot_name) return true; } +/* + * Append a suitably-quoted identifier or string literal to buf. + * "quote" should be either a double-quote or single-quote character. + * + * Caution: this quoting logic is sufficient for identifiers and literals + * in the replication grammar, but not always in regular SQL. Specifically, + * it'd fail for a string literal if standard_conforming_strings is off. + */ +void +AppendQuotedString(PQExpBuffer buf, const char *str, char quote) +{ + appendPQExpBufferChar(buf, quote); + while (*str) + { + char c = *str++; + + if (c == quote) + appendPQExpBufferChar(buf, c); + appendPQExpBufferChar(buf, c); + } + appendPQExpBufferChar(buf, quote); +} + /* * Append a "plain" option - one with no value - to a server command that * is being constructed. @@ -727,10 +755,13 @@ DropReplicationSlot(PGconn *conn, const char *slot_name) * write things like SOME_COMMAND OPTION1 OPTION2 'opt2value' OPTION3 42. The * new syntax uses a comma-separated list surrounded by parentheses, so the * equivalent is SOME_COMMAND (OPTION1, OPTION2 'optvalue', OPTION3 42). + * + * Note: we assume option names do not require quotes. Do not use this + * with option names coming from outside sources. */ void AppendPlainCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name) + const char *option_name) { if (buf->len > 0 && buf->data[buf->len - 1] != '(') { @@ -751,30 +782,26 @@ AppendPlainCommandOption(PQExpBuffer buf, bool use_new_option_syntax, */ void AppendStringCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, char *option_value) + const char *option_name, const char *option_value) { AppendPlainCommandOption(buf, use_new_option_syntax, option_name); if (option_value != NULL) { - size_t length = strlen(option_value); - char *escaped_value = palloc(1 + 2 * length); - - PQescapeStringConn(conn, escaped_value, option_value, length, NULL); - appendPQExpBuffer(buf, " '%s'", escaped_value); - pfree(escaped_value); + appendPQExpBufferChar(buf, ' '); + AppendQuotedLiteral(buf, option_value); } } /* - * Append an option with an associated integer value to a server command + * Append an option with an associated integer value to a server command that * is being constructed. * * See comments for AppendPlainCommandOption, above. */ void AppendIntegerCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, int32 option_value) + const char *option_name, int32 option_value) { AppendPlainCommandOption(buf, use_new_option_syntax, option_name); diff --git a/fetools/pg16/pg_basebackup/streamutil.h b/fetools/pg16/pg_basebackup/streamutil.h index 268c16321..973d0775b 100644 --- a/fetools/pg16/pg_basebackup/streamutil.h +++ b/fetools/pg16/pg_basebackup/streamutil.h @@ -42,15 +42,20 @@ extern bool RunIdentifySystem(PGconn *conn, char **sysid, XLogRecPtr *startpos, char **db_name); +extern void AppendQuotedString(PQExpBuffer buf, const char *str, char quote); +#define AppendQuotedIdentifier(b, s) AppendQuotedString(b, s, '"') +#define AppendQuotedLiteral(b, s) AppendQuotedString(b, s, '\'') extern void AppendPlainCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name); + const char *option_name); extern void AppendStringCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, char *option_value); + const char *option_name, + const char *option_value); extern void AppendIntegerCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, int32 option_value); + const char *option_name, + int32 option_value); extern bool GetSlotInformation(PGconn *conn, const char *slot_name, XLogRecPtr *restart_lsn, diff --git a/fetools/pg17/pg_basebackup/pg_createsubscriber.c b/fetools/pg17/pg_basebackup/pg_createsubscriber.c index eb90c23ca..b8f3b7111 100644 --- a/fetools/pg17/pg_basebackup/pg_createsubscriber.c +++ b/fetools/pg17/pg_basebackup/pg_createsubscriber.c @@ -1402,7 +1402,6 @@ drop_replication_slot(PGconn *conn, struct LogicalRepInfo *dbinfo, { pg_log_error("could not drop replication slot \"%s\" in database \"%s\": %s", slot_name, dbinfo->dbname, PQresultErrorMessage(res)); - dbinfo->made_replslot = false; /* don't try again. */ } PQclear(res); @@ -1665,7 +1664,6 @@ drop_publication(PGconn *conn, struct LogicalRepInfo *dbinfo) { pg_log_error("could not drop publication \"%s\" in database \"%s\": %s", dbinfo->pubname, dbinfo->dbname, PQresultErrorMessage(res)); - dbinfo->made_publication = false; /* don't try again. */ /* * Don't disconnect and exit here. This routine is used by primary diff --git a/fetools/pg17/pg_basebackup/pg_recvlogical.c b/fetools/pg17/pg_basebackup/pg_recvlogical.c index 3db520ed3..949cd6783 100644 --- a/fetools/pg17/pg_basebackup/pg_recvlogical.c +++ b/fetools/pg17/pg_basebackup/pg_recvlogical.c @@ -242,8 +242,9 @@ StreamLogicalLog(void) /* Initiate the replication stream at specified location */ query = createPQExpBuffer(); - appendPQExpBuffer(query, "START_REPLICATION SLOT \"%s\" LOGICAL %X/%X", - replication_slot, LSN_FORMAT_ARGS(startpos)); + appendPQExpBufferStr(query, "START_REPLICATION SLOT "); + AppendQuotedIdentifier(query, replication_slot); + appendPQExpBuffer(query, " LOGICAL %X/%X", LSN_FORMAT_ARGS(startpos)); /* print options if there are any */ if (noptions) @@ -256,11 +257,14 @@ StreamLogicalLog(void) appendPQExpBufferStr(query, ", "); /* write option name */ - appendPQExpBuffer(query, "\"%s\"", options[(i * 2)]); + AppendQuotedIdentifier(query, options[i * 2]); /* write option value if specified */ - if (options[(i * 2) + 1] != NULL) - appendPQExpBuffer(query, " '%s'", options[(i * 2) + 1]); + if (options[i * 2 + 1] != NULL) + { + appendPQExpBufferChar(query, ' '); + AppendQuotedLiteral(query, options[i * 2 + 1]); + } } if (noptions) @@ -340,7 +344,7 @@ StreamLogicalLog(void) outfd = fileno(stdout); else outfd = open(outfile, O_CREAT | O_APPEND | O_WRONLY | PG_BINARY, - S_IRUSR | S_IWUSR); + pg_file_create_mode); if (outfd == -1) { pg_log_error("could not open log file \"%s\": %m", outfile); diff --git a/fetools/pg17/pg_basebackup/receivelog.c b/fetools/pg17/pg_basebackup/receivelog.c index 3c6f9edf2..774126319 100644 --- a/fetools/pg17/pg_basebackup/receivelog.c +++ b/fetools/pg17/pg_basebackup/receivelog.c @@ -457,8 +457,7 @@ CheckServerVersionForStreaming(PGconn *conn) bool ReceiveXlogStream(PGconn *conn, StreamCtl *stream) { - char query[128]; - char slotcmd[128]; + PQExpBuffer query; PGresult *res; XLogRecPtr stoppos; @@ -483,7 +482,6 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) if (stream->replication_slot != NULL) { reportFlushPosition = true; - sprintf(slotcmd, "SLOT \"%s\" ", stream->replication_slot); } else { @@ -491,7 +489,6 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) reportFlushPosition = true; else reportFlushPosition = false; - slotcmd[0] = 0; } if (stream->sysidentifier != NULL) @@ -540,8 +537,10 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) */ if (!existsTimeLineHistoryFile(stream)) { - snprintf(query, sizeof(query), "TIMELINE_HISTORY %u", stream->timeline); - res = PQexec(conn, query); + query = createPQExpBuffer(); + appendPQExpBuffer(query, "TIMELINE_HISTORY %u", stream->timeline); + res = PQexec(conn, query->data); + destroyPQExpBuffer(query); if (PQresultStatus(res) != PGRES_TUPLES_OK) { /* FIXME: we might send it ok, but get an error */ @@ -577,11 +576,18 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) return true; /* Initiate the replication stream at specified location */ - snprintf(query, sizeof(query), "START_REPLICATION %s%X/%X TIMELINE %u", - slotcmd, - LSN_FORMAT_ARGS(stream->startpos), - stream->timeline); - res = PQexec(conn, query); + query = createPQExpBuffer(); + appendPQExpBufferStr(query, "START_REPLICATION"); + if (stream->replication_slot != NULL) + { + appendPQExpBufferStr(query, " SLOT "); + AppendQuotedIdentifier(query, stream->replication_slot); + } + appendPQExpBuffer(query, " %X/%X TIMELINE %u", + LSN_FORMAT_ARGS(stream->startpos), + stream->timeline); + res = PQexec(conn, query->data); + destroyPQExpBuffer(query); if (PQresultStatus(res) != PGRES_COPY_BOTH) { pg_log_error("could not send replication command \"%s\": %s", diff --git a/fetools/pg17/pg_basebackup/streamutil.c b/fetools/pg17/pg_basebackup/streamutil.c index dc604b153..d316f8427 100644 --- a/fetools/pg17/pg_basebackup/streamutil.c +++ b/fetools/pg17/pg_basebackup/streamutil.c @@ -572,7 +572,8 @@ GetSlotInformation(PGconn *conn, const char *slot_name, *restart_tli = tli_loc; query = createPQExpBuffer(); - appendPQExpBuffer(query, "READ_REPLICATION_SLOT %s", slot_name); + appendPQExpBufferStr(query, "READ_REPLICATION_SLOT "); + AppendQuotedIdentifier(query, slot_name); res = PQexec(conn, query->data); destroyPQExpBuffer(query); @@ -668,13 +669,17 @@ CreateReplicationSlot(PGconn *conn, const char *slot_name, const char *plugin, Assert(slot_name != NULL); /* Build base portion of query */ - appendPQExpBuffer(query, "CREATE_REPLICATION_SLOT \"%s\"", slot_name); + appendPQExpBufferStr(query, "CREATE_REPLICATION_SLOT "); + AppendQuotedIdentifier(query, slot_name); if (is_temporary) appendPQExpBufferStr(query, " TEMPORARY"); if (is_physical) appendPQExpBufferStr(query, " PHYSICAL"); else - appendPQExpBuffer(query, " LOGICAL \"%s\"", plugin); + { + appendPQExpBufferStr(query, " LOGICAL "); + AppendQuotedIdentifier(query, plugin); + } /* Add any requested options */ if (use_new_option_syntax) @@ -770,8 +775,8 @@ DropReplicationSlot(PGconn *conn, const char *slot_name) query = createPQExpBuffer(); /* Build query */ - appendPQExpBuffer(query, "DROP_REPLICATION_SLOT \"%s\"", - slot_name); + appendPQExpBufferStr(query, "DROP_REPLICATION_SLOT "); + AppendQuotedIdentifier(query, slot_name); res = PQexec(conn, query->data); if (PQresultStatus(res) != PGRES_COMMAND_OK) { @@ -799,6 +804,29 @@ DropReplicationSlot(PGconn *conn, const char *slot_name) return true; } +/* + * Append a suitably-quoted identifier or string literal to buf. + * "quote" should be either a double-quote or single-quote character. + * + * Caution: this quoting logic is sufficient for identifiers and literals + * in the replication grammar, but not always in regular SQL. Specifically, + * it'd fail for a string literal if standard_conforming_strings is off. + */ +void +AppendQuotedString(PQExpBuffer buf, const char *str, char quote) +{ + appendPQExpBufferChar(buf, quote); + while (*str) + { + char c = *str++; + + if (c == quote) + appendPQExpBufferChar(buf, c); + appendPQExpBufferChar(buf, c); + } + appendPQExpBufferChar(buf, quote); +} + /* * Append a "plain" option - one with no value - to a server command that * is being constructed. @@ -807,10 +835,13 @@ DropReplicationSlot(PGconn *conn, const char *slot_name) * write things like SOME_COMMAND OPTION1 OPTION2 'opt2value' OPTION3 42. The * new syntax uses a comma-separated list surrounded by parentheses, so the * equivalent is SOME_COMMAND (OPTION1, OPTION2 'optvalue', OPTION3 42). + * + * Note: we assume option names do not require quotes. Do not use this + * with option names coming from outside sources. */ void AppendPlainCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name) + const char *option_name) { if (buf->len > 0 && buf->data[buf->len - 1] != '(') { @@ -831,30 +862,26 @@ AppendPlainCommandOption(PQExpBuffer buf, bool use_new_option_syntax, */ void AppendStringCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, char *option_value) + const char *option_name, const char *option_value) { AppendPlainCommandOption(buf, use_new_option_syntax, option_name); if (option_value != NULL) { - size_t length = strlen(option_value); - char *escaped_value = palloc(1 + 2 * length); - - PQescapeStringConn(conn, escaped_value, option_value, length, NULL); - appendPQExpBuffer(buf, " '%s'", escaped_value); - pfree(escaped_value); + appendPQExpBufferChar(buf, ' '); + AppendQuotedLiteral(buf, option_value); } } /* - * Append an option with an associated integer value to a server command + * Append an option with an associated integer value to a server command that * is being constructed. * * See comments for AppendPlainCommandOption, above. */ void AppendIntegerCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, int32 option_value) + const char *option_name, int32 option_value) { AppendPlainCommandOption(buf, use_new_option_syntax, option_name); diff --git a/fetools/pg17/pg_basebackup/streamutil.h b/fetools/pg17/pg_basebackup/streamutil.h index 9b38e8c0f..f9bfe881c 100644 --- a/fetools/pg17/pg_basebackup/streamutil.h +++ b/fetools/pg17/pg_basebackup/streamutil.h @@ -44,15 +44,20 @@ extern bool RunIdentifySystem(PGconn *conn, char **sysid, XLogRecPtr *startpos, char **db_name); +extern void AppendQuotedString(PQExpBuffer buf, const char *str, char quote); +#define AppendQuotedIdentifier(b, s) AppendQuotedString(b, s, '"') +#define AppendQuotedLiteral(b, s) AppendQuotedString(b, s, '\'') extern void AppendPlainCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name); + const char *option_name); extern void AppendStringCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, char *option_value); + const char *option_name, + const char *option_value); extern void AppendIntegerCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, int32 option_value); + const char *option_name, + int32 option_value); extern bool GetSlotInformation(PGconn *conn, const char *slot_name, XLogRecPtr *restart_lsn, diff --git a/fetools/pg18/pg_basebackup/pg_createsubscriber.c b/fetools/pg18/pg_basebackup/pg_createsubscriber.c index 51d12aa7f..42a073e92 100644 --- a/fetools/pg18/pg_basebackup/pg_createsubscriber.c +++ b/fetools/pg18/pg_basebackup/pg_createsubscriber.c @@ -115,7 +115,7 @@ static void wait_for_end_recovery(const char *conninfo, const struct CreateSubscriberOptions *opt); static void create_publication(PGconn *conn, struct LogicalRepInfo *dbinfo); static void drop_publication(PGconn *conn, const char *pubname, - const char *dbname, bool *made_publication); + const char *dbname); static void check_and_drop_publications(PGconn *conn, struct LogicalRepInfo *dbinfo); static void create_subscription(PGconn *conn, const struct LogicalRepInfo *dbinfo); static void set_replication_progress(PGconn *conn, const struct LogicalRepInfo *dbinfo, @@ -203,8 +203,7 @@ cleanup_objects_atexit(void) if (conn != NULL) { if (dbinfo->made_publication) - drop_publication(conn, dbinfo->pubname, dbinfo->dbname, - &dbinfo->made_publication); + drop_publication(conn, dbinfo->pubname, dbinfo->dbname); if (dbinfo->made_replslot) drop_replication_slot(conn, dbinfo, dbinfo->replslotname); disconnect_database(conn, false); @@ -1465,7 +1464,6 @@ drop_replication_slot(PGconn *conn, struct LogicalRepInfo *dbinfo, { pg_log_error("could not drop replication slot \"%s\" in database \"%s\": %s", slot_name, dbinfo->dbname, PQresultErrorMessage(res)); - dbinfo->made_replslot = false; /* don't try again. */ } PQclear(res); @@ -1705,8 +1703,7 @@ create_publication(PGconn *conn, struct LogicalRepInfo *dbinfo) * Drop the specified publication in the given database. */ static void -drop_publication(PGconn *conn, const char *pubname, const char *dbname, - bool *made_publication) +drop_publication(PGconn *conn, const char *pubname, const char *dbname) { PQExpBuffer str = createPQExpBuffer(); PGresult *res; @@ -1736,7 +1733,6 @@ drop_publication(PGconn *conn, const char *pubname, const char *dbname, { pg_log_error("could not drop publication \"%s\" in database \"%s\": %s", pubname, dbname, PQresultErrorMessage(res)); - *made_publication = false; /* don't try again. */ /* * Don't disconnect and exit here. This routine is used by primary @@ -1786,8 +1782,7 @@ check_and_drop_publications(PGconn *conn, struct LogicalRepInfo *dbinfo) /* Drop each publication */ for (int i = 0; i < PQntuples(res); i++) - drop_publication(conn, PQgetvalue(res, i, 0), dbinfo->dbname, - &dbinfo->made_publication); + drop_publication(conn, PQgetvalue(res, i, 0), dbinfo->dbname); PQclear(res); } @@ -1797,8 +1792,7 @@ check_and_drop_publications(PGconn *conn, struct LogicalRepInfo *dbinfo) * those to provide necessary information to the user. */ if (!drop_all_pubs || dry_run) - drop_publication(conn, dbinfo->pubname, dbinfo->dbname, - &dbinfo->made_publication); + drop_publication(conn, dbinfo->pubname, dbinfo->dbname); } /* diff --git a/fetools/pg18/pg_basebackup/pg_recvlogical.c b/fetools/pg18/pg_basebackup/pg_recvlogical.c index fb7a6a1d0..113a43b81 100644 --- a/fetools/pg18/pg_basebackup/pg_recvlogical.c +++ b/fetools/pg18/pg_basebackup/pg_recvlogical.c @@ -244,8 +244,9 @@ StreamLogicalLog(void) /* Initiate the replication stream at specified location */ query = createPQExpBuffer(); - appendPQExpBuffer(query, "START_REPLICATION SLOT \"%s\" LOGICAL %X/%X", - replication_slot, LSN_FORMAT_ARGS(startpos)); + appendPQExpBufferStr(query, "START_REPLICATION SLOT "); + AppendQuotedIdentifier(query, replication_slot); + appendPQExpBuffer(query, " LOGICAL %X/%X", LSN_FORMAT_ARGS(startpos)); /* print options if there are any */ if (noptions) @@ -258,11 +259,14 @@ StreamLogicalLog(void) appendPQExpBufferStr(query, ", "); /* write option name */ - appendPQExpBuffer(query, "\"%s\"", options[(i * 2)]); + AppendQuotedIdentifier(query, options[i * 2]); /* write option value if specified */ - if (options[(i * 2) + 1] != NULL) - appendPQExpBuffer(query, " '%s'", options[(i * 2) + 1]); + if (options[i * 2 + 1] != NULL) + { + appendPQExpBufferChar(query, ' '); + AppendQuotedLiteral(query, options[i * 2 + 1]); + } } if (noptions) @@ -342,7 +346,7 @@ StreamLogicalLog(void) outfd = fileno(stdout); else outfd = open(outfile, O_CREAT | O_APPEND | O_WRONLY | PG_BINARY, - S_IRUSR | S_IWUSR); + pg_file_create_mode); if (outfd == -1) { pg_log_error("could not open log file \"%s\": %m", outfile); diff --git a/fetools/pg18/pg_basebackup/receivelog.c b/fetools/pg18/pg_basebackup/receivelog.c index eddb581d3..89710b19f 100644 --- a/fetools/pg18/pg_basebackup/receivelog.c +++ b/fetools/pg18/pg_basebackup/receivelog.c @@ -456,8 +456,7 @@ CheckServerVersionForStreaming(PGconn *conn) bool ReceiveXlogStream(PGconn *conn, StreamCtl *stream) { - char query[128]; - char slotcmd[128]; + PQExpBuffer query; PGresult *res; XLogRecPtr stoppos; @@ -482,7 +481,6 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) if (stream->replication_slot != NULL) { reportFlushPosition = true; - sprintf(slotcmd, "SLOT \"%s\" ", stream->replication_slot); } else { @@ -490,7 +488,6 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) reportFlushPosition = true; else reportFlushPosition = false; - slotcmd[0] = 0; } if (stream->sysidentifier != NULL) @@ -539,8 +536,10 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) */ if (!existsTimeLineHistoryFile(stream)) { - snprintf(query, sizeof(query), "TIMELINE_HISTORY %u", stream->timeline); - res = PQexec(conn, query); + query = createPQExpBuffer(); + appendPQExpBuffer(query, "TIMELINE_HISTORY %u", stream->timeline); + res = PQexec(conn, query->data); + destroyPQExpBuffer(query); if (PQresultStatus(res) != PGRES_TUPLES_OK) { /* FIXME: we might send it ok, but get an error */ @@ -576,11 +575,18 @@ ReceiveXlogStream(PGconn *conn, StreamCtl *stream) return true; /* Initiate the replication stream at specified location */ - snprintf(query, sizeof(query), "START_REPLICATION %s%X/%X TIMELINE %u", - slotcmd, - LSN_FORMAT_ARGS(stream->startpos), - stream->timeline); - res = PQexec(conn, query); + query = createPQExpBuffer(); + appendPQExpBufferStr(query, "START_REPLICATION"); + if (stream->replication_slot != NULL) + { + appendPQExpBufferStr(query, " SLOT "); + AppendQuotedIdentifier(query, stream->replication_slot); + } + appendPQExpBuffer(query, " %X/%X TIMELINE %u", + LSN_FORMAT_ARGS(stream->startpos), + stream->timeline); + res = PQexec(conn, query->data); + destroyPQExpBuffer(query); if (PQresultStatus(res) != PGRES_COPY_BOTH) { pg_log_error("could not send replication command \"%s\": %s", diff --git a/fetools/pg18/pg_basebackup/streamutil.c b/fetools/pg18/pg_basebackup/streamutil.c index c7b8a4c3a..a7b4fc084 100644 --- a/fetools/pg18/pg_basebackup/streamutil.c +++ b/fetools/pg18/pg_basebackup/streamutil.c @@ -501,7 +501,8 @@ GetSlotInformation(PGconn *conn, const char *slot_name, *restart_tli = tli_loc; query = createPQExpBuffer(); - appendPQExpBuffer(query, "READ_REPLICATION_SLOT %s", slot_name); + appendPQExpBufferStr(query, "READ_REPLICATION_SLOT "); + AppendQuotedIdentifier(query, slot_name); res = PQexec(conn, query->data); destroyPQExpBuffer(query); @@ -598,13 +599,17 @@ CreateReplicationSlot(PGconn *conn, const char *slot_name, const char *plugin, Assert(slot_name != NULL); /* Build base portion of query */ - appendPQExpBuffer(query, "CREATE_REPLICATION_SLOT \"%s\"", slot_name); + appendPQExpBufferStr(query, "CREATE_REPLICATION_SLOT "); + AppendQuotedIdentifier(query, slot_name); if (is_temporary) appendPQExpBufferStr(query, " TEMPORARY"); if (is_physical) appendPQExpBufferStr(query, " PHYSICAL"); else - appendPQExpBuffer(query, " LOGICAL \"%s\"", plugin); + { + appendPQExpBufferStr(query, " LOGICAL "); + AppendQuotedIdentifier(query, plugin); + } /* Add any requested options */ if (use_new_option_syntax) @@ -704,8 +709,8 @@ DropReplicationSlot(PGconn *conn, const char *slot_name) query = createPQExpBuffer(); /* Build query */ - appendPQExpBuffer(query, "DROP_REPLICATION_SLOT \"%s\"", - slot_name); + appendPQExpBufferStr(query, "DROP_REPLICATION_SLOT "); + AppendQuotedIdentifier(query, slot_name); res = PQexec(conn, query->data); if (PQresultStatus(res) != PGRES_COMMAND_OK) { @@ -733,6 +738,29 @@ DropReplicationSlot(PGconn *conn, const char *slot_name) return true; } +/* + * Append a suitably-quoted identifier or string literal to buf. + * "quote" should be either a double-quote or single-quote character. + * + * Caution: this quoting logic is sufficient for identifiers and literals + * in the replication grammar, but not always in regular SQL. Specifically, + * it'd fail for a string literal if standard_conforming_strings is off. + */ +void +AppendQuotedString(PQExpBuffer buf, const char *str, char quote) +{ + appendPQExpBufferChar(buf, quote); + while (*str) + { + char c = *str++; + + if (c == quote) + appendPQExpBufferChar(buf, c); + appendPQExpBufferChar(buf, c); + } + appendPQExpBufferChar(buf, quote); +} + /* * Append a "plain" option - one with no value - to a server command that * is being constructed. @@ -741,10 +769,13 @@ DropReplicationSlot(PGconn *conn, const char *slot_name) * write things like SOME_COMMAND OPTION1 OPTION2 'opt2value' OPTION3 42. The * new syntax uses a comma-separated list surrounded by parentheses, so the * equivalent is SOME_COMMAND (OPTION1, OPTION2 'optvalue', OPTION3 42). + * + * Note: we assume option names do not require quotes. Do not use this + * with option names coming from outside sources. */ void AppendPlainCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name) + const char *option_name) { if (buf->len > 0 && buf->data[buf->len - 1] != '(') { @@ -765,30 +796,26 @@ AppendPlainCommandOption(PQExpBuffer buf, bool use_new_option_syntax, */ void AppendStringCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, char *option_value) + const char *option_name, const char *option_value) { AppendPlainCommandOption(buf, use_new_option_syntax, option_name); if (option_value != NULL) { - size_t length = strlen(option_value); - char *escaped_value = palloc(1 + 2 * length); - - PQescapeStringConn(conn, escaped_value, option_value, length, NULL); - appendPQExpBuffer(buf, " '%s'", escaped_value); - pfree(escaped_value); + appendPQExpBufferChar(buf, ' '); + AppendQuotedLiteral(buf, option_value); } } /* - * Append an option with an associated integer value to a server command + * Append an option with an associated integer value to a server command that * is being constructed. * * See comments for AppendPlainCommandOption, above. */ void AppendIntegerCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, int32 option_value) + const char *option_name, int32 option_value) { AppendPlainCommandOption(buf, use_new_option_syntax, option_name); diff --git a/fetools/pg18/pg_basebackup/streamutil.h b/fetools/pg18/pg_basebackup/streamutil.h index 017b22730..ee2e98350 100644 --- a/fetools/pg18/pg_basebackup/streamutil.h +++ b/fetools/pg18/pg_basebackup/streamutil.h @@ -43,15 +43,20 @@ extern bool RunIdentifySystem(PGconn *conn, char **sysid, XLogRecPtr *startpos, char **db_name); +extern void AppendQuotedString(PQExpBuffer buf, const char *str, char quote); +#define AppendQuotedIdentifier(b, s) AppendQuotedString(b, s, '"') +#define AppendQuotedLiteral(b, s) AppendQuotedString(b, s, '\'') extern void AppendPlainCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name); + const char *option_name); extern void AppendStringCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, char *option_value); + const char *option_name, + const char *option_value); extern void AppendIntegerCommandOption(PQExpBuffer buf, bool use_new_option_syntax, - char *option_name, int32 option_value); + const char *option_name, + int32 option_value); extern bool GetSlotInformation(PGconn *conn, const char *slot_name, XLogRecPtr *restart_lsn, diff --git a/fetools/pg18/xlogreader.c b/fetools/pg18/xlogreader.c index f5cfccb68..fef6c093b 100644 --- a/fetools/pg18/xlogreader.c +++ b/fetools/pg18/xlogreader.c @@ -1596,9 +1596,6 @@ WALRead(XLogReaderState *state, #ifndef FRONTEND pgstat_report_wait_end(); - - pgstat_count_io_op_time(IOOBJECT_WAL, IOCONTEXT_NORMAL, IOOP_READ, - io_start, 1, readbytes); #endif if (readbytes <= 0) @@ -1611,6 +1608,11 @@ WALRead(XLogReaderState *state, return false; } +#ifndef FRONTEND + pgstat_count_io_op_time(IOOBJECT_WAL, IOCONTEXT_NORMAL, IOOP_READ, + io_start, 1, readbytes); +#endif + /* Update state for read */ recptr += readbytes; nbytes -= readbytes;