diff --git a/.editorconfig b/.editorconfig
new file mode 100644
index 0000000..8c727f3
--- /dev/null
+++ b/.editorconfig
@@ -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
\ No newline at end of file
diff --git a/pom.xml b/pom.xml
index c19bbfe..7fa2bbe 100644
--- a/pom.xml
+++ b/pom.xml
@@ -33,6 +33,10 @@
java
2.2.0.BUILD-SNAPSHOT
2.2.0.BUILD-SNAPSHOT
+ true
+
+ true
+
@@ -101,6 +105,14 @@
1.8
+
+ io.spring.javaformat
+ spring-javaformat-maven-plugin
+
+
+ org.apache.maven.plugins
+ maven-checkstyle-plugin
+
@@ -209,4 +221,4 @@
-
\ No newline at end of file
+
diff --git a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/AbstractCircuitBreakerFactory.java b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/AbstractCircuitBreakerFactory.java
index 2346f4c..c7b3296 100644
--- a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/AbstractCircuitBreakerFactory.java
+++ b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/AbstractCircuitBreakerFactory.java
@@ -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> {
+
private final ConcurrentHashMap 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 consumer, String ... ids) {
- for(String id : ids) {
+ public void configure(Consumer 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 getConfigurations() {
@@ -52,15 +53,16 @@ public abstract class AbstractCircuitBreakerFactory defaultConfiguration);
+
}
diff --git a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/CircuitBreaker.java b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/CircuitBreaker.java
index b320869..8cbdf79 100644
--- a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/CircuitBreaker.java
+++ b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/CircuitBreaker.java
@@ -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 {
diff --git a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/CircuitBreakerFactory.java b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/CircuitBreakerFactory.java
index 7fd5a7e..22d092e 100644
--- a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/CircuitBreakerFactory.java
+++ b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/CircuitBreakerFactory.java
@@ -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> extends AbstractCircuitBreakerFactory {
+public abstract class CircuitBreakerFactory>
+ extends AbstractCircuitBreakerFactory {
public abstract CircuitBreaker create(String id);
+
}
diff --git a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ConfigBuilder.java b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ConfigBuilder.java
index a053202..8957cb8 100644
--- a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ConfigBuilder.java
+++ b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ConfigBuilder.java
@@ -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 build();
-}
+ CONF build();
+
+}
diff --git a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/Customizer.java b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/Customizer.java
index c171382..6aa1641 100644
--- a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/Customizer.java
+++ b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/Customizer.java
@@ -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 {
+
void customize(TOCUSTOMIZE tocustomize);
+
}
diff --git a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/NoFallbackAvailableException.java b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/NoFallbackAvailableException.java
index 7861355..dee7f51 100644
--- a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/NoFallbackAvailableException.java
+++ b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/NoFallbackAvailableException.java
@@ -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);
}
+
}
diff --git a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ReactiveCircuitBreaker.java b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ReactiveCircuitBreaker.java
index c5f55c5..54c0fa6 100644
--- a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ReactiveCircuitBreaker.java
+++ b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ReactiveCircuitBreaker.java
@@ -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
*/
diff --git a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ReactiveCircuitBreakerFactory.java b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ReactiveCircuitBreakerFactory.java
index 5d4eeb0..1a2a356 100644
--- a/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ReactiveCircuitBreakerFactory.java
+++ b/spring-cloud-circuitbreaker-commons/src/main/java/org/springframework/cloud/circuitbreaker/commons/ReactiveCircuitBreakerFactory.java
@@ -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> extends AbstractCircuitBreakerFactory {
+public abstract class ReactiveCircuitBreakerFactory>
+ extends AbstractCircuitBreakerFactory {
public abstract ReactiveCircuitBreaker create(String id);
+
}
diff --git a/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/AbstractHystrixConfigBuilder.java b/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/AbstractHystrixConfigBuilder.java
index 2a1c14c..86926f8 100644
--- a/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/AbstractHystrixConfigBuilder.java
+++ b/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/AbstractHystrixConfigBuilder.java
@@ -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 implements ConfigBuilder {
+public abstract class AbstractHystrixConfigBuilder
+ implements ConfigBuilder {
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 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 implements ConfigBuild
}
protected HystrixCommandProperties.Setter getCommandPropertiesSetter() {
- return this.commandProperties != null
- ? this.commandProperties
+ return this.commandProperties != null ? this.commandProperties
: HystrixCommandProperties.Setter();
}
diff --git a/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreaker.java b/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreaker.java
index aa61df8..d779905 100644
--- a/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreaker.java
+++ b/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreaker.java
@@ -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();
}
+
}
diff --git a/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerAutoConfiguration.java b/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerAutoConfiguration.java
index d0fa0e1..1411972 100644
--- a/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerAutoConfiguration.java
+++ b/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerAutoConfiguration.java
@@ -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));
}
+
}
+
}
diff --git a/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerFactory.java b/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerFactory.java
index 9add3dd..1fce486 100644
--- a/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerFactory.java
+++ b/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerFactory.java
@@ -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 {
+public class HystrixCircuitBreakerFactory extends
+ CircuitBreakerFactory {
- private Function defaultConfiguration = id ->
- HystrixCommand.Setter.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id));
+ private Function defaultConfiguration = id -> HystrixCommand.Setter
+ .withGroupKey(HystrixCommandGroupKey.Factory.asKey(id));
-
- public void configureDefault(Function defaultConfiguration) {
+ public void configureDefault(
+ Function defaultConfiguration) {
this.defaultConfiguration = defaultConfiguration;
}
@@ -43,11 +46,13 @@ public class HystrixCircuitBreakerFactory extends CircuitBreakerFactory {
+ public static class HystrixConfigBuilder
+ extends AbstractHystrixConfigBuilder {
public HystrixConfigBuilder(String id) {
super(id);
@@ -55,9 +60,11 @@ public class HystrixCircuitBreakerFactory extends CircuitBreakerFactory 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 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 HystrixObservableCommand createCommand(Publisher toRun, Function fallback) {
+ private HystrixObservableCommand createCommand(Publisher toRun,
+ Function fallback) {
HystrixObservableCommand command = new HystrixObservableCommand(setter) {
@Override
protected Observable construct() {
@@ -67,12 +71,14 @@ public class ReactiveHystrixCircuitBreaker implements ReactiveCircuitBreaker {
@Override
protected Observable 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;
}
+
}
diff --git a/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerFactory.java b/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerFactory.java
index 9b19bf3..3b75ae0 100644
--- a/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerFactory.java
+++ b/spring-cloud-circuitbreaker-hystrix/src/main/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerFactory.java
@@ -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 {
+public class ReactiveHystrixCircuitBreakerFactory extends
+ ReactiveCircuitBreakerFactory {
- private Function defaultConfiguration = id ->
- HystrixObservableCommand.Setter.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id));
+ private Function 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 defaultConfiguration) {
+ public void configureDefault(
+ Function 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 {
+ public static class ReactiveHystrixConfigBuilder
+ extends AbstractHystrixConfigBuilder {
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());
}
+
}
+
}
diff --git a/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerIntegrationTest.java b/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerIntegrationTest.java
index 555de46..8ff19ad 100644
--- a/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerIntegrationTest.java
+++ b/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerIntegrationTest.java
@@ -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 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 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());
- }
}
diff --git a/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerTest.java b/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerTest.java
index f7d3d7c..e465ee6 100644
--- a/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerTest.java
+++ b/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/HystrixCircuitBreakerTest.java
@@ -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");
}
-}
\ No newline at end of file
+
+}
diff --git a/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerIntegrationTest.java b/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerIntegrationTest.java
index f6d1681..4f773b5 100644
--- a/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerIntegrationTest.java
+++ b/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerIntegrationTest.java
@@ -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 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 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 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 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());
- }
}
diff --git a/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerTest.java b/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerTest.java
index 1334707..9a04501 100644
--- a/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerTest.java
+++ b/spring-cloud-circuitbreaker-hystrix/src/test/java/org/springframework/cloud/circuitbreaker/hystrix/ReactiveHystrixCircuitBreakerTest.java
@@ -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 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 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 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" }));
}
-}
\ No newline at end of file
+
+}
diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfiguration.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfiguration.java
index 6d22e28..0612fcf 100644
--- a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfiguration.java
+++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfiguration.java
@@ -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> customizers = new ArrayList<>();
@Autowired(required = false)
- public ReactiveResilience4JCircuitBreakerFactory factory;
+ private List> customizers = new ArrayList<>();
+
+ @Autowired(required = false)
+ private ReactiveResilience4JCircuitBreakerFactory factory;
@PostConstruct
public void init() {
customizers.forEach(customizer -> customizer.customize(factory));
}
+
}
+
}
diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreaker.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreaker.java
index 66d2ec0..adb1f24 100644
--- a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreaker.java
+++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreaker.java
@@ -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> circuitBreakerCustomizer;
- public ReactiveResilience4JCircuitBreaker(String id, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config,
- CircuitBreakerRegistry circuitBreakerRegistry,
- Optional> circuitBreakerCustomizer) {
+ public ReactiveResilience4JCircuitBreaker(String id,
+ Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config,
+ CircuitBreakerRegistry circuitBreakerRegistry,
+ Optional> circuitBreakerCustomizer) {
this.id = id;
this.config = config;
this.registry = circuitBreakerRegistry;
@@ -51,26 +54,38 @@ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreake
@Override
public Mono run(Mono toRun, Function> fallback) {
- io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry.circuitBreaker(id, config.getCircuitBreakerConfig());
- circuitBreakerCustomizer.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
- Mono toReturn = toRun.transform(CircuitBreakerOperator.of(defaultCircuitBreaker))
+ io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry
+ .circuitBreaker(id, config.getCircuitBreakerConfig());
+ circuitBreakerCustomizer
+ .ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
+ Mono 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 Flux run(Flux toRun, Function> fallback) {
- io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry.circuitBreaker(id, config.getCircuitBreakerConfig());
- circuitBreakerCustomizer.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
- Flux toReturn = toRun.transform(CircuitBreakerOperator.of(defaultCircuitBreaker))
+ public Flux run(Flux toRun, Function> fallback) {
+ io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry
+ .circuitBreaker(id, config.getCircuitBreakerConfig());
+ circuitBreakerCustomizer
+ .ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
+ Flux 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;
diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreakerFactory.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreakerFactory.java
index 4167106..0cb737c 100644
--- a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreakerFactory.java
+++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreakerFactory.java
@@ -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 {
+public class ReactiveResilience4JCircuitBreakerFactory extends
+ ReactiveCircuitBreakerFactory {
- private Function defaultConfiguration = id ->
- new Resilience4JConfigBuilder(id)
- .circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
- .timeLimiterConfig(TimeLimiterConfig.ofDefaults())
- .build();
+ private Function defaultConfiguration = id -> new Resilience4JConfigBuilder(
+ id).circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
+ .timeLimiterConfig(TimeLimiterConfig.ofDefaults()).build();
+
+ private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry
+ .ofDefaults();
- private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry.ofDefaults();
private Map> 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 defaultConfiguration) {
+ public void configureDefault(
+ Function defaultConfiguration) {
this.defaultConfiguration = defaultConfiguration;
}
@@ -67,9 +70,11 @@ public class ReactiveResilience4JCircuitBreakerFactory extends ReactiveCircuitBr
this.circuitBreakerRegistry = registry;
}
- public void addCircuitBreakerCustomizer(Customizer customizer, String... ids) {
- for(String id : ids) {
+ public void addCircuitBreakerCustomizer(Customizer customizer,
+ String... ids) {
+ for (String id : ids) {
circuitBreakerCustomizers.put(id, customizer);
}
}
+
}
diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JAutoConfiguration.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JAutoConfiguration.java
index 9d972d2..ed69b27 100644
--- a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JAutoConfiguration.java
+++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JAutoConfiguration.java
@@ -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> customizers = new ArrayList<>();
@Autowired(required = false)
- public Resilience4JCircuitBreakerFactory factory;
+ private List> customizers = new ArrayList<>();
+
+ @Autowired(required = false)
+ private Resilience4JCircuitBreakerFactory factory;
@PostConstruct
public void init() {
customizers.forEach(customizer -> customizer.customize(factory));
}
+
}
-
-
}
diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreaker.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreaker.java
index 350a2fd..8d26498 100644
--- a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreaker.java
+++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreaker.java
@@ -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> circuitBreakerCustomizer;
- public Resilience4JCircuitBreaker(String id, io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig,
- TimeLimiterConfig timeLimiterConfig, CircuitBreakerRegistry circuitBreakerRegistry,
- ExecutorService executorService, Optional> circuitBreakerCustomizer) {
+ public Resilience4JCircuitBreaker(String id,
+ io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig,
+ TimeLimiterConfig timeLimiterConfig,
+ CircuitBreakerRegistry circuitBreakerRegistry,
+ ExecutorService executorService,
+ Optional> circuitBreakerCustomizer) {
this.id = id;
this.circuitBreakerConfig = circuitBreakerConfig;
this.registry = circuitBreakerRegistry;
@@ -59,13 +66,16 @@ public class Resilience4JCircuitBreaker implements CircuitBreaker {
public T run(Supplier toRun, Function fallback) {
TimeLimiter timeLimiter = TimeLimiter.of(timeLimiterConfig);
Supplier> 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 callable = io.github.resilience4j.circuitbreaker.CircuitBreaker
.decorateCallable(defaultCircuitBreaker, restrictedCall);
return Try.of(callable::call).recover(fallback).get();
}
+
}
diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerFactory.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerFactory.java
index 169bf47..3197455 100644
--- a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerFactory.java
+++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerFactory.java
@@ -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 {
+public class Resilience4JCircuitBreakerFactory extends
+ CircuitBreakerFactory {
- private Function defaultConfiguration = id ->
- new Resilience4JConfigBuilder(id)
- .circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
- .timeLimiterConfig(TimeLimiterConfig.ofDefaults())
- .build();
+ private Function 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> circuitBreakerCustomizers = new HashMap<>();
@Override
@@ -53,7 +55,8 @@ public class Resilience4JCircuitBreakerFactory extends CircuitBreakerFactory defaultConfiguration) {
+ public void configureDefault(
+ Function defaultConfiguration) {
this.defaultConfiguration = defaultConfiguration;
}
@@ -68,13 +71,16 @@ public class Resilience4JCircuitBreakerFactory extends CircuitBreakerFactory customizer, String... ids) {
- for(String id: ids) {
+ public void addCircuitBreakerCustomizer(Customizer customizer,
+ String... ids) {
+ for (String id : ids) {
circuitBreakerCustomizers.put(id, customizer);
}
}
diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JConfigBuilder.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JConfigBuilder.java
index 2fc38c8..eb7bfc1 100644
--- a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JConfigBuilder.java
+++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JConfigBuilder.java
@@ -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 {
+public class Resilience4JConfigBuilder implements
+ ConfigBuilder {
private String id;
+
private TimeLimiterConfig timeLimiterConfig;
+
private CircuitBreakerConfig circuitBreakerConfig;
public Resilience4JConfigBuilder(String id) {
@@ -38,7 +42,8 @@ public class Resilience4JConfigBuilder implements ConfigBuilder slowErrorConsumer;
+
@Mock
static EventConsumer slowSuccessConsumer;
+
@Mock
static EventConsumer normalErrorConsumer;
+
@Mock
static EventConsumer normalSuccessConsumer;
+
@Mock
static EventConsumer slowFluxErrorConsumer;
+
@Mock
static EventConsumer slowFluxSuccessConsumer;
+
@Mock
static EventConsumer normalFluxErrorConsumer;
+
@Mock
static EventConsumer 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 slow() {
return Mono.just("slow").delayElement(Duration.ofSeconds(3));
@@ -107,94 +142,91 @@ public class ReactiveResilience4JCircuitBreakerIntegrationTest {
@Bean
public Customizer 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 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 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 slowFlux() {
- return cbFactory.create("slowflux").run(WebClient.builder().baseUrl("http://localhost:" + port).build()
- .get().uri("/slowflux").retrieve().bodyToFlux(new ParameterizedTypeReference() { }), 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() {
+ }), t -> {
+ t.printStackTrace();
+ return Flux.just("fluxfallback");
+ });
}
public Flux 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());
}
}
diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreakerTest.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreakerTest.java
index b94132a..734e400 100644
--- a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreakerTest.java
+++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JCircuitBreakerTest.java
@@ -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"));
}
-}
\ No newline at end of file
+}
diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerIntegrationTest.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerIntegrationTest.java
index dc0ebd3..48bbe5e 100644
--- a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerIntegrationTest.java
+++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerIntegrationTest.java
@@ -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 slowErrorConsumer;
+
@Mock
static EventConsumer slowSuccessConsumer;
+
@Mock
static EventConsumer normalErrorConsumer;
+
@Mock
static EventConsumer 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 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());
- }
}
diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerTest.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerTest.java
index cbff809..225ac45 100644
--- a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerTest.java
+++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerTest.java
@@ -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");
}
-}
\ No newline at end of file
+
+}