Added support for Function binding
Added support for binding Suppliers, Function and Consumers as message producers and handlers for apps annotated with EnableBinding Changed '.name' property to '.definition' Resolves #1445
This commit is contained in:
@@ -44,8 +44,7 @@ import org.springframework.cloud.stream.binding.InputBindingLifecycle;
|
||||
import org.springframework.cloud.stream.binding.MessageChannelStreamListenerResultAdapter;
|
||||
import org.springframework.cloud.stream.binding.OutputBindingLifecycle;
|
||||
import org.springframework.cloud.stream.binding.StreamListenerAnnotationBeanPostProcessor;
|
||||
import org.springframework.cloud.stream.function.FunctionConfiguration;
|
||||
import org.springframework.cloud.stream.function.FunctionProperties;
|
||||
import org.springframework.cloud.stream.function.StreamFunctionProperties;
|
||||
import org.springframework.cloud.stream.micrometer.DestinationPublishingMetricsAutoConfiguration;
|
||||
import org.springframework.context.ApplicationListener;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
@@ -78,7 +77,7 @@ import org.springframework.util.Assert;
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
@Configuration
|
||||
@EnableConfigurationProperties({ BindingServiceProperties.class, SpringIntegrationProperties.class, FunctionProperties.class })
|
||||
@EnableConfigurationProperties({ BindingServiceProperties.class, SpringIntegrationProperties.class, StreamFunctionProperties.class })
|
||||
@Import({ DestinationPublishingMetricsAutoConfiguration.class, SpelExpressionConverterConfiguration.class })
|
||||
@Role(BeanDefinition.ROLE_INFRASTRUCTURE)
|
||||
@ConditionalOnBean(value = BinderTypeRegistry.class, search = SearchStrategy.CURRENT)
|
||||
|
||||
@@ -17,14 +17,10 @@
|
||||
package org.springframework.cloud.stream.function;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
|
||||
import org.springframework.beans.factory.support.BeanDefinitionRegistry;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||
import org.springframework.cloud.function.context.FunctionCatalog;
|
||||
import org.springframework.cloud.function.context.FunctionType;
|
||||
import org.springframework.cloud.function.context.catalog.FunctionInspector;
|
||||
import org.springframework.cloud.stream.binding.BindingBeanDefinitionRegistryUtils;
|
||||
import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory;
|
||||
import org.springframework.cloud.stream.messaging.Processor;
|
||||
import org.springframework.cloud.stream.messaging.Sink;
|
||||
@@ -32,9 +28,6 @@ import org.springframework.cloud.stream.messaging.Source;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.integration.dsl.IntegrationFlow;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.SubscribableChannel;
|
||||
import org.springframework.util.ClassUtils;
|
||||
|
||||
/**
|
||||
*
|
||||
@@ -43,7 +36,7 @@ import org.springframework.util.ClassUtils;
|
||||
* @since 2.1
|
||||
*/
|
||||
@Configuration
|
||||
@ConditionalOnProperty("spring.cloud.stream.function.name")
|
||||
@ConditionalOnProperty("spring.cloud.stream.function.definition")
|
||||
public class FunctionConfiguration {
|
||||
|
||||
@Autowired(required=false)
|
||||
@@ -55,13 +48,10 @@ public class FunctionConfiguration {
|
||||
@Autowired(required=false)
|
||||
private Sink sink;
|
||||
|
||||
@Autowired
|
||||
private ConfigurableListableBeanFactory registry;
|
||||
|
||||
@Bean
|
||||
public IntegrationFlowFunctionSupport functionSupport(FunctionCatalogWrapper functionCatalog,
|
||||
FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory,
|
||||
FunctionProperties functionProperties) {
|
||||
StreamFunctionProperties functionProperties) {
|
||||
|
||||
return new IntegrationFlowFunctionSupport(functionCatalog, functionInspector, messageConverterFactory,
|
||||
functionProperties);
|
||||
@@ -72,11 +62,13 @@ public class FunctionConfiguration {
|
||||
return new FunctionCatalogWrapper(catalog);
|
||||
}
|
||||
|
||||
|
||||
@ConditionalOnProperty("spring.cloud.stream.function.name")
|
||||
@ConditionalOnMissingBean
|
||||
/**
|
||||
* This configuration creates an instance of {@link IntegrationFlow} appropriate for binding declared using EnableBinding.
|
||||
* At the moment only Source, Processor and Sink are supported.
|
||||
*/
|
||||
@ConditionalOnMissingBean // starter apps typically already provide and instance of IntegrationFlow, so we don't need this one.
|
||||
@Bean
|
||||
public IntegrationFlow foo(IntegrationFlowFunctionSupport functionSupport) {
|
||||
public IntegrationFlow integrationFlowCreator(IntegrationFlowFunctionSupport functionSupport) {
|
||||
if (processor != null) {
|
||||
return functionSupport.integrationFlowForFunction(processor.input(), processor.output()).get();
|
||||
}
|
||||
@@ -86,14 +78,6 @@ public class FunctionConfiguration {
|
||||
else if (source != null) {
|
||||
return functionSupport.integrationFlowFromNamedSupplier().channel(this.source.output()).get();
|
||||
}
|
||||
|
||||
FunctionType ft = functionSupport.getCurrentFunctionType();
|
||||
BindingBeanDefinitionRegistryUtils.registerBindingTargetBeanDefinitions(Sink.class,
|
||||
Sink.class.getName(), (BeanDefinitionRegistry) registry);
|
||||
BindingBeanDefinitionRegistryUtils.registerBindingTargetsQualifiedBeanDefinitions(
|
||||
ClassUtils.resolveClassName(this.getClass().getName(), null), Sink.class,
|
||||
(BeanDefinitionRegistry) registry);
|
||||
return functionSupport.integrationFlowForFunction(registry.getBean("input", SubscribableChannel.class), null).get();
|
||||
//throw new UnsupportedOperationException("Not yet supotrted");
|
||||
throw new UnsupportedOperationException("Bindings other then Source, Processor and Sink are not currently supported");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -53,7 +53,7 @@ public class IntegrationFlowFunctionSupport {
|
||||
|
||||
private final CompositeMessageConverterFactory messageConverterFactory;
|
||||
|
||||
private final FunctionProperties functionProperties;
|
||||
private final StreamFunctionProperties functionProperties;
|
||||
|
||||
@Autowired
|
||||
private MessageChannel errorChannel;
|
||||
@@ -65,7 +65,7 @@ public class IntegrationFlowFunctionSupport {
|
||||
* @param functionProperties
|
||||
*/
|
||||
public IntegrationFlowFunctionSupport(FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector,
|
||||
CompositeMessageConverterFactory messageConverterFactory, FunctionProperties functionProperties) {
|
||||
CompositeMessageConverterFactory messageConverterFactory, StreamFunctionProperties functionProperties) {
|
||||
|
||||
Assert.notNull(functionCatalog, "'functionCatalog' must not be null");
|
||||
Assert.notNull(functionInspector, "'functionInspector' must not be null");
|
||||
@@ -78,18 +78,18 @@ public class IntegrationFlowFunctionSupport {
|
||||
}
|
||||
|
||||
public FunctionType getCurrentFunctionType() {
|
||||
return functionInspector.getRegistration(functionCatalog.lookup(this.functionProperties.getName())).getType();
|
||||
return functionInspector.getRegistration(functionCatalog.lookup(this.functionProperties.getDefinition())).getType();
|
||||
}
|
||||
|
||||
/**
|
||||
* Create an instance of the {@link IntegrationFlowBuilder} from a {@link Supplier} bean available in the context.
|
||||
* The name of the bean must be provided via `spring.cloud.stream.function.name` property.
|
||||
* The name of the bean must be provided via `spring.cloud.stream.function.definition` property.
|
||||
* @return instance of {@link IntegrationFlowBuilder}
|
||||
* @throws IllegalStateException if the named Supplier can not be located.
|
||||
*/
|
||||
public IntegrationFlowBuilder integrationFlowFromNamedSupplier() {
|
||||
if (StringUtils.hasText(this.functionProperties.getName())) {
|
||||
Supplier<?> supplier = functionCatalog.lookup(Supplier.class, this.functionProperties.getName());
|
||||
if (StringUtils.hasText(this.functionProperties.getDefinition())) {
|
||||
Supplier<?> supplier = functionCatalog.lookup(Supplier.class, this.functionProperties.getDefinition());
|
||||
if (supplier instanceof FluxSupplier) {
|
||||
supplier = ((FluxSupplier<?>)supplier).getTarget();
|
||||
}
|
||||
@@ -98,7 +98,7 @@ public class IntegrationFlowFunctionSupport {
|
||||
}
|
||||
|
||||
throw new IllegalStateException(
|
||||
"A Supplier is not specified in the 'spring.cloud.stream.function.name' property.");
|
||||
"A Supplier is not specified in the 'spring.cloud.stream.function.definition' property.");
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -135,7 +135,7 @@ public class IntegrationFlowFunctionSupport {
|
||||
|
||||
/**
|
||||
* Add a {@link Function} bean to the end of an integration flow.
|
||||
* The name of the bean must be provided via `spring.cloud.stream.function.name` property.
|
||||
* The name of the bean must be provided via `spring.cloud.stream.function.definition` property.
|
||||
* <p>
|
||||
* NOTE: If this method returns true, the integration flow is now represented
|
||||
* as a Reactive Streams {@link Publisher} bean.
|
||||
@@ -146,9 +146,9 @@ public class IntegrationFlowFunctionSupport {
|
||||
* @return true if {@link Function} was located and added and false if it wasn't.
|
||||
*/
|
||||
public <I,O> boolean andThenFunction(IntegrationFlowBuilder flowBuilder, MessageChannel outputChannel) {
|
||||
if (StringUtils.hasText(this.functionProperties.getName())) {
|
||||
if (StringUtils.hasText(this.functionProperties.getDefinition())) {
|
||||
FunctionInvoker<I,O> functionInvoker =
|
||||
new FunctionInvoker<>(this.functionProperties.getName(), this.functionCatalog,
|
||||
new FunctionInvoker<>(this.functionProperties.getDefinition(), this.functionCatalog,
|
||||
this.functionInspector, this.messageConverterFactory, this.errorChannel);
|
||||
|
||||
if (outputChannel != null) {
|
||||
|
||||
@@ -25,20 +25,21 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||
* @since 2.1
|
||||
*/
|
||||
@ConfigurationProperties("spring.cloud.stream.function")
|
||||
public class FunctionProperties {
|
||||
public class StreamFunctionProperties {
|
||||
|
||||
/**
|
||||
* Name of functions to bind. If several functions need to be composed into one, use pipes (e.g., 'fooFunc|barFunc')
|
||||
* Definition of functions to bind. If several functions need to be composed
|
||||
* into one, use pipes (e.g., 'fooFunc|barFunc')
|
||||
*/
|
||||
private String name;
|
||||
private String definition;
|
||||
|
||||
|
||||
public String getName() {
|
||||
return this.name;
|
||||
public String getDefinition() {
|
||||
return this.definition;
|
||||
}
|
||||
|
||||
public void setName(String name) {
|
||||
this.name = name;
|
||||
public void setDefinition(String definition) {
|
||||
this.definition = definition;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -84,7 +84,7 @@ public class FunctionInvokerTests {
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Function<Message, Message> messageToMessageNoType() {
|
||||
public Function<Message<?>, Message<?>> messageToMessageNoType() {
|
||||
return x -> MessageBuilder.withPayload(new Bar()).copyHeaders(x.getHeaders()).build();
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,127 @@
|
||||
/*
|
||||
* Copyright 2018 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
|
||||
*
|
||||
* http://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.cloud.stream.function;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.Date;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.boot.WebApplicationType;
|
||||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||||
import org.springframework.boot.builder.SpringApplicationBuilder;
|
||||
import org.springframework.cloud.stream.annotation.EnableBinding;
|
||||
import org.springframework.cloud.stream.binder.test.InputDestination;
|
||||
import org.springframework.cloud.stream.binder.test.OutputDestination;
|
||||
import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
|
||||
import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory;
|
||||
import org.springframework.cloud.stream.messaging.Processor;
|
||||
import org.springframework.cloud.stream.messaging.Sink;
|
||||
import org.springframework.cloud.stream.messaging.Source;
|
||||
import org.springframework.context.ConfigurableApplicationContext;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.PollableChannel;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
/**
|
||||
* This test validates proper function binding for applications where EnableBinding is declared.
|
||||
*
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
public class GreenfieldFunctionEnableBindingTests {
|
||||
|
||||
@Test
|
||||
public void testSourceFromSupplier() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SourceFromSupplier.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=date", "--spring.jmx.enabled=false")) {
|
||||
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
Message<byte[]> sourceMessage = target.receive(10000);
|
||||
Date date = (Date) new CompositeMessageConverterFactory().getMessageConverterForAllRegistered().fromMessage(sourceMessage, Date.class);
|
||||
assertThat(date).isEqualTo(new Date(12345L));
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testProcessorFromFunction() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(ProcessorFromFunction.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=toUpperCase", "--spring.jmx.enabled=false")) {
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
source.send(new GenericMessage<byte[]>("John Doe".getBytes()));
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
assertThat(target.receive(10000).getPayload()).isEqualTo("JOHN DOE".getBytes(StandardCharsets.UTF_8));
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSinkFromConsumer() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SinkFromConsumer.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=sink", "--spring.jmx.enabled=false")) {
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
PollableChannel result = context.getBean("result", PollableChannel.class);
|
||||
source.send(new GenericMessage<byte[]>("John Doe".getBytes()));
|
||||
assertThat(result.receive(10000).getPayload()).isEqualTo("John Doe");
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Source.class)
|
||||
public static class SourceFromSupplier {
|
||||
@Bean
|
||||
public Supplier<Date> date() {
|
||||
return () -> new Date(12345L);
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Processor.class)
|
||||
public static class ProcessorFromFunction {
|
||||
@Bean
|
||||
public Function<String, String> toUpperCase() {
|
||||
return s -> s.toUpperCase();
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Sink.class)
|
||||
public static class SinkFromConsumer {
|
||||
@Bean
|
||||
public PollableChannel result() {
|
||||
return new QueueChannel();
|
||||
}
|
||||
@Bean
|
||||
public Consumer<String> sink(PollableChannel result) {
|
||||
return s -> {
|
||||
result.send(new GenericMessage<String>(s));
|
||||
System.out.println(s);
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,132 +0,0 @@
|
||||
package org.springframework.cloud.stream.function;
|
||||
|
||||
import java.util.Date;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.springframework.boot.WebApplicationType;
|
||||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||||
import org.springframework.boot.builder.SpringApplicationBuilder;
|
||||
import org.springframework.cloud.stream.annotation.EnableBinding;
|
||||
import org.springframework.cloud.stream.binder.test.InputDestination;
|
||||
import org.springframework.cloud.stream.binder.test.OutputDestination;
|
||||
import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
|
||||
import org.springframework.cloud.stream.messaging.Processor;
|
||||
import org.springframework.cloud.stream.messaging.Sink;
|
||||
import org.springframework.cloud.stream.messaging.Source;
|
||||
import org.springframework.context.ConfigurableApplicationContext;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
public class NewSourceAsSupplierTests {
|
||||
|
||||
@Test
|
||||
public void testSourceFromSupplier() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SourceFromSupplier.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.name=date", "--spring.jmx.enabled=false")) {
|
||||
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
Message<byte[]> sourceMessage = target.receive(10000);
|
||||
System.out.println(sourceMessage);
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8));
|
||||
//etc
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testProcessorFromFunction() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(ProcessorFromFunction.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.name=toUpperCase", "--spring.jmx.enabled=false")) {
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
source.send(new GenericMessage<byte[]>("fopo".getBytes()));
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
Message<byte[]> targetMessage = target.receive(10000);
|
||||
System.out.println(new String(targetMessage.getPayload()));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8));
|
||||
//etc
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSinkFromConsumer() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SinkFromConsumer.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.name=sink", "--spring.jmx.enabled=false")) {
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
source.send(new GenericMessage<byte[]>("fopo".getBytes()));
|
||||
// OutputDestination target = context.getBean(OutputDestination.class);
|
||||
// Message<byte[]> targetMessage = target.receive(10000);
|
||||
// System.out.println(new String(targetMessage.getPayload()));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8));
|
||||
//etc
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSinkFromConsumerNoEnableBinding() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SinkFromConsumerNoEnableBinding.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.name=sink", "--spring.jmx.enabled=false")) {
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
source.send(new GenericMessage<byte[]>("Hello No Binding".getBytes()));
|
||||
// OutputDestination target = context.getBean(OutputDestination.class);
|
||||
// Message<byte[]> targetMessage = target.receive(10000);
|
||||
// System.out.println(new String(targetMessage.getPayload()));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8));
|
||||
// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8));
|
||||
//etc
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Source.class)
|
||||
public static class SourceFromSupplier {
|
||||
@Bean
|
||||
public Supplier<Date> date() {
|
||||
return () -> new Date();
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Processor.class)
|
||||
public static class ProcessorFromFunction {
|
||||
@Bean
|
||||
public Function<String, String> toUpperCase() {
|
||||
return s -> s.toUpperCase();
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Sink.class)
|
||||
public static class SinkFromConsumer {
|
||||
@Bean
|
||||
public Consumer<String> sink() {
|
||||
return s -> System.out.println(s);
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
// @EnableBinding(Sink.class)
|
||||
public static class SinkFromConsumerNoEnableBinding {
|
||||
@Bean
|
||||
public Consumer<String> sink() {
|
||||
return s -> System.out.println("==> " + s);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -74,7 +74,7 @@ public class ProcessorToFunctionsSupportTests {
|
||||
new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.name=toUpperCase", "--spring.jmx.enabled=false");
|
||||
.run("--spring.cloud.stream.function.definition=toUpperCase", "--spring.jmx.enabled=false");
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
@@ -88,7 +88,7 @@ public class ProcessorToFunctionsSupportTests {
|
||||
new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false");
|
||||
.run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false");
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
@@ -102,7 +102,7 @@ public class ProcessorToFunctionsSupportTests {
|
||||
new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(ConsumerConfiguration.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.name=log", "--spring.jmx.enabled=false");
|
||||
.run("--spring.cloud.stream.function.definition=log", "--spring.jmx.enabled=false");
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
|
||||
@@ -72,7 +72,7 @@ public class SourceToFunctionsSupportTests {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class)).web(
|
||||
WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.name=toUpperCase", "--spring.jmx.enabled=false")) {
|
||||
.run("--spring.cloud.stream.function.definition=toUpperCase", "--spring.jmx.enabled=false")) {
|
||||
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
assertThat(target.receive(1000).getPayload()).isEqualTo("HELLO FUNCTION".getBytes(StandardCharsets.UTF_8));
|
||||
@@ -84,7 +84,7 @@ public class SourceToFunctionsSupportTests {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class)).web(
|
||||
WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false")) {
|
||||
.run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false")) {
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
assertThat(target.receive(1000).getPayload()).isEqualTo(
|
||||
"HELLO FUNCTION:HELLO FUNCTION".getBytes(StandardCharsets.UTF_8));
|
||||
@@ -97,7 +97,7 @@ public class SourceToFunctionsSupportTests {
|
||||
new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfigurationNoConversionPossible.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf",
|
||||
.run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf",
|
||||
"--spring.jmx.enabled=false")) {
|
||||
PollableChannel errorChannel = context.getBean("errorChannel", PollableChannel.class);
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
@@ -112,7 +112,7 @@ public class SourceToFunctionsSupportTests {
|
||||
new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfigurationNoConversionPossible.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf",
|
||||
.run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf",
|
||||
"--spring.jmx.enabled=false")) {
|
||||
PollableChannel errorChannel = context.getBean("errorChannel", PollableChannel.class);
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
@@ -125,7 +125,7 @@ public class SourceToFunctionsSupportTests {
|
||||
public void testMessageSourceIsCreatedFromProvidedSupplier() {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web(
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.name=number", "--spring.jmx.enabled=false")) {
|
||||
WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=number", "--spring.jmx.enabled=false")) {
|
||||
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8));
|
||||
@@ -140,7 +140,7 @@ public class SourceToFunctionsSupportTests {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web(
|
||||
WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.name=number|concatWithSelf", "--spring.jmx.enabled=false")) {
|
||||
.run("--spring.cloud.stream.function.definition=number|concatWithSelf", "--spring.jmx.enabled=false")) {
|
||||
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
assertThat(target.receive(10000).getPayload()).isEqualTo("11".getBytes(StandardCharsets.UTF_8));
|
||||
@@ -155,7 +155,7 @@ public class SourceToFunctionsSupportTests {
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web(
|
||||
WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.name=number|concatWithSelf|multiplyByTwo",
|
||||
.run("--spring.cloud.stream.function.definition=number|concatWithSelf|multiplyByTwo",
|
||||
"--spring.jmx.enabled=false")) {
|
||||
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
@@ -177,7 +177,7 @@ public class SourceToFunctionsSupportTests {
|
||||
new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web(
|
||||
WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.function.name=doesNotExist", "--spring.jmx.enabled=false");
|
||||
.run("--spring.cloud.stream.function.definition=doesNotExist", "--spring.jmx.enabled=false");
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@@ -297,11 +297,11 @@ public class SourceToFunctionsSupportTests {
|
||||
private Source source;
|
||||
|
||||
@Autowired
|
||||
private FunctionProperties functionProperties;
|
||||
private StreamFunctionProperties functionProperties;
|
||||
|
||||
@Bean
|
||||
public IntegrationFlow messageSourceFlow(IntegrationFlowFunctionSupport functionSupport) {
|
||||
Assert.hasText(this.functionProperties.getName(), "Supplier name must be provided");
|
||||
Assert.hasText(this.functionProperties.getDefinition(), "Supplier name must be provided");
|
||||
|
||||
return functionSupport.integrationFlowFromNamedSupplier().channel(this.source.output()).get();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user