postgres.git / summary / log / commit / refs
commit 8f7617ff7af60f4d9b2be0ba2d89b2fedbc027cc
Author: Jeff Davis <jdavis@postgresql.org>
Date: Wed Aug 05 18:36:37 2026 +0000
Remove Subscription conninfo field; generate in caller.
After server-based subscriptions, conninfo became more than just a
catalog field. It has its own error paths, and it's important that
callers that don't need conninfo don't encounter errors related to it.
Reviewed-by: Amit Kapila <amit.kapila16@gmail.com>
Reviewed-by: Hayato Kuroda (Fujitsu) <kuroda.hayato@fujitsu.com>
Discussion: https://postgr.es/m/20260710195902.4f.noahmisch%40microsoft.com
Backpatch-through: 19
src/backend/catalog/pg_subscription.c | 105 ++++++++++++++-----------
src/backend/commands/subscriptioncmds.c | 43 +++++++---
src/backend/replication/logical/sequencesync.c | 2 +-
src/backend/replication/logical/tablesync.c | 2 +-
src/backend/replication/logical/worker.c | 24 +++++-
src/include/catalog/pg_subscription.h | 6 +-
src/include/replication/worker_internal.h | 1 +
7 files changed, 118 insertions(+), 65 deletions(-)
diff --git a/src/backend/catalog/pg_subscription.c b/src/backend/catalog/pg_subscription.c
index d1a110f1ff3..d1d5478b8fa 100644
--- a/src/backend/catalog/pg_subscription.c
+++ b/src/backend/catalog/pg_subscription.c
@@ -79,14 +79,10 @@ GetPublicationsStr(List *publications, StringInfo dest, bool quote_literal)
/*
* Fetch the subscription from the syscache.
*
- * If conninfo_needed is true, conninfo will be constructed, possibly
- * encountering errors in ForeignServerConnectionString(). Callers not
- * expecting such errors should pass false, in which case conninfo will be
- * NULL.
+ * Callers that need conninfo must call SubscriptionConninfo().
*/
Subscription *
-GetSubscription(Oid subid, bool missing_ok, bool conninfo_needed,
- bool conninfo_aclcheck)
+GetSubscription(Oid subid, bool missing_ok)
{
HeapTuple tup;
Subscription *sub;
@@ -96,8 +92,6 @@ GetSubscription(Oid subid, bool missing_ok, bool conninfo_needed,
MemoryContext cxt;
MemoryContext oldcxt;
- Assert(conninfo_needed || !conninfo_aclcheck);
-
tup = SearchSysCache1(SUBSCRIPTIONOID, ObjectIdGetDatum(subid));
if (!HeapTupleIsValid(tup))
@@ -139,42 +133,6 @@ GetSubscription(Oid subid, bool missing_ok, bool conninfo_needed,
sub->maxretention = subform->submaxretention;
sub->retentionactive = subform->subretentionactive;
- if (conninfo_needed)
- {
- if (OidIsValid(subform->subserver))
- {
- AclResult aclresult;
- ForeignServer *server;
-
- server = GetForeignServer(subform->subserver);
-
- if (conninfo_aclcheck)
- {
- /* recheck ACL if requested */
- aclresult = object_aclcheck(ForeignServerRelationId,
- subform->subserver,
- subform->subowner, ACL_USAGE);
-
- if (aclresult != ACLCHECK_OK)
- ereport(ERROR,
- (errcode(ERRCODE_INSUFFICIENT_PRIVILEGE),
- errmsg("subscription owner \"%s\" does not have permission on foreign server \"%s\"",
- GetUserNameFromId(subform->subowner, false),
- server->servername)));
- }
-
- sub->conninfo = ForeignServerConnectionString(subform->subowner,
- server);
- }
- else
- {
- datum = SysCacheGetAttrNotNull(SUBSCRIPTIONOID,
- tup,
- Anum_pg_subscription_subconninfo);
- sub->conninfo = TextDatumGetCString(datum);
- }
- }
-
/* Get slotname */
datum = SysCacheGetAttr(SUBSCRIPTIONOID,
tup,
@@ -219,6 +177,65 @@ GetSubscription(Oid subid, bool missing_ok, bool conninfo_needed,
return sub;
}
+/*
+ * Generate the connection string for a subscription.
+ *
+ * This is deliberately separate from GetSubscription() because resolving
+ * conninfo for a server-based subscription has its own error paths (foreign
+ * server USAGE, user mapping, ForeignServerConnectionString()). Keeping it
+ * separate lets a caller load the subscription and decide whether a
+ * connection is actually needed, and check things such as whether the
+ * subscription is enabled, before risking those errors. Callers that never
+ * connect thus never hit them, which matters during restore.
+ */
+char *
+SubscriptionConninfo(Subscription *sub, bool aclcheck)
+{
+ HeapTuple tup;
+ Form_pg_subscription subform;
+ Datum datum;
+ char *conninfo;
+
+ tup = SearchSysCache1(SUBSCRIPTIONOID, ObjectIdGetDatum(sub->oid));
+ if (!HeapTupleIsValid(tup))
+ elog(ERROR, "cache lookup failed for subscription %u", sub->oid);
+
+ subform = (Form_pg_subscription) GETSTRUCT(tup);
+
+ if (OidIsValid(subform->subserver))
+ {
+ ForeignServer *server;
+ AclResult aclresult;
+
+ server = GetForeignServer(subform->subserver);
+
+ if (aclcheck)
+ {
+ aclresult = object_aclcheck(ForeignServerRelationId,
+ subform->subserver,
+ sub->owner, ACL_USAGE);
+ if (aclresult != ACLCHECK_OK)
+ ereport(ERROR,
+ (errcode(ERRCODE_INSUFFICIENT_PRIVILEGE),
+ errmsg("subscription owner \"%s\" does not have permission on foreign server \"%s\"",
+ GetUserNameFromId(sub->owner, false),
+ server->servername)));
+ }
+
+ conninfo = ForeignServerConnectionString(sub->owner, server);
+ }
+ else
+ {
+ datum = SysCacheGetAttrNotNull(SUBSCRIPTIONOID, tup,
+ Anum_pg_subscription_subconninfo);
+ conninfo = TextDatumGetCString(datum);
+ }
+
+ ReleaseSysCache(tup);
+
+ return conninfo;
+}
+
/*
* Return number of subscriptions defined in given database.
* Used by dropdb() to check if database can indeed be dropped.
diff --git a/src/backend/commands/subscriptioncmds.c b/src/backend/commands/subscriptioncmds.c
index 9671541caf9..f553fe3294e 100644
--- a/src/backend/commands/subscriptioncmds.c
+++ b/src/backend/commands/subscriptioncmds.c
@@ -1028,7 +1028,7 @@ CreateSubscription(ParseState *pstate, CreateSubscriptionStmt *stmt,
static void
AlterSubscription_refresh(Subscription *sub, bool copy_data,
- List *validate_publications)
+ List *validate_publications, char *conninfo)
{
char *err;
List *pubrels = NIL;
@@ -1052,12 +1052,19 @@ AlterSubscription_refresh(Subscription *sub, bool copy_data,
WalReceiverConn *wrconn;
bool must_use_password;
+ /*
+ * Should not happen: CREATE/ALTER/DROP SUBSCRIPTION did not call
+ * SubscriptionConninfo() in a path where it's required.
+ */
+ if (!conninfo)
+ elog(ERROR, "no connection string provided for subscription");
+
/* Load the library providing us libpq calls. */
load_file("libpqwalreceiver", false);
/* Try to connect to the publisher. */
must_use_password = sub->passwordrequired && !sub->ownersuperuser;
- wrconn = walrcv_connect(sub->conninfo, true, true, must_use_password,
+ wrconn = walrcv_connect(conninfo, true, true, must_use_password,
sub->name, &err);
if (!wrconn)
ereport(ERROR,
@@ -1298,19 +1305,26 @@ AlterSubscription_refresh(Subscription *sub, bool copy_data,
* Marks all sequences with INIT state.
*/
static void
-AlterSubscription_refresh_seq(Subscription *sub)
+AlterSubscription_refresh_seq(Subscription *sub, char *conninfo)
{
char *err = NULL;
WalReceiverConn *wrconn;
bool must_use_password;
List *subrel_states;
+ /*
+ * Should not happen: CREATE/ALTER/DROP SUBSCRIPTION did not call
+ * SubscriptionConninfo() in a path where it's required.
+ */
+ if (!conninfo)
+ elog(ERROR, "no connection string provided for subscription");
+
/* Load the library providing us libpq calls. */
load_file("libpqwalreceiver", false);
/* Try to connect to the publisher. */
must_use_password = sub->passwordrequired && !sub->ownersuperuser;
- wrconn = walrcv_connect(sub->conninfo, true, true, must_use_password,
+ wrconn = walrcv_connect(conninfo, true, true, must_use_password,
sub->name, &err);
if (!wrconn)
ereport(ERROR,
@@ -1502,6 +1516,7 @@ AlterSubscription(ParseState *pstate, AlterSubscriptionStmt *stmt,
int max_retention;
bool retention_active;
char *new_conninfo = NULL;
+ char *orig_conninfo = NULL;
char *origin;
Subscription *sub;
Form_pg_subscription form;
@@ -1603,6 +1618,8 @@ AlterSubscription(ParseState *pstate, AlterSubscriptionStmt *stmt,
orig_conninfo_needed = false;
}
+ sub = GetSubscription(subid, false);
+
/*
* Skip ACL checks on the subscription's foreign server, if any. If
* changing the server (or replacing it with a raw connection), then the
@@ -1610,7 +1627,8 @@ AlterSubscription(ParseState *pstate, AlterSubscriptionStmt *stmt,
* there's no need to do an additional ACL check here; that will be done
* by the subscription worker.
*/
- sub = GetSubscription(subid, false, orig_conninfo_needed, false);
+ if (orig_conninfo_needed)
+ orig_conninfo = SubscriptionConninfo(sub, false);
retain_dead_tuples = sub->retaindeadtuples;
origin = sub->origin;
@@ -2065,7 +2083,8 @@ AlterSubscription(ParseState *pstate, AlterSubscriptionStmt *stmt,
sub->publications = stmt->publication;
AlterSubscription_refresh(sub, opts.copy_data,
- stmt->publication);
+ stmt->publication,
+ orig_conninfo);
}
break;
@@ -2120,7 +2139,8 @@ AlterSubscription(ParseState *pstate, AlterSubscriptionStmt *stmt,
sub->publications = publist;
AlterSubscription_refresh(sub, opts.copy_data,
- validate_publications);
+ validate_publications,
+ orig_conninfo);
}
break;
@@ -2159,7 +2179,8 @@ AlterSubscription(ParseState *pstate, AlterSubscriptionStmt *stmt,
PreventInTransactionBlock(isTopLevel, "ALTER SUBSCRIPTION ... REFRESH PUBLICATION");
- AlterSubscription_refresh(sub, opts.copy_data, NULL);
+ AlterSubscription_refresh(sub, opts.copy_data, NULL,
+ orig_conninfo);
break;
}
@@ -2172,7 +2193,7 @@ AlterSubscription(ParseState *pstate, AlterSubscriptionStmt *stmt,
errmsg("%s is not allowed for disabled subscriptions",
"ALTER SUBSCRIPTION ... REFRESH SEQUENCES"));
- AlterSubscription_refresh_seq(sub);
+ AlterSubscription_refresh_seq(sub, orig_conninfo);
break;
}
@@ -2244,7 +2265,7 @@ AlterSubscription(ParseState *pstate, AlterSubscriptionStmt *stmt,
char *err;
WalReceiverConn *wrconn;
- Assert(new_conninfo || orig_conninfo_needed);
+ Assert(new_conninfo || orig_conninfo);
/* Load the library providing us libpq calls. */
load_file("libpqwalreceiver", false);
@@ -2254,7 +2275,7 @@ AlterSubscription(ParseState *pstate, AlterSubscriptionStmt *stmt,
* available.
*/
must_use_password = sub->passwordrequired && !sub->ownersuperuser;
- wrconn = walrcv_connect(new_conninfo ? new_conninfo : sub->conninfo,
+ wrconn = walrcv_connect(new_conninfo ? new_conninfo : orig_conninfo,
true, true, must_use_password, sub->name,
&err);
if (!wrconn)
diff --git a/src/backend/replication/logical/sequencesync.c b/src/backend/replication/logical/sequencesync.c
index 9f0d2762dad..35286c3f727 100644
--- a/src/backend/replication/logical/sequencesync.c
+++ b/src/backend/replication/logical/sequencesync.c
@@ -815,7 +815,7 @@ LogicalRepSyncSequences(void)
* Establish the connection to the publisher for sequence synchronization.
*/
LogRepWorkerWalRcvConn =
- walrcv_connect(MySubscription->conninfo, true, true,
+ walrcv_connect(MySubscriptionConninfo, true, true,
must_use_password,
app_name.data, &err);
if (LogRepWorkerWalRcvConn == NULL)
diff --git a/src/backend/replication/logical/tablesync.c b/src/backend/replication/logical/tablesync.c
index a04b84ebc1d..e5101997cd3 100644
--- a/src/backend/replication/logical/tablesync.c
+++ b/src/backend/replication/logical/tablesync.c
@@ -1305,7 +1305,7 @@ LogicalRepSyncTableStart(XLogRecPtr *origin_startpos)
* so that synchronous replication can distinguish them.
*/
LogRepWorkerWalRcvConn =
- walrcv_connect(MySubscription->conninfo, true, true,
+ walrcv_connect(MySubscriptionConninfo, true, true,
must_use_password,
slotname, &err);
if (LogRepWorkerWalRcvConn == NULL)
diff --git a/src/backend/replication/logical/worker.c b/src/backend/replication/logical/worker.c
index 2dd421412d6..56949f49788 100644
--- a/src/backend/replication/logical/worker.c
+++ b/src/backend/replication/logical/worker.c
@@ -482,6 +482,7 @@ static MemoryContext LogicalStreamingContext = NULL;
WalReceiverConn *LogRepWorkerWalRcvConn = NULL;
Subscription *MySubscription = NULL;
+char *MySubscriptionConninfo = NULL;
static bool MySubscriptionValid = false;
static List *on_commit_wakeup_workers_subids = NIL;
@@ -5061,6 +5062,8 @@ void
maybe_reread_subscription(void)
{
Subscription *newsub;
+ char *old_conninfo;
+ char *new_conninfo;
bool started_tx = false;
/* When cache state is valid there is nothing to do here. */
@@ -5074,7 +5077,7 @@ maybe_reread_subscription(void)
started_tx = true;
}
- newsub = GetSubscription(MyLogicalRepWorker->subid, true, true, true);
+ newsub = GetSubscription(MyLogicalRepWorker->subid, true);
if (newsub)
{
@@ -5097,6 +5100,9 @@ maybe_reread_subscription(void)
proc_exit(0);
}
+ /* allocated in transaction context */
+ new_conninfo = SubscriptionConninfo(newsub, true);
+
/* Exit if the subscription was disabled. */
if (!newsub->enabled)
{
@@ -5120,7 +5126,7 @@ maybe_reread_subscription(void)
* 'parallel' to any other value or the server decides not to stream the
* in-progress transaction.
*/
- if (strcmp(newsub->conninfo, MySubscription->conninfo) != 0 ||
+ if (strcmp(new_conninfo, MySubscriptionConninfo) != 0 ||
strcmp(newsub->name, MySubscription->name) != 0 ||
strcmp(newsub->slotname, MySubscription->slotname) != 0 ||
newsub->binary != MySubscription->binary ||
@@ -5171,6 +5177,11 @@ maybe_reread_subscription(void)
MemoryContextDelete(MySubscription->cxt);
MySubscription = newsub;
+ /* copy to ApplyContext and update MySubscriptionConninfo */
+ old_conninfo = MySubscriptionConninfo;
+ MySubscriptionConninfo = MemoryContextStrdup(ApplyContext, new_conninfo);
+ pfree(old_conninfo);
+
/* Change synchronous commit according to the user's wishes */
SetConfigOption("synchronous_commit", MySubscription->synccommit,
PGC_BACKEND, PGC_S_OVERRIDE);
@@ -5718,7 +5729,7 @@ run_apply_worker(void)
must_use_password = MySubscription->passwordrequired &&
!MySubscription->ownersuperuser;
- LogRepWorkerWalRcvConn = walrcv_connect(MySubscription->conninfo, true,
+ LogRepWorkerWalRcvConn = walrcv_connect(MySubscriptionConninfo, true,
true, must_use_password,
MySubscription->name, &err);
@@ -5838,7 +5849,7 @@ InitializeLogRepWorker(void)
LockSharedObject(SubscriptionRelationId, MyLogicalRepWorker->subid, 0,
AccessShareLock);
- MySubscription = GetSubscription(MyLogicalRepWorker->subid, true, true, true);
+ MySubscription = GetSubscription(MyLogicalRepWorker->subid, true);
if (MySubscription)
{
@@ -5857,6 +5868,11 @@ InitializeLogRepWorker(void)
proc_exit(0);
}
+ /* build conninfo in transaction context and copy to ApplyContext */
+ MySubscriptionConninfo =
+ MemoryContextStrdup(ApplyContext,
+ SubscriptionConninfo(MySubscription, true));
+
MySubscriptionValid = true;
if (!MySubscription->enabled)
diff --git a/src/include/catalog/pg_subscription.h b/src/include/catalog/pg_subscription.h
index 48944201889..cbb46113c2d 100644
--- a/src/include/catalog/pg_subscription.h
+++ b/src/include/catalog/pg_subscription.h
@@ -164,7 +164,6 @@ typedef struct Subscription
* and the retention duration has not
* exceeded max_retention_duration, when
* defined */
- char *conninfo; /* Connection string to the publisher */
char *slotname; /* Name of the replication slot */
char *synccommit; /* Synchronous commit setting for worker */
char *walrcvtimeout; /* wal_receiver_timeout setting for worker */
@@ -212,9 +211,8 @@ typedef struct Subscription
#endif /* EXPOSE_TO_CLIENT_CODE */
-extern Subscription *GetSubscription(Oid subid, bool missing_ok,
- bool conninfo_needed,
- bool conninfo_aclcheck);
+extern Subscription *GetSubscription(Oid subid, bool missing_ok);
+extern char *SubscriptionConninfo(Subscription *sub, bool aclcheck);
extern void DisableSubscription(Oid subid);
extern int CountDBSubscriptions(Oid dbid);
diff --git a/src/include/replication/worker_internal.h b/src/include/replication/worker_internal.h
index 745b7d9e969..88cb7c1e252 100644
--- a/src/include/replication/worker_internal.h
+++ b/src/include/replication/worker_internal.h
@@ -247,6 +247,7 @@ extern PGDLLIMPORT struct WalReceiverConn *LogRepWorkerWalRcvConn;
/* Worker and subscription objects. */
extern PGDLLIMPORT Subscription *MySubscription;
+extern PGDLLIMPORT char *MySubscriptionConninfo;
extern PGDLLIMPORT LogicalRepWorker *MyLogicalRepWorker;
extern PGDLLIMPORT bool in_remote_transaction;
[parent: 55ea764269b6]