This is an automated email from the ASF dual-hosted git repository. quantranhong1999 pushed a commit to branch 3.9.x in repository https://gitbox.apache.org/repos/asf/james-project.git
commit f70bfa3e599fe5381c1be410d1b3de3a0c964097 Author: Benoit TELLIER <[email protected]> AuthorDate: Sun Aug 16 23:39:21 2026 +0700 [FIX] ClearMailRepository should count only once --- .../mailrepository/file/FileMailRepository.java | 11 ++- .../mailrepository/jpa/JPAMailRepository.java | 17 ++++ .../postgres/PostgresMailRepository.java | 8 +- .../postgres/PostgresMailRepositoryContentDAO.java | 5 +- .../james/mailrepository/api/MailRepository.java | 15 ++++ .../mailrepository/MailRepositoryContract.java | 22 +++++ .../blob/BlobMailRepositoryTest.java | 8 ++ .../cassandra/CassandraMailRepository.java | 9 +- .../memory/MemoryMailRepository.java | 10 +++ .../webadmin/service/ClearMailRepositoryTask.java | 32 +++++-- .../service/ClearMailRepositoryTaskTest.java | 97 ++++++++++++++++++++++ 11 files changed, 220 insertions(+), 14 deletions(-) diff --git a/server/data/data-file/src/main/java/org/apache/james/mailrepository/file/FileMailRepository.java b/server/data/data-file/src/main/java/org/apache/james/mailrepository/file/FileMailRepository.java index a53371d291..b2e89daa5c 100644 --- a/server/data/data-file/src/main/java/org/apache/james/mailrepository/file/FileMailRepository.java +++ b/server/data/data-file/src/main/java/org/apache/james/mailrepository/file/FileMailRepository.java @@ -27,6 +27,7 @@ import java.util.Collections; import java.util.HashSet; import java.util.Iterator; import java.util.Optional; +import java.util.function.Consumer; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -293,8 +294,16 @@ public class FileMailRepository implements MailRepository, Configurable, Initial @Override public void removeAll() { + removeAll(any -> { }); + } + + @Override + public void removeAll(Consumer<MailKey> progressCallback) { listStream() - .forEach(Throwing.<MailKey>consumer(this::remove).sneakyThrow()); + .forEach(Throwing.<MailKey>consumer(key -> { + remove(key); + progressCallback.accept(key); + }).sneakyThrow()); } private void internalRemove(MailKey key) { diff --git a/server/data/data-jpa/src/main/java/org/apache/james/mailrepository/jpa/JPAMailRepository.java b/server/data/data-jpa/src/main/java/org/apache/james/mailrepository/jpa/JPAMailRepository.java index dc35e4fd72..4426013eff 100644 --- a/server/data/data-jpa/src/main/java/org/apache/james/mailrepository/jpa/JPAMailRepository.java +++ b/server/data/data-jpa/src/main/java/org/apache/james/mailrepository/jpa/JPAMailRepository.java @@ -30,6 +30,7 @@ import java.util.Map; import java.util.Objects; import java.util.Optional; import java.util.StringTokenizer; +import java.util.function.Consumer; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -77,6 +78,7 @@ import com.fasterxml.jackson.databind.node.ObjectNode; import com.github.fge.lambdas.Throwing; import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableMap; +import com.google.common.collect.UnmodifiableIterator; /** * Implementation of a MailRepository on a database via JPA. @@ -84,6 +86,7 @@ import com.google.common.collect.ImmutableMap; public class JPAMailRepository implements MailRepository, Configurable, Initializable { private static final Logger LOGGER = LoggerFactory.getLogger(JPAMailRepository.class); private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); + private static final int DELETION_BATCH_SIZE = 100; private String repositoryName; @@ -377,6 +380,20 @@ public class JPAMailRepository implements MailRepository, Configurable, Initiali } } + /** + * Deletes by batches rather than in one go, so that the caller gets to observe the progress being made. + */ + @Override + public void removeAll(Consumer<MailKey> progressCallback) throws MessagingException { + UnmodifiableIterator<List<MailKey>> batches = com.google.common.collect.Iterators.partition(list(), DELETION_BATCH_SIZE); + + while (batches.hasNext()) { + List<MailKey> batch = batches.next(); + remove(batch); + batch.forEach(progressCallback); + } + } + @Override public void removeAll() throws MessagingException { EntityManager entityManager = entityManager(); diff --git a/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepository.java b/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepository.java index cc640377a1..510df2783e 100644 --- a/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepository.java +++ b/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepository.java @@ -21,6 +21,7 @@ package org.apache.james.mailrepository.postgres; import java.util.Collection; import java.util.Iterator; +import java.util.function.Consumer; import jakarta.inject.Inject; import jakarta.mail.MessagingException; @@ -80,6 +81,11 @@ public class PostgresMailRepository implements MailRepository { @Override public void removeAll() { - postgresMailRepositoryContentDAO.removeAll(url); + postgresMailRepositoryContentDAO.removeAll(url, any -> { }); + } + + @Override + public void removeAll(Consumer<MailKey> progressCallback) { + postgresMailRepositoryContentDAO.removeAll(url, progressCallback); } } diff --git a/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepositoryContentDAO.java b/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepositoryContentDAO.java index c90b0392f6..4524e5dc67 100644 --- a/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepositoryContentDAO.java +++ b/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepositoryContentDAO.java @@ -339,9 +339,10 @@ public class PostgresMailRepositoryContentDAO { .block(); } - public void removeAll(MailRepositoryUrl url) { + public void removeAll(MailRepositoryUrl url, Consumer<MailKey> progressCallback) { listMailKeys(url) - .flatMap(mailKey -> removeReactive(mailKey, url), DEFAULT_CONCURRENCY) + .flatMap(mailKey -> removeReactive(mailKey, url).thenReturn(mailKey), DEFAULT_CONCURRENCY) + .doOnNext(progressCallback) .then() .block(); } diff --git a/server/mailrepository/mailrepository-api/src/main/java/org/apache/james/mailrepository/api/MailRepository.java b/server/mailrepository/mailrepository-api/src/main/java/org/apache/james/mailrepository/api/MailRepository.java index 7013a5e893..ed1f55b508 100644 --- a/server/mailrepository/mailrepository-api/src/main/java/org/apache/james/mailrepository/api/MailRepository.java +++ b/server/mailrepository/mailrepository-api/src/main/java/org/apache/james/mailrepository/api/MailRepository.java @@ -22,6 +22,7 @@ package org.apache.james.mailrepository.api; import java.time.Instant; import java.util.Collection; import java.util.Iterator; +import java.util.function.Consumer; import java.util.function.Predicate; import jakarta.mail.MessagingException; @@ -212,4 +213,18 @@ public interface MailRepository { */ void removeAll() throws MessagingException; + /** + * Removes all mails from this repository, notifying the caller of the mails being removed. + * + * Implementations deleting mails one by one should override this in order to report their progress: this + * allows long running deletions to be monitored without counting the remaining mails over and over, which + * is a costly operation on large repositories. + * + * Implementations relying on a bulk deletion can keep the default: the caller then only learns about the + * progress once the deletion completed. + */ + default void removeAll(Consumer<MailKey> progressCallback) throws MessagingException { + removeAll(); + } + } diff --git a/server/mailrepository/mailrepository-api/src/test/java/org/apache/james/mailrepository/MailRepositoryContract.java b/server/mailrepository/mailrepository-api/src/test/java/org/apache/james/mailrepository/MailRepositoryContract.java index 6561119efc..8ff5f5763f 100644 --- a/server/mailrepository/mailrepository-api/src/test/java/org/apache/james/mailrepository/MailRepositoryContract.java +++ b/server/mailrepository/mailrepository-api/src/test/java/org/apache/james/mailrepository/MailRepositoryContract.java @@ -269,6 +269,28 @@ public interface MailRepositoryContract { testee.removeAll(); } + @Test + default void removeAllShouldReportTheRemovedMails() throws Exception { + MailRepository testee = retrieveRepository(); + testee.store(createMail(MAIL_1)); + testee.store(createMail(MAIL_2)); + + ImmutableList.Builder<MailKey> removed = ImmutableList.builder(); + testee.removeAll(removed::add); + + assertThat(removed.build()).containsExactlyInAnyOrder(MAIL_1, MAIL_2); + } + + @Test + default void removeAllShouldReportNothingWhenEmpty() throws Exception { + MailRepository testee = retrieveRepository(); + + ImmutableList.Builder<MailKey> removed = ImmutableList.builder(); + testee.removeAll(removed::add); + + assertThat(removed.build()).isEmpty(); + } + @Test default void retrieveShouldGetStoredEmojiMail() throws Exception { MailRepository testee = retrieveRepository(); diff --git a/server/mailrepository/mailrepository-blob/src/test/java/org/apache/james/mailrepository/blob/BlobMailRepositoryTest.java b/server/mailrepository/mailrepository-blob/src/test/java/org/apache/james/mailrepository/blob/BlobMailRepositoryTest.java index fbf0f8036e..72adeae0e1 100644 --- a/server/mailrepository/mailrepository-blob/src/test/java/org/apache/james/mailrepository/blob/BlobMailRepositoryTest.java +++ b/server/mailrepository/mailrepository-blob/src/test/java/org/apache/james/mailrepository/blob/BlobMailRepositoryTest.java @@ -30,6 +30,8 @@ import org.apache.james.mailrepository.api.MailRepositoryUrl; import org.apache.james.mailrepository.api.Protocol; import org.jetbrains.annotations.NotNull; import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Disabled; +import org.junit.jupiter.api.Test; class BlobMailRepositoryTest implements MailRepositoryContract { @@ -67,4 +69,10 @@ class BlobMailRepositoryTest implements MailRepositoryContract { public MailRepository retrieveRepository(MailRepositoryPath path) { return buildBlobMailRepository(path); } + + @Test + @Disabled("This repository removes blobs rather than mails: it can not tell which of them back a mail") + @Override + public void removeAllShouldReportTheRemovedMails() { + } } \ No newline at end of file diff --git a/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepository.java b/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepository.java index 8051f1122c..15dbe63d64 100644 --- a/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepository.java +++ b/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepository.java @@ -24,6 +24,7 @@ import static org.apache.james.util.ReactorUtils.publishIfPresent; import java.util.Iterator; import java.util.Optional; +import java.util.function.Consumer; import jakarta.inject.Inject; import jakarta.mail.MessagingException; @@ -178,8 +179,14 @@ public class CassandraMailRepository implements MailRepository { @Override public void removeAll() { + removeAll(any -> { }); + } + + @Override + public void removeAll(Consumer<MailKey> progressCallback) { keysDAO.list(url) - .flatMap(this::removeAsync, DEFAULT_CONCURRENCY) + .flatMap(key -> removeAsync(key).thenReturn(key), DEFAULT_CONCURRENCY) + .doOnNext(progressCallback) .then() .block(); } diff --git a/server/mailrepository/mailrepository-memory/src/main/java/org/apache/james/mailrepository/memory/MemoryMailRepository.java b/server/mailrepository/mailrepository-memory/src/main/java/org/apache/james/mailrepository/memory/MemoryMailRepository.java index 17b53002af..d535e6bdf1 100644 --- a/server/mailrepository/mailrepository-memory/src/main/java/org/apache/james/mailrepository/memory/MemoryMailRepository.java +++ b/server/mailrepository/mailrepository-memory/src/main/java/org/apache/james/mailrepository/memory/MemoryMailRepository.java @@ -22,6 +22,7 @@ package org.apache.james.mailrepository.memory; import java.util.Iterator; import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Consumer; import jakarta.mail.MessagingException; import jakarta.mail.internet.MimeMessage; @@ -89,6 +90,15 @@ public class MemoryMailRepository implements MailRepository { mails.clear(); } + @Override + public void removeAll(Consumer<MailKey> progressCallback) { + mails.keySet() + .forEach(key -> { + mails.remove(key); + progressCallback.accept(key); + }); + } + private Mail cloneMail(Mail mail) { try { Mail newMail = mail.duplicate(); diff --git a/server/protocols/webadmin/webadmin-mailrepository/src/main/java/org/apache/james/webadmin/service/ClearMailRepositoryTask.java b/server/protocols/webadmin/webadmin-mailrepository/src/main/java/org/apache/james/webadmin/service/ClearMailRepositoryTask.java index 877935e91a..5675648df6 100644 --- a/server/protocols/webadmin/webadmin-mailrepository/src/main/java/org/apache/james/webadmin/service/ClearMailRepositoryTask.java +++ b/server/protocols/webadmin/webadmin-mailrepository/src/main/java/org/apache/james/webadmin/service/ClearMailRepositoryTask.java @@ -22,6 +22,7 @@ package org.apache.james.webadmin.service; import java.time.Clock; import java.time.Instant; import java.util.Optional; +import java.util.concurrent.atomic.AtomicLong; import jakarta.inject.Inject; @@ -31,7 +32,6 @@ import org.apache.james.mailrepository.api.MailRepositoryStore; import org.apache.james.task.Task; import org.apache.james.task.TaskExecutionDetails; import org.apache.james.task.TaskType; -import org.reactivestreams.Publisher; import com.github.fge.lambdas.Throwing; @@ -102,18 +102,24 @@ public class ClearMailRepositoryTask implements Task { private final MailRepositoryStore mailRepositoryStore; private final MailRepositoryPath mailRepositoryPath; - private long initialCount = 0; + private final AtomicLong initialCount; + private final AtomicLong removedCount; public ClearMailRepositoryTask(MailRepositoryStore mailRepositoryStore, MailRepositoryPath path) { this.mailRepositoryStore = mailRepositoryStore; this.mailRepositoryPath = path; + this.initialCount = new AtomicLong(0); + this.removedCount = new AtomicLong(0); } @Override public Result run() { - initialCount = getRemainingSize().block(); + initialCount.set(getRemainingSize().block()); try { removeAllInAllRepositories(); + // Repositories deleting in bulk do not report any progress: only their completion tells us the + // repository is now empty. + removedCount.set(initialCount.get()); return Result.COMPLETED; } catch (MailRepositoryStore.MailRepositoryStoreException e) { LOGGER.error("Encountered error while clearing repository", e); @@ -123,7 +129,8 @@ public class ClearMailRepositoryTask implements Task { private void removeAllInAllRepositories() throws MailRepositoryStore.MailRepositoryStoreException { mailRepositoryStore.getByPath(mailRepositoryPath) - .forEach(Throwing.consumer(MailRepository::removeAll).sneakyThrow()); + .forEach(Throwing.<MailRepository>consumer(repository -> repository.removeAll(any -> removedCount.incrementAndGet())) + .sneakyThrow()); } @Override @@ -135,14 +142,21 @@ public class ClearMailRepositoryTask implements Task { return mailRepositoryPath; } + /** + * Progress is tracked in memory rather than by counting the remaining mails: counting is a full partition + * scan on distributed implementations, and running it upon each and every progress update is prohibitive on + * large repositories. Critically, it also used to make the task unable to reach a terminal state whenever + * that count failed, as recording completion, failure or cancellation all require these details. + */ @Override - public Publisher<Optional<TaskExecutionDetails.AdditionalInformation>> detailsReactive() { - return getRemainingSize() - .map(remainingSize -> new AdditionalInformation(mailRepositoryPath, initialCount, remainingSize, Clock.systemUTC().instant())) - .map(Optional::of); + public Optional<TaskExecutionDetails.AdditionalInformation> details() { + return Optional.of(new AdditionalInformation(mailRepositoryPath, + initialCount.get(), + Math.max(initialCount.get() - removedCount.get(), 0), + Clock.systemUTC().instant())); } - public Mono<Long> getRemainingSize() { + private Mono<Long> getRemainingSize() { try { return Flux.fromStream(mailRepositoryStore.getByPath(mailRepositoryPath)) .flatMap(MailRepository::sizeReactive) diff --git a/server/protocols/webadmin/webadmin-mailrepository/src/test/java/org/apache/james/webadmin/service/ClearMailRepositoryTaskTest.java b/server/protocols/webadmin/webadmin-mailrepository/src/test/java/org/apache/james/webadmin/service/ClearMailRepositoryTaskTest.java index 51e167a93a..f73b730b3e 100644 --- a/server/protocols/webadmin/webadmin-mailrepository/src/test/java/org/apache/james/webadmin/service/ClearMailRepositoryTaskTest.java +++ b/server/protocols/webadmin/webadmin-mailrepository/src/test/java/org/apache/james/webadmin/service/ClearMailRepositoryTaskTest.java @@ -19,16 +19,31 @@ package org.apache.james.webadmin.service; +import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; import java.time.Instant; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.Consumer; +import java.util.stream.Stream; + +import jakarta.mail.MessagingException; import org.apache.james.JsonSerializationVerifier; +import org.apache.james.mailrepository.api.MailKey; +import org.apache.james.mailrepository.api.MailRepository; import org.apache.james.mailrepository.api.MailRepositoryPath; import org.apache.james.mailrepository.api.MailRepositoryStore; import org.apache.james.server.task.json.JsonTaskSerializer; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentMatchers; + +import reactor.core.publisher.Mono; class ClearMailRepositoryTaskTest { @@ -59,6 +74,88 @@ class ClearMailRepositoryTaskTest { .isInstanceOf(ClearMailRepositoryTask.InvalidMailRepositoryPathDeserializationException.class); } + @Test + void detailsShouldNotQueryTheRepository() throws Exception { + MailRepository repository = mock(MailRepository.class); + MailRepositoryStore store = mock(MailRepositoryStore.class); + when(store.getByPath(MAIL_REPOSITORY_PATH)).thenAnswer(invocation -> Stream.of(repository)); + + ClearMailRepositoryTask task = new ClearMailRepositoryTask(store, MAIL_REPOSITORY_PATH); + + assertThat(task.details()).isPresent(); + verify(repository, never()).size(); + verify(repository, never()).sizeReactive(); + } + + @Test + void detailsShouldReportProgressWhileMailsAreBeingRemoved() throws Exception { + MailRepository repository = mock(MailRepository.class); + when(repository.sizeReactive()).thenReturn(Mono.just(3L)); + MailRepositoryStore store = mock(MailRepositoryStore.class); + when(store.getByPath(MAIL_REPOSITORY_PATH)).thenAnswer(invocation -> Stream.of(repository)); + + ClearMailRepositoryTask task = new ClearMailRepositoryTask(store, MAIL_REPOSITORY_PATH); + AtomicLong remainingHalfWayThrough = new AtomicLong(); + + doAnswer(invocation -> { + Consumer<MailKey> progressCallback = invocation.getArgument(0); + progressCallback.accept(new MailKey("mail1")); + progressCallback.accept(new MailKey("mail2")); + remainingHalfWayThrough.set(remainingCount(task)); + return null; + }).when(repository).removeAll(ArgumentMatchers.<Consumer<MailKey>>any()); + + task.run(); + + assertThat(remainingHalfWayThrough.get()).isEqualTo(1); + assertThat(remainingCount(task)).isEqualTo(0); + } + + /** + * Repositories deleting in bulk report no intermediate progress: only their completion tells us the + * repository is now empty. + */ + @Test + void detailsShouldReportAnEmptyRepositoryUponCompletionOfABulkRemoval() throws Exception { + MailRepository repository = mock(MailRepository.class); + when(repository.sizeReactive()).thenReturn(Mono.just(3L)); + MailRepositoryStore store = mock(MailRepositoryStore.class); + when(store.getByPath(MAIL_REPOSITORY_PATH)).thenAnswer(invocation -> Stream.of(repository)); + + ClearMailRepositoryTask task = new ClearMailRepositoryTask(store, MAIL_REPOSITORY_PATH); + + task.run(); + + assertThat(remainingCount(task)).isEqualTo(0); + } + + @Test + void detailsShouldNotReportANegativeRemainingCountWhenMailsGetAddedDuringTheRemoval() throws Exception { + MailRepository repository = mock(MailRepository.class); + when(repository.sizeReactive()).thenReturn(Mono.just(1L)); + MailRepositoryStore store = mock(MailRepositoryStore.class); + when(store.getByPath(MAIL_REPOSITORY_PATH)).thenAnswer(invocation -> Stream.of(repository)); + + ClearMailRepositoryTask task = new ClearMailRepositoryTask(store, MAIL_REPOSITORY_PATH); + AtomicLong remainingHalfWayThrough = new AtomicLong(); + + doAnswer(invocation -> { + Consumer<MailKey> progressCallback = invocation.getArgument(0); + progressCallback.accept(new MailKey("mail1")); + progressCallback.accept(new MailKey("mail2")); + remainingHalfWayThrough.set(remainingCount(task)); + return null; + }).when(repository).removeAll(ArgumentMatchers.<Consumer<MailKey>>any()); + + task.run(); + + assertThat(remainingHalfWayThrough.get()).isEqualTo(0); + } + + private long remainingCount(ClearMailRepositoryTask task) throws MessagingException { + return ((ClearMailRepositoryTask.AdditionalInformation) task.details().get()).getRemainingCount(); + } + @Test void additionalInformationShouldBeSerializable() throws Exception { JsonSerializationVerifier.dtoModule(ClearMailRepositoryTaskAdditionalInformationDTO.module()) --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
