diff --git a/docs/modules/ROOT/pages/spring-cloud-circuitbreaker-resilience4j/bulkhead-pattern-supporting.adoc b/docs/modules/ROOT/pages/spring-cloud-circuitbreaker-resilience4j/bulkhead-pattern-supporting.adoc index ba915f2..cbccea7 100644 --- a/docs/modules/ROOT/pages/spring-cloud-circuitbreaker-resilience4j/bulkhead-pattern-supporting.adoc +++ b/docs/modules/ROOT/pages/spring-cloud-circuitbreaker-resilience4j/bulkhead-pattern-supporting.adoc @@ -28,3 +28,41 @@ public Customizer defaultBulkheadCustomizer() { } ---- +== Reactive Bulkhead Pattern Supporting + +If you are using reactive programming with Spring Cloud CircuitBreaker, you can leverage the `ReactiveResilience4jBulkheadProvider` to support the Bulkhead pattern in reactive pipelines. +This provider decorates `Mono` and `Flux` instances to ensure bulkhead constraints are applied during reactive operations. + +Spring Cloud CircuitBreaker Resilience4j reactive support only uses the `SemaphoreBulkhead`. +If the property `spring.cloud.circuitbreaker.resilience4j.enableSemaphoreDefaultBulkhead` is set to `false`, a warning will be logged, and the `ReactiveResilience4jBulkheadProvider` will still use the `SemaphoreBulkhead`. + +== Configuring Reactive Bulkhead + +The `ReactiveResilience4jBulkheadProvider` can be customized using a `Customizer` bean, as shown below: + +[source,java] +---- +@Bean +public Customizer reactiveBulkheadCustomizer() { + return provider -> provider.configureDefault(id -> new Resilience4jBulkheadConfigurationBuilder() + .bulkheadConfig(BulkheadConfig.custom().maxConcurrentCalls(4).build()) + .build()); +} +---- + +You can also add individual bulkhead configurations for specific use cases: + +[source,java] +---- +@Bean +public Customizer reactiveSpecificBulkheadCustomizer() { + return provider -> provider.configure(builder -> { + builder.bulkheadConfig(BulkheadConfig.custom() + .maxConcurrentCalls(2) + .build()); + }, "serviceBulkhead"); +} +---- + +For more details, see the https://resilience4j.readme.io/docs/examples-1#decorate-mono-or-flux-with-a-bulkhead[Resilience4j Reactive Bulkhead Examples]. + 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 58c0a9a..af4a45b 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 @@ -19,14 +19,19 @@ package org.springframework.cloud.circuitbreaker.resilience4j; import java.util.ArrayList; import java.util.List; +import io.github.resilience4j.bulkhead.Bulkhead; +import io.github.resilience4j.bulkhead.BulkheadRegistry; import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; +import io.github.resilience4j.micrometer.tagged.TaggedBulkheadMetrics; import io.github.resilience4j.micrometer.tagged.TaggedCircuitBreakerMetrics; import io.github.resilience4j.micrometer.tagged.TaggedCircuitBreakerMetricsPublisher; import io.github.resilience4j.timelimiter.TimeLimiterRegistry; import io.micrometer.core.instrument.MeterRegistry; import jakarta.annotation.PostConstruct; +import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; @@ -41,6 +46,7 @@ import org.springframework.context.annotation.Configuration; * @author Ryan Baxter * @author Eric Bussieres * @author Thomas Vitale + * @author Yavor Chamov */ @Configuration(proxyBeanMethods = false) @ConditionalOnClass(name = { "reactor.core.publisher.Mono", "reactor.core.publisher.Flux", @@ -57,24 +63,55 @@ public class ReactiveResilience4JAutoConfiguration { @ConditionalOnMissingBean(ReactiveCircuitBreakerFactory.class) public ReactiveResilience4JCircuitBreakerFactory reactiveResilience4JCircuitBreakerFactory( CircuitBreakerRegistry circuitBreakerRegistry, TimeLimiterRegistry timeLimiterRegistry, + @Autowired(required = false) ReactiveResilience4jBulkheadProvider bulkheadProvider, Resilience4JConfigurationProperties resilience4JConfigurationProperties) { ReactiveResilience4JCircuitBreakerFactory factory = new ReactiveResilience4JCircuitBreakerFactory( - circuitBreakerRegistry, timeLimiterRegistry, resilience4JConfigurationProperties); + circuitBreakerRegistry, timeLimiterRegistry, bulkheadProvider, resilience4JConfigurationProperties); customizers.forEach(customizer -> customizer.customize(factory)); return factory; } @Configuration(proxyBeanMethods = false) - @ConditionalOnClass(name = { "reactor.core.publisher.Mono", "reactor.core.publisher.Flux", + @ConditionalOnClass(Bulkhead.class) + @ConditionalOnProperty(value = "spring.cloud.circuitbreaker.bulkhead.resilience4j.enabled", matchIfMissing = true) + public static class Resilience4jBulkheadConfiguration { + + @Autowired(required = false) + private List> bulkheadCustomizers = new ArrayList<>(); + + @Value("${spring.cloud.circuitbreaker.resilience4j.enableSemaphoreDefaultBulkhead:true}") + private boolean enableSemaphoreDefaultBulkhead; + + @Bean + public ReactiveResilience4jBulkheadProvider reactiveBulkheadProvider(BulkheadRegistry bulkheadRegistry) { + + if (!enableSemaphoreDefaultBulkhead) { + LoggerFactory.getLogger(Resilience4jBulkheadConfiguration.class) + .warn("Ignoring 'spring.cloud.circuitbreaker.resilience4j.enableSemaphoreDefaultBulkhead=false'. " + + "ReactiveResilience4jBulkheadProvider only supports SemaphoreBulkhead."); + } + + ReactiveResilience4jBulkheadProvider reactiveResilience4JCircuitBreaker = + new ReactiveResilience4jBulkheadProvider(bulkheadRegistry); + bulkheadCustomizers.forEach(customizer -> customizer.customize(reactiveResilience4JCircuitBreaker)); + return reactiveResilience4JCircuitBreaker; + } + } + + @Configuration(proxyBeanMethods = false) + @ConditionalOnClass(name = {"reactor.core.publisher.Mono", "reactor.core.publisher.Flux", "io.github.resilience4j.micrometer.tagged.TaggedCircuitBreakerMetrics", - "io.github.resilience4j.micrometer.tagged.TaggedCircuitBreakerMetricsPublisher" }) - @ConditionalOnBean({ MeterRegistry.class }) - @ConditionalOnMissingBean({ TaggedCircuitBreakerMetricsPublisher.class }) + "io.github.resilience4j.micrometer.tagged.TaggedCircuitBreakerMetricsPublisher"}) + @ConditionalOnBean({MeterRegistry.class}) + @ConditionalOnMissingBean({TaggedCircuitBreakerMetricsPublisher.class}) public static class MicrometerReactiveResilience4JCustomizerConfiguration { @Autowired(required = false) private ReactiveResilience4JCircuitBreakerFactory factory; + @Autowired(required = false) + private ReactiveResilience4jBulkheadProvider bulkheadProvider; + @Autowired(required = false) private TaggedCircuitBreakerMetrics taggedCircuitBreakerMetrics; @@ -90,6 +127,9 @@ public class ReactiveResilience4JAutoConfiguration { } taggedCircuitBreakerMetrics.bindTo(meterRegistry); } + if (bulkheadProvider != null) { + TaggedBulkheadMetrics.ofBulkheadRegistry(bulkheadProvider.getBulkheadRegistry()).bindTo(meterRegistry); + } } } 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 f4966dd..db6c9cb 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 @@ -43,6 +43,7 @@ import static org.springframework.cloud.circuitbreaker.resilience4j.Resilience4J * @author Ryan Baxter * @author Thomas Vitale * @author 荒 + * @author Yavor Chamov */ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreaker { @@ -50,6 +51,8 @@ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreake private final String groupName; + private final ReactiveResilience4jBulkheadProvider bulkheadProvider; + private final io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig; private final CircuitBreakerRegistry circuitBreakerRegistry; @@ -70,10 +73,19 @@ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreake this(id, groupName, config, circuitBreakerRegistry, timeLimiterRegistry, circuitBreakerCustomizer, false); } + @Deprecated public ReactiveResilience4JCircuitBreaker(String id, String groupName, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config, CircuitBreakerRegistry circuitBreakerRegistry, TimeLimiterRegistry timeLimiterRegistry, Optional> circuitBreakerCustomizer, boolean disableTimeLimiter) { + this(id, groupName, config, circuitBreakerRegistry, timeLimiterRegistry, circuitBreakerCustomizer, null, disableTimeLimiter); + } + + public ReactiveResilience4JCircuitBreaker(String id, String groupName, + Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config, + CircuitBreakerRegistry circuitBreakerRegistry, TimeLimiterRegistry timeLimiterRegistry, + Optional> circuitBreakerCustomizer, + ReactiveResilience4jBulkheadProvider bulkheadProvider, boolean disableTimeLimiter) { this.id = id; this.groupName = groupName; this.circuitBreakerConfig = config.getCircuitBreakerConfig(); @@ -81,13 +93,22 @@ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreake this.circuitBreakerCustomizer = circuitBreakerCustomizer; this.timeLimiterConfig = config.getTimeLimiterConfig(); this.timeLimiterRegistry = timeLimiterRegistry; + this.bulkheadProvider = bulkheadProvider; this.disableTimeLimiter = disableTimeLimiter; } @Override public Mono run(Mono toRun, Function> fallback) { + final Map tags = Map.of(CIRCUIT_BREAKER_GROUP_TAG, this.groupName); Tuple2> tuple = buildCircuitBreakerAndTimeLimiter(); - Mono toReturn = toRun.transform(CircuitBreakerOperator.of(tuple.getT1())); + Mono toReturn; + if (bulkheadProvider != null) { + toReturn = bulkheadProvider.decorateMono(groupName, tags, toRun); + } + else { + toReturn = toRun; + } + toReturn = toReturn.transform(CircuitBreakerOperator.of(tuple.getT1())); if (tuple.getT2().isPresent()) { final Duration timeoutDuration = tuple.getT2().get().getTimeLimiterConfig().getTimeoutDuration(); toReturn = toReturn.timeout(timeoutDuration) @@ -105,8 +126,16 @@ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreake @Override public Flux run(Flux toRun, Function> fallback) { + final Map tags = Map.of(CIRCUIT_BREAKER_GROUP_TAG, this.groupName); Tuple2> tuple = buildCircuitBreakerAndTimeLimiter(); - Flux toReturn = toRun.transform(CircuitBreakerOperator.of(tuple.getT1())); + Flux toReturn; + if (bulkheadProvider != null) { + toReturn = bulkheadProvider.decorateFlux(groupName, tags, toRun); + } + else { + toReturn = toRun; + } + toReturn = toReturn.transform(CircuitBreakerOperator.of(tuple.getT1())); if (tuple.getT2().isPresent()) { final Duration timeoutDuration = tuple.getT2().get().getTimeLimiterConfig().getTimeoutDuration(); toReturn = toReturn.timeout(timeoutDuration) 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 de3501a..03e78ad 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 @@ -36,10 +36,13 @@ import org.springframework.util.Assert; * @author Ryan Baxter * @author Thomas Vitale * @author 荒 + * @author Yavor Chamov */ public class ReactiveResilience4JCircuitBreakerFactory extends ReactiveCircuitBreakerFactory { + private final ReactiveResilience4jBulkheadProvider bulkheadProvider; + private Function defaultConfiguration; private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry.ofDefaults(); @@ -53,13 +56,21 @@ public class ReactiveResilience4JCircuitBreakerFactory extends @Deprecated public ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry circuitBreakerRegistry, TimeLimiterRegistry timeLimiterRegistry) { - this(circuitBreakerRegistry, timeLimiterRegistry, null); + this(circuitBreakerRegistry, timeLimiterRegistry, null, null); } + @Deprecated public ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry circuitBreakerRegistry, TimeLimiterRegistry timeLimiterRegistry, Resilience4JConfigurationProperties resilience4JConfigurationProperties) { + this(circuitBreakerRegistry, timeLimiterRegistry, null, resilience4JConfigurationProperties); + } + + public ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry circuitBreakerRegistry, + TimeLimiterRegistry timeLimiterRegistry, ReactiveResilience4jBulkheadProvider bulkheadProvider, + Resilience4JConfigurationProperties resilience4JConfigurationProperties) { this.circuitBreakerRegistry = circuitBreakerRegistry; + this.bulkheadProvider = bulkheadProvider; this.timeLimiterRegistry = timeLimiterRegistry; this.defaultConfiguration = id -> new Resilience4JConfigBuilder(id) .circuitBreakerConfig(this.circuitBreakerRegistry.getDefaultConfig()) @@ -106,7 +117,8 @@ public class ReactiveResilience4JCircuitBreakerFactory extends boolean isDisableTimeLimiter = ConfigurationPropertiesUtils .isDisableTimeLimiter(this.resilience4JConfigurationProperties, id, groupName); return new ReactiveResilience4JCircuitBreaker(id, groupName, config, circuitBreakerRegistry, - timeLimiterRegistry, Optional.ofNullable(circuitBreakerCustomizers.get(id)), isDisableTimeLimiter); + timeLimiterRegistry, Optional.ofNullable(circuitBreakerCustomizers.get(id)), + bulkheadProvider, isDisableTimeLimiter); } @Override diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4jBulkheadProvider.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4jBulkheadProvider.java new file mode 100644 index 0000000..dcf9fcd --- /dev/null +++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4jBulkheadProvider.java @@ -0,0 +1,104 @@ +/* + * 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.util.Map; +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Consumer; +import java.util.function.Function; + +import io.github.resilience4j.bulkhead.Bulkhead; +import io.github.resilience4j.bulkhead.BulkheadConfig; +import io.github.resilience4j.bulkhead.BulkheadRegistry; +import io.github.resilience4j.reactor.bulkhead.operator.BulkheadOperator; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +import org.springframework.lang.NonNull; +import org.springframework.util.Assert; + +/** + * @author Yavor Chamov + */ +public class ReactiveResilience4jBulkheadProvider { + + private final BulkheadRegistry bulkheadRegistry; + + private final ConcurrentHashMap + configurations = new ConcurrentHashMap<>(); + + private Function defaultConfiguration; + + public ReactiveResilience4jBulkheadProvider(BulkheadRegistry bulkheadRegistry) { + this.bulkheadRegistry = bulkheadRegistry; + this.defaultConfiguration = id -> new Resilience4jBulkheadConfigurationBuilder() + .bulkheadConfig(this.bulkheadRegistry.getDefaultConfig()) + .build(); + } + + public void configureDefault( + @NonNull Function defaultConfiguration) { + Assert.notNull(defaultConfiguration, "Default configuration must not be null"); + this.defaultConfiguration = defaultConfiguration; + } + + public void configure(Consumer consumer, String... ids) { + for (String id : ids) { + Resilience4jBulkheadConfigurationBuilder builder = new Resilience4jBulkheadConfigurationBuilder(); + consumer.accept(builder); + Resilience4jBulkheadConfigurationBuilder.BulkheadConfiguration configuration = builder.build(); + configurations.put(id, configuration); + } + } + + public void addBulkheadCustomizer(Consumer customizer, String... ids) { + for (String id : ids) { + Resilience4jBulkheadConfigurationBuilder.BulkheadConfiguration configuration = configurations + .computeIfAbsent(id, defaultConfiguration); + Bulkhead bulkhead = bulkheadRegistry.bulkhead(id, configuration.getBulkheadConfig()); + customizer.accept(bulkhead); + } + } + + public BulkheadRegistry getBulkheadRegistry() { + return bulkheadRegistry; + } + + public Mono decorateMono(String id, Map tags, Mono mono) { + Resilience4jBulkheadConfigurationBuilder.BulkheadConfiguration configuration = configurations + .computeIfAbsent(id, this::getConfiguration); + Bulkhead bulkhead = bulkheadRegistry.bulkhead(id, configuration.getBulkheadConfig(), tags); + return mono.transformDeferred(BulkheadOperator.of(bulkhead)); + } + + public Flux decorateFlux(String id, Map tags, Flux flux) { + Resilience4jBulkheadConfigurationBuilder.BulkheadConfiguration configuration = configurations + .computeIfAbsent(id, this::getConfiguration); + Bulkhead bulkhead = bulkheadRegistry.bulkhead(id, configuration.getBulkheadConfig(), tags); + return flux.transformDeferred(BulkheadOperator.of(bulkhead)); + } + + private Resilience4jBulkheadConfigurationBuilder.BulkheadConfiguration getConfiguration(String id) { + Resilience4jBulkheadConfigurationBuilder builder = new Resilience4jBulkheadConfigurationBuilder(); + Resilience4jBulkheadConfigurationBuilder.BulkheadConfiguration defaultConfiguration = + this.defaultConfiguration.apply(id); + Optional bulkheadConfiguration = bulkheadRegistry.getConfiguration(id); + builder.bulkheadConfig(bulkheadConfiguration.orElse(defaultConfiguration.getBulkheadConfig())); + return builder.build(); + } +} diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/BulkheadConfigurationTests.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/BulkheadConfigurationTests.java index 05ae3fb..67c7901 100644 --- a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/BulkheadConfigurationTests.java +++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/BulkheadConfigurationTests.java @@ -16,12 +16,18 @@ package org.springframework.cloud.circuitbreaker.resilience4j; +import java.time.Duration; +import java.util.Optional; + +import io.github.resilience4j.bulkhead.Bulkhead; import io.github.resilience4j.bulkhead.BulkheadConfig; import io.github.resilience4j.bulkhead.BulkheadRegistry; import io.github.resilience4j.bulkhead.ThreadPoolBulkhead; import io.github.resilience4j.bulkhead.ThreadPoolBulkheadConfig; import io.github.resilience4j.bulkhead.ThreadPoolBulkheadRegistry; import org.junit.jupiter.api.Test; +import reactor.core.publisher.Mono; +import reactor.test.StepVerifier; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.test.context.runner.WebApplicationContextRunner; @@ -31,6 +37,7 @@ import static org.assertj.core.api.Assertions.assertThat; /** * @author Ryan Baxter + * @author Yavor Chamov */ public class BulkheadConfigurationTests { @@ -103,6 +110,43 @@ public class BulkheadConfigurationTests { } + @Test + void testReactiveInstanceConfigurationOverridesConfigAndCustomizerProperties() { + new WebApplicationContextRunner() + .withUserConfiguration(Application.class) + .withPropertyValues( + "resilience4j.bulkhead.configs.testme.max-concurrent-calls=30", + "resilience4j.bulkhead.instances.testme.max-concurrent-calls=40" + ) + .run(context -> { + final String id = "testme"; + + ReactiveResilience4JCircuitBreakerFactory resilience4JCircuitBreakerFactory = context + .getBean(ReactiveResilience4JCircuitBreakerFactory.class); + + ReactiveResilience4jBulkheadProvider bulkheadProvider = context.getBean(ReactiveResilience4jBulkheadProvider.class); + + bulkheadProvider.configure(builder -> { + BulkheadConfig bulkheadConfig = BulkheadConfig.custom() + .maxConcurrentCalls(50) + .build(); + builder.bulkheadConfig(bulkheadConfig); + }, id); + + BulkheadRegistry bulkheadRegistry = context.getBean(BulkheadRegistry.class); + + Mono result = resilience4JCircuitBreakerFactory.create(id).run(Mono.just("test")); + + StepVerifier.create(result) + .expectNext("test") + .verifyComplete(); + + Optional bulkheadOptional = bulkheadRegistry.find(id); + assertThat(bulkheadOptional).isPresent(); + assertThat(bulkheadOptional.get().getBulkheadConfig().getMaxConcurrentCalls()).isEqualTo(40); + }); + } + @Test void testCustomizerConfigurationOverridesConfigProperties() { new WebApplicationContextRunner().withUserConfiguration(Application.class) @@ -130,6 +174,42 @@ public class BulkheadConfigurationTests { } + @Test + void testReactiveCustomizerConfigurationOverridesConfigProperties() { + new WebApplicationContextRunner() + .withUserConfiguration(Application.class) + .withPropertyValues( + "resilience4j.bulkhead.configs.testme.max-concurrent-calls=30" + ) + .run(context -> { + final String id = "testme"; + + ReactiveResilience4JCircuitBreakerFactory resilience4JCircuitBreakerFactory = context + .getBean(ReactiveResilience4JCircuitBreakerFactory.class); + + ReactiveResilience4jBulkheadProvider bulkheadProvider = context.getBean(ReactiveResilience4jBulkheadProvider.class); + + bulkheadProvider.configure(builder -> { + BulkheadConfig bulkheadConfig = BulkheadConfig.custom() + .maxConcurrentCalls(50) + .build(); + builder.bulkheadConfig(bulkheadConfig); + }, id); + + BulkheadRegistry bulkheadRegistry = context.getBean(BulkheadRegistry.class); + + Mono result = resilience4JCircuitBreakerFactory.create(id).run(Mono.just("test")); + + StepVerifier.create(result) + .expectNext("test") + .verifyComplete(); + + Optional bulkheadOptional = bulkheadRegistry.find(id); + assertThat(bulkheadOptional).isPresent(); + assertThat(bulkheadOptional.get().getBulkheadConfig().getMaxConcurrentCalls()).isEqualTo(50); + }); + } + @Test void useThreadPoolWithSemaphorePropertySet() { new WebApplicationContextRunner().withUserConfiguration(Application.class) @@ -202,6 +282,36 @@ public class BulkheadConfigurationTests { } + @Test + void configureDefaultOverridesPropertyDefaultForReactiveBulkhead() { + new WebApplicationContextRunner().withUserConfiguration(Application.class) + .withPropertyValues("resilience4j.bulkhead.config.default.max-concurrent-calls=30", + "resilience4j.threadpool.config.default.queueCapacity=30"/* + * , + * "resilience4j.bulkhead.instances.testme.max-wait-duration=30s", + * "resilience4j.threadPoolBulkHead.instances.testme.core-threadpool-size=1" + */) + .run(context -> { + final String id = "testme"; + ReactiveResilience4JCircuitBreakerFactory resilience4JCircuitBreakerFactory = context + .getBean(ReactiveResilience4JCircuitBreakerFactory.class); + ReactiveResilience4jBulkheadProvider bulkheadProvider = context.getBean(ReactiveResilience4jBulkheadProvider.class); + bulkheadProvider.configureDefault(bulkheadId -> new Resilience4jBulkheadConfigurationBuilder() + .bulkheadConfig(BulkheadConfig.custom().maxConcurrentCalls(50).maxWaitDuration(Duration.ofSeconds(10)).build()) + .build()); + BulkheadRegistry bulkheadRegistry = context.getBean(BulkheadRegistry.class); + Mono result = resilience4JCircuitBreakerFactory.create(id).run(Mono.just("test")); + StepVerifier.create(result) + .expectNext("test") + .verifyComplete(); + Optional bulkheadOptional = bulkheadRegistry.find(id); + assertThat(bulkheadOptional).isPresent(); + assertThat(bulkheadOptional.get().getBulkheadConfig().getMaxConcurrentCalls()).isEqualTo(50); + assertThat(bulkheadOptional.get().getBulkheadConfig().getMaxWaitDuration()).isEqualTo(Duration.ofSeconds(10)); + }); + + } + @Test void configureDefaultOverridesPropertyDefaultForSemaphore() { new WebApplicationContextRunner().withUserConfiguration(Application.class) @@ -230,6 +340,39 @@ public class BulkheadConfigurationTests { } + @Test + void configureDefaultOverridesPropertyDefaultForReactiveSemaphore() { + new WebApplicationContextRunner() + .withUserConfiguration(Application.class) + .withPropertyValues( + "resilience4j.bulkhead.config.default.max-concurrent-calls=30", + "spring.cloud.circuitbreaker.resilience4j.enableSemaphoreDefaultBulkhead=true" + ) + .run(context -> { + final String id = "testme"; + + ReactiveResilience4JCircuitBreakerFactory resilience4JCircuitBreakerFactory = context + .getBean(ReactiveResilience4JCircuitBreakerFactory.class); + ReactiveResilience4jBulkheadProvider bulkheadProvider = context.getBean(ReactiveResilience4jBulkheadProvider.class); + + bulkheadProvider.configureDefault(bulkheadId -> new Resilience4jBulkheadConfigurationBuilder() + .bulkheadConfig(BulkheadConfig.custom().maxConcurrentCalls(40).build()) + .build()); + + BulkheadRegistry bulkheadRegistry = context.getBean(BulkheadRegistry.class); + + Mono result = resilience4JCircuitBreakerFactory.create(id).run(Mono.just("test")); + + StepVerifier.create(result) + .expectNext("test") + .verifyComplete(); + + Optional semaphoreBulkhead = bulkheadRegistry.find(id); + assertThat(semaphoreBulkhead).isPresent(); + assertThat(semaphoreBulkhead.get().getBulkheadConfig().getMaxConcurrentCalls()).isEqualTo(40); + }); + } + @Configuration(proxyBeanMethods = false) @EnableAutoConfiguration protected static class Application { diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationWithoutBulkheadTest.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationWithoutBulkheadTest.java new file mode 100644 index 0000000..4c617eb --- /dev/null +++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JAutoConfigurationWithoutBulkheadTest.java @@ -0,0 +1,55 @@ +/* + * 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 org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.boot.SpringBootConfiguration; +import org.springframework.boot.WebApplicationType; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.cloud.test.ClassPathExclusions; +import org.springframework.cloud.test.ModifiedClassPathRunner; +import org.springframework.context.ConfigurableApplicationContext; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * @author Yavor Chamov + */ +@RunWith(ModifiedClassPathRunner.class) +@ClassPathExclusions("resilience4j-bulkhead-*.jar") +public class ReactiveResilience4JAutoConfigurationWithoutBulkheadTest { + + @Test + public void testWithoutBulkhead() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder().web(WebApplicationType.NONE) + .sources(TestApp.class) + .run() + ) { + assertThat(context.containsBean("reactiveBulkheadProvider")).isFalse(); + } + } + + @SpringBootConfiguration + @EnableAutoConfiguration + protected static class TestApp { + + } + +} 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 7d7fd16..97ba357 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 @@ -30,6 +30,7 @@ import org.springframework.cloud.test.ModifiedClassPathRunner; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; +import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -37,14 +38,15 @@ import static org.mockito.Mockito.verify; /** * @author Ryan Baxter * @author Thomas Vitale + * @author Yavor Chamov */ @RunWith(ModifiedClassPathRunner.class) @ClassPathExclusions({ "micrometer-core-*.jar", "resilience4j-micrometer-*.jar" }) public class ReactiveResilience4JAutoConfigurationWithoutMetricsTest { static ReactiveResilience4JCircuitBreakerFactory circuitBreakerFactory = spy( - new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - TimeLimiterRegistry.ofDefaults(), new Resilience4JConfigurationProperties())); + new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults(), new Resilience4JConfigurationProperties())); @Test public void testWithoutMetrics() { @@ -52,6 +54,19 @@ public class ReactiveResilience4JAutoConfigurationWithoutMetricsTest { .sources(TestApp.class) .run()) { verify(circuitBreakerFactory, times(0)).getCircuitBreakerRegistry(); + assertThat(context.containsBean("circuitBreakerFactory")).isTrue(); + assertThat(context.containsBean("reactiveBulkheadProvider")).isTrue(); + } + } + + @Test + public void testProviderCreatedWhenEnableSemaphoreDefaultBulkheadFalse() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder() + .web(WebApplicationType.NONE) + .sources(TestApp.class) + .properties("spring.cloud.circuitbreaker.resilience4j.enableSemaphoreDefaultBulkhead=false") + .run()) { + assertThat(context.containsBean("reactiveBulkheadProvider")).isTrue(); } } diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JBulkheadAndTimeLimiterIntegrationTest.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JBulkheadAndTimeLimiterIntegrationTest.java new file mode 100644 index 0000000..2f7d931 --- /dev/null +++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4JBulkheadAndTimeLimiterIntegrationTest.java @@ -0,0 +1,116 @@ +/* + * Copyright 2013-2023 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 java.util.concurrent.TimeoutException; + +import io.github.resilience4j.bulkhead.BulkheadFullException; +import io.github.resilience4j.timelimiter.TimeLimiterConfig; +import org.junit.Test; +import org.junit.runner.RunWith; +import reactor.core.publisher.Mono; +import reactor.test.StepVerifier; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.cloud.client.circuitbreaker.Customizer; +import org.springframework.cloud.client.circuitbreaker.ReactiveCircuitBreaker; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.stereotype.Service; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.junit4.SpringRunner; + +/** + * @author Yavor Chamov + */ +@RunWith(SpringRunner.class) +@SpringBootTest(classes = ReactiveResilience4JBulkheadAndTimeLimiterIntegrationTest.Application.class) +@DirtiesContext +public class ReactiveResilience4JBulkheadAndTimeLimiterIntegrationTest { + + /** + * Slow bulkhead name. + */ + public static final String SLOW_BULKHEAD = "slowBulkhead"; + + @Autowired + Application.DemoReactiveService service; + + @Test + public void testBulkheadThreadInterrupted() { + + StepVerifier.create(service.bulkheadWithDelay(1000)) + .expectNext(Application.CompletionStatus.INTERRUPTED) + .verifyComplete(); + } + + @Test + public void testBulkheadFastCallNotInterrupted() { + StepVerifier.create(service.bulkheadWithDelay(300)) + .expectNext(Application.CompletionStatus.SUCCESS) + .verifyComplete(); + } + + @Configuration(proxyBeanMethods = false) + @EnableAutoConfiguration + protected static class Application { + + @Bean + public Customizer reactiveSlowBulkheadCustomizer() { + TimeLimiterConfig timeLimiterConfig = TimeLimiterConfig.custom() + .timeoutDuration(Duration.ofMillis(500)) + .build(); + return circuitBreakerFactory -> circuitBreakerFactory + .configure(builder -> builder.timeLimiterConfig(timeLimiterConfig), SLOW_BULKHEAD); + } + + @Bean + public Customizer reactiveBulkheadProviderCustomizer() { + return provider -> provider.addBulkheadCustomizer(bulkhead -> { }, SLOW_BULKHEAD); + } + + enum CompletionStatus { + SUCCESS, INTERRUPTED + } + + @Service + public static class DemoReactiveService { + + private final ReactiveCircuitBreaker slowBulkheadCircuitBreaker; + + DemoReactiveService(ReactiveResilience4JCircuitBreakerFactory circuitBreakerFactory) { + this.slowBulkheadCircuitBreaker = circuitBreakerFactory.create(SLOW_BULKHEAD); + } + + public Mono bulkheadWithDelay(long delay) { + return slowBulkheadCircuitBreaker.run( + Mono.just(CompletionStatus.SUCCESS) + .delayElement(Duration.ofMillis(delay)), + throwable -> { + if (throwable instanceof TimeoutException || throwable instanceof BulkheadFullException) { + return Mono.just(CompletionStatus.INTERRUPTED); + } + return Mono.error(throwable); + } + ); + } + } + } +} 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 37d91d7..bb4d8f6 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 @@ -20,6 +20,7 @@ import java.util.Arrays; import java.util.Collections; import java.util.concurrent.TimeUnit; +import io.github.resilience4j.bulkhead.BulkheadRegistry; import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; import io.github.resilience4j.timelimiter.TimeLimiterRegistry; import org.junit.Assert; @@ -36,13 +37,15 @@ import static org.assertj.core.api.Assertions.assertThat; /** * @author Ryan Baxter * @author Thomas Vitale + * @author Yavor Chamov */ public class ReactiveResilience4JCircuitBreakerTest { @Test public void runMono() { ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - TimeLimiterRegistry.ofDefaults(), new Resilience4JConfigurationProperties()) + TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), + new Resilience4JConfigurationProperties()) .create("foo"); assertThat(Mono.just("foobar").transform(cb::run).block()).isEqualTo("foobar"); } @@ -50,7 +53,8 @@ public class ReactiveResilience4JCircuitBreakerTest { @Test public void runMonoWithFallback() { ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - TimeLimiterRegistry.ofDefaults(), new Resilience4JConfigurationProperties()) + TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), + new Resilience4JConfigurationProperties()) .create("foo"); assertThat(Mono.error(new RuntimeException("boom")) .transform(it -> cb.run(it, t -> Mono.just("fallback"))) @@ -60,7 +64,8 @@ public class ReactiveResilience4JCircuitBreakerTest { @Test public void runFlux() { ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - TimeLimiterRegistry.ofDefaults(), new Resilience4JConfigurationProperties()) + TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), + new Resilience4JConfigurationProperties()) .create("foo"); assertThat(Flux.just("foobar", "hello world").transform(cb::run).collectList().block()) .isEqualTo(Arrays.asList("foobar", "hello world")); @@ -69,7 +74,8 @@ public class ReactiveResilience4JCircuitBreakerTest { @Test public void runFluxWithFallback() { ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - TimeLimiterRegistry.ofDefaults(), new Resilience4JConfigurationProperties()) + TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), + new Resilience4JConfigurationProperties()) .create("foo"); assertThat(Flux.error(new RuntimeException("boom")) .transform(it -> cb.run(it, t -> Flux.just("fallback"))) @@ -85,7 +91,8 @@ public class ReactiveResilience4JCircuitBreakerTest { public void runWithDefaultTimeLimiter() { final TimeLimiterRegistry timeLimiterRegistry = TimeLimiterRegistry.ofDefaults(); ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - timeLimiterRegistry, new Resilience4JConfigurationProperties()) + timeLimiterRegistry, new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), + new Resilience4JConfigurationProperties()) .create("foo"); assertThat(Mono.fromCallable(() -> { @@ -110,7 +117,7 @@ public class ReactiveResilience4JCircuitBreakerTest { public void runWithDefaultTimeLimiterTooSlow() { final TimeLimiterRegistry timeLimiterRegistry = TimeLimiterRegistry.ofDefaults(); ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - timeLimiterRegistry, new Resilience4JConfigurationProperties()) + timeLimiterRegistry, new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), new Resilience4JConfigurationProperties()) .create("foo"); Mono.fromCallable(() -> { @@ -141,7 +148,7 @@ public class ReactiveResilience4JCircuitBreakerTest { final Resilience4JConfigurationProperties resilience4JConfigurationProperties = new Resilience4JConfigurationProperties(); resilience4JConfigurationProperties.setDisableTimeLimiter(true); ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - timeLimiterRegistry, resilience4JConfigurationProperties) + timeLimiterRegistry, new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), resilience4JConfigurationProperties) .create("foo"); assertThat(Mono.fromCallable(() -> { @@ -158,4 +165,131 @@ public class ReactiveResilience4JCircuitBreakerTest { }).subscribeOn(Schedulers.single()).transform(cb::run).block()).isEqualTo("foobar"); } + @Test + public void runMonoWithBulkheadProvider() { + ReactiveResilience4jBulkheadProvider bulkheadProvider = new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()); + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults(), + bulkheadProvider, + new Resilience4JConfigurationProperties() + ).create("foo"); + + assertThat(Mono.just("bulkheadMono") + .transform(cb::run) + .block()) + .isEqualTo("bulkheadMono"); + } + + @Test + public void runMonoWithBulkheadProviderAndFallback() { + ReactiveResilience4jBulkheadProvider bulkheadProvider = new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()); + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults(), + bulkheadProvider, + new Resilience4JConfigurationProperties() + ).create("foo"); + + assertThat(Mono.error(new RuntimeException("exception")) + .transform(it -> cb.run(it, t -> Mono.just("bulkheadFallback"))) + .block()) + .isEqualTo("bulkheadFallback"); + } + + @Test + public void runFluxWithBulkheadProvider() { + ReactiveResilience4jBulkheadProvider bulkheadProvider = new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()); + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults(), + bulkheadProvider, + new Resilience4JConfigurationProperties() + ).create("foo"); + + assertThat(Flux.just("bulkheadFlux1", "bulkheadFlux2") + .transform(cb::run) + .collectList() + .block()) + .isEqualTo(Arrays.asList("bulkheadFlux1", "bulkheadFlux2")); + } + + @Test + public void runFluxWithBulkheadProviderAndFallback() { + ReactiveResilience4jBulkheadProvider bulkheadProvider = new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()); + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults(), + bulkheadProvider, + new Resilience4JConfigurationProperties() + ).create("foo"); + + assertThat(Flux.error(new RuntimeException("exception")) + .transform(it -> cb.run(it, t -> Flux.just("bulkheadFallbackFlux"))) + .collectList() + .block()) + .isEqualTo(Collections.singletonList("bulkheadFallbackFlux")); + } + + @Test + public void runMonoWithoutBulkheadProvider() { + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults(), + null, + new Resilience4JConfigurationProperties() + ).create("foo"); + + assertThat(Mono.just("noBulkheadMono") + .transform(cb::run) + .block()) + .isEqualTo("noBulkheadMono"); + } + + @Test + public void runMonoWithoutBulkheadProviderWithFallback() { + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults(), + null, + new Resilience4JConfigurationProperties() + ).create("foo"); + + assertThat(Mono.error(new RuntimeException("exception")) + .transform(it -> cb.run(it, t -> Mono.just("noBulkheadFallback"))) + .block()) + .isEqualTo("noBulkheadFallback"); + } + + @Test + public void runFluxWithoutBulkheadProvider() { + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults(), + null, + new Resilience4JConfigurationProperties() + ).create("foo"); + + assertThat(Flux.just("noBulkheadFlux1", "noBulkheadFlux2") + .transform(cb::run) + .collectList() + .block()) + .isEqualTo(Arrays.asList("noBulkheadFlux1", "noBulkheadFlux2")); + } + + @Test + public void runFluxWithoutBulkheadProviderWithFallback() { + ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory( + CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults(), + null, + new Resilience4JConfigurationProperties() + ).create("foo"); + + assertThat(Flux.error(new RuntimeException("boom")) + .transform(it -> cb.run(it, t -> Flux.just("noBulkheadFallbackFlux"))) + .collectList() + .block()) + .isEqualTo(Collections.singletonList("noBulkheadFallbackFlux")); + } } diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4jBulkheadIntegrationTest.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4jBulkheadIntegrationTest.java new file mode 100644 index 0000000..f0e8c05 --- /dev/null +++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/ReactiveResilience4jBulkheadIntegrationTest.java @@ -0,0 +1,234 @@ +/* + * Copyright 2013-2018 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 java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Map; + +import io.github.resilience4j.bulkhead.BulkheadConfig; +import io.github.resilience4j.bulkhead.BulkheadFullException; +import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig; +import io.github.resilience4j.timelimiter.TimeLimiterConfig; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.junit.Test; +import org.junit.runner.RunWith; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; +import reactor.test.StepVerifier; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.web.client.TestRestTemplate; +import org.springframework.cloud.client.circuitbreaker.Customizer; +import org.springframework.cloud.client.circuitbreaker.ReactiveCircuitBreaker; +import org.springframework.cloud.client.circuitbreaker.ReactiveCircuitBreakerFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.http.HttpHeaders; +import org.springframework.stereotype.Service; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.junit4.SpringRunner; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.RequestHeader; +import org.springframework.web.bind.annotation.RestController; + +import static org.assertj.core.api.AssertionsForInterfaceTypes.assertThat; +import static org.springframework.boot.test.context.SpringBootTest.WebEnvironment.RANDOM_PORT; + +/** + * @author Yavor Chamov + */ +@RunWith(SpringRunner.class) +@SpringBootTest(webEnvironment = RANDOM_PORT, classes = ReactiveResilience4jBulkheadIntegrationTest.Application.class, + properties = {"management.endpoints.web.exposure.include=*"}) +@DirtiesContext +public class ReactiveResilience4jBulkheadIntegrationTest { + + @Autowired + Application.DemoControllerService service; + + @Autowired + private TestRestTemplate rest; + + @Test + public void testSlow() { + StepVerifier.create(service.slow()) + .expectNext("fallback") + .verifyComplete(); + } + + @Test + public void testNormal() { + StepVerifier.create(service.normal()) + .expectNext("normal") + .verifyComplete(); + } + + @Test + public void testSlowResponsesDontFailSubsequentGoodRequests() { + StepVerifier.create(service.slowOnDemand(5000)) + .expectNext("fallback") + .verifyComplete(); + + StepVerifier.create(service.slowOnDemand(0)) + .expectNext("normal") + .verifyComplete(); + } + + @Test + public void testBulkheadConcurrentCallLimit() { + List results = Collections.synchronizedList(new ArrayList<>()); + + Flux.merge( + Mono.defer(() -> service.slowBulkhead() + .doOnNext(results::add)), + Mono.defer(() -> service.slowBulkhead() + .doOnNext(results::add)) + ).blockLast(); + + assertThat(results).containsExactlyInAnyOrder("slowBulkhead", "rejected"); + } + + @Test + public void testResilience4JMetricsAvailable() { + StepVerifier.create(service.normal()) + .expectNext("normal") + .verifyComplete(); + + @SuppressWarnings("unchecked") + Map metrics = rest.getForObject("/actuator/metrics", Map.class); + List metricNames = (List) metrics.get("names"); + + assertThat(metricNames).contains("resilience4j.bulkhead.max.allowed.concurrent.calls"); + assertThat(rest + .getForObject("/actuator/metrics/resilience4j.bulkhead.max.allowed.concurrent.calls", Map.class) + .get("availableTags")).isNotNull(); + } + + @Configuration(proxyBeanMethods = false) + @EnableAutoConfiguration + @RestController + protected static class Application { + + private static final Log LOG = LogFactory.getLog(Application.class); + + @GetMapping("/slow") + public Mono slow() { + return Mono.delay(Duration.ofSeconds(3)).thenReturn("slow"); + } + + @GetMapping("/normal") + public Mono normal() { + return Mono.just("normal"); + } + + @GetMapping("/slowOnDemand") + public Mono slowOnDemand(@RequestHeader HttpHeaders headers) { + if (headers.containsKey("delayInMilliseconds")) { + String delayString = headers.getFirst("delayInMilliseconds"); + LOG.info("delay header: " + delayString); + if (delayString != null) { + try { + long delay = Long.parseLong(delayString); + return Mono.delay(Duration.ofMillis(delay)).thenReturn("normal"); + } + catch (NumberFormatException e) { + e.printStackTrace(); + } + } + } + return Mono.just("normal"); + } + + @GetMapping("/slowBulkhead") + public Mono slowBulkhead() { + return Mono.delay(Duration.ofSeconds(3)).thenReturn("slowBulkhead"); + } + + @Bean + public Customizer slowCustomizer() { + return factory -> { + factory.configure(builder -> builder.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults()) + .timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(3)).build()), + "slow"); + factory.configureDefault(id -> new Resilience4JConfigBuilder(id) + .timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(4)).build()) + .circuitBreakerConfig(CircuitBreakerConfig.ofDefaults()) + .build()); + }; + } + + @Bean + public Customizer bulkheadCustomizer() { + return provider -> { + provider.configure(builder -> builder.bulkheadConfig( + BulkheadConfig.custom().maxConcurrentCalls(1).build()), "slowBulkhead"); + }; + } + + @Service + public static class DemoControllerService { + + private final ReactiveCircuitBreakerFactory cbFactory; + + private final ReactiveCircuitBreaker circuitBreakerSlow; + + DemoControllerService(ReactiveCircuitBreakerFactory cbFactory) { + this.cbFactory = cbFactory; + this.circuitBreakerSlow = cbFactory.create("slow"); + } + + public Mono slow() { + return circuitBreakerSlow.run( + Mono.delay(Duration.ofSeconds(4)) + .thenReturn("slow"), + t -> Mono.just("fallback") + ); + } + + public Mono normal() { + return cbFactory.create("normal").run(Mono.just("normal"), t -> Mono.just("fallback")); + } + + public Mono slowOnDemand(int delayInMilliseconds) { + LOG.info("delay: " + delayInMilliseconds); + return circuitBreakerSlow.run( + Mono.delay(Duration.ofMillis(delayInMilliseconds)).thenReturn("normal"), + t -> Mono.just("fallback") + ); + } + + public Mono slowBulkhead() { + return cbFactory.create("slowBulkhead") + .run( + Mono.delay(Duration.ofSeconds(2)).thenReturn("slowBulkhead"), + throwable -> { + if (throwable instanceof BulkheadFullException) { + return Mono.just("rejected"); + } + return Mono.just("fallback"); + } + ); + } + } + } +} diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JAutoConfigurationMetricsConfigRun.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JAutoConfigurationMetricsConfigRun.java index d59e9c6..163f0af 100644 --- a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JAutoConfigurationMetricsConfigRun.java +++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JAutoConfigurationMetricsConfigRun.java @@ -41,12 +41,12 @@ import static org.mockito.Mockito.verify; public class Resilience4JAutoConfigurationMetricsConfigRun { static ReactiveResilience4JCircuitBreakerFactory reactiveCircuitBreakerFactory = spy( - new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - TimeLimiterRegistry.ofDefaults(), new Resilience4JConfigurationProperties())); + new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), + TimeLimiterRegistry.ofDefaults(), new Resilience4JConfigurationProperties())); static Resilience4JCircuitBreakerFactory circuitBreakerFactory = spy( - new Resilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), TimeLimiterRegistry.ofDefaults(), - mock(Resilience4jBulkheadProvider.class), new Resilience4JConfigurationProperties())); + new Resilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), TimeLimiterRegistry.ofDefaults(), + mock(Resilience4jBulkheadProvider.class), new Resilience4JConfigurationProperties())); @Test public void testWithMetricsConfigReactive() {