diff --git a/src/main/java/org/springframework/data/gemfire/PeerRegionFactoryBean.java b/src/main/java/org/springframework/data/gemfire/PeerRegionFactoryBean.java index c113e939..9054d64c 100644 --- a/src/main/java/org/springframework/data/gemfire/PeerRegionFactoryBean.java +++ b/src/main/java/org/springframework/data/gemfire/PeerRegionFactoryBean.java @@ -15,12 +15,12 @@ */ package org.springframework.data.gemfire; -import static java.util.Arrays.stream; -import static org.springframework.data.gemfire.util.ArrayUtils.nullSafeArray; import static org.springframework.data.gemfire.util.RuntimeExceptionFactory.newIllegalArgumentException; +import java.util.ArrayList; import java.util.Arrays; import java.util.HashSet; +import java.util.List; import java.util.Objects; import java.util.Optional; import java.util.Set; @@ -48,7 +48,6 @@ import org.apache.geode.cache.Scope; import org.apache.geode.cache.asyncqueue.AsyncEventQueue; import org.apache.geode.cache.wan.GatewaySender; import org.apache.geode.compression.Compressor; -import org.apache.geode.internal.cache.UserSpecifiedRegionAttributes; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.FactoryBean; @@ -57,19 +56,22 @@ import org.springframework.core.io.Resource; import org.springframework.data.gemfire.client.ClientRegionFactoryBean; import org.springframework.data.gemfire.eviction.EvictingRegionFactoryBean; import org.springframework.data.gemfire.expiration.ExpiringRegionFactoryBean; +import org.springframework.data.gemfire.util.ArrayUtils; import org.springframework.data.gemfire.util.CollectionUtils; import org.springframework.data.gemfire.util.RegionUtils; +import org.springframework.data.gemfire.util.SpringUtils; +import org.springframework.lang.NonNull; +import org.springframework.lang.Nullable; import org.springframework.util.Assert; -import org.springframework.util.ObjectUtils; import org.springframework.util.ReflectionUtils; import org.springframework.util.StringUtils; /** - * Abstract Spring {@link FactoryBean} base class extended by other SDG {@link FactoryBean FactoryBeans} used to - * construct, configure and initialize peer {@link Region Regions}. + * Spring {@link FactoryBean} and abstract base class extended by other SDG {@link FactoryBean FactoryBeans} + * used to construct, configure and initialize {@literal peer} {@link Region Regions}. * - * This {@link FactoryBean} allows for very easy and flexible creation of peer {@link Region}. - * For client {@link Region Regions}, however, see the {@link ClientRegionFactoryBean}. + * This {@link FactoryBean} allows for very easy and flexible creation of {@literal peer} {@link Region Regions}. + * For {@literal client} {@link Region Regions}, see the {@link ClientRegionFactoryBean}. * * @author Costin Leau * @author David Turanski @@ -78,8 +80,11 @@ import org.springframework.util.StringUtils; * @see org.apache.geode.cache.CacheListener * @see org.apache.geode.cache.CacheLoader * @see org.apache.geode.cache.CacheWriter + * @see org.apache.geode.cache.CustomExpiry * @see org.apache.geode.cache.DataPolicy + * @see org.apache.geode.cache.DiskStore * @see org.apache.geode.cache.EvictionAttributes + * @see org.apache.geode.cache.ExpirationAttributes * @see org.apache.geode.cache.GemFireCache * @see org.apache.geode.cache.PartitionAttributes * @see org.apache.geode.cache.Region @@ -88,10 +93,15 @@ import org.springframework.util.StringUtils; * @see org.apache.geode.cache.RegionShortcut * @see org.apache.geode.cache.Scope * @see org.apache.geode.cache.asyncqueue.AsyncEventQueue + * @see org.apache.geode.cache.wan.GatewaySender + * @see org.apache.geode.compression.Compressor * @see org.springframework.beans.factory.DisposableBean * @see org.springframework.context.SmartLifecycle + * @see org.springframework.data.gemfire.ConfigurableRegionFactoryBean * @see org.springframework.data.gemfire.ResolvableRegionFactoryBean * @see org.springframework.data.gemfire.client.ClientRegionFactoryBean + * @see org.springframework.data.gemfire.eviction.EvictingRegionFactoryBean + * @see org.springframework.data.gemfire.expiration.ExpiringRegionFactoryBean * @see org.springframework.data.gemfire.config.annotation.RegionConfigurer */ @SuppressWarnings("unused") @@ -106,8 +116,6 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact private Boolean persistent; private Boolean statisticsEnabled; - private AsyncEventQueue[] asyncEventQueues; - private CacheListener[] cacheListeners; private CacheLoader cacheLoader; @@ -131,7 +139,10 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact private ExpirationAttributes regionIdleTimeout; private ExpirationAttributes regionTimeToLive; - private GatewaySender[] gatewaySenders; + private List asyncEventQueues = new ArrayList<>(); + private List gatewaySenders = new ArrayList<>(); + private List asyncEventQueueIds = new ArrayList<>(); + private List gatewaySenderIds = new ArrayList<>(); private RegionAttributes attributes; @@ -143,9 +154,6 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact private String diskStoreName; - private String[] asyncEventQueueIds; - private String[] gatewaySenderIds; - /** * Creates a new {@link Region} with the given {@link String name}. * @@ -168,33 +176,38 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact Region region = newRegion(regionFactory, getParent(), regionName); - return enableAsLockGrantor(region); - } - - private Region enableAsLockGrantor(Region region) { - - Optional.ofNullable(region) - .filter(it -> it.getAttributes().isLockGrantor()) - .ifPresent(Region::becomeLockGrantor); + region = becomeLockGrantor(region); return region; } + private Region becomeLockGrantor(Region region) { + + if (isLockGrantor(region)) { + region.becomeLockGrantor(); + } + + return region; + } + + private boolean isLockGrantor(@Nullable Region region) { + return region != null && region.getAttributes() != null && region.getAttributes().isLockGrantor(); + } + private Region newRegion(RegionFactory regionFactory, Region parentRegion, String regionName) { - return Optional.ofNullable(parentRegion) - .map(parent -> { + if (parentRegion != null) { - logInfo("Creating Subregion [%1$s] with parent Region [%2$s]", regionName, parent.getName()); + logInfo("Creating Subregion [%1$s] with parent Region [%2$s]", regionName, parentRegion.getName()); - return regionFactory.createSubregion(parent, regionName); - }) - .orElseGet(() -> { + return regionFactory.createSubregion(parentRegion, regionName); + } + else { - logInfo("Created Region [%s]", regionName); + logInfo("Created Region [%s]", regionName); - return regionFactory.create(regionName); - }); + return regionFactory.create(regionName); + } } private Cache resolveCache(GemFireCache gemfireCache) { @@ -262,7 +275,8 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact getConfiguredAsyncEventQueueIds().forEach(regionFactory::addAsyncEventQueueId); - stream(nullSafeArray(this.cacheListeners, CacheListener.class)).forEach(regionFactory::addCacheListener); + Arrays.stream(ArrayUtils.nullSafeArray(this.cacheListeners, CacheListener.class)) + .forEach(regionFactory::addCacheListener); Optional.ofNullable(this.cacheLoader).ifPresent(regionFactory::setCacheLoader); @@ -322,11 +336,12 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact Set asyncEventQueueIds = new HashSet<>(); - Arrays.stream(nullSafeArray(this.asyncEventQueues, AsyncEventQueue.class)) + CollectionUtils.nullSafeList(this.asyncEventQueues).stream() + .filter(Objects::nonNull) .map(AsyncEventQueue::getId) .collect(Collectors.toCollection(() -> asyncEventQueueIds)); - Arrays.stream(nullSafeArray(this.asyncEventQueueIds, String.class)) + CollectionUtils.nullSafeList(this.asyncEventQueueIds).stream() .filter(StringUtils::hasText) .map(String::trim) .collect(Collectors.toCollection(() -> asyncEventQueueIds)); @@ -338,11 +353,12 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact Set gatewaySenderIds = new HashSet<>(); - Arrays.stream(nullSafeArray(this.gatewaySenders, GatewaySender.class)) + CollectionUtils.nullSafeList(this.gatewaySenders).stream() + .filter(Objects::nonNull) .map(GatewaySender::getId) .collect(Collectors.toCollection(() -> gatewaySenderIds)); - Arrays.stream(nullSafeArray(this.gatewaySenderIds, String.class)) + CollectionUtils.nullSafeList(this.gatewaySenderIds).stream() .filter(StringUtils::hasText) .map(String::trim) .collect(Collectors.toCollection(() -> gatewaySenderIds)); @@ -351,8 +367,6 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact } /* - * (non-Javadoc) - * * This method is not considered part of the PeerRegionFactoryBean API and is strictly used for testing purposes! * * NOTE: Cannot pass RegionAttributes.class as the "targetType" in the second invocation of getFieldValue(..) @@ -365,7 +379,7 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact * @see org.apache.geode.cache.RegionAttributes#getDataPolicy * @see org.apache.geode.cache.DataPolicy */ - @SuppressWarnings({ "deprecation", "unchecked" }) + @SuppressWarnings({ "deprecation", "rawtypes", "unchecked" }) DataPolicy getDataPolicy(RegionFactory regionFactory, RegionShortcut regionShortcut) { return getFieldValue(regionFactory, "attrsFactory", AttributesFactory.class) @@ -430,7 +444,7 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact regionFactory.setEntryIdleTimeout(regionAttributes.getEntryIdleTimeout()); regionFactory.setEntryTimeToLive(regionAttributes.getEntryTimeToLive()); - // NOTE: EvictionAttributes are created by certain RegionShortcuts; need the null check! + // NOTE: EvictionAttributes are created by certain RegionShortcuts; null check needed! if (isUserSpecifiedEvictionAttributes(regionAttributes)) { regionFactory.setEvictionAttributes(regionAttributes.getEvictionAttributes()); } @@ -471,6 +485,7 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact * @see org.apache.geode.cache.RegionAttributes * @see org.apache.geode.cache.RegionFactory */ + @SuppressWarnings("rawtypes") protected void mergePartitionAttributes(RegionFactory regionFactory, RegionAttributes regionAttributes) { @@ -510,35 +525,40 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact allow &= getDataPolicy().withPersistence() || (getAttributes() != null && getAttributes().getEvictionAttributes() != null - && EvictionAction.OVERFLOW_TO_DISK.equals(attributes.getEvictionAttributes().getAction())); + && EvictionAction.OVERFLOW_TO_DISK.equals(this.attributes.getEvictionAttributes().getAction())); return allow; } /* - * (non-Javadoc) - * * This method is not part of the PeerRegionFactoryBean API and is strictly used for testing purposes! * +<<<<<<< HEAD * NOTE unfortunately, must resort to using a Pivotal GemFire internal class, ugh! * +======= +>>>>>>> f7643fe4a... DATAGEODE-368 - Add API to attach additional AsyncEventQueues and GatewaySenders to peer Regions. * @see org.apache.geode.internal.cache.UserSpecifiedRegionAttributes#hasEvictionAttributes */ + boolean isUserSpecifiedEvictionAttributes(RegionAttributes regionAttributes) { - boolean isUserSpecifiedEvictionAttributes(final RegionAttributes regionAttributes) { + SpringUtils.ValueReturningThrowableOperation hasEvictionAttributes = () -> + Optional.ofNullable(regionAttributes) + .map(Object::getClass) + .map(type -> ReflectionUtils.findMethod(type, "hasEvictionAttributes")) + .map(method -> ReflectionUtils.invokeMethod(method, regionAttributes)) + .map(Boolean.TRUE::equals) + .orElse(false); - return regionAttributes instanceof UserSpecifiedRegionAttributes - && ((UserSpecifiedRegionAttributes) regionAttributes).hasEvictionAttributes(); + return SpringUtils.safeGetValue(hasEvictionAttributes, false); } /* - * (non-Javadoc) - * * This method is not part of the PeerRegionFactoryBean API and is strictly used for testing purposes! * * @see org.apache.geode.cache.AttributesFactory#validateAttributes(:RegionAttributes) */ - @SuppressWarnings("deprecation") + @SuppressWarnings({ "deprecation", "rawtypes" }) void validateRegionAttributes(RegionAttributes regionAttributes) { org.apache.geode.cache.AttributesFactory.validateAttributes(regionAttributes); } @@ -631,6 +651,7 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact } } + @SuppressWarnings("rawtypes") private DataPolicy getDataPolicy(RegionAttributes regionAttributes, DataPolicy defaultDataPolicy) { return Optional.ofNullable(regionAttributes) @@ -639,7 +660,7 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact } /** - * Closes and destroys the {@link Region}. + * Closes and destroys this {@link Region}. * * @throws Exception if {@code destroy()} fails. * @see org.springframework.beans.factory.DisposableBean @@ -662,25 +683,82 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact } /** - * Configures an array of {@link AsyncEventQueue AsyncEventQueues} for this {@link Region} used to perform - * asynchronous data access operations, e.g. {@literal asynchronous write-behind}. + * Configures an array of {@link AsyncEventQueue AsyncEventQueues} for this {@link Region}, which are used + * to perform asynchronous data access operations, e.g. {@literal asynchronous, write-behind operations}. * - * @param asyncEventQueues array of {@link AsyncEventQueue AsyncEventQueues} used by this {@link Region} - * to perform asynchronous data access operations. + * This method clears any existing, registered {@link AsyncEventQueue AsyncEventQueues} (AEQ) already associated + * with this {@link Region}. Use {@link #addAsyncEventQueues(AsyncEventQueue[])} + * or {@link #addAsyncEventQueueIds(String[])} to append to the existing AEQs already registered instead. + * + * @param asyncEventQueues array of {@link AsyncEventQueue AsyncEventQueues} registered with and used by + * this {@link Region} to perform asynchronous data access operations. * @see org.apache.geode.cache.asyncqueue.AsyncEventQueue + * @see #addAsyncEventQueues(AsyncEventQueue[]) + * @see #addAsyncEventQueueIds(String[]) + * @see #setAsyncEventQueueIds(String[]) */ - public void setAsyncEventQueues(AsyncEventQueue[] asyncEventQueues) { - this.asyncEventQueues = asyncEventQueues; + public void setAsyncEventQueues(@NonNull AsyncEventQueue[] asyncEventQueues) { + this.asyncEventQueues.clear(); + addAsyncEventQueues(asyncEventQueues); } - public void setAsyncEventQueueIds(String[] asyncEventQueueIds) { - this.asyncEventQueueIds = asyncEventQueueIds; + /** + * Configures an array of {@link AsyncEventQueue AsyncEventQueues} (AEQ) for this {@link Region} + * by {@link String AEQ ID}. + * + * This method clears any existing, registered {@link AsyncEventQueue AsyncEventQueues} (AEQ) already associated + * with this {@link Region} by {@literal AEQ ID}. Use {@link #addAsyncEventQueues(AsyncEventQueue[])} + * or {@link #addAsyncEventQueueIds(String[])} to append to the existing AEQs already registered instead. + * + * or {@link #addAsyncEventQueueIds(String[])} to append to the existing AEQs already registered instead. + * @param asyncEventQueueIds array of {@link String Strings} specifying {@link String AEQ IDs} to be registered + * with this {@link Region}. + * @see #addAsyncEventQueues(AsyncEventQueue[]) + * @see #setAsyncEventQueues(AsyncEventQueue[]) + * @see #setAsyncEventQueueIds(String[]) + */ + public void setAsyncEventQueueIds(@NonNull String[] asyncEventQueueIds) { + this.asyncEventQueueIds.clear(); + addAsyncEventQueueIds(asyncEventQueueIds); + } + + /** + * Registers the array of {@link AsyncEventQueue AsyncEventQueues} (AEQ) with this {@link Region} by appending to + * the already existing, registered AEQs for this {@link Region}. + * + * @param asyncEventQueues array of {@link AsyncEventQueue AsyncEventQueues} to register with this {@link Region}. + * @see #addAsyncEventQueueIds(String[]) + * @see #setAsyncEventQueues(AsyncEventQueue[]) + * @see #setAsyncEventQueueIds(String[]) + */ + public void addAsyncEventQueues(@NonNull AsyncEventQueue[] asyncEventQueues) { + + Arrays.stream(ArrayUtils.nullSafeArray(asyncEventQueues, AsyncEventQueue.class)) + .filter(Objects::nonNull) + .forEach(this.asyncEventQueues::add); + } + + /** + * Registers the array of {@link AsyncEventQueue AsyncEventQueues} (AEQ) with this {@link Region} + * by {@link String ID} by appending to the already existing, registered AEQs for this {@link Region}. + * + * @param asyncEventQueueIds array of {@link AsyncEventQueue AsyncEventQueue} {@link String IDs} to register with + * this {@link Region}. + * @see #addAsyncEventQueues(AsyncEventQueue[]) + * @see #setAsyncEventQueues(AsyncEventQueue[]) + * @see #setAsyncEventQueueIds(String[]) + */ + public void addAsyncEventQueueIds(@NonNull String[] asyncEventQueueIds) { + + Arrays.stream(ArrayUtils.nullSafeArray(asyncEventQueueIds, String.class)) + .filter(StringUtils::hasText) + .forEach(this.asyncEventQueueIds::add); } /** * Sets the {@link RegionAttributes} used to configure this {@link Region}. * - * Specifying {@link RegionAttributes} allows maximum control in specifying various {@link Region} settings. + * Specifying {@link RegionAttributes} allows full control in configuring various {@link Region} settings. * Used only when the {@link Region} is created and not when the {@link Region} is looked up. * * @param attributes {@link RegionAttributes} used to configure this {@link Region}. @@ -697,7 +775,10 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact * @see org.apache.geode.cache.RegionAttributes */ public RegionAttributes getAttributes() { - return Optional.ofNullable(getRegion()).map(Region::getAttributes).orElse(this.attributes); + + return Optional.ofNullable(getRegion()) + .map(Region::getAttributes) + .orElse(this.attributes); } /** @@ -841,24 +922,83 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact /** * Configures the {@link GatewaySender GatewaySenders} used to send data and events from this {@link Region} - * to a corresponding {@link Region} in a remote cluster/site. + * to a matching {@link Region} in a remote cluster. + * + * This method clears all existing, registered {@link GatewaySender GatewaySenders} already associated with + * this {@link Region}. Use {@link #addGatewaySenders(GatewaySender[])} + * or {@link #addGatewaySendersIds(String[])} to append to the existing, registered + * {@link GatewaySender GatewaySenders} for this {@link Region}. * * @param gatewaySenders {@link GatewaySender GatewaySenders} used to send data and events from this {@link Region} - * to a corresponding {@link Region} in a remote cluster/site. + * to a matching {@link Region} in a remote cluster. * @see org.apache.geode.cache.wan.GatewaySender + * @see #addGatewaySenders(GatewaySender[]) + * @see #addGatewaySendersIds(String[]) + * @see #setGatewaySenderIds(String[]) */ - public void setGatewaySenders(GatewaySender[] gatewaySenders) { - this.gatewaySenders = gatewaySenders; - } - - public void setGatewaySenderIds(String[] gatewaySenderIds) { - this.gatewaySenderIds = gatewaySenderIds; + public void setGatewaySenders(@NonNull GatewaySender[] gatewaySenders) { + this.gatewaySenders.clear(); + addGatewaySenders(gatewaySenders); } /** - * Configures whether to enable this {@link Region} with the ability to store data in {@literal off-heap memory}. + * Configures the {@link GatewaySender GatewaySenders} by {@link String ID} used to send data and events from + * this {@link Region} to a matching {@link Region} in a remote cluster. * - * @param offHeap {@link Boolean} value indicating whether to enable {@literal off-heap memory} + * This method clears all existing, registered {@link GatewaySender GatewaySenders} already associated with + * this {@link Region}. Use {@link #addGatewaySenders(GatewaySender[])} + * or {@link #addGatewaySendersIds(String[])} to append to the existing, registered + * {@link GatewaySender GatewaySenders} for this {@link Region}. + * + * @param gatewaySenderIds {@link String} array containing {@link GatewaySender} {@link String IDs} to register + * with this {@link Region}. + * @see #addGatewaySenders(GatewaySender[]) + * @see #addGatewaySendersIds(String[]) + * @see #setGatewaySenders(GatewaySender[]) + */ + public void setGatewaySenderIds(@NonNull String[] gatewaySenderIds) { + this.gatewaySenderIds.clear(); + addGatewaySendersIds(gatewaySenderIds); + } + + /** + * Registers the array of {@link GatewaySender GatewaySenders} with this {@link Region} by appending to the already + * existing, registered {@link GatewaySender GatewaySenders} for this {@link Region}. + * + * @param gatewaySenders array of {@link GatewaySender GatewaySenders} to register with this {@link Region}. + * @see org.apache.geode.cache.wan.GatewaySender + * @see #addGatewaySendersIds(String[]) + * @see #setGatewaySenders(GatewaySender[]) + * @see #setGatewaySenderIds(String[]) + */ + public void addGatewaySenders(@NonNull GatewaySender[] gatewaySenders) { + + Arrays.stream(ArrayUtils.nullSafeArray(gatewaySenders, GatewaySender.class)) + .filter(Objects::nonNull) + .forEach(this.gatewaySenders::add); + } + + /** + * Registers the array of {@link GatewaySender} {@link String IDs} with this {@link Region} by appending to + * the already existing, registered {@link GatewaySender GatewaySenders} for this {@link Region}. + * + * @param gatewaySenderIds array of {@link GatewaySender} {@link String IDs} to register with this {@link Region}. + * @see org.apache.geode.cache.wan.GatewaySender + * @see #addGatewaySenders(GatewaySender[]) + * @see #setGatewaySenders(GatewaySender[]) + * @see #setGatewaySenderIds(String[]) + */ + public void addGatewaySendersIds(@NonNull String[] gatewaySenderIds) { + + Arrays.stream(ArrayUtils.nullSafeArray(gatewaySenderIds, String.class)) + .filter(StringUtils::hasText) + .forEach(this.gatewaySenderIds::add); + } + + /** + * Configures this {@link Region} with the capability to store data in {@literal off-heap memory}. + * + * @param offHeap {@link Boolean} value indicating whether to enable the use of {@literal off-heap memory} * for this {@link Region}. * @see org.apache.geode.cache.RegionFactory#setOffHeap(boolean) */ @@ -869,20 +1009,21 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact /** * Returns a {@link Boolean} value indicating whether {@literal off-heap memory} was enabled for this {@link Region}. * - * {@literal Off-heap memory} will be enabled if this method returns a {@literal non-null} {@link Boolean} value - * evaluating to {@literal true}. + * {@literal Off-heap memory} will be enabled for this {@link Region} if this method returns a {@literal non-null}, + * {@link Boolean} value evaluating to {@literal true}. * - * @return a {@link Boolean} value indicating whether {@literal off-heap memory} is enabled for this {@link Region}. + * @return a {@link Boolean} value indicating whether {@literal off-heap memory} use is enabled for + * this {@link Region}. */ public Boolean getOffHeap() { return this.offHeap; } /** - * Returns a boolean value indicating whether {@literal off-heap memory} has been enabled for this {@link Region}. + * Returns a boolean value indicating whether {@literal off-heap memory} use was enabled for this {@link Region}. * - * @return a {@literal boolean} value indicating whether {@literal off-heap memory} has been enabled - * for this {@link Region}. + * @return a {@literal boolean} value indicating whether {@literal off-heap memory} use was enabled for + * this {@link Region}. * @see #getOffHeap() */ public boolean isOffHeap() { @@ -1008,9 +1149,9 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact @SuppressWarnings("all") public void start() { - if (!ObjectUtils.isEmpty(this.gatewaySenders)) { + if (!this.gatewaySenders.isEmpty()) { synchronized(this.gatewaySenders) { - Arrays.stream(this.gatewaySenders) + this.gatewaySenders.stream() .filter(Objects::nonNull) .filter(gatewaySender -> !gatewaySender.isManualStart()) .filter(gatewaySender -> !gatewaySender.isRunning()) @@ -1031,11 +1172,9 @@ public abstract class PeerRegionFactoryBean extends ConfigurableRegionFact @SuppressWarnings("all") public void stop() { - if (!ObjectUtils.isEmpty(this.gatewaySenders)) { + if (!this.gatewaySenders.isEmpty()) { synchronized (this.gatewaySenders) { - for (GatewaySender gatewaySender : this.gatewaySenders) { - gatewaySender.stop(); - } + this.gatewaySenders.forEach(GatewaySender::stop); } } diff --git a/src/test/java/org/springframework/data/gemfire/config/xml/GemfireV7GatewayNamespaceTest.java b/src/test/java/org/springframework/data/gemfire/config/xml/GemfireV7GatewayNamespaceTest.java index 8c2a478d..70d56a1f 100644 --- a/src/test/java/org/springframework/data/gemfire/config/xml/GemfireV7GatewayNamespaceTest.java +++ b/src/test/java/org/springframework/data/gemfire/config/xml/GemfireV7GatewayNamespaceTest.java @@ -15,18 +15,16 @@ */ package org.springframework.data.gemfire.config.xml; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertFalse; -import static org.junit.Assert.assertNotNull; -import static org.junit.Assert.assertSame; -import static org.junit.Assert.assertTrue; -import static org.springframework.data.gemfire.util.ArrayUtils.nullSafeArray; +import static org.assertj.core.api.Assertions.assertThat; import java.io.File; import java.io.InputStream; import java.io.OutputStream; import java.util.List; +import org.junit.AfterClass; +import org.junit.Test; + import org.apache.geode.cache.Region; import org.apache.geode.cache.asyncqueue.AsyncEvent; import org.apache.geode.cache.asyncqueue.AsyncEventListener; @@ -38,9 +36,6 @@ import org.apache.geode.cache.wan.GatewaySender; import org.apache.geode.cache.wan.GatewaySender.OrderPolicy; import org.apache.geode.cache.wan.GatewayTransportFilter; -import org.junit.AfterClass; -import org.junit.Test; - import org.springframework.data.gemfire.PeerRegionFactoryBean; import org.springframework.data.gemfire.RecreatingContextTest; import org.springframework.data.gemfire.TestUtils; @@ -48,20 +43,25 @@ import org.springframework.data.gemfire.test.GemfireTestBeanPostProcessor; import org.springframework.data.gemfire.wan.GatewaySenderFactoryBean; /** - * This test is only valid for GF 7.0 and above + * Integration Tests testing and asserting GemFire 7.0 WAN functionality and configuration. * * @author David Turanski * @author John Blum + * @see org.junit.Test + * @see org.apache.geode.cache.asyncqueue.AsyncEventQueue + * @see org.apache.geode.cache.wan.GatewayReceiver + * @see org.apache.geode.cache.wan.GatewaySender + * @see org.springframework.data.gemfire.RecreatingContextTest */ +@SuppressWarnings("unused") public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest { - /* - * (non-Javadoc) - * @see org.springframework.data.gemfire.RecreatingContextTest#location() - */ - @Override - protected String location() { - return "/org/springframework/data/gemfire/config/xml/gateway-v7-ns.xml"; + @AfterClass + public static void tearDown() { + + for (String name : new File(".").list((file, filename) -> filename.startsWith("BACKUP"))) { + new File(name).delete(); + } } @Override @@ -69,163 +69,166 @@ public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest { this.applicationContext.getBeanFactory().addBeanPostProcessor(new GemfireTestBeanPostProcessor()); } - @AfterClass - @SuppressWarnings("all") - public static void tearDown() { - - for (String name : nullSafeArray(new File(".") - .list((file, filename) -> filename.startsWith("BACKUP")), String.class)) { - - new File(name).delete(); - } + @Override + protected String location() { + return "/org/springframework/data/gemfire/config/xml/gateway-v7-ns.xml"; } @Test - public void testAsyncEventQueue() { + public void asyncEventQueueConfigurationIsCorrect() { AsyncEventQueue asyncEventQueue = this.applicationContext.getBean("async-event-queue", AsyncEventQueue.class); - assertNotNull(asyncEventQueue); - assertTrue(asyncEventQueue.isBatchConflationEnabled()); - assertEquals(10, asyncEventQueue.getBatchSize()); - assertEquals(3, asyncEventQueue.getBatchTimeInterval()); - assertEquals("diskstore", asyncEventQueue.getDiskStoreName()); - assertTrue(asyncEventQueue.isDiskSynchronous()); - assertEquals(50, asyncEventQueue.getMaximumQueueMemory()); - assertEquals(OrderPolicy.KEY, asyncEventQueue.getOrderPolicy()); - assertFalse(asyncEventQueue.isParallel()); - assertTrue(asyncEventQueue.isPersistent()); + assertThat(asyncEventQueue).isNotNull(); + assertThat(asyncEventQueue.isBatchConflationEnabled()).isTrue(); + assertThat(asyncEventQueue.getBatchSize()).isEqualTo(10); + assertThat(asyncEventQueue.getBatchTimeInterval()).isEqualTo(3); + assertThat(asyncEventQueue.getDiskStoreName()).isEqualTo("diskstore"); + assertThat(asyncEventQueue.isDiskSynchronous()).isTrue(); + assertThat(asyncEventQueue.getMaximumQueueMemory()).isEqualTo(50); + assertThat(asyncEventQueue.getOrderPolicy()).isEqualTo(OrderPolicy.KEY); + assertThat(asyncEventQueue.isParallel()).isFalse(); + assertThat(asyncEventQueue.isPersistent()).isTrue(); } @Test - public void testGatewaySender() throws Exception { + public void gatewaySenderFactoryBeanConfigurationIsCorrect() throws Exception { GatewaySenderFactoryBean gatewaySenderFactoryBean = this.applicationContext.getBean("&gateway-sender", GatewaySenderFactoryBean.class); - assertNotNull(gatewaySenderFactoryBean); - assertNotNull(TestUtils.readField("cache", gatewaySenderFactoryBean)); - assertEquals(2, - TestUtils.readField("remoteDistributedSystemId", gatewaySenderFactoryBean).longValue()); - assertEquals(10, TestUtils.readField("alertThreshold", gatewaySenderFactoryBean).longValue()); - assertTrue(Boolean.TRUE.equals(TestUtils.readField("batchConflationEnabled", gatewaySenderFactoryBean))); - assertEquals(11, TestUtils.readField("batchSize", gatewaySenderFactoryBean).intValue()); - assertEquals(12, TestUtils.readField("dispatcherThreads", gatewaySenderFactoryBean).intValue()); - assertEquals(false, TestUtils.readField("diskSynchronous", gatewaySenderFactoryBean)); - assertEquals(true, TestUtils.readField("manualStart", gatewaySenderFactoryBean)); + assertThat(gatewaySenderFactoryBean).isNotNull(); + assertThat(gatewaySenderFactoryBean.getCache()).isNotNull(); + assertThat(gatewaySenderFactoryBean.getRemoteDistributedSystemId()).isEqualTo(2); + assertThat(gatewaySenderFactoryBean.getAlertThreshold()).isEqualTo(10); + assertThat(gatewaySenderFactoryBean.getBatchConflationEnabled()).isTrue(); + assertThat(gatewaySenderFactoryBean.getBatchSize()).isEqualTo(11); + assertThat(gatewaySenderFactoryBean.getDispatcherThreads()).isEqualTo(12); + assertThat(gatewaySenderFactoryBean.getDiskSynchronous()).isFalse(); + assertThat(gatewaySenderFactoryBean.isManualStart()).isTrue(); List eventFilters = TestUtils.readField("eventFilters", gatewaySenderFactoryBean); - assertNotNull(eventFilters); - assertEquals(2, eventFilters.size()); - assertTrue(eventFilters.get(0) instanceof TestEventFilter); + assertThat(eventFilters).isNotNull(); + assertThat(eventFilters.size()).isEqualTo(2); + assertThat(eventFilters.get(0)).isInstanceOf(TestEventFilter.class); List transportFilters = TestUtils .readField("transportFilters", gatewaySenderFactoryBean); - assertNotNull(transportFilters); - assertEquals(2, transportFilters.size()); - assertTrue(transportFilters.get(0) instanceof TestTransportFilter); + assertThat(transportFilters).isNotNull(); + assertThat(transportFilters.size()).isEqualTo(2); + assertThat(transportFilters.get(0)).isInstanceOf(TestTransportFilter.class); } @Test @SuppressWarnings("rawtypes") - public void testInnerGatewaySender() throws Exception { + public void nestedGatewaySenderConfigurationIsCorrect() throws Exception { - Region region = this.applicationContext.getBean("region-inner-gateway-sender", Region.class); + Region region = this.applicationContext.getBean("region-with-nested-gateway-sender", Region.class); - assertNotNull(region.getAttributes().getGatewaySenderIds()); - assertEquals(2, region.getAttributes().getGatewaySenderIds().size()); + assertThat(region).isNotNull(); + assertThat(region.getAttributes()).isNotNull(); + assertThat(region.getAttributes().getGatewaySenderIds()).isNotNull(); + assertThat(region.getAttributes().getGatewaySenderIds()).hasSize(2); - PeerRegionFactoryBean regionFactoryBean = - this.applicationContext.getBean("®ion-inner-gateway-sender", PeerRegionFactoryBean.class); + PeerRegionFactoryBean regionFactoryBean = applicationContext.getBean("®ion-with-nested-gateway-sender", PeerRegionFactoryBean.class); - Object[] gatewaySenders = TestUtils.readField("gatewaySenders", regionFactoryBean); + List gatewaySenders = TestUtils.readField("gatewaySenders", regionFactoryBean); - assertNotNull(gatewaySenders); - assertEquals(2, gatewaySenders.length); + assertThat(gatewaySenders).isNotNull(); + assertThat(gatewaySenders).hasSize(2); - GatewaySender gatewaySender = (GatewaySender) gatewaySenders[0]; + GatewaySender gatewaySender = gatewaySenders.get(0); - assertNotNull(gatewaySender); - assertEquals(1, gatewaySender.getRemoteDSId()); - assertEquals(false, gatewaySender.isManualStart()); - assertEquals(true, gatewaySender.isRunning()); - assertEquals(10, gatewaySender.getAlertThreshold()); - assertEquals(11, gatewaySender.getBatchSize()); - assertEquals(3000, gatewaySender.getBatchTimeInterval()); - assertEquals(2, gatewaySender.getDispatcherThreads()); - assertEquals("diskstore", gatewaySender.getDiskStoreName()); - assertEquals(true, gatewaySender.isDiskSynchronous()); - assertTrue(gatewaySender.isBatchConflationEnabled()); - assertEquals(50, gatewaySender.getMaximumQueueMemory()); - assertEquals(OrderPolicy.THREAD, gatewaySender.getOrderPolicy()); - assertTrue(gatewaySender.isPersistenceEnabled()); - assertFalse(gatewaySender.isParallel()); - assertEquals(16536, gatewaySender.getSocketBufferSize()); - assertEquals(3000, gatewaySender.getSocketReadTimeout()); + assertThat(gatewaySender).isNotNull(); + assertThat(gatewaySender.getRemoteDSId()).isEqualTo(1); + assertThat(gatewaySender.isManualStart()).isFalse(); + assertThat(gatewaySender.isRunning()).isTrue(); + assertThat(gatewaySender.getAlertThreshold()).isEqualTo(10); + assertThat(gatewaySender.getBatchSize()).isEqualTo(11); + assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(3000); + assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(2); + assertThat(gatewaySender.getDiskStoreName()).isEqualTo("diskstore"); + assertThat(gatewaySender.isDiskSynchronous()).isEqualTo(true); + assertThat(gatewaySender.isBatchConflationEnabled()).isTrue(); + assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(50); + assertThat(gatewaySender.getOrderPolicy()).isEqualTo(OrderPolicy.THREAD); + assertThat(gatewaySender.isPersistenceEnabled()).isTrue(); + assertThat(gatewaySender.isParallel()).isFalse(); + assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(16536); + assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(3000); List eventFilters = gatewaySender.getGatewayEventFilters(); - assertNotNull(eventFilters); - assertEquals(1, eventFilters.size()); - assertTrue(eventFilters.get(0) instanceof TestEventFilter); + assertThat(eventFilters).isNotNull(); + assertThat(eventFilters).hasSize(1); + assertThat(eventFilters.get(0)).isInstanceOf(TestEventFilter.class); List transportFilters = gatewaySender.getGatewayTransportFilters(); - assertNotNull(transportFilters); - assertEquals(1, transportFilters.size()); - assertTrue(transportFilters.get(0) instanceof TestTransportFilter); + assertThat(transportFilters).isNotNull(); + assertThat(transportFilters).hasSize(1); + assertThat(transportFilters.get(0)).isInstanceOf(TestTransportFilter.class); } @Test - public void testGatewaySenderWithEventTransportFilterRefs() throws Exception { + public void gatewaySenderWithEventTransportFilterRefsConfigurationIsCorrect() throws Exception { GatewaySenderFactoryBean gatewaySenderFactoryBean = this.applicationContext.getBean("&gateway-sender-with-event-transport-filter-refs", GatewaySenderFactoryBean.class); - assertNotNull(gatewaySenderFactoryBean); - assertNotNull(TestUtils.readField("cache", gatewaySenderFactoryBean)); - assertEquals(3, - TestUtils.readField("remoteDistributedSystemId", gatewaySenderFactoryBean).intValue()); - assertTrue(Boolean.TRUE.equals(TestUtils.readField("batchConflationEnabled", gatewaySenderFactoryBean))); - assertEquals(50, TestUtils.readField("batchSize", gatewaySenderFactoryBean).intValue()); - assertEquals(10, TestUtils.readField("dispatcherThreads", gatewaySenderFactoryBean).intValue()); - assertEquals(true, TestUtils.readField("manualStart", gatewaySenderFactoryBean)); + assertThat(gatewaySenderFactoryBean).isNotNull(); + assertThat(gatewaySenderFactoryBean.getCache()).isNotNull(); + assertThat(gatewaySenderFactoryBean.getRemoteDistributedSystemId()).isEqualTo(3); + assertThat(gatewaySenderFactoryBean.getBatchConflationEnabled()).isTrue(); + assertThat(gatewaySenderFactoryBean.getBatchSize()).isEqualTo(50); + assertThat(gatewaySenderFactoryBean.getDispatcherThreads()).isEqualTo(10); + assertThat(gatewaySenderFactoryBean.isManualStart()).isTrue(); List eventFilters = TestUtils.readField("eventFilters", gatewaySenderFactoryBean); - assertNotNull(eventFilters); - assertEquals(1, eventFilters.size()); - assertTrue(eventFilters.get(0) instanceof TestEventFilter); - assertSame(applicationContext.getBean("event-filter"), eventFilters.get(0)); + assertThat(eventFilters).isNotNull(); + assertThat(eventFilters).hasSize(1); + assertThat(eventFilters.get(0)).isInstanceOf(TestEventFilter.class); + assertThat(eventFilters.get(0)).isSameAs(applicationContext.getBean("event-filter")); - List transportFilters = TestUtils - .readField("transportFilters", gatewaySenderFactoryBean); + List transportFilters = + TestUtils.readField("transportFilters", gatewaySenderFactoryBean); - assertNotNull(transportFilters); - assertEquals(1, transportFilters.size()); - assertTrue(transportFilters.get(0) instanceof TestTransportFilter); - assertSame(applicationContext.getBean("transport-filter"), transportFilters.get(0)); + assertThat(transportFilters).isNotNull(); + assertThat(transportFilters).hasSize(1); + assertThat(transportFilters.get(0)).isInstanceOf(TestTransportFilter.class); + assertThat(transportFilters.get(0)).isSameAs(applicationContext.getBean("transport-filter")); } @Test - public void testGatewayReceiver() { + public void gatewayReceiverConfigurationIsCorrect() { GatewayReceiver gatewayReceiver = this.applicationContext.getBean("gateway-receiver", GatewayReceiver.class); - assertNotNull(gatewayReceiver); - assertEquals("192.168.0.1", gatewayReceiver.getBindAddress()); - assertEquals(12345, gatewayReceiver.getStartPort()); - assertEquals(23456, gatewayReceiver.getEndPort()); - assertEquals(3000, gatewayReceiver.getMaximumTimeBetweenPings()); - assertEquals(16536, gatewayReceiver.getSocketBufferSize()); + assertThat(gatewayReceiver).isNotNull(); + assertThat(gatewayReceiver.getBindAddress()).isEqualTo("192.168.0.1"); + assertThat(gatewayReceiver.getStartPort()).isEqualTo(12345); + assertThat(gatewayReceiver.getEndPort()).isEqualTo(23456); + assertThat(gatewayReceiver.getMaximumTimeBetweenPings()).isEqualTo(3000); + assertThat(gatewayReceiver.getSocketBufferSize()).isEqualTo(16536); + } + + public static class TestAsyncEventListener implements AsyncEventListener { + + @Override + public void close() { } + + @Override + public boolean processEvents(List events) { + return false; + } } - @SuppressWarnings("rawtypes") public static class TestEventFilter implements GatewayEventFilter { @Override @@ -233,16 +236,15 @@ public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest { } @Override - public void afterAcknowledgement(GatewayQueueEvent arg0) { - } + public void afterAcknowledgement(GatewayQueueEvent event) { } @Override - public boolean beforeEnqueue(GatewayQueueEvent arg0) { + public boolean beforeEnqueue(GatewayQueueEvent event) { return false; } @Override - public boolean beforeTransmit(GatewayQueueEvent arg0) { + public boolean beforeTransmit(GatewayQueueEvent event) { return false; } } @@ -254,26 +256,13 @@ public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest { } @Override - public InputStream getInputStream(InputStream arg0) { + public InputStream getInputStream(InputStream in) { return null; } @Override - public OutputStream getOutputStream(OutputStream arg0) { + public OutputStream getOutputStream(OutputStream out) { return null; } } - - @SuppressWarnings("rawtypes") - public static class TestAsyncEventListener implements AsyncEventListener { - - @Override - public void close() { - } - - @Override - public boolean processEvents(List arg0) { - return false; - } - } } diff --git a/src/test/resources/org/springframework/data/gemfire/config/xml/gateway-v7-ns.xml b/src/test/resources/org/springframework/data/gemfire/config/xml/gateway-v7-ns.xml index 976b3ba5..6b52e960 100644 --- a/src/test/resources/org/springframework/data/gemfire/config/xml/gateway-v7-ns.xml +++ b/src/test/resources/org/springframework/data/gemfire/config/xml/gateway-v7-ns.xml @@ -16,7 +16,25 @@ - + + + + + + + + + - - - - - - - -