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 b19553cd8be6a247a315b687d45f6b926aff6b02
Author: Quan Tran <[email protected]>
AuthorDate: Sat Sep 26 21:45:07 2026 +0700

    JAMES-2586 Reproduce a pooled Postgres connection staying busy after a 
timed out query
    
    Cancelling the reactive pipeline upon `jooq.reactive.timeout` does not
    stop the query on the Postgres server: the connection keeps processing
    it and is handed back to the pool in that state. The next query
    acquiring that connection queues behind the still running statement
    and times out in turn, which cascades into the "everything needing
    Postgres fails, then sometimes recovers" behaviour of
    linagora/tmail-backend#1599.
    
    PostgresExecutorTimeoutTest reproduces it with a pool of a single
    connection and a 60 seconds `pg_sleep`:
     - a query stalling on the database side must time out,
     - right after the timeout, a query on that connection must succeed.
    The latter fails before the fix.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../postgres/PostgresExecutorTimeoutTest.java      | 129 +++++++++++++++++++++
 1 file changed, 129 insertions(+)

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
new file mode 100644
index 0000000000..a8fbf852f8
--- /dev/null
+++ 
b/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExecutorTimeoutTest.java
@@ -0,0 +1,129 @@
+/****************************************************************
+ * Licensed to the Apache Software Foundation (ASF) under one   *
+ * or more contributor license agreements.  See the NOTICE file *
+ * distributed with this work for additional information        *
+ * regarding copyright ownership.  The ASF licenses this file   *
+ * to you under the Apache License, Version 2.0 (the            *
+ * "License"); you may not use this file except in compliance   *
+ * with the License.  You may obtain a copy of the License at   *
+ *                                                              *
+ *   http://www.apache.org/licenses/LICENSE-2.0                 *
+ *                                                              *
+ * Unless required by applicable law or agreed to in writing,   *
+ * software distributed under the License is distributed on an  *
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY       *
+ * KIND, either express or implied.  See the License for the    *
+ * specific language governing permissions and limitations      *
+ * under the License.                                           *
+ ****************************************************************/
+
+package org.apache.james.backends.postgres;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+import java.time.Duration;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.TimeoutException;
+
+import org.apache.james.backends.postgres.utils.JamesPostgresConnectionFactory;
+import 
org.apache.james.backends.postgres.utils.PoolBackedPostgresConnectionFactory;
+import org.apache.james.backends.postgres.utils.PostgresExecutor;
+import org.apache.james.metrics.tests.RecordingMetricFactory;
+import org.jooq.Field;
+import org.jooq.Record;
+import org.jooq.Table;
+import org.jooq.impl.DSL;
+import org.jooq.impl.SQLDataType;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+class PostgresExecutorTimeoutTest {
+    private static final Duration JOOQ_REACTIVE_TIMEOUT = 
Duration.ofMillis(500);
+    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 Table<Record> TABLE = DSL.table("timeout_test");
+    private static final Field<Integer> ID = DSL.field("id", 
SQLDataType.INTEGER);
+
+    @RegisterExtension
+    static PostgresExtension postgresExtension = PostgresExtension.empty();
+
+    private static JamesPostgresConnectionFactory singleConnectionFactory;
+    private static PostgresExecutor postgresExecutor;
+
+    @BeforeAll
+    static void beforeAll() {
+        PostgresConfiguration extensionConfiguration = 
postgresExtension.getPostgresConfiguration();
+        PostgresConfiguration shortTimeoutConfiguration = 
PostgresConfiguration.builder()
+            .databaseName(extensionConfiguration.getDatabaseName())
+            .databaseSchema(extensionConfiguration.getDatabaseSchema())
+            .host(extensionConfiguration.getHost())
+            .port(extensionConfiguration.getPort())
+            
.username(extensionConfiguration.getDefaultCredential().getUsername())
+            
.password(extensionConfiguration.getDefaultCredential().getPassword())
+            
.byPassRLSUser(extensionConfiguration.getByPassRLSCredential().getUsername())
+            
.byPassRLSPassword(extensionConfiguration.getByPassRLSCredential().getPassword())
+            .rowLevelSecurityEnabled(false)
+            .jooqReactiveTimeout(Optional.of(JOOQ_REACTIVE_TIMEOUT))
+            .build();
+        singleConnectionFactory = new 
PoolBackedPostgresConnectionFactory(RowLevelSecurity.DISABLED,
+            SINGLE_CONNECTION_POOL, SINGLE_CONNECTION_POOL, 
postgresExtension.getConnectionFactory());
+        postgresExecutor = new 
PostgresExecutor.Factory(singleConnectionFactory, shortTimeoutConfiguration, 
new RecordingMetricFactory())
+            .create();
+    }
+
+    @AfterAll
+    static void afterAll() {
+        singleConnectionFactory.close().block();
+    }
+
+    @BeforeEach
+    void beforeEach() {
+        PostgresExecutor setUpExecutor = 
postgresExtension.getDefaultPostgresExecutor();
+        setUpExecutor.executeVoid(dslContext -> 
Mono.from(dslContext.createTableIfNotExists(TABLE)
+                .column(ID)))
+            .block();
+        Flux.range(1, ROW_COUNT)
+            .concatMap(id -> setUpExecutor.executeVoid(dslContext -> 
Mono.from(dslContext.insertInto(TABLE, ID).values(id))))
+            .blockLast();
+    }
+
+    @AfterEach
+    void afterEach() {
+        postgresExtension.getDefaultPostgresExecutor()
+            .executeVoid(dslContext -> 
Mono.from(dslContext.dropTableIfExists(TABLE)))
+            .block();
+    }
+
+    @Test
+    void 
executeRowsShouldTimeoutWhenTheDatabaseDoesNotAnswerWhileRowsAreAwaited() {
+        assertThatThrownBy(() -> 
sleepOnTheDatabaseSide().collectList().block())
+            .hasCauseInstanceOf(TimeoutException.class);
+    }
+
+    @Test
+    void connectionShouldBeUsableRightAfterATimeout() {
+        assertThatThrownBy(() -> 
sleepOnTheDatabaseSide().collectList().block())
+            .hasCauseInstanceOf(TimeoutException.class);
+
+        List<Integer> ids = postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(ID).from(TABLE).orderBy(ID)))
+            .map(record -> record.get(ID))
+            .collectList()
+            .block();
+
+        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