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 b05fc0e5..3307422d 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 @@ -62,8 +62,9 @@ class RabbitEventExternalizerConfiguration { return new DelegatingEventExternalizer(configuration, (target, payload) -> { var routing = BrokerRouting.of(target, context); + var headers = configuration.getHeadersFor(payload); - operations.convertAndSend(routing.getTarget(), routing.getKey(payload), payload); + operations.convertAndSend(routing.getTarget(), routing.getKey(payload), payload, headers); return CompletableFuture.completedFuture(null); }); diff --git a/spring-modulith-events/spring-modulith-events-api/src/main/java/org/springframework/modulith/events/DefaultEventExternalizationConfiguration.java b/spring-modulith-events/spring-modulith-events-api/src/main/java/org/springframework/modulith/events/DefaultEventExternalizationConfiguration.java index 1363f30e..d1620ca9 100644 --- a/spring-modulith-events/spring-modulith-events-api/src/main/java/org/springframework/modulith/events/DefaultEventExternalizationConfiguration.java +++ b/spring-modulith-events/spring-modulith-events-api/src/main/java/org/springframework/modulith/events/DefaultEventExternalizationConfiguration.java @@ -15,6 +15,7 @@ */ package org.springframework.modulith.events; +import java.util.Map; import java.util.function.Function; import java.util.function.Predicate; @@ -31,6 +32,7 @@ class DefaultEventExternalizationConfiguration implements EventExternalizationCo private final Predicate filter; private final Function mapper; private final Function router; + private final Function> headers; /** * Creates a new {@link DefaultEventExternalizationConfiguration} @@ -38,9 +40,10 @@ class DefaultEventExternalizationConfiguration implements EventExternalizationCo * @param filter must not be {@literal null}. * @param mapper must not be {@literal null}. * @param router must not be {@literal null}. + * @param headers must not be {@literal null}. */ DefaultEventExternalizationConfiguration(Predicate filter, Function mapper, - Function router) { + Function router, Function> headers) { Assert.notNull(filter, "Filter must not be null!"); Assert.notNull(mapper, "Mapper must not be null!"); @@ -49,6 +52,7 @@ class DefaultEventExternalizationConfiguration implements EventExternalizationCo this.filter = filter; this.mapper = mapper; this.router = router; + this.headers = headers; } /** @@ -95,4 +99,16 @@ class DefaultEventExternalizationConfiguration implements EventExternalizationCo return router.apply(event).verify(); } + + /* + * (non-Javadoc) + * @see org.springframework.modulith.events.EventExternalizationConfiguration#getHeadersFor(java.lang.Object) + */ + @Override + public Map getHeadersFor(Object event) { + + Assert.notNull(event, "Event must not be null!"); + + return headers.apply(event); + } } diff --git a/spring-modulith-events/spring-modulith-events-api/src/main/java/org/springframework/modulith/events/EventExternalizationConfiguration.java b/spring-modulith-events/spring-modulith-events-api/src/main/java/org/springframework/modulith/events/EventExternalizationConfiguration.java index 787c11db..8c4436d2 100644 --- a/spring-modulith-events/spring-modulith-events-api/src/main/java/org/springframework/modulith/events/EventExternalizationConfiguration.java +++ b/spring-modulith-events/spring-modulith-events-api/src/main/java/org/springframework/modulith/events/EventExternalizationConfiguration.java @@ -19,7 +19,9 @@ import static org.springframework.core.annotation.AnnotatedElementUtils.*; import java.lang.annotation.Annotation; import java.util.Collection; +import java.util.Collections; import java.util.List; +import java.util.Map; import java.util.Optional; import java.util.function.BiFunction; import java.util.function.BiPredicate; @@ -188,6 +190,15 @@ public interface EventExternalizationConfiguration { */ RoutingTarget determineTarget(Object event); + /** + * Returns the headers to be attached to the message sent out for the given event. + * + * @param event must not be {@literal null}. + * @return will never be {@literal null}. + * @since 1.3 + */ + Map getHeadersFor(Object event); + /** * API to define which events are supposed to be selected for externalization. * @@ -367,6 +378,7 @@ public interface EventExternalizationConfiguration { private final Predicate filter; private final Function mapper; private final Function router; + private final Function> headers; /** * Creates a new {@link Router} for the given selector {@link Predicate} and mapper and router {@link Function}s. @@ -374,16 +386,20 @@ public interface EventExternalizationConfiguration { * @param filter must not be {@literal null}. * @param mapper must not be {@literal null}. * @param router must not be {@literal null}. + * @param headers must not be {@literal null}. */ - Router(Predicate filter, Function mapper, Function router) { + Router(Predicate filter, Function mapper, Function router, + Function> headers) { Assert.notNull(filter, "Selector must not be null!"); Assert.notNull(mapper, "Mapper must not be null!"); Assert.notNull(router, "Router must not be null!"); + Assert.notNull(headers, "Headers extractor must not be null!"); this.filter = filter; this.mapper = mapper; this.router = router; + this.headers = headers; } /** @@ -392,7 +408,7 @@ public interface EventExternalizationConfiguration { * @param filter must not be {@literal null}. */ Router(Predicate filter) { - this(filter, Function.identity(), DEFAULT_ROUTER); + this(filter, Function.identity(), DEFAULT_ROUTER, it -> Collections.emptyMap()); } /** @@ -406,7 +422,7 @@ public interface EventExternalizationConfiguration { Assert.notNull(mapper, "Mapper must not be null!"); - return new Router(filter, mapper, router); + return new Router(filter, mapper, router, headers); } /** @@ -428,7 +444,42 @@ public interface EventExternalizationConfiguration { .map(mapper::apply) .orElse(it); - return new Router(filter, this.mapper.compose(combined), router); + return new Router(filter, this.mapper.compose(combined), router, headers); + } + + /** + * Registers the given function to extract headers from the events to be externalized. Will reset the entire header + * extractor arrangement. For type-specific extractions, see {@link #headers(Class, Function)}. + * + * @param extractor must not be {@literal null}. + * @return will never be {@literal null}. + * @see #headers(Class, Function) + * @since 1.3 + */ + public Router headers(Function> extractor) { + + Assert.notNull(extractor, "Headers extractor must not be null!"); + + return new Router(filter, mapper, router, extractor); + } + + /** + * Registers the given type-specific function to extract headers from the events to be externalized. + * + * @param extractor must not be {@literal null}. + * @return will never be {@literal null}. + * @since 1.3 + */ + public Router headers(Class type, Function> extractor) { + + Assert.notNull(type, "Type must not be null!"); + Assert.notNull(extractor, "Headers extractor must not be null!"); + + Function> combined = it -> toOptional(type, it) + .map(extractor::apply) + .orElseGet(() -> this.headers.apply(it)); + + return new Router(filter, mapper, router, combined); } /** @@ -437,7 +488,7 @@ public interface EventExternalizationConfiguration { * @return will never be {@literal null}. */ public Router routeMapped() { - return new Router(filter, mapper, router.compose(mapper)); + return new Router(filter, mapper, router.compose(mapper), headers); } /** @@ -450,7 +501,7 @@ public interface EventExternalizationConfiguration { Assert.notNull(router, "Router must not be null!"); - return new Router(filter, mapper, router); + return new Router(filter, mapper, router, headers); } /** @@ -466,9 +517,11 @@ public interface EventExternalizationConfiguration { Assert.notNull(type, "Type must not be null!"); Assert.notNull(router, "Router must not be null!"); - return new Router(filter, mapper, it -> toOptional(type, it) + Function adapted = it -> toOptional(type, it) .map(router::apply) - .orElseGet(() -> this.router.apply(it))); + .orElseGet(() -> this.router.apply(it)); + + return new Router(filter, mapper, adapted, headers); } /** @@ -487,9 +540,11 @@ public interface EventExternalizationConfiguration { Assert.notNull(type, "Type must not be null!"); Assert.notNull(extractor, "Extractor must not be null!"); - return new Router(filter, mapper, it -> toOptional(type, it) + Function adapted = it -> toOptional(type, it) .map(t -> this.router.apply(t).withKey(extractor.apply(t))) - .orElseGet(() -> this.router.apply(it))); + .orElseGet(() -> this.router.apply(it)); + + return new Router(filter, mapper, adapted, headers); } /** @@ -503,9 +558,9 @@ public interface EventExternalizationConfiguration { Assert.notNull(router, "Router must not be null!"); - return new Router(filter, mapper, it -> router.apply(it) - .orElseGet(() -> this.router.apply(it))) - .build(); + Function adapted = it -> router.apply(it).orElseGet(() -> this.router.apply(it)); + + return new Router(filter, mapper, adapted, headers).build(); } /** @@ -533,16 +588,16 @@ public interface EventExternalizationConfiguration { Assert.notNull(router, "Router must not be null!"); - return new Router(filter, mapper, it -> router.apply(it.getClass())); + return new Router(filter, mapper, it -> router.apply(it.getClass()), headers); } /** - * Creates a new {@link EventExternalizationConfiguration} refelcting the current configuration. + * Creates a new {@link EventExternalizationConfiguration} reflecting the current configuration. * * @return will never be {@literal null}. */ public EventExternalizationConfiguration build() { - return new DefaultEventExternalizationConfiguration(filter, mapper, router); + return new DefaultEventExternalizationConfiguration(filter, mapper, router, headers); } private static Optional toOptional(Class type, Object source) { diff --git a/spring-modulith-events/spring-modulith-events-api/src/test/java/org/springframework/modulith/events/EventExternalizationConfigurationUnitTests.java b/spring-modulith-events/spring-modulith-events-api/src/test/java/org/springframework/modulith/events/EventExternalizationConfigurationUnitTests.java index 08ebcd92..556f3f47 100644 --- a/spring-modulith-events/spring-modulith-events-api/src/test/java/org/springframework/modulith/events/EventExternalizationConfigurationUnitTests.java +++ b/spring-modulith-events/spring-modulith-events-api/src/test/java/org/springframework/modulith/events/EventExternalizationConfigurationUnitTests.java @@ -22,6 +22,8 @@ import lombok.RequiredArgsConstructor; import java.lang.annotation.Retention; import java.lang.annotation.RetentionPolicy; +import java.util.List; +import java.util.Map; import org.junit.jupiter.api.Test; @@ -130,6 +132,21 @@ class EventExternalizationConfigurationUnitTests { assertThat(target.getKey()).isNull(); } + @Test // GH-855 + void registersHeaderExtractor() { + + var configuration = defaults("org.springframework.modulith") + .headers(AnotherSampleEvent.class, it -> Map.of("another", "anotherValue")) + .headers(SampleEvent.class, it -> Map.of("sample", "value")) + .build(); + + assertThat(configuration.getHeadersFor(new SampleEvent())) + .containsEntry("sample", "value"); + + assertThat(configuration.getHeadersFor(new AnotherSampleEvent())) + .containsEntry("another", "anotherValue"); + } + @Retention(RetentionPolicy.RUNTIME) @interface CustomExternalized { String value() default ""; 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 33bae675..655c9fab 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 @@ -27,6 +27,9 @@ import org.springframework.context.expression.BeanFactoryResolver; import org.springframework.expression.spel.support.StandardEvaluationContext; import org.springframework.kafka.core.KafkaOperations; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.support.KafkaHeaders; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.MessageBuilder; import org.springframework.modulith.events.EventExternalizationConfiguration; import org.springframework.modulith.events.config.EventExternalizationAutoConfiguration; import org.springframework.modulith.events.support.BrokerRouting; @@ -61,7 +64,17 @@ class KafkaEventExternalizerConfiguration { return new DelegatingEventExternalizer(configuration, (target, payload) -> { var routing = BrokerRouting.of(target, context); - return operations.send(routing.getTarget(), routing.getKey(payload), payload); + + var builder = payload instanceof Message message + ? MessageBuilder.fromMessage(message) + : MessageBuilder.withPayload(payload).copyHeaders(configuration.getHeadersFor(payload)); + + var message = builder + .setHeaderIfAbsent(KafkaHeaders.KEY, routing.getKey(payload)) + .setHeaderIfAbsent(KafkaHeaders.TOPIC, routing.getTarget()) + .build(); + + return operations.send(message); }); } } diff --git a/spring-modulith-events/spring-modulith-events-kafka/src/main/java/org/springframework/modulith/events/kafka/KafkaJacksonConfiguration.java b/spring-modulith-events/spring-modulith-events-kafka/src/main/java/org/springframework/modulith/events/kafka/KafkaJacksonConfiguration.java index 24947692..144ed793 100644 --- a/spring-modulith-events/spring-modulith-events-kafka/src/main/java/org/springframework/modulith/events/kafka/KafkaJacksonConfiguration.java +++ b/spring-modulith-events/spring-modulith-events-kafka/src/main/java/org/springframework/modulith/events/kafka/KafkaJacksonConfiguration.java @@ -22,6 +22,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.PropertySource; +import org.springframework.kafka.support.converter.ByteArrayJsonMessageConverter; import org.springframework.kafka.support.converter.JsonMessageConverter; import com.fasterxml.jackson.databind.ObjectMapper; @@ -41,7 +42,7 @@ class KafkaJacksonConfiguration { @Bean @ConditionalOnBean(ObjectMapper.class) @ConditionalOnMissingBean(JsonMessageConverter.class) - JsonMessageConverter jsonMessageConverter(ObjectMapper mapper) { - return new JsonMessageConverter(mapper); + ByteArrayJsonMessageConverter jsonMessageConverter(ObjectMapper mapper) { + return new ByteArrayJsonMessageConverter(mapper); } } diff --git a/spring-modulith-events/spring-modulith-events-kafka/src/main/resources/kafka-json.properties b/spring-modulith-events/spring-modulith-events-kafka/src/main/resources/kafka-json.properties index 6cae2776..261ea743 100644 --- a/spring-modulith-events/spring-modulith-events-kafka/src/main/resources/kafka-json.properties +++ b/spring-modulith-events/spring-modulith-events-kafka/src/main/resources/kafka-json.properties @@ -1,2 +1,2 @@ -spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer +spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.ByteArraySerializer spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.ByteArrayDeserializer 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 a02a67f5..3c0e6601 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 @@ -18,11 +18,20 @@ package org.springframework.modulith.events.kafka; import static org.assertj.core.api.Assertions.*; import static org.mockito.Mockito.*; +import java.util.Map; +import java.util.function.Consumer; +import java.util.function.Supplier; + import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; import org.springframework.boot.autoconfigure.AutoConfigurations; import org.springframework.boot.test.context.runner.ApplicationContextRunner; import org.springframework.kafka.core.KafkaOperations; +import org.springframework.lang.Nullable; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.MessageBuilder; import org.springframework.modulith.events.EventExternalizationConfiguration; +import org.springframework.modulith.events.Externalized; import org.springframework.modulith.events.support.DelegatingEventExternalizer; /** @@ -33,6 +42,8 @@ import org.springframework.modulith.events.support.DelegatingEventExternalizer; */ class KafkaEventExternalizerConfigurationIntegrationTests { + private KafkaOperations operations = mock(KafkaOperations.class); + @Test // GH-342 void registersExternalizerByDefault() { @@ -52,12 +63,61 @@ class KafkaEventExternalizerConfigurationIntegrationTests { }); } + @Test // GH-855 + void addsHeadersIfConfigured() { + + var config = EventExternalizationConfiguration.defaults("org") + .headers(Sample.class, __ -> Map.of("key", "value")) + .build(); + + assertMessage(config, it -> { + assertThat(it.getHeaders()).containsKey("key"); + }); + } + + @Test // GH-855 + void sendsMessageAsIsIfMappingTarget() { + + var config = EventExternalizationConfiguration.defaults("org") + .mapping(Sample.class, it -> MessageBuilder.withPayload(it).setHeader("key", "value").build()) + .build(); + + assertMessage(config, it -> { + assertThat(it.getHeaders()).contains(Map.entry("key", "value")); + }); + } + + private void assertMessage(EventExternalizationConfiguration configuration, Consumer> assertions) { + + basicSetup(configuration) + .run(ctxt -> { + + ctxt.getBean(DelegatingEventExternalizer.class).externalize(new Sample()); + + var captor = ArgumentCaptor.forClass(Message.class); + verify(operations).send(captor.capture()); + + assertions.accept(captor.getValue()); + + }); + } + private ApplicationContextRunner basicSetup() { + return basicSetup(null); + } + + private ApplicationContextRunner basicSetup(@Nullable EventExternalizationConfiguration config) { + + Supplier configProvider = () -> config == null + ? EventExternalizationConfiguration.disabled() + : config; return new ApplicationContextRunner() - .withConfiguration( - AutoConfigurations.of(KafkaEventExternalizerConfiguration.class)) - .withBean(EventExternalizationConfiguration.class, () -> EventExternalizationConfiguration.disabled()) - .withBean(KafkaOperations.class, () -> mock(KafkaOperations.class)); + .withConfiguration(AutoConfigurations.of(KafkaEventExternalizerConfiguration.class)) + .withBean(EventExternalizationConfiguration.class, configProvider) + .withBean(KafkaOperations.class, () -> operations); } + + @Externalized + record Sample() {} } 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 46f64fa3..91599fa5 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 @@ -66,6 +66,7 @@ class SpringMessagingEventExternalizerConfiguration { var message = MessageBuilder .withPayload(payload) .setHeader(MODULITH_ROUTING_HEADER, target.toString()) + .copyHeadersIfAbsent(configuration.getHeadersFor(payload)) .build(); if (logger.isDebugEnabled()) { diff --git a/spring-modulith-examples/spring-modulith-example-kafka/src/test/java/example/TestApplication.java b/spring-modulith-examples/spring-modulith-example-kafka/src/test/java/example/TestApplication.java index 43db1707..e3cdb4b4 100644 --- a/spring-modulith-examples/spring-modulith-example-kafka/src/test/java/example/TestApplication.java +++ b/spring-modulith-examples/spring-modulith-example-kafka/src/test/java/example/TestApplication.java @@ -28,6 +28,7 @@ import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Primary; import org.springframework.kafka.core.KafkaOperations; +import org.springframework.messaging.Message; /** * @author Oliver Drotbohm @@ -44,9 +45,9 @@ public class TestApplication { var mock = mock(KafkaOperations.class); - when(mock.send(any(), any())).then(invocation -> { + when(mock.send(any(Message.class))).then(invocation -> { - logger.info("Sending message {} to {}.", invocation.getArguments()[1], invocation.getArguments()[0]); + logger.info("Sending message {}.", invocation.getArguments()[0]); return null; });