Resolve merge conflicts and cherry pick #108
This commit is contained in:
@@ -94,6 +94,39 @@ public Customizer<ReactiveResilience4JCircuitBreakerFactory> 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`.
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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<Customizer<CircuitBreaker>> circuitBreakerCustomizer;
|
||||
private final TimeLimiterConfig timeLimiterConfig;
|
||||
|
||||
private final TimeLimiterRegistry timeLimiterRegistry;
|
||||
|
||||
private final Optional<Customizer<CircuitBreaker>> circuitBreakerCustomizer;
|
||||
|
||||
@Deprecated
|
||||
public ReactiveResilience4JCircuitBreaker(String id,
|
||||
Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config,
|
||||
CircuitBreakerRegistry circuitBreakerRegistry,
|
||||
Optional<Customizer<CircuitBreaker>> 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<Customizer<CircuitBreaker>> circuitBreakerCustomizer) {
|
||||
this.id = id;
|
||||
this.circuitBreakerConfig = config.getCircuitBreakerConfig();
|
||||
this.circuitBreakerRegistry = circuitBreakerRegistry;
|
||||
this.circuitBreakerCustomizer = circuitBreakerCustomizer;
|
||||
this.timeLimiterConfig = config.getTimeLimiterConfig();
|
||||
this.timeLimiterRegistry = timeLimiterRegistry;
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> Mono<T> run(Mono<T> toRun, Function<Throwable, Mono<T>> fallback) {
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry
|
||||
.circuitBreaker(id, config.getCircuitBreakerConfig());
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = circuitBreakerRegistry
|
||||
.circuitBreaker(id, circuitBreakerConfig);
|
||||
circuitBreakerCustomizer
|
||||
.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
TimeLimiter timeLimiter = timeLimiterRegistry.timeLimiter(id, timeLimiterConfig);
|
||||
Mono<T> 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 <T> Flux<T> run(Flux<T> toRun, Function<Throwable, Flux<T>> 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<T> 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);
|
||||
}
|
||||
|
||||
@@ -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<Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration, Resilience4JConfigBuilder> {
|
||||
|
||||
private Function<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration = id -> new Resilience4JConfigBuilder(
|
||||
id).circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.ofDefaults()).build();
|
||||
private Function<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration;
|
||||
|
||||
private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry
|
||||
.ofDefaults();
|
||||
|
||||
private TimeLimiterRegistry timeLimiterRegistry = TimeLimiterRegistry.ofDefaults();
|
||||
|
||||
private Map<String, Customizer<CircuitBreaker>> 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<String, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration> defaultConfiguration) {
|
||||
|
||||
@@ -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));
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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() {
|
||||
|
||||
@@ -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"));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user