Re-write integration test logic verifying subscription queue connections have been established between client and server.

This commit is contained in:
John Blum
2018-05-26 18:37:33 -07:00
parent 79769628e7
commit 89731d7d31

View File

@@ -16,14 +16,11 @@
package org.springframework.data.gemfire.tests.integration.config;
import static org.springframework.data.gemfire.util.CollectionUtils.nullSafeSet;
import static org.springframework.data.gemfire.util.CollectionUtils.nullSafeMap;
import java.lang.reflect.Method;
import java.net.InetSocketAddress;
import java.util.Collections;
import java.util.Objects;
import java.util.Collection;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
@@ -32,15 +29,16 @@ import org.apache.geode.cache.client.ClientRegionShortcut;
import org.apache.geode.cache.client.Pool;
import org.apache.geode.cache.client.PoolManager;
import org.apache.geode.cache.client.internal.PoolImpl;
import org.apache.geode.internal.concurrent.ConcurrentHashSet;
import org.apache.geode.management.membership.ClientMembership;
import org.apache.geode.management.membership.ClientMembershipEvent;
import org.apache.geode.management.membership.ClientMembershipListenerAdapter;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.ListableBeanFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.beans.factory.config.BeanPostProcessor;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.gemfire.client.ClientCacheFactoryBean;
import org.springframework.data.gemfire.client.ClientRegionFactoryBean;
import org.springframework.data.gemfire.client.ClientRegionShortcutWrapper;
import org.springframework.data.gemfire.client.PoolFactoryBean;
@@ -49,6 +47,7 @@ import org.springframework.data.gemfire.config.xml.GemfireConstants;
import org.springframework.data.gemfire.listener.ContinuousQueryListenerContainer;
import org.springframework.data.gemfire.tests.integration.ClientServerIntegrationTestsSupport;
import org.springframework.data.gemfire.tests.util.ObjectUtils;
import org.springframework.lang.Nullable;
import org.springframework.util.Assert;
import org.springframework.util.ReflectionUtils;
import org.springframework.util.StringUtils;
@@ -66,11 +65,15 @@ import org.springframework.util.StringUtils;
* @see org.apache.geode.cache.client.PoolManager
* @see org.apache.geode.management.membership.ClientMembership
* @see org.apache.geode.management.membership.ClientMembershipListenerAdapter
* @see org.springframework.beans.factory.ListableBeanFactory
* @see org.springframework.beans.factory.config.BeanPostProcessor
* @see org.springframework.context.annotation.Bean
* @see org.springframework.context.annotation.Configuration
* @see org.springframework.data.gemfire.client.ClientCacheFactoryBean
* @see org.springframework.data.gemfire.client.ClientRegionFactoryBean
* @see org.springframework.data.gemfire.client.PoolFactoryBean
* @see org.springframework.data.gemfire.config.annotation.ClientCacheConfigurer
* @see org.springframework.data.gemfire.listener.ContinuousQueryListenerContainer
* @see org.springframework.data.gemfire.tests.integration.ClientServerIntegrationTestsSupport
* @since 1.0.0
*/
@@ -90,27 +93,23 @@ public class SubscriptionEnabledClientServerIntegrationTestConfiguration {
private static final String LOCALHOST = ClientServerIntegrationTestsSupport.DEFAULT_HOSTNAME;
protected Set<String> getTargetRegionBeans() {
return Collections.emptySet();
}
@Bean
BeanPostProcessor clientServerReadyBeanPostProcessor(
BeanPostProcessor clientServerReadyBeanPostProcessor(ListableBeanFactory beanFactory,
@Value("${" + GEMFIRE_CACHE_SERVER_PORT_PROPERTY + ":40404}") int port) {
return new BeanPostProcessor() {
private final AtomicBoolean verifyGemFireServerIsRunning = new AtomicBoolean(true);
private final Set<String> poolNames = new ConcurrentHashSet<>();
@Nullable
@Override
@SuppressWarnings("all")
public Object postProcessBeforeInitialization(Object bean, String beanName) throws BeansException {
public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
if (shouldVerifyGemFireServerIsRunning(bean, beanName)) {
try {
verifyClientCacheSubscriptionQueueConnectionsEstablished();
verifyClientCacheNotified();
verifyClientCacheSubscriptionQueueConnectionEstablished();
}
catch (InterruptedException cause) {
Thread.currentThread().interrupt();
@@ -127,72 +126,31 @@ public class SubscriptionEnabledClientServerIntegrationTestConfiguration {
}
private boolean isBeanOfImportance(Object bean, String beanName) {
return isContinuousQueryListenerContainer(bean, beanName) || isPool(bean, beanName)
|| isRegion(bean, beanName);
return isContinuousQueryListenerContainer(bean) || isProxyClientRegion(bean);
}
private boolean isContinuousQueryListenerContainer(Object bean, String beanName) {
if (bean instanceof ContinuousQueryListenerContainer) {
ContinuousQueryListenerContainer continuousQueryListenerContainer =
(ContinuousQueryListenerContainer) bean;
return Optional.ofNullable(continuousQueryListenerContainer.getPoolName())
.filter(StringUtils::hasText)
.map(this.poolNames::add)
.orElseGet(() -> {
this.poolNames.add(GEMFIRE_DEFAULT_POOL_NAME);
return true;
});
}
return false;
private boolean isClientCache(Object bean) {
return bean instanceof ClientCacheFactoryBean;
}
private boolean isPool(Object bean, String beanName) {
if (bean instanceof PoolFactoryBean) {
// TODO: uncomment in SD Lovelace and replace this.poolNames.add(beanName)
//this.poolNames.add(((PoolFactoryBean) bean).getName());
this.poolNames.add(beanName);
return true;
}
return false;
private boolean isContinuousQueryListenerContainer(Object bean) {
return bean instanceof ContinuousQueryListenerContainer;
}
private boolean isRegion(Object bean, String beanName) {
return isTargetRegionBean(beanName) || isProxyClientRegion(bean, beanName);
}
private boolean isTargetRegionBean(String beanName) {
return nullSafeSet(getTargetRegionBeans()).contains(beanName);
}
private boolean isProxyClientRegion(Object bean, String beanName) {
private boolean isProxyClientRegion(Object bean) {
if (bean instanceof ClientRegionFactoryBean) {
ClientRegionFactoryBean<?, ?> clientRegionFactoryBean = (ClientRegionFactoryBean<?, ?>) bean;
ClientRegionFactoryBean<?, ?> clientRegionFactoryBean = (ClientRegionFactoryBean) bean;
Optional<String> poolName = clientRegionFactoryBean.getPoolName()
.filter(this::isNotDefaultPool);
if (poolName.isPresent()) {
this.poolNames.add(poolName.get());
return true;
}
.filter(StringUtils::hasText);
Optional<ClientRegionShortcut> clientRegionShortcut =
resolveClientRegionShortcut(clientRegionFactoryBean)
.filter(this::isProxyClientRegion);
if (clientRegionShortcut.isPresent()) {
this.poolNames.add(GEMFIRE_DEFAULT_POOL_NAME);
return true;
}
return poolName.isPresent() || clientRegionShortcut.isPresent();
}
return false;
@@ -202,13 +160,9 @@ public class SubscriptionEnabledClientServerIntegrationTestConfiguration {
return ClientRegionShortcutWrapper.valueOf(clientRegionShortcut).isProxy();
}
private boolean isNotDefaultPool(String poolName) {
return !GEMFIRE_DEFAULT_POOL_NAME.equals(poolName);
}
@SuppressWarnings("unchecked")
private Optional<ClientRegionShortcut> resolveClientRegionShortcut(
ClientRegionFactoryBean<?, ?> clientRegionFactoryBean) {
ClientRegionFactoryBean<?, ?> clientRegionFactoryBean) {
try {
@@ -232,65 +186,42 @@ public class SubscriptionEnabledClientServerIntegrationTestConfiguration {
}
@SuppressWarnings("all")
private void verifyClientCacheSubscriptionQueueConnectionEstablished() throws InterruptedException {
private void verifyClientCacheSubscriptionQueueConnectionsEstablished() {
resolvePoolNames().stream()
.map(this::resolvePool)
resolvePools().stream()
.filter(pool -> pool instanceof PoolImpl)
.map(pool -> (PoolImpl) pool)
.forEach(pool -> {
boolean clientCacheSubscriptionQueueConnectionEstablished = false;
if (pool instanceof PoolImpl) {
long timeout = System.currentTimeMillis() + DEFAULT_TIMEOUT;
while (System.currentTimeMillis() < timeout && !((PoolImpl) pool).isPrimaryUpdaterAlive()) {
synchronized (pool) {
ObjectUtils.doOperationSafely(() -> {
TimeUnit.MILLISECONDS.timedWait(pool, 500L);
return null;
});
}
long timeout = System.currentTimeMillis() + DEFAULT_TIMEOUT;
while (System.currentTimeMillis() < timeout && !pool.isPrimaryUpdaterAlive()) {
synchronized (pool) {
ObjectUtils.doOperationSafely(() -> {
TimeUnit.MILLISECONDS.timedWait(pool, 500L);
return null;
});
}
clientCacheSubscriptionQueueConnectionEstablished =
((PoolImpl) pool).isPrimaryUpdaterAlive();
}
Assert.state(clientCacheSubscriptionQueueConnectionEstablished,
Assert.state(pool.isPrimaryUpdaterAlive(),
String.format("ClientCache subscription queue connection not established;"
+ " Pool [%s] has configuration [locators = %s, servers = %s]",
pool, resolvePoolLocators(pool), resolvePoolServers(pool)));
pool, pool.getLocators(), pool.getServers()));
});
}
private Pool resolvePool(String poolName) {
private Collection<Pool> resolvePools() {
return Optional.ofNullable(poolName)
.filter(StringUtils::hasText)
.map(PoolManager::find)
.filter(Objects::nonNull)
.orElse(null);
eagerlyInitializeSpringManagedPoolBeans();
return nullSafeMap(PoolManager.getAll()).values();
}
private Iterable<InetSocketAddress> resolvePoolLocators(Pool pool) {
return pool != null ? pool.getLocators() : Collections.emptyList();
}
private void eagerlyInitializeSpringManagedPoolBeans() {
private Set<String> resolvePoolNames() {
if (this.poolNames.isEmpty()) {
this.poolNames.add(GEMFIRE_DEFAULT_POOL_NAME);
}
return this.poolNames;
}
private Iterable<InetSocketAddress> resolvePoolServers(Pool pool) {
return pool != null ? pool.getServers() : Collections.emptyList();
beanFactory.getBeansOfType(PoolFactoryBean.class).keySet()
.forEach(beanName -> beanFactory.getBean(beanName, Pool.class));
}
};
}