Add support for checkstyle rules and fix checkstyle failures. Fixes #14
This commit is contained in:
18
.editorconfig
Normal file
18
.editorconfig
Normal file
@@ -0,0 +1,18 @@
|
||||
# EditorConfig is awesome: http://EditorConfig.org
|
||||
|
||||
# top-most EditorConfig file
|
||||
root = true
|
||||
|
||||
[*]
|
||||
indent_style = tab
|
||||
indent_size = 4
|
||||
end_of_line = lf
|
||||
insert_final_newline = true
|
||||
|
||||
[*.yml]
|
||||
indent_style = space
|
||||
indent_size = 2
|
||||
|
||||
[*.yaml]
|
||||
indent_style = space
|
||||
indent_size = 2
|
||||
14
pom.xml
14
pom.xml
@@ -33,6 +33,10 @@
|
||||
<sonar.language>java</sonar.language>
|
||||
<spring-cloud-commons.version>2.2.0.BUILD-SNAPSHOT</spring-cloud-commons.version>
|
||||
<spring-cloud-netflix.version>2.2.0.BUILD-SNAPSHOT</spring-cloud-netflix.version>
|
||||
<maven-checkstyle-plugin.includeTestSourceDirectory>true
|
||||
</maven-checkstyle-plugin.includeTestSourceDirectory>
|
||||
<maven-checkstyle-plugin.failOnViolation>true
|
||||
</maven-checkstyle-plugin.failOnViolation>
|
||||
</properties>
|
||||
<dependencyManagement>
|
||||
<dependencies>
|
||||
@@ -101,6 +105,14 @@
|
||||
<target>1.8</target>
|
||||
</configuration>
|
||||
</plugin>
|
||||
<plugin>
|
||||
<groupId>io.spring.javaformat</groupId>
|
||||
<artifactId>spring-javaformat-maven-plugin</artifactId>
|
||||
</plugin>
|
||||
<plugin>
|
||||
<groupId>org.apache.maven.plugins</groupId>
|
||||
<artifactId>maven-checkstyle-plugin</artifactId>
|
||||
</plugin>
|
||||
</plugins>
|
||||
</build>
|
||||
<profiles>
|
||||
@@ -209,4 +221,4 @@
|
||||
</profile>
|
||||
</profiles>
|
||||
|
||||
</project>
|
||||
</project>
|
||||
|
||||
@@ -19,23 +19,24 @@ package org.springframework.cloud.circuitbreaker.commons;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
/**
|
||||
* Base class for factories which produce circuit breakers.
|
||||
*
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public abstract class AbstractCircuitBreakerFactory<CONF, CONFB extends ConfigBuilder<CONF>> {
|
||||
|
||||
private final ConcurrentHashMap<String, CONF> configurations = new ConcurrentHashMap<>();
|
||||
|
||||
/**
|
||||
* Adds configurations for circuit breakers
|
||||
* Adds configurations for circuit breakers.
|
||||
* @param ids The id of the circuit breaker
|
||||
* @param consumer A configuration builder consumer,
|
||||
* allows consumers to customize the builder before the configuration is built
|
||||
* @param consumer A configuration builder consumer, allows consumers to customize the
|
||||
* builder before the configuration is built
|
||||
*/
|
||||
public void configure(Consumer<CONFB> consumer, String ... ids) {
|
||||
for(String id : ids) {
|
||||
public void configure(Consumer<CONFB> consumer, String... ids) {
|
||||
for (String id : ids) {
|
||||
CONFB builder = configBuilder(id);
|
||||
consumer.accept(builder);
|
||||
CONF conf = builder.build();
|
||||
@@ -44,7 +45,7 @@ public abstract class AbstractCircuitBreakerFactory<CONF, CONFB extends ConfigBu
|
||||
}
|
||||
|
||||
/**
|
||||
* Gets the configurations for the circuit breakers
|
||||
* Gets the configurations for the circuit breakers.
|
||||
* @return The configurations
|
||||
*/
|
||||
protected ConcurrentHashMap<String, CONF> getConfigurations() {
|
||||
@@ -52,15 +53,16 @@ public abstract class AbstractCircuitBreakerFactory<CONF, CONFB extends ConfigBu
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a configuration builder for the given id
|
||||
* Creates a configuration builder for the given id.
|
||||
* @param id The id of the circuit breaker
|
||||
* @return The configuration builder
|
||||
*/
|
||||
protected abstract CONFB configBuilder(String id);
|
||||
|
||||
/**
|
||||
* Sets the default configuration for circuit breakers
|
||||
* Sets the default configuration for circuit breakers.
|
||||
* @param defaultConfiguration A function that returns the default configuration
|
||||
*/
|
||||
public abstract void configureDefault(Function<String, CONF> defaultConfiguration);
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2018 the original author or authors.
|
||||
* Copyright 2013-2019 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.
|
||||
@@ -13,6 +13,7 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.commons;
|
||||
|
||||
import java.util.function.Function;
|
||||
@@ -20,6 +21,7 @@ import java.util.function.Supplier;
|
||||
|
||||
/**
|
||||
* Spring Cloud circuit breaker.
|
||||
*
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public interface CircuitBreaker {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2018 the original author or authors.
|
||||
* Copyright 2013-2019 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.
|
||||
@@ -13,6 +13,7 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.commons;
|
||||
|
||||
/**
|
||||
@@ -20,7 +21,9 @@ package org.springframework.cloud.circuitbreaker.commons;
|
||||
*
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public abstract class CircuitBreakerFactory<CONF, CONFB extends ConfigBuilder<CONF>> extends AbstractCircuitBreakerFactory<CONF, CONFB> {
|
||||
public abstract class CircuitBreakerFactory<CONF, CONFB extends ConfigBuilder<CONF>>
|
||||
extends AbstractCircuitBreakerFactory<CONF, CONFB> {
|
||||
|
||||
public abstract CircuitBreaker create(String id);
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2018 the original author or authors.
|
||||
* Copyright 2013-2019 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.
|
||||
@@ -13,14 +13,16 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.commons;
|
||||
|
||||
/**
|
||||
* A builder for circuit breaker configurations
|
||||
* A builder for circuit breaker configurations.
|
||||
*
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public interface ConfigBuilder<CONF> {
|
||||
CONF build();
|
||||
}
|
||||
|
||||
CONF build();
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2018 the original author or authors.
|
||||
* Copyright 2013-2019 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.
|
||||
@@ -13,13 +13,16 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.commons;
|
||||
|
||||
/**
|
||||
* Customizes the parameterized class
|
||||
* Customizes the parameterized class.
|
||||
*
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public interface Customizer<TOCUSTOMIZE> {
|
||||
|
||||
void customize(TOCUSTOMIZE tocustomize);
|
||||
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.commons;
|
||||
|
||||
/**
|
||||
@@ -25,4 +26,5 @@ public class NoFallbackAvailableException extends RuntimeException {
|
||||
public NoFallbackAvailableException(String message, Throwable cause) {
|
||||
super(message, cause);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2018 the original author or authors.
|
||||
* Copyright 2013-2019 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.
|
||||
@@ -13,15 +13,16 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.commons;
|
||||
|
||||
import java.util.function.Function;
|
||||
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import java.util.function.Function;
|
||||
|
||||
/**
|
||||
* Spring Cloud reactive circuit breaker API
|
||||
* Spring Cloud reactive circuit breaker API.
|
||||
*
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
|
||||
@@ -13,14 +13,17 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.commons;
|
||||
|
||||
/**
|
||||
* Creates reactive circuit breakers
|
||||
* Creates reactive circuit breakers.
|
||||
*
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public abstract class ReactiveCircuitBreakerFactory<CONF, CONFB extends ConfigBuilder<CONF>> extends AbstractCircuitBreakerFactory<CONF, CONFB> {
|
||||
public abstract class ReactiveCircuitBreakerFactory<CONF, CONFB extends ConfigBuilder<CONF>>
|
||||
extends AbstractCircuitBreakerFactory<CONF, CONFB> {
|
||||
|
||||
public abstract ReactiveCircuitBreaker create(String id);
|
||||
|
||||
}
|
||||
|
||||
@@ -13,28 +13,32 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.hystrix;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.ConfigBuilder;
|
||||
import org.springframework.util.StringUtils;
|
||||
import com.netflix.hystrix.HystrixCommandGroupKey;
|
||||
import com.netflix.hystrix.HystrixCommandKey;
|
||||
import com.netflix.hystrix.HystrixCommandProperties;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.ConfigBuilder;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public abstract class AbstractHystrixConfigBuilder<CONFB> implements ConfigBuilder<CONFB> {
|
||||
public abstract class AbstractHystrixConfigBuilder<CONFB>
|
||||
implements ConfigBuilder<CONFB> {
|
||||
|
||||
private final String commandName;
|
||||
|
||||
protected String groupName;
|
||||
|
||||
protected HystrixCommandProperties.Setter commandProperties;
|
||||
|
||||
public AbstractHystrixConfigBuilder(String id) {
|
||||
this.commandName = id;
|
||||
}
|
||||
|
||||
|
||||
public AbstractHystrixConfigBuilder groupName(String groupName) {
|
||||
this.groupName = groupName;
|
||||
return this;
|
||||
@@ -50,7 +54,8 @@ public abstract class AbstractHystrixConfigBuilder<CONFB> implements ConfigBuild
|
||||
String groupNameToUse;
|
||||
if (StringUtils.hasText(this.groupName)) {
|
||||
groupNameToUse = this.groupName;
|
||||
} else {
|
||||
}
|
||||
else {
|
||||
groupNameToUse = commandName + "group";
|
||||
}
|
||||
return HystrixCommandGroupKey.Factory.asKey(groupNameToUse);
|
||||
@@ -61,8 +66,7 @@ public abstract class AbstractHystrixConfigBuilder<CONFB> implements ConfigBuild
|
||||
}
|
||||
|
||||
protected HystrixCommandProperties.Setter getCommandPropertiesSetter() {
|
||||
return this.commandProperties != null
|
||||
? this.commandProperties
|
||||
return this.commandProperties != null ? this.commandProperties
|
||||
: HystrixCommandProperties.Setter();
|
||||
}
|
||||
|
||||
|
||||
@@ -19,11 +19,12 @@ package org.springframework.cloud.circuitbreaker.hystrix;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.CircuitBreaker;
|
||||
import com.netflix.hystrix.HystrixCommand;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.CircuitBreaker;
|
||||
|
||||
/**
|
||||
* Hystrix implementation of {@link CircuitBreaker}
|
||||
* Hystrix implementation of {@link CircuitBreaker}.
|
||||
*
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
@@ -51,4 +52,5 @@ public class HystrixCircuitBreaker implements CircuitBreaker {
|
||||
};
|
||||
return command.execute();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2018 the original author or authors.
|
||||
* Copyright 2013-2019 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.
|
||||
@@ -12,14 +12,17 @@
|
||||
* 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.circuitbreaker.hystrix;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
import javax.annotation.PostConstruct;
|
||||
|
||||
import com.netflix.hystrix.Hystrix;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
|
||||
@@ -28,7 +31,6 @@ import org.springframework.cloud.circuitbreaker.commons.Customizer;
|
||||
import org.springframework.cloud.circuitbreaker.commons.ReactiveCircuitBreakerFactory;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import com.netflix.hystrix.Hystrix;
|
||||
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
@@ -45,7 +47,8 @@ public class HystrixCircuitBreakerAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean(ReactiveCircuitBreakerFactory.class)
|
||||
@ConditionalOnClass(name = {"reactor.core.publisher.Mono", "reactor.core.publisher.Flux"})
|
||||
@ConditionalOnClass(name = { "reactor.core.publisher.Mono",
|
||||
"reactor.core.publisher.Flux" })
|
||||
public ReactiveHystrixCircuitBreakerFactory reactiveHystrixCircuitBreakerFactory() {
|
||||
return new ReactiveHystrixCircuitBreakerFactory();
|
||||
}
|
||||
@@ -63,10 +66,12 @@ public class HystrixCircuitBreakerAutoConfiguration {
|
||||
public void init() {
|
||||
customizers.forEach(customizer -> customizer.customize(factory));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@ConditionalOnClass(name = {"reactor.core.publisher.Mono", "reactor.core.publisher.Flux"})
|
||||
@ConditionalOnClass(name = { "reactor.core.publisher.Mono",
|
||||
"reactor.core.publisher.Flux" })
|
||||
protected static class ReactiveHystrixCircuitBreakerCustomizerConfiguration {
|
||||
|
||||
@Autowired(required = false)
|
||||
@@ -79,5 +84,7 @@ public class HystrixCircuitBreakerAutoConfiguration {
|
||||
public void init() {
|
||||
customizers.forEach(customizer -> customizer.customize(factory));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -17,23 +17,26 @@
|
||||
package org.springframework.cloud.circuitbreaker.hystrix;
|
||||
|
||||
import java.util.function.Function;
|
||||
import org.springframework.cloud.circuitbreaker.commons.CircuitBreakerFactory;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
import com.netflix.hystrix.HystrixCommand;
|
||||
import com.netflix.hystrix.HystrixCommandGroupKey;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.CircuitBreakerFactory;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
* Builds Hystrix circuit breakers.
|
||||
*
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public class HystrixCircuitBreakerFactory extends CircuitBreakerFactory<HystrixCommand.Setter, HystrixCircuitBreakerFactory.HystrixConfigBuilder> {
|
||||
public class HystrixCircuitBreakerFactory extends
|
||||
CircuitBreakerFactory<HystrixCommand.Setter, HystrixCircuitBreakerFactory.HystrixConfigBuilder> {
|
||||
|
||||
private Function<String, HystrixCommand.Setter> defaultConfiguration = id ->
|
||||
HystrixCommand.Setter.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id));
|
||||
private Function<String, HystrixCommand.Setter> defaultConfiguration = id -> HystrixCommand.Setter
|
||||
.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id));
|
||||
|
||||
|
||||
public void configureDefault(Function<String, HystrixCommand.Setter> defaultConfiguration) {
|
||||
public void configureDefault(
|
||||
Function<String, HystrixCommand.Setter> defaultConfiguration) {
|
||||
this.defaultConfiguration = defaultConfiguration;
|
||||
}
|
||||
|
||||
@@ -43,11 +46,13 @@ public class HystrixCircuitBreakerFactory extends CircuitBreakerFactory<HystrixC
|
||||
|
||||
public HystrixCircuitBreaker create(String id) {
|
||||
Assert.hasText(id, "A CircuitBreaker must have an id.");
|
||||
HystrixCommand.Setter setter = getConfigurations().computeIfAbsent(id, defaultConfiguration);
|
||||
HystrixCommand.Setter setter = getConfigurations().computeIfAbsent(id,
|
||||
defaultConfiguration);
|
||||
return new HystrixCircuitBreaker(setter);
|
||||
}
|
||||
|
||||
public static class HystrixConfigBuilder extends AbstractHystrixConfigBuilder<HystrixCommand.Setter> {
|
||||
public static class HystrixConfigBuilder
|
||||
extends AbstractHystrixConfigBuilder<HystrixCommand.Setter> {
|
||||
|
||||
public HystrixConfigBuilder(String id) {
|
||||
super(id);
|
||||
@@ -55,9 +60,11 @@ public class HystrixCircuitBreakerFactory extends CircuitBreakerFactory<HystrixC
|
||||
|
||||
@Override
|
||||
public HystrixCommand.Setter build() {
|
||||
return HystrixCommand.Setter.withGroupKey(getGroupKey()).andCommandKey(getCommandKey())
|
||||
return HystrixCommand.Setter.withGroupKey(getGroupKey())
|
||||
.andCommandKey(getCommandKey())
|
||||
.andCommandPropertiesDefaults(getCommandPropertiesSetter());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -16,16 +16,17 @@
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.hystrix;
|
||||
|
||||
import java.util.function.Function;
|
||||
|
||||
import com.netflix.hystrix.HystrixObservableCommand;
|
||||
import org.reactivestreams.Publisher;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
import rx.Observable;
|
||||
import rx.RxReactiveStreams;
|
||||
import rx.Subscription;
|
||||
|
||||
import java.util.function.Function;
|
||||
import org.reactivestreams.Publisher;
|
||||
import org.springframework.cloud.circuitbreaker.commons.ReactiveCircuitBreaker;
|
||||
import com.netflix.hystrix.HystrixObservableCommand;
|
||||
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
@@ -43,7 +44,8 @@ public class ReactiveHystrixCircuitBreaker implements ReactiveCircuitBreaker {
|
||||
HystrixObservableCommand<T> command = createCommand(toRun, fallback);
|
||||
|
||||
return Mono.create(s -> {
|
||||
Subscription sub = command.toObservable().subscribe(s::success, s::error, s::success);
|
||||
Subscription sub = command.toObservable().subscribe(s::success, s::error,
|
||||
s::success);
|
||||
s.onCancel(sub::unsubscribe);
|
||||
});
|
||||
}
|
||||
@@ -53,12 +55,14 @@ public class ReactiveHystrixCircuitBreaker implements ReactiveCircuitBreaker {
|
||||
HystrixObservableCommand<T> command = createCommand(toRun, fallback);
|
||||
|
||||
return Flux.create(s -> {
|
||||
Subscription sub = command.toObservable().subscribe(s::next, s::error, s::complete);
|
||||
Subscription sub = command.toObservable().subscribe(s::next, s::error,
|
||||
s::complete);
|
||||
s.onCancel(sub::unsubscribe);
|
||||
});
|
||||
}
|
||||
|
||||
private <T> HystrixObservableCommand<T> createCommand(Publisher<T> toRun, Function fallback) {
|
||||
private <T> HystrixObservableCommand<T> createCommand(Publisher<T> toRun,
|
||||
Function fallback) {
|
||||
HystrixObservableCommand<T> command = new HystrixObservableCommand<T>(setter) {
|
||||
@Override
|
||||
protected Observable<T> construct() {
|
||||
@@ -67,12 +71,14 @@ public class ReactiveHystrixCircuitBreaker implements ReactiveCircuitBreaker {
|
||||
|
||||
@Override
|
||||
protected Observable<T> resumeWithFallback() {
|
||||
if(fallback == null) {
|
||||
if (fallback == null) {
|
||||
super.resumeWithFallback();
|
||||
}
|
||||
return RxReactiveStreams.toObservable((Publisher)fallback.apply(this.getExecutionException()));
|
||||
return RxReactiveStreams.toObservable(
|
||||
(Publisher) fallback.apply(this.getExecutionException()));
|
||||
}
|
||||
};
|
||||
return command;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -17,20 +17,22 @@
|
||||
package org.springframework.cloud.circuitbreaker.hystrix;
|
||||
|
||||
import java.util.function.Function;
|
||||
|
||||
import com.netflix.hystrix.HystrixCommandGroupKey;
|
||||
import com.netflix.hystrix.HystrixObservableCommand;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.ReactiveCircuitBreaker;
|
||||
import org.springframework.cloud.circuitbreaker.commons.ReactiveCircuitBreakerFactory;
|
||||
import org.springframework.util.Assert;
|
||||
import com.netflix.hystrix.HystrixCommandGroupKey;
|
||||
import com.netflix.hystrix.HystrixObservableCommand;
|
||||
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public class ReactiveHystrixCircuitBreakerFactory extends ReactiveCircuitBreakerFactory<HystrixObservableCommand.Setter,
|
||||
ReactiveHystrixCircuitBreakerFactory.ReactiveHystrixConfigBuilder> {
|
||||
public class ReactiveHystrixCircuitBreakerFactory extends
|
||||
ReactiveCircuitBreakerFactory<HystrixObservableCommand.Setter, ReactiveHystrixCircuitBreakerFactory.ReactiveHystrixConfigBuilder> {
|
||||
|
||||
private Function<String, HystrixObservableCommand.Setter> defaultConfiguration = id ->
|
||||
HystrixObservableCommand.Setter.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id));
|
||||
private Function<String, HystrixObservableCommand.Setter> defaultConfiguration = id -> HystrixObservableCommand.Setter
|
||||
.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id));
|
||||
|
||||
@Override
|
||||
protected ReactiveHystrixConfigBuilder configBuilder(String id) {
|
||||
@@ -38,18 +40,21 @@ public class ReactiveHystrixCircuitBreakerFactory extends ReactiveCircuitBreaker
|
||||
}
|
||||
|
||||
@Override
|
||||
public void configureDefault(Function<String, HystrixObservableCommand.Setter> defaultConfiguration) {
|
||||
public void configureDefault(
|
||||
Function<String, HystrixObservableCommand.Setter> defaultConfiguration) {
|
||||
this.defaultConfiguration = defaultConfiguration;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ReactiveCircuitBreaker create(String id) {
|
||||
Assert.hasText(id, "A CircuitBreaker must have an id.");
|
||||
HystrixObservableCommand.Setter setter = getConfigurations().computeIfAbsent(id, defaultConfiguration);
|
||||
HystrixObservableCommand.Setter setter = getConfigurations().computeIfAbsent(id,
|
||||
defaultConfiguration);
|
||||
return new ReactiveHystrixCircuitBreaker(setter);
|
||||
}
|
||||
|
||||
public static class ReactiveHystrixConfigBuilder extends AbstractHystrixConfigBuilder<HystrixObservableCommand.Setter> {
|
||||
public static class ReactiveHystrixConfigBuilder
|
||||
extends AbstractHystrixConfigBuilder<HystrixObservableCommand.Setter> {
|
||||
|
||||
public ReactiveHystrixConfigBuilder(String id) {
|
||||
super(id);
|
||||
@@ -57,8 +62,11 @@ public class ReactiveHystrixCircuitBreakerFactory extends ReactiveCircuitBreaker
|
||||
|
||||
@Override
|
||||
public HystrixObservableCommand.Setter build() {
|
||||
return HystrixObservableCommand.Setter.withGroupKey(getGroupKey()).andCommandKey(getCommandKey())
|
||||
return HystrixObservableCommand.Setter.withGroupKey(getGroupKey())
|
||||
.andCommandKey(getCommandKey())
|
||||
.andCommandPropertiesDefaults(getCommandPropertiesSetter());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -16,8 +16,12 @@
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.hystrix;
|
||||
|
||||
import com.netflix.hystrix.HystrixCommand;
|
||||
import com.netflix.hystrix.HystrixCommandGroupKey;
|
||||
import com.netflix.hystrix.HystrixCommandProperties;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
@@ -32,11 +36,8 @@ import org.springframework.test.context.junit4.SpringRunner;
|
||||
import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.RequestMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
import com.netflix.hystrix.HystrixCommand;
|
||||
import com.netflix.hystrix.HystrixCommandGroupKey;
|
||||
import com.netflix.hystrix.HystrixCommandProperties;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.springframework.boot.test.context.SpringBootTest.WebEnvironment.RANDOM_PORT;
|
||||
|
||||
/**
|
||||
@@ -47,10 +48,24 @@ import static org.springframework.boot.test.context.SpringBootTest.WebEnvironmen
|
||||
@DirtiesContext
|
||||
public class HystrixCircuitBreakerIntegrationTest {
|
||||
|
||||
@Autowired
|
||||
Application.DemoControllerService service;
|
||||
|
||||
@Test
|
||||
public void testSlow() {
|
||||
assertThat(service.slow()).isEqualTo("fallback");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testNormal() {
|
||||
assertThat(service.normal()).isEqualTo("normal");
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableAutoConfiguration
|
||||
@RestController
|
||||
protected static class Application {
|
||||
|
||||
@RequestMapping("/slow")
|
||||
public String slow() throws InterruptedException {
|
||||
Thread.sleep(3000);
|
||||
@@ -62,50 +77,49 @@ public class HystrixCircuitBreakerIntegrationTest {
|
||||
return "normal";
|
||||
}
|
||||
|
||||
|
||||
@Bean
|
||||
public Customizer<HystrixCircuitBreakerFactory> customizer() {
|
||||
return factory -> factory.configure(builder -> builder.commandProperties(
|
||||
HystrixCommandProperties.Setter().withExecutionTimeoutInMilliseconds(2000)), "slow");
|
||||
return factory -> factory
|
||||
.configure(
|
||||
builder -> builder.commandProperties(HystrixCommandProperties
|
||||
.Setter().withExecutionTimeoutInMilliseconds(2000)),
|
||||
"slow");
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Customizer<HystrixCircuitBreakerFactory> defaultConfig() {
|
||||
return factory -> factory.configureDefault(id -> HystrixCommand.Setter
|
||||
.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id))
|
||||
.andCommandPropertiesDefaults(HystrixCommandProperties.Setter().withExecutionTimeoutInMilliseconds(4000)));
|
||||
.andCommandPropertiesDefaults(HystrixCommandProperties.Setter()
|
||||
.withExecutionTimeoutInMilliseconds(4000)));
|
||||
}
|
||||
|
||||
@Service
|
||||
public static class DemoControllerService {
|
||||
|
||||
private TestRestTemplate rest;
|
||||
|
||||
private CircuitBreakerFactory cbFactory;
|
||||
|
||||
public DemoControllerService(TestRestTemplate rest, CircuitBreakerFactory cbBuilder) {
|
||||
DemoControllerService(TestRestTemplate rest,
|
||||
CircuitBreakerFactory cbBuilder) {
|
||||
this.rest = rest;
|
||||
this.cbFactory = cbBuilder;
|
||||
}
|
||||
|
||||
public String slow() {
|
||||
return cbFactory.create("slow").run(() -> rest.getForObject("/slow", String.class), t -> "fallback");
|
||||
return cbFactory.create("slow").run(
|
||||
() -> rest.getForObject("/slow", String.class), t -> "fallback");
|
||||
}
|
||||
|
||||
public String normal() {
|
||||
return cbFactory.create("normal").run(() -> rest.getForObject("/normal", String.class), t -> "fallback");
|
||||
return cbFactory.create("normal").run(
|
||||
() -> rest.getForObject("/normal", String.class),
|
||||
t -> "fallback");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Autowired
|
||||
Application.DemoControllerService service;
|
||||
|
||||
@Test
|
||||
public void testSlow() {
|
||||
assertEquals("fallback", service.slow());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testNormal() {
|
||||
assertEquals("normal", service.normal());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,10 +17,10 @@
|
||||
package org.springframework.cloud.circuitbreaker.hystrix;
|
||||
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.CircuitBreaker;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
@@ -31,14 +31,15 @@ public class HystrixCircuitBreakerTest {
|
||||
public void run() {
|
||||
CircuitBreaker cb = new HystrixCircuitBreakerFactory().create("foo");
|
||||
String s = cb.run(() -> "foobar", t -> "fallback");
|
||||
assertEquals("foobar", cb.run(() -> "foobar", t -> "fallback"));
|
||||
assertThat(cb.run(() -> "foobar", t -> "fallback")).isEqualTo("foobar");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void fallback() {
|
||||
CircuitBreaker cb = new HystrixCircuitBreakerFactory().create("foo");
|
||||
assertEquals("fallback", cb.run(() -> {
|
||||
assertThat((String) cb.run(() -> {
|
||||
throw new RuntimeException("Boom");
|
||||
}, t -> "fallback"));
|
||||
}, t -> "fallback")).isEqualTo("fallback");
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -16,12 +16,16 @@
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.hystrix;
|
||||
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
import com.netflix.hystrix.HystrixCommandGroupKey;
|
||||
import com.netflix.hystrix.HystrixCommandProperties;
|
||||
import com.netflix.hystrix.HystrixObservableCommand;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
@@ -37,11 +41,8 @@ import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.RequestMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
import org.springframework.web.reactive.function.client.WebClient;
|
||||
import com.netflix.hystrix.HystrixCommandGroupKey;
|
||||
import com.netflix.hystrix.HystrixCommandProperties;
|
||||
import com.netflix.hystrix.HystrixObservableCommand;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.springframework.boot.test.context.SpringBootTest.WebEnvironment.RANDOM_PORT;
|
||||
|
||||
/**
|
||||
@@ -55,6 +56,23 @@ public class ReactiveHystrixCircuitBreakerIntegrationTest {
|
||||
@LocalServerPort
|
||||
int port = 0;
|
||||
|
||||
@Autowired
|
||||
ReactiveHystrixCircuitBreakerIntegrationTest.Application.DemoControllerService service;
|
||||
|
||||
@Before
|
||||
public void setup() {
|
||||
service.setPort(port);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSlow() {
|
||||
assertThat(service.slow().block()).isEqualTo("fallback");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testNormal() {
|
||||
assertThat(service.normal().block()).isEqualTo("normal");
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableAutoConfiguration
|
||||
@@ -73,65 +91,59 @@ public class ReactiveHystrixCircuitBreakerIntegrationTest {
|
||||
|
||||
@Bean
|
||||
public Customizer<ReactiveHystrixCircuitBreakerFactory> customizer() {
|
||||
return factory -> factory.configure(builder -> builder.commandProperties(
|
||||
HystrixCommandProperties.Setter().withExecutionTimeoutInMilliseconds(2000)), "slow");
|
||||
return factory -> factory
|
||||
.configure(
|
||||
builder -> builder.commandProperties(HystrixCommandProperties
|
||||
.Setter().withExecutionTimeoutInMilliseconds(2000)),
|
||||
"slow");
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Customizer<ReactiveHystrixCircuitBreakerFactory> defaultConfig() {
|
||||
return factory -> factory.configureDefault(id -> HystrixObservableCommand.Setter
|
||||
.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id))
|
||||
.andCommandPropertiesDefaults(HystrixCommandProperties.Setter()
|
||||
.withExecutionTimeoutInMilliseconds(4000)));
|
||||
return factory -> factory
|
||||
.configureDefault(id -> HystrixObservableCommand.Setter
|
||||
.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id))
|
||||
.andCommandPropertiesDefaults(HystrixCommandProperties
|
||||
.Setter().withExecutionTimeoutInMilliseconds(4000)));
|
||||
}
|
||||
|
||||
@Service
|
||||
public static class DemoControllerService {
|
||||
|
||||
private int port = 0;
|
||||
|
||||
private ReactiveCircuitBreakerFactory cbFactory;
|
||||
|
||||
|
||||
public DemoControllerService(ReactiveCircuitBreakerFactory cbBuilder) {
|
||||
DemoControllerService(ReactiveCircuitBreakerFactory cbBuilder) {
|
||||
this.cbFactory = cbBuilder;
|
||||
}
|
||||
|
||||
public Mono<String> slow() {
|
||||
return cbFactory.create("slow").run(WebClient.builder().baseUrl("http://localhost:" + port).build()
|
||||
.get().uri("/slow").retrieve().bodyToMono(String.class), t -> {
|
||||
t.printStackTrace();
|
||||
return Mono.just("fallback");
|
||||
});
|
||||
return cbFactory.create("slow").run(
|
||||
WebClient.builder().baseUrl("http://localhost:" + port).build()
|
||||
.get().uri("/slow").retrieve().bodyToMono(String.class),
|
||||
t -> {
|
||||
t.printStackTrace();
|
||||
return Mono.just("fallback");
|
||||
});
|
||||
}
|
||||
|
||||
public Mono<String> normal() {
|
||||
return cbFactory.create("normal").run(WebClient.builder().baseUrl("http://localhost:" + port).build()
|
||||
.get().uri("/normal").retrieve().bodyToMono(String.class), t -> {
|
||||
t.printStackTrace();
|
||||
return Mono.just("fallback");
|
||||
});
|
||||
return cbFactory.create("normal").run(
|
||||
WebClient.builder().baseUrl("http://localhost:" + port).build()
|
||||
.get().uri("/normal").retrieve().bodyToMono(String.class),
|
||||
t -> {
|
||||
t.printStackTrace();
|
||||
return Mono.just("fallback");
|
||||
});
|
||||
}
|
||||
|
||||
public void setPort(int port) {
|
||||
this.port = port;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Autowired
|
||||
ReactiveHystrixCircuitBreakerIntegrationTest.Application.DemoControllerService service;
|
||||
|
||||
@Before
|
||||
public void setup() {
|
||||
service.setPort(port);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSlow() {
|
||||
assertEquals("fallback", service.slow().block());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testNormal() {
|
||||
assertEquals("normal", service.normal().block());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -16,14 +16,14 @@
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.hystrix;
|
||||
|
||||
import org.assertj.core.util.Arrays;
|
||||
import org.junit.Test;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import org.assertj.core.util.Arrays;
|
||||
import org.junit.Test;
|
||||
import org.springframework.cloud.circuitbreaker.commons.ReactiveCircuitBreaker;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
@@ -32,27 +32,39 @@ public class ReactiveHystrixCircuitBreakerTest {
|
||||
|
||||
@Test
|
||||
public void monoRun() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveHystrixCircuitBreakerFactory().create("foo");
|
||||
ReactiveCircuitBreaker cb = new ReactiveHystrixCircuitBreakerFactory()
|
||||
.create("foo");
|
||||
Mono<String> s = cb.run(Mono.just("foobar"), t -> Mono.just("fallback"));
|
||||
assertEquals("foobar", s.block());
|
||||
assertThat(s.block()).isEqualTo("foobar");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void monoFallback() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveHystrixCircuitBreakerFactory().create("foo");
|
||||
assertEquals("fallback", cb.run(Mono.error(new RuntimeException("boom")), t -> Mono.just("fallback")).block());
|
||||
ReactiveCircuitBreaker cb = new ReactiveHystrixCircuitBreakerFactory()
|
||||
.create("foo");
|
||||
assertThat(cb
|
||||
.run(Mono.error(new RuntimeException("boom")), t -> Mono.just("fallback"))
|
||||
.block()).isEqualTo("fallback");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void fluxRun() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveHystrixCircuitBreakerFactory().create("foo");
|
||||
Flux<String> s = cb.run(Flux.just("foobar","hello world"), t -> Flux.just("fallback"));
|
||||
assertEquals(Arrays.asList(new String[]{"foobar", "hello world"}), s.collectList().block());
|
||||
ReactiveCircuitBreaker cb = new ReactiveHystrixCircuitBreakerFactory()
|
||||
.create("foo");
|
||||
Flux<String> s = cb.run(Flux.just("foobar", "hello world"),
|
||||
t -> Flux.just("fallback"));
|
||||
assertThat(s.collectList().block())
|
||||
.isEqualTo(Arrays.asList(new String[] { "foobar", "hello world" }));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void fluxFallback() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveHystrixCircuitBreakerFactory().create("foo");
|
||||
assertEquals(Arrays.asList(new String[]{"fallback"}), cb.run(Flux.error(new RuntimeException("boom")), t -> Flux.just("fallback")).collectList().block());
|
||||
ReactiveCircuitBreaker cb = new ReactiveHystrixCircuitBreakerFactory()
|
||||
.create("foo");
|
||||
assertThat(cb
|
||||
.run(Flux.error(new RuntimeException("boom")), t -> Flux.just("fallback"))
|
||||
.collectList().block())
|
||||
.isEqualTo(Arrays.asList(new String[] { "fallback" }));
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -20,7 +20,6 @@ import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
import javax.annotation.PostConstruct;
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreaker;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
|
||||
@@ -34,7 +33,8 @@ import org.springframework.context.annotation.Configuration;
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
@Configuration
|
||||
@ConditionalOnClass(name = {"reactor.core.publisher.Mono", "reactor.core.publisher.Flux"})
|
||||
@ConditionalOnClass(name = { "reactor.core.publisher.Mono",
|
||||
"reactor.core.publisher.Flux" })
|
||||
public class ReactiveResilience4JAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
@@ -44,17 +44,21 @@ public class ReactiveResilience4JAutoConfiguration {
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@ConditionalOnClass(name = {"reactor.core.publisher.Mono", "reactor.core.publisher.Flux"})
|
||||
@ConditionalOnClass(name = { "reactor.core.publisher.Mono",
|
||||
"reactor.core.publisher.Flux" })
|
||||
public static class ReactiveResilience4JCustomizerConfiguration {
|
||||
@Autowired(required = false)
|
||||
public List<Customizer<ReactiveResilience4JCircuitBreakerFactory>> customizers = new ArrayList<>();
|
||||
|
||||
@Autowired(required = false)
|
||||
public ReactiveResilience4JCircuitBreakerFactory factory;
|
||||
private List<Customizer<ReactiveResilience4JCircuitBreakerFactory>> customizers = new ArrayList<>();
|
||||
|
||||
@Autowired(required = false)
|
||||
private ReactiveResilience4JCircuitBreakerFactory factory;
|
||||
|
||||
@PostConstruct
|
||||
public void init() {
|
||||
customizers.forEach(customizer -> customizer.customize(factory));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2018 the original author or authors.
|
||||
* Copyright 2013-2019 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.
|
||||
@@ -13,36 +13,39 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
import java.util.function.Function;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreaker;
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry;
|
||||
import io.github.resilience4j.reactor.circuitbreaker.operator.CircuitBreakerOperator;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
import java.util.function.Function;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.Customizer;
|
||||
import org.springframework.cloud.circuitbreaker.commons.ReactiveCircuitBreaker;
|
||||
|
||||
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreaker {
|
||||
|
||||
private String id;
|
||||
|
||||
private Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config;
|
||||
|
||||
private CircuitBreakerRegistry registry;
|
||||
|
||||
private Optional<Customizer<CircuitBreaker>> circuitBreakerCustomizer;
|
||||
|
||||
public ReactiveResilience4JCircuitBreaker(String id, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config,
|
||||
CircuitBreakerRegistry circuitBreakerRegistry,
|
||||
Optional<Customizer<CircuitBreaker>> circuitBreakerCustomizer) {
|
||||
public ReactiveResilience4JCircuitBreaker(String id,
|
||||
Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config,
|
||||
CircuitBreakerRegistry circuitBreakerRegistry,
|
||||
Optional<Customizer<CircuitBreaker>> circuitBreakerCustomizer) {
|
||||
this.id = id;
|
||||
this.config = config;
|
||||
this.registry = circuitBreakerRegistry;
|
||||
@@ -51,26 +54,38 @@ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreake
|
||||
|
||||
@Override
|
||||
public <T> Mono<T> run(Mono<T> toRun, Function<Throwable, Mono<T>> fallback) {
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry.circuitBreaker(id, config.getCircuitBreakerConfig());
|
||||
circuitBreakerCustomizer.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
Mono<T> toReturn = toRun.transform(CircuitBreakerOperator.of(defaultCircuitBreaker))
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry
|
||||
.circuitBreaker(id, config.getCircuitBreakerConfig());
|
||||
circuitBreakerCustomizer
|
||||
.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
Mono<T> toReturn = toRun
|
||||
.transform(CircuitBreakerOperator.of(defaultCircuitBreaker))
|
||||
.timeout(config.getTimeLimiterConfig().getTimeoutDuration())
|
||||
// Since we are using the Mono timeout we need to tell the circuit breaker about the error
|
||||
.doOnError(TimeoutException.class, t -> defaultCircuitBreaker.onError(config.getTimeLimiterConfig().getTimeoutDuration().toMillis(), t));
|
||||
if(fallback != null) {
|
||||
// Since we are using the Mono timeout we need to tell the circuit breaker
|
||||
// about the error
|
||||
.doOnError(TimeoutException.class, t -> defaultCircuitBreaker.onError(
|
||||
config.getTimeLimiterConfig().getTimeoutDuration().toMillis(),
|
||||
t));
|
||||
if (fallback != null) {
|
||||
toReturn = toReturn.onErrorResume(t -> fallback.apply(t));
|
||||
}
|
||||
return toReturn;
|
||||
}
|
||||
|
||||
public<T> Flux<T> run(Flux<T> toRun, Function<Throwable, Flux<T>> fallback) {
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry.circuitBreaker(id, config.getCircuitBreakerConfig());
|
||||
circuitBreakerCustomizer.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
Flux<T> toReturn = toRun.transform(CircuitBreakerOperator.of(defaultCircuitBreaker))
|
||||
public <T> Flux<T> run(Flux<T> toRun, Function<Throwable, Flux<T>> fallback) {
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry
|
||||
.circuitBreaker(id, config.getCircuitBreakerConfig());
|
||||
circuitBreakerCustomizer
|
||||
.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
Flux<T> toReturn = toRun
|
||||
.transform(CircuitBreakerOperator.of(defaultCircuitBreaker))
|
||||
.timeout(config.getTimeLimiterConfig().getTimeoutDuration())
|
||||
// Since we are using the Flux timeout we need to tell the circuit breaker about the error
|
||||
.doOnError(TimeoutException.class, t -> defaultCircuitBreaker.onError(config.getTimeLimiterConfig().getTimeoutDuration().toMillis(), t));
|
||||
if(fallback != null) {
|
||||
// Since we are using the Flux timeout we need to tell the circuit breaker
|
||||
// about the error
|
||||
.doOnError(TimeoutException.class, t -> defaultCircuitBreaker.onError(
|
||||
config.getTimeLimiterConfig().getTimeoutDuration().toMillis(),
|
||||
t));
|
||||
if (fallback != null) {
|
||||
toReturn = toReturn.onErrorResume(t -> fallback.apply(t));
|
||||
}
|
||||
return toReturn;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2018 the original author or authors.
|
||||
* Copyright 2013-2019 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.
|
||||
@@ -13,19 +13,19 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.function.Function;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreaker;
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiterConfig;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.function.Function;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.Customizer;
|
||||
import org.springframework.cloud.circuitbreaker.commons.ReactiveCircuitBreaker;
|
||||
import org.springframework.cloud.circuitbreaker.commons.ReactiveCircuitBreakerFactory;
|
||||
@@ -34,23 +34,25 @@ import org.springframework.util.Assert;
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public class ReactiveResilience4JCircuitBreakerFactory extends ReactiveCircuitBreakerFactory<Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration, Resilience4JConfigBuilder> {
|
||||
public class ReactiveResilience4JCircuitBreakerFactory extends
|
||||
ReactiveCircuitBreakerFactory<Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration, Resilience4JConfigBuilder> {
|
||||
|
||||
private Function<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration = id ->
|
||||
new Resilience4JConfigBuilder(id)
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.ofDefaults())
|
||||
.build();
|
||||
private Function<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration = id -> new Resilience4JConfigBuilder(
|
||||
id).circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.ofDefaults()).build();
|
||||
|
||||
private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry
|
||||
.ofDefaults();
|
||||
|
||||
private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry.ofDefaults();
|
||||
private Map<String, Customizer<CircuitBreaker>> circuitBreakerCustomizers = new HashMap<>();
|
||||
|
||||
@Override
|
||||
public ReactiveCircuitBreaker create(String id) {
|
||||
Assert.hasText(id, "A CircuitBreaker must have an id.");
|
||||
Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config = getConfigurations().computeIfAbsent(id, defaultConfiguration);
|
||||
return new ReactiveResilience4JCircuitBreaker(id, config,
|
||||
circuitBreakerRegistry, Optional.ofNullable(circuitBreakerCustomizers.get(id)));
|
||||
Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config = getConfigurations()
|
||||
.computeIfAbsent(id, defaultConfiguration);
|
||||
return new ReactiveResilience4JCircuitBreaker(id, config, circuitBreakerRegistry,
|
||||
Optional.ofNullable(circuitBreakerCustomizers.get(id)));
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -59,7 +61,8 @@ public class ReactiveResilience4JCircuitBreakerFactory extends ReactiveCircuitBr
|
||||
}
|
||||
|
||||
@Override
|
||||
public void configureDefault(Function<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration) {
|
||||
public void configureDefault(
|
||||
Function<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration) {
|
||||
this.defaultConfiguration = defaultConfiguration;
|
||||
}
|
||||
|
||||
@@ -67,9 +70,11 @@ public class ReactiveResilience4JCircuitBreakerFactory extends ReactiveCircuitBr
|
||||
this.circuitBreakerRegistry = registry;
|
||||
}
|
||||
|
||||
public void addCircuitBreakerCustomizer(Customizer<CircuitBreaker> customizer, String... ids) {
|
||||
for(String id : ids) {
|
||||
public void addCircuitBreakerCustomizer(Customizer<CircuitBreaker> customizer,
|
||||
String... ids) {
|
||||
for (String id : ids) {
|
||||
circuitBreakerCustomizers.put(id, customizer);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -13,20 +13,18 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
import javax.annotation.PostConstruct;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreaker;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
|
||||
import org.springframework.cloud.circuitbreaker.commons.CircuitBreakerFactory;
|
||||
import org.springframework.cloud.circuitbreaker.commons.Customizer;
|
||||
import org.springframework.cloud.circuitbreaker.commons.ReactiveCircuitBreakerFactory;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
@@ -44,18 +42,18 @@ public class Resilience4JAutoConfiguration {
|
||||
|
||||
@Configuration
|
||||
public static class Resilience4JCustomizerConfiguration {
|
||||
@Autowired(required = false)
|
||||
public List<Customizer<Resilience4JCircuitBreakerFactory>> customizers = new ArrayList<>();
|
||||
|
||||
@Autowired(required = false)
|
||||
public Resilience4JCircuitBreakerFactory factory;
|
||||
private List<Customizer<Resilience4JCircuitBreakerFactory>> customizers = new ArrayList<>();
|
||||
|
||||
@Autowired(required = false)
|
||||
private Resilience4JCircuitBreakerFactory factory;
|
||||
|
||||
@PostConstruct
|
||||
public void init() {
|
||||
customizers.forEach(customizer -> customizer.customize(factory));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
|
||||
}
|
||||
|
||||
@@ -16,37 +16,44 @@
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiter;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiterConfig;
|
||||
import io.vavr.control.Try;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiter;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiterConfig;
|
||||
import io.vavr.control.Try;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.CircuitBreaker;
|
||||
import org.springframework.cloud.circuitbreaker.commons.Customizer;
|
||||
|
||||
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public class Resilience4JCircuitBreaker implements CircuitBreaker {
|
||||
|
||||
private String id;
|
||||
|
||||
private io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig;
|
||||
|
||||
private CircuitBreakerRegistry registry;
|
||||
|
||||
private TimeLimiterConfig timeLimiterConfig;
|
||||
|
||||
private ExecutorService executorService;
|
||||
|
||||
private Optional<Customizer<io.github.resilience4j.circuitbreaker.CircuitBreaker>> circuitBreakerCustomizer;
|
||||
|
||||
public Resilience4JCircuitBreaker(String id, io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig,
|
||||
TimeLimiterConfig timeLimiterConfig, CircuitBreakerRegistry circuitBreakerRegistry,
|
||||
ExecutorService executorService, Optional<Customizer<io.github.resilience4j.circuitbreaker.CircuitBreaker>> circuitBreakerCustomizer) {
|
||||
public Resilience4JCircuitBreaker(String id,
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig,
|
||||
TimeLimiterConfig timeLimiterConfig,
|
||||
CircuitBreakerRegistry circuitBreakerRegistry,
|
||||
ExecutorService executorService,
|
||||
Optional<Customizer<io.github.resilience4j.circuitbreaker.CircuitBreaker>> circuitBreakerCustomizer) {
|
||||
this.id = id;
|
||||
this.circuitBreakerConfig = circuitBreakerConfig;
|
||||
this.registry = circuitBreakerRegistry;
|
||||
@@ -59,13 +66,16 @@ public class Resilience4JCircuitBreaker implements CircuitBreaker {
|
||||
public <T> T run(Supplier<T> toRun, Function<Throwable, T> fallback) {
|
||||
TimeLimiter timeLimiter = TimeLimiter.of(timeLimiterConfig);
|
||||
Supplier<Future<T>> futureSupplier = () -> executorService.submit(toRun::get);
|
||||
Callable restrictedCall = TimeLimiter
|
||||
.decorateFutureSupplier(timeLimiter, futureSupplier);
|
||||
Callable restrictedCall = TimeLimiter.decorateFutureSupplier(timeLimiter,
|
||||
futureSupplier);
|
||||
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry.circuitBreaker(id, circuitBreakerConfig);
|
||||
circuitBreakerCustomizer.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry
|
||||
.circuitBreaker(id, circuitBreakerConfig);
|
||||
circuitBreakerCustomizer
|
||||
.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
Callable<T> callable = io.github.resilience4j.circuitbreaker.CircuitBreaker
|
||||
.decorateCallable(defaultCircuitBreaker, restrictedCall);
|
||||
return Try.of(callable::call).recover(fallback).get();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -16,18 +16,18 @@
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreaker;
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiterConfig;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.function.Function;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreaker;
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiterConfig;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.CircuitBreakerFactory;
|
||||
import org.springframework.cloud.circuitbreaker.commons.Customizer;
|
||||
import org.springframework.util.Assert;
|
||||
@@ -35,16 +35,18 @@ import org.springframework.util.Assert;
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public class Resilience4JCircuitBreakerFactory extends CircuitBreakerFactory<Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration, Resilience4JConfigBuilder> {
|
||||
public class Resilience4JCircuitBreakerFactory extends
|
||||
CircuitBreakerFactory<Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration, Resilience4JConfigBuilder> {
|
||||
|
||||
private Function<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration = id ->
|
||||
new Resilience4JConfigBuilder(id)
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.ofDefaults())
|
||||
.build();
|
||||
private Function<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration = id -> new Resilience4JConfigBuilder(
|
||||
id).circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.ofDefaults()).build();
|
||||
|
||||
private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry
|
||||
.ofDefaults();
|
||||
|
||||
private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry.ofDefaults();
|
||||
private ExecutorService executorService = Executors.newSingleThreadExecutor();
|
||||
|
||||
private Map<String, Customizer<CircuitBreaker>> circuitBreakerCustomizers = new HashMap<>();
|
||||
|
||||
@Override
|
||||
@@ -53,7 +55,8 @@ public class Resilience4JCircuitBreakerFactory extends CircuitBreakerFactory<Res
|
||||
}
|
||||
|
||||
@Override
|
||||
public void configureDefault(Function<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration) {
|
||||
public void configureDefault(
|
||||
Function<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration) {
|
||||
this.defaultConfiguration = defaultConfiguration;
|
||||
}
|
||||
|
||||
@@ -68,13 +71,16 @@ public class Resilience4JCircuitBreakerFactory extends CircuitBreakerFactory<Res
|
||||
@Override
|
||||
public Resilience4JCircuitBreaker create(String id) {
|
||||
Assert.hasText(id, "A CircuitBreaker must have an id.");
|
||||
Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config = getConfigurations().computeIfAbsent(id, defaultConfiguration);
|
||||
return new Resilience4JCircuitBreaker(id, config.getCircuitBreakerConfig(), config.getTimeLimiterConfig(),
|
||||
circuitBreakerRegistry, executorService, Optional.ofNullable(circuitBreakerCustomizers.get(id)));
|
||||
Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config = getConfigurations()
|
||||
.computeIfAbsent(id, defaultConfiguration);
|
||||
return new Resilience4JCircuitBreaker(id, config.getCircuitBreakerConfig(),
|
||||
config.getTimeLimiterConfig(), circuitBreakerRegistry, executorService,
|
||||
Optional.ofNullable(circuitBreakerCustomizers.get(id)));
|
||||
}
|
||||
|
||||
public void addCircuitBreakerCustomizer(Customizer<CircuitBreaker> customizer, String... ids) {
|
||||
for(String id: ids) {
|
||||
public void addCircuitBreakerCustomizer(Customizer<CircuitBreaker> customizer,
|
||||
String... ids) {
|
||||
for (String id : ids) {
|
||||
circuitBreakerCustomizers.put(id, customizer);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2018 the original author or authors.
|
||||
* Copyright 2013-2019 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.
|
||||
@@ -13,6 +13,7 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
|
||||
@@ -23,10 +24,13 @@ import org.springframework.cloud.circuitbreaker.commons.ConfigBuilder;
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
*/
|
||||
public class Resilience4JConfigBuilder implements ConfigBuilder<Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> {
|
||||
public class Resilience4JConfigBuilder implements
|
||||
ConfigBuilder<Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> {
|
||||
|
||||
private String id;
|
||||
|
||||
private TimeLimiterConfig timeLimiterConfig;
|
||||
|
||||
private CircuitBreakerConfig circuitBreakerConfig;
|
||||
|
||||
public Resilience4JConfigBuilder(String id) {
|
||||
@@ -38,7 +42,8 @@ public class Resilience4JConfigBuilder implements ConfigBuilder<Resilience4JConf
|
||||
return this;
|
||||
}
|
||||
|
||||
public Resilience4JConfigBuilder circuitBreakerConfig(CircuitBreakerConfig circuitBreakerConfig) {
|
||||
public Resilience4JConfigBuilder circuitBreakerConfig(
|
||||
CircuitBreakerConfig circuitBreakerConfig) {
|
||||
this.circuitBreakerConfig = circuitBreakerConfig;
|
||||
return this;
|
||||
}
|
||||
@@ -47,15 +52,18 @@ public class Resilience4JConfigBuilder implements ConfigBuilder<Resilience4JConf
|
||||
public Resilience4JCircuitBreakerConfiguration build() {
|
||||
Resilience4JCircuitBreakerConfiguration config = new Resilience4JCircuitBreakerConfiguration();
|
||||
config.setId(id);
|
||||
//TODO null checks?
|
||||
// TODO null checks?
|
||||
config.setCircuitBreakerConfig(circuitBreakerConfig);
|
||||
config.setTimeLimiterConfig(timeLimiterConfig);
|
||||
return config;
|
||||
}
|
||||
|
||||
public static class Resilience4JCircuitBreakerConfiguration {
|
||||
|
||||
private String id;
|
||||
|
||||
private TimeLimiterConfig timeLimiterConfig;
|
||||
|
||||
private CircuitBreakerConfig circuitBreakerConfig;
|
||||
|
||||
public String getId() {
|
||||
@@ -81,5 +89,7 @@ public class Resilience4JConfigBuilder implements ConfigBuilder<Resilience4JConf
|
||||
public void setCircuitBreakerConfig(CircuitBreakerConfig circuitBreakerConfig) {
|
||||
this.circuitBreakerConfig = circuitBreakerConfig;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -13,22 +13,22 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreaker;
|
||||
import java.time.Duration;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
|
||||
import io.github.resilience4j.circuitbreaker.event.CircuitBreakerOnErrorEvent;
|
||||
import io.github.resilience4j.circuitbreaker.event.CircuitBreakerOnSuccessEvent;
|
||||
import io.github.resilience4j.core.EventConsumer;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiterConfig;
|
||||
import org.mockito.Mock;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import java.time.Duration;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.mockito.Mock;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
import reactor.test.StepVerifier;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
@@ -47,7 +47,7 @@ import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
import org.springframework.web.reactive.function.client.WebClient;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.internal.verification.VerificationModeFactory.times;
|
||||
@@ -60,30 +60,65 @@ import static org.springframework.boot.test.context.SpringBootTest.WebEnvironmen
|
||||
@SpringBootTest(webEnvironment = RANDOM_PORT, classes = ReactiveResilience4JCircuitBreakerIntegrationTest.Application.class)
|
||||
@DirtiesContext
|
||||
public class ReactiveResilience4JCircuitBreakerIntegrationTest {
|
||||
|
||||
@LocalServerPort
|
||||
int port = 0;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnErrorEvent> slowErrorConsumer;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnSuccessEvent> slowSuccessConsumer;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnErrorEvent> normalErrorConsumer;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnSuccessEvent> normalSuccessConsumer;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnErrorEvent> slowFluxErrorConsumer;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnSuccessEvent> slowFluxSuccessConsumer;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnErrorEvent> normalFluxErrorConsumer;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnSuccessEvent> normalFluxSuccessConsumer;
|
||||
|
||||
@Autowired
|
||||
ReactiveResilience4JCircuitBreakerIntegrationTest.Application.DemoControllerService service;
|
||||
|
||||
@Before
|
||||
public void setup() {
|
||||
service.setPort(port);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
assertThat(service.normal().block()).isEqualTo("normal");
|
||||
verify(normalErrorConsumer, times(0)).consumeEvent(any());
|
||||
verify(normalSuccessConsumer, times(1)).consumeEvent(any());
|
||||
assertThat(service.slow().block()).isEqualTo("fallback");
|
||||
verify(slowErrorConsumer, times(1)).consumeEvent(any());
|
||||
verify(slowSuccessConsumer, times(0)).consumeEvent(any());
|
||||
StepVerifier.create(service.normalFlux()).expectNext("normalflux")
|
||||
.verifyComplete();
|
||||
verify(normalFluxErrorConsumer, times(0)).consumeEvent(any());
|
||||
verify(normalFluxSuccessConsumer, times(1)).consumeEvent(any());
|
||||
StepVerifier.create(service.slowFlux()).expectNext("fluxfallback")
|
||||
.verifyComplete();
|
||||
verify(slowFluxErrorConsumer, times(1)).consumeEvent(any());
|
||||
verify(slowSuccessConsumer, times(0)).consumeEvent(any());
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableAutoConfiguration
|
||||
@RestController
|
||||
protected static class Application {
|
||||
|
||||
@GetMapping("/slow")
|
||||
public Mono<String> slow() {
|
||||
return Mono.just("slow").delayElement(Duration.ofSeconds(3));
|
||||
@@ -107,94 +142,91 @@ public class ReactiveResilience4JCircuitBreakerIntegrationTest {
|
||||
@Bean
|
||||
public Customizer<ReactiveResilience4JCircuitBreakerFactory> slowCusomtizer() {
|
||||
return factory -> {
|
||||
factory.configureDefault(id -> new Resilience4JConfigBuilder(id)
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(4)).build()).build());
|
||||
factory.configure(builder -> builder
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(2)).build())
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults()), "slow", "slowflux");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker ->
|
||||
circuitBreaker.getEventPublisher().onError(slowErrorConsumer).onSuccess(slowSuccessConsumer),
|
||||
"slow" );
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker.getEventPublisher().onError(normalErrorConsumer).onSuccess(normalSuccessConsumer),
|
||||
"normal");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker.getEventPublisher().onError(slowFluxErrorConsumer).onSuccess(slowFluxSuccessConsumer),
|
||||
"slowflux");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker.getEventPublisher().onError(normalFluxErrorConsumer).onSuccess(normalFluxSuccessConsumer),
|
||||
"normalflux");
|
||||
factory.configureDefault(
|
||||
id -> new Resilience4JConfigBuilder(id)
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom()
|
||||
.timeoutDuration(Duration.ofSeconds(4)).build())
|
||||
.build());
|
||||
factory.configure(
|
||||
builder -> builder
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom()
|
||||
.timeoutDuration(Duration.ofSeconds(2)).build())
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults()),
|
||||
"slow", "slowflux");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker
|
||||
.getEventPublisher().onError(slowErrorConsumer)
|
||||
.onSuccess(slowSuccessConsumer), "slow");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker
|
||||
.getEventPublisher().onError(normalErrorConsumer)
|
||||
.onSuccess(normalSuccessConsumer), "normal");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker
|
||||
.getEventPublisher().onError(slowFluxErrorConsumer)
|
||||
.onSuccess(slowFluxSuccessConsumer), "slowflux");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker
|
||||
.getEventPublisher().onError(normalFluxErrorConsumer)
|
||||
.onSuccess(normalFluxSuccessConsumer), "normalflux");
|
||||
};
|
||||
}
|
||||
|
||||
@Service
|
||||
public static class DemoControllerService {
|
||||
|
||||
private int port = 0;
|
||||
|
||||
private ReactiveCircuitBreakerFactory cbFactory;
|
||||
|
||||
|
||||
public DemoControllerService(ReactiveCircuitBreakerFactory cbFactory) {
|
||||
DemoControllerService(ReactiveCircuitBreakerFactory cbFactory) {
|
||||
this.cbFactory = cbFactory;
|
||||
}
|
||||
|
||||
public Mono<String> slow() {
|
||||
return cbFactory.create("slow").run(WebClient.builder().baseUrl("http://localhost:" + port).build()
|
||||
.get().uri("/slow").retrieve().bodyToMono(String.class), t -> {
|
||||
t.printStackTrace();
|
||||
return Mono.just("fallback");
|
||||
});
|
||||
return cbFactory.create("slow").run(
|
||||
WebClient.builder().baseUrl("http://localhost:" + port).build()
|
||||
.get().uri("/slow").retrieve().bodyToMono(String.class),
|
||||
t -> {
|
||||
t.printStackTrace();
|
||||
return Mono.just("fallback");
|
||||
});
|
||||
}
|
||||
|
||||
public Mono<String> normal() {
|
||||
return cbFactory.create("normal").run(WebClient.builder().baseUrl("http://localhost:" + port).build()
|
||||
.get().uri("/normal").retrieve().bodyToMono(String.class), t -> {
|
||||
t.printStackTrace();
|
||||
return Mono.just("fallback");
|
||||
});
|
||||
return cbFactory.create("normal").run(
|
||||
WebClient.builder().baseUrl("http://localhost:" + port).build()
|
||||
.get().uri("/normal").retrieve().bodyToMono(String.class),
|
||||
t -> {
|
||||
t.printStackTrace();
|
||||
return Mono.just("fallback");
|
||||
});
|
||||
}
|
||||
|
||||
public Flux<String> slowFlux() {
|
||||
return cbFactory.create("slowflux").run(WebClient.builder().baseUrl("http://localhost:" + port).build()
|
||||
.get().uri("/slowflux").retrieve().bodyToFlux(new ParameterizedTypeReference<String>() { }), t -> {
|
||||
t.printStackTrace();
|
||||
return Flux.just("fluxfallback");
|
||||
});
|
||||
return cbFactory.create("slowflux")
|
||||
.run(WebClient.builder().baseUrl("http://localhost:" + port)
|
||||
.build().get().uri("/slowflux").retrieve()
|
||||
.bodyToFlux(new ParameterizedTypeReference<String>() {
|
||||
}), t -> {
|
||||
t.printStackTrace();
|
||||
return Flux.just("fluxfallback");
|
||||
});
|
||||
}
|
||||
|
||||
public Flux<String> normalFlux() {
|
||||
return cbFactory.create("normalflux").run(WebClient.builder().baseUrl("http://localhost:" + port).build()
|
||||
.get().uri("/normalflux").retrieve().bodyToFlux(String.class), t -> {
|
||||
t.printStackTrace();
|
||||
return Flux.just("fluxfallback");
|
||||
});
|
||||
return cbFactory.create("normalflux")
|
||||
.run(WebClient.builder().baseUrl("http://localhost:" + port)
|
||||
.build().get().uri("/normalflux").retrieve()
|
||||
.bodyToFlux(String.class), t -> {
|
||||
t.printStackTrace();
|
||||
return Flux.just("fluxfallback");
|
||||
});
|
||||
}
|
||||
|
||||
public void setPort(int port) {
|
||||
this.port = port;
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@Autowired
|
||||
ReactiveResilience4JCircuitBreakerIntegrationTest.Application.DemoControllerService service;
|
||||
|
||||
@Before
|
||||
public void setup() {
|
||||
service.setPort(port);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
assertEquals("normal", service.normal().block());
|
||||
verify(normalErrorConsumer, times(0)).consumeEvent(any());
|
||||
verify(normalSuccessConsumer, times(1)).consumeEvent(any());
|
||||
assertEquals("fallback", service.slow().block());
|
||||
verify(slowErrorConsumer, times(1)).consumeEvent(any());
|
||||
verify(slowSuccessConsumer, times(0)).consumeEvent(any());
|
||||
StepVerifier.create(service.normalFlux()).expectNext("normalflux").verifyComplete();
|
||||
verify(normalFluxErrorConsumer, times(0)).consumeEvent(any());
|
||||
verify(normalFluxSuccessConsumer, times(1)).consumeEvent(any());
|
||||
StepVerifier.create(service.slowFlux()).expectNext("fluxfallback").verifyComplete();
|
||||
verify(slowFluxErrorConsumer, times(1)).consumeEvent(any());
|
||||
verify(slowSuccessConsumer, times(0)).consumeEvent(any());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -13,18 +13,18 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
|
||||
import java.util.Arrays;
|
||||
|
||||
import org.junit.Test;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import org.junit.Test;
|
||||
import org.springframework.cloud.circuitbreaker.commons.ReactiveCircuitBreaker;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
@@ -33,26 +33,35 @@ public class ReactiveResilience4JCircuitBreakerTest {
|
||||
|
||||
@Test
|
||||
public void runMono() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals("foobar", cb.run(Mono.just("foobar")).block());
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory()
|
||||
.create("foo");
|
||||
assertThat(cb.run(Mono.just("foobar")).block()).isEqualTo("foobar");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runMonoWithFallback() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals("fallback", cb.run(Mono.error(new RuntimeException("boom")), t -> Mono.just("fallback")).block());
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory()
|
||||
.create("foo");
|
||||
assertThat(cb
|
||||
.run(Mono.error(new RuntimeException("boom")), t -> Mono.just("fallback"))
|
||||
.block()).isEqualTo("fallback");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runFlux() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals(Arrays.asList("foobar", "hello world"), cb.run(Flux.just("foobar", "hello world")).collectList().block());
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory()
|
||||
.create("foo");
|
||||
assertThat(cb.run(Flux.just("foobar", "hello world")).collectList().block())
|
||||
.isEqualTo(Arrays.asList("foobar", "hello world"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runFluxWithFallback() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals(Arrays.asList("fallback"), cb.run(Flux.error(new RuntimeException("boom")), t -> Flux.just("fallback")).collectList().block());
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory()
|
||||
.create("foo");
|
||||
assertThat(cb
|
||||
.run(Flux.error(new RuntimeException("boom")), t -> Flux.just("fallback"))
|
||||
.collectList().block()).isEqualTo(Arrays.asList("fallback"));
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -16,13 +16,13 @@
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
|
||||
import io.github.resilience4j.circuitbreaker.event.CircuitBreakerOnErrorEvent;
|
||||
import io.github.resilience4j.circuitbreaker.event.CircuitBreakerOnSuccessEvent;
|
||||
import io.github.resilience4j.core.EventConsumer;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiterConfig;
|
||||
|
||||
import java.time.Duration;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.mockito.Mock;
|
||||
@@ -41,7 +41,7 @@ import org.springframework.test.context.junit4.SpringRunner;
|
||||
import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.internal.verification.VerificationModeFactory.times;
|
||||
@@ -55,20 +55,40 @@ import static org.springframework.boot.test.context.SpringBootTest.WebEnvironmen
|
||||
@DirtiesContext
|
||||
public class Resilience4JCircuitBreakerIntegrationTest {
|
||||
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnErrorEvent> slowErrorConsumer;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnSuccessEvent> slowSuccessConsumer;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnErrorEvent> normalErrorConsumer;
|
||||
|
||||
@Mock
|
||||
static EventConsumer<CircuitBreakerOnSuccessEvent> normalSuccessConsumer;
|
||||
|
||||
@Autowired
|
||||
Application.DemoControllerService service;
|
||||
|
||||
@Test
|
||||
public void testSlow() {
|
||||
assertThat(service.slow()).isEqualTo("fallback");
|
||||
verify(slowErrorConsumer, times(1)).consumeEvent(any());
|
||||
verify(slowSuccessConsumer, times(0)).consumeEvent(any());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testNormal() {
|
||||
assertThat(service.normal()).isEqualTo("normal");
|
||||
verify(normalErrorConsumer, times(0)).consumeEvent(any());
|
||||
verify(normalSuccessConsumer, times(1)).consumeEvent(any());
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableAutoConfiguration
|
||||
@RestController
|
||||
protected static class Application {
|
||||
|
||||
@GetMapping("/slow")
|
||||
public String slow() throws InterruptedException {
|
||||
Thread.sleep(3000);
|
||||
@@ -83,55 +103,51 @@ public class Resilience4JCircuitBreakerIntegrationTest {
|
||||
@Bean
|
||||
public Customizer<Resilience4JCircuitBreakerFactory> slowCustomizer() {
|
||||
return factory -> {
|
||||
factory.configure(builder -> builder.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(2)).build()), "slow");
|
||||
factory.configure(
|
||||
builder -> builder
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom()
|
||||
.timeoutDuration(Duration.ofSeconds(2)).build()),
|
||||
"slow");
|
||||
factory.configureDefault(id -> new Resilience4JConfigBuilder(id)
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(4)).build())
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.build());
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker ->
|
||||
circuitBreaker.getEventPublisher().onError(slowErrorConsumer).onSuccess(slowSuccessConsumer),
|
||||
"slow" );
|
||||
factory.addCircuitBreakerCustomizer(
|
||||
circuitBreaker -> circuitBreaker.getEventPublisher().onError(normalErrorConsumer).onSuccess(normalSuccessConsumer),
|
||||
"normal");
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom()
|
||||
.timeoutDuration(Duration.ofSeconds(4)).build())
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults()).build());
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker
|
||||
.getEventPublisher().onError(slowErrorConsumer)
|
||||
.onSuccess(slowSuccessConsumer), "slow");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker
|
||||
.getEventPublisher().onError(normalErrorConsumer)
|
||||
.onSuccess(normalSuccessConsumer), "normal");
|
||||
};
|
||||
}
|
||||
|
||||
@Service
|
||||
public static class DemoControllerService {
|
||||
|
||||
private TestRestTemplate rest;
|
||||
|
||||
private CircuitBreakerFactory cbFactory;
|
||||
|
||||
public DemoControllerService(TestRestTemplate rest, CircuitBreakerFactory cbFactory) {
|
||||
DemoControllerService(TestRestTemplate rest,
|
||||
CircuitBreakerFactory cbFactory) {
|
||||
this.rest = rest;
|
||||
this.cbFactory = cbFactory;
|
||||
}
|
||||
|
||||
public String slow() {
|
||||
return cbFactory.create("slow").run(() -> rest.getForObject("/slow", String.class), t -> "fallback");
|
||||
return cbFactory.create("slow").run(
|
||||
() -> rest.getForObject("/slow", String.class), t -> "fallback");
|
||||
}
|
||||
|
||||
public String normal() {
|
||||
return cbFactory.create("normal").run(() -> rest.getForObject("/normal", String.class), t -> "fallback");
|
||||
return cbFactory.create("normal").run(
|
||||
() -> rest.getForObject("/normal", String.class),
|
||||
t -> "fallback");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Autowired
|
||||
Application.DemoControllerService service;
|
||||
|
||||
@Test
|
||||
public void testSlow() {
|
||||
assertEquals("fallback", service.slow());
|
||||
verify(slowErrorConsumer, times(1)).consumeEvent(any());
|
||||
verify(slowSuccessConsumer, times(0)).consumeEvent(any());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testNormal() {
|
||||
assertEquals("normal", service.normal());
|
||||
verify(normalErrorConsumer, times(0)).consumeEvent(any());
|
||||
verify(normalSuccessConsumer, times(1)).consumeEvent(any());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -16,12 +16,11 @@
|
||||
|
||||
package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
|
||||
import java.util.Collections;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.CircuitBreaker;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
/**
|
||||
* @author Ryan Baxter
|
||||
@@ -31,14 +30,15 @@ public class Resilience4JCircuitBreakerTest {
|
||||
@Test
|
||||
public void run() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals("foobar", cb.run(() -> "foobar"));
|
||||
assertThat(cb.run(() -> "foobar")).isEqualTo("foobar");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runWithFallback() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals("fallback", cb.run(() -> {
|
||||
assertThat((String) cb.run(() -> {
|
||||
throw new RuntimeException("boom");
|
||||
}, t -> "fallback"));
|
||||
}, t -> "fallback")).isEqualTo("fallback");
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user