From 93dc27b588ee7be98057d5468d3d60cf7a2cb79f Mon Sep 17 00:00:00 2001 From: Yavor Chamov Date: Sun, 22 Dec 2024 14:31:22 +0200 Subject: [PATCH] Add support for reactive Resilience4jBulkheadProvider to provide bulkhead support for reactive operations (Mono and Flux). --- ...nce4JAutoConfigurationWithoutBulkheadTest.java | 2 -- ...ce4JBulkheadAndTimeLimiterIntegrationTest.java | 3 +++ .../ReactiveResilience4JCircuitBreakerTest.java | 15 ++++++++++----- ...activeResilience4jBulkheadIntegrationTest.java | 3 +++ 4 files changed, 16 insertions(+), 7 deletions(-) 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 index ae655cf..4c617eb 100644 --- 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 @@ -30,8 +30,6 @@ import org.springframework.context.ConfigurableApplicationContext; import static org.assertj.core.api.Assertions.assertThat; /** - * @author Ryan Baxter - * @author Thomas Vitale * @author Yavor Chamov */ @RunWith(ModifiedClassPathRunner.class) 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 index 55d2cc4..2f7d931 100644 --- 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 @@ -37,6 +37,9 @@ 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 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 f0f2afa..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 @@ -44,7 +44,8 @@ public class ReactiveResilience4JCircuitBreakerTest { @Test public void runMono() { ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), new Resilience4JConfigurationProperties()) + TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), + new Resilience4JConfigurationProperties()) .create("foo"); assertThat(Mono.just("foobar").transform(cb::run).block()).isEqualTo("foobar"); } @@ -52,7 +53,8 @@ public class ReactiveResilience4JCircuitBreakerTest { @Test public void runMonoWithFallback() { ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.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"))) @@ -62,7 +64,8 @@ public class ReactiveResilience4JCircuitBreakerTest { @Test public void runFlux() { ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.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")); @@ -71,7 +74,8 @@ public class ReactiveResilience4JCircuitBreakerTest { @Test public void runFluxWithFallback() { ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.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"))) @@ -87,7 +91,8 @@ public class ReactiveResilience4JCircuitBreakerTest { public void runWithDefaultTimeLimiter() { final TimeLimiterRegistry timeLimiterRegistry = TimeLimiterRegistry.ofDefaults(); ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(), - timeLimiterRegistry, new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), new Resilience4JConfigurationProperties()) + timeLimiterRegistry, new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), + new Resilience4JConfigurationProperties()) .create("foo"); assertThat(Mono.fromCallable(() -> { 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 index 7fc3627..f0e8c05 100644 --- 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 @@ -54,6 +54,9 @@ 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=*"})