Additional enhancements to dynamic binding test feature

This commit is contained in:
Oleg Zhurakousky
2020-02-03 13:23:43 +01:00
parent ce221adf0d
commit 2fca323654
3 changed files with 71 additions and 28 deletions

View File

@@ -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.
*/

View File

@@ -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<MessageChannel> 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<MessageChannel> bindConsumer = binder.bindConsumer(inputProperties.getDestination(), null,
applicationContext.getBean(inputBindingName, MessageChannel.class),
consumerProperties == null ? new ConsumerProperties() : consumerProperties);

View File

@@ -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<byte[]>("hello".getBytes()));
fail(); // it should since there are no functions and no bindings
}
catch (Exception e) {
// good, we expected it
}
Function<byte[], byte[]> 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<byte[]>("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<byte[]>("hello".getBytes()));
fail(); // it should since there are no functions and no bindings
}
catch (Exception e) {
// good, we expected it
}
Function<String, String> function = v -> v.toUpperCase();
FunctionBindingTestUtils.bind(context, function);
InputDestination input = context.getBean(InputDestination.class);
try {
input.send(new GenericMessage<byte[]>("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<String, String> function = v -> v.toUpperCase();
FunctionBindingTestUtils.bind(context, function);
input.send(new GenericMessage<byte[]>("hello".getBytes()));
OutputDestination output = context.getBean(OutputDestination.class);
assertThat(new String(output.receive().getPayload())).isEqualTo("HELLO");
}
@Test