From 5614a7d7652b3c31f8b782a3aa6f4cd8af9e70d2 Mon Sep 17 00:00:00 2001 From: Marcus Da Coregio Date: Mon, 16 Oct 2023 10:35:26 -0300 Subject: [PATCH] Fix Reactive Redis Sessions Never Expiring Prior to this commit, we were using Mono#and(Mono) to execute operations sequentially on Redis. However, the and(...) operator does not guarantee that the operations will run sequentially. This commit changes from Mono#and to Mono#flatMap Closes gh-2464 --- .../ReactiveRedisSessionRepositoryITests.java | 40 ++++++++++++++++++- .../redis/ReactiveRedisSessionRepository.java | 7 ++-- 2 files changed, 43 insertions(+), 4 deletions(-) diff --git a/spring-session-data-redis/src/integration-test/java/org/springframework/session/data/redis/ReactiveRedisSessionRepositoryITests.java b/spring-session-data-redis/src/integration-test/java/org/springframework/session/data/redis/ReactiveRedisSessionRepositoryITests.java index ded40a9b..086febad 100644 --- a/spring-session-data-redis/src/integration-test/java/org/springframework/session/data/redis/ReactiveRedisSessionRepositoryITests.java +++ b/spring-session-data-redis/src/integration-test/java/org/springframework/session/data/redis/ReactiveRedisSessionRepositoryITests.java @@ -16,14 +16,17 @@ package org.springframework.session.data.redis; +import java.time.Duration; import java.time.Instant; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import reactor.core.publisher.Mono; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Configuration; +import org.springframework.data.redis.core.ReactiveHashOperations; import org.springframework.data.redis.core.ReactiveRedisOperations; import org.springframework.session.Session; import org.springframework.session.data.redis.ReactiveRedisSessionRepository.RedisSession; @@ -35,8 +38,11 @@ import org.springframework.test.util.ReflectionTestUtils; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatIllegalStateException; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.endsWith; import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.willAnswer; import static org.mockito.Mockito.reset; import static org.mockito.Mockito.spy; @@ -53,6 +59,13 @@ class ReactiveRedisSessionRepositoryITests extends AbstractRedisITests { @Autowired private ReactiveRedisSessionRepository repository; + private ReactiveRedisOperations sessionRedisOperations; + + @BeforeEach + void setup() { + this.sessionRedisOperations = this.repository.getSessionRedisOperations(); + } + @Test void saves() { RedisSession toSave = this.repository.createSession().block(); @@ -75,7 +88,8 @@ class ReactiveRedisSessionRepositoryITests extends AbstractRedisITests { assertThat(this.repository.findById(toSave.getId()).block()).isNull(); } - @Test // gh-1399 + @Test + // gh-1399 void saveMultipleTimes() { RedisSession session = this.repository.createSession().block(); session.setAttribute("attribute1", "value1"); @@ -245,6 +259,30 @@ class ReactiveRedisSessionRepositoryITests extends AbstractRedisITests { reset(spyOperations); } + // gh-2464 + @Test + void saveWhenPutAllIsDelayedThenExpireShouldBeSet() { + ReactiveRedisOperations spy = spy(this.sessionRedisOperations); + ReflectionTestUtils.setField(this.repository, "sessionRedisOperations", spy); + ReactiveHashOperations opsForHash = spy(this.sessionRedisOperations.opsForHash()); + given(spy.opsForHash()).willReturn(opsForHash); + willAnswer((invocation) -> Mono.delay(Duration.ofSeconds(1)).then((Mono) invocation.callRealMethod())) + .given(opsForHash).putAll(anyString(), any()); + RedisSession toSave = this.repository.createSession().block(); + + String expectedAttributeName = "a"; + String expectedAttributeValue = "b"; + + toSave.setAttribute(expectedAttributeName, expectedAttributeValue); + this.repository.save(toSave).block(); + + String id = toSave.getId(); + Duration expireDuration = this.sessionRedisOperations.getExpire("spring:session:sessions:" + id).block(); + + assertThat(expireDuration).isNotEqualTo(Duration.ZERO); + reset(spy); + } + @Configuration @EnableRedisWebSession static class Config extends BaseConfig { diff --git a/spring-session-data-redis/src/main/java/org/springframework/session/data/redis/ReactiveRedisSessionRepository.java b/spring-session-data-redis/src/main/java/org/springframework/session/data/redis/ReactiveRedisSessionRepository.java index 13fc623f..c90947d3 100644 --- a/spring-session-data-redis/src/main/java/org/springframework/session/data/redis/ReactiveRedisSessionRepository.java +++ b/spring-session-data-redis/src/main/java/org/springframework/session/data/redis/ReactiveRedisSessionRepository.java @@ -286,10 +286,11 @@ public class ReactiveRedisSessionRepository setTtl = ReactiveRedisSessionRepository.this.sessionRedisOperations.persist(sessionKey); } - return update.and(setTtl).and((s) -> { + Mono clearDelta = Mono.fromDirect((s) -> { this.delta.clear(); s.onComplete(); - }).then(); + }); + return update.flatMap((unused) -> setTtl).flatMap((unused) -> clearDelta).then(); } private Mono saveChangeSessionId() { @@ -312,7 +313,7 @@ public class ReactiveRedisSessionRepository String sessionKey = getSessionKey(sessionId); return ReactiveRedisSessionRepository.this.sessionRedisOperations.rename(originalSessionKey, sessionKey) - .and(replaceSessionId).onErrorResume((ex) -> { + .flatMap((unused) -> Mono.fromDirect(replaceSessionId)).onErrorResume((ex) -> { String message = NestedExceptionUtils.getMostSpecificCause(ex).getMessage(); return StringUtils.startsWithIgnoreCase(message, "ERR no such key"); }, (ex) -> Mono.empty());