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

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

commit 9f898f0796a5fc29f7fdaa0a3572a3d137df2980
Author: Quan Tran <[email protected]>
AuthorDate: Sat Sep 26 21:45:56 2026 +0700

    JAMES-2586 Cancel the Postgres query upon jOOQ reactive timeout
    
    Upon TimeoutException, `PostgresExecutor` now asks Postgres to cancel
    the running query through r2dbc-postgresql `cancelRequest()` (the pooled
    connection being unwrapped to the driver connection) before propagating
    the error and releasing the connection. The cancel request goes through
    a dedicated connection, so it does not depend on the busy one.
    
    The connection thus becomes usable again within milliseconds instead of
    staying busy until the query completes on its own, which prevents a
    single slow query from cascading timeouts onto unrelated queries.
    
    This applies to every execution method, including each page of
    `executeRowsPaginated`.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../backends/postgres/utils/PostgresExecutor.java  | 45 +++++++++++++++++++---
 1 file changed, 39 insertions(+), 6 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 baee1f95e6..dc7db5c768 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
@@ -51,8 +51,10 @@ import org.slf4j.LoggerFactory;
 
 import com.google.common.annotations.VisibleForTesting;
 
+import io.r2dbc.postgresql.api.PostgresqlConnection;
 import io.r2dbc.spi.Connection;
 import io.r2dbc.spi.R2dbcBadGrammarException;
+import io.r2dbc.spi.Wrapped;
 import reactor.core.publisher.Flux;
 import reactor.core.publisher.Mono;
 import reactor.util.retry.Retry;
@@ -121,7 +123,7 @@ public class PostgresExecutor {
                 connection -> dslContext(connection)
                     .flatMap(queryFunction)
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
-                    .doOnError(TimeoutException.class, e -> 
LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e))
+                    .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
                         .filter(preparedStatementConflictException()))
                     .then(),
@@ -153,7 +155,7 @@ public class PostgresExecutor {
                     Flux<Record> recordFlux = dslContext(connection)
                         .flatMapMany(queryFunction)
                         
.timeout(postgresConfiguration.getJooqReactiveTimeout())
-                        .doOnError(TimeoutException.class, e -> 
LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e))
+                        .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                         .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, 
MIN_BACKOFF)
                             .filter(preparedStatementConflictException()));
 
@@ -210,7 +212,7 @@ public class PostgresExecutor {
                 connection -> dslContext(connection)
                     .flatMapMany(queryFunction)
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
-                    .doOnError(TimeoutException.class, e -> 
LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e))
+                    .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
                         .filter(preparedStatementConflictException())),
                 jamesPostgresConnectionFactory::closeConnection)));
@@ -222,7 +224,7 @@ public class PostgresExecutor {
                 connection -> dslContext(connection)
                     .flatMap(queryFunction.andThen(Mono::from))
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
-                    .doOnError(TimeoutException.class, e -> 
LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e))
+                    .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
                         .filter(preparedStatementConflictException())),
                 jamesPostgresConnectionFactory::closeConnection)));
@@ -240,7 +242,7 @@ public class PostgresExecutor {
                 connection -> dslContext(connection)
                     .flatMap(queryFunction)
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
-                    .doOnError(TimeoutException.class, e -> 
LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e))
+                    .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
                         .filter(preparedStatementConflictException()))
                     .map(Record1::value1),
@@ -258,7 +260,7 @@ public class PostgresExecutor {
                 connection -> dslContext(connection)
                     .flatMap(queryFunction)
                     .timeout(postgresConfiguration.getJooqReactiveTimeout())
-                    .doOnError(TimeoutException.class, e -> 
LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e))
+                    .onErrorResume(TimeoutException.class, e -> 
handleTimeout(connection, e))
                     .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF)
                         .filter(preparedStatementConflictException())),
                 jamesPostgresConnectionFactory::closeConnection)));
@@ -288,6 +290,37 @@ public class PostgresExecutor {
         jamesPostgresConnectionFactory.close().block();
     }
 
+    private <T> Mono<T> handleTimeout(Connection connection, TimeoutException 
timeoutException) {
+        LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, timeoutException);
+        return cancelRunningQuery(connection)
+            .then(Mono.error(timeoutException));
+    }
+
+    /**
+     * 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.
+     */
+    private Mono<Void> cancelRunningQuery(Connection connection) {
+        return unwrapPostgresqlConnection(connection)
+            .map(postgresqlConnection -> postgresqlConnection.cancelRequest()
+                .onErrorResume(e -> {
+                    LOGGER.warn("Failed to cancel the timed out Postgres 
query", e);
+                    return Mono.empty();
+                }))
+            .orElse(Mono.empty());
+    }
+
+    private Optional<PostgresqlConnection> 
unwrapPostgresqlConnection(Connection connection) {
+        if (connection instanceof PostgresqlConnection postgresqlConnection) {
+            return Optional.of(postgresqlConnection);
+        }
+        if (connection instanceof Wrapped<?> wrapped && wrapped.unwrap() 
instanceof PostgresqlConnection postgresqlConnection) {
+            return Optional.of(postgresqlConnection);
+        }
+        return Optional.empty();
+    }
+
     private Predicate<Throwable> preparedStatementConflictException() {
         return throwable -> throwable.getCause() instanceof 
R2dbcBadGrammarException
             && throwable.getMessage().contains("prepared statement")


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

Reply via email to