diff --git a/contrib/postgres_fdw/connection.c b/contrib/postgres_fdw/connection.c
index b5d4cf3dccc..00c264d0276 100644
--- a/contrib/postgres_fdw/connection.c
+++ b/contrib/postgres_fdw/connection.c
@@ -112,6 +112,20 @@ static uint32 pgfdw_we_get_result = 0;
*/
#define RETRY_CANCEL_TIMEOUT 1000
+/*
+ * Macro for constructing commit command to be sent
+ *
+ * We synchronize the read/write mode before committing remote transactions
+ * so deferred triggers on remote servers can run in the right mode.
+ */
+#define CONSTRUCT_COMMIT_COMMAND(sql, entry) \
+ do { \
+ if ((read_only_level > 0) && !(entry)->xact_read_only) \
+ strcpy((sql), "SET TRANSACTION READ ONLY; COMMIT TRANSACTION"); \
+ else \
+ strcpy((sql), "COMMIT TRANSACTION"); \
+ } while(0)
+
/* Macro for constructing abort command to be sent */
#define CONSTRUCT_ABORT_COMMAND(sql, entry, toplevel) \
do { \
@@ -397,7 +411,8 @@ make_new_connection(ConnCacheEntry *entry, UserMapping *user)
entry->mapping_hashvalue =
GetSysCacheHashValue1(USERMAPPINGOID,
ObjectIdGetDatum(user->umid));
- memset(&entry->state, 0, sizeof(entry->state));
+ entry->state.pendingAreq = NULL;
+ entry->state.entry = entry;
/*
* Determine whether to keep the connection that we're about to make here
@@ -929,6 +944,7 @@ begin_remote_xact(ConnCacheEntry *entry)
*/
StringInfoData sql;
bool ro = (read_only_level == 1);
+ int remoteversion = PQserverVersion(entry->conn);
elog(DEBUG3, "starting remote transaction on connection %p",
entry->conn);
@@ -941,8 +957,15 @@ begin_remote_xact(ConnCacheEntry *entry)
appendStringInfoString(&sql, "REPEATABLE READ");
if (ro)
appendStringInfoString(&sql, " READ ONLY");
- if (XactDeferrable)
- appendStringInfoString(&sql, " DEFERRABLE");
+ else
+ appendStringInfoString(&sql, " READ WRITE");
+ if (remoteversion >= 90100)
+ {
+ if (XactDeferrable)
+ appendStringInfoString(&sql, " DEFERRABLE");
+ else
+ appendStringInfoString(&sql, " NOT DEFERRABLE");
+ }
entry->changing_xact_state = true;
do_sql_command(entry->conn, sql.data);
entry->xact_depth = 1;
@@ -971,7 +994,7 @@ begin_remote_xact(ConnCacheEntry *entry)
if (entry->xact_depth == read_only_level)
{
entry->changing_xact_state = true;
- do_sql_command(entry->conn, "SET transaction_read_only = on");
+ do_sql_command(entry->conn, "SET TRANSACTION READ ONLY");
entry->xact_read_only = true;
entry->changing_xact_state = false;
}
@@ -1004,7 +1027,7 @@ begin_remote_xact(ConnCacheEntry *entry)
initStringInfo(&sql);
appendStringInfo(&sql, "SAVEPOINT s%d", entry->xact_depth + 1);
if (ro)
- appendStringInfoString(&sql, "; SET transaction_read_only = on");
+ appendStringInfoString(&sql, "; SET TRANSACTION READ ONLY");
entry->changing_xact_state = true;
do_sql_command(entry->conn, sql.data);
entry->xact_depth++;
@@ -1061,6 +1084,20 @@ GetPrepStmtNumber(PGconn *conn)
return ++prep_stmt_number;
}
+/*
+ * Exported version of begin_remote_xact().
+ *
+ * This can be called for connections on which begin_remote_xact() has started
+ * a remote transaction.
+ */
+void
+pgfdw_begin_remote_xact(ConnCacheEntry *entry)
+{
+ Assert(entry);
+ Assert(entry->xact_depth > 0);
+ begin_remote_xact(entry);
+}
+
/*
* Submit a query and wait for the result.
*
@@ -1077,6 +1114,14 @@ pgfdw_exec_query(PGconn *conn, const char *query, PgFdwConnState *state)
if (state && state->pendingAreq)
process_pending_request(state->pendingAreq);
+ /*
+ * Ensure the local and remote (sub)transactions are synchronized. Note
+ * that we need to do this because this function can be called from open
+ * cursors, bypassing begin_remote_xact().
+ */
+ if (state)
+ pgfdw_begin_remote_xact(state->entry);
+
if (!PQsendQuery(conn, query))
return NULL;
return pgfdw_get_result(conn);
@@ -1184,6 +1229,19 @@ pgfdw_xact_callback(XactEvent event, void *arg)
if (!xact_got_connection)
return;
+ /*
+ * The local transaction may have become read-only since the last remote
+ * operation, so ensure read_only_level is set for later processing.
+ */
+ if (XactReadOnly)
+ {
+ if (read_only_level == 0)
+ read_only_level = 1;
+ Assert(read_only_level == 1);
+ }
+ else
+ Assert(read_only_level == 0);
+
/*
* Scan all connection cache entries to find open remote transactions, and
* close them.
@@ -1200,6 +1258,8 @@ pgfdw_xact_callback(XactEvent event, void *arg)
/* If it has an open remote transaction, try to close it */
if (entry->xact_depth > 0)
{
+ char sql[100];
+
elog(DEBUG3, "closing remote transaction on connection %p",
entry->conn);
@@ -1215,14 +1275,17 @@ pgfdw_xact_callback(XactEvent event, void *arg)
pgfdw_reject_incomplete_xact_state_change(entry);
/* Commit all remote transactions during pre-commit */
+ CONSTRUCT_COMMIT_COMMAND(sql, entry);
entry->changing_xact_state = true;
if (entry->parallel_commit)
{
- do_sql_command_begin(entry->conn, "COMMIT TRANSACTION");
+ do_sql_command_begin(entry->conn, sql);
pending_entries = lappend(pending_entries, entry);
continue;
}
- do_sql_command(entry->conn, "COMMIT TRANSACTION");
+ do_sql_command(entry->conn, sql);
+ if ((read_only_level > 0) && !entry->xact_read_only)
+ entry->xact_read_only = true;
entry->changing_xact_state = false;
/*
@@ -1929,10 +1992,10 @@ pgfdw_abort_cleanup(ConnCacheEntry *entry, bool toplevel)
* If pendingAreq of the per-connection state is not NULL, it means that
* an asynchronous fetch begun by fetch_more_data_begin() was not done
* successfully and thus the per-connection state was not reset in
- * fetch_more_data(); in that case reset the per-connection state here.
+ * fetch_more_data(); in that case reset pendingAreq here.
*/
if (entry->state.pendingAreq)
- memset(&entry->state, 0, sizeof(entry->state));
+ entry->state.pendingAreq = NULL;
/* Disarm changing_xact_state if it all worked */
entry->changing_xact_state = false;
@@ -2019,6 +2082,8 @@ pgfdw_finish_pre_commit_cleanup(List *pending_entries)
*/
foreach(lc, pending_entries)
{
+ char sql[100];
+
entry = (ConnCacheEntry *) lfirst(lc);
Assert(entry->changing_xact_state);
@@ -2027,7 +2092,10 @@ pgfdw_finish_pre_commit_cleanup(List *pending_entries)
* We might already have received the result on the socket, so pass
* consume_input=true to try to consume it first
*/
- do_sql_command_end(entry->conn, "COMMIT TRANSACTION", true);
+ CONSTRUCT_COMMIT_COMMAND(sql, entry);
+ do_sql_command_end(entry->conn, sql, true);
+ if ((read_only_level > 0) && !(entry)->xact_read_only)
+ entry->xact_read_only = true;
entry->changing_xact_state = false;
/* Do a DEALLOCATE ALL in parallel if needed */
@@ -2221,9 +2289,9 @@ pgfdw_finish_abort_cleanup(List *pending_entries, List *cancel_requested,
entry->have_error = false;
}
- /* Reset the per-connection state if needed */
+ /* Reset pendingAreq here if any */
if (entry->state.pendingAreq)
- memset(&entry->state, 0, sizeof(entry->state));
+ entry->state.pendingAreq = NULL;
/* We're done with this entry; unset the changing_xact_state flag */
entry->changing_xact_state = false;
@@ -2266,9 +2334,9 @@ pgfdw_finish_abort_cleanup(List *pending_entries, List *cancel_requested,
entry->have_prep_stmt = false;
entry->have_error = false;
- /* Reset the per-connection state if needed */
+ /* Reset pendingAreq here if any */
if (entry->state.pendingAreq)
- memset(&entry->state, 0, sizeof(entry->state));
+ entry->state.pendingAreq = NULL;
/* We're done with this entry; unset the changing_xact_state flag */
entry->changing_xact_state = false;
diff --git a/contrib/postgres_fdw/postgres_fdw.c b/contrib/postgres_fdw/postgres_fdw.c
index 2bcff4b26b4..445eb3d4d79 100644
--- a/contrib/postgres_fdw/postgres_fdw.c
+++ b/contrib/postgres_fdw/postgres_fdw.c
@@ -4058,6 +4058,13 @@ create_cursor(ForeignScanState *node)
if (fsstate->conn_state->pendingAreq)
process_pending_request(fsstate->conn_state->pendingAreq);
+ /*
+ * Ensure the local and remote (sub)transactions are synchronized. Note
+ * that we need to do this because this function can be called from open
+ * cursors, bypassing begin_remote_xact().
+ */
+ pgfdw_begin_remote_xact(fsstate->conn_state->entry);
+
/*
* Construct array of query parameter values in text format. We do the
* conversions in the short-lived per-tuple context, so as not to cause a
@@ -4147,7 +4154,7 @@ fetch_more_data(ForeignScanState *node)
if (PQresultStatus(res) != PGRES_TUPLES_OK)
pgfdw_report_error(res, conn, fsstate->query);
- /* Reset per-connection state */
+ /* Reset the pending asynchronous request */
fsstate->conn_state->pendingAreq = NULL;
}
else
@@ -8840,6 +8847,13 @@ fetch_more_data_begin(AsyncRequest *areq)
Assert(!fsstate->conn_state->pendingAreq);
+ /*
+ * Ensure the local and remote (sub)transactions are synchronized. Note
+ * that we need to do this because this function can be called from open
+ * cursors, bypassing begin_remote_xact().
+ */
+ pgfdw_begin_remote_xact(fsstate->conn_state->entry);
+
/* Create the cursor synchronously. */
if (!fsstate->cursor_exists)
create_cursor(node);
diff --git a/contrib/postgres_fdw/postgres_fdw.h b/contrib/postgres_fdw/postgres_fdw.h
index da7da1c2ea9..b9b460141f9 100644
--- a/contrib/postgres_fdw/postgres_fdw.h
+++ b/contrib/postgres_fdw/postgres_fdw.h
@@ -146,7 +146,8 @@ typedef struct PgFdwRelationInfo
*/
typedef struct PgFdwConnState
{
- AsyncRequest *pendingAreq; /* pending async request */
+ AsyncRequest *pendingAreq; /* pending async request */
+ struct ConnCacheEntry *entry; /* link to containing ConnCacheEntry */
} PgFdwConnState;
/*
@@ -173,6 +174,7 @@ extern void ReleaseConnection(PGconn *conn);
extern unsigned int GetCursorNumber(PGconn *conn);
extern unsigned int GetPrepStmtNumber(PGconn *conn);
extern void do_sql_command(PGconn *conn, const char *sql);
+extern void pgfdw_begin_remote_xact(struct ConnCacheEntry *entry);
extern PGresult *pgfdw_get_result(PGconn *conn);
extern PGresult *pgfdw_exec_query(PGconn *conn, const char *query,
PgFdwConnState *state);
diff --git a/doc/src/sgml/postgres-fdw.sgml b/doc/src/sgml/postgres-fdw.sgml
index fe4e6478e28..24fd5e342b4 100644
--- a/doc/src/sgml/postgres-fdw.sgml
+++ b/doc/src/sgml/postgres-fdw.sgml
@@ -1162,6 +1162,7 @@ CREATE SUBSCRIPTION my_subscription SERVER subscription_server PUBLICATION testp
local transaction: if the local transaction is DEFERRABLE,
the remote transaction is opened in DEFERRABLE mode,
otherwise it is opened in NOT DEFERRABLE mode.
+ (This rule is only applied to remote servers 9.1 and newer.)