This is an automated email from the ASF dual-hosted git repository.

Arsnael pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/james-project.git


The following commit(s) were added to refs/heads/master by this push:
     new d2d6cdf343 [FIX] Improve PostgresExecutor cancelation support
d2d6cdf343 is described below

commit d2d6cdf343ba9fe2ef9666f3566e613f93d11f04
Author: Benoit TELLIER <[email protected]>
AuthorDate: Fri Oct 2 22:11:56 2026 +0200

    [FIX] Improve PostgresExecutor cancelation support
---
 .../backends/postgres/utils/PostgresExecutor.java  | 81 +++++++++++++++-------
 .../postgres/PostgresExecutorTimeoutTest.java      | 36 ++++++++++
 2 files changed, 93 insertions(+), 24 deletions(-)

diff --git 
a/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java
 
b/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java
index b5fbd1c128..3bed6f8858 100644
--- 
a/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java
+++ 
b/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java
@@ -119,15 +119,14 @@ public class PostgresExecutor {
 
     public Mono<Void> executeVoid(Function<DSLContext, Mono<?>> queryFunction) 
{
         return 
Mono.from(metricFactory.decoratePublisherWithTimerMetric("postgres-execution",
-            Mono.usingWhen(getConnection(domain),
+            usingConnection(
                 connection -> dslContext(connection)
                     .flatMap(queryFunction)
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
                     .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
                         .filter(preparedStatementConflictException()))
-                    .then(),
-                jamesPostgresConnectionFactory::closeConnection)));
+                    .then())));
     }
 
     public Flux<Record> executeRows(Function<DSLContext, Flux<Record>> 
queryFunction) {
@@ -150,7 +149,7 @@ public class PostgresExecutor {
      */
     public Flux<Record> executeRows(Function<DSLContext, Flux<Record>> 
queryFunction, boolean isEagerFetch) {
         return 
Flux.from(metricFactory.decoratePublisherWithTimerMetric("postgres-execution",
-            Flux.usingWhen(getConnection(domain),
+            usingConnectionMany(
                 connection -> {
                     Flux<Record> recordFlux = dslContext(connection)
                         .flatMapMany(queryFunction)
@@ -164,8 +163,7 @@ public class PostgresExecutor {
                     } else {
                         return recordFlux;
                     }
-                },
-                jamesPostgresConnectionFactory::closeConnection)));
+                })));
     }
 
     /**
@@ -208,26 +206,24 @@ public class PostgresExecutor {
 
     public Flux<Record> executeDeleteAndReturnList(Function<DSLContext, 
DeleteResultStep<Record>> queryFunction) {
         return 
Flux.from(metricFactory.decoratePublisherWithTimerMetric("postgres-execution",
-            Flux.usingWhen(getConnection(domain),
+            usingConnectionMany(
                 connection -> dslContext(connection)
                     .flatMapMany(queryFunction)
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
                     .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
-                        .filter(preparedStatementConflictException())),
-                jamesPostgresConnectionFactory::closeConnection)));
+                        .filter(preparedStatementConflictException())))));
     }
 
     public Mono<Record> executeRow(Function<DSLContext, Publisher<Record>> 
queryFunction) {
         return 
Mono.from(metricFactory.decoratePublisherWithTimerMetric("postgres-execution",
-            Mono.usingWhen(getConnection(domain),
+            usingConnection(
                 connection -> dslContext(connection)
                     .flatMap(queryFunction.andThen(Mono::from))
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
                     .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
-                        .filter(preparedStatementConflictException())),
-                jamesPostgresConnectionFactory::closeConnection)));
+                        .filter(preparedStatementConflictException())))));
     }
 
     public Mono<Optional<Record>> 
executeSingleRowOptional(Function<DSLContext, Publisher<Record>> queryFunction) 
{
@@ -238,15 +234,14 @@ public class PostgresExecutor {
 
     public Mono<Integer> executeCount(Function<DSLContext, 
Mono<Record1<Integer>>> queryFunction) {
         return 
Mono.from(metricFactory.decoratePublisherWithTimerMetric("postgres-execution",
-            Mono.usingWhen(getConnection(domain),
+            usingConnection(
                 connection -> dslContext(connection)
                     .flatMap(queryFunction)
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
                     .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
                         .filter(preparedStatementConflictException()))
-                    .map(Record1::value1),
-                jamesPostgresConnectionFactory::closeConnection)));
+                    .map(Record1::value1))));
     }
 
     public Mono<Boolean> executeExists(Function<DSLContext, 
SelectConditionStep<?>> queryFunction) {
@@ -256,19 +251,18 @@ public class PostgresExecutor {
 
     public Mono<Integer> executeReturnAffectedRowsCount(Function<DSLContext, 
Mono<Integer>> queryFunction) {
         return 
Mono.from(metricFactory.decoratePublisherWithTimerMetric("postgres-execution",
-            Mono.usingWhen(getConnection(domain),
+            usingConnection(
                 connection -> dslContext(connection)
                     .flatMap(queryFunction)
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
                     .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
-                        .filter(preparedStatementConflictException())),
-                jamesPostgresConnectionFactory::closeConnection)));
+                        .filter(preparedStatementConflictException())))));
     }
 
     public <T> Mono<T> executeTransaction(Function<DSLContext, Mono<T>> 
transactionFunction) {
         return 
Mono.from(metricFactory.decoratePublisherWithTimerMetric("postgres-transaction-execution",
-            Mono.usingWhen(getConnection(domain),
+            usingConnection(this::cancelRunningQueryThenRollbackAndRelease,
                 connection -> Mono.from(connection.beginTransaction())
                     .then(dslContext(connection)
                         .flatMap(transactionFunction)
@@ -277,8 +271,7 @@ public class PostgresExecutor {
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
                     .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
-                        .filter(preparedStatementConflictException())),
-                jamesPostgresConnectionFactory::closeConnection)));
+                        .filter(preparedStatementConflictException())))));
     }
 
     public JamesPostgresConnectionFactory connectionFactory() {
@@ -290,6 +283,44 @@ public class PostgresExecutor {
         jamesPostgresConnectionFactory.close().block();
     }
 
+    private <T> Mono<T> usingConnection(Function<Connection, Mono<T>> closure) 
{
+        return usingConnection(this::cancelRunningQueryAndRelease, closure);
+    }
+
+    /**
+     * Releasing the connection upon cancellation is not enough: see {@link 
#cancelRunningQuery(Connection)}.
+     */
+    private <T> Mono<T> usingConnection(Function<Connection, Mono<Void>> 
onCancel, Function<Connection, Mono<T>> closure) {
+        return Mono.usingWhen(getConnection(domain),
+            closure,
+            jamesPostgresConnectionFactory::closeConnection,
+            (connection, error) -> 
jamesPostgresConnectionFactory.closeConnection(connection),
+            onCancel);
+    }
+
+    private <T> Flux<T> usingConnectionMany(Function<Connection, Flux<T>> 
closure) {
+        return Flux.usingWhen(getConnection(domain),
+            closure,
+            jamesPostgresConnectionFactory::closeConnection,
+            (connection, error) -> 
jamesPostgresConnectionFactory.closeConnection(connection),
+            this::cancelRunningQueryAndRelease);
+    }
+
+    private Mono<Void> cancelRunningQueryAndRelease(Connection connection) {
+        return cancelRunningQuery(connection)
+            .then(jamesPostgresConnectionFactory.closeConnection(connection));
+    }
+
+    private Mono<Void> cancelRunningQueryThenRollbackAndRelease(Connection 
connection) {
+        return cancelRunningQuery(connection)
+            .then(Mono.from(connection.rollbackTransaction())
+                .onErrorResume(e -> {
+                    LOGGER.warn("Failed to rollback the cancelled Postgres 
transaction", e);
+                    return Mono.empty();
+                }))
+            .then(jamesPostgresConnectionFactory.closeConnection(connection));
+    }
+
     private <T> Mono<T> handleTimeout(Connection connection, TimeoutException 
timeoutException) {
         LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, timeoutException);
         return cancelRunningQuery(connection)
@@ -297,9 +328,11 @@ public class PostgresExecutor {
     }
 
     /**
-     * Cancelling the reactive pipeline does not stop the query on the 
Postgres server side: the connection stays busy
-     * until the query completes, and is handed back to the pool in that 
state. Asking Postgres to cancel the running query
-     * ensures the connection is quickly usable again.
+     * Cancelling the reactive pipeline (timeout, or a downstream 
short-circuit like `any`, `next`, `take`...) does not stop
+     * the query on the Postgres server side: the connection stays busy until 
the query completes, and is handed back to
+     * the pool in that state. Asking Postgres to cancel the running query 
ensures the connection is quickly usable again.
+     * <p>
+     * Asking Postgres to cancel a query that already completed is a no-op.
      */
     private Mono<Void> cancelRunningQuery(Connection connection) {
         return unwrapPostgresqlConnection(connection)
diff --git 
a/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExecutorTimeoutTest.java
 
b/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExecutorTimeoutTest.java
index be40e6f4b9..594a349bae 100644
--- 
a/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExecutorTimeoutTest.java
+++ 
b/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExecutorTimeoutTest.java
@@ -53,6 +53,7 @@ class PostgresExecutorTimeoutTest {
     private static final int ROW_COUNT = 3;
     private static final int ROW_INSERTED_BY_TIMED_OUT_TRANSACTION = 42;
     private static final Duration WELL_BEFORE_THE_DATABASE_SLEEP_ENDS = 
Duration.ofSeconds(10);
+    private static final Duration BEFORE_THE_JOOQ_REACTIVE_TIMEOUT = 
Duration.ofMillis(100);
     private static final Table<Record> TABLE = DSL.table("timeout_test");
     private static final Field<Integer> ID = DSL.field("id", 
SQLDataType.INTEGER);
 
@@ -141,6 +142,41 @@ class PostgresExecutorTimeoutTest {
         assertThat(ids).containsExactly(1, 2, 3);
     }
 
+    @Test
+    void connectionShouldBeUsableRightAfterCancellingExecuteRows() {
+        sleepOnTheDatabaseSide()
+            .take(BEFORE_THE_JOOQ_REACTIVE_TIMEOUT)
+            .collectList()
+            .block();
+
+        
assertThat(readIds().block(WELL_BEFORE_THE_DATABASE_SLEEP_ENDS)).containsExactly(1,
 2, 3);
+    }
+
+    @Test
+    void connectionShouldBeUsableRightAfterCancellingExecuteRow() {
+        postgresExecutor.executeRow(dslContext -> 
Mono.from(dslContext.select(DSL.field("pg_sleep(" + 
LONGER_THAN_TIMEOUT_IN_SECONDS.toSeconds() + ")"))))
+            .take(BEFORE_THE_JOOQ_REACTIVE_TIMEOUT)
+            .blockOptional();
+
+        
assertThat(readIds().block(WELL_BEFORE_THE_DATABASE_SLEEP_ENDS)).containsExactly(1,
 2, 3);
+    }
+
+    @Test
+    void 
executeTransactionShouldRollbackAndLeaveTheConnectionUsableRightAfterACancellation()
 {
+        postgresExecutor.executeTransaction(dslContext -> 
Mono.from(dslContext.insertInto(TABLE, 
ID).values(ROW_INSERTED_BY_TIMED_OUT_TRANSACTION))
+                .then(Mono.from(dslContext.select(DSL.field("pg_sleep(" + 
LONGER_THAN_TIMEOUT_IN_SECONDS.toSeconds() + ")")))))
+            .take(BEFORE_THE_JOOQ_REACTIVE_TIMEOUT)
+            .blockOptional();
+
+        
assertThat(readIds().block(WELL_BEFORE_THE_DATABASE_SLEEP_ENDS)).containsExactly(1,
 2, 3);
+    }
+
+    private Mono<List<Integer>> readIds() {
+        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(ID).from(TABLE).orderBy(ID)))
+            .map(record -> record.get(ID))
+            .collectList();
+    }
+
     private Flux<Record> sleepOnTheDatabaseSide() {
         return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(DSL.field("pg_sleep(" + 
LONGER_THAN_TIMEOUT_IN_SECONDS.toSeconds() + ")"))));
     }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to