diff --git a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 3a618e32e..b7b3b7513 100644 --- a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2022 the original author or authors. + * Copyright 2019-2023 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. @@ -25,6 +25,7 @@ import java.util.Collection; import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.function.BiConsumer; import java.util.function.Consumer; import java.util.function.Function; import java.util.function.Supplier; @@ -34,6 +35,7 @@ import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; +import org.springframework.cloud.function.context.FunctionCatalog; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.publisher.Sinks; @@ -78,6 +80,7 @@ import static org.junit.Assert.fail; /** * * @author Oleg Zhurakousky + * @author Soby Chacko * */ public class ImplicitFunctionBindingTests { @@ -657,6 +660,18 @@ public class ImplicitFunctionBindingTests { } } + @Test + void functionInvocationWrapperReflectsBiConsumerTargetFunctionType() { + System.clearProperty("spring.cloud.function.definition"); + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(WrappedBiConsumerAutoConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false")) { + FunctionCatalog functionCatalog = context.getBean(FunctionCatalog.class); + FunctionInvocationWrapper functionWrapper = functionCatalog.lookup("testBiConsumer"); + assertThat(functionWrapper.isWrappedBiConsumer()).isTrue(); + } + } + @Test void testCollectionAndMapConversionDuringComposition() { System.clearProperty("spring.cloud.function.definition"); @@ -1617,6 +1632,16 @@ public class ImplicitFunctionBindingTests { } } + @EnableAutoConfiguration + public static class WrappedBiConsumerAutoConfiguration { + + @Bean + public BiConsumer> testBiConsumer() { + return (a, b) -> { }; + } + + } + public static class Person { private String name; private int id; diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index e8125269d..e0ed79821 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -839,7 +839,12 @@ public class FunctionConfiguration { } else { this.inputCount = FunctionTypeUtils.getInputCount(function); - this.outputCount = this.getOutputCount(function, false); + if (function.isWrappedBiConsumer()) { + this.outputCount = 0; + } + else { + this.outputCount = this.getOutputCount(function, false); + } } AtomicReference proxyFactory = new AtomicReference<>();