From d16acf6475f944c67ce07ed833af67ceba7c6064 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 23 Mar 2020 16:29:35 +0100 Subject: [PATCH] GH-1931 Enhance header propagation in StreamBridge Given that StreamBridge is just a bridge and should honor whetever user sends through it, this fix ensures that message headers set by user are propagated Resolves #1931 --- pom.xml | 12 +++++++++ .../cloud/stream/function/StreamBridge.java | 3 +-- .../stream/function/RoutingFunctionTests.java | 8 ++---- .../stream/function/StreamBridgeTests.java | 26 +++++++++++++++++++ 4 files changed, 41 insertions(+), 8 deletions(-) diff --git a/pom.xml b/pom.xml index db9da884c..b8faeea27 100644 --- a/pom.xml +++ b/pom.xml @@ -208,6 +208,18 @@ + + + spring-snapshots + Spring Snapshots + https://repo.spring.io/libs-snapshot-local + + + spring-milestones + Spring milestones + https://repo.spring.io/libs-milestone-local + + spring-snapshots diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java index 563df07d2..7e277986f 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java @@ -131,8 +131,7 @@ public final class StreamBridge implements SmartInitializingSingleton { // we're registering a dummy pass-through function to ensure that it goes through the // same process (type conversion, etc) as other function invocation. FunctionRegistration> fr = new FunctionRegistration<>(v -> v, channelEntry.getKey()); - fr.type(FunctionType.from(Object.class).to(Object.class)); - this.functionRegistry.register(fr); + this.functionRegistry.register(fr.type(FunctionType.from(Object.class).to(Object.class).message())); this.initialized = true; } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java index 05acfa37e..118733dd6 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java @@ -18,7 +18,6 @@ package org.springframework.cloud.stream.function; import java.util.function.Function; -import org.junit.After; import org.junit.Before; import org.junit.Test; import reactor.core.publisher.Flux; @@ -48,16 +47,13 @@ import static org.assertj.core.api.Assertions.assertThat; */ public class RoutingFunctionTests { - @After - public void after() { - System.getProperties().remove("spring.cloud.function.routing.enabled"); - System.getProperties().remove("spring.cloud.stream.function.definition"); - } @Before public void before() { System.getProperties().remove("spring.cloud.function.routing.enabled"); System.getProperties().remove("spring.cloud.stream.function.definition"); + System.getProperties().remove("spring.cloud.function.definition"); + System.getProperties().remove("spring.cloud.function.routing-expression"); } @Test diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java index a1b2053d1..0e7d1cdca 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java @@ -30,6 +30,7 @@ import org.springframework.cloud.stream.binder.test.TestChannelBinderConfigurati import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.messaging.Message; +import org.springframework.messaging.support.MessageBuilder; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.fail; @@ -76,6 +77,31 @@ public class StreamBridgeTests { } } + @Test + public void testBridgeFunctionsSendingMessagePreservingHeaders() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration + .getCompleteConfiguration(EmptyConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.cloud.stream.source=foo;bar", + "--spring.jmx.enabled=false")) { + + StreamBridge bridge = context.getBean(StreamBridge.class); + bridge.send("foo-out-0", MessageBuilder.withPayload("hello foo").setHeader("foo", "foo").build()); + bridge.send("bar-out-0", MessageBuilder.withPayload("hello bar").setHeader("bar", "bar").build()); + + + OutputDestination outputDestination = context.getBean(OutputDestination.class); + + + Message message = outputDestination.receive(100, "foo-out-0"); + assertThat(message.getPayload()).isEqualTo("hello foo".getBytes()); + assertThat(message.getHeaders().get("foo")).isEqualTo("foo"); + + message = outputDestination.receive(100, "bar-out-0"); + assertThat(message.getPayload()).isEqualTo("hello bar".getBytes()); + assertThat(message.getHeaders().get("bar")).isEqualTo("bar"); + } + } + @Test public void testBridgeFunctionsWitthPartitionInformation() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration