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.)