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]
