From 3ac51f5090771561431150c75e7b55543b323786 Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Wed, 23 Sep 2026 14:14:37 +0200 Subject: [PATCH 01/15] GH-1899 - Declare releaseTrain profile explicitly to deactivate profiles active by default. Maven only deactivates profiles active by default if any of the explicitly activated profiles is declared in the same pom that declares the one active by default. We now explicitly declare the releaseTrain profile so that the one named "default" is properly disabled. --- pom.xml | 6 ++++++ spring-modulith-events/pom.xml | 7 ++++++- 2 files changed, 12 insertions(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index 75c804903..8421c7b6f 100644 --- a/pom.xml +++ b/pom.xml @@ -138,6 +138,12 @@ limitations under the License. + + + + releaseTrain + + diff --git a/spring-modulith-events/pom.xml b/spring-modulith-events/pom.xml index 97f35404c..ebef555bd 100644 --- a/spring-modulith-events/pom.xml +++ b/spring-modulith-events/pom.xml @@ -55,11 +55,16 @@ - + + sonatype + + releaseTrain + + From e48d3f902ac049f58bded48944f327440d6af0e0 Mon Sep 17 00:00:00 2001 From: Roland Beisel Date: Sun, 20 Sep 2026 20:47:25 +0200 Subject: [PATCH 02/15] GH-1895 - Defer event mapping in Namastack outbox mode Store the original domain event in the Namastack outbox and apply the configured mapping when the record is delivered. This keeps selection and routing based on the original event and reports mapping or broker serialization failures as outbox failures. Original pull request: GH-1897 Signed-off-by: Roland Beisel --- ...rnalizerConfigurationIntegrationTests.java | 29 +++++++++++++++++++ .../NamastackOutboxEventRecorder.java | 5 ++-- ...NamastackOutboxEventRecorderUnitTests.java | 27 +++++++++++++++++ 3 files changed, 58 insertions(+), 3 deletions(-) diff --git a/spring-modulith-events/spring-modulith-events-kafka/src/test/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfigurationIntegrationTests.java b/spring-modulith-events/spring-modulith-events-kafka/src/test/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfigurationIntegrationTests.java index f03d71921..6eb46f0e0 100644 --- a/spring-modulith-events/spring-modulith-events-kafka/src/test/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfigurationIntegrationTests.java +++ b/spring-modulith-events/spring-modulith-events-kafka/src/test/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfigurationIntegrationTests.java @@ -151,6 +151,18 @@ void publishesEventExternalizedAfterNamastackExternalization() { assertEventExternalizedPublished(OutboxHandler.class, (transport, event) -> transport.handle(event, null)); } + @Test // GH-1895 + void appliesMappingForNamastackExternalization() { + + var config = EventExternalizationConfiguration.defaults("org") + .mapping(Sample.class, __ -> new MappedSample()) + .build(); + + assertOutboxMessage(config, it -> { + assertThat(it.getPayload()).isInstanceOf(MappedSample.class); + }); + } + private void assertMessage(EventExternalizationConfiguration configuration, Consumer> assertions) { basicSetup(configuration) @@ -165,6 +177,21 @@ private void assertMessage(EventExternalizationConfiguration configuration, Cons }); } + private void assertOutboxMessage(EventExternalizationConfiguration configuration, Consumer> assertions) { + + basicSetup(configuration) + .withPropertyValues(ExternalizationMode.PROPERTY + "=" + ExternalizationMode.OUTBOX) + .run(ctxt -> { + + ctxt.getBean(OutboxHandler.class).handle(new Sample(), null); + + var captor = ArgumentCaptor.forClass(Message.class); + verify(operations).send(captor.capture()); + + assertions.accept(captor.getValue()); + }); + } + private ApplicationContextRunner basicSetup() { return basicSetup(null); } @@ -208,4 +235,6 @@ private void assertEventExternalizedPublished(Class transportType, BiCons @Externalized record Sample() {} + + record MappedSample() {} } diff --git a/spring-modulith-events/spring-modulith-events-namastack/src/main/java/org/springframework/modulith/events/namastack/NamastackOutboxEventRecorder.java b/spring-modulith-events/spring-modulith-events-namastack/src/main/java/org/springframework/modulith/events/namastack/NamastackOutboxEventRecorder.java index cbc288371..317ccddba 100644 --- a/spring-modulith-events/spring-modulith-events-namastack/src/main/java/org/springframework/modulith/events/namastack/NamastackOutboxEventRecorder.java +++ b/spring-modulith-events/spring-modulith-events-namastack/src/main/java/org/springframework/modulith/events/namastack/NamastackOutboxEventRecorder.java @@ -81,16 +81,15 @@ public void onApplicationEvent(PayloadApplicationEvent event) { } var target = configuration.determineTarget(payload); - var mapped = configuration.map(payload); var routing = BrokerRouting.of(target, context); - var key = routing.getKey(mapped); + var key = routing.getKey(payload); if (LOGGER.isDebugEnabled()) { LOGGER.debug("Scheduling event of type {} to outbox for target {}.", payload.getClass().getName(), target.getTarget()); } - scheduleToOutbox(mapped, key); + scheduleToOutbox(payload, key); } /** diff --git a/spring-modulith-events/spring-modulith-events-namastack/src/test/java/org/springframework/modulith/events/namastack/NamastackOutboxEventRecorderUnitTests.java b/spring-modulith-events/spring-modulith-events-namastack/src/test/java/org/springframework/modulith/events/namastack/NamastackOutboxEventRecorderUnitTests.java index 7323fd296..b9f2c0a5b 100644 --- a/spring-modulith-events/spring-modulith-events-namastack/src/test/java/org/springframework/modulith/events/namastack/NamastackOutboxEventRecorderUnitTests.java +++ b/spring-modulith-events/spring-modulith-events-namastack/src/test/java/org/springframework/modulith/events/namastack/NamastackOutboxEventRecorderUnitTests.java @@ -15,10 +15,13 @@ */ package org.springframework.modulith.events.namastack; +import static org.assertj.core.api.Assertions.*; import static org.mockito.Mockito.*; import io.namastack.outbox.Outbox; +import java.util.concurrent.atomic.AtomicInteger; + import org.junit.jupiter.api.Test; import org.springframework.context.PayloadApplicationEvent; import org.springframework.expression.spel.support.StandardEvaluationContext; @@ -48,6 +51,30 @@ void evaluatesRoutingKeyExpressionAgainstEventPayload() { verify(outbox).schedule(payload, "value"); } + @Test // GH-1895 + void storesOriginalEventWithoutApplyingMapping() { + + var mappings = new AtomicInteger(); + var configuration = EventExternalizationConfiguration.externalizing() + .select(EventExternalizationConfiguration.annotatedAsExternalized()) + .mapping(SampleEvent.class, it -> { + mappings.incrementAndGet(); + return new MappedEvent(it.getValue()); + }) + .build(); + + var payload = new SampleEvent("value"); + var outbox = mock(Outbox.class); + var recorder = new NamastackOutboxEventRecorder(configuration, outbox, new StandardEvaluationContext()); + + recorder.onApplicationEvent(new PayloadApplicationEvent<>(this, payload)); + + verify(outbox).schedule(payload, "value"); + assertThat(mappings).hasValue(0); + } + + record MappedEvent(String value) {} + @Externalized("target::#{getValue()}") static class SampleEvent { From 9afdcf40815bfcfdc22ca6a26978f47c336a9f04 Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Wed, 23 Sep 2026 15:31:25 +0200 Subject: [PATCH 03/15] GH-1896 - Calculate Broker routing on original event instead of the mapped payload. --- .../RabbitEventExternalizerConfiguration.java | 47 +++++++------ ...EventExternalizationAutoConfiguration.java | 6 +- .../EventExternalizerModuleListener.java | 25 +++++++ .../support/EventExternalizerSupport.java | 67 +++++++++++++++++-- .../support/OutboxEventExternalizer.java | 23 +++++++ .../OutboxEventExternalizerFactory.java | 32 ++++++++- .../TransportAwareEventExternalizer.java | 28 +++++++- ...ntExternalizerModuleListenerUnitTests.java | 41 ++++++++++++ .../JmsEventExternalizerConfiguration.java | 2 +- .../KafkaEventExternalizerConfiguration.java | 2 +- ...rnalizerConfigurationIntegrationTests.java | 32 +++++++++ ...ssagingEventExternalizerConfiguration.java | 3 +- ...agingEventPublicationIntegrationTests.java | 12 +++- 13 files changed, 279 insertions(+), 41 deletions(-) diff --git a/spring-modulith-events/spring-modulith-events-amqp/src/main/java/org/springframework/modulith/events/amqp/RabbitEventExternalizerConfiguration.java b/spring-modulith-events/spring-modulith-events-amqp/src/main/java/org/springframework/modulith/events/amqp/RabbitEventExternalizerConfiguration.java index c7abde3b4..1c81032f1 100644 --- a/spring-modulith-events/spring-modulith-events-amqp/src/main/java/org/springframework/modulith/events/amqp/RabbitEventExternalizerConfiguration.java +++ b/spring-modulith-events/spring-modulith-events-amqp/src/main/java/org/springframework/modulith/events/amqp/RabbitEventExternalizerConfiguration.java @@ -68,6 +68,16 @@ class RabbitEventExternalizerConfiguration { private static final Logger logger = LoggerFactory.getLogger(RabbitEventExternalizerConfiguration.class); private static final String PUBLISHER_CONFIRM_TYPE = "spring.rabbitmq.publisher-confirm-type"; + private final EvaluationContext evaluationContext; + + RabbitEventExternalizerConfiguration(BeanFactory factory) { + + var context = new StandardEvaluationContext(); + context.setBeanResolver(new BeanFactoryResolver(factory)); + + this.evaluationContext = context; + } + @Bean @ConditionalOnProperty(name = ExternalizationMode.PROPERTY, havingValue = "module-listener", matchIfMissing = true) EventExternalizerModuleListener rabbitEventExternalizer(EventExternalizationConfiguration configuration, @@ -76,22 +86,22 @@ EventExternalizerModuleListener rabbitEventExternalizer(EventExternalizationConf logger.debug("Registering domain event externalization to RabbitMQ…"); return new EventExternalizerModuleListener(configuration, - createRabbitTransport(configuration, operations, factory)); + createRabbitTransport(configuration, operations), factory); } @AutoConfiguration @ConditionalOnProperty(name = ExternalizationMode.PROPERTY, havingValue = "outbox") - static class RabbitOutboxConfiguration { + class RabbitOutboxConfiguration { private final OutboxEventExternalizer externalizer; RabbitOutboxConfiguration(EventExternalizationConfiguration configuration, RabbitMessageOperations operations, - BeanFactory beanFactory, OutboxEventExternalizerFactory factory, Environment environment) { + OutboxEventExternalizerFactory factory, Environment environment) { Assert.state("correlated".equalsIgnoreCase(environment.getProperty(PUBLISHER_CONFIRM_TYPE)), () -> "RabbitMQ outbox event externalization requires " + PUBLISHER_CONFIRM_TYPE + "=correlated!"); - this.externalizer = factory.forTransport(createConfirmingRabbitTransport(configuration, operations, beanFactory)); + this.externalizer = factory.forTransport(createConfirmingRabbitTransport(configuration, operations)); } @AutoConfiguration @@ -121,28 +131,26 @@ JobRunrExternalizationTransport jobRunrRabbitOutboxExternalizer() { } } - private static EventExternalizationTransport createRabbitTransport( - EventExternalizationConfiguration configuration, RabbitMessageOperations operations, - BeanFactory factory) { + private EventExternalizationTransport createRabbitTransport( + EventExternalizationConfiguration configuration, RabbitMessageOperations operations) { return (payload, target) -> { - send(payload, target, configuration, Collections.emptyMap(), operations, factory); + send(payload, target, configuration, Collections.emptyMap(), operations); return CompletableFuture.completedFuture(null); }; } - private static EventExternalizationTransport createConfirmingRabbitTransport( - EventExternalizationConfiguration configuration, RabbitMessageOperations operations, - BeanFactory factory) { + private EventExternalizationTransport createConfirmingRabbitTransport(EventExternalizationConfiguration configuration, + RabbitMessageOperations operations) { return (payload, target) -> { var correlation = new CorrelationData(); var headers = Map.of(AmqpHeaders.PUBLISH_CONFIRM_CORRELATION, (Object) correlation); - send(payload, target, configuration, headers, operations, factory); + send(payload, target, configuration, headers, operations); return correlation.getFuture().thenAccept(confirm -> { @@ -154,22 +162,13 @@ private static EventExternalizationTransport createConfirmingRabbitTransport( }; } - private static void send(Object payload, RoutingTarget target, EventExternalizationConfiguration configuration, - Map additionalHeaders, RabbitMessageOperations operations, BeanFactory factory) { - - var routing = BrokerRouting.of(target, createContext(factory)); + private void send(Object payload, RoutingTarget target, EventExternalizationConfiguration configuration, + Map additionalHeaders, RabbitMessageOperations operations) { + var routing = BrokerRouting.of(target, evaluationContext); var headers = new HashMap<>(configuration.getHeadersFor(payload)); headers.putAll(additionalHeaders); operations.convertAndSend(routing.getTarget(payload), routing.getKey(payload), payload, additionalHeaders); } - - private static EvaluationContext createContext(BeanFactory factory) { - - var context = new StandardEvaluationContext(); - context.setBeanResolver(new BeanFactoryResolver(factory)); - - return context; - } } diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventExternalizationAutoConfiguration.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventExternalizationAutoConfiguration.java index aa7b35e5b..7b1896834 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventExternalizationAutoConfiguration.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventExternalizationAutoConfiguration.java @@ -26,7 +26,7 @@ import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; -import org.springframework.context.ApplicationEventPublisher; +import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationListener; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Role; @@ -89,8 +89,8 @@ static EventExternalizationConfiguration eventExternalizationConfiguration(BeanF @Role(BeanDefinition.ROLE_INFRASTRUCTURE) @ConditionalOnProperty(name = ExternalizationMode.PROPERTY, havingValue = "outbox") static OutboxEventExternalizerFactory outboxEventExternalizerFactory(EventExternalizationConfiguration configuration, - ApplicationEventPublisher publisher) { - return new OutboxEventExternalizerFactory(configuration, publisher); + ApplicationContext context) { + return new OutboxEventExternalizerFactory(configuration, context); } /** diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/EventExternalizerModuleListener.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/EventExternalizerModuleListener.java index 8116bec3f..86909c791 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/EventExternalizerModuleListener.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/EventExternalizerModuleListener.java @@ -17,6 +17,7 @@ import java.util.concurrent.CompletableFuture; +import org.springframework.beans.factory.BeanFactory; import org.springframework.modulith.events.ApplicationModuleListener; import org.springframework.modulith.events.EventExternalizationConfiguration; import org.springframework.modulith.events.core.ConditionalEventListener; @@ -47,7 +48,11 @@ public class EventExternalizerModuleListener extends TransportAwareEventExternal * * @param configuration must not be {@literal null}. * @param transport must not be {@literal null}. + * @deprecated since 2.2, 2.1.2, for removal in 2.3. Use + * {@link #EventExternalizerModuleListener(EventExternalizationConfiguration, EventExternalizationTransport, BeanFactory)} + * instead. */ + @Deprecated(since = "2.2, 2.1.2", forRemoval = true) public EventExternalizerModuleListener(EventExternalizationConfiguration configuration, EventExternalizationTransport transport) { @@ -58,6 +63,26 @@ public EventExternalizerModuleListener(EventExternalizationConfiguration configu this.configuration = configuration; } + /** + * Creates a new {@link EventExternalizerModuleListener} for the given {@link EventExternalizationConfiguration}, + * {@link EventExternalizationTransport} and {@link BeanFactory} to resolve routing target and key expressions (which + * may refer to beans) against the original event. + * + * @param configuration must not be {@literal null}. + * @param transport must not be {@literal null}. + * @param beanFactory must not be {@literal null}. + * @since 2.2, 2.1.2 + */ + public EventExternalizerModuleListener(EventExternalizationConfiguration configuration, + EventExternalizationTransport transport, BeanFactory beanFactory) { + + super(configuration, transport, beanFactory); + + Assert.notNull(transport, "EventExternalizationTransport must not be null!"); + + this.configuration = configuration; + } + /* * (non-Javadoc) * @see org.springframework.modulith.events.ConditionalEventListener#supports(java.lang.Object) diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/EventExternalizerSupport.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/EventExternalizerSupport.java index 49624c06d..31ba66101 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/EventExternalizerSupport.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/EventExternalizerSupport.java @@ -20,6 +20,10 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.BeanFactory; +import org.springframework.context.expression.BeanFactoryResolver; +import org.springframework.expression.EvaluationContext; +import org.springframework.expression.spel.support.StandardEvaluationContext; import org.springframework.modulith.events.EventExternalizationConfiguration; import org.springframework.modulith.events.EventExternalized; import org.springframework.modulith.events.RoutingTarget; @@ -38,18 +42,50 @@ abstract class EventExternalizerSupport { private static final Logger logger = LoggerFactory.getLogger(EventExternalizerSupport.class); private final EventExternalizationConfiguration configuration; + private final EvaluationContext context; private final Semaphore semaphore = new Semaphore(1); /** - * Creates a new {@link EventExternalizerSupport} for the given {@link EventExternalizationConfiguration}. + * Creates a new {@link EventExternalizerSupport} for the given {@link EventExternalizationConfiguration}, not + * resolving bean references in routing target and key expressions. * * @param configuration must not be {@literal null}. + * @deprecated since 2.2, 2.1.2, for removal in 2.3. Use + * {@link #EventExternalizerSupport(EventExternalizationConfiguration, BeanFactory)} instead. */ + @Deprecated(since = "2.2, 2.1.2", forRemoval = true) protected EventExternalizerSupport(EventExternalizationConfiguration configuration) { + this(configuration, new StandardEvaluationContext()); + } + + /** + * Creates a new {@link EventExternalizerSupport} for the given {@link EventExternalizationConfiguration}, resolving + * routing target and key expressions (which may refer to beans) against the original event using an + * {@link EvaluationContext} backed by the given {@link BeanFactory}. + * + * @param configuration must not be {@literal null}. + * @param beanFactory must not be {@literal null}. + * @since 2.2, 2.1.2 + */ + protected EventExternalizerSupport(EventExternalizationConfiguration configuration, BeanFactory beanFactory) { + + Assert.notNull(configuration, "EventExternalizationConfiguration must not be null!"); + Assert.notNull(beanFactory, "BeanFactory must not be null!"); + + var context = new StandardEvaluationContext(); + context.setBeanResolver(new BeanFactoryResolver(beanFactory)); + + this.configuration = configuration; + this.context = context; + } + + private EventExternalizerSupport(EventExternalizationConfiguration configuration, EvaluationContext context) { Assert.notNull(configuration, "EventExternalizationConfiguration must not be null!"); + Assert.notNull(context, "EvaluationContext must not be null!"); this.configuration = configuration; + this.context = context; } /** @@ -67,17 +103,18 @@ public CompletableFuture externalize(Object event) { } var target = configuration.determineTarget(event); + var resolved = resolve(target, event); var mapped = configuration.map(event); if (logger.isTraceEnabled()) { - logger.trace("Externalizing event of type {} to {}, payload: {}).", event.getClass(), target, mapped); + logger.trace("Externalizing event of type {} to {}, payload: {}).", event.getClass(), resolved, mapped); } else if (logger.isDebugEnabled()) { - logger.debug("Externalizing event of type {} to {}.", event.getClass(), target); + logger.debug("Externalizing event of type {} to {}.", event.getClass(), resolved); } return configuration.serializeExternalization() - ? doExternalizeSerialized(event, mapped, target) - : doExternalize(event, mapped, target); + ? doExternalizeSerialized(event, mapped, resolved) + : doExternalize(event, mapped, resolved); } /** @@ -110,4 +147,24 @@ private CompletableFuture doExternalize(Object event, Object mapped, RoutingT return externalize(mapped, target) .thenApply(it -> new EventExternalized<>(event, mapped, target, it)); } + + /** + * Resolves dynamic target and key expressions declared on the given {@link RoutingTarget} against the given + * (original, unmapped) event, so that a mapping applied downstream cannot affect routing. + * + * @param target must not be {@literal null}. + * @param event must not be {@literal null}. + * @return will never be {@literal null}. + */ + private RoutingTarget resolve(RoutingTarget target, Object event) { + + if (!target.hasExpression()) { + return target; + } + + var routing = BrokerRouting.of(target, context); + var resolved = RoutingTarget.forTarget(routing.getTarget(event)); + + return target.getKey() == null ? resolved.withoutKey() : resolved.andKey(routing.getKey(event)); + } } diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/OutboxEventExternalizer.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/OutboxEventExternalizer.java index 45147d973..ebb4354b6 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/OutboxEventExternalizer.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/OutboxEventExternalizer.java @@ -18,6 +18,7 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; +import org.springframework.beans.factory.BeanFactory; import org.springframework.context.ApplicationEventPublisher; import org.springframework.modulith.events.EventExternalizationConfiguration; @@ -43,7 +44,11 @@ public class OutboxEventExternalizer extends TransportAwareEventExternalizer { * * @param configuration must not be {@literal null}. * @param publisher must not be {@literal null}. + * @deprecated since 2.2, 2.1.2, for removal in 2.3. Use + * {@link #OutboxEventExternalizer(EventExternalizationConfiguration, ApplicationEventPublisher, EventExternalizationTransport, BeanFactory)} + * instead. */ + @Deprecated(since = "2.2, 2.1.2", forRemoval = true) public OutboxEventExternalizer(EventExternalizationConfiguration configuration, ApplicationEventPublisher publisher, EventExternalizationTransport transport) { @@ -52,6 +57,24 @@ public OutboxEventExternalizer(EventExternalizationConfiguration configuration, this.events = publisher; } + /** + * Creates a new {@link OutboxEventExternalizer} for the given {@link EventExternalizationConfiguration}, transport + * function and {@link BeanFactory} to resolve routing target and key expressions (which may refer to beans) against + * the original event. + * + * @param configuration must not be {@literal null}. + * @param publisher must not be {@literal null}. + * @param beanFactory must not be {@literal null}. + * @since 2.2, 2.1.2 + */ + public OutboxEventExternalizer(EventExternalizationConfiguration configuration, ApplicationEventPublisher publisher, + EventExternalizationTransport transport, BeanFactory beanFactory) { + + super(configuration, transport, beanFactory); + + this.events = publisher; + } + /* * (non-Javadoc) * @see org.springframework.modulith.events.support.EventExternalizerSupport#externalize(java.lang.Object) diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/OutboxEventExternalizerFactory.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/OutboxEventExternalizerFactory.java index cf427705b..40e224186 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/OutboxEventExternalizerFactory.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/OutboxEventExternalizerFactory.java @@ -15,6 +15,9 @@ */ package org.springframework.modulith.events.support; +import org.jspecify.annotations.Nullable; +import org.springframework.beans.factory.BeanFactory; +import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationEventPublisher; import org.springframework.modulith.events.EventExternalizationConfiguration; import org.springframework.util.Assert; @@ -31,22 +34,45 @@ public class OutboxEventExternalizerFactory { private final EventExternalizationConfiguration configuration; private final ApplicationEventPublisher publisher; + private final @Nullable BeanFactory beanFactory; /** * Creates a new {@link OutboxEventExternalizerFactory} for the given {@link EventExternalizationConfiguration} and - * {@link ApplicationEventPublisher}. + * {@link ApplicationEventPublisher}, not resolving bean references in routing target and key expressions. * * @param configuration must not be {@literal null}. * @param publisher must not be {@literal null}. + * @deprecated since 2.2, 2.1.2, for removal in 2.3. Use + * {@link #OutboxEventExternalizerFactory(EventExternalizationConfiguration, ApplicationContext)} instead. */ + @Deprecated(since = "2.2, 2.1.2", forRemoval = true) public OutboxEventExternalizerFactory(EventExternalizationConfiguration configuration, ApplicationEventPublisher publisher) { + this(configuration, publisher, null); + } + + /** + * Creates a new {@link OutboxEventExternalizerFactory} for the given {@link EventExternalizationConfiguration} and + * {@link ApplicationContext}, resolving routing target and key expressions (which may refer to beans) against the + * original event. + * + * @param configuration must not be {@literal null}. + * @param context must not be {@literal null}. + * @since 2.2, 2.1.2 + */ + public OutboxEventExternalizerFactory(EventExternalizationConfiguration configuration, ApplicationContext context) { + this(configuration, context, context); + } + + private OutboxEventExternalizerFactory(EventExternalizationConfiguration configuration, + ApplicationEventPublisher publisher, @Nullable BeanFactory beanFactory) { Assert.notNull(configuration, "EventExternalizationConfiguration must not be null!"); Assert.notNull(publisher, "ApplicationEventPublisher must not be null!"); this.configuration = configuration; this.publisher = publisher; + this.beanFactory = beanFactory; } /** @@ -59,6 +85,8 @@ public OutboxEventExternalizer forTransport(EventExternalizationTransport transp Assert.notNull(transport, "EventExternalizationTransport must not be null!"); - return new OutboxEventExternalizer(configuration, publisher, transport); + return beanFactory != null + ? new OutboxEventExternalizer(configuration, publisher, transport, beanFactory) + : new OutboxEventExternalizer(configuration, publisher, transport); } } diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/TransportAwareEventExternalizer.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/TransportAwareEventExternalizer.java index 2835f94e6..5865ffc8e 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/TransportAwareEventExternalizer.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/TransportAwareEventExternalizer.java @@ -17,6 +17,7 @@ import java.util.concurrent.CompletableFuture; +import org.springframework.beans.factory.BeanFactory; import org.springframework.modulith.events.EventExternalizationConfiguration; import org.springframework.modulith.events.RoutingTarget; import org.springframework.util.Assert; @@ -35,11 +36,16 @@ abstract class TransportAwareEventExternalizer extends EventExternalizerSupport /** * Creates a new {@link TransportAwareEventExternalizer} for the given {@link EventExternalizationConfiguration} and - * {@link EventExternalizationTransport} implementing the actual externalization. + * {@link EventExternalizationTransport} implementing the actual externalization, not resolving bean references in + * routing target and key expressions. * * @param configuration must not be {@literal null}. * @param transport must not be {@literal null}. + * @deprecated since 2.2, 2.1.2, for removal in 2.3. Use + * {@link #TransportAwareEventExternalizer(EventExternalizationConfiguration, EventExternalizationTransport, BeanFactory)} + * instead. */ + @Deprecated(since = "2.2, 2.1.2", forRemoval = true) public TransportAwareEventExternalizer(EventExternalizationConfiguration configuration, EventExternalizationTransport transport) { @@ -50,6 +56,26 @@ public TransportAwareEventExternalizer(EventExternalizationConfiguration configu this.transport = transport; } + /** + * Creates a new {@link TransportAwareEventExternalizer} for the given {@link EventExternalizationConfiguration}, + * {@link EventExternalizationTransport} implementing the actual externalization and {@link BeanFactory} to resolve + * routing target and key expressions (which may refer to beans) against the original event. + * + * @param configuration must not be {@literal null}. + * @param transport must not be {@literal null}. + * @param beanFactory must not be {@literal null}. + * @since 2.2, 2.1.2 + */ + public TransportAwareEventExternalizer(EventExternalizationConfiguration configuration, + EventExternalizationTransport transport, BeanFactory beanFactory) { + + super(configuration, beanFactory); + + Assert.notNull(transport, "EventExternalizationTransport must not be null!"); + + this.transport = transport; + } + /* * (non-Javadoc) * @see org.springframework.modulith.events.support.EventExternalizerSupport#externalize(org.springframework.modulith.events.RoutingTarget, java.lang.Object) diff --git a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/EventExternalizerModuleListenerUnitTests.java b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/EventExternalizerModuleListenerUnitTests.java index 9922ef426..0f14ba7be 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/EventExternalizerModuleListenerUnitTests.java +++ b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/EventExternalizerModuleListenerUnitTests.java @@ -28,6 +28,7 @@ import java.util.concurrent.atomic.AtomicInteger; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; import org.springframework.core.annotation.AnnotatedElementUtils; import org.springframework.modulith.events.Externalized; import org.springframework.modulith.events.RoutingTarget; @@ -114,6 +115,46 @@ public void testOnlyOneDownstreamExecutionAtATime() throws InterruptedException verify(mock, times(3)).externalize(any(Object.class), any(RoutingTarget.class)); } + @Test // GH-1896 + void evaluatesRoutingExpressionsAgainstOriginalEventNotMappedPayload() { + + var mock = mock(EventExternalizationTransport.class); + when(mock.externalize(any(), any())).thenReturn(CompletableFuture.completedFuture(null)); + + var configuration = externalizing() + .select(annotatedAsExternalized()) + .mapping(RoutedSample.class, __ -> new Mapped()) + .build(); + + var service = new EventExternalizerModuleListener(configuration, mock::externalize); + + service.externalize(new RoutedSample("42")).join(); + + var payload = ArgumentCaptor.forClass(Object.class); + var target = ArgumentCaptor.forClass(RoutingTarget.class); + + verify(mock).externalize(payload.capture(), target.capture()); + + assertThat(payload.getValue()).isInstanceOf(Mapped.class); + assertThat(target.getValue().getKey()).isEqualTo("42"); + } + @Externalized static class Sample {} + + @Externalized("target::#{getValue()}") + static class RoutedSample { + + private final String value; + + RoutedSample(String value) { + this.value = value; + } + + public String getValue() { + return value; + } + } + + record Mapped() {} } diff --git a/spring-modulith-events/spring-modulith-events-jms/src/main/java/org/springframework/modulith/events/jms/JmsEventExternalizerConfiguration.java b/spring-modulith-events/spring-modulith-events-jms/src/main/java/org/springframework/modulith/events/jms/JmsEventExternalizerConfiguration.java index 27649f32c..2fea8760f 100644 --- a/spring-modulith-events/spring-modulith-events-jms/src/main/java/org/springframework/modulith/events/jms/JmsEventExternalizerConfiguration.java +++ b/spring-modulith-events/spring-modulith-events-jms/src/main/java/org/springframework/modulith/events/jms/JmsEventExternalizerConfiguration.java @@ -66,7 +66,7 @@ EventExternalizerModuleListener jmsEventExternalizer(EventExternalizationConfigu logger.debug("Registering domain event externalization to JMS…"); return new EventExternalizerModuleListener(configuration, - createJmsTransport(operations, serializer, factory)); + createJmsTransport(operations, serializer, factory), factory); } @AutoConfiguration diff --git a/spring-modulith-events/spring-modulith-events-kafka/src/main/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfiguration.java b/spring-modulith-events/spring-modulith-events-kafka/src/main/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfiguration.java index b4a4129c8..94c65c5c8 100644 --- a/spring-modulith-events/spring-modulith-events-kafka/src/main/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfiguration.java +++ b/spring-modulith-events/spring-modulith-events-kafka/src/main/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfiguration.java @@ -69,7 +69,7 @@ EventExternalizerModuleListener kafkaEventExternalizer(EventExternalizationConfi logger.debug("Registering domain event externalization to Kafka…"); return new EventExternalizerModuleListener(configuration, - createKafkaTransport(configuration, operations, factory)); + createKafkaTransport(configuration, operations, factory), factory); } @AutoConfiguration diff --git a/spring-modulith-events/spring-modulith-events-kafka/src/test/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfigurationIntegrationTests.java b/spring-modulith-events/spring-modulith-events-kafka/src/test/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfigurationIntegrationTests.java index 6eb46f0e0..f12972603 100644 --- a/spring-modulith-events/spring-modulith-events-kafka/src/test/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfigurationIntegrationTests.java +++ b/spring-modulith-events/spring-modulith-events-kafka/src/test/java/org/springframework/modulith/events/kafka/KafkaEventExternalizerConfigurationIntegrationTests.java @@ -31,6 +31,7 @@ import org.springframework.boot.test.context.FilteredClassLoader; import org.springframework.boot.test.context.runner.ApplicationContextRunner; import org.springframework.kafka.core.KafkaOperations; +import org.springframework.kafka.support.KafkaHeaders; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import org.springframework.modulith.events.EventExternalizationConfiguration; @@ -163,6 +164,29 @@ void appliesMappingForNamastackExternalization() { }); } + @Test // GH-1896 + void evaluatesRoutingKeyAgainstOriginalEventWhenMapped() { + + var config = EventExternalizationConfiguration.defaults("org") + .mapping(RoutedSample.class, __ -> new MappedSample()) + .build(); + + basicSetup(config) + .run(ctxt -> { + + ctxt.getBean(EventExternalizerModuleListener.class).externalize(new RoutedSample("42")); + + var captor = ArgumentCaptor.forClass(Message.class); + verify(operations).send(captor.capture()); + + var message = captor.getValue(); + + assertThat(message.getPayload()).isInstanceOf(MappedSample.class); + assertThat(message.getHeaders().get(KafkaHeaders.KEY)).isEqualTo("42"); + assertThat(message.getHeaders().get(KafkaHeaders.TOPIC)).isEqualTo("orders"); + }); + } + private void assertMessage(EventExternalizationConfiguration configuration, Consumer> assertions) { basicSetup(configuration) @@ -237,4 +261,12 @@ private void assertEventExternalizedPublished(Class transportType, BiCons record Sample() {} record MappedSample() {} + + @Externalized("orders::#{getValue()}") + record RoutedSample(String value) { + + public String getValue() { + return value; + } + } } diff --git a/spring-modulith-events/spring-modulith-events-messaging/src/main/java/org/springframework/modulith/events/messaging/SpringMessagingEventExternalizerConfiguration.java b/spring-modulith-events/spring-modulith-events-messaging/src/main/java/org/springframework/modulith/events/messaging/SpringMessagingEventExternalizerConfiguration.java index 1510fdc62..085240357 100644 --- a/spring-modulith-events/spring-modulith-events-messaging/src/main/java/org/springframework/modulith/events/messaging/SpringMessagingEventExternalizerConfiguration.java +++ b/spring-modulith-events/spring-modulith-events-messaging/src/main/java/org/springframework/modulith/events/messaging/SpringMessagingEventExternalizerConfiguration.java @@ -67,7 +67,8 @@ EventExternalizerModuleListener springMessagingEventExternalizer(EventExternaliz logger.debug("Registering domain event externalization for Spring Messaging…"); - return new EventExternalizerModuleListener(configuration, createMessagingTransport(configuration, factory)); + return new EventExternalizerModuleListener(configuration, createMessagingTransport(configuration, factory), + factory); } @AutoConfiguration diff --git a/spring-modulith-events/spring-modulith-events-messaging/src/test/java/org/springframework/modulith/events/messaging/SpringMessagingEventPublicationIntegrationTests.java b/spring-modulith-events/spring-modulith-events-messaging/src/test/java/org/springframework/modulith/events/messaging/SpringMessagingEventPublicationIntegrationTests.java index 258bf2ae9..60164b5c7 100644 --- a/spring-modulith-events/spring-modulith-events-messaging/src/test/java/org/springframework/modulith/events/messaging/SpringMessagingEventPublicationIntegrationTests.java +++ b/spring-modulith-events/spring-modulith-events-messaging/src/test/java/org/springframework/modulith/events/messaging/SpringMessagingEventPublicationIntegrationTests.java @@ -46,7 +46,8 @@ @SpringBootTest class SpringMessagingEventPublicationIntegrationTests { - private static final String TARGET = "target::#{someExpression}"; + private static final String TARGET = "target::#{getSomeExpression()}"; + private static final String RESOLVED_KEY = "resolved"; private static final String CHANNEL_NAME = "target"; private static final AtomicInteger COUNTER = new AtomicInteger(); @@ -69,7 +70,7 @@ IntegrationFlow inboundIntegrationFlow(@Qualifier(CHANNEL_NAME) MessageChannel i .handle((__, headers) -> { assertThat(headers.get(SpringMessagingEventExternalizerConfiguration.MODULITH_ROUTING_HEADER)) - .isEqualTo(TARGET); + .isEqualTo("target::" + RESOLVED_KEY); COUNTER.incrementAndGet(); return null; @@ -100,7 +101,12 @@ void publishesEventToSpringMessaging() throws Exception { } @Externalized(TARGET) - static class TestEvent {} + static class TestEvent { + + public String getSomeExpression() { + return RESOLVED_KEY; + } + } @RequiredArgsConstructor static class TestPublisher { From d648975880a656c9b88c22d80a4f07161d76d112 Mon Sep 17 00:00:00 2001 From: cwjohnpark Date: Mon, 28 Sep 2026 13:38:33 +0900 Subject: [PATCH 04/15] GH-1901 - Use type name for array types in FormattableType.of(ResolvableType). FormattableType.of(ResolvableType) used Class#getName(), which renders array types by their JVM binary names (e.g. [B or [Ljava.lang.String;). As the instance is cached under the same key that FormattableType.of(Class) uses, the broken instance also leaked into Class-based lookups, so the result depended on call order. We now use Class#getTypeName() in line with FormattableType.of(Class). Original pull request: GH-1902 Signed-off-by: cwjohnpark --- .../modulith/core/FormattableType.java | 2 +- .../core/FormattableTypeUnitTests.java | 28 +++++++++++++++++++ 2 files changed, 29 insertions(+), 1 deletion(-) diff --git a/spring-modulith-core/src/main/java/org/springframework/modulith/core/FormattableType.java b/spring-modulith-core/src/main/java/org/springframework/modulith/core/FormattableType.java index 47baca8d3..718ca2ae3 100644 --- a/spring-modulith-core/src/main/java/org/springframework/modulith/core/FormattableType.java +++ b/spring-modulith-core/src/main/java/org/springframework/modulith/core/FormattableType.java @@ -152,7 +152,7 @@ public static FormattableType of(ResolvableType type) { .toList(); return resolved != null - ? CACHE.computeIfAbsent(type.toString(), __ -> new FormattableType(resolved.getName(), generics)) + ? CACHE.computeIfAbsent(type.toString(), __ -> new FormattableType(resolved.getTypeName(), generics)) : WILDCARD; } diff --git a/spring-modulith-core/src/test/java/org/springframework/modulith/core/FormattableTypeUnitTests.java b/spring-modulith-core/src/test/java/org/springframework/modulith/core/FormattableTypeUnitTests.java index 71d476fec..8795d86e3 100644 --- a/spring-modulith-core/src/test/java/org/springframework/modulith/core/FormattableTypeUnitTests.java +++ b/spring-modulith-core/src/test/java/org/springframework/modulith/core/FormattableTypeUnitTests.java @@ -128,6 +128,32 @@ void rendersDeclaredGenerics() throws Exception { assertThat(type.getAbbreviatedFullName()).isEqualTo("j.u.List>"); } + @Test // GH-1901 + void handlesArrayTypesFromResolvableType() throws Exception { + + var method = Sample.class.getMethod("arrays", Integer[].class, byte[].class, Integer[][].class, List.class); + + assertThat(FormattableType.of(ResolvableType.forMethodParameter(method, 0)).getAbbreviatedFullName()) + .isEqualTo("j.l.Integer[]"); + assertThat(FormattableType.of(ResolvableType.forMethodParameter(method, 1)).getAbbreviatedFullName()) + .isEqualTo("byte[]"); + assertThat(FormattableType.of(ResolvableType.forMethodParameter(method, 2)).getFullName()) + .isEqualTo("java.lang.Integer[][]"); + assertThat(FormattableType.of(ResolvableType.forMethodParameter(method, 3)).getAbbreviatedFullName()) + .isEqualTo("j.u.List"); + } + + @Test // GH-1901 + void rendersArrayTypesConsistentlyIndependentOfLookupOrder() { + + var fromResolvableType = FormattableType.of(ResolvableType.forClass(Sample[].class)); + var fromClass = FormattableType.of(Sample[].class); + + assertThat(fromClass.getFullName()) + .isEqualTo(fromResolvableType.getFullName()) + .isEqualTo("org.springframework.modulith.core.FormattableTypeUnitTests.Sample[]"); + } + interface Sample { List wildcarded(List parameterized); @@ -135,5 +161,7 @@ interface Sample { void classWildcard(Class type); List> genericReturnType(); + + void arrays(Integer[] integers, byte[] bytes, Integer[][] matrix, List list); } } From e75d14ab6cca966fc0531624d2b8d11bccc4cee3 Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Thu, 1 Oct 2026 10:50:53 +0200 Subject: [PATCH 05/15] GH-1901 - Polishing. Refactored test cases. --- .../core/FormattableTypeUnitTests.java | 41 ++++++++++++++----- 1 file changed, 31 insertions(+), 10 deletions(-) diff --git a/spring-modulith-core/src/test/java/org/springframework/modulith/core/FormattableTypeUnitTests.java b/spring-modulith-core/src/test/java/org/springframework/modulith/core/FormattableTypeUnitTests.java index 8795d86e3..e7ebf44d6 100644 --- a/spring-modulith-core/src/test/java/org/springframework/modulith/core/FormattableTypeUnitTests.java +++ b/spring-modulith-core/src/test/java/org/springframework/modulith/core/FormattableTypeUnitTests.java @@ -17,10 +17,15 @@ import static org.assertj.core.api.Assertions.*; +import java.lang.reflect.Method; import java.util.List; import java.util.Map; +import java.util.stream.Stream; +import org.junit.jupiter.api.DynamicTest; +import org.junit.jupiter.api.NamedExecutable; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestFactory; import org.springframework.core.ResolvableType; import org.springframework.modulith.core.FormattableType.NonModuleTypeAbbreviation; @@ -128,19 +133,18 @@ void rendersDeclaredGenerics() throws Exception { assertThat(type.getAbbreviatedFullName()).isEqualTo("j.u.List>"); } - @Test // GH-1901 - void handlesArrayTypesFromResolvableType() throws Exception { + @TestFactory // GH-1901 + Stream handlesArrayTypesFromResolvableType() throws Exception { var method = Sample.class.getMethod("arrays", Integer[].class, byte[].class, Integer[][].class, List.class); - assertThat(FormattableType.of(ResolvableType.forMethodParameter(method, 0)).getAbbreviatedFullName()) - .isEqualTo("j.l.Integer[]"); - assertThat(FormattableType.of(ResolvableType.forMethodParameter(method, 1)).getAbbreviatedFullName()) - .isEqualTo("byte[]"); - assertThat(FormattableType.of(ResolvableType.forMethodParameter(method, 2)).getFullName()) - .isEqualTo("java.lang.Integer[][]"); - assertThat(FormattableType.of(ResolvableType.forMethodParameter(method, 3)).getAbbreviatedFullName()) - .isEqualTo("j.u.List"); + var tests = Stream.of( + new $(method, 0, "j.l.Integer[]"), + new $(method, 1, "byte[]"), + new $(method, 2, "j.l.Integer[][]"), + new $(method, 3, "j.u.List")); + + return DynamicTest.stream(tests); } @Test // GH-1901 @@ -164,4 +168,21 @@ interface Sample { void arrays(Integer[] integers, byte[] bytes, Integer[][] matrix, List list); } + + record $(ResolvableType type, String expected) implements NamedExecutable { + + public $(Method method, int parameterIndex, String expected) { + this(ResolvableType.forMethodParameter(method, parameterIndex), expected); + } + + @Override + public final String toString() { + return "%s renders as %s".formatted(type, expected); + } + + @Override + public void execute() throws Throwable { + assertThat(FormattableType.of(type).getAbbreviatedFullName()).isEqualTo(expected); + } + } } From 5b319b4554d50d48826ef18fc829a02db47672a6 Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Thu, 1 Oct 2026 14:19:23 +0200 Subject: [PATCH 06/15] GH-1903 - Honor registry trigger annotation for publication completion. CompletionRegisteringAdvisor used to decorate all AFTER_COMMIT @TransactionalEventListener methods, no matter whether spring.modulith.events.registry-trigger-annotation excluded them. For an excluded listener no publication is ever registered, so the completion fell back to the payload-based lookup in the repository, which serializes the event and thus failed for events that cannot be serialized. The decision which event listener methods are considered by the Event Publication Registry is now encapsulated in EventListenerMethodMetadata, used by both the advisor and TransactionalEventListeners. --- .../EventPublicationAutoConfiguration.java | 5 +- .../config/EventPublicationConfiguration.java | 7 +- .../core/EventListenerMethodMetadata.java | 178 ++++++++++++++++++ .../core/TransactionalEventListeners.java | 66 +------ .../support/CompletionRegisteringAdvisor.java | 45 +++-- .../EventListenerMethodMetadataUnitTests.java | 115 +++++++++++ .../TransactionalEventListenersUnitTests.java | 28 +-- ...ionRegisteringAdvisorIntegrationTests.java | 7 +- ...CompletionRegisteringAdvisorUnitTests.java | 33 +++- .../antora/modules/ROOT/pages/events.adoc | 3 +- 10 files changed, 376 insertions(+), 111 deletions(-) create mode 100644 spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/EventListenerMethodMetadata.java create mode 100644 spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/EventListenerMethodMetadataUnitTests.java diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventPublicationAutoConfiguration.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventPublicationAutoConfiguration.java index abf9d6602..49cf3472e 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventPublicationAutoConfiguration.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventPublicationAutoConfiguration.java @@ -113,8 +113,9 @@ static DefaultFailedEventPublications failedEventPublications( @Bean @Role(BeanDefinition.ROLE_INFRASTRUCTURE) @ConditionalOnBean(EventPublicationRegistry.class) - static CompletionRegisteringAdvisor completionRegisteringAdvisor(ObjectFactory registry) { - return EventPublicationConfiguration.completionRegisteringAdvisor(registry); + static CompletionRegisteringAdvisor completionRegisteringAdvisor(ObjectFactory registry, + ObjectFactory environment) { + return EventPublicationConfiguration.completionRegisteringAdvisor(registry, environment); } @Bean diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventPublicationConfiguration.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventPublicationConfiguration.java index 99db6221a..8c76b02b0 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventPublicationConfiguration.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/EventPublicationConfiguration.java @@ -28,6 +28,7 @@ import org.springframework.modulith.events.core.DefaultAbandonedEventPublications; import org.springframework.modulith.events.core.DefaultEventPublicationRegistry; import org.springframework.modulith.events.core.DefaultFailedEventPublications; +import org.springframework.modulith.events.core.EventListenerMethodMetadata; import org.springframework.modulith.events.core.EventPublicationRegistry; import org.springframework.modulith.events.core.EventPublicationRepository; import org.springframework.modulith.events.support.CompletionRegisteringAdvisor; @@ -81,7 +82,9 @@ static DefaultFailedEventPublications failedEventPublications( @Bean @Role(BeanDefinition.ROLE_INFRASTRUCTURE) - static CompletionRegisteringAdvisor completionRegisteringAdvisor(ObjectFactory registry) { - return new CompletionRegisteringAdvisor(registry::getObject); + static CompletionRegisteringAdvisor completionRegisteringAdvisor(ObjectFactory registry, + ObjectFactory environment) { + return new CompletionRegisteringAdvisor(registry::getObject, + EventListenerMethodMetadata.of(environment::getObject)); } } diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/EventListenerMethodMetadata.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/EventListenerMethodMetadata.java new file mode 100644 index 000000000..d58f0fefe --- /dev/null +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/EventListenerMethodMetadata.java @@ -0,0 +1,178 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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.springframework.modulith.events.core; + +import java.lang.annotation.Annotation; +import java.lang.reflect.Method; +import java.util.function.Supplier; + +import org.jspecify.annotations.Nullable; +import org.springframework.core.annotation.AnnotatedElementUtils; +import org.springframework.core.env.Environment; +import org.springframework.transaction.event.TransactionPhase; +import org.springframework.transaction.event.TransactionalApplicationListener; +import org.springframework.transaction.event.TransactionalApplicationListenerMethodAdapter; +import org.springframework.transaction.event.TransactionalEventListener; +import org.springframework.util.Assert; +import org.springframework.util.ClassUtils; +import org.springframework.util.ReflectionUtils; +import org.springframework.util.StringUtils; +import org.springframework.util.function.SingletonSupplier; + +/** + * Metadata about event listener methods, in particular which of them are to be considered by the Event Publication + * Registry, as configured in the {@link Environment}. By default, all listeners (meta-)annotated with + * {@link TransactionalEventListener} and listening for {@link TransactionPhase#AFTER_COMMIT} trigger an entry in + * the registry. + * Configuring an annotation via {@code spring.modulith.events.registry-trigger-annotation} narrows that down to the + * listeners carrying it, which in turn causes all other listeners to be left alone by the registry entirely. + * + * @author Oliver Drotbohm + * @since 2.2, 2.1.2 + */ +public class EventListenerMethodMetadata { + + static final String TRIGGER_ANNOTATION_PROPERTY = "spring.modulith.events.registry-trigger-annotation"; + + private static final Method GET_TARGET_METHOD; + + static { + + GET_TARGET_METHOD = ReflectionUtils + .findMethod(TransactionalApplicationListenerMethodAdapter.class, "getTargetMethod"); + ReflectionUtils.makeAccessible(GET_TARGET_METHOD); + } + + private final Supplier<@Nullable Class> triggerAnnotation; + + private EventListenerMethodMetadata(Supplier<@Nullable Class> triggerAnnotation) { + this.triggerAnnotation = triggerAnnotation; + } + + /** + * Creates a new {@link EventListenerMethodMetadata} for the given {@link Environment}. The latter is only accessed on + * first use and the trigger annotation resolved from it cached. + * + * @param environment must not be {@literal null}. + * @return will never be {@literal null}. + */ + public static EventListenerMethodMetadata of(Supplier environment) { + + Assert.notNull(environment, "Environment must not be null!"); + + return new EventListenerMethodMetadata( + SingletonSupplier.ofNullable(() -> findTriggerAnnotation(environment.get()))); + } + + /** + * Creates a new {@link EventListenerMethodMetadata} considering all transactional event listener methods, i.e. the + * arrangement in place if no trigger annotation is configured. + * + * @return will never be {@literal null}. + */ + public static EventListenerMethodMetadata all() { + return new EventListenerMethodMetadata(() -> null); + } + + /** + * Returns whether the given {@link Method} is supposed to trigger an entry in the Event Publication Registry, i.e. + * whether it is an {@link TransactionPhase#AFTER_COMMIT} {@link TransactionalEventListener} carrying the configured + * trigger annotation, if any. + * + * @param method must not be {@literal null}. + */ + public boolean triggersRegistry(Method method) { + + Assert.notNull(method, "Method must not be null!"); + + var annotation = AnnotatedElementUtils.findMergedAnnotation(method, TransactionalEventListener.class); + + return annotation != null + && annotation.phase().equals(TransactionPhase.AFTER_COMMIT) + && hasTriggerAnnotation(method); + } + + /** + * Returns whether the given {@link TransactionalApplicationListener} is supposed to trigger an entry in the Event + * Publication Registry, i.e. whether it listens for {@link TransactionPhase#AFTER_COMMIT} and is backed by a method + * carrying the configured trigger annotation, if any. Listeners not backed by an event listener method are excluded + * as soon as a trigger annotation is configured, as their qualification cannot be established. + * + * @param listener must not be {@literal null}. + */ + boolean triggersRegistry(TransactionalApplicationListener listener) { + + Assert.notNull(listener, "TransactionalApplicationListener must not be null!"); + + if (!listener.getTransactionPhase().equals(TransactionPhase.AFTER_COMMIT)) { + return false; + } + + if (triggerAnnotation.get() == null) { + return true; + } + + if (!(listener instanceof TransactionalApplicationListenerMethodAdapter adapter)) { + return false; + } + + var method = (Method) ReflectionUtils.invokeMethod(GET_TARGET_METHOD, adapter); + + return method != null && hasTriggerAnnotation(method); + } + + /** + * Returns whether the given {@link Method} carries the configured trigger annotation, if any. + * + * @param method must not be {@literal null}. + */ + private boolean hasTriggerAnnotation(Method method) { + + var annotationType = triggerAnnotation.get(); + + return annotationType == null || AnnotatedElementUtils.hasAnnotation(method, annotationType); + } + + /** + * Returns the annotation type configured via {@value #TRIGGER_ANNOTATION_PROPERTY} or {@literal null} if none is + * configured. + * + * @param environment must not be {@literal null}. + */ + @SuppressWarnings("unchecked") + private static @Nullable Class findTriggerAnnotation(Environment environment) { + + var annotationName = environment.getProperty(TRIGGER_ANNOTATION_PROPERTY); + + if (!StringUtils.hasText(annotationName)) { + return null; + } + + try { + + var annotationType = ClassUtils.forName(annotationName, EventListenerMethodMetadata.class.getClassLoader()); + + if (!annotationType.isAnnotation()) { + throw new IllegalStateException("Configured type is not an annotation!"); + } + + return (Class) annotationType; + + } catch (ClassNotFoundException o_O) { + throw new IllegalStateException(o_O); + } + } +} diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/TransactionalEventListeners.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/TransactionalEventListeners.java index 39ea6c152..3be55f29c 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/TransactionalEventListeners.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/TransactionalEventListeners.java @@ -15,8 +15,6 @@ */ package org.springframework.modulith.events.core; -import java.lang.annotation.Annotation; -import java.lang.reflect.Method; import java.util.Collection; import java.util.List; import java.util.function.Consumer; @@ -27,16 +25,10 @@ import org.springframework.context.ApplicationEvent; import org.springframework.context.ApplicationListener; -import org.springframework.core.annotation.AnnotatedElementUtils; import org.springframework.core.annotation.AnnotationAwareOrderComparator; import org.springframework.core.env.Environment; -import org.springframework.transaction.event.TransactionPhase; import org.springframework.transaction.event.TransactionalApplicationListener; -import org.springframework.transaction.event.TransactionalApplicationListenerMethodAdapter; import org.springframework.util.Assert; -import org.springframework.util.ClassUtils; -import org.springframework.util.ReflectionUtils; -import org.springframework.util.StringUtils; /** * First-class collection to work with transactional event listeners, i.e. {@link ApplicationListener} instances that @@ -50,17 +42,6 @@ */ public class TransactionalEventListeners { - static final String TRIGGER_ANNOTATION_PROPERTY = "spring.modulith.events.registry-trigger-annotation"; - - private static final Method GET_TARGET_METHOD; - - static { - - GET_TARGET_METHOD = ReflectionUtils - .findMethod(TransactionalApplicationListenerMethodAdapter.class, "getTargetMethod"); - ReflectionUtils.makeAccessible(GET_TARGET_METHOD); - } - private final List> listeners; /** @@ -68,6 +49,7 @@ public class TransactionalEventListeners { * {@link TransactionalApplicationListener}. * * @param listeners must not be {@literal null}. + * @param environment must not be {@literal null}. */ @SuppressWarnings({ "rawtypes", "unchecked" }) public TransactionalEventListeners(Collection> listeners, @@ -75,11 +57,12 @@ public TransactionalEventListeners(Collection> listeners, Assert.notNull(listeners, "ApplicationListeners must not be null!"); + var metadata = EventListenerMethodMetadata.of(environment); + this.listeners = (List) listeners.stream() .filter(TransactionalApplicationListener.class::isInstance) .map(TransactionalApplicationListener.class::cast) - .filter(it -> it.getTransactionPhase().equals(TransactionPhase.AFTER_COMMIT)) - .filter(byAnnotationFilter(environment)) + .filter(metadata::triggersRegistry) .sorted(AnnotationAwareOrderComparator.INSTANCE) .toList(); } @@ -137,45 +120,4 @@ Stream> stream() { return listeners.stream(); } - /** - * Returns a {@link Predicate} filtering the listeners by the trigger annotation configured in - * {@code spring.modulith.events.annotation}. - * - * @param environment must not be {@literal null}. - * @return will never be {@literal null}. - * @since 2.1 - */ - @SuppressWarnings({ "rawtypes", "unchecked" }) - private static Predicate byAnnotationFilter( - Supplier environment) { - - return listener -> { - - var annotationName = environment.get().getProperty(TRIGGER_ANNOTATION_PROPERTY); - - if (!StringUtils.hasText(annotationName)) { - return true; - } - - try { - - var annotationType = ClassUtils.forName(annotationName, TransactionalEventListeners.class.getClassLoader()); - - if (!annotationType.isAnnotation()) { - throw new IllegalStateException("Configured type is not an annotation!"); - } - - if (!(listener instanceof TransactionalApplicationListenerMethodAdapter)) { - return false; - } - - var method = (Method) ReflectionUtils.invokeMethod(GET_TARGET_METHOD, listener); - - return AnnotatedElementUtils.hasAnnotation(method, (Class) annotationType); - - } catch (ClassNotFoundException o_O) { - throw new IllegalStateException(o_O); - } - }; - } } diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisor.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisor.java index af29a313d..d16746251 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisor.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisor.java @@ -32,10 +32,9 @@ import org.springframework.aop.support.StaticMethodMatcher; import org.springframework.aop.support.annotation.AnnotationMatchingPointcut; import org.springframework.core.Ordered; -import org.springframework.core.annotation.AnnotatedElementUtils; +import org.springframework.modulith.events.core.EventListenerMethodMetadata; import org.springframework.modulith.events.core.EventPublicationRegistry; import org.springframework.modulith.events.core.PublicationTargetIdentifier; -import org.springframework.transaction.event.TransactionPhase; import org.springframework.transaction.event.TransactionalApplicationListenerMethodAdapter; import org.springframework.transaction.event.TransactionalEventListener; import org.springframework.util.Assert; @@ -55,13 +54,33 @@ public class CompletionRegisteringAdvisor extends AbstractPointcutAdvisor { private final Advice advice; /** - * Creates a new {@link CompletionRegisteringAdvisor} for the given {@link EventPublicationRegistry}. + * Creates a new {@link CompletionRegisteringAdvisor} for the given {@link EventPublicationRegistry}, decorating all + * {@link TransactionalEventListener} annotated methods, no matter whether they are selected via + * {@code spring.modulith.events.registry-trigger-annotation}. * * @param registry must not be {@literal null}. + * @deprecated since 2.2, 2.1.2, for removal in 2.3. Use + * {@link #CompletionRegisteringAdvisor(Supplier, EventListenerMethodMetadata)} instead to honor the + * configured registry trigger annotation. */ + @Deprecated(since = "2.2, 2.1.2", forRemoval = true) public CompletionRegisteringAdvisor(Supplier registry) { + this(registry, EventListenerMethodMetadata.all()); + } + + /** + * Creates a new {@link CompletionRegisteringAdvisor} for the given {@link EventPublicationRegistry}, only decorating + * methods that actually trigger an entry in the former according to the given {@link EventListenerMethodMetadata}. + * + * @param registry must not be {@literal null}. + * @param metadata must not be {@literal null}. + * @since 2.2, 2.1.2 + */ + public CompletionRegisteringAdvisor(Supplier registry, + EventListenerMethodMetadata metadata) { Assert.notNull(registry, "EventPublicationRegistry must not be null!"); + Assert.notNull(metadata, "EventListenerMethodMetadata must not be null!"); this.pointcut = new AnnotationMatchingPointcut(null, TransactionalEventListener.class, true) { @@ -71,7 +90,7 @@ public CompletionRegisteringAdvisor(Supplier registry) */ @Override public MethodMatcher getMethodMatcher() { - return new CommitListenerMethodMatcher(super.getMethodMatcher()); + return new CommitListenerMethodMatcher(super.getMethodMatcher(), metadata); } }; @@ -103,14 +122,19 @@ public Advice getAdvice() { private static class CommitListenerMethodMatcher extends StaticMethodMatcher { private final MethodMatcher delegate; + private final EventListenerMethodMetadata metadata; /** - * Creates a new {@link CommitListenerMethodMatcher} with the given delegate {@link MethodMatcher}. + * Creates a new {@link CommitListenerMethodMatcher} with the given delegate {@link MethodMatcher} and + * {@link EventListenerMethodMetadata}. * * @param delegate must not be {@literal null}. + * @param metadata must not be {@literal null}. */ - public CommitListenerMethodMatcher(MethodMatcher delegate) { + public CommitListenerMethodMatcher(MethodMatcher delegate, EventListenerMethodMetadata metadata) { + this.delegate = delegate; + this.metadata = metadata; } /* @@ -119,14 +143,7 @@ public CommitListenerMethodMatcher(MethodMatcher delegate) { */ @Override public boolean matches(Method method, Class targetClass) { - - if (!delegate.matches(method, targetClass)) { - return false; - } - - var annotation = AnnotatedElementUtils.findMergedAnnotation(method, TransactionalEventListener.class); - - return annotation != null && annotation.phase().equals(TransactionPhase.AFTER_COMMIT); + return delegate.matches(method, targetClass) && metadata.triggersRegistry(method); } } diff --git a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/EventListenerMethodMetadataUnitTests.java b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/EventListenerMethodMetadataUnitTests.java new file mode 100644 index 000000000..c3567183d --- /dev/null +++ b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/EventListenerMethodMetadataUnitTests.java @@ -0,0 +1,115 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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.springframework.modulith.events.core; + +import static org.assertj.core.api.Assertions.*; +import static org.springframework.modulith.events.core.EventListenerMethodMetadata.*; + +import java.lang.reflect.Method; + +import org.junit.jupiter.api.Test; +import org.springframework.context.event.EventListener; +import org.springframework.mock.env.MockEnvironment; +import org.springframework.modulith.events.ApplicationModuleListener; +import org.springframework.transaction.event.TransactionPhase; +import org.springframework.transaction.event.TransactionalEventListener; +import org.springframework.util.ReflectionUtils; + +/** + * Unit tests for {@link EventListenerMethodMetadata}. + * + * @author Oliver Drotbohm + */ +class EventListenerMethodMetadataUnitTests { + + MockEnvironment environment = new MockEnvironment(); + + @Test // GH-1630 + void rejectsNotLoadableTriggerAnnotation() { + + environment.setProperty(TRIGGER_ANNOTATION_PROPERTY, "some.non.loadable.Type"); + + var metadata = EventListenerMethodMetadata.of(() -> environment); + + assertThatIllegalStateException().isThrownBy(() -> metadata.triggersRegistry(method("plain"))); + } + + @Test // GH-1630 + void rejectsNonAnnotationTypeForTriggerAnnotation() { + + environment.setProperty(TRIGGER_ANNOTATION_PROPERTY, "java.lang.String"); + + var metadata = EventListenerMethodMetadata.of(() -> environment); + + assertThatIllegalStateException().isThrownBy(() -> metadata.triggersRegistry(method("plain"))); + } + + @Test // GH-1903 + void considersAllListenerMethodsByDefault() { + + var metadata = EventListenerMethodMetadata.of(() -> environment); + + assertThat(metadata.triggersRegistry(method("plain"))).isTrue(); + assertThat(metadata.triggersRegistry(method("moduleListener"))).isTrue(); + } + + @Test // GH-1903 + void considersTriggerAnnotation() { + + environment.setProperty(TRIGGER_ANNOTATION_PROPERTY, ApplicationModuleListener.class.getName()); + + var metadata = EventListenerMethodMetadata.of(() -> environment); + + assertThat(metadata.triggersRegistry(method("plain"))).isFalse(); + assertThat(metadata.triggersRegistry(method("moduleListener"))).isTrue(); + } + + @Test // GH-1903 + void onlyConsidersAfterCommitListenerMethods() { + + var metadata = EventListenerMethodMetadata.all(); + + assertThat(metadata.triggersRegistry(method("afterRollback"))).isFalse(); + assertThat(metadata.triggersRegistry(method("plainEventListener"))).isFalse(); + assertThat(metadata.triggersRegistry(method("nonListener"))).isFalse(); + } + + private static Method method(String name) { + + var method = ReflectionUtils.findMethod(SampleListener.class, name, Object.class); + + assertThat(method).isNotNull(); + + return method; + } + + static class SampleListener { + + @TransactionalEventListener + void plain(Object event) {} + + @ApplicationModuleListener + void moduleListener(Object event) {} + + @TransactionalEventListener(phase = TransactionPhase.AFTER_ROLLBACK) + void afterRollback(Object event) {} + + @EventListener + void plainEventListener(Object event) {} + + void nonListener(Object event) {} + } +} diff --git a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/TransactionalEventListenersUnitTests.java b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/TransactionalEventListenersUnitTests.java index 103957b6d..d60484fde 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/TransactionalEventListenersUnitTests.java +++ b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/TransactionalEventListenersUnitTests.java @@ -16,7 +16,7 @@ package org.springframework.modulith.events.core; import static org.assertj.core.api.Assertions.*; -import static org.springframework.modulith.events.core.TransactionalEventListeners.*; +import static org.springframework.modulith.events.core.EventListenerMethodMetadata.*; import java.util.List; @@ -69,32 +69,6 @@ void considersTriggerAnnotation() { assertThat(listeners.stream()).containsExactly(second); } - @Test // GH-1630 - void rejectsNotLoadableTriggerAnnotation() { - - var environment = new MockEnvironment(); - environment.setProperty(TRIGGER_ANNOTATION_PROPERTY, "some.non.loadable.Type"); - - var second = getAdapter(ModuleListener.class, "on", SampleEvent.class); - - assertThatIllegalStateException().isThrownBy(() -> { - new TransactionalEventListeners(List.of(second), () -> environment); - }); - } - - @Test // GH-1630 - void rejectsNonAnnotationTypeForTriggerAnnotation() { - - var environment = new MockEnvironment(); - environment.setProperty(TRIGGER_ANNOTATION_PROPERTY, "java.lang.String"); - - var second = getAdapter(ModuleListener.class, "on", SampleEvent.class); - - assertThatIllegalStateException().isThrownBy(() -> { - new TransactionalEventListeners(List.of(second), () -> environment); - }); - } - private static TransactionalApplicationListenerMethodAdapter getAdapter(Class type, String methodName, Class parameter) { diff --git a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisorIntegrationTests.java b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisorIntegrationTests.java index 0ad28bd38..655fa501f 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisorIntegrationTests.java +++ b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisorIntegrationTests.java @@ -27,6 +27,8 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Role; +import org.springframework.core.env.Environment; +import org.springframework.modulith.events.core.EventListenerMethodMetadata; import org.springframework.modulith.events.core.EventPublicationRegistry; import org.springframework.modulith.events.support.CompletionRegisteringAdvisor.CompletionRegisteringMethodInterceptor; import org.springframework.scheduling.annotation.AnnotationAsyncExecutionInterceptor; @@ -64,11 +66,12 @@ PlatformTransactionManager transactionManager() { @Bean @Role(BeanDefinition.ROLE_INFRASTRUCTURE) - static CompletionRegisteringAdvisor completionRegisteringAdvisor() { + static CompletionRegisteringAdvisor completionRegisteringAdvisor(Environment environment) { var publicationRegistry = mock(EventPublicationRegistry.class); - return new CompletionRegisteringAdvisor(() -> publicationRegistry); + return new CompletionRegisteringAdvisor(() -> publicationRegistry, + EventListenerMethodMetadata.of(() -> environment)); } } diff --git a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisorUnitTests.java b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisorUnitTests.java index a8a0a12ef..136da788c 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisorUnitTests.java +++ b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/support/CompletionRegisteringAdvisorUnitTests.java @@ -26,6 +26,9 @@ import org.springframework.aop.framework.Advised; import org.springframework.aop.framework.ProxyFactory; import org.springframework.context.event.EventListener; +import org.springframework.mock.env.MockEnvironment; +import org.springframework.modulith.events.ApplicationModuleListener; +import org.springframework.modulith.events.core.EventListenerMethodMetadata; import org.springframework.modulith.events.core.EventPublicationRegistry; import org.springframework.scheduling.annotation.Async; import org.springframework.transaction.event.TransactionPhase; @@ -38,8 +41,11 @@ */ class CompletionRegisteringAdvisorUnitTests { + static final String TRIGGER_ANNOTATION_PROPERTY = "spring.modulith.events.registry-trigger-annotation"; + EventPublicationRegistry registry = mock(EventPublicationRegistry.class); SomeEventListener bean = new SomeEventListener(); + MockEnvironment environment = new MockEnvironment(); @Test void triggersCompletionForAfterCommitEventListener() throws Exception { @@ -97,6 +103,27 @@ void exposesResultForCompletableFuture() throws Exception { assertThat(future.get()).isNotNull(); } + @Test // GH-1903 + void doesNotTriggerCompletionForListenerMissingTriggerAnnotation() { + + environment.setProperty(TRIGGER_ANNOTATION_PROPERTY, ApplicationModuleListener.class.getName()); + + createProxyFor(bean).onAfterCommit(new Object()); + + verify(registry, never()).markProcessing(any(), any()); + verify(registry, never()).markCompleted(any(), any()); + } + + @Test // GH-1903 + void triggersCompletionForListenerCarryingTriggerAnnotation() { + + environment.setProperty(TRIGGER_ANNOTATION_PROPERTY, ApplicationModuleListener.class.getName()); + + createProxyFor(bean).onModuleEvent(new Object()); + + verify(registry).markCompleted(any(), any()); + } + private void assertCompletion(BiConsumer consumer) { assertCompletion(consumer, true); } @@ -120,7 +147,8 @@ private void assertCompletion(BiConsumer consumer, bo private T createProxyFor(T bean) { ProxyFactory factory = new ProxyFactory(bean); - factory.addAdvisor(new CompletionRegisteringAdvisor(() -> registry)); + factory.addAdvisor(new CompletionRegisteringAdvisor(() -> registry, + EventListenerMethodMetadata.of(() -> environment))); return (T) factory.getProxy(); } @@ -132,6 +160,9 @@ void onAfterCommit(Object event) {} @TransactionalEventListener(phase = TransactionPhase.AFTER_ROLLBACK) void onAfterRollback(Object object) {} + @ApplicationModuleListener + void onModuleEvent(Object event) {} + @EventListener void simpleEventListener(Object object) {} diff --git a/src/docs/antora/modules/ROOT/pages/events.adoc b/src/docs/antora/modules/ROOT/pages/events.adoc index feba439cd..31486aa9b 100644 --- a/src/docs/antora/modules/ROOT/pages/events.adoc +++ b/src/docs/antora/modules/ROOT/pages/events.adoc @@ -206,7 +206,8 @@ If you want to customize this, check out the xref:appendix.adoc#configuration-pr .The transactional event listener arrangement before execution image::event-publication-registry-start.png[] -Each transactional event listener is wrapped into an aspect that marks that log entry as completed if the execution of the listener succeeds. +Each transactional event listener that triggers such a log entry is wrapped into an aspect that marks the entry as completed if the execution of the listener succeeds. +Listeners excluded via the trigger annotation described above are left untouched entirely. In case the listener fails, the log entry stays untouched so that retry mechanisms can be deployed depending on the application's needs. Automatic re-publication of the events can be enabled via the xref:appendix.adoc#configuration-properties[`spring.modulith.events.republish-outstanding-events-on-restart`] property. From 6acf573d78d91c2dccbd88a502203c3482066ea0 Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Fri, 2 Oct 2026 12:55:57 +0200 Subject: [PATCH 07/15] GH-1917 - Skip stale-publication handling for statuses without configured staleness. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit DefaultEventPublicationRegistry now only marks publications as failed for statuses Staleness.isMonitored(…) reports as monitored, so configuring only the resubmitted timeout no longer fails PUBLISHED and PROCESSING publications. StalenessProperties rejects zero and negative durations and refuses staleness lookups for statuses that cannot become stale. --- .../events/config/StalenessProperties.java | 25 ++-- .../core/DefaultEventPublicationRegistry.java | 4 + .../modulith/events/core/Staleness.java | 44 ++++++- .../config/StalenessPropertiesUnitTests.java | 119 ++++++++++++++++++ ...aultEventPublicationRegistryUnitTests.java | 19 +++ 5 files changed, 202 insertions(+), 9 deletions(-) create mode 100644 spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/config/StalenessPropertiesUnitTests.java diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/StalenessProperties.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/StalenessProperties.java index b9fcfb6df..a4cb75663 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/StalenessProperties.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/config/StalenessProperties.java @@ -65,17 +65,18 @@ public class StalenessProperties implements Staleness { @Nullable Duration resubmitted, @Nullable Duration checkInterval) { - this.published = published == null ? Duration.ZERO : published; - this.processing = processing == null ? Duration.ZERO : processing; - this.resubmitted = resubmitted == null ? Duration.ZERO : resubmitted; - this.checkInterval = checkInterval == null ? Duration.ofMinutes(1) : checkInterval; + this.published = published == null ? UNCONFIGURED_DURATION : requirePositive(published, "published"); + this.processing = processing == null ? UNCONFIGURED_DURATION : requirePositive(processing, "processing"); + this.resubmitted = resubmitted == null ? UNCONFIGURED_DURATION : requirePositive(resubmitted, "resubmitted"); + this.checkInterval = checkInterval == null ? Duration.ofMinutes(1) + : requirePositive(checkInterval, "checkInterval"); } boolean monitorStaleness() { - return !published.equals(Duration.ZERO) - || !processing.equals(Duration.ZERO) - || !resubmitted.equals(Duration.ZERO); + return !published.equals(UNCONFIGURED_DURATION) + || !processing.equals(UNCONFIGURED_DURATION) + || !resubmitted.equals(UNCONFIGURED_DURATION); } Duration getCheckInterval() { @@ -95,7 +96,15 @@ public Duration getStaleness(Status status) { case PUBLISHED -> published; case PROCESSING -> processing; case RESUBMITTED -> resubmitted; - default -> Duration.ZERO; + default -> throw new IllegalArgumentException("Unsupported status: " + status); }; } + + private static Duration requirePositive(Duration duration, String name) { + + Assert.isTrue(!duration.isZero() && !duration.isNegative(), + () -> "Staleness property '%s' must be greater than zero but was %s!".formatted(name, duration)); + + return duration; + } } diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/DefaultEventPublicationRegistry.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/DefaultEventPublicationRegistry.java index b3e0bbe72..90dfcf491 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/DefaultEventPublicationRegistry.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/DefaultEventPublicationRegistry.java @@ -409,6 +409,10 @@ private void propagateStateTransition(Object event, PublicationTargetIdentifier private void markFailed(Status status, Staleness staleness) { + if (!staleness.isMonitored(status)) { + return; + } + var duration = staleness.getStaleness(status); var reference = clock.instant().minus(duration); var result = events.findByStatus(status).stream() diff --git a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/Staleness.java b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/Staleness.java index cbf94c0f1..21f8202c7 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/Staleness.java +++ b/spring-modulith-events/spring-modulith-events-core/src/main/java/org/springframework/modulith/events/core/Staleness.java @@ -28,12 +28,54 @@ */ public interface Staleness { + /** + * The {@link Duration} to be used to indicate that no staleness has been configured for a particular {@link Status}, + * i.e. that publications in that state are not supposed to be checked for staleness. + * + * @since 2.0.9, 2.1.2, 2.2 + */ + Duration UNCONFIGURED_DURATION = Duration.ZERO; + + /** + * Returns whether {@link org.springframework.modulith.events.EventPublication}s with the given {@link Status} are + * supposed to be checked for staleness at all. Callers are expected to inspect this before calling + * {@link #getStaleness(Status)}. Returns {@literal false}, and does not fail, for {@link Status}es that cannot become + * stale. + * + * @param status must not be {@literal null}. + * @return whether staleness needs to be handled for the given {@link Status}. + * @since 2.0.9, 2.1.2, 2.2 + */ + default boolean isMonitored(Status status) { + return isStaleable(status) && !UNCONFIGURED_DURATION.equals(getStaleness(status)); + } + + /** + * Returns whether {@link org.springframework.modulith.events.EventPublication}s with the given {@link Status} can + * become stale at all. + * + * @param status must not be {@literal null}. + * @return whether the given {@link Status} can become stale. + * @since 2.0.9, 2.1.2, 2.2 + */ + default boolean isStaleable(Status status) { + + return switch (status) { + case PUBLISHED, PROCESSING, RESUBMITTED -> true; + default -> false; + }; + } + /** * Returns the {@link Duration} after which {@link org.springframework.modulith.events.EventPublication}s with a - * certain {@link Status} are considered stale. + * certain {@link Status} are considered stale. Guaranteed to be positive if {@link #isMonitored(Status)} returned + * {@literal true} before. Rejects {@link Status}es not considered {@link #isStaleable(Status) staleable}. * * @param status must not be {@literal null}. * @return will never be {@literal null}. + * @throws IllegalArgumentException in case the given {@link Status} cannot become stale. + * @see #isMonitored(Status) + * @see #isStaleable(Status) */ Duration getStaleness(Status status); } diff --git a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/config/StalenessPropertiesUnitTests.java b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/config/StalenessPropertiesUnitTests.java new file mode 100644 index 000000000..3488620a4 --- /dev/null +++ b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/config/StalenessPropertiesUnitTests.java @@ -0,0 +1,119 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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.springframework.modulith.events.config; + +import static org.assertj.core.api.Assertions.*; +import static org.springframework.modulith.events.EventPublication.Status.*; + +import java.time.Duration; +import java.util.function.Supplier; +import java.util.stream.Stream; + +import org.junit.jupiter.api.DynamicTest; +import org.junit.jupiter.api.NamedExecutable; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestFactory; + +/** + * Unit tests for {@link StalenessProperties}. + * + * @author Oliver Drotbohm + */ +class StalenessPropertiesUnitTests { + + @TestFactory // GH-1917 + Stream rejectsZeroAndNegativeDurations() { + + var negative = Duration.ofSeconds(-1); + + var tests = Stream.of( + new $("published", Duration.ZERO), + new $("published", negative), + new $("processing", Duration.ZERO), + new $("processing", negative), + new $("resubmitted", Duration.ZERO), + new $("resubmitted", negative), + new $("checkInterval", Duration.ZERO), + new $("checkInterval", negative)); + + return DynamicTest.stream(tests); + } + + @Test // GH-1917 + void acceptsPositiveAndUnsetDurations() { + + var properties = new StalenessProperties(null, null, Duration.ofMinutes(5), null); + + assertThat(properties.monitorStaleness()).isTrue(); + assertThat(properties.getStaleness(PROCESSING)).isEqualTo(Duration.ZERO); + } + + @Test // GH-1917 + void onlyMonitorsStatusesWithConfiguredStaleness() { + + var properties = new StalenessProperties(null, null, Duration.ofMinutes(5), null); + + assertThat(properties.isMonitored(RESUBMITTED)).isTrue(); + assertThat(properties.isMonitored(PUBLISHED)).isFalse(); + assertThat(properties.isMonitored(PROCESSING)).isFalse(); + } + + @Test // GH-1917 + void doesNotMonitorStatusesThatCannotBecomeStale() { + + var properties = new StalenessProperties(Duration.ofMinutes(1), Duration.ofMinutes(2), Duration.ofMinutes(3), + null); + + assertThat(properties.isMonitored(COMPLETED)).isFalse(); + assertThat(properties.isMonitored(FAILED)).isFalse(); + } + + @Test // GH-1917 + void rejectsStalenessLookupForStatusesThatCannotBecomeStale() { + + var properties = new StalenessProperties(Duration.ofMinutes(1), Duration.ofMinutes(2), Duration.ofMinutes(3), + null); + + assertThatIllegalArgumentException().isThrownBy(() -> properties.getStaleness(COMPLETED)); + assertThatIllegalArgumentException().isThrownBy(() -> properties.getStaleness(FAILED)); + } + + record $(String property, Duration value, Supplier factory) implements NamedExecutable { + + $(String property, Duration value) { + this(property, value, switch (property) { + case "published" -> () -> new StalenessProperties(value, null, null, null); + case "processing" -> () -> new StalenessProperties(null, value, null, null); + case "resubmitted" -> () -> new StalenessProperties(null, null, value, null); + case "checkInterval" -> () -> new StalenessProperties(null, null, null, value); + default -> throw new IllegalArgumentException("Unknown property " + property); + }); + } + + @Override + public final String toString() { + return "Rejects %s for %s".formatted(value, property); + } + + @Override + public void execute() { + + assertThatIllegalArgumentException() + .isThrownBy(factory::get) + .withMessageContaining(property); + } + } +} diff --git a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/DefaultEventPublicationRegistryUnitTests.java b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/DefaultEventPublicationRegistryUnitTests.java index 2b4cfd026..7712ab8b1 100644 --- a/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/DefaultEventPublicationRegistryUnitTests.java +++ b/spring-modulith-events/spring-modulith-events-core/src/test/java/org/springframework/modulith/events/core/DefaultEventPublicationRegistryUnitTests.java @@ -241,6 +241,25 @@ void doesNotConsiderRecentlyResubmittedPublicationStaleBasedOnOriginalPublicatio verify(repository, never()).markFailed(any()); } + @Test // GH-1917 + void doesNotMarkPublicationsFailedForStatusWithZeroStaleness() { + + var now = Instant.now(); + + var publication = mock(TargetEventPublication.class, CALLS_REAL_METHODS); + lenient().when(publication.getPublicationDate()).thenReturn(now.minusSeconds(1)); + lenient().when(publication.getLastResubmissionDate()).thenReturn(null); + + lenient().when(repository.findByStatus(Status.PUBLISHED)).thenReturn(List.of(publication)); + lenient().when(repository.findByStatus(Status.PROCESSING)).thenReturn(List.of(publication)); + when(repository.findByStatus(Status.RESUBMITTED)).thenReturn(Collections.emptyList()); + + createRegistry(now).markStalePublicationsFailed( + status -> status == Status.RESUBMITTED ? Duration.ofMinutes(5) : Duration.ZERO); + + verify(repository, never()).markFailed(any()); + } + private DefaultEventPublicationRegistry createRegistry(Instant instant) { var clock = Clock.fixed(instant, ZoneId.systemDefault()); From 67989011c3107493a84db93eeb71b11644a89e7c Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Fri, 2 Oct 2026 13:56:29 +0200 Subject: [PATCH 08/15] GH-1920 - Upgrade to Develocity conventions 0.0.27. --- .mvn/extensions.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.mvn/extensions.xml b/.mvn/extensions.xml index 25506efe4..e871cd89e 100644 --- a/.mvn/extensions.xml +++ b/.mvn/extensions.xml @@ -3,6 +3,6 @@ io.spring.develocity.conventions develocity-conventions-maven-extension - 0.0.26 + 0.0.27 From c103395eedd6909eab57648c6d55246808723b48 Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Fri, 2 Oct 2026 13:57:14 +0200 Subject: [PATCH 09/15] GH-1921 - Upgrade Maven Wrapper to 3.10. --- .mvn/wrapper/maven-wrapper.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.mvn/wrapper/maven-wrapper.properties b/.mvn/wrapper/maven-wrapper.properties index 216df0589..bb58b3c26 100644 --- a/.mvn/wrapper/maven-wrapper.properties +++ b/.mvn/wrapper/maven-wrapper.properties @@ -1,3 +1,3 @@ wrapperVersion=3.3.4 distributionType=only-script -distributionUrl=https://repo.maven.apache.org/maven2/org/apache/maven/apache-maven/3.9.16/apache-maven-3.9.16-bin.zip +distributionUrl=https://repo.maven.apache.org/maven2/org/apache/maven/apache-maven/3.10.0/apache-maven-3.10.0-bin.zip From 5be012354421eb8fd3dc7c9eecad5b4783b3c6b8 Mon Sep 17 00:00:00 2001 From: char-yb Date: Wed, 7 Oct 2026 10:33:28 +0900 Subject: [PATCH 10/15] GH-1925 - Prevent completed event publications from being resubmitted or marked as failed. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The status transitions in the JPA, JDBC and Neo4j event publication repositories only checked that a publication was not already in the target state. A publication that completed after having been looked up as failed or stale could thus still be claimed for resubmission or be flipped back to FAILED, causing the listener to be invoked again for an already processed event. The statements used by markFailed(…), markProcessing(…) and markResubmitted(…) now also require the completion date to be unset, in line with the MongoDB repository. Signed-off-by: char-yb --- .../JdbcEventPublicationRepositoryV2.java | 2 ++ ...blicationRepositoryV2IntegrationTests.java | 21 +++++++++++++++ .../jpa/JpaEventPublicationRepository.java | 2 ++ ...PublicationRepositoryIntegrationTests.java | 27 +++++++++++++++++++ .../Neo4jEventPublicationRepository.java | 2 ++ .../Neo4jEventPublicationRepositoryTest.java | 21 +++++++++++++++ 6 files changed, 75 insertions(+) diff --git a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2.java b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2.java index ca482cc7d..9cea083ca 100644 --- a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2.java +++ b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2.java @@ -109,6 +109,7 @@ private static String getUpdateSql(String table, Status status) { STATUS = '%s' WHERE ID = ? + AND COMPLETION_DATE IS NULL AND (STATUS IS NULL OR STATUS != '%s') """.formatted(table, status.name(), status.name())); } @@ -418,6 +419,7 @@ public boolean markResubmitted(UUID identifier, Instant instant) { LAST_RESUBMISSION_DATE = ? WHERE ID = ? + AND COMPLETION_DATE IS NULL AND (STATUS IS NULL OR STATUS != 'RESUBMITTED') """.formatted(settings.getTable())); diff --git a/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2IntegrationTests.java b/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2IntegrationTests.java index 25735a065..d35742949 100644 --- a/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2IntegrationTests.java +++ b/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2IntegrationTests.java @@ -661,6 +661,27 @@ void marksPublicationWithNullStatusColumnAsFailed() { .containsExactly(publication.getIdentifier()); } + @Test // GH-1925 + void doesNotResubmitCompletedPublication() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + + assertThat(repository.markResubmitted(publication.getIdentifier(), Instant.now())).isFalse(); + } + + @Test // GH-1925 + void doesNotMarkCompletedPublicationFailed() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + repository.markFailed(publication.getIdentifier()); + + assertThat(repository.findByStatus(Status.FAILED)).isEmpty(); + } + /** * Simulates a publication persisted by a schema version that predates the {@code STATUS} column, i.e. one for * which the column was never backfilled and is {@literal null} rather than defaulted. diff --git a/spring-modulith-events/spring-modulith-events-jpa/src/main/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepository.java b/spring-modulith-events/spring-modulith-events-jpa/src/main/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepository.java index 5651170de..7f5e0ab34 100644 --- a/spring-modulith-events/spring-modulith-events-jpa/src/main/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepository.java +++ b/spring-modulith-events/spring-modulith-events-jpa/src/main/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepository.java @@ -154,6 +154,7 @@ class JpaEventPublicationRepository implements EventPublicationRepository { update DefaultJpaEventPublication p set p.status = ?1 where p.id = ?2 + and p.completionDate is null and (status is null or status != ?1) """; @@ -163,6 +164,7 @@ class JpaEventPublicationRepository implements EventPublicationRepository { p.completionAttempts = p.completionAttempts + 1, p.lastResubmissionDate = ?1 where p.id = ?2 + and p.completionDate is null and (p.status is null or p.status != org.springframework.modulith.events.EventPublication$Status.RESUBMITTED) """; diff --git a/spring-modulith-events/spring-modulith-events-jpa/src/test/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepositoryIntegrationTests.java b/spring-modulith-events/spring-modulith-events-jpa/src/test/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepositoryIntegrationTests.java index fdfb0671c..b4ff8b18b 100644 --- a/spring-modulith-events/spring-modulith-events-jpa/src/test/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepositoryIntegrationTests.java +++ b/spring-modulith-events/spring-modulith-events-jpa/src/test/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepositoryIntegrationTests.java @@ -619,6 +619,33 @@ void marksPublicationWithNullStatusColumnAsFailed() { .containsExactly(publication.getIdentifier()); } + @Test // GH-1925 + void doesNotResubmitCompletedPublication() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + em.flush(); + em.clear(); + + assertThat(repository.markResubmitted(publication.getIdentifier(), Instant.now())).isFalse(); + } + + @Test // GH-1925 + void doesNotMarkCompletedPublicationFailed() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + em.flush(); + em.clear(); + + repository.markFailed(publication.getIdentifier()); + em.clear(); + + assertThat(repository.findByStatus(Status.FAILED)).isEmpty(); + } + /** * Simulates a publication persisted by a schema version that predates the {@code status} column, i.e. one for * which the column was never backfilled and is {@literal null} rather than defaulted. diff --git a/spring-modulith-events/spring-modulith-events-neo4j/src/main/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepository.java b/spring-modulith-events/spring-modulith-events-neo4j/src/main/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepository.java index 3fa70401a..cf33cb38e 100644 --- a/spring-modulith-events/spring-modulith-events-neo4j/src/main/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepository.java +++ b/spring-modulith-events/spring-modulith-events-neo4j/src/main/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepository.java @@ -160,6 +160,7 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository { private static final Statement UPDATE_STATUS_STATEMENT = match(EVENT_PUBLICATION_NODE) .where(EVENT_PUBLICATION_NODE.property(ID).eq(parameter(ID))) + .and(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNull()) .and(EVENT_PUBLICATION_NODE.property(STATUS).isNull() .or(EVENT_PUBLICATION_NODE.property(STATUS).ne(parameter(STATUS)))) .set(EVENT_PUBLICATION_NODE.property(STATUS).to(parameter(STATUS))) @@ -167,6 +168,7 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository { private static final Statement RESUBMIT_STATEMENT = match(EVENT_PUBLICATION_NODE) .where(EVENT_PUBLICATION_NODE.property(ID).eq(parameter(ID))) + .and(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNull()) .and(EVENT_PUBLICATION_NODE.property(STATUS).isNull() .or(EVENT_PUBLICATION_NODE.property(STATUS).ne(literalOf(Status.RESUBMITTED.name())))) .set(EVENT_PUBLICATION_NODE.property(STATUS).to(parameter(STATUS))) diff --git a/spring-modulith-events/spring-modulith-events-neo4j/src/test/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepositoryTest.java b/spring-modulith-events/spring-modulith-events-neo4j/src/test/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepositoryTest.java index 00b72d343..d88109c99 100644 --- a/spring-modulith-events/spring-modulith-events-neo4j/src/test/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepositoryTest.java +++ b/spring-modulith-events/spring-modulith-events-neo4j/src/test/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepositoryTest.java @@ -581,6 +581,27 @@ void marksPublicationWithNullStatusPropertyAsFailed() { .containsExactly(publication.getIdentifier()); } + @Test // GH-1925 + void doesNotResubmitCompletedPublication() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + + assertThat(repository.markResubmitted(publication.getIdentifier(), Instant.now())).isFalse(); + } + + @Test // GH-1925 + void doesNotMarkCompletedPublicationFailed() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + repository.markFailed(publication.getIdentifier()); + + assertThat(repository.findByStatus(EventPublication.Status.FAILED)).isEmpty(); + } + /** * Simulates a publication persisted by a schema version that predates the {@code status} property, i.e. one for * which the property is missing rather than defaulted. From 2b217a0f59e3f55228334a7fb91bf1c59f5756aa Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Fri, 9 Oct 2026 01:37:12 +0530 Subject: [PATCH 11/15] GH-1927 - Shared modules do not list themselves as allowed dependencies. --- .../modulith/core/ApplicationModule.java | 1 + .../modulith/core/ModuleUnitTest.java | 15 ++++++++++++ .../modulith/core/TestUtils.java | 15 ++++++++++++ .../java/reproducers/gh1927/Application.java | 21 +++++++++++++++++ .../java/reproducers/gh1927/order/Order.java | 18 +++++++++++++++ .../reproducers/gh1927/shared/Shared.java | 23 +++++++++++++++++++ .../gh1927/shared/package-info.java | 2 ++ 7 files changed, 95 insertions(+) create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1927/Application.java create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1927/order/Order.java create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1927/shared/Shared.java create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1927/shared/package-info.java diff --git a/spring-modulith-core/src/main/java/org/springframework/modulith/core/ApplicationModule.java b/spring-modulith-core/src/main/java/org/springframework/modulith/core/ApplicationModule.java index 050672a94..104228355 100644 --- a/spring-modulith-core/src/main/java/org/springframework/modulith/core/ApplicationModule.java +++ b/spring-modulith-core/src/main/java/org/springframework/modulith/core/ApplicationModule.java @@ -548,6 +548,7 @@ public AllowedDependencies getAllowedDependencies(ApplicationModules modules) { .map(it -> AllowedDependency.of(it, this, modules)); var sharedDependencies = modules.getSharedModules().stream() + .filter(it -> !it.equals(this)) .map(AllowedDependency::to); return Stream.concat(explicitlyDeclaredModules, sharedDependencies) // diff --git a/spring-modulith-core/src/test/java/org/springframework/modulith/core/ModuleUnitTest.java b/spring-modulith-core/src/test/java/org/springframework/modulith/core/ModuleUnitTest.java index c2b296721..eb93772ab 100644 --- a/spring-modulith-core/src/test/java/org/springframework/modulith/core/ModuleUnitTest.java +++ b/spring-modulith-core/src/test/java/org/springframework/modulith/core/ModuleUnitTest.java @@ -127,4 +127,19 @@ void obtainsDependenciesForCyclicArrangement() { .contains(it.getIdentifier()); }); } + + @Test // GH-1927 + void doesNotListSharedModuleAsOwnAllowedDependency() { + + var modules = TestUtils.of(reproducers.gh1927.Application.class); + var shared = modules.getModuleByName("shared").orElseThrow(); + + assertThat(shared.getAllowedDependencies(modules)).isEmpty(); + + assertThatExceptionOfType(Violations.class) // + .isThrownBy(modules::verify) // + .satisfies(ex -> { + assertThat(ex.getMessages()).anySatisfy(message -> assertThat(message).contains("Allowed targets: none")); + }); + } } diff --git a/spring-modulith-core/src/test/java/org/springframework/modulith/core/TestUtils.java b/spring-modulith-core/src/test/java/org/springframework/modulith/core/TestUtils.java index 83bf9f12f..e0dd2f0f2 100644 --- a/spring-modulith-core/src/test/java/org/springframework/modulith/core/TestUtils.java +++ b/spring-modulith-core/src/test/java/org/springframework/modulith/core/TestUtils.java @@ -57,6 +57,21 @@ public static ApplicationModules of(String basePackage, String... ignoredPackage return of(ModulithMetadata.of(basePackage), JavaClass.Predicates.resideInAnyPackage(ignoredPackages)); } + /** + * Creates an {@link ApplicationModules} instance from the given modulith type but only inspecting the test code. + * Contrary to {@link #of(String, String...)}, this honors the configuration declared on the type, e.g. shared modules + * via {@link org.springframework.modulith.Modulithic}. + * + * @param modulithType must not be {@literal null}. + * @return will never be {@literal null}. + */ + public static ApplicationModules of(Class modulithType) { + + Assert.notNull(modulithType, "Modulith type must not be null!"); + + return ApplicationModules.of(modulithType, new ImportOption.OnlyIncludeTests()); + } + /** * Returns all {@link Classes} of this module. * diff --git a/spring-modulith-core/src/test/java/reproducers/gh1927/Application.java b/spring-modulith-core/src/test/java/reproducers/gh1927/Application.java new file mode 100644 index 000000000..73c1942c8 --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1927/Application.java @@ -0,0 +1,21 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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 reproducers.gh1927; + +import org.springframework.modulith.Modulithic; + +@Modulithic(sharedModules = "shared") +public class Application {} diff --git a/spring-modulith-core/src/test/java/reproducers/gh1927/order/Order.java b/spring-modulith-core/src/test/java/reproducers/gh1927/order/Order.java new file mode 100644 index 000000000..a33109c3e --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1927/order/Order.java @@ -0,0 +1,18 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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 reproducers.gh1927.order; + +public class Order {} diff --git a/spring-modulith-core/src/test/java/reproducers/gh1927/shared/Shared.java b/spring-modulith-core/src/test/java/reproducers/gh1927/shared/Shared.java new file mode 100644 index 000000000..18a3df07d --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1927/shared/Shared.java @@ -0,0 +1,23 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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 reproducers.gh1927.shared; + +import reproducers.gh1927.order.Order; + +class Shared { + + Order order; +} diff --git a/spring-modulith-core/src/test/java/reproducers/gh1927/shared/package-info.java b/spring-modulith-core/src/test/java/reproducers/gh1927/shared/package-info.java new file mode 100644 index 000000000..3981b8ba1 --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1927/shared/package-info.java @@ -0,0 +1,2 @@ +@org.springframework.modulith.ApplicationModule(allowedDependencies = {}) +package reproducers.gh1927.shared; From 93f6c53d75e4e0ba78bb5444f836cd1c4f601ced Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Fri, 9 Oct 2026 11:32:51 +0200 Subject: [PATCH 12/15] GH-1910 - Let backport script also change ticket references. --- etc/backport-ticket.sh | 4 ++++ etc/functions.sh | 26 ++++++++++++++++++++++++++ 2 files changed, 30 insertions(+) diff --git a/etc/backport-ticket.sh b/etc/backport-ticket.sh index 1ce235b35..943bb9e7c 100755 --- a/etc/backport-ticket.sh +++ b/etc/backport-ticket.sh @@ -115,6 +115,10 @@ do updateCommitMessage "$sourceGh" "$targetGh" echo "Updated commit message" + # Replace @Test ticket references in the cherry-picked test sources with the new one + updateTestReferences "$sourceGh" "$targetGh" || _exit 1 "Failed to update test references for commit $sha" + echo "Updated test references" + done <<< "$shas" done diff --git a/etc/functions.sh b/etc/functions.sh index e7da2ff43..a41082d38 100755 --- a/etc/functions.sh +++ b/etc/functions.sh @@ -214,3 +214,29 @@ updateCommitMessage() { return 1 fi } + +# Replaces references to the source ticket on @Test lines (e.g. "@Test // GH-1234") in the Java files touched by +# the HEAD commit with the target ticket and amends the commit if anything changed. +updateTestReferences() { + local source="$1" + local target="$2" + local files + local changed=false + + files=$(git diff-tree --no-commit-id --name-only -r --diff-filter=AM HEAD | grep '\.java$') + + while IFS= read -r file; do + [ -f "$file" ] || continue + SOURCE_GH="$source" TARGET_GH="$target" perl -i -pe 's/\b\Q$ENV{SOURCE_GH}\E\b/$ENV{TARGET_GH}/g if /\@Test\b/' "$file" + if ! git diff --quiet -- "$file"; then + git add -- "$file" + changed=true + fi + done <<< "$files" + + if [ "$changed" == "true" ]; then + git commit --amend --no-edit || return 1 + fi + + return 0 +} From 0ffe07e62acd995a7ffe7c3000c876ec155c5aa7 Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Fri, 9 Oct 2026 14:12:42 +0200 Subject: [PATCH 13/15] GH-1924 - Support to explicitly depend on everything an ApplicationModule exposes via module :: **. --- .../modulith/ApplicationModule.java | 4 +- .../modulith/core/ApplicationModule.java | 238 ++++++++++++++---- .../modulith/core/NamedInterfaces.java | 39 +++ .../modulith/core/ModuleUnitTest.java | 29 ++- .../java/reproducers/gh1924/Application.java | 21 ++ .../gh1924/inventory/Inventory.java | 18 ++ .../gh1924/inventory/package-info.java | 19 ++ .../reproducers/gh1924/order/OrderId.java | 18 ++ .../gh1924/order/events/OrderShipped.java | 18 ++ .../gh1924/order/events/package-info.java | 19 ++ .../gh1924/order/internal/Hidden.java | 18 ++ .../modules/ROOT/pages/fundamentals.adoc | 34 +++ 12 files changed, 421 insertions(+), 54 deletions(-) create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1924/Application.java create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1924/inventory/Inventory.java create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1924/inventory/package-info.java create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1924/order/OrderId.java create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1924/order/events/OrderShipped.java create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1924/order/events/package-info.java create mode 100644 spring-modulith-core/src/test/java/reproducers/gh1924/order/internal/Hidden.java diff --git a/spring-modulith-api/src/main/java/org/springframework/modulith/ApplicationModule.java b/spring-modulith-api/src/main/java/org/springframework/modulith/ApplicationModule.java index 99e7f24b3..f2aa8e6c7 100644 --- a/spring-modulith-api/src/main/java/org/springframework/modulith/ApplicationModule.java +++ b/spring-modulith-api/src/main/java/org/springframework/modulith/ApplicationModule.java @@ -51,7 +51,9 @@ * {@link Modulith}/{@link Modulithic} will be allowed, too. Names listed are local ones, unless the application has * configured {@link Modulithic#useFullyQualifiedModuleNames()} to {@literal true}. Explicit references to * {@link NamedInterface}s need to be separated by a double colon {@code ::}, e.g. {@code module::API} if - * {@code module} is the logical module name and {@code API} is the name of the named interface. + * {@code module} is the logical module name and {@code API} is the name of the named interface. Use {@code module::*} + * to refer to all explicitly declared named interfaces of the module (not including its base package) and + * {@code module::**} to refer to everything the module exposes, i.e. the base package and all named interfaces. *

* Declaring an empty array will allow no dependencies to other modules. To not restrict the dependencies at all, * leave the attribute at its default value. diff --git a/spring-modulith-core/src/main/java/org/springframework/modulith/core/ApplicationModule.java b/spring-modulith-core/src/main/java/org/springframework/modulith/core/ApplicationModule.java index 104228355..78ffb16a4 100644 --- a/spring-modulith-core/src/main/java/org/springframework/modulith/core/ApplicationModule.java +++ b/spring-modulith-core/src/main/java/org/springframework/modulith/core/ApplicationModule.java @@ -540,20 +540,13 @@ public AllowedDependencies getAllowedDependencies(ApplicationModules modules) { var allowedDependencyNames = information.getDeclaredDependencies(); - if (AllowedDependencies.isOpen(allowedDependencyNames)) { - return AllowedDependencies.open(); - } - - var explicitlyDeclaredModules = allowedDependencyNames.stream() // - .map(it -> AllowedDependency.of(it, this, modules)); - - var sharedDependencies = modules.getSharedModules().stream() - .filter(it -> !it.equals(this)) - .map(AllowedDependency::to); - - return Stream.concat(explicitlyDeclaredModules, sharedDependencies) // - .distinct() // - .collect(Collectors.collectingAndThen(Collectors.toList(), AllowedDependencies::closed)); + return AllowedDependencies.isOpen(allowedDependencyNames) + ? AllowedDependencies.open() + : AllowedDependencies.NONE + // Explicitly declared + .and(AllowedDependencies.of(this, allowedDependencyNames, modules)) + // All shared modules except itself + .and(AllowedDependencies.to(modules.getSharedModules().stream().filter(Predicate.not(this::equals)))); } /** @@ -933,43 +926,50 @@ public static class AllowedDependency { private static final String INVALID_EXPLICIT_MODULE_DEPENDENCY = "Invalid explicit module dependency in %s! No module found with name '%s'."; private static final String INVALID_NAMED_INTERFACE_DECLARATION = "No named interface named '%s' found! Original dependency declaration: %s -> %s."; - private static final String WILDCARD = "*"; + private static final String ALL_NAMED_INTERFACES = "*"; + private static final String EVERYTHING = "**"; private final ApplicationModule target; - private final @Nullable NamedInterface namedInterface; + private final NamedInterface namedInterface; /** * Creates a new {@link AllowedDependency} for the given {@link ApplicationModule} and {@link NamedInterface}. * * @param target must not be {@literal null}. - * @param namedInterface can be {@literal null}. + * @param namedInterface must not be {@literal null}. */ - private AllowedDependency(ApplicationModule target, @Nullable NamedInterface namedInterface) { + private AllowedDependency(ApplicationModule target, NamedInterface namedInterface) { Assert.notNull(target, "Target ApplicationModule must not be null!"); + Assert.notNull(namedInterface, "NamedInterface must not be null!"); this.target = target; this.namedInterface = namedInterface; } /** - * Creates an {@link AllowedDependency} to the module and optionally named interface defined by the given - * identifier. + * Creates {@link AllowedDependency} instances to the module and optionally named interface defined by the given + * expression. The wildcard variants are expanded into the individual named interfaces they stand for. * - * @param identifier must not be {@literal null} or empty. Follows the - * {@code ${moduleName}(::${namedInterfaceName})} pattern. * @param source the source module of the dependency, must not be {@literal null}. + * @param expression must not be {@literal null} or empty. Follows the + * {@code ${moduleName}(::${namedInterfaceName})} pattern, where the interface name can be {@code *} to + * refer to all explicitly declared named interfaces or {@code **} to additionally include the unnamed + * interface, i.e. the module's base package. * @param modules must not be {@literal null}. * @return will never be {@literal null}. * @throws IllegalArgumentException in case the given identifier is invalid, i.e. does not refer to an existing * module or named interface. + * @see org.springframework.modulith.ApplicationModule#allowedDependencies() */ - static AllowedDependency of(String identifier, ApplicationModule source, + static AllowedDependencies of(ApplicationModule source, String expression, ApplicationModules modules) { - Assert.hasText(identifier, "Module dependency identifier must not be null or empty!"); + Assert.notNull(source, "Source ApplicationModule must not be null!"); + Assert.hasText(expression, "Module dependency expression must not be null or empty!"); + Assert.notNull(modules, "ApplicationModules must not be null!"); - var segments = identifier.split("::"); + var segments = expression.split("::"); var targetModuleName = segments[0].trim(); var namedInterfaceName = segments.length > 1 ? segments[1].trim() : null; @@ -977,19 +977,28 @@ static AllowedDependency of(String identifier, ApplicationModule source, .orElseThrow(() -> new IllegalArgumentException( INVALID_EXPLICIT_MODULE_DEPENDENCY.formatted(source.getIdentifier(), targetModuleName))); - if (WILDCARD.equals(namedInterfaceName)) { - return new AllowedDependency(target, null); + var namedInterfaces = target.getNamedInterfaces(); + + if (ALL_NAMED_INTERFACES.equals(namedInterfaceName)) { + + return AllowedDependencies.closed(namedInterfaces.namedOnly().stream() + .map(it -> toNamedInterface(target, it)).toList()); + } + + if (EVERYTHING.equals(namedInterfaceName)) { + + return AllowedDependencies.closed(namedInterfaces.stream() + .map(it -> toNamedInterface(target, it)).toList()); } - var namedInterfaces = target.getNamedInterfaces(); var namedInterface = namedInterfaceName == null ? namedInterfaces.getUnnamedInterface() : namedInterfaces.getByName(namedInterfaceName) .orElseThrow(() -> new IllegalArgumentException( INVALID_NAMED_INTERFACE_DECLARATION.formatted(namedInterfaceName, source.getIdentifier(), - identifier))); + expression))); - return new AllowedDependency(target, namedInterface); + return AllowedDependencies.closed(List.of(toNamedInterface(target, namedInterface))); } /** @@ -1002,7 +1011,21 @@ static AllowedDependency to(ApplicationModule module) { Assert.notNull(module, "ApplicationModule must not be null!"); - return new AllowedDependency(module, module.getNamedInterfaces().getUnnamedInterface()); + return toNamedInterface(module, module.getNamedInterfaces().getUnnamedInterface()); + } + + /** + * Creates a new {@link AllowedDependency} to the given {@link NamedInterface} of the given + * {@link ApplicationModule}. Passing the module's unnamed interface refers to its base package, which corresponds + * to a plain {@code module} declaration. + * + * @param module must not be {@literal null}. + * @param namedInterface must not be {@literal null}. + * @return will never be {@literal null}. + * @since 2.2 + */ + private static AllowedDependency toNamedInterface(ApplicationModule module, NamedInterface namedInterface) { + return new AllowedDependency(module, namedInterface); } /** @@ -1015,14 +1038,24 @@ public ApplicationModule getTargetModule() { } /** - * Returns the {@link NamedInterface} declared as valid target if declared. + * Returns the {@link NamedInterface} declared as valid target. * - * @return can be {@literal null}. + * @return will never be {@literal null}. */ - public @Nullable NamedInterface getTargetNamedInterface() { + public NamedInterface getTargetNamedInterface() { return namedInterface; } + /** + * Returns whether the {@link AllowedDependency} refers to the given {@link ApplicationModule}. + * + * @param module must not be {@literal null}. + * @since 2.2 + */ + boolean refersTo(ApplicationModule module) { + return target.equals(module); + } + /** * Returns whether the {@link AllowedDependency} contains the given {@link JavaClass}. * @@ -1033,9 +1066,7 @@ boolean contains(JavaClass type) { Assert.notNull(type, "Type must not be null!"); - return namedInterface == null - ? target.getNamedInterfaces().containsInExplicitInterface(type) - : namedInterface.contains(type); + return namedInterface.contains(type); } /** @@ -1048,9 +1079,7 @@ boolean contains(Class type) { Assert.notNull(type, "Type must not be null!"); - return namedInterface == null - ? target.getNamedInterfaces().containsInExplicitInterface(type) - : namedInterface.contains(type); + return namedInterface.contains(type); } /* @@ -1061,17 +1090,12 @@ boolean contains(Class type) { public String toString() { var result = target.getIdentifier().toString(); - var ni = namedInterface; - if (ni == null) { - return result + " :: " + WILDCARD; - } - - if (ni.isUnnamed()) { + if (namedInterface.isUnnamed()) { return result; } - return result + " :: " + ni.getName(); + return result + " :: " + namedInterface.getName(); } /* @@ -1112,22 +1136,88 @@ public int hashCode() { public static class AllowedDependencies implements Iterable { private static final String OPEN_TOKEN = "¯\\_(ツ)_/¯"; + private static final AllowedDependencies NONE = AllowedDependencies.closed(Collections.emptyList()); private final List dependencies; private final boolean closed; + /** + * Returns whether the given declared dependencies consist of nothing but the open token, i.e. whether the + * dependencies are not restricted. + * + * @param AllowedDependencies must not be {@literal null}. + */ static boolean isOpen(List AllowedDependencies) { return AllowedDependencies.size() == 1 && AllowedDependencies.get(0).equals(OPEN_TOKEN); } + /** + * Creates an {@link AllowedDependencies} instance that does not restrict dependencies, i.e. every dependency is + * considered allowed. Used for modules that do not declare any allowed dependencies explicitly. + * + * @return will never be {@literal null}. + */ static AllowedDependencies open() { return new AllowedDependencies(Collections.emptyList(), false); } + /** + * Creates an {@link AllowedDependencies} instance that only allows the given {@link AllowedDependency} instances. + * An empty {@link List} means no dependency is allowed at all. + * + * @param dependencies must not be {@literal null}. + * @return will never be {@literal null}. + */ static AllowedDependencies closed(List dependencies) { return new AllowedDependencies(dependencies, true); } + /** + * Creates a new AllowedDependecies from the given {@link ApplicationModule} to the targets defined in the + * expression. + * + * @param source the source {@link ApplicationModule}, must not be {@literal null}. + * @param expression the expression identifying the targets, must not be {@literal null} or empty. + * @param modules all {@link ApplicationModules} available, must not be {@literal null}.. + * @return will never be {@literal null}. + * @since 2.2 + */ + static AllowedDependencies of(ApplicationModule source, String expression, ApplicationModules modules) { + return AllowedDependency.of(source, expression, modules); + } + + /** + * Creates a new {@link AllowedDependencies} from the given {@link ApplicationModule} to the targets defined by all + * of the given expressions. Targets referred to by multiple expressions are only included once. + * + * @param source the source {@link ApplicationModule}, must not be {@literal null}. + * @param expression the expressions identifying the targets, must not be {@literal null}. + * @param modules all {@link ApplicationModules} available, must not be {@literal null}. + * @return will never be {@literal null}. + * @throws IllegalArgumentException in case any of the expressions is invalid. + * @since 2.2 + * @see #of(ApplicationModule, String, ApplicationModules) + */ + private static AllowedDependencies of(ApplicationModule source, Collection expression, + ApplicationModules modules) { + + return expression.stream() + .map(it -> AllowedDependencies.of(source, it, modules)) + .reduce(AllowedDependencies.NONE, AllowedDependencies::and); + } + + /** + * Creates a new {@link AllowedDependencies} to the unnamed interface of each of the given + * {@link ApplicationModule}s. + * + * @param modules must not be {@literal null}. + * @return will never be {@literal null}. + * @since 2.2 + */ + private static AllowedDependencies to(Stream modules) { + return AllowedDependencies.closed(modules.map(AllowedDependency::to).toList()); + } + /** * Creates a new {@link AllowedDependencies} for the given {@link List} of {@link AllowedDependency}. * @@ -1157,6 +1247,60 @@ public Stream stream() { return dependencies.stream(); } + /** + * Creates a new {@link AllowedDependencies} adding the given {@link AllowedDependency} instance unless they're + * already present. + * + * @param others must not be {@literal null}. + * @return will never be {@literal null}. + * @since 2.2 + */ + public AllowedDependencies and(Iterable others) { + + Assert.notNull(others, "AllowedDependencies must not be null!"); + + var result = new ArrayList<>(dependencies); + + for (var dependency : others) { + if (!result.contains(dependency)) { + result.add(dependency); + } + } + + return AllowedDependencies.closed(result); + } + + /** + * Returns all {@link NamedInterfaces} targeted by the {@link AllowedDependency} instances. + * + * @return will never be {@literal null}. + * @since 2.2 + */ + NamedInterfaces getTargetNamedInterfaces() { + + return dependencies.stream() + .map(it -> it.getTargetNamedInterface()) + .collect(NamedInterfaces.collector()); + } + + /** + * Returns whether the {@link AllowedDependency} instances refer to exactly the given {@link NamedInterface}s of the + * given {@link ApplicationModule}, i.e. there is a dependency for each of them and no dependency refers to anything + * else. + * + * @param module must not be {@literal null}. + * @param namedInterfaces must not be {@literal null}. + * @since 2.2 + */ + boolean referTo(ApplicationModule module, NamedInterfaces namedInterfaces) { + + Assert.notNull(module, "ApplicationModule must not be null!"); + Assert.notNull(namedInterfaces, "NamedInterfaces must not be null!"); + + return dependencies.stream().allMatch(it -> it.refersTo(module)) + && getTargetNamedInterfaces().equals(namedInterfaces); + } + /** * Returns whether the given {@link JavaClass} is a valid dependency. * diff --git a/spring-modulith-core/src/main/java/org/springframework/modulith/core/NamedInterfaces.java b/spring-modulith-core/src/main/java/org/springframework/modulith/core/NamedInterfaces.java index 9cf3e4420..6be2f3186 100644 --- a/spring-modulith-core/src/main/java/org/springframework/modulith/core/NamedInterfaces.java +++ b/spring-modulith-core/src/main/java/org/springframework/modulith/core/NamedInterfaces.java @@ -25,9 +25,11 @@ import java.util.List; import java.util.Optional; import java.util.function.Predicate; +import java.util.stream.Collector; import java.util.stream.Collectors; import java.util.stream.Stream; +import org.jspecify.annotations.Nullable; import org.springframework.core.annotation.AnnotatedElementUtils; import org.springframework.modulith.core.Types.JavaTypes; import org.springframework.util.Assert; @@ -287,6 +289,19 @@ boolean containsInExplicitInterface(Class type) { .anyMatch(NamedInterface::isNamed); } + /** + * Returns only explicitly named {@link NamedInterface}s. + * + * @return will never be {@literal null}. + * @since 2.2 + */ + NamedInterfaces namedOnly() { + + return namedInterfaces.stream() + .filter(NamedInterface::isNamed) + .collect(NamedInterfaces.collector()); + } + /* * (non-Javadoc) * @see java.lang.Object#toString() @@ -299,6 +314,30 @@ public String toString() { .collect(Collectors.joining(System.lineSeparator())); } + /* + * (non-Javadoc) + * @see java.lang.Object#equals(java.lang.Object) + */ + @Override + public boolean equals(@Nullable Object obj) { + + return this == obj + || obj instanceof NamedInterfaces that && this.namedInterfaces.equals(that.namedInterfaces); + } + + /* + * (non-Javadoc) + * @see java.lang.Object#hashCode() + */ + @Override + public int hashCode() { + return this.namedInterfaces.hashCode(); + } + + static Collector collector() { + return Collectors.collectingAndThen(Collectors.toList(), NamedInterfaces::new); + } + private static NamedInterfaces of(NamedInterface interfaces) { return new NamedInterfaces(List.of(interfaces)); } diff --git a/spring-modulith-core/src/test/java/org/springframework/modulith/core/ModuleUnitTest.java b/spring-modulith-core/src/test/java/org/springframework/modulith/core/ModuleUnitTest.java index eb93772ab..956a951f8 100644 --- a/spring-modulith-core/src/test/java/org/springframework/modulith/core/ModuleUnitTest.java +++ b/spring-modulith-core/src/test/java/org/springframework/modulith/core/ModuleUnitTest.java @@ -17,8 +17,9 @@ import static org.assertj.core.api.Assertions.*; -import example.ni.api.ApiType; -import example.ni.spi.SpiType; +import reproducers.gh1924.order.OrderId; +import reproducers.gh1924.order.events.OrderShipped; +import reproducers.gh1924.order.internal.Hidden; import java.util.List; @@ -27,7 +28,7 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.TestInstance; import org.junit.jupiter.api.TestInstance.Lifecycle; -import org.springframework.modulith.core.ApplicationModule.AllowedDependency; +import org.springframework.modulith.core.ApplicationModule.AllowedDependencies; import com.acme.withatbean.SampleAggregate; import com.acme.withatbean.TestEvents.JMoleculesAnnotated; @@ -109,10 +110,10 @@ void wildcardedDeclaredDependencyAllowsDependenciesToAllNamedInterfaces() { var modules = TestUtils.of("example", "example.ninvalid"); var module = modules.getModuleByName("ni").orElseThrow(); - var dependency = AllowedDependency.of("ni :: *", module, modules); + var namedInterfaces = module.getNamedInterfaces().namedOnly(); + var dependencies = AllowedDependencies.of(module, "ni :: *", modules); - assertThat(dependency.contains(SpiType.class)).isTrue(); - assertThat(dependency.contains(ApiType.class)).isTrue(); + assertThat(dependencies.referTo(module, namedInterfaces)).isTrue(); } @Test // GH-1299 @@ -142,4 +143,20 @@ void doesNotListSharedModuleAsOwnAllowedDependency() { assertThat(ex.getMessages()).anySatisfy(message -> assertThat(message).contains("Allowed targets: none")); }); } + + @Test // GH-1924 + void allowsAccessToEverythingExposedViaDoubleAsterisk() { + + var modules = TestUtils.of(reproducers.gh1924.Application.class); + + var inventory = modules.getModuleByName("inventory").orElseThrow(); + var order = modules.getModuleByName("order").orElseThrow(); + + var allowed = inventory.getAllowedDependencies(modules); + + assertThat(allowed.referTo(order, order.getNamedInterfaces())).isTrue(); + assertThat(allowed.isAllowedDependency(OrderId.class)).isTrue(); + assertThat(allowed.isAllowedDependency(OrderShipped.class)).isTrue(); + assertThat(allowed.isAllowedDependency(Hidden.class)).isFalse(); + } } diff --git a/spring-modulith-core/src/test/java/reproducers/gh1924/Application.java b/spring-modulith-core/src/test/java/reproducers/gh1924/Application.java new file mode 100644 index 000000000..67cdf9049 --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1924/Application.java @@ -0,0 +1,21 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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 reproducers.gh1924; + +import org.springframework.modulith.Modulithic; + +@Modulithic +public class Application {} diff --git a/spring-modulith-core/src/test/java/reproducers/gh1924/inventory/Inventory.java b/spring-modulith-core/src/test/java/reproducers/gh1924/inventory/Inventory.java new file mode 100644 index 000000000..e4b565ac3 --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1924/inventory/Inventory.java @@ -0,0 +1,18 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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 reproducers.gh1924.inventory; + +class Inventory {} diff --git a/spring-modulith-core/src/test/java/reproducers/gh1924/inventory/package-info.java b/spring-modulith-core/src/test/java/reproducers/gh1924/inventory/package-info.java new file mode 100644 index 000000000..5e06cf1ca --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1924/inventory/package-info.java @@ -0,0 +1,19 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ +@ApplicationModule(allowedDependencies = "order :: **") +package reproducers.gh1924.inventory; + +import org.springframework.modulith.ApplicationModule; diff --git a/spring-modulith-core/src/test/java/reproducers/gh1924/order/OrderId.java b/spring-modulith-core/src/test/java/reproducers/gh1924/order/OrderId.java new file mode 100644 index 000000000..667e6c2aa --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1924/order/OrderId.java @@ -0,0 +1,18 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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 reproducers.gh1924.order; + +public class OrderId {} diff --git a/spring-modulith-core/src/test/java/reproducers/gh1924/order/events/OrderShipped.java b/spring-modulith-core/src/test/java/reproducers/gh1924/order/events/OrderShipped.java new file mode 100644 index 000000000..bf4c108fc --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1924/order/events/OrderShipped.java @@ -0,0 +1,18 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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 reproducers.gh1924.order.events; + +public class OrderShipped {} diff --git a/spring-modulith-core/src/test/java/reproducers/gh1924/order/events/package-info.java b/spring-modulith-core/src/test/java/reproducers/gh1924/order/events/package-info.java new file mode 100644 index 000000000..058d92c45 --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1924/order/events/package-info.java @@ -0,0 +1,19 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ +@NamedInterface("events") +package reproducers.gh1924.order.events; + +import org.springframework.modulith.NamedInterface; diff --git a/spring-modulith-core/src/test/java/reproducers/gh1924/order/internal/Hidden.java b/spring-modulith-core/src/test/java/reproducers/gh1924/order/internal/Hidden.java new file mode 100644 index 000000000..836fed563 --- /dev/null +++ b/spring-modulith-core/src/test/java/reproducers/gh1924/order/internal/Hidden.java @@ -0,0 +1,18 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed 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 + * + * https://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 reproducers.gh1924.order.internal; + +public class Hidden {} diff --git a/src/docs/antora/modules/ROOT/pages/fundamentals.adoc b/src/docs/antora/modules/ROOT/pages/fundamentals.adoc index 8cd753257..3d7b55adc 100644 --- a/src/docs/antora/modules/ROOT/pages/fundamentals.adoc +++ b/src/docs/antora/modules/ROOT/pages/fundamentals.adoc @@ -412,6 +412,40 @@ class ModuleMetadata {} ---- ====== +Note that the asterisk does *not* cover the application module's base package. +To allow access to everything an application module exposes -- its base package *and* all explicitly declared named interfaces -- use a double asterisk (``**``): + +.Using the double asterisk to declare allowed dependencies to everything a module exposes +[tabs] +====== +Java:: ++ +[source, java, role="primary", chomp="none"] +---- +@org.springframework.modulith.ApplicationModule( + allowedDependencies = "order :: **" +) +package example.inventory; +---- +Kotlin:: ++ +[source, kotlin, role="secondary", chomp="none"] +---- +package example.inventory + +import org.springframework.modulith.ApplicationModule +import org.springframework.modulith.PackageInfo + +@ApplicationModule( + allowedDependencies = "order :: **" +) +@PackageInfo +class ModuleMetadata {} +---- +====== + +This is equivalent to declaring `{ "order", "order :: *" }`. + If you require more generic control about the named interfaces of an application module, check out xref:fundamentals.adoc#customizing-named-interfaces[the customization section]. [[customizing-modules-arrangement]] From 0a4e2669d5cfe378e615206b4bf8482404ae3a28 Mon Sep 17 00:00:00 2001 From: char-yb Date: Wed, 7 Oct 2026 10:33:28 +0900 Subject: [PATCH 14/15] GH-1925 - Prevent completed event publications from being resubmitted or marked as failed. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The status transitions in the JPA, JDBC and Neo4j event publication repositories only checked that a publication was not already in the target state. A publication that completed after having been looked up as failed or stale could thus still be claimed for resubmission or be flipped back to FAILED, causing the listener to be invoked again for an already processed event. The statements used by markFailed(…), markProcessing(…) and markResubmitted(…) now also require the completion date to be unset, in line with the MongoDB repository. Signed-off-by: char-yb --- .../JdbcEventPublicationRepositoryV2.java | 2 ++ ...blicationRepositoryV2IntegrationTests.java | 21 +++++++++++++++ .../jpa/JpaEventPublicationRepository.java | 2 ++ ...PublicationRepositoryIntegrationTests.java | 27 +++++++++++++++++++ .../Neo4jEventPublicationRepository.java | 2 ++ .../Neo4jEventPublicationRepositoryTest.java | 21 +++++++++++++++ 6 files changed, 75 insertions(+) diff --git a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2.java b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2.java index ca482cc7d..9cea083ca 100644 --- a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2.java +++ b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2.java @@ -109,6 +109,7 @@ private static String getUpdateSql(String table, Status status) { STATUS = '%s' WHERE ID = ? + AND COMPLETION_DATE IS NULL AND (STATUS IS NULL OR STATUS != '%s') """.formatted(table, status.name(), status.name())); } @@ -418,6 +419,7 @@ public boolean markResubmitted(UUID identifier, Instant instant) { LAST_RESUBMISSION_DATE = ? WHERE ID = ? + AND COMPLETION_DATE IS NULL AND (STATUS IS NULL OR STATUS != 'RESUBMITTED') """.formatted(settings.getTable())); diff --git a/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2IntegrationTests.java b/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2IntegrationTests.java index 25735a065..d35742949 100644 --- a/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2IntegrationTests.java +++ b/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepositoryV2IntegrationTests.java @@ -661,6 +661,27 @@ void marksPublicationWithNullStatusColumnAsFailed() { .containsExactly(publication.getIdentifier()); } + @Test // GH-1925 + void doesNotResubmitCompletedPublication() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + + assertThat(repository.markResubmitted(publication.getIdentifier(), Instant.now())).isFalse(); + } + + @Test // GH-1925 + void doesNotMarkCompletedPublicationFailed() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + repository.markFailed(publication.getIdentifier()); + + assertThat(repository.findByStatus(Status.FAILED)).isEmpty(); + } + /** * Simulates a publication persisted by a schema version that predates the {@code STATUS} column, i.e. one for * which the column was never backfilled and is {@literal null} rather than defaulted. diff --git a/spring-modulith-events/spring-modulith-events-jpa/src/main/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepository.java b/spring-modulith-events/spring-modulith-events-jpa/src/main/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepository.java index 5651170de..7f5e0ab34 100644 --- a/spring-modulith-events/spring-modulith-events-jpa/src/main/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepository.java +++ b/spring-modulith-events/spring-modulith-events-jpa/src/main/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepository.java @@ -154,6 +154,7 @@ class JpaEventPublicationRepository implements EventPublicationRepository { update DefaultJpaEventPublication p set p.status = ?1 where p.id = ?2 + and p.completionDate is null and (status is null or status != ?1) """; @@ -163,6 +164,7 @@ class JpaEventPublicationRepository implements EventPublicationRepository { p.completionAttempts = p.completionAttempts + 1, p.lastResubmissionDate = ?1 where p.id = ?2 + and p.completionDate is null and (p.status is null or p.status != org.springframework.modulith.events.EventPublication$Status.RESUBMITTED) """; diff --git a/spring-modulith-events/spring-modulith-events-jpa/src/test/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepositoryIntegrationTests.java b/spring-modulith-events/spring-modulith-events-jpa/src/test/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepositoryIntegrationTests.java index fdfb0671c..b4ff8b18b 100644 --- a/spring-modulith-events/spring-modulith-events-jpa/src/test/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepositoryIntegrationTests.java +++ b/spring-modulith-events/spring-modulith-events-jpa/src/test/java/org/springframework/modulith/events/jpa/JpaEventPublicationRepositoryIntegrationTests.java @@ -619,6 +619,33 @@ void marksPublicationWithNullStatusColumnAsFailed() { .containsExactly(publication.getIdentifier()); } + @Test // GH-1925 + void doesNotResubmitCompletedPublication() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + em.flush(); + em.clear(); + + assertThat(repository.markResubmitted(publication.getIdentifier(), Instant.now())).isFalse(); + } + + @Test // GH-1925 + void doesNotMarkCompletedPublicationFailed() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + em.flush(); + em.clear(); + + repository.markFailed(publication.getIdentifier()); + em.clear(); + + assertThat(repository.findByStatus(Status.FAILED)).isEmpty(); + } + /** * Simulates a publication persisted by a schema version that predates the {@code status} column, i.e. one for * which the column was never backfilled and is {@literal null} rather than defaulted. diff --git a/spring-modulith-events/spring-modulith-events-neo4j/src/main/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepository.java b/spring-modulith-events/spring-modulith-events-neo4j/src/main/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepository.java index 3fa70401a..cf33cb38e 100644 --- a/spring-modulith-events/spring-modulith-events-neo4j/src/main/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepository.java +++ b/spring-modulith-events/spring-modulith-events-neo4j/src/main/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepository.java @@ -160,6 +160,7 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository { private static final Statement UPDATE_STATUS_STATEMENT = match(EVENT_PUBLICATION_NODE) .where(EVENT_PUBLICATION_NODE.property(ID).eq(parameter(ID))) + .and(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNull()) .and(EVENT_PUBLICATION_NODE.property(STATUS).isNull() .or(EVENT_PUBLICATION_NODE.property(STATUS).ne(parameter(STATUS)))) .set(EVENT_PUBLICATION_NODE.property(STATUS).to(parameter(STATUS))) @@ -167,6 +168,7 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository { private static final Statement RESUBMIT_STATEMENT = match(EVENT_PUBLICATION_NODE) .where(EVENT_PUBLICATION_NODE.property(ID).eq(parameter(ID))) + .and(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNull()) .and(EVENT_PUBLICATION_NODE.property(STATUS).isNull() .or(EVENT_PUBLICATION_NODE.property(STATUS).ne(literalOf(Status.RESUBMITTED.name())))) .set(EVENT_PUBLICATION_NODE.property(STATUS).to(parameter(STATUS))) diff --git a/spring-modulith-events/spring-modulith-events-neo4j/src/test/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepositoryTest.java b/spring-modulith-events/spring-modulith-events-neo4j/src/test/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepositoryTest.java index 00b72d343..d88109c99 100644 --- a/spring-modulith-events/spring-modulith-events-neo4j/src/test/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepositoryTest.java +++ b/spring-modulith-events/spring-modulith-events-neo4j/src/test/java/org/springframework/modulith/events/neo4j/Neo4jEventPublicationRepositoryTest.java @@ -581,6 +581,27 @@ void marksPublicationWithNullStatusPropertyAsFailed() { .containsExactly(publication.getIdentifier()); } + @Test // GH-1925 + void doesNotResubmitCompletedPublication() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + + assertThat(repository.markResubmitted(publication.getIdentifier(), Instant.now())).isFalse(); + } + + @Test // GH-1925 + void doesNotMarkCompletedPublicationFailed() { + + var publication = createPublication(new TestEvent("completed")); + + repository.markCompleted(publication.getIdentifier(), Instant.now()); + repository.markFailed(publication.getIdentifier()); + + assertThat(repository.findByStatus(EventPublication.Status.FAILED)).isEmpty(); + } + /** * Simulates a publication persisted by a schema version that predates the {@code status} property, i.e. one for * which the property is missing rather than defaulted. From 8973a9c81ad56711c744a5a61b91d5da786311ee Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Fri, 9 Oct 2026 14:32:05 +0200 Subject: [PATCH 15/15] GH-1412 - Avoid $merge when archiving MongoDB event publications. MongoDbEventPublicationRepository previously archived completed or abandoned publications via an aggregation using $merge, which MongoDB rejects inside transactions (error 263). We now read the affected publications and write them to the archive collection through a single unordered bulk upsert using $setOnInsert, which is transaction-safe and still keeps already existing archive entries. --- .../MongoDbEventPublicationRepository.java | 40 ++++++++++------- ...MongoDbEventPublicationRepositoryTest.java | 44 +++++++++++++++++++ 2 files changed, 69 insertions(+), 15 deletions(-) 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 f78887729..dcc3fdcfe 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 @@ -15,7 +15,6 @@ */ package org.springframework.modulith.events.mongodb; -import static org.springframework.data.mongodb.core.aggregation.Aggregation.*; import static org.springframework.data.mongodb.core.query.Criteria.*; import static org.springframework.data.mongodb.core.query.Query.*; @@ -31,10 +30,10 @@ import org.springframework.data.annotation.Id; import org.springframework.data.core.TypeInformation; import org.springframework.data.domain.Sort; +import org.springframework.data.mongodb.core.BulkOperations.BulkMode; import org.springframework.data.mongodb.core.FindAndModifyOptions; import org.springframework.data.mongodb.core.MongoTemplate; import org.springframework.data.mongodb.core.aggregation.Fields; -import org.springframework.data.mongodb.core.aggregation.MergeOperation.WhenDocumentsMatch; import org.springframework.data.mongodb.core.query.Criteria; import org.springframework.data.mongodb.core.query.Query; import org.springframework.data.mongodb.core.query.Update; @@ -455,23 +454,34 @@ private void archiveTerminal(Collection identifiers, Instant now, Status s return; } - var aggregation = newAggregation(MongoDbEventPublication.class, + // Avoid $merge as it's not supported in transactions. Upserting via $setOnInsert keeps existing archive entries. + var publications = mongoTemplate.find(query(where(ID).in(identifiers).and(COMPLETION_DATE).isNull()), + MongoDbEventPublication.class, collection); - match(where(ID).in(identifiers).and(COMPLETION_DATE).isNull()), + if (!publications.isEmpty()) { - addFields() - .addFieldWithValue(COMPLETION_DATE, now) - .addFieldWithValue(STATUS, status.name()) - .build(), + var operations = mongoTemplate.bulkOps(BulkMode.UNORDERED, MongoDbEventPublication.class, archiveCollection); + var converter = mongoTemplate.getConverter(); - merge() - .intoCollection(archiveCollection) - .on(ID) - .whenMatched(WhenDocumentsMatch.keepExistingDocument()) - .build()) - .withOptions(newAggregationOptions().skipOutput().build()); + for (var publication : publications) { + + publication.completionDate = now; + publication.status = status; + + var document = new Document(); + converter.write(publication, document); + + var update = new Update(); + document.entrySet().stream() + .filter(it -> !it.getKey().equals(Fields.UNDERSCORE_ID)) + .forEach(it -> update.setOnInsert(it.getKey(), it.getValue())); + + operations.upsert(query(where(ID).is(publication.id)), update); + } + + operations.execute(); + } - mongoTemplate.aggregate(aggregation, collection, Document.class); mongoTemplate.remove(query(where(ID).in(identifiers)), MongoDbEventPublication.class, collection); } 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 093e611f3..baf13bed0 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 @@ -39,6 +39,7 @@ 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.MongoTransactionManager; import org.springframework.data.mongodb.core.MongoTemplate; import org.springframework.modulith.events.EventPublication.Status; import org.springframework.modulith.events.ResubmissionOptions; @@ -51,6 +52,7 @@ import org.springframework.modulith.testapp.TestApplication; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.TestPropertySource; +import org.springframework.transaction.support.TransactionTemplate; import org.testcontainers.junit.jupiter.Testcontainers; /** @@ -617,6 +619,48 @@ void doesNotAbandonPublicationThatHasBeenConcurrentlyResubmitted() { assertThat(repository.findByStatus(Status.ABANDONED)).isEmpty(); } + @Test // GH-1412 + void archivesPublicationsInsideTransaction() { + + assumeTrue(completionMode == CompletionMode.ARCHIVE); + + var first = createPublication(new TestEvent("first")); + var second = createPublication(new TestEvent("second")); + var template = new TransactionTemplate(new MongoTransactionManager(mongoTemplate.getMongoDatabaseFactory())); + + template.executeWithoutResult(__ -> { + repository.markCompleted(first.getIdentifier(), Instant.now()); + repository.markCompleted(second.getIdentifier(), Instant.now()); + }); + + assertThat(repository.findIncompletePublications()).isEmpty(); + assertThat(repository.findCompletedPublications()) + .extracting(TargetEventPublication::getIdentifier) + .containsExactlyInAnyOrder(first.getIdentifier(), second.getIdentifier()); + } + + @Test // GH-1412 + void keepsExistingArchiveEntryWhenArchivingAgain() { + + assumeTrue(completionMode == CompletionMode.ARCHIVE); + + var publication = createPublication(new TestEvent("first")); + var firstCompletion = Instant.now().truncatedTo(ChronoUnit.MILLIS); + + repository.markCompleted(publication.getIdentifier(), firstCompletion); + + // Re-create the active entry with the same identifier to simulate a concurrent archival + mongoTemplate.save(new MongoDbEventPublication(publication.getIdentifier(), Instant.now(), "listener", + new TestEvent("first"), null, Status.PUBLISHED, null, 1)); + + repository.markCompleted(publication.getIdentifier(), firstCompletion.plusSeconds(60)); + + assertThat(mongoTemplate.findAll(MongoDbEventPublication.class, archiveCollection)).singleElement() + .satisfies(it -> assertThat(it.completionDate).isEqualTo(firstCompletion)); + assertThat(mongoTemplate.findAll(MongoDbEventPublication.class, + mongoTemplate.getCollectionName(MongoDbEventPublication.class))).isEmpty(); + } + private TargetEventPublication createPublication(Object event) { return createPublication(event, TARGET_IDENTIFIER); }