diff --git a/contrib/postgres_fdw/connection.c b/contrib/postgres_fdw/connection.c index b5d4cf3dccc..652fa4a943d 100644 --- a/contrib/postgres_fdw/connection.c +++ b/contrib/postgres_fdw/connection.c @@ -397,7 +397,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 @@ -1061,6 +1062,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 +1092,13 @@ pgfdw_exec_query(PGconn *conn, const char *query, PgFdwConnState *state) if (state && state->pendingAreq) process_pending_request(state->pendingAreq); + /* + * Second, synchronize the local/remote transactions. Note that we need + * to do this because this function can be called from open cursors. + */ + if (state) + pgfdw_begin_remote_xact(state->entry); + if (!PQsendQuery(conn, query)) return NULL; return pgfdw_get_result(conn); @@ -1929,10 +1951,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; @@ -2221,9 +2243,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 +2288,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/expected/postgres_fdw.out b/contrib/postgres_fdw/expected/postgres_fdw.out index 739f43af7bb..3aab56b0642 100644 --- a/contrib/postgres_fdw/expected/postgres_fdw.out +++ b/contrib/postgres_fdw/expected/postgres_fdw.out @@ -5332,6 +5332,8 @@ ALTER FOREIGN TABLE ft1 ALTER COLUMN c8 TYPE user_enum; -- =================================================================== -- subtransaction -- + local/remote error doesn't break cursor +-- + cursors opened before a savepoint are disallowed to be first +-- fetched within it -- =================================================================== BEGIN; DECLARE c CURSOR FOR SELECT * FROM ft1 ORDER BY c1; @@ -5371,6 +5373,12 @@ SELECT * FROM ft1 ORDER BY c1 LIMIT 1; (1 row) COMMIT; +BEGIN; +DECLARE c CURSOR FOR SELECT * FROM ft1 ORDER BY c1; +SAVEPOINT s; +FETCH c; +ERROR: cannot perform the first fetch of a cursor within a deeper subtransaction than it was created in +ABORT; -- =================================================================== -- test handling of collations -- =================================================================== diff --git a/contrib/postgres_fdw/postgres_fdw.c b/contrib/postgres_fdw/postgres_fdw.c index 2bcff4b26b4..b1d75487c67 100644 --- a/contrib/postgres_fdw/postgres_fdw.c +++ b/contrib/postgres_fdw/postgres_fdw.c @@ -17,6 +17,7 @@ #include "access/htup_details.h" #include "access/sysattr.h" #include "access/table.h" +#include "access/xact.h" #include "catalog/pg_opfamily.h" #include "commands/defrem.h" #include "commands/explain_format.h" @@ -189,6 +190,7 @@ typedef struct PgFdwScanState FmgrInfo *param_flinfo; /* output conversion functions for them */ List *param_exprs; /* executable expressions for param values */ const char **param_values; /* textual values of query parameters */ + int created_at; /* xact depth at which the scan was created */ /* for storing result tuples */ HeapTuple *tuples; /* array of currently-retrieved tuples */ @@ -1763,6 +1765,9 @@ postgresBeginForeignScan(ForeignScanState *node, int eflags) fsstate->cursor_number = GetCursorNumber(fsstate->conn); fsstate->cursor_exists = false; + /* Get the current local transaction's nesting depth */ + fsstate->created_at = GetCurrentTransactionNestLevel(); + /* Get private info created by planner functions. */ fsstate->query = strVal(list_nth(fsplan->fdw_private, FdwScanPrivateSelectSql)); @@ -4054,10 +4059,21 @@ create_cursor(ForeignScanState *node) StringInfoData buf; PGresult *res; + if (fsstate->created_at < GetCurrentTransactionNestLevel()) + ereport(ERROR, + (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), + errmsg("cannot perform the first fetch of a cursor within a deeper subtransaction than it was created in"))); + /* First, process a pending asynchronous request, if any. */ if (fsstate->conn_state->pendingAreq) process_pending_request(fsstate->conn_state->pendingAreq); + /* + * Second, synchronize the local/remote transactions. Note that we need + * to do this because this function can be called from open cursors. + */ + 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 +4163,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 @@ -4394,6 +4410,9 @@ create_foreign_modify(EState *estate, * result if any. (This is the shared guts of postgresExecForeignInsert, * postgresExecForeignBatchInsert, postgresExecForeignUpdate, and * postgresExecForeignDelete.) + * + * Note: this function is never called from open cursors, thus no need to + * synchronize the local/remote transactions. */ static TupleTableSlot ** execute_foreign_modify(EState *estate, @@ -4832,6 +4851,9 @@ rebuild_fdw_scan_tlist(ForeignScan *fscan, List *tlist) /* * Execute a direct UPDATE/DELETE statement. + * + * Note: this function is never called from open cursors, thus no need to + * synchronize the local/remote transactions. */ static void execute_dml_stmt(ForeignScanState *node) @@ -8840,9 +8862,14 @@ fetch_more_data_begin(AsyncRequest *areq) Assert(!fsstate->conn_state->pendingAreq); - /* Create the cursor synchronously. */ + /* + * Create the cursor synchronously if not already done. Otherwise, + * synchronize the local/remote transactions before the data fetch. + */ if (!fsstate->cursor_exists) create_cursor(node); + else + pgfdw_begin_remote_xact(fsstate->conn_state->entry); /* We will send this query, but not wait for the response. */ snprintf(sql, sizeof(sql), "FETCH %d FROM c%u", 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/contrib/postgres_fdw/sql/postgres_fdw.sql b/contrib/postgres_fdw/sql/postgres_fdw.sql index f1ca3204382..9c271953206 100644 --- a/contrib/postgres_fdw/sql/postgres_fdw.sql +++ b/contrib/postgres_fdw/sql/postgres_fdw.sql @@ -1635,6 +1635,8 @@ ALTER FOREIGN TABLE ft1 ALTER COLUMN c8 TYPE user_enum; -- =================================================================== -- subtransaction -- + local/remote error doesn't break cursor +-- + cursors opened before a savepoint are disallowed to be first +-- fetched within it -- =================================================================== BEGIN; DECLARE c CURSOR FOR SELECT * FROM ft1 ORDER BY c1; @@ -1650,6 +1652,12 @@ FETCH c; SELECT * FROM ft1 ORDER BY c1 LIMIT 1; COMMIT; +BEGIN; +DECLARE c CURSOR FOR SELECT * FROM ft1 ORDER BY c1; +SAVEPOINT s; +FETCH c; +ABORT; + -- =================================================================== -- test handling of collations -- ===================================================================