Change RouteReader.getRoutes() from List<Route> to Flux<Route>

This commit is contained in:
Spencer Gibb
2017-01-27 17:06:42 -07:00
parent 8f368db463
commit c0ec9cdd16
6 changed files with 63 additions and 32 deletions

View File

@@ -16,6 +16,8 @@ import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Mono;
/**
* @author Spencer Gibb
*/
@@ -70,22 +72,23 @@ public class GatewayEndpoint {/*extends AbstractEndpoint<Map<String, Object>> {*
}
@GetMapping("/routes")
public List<Route> routes() {
return this.routeReader.getRoutes();
public Mono<List<Route>> routes() {
return this.routeReader.getRoutes().collectList();
}
@GetMapping("/routes/{id}")
public Route route(@PathVariable String id) {
return this.routeReader.getRoutes().stream()
public Mono<Route> route(@PathVariable String id) {
return this.routeReader.getRoutes()
.filter(route -> route.getId().equals(id))
.findFirst().get();
.singleOrEmpty();
}
@GetMapping("/routes/{id}/combinedfilters")
public Map<String, Object> combinedfilters(@PathVariable String id) {
final Optional<Route> route = this.routeReader.getRoutes().stream()
Mono<Route> route = this.routeReader.getRoutes()
.filter(r -> r.getId().equals(id))
.findFirst();
return getNamesToOrders(this.filteringWebHandler.combineFiltersForRoute(route));
.singleOrEmpty();
Optional<Route> optional = Optional.ofNullable(route.block()); //TODO: remove block();
return getNamesToOrders(this.filteringWebHandler.combineFiltersForRoute(optional));
}
}

View File

@@ -0,0 +1,22 @@
package org.springframework.cloud.gateway.api;
import org.springframework.cloud.gateway.config.Route;
import reactor.core.publisher.Flux;
/**
* @author Spencer Gibb
*/
public class CompositeRouteReader implements RouteReader {
private final Flux<RouteReader> delegates;
public CompositeRouteReader(Flux<RouteReader> delegates) {
this.delegates = delegates;
}
@Override
public Flux<Route> getRoutes() {
return this.delegates.flatMap(RouteReader::getRoutes);
}
}

View File

@@ -2,12 +2,12 @@ package org.springframework.cloud.gateway.api;
import org.springframework.cloud.gateway.config.Route;
import java.util.List;
import reactor.core.publisher.Flux;
/**
* @author Spencer Gibb
*/
public interface RouteReader {
List<Route> getRoutes();
Flux<Route> getRoutes();
}

View File

@@ -2,7 +2,7 @@ package org.springframework.cloud.gateway.config;
import org.springframework.cloud.gateway.api.RouteReader;
import java.util.List;
import reactor.core.publisher.Flux;
/**
* @author Spencer Gibb
@@ -16,7 +16,7 @@ public class PropertiesRouteReader implements RouteReader {
}
@Override
public List<Route> getRoutes() {
return this.properties.getRoutes();
public Flux<Route> getRoutes() {
return Flux.fromIterable(this.properties.getRoutes());
}
}

View File

@@ -16,6 +16,7 @@ import org.springframework.web.reactive.handler.AbstractHandlerMapping;
import org.springframework.web.server.ServerWebExchange;
import org.springframework.web.server.WebHandler;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.GATEWAY_HANDLER_MAPPER_ATTR;
@@ -60,7 +61,8 @@ public class RoutePredicateHandlerMapping extends AbstractHandlerMapping {
@Override
protected void initApplicationContext() throws BeansException {
super.initApplicationContext();
registerHandlers(this.routeReader.getRoutes());
Flux<Route> routes = this.routeReader.getRoutes();
registerHandlers(routes.collectList().block()); //TODO: convert rest of class to Reactive
}
protected void registerHandlers(List<Route> routes) {

View File

@@ -44,7 +44,7 @@ public class GatewayIntegrationTests {
private static final String HANDLER_MAPPER_HEADER = "X-Gateway-Handler-Mapper-Class";
private static final String ROUTE_ID_HEADER = "X-Gateway-Route-Id";
public static final Duration DURATION = Duration.ofSeconds(3);
public static final Duration DURATION = Duration.ofSeconds(5);
@LocalServerPort
private int port;
@@ -65,7 +65,7 @@ public class GatewayIntegrationTests {
@Test
public void addRequestHeaderFilterWorks() {
Mono<Map> result = webClient.exchange(
GET("http://localhost:" + port + "/headers")
GET(baseUrl() + "/headers")
.header("Host", "www.addrequestheader.org")
.build()
).then(response -> response.body(toMono(Map.class)));
@@ -93,7 +93,7 @@ public class GatewayIntegrationTests {
private void testRequestParameterFilter(String query) {
Mono<Map> result = webClient.exchange(
GET("http://localhost:" + port + "/get" + query)
GET(baseUrl() + "/get" + query)
.header("Host", "www.addrequestparameter.org")
.build()
).then(response -> response.body(toMono(Map.class)));
@@ -112,7 +112,7 @@ public class GatewayIntegrationTests {
@Test
public void addResponseHeaderFilterWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/headers")
GET(baseUrl() + "/headers")
.header("Host", "www.addresponseheader.org")
.build()
);
@@ -131,7 +131,7 @@ public class GatewayIntegrationTests {
@Test
public void compositeRouteWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/headers?foo=bar&baz")
GET(baseUrl() + "/headers?foo=bar&baz")
.header("Host", "www.foo.org")
.header("X-Request-Id", "123")
.cookie("chocolate", "chip")
@@ -158,7 +158,7 @@ public class GatewayIntegrationTests {
@Test
public void hostRouteWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/get")
GET(baseUrl() + "/get")
.header("Host", "www.example.org")
.build()
);
@@ -181,7 +181,7 @@ public class GatewayIntegrationTests {
@Test
public void hystrixFilterWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/get")
GET(baseUrl() + "/get")
.header("Host", "www.hystrixsuccess.org")
.build()
);
@@ -202,7 +202,7 @@ public class GatewayIntegrationTests {
@Test
public void hystrixFilterTimesout() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/delay/3")
GET(baseUrl() + "/delay/3")
.header("Host", "www.hystrixfailure.org")
.build()
);
@@ -215,7 +215,7 @@ public class GatewayIntegrationTests {
@Test
public void loadBalancerFilterWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/get")
GET(baseUrl() + "/get")
.header("Host", "www.loadbalancerclient.org")
.build()
);
@@ -235,7 +235,7 @@ public class GatewayIntegrationTests {
@Test
public void postWorks() {
ClientRequest<Mono<String>> request = POST("http://localhost:" + port + "/post")
ClientRequest<Mono<String>> request = POST(baseUrl() + "/post")
.header("Host", "www.example.org")
.body(Mono.just("testdata"), String.class);
@@ -251,7 +251,7 @@ public class GatewayIntegrationTests {
@Test
public void redirectToFilterWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port)
GET(baseUrl())
.header("Host", "www.redirectto.org")
.build()
);
@@ -273,7 +273,7 @@ public class GatewayIntegrationTests {
@SuppressWarnings("unchecked")
public void removeRequestHeaderFilterWorks() {
Mono<Map> result = webClient.exchange(
GET("http://localhost:" + port + "/headers")
GET(baseUrl() + "/headers")
.header("Host", "www.removerequestheader.org")
.header("X-Request-Foo", "Bar")
.build()
@@ -293,7 +293,7 @@ public class GatewayIntegrationTests {
@Test
public void removeResponseHeaderFilterWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/headers")
GET(baseUrl() + "/headers")
.header("Host", "www.removereresponseheader.org")
.build()
);
@@ -311,7 +311,7 @@ public class GatewayIntegrationTests {
@Test
public void rewritePathFilterWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/foo/get")
GET(baseUrl() + "/foo/get")
.header("Host", "www.baz.org")
.build()
);
@@ -329,7 +329,7 @@ public class GatewayIntegrationTests {
@Test
public void setPathFilterWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/foo/get")
GET(baseUrl() + "/foo/get")
.header("Host", "www.setpath.org")
.build()
);
@@ -347,7 +347,7 @@ public class GatewayIntegrationTests {
@Test
public void setResponseHeaderFilterWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/headers")
GET(baseUrl() + "/headers")
.header("Host", "www.setreresponseheader.org")
.build()
);
@@ -375,7 +375,7 @@ public class GatewayIntegrationTests {
private void setStatusStringTest(String host, HttpStatus status) {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/headers")
GET(baseUrl() + "/headers")
.header("Host", host)
.build()
);
@@ -393,7 +393,7 @@ public class GatewayIntegrationTests {
@Test
public void urlRouteWorks() {
Mono<ClientResponse> result = webClient.exchange(
GET("http://localhost:" + port + "/get").build()
GET(baseUrl() + "/get").build()
);
StepVerifier.create(result)
@@ -411,6 +411,10 @@ public class GatewayIntegrationTests {
.verify(DURATION);
}
private String baseUrl() {
return "http://localhost:" + port;
}
@EnableAutoConfiguration
@SpringBootConfiguration
public static class TestConfig {