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 0baecc4b26bc89cdc74e475b13df38dd58770ce5
Author: Quan Tran <[email protected]>
AuthorDate: Sat Sep 26 21:41:22 2026 +0700

    JAMES-2586 Regression test for throttled full re-indexing on Postgres
    
    A throttled full re-indexing on the Postgres app
    (`/mailboxes?task=reIndex&messagesPerSecond=N`) used to fail after exactly
    `jooq.reactive.timeout` (10 seconds by default) with:
    
      PostgresExecutor - Time out executing Postgres query. May need to check
      either jOOQ reactive issue or Postgres DB performance.
      java.util.concurrent.TimeoutException: Did not observe any item or
      terminal signal within 10000ms in 'flatMapMany'
    
    The re-indexing lists every mailbox with a mailbox concurrency of 1 and
    throttles the messages of each mailbox, so a streamed listing query was
    kept open while the first mailbox was slowly re-indexed, and the reactive
    timeout fired while Postgres only waited for downstream demand.
    
    Paginated listings (see "[FIX] PG: Generalise streaming with paging") no
    longer hold a query open during the consumption. This test locks in the
    reported scenario on a Guice Postgres server, with a 2 seconds timeout
    and 1 message per second over two mailboxes.
    
    `PostgresExtension.withJooqReactiveTimeout` allows tests to shorten the
    timeout.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../james/backends/postgres/PostgresExtension.java |   9 +-
 .../PostgresReIndexingIntegrationTest.java         | 151 +++++++++++++++++++++
 2 files changed, 159 insertions(+), 1 deletion(-)

diff --git 
a/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExtension.java
 
b/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExtension.java
index 84143e5b80..1db077b0db 100644
--- 
a/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExtension.java
+++ 
b/backends-common/postgres/src/test/java/org/apache/james/backends/postgres/PostgresExtension.java
@@ -87,11 +87,13 @@ public class PostgresExtension implements 
GuiceModuleTestExtension {
     }
 
     public static final PoolSize DEFAULT_POOL_SIZE = PoolSize.SMALL;
+    public static final Duration DEFAULT_JOOQ_REACTIVE_TIMEOUT = 
Duration.ofSeconds(20L);
     public static PostgreSQLContainer<?> PG_CONTAINER = 
DockerPostgresSingleton.SINGLETON;
     private final PostgresDataDefinition postgresDataDefinition;
     private final RowLevelSecurity rowLevelSecurity;
     private final PostgresFixture.Database selectedDatabase;
     private PoolSize poolSize;
+    private Duration jooqReactiveTimeout = DEFAULT_JOOQ_REACTIVE_TIMEOUT;
     private PostgresConfiguration postgresConfiguration;
     private PostgresExecutor defaultPostgresExecutor;
     private PostgresExecutor byPassRLSPostgresExecutor;
@@ -100,6 +102,11 @@ public class PostgresExtension implements 
GuiceModuleTestExtension {
     private PostgresExecutor.Factory executorFactory;
     private PostgresTableManager postgresTableManager;
 
+    public PostgresExtension withJooqReactiveTimeout(Duration 
jooqReactiveTimeout) {
+        this.jooqReactiveTimeout = jooqReactiveTimeout;
+        return this;
+    }
+
     public void pause() {
         
PG_CONTAINER.getDockerClient().pauseContainerCmd(PG_CONTAINER.getContainerId())
             .exec();
@@ -161,7 +168,7 @@ public class PostgresExtension implements 
GuiceModuleTestExtension {
             .byPassRLSUser(DEFAULT_DATABASE.dbUser())
             .byPassRLSPassword(DEFAULT_DATABASE.dbPassword())
             
.rowLevelSecurityEnabled(rowLevelSecurity.isRowLevelSecurityEnabled())
-            .jooqReactiveTimeout(Optional.of(Duration.ofSeconds(20L)))
+            .jooqReactiveTimeout(Optional.of(jooqReactiveTimeout))
             .build();
 
         Function<PostgresConfiguration.Credential, 
PostgresqlConnectionConfiguration> postgresqlConnectionConfigurationFunction = 
credential ->
diff --git 
a/server/protocols/webadmin-integration-test/postgres-webadmin-integration-test/src/test/java/org/apache/james/webadmin/integration/postgres/PostgresReIndexingIntegrationTest.java
 
b/server/protocols/webadmin-integration-test/postgres-webadmin-integration-test/src/test/java/org/apache/james/webadmin/integration/postgres/PostgresReIndexingIntegrationTest.java
new file mode 100644
index 0000000000..a111cf8030
--- /dev/null
+++ 
b/server/protocols/webadmin-integration-test/postgres-webadmin-integration-test/src/test/java/org/apache/james/webadmin/integration/postgres/PostgresReIndexingIntegrationTest.java
@@ -0,0 +1,151 @@
+/****************************************************************
+ * 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.webadmin.integration.postgres;
+
+import static io.restassured.RestAssured.given;
+import static io.restassured.RestAssured.with;
+import static 
org.apache.james.data.UsersRepositoryModuleChooser.Implementation.DEFAULT;
+import static org.hamcrest.Matchers.is;
+
+import java.time.Duration;
+
+import org.apache.james.FakeMessageSearchIndex;
+import org.apache.james.GuiceJamesServer;
+import org.apache.james.JamesServerBuilder;
+import org.apache.james.JamesServerExtension;
+import org.apache.james.PostgresJamesConfiguration;
+import org.apache.james.PostgresJamesServerMain;
+import org.apache.james.SearchConfiguration;
+import org.apache.james.backends.postgres.PostgresExtension;
+import org.apache.james.core.Username;
+import org.apache.james.mailbox.MailboxSession;
+import org.apache.james.mailbox.MessageManager.AppendCommand;
+import org.apache.james.mailbox.model.Mailbox;
+import org.apache.james.mailbox.model.MailboxId;
+import org.apache.james.mailbox.model.MailboxPath;
+import org.apache.james.mailbox.store.mail.model.MailboxMessage;
+import org.apache.james.mailbox.store.search.ListeningMessageSearchIndex;
+import org.apache.james.modules.MailboxProbeImpl;
+import org.apache.james.probe.DataProbe;
+import org.apache.james.utils.DataProbeImpl;
+import org.apache.james.utils.WebAdminGuiceProbe;
+import org.apache.james.webadmin.WebAdminUtils;
+import org.apache.james.webadmin.routes.TasksRoutes;
+import org.apache.mailbox.tools.indexer.FullReindexingTask;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import io.restassured.RestAssured;
+import reactor.core.publisher.Mono;
+
+/**
+ * Regression test for the "Time out executing Postgres query" error reported 
when running a throttled full re-indexing
+ * against the Postgres app: listing mailboxes and messages must not hold a 
Postgres query open while the messages are
+ * slowly re-indexed, otherwise the jOOQ reactive timeout fails the task.
+ */
+class PostgresReIndexingIntegrationTest {
+    private static class AcceptingMessageSearchIndex extends 
FakeMessageSearchIndex {
+        @Override
+        public Mono<Void> add(MailboxSession session, Mailbox mailbox, 
MailboxMessage message) {
+            return Mono.empty();
+        }
+
+        @Override
+        public Mono<Void> deleteAll(MailboxSession session, MailboxId 
mailboxId) {
+            return Mono.empty();
+        }
+
+        @Override
+        public void postReindexing() {
+
+        }
+    }
+
+    private static final Duration JOOQ_REACTIVE_TIMEOUT = 
Duration.ofSeconds(2);
+    private static final int ONE_MESSAGE_PER_SECOND = 1;
+    private static final int INBOX_MESSAGE_COUNT = 4;
+    private static final int SENT_MESSAGE_COUNT = 1;
+    private static final String DOMAIN = "domain.tld";
+    private static final Username BOB = Username.of("bob@" + DOMAIN);
+    private static final String PASSWORD = "password";
+    private static final MailboxPath BOB_INBOX = MailboxPath.inbox(BOB);
+    private static final MailboxPath BOB_SENT = MailboxPath.forUser(BOB, 
"Sent");
+
+    @RegisterExtension
+    static JamesServerExtension jamesServerExtension = new 
JamesServerBuilder<PostgresJamesConfiguration>(tmpDir ->
+        PostgresJamesConfiguration.builder()
+            .workingDirectory(tmpDir)
+            .configurationFromClasspath()
+            .searchConfiguration(SearchConfiguration.scanning())
+            .usersRepository(DEFAULT)
+            .eventBusImpl(PostgresJamesConfiguration.EventBusImpl.IN_MEMORY)
+            .build())
+        
.extension(PostgresExtension.empty().withJooqReactiveTimeout(JOOQ_REACTIVE_TIMEOUT))
+        .server(configuration -> 
PostgresJamesServerMain.createServer(configuration)
+            .overrideWith(binder -> 
binder.bind(ListeningMessageSearchIndex.class).toInstance(new 
AcceptingMessageSearchIndex())))
+        .build();
+
+    private MailboxProbeImpl mailboxProbe;
+
+    @BeforeEach
+    void setUp(GuiceJamesServer guiceJamesServer) throws Exception {
+        DataProbe dataProbe = guiceJamesServer.getProbe(DataProbeImpl.class);
+        mailboxProbe = guiceJamesServer.getProbe(MailboxProbeImpl.class);
+        WebAdminGuiceProbe webAdminGuiceProbe = 
guiceJamesServer.getProbe(WebAdminGuiceProbe.class);
+        RestAssured.requestSpecification = 
WebAdminUtils.buildRequestSpecification(webAdminGuiceProbe.getWebAdminPort())
+            .build();
+
+        dataProbe.addDomain(DOMAIN);
+        dataProbe.addUser(BOB.asString(), PASSWORD);
+        mailboxProbe.createMailbox(BOB_INBOX);
+        mailboxProbe.createMailbox(BOB_SENT);
+        appendMessages(BOB_INBOX, INBOX_MESSAGE_COUNT);
+        appendMessages(BOB_SENT, SENT_MESSAGE_COUNT);
+    }
+
+    @Test
+    void 
throttledFullReIndexingShouldNotTimeoutWhenIndexingAMailboxTakesLongerThanTheJooqReactiveTimeout()
 {
+        String taskId = with()
+            .queryParam("task", "reIndex")
+            .queryParam("messagesPerSecond", ONE_MESSAGE_PER_SECOND)
+            .post("/mailboxes")
+            .jsonPath()
+            .get("taskId");
+
+        given()
+            .basePath(TasksRoutes.BASE)
+        .when()
+            .get(taskId + "/await")
+        .then()
+            .body("status", is("completed"))
+            .body("type", is(FullReindexingTask.FULL_RE_INDEXING.asString()))
+            .body("additionalInformation.successfullyReprocessedMailCount", 
is(INBOX_MESSAGE_COUNT + SENT_MESSAGE_COUNT))
+            .body("additionalInformation.failedReprocessedMailCount", is(0))
+            .body("additionalInformation.runningOptions.messagesPerSecond", 
is(ONE_MESSAGE_PER_SECOND));
+    }
+
+    private void appendMessages(MailboxPath mailboxPath, int count) throws 
Exception {
+        for (int i = 0; i < count; i++) {
+            mailboxProbe.appendMessage(BOB.asString(), mailboxPath,
+                AppendCommand.builder().build("header: value\r\n\r\nbody " + 
i));
+        }
+    }
+}


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

Reply via email to