Fix imperative function|consumer composition
This commit is contained in:
@@ -83,12 +83,12 @@ import org.springframework.core.env.Environment;
|
||||
import org.springframework.core.type.MethodMetadata;
|
||||
import org.springframework.integration.channel.AbstractMessageChannel;
|
||||
import org.springframework.integration.channel.AbstractSubscribableChannel;
|
||||
import org.springframework.integration.channel.MessageChannelReactiveUtils;
|
||||
import org.springframework.integration.dsl.IntegrationFlow;
|
||||
import org.springframework.integration.dsl.IntegrationFlowBuilder;
|
||||
import org.springframework.integration.dsl.IntegrationFlows;
|
||||
import org.springframework.integration.handler.ServiceActivatingHandler;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.integration.util.IntegrationReactiveUtils;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
@@ -364,7 +364,7 @@ public class FunctionConfiguration {
|
||||
if (isReactiveOrMultipleInputOutput(bindableProxyFactory, functionType)) {
|
||||
Publisher[] inputPublishers = inputBindingNames.stream().map(inputBindingName -> {
|
||||
SubscribableChannel inputChannel = this.applicationContext.getBean(inputBindingName, SubscribableChannel.class);
|
||||
return MessageChannelReactiveUtils.toPublisher(inputChannel);
|
||||
return IntegrationReactiveUtils.messageChannelToFlux(inputChannel);
|
||||
}).toArray(Publisher[]::new);
|
||||
|
||||
Function functionToInvoke = function;
|
||||
@@ -393,7 +393,9 @@ public class FunctionConfiguration {
|
||||
}
|
||||
else {
|
||||
String outputDestinationName = this.determineOutputDestinationName(0, bindableProxyFactory, functionType);
|
||||
this.adjustFunctionForNativeEncodingIfNecessary(outputDestinationName, function, 0);
|
||||
if (StringUtils.hasText(outputDestinationName)) {
|
||||
this.adjustFunctionForNativeEncodingIfNecessary(outputDestinationName, function, 0);
|
||||
}
|
||||
String inputDestinationName = inputBindingNames.iterator().next();
|
||||
Object inputDestination = this.applicationContext.getBean(inputDestinationName);
|
||||
if (inputDestination != null && inputDestination instanceof SubscribableChannel) {
|
||||
|
||||
@@ -298,6 +298,26 @@ public class ImplicitFunctionBindingTests {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void fooFunctionComposedWithConsumerNonReactive() {
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionConsumerCopositionConfiguration.class))
|
||||
.web(WebApplicationType.NONE).run("--spring.jmx.enabled=false", "--spring.cloud.function.definition=echo|consumer")) {
|
||||
|
||||
assertThat(context.containsBean("echoconsumer-out-0")).isFalse();
|
||||
|
||||
InputDestination inputDestination = context.getBean(InputDestination.class);
|
||||
Message<byte[]> inputMessage = MessageBuilder.withPayload("Hello".getBytes()).build();
|
||||
inputDestination.send(inputMessage);
|
||||
|
||||
assertThat(System.getProperty("FunctionConsumerCopositionConfiguration")).isEqualTo("Hello");
|
||||
System.clearProperty("FunctionConsumerCopositionConfiguration");
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
@Test
|
||||
public void testReactiveConsumerWithoutDefinitionProperty() {
|
||||
System.clearProperty("spring.cloud.function.definition");
|
||||
@@ -781,6 +801,26 @@ public class ImplicitFunctionBindingTests {
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
public static class FunctionConsumerCopositionConfiguration {
|
||||
|
||||
@Bean
|
||||
public Consumer<String> consumer() {
|
||||
return v -> {
|
||||
System.out.println("==== Consuming " + v);
|
||||
System.setProperty("FunctionConsumerCopositionConfiguration", v);
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Function<String, String> echo() {
|
||||
return v -> {
|
||||
System.out.println("==> Echo " + v);
|
||||
return v;
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
public static class ReactiveFunctionConfiguration {
|
||||
|
||||
|
||||
Reference in New Issue
Block a user