diff --git a/spring-context/src/test/java/org/springframework/cache/annotation/ReactiveCachingTests.java b/spring-context/src/test/java/org/springframework/cache/annotation/ReactiveCachingTests.java index 0ba7898253c..8e5a1979e35 100644 --- a/spring-context/src/test/java/org/springframework/cache/annotation/ReactiveCachingTests.java +++ b/spring-context/src/test/java/org/springframework/cache/annotation/ReactiveCachingTests.java @@ -16,10 +16,12 @@ package org.springframework.cache.annotation; +import java.time.Duration; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; import java.util.function.Supplier; @@ -31,6 +33,7 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.test.StepVerifier; +import org.springframework.aop.framework.AopProxyUtils; import org.springframework.cache.Cache; import org.springframework.cache.CacheManager; import org.springframework.cache.concurrent.ConcurrentMapCache; @@ -255,6 +258,62 @@ class ReactiveCachingTests { .verify(); } + @Test // gh-37403 + void cacheableFluxCanBeResubscribedAfterError() { + AnnotationConfigApplicationContext ctx = new AnnotationConfigApplicationContext( + EarlyCacheHitDeterminationConfig.class, ReactiveResubscriptionService.class); + ReactiveResubscriptionService service = ctx.getBean(ReactiveResubscriptionService.class); + ReactiveResubscriptionService target = + (ReactiveResubscriptionService) AopProxyUtils.getSingletonTarget(service); + + // The 1st subscription fails and the 2nd subscription (retry) succeeds. + StepVerifier.create(service.cacheable("key").retry(1)) + .expectNext("a", "b") + .expectComplete() + .verify(Duration.ofSeconds(1)); + + assertThat(target.cacheableSubscriptions).hasValue(2); + assertThat(ctx.getBean(CacheManager.class).getCache("first").get("key").get()) + .isEqualTo(List.of("a", "b")); + } + + @Test // gh-37403 + void cachePutFluxCanBeResubscribedAfterError() { + AnnotationConfigApplicationContext ctx = new AnnotationConfigApplicationContext( + EarlyCacheHitDeterminationConfig.class, ReactiveResubscriptionService.class); + ReactiveResubscriptionService service = ctx.getBean(ReactiveResubscriptionService.class); + ReactiveResubscriptionService target = + (ReactiveResubscriptionService) AopProxyUtils.getSingletonTarget(service); + + StepVerifier.create(service.cachePut("key").retry(1)) + .expectNext("a", "b") + .expectComplete() + .verify(Duration.ofSeconds(1)); + + assertThat(target.cachePutSubscriptions).hasValue(2); + assertThat(ctx.getBean(CacheManager.class).getCache("first").get("key").get()) + .isEqualTo(List.of("a", "b")); + } + + @Test // gh-37403 + void cacheEvictFluxCanBeResubscribedAfterError() { + AnnotationConfigApplicationContext ctx = new AnnotationConfigApplicationContext( + EarlyCacheHitDeterminationConfig.class, ReactiveResubscriptionService.class); + ReactiveResubscriptionService service = ctx.getBean(ReactiveResubscriptionService.class); + ReactiveResubscriptionService target = + (ReactiveResubscriptionService) AopProxyUtils.getSingletonTarget(service); + Cache cache = ctx.getBean(CacheManager.class).getCache("first"); + cache.put("key", List.of("a", "b")); + + StepVerifier.create(service.cacheEvict("key").retry(1)) + .expectNext("a", "b") + .expectComplete() + .verify(Duration.ofSeconds(1)); + + assertThat(target.cacheEvictSubscriptions).hasValue(2); + assertThat(cache.get("key")).isNull(); + } + @CacheConfig("first") static class ReactiveCacheableService { @@ -374,6 +433,38 @@ class ReactiveCachingTests { } + @CacheConfig("first") + static class ReactiveResubscriptionService { + + final AtomicInteger cacheableSubscriptions = new AtomicInteger(); + + final AtomicInteger cachePutSubscriptions = new AtomicInteger(); + + final AtomicInteger cacheEvictSubscriptions = new AtomicInteger(); + + @Cacheable + Flux cacheable(String key) { + return failOnFirstSubscription(this.cacheableSubscriptions); + } + + @CachePut + Flux cachePut(String key) { + return failOnFirstSubscription(this.cachePutSubscriptions); + } + + @CacheEvict + Flux cacheEvict(String key) { + return failOnFirstSubscription(this.cacheEvictSubscriptions); + } + + private static Flux failOnFirstSubscription(AtomicInteger subscriptions) { + return Flux.defer(() -> (subscriptions.incrementAndGet() == 1 ? + Flux.error(new IllegalStateException("1st subscription fails")) : + Flux.just("a", "b"))); + } + } + + @Configuration(proxyBeanMethods = false) @EnableCaching static class EarlyCacheHitDeterminationConfig { diff --git a/spring-context/src/test/java/org/springframework/resilience/ReactiveRetryInterceptorTests.java b/spring-context/src/test/java/org/springframework/resilience/ReactiveRetryInterceptorTests.java index e4e7a8140d0..41f3ec86baf 100644 --- a/spring-context/src/test/java/org/springframework/resilience/ReactiveRetryInterceptorTests.java +++ b/spring-context/src/test/java/org/springframework/resilience/ReactiveRetryInterceptorTests.java @@ -208,7 +208,7 @@ class ReactiveRetryInterceptorTests { CacheableAnnotatedBean target = (CacheableAnnotatedBean) AopProxyUtils.getSingletonTarget(proxy); // Simulates a failure on the 1st subscription to the target Flux. - List result = proxy.retryOperation().collectList().block(Duration.ofSeconds(5)); + List result = proxy.retryOperation().collectList().block(Duration.ofSeconds(1)); assertThat(result).containsExactly("a", "b"); // Subscribed twice: the 2nd attempt re-subscribed to the Flux returned from // the cache interceptor and reached the target Flux again instead of hanging. @@ -232,7 +232,7 @@ class ReactiveRetryInterceptorTests { cache.put(SimpleKey.EMPTY, List.of("a", "b")); // Simulates a failure on the 1st subscription to the target Flux. - List result = proxy.retryOperation().collectList().block(Duration.ofSeconds(5)); + List result = proxy.retryOperation().collectList().block(Duration.ofSeconds(1)); assertThat(result).containsExactly("a", "b"); // Subscribed twice: the 2nd attempt re-subscribed to the Flux returned from // the cache interceptor and reached the target Flux again instead of hanging.