Polish gh-3224

This commit is contained in:
sgibb
2024-03-08 13:40:20 -05:00
parent 9eb453b759
commit e802d2863b

View File

@@ -26,6 +26,7 @@ import org.apache.commons.logging.LogFactory;
import reactor.cache.CacheFlux;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.core.publisher.Signal;
import org.springframework.cloud.gateway.event.RefreshRoutesEvent;
import org.springframework.cloud.gateway.event.RefreshRoutesResultEvent;
@@ -90,20 +91,16 @@ public class CachingRouteLocator
scopedRoutes.subscribe(scopedRoutesList -> {
Flux.concat(Flux.fromIterable(scopedRoutesList), getNonScopedRoutes(event))
.sort(AnnotationAwareOrderComparator.INSTANCE).materialize().collect(Collectors.toList())
.subscribe(signals -> {
cache.put(CACHE_KEY, signals);
applicationEventPublisher.publishEvent(new RefreshRoutesResultEvent(this));
}, this::handleRefreshError);
.subscribe(this::publishRefreshEvent, this::handleRefreshError);
}, this::handleRefreshError);
}
else {
final Mono<List<Route>> allRoutes = fetch().collect(Collectors.toList());
allRoutes.subscribe(list -> Flux.fromIterable(list).materialize().collect(Collectors.toList())
.subscribe(signals -> {
cache.put(CACHE_KEY, signals);
applicationEventPublisher.publishEvent(new RefreshRoutesResultEvent(this));
}, this::handleRefreshError), this::handleRefreshError);
allRoutes.subscribe(
list -> Flux.fromIterable(list).materialize().collect(Collectors.toList())
.subscribe(this::publishRefreshEvent, this::handleRefreshError),
this::handleRefreshError);
}
}
catch (Throwable e) {
@@ -111,6 +108,11 @@ public class CachingRouteLocator
}
}
private void publishRefreshEvent(List<Signal<Route>> signals) {
cache.put(CACHE_KEY, signals);
applicationEventPublisher.publishEvent(new RefreshRoutesResultEvent(this));
}
private Flux<Route> getNonScopedRoutes(RefreshRoutesEvent scopedEvent) {
return this.getRoutes()
.filter(route -> !RouteLocator.matchMetadata(route.getMetadata(), scopedEvent.getMetadata()));