From fff6bce9d08f6abac2bc020d651576025bbd6231 Mon Sep 17 00:00:00 2001 From: John Blum Date: Thu, 28 Apr 2022 14:14:48 -0700 Subject: [PATCH] Fix race condition in concurrent Session access Integeration Tests. The race condition involved the absence/presence of the Session ID in GemFire/Geode client/server Integration Tests when using client PROXY or CACHING_PROXY Regions to persist Session state. When the remote Region operation (e.g. put(..), when storing the new Session in the cache Region) occurred, the current Thread could then be blocked on an IO operation, which would then free up another MultithreadedTC test framework thread to continue, but the Session creating Thread may not have set theh Session ID required by the other threads in a timely manner. --- ...rentSessionOperationsIntegrationTests.java | 53 +++++++++++------- ...entCachingProxyRegionIntegrationTests.java | 55 ++++++++----------- ...singClientLocalRegionIntegrationTests.java | 22 +++----- ...singClientProxyRegionIntegrationTests.java | 9 +-- 4 files changed, 67 insertions(+), 72 deletions(-) diff --git a/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/AbstractConcurrentSessionOperationsIntegrationTests.java b/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/AbstractConcurrentSessionOperationsIntegrationTests.java index ece6829..968ffb4 100644 --- a/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/AbstractConcurrentSessionOperationsIntegrationTests.java +++ b/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/AbstractConcurrentSessionOperationsIntegrationTests.java @@ -23,6 +23,7 @@ import static org.springframework.data.gemfire.util.RuntimeExceptionFactory.newI import static org.springframework.session.data.gemfire.AbstractGemFireOperationsSessionRepository.GemFireSession; import java.time.Instant; +import java.util.Objects; import java.util.Optional; import java.util.concurrent.atomic.AtomicReference; @@ -64,6 +65,8 @@ public abstract class AbstractConcurrentSessionOperationsIntegrationTests extend private final AbstractConcurrentSessionOperationsIntegrationTests testInstance; + private final AtomicReference sessionId = new AtomicReference<>(null); + private final GemFireOperationsSessionRepository sessionRepository; protected AbstractConcurrentSessionOperationsTestCase( @@ -83,44 +86,51 @@ public abstract class AbstractConcurrentSessionOperationsIntegrationTests extend GemFireOperationsSessionRepository.class.getName(), ObjectUtils.nullSafeClassName(sessionRepository))); } - @NonNull @SuppressWarnings("unused") - protected AbstractConcurrentSessionOperationsIntegrationTests getTestInstance() { + @SuppressWarnings("unused") + protected @NonNull AbstractConcurrentSessionOperationsIntegrationTests getTestInstance() { return this.testInstance; } - @NonNull - protected GemFireOperationsSessionRepository getSessionRepository() { + protected @NonNull GemFireOperationsSessionRepository getSessionRepository() { return this.sessionRepository; } - @Nullable - protected Session findById(String id) { + protected @NonNull String getSessionId() { + return this.sessionId.get(); + } + + protected void setSessionId(@Nullable String sessionId) { + this.sessionId.set(sessionId); + } + + protected @Nullable Session findById(@NonNull String id) { return getSessionRepository().findById(id); } - @NonNull - protected Session newSession() { + protected @NonNull Session newSession() { return getSessionRepository().createSession(); } - @Nullable - protected T save(@Nullable T session) { + protected @Nullable T save(@Nullable T session) { getSessionRepository().save(session); return session; } + + protected void waitOnAvailableSessionId() { + AbstractConcurrentSessionOperationsIntegrationTests.waitOn(() -> Objects.nonNull(this.sessionId.get())); + } } @SuppressWarnings("unused") public static class ConcurrentSessionOperationsTestCase extends AbstractConcurrentSessionOperationsTestCase { private final AtomicReference lastAccessedTime = new AtomicReference<>(null); - private final AtomicReference sessionId = new AtomicReference<>(null); public ConcurrentSessionOperationsTestCase(AbstractConcurrentSessionOperationsIntegrationTests testInstance) { super(testInstance); } - // Creator Thread + // Session Creator Thread @SuppressWarnings("rawtypes") public void thread1() { @@ -142,7 +152,7 @@ public abstract class AbstractConcurrentSessionOperationsIntegrationTests extend save(session); - this.sessionId.set(session.getId()); + setSessionId(session.getId()); waitForTick(4); assertTick(4); @@ -153,18 +163,19 @@ public abstract class AbstractConcurrentSessionOperationsIntegrationTests extend save(session); } - // Modifier (Attribute) Thread + // Session Attribute Modifier Thread public void thread2() { Thread.currentThread().setName("User Session Two"); waitForTick(1); assertTick(1); + waitOnAvailableSessionId(); - Session session = findById(this.sessionId.get()); + Session session = findById(getSessionId()); assertThat(session).isNotNull(); - assertThat(session.getId()).isEqualTo(this.sessionId.get()); + assertThat(session.getId()).isEqualTo(getSessionId()); assertThat(session.isExpired()).isFalse(); assertThat(session.getAttributeNames()).containsOnly("attributeOne", "attributeTwo"); assertThat(session.getAttribute("attributeOne")).isEqualTo("testOne"); @@ -181,7 +192,7 @@ public abstract class AbstractConcurrentSessionOperationsIntegrationTests extend save(session); } - // Modifier (Timestamp) Thread + // Session Timestamp Modifier Thread public void thread3() { Thread.currentThread().setName("User Session Three"); @@ -189,10 +200,10 @@ public abstract class AbstractConcurrentSessionOperationsIntegrationTests extend waitForTick(1); assertTick(1); - Session session = findById(this.sessionId.get()); + Session session = findById(getSessionId()); assertThat(session).isNotNull(); - assertThat(session.getId()).isEqualTo(this.sessionId.get()); + assertThat(session.getId()).isEqualTo(getSessionId()); assertThat(session.isExpired()).isFalse(); assertThat(session.getAttributeNames()).containsOnly("attributeOne", "attributeTwo"); assertThat(session.getAttribute("attributeOne")).isEqualTo("testOne"); @@ -211,10 +222,10 @@ public abstract class AbstractConcurrentSessionOperationsIntegrationTests extend super.finish(); - Session session = findById(this.sessionId.get()); + Session session = findById(getSessionId()); assertThat(session).isNotNull(); - assertThat(session.getId()).isEqualTo(this.sessionId.get()); + assertThat(session.getId()).isEqualTo(getSessionId()); assertThat(session.getAttributeNames()).containsOnly("attributeOne", "attributeTwo", "attributeThree"); assertThat(session.getAttribute("attributeOne")).isEqualTo("testOne"); assertThat(session.getAttribute("attributeTwo")).isEqualTo("testTwo"); diff --git a/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegrationTests.java b/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegrationTests.java index 6bd0fb3..7ba813d 100644 --- a/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegrationTests.java +++ b/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegrationTests.java @@ -30,7 +30,6 @@ import java.io.IOException; import java.time.Instant; import java.util.Arrays; import java.util.Objects; -import java.util.concurrent.atomic.AtomicReference; import java.util.function.Predicate; import edu.umd.cs.mtc.TestFramework; @@ -54,6 +53,7 @@ import org.apache.geode.internal.InternalDataSerializer; import org.springframework.context.annotation.AnnotationConfigApplicationContext; import org.springframework.data.gemfire.config.annotation.CacheServerApplication; import org.springframework.data.gemfire.config.annotation.ClientCacheApplication; +import org.springframework.lang.NonNull; import org.springframework.session.Session; import org.springframework.session.data.gemfire.config.annotation.web.http.EnableGemFireHttpSession; import org.springframework.session.data.gemfire.config.annotation.web.http.GemFireHttpSessionConfiguration; @@ -90,8 +90,6 @@ import org.springframework.test.context.junit4.SpringRunner; public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegrationTests extends AbstractConcurrentSessionOperationsIntegrationTests { - private static final String GEMFIRE_LOG_LEVEL = "error"; - @Before public void setup() { @@ -128,10 +126,8 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration // Tests that 2 Threads share the same Session object reference and therefore see's each other's changes. public static class ConcurrentCachedSessionOperationsTestCase extends AbstractConcurrentSessionOperationsTestCase { - private final AtomicReference sessionId = new AtomicReference<>(null); - public ConcurrentCachedSessionOperationsTestCase( - ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegrationTests testInstance) { + @NonNull ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegrationTests testInstance) { super(testInstance); } @@ -153,7 +149,7 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration save(session); - this.sessionId.set(session.getId()); + setSessionId(session.getId()); } public void thread1() { @@ -164,10 +160,10 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration Instant beforeLastAccessedTime = Instant.now(); - Session session = findById(this.sessionId.get()); + Session session = findById(getSessionId()); assertThat(session).isNotNull(); - assertThat(session.getId()).isEqualTo(this.sessionId.get()); + assertThat(session.getId()).isEqualTo(getSessionId()); assertThat(session.getLastAccessedTime()).isAfterOrEqualTo(beforeLastAccessedTime); assertThat(session.getLastAccessedTime()).isBeforeOrEqualTo(Instant.now()); assertThat(session.isExpired()).isFalse(); @@ -192,10 +188,10 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration Instant beforeLastAccessedTime = Instant.now(); - Session session = findById(this.sessionId.get()); + Session session = findById(getSessionId()); assertThat(session).isNotNull(); - assertThat(session.getId()).isEqualTo(this.sessionId.get()); + assertThat(session.getId()).isEqualTo(getSessionId()); assertThat(session.getLastAccessedTime()).isAfterOrEqualTo(beforeLastAccessedTime); assertThat(session.getLastAccessedTime()).isBeforeOrEqualTo(Instant.now()); assertThat(session.isExpired()).isFalse(); @@ -217,8 +213,6 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration private static final String DATA_SERIALIZER_NOT_FOUND_EXCEPTION_MESSAGE = "No DataSerializer was found capable of de/serializing Sessions"; - private final AtomicReference sessionId = new AtomicReference<>(null); - private final DataSerializer sessionSerializer; private final Region sessions; @@ -232,6 +226,14 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration this.sessionSerializer = reregisterDataSerializer(resolveDataSerializer()); } + private DataSerializer reregisterDataSerializer(DataSerializer dataSerializer) { + + InternalDataSerializer.unregister(dataSerializer.getId()); + InternalDataSerializer._register(dataSerializer, false); + + return dataSerializer; + } + private DataSerializer resolveDataSerializer() { return Arrays.stream(nullSafeArray(InternalDataSerializer.getSerializers(), DataSerializer.class)) @@ -258,14 +260,6 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration }; } - private DataSerializer reregisterDataSerializer(DataSerializer dataSerializer) { - - InternalDataSerializer.unregister(dataSerializer.getId()); - InternalDataSerializer._register(dataSerializer, false); - - return dataSerializer; - } - private Session get(String id) { return this.sessions.get(id); } @@ -310,7 +304,7 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration getSessionRepository().commit(loadedSession); - this.sessionId.set(session.getId()); + setSessionId(session.getId()); } public void thread2() { @@ -319,11 +313,12 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration waitForTick(1); assertTick(1); + waitOnAvailableSessionId(); - Session session = get(this.sessionId.get()); + Session session = get(getSessionId()); assertThat(session).isInstanceOf(GemFireSession.class); - assertThat(session.getId()).isEqualTo(this.sessionId.get()); + assertThat(session.getId()).isEqualTo(getSessionId()); assertThat(session.isExpired()).isFalse(); assertThat(session.getAttributeNames()).containsOnly("attributeOne", "attributeTwo"); assertThat(session.getAttribute("attributeOne")).isEqualTo("testOne"); @@ -336,10 +331,10 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration @Override public void finish() { - Session session = get(this.sessionId.get()); + Session session = get(getSessionId()); assertThat(session).isNotNull(); - assertThat(session.getId()).isEqualTo(this.sessionId.get()); + assertThat(session.getId()).isEqualTo(getSessionId()); assertThat(session.isExpired()).isFalse(); assertThat(session.getAttributeNames()).containsOnly("attributeOne", "attributeTwo"); assertThat(session.getAttribute("attributeOne")).isEqualTo("testOne"); @@ -352,6 +347,7 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration verify(this.sessionSerializer, times(2)) .toData(isA(GemFireSession.class), isA(DataOutput.class)); + } catch (ClassNotFoundException | IOException ignore) { } } @@ -364,7 +360,7 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration // Tests fail when 'copyOnRead' is set to 'true'! //@ClientCacheApplication(copyOnRead = true, logLevel = GEMFIRE_LOG_LEVEL, subscriptionEnabled = true) - @ClientCacheApplication(logLevel = GEMFIRE_LOG_LEVEL, subscriptionEnabled = true) + @ClientCacheApplication(subscriptionEnabled = true) @EnableGemFireHttpSession( clientRegionShortcut = ClientRegionShortcut.CACHING_PROXY, poolName = "DEFAULT", @@ -373,10 +369,7 @@ public class ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegration ) static class GemFireClientConfiguration { } - @CacheServerApplication( - name = "ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegrationTests", - logLevel = GEMFIRE_LOG_LEVEL - ) + @CacheServerApplication(name = "ConcurrentSessionOperationsUsingClientCachingProxyRegionIntegrationTests") @EnableGemFireHttpSession( regionName = "Sessions", sessionSerializerBeanName = GemFireHttpSessionConfiguration.SESSION_DATA_SERIALIZER_BEAN_NAME diff --git a/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests.java b/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests.java index 73559b9..60d5fa1 100644 --- a/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests.java +++ b/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests.java @@ -17,18 +17,17 @@ package org.springframework.session.data.gemfire; import static org.assertj.core.api.Assertions.assertThat; -import java.util.concurrent.atomic.AtomicReference; - -import edu.umd.cs.mtc.TestFramework; - import org.junit.Test; import org.junit.runner.RunWith; +import edu.umd.cs.mtc.TestFramework; + import org.apache.geode.cache.Region; import org.apache.geode.cache.client.ClientCache; import org.apache.geode.cache.client.ClientRegionShortcut; import org.springframework.data.gemfire.config.annotation.ClientCacheApplication; +import org.springframework.lang.NonNull; import org.springframework.session.Session; import org.springframework.session.data.gemfire.config.annotation.web.http.EnableGemFireHttpSession; import org.springframework.session.data.gemfire.config.annotation.web.http.GemFireHttpSessionConfiguration; @@ -56,8 +55,6 @@ import org.springframework.test.context.junit4.SpringRunner; public class ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests extends AbstractConcurrentSessionOperationsIntegrationTests { - private static final String GEMFIRE_LOG_LEVEL = "error"; - @Test public void concurrentLocalSessionAccessIsCorrect() throws Throwable { TestFramework.runOnce(new ConcurrentLocalSessionAccessTestCase(this)); @@ -66,10 +63,8 @@ public class ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests @SuppressWarnings("unused") public static class ConcurrentLocalSessionAccessTestCase extends AbstractConcurrentSessionOperationsTestCase { - private final AtomicReference sessionId = new AtomicReference<>(null); - public ConcurrentLocalSessionAccessTestCase( - ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests testInstance) { + @NonNull ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests testInstance) { super(testInstance); } @@ -89,7 +84,7 @@ public class ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests save(session); - this.sessionId.set(session.getId()); + setSessionId(session.getId()); waitForTick(2); assertTick(2); @@ -105,11 +100,12 @@ public class ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests waitForTick(1); assertTick(1); + waitOnAvailableSessionId(); - Session session = findById(this.sessionId.get()); + Session session = findById(getSessionId()); assertThat(session).isNotNull(); - assertThat(session.getId()).isEqualTo(this.sessionId.get()); + assertThat(session.getId()).isEqualTo(getSessionId()); assertThat(session.isExpired()).isFalse(); assertThat(session.getAttributeNames()).isEmpty(); @@ -122,7 +118,7 @@ public class ConcurrentSessionOperationsUsingClientLocalRegionIntegrationTests } } - @ClientCacheApplication(logLevel = GEMFIRE_LOG_LEVEL) + @ClientCacheApplication @EnableGemFireHttpSession( clientRegionShortcut = ClientRegionShortcut.LOCAL, poolName = "DEFAULT", diff --git a/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientProxyRegionIntegrationTests.java b/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientProxyRegionIntegrationTests.java index bbfdd8c..51c9068 100644 --- a/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientProxyRegionIntegrationTests.java +++ b/spring-session-data-geode/src/integration-test/java/org/springframework/session/data/gemfire/ConcurrentSessionOperationsUsingClientProxyRegionIntegrationTests.java @@ -57,14 +57,12 @@ import org.springframework.test.context.junit4.SpringRunner; public class ConcurrentSessionOperationsUsingClientProxyRegionIntegrationTests extends AbstractConcurrentSessionOperationsIntegrationTests { - private static final String GEMFIRE_LOG_LEVEL = "error"; - @BeforeClass public static void startGemFireServer() throws IOException { startGemFireServer(GemFireServerConfiguration.class); } - @ClientCacheApplication(logLevel = GEMFIRE_LOG_LEVEL, subscriptionEnabled = true) + @ClientCacheApplication(subscriptionEnabled = true) @EnableGemFireHttpSession( clientRegionShortcut = ClientRegionShortcut.PROXY, poolName = "DEFAULT", @@ -72,10 +70,7 @@ public class ConcurrentSessionOperationsUsingClientProxyRegionIntegrationTests ) static class GemFireClientConfiguration { } - @CacheServerApplication( - name = "ConcurrentSessionOperationsUsingClientProxyRegionIntegrationTests", - logLevel = GEMFIRE_LOG_LEVEL - ) + @CacheServerApplication(name = "ConcurrentSessionOperationsUsingClientProxyRegionIntegrationTests") @EnableGemFireHttpSession( sessionSerializerBeanName = GemFireHttpSessionConfiguration.SESSION_DATA_SERIALIZER_BEAN_NAME )