From e5ec2428d7adf3f852db3ccbc316253f3348e7f2 Mon Sep 17 00:00:00 2001 From: Ryan Baxter Date: Wed, 26 May 2021 14:27:34 -0400 Subject: [PATCH] Resolve merge conflicts and cherry pick #108 --- ...ing-cloud-circuitbreaker-resilience4j.adoc | 33 +++++++ ...ReactiveResilience4JAutoConfiguration.java | 12 ++- .../ReactiveResilience4JCircuitBreaker.java | 68 ++++++++++---- ...tiveResilience4JCircuitBreakerFactory.java | 34 +++++-- ...lience4JAutoConfigurationPropertyTest.java | 88 +++++++++++++++++++ ...4JAutoConfigurationWithoutMetricsTest.java | 9 +- ...eactiveResilience4JCircuitBreakerTest.java | 35 +++++--- 7 files changed, 236 insertions(+), 43 deletions(-) create mode 100644 spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationPropertyTest.java diff --git a/docs/src/main/asciidoc/spring-cloud-circuitbreaker-resilience4j.adoc b/docs/src/main/asciidoc/spring-cloud-circuitbreaker-resilience4j.adoc index cf05d67..4e2853e 100644 --- a/docs/src/main/asciidoc/spring-cloud-circuitbreaker-resilience4j.adoc +++ b/docs/src/main/asciidoc/spring-cloud-circuitbreaker-resilience4j.adoc @@ -94,6 +94,39 @@ public Customizer slowCusomtizer() { ---- ==== +==== Circuit Breaker Properties Configuration + +You can configure `CircuitBreaker` and `TimeLimiter` instances in your application's configuration properties file. +Property configuration has higher priority than Java `Customizer` configuration. + +==== +[source] +---- +resilience4j.circuitbreaker: + instances: + backendA: + registerHealthIndicator: true + slidingWindowSize: 100 + backendB: + registerHealthIndicator: true + slidingWindowSize: 10 + permittedNumberOfCallsInHalfOpenState: 3 + slidingWindowType: TIME_BASED + recordFailurePredicate: io.github.robwin.exception.RecordFailurePredicate + +resilience4j.timelimiter: + instances: + backendA: + timeoutDuration: 2s + cancelRunningFuture: true + backendB: + timeoutDuration: 1s + cancelRunningFuture: false +---- +==== + +For more information on Resilience4j property configuration, see https://resilience4j.readme.io/docs/getting-started-3#configuration[Resilience4J Spring Boot 2 Configuration]. + ==== Bulkhead pattern supporting If `resilience4j-bulkhead` is on the classpath, Spring Cloud CircuitBreaker will wrap all methods with a Resilience4j Bulkhead. You can disable the Resilience4j Bulkhead by setting `spring.cloud.circuitbreaker.bulkhead.resilience4j.enabled` to `false`. 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 a09184b..768268d 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 @@ -1,5 +1,5 @@ /* - * Copyright 2013-2019 the original author or authors. + * Copyright 2013-2021 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. @@ -21,7 +21,9 @@ import java.util.List; import javax.annotation.PostConstruct; +import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; import io.github.resilience4j.micrometer.tagged.TaggedCircuitBreakerMetrics; +import io.github.resilience4j.timelimiter.TimeLimiterRegistry; import io.micrometer.core.instrument.MeterRegistry; import org.springframework.beans.factory.annotation.Autowired; @@ -37,6 +39,7 @@ import org.springframework.context.annotation.Configuration; /** * @author Ryan Baxter * @author Eric Bussieres + * @author Thomas Vitale */ @Configuration(proxyBeanMethods = false) @ConditionalOnClass(name = { "reactor.core.publisher.Mono", "reactor.core.publisher.Flux", @@ -50,8 +53,11 @@ public class ReactiveResilience4JAutoConfiguration { @Bean @ConditionalOnMissingBean(ReactiveCircuitBreakerFactory.class) - public ReactiveResilience4JCircuitBreakerFactory reactiveResilience4JCircuitBreakerFactory() { - ReactiveResilience4JCircuitBreakerFactory factory = new ReactiveResilience4JCircuitBreakerFactory(); + public ReactiveResilience4JCircuitBreakerFactory reactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry circuitBreakerRegistry, + TimeLimiterRegistry timeLimiterRegistry) { + ReactiveResilience4JCircuitBreakerFactory factory = new ReactiveResilience4JCircuitBreakerFactory( + circuitBreakerRegistry, timeLimiterRegistry); customizers.forEach(customizer -> customizer.customize(factory)); return 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 bd50961..05a1d1a 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-2019 the original author or authors. + * Copyright 2013-2021 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. @@ -24,6 +24,9 @@ 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 io.github.resilience4j.timelimiter.TimeLimiter; +import io.github.resilience4j.timelimiter.TimeLimiterConfig; +import io.github.resilience4j.timelimiter.TimeLimiterRegistry; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -32,42 +35,66 @@ import org.springframework.cloud.client.circuitbreaker.ReactiveCircuitBreaker; /** * @author Ryan Baxter + * @author Thomas Vitale */ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreaker { - private String id; + private final String id; - private Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config; + private final io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig; - private CircuitBreakerRegistry registry; + private final CircuitBreakerRegistry circuitBreakerRegistry; - private Optional> circuitBreakerCustomizer; + private final TimeLimiterConfig timeLimiterConfig; + private final TimeLimiterRegistry timeLimiterRegistry; + + private final Optional> circuitBreakerCustomizer; + + @Deprecated public ReactiveResilience4JCircuitBreaker(String id, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config, CircuitBreakerRegistry circuitBreakerRegistry, Optional> circuitBreakerCustomizer) { this.id = id; - this.config = config; - this.registry = circuitBreakerRegistry; + this.circuitBreakerConfig = config.getCircuitBreakerConfig(); + this.circuitBreakerRegistry = circuitBreakerRegistry; this.circuitBreakerCustomizer = circuitBreakerCustomizer; + this.timeLimiterConfig = config.getTimeLimiterConfig(); + this.timeLimiterRegistry = TimeLimiterRegistry.ofDefaults(); + } + + public ReactiveResilience4JCircuitBreaker(String id, + Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config, + CircuitBreakerRegistry circuitBreakerRegistry, + TimeLimiterRegistry timeLimiterRegistry, + Optional> circuitBreakerCustomizer) { + this.id = id; + this.circuitBreakerConfig = config.getCircuitBreakerConfig(); + this.circuitBreakerRegistry = circuitBreakerRegistry; + this.circuitBreakerCustomizer = circuitBreakerCustomizer; + this.timeLimiterConfig = config.getTimeLimiterConfig(); + this.timeLimiterRegistry = timeLimiterRegistry; } @Override public Mono run(Mono toRun, Function> fallback) { - io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry - .circuitBreaker(id, config.getCircuitBreakerConfig()); + io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = circuitBreakerRegistry + .circuitBreaker(id, circuitBreakerConfig); circuitBreakerCustomizer .ifPresent(customizer -> customizer.customize(defaultCircuitBreaker)); + TimeLimiter timeLimiter = timeLimiterRegistry.timeLimiter(id, timeLimiterConfig); Mono toReturn = toRun .transform(CircuitBreakerOperator.of(defaultCircuitBreaker)) - .timeout(config.getTimeLimiterConfig().getTimeoutDuration()) + .timeout(timeLimiter.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(), TimeUnit.MILLISECONDS, - t)); + t -> defaultCircuitBreaker + .onError( + timeLimiter.getTimeLimiterConfig() + .getTimeoutDuration().toMillis(), + TimeUnit.MILLISECONDS, t)); if (fallback != null) { toReturn = toReturn.onErrorResume(fallback); } @@ -75,19 +102,22 @@ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreake } public Flux run(Flux toRun, Function> fallback) { - io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry - .circuitBreaker(id, config.getCircuitBreakerConfig()); + io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = circuitBreakerRegistry + .circuitBreaker(id, circuitBreakerConfig); circuitBreakerCustomizer .ifPresent(customizer -> customizer.customize(defaultCircuitBreaker)); + TimeLimiter timeLimiter = timeLimiterRegistry.timeLimiter(id, timeLimiterConfig); Flux toReturn = toRun .transform(CircuitBreakerOperator.of(defaultCircuitBreaker)) - .timeout(config.getTimeLimiterConfig().getTimeoutDuration()) + .timeout(timeLimiter.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(), TimeUnit.MILLISECONDS, - t)); + t -> defaultCircuitBreaker + .onError( + timeLimiter.getTimeLimiterConfig() + .getTimeoutDuration().toMillis(), + TimeUnit.MILLISECONDS, t)); if (fallback != null) { toReturn = toReturn.onErrorResume(fallback); } 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 d2a31d1..6cc833f 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-2019 the original author or authors. + * Copyright 2013-2021 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. @@ -25,6 +25,7 @@ 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 io.github.resilience4j.timelimiter.TimeLimiterRegistry; import org.springframework.cloud.client.circuitbreaker.Customizer; import org.springframework.cloud.client.circuitbreaker.ReactiveCircuitBreaker; @@ -33,25 +34,44 @@ import org.springframework.util.Assert; /** * @author Ryan Baxter + * @author Thomas Vitale */ public class ReactiveResilience4JCircuitBreakerFactory extends ReactiveCircuitBreakerFactory { - private Function defaultConfiguration = id -> new Resilience4JConfigBuilder( - id).circuitBreakerConfig(CircuitBreakerConfig.ofDefaults()) - .timeLimiterConfig(TimeLimiterConfig.ofDefaults()).build(); + private Function defaultConfiguration; private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry .ofDefaults(); + private TimeLimiterRegistry timeLimiterRegistry = TimeLimiterRegistry.ofDefaults(); + private Map> circuitBreakerCustomizers = new HashMap<>(); + @Deprecated + public ReactiveResilience4JCircuitBreakerFactory() { + this.defaultConfiguration = id -> new Resilience4JConfigBuilder(id) + .circuitBreakerConfig(CircuitBreakerConfig.ofDefaults()) + .timeLimiterConfig(TimeLimiterConfig.ofDefaults()).build(); + } + + public ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry circuitBreakerRegistry, + TimeLimiterRegistry timeLimiterRegistry) { + this.circuitBreakerRegistry = circuitBreakerRegistry; + this.timeLimiterRegistry = timeLimiterRegistry; + this.defaultConfiguration = id -> new Resilience4JConfigBuilder(id) + .circuitBreakerConfig(this.circuitBreakerRegistry.getDefaultConfig()) + .timeLimiterConfig(this.timeLimiterRegistry.getDefaultConfig()).build(); + } + @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, + timeLimiterRegistry, Optional.ofNullable(circuitBreakerCustomizers.get(id))); } @@ -60,10 +80,14 @@ public class ReactiveResilience4JCircuitBreakerFactory extends return new Resilience4JConfigBuilder(id); } - CircuitBreakerRegistry getCircuitBreakerRegistry() { + public CircuitBreakerRegistry getCircuitBreakerRegistry() { return circuitBreakerRegistry; } + public TimeLimiterRegistry getTimeLimiterRegistry() { + return timeLimiterRegistry; + } + @Override public void configureDefault( Function defaultConfiguration) { diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationPropertyTest.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationPropertyTest.java new file mode 100644 index 0000000..fdc067d --- /dev/null +++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationPropertyTest.java @@ -0,0 +1,88 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * 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.resilience4j; + +import java.time.Duration; + +import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; +import io.github.resilience4j.timelimiter.TimeLimiterRegistry; +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; +import org.springframework.test.context.ActiveProfiles; +import org.springframework.test.context.junit4.SpringRunner; + +import static org.assertj.core.api.AssertionsForClassTypes.assertThat; + +/** + * @author Thomas Vitale + */ +@RunWith(SpringRunner.class) +@SpringBootTest(classes = ReactiveResilience4JAutoConfigurationPropertyTest.class) +@EnableAutoConfiguration +@ActiveProfiles(profiles = "test-properties") +public class ReactiveResilience4JAutoConfigurationPropertyTest { + + @Autowired + ReactiveResilience4JCircuitBreakerFactory factory; + + @Test + public void testCircuitBreakerPropertiesPopulated() { + CircuitBreakerRegistry circuitBreakerRegistry = factory + .getCircuitBreakerRegistry(); + assertThat(circuitBreakerRegistry).isNotNull(); + assertThat(circuitBreakerRegistry.find("test_circuit")).isPresent(); + assertThat(circuitBreakerRegistry.find("test_circuit").get() + .getCircuitBreakerConfig().getMinimumNumberOfCalls()).isEqualTo(5); + } + + @Test + public void testTimeLimiterPropertiesPopulated() { + TimeLimiterRegistry timeLimiterRegistry = factory.getTimeLimiterRegistry(); + assertThat(timeLimiterRegistry).isNotNull(); + assertThat(timeLimiterRegistry.find("test_circuit")).isPresent(); + assertThat(timeLimiterRegistry.find("test_circuit").get().getTimeLimiterConfig() + .getTimeoutDuration()).isEqualTo(Duration.ofSeconds(18)); + } + + @Test + public void testDefaultCircuitBreakerPropertiesPopulated() { + factory.create("default_circuitBreaker").run(Mono.just("result")); + CircuitBreakerRegistry circuitBreakerRegistry = factory + .getCircuitBreakerRegistry(); + assertThat(circuitBreakerRegistry).isNotNull(); + assertThat(circuitBreakerRegistry.find("default_circuitBreaker")).isPresent(); + assertThat(circuitBreakerRegistry.find("default_circuitBreaker").get() + .getCircuitBreakerConfig().getMinimumNumberOfCalls()).isEqualTo(20); + } + + @Test + public void testDefaultTimeLimiterPropertiesPopulated() { + factory.create("default_circuitBreaker").run(Mono.just("result")); + TimeLimiterRegistry timeLimiterRegistry = factory.getTimeLimiterRegistry(); + assertThat(timeLimiterRegistry).isNotNull(); + assertThat(timeLimiterRegistry.find("default_circuitBreaker")).isPresent(); + assertThat(timeLimiterRegistry.find("default_circuitBreaker").get() + .getTimeLimiterConfig().getTimeoutDuration()) + .isEqualTo(Duration.ofMillis(150)); + } + +} diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationWithoutMetricsTest.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationWithoutMetricsTest.java index b4d7ebe..a4cac1c 100644 --- a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationWithoutMetricsTest.java +++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationWithoutMetricsTest.java @@ -1,5 +1,5 @@ /* - * Copyright 2013-2019 the original author or authors. + * Copyright 2013-2021 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. @@ -16,6 +16,8 @@ package org.springframework.cloud.circuitbreaker.resilience4j; +import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; +import io.github.resilience4j.timelimiter.TimeLimiterRegistry; import org.junit.Test; import org.junit.runner.RunWith; @@ -34,13 +36,16 @@ import static org.mockito.Mockito.verify; /** * @author Ryan Baxter + * @author Thomas Vitale */ @RunWith(ModifiedClassPathRunner.class) @ClassPathExclusions({ "micrometer-core-*.jar", "resilience4j-micrometer-*.jar" }) public class ReactiveResilience4JAutoConfigurationWithoutMetricsTest { static ReactiveResilience4JCircuitBreakerFactory circuitBreakerFactory = spy( - new ReactiveResilience4JCircuitBreakerFactory()); + new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults())); @Test public void testWithoutMetrics() { 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 ca7ce0b..4726e70 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 @@ -1,5 +1,5 @@ /* - * Copyright 2013-2018 the original author or authors. + * Copyright 2013-2021 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. @@ -17,7 +17,10 @@ package org.springframework.cloud.circuitbreaker.resilience4j; import java.util.Arrays; +import java.util.Collections; +import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; +import io.github.resilience4j.timelimiter.TimeLimiterRegistry; import org.junit.Test; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -28,21 +31,23 @@ import static org.assertj.core.api.Assertions.assertThat; /** * @author Ryan Baxter + * @author Thomas Vitale */ public class ReactiveResilience4JCircuitBreakerTest { @Test public void runMono() { - ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory() - .create("foo"); - assertThat(Mono.just("foobar").transform(it -> cb.run(it)).block()) - .isEqualTo("foobar"); + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), TimeLimiterRegistry.ofDefaults()) + .create("foo"); + assertThat(Mono.just("foobar").transform(cb::run).block()).isEqualTo("foobar"); } @Test public void runMonoWithFallback() { - ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory() - .create("foo"); + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), TimeLimiterRegistry.ofDefaults()) + .create("foo"); assertThat(Mono.error(new RuntimeException("boom")) .transform(it -> cb.run(it, t -> Mono.just("fallback"))).block()) .isEqualTo("fallback"); @@ -50,19 +55,21 @@ public class ReactiveResilience4JCircuitBreakerTest { @Test public void runFlux() { - ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory() - .create("foo"); - assertThat(Flux.just("foobar", "hello world").transform(it -> cb.run(it)) - .collectList().block()).isEqualTo(Arrays.asList("foobar", "hello world")); + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), TimeLimiterRegistry.ofDefaults()) + .create("foo"); + assertThat(Flux.just("foobar", "hello world").transform(cb::run).collectList() + .block()).isEqualTo(Arrays.asList("foobar", "hello world")); } @Test public void runFluxWithFallback() { - ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory() - .create("foo"); + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), TimeLimiterRegistry.ofDefaults()) + .create("foo"); assertThat(Flux.error(new RuntimeException("boom")) .transform(it -> cb.run(it, t -> Flux.just("fallback"))).collectList() - .block()).isEqualTo(Arrays.asList("fallback")); + .block()).isEqualTo(Collections.singletonList("fallback")); } }