From 899f295fb5e3a223a176d77edd923ac455b41cc4 Mon Sep 17 00:00:00 2001 From: Ramesh Malla Date: Mon, 31 Aug 2026 01:24:01 +0200 Subject: [PATCH 1/2] failsafe composition fix Signed-off-by: Ramesh Malla --- .../riptide/failsafe/FailsafePlugin.java | 4 + .../DefaultRiptideRegistrar.java | 97 ++++++--------- .../riptide/autoconfigure/Defaulting.java | 9 ++ .../autoconfigure/FailsafePluginFactory.java | 110 ++++++++---------- .../LegacyFailsafeThreadsException.java | 25 ++++ .../autoconfigure/RiptideProperties.java | 42 +++++++ ...FailSafeExecutorAutoConfigurationTest.java | 89 ++++++++++---- .../riptide/autoconfigure/PluginTest.java | 30 ++++- .../test/resources/application-default.yml | 22 +--- 9 files changed, 263 insertions(+), 165 deletions(-) create mode 100644 riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/LegacyFailsafeThreadsException.java diff --git a/riptide-failsafe/src/main/java/org/zalando/riptide/failsafe/FailsafePlugin.java b/riptide-failsafe/src/main/java/org/zalando/riptide/failsafe/FailsafePlugin.java index 22df68f4c..19098d6f2 100644 --- a/riptide-failsafe/src/main/java/org/zalando/riptide/failsafe/FailsafePlugin.java +++ b/riptide-failsafe/src/main/java/org/zalando/riptide/failsafe/FailsafePlugin.java @@ -52,6 +52,10 @@ public FailsafePlugin withPolicy(final RequestPolicy policy) { return new FailsafePlugin(policies.append(policy), decorators, executorService); } + public FailsafePlugin withPolicies(final List requestPolicies) { + return new FailsafePlugin(policies.concat(requestPolicies), decorators, executorService); + } + public FailsafePlugin withExecutor(@Nullable final ExecutorService executorService) { if (executorService instanceof ThreadPoolExecutor && ((ThreadPoolExecutor) executorService).getCorePoolSize() == 1) { diff --git a/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/DefaultRiptideRegistrar.java b/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/DefaultRiptideRegistrar.java index ee6fa7e36..c16519d2b 100644 --- a/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/DefaultRiptideRegistrar.java +++ b/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/DefaultRiptideRegistrar.java @@ -4,7 +4,7 @@ import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableMap; import dev.failsafe.CircuitBreaker; -import dev.failsafe.Timeout; +import dev.failsafe.Policy; import io.micrometer.core.instrument.Tag; import io.micrometer.core.instrument.Tags; import io.opentracing.contrib.concurrent.TracedExecutorService; @@ -37,7 +37,6 @@ import org.zalando.riptide.chaos.Probability; import org.zalando.riptide.compatibility.HttpOperations; import org.zalando.riptide.compression.RequestCompressionPlugin; -import org.zalando.riptide.failsafe.BackupRequest; import org.zalando.riptide.failsafe.CircuitBreakerListener; import org.zalando.riptide.failsafe.FailsafePlugin; import org.zalando.riptide.httpclient.ApacheClientHttpRequestFactory; @@ -251,11 +250,8 @@ private List registerPlugins(final String id, final Client client registerLogbookPlugin(id, client), registerOpenTracingPlugin(id, client), registerOpenTelemetryPlugin(id, client), - registerCircuitBreakerFailsafePlugin(id, client), - registerRetryPolicyFailsafePlugin(id, client), registerAuthorizationPlugin(id, client), - registerBackupRequestFailsafePlugin(id, client), - registerTimeoutFailsafePlugin(id, client), + registerFailsafePluginWithConfiguredPolicies(id, client), registerOriginalStackTracePlugin(id, client), registerCustomPlugin(id)); @@ -409,38 +405,6 @@ private Optional registerOpenTelemetryPlugin(final String id, final Clie return Optional.empty(); } - private Optional registerCircuitBreakerFailsafePlugin(final String id, final Client client) { - if (client.getCircuitBreaker().getEnabled()) { - final String pluginId = registry.registerIfAbsent(name(id, CircuitBreaker.class, FailsafePlugin.class), - () -> { - log.debug("Client [{}]: Registering [CircuitBreakerFailsafePlugin]", id); - return genericBeanDefinition(FailsafePluginFactory.class) - .setFactoryMethod("createCircuitBreakerPlugin") - .addConstructorArgValue(registerCircuitBreaker(id, client)) - .addConstructorArgValue(createTaskDecorators(id, client)) - .addConstructorArgValue(createExecutor(id + "-circuit-breaker", "failsafe.circuitbreaker.executor", client, client.getCircuitBreaker().getThreads())); - }); - return Optional.of(pluginId); - } - return Optional.empty(); - } - - private Optional registerRetryPolicyFailsafePlugin(final String id, final Client client) { - if (client.getRetry().getEnabled()) { - - final String pluginId = registry.registerIfAbsent(name(id, "RetryPolicy", FailsafePlugin.class), () -> { - log.debug("Client [{}]: Registering [RetryPolicyFailsafePlugin]", id); - return genericBeanDefinition(FailsafePluginFactory.class) - .setFactoryMethod("createRetryFailsafePlugin") - .addConstructorArgValue(client) - .addConstructorArgValue(createTaskDecorators(id, client)) - .addConstructorArgValue(createExecutor(id + "-retry-policy", "failsafe.retry.executor", client, client.getRetry().getThreads())); - }); - return Optional.of(pluginId); - } - return Optional.empty(); - } - private Optional registerAuthorizationPlugin(final String id, final Client client) { if (client.getAuth().getEnabled()) { final String pluginId = registry.registerIfAbsent(id, AuthorizationPlugin.class, () -> { @@ -453,37 +417,52 @@ private Optional registerAuthorizationPlugin(final String id, final Clie return Optional.empty(); } - private Optional registerBackupRequestFailsafePlugin(final String id, final Client client) { - if (client.getBackupRequest().getEnabled()) { - final String pluginId = registry.registerIfAbsent(name(id, BackupRequest.class, FailsafePlugin.class), - () -> { - log.debug("Client [{}]: Registering [BackupRequestFailsafePlugin]", id); - return genericBeanDefinition(FailsafePluginFactory.class) - .setFactoryMethod("createBackupRequestPlugin") - .addConstructorArgValue(client) - .addConstructorArgValue(createTaskDecorators(id, client)) - .addConstructorArgValue(createExecutor(id + "-backup-request", "failsafe.backuprequest.executor", client, client.getBackupRequest().getThreads())); - }); - return Optional.of(pluginId); - } - return Optional.empty(); - } + private Optional registerFailsafePluginWithConfiguredPolicies(final String id, final Client client) { + if(client.getTimeouts().getEnabled() || client.getRetry().getEnabled() || client.getBackupRequest().getEnabled() || client.getCircuitBreaker().getEnabled()){ + final Object executor = createExecutor(id + "-failsafe", "failsafe.executor", client, + resolveFailsafeThreads(id, client)); + + final BeanMetadataElement circuitBreaker = client.getCircuitBreaker().getEnabled() + ? registerCircuitBreaker(id, client) + : null; - private Optional registerTimeoutFailsafePlugin(final String id, final Client client) { - if (client.getTimeouts().getEnabled()) { - final String pluginId = registry.registerIfAbsent(name(id, Timeout.class, FailsafePlugin.class), () -> { - log.debug("Client [{}]: Registering [TimeoutFailsafePlugin]", id); + final String pluginId = registry.registerIfAbsent(name(id, Policy.class, FailsafePlugin.class), () -> { + log.debug("Client [{}]: Registering [FailsafePlugin]", id); return genericBeanDefinition(FailsafePluginFactory.class) - .setFactoryMethod("createTimeoutPlugin") + .setFactoryMethod("createFailsafePlugin") .addConstructorArgValue(client) .addConstructorArgValue(createTaskDecorators(id, client)) - .addConstructorArgValue(createExecutor(id + "-timeout","failsafe.timeout.executor", client, client.getTimeouts().getThreads())); + .addConstructorArgValue(executor) + .addConstructorArgValue(circuitBreaker); }); return Optional.of(pluginId); } return Optional.empty(); } + /** + * Resolves the single, shared thread pool configuration for a client's merged + * {@link FailsafePlugin}. The deprecated per-policy {@code *.threads} configs are no longer + * supported; using one throws {@link LegacyFailsafeThreadsException}. + */ + @Nullable + private RiptideProperties.Threads resolveFailsafeThreads(final String id, final Client client) { + rejectLegacyThreads(id, "retry.threads", client.getRetry().getEnabled(), client.getRetry().getThreads()); + rejectLegacyThreads(id, "circuit-breaker.threads", client.getCircuitBreaker().getEnabled(), client.getCircuitBreaker().getThreads()); + rejectLegacyThreads(id, "backup-request.threads", client.getBackupRequest().getEnabled(), client.getBackupRequest().getThreads()); + rejectLegacyThreads(id, "timeouts.threads", client.getTimeouts().getEnabled(), client.getTimeouts().getThreads()); + + final RiptideProperties.Threads failsafeThreads = client.getFailsafe().getThreads(); + return failsafeThreads != null && failsafeThreads.getEnabled() ? failsafeThreads : null; + } + + private static void rejectLegacyThreads(final String id, final String property, + final boolean policyEnabled, @Nullable final RiptideProperties.Threads threads) { + if (policyEnabled && threads != null && Boolean.TRUE.equals(threads.getEnabled())) { + throw new LegacyFailsafeThreadsException(id, property); + } + } + private Optional registerOriginalStackTracePlugin(final String id, final Client client) { if (client.getStackTracePreservation().getEnabled()) { final String pluginId = registry.registerIfAbsent(id, OriginalStackTracePlugin.class, () -> { diff --git a/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/Defaulting.java b/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/Defaulting.java index e2b013e9c..738c3d92b 100644 --- a/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/Defaulting.java +++ b/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/Defaulting.java @@ -39,6 +39,7 @@ import static org.zalando.riptide.autoconfigure.RiptideProperties.CircuitBreaker; import static org.zalando.riptide.autoconfigure.RiptideProperties.Client; import static org.zalando.riptide.autoconfigure.RiptideProperties.Defaults; +import static org.zalando.riptide.autoconfigure.RiptideProperties.Failsafe; import static org.zalando.riptide.autoconfigure.RiptideProperties.Retry; import static org.zalando.riptide.autoconfigure.RiptideProperties.Threads; @@ -74,6 +75,7 @@ private static Defaults merge(final Defaults defaults) { defaults.getCircuitBreaker(), defaults.getBackupRequest(), defaults.getTimeouts(), + defaults.getFailsafe(), defaults.getRequestCompression(), defaults.getCertificatePinning(), defaults.getCaching(), @@ -117,6 +119,7 @@ private static Client merge(final Client base, final Defaults defaults) { merge(base.getCircuitBreaker(), defaults.getCircuitBreaker(), Defaulting::merge), merge(base.getBackupRequest(), defaults.getBackupRequest(), Defaulting::merge), merge(base.getTimeouts(), defaults.getTimeouts(), Defaulting::merge), + merge(base.getFailsafe(), defaults.getFailsafe(), Defaulting::merge), merge(base.getRequestCompression(), defaults.getRequestCompression(), Defaulting::merge), merge(base.getCertificatePinning(), defaults.getCertificatePinning(), Defaulting::merge), merge(base.getCaching(), defaults.getCaching(), Defaulting::merge), @@ -245,6 +248,12 @@ private static Timeouts merge(final Timeouts base, final Timeouts defaults) { ); } + private static Failsafe merge(final Failsafe base, final Failsafe defaults) { + return new Failsafe( + either(base.getThreads(), defaults.getThreads()) + ); + } + private static RequestCompression merge(final RequestCompression base, final RequestCompression defaults) { return new RequestCompression( either(base.getEnabled(), defaults.getEnabled()) diff --git a/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/FailsafePluginFactory.java b/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/FailsafePluginFactory.java index 69aec5760..bebba85ab 100644 --- a/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/FailsafePluginFactory.java +++ b/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/FailsafePluginFactory.java @@ -17,6 +17,7 @@ import org.zalando.riptide.failsafe.FailsafePlugin; import org.zalando.riptide.failsafe.RateLimitResetDelayFunction; import org.zalando.riptide.failsafe.RequestPolicies; +import org.zalando.riptide.failsafe.RequestPolicy; import org.zalando.riptide.failsafe.RetryAfterDelayFunction; import org.zalando.riptide.failsafe.RetryException; import org.zalando.riptide.failsafe.RetryRequestPolicy; @@ -25,6 +26,7 @@ import javax.annotation.Nullable; import java.time.Duration; +import java.util.ArrayList; import java.util.Arrays; import java.util.List; import java.util.Optional; @@ -48,17 +50,6 @@ private FailsafePluginFactory() { } - public static Plugin createCircuitBreakerPlugin( - final CircuitBreaker breaker, - final List decorators, - @Nullable final ExecutorService executorService) { - - return new FailsafePlugin() - .withExecutor(executorService) - .withPolicy(breaker) - .withDecorator(composite(decorators)); - } - public static CircuitBreaker createCircuitBreaker( final Client client, final CircuitBreakerListener listener) { @@ -85,30 +76,53 @@ public static CircuitBreaker createCircuitBreaker( return breakerBuilder.build(); } - public static Plugin createRetryFailsafePlugin( + + public static Plugin createFailsafePlugin( final Client client, final List decorators, - @Nullable final ExecutorService executorService) { - - if (client.getTransientFaultDetection().getEnabled()) { - return new FailsafePlugin() - .withExecutor(executorService) - .withPolicy(new RetryRequestPolicy(getRetryPolicyBuilder(client) - .handleIf(toCheckedPredicate(transientSocketFaults())) - .build()) - .withPredicate(new IdempotencyPredicate())) - .withPolicy(new RetryRequestPolicy(getRetryPolicyBuilder(client) - .handleIf(toCheckedPredicate(transientConnectionFaults())) - .build()) - .withPredicate(alwaysTrue())) - .withPolicy(new RetryRequestPolicy(getRetryPolicyBuilder(client).handle(RetryException.class).build())) - .withDecorator(composite(decorators)); - } else { - return new FailsafePlugin() - .withExecutor(executorService) - .withPolicy(new RetryRequestPolicy(getRetryPolicyBuilder(client).handle(RetryException.class).build())) - .withDecorator(composite(decorators)); + @Nullable final ExecutorService executorService, + @Nullable final CircuitBreaker circuitBreaker + ){ + final List requestPolicies = new ArrayList<>(); + + if(client.getTimeouts().getEnabled()){ + final Duration timeout = client.getTimeouts().getGlobal().toDuration(); + requestPolicies.add(RequestPolicies.of(Timeout.builder(timeout) + .withInterrupt() + .build())); + } + + if(client.getBackupRequest().getEnabled()){ + final TimeSpan delay = client.getBackupRequest().getDelay(); + requestPolicies.add(RequestPolicies.of( + new BackupRequest<>(delay.getAmount(), delay.getUnit()), + new IdempotencyPredicate())); + } + + if (client.getRetry().getEnabled()) { + if (client.getTransientFaultDetection().getEnabled()) { + requestPolicies.add(new RetryRequestPolicy(getRetryPolicyBuilder(client) + .handleIf(toCheckedPredicate(transientSocketFaults())) + .build()) + .withPredicate(new IdempotencyPredicate())); + requestPolicies.add(new RetryRequestPolicy(getRetryPolicyBuilder(client) + .handleIf(toCheckedPredicate(transientConnectionFaults())) + .build()) + .withPredicate(alwaysTrue())); + requestPolicies.add(new RetryRequestPolicy(getRetryPolicyBuilder(client).handle(RetryException.class).build())); + } else { + requestPolicies.add((new RetryRequestPolicy(getRetryPolicyBuilder(client).handle(RetryException.class).build()))); + } + } + + if (circuitBreaker != null) { + requestPolicies.add(RequestPolicies.of(circuitBreaker)); } + + return new FailsafePlugin() + .withExecutor(executorService) + .withPolicies(requestPolicies) + .withDecorator(composite(decorators)); } private static RetryPolicyBuilder getRetryPolicyBuilder(Client client) { @@ -151,38 +165,6 @@ private static RetryPolicyBuilder getRetryPolicyBuilder(Clie return policyBuilder; } - public static Plugin createBackupRequestPlugin( - final Client client, - final List decorators, - @Nullable final ExecutorService executorService) { - - final TimeSpan delay = client.getBackupRequest().getDelay(); - - return new FailsafePlugin() - .withExecutor(executorService) - .withPolicy(RequestPolicies.of( - new BackupRequest<>(delay.getAmount(), delay.getUnit()), - new IdempotencyPredicate())) - .withDecorator(composite(decorators)); - } - - public static Plugin createTimeoutPlugin( - final Client client, - final List decorators, - @Nullable final ExecutorService executorService) { - - final Duration timeout = client.getTimeouts().getGlobal().toDuration(); - - return new FailsafePlugin() - .withExecutor(executorService) - .withPolicy( - Timeout.builder(timeout) - .withInterrupt() - .build() - ) - .withDecorator(composite(decorators)); - } - private static ContextualSupplier delayFunction() { return new CompositeDelayFunction<>(Arrays.asList( new RetryAfterDelayFunction(systemUTC()), diff --git a/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/LegacyFailsafeThreadsException.java b/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/LegacyFailsafeThreadsException.java new file mode 100644 index 000000000..0c7b6916b --- /dev/null +++ b/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/LegacyFailsafeThreadsException.java @@ -0,0 +1,25 @@ +package org.zalando.riptide.autoconfigure; + +import lombok.Getter; +import org.springframework.core.NestedRuntimeException; + +@Getter +public class LegacyFailsafeThreadsException extends NestedRuntimeException { + + private final String clientId; + private final String property; + + public LegacyFailsafeThreadsException(final String clientId, final String property) { + super(createMessage(clientId, property)); + this.clientId = clientId; + this.property = property; + } + + private static String createMessage(final String clientId, final String property) { + return String.format( + "Client [%s]: [riptide.clients.%s.%s] is no longer supported, configure " + + "[riptide.clients.%s.failsafe.threads] instead", + clientId, clientId, property, clientId); + } + +} diff --git a/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/RiptideProperties.java b/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/RiptideProperties.java index 430dfaf06..046a5101a 100644 --- a/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/RiptideProperties.java +++ b/riptide-spring-boot-autoconfigure/src/main/java/org/zalando/riptide/autoconfigure/RiptideProperties.java @@ -108,6 +108,9 @@ public static final class Defaults { @NestedConfigurationProperty private Timeouts timeouts = new Timeouts(false, null, null); + @NestedConfigurationProperty + private Failsafe failsafe = new Failsafe(new Threads(false, null, null, null, null)); + @NestedConfigurationProperty private RequestCompression requestCompression = new RequestCompression(false); @@ -196,6 +199,9 @@ public static final class Client { @NestedConfigurationProperty private Timeouts timeouts; + @NestedConfigurationProperty + private Failsafe failsafe; + @NestedConfigurationProperty private RequestCompression requestCompression; @@ -306,6 +312,13 @@ public static final class Retry { private TimeSpan maxDuration; private Double jitterFactor; private TimeSpan jitter; + + /** + * @deprecated no longer supported; configure {@code riptide.clients..failsafe.threads} + * instead. Setting this field now throws {@link LegacyFailsafeThreadsException} at + * startup. + */ + @Deprecated private Threads threads; @Getter @@ -330,6 +343,13 @@ public static final class CircuitBreaker { private RatioInTimeSpan failureRateThreshold; private TimeSpan delay; private Ratio successThreshold; + + /** + * @deprecated no longer supported; configure {@code riptide.clients..failsafe.threads} + * instead. Setting this field now throws {@link LegacyFailsafeThreadsException} at + * startup. + */ + @Deprecated private Threads threads; } @@ -340,6 +360,13 @@ public static final class CircuitBreaker { public static final class BackupRequest { private Boolean enabled; private TimeSpan delay; + + /** + * @deprecated no longer supported; configure {@code riptide.clients..failsafe.threads} + * instead. Setting this field now throws {@link LegacyFailsafeThreadsException} at + * startup. + */ + @Deprecated private Threads threads; } @@ -350,6 +377,21 @@ public static final class BackupRequest { public static final class Timeouts { private Boolean enabled; private TimeSpan global; + + /** + * @deprecated no longer supported; configure {@code riptide.clients..failsafe.threads} + * instead. Setting this field now throws {@link LegacyFailsafeThreadsException} at + * startup. + */ + @Deprecated + private Threads threads; + } + + @Getter + @Setter + @NoArgsConstructor + @AllArgsConstructor + public static final class Failsafe { private Threads threads; } diff --git a/riptide-spring-boot-autoconfigure/src/test/java/org/zalando/riptide/autoconfigure/FailSafeExecutorAutoConfigurationTest.java b/riptide-spring-boot-autoconfigure/src/test/java/org/zalando/riptide/autoconfigure/FailSafeExecutorAutoConfigurationTest.java index 2f03996f1..bdb052b6e 100644 --- a/riptide-spring-boot-autoconfigure/src/test/java/org/zalando/riptide/autoconfigure/FailSafeExecutorAutoConfigurationTest.java +++ b/riptide-spring-boot-autoconfigure/src/test/java/org/zalando/riptide/autoconfigure/FailSafeExecutorAutoConfigurationTest.java @@ -1,19 +1,34 @@ package org.zalando.riptide.autoconfigure; +import com.google.common.collect.ImmutableMap; import lombok.extern.slf4j.Slf4j; import org.assertj.core.api.Assertions; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.support.SimpleBeanDefinitionRegistry; import org.springframework.boot.autoconfigure.ImportAutoConfiguration; import org.springframework.boot.autoconfigure.jackson.JacksonAutoConfiguration; import org.springframework.context.ApplicationContext; import org.springframework.context.annotation.Configuration; import org.springframework.test.context.ActiveProfiles; import org.zalando.logbook.autoconfigure.LogbookAutoConfiguration; +import org.zalando.riptide.autoconfigure.RiptideProperties.BackupRequest; +import org.zalando.riptide.autoconfigure.RiptideProperties.CircuitBreaker; +import org.zalando.riptide.autoconfigure.RiptideProperties.Client; +import org.zalando.riptide.autoconfigure.RiptideProperties.Defaults; +import org.zalando.riptide.autoconfigure.RiptideProperties.Retry; +import org.zalando.riptide.autoconfigure.RiptideProperties.Threads; +import org.zalando.riptide.autoconfigure.RiptideProperties.Timeouts; +import java.net.URI; import java.util.concurrent.ExecutorService; +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.containsString; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + @RiptideClientTest @ActiveProfiles("default") @Slf4j @@ -33,27 +48,61 @@ static class ContextConfiguration { private ApplicationContext applicationContext; @Test - public void shouldContainExecutorsConfiguredForFailSafePolicies() { - final var customExecutorTestRetryPolicyExecutorService = applicationContext.getBean("customExecutorTestRetryPolicyExecutorService"); - final var customExecutorTestCircuitBreakerExecutorService = applicationContext.getBean("customExecutorTestCircuitBreakerExecutorService", ExecutorService.class); - final var customExecutorTestBackupRequestExecutorService = applicationContext.getBean("customExecutorTestBackupRequestExecutorService", ExecutorService.class); - final var customExecutorTestTimeoutExecutorService = applicationContext.getBean("customExecutorTestTimeoutExecutorService", ExecutorService.class); - Assertions.assertThat(customExecutorTestRetryPolicyExecutorService) - .isNotNull() - .hasFieldOrPropertyWithValue("corePoolSize",2) - .hasFieldOrPropertyWithValue("maximumPoolSize", 13); - Assertions.assertThat(customExecutorTestCircuitBreakerExecutorService) - .isNotNull() - .hasFieldOrPropertyWithValue("corePoolSize",2) - .hasFieldOrPropertyWithValue("maximumPoolSize", 10); - Assertions.assertThat(customExecutorTestBackupRequestExecutorService) - .isNotNull() - .hasFieldOrPropertyWithValue("corePoolSize",2) - .hasFieldOrPropertyWithValue("maximumPoolSize", 12); - Assertions.assertThat(customExecutorTestTimeoutExecutorService) + public void shouldContainSingleExecutorConfiguredViaFailsafeThreads() { + final var customExecutorTestFailsafeExecutorService = applicationContext.getBean("customExecutorTestFailsafeExecutorService", ExecutorService.class); + Assertions.assertThat(customExecutorTestFailsafeExecutorService) .isNotNull() - .hasFieldOrPropertyWithValue("corePoolSize",2) - .hasFieldOrPropertyWithValue("maximumPoolSize", 11); + .hasFieldOrPropertyWithValue("corePoolSize", 2) + .hasFieldOrPropertyWithValue("maximumPoolSize", 14); + } + + @Test + void shouldRejectLegacyRetryThreads() { + final Client client = new Client(); + client.setRetry(new Retry(true, null, null, null, null, null, null, new Threads(true, 2, 4, null, null))); + + assertRejected("legacy-retry", client, "retry.threads"); + } + + @Test + void shouldRejectLegacyCircuitBreakerThreads() { + final Client client = new Client(); + client.setCircuitBreaker(new CircuitBreaker(true, null, null, null, null, new Threads(true, 2, 4, null, null))); + + assertRejected("legacy-circuit-breaker", client, "circuit-breaker.threads"); + } + + @Test + void shouldRejectLegacyBackupRequestThreads() { + final Client client = new Client(); + client.setBackupRequest(new BackupRequest(true, null, new Threads(true, 2, 4, null, null))); + + assertRejected("legacy-backup-request", client, "backup-request.threads"); + } + + @Test + void shouldRejectLegacyTimeoutsThreads() { + final Client client = new Client(); + client.setTimeouts(new Timeouts(true, null, new Threads(true, 2, 4, null, null))); + + assertRejected("legacy-timeouts", client, "timeouts.threads"); + } + + private void assertRejected(final String id, final Client client, final String expectedProperty) { + final RiptideProperties properties = Defaulting.withDefaults( + new RiptideProperties(new Defaults(), ImmutableMap.of(id, client))); + properties.getClients().get(id).setBaseUrl(URI.create("http://example.com")); + + final DefaultRiptideRegistrar registrar = new DefaultRiptideRegistrar( + new Registry(new SimpleBeanDefinitionRegistry()), properties); + + final LegacyFailsafeThreadsException exception = assertThrows( + LegacyFailsafeThreadsException.class, registrar::register); + + assertEquals(id, exception.getClientId()); + assertEquals(expectedProperty, exception.getProperty()); + assertThat(exception.getMessage(), containsString("riptide.clients." + id + "." + expectedProperty)); + assertThat(exception.getMessage(), containsString("riptide.clients." + id + ".failsafe.threads")); } } diff --git a/riptide-spring-boot-autoconfigure/src/test/java/org/zalando/riptide/autoconfigure/PluginTest.java b/riptide-spring-boot-autoconfigure/src/test/java/org/zalando/riptide/autoconfigure/PluginTest.java index 20148f504..a143914bc 100644 --- a/riptide-spring-boot-autoconfigure/src/test/java/org/zalando/riptide/autoconfigure/PluginTest.java +++ b/riptide-spring-boot-autoconfigure/src/test/java/org/zalando/riptide/autoconfigure/PluginTest.java @@ -13,6 +13,7 @@ import org.zalando.riptide.OriginalStackTracePlugin; import org.zalando.riptide.Plugin; import org.zalando.riptide.failsafe.FailsafePlugin; +import org.zalando.riptide.failsafe.RequestPolicy; import org.zalando.riptide.logbook.LogbookPlugin; import org.zalando.riptide.micrometer.MicrometerPlugin; import org.zalando.riptide.opentelemetry.OpenTelemetryPlugin; @@ -25,7 +26,9 @@ import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.Matchers.contains; import static org.hamcrest.Matchers.hasItem; +import static org.hamcrest.Matchers.hasSize; import static org.hamcrest.Matchers.instanceOf; +import static org.hamcrest.Matchers.is; import static org.springframework.boot.test.context.SpringBootTest.WebEnvironment.NONE; @SpringBootTest(classes = PluginTest.TestConfiguration.class, webEnvironment = NONE) @@ -83,13 +86,27 @@ void shouldUseFailsafePlugin() throws Exception { @Test void shouldUseBackupRequestPlugin() throws Exception { - assertThat(getPlugins(baz), contains(asList( + final List plugins = getPlugins(baz); + assertThat(plugins, contains(asList( instanceOf(Plugin.class), // internal plugin instanceOf(Plugin.class), // internal plugin instanceOf(Plugin.class), // internal plugin instanceOf(MicrometerPlugin.class), - instanceOf(FailsafePlugin.class), // backup requests - instanceOf(FailsafePlugin.class)))); // timeouts + instanceOf(FailsafePlugin.class)))); + + final List failsafePlugins = plugins.stream() + .filter(plugin -> plugin instanceof FailsafePlugin) + .map(plugin -> (FailsafePlugin) plugin) + .toList(); + + assertThat("There should be exactly one FailsafePlugin", failsafePlugins, hasSize(1)); + + final List requestPolicies = getRequestPolicies(failsafePlugins.get(0)); + assertThat(requestPolicies, contains(asList( + instanceOf(RequestPolicy.class), // DefaultRequestPolicy wrapping Timeout + instanceOf(RequestPolicy.class)))); // ConditionalRequestPolicy wrapping BackupRequest + assertThat(requestPolicies.get(0).getClass().getSimpleName(), is("DefaultRequestPolicy")); + assertThat(requestPolicies.get(1).getClass().getSimpleName(), is("ConditionalRequestPolicy")); } @Test @@ -109,6 +126,13 @@ void shouldUseOpenTelemetryPlugin() throws Exception { assertThat(getPlugins(github), hasItem(instanceOf(OpenTelemetryPlugin.class))); } + @SuppressWarnings("unchecked") + private List getRequestPolicies(final FailsafePlugin failsafePlugin) throws Exception { + final Field field = FailsafePlugin.class.getDeclaredField("policies"); + field.setAccessible(true); + return (List) field.get(failsafePlugin); + } + private List getPlugins(final Http http) throws Exception { final Field field = http.getClass().getDeclaredField("plugin"); field.setAccessible(true); diff --git a/riptide-spring-boot-autoconfigure/src/test/resources/application-default.yml b/riptide-spring-boot-autoconfigure/src/test/resources/application-default.yml index 9a274c5d1..5d69ac9e2 100644 --- a/riptide-spring-boot-autoconfigure/src/test/resources/application-default.yml +++ b/riptide-spring-boot-autoconfigure/src/test/resources/application-default.yml @@ -92,6 +92,7 @@ riptide: failure-threshold: 1 success-threshold: 1 baz: + base-url: http://baz backup-request: enabled: true delay: 100 milliseconds @@ -110,36 +111,19 @@ riptide: retry: enabled: true max-retries: 2 - threads: - max-size: 13 - min-size: 2 - enabled: true - keep-alive: 5 minutes - queue-size: 10 circuit-breaker: enabled: true failure-rate-threshold: 3 in 5 seconds success-threshold: 1 - threads: - max-size: 10 - min-size: 2 - enabled: true - keep-alive: 5 minutes - queue-size: 10 timeouts: enabled: true global: 5 seconds - threads: - max-size: 11 - min-size: 2 - enabled: true - keep-alive: 5 minutes - queue-size: 10 backup-request: enabled: true delay: 75 milliseconds + failsafe: threads: - max-size: 12 + max-size: 14 min-size: 2 enabled: true keep-alive: 5 minutes From bbaa151a3e6d680cad699cc5fcabd7ea5ecce848 Mon Sep 17 00:00:00 2001 From: Ramesh Malla Date: Mon, 31 Aug 2026 13:46:13 +0200 Subject: [PATCH 2/2] Update readme Signed-off-by: Ramesh Malla --- riptide-spring-boot-autoconfigure/README.md | 75 ++++++++------------- 1 file changed, 27 insertions(+), 48 deletions(-) diff --git a/riptide-spring-boot-autoconfigure/README.md b/riptide-spring-boot-autoconfigure/README.md index 4c279be4e..840ddbe2a 100644 --- a/riptide-spring-boot-autoconfigure/README.md +++ b/riptide-spring-boot-autoconfigure/README.md @@ -27,17 +27,12 @@ riptide.clients: enabled: true fixed-delay: 50 milliseconds max-retries: 5 - threads: - max-size: 10 - min-size: 2 - enabled: true - keep-alive: 5 minutes - queue-size: 10 circuit-breaker: enabled: true failure-threshold: 3 out of 5 delay: 30 seconds success-threshold: 5 out of 5 + failsafe: threads: max-size: 10 min-size: 2 @@ -298,36 +293,19 @@ riptide: max-retries: 5 max-duration: 2 seconds jitter: 25 milliseconds - threads: - max-size: 10 - min-size: 2 - enabled: true - keep-alive: 5 minutes - queue-size: 10 circuit-breaker: enabled: true failure-threshold: 3 out of 5 failure-rate-threshold: 3 out of 5 in 5 seconds delay: 30 seconds success-threshold: 5 out of 5 - threads: - max-size: 10 - min-size: 2 - enabled: true - keep-alive: 5 minutes - queue-size: 10 backup-request: enabled: true delay: 75 milliseconds - threads: - max-size: 10 - min-size: 2 - enabled: true - keep-alive: 5 minutes - queue-size: 10 timeouts: enabled: true global: 500 milliseconds + failsafe: threads: max-size: 10 min-size: 2 @@ -413,6 +391,13 @@ For a complete overview of available properties, they type and default value ple | `│   │   ├── max-per-route` | `int` | `20` | | `│   │   ├── max-total` | `int` | `20` (or at least `max-per-route`); a warning is logged if the configured value is overridden | | `│   │   └── mode` | `String` | `streaming` (alternative is `buffering`) | +| `│ ├── failsafe` | | shared thread pool for `retry`/`circuit-breaker`/`backup-request`/`timeouts`; see [custom executor](#customization) | +| `│ │ └── threads` | | | +| `│ │ ├── enabled` | `boolean` | `false` | +| `│ │ ├── min-size` | `int` | `1` | +| `│ │ ├── max-size` | `int` | `1` | +| `│ │ ├── keep-alive` | `TimeSpan` | `1 minute` | +| `│ │ └── queue-size` | `int` | `0` (no queue) | | `│   ├── logging` | | | | `│   │   └── enabled` | `boolean` | `false` | | `│   ├── metrics` | | | @@ -501,6 +486,13 @@ For a complete overview of available properties, they type and default value ple | `        │   ├── time-to-live` | `TimeSpan` | see `defaults` | | `        │   ├── max-per-route` | `int` | see `defaults` | | `        │   └── max-total` | `int` | see `defaults` | +| ` ├── failsafe` | | shared thread pool for `retry`/`circuit-breaker`/`backup-request`/`timeouts` | +| ` │ └── threads` | | | +| ` │ ├── enabled` | `boolean` | see `defaults` | +| ` │ ├── min-size` | `int` | see `defaults` | +| ` │ ├── max-size` | `int` | see `defaults` | +| ` │ ├── keep-alive` | `TimeSpan` | see `defaults` | +| ` │ └── queue-size` | `int` | see `defaults` | | `        ├── logging` | | | | `        │   └── enabled` | `boolean` | see `defaults` | | `        ├── metrics` | | | @@ -636,7 +628,7 @@ public ClientHttpMessageConverters exampleHttpMessageConverters() { The following code can be used if you cannot use your client name in the method name (e.g. your client name is `my-client`): ```java -@Bean(name = "my-clientCircuitBreakerExecutorService") +@Bean(name = "my-clientHttpMessageConverters") public ClientHttpMessageConverters httpMessageConverters() { return new ClientHttpMessageConverters(singletonList(new Jaxb2RootElementHttpMessageConverter())); } @@ -664,15 +656,19 @@ The following table shows all beans with their respective name (for the `example | `exampleFaultClassifier` | `FaultClassifier` | | `exampleCircuitBreakerListener` | `CircuitBreakerListener` | | `exampleAuthorizationProvider` | `AuthorizationProvider` | -| `exampleRetryPolicyExecutorService` | `ExecutorService` | -| `exampleCircuitBreakerExecutorService` | `ExecutorService` | -| `exampleBackupRequestExecutorService` | `ExecutorService` | -| `exampleTimeoutExecutorService` | `ExecutorService` | +| `exampleFailsafeExecutorService` | `ExecutorService` (shared by retry/circuit-breaker/backup-request/timeouts, only if `failsafe.threads.enabled`) | If you override a bean then all of its dependencies (see the [graph](#customization)), will **not** be registered, unless required by some other bean. -Riptide uses Failsafe underneath to manage resiliency flows, and Failsafe supports custom thread pool executors. For more details, refer to the [riptide-failsafe](https://github.com/zalando/riptide/tree/main/riptide-failsafe#custom-executor) documentation. To configure a custom thread pool executor for retry, circuit breaker, backup requests, and timeout features, follow the configuration steps below. +Riptide uses Failsafe underneath to manage resiliency flows. `retry`, `circuit-breaker`, `backup-request` and +`timeouts` are merged into a single `FailsafePlugin` per client and share one thread pool executor, configured via +`failsafe.threads`. For more details, refer to the [riptide-failsafe](https://github.com/zalando/riptide/tree/main/riptide-failsafe#custom-executor) documentation. + +> The previous per-policy `retry.threads`/`circuit-breaker.threads`/`backup-request.threads`/`timeouts.threads` +> settings are no longer supported. Configuring any of them now throws a `LegacyFailsafeThreadsException` at +> startup; migrate to `failsafe.threads` below. + ```yaml retry: enabled: true @@ -680,36 +676,19 @@ Riptide uses Failsafe underneath to manage resiliency flows, and Failsafe suppor max-retries: 5 max-duration: 2 seconds jitter: 25 milliseconds - threads: - max-size: 10 - min-size: 2 - enabled: true - keep-alive: 5 minutes - queue-size: 10 circuit-breaker: enabled: true failure-threshold: 3 out of 5 failure-rate-threshold: 3 out of 5 in 5 seconds delay: 30 seconds success-threshold: 5 out of 5 - threads: - max-size: 10 - min-size: 2 - enabled: true - keep-alive: 5 minutes - queue-size: 10 backup-request: enabled: true delay: 75 milliseconds - threads: - max-size: 10 - min-size: 2 - enabled: true - keep-alive: 5 minutes - queue-size: 10 timeouts: enabled: true global: 500 milliseconds + failsafe: threads: max-size: 10 min-size: 2