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]
