From 83dbec58be57e2cafe05e264be77788a290bd7aa Mon Sep 17 00:00:00 2001 From: John Blum Date: Thu, 24 May 2018 17:51:50 -0700 Subject: [PATCH] Add proper handing and support for subscription-enabled, client/server integration tests configuration. --- .../ClientServerIntegrationTestsSupport.java | 5 +- .../integration/IntegrationTestsSupport.java | 12 +- ...entServerIntegrationTestConfiguration.java | 176 ++++++++++++++++++ 3 files changed, 187 insertions(+), 6 deletions(-) create mode 100644 spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/config/SubscriptionEnabledClientServerIntegrationTestConfiguration.java diff --git a/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/ClientServerIntegrationTestsSupport.java b/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/ClientServerIntegrationTestsSupport.java index 7f92ae3..3491b59 100644 --- a/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/ClientServerIntegrationTestsSupport.java +++ b/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/ClientServerIntegrationTestsSupport.java @@ -58,11 +58,12 @@ import org.springframework.data.gemfire.util.CollectionUtils; @SuppressWarnings("unused") public abstract class ClientServerIntegrationTestsSupport extends IntegrationTestsSupport { + public static final String DEFAULT_HOSTNAME = "localhost"; + public static final String GEMFIRE_CACHE_SERVER_PORT_PROPERTY = "spring.data.gemfire.cache.server.port"; + protected static final String DEBUG_ENDPOINT = "-agentlib:jdwp=transport=dt_socket,server=y,suspend=y,address=5005"; protected static final String DEBUGGING_ENABLED_PROPERTY = "spring.data.gemfire.debugging.enabled"; - protected static final String DEFAULT_HOSTNAME = "localhost"; protected static final String DIRECTORY_DELETE_ON_EXIT_PROPERTY = "spring.data.gemfire.directory.delete-on-exit"; - protected static final String GEMFIRE_CACHE_SERVER_PORT_PROPERTY = "spring.data.gemfire.cache.server.port"; protected static final String PROCESS_RUN_MANUAL_PROPERTY = "spring.data.gemfire.process.run-manual"; protected static final String SYSTEM_PROPERTIES_LOG_FILE = "system-properties.log"; diff --git a/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/IntegrationTestsSupport.java b/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/IntegrationTestsSupport.java index 9064772..6654f78 100644 --- a/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/IntegrationTestsSupport.java +++ b/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/IntegrationTestsSupport.java @@ -39,9 +39,9 @@ public abstract class IntegrationTestsSupport { protected static final String GEMFIRE_LOG_FILE = "gemfire-server.log"; protected static final String GEMFIRE_LOG_FILE_PROPERTY = "spring.data.gemfire.log.file"; - protected static final String GEMFIRE_LOG_LEVEL = "warning"; + protected static final String GEMFIRE_LOG_LEVEL = "error"; protected static final String GEMFIRE_LOG_LEVEL_PROPERTY = "spring.data.gemfire.log.level"; - protected static final String TEST_GEMFIRE_LOG_LEVEL = "warning"; + protected static final String TEST_GEMFIRE_LOG_LEVEL = "error"; @BeforeClass public static void closeAnyExistingGemFireCacheInstanceBeforeTestExecution() { @@ -73,8 +73,12 @@ public abstract class IntegrationTestsSupport { } private static GemFireCache close(GemFireCache cache) { - cache.close(); - return cache; + + return Optional.ofNullable(cache) + .map(it -> { + cache.close(); + return cache; + }).orElse(cache); } protected static boolean waitOn(Condition condition) { 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 new file mode 100644 index 0000000..8919c60 --- /dev/null +++ b/spring-test-data-geode/src/main/java/org/springframework/data/gemfire/tests/integration/config/SubscriptionEnabledClientServerIntegrationTestConfiguration.java @@ -0,0 +1,176 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing + * permissions and limitations under the License. + */ + +package org.springframework.data.gemfire.tests.integration.config; + +import static org.springframework.data.gemfire.util.ArrayUtils.nullSafeArray; + +import java.util.Arrays; +import java.util.Objects; +import java.util.Optional; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.geode.cache.Region; +import org.apache.geode.cache.client.ClientCache; +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.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.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.config.annotation.ClientCacheConfigurer; +import org.springframework.data.gemfire.config.xml.GemfireConstants; +import org.springframework.data.gemfire.tests.integration.ClientServerIntegrationTestsSupport; +import org.springframework.util.Assert; + +/** + * The SubscriptionEnabledClientServerIntegrationTestConfiguration class... + * + * @author John Blum + * @since 1.0.0 + */ +@Configuration +@SuppressWarnings("unused") +public class SubscriptionEnabledClientServerIntegrationTestConfiguration { + + private static final long DEFAULT_TIMEOUT = TimeUnit.SECONDS.toMillis(60); + + private static final CountDownLatch LATCH = new CountDownLatch(1); + + private static final String GEMFIRE_CACHE_SERVER_PORT_PROPERTY = + ClientServerIntegrationTestsSupport.GEMFIRE_CACHE_SERVER_PORT_PROPERTY; + + private static final String GEMFIRE_DEFAULT_POOL_NAME = "DEFAULT"; + + private static final String LOCALHOST = ClientServerIntegrationTestsSupport.DEFAULT_HOSTNAME; + + @Bean + BeanPostProcessor clientServerReadyBeanPostProcessor( + @Value("${" + GEMFIRE_CACHE_SERVER_PORT_PROPERTY + ":40404}") int port) { + + return new BeanPostProcessor() { + + private final AtomicBoolean checkGemFireServerIsRunning = new AtomicBoolean(true); + + private final AtomicReference defaultPool = new AtomicReference<>(null); + + @SuppressWarnings("all") + public Object postProcessBeforeInitialization(Object bean, String beanName) throws BeansException { + + if (shouldCheckWhetherGemFireServerIsRunning(bean, beanName)) { + try { + validateClientCacheNotified(); + validateClientCacheSubscriptionQueueConnectionEstablished(); + } + catch (InterruptedException cause) { + Thread.currentThread().interrupt(); + } + } + + return bean; + } + + private boolean shouldCheckWhetherGemFireServerIsRunning(Object bean, String beanName) { + + return isGemFireRegion(bean, beanName) + ? checkGemFireServerIsRunning.compareAndSet(true, false) + : whenGemFireCache(bean, beanName); + } + + private boolean isGemFireRegion(Object bean, String beanName) { + return bean instanceof Region; + } + + private boolean whenGemFireCache(Object bean, String beanName) { + + if (bean instanceof ClientCache) { + defaultPool.compareAndSet(null, ((ClientCache) bean).getDefaultPool()); + } + + return false; + } + + private void validateClientCacheNotified() throws InterruptedException { + + boolean didNotTimeout = LATCH.await(DEFAULT_TIMEOUT, TimeUnit.MILLISECONDS); + + Assert.state(didNotTimeout, String.format( + "Apache Geode CacheServer failed to start on host [%s] and port [%d]", LOCALHOST, port)); + } + + @SuppressWarnings("all") + private void validateClientCacheSubscriptionQueueConnectionEstablished() throws InterruptedException { + + boolean clientCacheSubscriptionQueueConnectionEstablished = false; + + Pool pool = resolvePool(this.defaultPool.get(), + GemfireConstants.DEFAULT_GEMFIRE_POOL_NAME, GEMFIRE_DEFAULT_POOL_NAME); + + if (pool instanceof PoolImpl) { + + long timeout = System.currentTimeMillis() + DEFAULT_TIMEOUT; + + while (System.currentTimeMillis() < timeout && !((PoolImpl) pool).isPrimaryUpdaterAlive()) { + synchronized (pool) { + TimeUnit.MILLISECONDS.timedWait(pool, 500L); + } + + } + + clientCacheSubscriptionQueueConnectionEstablished |= ((PoolImpl) pool).isPrimaryUpdaterAlive(); + } + + Assert.state(clientCacheSubscriptionQueueConnectionEstablished, + String.format("ClientCache subscription queue connection not established;" + + " Apache Geode Pool was [%s];" + + " Apache Geode Pool configuration was [locators = %s, servers = %s]", + pool, pool.getLocators(), pool.getServers())); + } + + private Pool resolvePool(Pool pool, String... poolNames) { + + return Optional.ofNullable(pool) + .orElseGet(() -> Arrays.stream(nullSafeArray(poolNames, String.class)) + .map(PoolManager::find) + .filter(Objects::nonNull) + .findFirst() + .orElse(null)); + } + }; + } + + @Bean + ClientCacheConfigurer registerClientMembershipListener() { + + return (beanName, bean) -> + + ClientMembership.registerClientMembershipListener(new ClientMembershipListenerAdapter() { + + @Override + public void memberJoined(ClientMembershipEvent event) { + LATCH.countDown(); + } + }); + } +}