diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/MessageConverterUtils.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/MessageConverterUtils.java index 4422e1455..4f9bf0880 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/MessageConverterUtils.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/MessageConverterUtils.java @@ -30,12 +30,6 @@ import org.springframework.util.StringUtils; */ public abstract class MessageConverterUtils { - /** - * An MimeType specifying a {@link Tuple}. - */ - public static final MimeType X_SPRING_TUPLE = MimeType - .valueOf("application/x-spring-tuple"); - /** * A general MimeType for Java Types. */ diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/FunctionBindingTestUtils.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/FunctionBindingTestUtils.java index 767ef498a..704721e6b 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/FunctionBindingTestUtils.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/FunctionBindingTestUtils.java @@ -22,6 +22,9 @@ import java.util.function.Function; import java.util.function.Supplier; import org.springframework.beans.factory.InitializingBean; +import org.springframework.cloud.function.context.FunctionCatalog; +import org.springframework.cloud.function.context.FunctionRegistration; +import org.springframework.cloud.function.context.catalog.BeanFactoryAwareFunctionRegistry.FunctionInvocationWrapper; import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.ProducerProperties; @@ -31,14 +34,30 @@ import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.messaging.MessageChannel; +/** + * A utility class to assist with just-in-time bindings. + * It is intended for internal framework testing and is NOTfor public use. Breaking changes are likely!!! + * + * @author Oleg Zhurakousky + * @since 3.0.2 + * + */ public class FunctionBindingTestUtils { + @SuppressWarnings("rawtypes") public static void bind(ConfigurableApplicationContext applicationContext, Object function) { try { - String functionName = function instanceof Function ? "function" : (function instanceof Consumer ? "consumer" : "supplier"); + Object targetFunction = function; + if (function instanceof FunctionRegistration) { + targetFunction = ((FunctionRegistration) function).getTarget(); + } + String functionName = targetFunction instanceof Function ? "function" : (targetFunction instanceof Consumer ? "consumer" : "supplier"); + System.setProperty("spring.cloud.function.definition", functionName); applicationContext.getBeanFactory().registerSingleton(functionName, function); + Object actualFunction = ((FunctionInvocationWrapper) applicationContext + .getBean(FunctionCatalog.class).lookup(functionName)).getTarget(); InitializingBean functionBindingRegistrar = applicationContext.getBean("functionBindingRegistrar", InitializingBean.class); functionBindingRegistrar.afterPropertiesSet(); @@ -58,15 +77,14 @@ public class FunctionBindingTestUtils { ConsumerProperties consumerProperties = inputProperties.getConsumer(); ProducerProperties producerProperties = outputProperties.getProducer(); - TestChannelBinder binder = applicationContext.getBean(TestChannelBinder.class); - if (function instanceof Supplier || function instanceof Function) { + if (actualFunction instanceof Supplier || actualFunction instanceof Function) { Binding bindProducer = binder.bindProducer(outputProperties.getDestination(), applicationContext.getBean(outputBindingName, MessageChannel.class), producerProperties == null ? new ProducerProperties() : producerProperties); bindProducer.start(); } - if (function instanceof Consumer || function instanceof Function) { + if (actualFunction instanceof Consumer || actualFunction instanceof Function) { Binding bindConsumer = binder.bindConsumer(inputProperties.getDestination(), null, applicationContext.getBean(inputBindingName, MessageChannel.class), consumerProperties == null ? new ConsumerProperties() : consumerProperties); 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 105bf89e7..26645855a 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 @@ -33,6 +33,8 @@ 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.function.context.FunctionRegistration; +import org.springframework.cloud.function.context.FunctionType; import org.springframework.cloud.function.context.config.ContextFunctionCatalogAutoConfiguration; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.StreamListener; @@ -65,31 +67,60 @@ public class ImplicitFunctionBindingTests { @After public void after() { System.clearProperty("spring.cloud.function.definition"); - System.clearProperty("spring.cloud.function.definition"); + } + + @SuppressWarnings({ "unchecked", "rawtypes" }) + @Test + public void dynamicBindingTestWithFunctionRegistrationAndExplicitDestination() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(EmptyConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.function-in-0.destination=input")) { + + InputDestination input = context.getBean(InputDestination.class); + try { + input.send(new GenericMessage("hello".getBytes())); + fail(); // it should since there are no functions and no bindings + } + catch (Exception e) { + // good, we expected it + } + + Function function = v -> v; + FunctionRegistration functionRegistration = new FunctionRegistration(function, "function"); + functionRegistration = functionRegistration.type(FunctionType.from(byte[].class).to(byte[].class)); + FunctionBindingTestUtils.bind(context, functionRegistration); + + input.send(new GenericMessage("hello".getBytes()), "input"); + + OutputDestination output = context.getBean(OutputDestination.class); + assertThat(output.receive(1000, "function-out-0").getPayload()).isEqualTo("hello".getBytes()); + } } @Test - public void marcinsTest() { - ConfigurableApplicationContext context = new SpringApplicationBuilder( + public void dynamicBindingTestWithFunction() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(EmptyConfiguration.class)) - .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false"); + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false")) { + InputDestination input = context.getBean(InputDestination.class); + try { + input.send(new GenericMessage("hello".getBytes())); + fail(); // it should since there are no functions and no bindings + } + catch (Exception e) { + // good, we expected it + } + + Function function = v -> v.toUpperCase(); + FunctionBindingTestUtils.bind(context, function); - InputDestination input = context.getBean(InputDestination.class); - try { input.send(new GenericMessage("hello".getBytes())); - fail(); // it should since there are no functions and no bindings + + OutputDestination output = context.getBean(OutputDestination.class); + assertThat(new String(output.receive().getPayload())).isEqualTo("HELLO"); } - catch (Exception e) { - // good, we expected it - } - - Function function = v -> v.toUpperCase(); - FunctionBindingTestUtils.bind(context, function); - - input.send(new GenericMessage("hello".getBytes())); - - OutputDestination output = context.getBean(OutputDestination.class); - assertThat(new String(output.receive().getPayload())).isEqualTo("HELLO"); } @Test