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]