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 )