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