@@ -64,6 +64,7 @@ import org.springframework.web.reactive.function.server.RouterFunction;
|
|||||||
import org.springframework.web.reactive.function.server.RouterFunctions;
|
import org.springframework.web.reactive.function.server.RouterFunctions;
|
||||||
import org.springframework.web.reactive.function.server.ServerRequest;
|
import org.springframework.web.reactive.function.server.ServerRequest;
|
||||||
import org.springframework.web.reactive.function.server.ServerResponse;
|
import org.springframework.web.reactive.function.server.ServerResponse;
|
||||||
|
import org.springframework.web.reactive.function.server.ServerResponse.BodyBuilder;
|
||||||
import org.springframework.web.server.WebExceptionHandler;
|
import org.springframework.web.server.WebExceptionHandler;
|
||||||
import org.springframework.web.server.adapter.HttpWebHandlerAdapter;
|
import org.springframework.web.server.adapter.HttpWebHandlerAdapter;
|
||||||
import org.springframework.web.server.adapter.WebHttpHandlerBuilder;
|
import org.springframework.web.server.adapter.WebHttpHandlerBuilder;
|
||||||
@@ -244,8 +245,13 @@ class FunctionEndpointFactory {
|
|||||||
Mono<ResponseEntity<?>> stream = request.bodyToMono(String.class)
|
Mono<ResponseEntity<?>> stream = request.bodyToMono(String.class)
|
||||||
.flatMap(content -> this.processor.post(wrapper, content, false));
|
.flatMap(content -> this.processor.post(wrapper, content, false));
|
||||||
return stream.flatMap(entity -> {
|
return stream.flatMap(entity -> {
|
||||||
return status(entity.getStatusCode()).headers(headers -> headers.addAll(entity.getHeaders()))
|
BodyBuilder builder = status(entity.getStatusCode()).headers(headers -> headers.addAll(entity.getHeaders()));
|
||||||
.body(entity.hasBody() ? Mono.just((T) entity.getBody()) : Mono.empty(), outputType);
|
if (outputType == null) { // consumer
|
||||||
|
return builder.build();
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
return builder.body(entity != null && entity.hasBody() ? Mono.just((T) entity.getBody()) : Mono.empty(), outputType);
|
||||||
|
}
|
||||||
});
|
});
|
||||||
}).andRoute(GET("/**"), request -> {
|
}).andRoute(GET("/**"), request -> {
|
||||||
FunctionInvocationWrapper funcWrapper = extract(request);
|
FunctionInvocationWrapper funcWrapper = extract(request);
|
||||||
|
|||||||
@@ -17,6 +17,7 @@
|
|||||||
package org.springframework.cloud.function.web.function;
|
package org.springframework.cloud.function.web.function;
|
||||||
|
|
||||||
import java.net.URI;
|
import java.net.URI;
|
||||||
|
import java.util.function.Consumer;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
import java.util.function.Supplier;
|
import java.util.function.Supplier;
|
||||||
|
|
||||||
@@ -67,6 +68,18 @@ public class FunctionEndpointInitializerTests {
|
|||||||
assertThat(response.getStatusCode()).isEqualTo(HttpStatus.NOT_FOUND);
|
assertThat(response.getStatusCode()).isEqualTo(HttpStatus.NOT_FOUND);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void testConsumerMapping() throws Exception {
|
||||||
|
FunctionalSpringApplication.run(ConsumerConfiguration.class);
|
||||||
|
TestRestTemplate testRestTemplate = new TestRestTemplate();
|
||||||
|
String port = System.getProperty("server.port");
|
||||||
|
Thread.sleep(200);
|
||||||
|
ResponseEntity<String> response = testRestTemplate
|
||||||
|
.postForEntity(new URI("http://localhost:" + port + "/uppercase"), "stressed", String.class);
|
||||||
|
assertThat(response.getBody()).isNull();
|
||||||
|
assertThat(response.getStatusCode()).isEqualTo(HttpStatus.ACCEPTED);
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void testSingleFunctionMapping() throws Exception {
|
public void testSingleFunctionMapping() throws Exception {
|
||||||
FunctionalSpringApplication.run(ApplicationConfiguration.class);
|
FunctionalSpringApplication.run(ApplicationConfiguration.class);
|
||||||
@@ -114,6 +127,23 @@ public class FunctionEndpointInitializerTests {
|
|||||||
assertThat(response.getBody()).isEqualTo("Jim Lahey");
|
assertThat(response.getBody()).isEqualTo("Jim Lahey");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@SpringBootConfiguration
|
||||||
|
protected static class ConsumerConfiguration
|
||||||
|
implements ApplicationContextInitializer<GenericApplicationContext> {
|
||||||
|
|
||||||
|
public Consumer<String> consume() {
|
||||||
|
return v -> System.out.println(v);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void initialize(GenericApplicationContext applicationContext) {
|
||||||
|
applicationContext.registerBean("consume", FunctionRegistration.class,
|
||||||
|
() -> new FunctionRegistration<>(consume())
|
||||||
|
.type(FunctionType.consumer(String.class)));
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
@SpringBootConfiguration
|
@SpringBootConfiguration
|
||||||
protected static class ApplicationConfiguration
|
protected static class ApplicationConfiguration
|
||||||
|
|||||||
Reference in New Issue
Block a user