GH-1891 Added test to validate propper behavior of stateful functions
While the original issue was addressed with the reversak of a5e990b2d0 commit
this commit simply adds a test to validate the expected behavior
Resolves #1891
This commit is contained in:
@@ -17,6 +17,7 @@
|
||||
package org.springframework.cloud.stream.function;
|
||||
|
||||
import java.io.Serializable;
|
||||
import java.time.Duration;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
@@ -59,7 +60,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@After
|
||||
public void after() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
}
|
||||
|
||||
@@ -86,7 +87,7 @@ public class ImplicitFunctionBindingTests {
|
||||
NoEnableBindingConfiguration.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.jmx.enabled=false",
|
||||
"--spring.cloud.stream.function.definition=func")) {
|
||||
"--spring.cloud.function.definition=func")) {
|
||||
|
||||
InputDestination inputDestination = context.getBean(InputDestination.class);
|
||||
OutputDestination outputDestination = context
|
||||
@@ -102,6 +103,36 @@ public class ImplicitFunctionBindingTests {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testReactiveFunctionWithState() {
|
||||
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
NoEnableBindingConfiguration.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.jmx.enabled=false",
|
||||
"--spring.cloud.function.definition=aggregate")) {
|
||||
|
||||
InputDestination inputDestination = context.getBean(InputDestination.class);
|
||||
OutputDestination outputDestination = context
|
||||
.getBean(OutputDestination.class);
|
||||
|
||||
Message<byte[]> inputMessage = MessageBuilder
|
||||
.withPayload("Hello".getBytes()).build();
|
||||
inputDestination.send(inputMessage);
|
||||
inputDestination.send(inputMessage);
|
||||
inputDestination.send(inputMessage);
|
||||
assertThat(new String(outputDestination.receive(2000).getPayload())).isEqualTo("HelloHelloHello");
|
||||
assertThat(new String(outputDestination.receive(2000).getPayload())).isEqualTo("");
|
||||
|
||||
inputDestination.send(inputMessage);
|
||||
inputDestination.send(inputMessage);
|
||||
inputDestination.send(inputMessage);
|
||||
inputDestination.send(inputMessage);
|
||||
assertThat(new String(outputDestination.receive(2000).getPayload())).isEqualTo("HelloHelloHelloHello");
|
||||
}
|
||||
}
|
||||
|
||||
@SuppressWarnings("rawtypes")
|
||||
@Test
|
||||
public void testFunctionWithUseNativeEncoding() {
|
||||
@@ -152,7 +183,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test
|
||||
public void testSimpleFunctionWithoutDefinitionProperty() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
SingleFunctionConfiguration.class))
|
||||
@@ -175,7 +206,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test
|
||||
public void testSimpleConsumerWithoutDefinitionProperty() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
SingleConsumerConfiguration.class))
|
||||
@@ -194,7 +225,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test
|
||||
public void testReactiveConsumerWithoutDefinitionProperty() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
SingleReactiveConsumerConfiguration.class))
|
||||
@@ -217,7 +248,7 @@ public class ImplicitFunctionBindingTests {
|
||||
TestChannelBinderConfiguration
|
||||
.getCompleteConfiguration(SingleConsumerConfiguration.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.definition=consumer",
|
||||
.run("--spring.cloud.function.definition=consumer",
|
||||
"--spring.jmx.enabled=false",
|
||||
"--spring.cloud.stream.bindings.input.content-type=text/plain",
|
||||
"--spring.cloud.stream.bindings.input.consumer.use-native-decoding=true")) {
|
||||
@@ -229,7 +260,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test
|
||||
public void testBindingWithReactiveFunction() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
ReactiveFunctionConfiguration.class))
|
||||
@@ -256,7 +287,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test
|
||||
public void testFunctionConfigDisabledIfStreamListenerIsUsed() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
LegacyConfiguration.class))
|
||||
@@ -270,7 +301,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test(expected = Exception.class)
|
||||
public void testDeclaredTypeVsActualInstance() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
SCF_GH_409Configuration.class))
|
||||
@@ -288,7 +319,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test
|
||||
public void testWithContextTypeApplicationProperty() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
SingleFunctionConfiguration.class))
|
||||
@@ -317,7 +348,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test
|
||||
public void testWithIntegrationFlowAsFunction() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
FunctionSampleSpringIntegrationConfiguration.class))
|
||||
@@ -338,7 +369,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test
|
||||
public void testSupplierWithCustomPoller() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
SupplierWithExplicitPollerConfiguration.class))
|
||||
@@ -358,14 +389,14 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test
|
||||
public void testSupplierWithCustomPollerAndMappedOutput() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
SupplierWithExplicitPollerConfiguration.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.jmx.enabled=false",
|
||||
"--spring.cloud.stream.poller.fixed-delay=2000",
|
||||
"--spring.cloud.stream.function.bindings.supplier-out-0=output")) {
|
||||
"--spring.cloud.function.bindings.supplier-out-0=output")) {
|
||||
|
||||
OutputDestination outputDestination = context.getBean(OutputDestination.class);
|
||||
|
||||
@@ -380,7 +411,7 @@ public class ImplicitFunctionBindingTests {
|
||||
|
||||
@Test
|
||||
public void testNoFunctionEnabledConfiguration() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(NoFunctionEnabledConfiguration.class))
|
||||
.web(WebApplicationType.NONE).run("--spring.jmx.enabled=false")) {
|
||||
@@ -408,6 +439,15 @@ public class ImplicitFunctionBindingTests {
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Function<Flux<String>, Flux<String>> aggregate() {
|
||||
return inbound -> inbound.
|
||||
log()
|
||||
.window(Duration.ofSeconds(1))
|
||||
.flatMap(w -> w.reduce("", (s1,s2)->s1+s2))
|
||||
.log();
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Consumer<String> cons() {
|
||||
return x -> {
|
||||
@@ -487,6 +527,7 @@ public class ImplicitFunctionBindingTests {
|
||||
return new Foo();
|
||||
}
|
||||
|
||||
@SuppressWarnings("serial")
|
||||
private static class Foo implements Supplier<Object>, Serializable {
|
||||
|
||||
@Override
|
||||
@@ -501,6 +542,7 @@ public class ImplicitFunctionBindingTests {
|
||||
@EnableAutoConfiguration
|
||||
public static class FunctionSampleSpringIntegrationConfiguration {
|
||||
|
||||
@SuppressWarnings("deprecation")
|
||||
@Bean
|
||||
public IntegrationFlow uppercaseFlow() {
|
||||
return IntegrationFlows.from(MessageFunction.class, "uppercase")
|
||||
|
||||
Reference in New Issue
Block a user