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 ac241179b804fa36bd43f7e099530dec09f564a2
Author: Quan Tran <[email protected]>
AuthorDate: Tue Sep 29 14:22:32 2026 +0700

    JAMES-2586 Cancel the Postgres query upon jOOQ reactive timeout in 
executeTransaction
    
    `executeTransaction`, introduced by #3203, only logged the jOOQ reactive
    timeout. The timed out query kept running on the Postgres server, and the
    error only reached the caller once the connection was released, that is
    once the query completed on its own: the caller and the pooled
    connection stayed blocked for as long as the query ran.
    
    Handle its timeout like the other `PostgresExecutor` methods: cancel the
    running query through r2dbc-postgresql `cancelRequest()` before
    propagating the error and releasing the connection. The transaction is
    rolled back upon connection release.
    
    `PostgresExecutorTimeoutTest` covers it with a transaction inserting a
    row then sleeping on the database side, on a single connection pool:
    the timeout is surfaced right away, the connection is usable right after,
    and the row is not committed.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../backends/postgres/utils/PostgresExecutor.java      |  2 +-
 .../backends/postgres/PostgresExecutorTimeoutTest.java | 18 ++++++++++++++++++
 2 files changed, 19 insertions(+), 1 deletion(-)

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 dc7db5c768..b5fbd1c128 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
@@ -275,7 +275,7 @@ public class PostgresExecutor {
                         .flatMap(result -> 
Mono.from(connection.commitTransaction()).thenReturn(result))
                         .onErrorResume(throwable -> 
Mono.from(connection.rollbackTransaction()).then(Mono.error(throwable))))
                     .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)));
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 a8fbf852f8..be40e6f4b9 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
@@ -51,6 +51,8 @@ class PostgresExecutorTimeoutTest {
     private static final Duration LONGER_THAN_TIMEOUT_IN_SECONDS = 
Duration.ofSeconds(60);
     private static final int SINGLE_CONNECTION_POOL = 1;
     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 Table<Record> TABLE = DSL.table("timeout_test");
     private static final Field<Integer> ID = DSL.field("id", 
SQLDataType.INTEGER);
 
@@ -123,6 +125,22 @@ class PostgresExecutorTimeoutTest {
         assertThat(ids).containsExactly(1, 2, 3);
     }
 
+    @Test
+    void 
executeTransactionShouldRollbackAndLeaveTheConnectionUsableRightAfterATimeout() 
{
+        assertThatThrownBy(() -> 
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() + ")")))))
+                .block(WELL_BEFORE_THE_DATABASE_SLEEP_ENDS))
+            .hasCauseInstanceOf(TimeoutException.class)
+            .hasMessageContaining("Did not observe any item or terminal 
signal");
+
+        List<Integer> ids = postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(ID).from(TABLE).orderBy(ID)))
+            .map(record -> record.get(ID))
+            .collectList()
+            .block(WELL_BEFORE_THE_DATABASE_SLEEP_ENDS);
+
+        assertThat(ids).containsExactly(1, 2, 3);
+    }
+
     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