Revise contribution

- Add tests for @⁠CachePut, @⁠Cacheable, and
  @⁠CacheEvict with Flux resubscription in
  ReactiveCachingTests, using plain Reactor retry() instead of
  @⁠Retryable.
- Reduce the block() timeouts in ReactiveRetryInterceptorTests from 5
  seconds to 1 second so that regressions fail faster.

See gh-37403
This commit is contained in:
Sam Brannen
2026-10-06 16:09:00 +02:00
parent 440139fec1
commit 1ef3d070c9
2 changed files with 93 additions and 2 deletions
@@ -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<String> cacheable(String key) {
return failOnFirstSubscription(this.cacheableSubscriptions);
}
@CachePut
Flux<String> cachePut(String key) {
return failOnFirstSubscription(this.cachePutSubscriptions);
}
@CacheEvict
Flux<String> cacheEvict(String key) {
return failOnFirstSubscription(this.cacheEvictSubscriptions);
}
private static Flux<String> 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 {
@@ -208,7 +208,7 @@ class ReactiveRetryInterceptorTests {
CacheableAnnotatedBean target = (CacheableAnnotatedBean) AopProxyUtils.getSingletonTarget(proxy);
// Simulates a failure on the 1st subscription to the target Flux.
List<String> result = proxy.retryOperation().collectList().block(Duration.ofSeconds(5));
List<String> 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<String> result = proxy.retryOperation().collectList().block(Duration.ofSeconds(5));
List<String> 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.