GH-1794 Disabled function initialization for legacy cases
Fix bypassing of the functional configuration for cases where legacy configuration elements are present (e.g., EnableBinding) Resolves #1794
This commit is contained in:
@@ -41,6 +41,7 @@ import org.springframework.cloud.function.context.PollableSupplier;
|
||||
import org.springframework.cloud.function.context.catalog.BeanFactoryAwareFunctionRegistry.FunctionInvocationWrapper;
|
||||
import org.springframework.cloud.function.context.catalog.FunctionInspector;
|
||||
import org.springframework.cloud.function.context.catalog.FunctionTypeUtils;
|
||||
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;
|
||||
@@ -82,39 +83,42 @@ public class FunctionConfiguration {
|
||||
|
||||
@Bean
|
||||
public InitializingBean functionChannelBindingInitializer(FunctionCatalog functionCatalog, FunctionInspector functionInspector,
|
||||
StreamFunctionProperties functionProperties, @Nullable BindableProxyFactory[] bindableProxyFactory) {
|
||||
StreamFunctionProperties functionProperties, @Nullable BindableProxyFactory[] bindableProxyFactory, GenericApplicationContext context) {
|
||||
return new FunctionChannelBindingInitializer(functionCatalog, functionInspector, functionProperties,
|
||||
ObjectUtils.isEmpty(bindableProxyFactory) ? null : bindableProxyFactory[0]);
|
||||
ObjectUtils.isEmpty(bindableProxyFactory) ? null : bindableProxyFactory[0]);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public IntegrationFlow standAloneSupplierFlow(FunctionCatalog functionCatalog, FunctionInspector functionInspector,
|
||||
StreamFunctionProperties functionProperties, GenericApplicationContext context) {
|
||||
|
||||
IntegrationFlow integrationFlow = null;
|
||||
FunctionInvocationWrapper functionWrapper = functionCatalog.lookup(functionProperties.getDefinition());
|
||||
if (functionWrapper != null) {
|
||||
AtomicReference<MonoSink<Object>> triggerRef = new AtomicReference<>();
|
||||
Publisher<Object> beginPublishingTrigger = Mono.create(emmiter -> {
|
||||
triggerRef.set(emmiter);
|
||||
});
|
||||
context.addApplicationListener(event -> {
|
||||
if (event instanceof BindingCreatedEvent) {
|
||||
if (triggerRef.get() != null) {
|
||||
triggerRef.get().success();
|
||||
if (functionCatalog != null && ObjectUtils.isEmpty(context.getBeanNamesForAnnotation(EnableBinding.class))) {
|
||||
FunctionInvocationWrapper functionWrapper = functionCatalog.lookup(functionProperties.getDefinition());
|
||||
if (functionWrapper != null) {
|
||||
AtomicReference<MonoSink<Object>> triggerRef = new AtomicReference<>();
|
||||
Publisher<Object> beginPublishingTrigger = Mono.create(emmiter -> {
|
||||
triggerRef.set(emmiter);
|
||||
});
|
||||
context.addApplicationListener(event -> {
|
||||
if (event instanceof BindingCreatedEvent) {
|
||||
if (triggerRef.get() != null) {
|
||||
triggerRef.get().success();
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
RootBeanDefinition bd = (RootBeanDefinition) context.getBeanDefinition(functionProperties.getParsedDefinition()[0]);
|
||||
Method factoryMethod = bd.getResolvedFactoryMethod();
|
||||
PollableSupplier pollable = factoryMethod.getReturnType().isAssignableFrom(Supplier.class)
|
||||
? AnnotationUtils.findAnnotation(factoryMethod, PollableSupplier.class)
|
||||
: null;
|
||||
|
||||
if (!functionProperties.isComposeFrom() && !functionProperties.isComposeTo() && functionWrapper.isSupplier()) {
|
||||
integrationFlow = this.integrationFlowFromProvidedSupplier(functionWrapper, functionInspector, beginPublishingTrigger, pollable)
|
||||
.channel("output").get();
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
RootBeanDefinition bd = (RootBeanDefinition) context.getBeanDefinition(functionProperties.getParsedDefinition()[0]);
|
||||
Method factoryMethod = bd.getResolvedFactoryMethod();
|
||||
PollableSupplier pollable = factoryMethod.getReturnType().isAssignableFrom(Supplier.class)
|
||||
? AnnotationUtils.findAnnotation(factoryMethod, PollableSupplier.class)
|
||||
: null;
|
||||
|
||||
if (!functionProperties.isComposeFrom() && !functionProperties.isComposeTo() && functionWrapper.isSupplier()) {
|
||||
integrationFlow = this.integrationFlowFromProvidedSupplier(functionWrapper, functionInspector, beginPublishingTrigger, pollable)
|
||||
.channel("output").get();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -166,7 +166,6 @@ public class GreenfieldFunctionEnableBindingTests {
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Source.class)
|
||||
public static class SourceFromSupplier {
|
||||
|
||||
@Bean
|
||||
|
||||
@@ -26,9 +26,12 @@ import reactor.core.publisher.Flux;
|
||||
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.annotation.StreamListener;
|
||||
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.Sink;
|
||||
import org.springframework.context.ConfigurableApplicationContext;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
@@ -178,7 +181,19 @@ public class ImplicitFunctionBindingTests {
|
||||
assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes());
|
||||
outputMessage = outputDestination.receive();
|
||||
assertThat(outputMessage.getPayload()).isEqualTo("Hello Again".getBytes());
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFunctionConfigDisabledIfStreamListenerIsUsed() {
|
||||
System.clearProperty("spring.cloud.stream.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
LegacyConfiguration.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.jmx.enabled=false")) {
|
||||
|
||||
assertThat(context.getBean("standAloneSupplierFlow")).isEqualTo(null);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -237,4 +252,14 @@ public class ImplicitFunctionBindingTests {
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@EnableBinding(Sink.class)
|
||||
public static class LegacyConfiguration {
|
||||
|
||||
@StreamListener(Sink.INPUT)
|
||||
public void handle(String value) {
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -232,6 +232,9 @@ public class SourceToFunctionsSupportTests {
|
||||
assertThat(new String(target.receive(2000).getPayload())).isEqualTo("4");
|
||||
assertThat(new String(target.receive(2000).getPayload())).isEqualTo("5");
|
||||
assertThat(new String(target.receive(2000).getPayload())).isEqualTo("6");
|
||||
|
||||
assertThat(context.getBean("standAloneSupplierFlow")).isNotEqualTo(null);
|
||||
assertThat(context.getBean("functionChannelBindingInitializer")).isNotEqualTo(null);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user