diff --git a/spring-modulith-events/spring-modulith-events-mongodb/src/main/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepository.java b/spring-modulith-events/spring-modulith-events-mongodb/src/main/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepository.java index e6c73d2bd..ef87fea2f 100644 --- a/spring-modulith-events/spring-modulith-events-mongodb/src/main/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepository.java +++ b/spring-modulith-events/spring-modulith-events-mongodb/src/main/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepository.java @@ -109,7 +109,7 @@ public void markCompleted(Object event, PublicationTargetIdentifier identifier, var criteria = byEventAndListenerId(event, identifier); var query = defaultQuery(criteria); - var update = Update.update(COMPLETION_DATE, completionDate); + var update = Update.update(COMPLETION_DATE, completionDate).set(STATUS, Status.COMPLETED); if (completionMode == CompletionMode.DELETE) { @@ -134,7 +134,7 @@ public void markCompleted(UUID identifier, Instant completionDate) { var criteria = where(ID).is(identifier).and(COMPLETION_DATE).isNull(); var query = query(criteria); - var update = Update.update(COMPLETION_DATE, completionDate); + var update = Update.update(COMPLETION_DATE, completionDate).set(STATUS, Status.COMPLETED); if (completionMode == CompletionMode.DELETE) { @@ -156,7 +156,7 @@ public void markCompleted(UUID identifier, Instant completionDate) { @Override public void markFailed(UUID identifier) { - var query = query(where(ID).is(identifier).and(STATUS).ne(Status.FAILED)); + var query = query(where(ID).is(identifier).and(STATUS).ne(Status.FAILED).and(COMPLETION_DATE).isNull()); var update = Update.update(STATUS, Status.FAILED); mongoTemplate.findAndModify(query, update, MongoDbEventPublication.class, collection); @@ -169,7 +169,7 @@ public void markFailed(UUID identifier) { @Override public boolean markResubmitted(UUID identifier, Instant resubmissionDate) { - var query = query(where(ID).is(identifier).and(STATUS).ne(Status.RESUBMITTED)); + var query = query(where(ID).is(identifier).and(STATUS).ne(Status.RESUBMITTED).and(COMPLETION_DATE).isNull()); var update = Update.update(STATUS, Status.RESUBMITTED) .inc(COMPLETION_ATTEMPTS, 1) .set(LAST_RESUBMISSION_DATE, resubmissionDate); @@ -230,7 +230,7 @@ public List findCompletedPublications() { @Override public List findFailedPublications(FailedCriteria criteria) { - var statusFailed = where(STATUS).is(Status.FAILED); + var statusFailed = byStatus(Status.FAILED); var noStatusAndCompletionDate = where(STATUS).isNull().and(COMPLETION_DATE).isNull(); var baseCriteria = new Criteria().orOperator(statusFailed, noStatusAndCompletionDate); @@ -253,6 +253,18 @@ public List findFailedPublications(FailedCriteria criter return readMapped(limit != -1 ? query.limit((int) limit) : query); } + /* + * (non-Javadoc) + * @see org.springframework.modulith.events.core.EventPublicationRepository#findByStatus(org.springframework.modulith.events.EventPublication.Status) + */ + @Override + public List findByStatus(Status status) { + + var collection = status == Status.COMPLETED ? archiveCollection : this.collection; + + return readMapped(defaultQuery(byStatus(status)), collection); + } + /* * (non-Javadoc) * @see org.springframework.modulith.events.core.EventPublicationRepository#countByStatus(org.springframework.modulith.events.EventPublication.Status) @@ -260,11 +272,9 @@ public List findFailedPublications(FailedCriteria criter @Override public int countByStatus(Status status) { - var collection = status == Status.COMPLETED && completionMode == CompletionMode.ARCHIVE - ? archiveCollection - : this.collection; + var collection = status == Status.COMPLETED ? archiveCollection : this.collection; - return (int) mongoTemplate.count(query(where(STATUS).is(status)), MongoDbEventPublication.class, collection); + return (int) mongoTemplate.count(query(byStatus(status)), MongoDbEventPublication.class, collection); } /* @@ -323,6 +333,16 @@ private Criteria byEventAndListenerId(Object event, PublicationTargetIdentifier .and(COMPLETION_DATE).isNull(); } + private static Criteria byStatus(Status status) { + + Assert.notNull(status, "Status must not be null!"); + + // Older publications can have a completion date without a corresponding status update. + return status == Status.COMPLETED + ? where(COMPLETION_DATE).ne(null) + : where(STATUS).is(status).and(COMPLETION_DATE).isNull(); + } + private static MongoDbEventPublication domainToDocument(TargetEventPublication publication) { return new MongoDbEventPublication( // @@ -352,6 +372,7 @@ private void markCompleted(Criteria lookup, Instant now) { addFields() .addFieldWithValue(COMPLETION_DATE, now) + .addFieldWithValue(STATUS, Status.COMPLETED.name()) .build(), merge() @@ -407,11 +428,11 @@ public void markCompleted(Instant instant) { @Override public Status getStatus() { - if (publication.status != null) { - return publication.status; + if (publication.completionDate != null) { + return Status.COMPLETED; } - return publication.completionDate != null ? Status.COMPLETED : Status.PUBLISHED; + return publication.status != null ? publication.status : Status.PUBLISHED; } @Override diff --git a/spring-modulith-events/spring-modulith-events-mongodb/src/test/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepositoryTest.java b/spring-modulith-events/spring-modulith-events-mongodb/src/test/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepositoryTest.java index 84f6b5de4..75177159f 100644 --- a/spring-modulith-events/spring-modulith-events-mongodb/src/test/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepositoryTest.java +++ b/spring-modulith-events/spring-modulith-events-mongodb/src/test/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepositoryTest.java @@ -18,10 +18,13 @@ import static org.assertj.core.api.Assertions.*; import static org.junit.jupiter.api.Assumptions.*; +import java.time.Clock; +import java.time.Duration; import java.time.Instant; import java.time.LocalDateTime; import java.time.ZoneOffset; import java.time.temporal.ChronoUnit; +import java.util.ArrayList; import java.util.Comparator; import java.util.List; import java.util.UUID; @@ -30,12 +33,16 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Nested; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.data.mongodb.test.autoconfigure.DataMongoTest; import org.springframework.context.annotation.Import; import org.springframework.core.env.Environment; import org.springframework.data.mongodb.core.MongoTemplate; import org.springframework.modulith.events.EventPublication.Status; +import org.springframework.modulith.events.ResubmissionOptions; +import org.springframework.modulith.events.core.DefaultEventPublicationRegistry; import org.springframework.modulith.events.core.EventPublicationRepository.FailedCriteria; import org.springframework.modulith.events.core.PublicationTargetIdentifier; import org.springframework.modulith.events.core.TargetEventPublication; @@ -328,19 +335,136 @@ void claimsResubmissionOnce() { assertThat(repository.markResubmitted(publication.getIdentifier(), now)).isFalse(); } - @Test // GH-1336 + @Test // GH-1336, GH-1855 void countsByStatus() { var event = new TestEvent("first"); var publication = createPublication(event); - assertOneByStatus(Status.PUBLISHED); + assertByStatus(Status.PUBLISHED, publication.getIdentifier()); repository.markFailed(publication.getIdentifier()); - assertOneByStatus(Status.FAILED); + assertByStatus(Status.FAILED, publication.getIdentifier()); repository.markResubmitted(publication.getIdentifier(), Instant.now()); - assertOneByStatus(Status.RESUBMITTED); + assertByStatus(Status.RESUBMITTED, publication.getIdentifier()); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + assertCompleted(publication.getIdentifier()); + } + + @Test // GH-1855 + void marksPublicationAsCompletedByEventAndTargetIdentifier() { + + var publication = createPublication(new TestEvent("first")); + + repository.markCompleted(publication.getEvent(), TARGET_IDENTIFIER, Instant.now()); + + assertCompleted(publication.getIdentifier()); + } + + @ParameterizedTest // GH-1855 + @EnumSource(Status.class) + void findsPublicationsByStatusInPublicationOrder(Status status) { + + var now = Instant.parse("2026-01-01T12:00:00Z"); + var second = savePublicationAt(now, status); + var first = savePublicationAt(now.minusSeconds(1), status); + savePublicationAt(now.minusSeconds(2), status == Status.FAILED ? Status.PUBLISHED : Status.FAILED); + + assertByStatusResult(status, first.id, second.id); + } + + @ParameterizedTest // GH-1855 + @EnumSource(value = Status.class, names = "COMPLETED", mode = EnumSource.Mode.EXCLUDE) + void recognizesPreviouslyCompletedPublications(Status status) { + + var now = Instant.parse("2026-01-01T12:00:00Z"); + var collection = completionMode == CompletionMode.ARCHIVE + ? archiveCollection + : mongoTemplate.getCollectionName(MongoDbEventPublication.class); + + // Older versions recorded completion without updating the stored status. + var publication = new MongoDbEventPublication(UUID.randomUUID(), now.minusSeconds(60), "listener", + new TestEvent("completed"), now, status, null, 1); + + mongoTemplate.save(publication, collection); + + assertByStatus(Status.COMPLETED, publication.id); + assertThat(repository.findCompletedPublications()).singleElement() + .satisfies(it -> assertThat(it.getStatus()).isEqualTo(Status.COMPLETED)); + assertThat(repository.findFailedPublications(FailedCriteria.ALL)).isEmpty(); + + repository.markFailed(publication.id); + assertThat(repository.markResubmitted(publication.id, now.plusSeconds(1))).isFalse(); + + assertThat(mongoTemplate.findAll(MongoDbEventPublication.class, collection)).singleElement() + .satisfies(it -> assertThat(it.status).isEqualTo(status)); + } + + @Test // GH-1855 + void doesNotFailPublicationCompletedAfterStatusLookup() { + + var publication = createPublication(new TestEvent("first")); + var candidates = repository.findByStatus(Status.PUBLISHED); + + assertThat(candidates).hasSize(1); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + candidates.forEach(it -> repository.markFailed(it.getIdentifier())); + + assertCompleted(publication.getIdentifier()); + assertThat(repository.findFailedPublications(FailedCriteria.ALL)).isEmpty(); + } + + @Test // GH-1855 + void doesNotResubmitPublicationCompletedAfterFailedLookup() { + + var publication = createPublication(new TestEvent("first")); + repository.markFailed(publication.getIdentifier()); + + var candidates = repository.findFailedPublications(FailedCriteria.ALL); + + assertThat(candidates).hasSize(1); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + candidates.forEach(it -> assertThat(repository.markResubmitted(it.getIdentifier(), Instant.now())).isFalse()); + + assertCompleted(publication.getIdentifier()); + } + + @Test // GH-1855 + void resubmitsOnlyStaleIncompletePublications() { + + var now = Instant.parse("2026-01-01T12:00:00Z"); + var old = now.minusSeconds(120); + var published = savePublicationAt(old, Status.PUBLISHED); + var processing = savePublicationAt(old.plusSeconds(1), Status.PROCESSING); + var resubmitted = savePublicationAt(old.plusSeconds(2), Status.PUBLISHED); + var recent = savePublicationAt(old.plusSeconds(3), Status.PUBLISHED); + var completed = savePublicationAt(old.plusSeconds(4), Status.PUBLISHED); + + assertThat(repository.markResubmitted(resubmitted.id, old.plusSeconds(30))).isTrue(); + assertThat(repository.markResubmitted(recent.id, now.minusSeconds(5))).isTrue(); + repository.markCompleted(completed.id, now.minusSeconds(30)); + + var registry = new DefaultEventPublicationRegistry(repository, Clock.fixed(now, ZoneOffset.UTC)); + + registry.markStalePublicationsFailed(__ -> Duration.ofSeconds(60)); + + assertByStatusResult(Status.FAILED, published.id, processing.id, resubmitted.id); + assertByStatusResult(Status.RESUBMITTED, recent.id); + + var resubmittedIdentifiers = new ArrayList(); + + registry.processFailedPublications(ResubmissionOptions.defaults(), + it -> resubmittedIdentifiers.add(it.getIdentifier())); + + assertThat(resubmittedIdentifiers).containsExactly(published.id, processing.id, resubmitted.id); + assertByStatusResult(Status.FAILED); + assertByStatusResult(Status.RESUBMITTED, published.id, processing.id, resubmitted.id, recent.id); + assertByStatusResult(Status.COMPLETED, + completionMode == CompletionMode.DELETE ? new UUID[0] : new UUID[] { completed.id }); } @Test // GH-1336 @@ -399,12 +523,56 @@ private void savePublicationAt(LocalDateTime date) { mongoTemplate.save(publication); } - private void assertOneByStatus(Status reference) { + private MongoDbEventPublication savePublicationAt(Instant date, Status status) { + + var completed = status == Status.COMPLETED; + var publication = new MongoDbEventPublication(UUID.randomUUID(), date, "listener", new TestEvent("event"), + completed ? date.plusSeconds(1) : null, status, null, 1); + + return completed && completionMode == CompletionMode.ARCHIVE + ? mongoTemplate.save(publication, archiveCollection) + : mongoTemplate.save(publication); + } + + private void assertCompleted(UUID identifier) { + + if (completionMode == CompletionMode.DELETE) { + + assertByStatus(Status.COMPLETED); + assertThat(mongoTemplate.findAll(MongoDbEventPublication.class)).isEmpty(); + + } else { + + assertByStatus(Status.COMPLETED, identifier); + + var collection = completionMode == CompletionMode.ARCHIVE + ? archiveCollection + : mongoTemplate.getCollectionName(MongoDbEventPublication.class); + + assertThat(mongoTemplate.findAll(MongoDbEventPublication.class, collection)).singleElement().satisfies(it -> { + assertThat(it.id).isEqualTo(identifier); + assertThat(it.status).isEqualTo(Status.COMPLETED); + assertThat(it.completionDate).isNotNull(); + }); + } + } + + private void assertByStatus(Status reference, UUID... identifiers) { for (var status : Status.values()) { - assertThat(repository.countByStatus(status)).isEqualTo(status == reference ? 1 : 0); + assertByStatusResult(status, status == reference ? identifiers : new UUID[0]); } } + + private void assertByStatusResult(Status status, UUID... identifiers) { + + assertThat(repository.countByStatus(status)).isEqualTo(identifiers.length); + assertThat(repository.findByStatus(status)) + .extracting(TargetEventPublication::getIdentifier) + .containsExactly(identifiers); + assertThat(repository.findByStatus(status)) + .allSatisfy(it -> assertThat(it.getStatus()).isEqualTo(status)); + } } @Nested