From 043bf8e8c6d161885e2e47fd2bbf841c71c280c2 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 10 Sep 2019 15:59:19 +0200 Subject: [PATCH] GH-1803 Propagate BindingProperties to FunctionConfiguration Resolves #1803 --- .../function/FunctionConfiguration.java | 24 ++++++++++----- .../ImplicitFunctionBindingTests.java | 29 +++++++++++++++++++ 2 files changed, 45 insertions(+), 8 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index 62cab49ac..6d3f18f87 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -46,8 +46,12 @@ import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.BindingCreatedEvent; import org.springframework.cloud.stream.binding.BindableProxyFactory; import org.springframework.cloud.stream.config.BinderFactoryAutoConfiguration; +import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceConfiguration; +import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel; +import org.springframework.cloud.stream.messaging.Sink; +import org.springframework.cloud.stream.messaging.Source; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; import org.springframework.context.annotation.Bean; @@ -87,9 +91,9 @@ class FunctionConfiguration { @Bean public InitializingBean functionChannelBindingInitializer(FunctionCatalog functionCatalog, FunctionInspector functionInspector, - StreamFunctionProperties functionProperties, @Nullable BindableProxyFactory[] bindableProxyFactory, GenericApplicationContext context) { + StreamFunctionProperties functionProperties, @Nullable BindableProxyFactory[] bindableProxyFactory, BindingServiceProperties serviceProperties) { return new FunctionChannelBindingInitializer(functionCatalog, functionInspector, functionProperties, - ObjectUtils.isEmpty(bindableProxyFactory) ? null : bindableProxyFactory[0]); + ObjectUtils.isEmpty(bindableProxyFactory) ? null : bindableProxyFactory[0], serviceProperties); } @Bean @@ -185,16 +189,18 @@ class FunctionConfiguration { private final BindableProxyFactory bindableProxyFactory; + private final BindingServiceProperties serviceProperties; + private GenericApplicationContext context; FunctionChannelBindingInitializer(FunctionCatalog functionCatalog, FunctionInspector functionInspector, - StreamFunctionProperties functionProperties, BindableProxyFactory bindableProxyFactory) { + StreamFunctionProperties functionProperties, BindableProxyFactory bindableProxyFactory, BindingServiceProperties serviceProperties) { this.functionCatalog = functionCatalog; this.functionInspector = functionInspector; this.functionProperties = functionProperties; this.bindableProxyFactory = bindableProxyFactory; - + this.serviceProperties = serviceProperties; } @Override @@ -230,10 +236,11 @@ class FunctionConfiguration { if (functionProperties.isComposeTo() && messageChannel instanceof SubscribableChannel && "input".equals(channelName)) { throw new UnsupportedOperationException("Composing at tail is not currently supported"); } - else if (functionProperties.isComposeFrom() && "output".equals(channelName)) { + else if (functionProperties.isComposeFrom() && Source.OUTPUT.equals(channelName)) { Assert.notNull(this.bindableProxyFactory, "Can not compose function into the existing app since `bindableProxyFactory` is null."); logger.info("Composing at the head of 'output' channel"); - FunctionInvocationWrapper function = functionCatalog.lookup(functionProperties.getDefinition(), "application/json"); + BindingProperties properties = this.serviceProperties.getBindings().get(Source.OUTPUT); + FunctionInvocationWrapper function = functionCatalog.lookup(functionProperties.getDefinition(), properties.getContentType()); ServiceActivatingHandler handler = new ServiceActivatingHandler(new FunctionWrapper(function)); handler.setBeanFactory(context); handler.afterPropertiesSet(); @@ -249,8 +256,9 @@ class FunctionConfiguration { subscribeChannel.subscribe(handler); } else { - FunctionInvocationWrapper function = functionCatalog.lookup(functionProperties.getDefinition(), "application/json"); - if (/*!function.isSupplier() && */"input".equals(channelName)) { + if (Sink.INPUT.equals(channelName)) { + BindingProperties properties = this.serviceProperties.getBindings().get(Sink.INPUT); + FunctionInvocationWrapper function = functionCatalog.lookup(functionProperties.getDefinition(), properties.getContentType()); this.postProcessForStandAloneFunction(function, messageChannel); } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 5bd9b4328..481226a8c 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -197,6 +197,35 @@ public class ImplicitFunctionBindingTests { } } + + @Test + public void testWithContextTypeApplicationProperty() { + System.clearProperty("spring.cloud.stream.function.definition"); + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + SingleFunctionConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.stream.bindings.input.content-type=text/plain")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context + .getBean(OutputDestination.class); + + Message inputMessageOne = MessageBuilder + .withPayload("Hello".getBytes()).build(); + Message inputMessageTwo = MessageBuilder + .withPayload("Hello Again".getBytes()).build(); + inputDestination.send(inputMessageOne); + inputDestination.send(inputMessageTwo); + + Message outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes()); + outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("Hello Again".getBytes()); + } + } + @EnableAutoConfiguration public static class NoEnableBindingConfiguration {