postgres.git / summary / log / commit / refs
commit 7d74b26447229cccc2005d8b48b36bcd8bd17f33
Author: Etsuro Fujita <efujita@postgresql.org>
Date: Fri Aug 07 08:30:01 2026 +0000
Drain pending asynchronous requests during ExecReScanAppend.
The logic for asynchronous Append assumes that pending requests made for
subplans of an Append are drained during ExecReScanAppend. To ensure
that, commit 9e283fc85 modified postgresReScanForeignScan to drain such
a request if any, but failed to take into account that if such a request
was made for a subplan that is re-scanned with parameter changes or
pruned in the next round by runtime pruning, the postgres_fdw callback
function is called after ExecReScanAppend or never called, respectively.
This would cause such a request to remain even after ExecReScanAppend,
leading to incorrect results, an infinite loop, or an assertion failure.
To fix, modify ExecReScanAppend to, for each of the pending requests,
give the FDW a chance to drain that request using the existing
ForeignAsyncConfigureWait/ForeignAsyncNotify callback functions. This
makes the change made to postgresReScanForeignScan useless, so remove it
as well.
Back-patch to v14 where asynchronous Append was added.
Reported-by: Alexander Korotkov <aekorotkov@gmail.com>
Co-authored-by: Alexander Korotkov <aekorotkov@gmail.com>
Co-authored-by: Gleb Kashkin <g.kashkin@postgrespro.ru>
Co-authored-by: Etsuro Fujita <etsuro.fujita@gmail.com>
Reviewed-by: Alexander Pyhalov <a.pyhalov@postgrespro.ru>
Reviewed-by: Gleb Kashkin <g.kashkin@postgrespro.ru>
Discussion: https://postgr.es/m/CAPpHfduMOTnV5Zj2KGJ7zanL_10QvccZHtPUaDfJvBhsh9axnQ%40mail.gmail.com
Backpatch-through: 14
contrib/postgres_fdw/expected/postgres_fdw.out | 70 ++++++++++++++++++-
contrib/postgres_fdw/postgres_fdw.c | 15 ++--
contrib/postgres_fdw/sql/postgres_fdw.sql | 14 +++-
src/backend/executor/nodeAppend.c | 97 +++++++++++++++++++++-----
4 files changed, 163 insertions(+), 33 deletions(-)
diff --git a/contrib/postgres_fdw/expected/postgres_fdw.out b/contrib/postgres_fdw/expected/postgres_fdw.out
index d19121b05da..5a77671d35f 100644
--- a/contrib/postgres_fdw/expected/postgres_fdw.out
+++ b/contrib/postgres_fdw/expected/postgres_fdw.out
@@ -11836,6 +11836,72 @@ SELECT * FROM result_tbl ORDER BY a;
(3 rows)
DELETE FROM result_tbl;
+-- Test ExecAppendAsyncReset code path that drains outstanding async requests
+-- (case where subplans are re-scanned with parameter changes)
+EXPLAIN (VERBOSE, COSTS OFF)
+SELECT o.x FROM (VALUES (2505), (3505)) o(x), LATERAL (SELECT a FROM async_pt WHERE a = o.x OR a = 1505 LIMIT 1) s ORDER BY o.x;
+ QUERY PLAN
+-------------------------------------------------------------------------------------------------------------------
+ Sort
+ Output: "*VALUES*".column1
+ Sort Key: "*VALUES*".column1
+ -> Nested Loop
+ Output: "*VALUES*".column1
+ -> Values Scan on "*VALUES*"
+ Output: "*VALUES*".column1
+ -> Limit
+ Output: NULL::integer
+ -> Append
+ -> Async Foreign Scan on public.async_p1 async_pt_1
+ Output: NULL::integer
+ Remote SQL: SELECT NULL FROM public.base_tbl1 WHERE (((a = $1::integer) OR (a = 1505)))
+ -> Async Foreign Scan on public.async_p2 async_pt_2
+ Output: NULL::integer
+ Remote SQL: SELECT NULL FROM public.base_tbl2 WHERE (((a = $1::integer) OR (a = 1505)))
+ -> Async Foreign Scan on public.async_p3 async_pt_3
+ Output: NULL::integer
+ Remote SQL: SELECT NULL FROM public.base_tbl3 WHERE (((a = $1::integer) OR (a = 1505)))
+(19 rows)
+
+SELECT o.x FROM (VALUES (2505), (3505)) o(x), LATERAL (SELECT a FROM async_pt WHERE a = o.x OR a = 1505 LIMIT 1) s ORDER BY o.x;
+ x
+------
+ 2505
+ 3505
+(2 rows)
+
+EXPLAIN (VERBOSE, COSTS OFF)
+SELECT o.x FROM (VALUES (2505), (3505)) o(x), LATERAL (SELECT a FROM async_pt WHERE a = o.x OR (a = 1505 AND o.x = 2505) LIMIT 1) s ORDER BY o.x;
+ QUERY PLAN
+----------------------------------------------------------------------------------------------------------------------------------------------
+ Sort
+ Output: "*VALUES*".column1
+ Sort Key: "*VALUES*".column1
+ -> Nested Loop
+ Output: "*VALUES*".column1
+ -> Values Scan on "*VALUES*"
+ Output: "*VALUES*".column1
+ -> Limit
+ Output: NULL::integer
+ -> Append
+ -> Async Foreign Scan on public.async_p1 async_pt_1
+ Output: NULL::integer
+ Remote SQL: SELECT NULL FROM public.base_tbl1 WHERE (((a = $1::integer) OR ((a = 1505) AND ($1::integer = 2505))))
+ -> Async Foreign Scan on public.async_p2 async_pt_2
+ Output: NULL::integer
+ Remote SQL: SELECT NULL FROM public.base_tbl2 WHERE (((a = $1::integer) OR ((a = 1505) AND ($1::integer = 2505))))
+ -> Async Foreign Scan on public.async_p3 async_pt_3
+ Output: NULL::integer
+ Remote SQL: SELECT NULL FROM public.base_tbl3 WHERE (((a = $1::integer) OR ((a = 1505) AND ($1::integer = 2505))))
+(19 rows)
+
+SELECT o.x FROM (VALUES (2505), (3505)) o(x), LATERAL (SELECT a FROM async_pt WHERE a = o.x OR (a = 1505 AND o.x = 2505) LIMIT 1) s ORDER BY o.x;
+ x
+------
+ 2505
+ 3505
+(2 rows)
+
-- Test COPY TO when foreign table is partition
COPY async_pt TO stdout; --error
ERROR: cannot copy from foreign table "async_p1"
@@ -12621,8 +12687,8 @@ DROP TABLE base_tbl1;
DROP TABLE base_tbl2;
DROP TABLE result_tbl;
DROP TABLE join_tbl;
--- Test that an asynchronous fetch is processed before restarting the scan in
--- ReScanForeignScan
+-- Test ExecAppendAsyncReset code path that drains outstanding async requests
+-- (case where subplans are re-scanned without parameter changes)
CREATE TABLE base_tbl (a int, b int);
INSERT INTO base_tbl VALUES (1, 11), (2, 22), (3, 33);
CREATE FOREIGN TABLE foreign_tbl (b int)
diff --git a/contrib/postgres_fdw/postgres_fdw.c b/contrib/postgres_fdw/postgres_fdw.c
index 8b660a6c02c..b9739610131 100644
--- a/contrib/postgres_fdw/postgres_fdw.c
+++ b/contrib/postgres_fdw/postgres_fdw.c
@@ -1751,16 +1751,11 @@ postgresReScanForeignScan(ForeignScanState *node)
return;
/*
- * If the node is async-capable, and an asynchronous fetch for it has
- * begun, the asynchronous fetch might not have yet completed. Check if
- * the node is async-capable, and an asynchronous fetch for it is still in
- * progress; if so, complete the asynchronous fetch before restarting the
- * scan.
- */
- if (fsstate->async_capable &&
- fsstate->conn_state->pendingAreq &&
- fsstate->conn_state->pendingAreq->requestee == (PlanState *) node)
- fetch_more_data(node);
+ * If the node is async-capable, any asynchronous fetch made for it should
+ * have been processed before we get here (see ExecAppendAsyncReset()).
+ */
+ Assert(!fsstate->async_capable || !fsstate->conn_state->pendingAreq ||
+ fsstate->conn_state->pendingAreq->requestee != (PlanState *) node);
/*
* If any internal parameters affecting this node have changed, we'd
diff --git a/contrib/postgres_fdw/sql/postgres_fdw.sql b/contrib/postgres_fdw/sql/postgres_fdw.sql
index e7019952173..54d09040d0d 100644
--- a/contrib/postgres_fdw/sql/postgres_fdw.sql
+++ b/contrib/postgres_fdw/sql/postgres_fdw.sql
@@ -4090,6 +4090,16 @@ INSERT INTO result_tbl SELECT * FROM async_pt WHERE b === 505;
SELECT * FROM result_tbl ORDER BY a;
DELETE FROM result_tbl;
+-- Test ExecAppendAsyncReset code path that drains outstanding async requests
+-- (case where subplans are re-scanned with parameter changes)
+EXPLAIN (VERBOSE, COSTS OFF)
+SELECT o.x FROM (VALUES (2505), (3505)) o(x), LATERAL (SELECT a FROM async_pt WHERE a = o.x OR a = 1505 LIMIT 1) s ORDER BY o.x;
+SELECT o.x FROM (VALUES (2505), (3505)) o(x), LATERAL (SELECT a FROM async_pt WHERE a = o.x OR a = 1505 LIMIT 1) s ORDER BY o.x;
+
+EXPLAIN (VERBOSE, COSTS OFF)
+SELECT o.x FROM (VALUES (2505), (3505)) o(x), LATERAL (SELECT a FROM async_pt WHERE a = o.x OR (a = 1505 AND o.x = 2505) LIMIT 1) s ORDER BY o.x;
+SELECT o.x FROM (VALUES (2505), (3505)) o(x), LATERAL (SELECT a FROM async_pt WHERE a = o.x OR (a = 1505 AND o.x = 2505) LIMIT 1) s ORDER BY o.x;
+
-- Test COPY TO when foreign table is partition
COPY async_pt TO stdout; --error
@@ -4329,8 +4339,8 @@ DROP TABLE base_tbl2;
DROP TABLE result_tbl;
DROP TABLE join_tbl;
--- Test that an asynchronous fetch is processed before restarting the scan in
--- ReScanForeignScan
+-- Test ExecAppendAsyncReset code path that drains outstanding async requests
+-- (case where subplans are re-scanned without parameter changes)
CREATE TABLE base_tbl (a int, b int);
INSERT INTO base_tbl VALUES (1, 11), (2, 22), (3, 33);
CREATE FOREIGN TABLE foreign_tbl (b int)
diff --git a/src/backend/executor/nodeAppend.c b/src/backend/executor/nodeAppend.c
index 987358e27fa..51654b31ad1 100644
--- a/src/backend/executor/nodeAppend.c
+++ b/src/backend/executor/nodeAppend.c
@@ -95,6 +95,7 @@ static void ExecAppendAsyncBegin(AppendState *node);
static bool ExecAppendAsyncGetNext(AppendState *node, TupleTableSlot **result);
static bool ExecAppendAsyncRequest(AppendState *node, TupleTableSlot **result);
static void ExecAppendAsyncEventWait(AppendState *node);
+static void ExecAppendAsyncReset(AppendState *node);
static void classify_matching_subplans(AppendState *node);
/* ----------------------------------------------------------------
@@ -426,6 +427,10 @@ ExecReScanAppend(AppendState *node)
int nasyncplans = node->as_nasyncplans;
int i;
+ /* If there are any async subplans, reset async requests made for them. */
+ if (nasyncplans > 0)
+ ExecAppendAsyncReset(node);
+
/*
* If any PARAM_EXEC Params used in pruning expressions have changed, then
* we'd better unset the valid subplans so that they are reselected for
@@ -461,25 +466,6 @@ ExecReScanAppend(AppendState *node)
ExecReScan(subnode);
}
- /* Reset async state */
- if (nasyncplans > 0)
- {
- i = -1;
- while ((i = bms_next_member(node->as_asyncplans, i)) >= 0)
- {
- AsyncRequest *areq = node->as_asyncrequests[i];
-
- areq->callback_pending = false;
- areq->request_complete = false;
- areq->result = NULL;
- }
-
- node->as_nasyncresults = 0;
- node->as_nasyncremain = 0;
- bms_free(node->as_needrequest);
- node->as_needrequest = NULL;
- }
-
/* Let choose_next_subplan_* function handle setting the first subplan */
node->as_whichplan = INVALID_SUBPLAN_INDEX;
node->as_syncdone = false;
@@ -1135,6 +1121,79 @@ ExecAppendAsyncEventWait(AppendState *node)
}
}
+/* ----------------------------------------------------------------
+ * ExecAppendAsyncReset
+ *
+ * Reset asynchronous requests made for async-capable subplans.
+ * ----------------------------------------------------------------
+ */
+static void
+ExecAppendAsyncReset(AppendState *node)
+{
+ int i;
+
+ /* We should never be called when there are no async subplans. */
+ Assert(node->as_nasyncplans > 0);
+
+ /*
+ * Drain pending async requests if any. We force the as_syncdone flag to
+ * be true so that ExecAppendAsyncEventWait() waits until at least one
+ * event occurs.
+ */
+ node->as_syncdone = true;
+ for (;;)
+ {
+ bool found = false;
+
+ /*
+ * When called from ExecAppendAsyncEventWait(), postgres_fdw (and
+ * possibly other FDWs) will skip configuration of events for pending
+ * requests in some cases if as_needrequest isn't empty. To avoid
+ * that, discard results we already have. Note that we need to do
+ * this on every iteration, as the call to that function may produce
+ * new results.
+ */
+ node->as_nasyncresults = 0;
+ bms_free(node->as_needrequest);
+ node->as_needrequest = NULL;
+
+ i = -1;
+ while ((i = bms_next_member(node->as_asyncplans, i)) >= 0)
+ {
+ AsyncRequest *areq = node->as_asyncrequests[i];
+
+ if (areq->callback_pending)
+ {
+ found = true;
+ break;
+ }
+ }
+ if (!found)
+ break;
+
+ CHECK_FOR_INTERRUPTS();
+
+ /* Wait or poll for async events. */
+ ExecAppendAsyncEventWait(node);
+ }
+
+ /* Reset async requests. */
+ i = -1;
+ while ((i = bms_next_member(node->as_asyncplans, i)) >= 0)
+ {
+ AsyncRequest *areq = node->as_asyncrequests[i];
+
+ Assert(!areq->callback_pending);
+ areq->request_complete = false;
+ areq->result = NULL;
+ }
+
+ /* Reset state variables. */
+ Assert(node->as_nasyncresults == 0);
+ Assert(node->as_needrequest == NULL);
+ node->as_nasyncremain = 0;
+}
+
/* ----------------------------------------------------------------
* ExecAsyncAppendResponse
*
[parent: acdfaee1b947]