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


The following commit(s) were added to refs/heads/master by this push:
     new 51867068cf [FIX] PG pagination: buffer next page locally (#3215)
51867068cf is described below

commit 51867068cf03c006836274f9baa8451a99ee2cde
Author: Benoit TELLIER <[email protected]>
AuthorDate: Tue Sep 29 12:05:50 2026 +0200

    [FIX] PG pagination: buffer next page locally (#3215)
---
 .../backends/postgres/utils/PostgresExecutor.java  |   5 +-
 .../backends/postgres/PostgresExecutorTest.java    | 114 +++++++++++++++++++++
 2 files changed, 118 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 c8e6cd8986..baee1f95e6 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
@@ -198,7 +198,10 @@ public class PostgresExecutor {
 
     private Mono<List<Record>> executePage(BiFunction<DSLContext, 
Optional<Record>, SelectLimitStep<? extends Record>> pageQuery, 
Optional<Record> lastRecord, int pageSize) {
         return executeRows(dslContext -> Flux.from(pageQuery.apply(dslContext, 
lastRecord).limit(pageSize)))
-            .collectList();
+            .collectList()
+            // expand subscribes to the next page without downstream demand, 
and collectList only requests upon demand:
+            // cache forces the page to be read right away so that the 
connection is released, and replays it later
+            .cache();
     }
 
     public Flux<Record> executeDeleteAndReturnList(Function<DSLContext, 
DeleteResultStep<Record>> queryFunction) {
diff --git 
a/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExecutorTest.java
 
b/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExecutorTest.java
new file mode 100644
index 0000000000..3954bba3f7
--- /dev/null
+++ 
b/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExecutorTest.java
@@ -0,0 +1,114 @@
+/****************************************************************
+ * 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 java.time.Duration;
+import java.util.List;
+import java.util.Optional;
+import java.util.stream.IntStream;
+
+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.impl.DSL;
+import org.jooq.impl.SQLDataType;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import com.google.common.collect.ImmutableList;
+
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+class PostgresExecutorTest {
+    private static final Duration JOOQ_REACTIVE_TIMEOUT = 
Duration.ofSeconds(1);
+    private static final Field<Integer> ID = DSL.field("id", 
SQLDataType.INTEGER.notNull());
+    private static final int PAGE_SIZE = 10;
+    private static final int ROW_COUNT = 25;
+
+    @RegisterExtension
+    static PostgresExtension postgresExtension = PostgresExtension.empty();
+
+    private PostgresExecutor postgresExecutor;
+
+    @BeforeEach
+    void beforeEach() {
+        PostgresConfiguration configuration = PostgresConfiguration.builder()
+            .username("james")
+            .password("secret")
+            .jooqReactiveTimeout(Optional.of(JOOQ_REACTIVE_TIMEOUT))
+            .build();
+        postgresExecutor = new PostgresExecutor.Factory(
+            new PoolBackedPostgresConnectionFactory(RowLevelSecurity.DISABLED, 
postgresExtension.getConnectionFactory()),
+            configuration,
+            new RecordingMetricFactory())
+            .create();
+
+        postgresExecutor.executeVoid(dslContext -> 
Mono.from(dslContext.createTableIfNotExists("paginated")
+                .column(ID)
+                .constraints(DSL.constraint().primaryKey(ID))))
+            .block();
+        Flux.fromStream(IntStream.range(0, ROW_COUNT).boxed())
+            .concatMap(id -> postgresExecutor.executeVoid(dslContext -> 
Mono.from(dslContext.insertInto(DSL.table("paginated"), ID)
+                .values(id))))
+            .then()
+            .block();
+    }
+
+    @AfterEach
+    void afterEach() {
+        postgresExecutor.executeVoid(dslContext -> 
Mono.from(dslContext.dropTableIfExists("paginated")))
+            .block();
+    }
+
+    @Test
+    void executeRowsPaginatedShouldReturnAllRows() {
+        List<Integer> ids = postgresExecutor.executeRowsPaginated((dslContext, 
lastRecord) -> dslContext.select(ID)
+                    .from(DSL.table("paginated"))
+                    .where(lastRecord.map(record -> 
ID.greaterThan(record.get(ID))).orElseGet(DSL::noCondition))
+                    .orderBy(ID), PAGE_SIZE)
+            .map(record -> record.get(ID))
+            .collectList()
+            .block();
+
+        assertThat(ids).containsExactlyElementsOf(IntStream.range(0, 
ROW_COUNT).boxed().collect(ImmutableList.toImmutableList()));
+    }
+
+    @Test
+    void 
executeRowsPaginatedShouldNotTimeoutWhenConsumerIsSlowerThanTheReactiveTimeout()
 {
+        // Consuming a page takes 2 seconds while the reactive timeout is 1 
second: the next page must not wait
+        // for downstream demand while holding its connection
+        List<Integer> ids = postgresExecutor.executeRowsPaginated((dslContext, 
lastRecord) -> dslContext.select(ID)
+                    .from(DSL.table("paginated"))
+                    .where(lastRecord.map(record -> 
ID.greaterThan(record.get(ID))).orElseGet(DSL::noCondition))
+                    .orderBy(ID), PAGE_SIZE)
+            .map(record -> record.get(ID))
+            .delayElements(Duration.ofMillis(200))
+            .collectList()
+            .block();
+
+        assertThat(ids).containsExactlyElementsOf(IntStream.range(0, 
ROW_COUNT).boxed().collect(ImmutableList.toImmutableList()));
+    }
+}


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

Reply via email to