diff --git a/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/config/SubscriptionEnabledClientServerIntegrationTestConfiguration.java b/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/config/SubscriptionEnabledClientServerIntegrationTestConfiguration.java index ad02a0e..33a17d9 100644 --- a/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/config/SubscriptionEnabledClientServerIntegrationTestConfiguration.java +++ b/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/config/SubscriptionEnabledClientServerIntegrationTestConfiguration.java @@ -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 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 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 poolName = clientRegionFactoryBean.getPoolName() - .filter(this::isNotDefaultPool); - - if (poolName.isPresent()) { - this.poolNames.add(poolName.get()); - return true; - } + .filter(StringUtils::hasText); Optional 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 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 resolvePools() { - return Optional.ofNullable(poolName) - .filter(StringUtils::hasText) - .map(PoolManager::find) - .filter(Objects::nonNull) - .orElse(null); + eagerlyInitializeSpringManagedPoolBeans(); + + return nullSafeMap(PoolManager.getAll()).values(); } - private Iterable resolvePoolLocators(Pool pool) { - return pool != null ? pool.getLocators() : Collections.emptyList(); - } + private void eagerlyInitializeSpringManagedPoolBeans() { - private Set resolvePoolNames() { - - if (this.poolNames.isEmpty()) { - this.poolNames.add(GEMFIRE_DEFAULT_POOL_NAME); - } - - return this.poolNames; - } - - private Iterable resolvePoolServers(Pool pool) { - return pool != null ? pool.getServers() : Collections.emptyList(); + beanFactory.getBeansOfType(PoolFactoryBean.class).keySet() + .forEach(beanName -> beanFactory.getBean(beanName, Pool.class)); } }; }