From 7ea953dfa26a4dddf5cb2119db3b48e2847caef2 Mon Sep 17 00:00:00 2001 From: Ryan Baxter Date: Tue, 12 Mar 2019 20:04:07 -0400 Subject: [PATCH] Add support for checkstyle rules and fix checkstyle failures. Fixes #14 --- .editorconfig | 18 ++ pom.xml | 14 +- .../AbstractCircuitBreakerFactory.java | 20 ++- .../commons/CircuitBreaker.java | 4 +- .../commons/CircuitBreakerFactory.java | 7 +- .../circuitbreaker/commons/ConfigBuilder.java | 10 +- .../circuitbreaker/commons/Customizer.java | 7 +- .../commons/NoFallbackAvailableException.java | 2 + .../commons/ReactiveCircuitBreaker.java | 9 +- .../ReactiveCircuitBreakerFactory.java | 7 +- .../hystrix/AbstractHystrixConfigBuilder.java | 18 +- .../hystrix/HystrixCircuitBreaker.java | 6 +- ...ystrixCircuitBreakerAutoConfiguration.java | 17 +- .../hystrix/HystrixCircuitBreakerFactory.java | 29 +-- .../ReactiveHystrixCircuitBreaker.java | 22 ++- .../ReactiveHystrixCircuitBreakerFactory.java | 28 +-- .../HystrixCircuitBreakerIntegrationTest.java | 60 ++++--- .../hystrix/HystrixCircuitBreakerTest.java | 13 +- ...eHystrixCircuitBreakerIntegrationTest.java | 94 +++++----- .../ReactiveHystrixCircuitBreakerTest.java | 38 ++-- ...ReactiveResilience4JAutoConfiguration.java | 16 +- .../ReactiveResilience4JCircuitBreaker.java | 61 ++++--- ...tiveResilience4JCircuitBreakerFactory.java | 45 ++--- .../Resilience4JAutoConfiguration.java | 16 +- .../Resilience4JCircuitBreaker.java | 38 ++-- .../Resilience4JCircuitBreakerFactory.java | 44 +++-- .../Resilience4JConfigBuilder.java | 18 +- ...lience4JCircuitBreakerIntegrationTest.java | 166 +++++++++++------- ...eactiveResilience4JCircuitBreakerTest.java | 37 ++-- ...lience4JCircuitBreakerIntegrationTest.java | 84 +++++---- .../Resilience4JCircuitBreakerTest.java | 14 +- 31 files changed, 594 insertions(+), 368 deletions(-) create mode 100644 .editorconfig 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 + +}