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.
This commit is contained in:
John Blum
2022-04-28 14:14:48 -07:00
parent 875c6411f3
commit fff6bce9d0
4 changed files with 67 additions and 72 deletions

View File

@@ -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<String> 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 extends Session> T save(@Nullable T session) {
protected @Nullable <T extends Session> 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<Instant> lastAccessedTime = new AtomicReference<>(null);
private final AtomicReference<String> 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.<String>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.<String>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.<String>getAttribute("attributeOne")).isEqualTo("testOne");
assertThat(session.<String>getAttribute("attributeTwo")).isEqualTo("testTwo");

View File

@@ -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<String> 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<String> sessionId = new AtomicReference<>(null);
private final DataSerializer sessionSerializer;
private final Region<Object, Session> 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.<String>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.<String>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

View File

@@ -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<String> 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",

View File

@@ -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
)