Add support for detecting FunctionRegistration or Function
User can now provide a Function or an ApplicationInitializer. Also the initializer can create a FunctionRegistration with the handler name instead of a bean with the handler name. Better control of input and output types that way. Fixes gh-231
This commit is contained in:
@@ -5,7 +5,7 @@
|
|||||||
|
|
||||||
<groupId>com.example</groupId>
|
<groupId>com.example</groupId>
|
||||||
<artifactId>flux-sample</artifactId>
|
<artifactId>flux-sample</artifactId>
|
||||||
<version>1.0.0.M1</version>
|
<version>1.0.0.RC1</version>
|
||||||
<packaging>jar</packaging>
|
<packaging>jar</packaging>
|
||||||
|
|
||||||
<parent>
|
<parent>
|
||||||
@@ -18,7 +18,7 @@
|
|||||||
<properties>
|
<properties>
|
||||||
<java.version>1.8</java.version>
|
<java.version>1.8</java.version>
|
||||||
<spring-cloud-function.version>2.0.0.BUILD-SNAPSHOT</spring-cloud-function.version>
|
<spring-cloud-function.version>2.0.0.BUILD-SNAPSHOT</spring-cloud-function.version>
|
||||||
<wrapper.version>1.0.15.RELEASE</wrapper.version>
|
<wrapper.version>1.0.17.RELEASE</wrapper.version>
|
||||||
</properties>
|
</properties>
|
||||||
|
|
||||||
<dependencies>
|
<dependencies>
|
||||||
|
|||||||
@@ -18,7 +18,7 @@
|
|||||||
<properties>
|
<properties>
|
||||||
<java.version>1.8</java.version>
|
<java.version>1.8</java.version>
|
||||||
<spring-cloud-function.version>2.0.0.BUILD-SNAPSHOT</spring-cloud-function.version>
|
<spring-cloud-function.version>2.0.0.BUILD-SNAPSHOT</spring-cloud-function.version>
|
||||||
<wrapper.version>1.0.12.RELEASE</wrapper.version>
|
<wrapper.version>1.0.17.RELEASE</wrapper.version>
|
||||||
</properties>
|
</properties>
|
||||||
|
|
||||||
<dependencies>
|
<dependencies>
|
||||||
|
|||||||
@@ -18,13 +18,10 @@ package org.springframework.cloud.function.deployer;
|
|||||||
|
|
||||||
import java.lang.reflect.Field;
|
import java.lang.reflect.Field;
|
||||||
import java.net.URL;
|
import java.net.URL;
|
||||||
import java.util.Collections;
|
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
||||||
import org.springframework.beans.BeanUtils;
|
|
||||||
import org.springframework.boot.SpringApplication;
|
import org.springframework.boot.SpringApplication;
|
||||||
import org.springframework.context.ApplicationContext;
|
import org.springframework.cloud.function.context.FunctionalSpringApplication;
|
||||||
import org.springframework.context.ApplicationContextInitializer;
|
|
||||||
import org.springframework.context.ConfigurableApplicationContext;
|
import org.springframework.context.ConfigurableApplicationContext;
|
||||||
import org.springframework.core.env.MapPropertySource;
|
import org.springframework.core.env.MapPropertySource;
|
||||||
import org.springframework.core.env.StandardEnvironment;
|
import org.springframework.core.env.StandardEnvironment;
|
||||||
@@ -59,16 +56,7 @@ public class ContextRunner {
|
|||||||
new MapPropertySource("appDeployer", properties));
|
new MapPropertySource("appDeployer", properties));
|
||||||
running = true;
|
running = true;
|
||||||
Class<?> sourceClass = ClassUtils.resolveClassName(source, null);
|
Class<?> sourceClass = ClassUtils.resolveClassName(source, null);
|
||||||
ApplicationContextInitializer<?> initializer = null;
|
|
||||||
if (ApplicationContextInitializer.class.isAssignableFrom(sourceClass)) {
|
|
||||||
initializer = BeanUtils.instantiateClass(sourceClass, ApplicationContextInitializer.class);
|
|
||||||
sourceClass = Dummy.class;
|
|
||||||
}
|
|
||||||
SpringApplication builder = builder(sourceClass);
|
SpringApplication builder = builder(sourceClass);
|
||||||
if (initializer!=null) {
|
|
||||||
builder.addInitializers(initializer);
|
|
||||||
builder.setDefaultProperties(Collections.singletonMap("spring.functional.enabled", "true"));
|
|
||||||
}
|
|
||||||
builder.setEnvironment(environment);
|
builder.setEnvironment(environment);
|
||||||
builder.setRegisterShutdownHook(false);
|
builder.setRegisterShutdownHook(false);
|
||||||
context = builder.run(args);
|
context = builder.run(args);
|
||||||
@@ -131,19 +119,8 @@ public class ContextRunner {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private static SpringApplication builder(Class<?> type) {
|
private static SpringApplication builder(Class<?> type) {
|
||||||
if (type==Dummy.class) {
|
SpringApplication application = new FunctionalSpringApplication(type);
|
||||||
SpringApplication application = new SpringApplication() {
|
return application;
|
||||||
@Override
|
|
||||||
protected void load(ApplicationContext context, Object[] sources) {
|
|
||||||
}
|
|
||||||
};
|
|
||||||
// Boot doesn't allow null sources
|
|
||||||
application.setSources(Collections.singleton(Dummy.class.getName()));
|
|
||||||
return application;
|
|
||||||
}
|
|
||||||
return new SpringApplication(type);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private class Dummy {}
|
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -51,10 +51,12 @@ import org.springframework.boot.loader.archive.Archive;
|
|||||||
import org.springframework.boot.loader.archive.ExplodedArchive;
|
import org.springframework.boot.loader.archive.ExplodedArchive;
|
||||||
import org.springframework.boot.loader.archive.JarFileArchive;
|
import org.springframework.boot.loader.archive.JarFileArchive;
|
||||||
import org.springframework.cloud.deployer.resource.support.DelegatingResourceLoader;
|
import org.springframework.cloud.deployer.resource.support.DelegatingResourceLoader;
|
||||||
|
import org.springframework.cloud.function.context.FunctionCatalog;
|
||||||
import org.springframework.cloud.function.context.FunctionRegistration;
|
import org.springframework.cloud.function.context.FunctionRegistration;
|
||||||
import org.springframework.cloud.function.context.FunctionRegistry;
|
import org.springframework.cloud.function.context.FunctionRegistry;
|
||||||
import org.springframework.cloud.function.context.FunctionType;
|
import org.springframework.cloud.function.context.FunctionType;
|
||||||
import org.springframework.cloud.function.context.catalog.FunctionInspector;
|
import org.springframework.cloud.function.context.catalog.FunctionInspector;
|
||||||
|
import org.springframework.cloud.function.core.FluxFunction;
|
||||||
import org.springframework.context.ConfigurableApplicationContext;
|
import org.springframework.context.ConfigurableApplicationContext;
|
||||||
import org.springframework.context.annotation.Configuration;
|
import org.springframework.context.annotation.Configuration;
|
||||||
import org.springframework.core.env.Environment;
|
import org.springframework.core.env.Environment;
|
||||||
@@ -381,6 +383,32 @@ class FunctionCreatorConfiguration {
|
|||||||
Object result = null;
|
Object result = null;
|
||||||
if (this.runner != null) {
|
if (this.runner != null) {
|
||||||
result = this.runner.getBean(type);
|
result = this.runner.getBean(type);
|
||||||
|
if (result == null) {
|
||||||
|
if (this.runner.containsBean(FunctionCatalog.class.getName())) {
|
||||||
|
Object catalog = this.runner
|
||||||
|
.getBean(FunctionCatalog.class.getName());
|
||||||
|
result = this.runner.evaluate("lookup(#function).getTarget()",
|
||||||
|
catalog, "function", type);
|
||||||
|
if (result != null) {
|
||||||
|
logger.info("Located registration: " + type + " of type "
|
||||||
|
+ result.getClass());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
logger.info("Located bean: " + type + " of type "
|
||||||
|
+ result.getClass());
|
||||||
|
if (result.getClass().getName()
|
||||||
|
.equals(FunctionRegistration.class.getName())) {
|
||||||
|
result = this.runner.evaluate("getTarget()", result);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (result != null) {
|
||||||
|
if (result.getClass().getName()
|
||||||
|
.equals(FluxFunction.class.getName())) {
|
||||||
|
result = this.runner.evaluate("getTarget()", result);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
if (result == null) {
|
if (result == null) {
|
||||||
logger.info("No bean found. Instantiating: " + type);
|
logger.info("No bean found. Instantiating: " + type);
|
||||||
@@ -390,7 +418,6 @@ class FunctionCreatorConfiguration {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (result != null) {
|
if (result != null) {
|
||||||
logger.info("Located bean: " + type);
|
|
||||||
return result;
|
return result;
|
||||||
}
|
}
|
||||||
throw new IllegalStateException("Cannot create bean for: " + type);
|
throw new IllegalStateException("Cannot create bean for: " + type);
|
||||||
|
|||||||
@@ -21,7 +21,6 @@ import java.util.function.Supplier;
|
|||||||
|
|
||||||
import org.junit.Assume;
|
import org.junit.Assume;
|
||||||
import org.junit.BeforeClass;
|
import org.junit.BeforeClass;
|
||||||
import org.junit.Ignore;
|
|
||||||
import org.junit.Rule;
|
import org.junit.Rule;
|
||||||
import org.junit.Test;
|
import org.junit.Test;
|
||||||
import org.junit.runner.RunWith;
|
import org.junit.runner.RunWith;
|
||||||
@@ -97,13 +96,45 @@ public abstract class FunctionCreatorConfigurationTests {
|
|||||||
@EnableAutoConfiguration
|
@EnableAutoConfiguration
|
||||||
@TestPropertySource(properties = {
|
@TestPropertySource(properties = {
|
||||||
"function.location=app:classpath,file:target/test-classes,file:target/test-classes/app",
|
"function.location=app:classpath,file:target/test-classes,file:target/test-classes/app",
|
||||||
"function.bean=myDoubler",
|
"function.bean=doubler",
|
||||||
"function.main=org.springframework.cloud.function.test.FunctionRegistrar" })
|
"function.main=org.springframework.cloud.function.test.FunctionRegistrar" })
|
||||||
public static class SingleFunctionWithRegistrarTests
|
public static class SingleFunctionWithRegistrarTests
|
||||||
extends FunctionCreatorConfigurationTests {
|
extends FunctionCreatorConfigurationTests {
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
@Ignore // related to boot 2.1 no bean override change
|
public void testDouble() {
|
||||||
|
Function<Flux<Integer>, Flux<Integer>> function = catalog
|
||||||
|
.lookup(Function.class, "function0");
|
||||||
|
assertThat(function.apply(Flux.just(2)).blockFirst()).isEqualTo(4);
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
@EnableAutoConfiguration
|
||||||
|
@TestPropertySource(properties = {
|
||||||
|
"function.location=app:classpath,file:target/test-classes,file:target/test-classes/app",
|
||||||
|
"function.bean=frenchizer",
|
||||||
|
"function.main=org.springframework.cloud.function.test.FunctionRegistrar" })
|
||||||
|
public static class SingleFunctionWithRegistrarAndRegistrationTests
|
||||||
|
extends FunctionCreatorConfigurationTests {
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void testFrenchize() {
|
||||||
|
Function<Flux<Integer>, Flux<String>> function = catalog
|
||||||
|
.lookup(Function.class, "function0");
|
||||||
|
assertThat(function.apply(Flux.just(2)).blockFirst()).isEqualTo("deux");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@EnableAutoConfiguration
|
||||||
|
@TestPropertySource(properties = {
|
||||||
|
"function.location=app:classpath,file:target/test-classes,file:target/test-classes/app",
|
||||||
|
"function.bean=myDoubler",
|
||||||
|
"function.main=org.springframework.cloud.function.test.FunctionInitializer" })
|
||||||
|
public static class SingleFunctionWithInitializerTests
|
||||||
|
extends FunctionCreatorConfigurationTests {
|
||||||
|
|
||||||
|
@Test
|
||||||
public void testDouble() {
|
public void testDouble() {
|
||||||
Function<Flux<Integer>, Flux<Integer>> function = catalog
|
Function<Flux<Integer>, Flux<Integer>> function = catalog
|
||||||
.lookup(Function.class, "function0");
|
.lookup(Function.class, "function0");
|
||||||
|
|||||||
@@ -0,0 +1,53 @@
|
|||||||
|
/*
|
||||||
|
* Copyright 2017 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.
|
||||||
|
* You may obtain a copy of the License at
|
||||||
|
*
|
||||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
*
|
||||||
|
* Unless required by applicable law or agreed to in writing, software
|
||||||
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
* See the License for the specific language governing permissions and
|
||||||
|
* limitations under the License.
|
||||||
|
*/
|
||||||
|
|
||||||
|
package org.springframework.cloud.function.test;
|
||||||
|
|
||||||
|
import org.springframework.boot.SpringApplication;
|
||||||
|
import org.springframework.cloud.function.context.FunctionalSpringApplication;
|
||||||
|
import org.springframework.context.ApplicationContextInitializer;
|
||||||
|
import org.springframework.context.annotation.Bean;
|
||||||
|
import org.springframework.context.support.GenericApplicationContext;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @author Dave Syer
|
||||||
|
*/
|
||||||
|
public class FunctionInitializer
|
||||||
|
implements ApplicationContextInitializer<GenericApplicationContext> {
|
||||||
|
|
||||||
|
@Bean
|
||||||
|
public Doubler myDoubler() {
|
||||||
|
return new Doubler();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Bean
|
||||||
|
public Frenchizer myFrenchizer() {
|
||||||
|
return new Frenchizer();
|
||||||
|
}
|
||||||
|
|
||||||
|
public static void main(String[] args) throws Exception {
|
||||||
|
SpringApplication application = new FunctionalSpringApplication(
|
||||||
|
FunctionInitializer.class);
|
||||||
|
application.run(args);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void initialize(GenericApplicationContext context) {
|
||||||
|
// TODO: support for FunctionRegistration
|
||||||
|
context.registerBean("myDoubler", Doubler.class, () -> myDoubler());
|
||||||
|
context.registerBean("myFrenchizer", Frenchizer.class, () -> myFrenchizer());
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -16,10 +16,10 @@
|
|||||||
|
|
||||||
package org.springframework.cloud.function.test;
|
package org.springframework.cloud.function.test;
|
||||||
|
|
||||||
import java.util.Collections;
|
|
||||||
|
|
||||||
import org.springframework.boot.SpringApplication;
|
import org.springframework.boot.SpringApplication;
|
||||||
import org.springframework.context.ApplicationContext;
|
import org.springframework.cloud.function.context.FunctionRegistration;
|
||||||
|
import org.springframework.cloud.function.context.FunctionType;
|
||||||
|
import org.springframework.cloud.function.context.FunctionalSpringApplication;
|
||||||
import org.springframework.context.ApplicationContextInitializer;
|
import org.springframework.context.ApplicationContextInitializer;
|
||||||
import org.springframework.context.annotation.Bean;
|
import org.springframework.context.annotation.Bean;
|
||||||
import org.springframework.context.support.GenericApplicationContext;
|
import org.springframework.context.support.GenericApplicationContext;
|
||||||
@@ -41,21 +41,21 @@ public class FunctionRegistrar
|
|||||||
}
|
}
|
||||||
|
|
||||||
public static void main(String[] args) throws Exception {
|
public static void main(String[] args) throws Exception {
|
||||||
SpringApplication application = new SpringApplication(Object.class) {
|
SpringApplication application = new FunctionalSpringApplication(
|
||||||
@Override
|
FunctionRegistrar.class);
|
||||||
protected void load(ApplicationContext context, Object[] sources) {
|
|
||||||
}
|
|
||||||
};
|
|
||||||
application.addInitializers(new FunctionRegistrar());
|
|
||||||
application.setDefaultProperties(
|
|
||||||
Collections.singletonMap("spring.functional.enabled", "true"));
|
|
||||||
application.run(args);
|
application.run(args);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void initialize(GenericApplicationContext context) {
|
public void initialize(GenericApplicationContext context) {
|
||||||
// TODO: support for FunctionRegistration
|
context.registerBean("theDoubler", FunctionRegistration.class,
|
||||||
context.registerBean("myDoubler", Doubler.class, () -> myDoubler());
|
() -> new FunctionRegistration<>(myDoubler(), "doubler")
|
||||||
context.registerBean("myFrenchizer", Frenchizer.class, () -> myFrenchizer());
|
.type(FunctionType.of((Doubler.class))));
|
||||||
|
context.registerBean("frenchizer", FunctionRegistration.class, () -> {
|
||||||
|
Frenchizer function = myFrenchizer();
|
||||||
|
function.init();
|
||||||
|
return new FunctionRegistration<>(function, "theFrenchizer")
|
||||||
|
.type(FunctionType.of((Frenchizer.class)));
|
||||||
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -16,10 +16,13 @@
|
|||||||
|
|
||||||
package org.springframework.cloud.function.web;
|
package org.springframework.cloud.function.web;
|
||||||
|
|
||||||
|
import java.lang.reflect.Method;
|
||||||
import java.lang.reflect.ParameterizedType;
|
import java.lang.reflect.ParameterizedType;
|
||||||
import java.lang.reflect.Type;
|
import java.lang.reflect.Type;
|
||||||
|
import java.util.Arrays;
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
@@ -38,15 +41,29 @@ import org.springframework.cloud.function.context.message.MessageUtils;
|
|||||||
import org.springframework.cloud.function.core.FluxWrapper;
|
import org.springframework.cloud.function.core.FluxWrapper;
|
||||||
import org.springframework.cloud.function.json.JsonMapper;
|
import org.springframework.cloud.function.json.JsonMapper;
|
||||||
import org.springframework.cloud.function.web.util.HeaderUtils;
|
import org.springframework.cloud.function.web.util.HeaderUtils;
|
||||||
|
import org.springframework.core.MethodParameter;
|
||||||
|
import org.springframework.core.ReactiveAdapter;
|
||||||
|
import org.springframework.core.ReactiveAdapterRegistry;
|
||||||
import org.springframework.core.ResolvableType;
|
import org.springframework.core.ResolvableType;
|
||||||
|
import org.springframework.core.codec.DecodingException;
|
||||||
|
import org.springframework.core.codec.Hints;
|
||||||
import org.springframework.http.HttpHeaders;
|
import org.springframework.http.HttpHeaders;
|
||||||
import org.springframework.http.HttpStatus;
|
import org.springframework.http.HttpStatus;
|
||||||
|
import org.springframework.http.MediaType;
|
||||||
import org.springframework.http.ResponseEntity;
|
import org.springframework.http.ResponseEntity;
|
||||||
import org.springframework.http.ResponseEntity.BodyBuilder;
|
import org.springframework.http.ResponseEntity.BodyBuilder;
|
||||||
|
import org.springframework.http.codec.HttpMessageReader;
|
||||||
|
import org.springframework.http.codec.ServerCodecConfigurer;
|
||||||
|
import org.springframework.http.server.reactive.ServerHttpRequest;
|
||||||
|
import org.springframework.http.server.reactive.ServerHttpResponse;
|
||||||
import org.springframework.messaging.Message;
|
import org.springframework.messaging.Message;
|
||||||
import org.springframework.util.LinkedMultiValueMap;
|
import org.springframework.util.LinkedMultiValueMap;
|
||||||
import org.springframework.util.MultiValueMap;
|
import org.springframework.util.MultiValueMap;
|
||||||
|
import org.springframework.util.ReflectionUtils;
|
||||||
import org.springframework.util.StringUtils;
|
import org.springframework.util.StringUtils;
|
||||||
|
import org.springframework.web.server.ServerWebExchange;
|
||||||
|
import org.springframework.web.server.ServerWebInputException;
|
||||||
|
import org.springframework.web.server.UnsupportedMediaTypeStatusException;
|
||||||
|
|
||||||
import reactor.core.publisher.Flux;
|
import reactor.core.publisher.Flux;
|
||||||
import reactor.core.publisher.Mono;
|
import reactor.core.publisher.Mono;
|
||||||
@@ -65,11 +82,16 @@ public class RequestProcessor {
|
|||||||
|
|
||||||
private final JsonMapper mapper;
|
private final JsonMapper mapper;
|
||||||
|
|
||||||
public RequestProcessor(FunctionInspector inspector, ObjectProvider<JsonMapper> mapper,
|
private final List<HttpMessageReader<?>> messageReaders;
|
||||||
StringConverter converter) {
|
|
||||||
|
public RequestProcessor(FunctionInspector inspector,
|
||||||
|
ObjectProvider<JsonMapper> mapper, StringConverter converter,
|
||||||
|
ObjectProvider<ServerCodecConfigurer> codecs) {
|
||||||
this.mapper = mapper.getIfAvailable();
|
this.mapper = mapper.getIfAvailable();
|
||||||
this.inspector = inspector;
|
this.inspector = inspector;
|
||||||
this.converter = converter;
|
this.converter = converter;
|
||||||
|
ServerCodecConfigurer source = codecs.getIfAvailable();
|
||||||
|
this.messageReaders = source == null ? null : source.getReaders();
|
||||||
}
|
}
|
||||||
|
|
||||||
public static FunctionWrapper wrapper(
|
public static FunctionWrapper wrapper(
|
||||||
@@ -95,6 +117,12 @@ public class RequestProcessor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public Mono<ResponseEntity<?>> post(FunctionWrapper wrapper,
|
||||||
|
ServerWebExchange exchange) {
|
||||||
|
return Mono.from(body(wrapper.handler(), exchange))
|
||||||
|
.flatMap(body -> response(wrapper, body, false));
|
||||||
|
}
|
||||||
|
|
||||||
public Mono<ResponseEntity<?>> post(FunctionWrapper wrapper, String body,
|
public Mono<ResponseEntity<?>> post(FunctionWrapper wrapper, String body,
|
||||||
boolean stream) {
|
boolean stream) {
|
||||||
Object function = wrapper.handler();
|
Object function = wrapper.handler();
|
||||||
@@ -102,12 +130,15 @@ public class RequestProcessor {
|
|||||||
Type itemType = getItemType(function);
|
Type itemType = getItemType(function);
|
||||||
|
|
||||||
Object input = body;
|
Object input = body;
|
||||||
if (StringUtils.hasText(body) && this.mapper!=null) {
|
if (StringUtils.hasText(body) && this.mapper != null) {
|
||||||
if (body.startsWith("[")) {
|
if (body.startsWith("[")) {
|
||||||
Class<?> collectionType = Collection.class.isAssignableFrom(inputType) ? inputType : Collection.class;
|
Class<?> collectionType = Collection.class.isAssignableFrom(inputType)
|
||||||
input = mapper.toObject(body, ResolvableType
|
? inputType
|
||||||
.forClassWithGenerics(collectionType, (Class<?>) itemType)
|
: Collection.class;
|
||||||
.getType());
|
input = mapper.toObject(body,
|
||||||
|
ResolvableType
|
||||||
|
.forClassWithGenerics(collectionType, (Class<?>) itemType)
|
||||||
|
.getType());
|
||||||
}
|
}
|
||||||
else {
|
else {
|
||||||
if (inputType == String.class) {
|
if (inputType == String.class) {
|
||||||
@@ -124,7 +155,7 @@ public class RequestProcessor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return post(wrapper, input, null, stream);
|
return response(wrapper, input, stream);
|
||||||
}
|
}
|
||||||
|
|
||||||
public Mono<ResponseEntity<?>> stream(FunctionWrapper request) {
|
public Mono<ResponseEntity<?>> stream(FunctionWrapper request) {
|
||||||
@@ -134,19 +165,17 @@ public class RequestProcessor {
|
|||||||
return stream(request, result);
|
return stream(request, result);
|
||||||
}
|
}
|
||||||
|
|
||||||
private Mono<ResponseEntity<?>> post(FunctionWrapper wrapper, Object body,
|
private Mono<ResponseEntity<?>> response(FunctionWrapper wrapper, Object body,
|
||||||
MultiValueMap<String, String> params, boolean stream) {
|
boolean stream) {
|
||||||
|
|
||||||
Iterable<?> iterable = body instanceof Collection ? (Collection<?>) body
|
Iterable<?> iterable = body instanceof Collection ? (Collection<?>) body
|
||||||
: (body instanceof Set ? Collections.singleton(body) : Collections.singletonList(body));
|
: (body instanceof Set ? Collections.singleton(body)
|
||||||
|
: Collections.singletonList(body));
|
||||||
|
|
||||||
Function<Publisher<?>, Publisher<?>> function = wrapper.function();
|
Function<Publisher<?>, Publisher<?>> function = wrapper.function();
|
||||||
Consumer<Publisher<?>> consumer = wrapper.consumer();
|
Consumer<Publisher<?>> consumer = wrapper.consumer();
|
||||||
|
|
||||||
MultiValueMap<String, String> form = wrapper.params();
|
MultiValueMap<String, String> form = wrapper.params();
|
||||||
if (params != null) {
|
|
||||||
form.putAll(params);
|
|
||||||
}
|
|
||||||
|
|
||||||
boolean inputIsCollection = Collection.class
|
boolean inputIsCollection = Collection.class
|
||||||
.isAssignableFrom(inspector.getInputType(wrapper.handler()));
|
.isAssignableFrom(inspector.getInputType(wrapper.handler()));
|
||||||
@@ -250,6 +279,98 @@ public class RequestProcessor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private Publisher<?> body(Object handler, ServerWebExchange exchange) {
|
||||||
|
|
||||||
|
ResolvableType elementType = ResolvableType
|
||||||
|
.forClass(this.inspector.getInputType(handler));
|
||||||
|
ResolvableType actualType = elementType;
|
||||||
|
Class<?> resolvedType = elementType.resolve();
|
||||||
|
ReactiveAdapter adapter = (resolvedType != null
|
||||||
|
? getAdapterRegistry().getAdapter(resolvedType)
|
||||||
|
: null);
|
||||||
|
|
||||||
|
ServerHttpRequest request = exchange.getRequest();
|
||||||
|
ServerHttpResponse response = exchange.getResponse();
|
||||||
|
|
||||||
|
MediaType contentType = request.getHeaders().getContentType();
|
||||||
|
MediaType mediaType = (contentType != null ? contentType
|
||||||
|
: MediaType.APPLICATION_OCTET_STREAM);
|
||||||
|
|
||||||
|
if (logger.isDebugEnabled()) {
|
||||||
|
logger.debug(exchange.getLogPrefix() + (contentType != null
|
||||||
|
? "Content-Type:" + contentType
|
||||||
|
: "No Content-Type, using " + MediaType.APPLICATION_OCTET_STREAM));
|
||||||
|
}
|
||||||
|
boolean isBodyRequired = (adapter != null && !adapter.supportsEmpty());
|
||||||
|
|
||||||
|
MethodParameter bodyParam = new MethodParameter(handlerMethod(handler), 0);
|
||||||
|
for (HttpMessageReader<?> reader : getMessageReaders()) {
|
||||||
|
if (reader.canRead(elementType, mediaType)) {
|
||||||
|
Map<String, Object> readHints = Hints.from(Hints.LOG_PREFIX_HINT,
|
||||||
|
exchange.getLogPrefix());
|
||||||
|
if (adapter != null && adapter.isMultiValue()) {
|
||||||
|
if (logger.isDebugEnabled()) {
|
||||||
|
logger.debug(
|
||||||
|
exchange.getLogPrefix() + "0..N [" + elementType + "]");
|
||||||
|
}
|
||||||
|
Flux<?> flux = reader.read(actualType, elementType, request, response,
|
||||||
|
readHints);
|
||||||
|
flux = flux.onErrorResume(
|
||||||
|
ex -> Flux.error(handleReadError(bodyParam, ex)));
|
||||||
|
if (isBodyRequired) {
|
||||||
|
flux = flux.switchIfEmpty(
|
||||||
|
Flux.error(() -> handleMissingBody(bodyParam)));
|
||||||
|
}
|
||||||
|
return Mono.just(adapter.fromPublisher(flux));
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
// Single-value (with or without reactive type wrapper)
|
||||||
|
if (logger.isDebugEnabled()) {
|
||||||
|
logger.debug(
|
||||||
|
exchange.getLogPrefix() + "0..1 [" + elementType + "]");
|
||||||
|
}
|
||||||
|
Mono<?> mono = reader.readMono(actualType, elementType, request,
|
||||||
|
response, readHints);
|
||||||
|
mono = mono.onErrorResume(
|
||||||
|
ex -> Mono.error(handleReadError(bodyParam, ex)));
|
||||||
|
if (isBodyRequired) {
|
||||||
|
mono = mono.switchIfEmpty(
|
||||||
|
Mono.error(() -> handleMissingBody(bodyParam)));
|
||||||
|
}
|
||||||
|
return (adapter != null ? Mono.just(adapter.fromPublisher(mono))
|
||||||
|
: Mono.from(mono));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return Mono.error(new UnsupportedMediaTypeStatusException(mediaType,
|
||||||
|
Arrays.asList(MediaType.APPLICATION_JSON), elementType));
|
||||||
|
}
|
||||||
|
|
||||||
|
private Method handlerMethod(Object handler) {
|
||||||
|
return ReflectionUtils.findMethod(handler.getClass(), "apply", (Class<?>[]) null);
|
||||||
|
}
|
||||||
|
|
||||||
|
public List<HttpMessageReader<?>> getMessageReaders() {
|
||||||
|
return this.messageReaders;
|
||||||
|
}
|
||||||
|
|
||||||
|
private Throwable handleReadError(MethodParameter parameter, Throwable ex) {
|
||||||
|
return (ex instanceof DecodingException
|
||||||
|
? new ServerWebInputException("Failed to read HTTP message", parameter,
|
||||||
|
ex)
|
||||||
|
: ex);
|
||||||
|
}
|
||||||
|
|
||||||
|
private ServerWebInputException handleMissingBody(MethodParameter param) {
|
||||||
|
return new ServerWebInputException(
|
||||||
|
"Request body is missing: " + param.getExecutable().toGenericString());
|
||||||
|
}
|
||||||
|
|
||||||
|
private ReactiveAdapterRegistry getAdapterRegistry() {
|
||||||
|
return ReactiveAdapterRegistry.getSharedInstance();
|
||||||
|
}
|
||||||
|
|
||||||
private Publisher<?> value(Function<Publisher<?>, Publisher<?>> function,
|
private Publisher<?> value(Function<Publisher<?>, Publisher<?>> function,
|
||||||
Publisher<String> value) {
|
Publisher<String> value) {
|
||||||
Flux<?> input = Flux.from(value).map(body -> converter.convert(function, body));
|
Flux<?> input = Flux.from(value).map(body -> converter.convert(function, body));
|
||||||
|
|||||||
@@ -83,6 +83,13 @@ public class FunctionController {
|
|||||||
return map;
|
return map;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@PostMapping(path = "/**", consumes = MediaType.APPLICATION_OCTET_STREAM_VALUE)
|
||||||
|
@ResponseBody
|
||||||
|
public Mono<ResponseEntity<?>> post(ServerWebExchange request) {
|
||||||
|
FunctionWrapper wrapper = wrapper(request);
|
||||||
|
return processor.post(wrapper, request);
|
||||||
|
}
|
||||||
|
|
||||||
@PostMapping(path = "/**")
|
@PostMapping(path = "/**")
|
||||||
@ResponseBody
|
@ResponseBody
|
||||||
public Mono<ResponseEntity<?>> post(ServerWebExchange request,
|
public Mono<ResponseEntity<?>> post(ServerWebExchange request,
|
||||||
|
|||||||
@@ -106,7 +106,8 @@ class FunctionEndpointInitializer
|
|||||||
context.registerBean(RequestProcessor.class,
|
context.registerBean(RequestProcessor.class,
|
||||||
() -> new RequestProcessor(context.getBean(FunctionInspector.class),
|
() -> new RequestProcessor(context.getBean(FunctionInspector.class),
|
||||||
context.getBeanProvider(JsonMapper.class),
|
context.getBeanProvider(JsonMapper.class),
|
||||||
context.getBean(StringConverter.class)));
|
context.getBean(StringConverter.class),
|
||||||
|
context.getBeanProvider(ServerCodecConfigurer.class)));
|
||||||
context.registerBean(FunctionEndpointFactory.class,
|
context.registerBean(FunctionEndpointFactory.class,
|
||||||
() -> new FunctionEndpointFactory(context.getBean(FunctionCatalog.class),
|
() -> new FunctionEndpointFactory(context.getBean(FunctionCatalog.class),
|
||||||
context.getBean(FunctionInspector.class),
|
context.getBean(FunctionInspector.class),
|
||||||
@@ -226,7 +227,8 @@ class FunctionEndpointFactory {
|
|||||||
logger.info("Found functions: " + names);
|
logger.info("Found functions: " + names);
|
||||||
if (handler != null) {
|
if (handler != null) {
|
||||||
logger.info("Configured function: " + handler);
|
logger.info("Configured function: " + handler);
|
||||||
Assert.isTrue(names.contains(handler), "Cannot locate function: " + handler);
|
Assert.isTrue(names.contains(handler),
|
||||||
|
"Cannot locate function: " + handler);
|
||||||
return catalog.lookup(Function.class, handler);
|
return catalog.lookup(Function.class, handler);
|
||||||
}
|
}
|
||||||
return catalog.lookup(Function.class, names.iterator().next());
|
return catalog.lookup(Function.class, names.iterator().next());
|
||||||
|
|||||||
Reference in New Issue
Block a user