diff --git a/src/main/java/org/springframework/data/gemfire/config/xml/AsyncEventQueueParser.java b/src/main/java/org/springframework/data/gemfire/config/xml/AsyncEventQueueParser.java index 411e098b..b53567e2 100644 --- a/src/main/java/org/springframework/data/gemfire/config/xml/AsyncEventQueueParser.java +++ b/src/main/java/org/springframework/data/gemfire/config/xml/AsyncEventQueueParser.java @@ -20,7 +20,6 @@ import org.springframework.beans.factory.config.RuntimeBeanReference; import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser; import org.springframework.beans.factory.xml.ParserContext; -import org.springframework.data.gemfire.GemfireUtils; import org.springframework.data.gemfire.util.SpringUtils; import org.springframework.data.gemfire.wan.AsyncEventQueueFactoryBean; import org.springframework.util.StringUtils; @@ -50,34 +49,50 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser { */ @Override protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { + builder.setLazyInit(false); parseAsyncEventListener(element, parserContext, builder); parseCache(element, builder); parseDiskStore(element, builder); + ParsingUtils.setPropertyValue(element, builder, "enable-batch-conflation", "batchConflationEnabled"); + ParsingUtils.setPropertyValue(element, builder, "batch-conflation-enabled"); ParsingUtils.setPropertyValue(element, builder, "batch-size"); + ParsingUtils.setPropertyValue(element, builder, "batch-time-interval"); + ParsingUtils.setPropertyValue(element, builder, "disk-synchronous"); + ParsingUtils.setPropertyValue(element, builder, "dispatcher-threads"); + ParsingUtils.setPropertyValue(element, builder, "forward-expiration-destroy"); ParsingUtils.setPropertyValue(element, builder, "maximum-queue-memory"); + ParsingUtils.setPropertyValue(element, builder, "order-policy"); ParsingUtils.setPropertyValue(element, builder, "parallel"); ParsingUtils.setPropertyValue(element, builder, "persistent"); - if (GemfireUtils.GEMFIRE_VERSION.compareTo("7.0.1") >= 0) { - ParsingUtils.setPropertyValue(element, builder, "enable-batch-conflation", "batchConflationEnabled"); - ParsingUtils.setPropertyValue(element, builder, "batch-conflation-enabled"); - ParsingUtils.setPropertyValue(element, builder, "batch-time-interval"); - ParsingUtils.setPropertyValue(element, builder, "disk-synchronous"); - ParsingUtils.setPropertyValue(element, builder, "dispatcher-threads"); - ParsingUtils.setPropertyValue(element, builder, "order-policy"); + Element eventFilterElement = DomUtils.getChildElementByTagName(element, "event-filter"); + + if (eventFilterElement != null) { + builder.addPropertyValue("gatewayEventFilters", + ParsingUtils.parseRefOrNestedBeanDeclaration(eventFilterElement, parserContext, builder)); + } + + Element eventSubstitutionFilterElement = + DomUtils.getChildElementByTagName(element, "event-substitution-filter"); + + if (eventSubstitutionFilterElement != null) { + builder.addPropertyValue("gatewayEventSubstitutionFilter", + ParsingUtils.parseRefOrSingleNestedBeanDeclaration(eventSubstitutionFilterElement, parserContext, builder)); } ParsingUtils.setPropertyValue(element, builder, NAME_ATTRIBUTE); if (!StringUtils.hasText(element.getAttribute(NAME_ATTRIBUTE))) { if (element.getParentNode().getNodeName().endsWith("region")) { + Element region = (Element) element.getParentNode(); String regionName = StringUtils.hasText(region.getAttribute(NAME_ATTRIBUTE)) - ? region.getAttribute(NAME_ATTRIBUTE) : region.getAttribute(ID_ATTRIBUTE); + ? region.getAttribute(NAME_ATTRIBUTE) + : region.getAttribute(ID_ATTRIBUTE); int index = 0; @@ -92,13 +107,12 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser { } } - /* (non-Javadoc) */ private void parseAsyncEventListener(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { + Element asyncEventListenerElement = DomUtils.getChildElementByTagName(element, "async-event-listener"); - Object asyncEventListener = ParsingUtils.parseRefOrSingleNestedBeanDeclaration(asyncEventListenerElement, - parserContext, - builder); + Object asyncEventListener = + ParsingUtils.parseRefOrSingleNestedBeanDeclaration(asyncEventListenerElement, parserContext, builder); builder.addPropertyValue("asyncEventListener", asyncEventListener); @@ -109,15 +123,16 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser { /* (non-Javadoc) */ private void parseCache(Element element, BeanDefinitionBuilder builder) { + String cacheRefAttribute = element.getAttribute("cache-ref"); - String cacheName = SpringUtils.defaultIfEmpty(cacheRefAttribute, - GemfireConstants.DEFAULT_GEMFIRE_CACHE_NAME); + String cacheName = SpringUtils.defaultIfEmpty(cacheRefAttribute, GemfireConstants.DEFAULT_GEMFIRE_CACHE_NAME); builder.addConstructorArgReference(cacheName); } /* (non-Javadoc) */ private void parseDiskStore(Element element, BeanDefinitionBuilder builder) { + ParsingUtils.setPropertyValue(element, builder, "disk-store-ref"); String diskStoreRef = element.getAttribute("disk-store-ref"); diff --git a/src/main/java/org/springframework/data/gemfire/wan/AbstractWANComponentFactoryBean.java b/src/main/java/org/springframework/data/gemfire/wan/AbstractWANComponentFactoryBean.java index 1e9561a8..fdbc1315 100644 --- a/src/main/java/org/springframework/data/gemfire/wan/AbstractWANComponentFactoryBean.java +++ b/src/main/java/org/springframework/data/gemfire/wan/AbstractWANComponentFactoryBean.java @@ -13,10 +13,8 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.data.gemfire.wan; -import java.util.Arrays; -import java.util.List; +package org.springframework.data.gemfire.wan; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -29,35 +27,34 @@ import org.springframework.util.Assert; import org.springframework.util.StringUtils; /** - * Base class for GemFire WAN Gateway component factory beans. + * Abstract base class for WAN Gateway components. * * @author David Turanski * @author John Blum + * @see org.apache.geode.cache.Cache * @see org.springframework.beans.factory.BeanNameAware * @see org.springframework.beans.factory.DisposableBean * @see org.springframework.beans.factory.FactoryBean * @see org.springframework.beans.factory.InitializingBean */ -public abstract class AbstractWANComponentFactoryBean implements BeanNameAware, FactoryBean, - InitializingBean, DisposableBean { - - protected static final List VALID_ORDER_POLICIES = Arrays.asList("KEY", "PARTITION", "THREAD"); - - protected Log log = LogFactory.getLog(getClass()); +public abstract class AbstractWANComponentFactoryBean + implements BeanNameAware, DisposableBean, FactoryBean, InitializingBean { protected final Cache cache; + protected final Log log = LogFactory.getLog(getClass()); + protected Object factory; private String beanName; private String name; - protected AbstractWANComponentFactoryBean(final Cache cache) { + protected AbstractWANComponentFactoryBean(Cache cache) { this.cache = cache; } @Override - public final void setBeanName(final String beanName) { + public void setBeanName(String beanName) { this.beanName = beanName; } @@ -65,20 +62,14 @@ public abstract class AbstractWANComponentFactoryBean implements BeanNameAwar this.factory = factory; } - public void setName(final String name) { + public void setName(String name) { this.name = name; } public String getName() { - return (StringUtils.hasText(name) ? name : beanName); + return StringUtils.hasText(this.name) ? this.name : this.beanName; } - @Override - public abstract T getObject() throws Exception; - - @Override - public abstract Class getObjectType(); - @Override public final boolean isSingleton() { return true; @@ -86,15 +77,16 @@ public abstract class AbstractWANComponentFactoryBean implements BeanNameAwar @Override public final void afterPropertiesSet() throws Exception { - Assert.notNull(cache, "Cache must not be null."); - Assert.notNull(getName(), "Name must not be null."); + + Assert.notNull(this.cache, "Cache must not be null"); + Assert.notNull(getName(), "Name must not be null"); + doInit(); } protected abstract void doInit() throws Exception; @Override - public void destroy() throws Exception { - } + public void destroy() throws Exception { } } diff --git a/src/main/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBean.java b/src/main/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBean.java index fe0e4f7f..4b8ce27c 100644 --- a/src/main/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBean.java +++ b/src/main/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBean.java @@ -13,21 +13,36 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.data.gemfire.wan; +import static org.springframework.data.gemfire.util.CollectionUtils.nullSafeList; + +import java.util.List; +import java.util.Optional; + import org.apache.geode.cache.Cache; import org.apache.geode.cache.CacheClosedException; +import org.apache.geode.cache.Region; import org.apache.geode.cache.asyncqueue.AsyncEventListener; import org.apache.geode.cache.asyncqueue.AsyncEventQueue; import org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory; +import org.apache.geode.cache.wan.GatewayEventFilter; +import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter; import org.apache.geode.cache.wan.GatewaySender; +import org.springframework.beans.factory.FactoryBean; import org.springframework.util.Assert; /** - * FactoryBean for creating GemFire {@link AsyncEventQueue}s. + * Spring {@link FactoryBean} for creating Apache Geode/Pivotal GemFire {@link AsyncEventQueue AsyncEventQueues}. * * @author David Turanski * @author John Blum + * @see org.apache.geode.cache.Cache + * @see org.apache.geode.cache.Region + * @see org.apache.geode.cache.asyncqueue.AsyncEventListener + * @see org.apache.geode.cache.asyncqueue.AsyncEventQueue + * @see org.springframework.data.gemfire.wan.AbstractWANComponentFactoryBean */ @SuppressWarnings("unused") public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean { @@ -38,6 +53,7 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< private Boolean batchConflationEnabled; private Boolean diskSynchronous; + private Boolean forwardExpirationDestroy; private Boolean parallel; private Boolean persistent; @@ -46,8 +62,13 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< private Integer dispatcherThreads; private Integer maximumQueueMemory; + private GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter; + + private GatewaySender.OrderPolicy orderPolicy; + + private List gatewayEventFilters; + private String diskStoreReference; - private String orderPolicy; /** * Constructs an instance of the AsyncEventQueueFactoryBean for creating an GemFire AsyncEventQueue. @@ -65,77 +86,60 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< * @param cache the GemFire Cache reference. * @param asyncEventListener required {@link AsyncEventListener} */ - public AsyncEventQueueFactoryBean(final Cache cache, final AsyncEventListener asyncEventListener) { + public AsyncEventQueueFactoryBean(Cache cache, AsyncEventListener asyncEventListener) { + super(cache); + setAsyncEventListener(asyncEventListener); } @Override public AsyncEventQueue getObject() throws Exception { - return asyncEventQueue; + return this.asyncEventQueue; } @Override public Class getObjectType() { - return AsyncEventQueue.class; + return this.asyncEventQueue != null ? this.asyncEventQueue.getClass() : AsyncEventQueue.class; } @Override protected void doInit() { - Assert.notNull(this.asyncEventListener, "The AsyncEventListener cannot be null."); - AsyncEventQueueFactory asyncEventQueueFactory = (this.factory != null ? (AsyncEventQueueFactory) factory - : cache.createAsyncEventQueueFactory()); + Assert.state(this.asyncEventListener != null, "AsyncEventListener must not be null"); - if (batchSize != null) { - asyncEventQueueFactory.setBatchSize(batchSize); - } + AsyncEventQueueFactory asyncEventQueueFactory = + this.factory != null ? (AsyncEventQueueFactory) this.factory : this.cache.createAsyncEventQueueFactory(); - if (batchTimeInterval != null) { - asyncEventQueueFactory.setBatchTimeInterval(batchTimeInterval); - } - - if (batchConflationEnabled != null) { - asyncEventQueueFactory.setBatchConflationEnabled(batchConflationEnabled); - } - - if (dispatcherThreads != null) { - asyncEventQueueFactory.setDispatcherThreads(dispatcherThreads); - } - - if (diskStoreReference != null) { - asyncEventQueueFactory.setDiskStoreName(diskStoreReference); - } - - if (diskSynchronous != null) { - asyncEventQueueFactory.setDiskSynchronous(diskSynchronous); - } - - if (maximumQueueMemory != null) { - asyncEventQueueFactory.setMaximumQueueMemory(maximumQueueMemory); - } + Optional.ofNullable(this.batchConflationEnabled).ifPresent(asyncEventQueueFactory::setBatchConflationEnabled); + Optional.ofNullable(this.batchSize).ifPresent(asyncEventQueueFactory::setBatchSize); + Optional.ofNullable(this.batchTimeInterval).ifPresent(asyncEventQueueFactory::setBatchTimeInterval); + Optional.ofNullable(this.diskStoreReference).ifPresent(asyncEventQueueFactory::setDiskStoreName); + Optional.ofNullable(this.diskSynchronous).ifPresent(asyncEventQueueFactory::setDiskSynchronous); + Optional.ofNullable(this.dispatcherThreads).ifPresent(asyncEventQueueFactory::setDispatcherThreads); + Optional.ofNullable(this.forwardExpirationDestroy).ifPresent(asyncEventQueueFactory::setForwardExpirationDestroy); + Optional.ofNullable(this.gatewayEventSubstitutionFilter).ifPresent(asyncEventQueueFactory::setGatewayEventSubstitutionListener); + Optional.ofNullable(this.maximumQueueMemory).ifPresent(asyncEventQueueFactory::setMaximumQueueMemory); + Optional.ofNullable(this.persistent).ifPresent(asyncEventQueueFactory::setPersistent); asyncEventQueueFactory.setParallel(isParallelEventQueue()); - if (orderPolicy != null) { - Assert.isTrue(isSerialEventQueue(), "Order Policy cannot be used with a Parallel Event Queue."); + nullSafeList(this.gatewayEventFilters).forEach(asyncEventQueueFactory::addGatewayEventFilter); - Assert.isTrue(VALID_ORDER_POLICIES.contains(orderPolicy.toUpperCase()), String.format( - "The value of Order Policy '$1%s' is invalid.", orderPolicy)); + if (this.orderPolicy != null) { - asyncEventQueueFactory.setOrderPolicy(GatewaySender.OrderPolicy.valueOf(orderPolicy.toUpperCase())); + Assert.state(isSerialEventQueue(), "OrderPolicy cannot be used with a Parallel AsyncEventQueue"); + + asyncEventQueueFactory.setOrderPolicy(this.orderPolicy); } - if (persistent != null) { - asyncEventQueueFactory.setPersistent(persistent); - } - - asyncEventQueue = asyncEventQueueFactory.create(getName(), this.asyncEventListener); + setAsyncEventQueue(asyncEventQueueFactory.create(getName(), this.asyncEventListener)); } @Override public void destroy() throws Exception { - if (!cache.isClosed()) { + + if (!this.cache.isClosed()) { try { this.asyncEventListener.close(); } @@ -145,22 +149,28 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< } public final void setAsyncEventListener(AsyncEventListener listener) { + Assert.state(this.asyncEventQueue == null, - "Setting an AsyncEventListener is not allowed once the AsyncEventQueue has been created."); + "Setting an AsyncEventListener is not allowed once the AsyncEventQueue has been created"); + this.asyncEventListener = listener; } /** - * @param asyncEventQueue overrides Async Event Queue returned by this FactoryBean. + * Configures the {@link AsyncEventQueue} returned by this {@link FactoryBean}. + * + * @param asyncEventQueue overrides {@link AsyncEventQueue} returned by this {@link FactoryBean}. + * @see org.apache.geode.cache.asyncqueue.AsyncEventQueue */ public void setAsyncEventQueue(AsyncEventQueue asyncEventQueue) { this.asyncEventQueue = asyncEventQueue; } /** - * Enable or disable the Async Event Queue's (AEQ) should conflate messages. + * Enable or disable {@link AsyncEventQueue} (AEQ) message conflation. * - * @param batchConflationEnabled a boolean value indicating whether to conflate queued events. + * @param batchConflationEnabled {@link Boolean} indicating whether to conflate queued events. + * @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setBatchConflationEnabled(boolean) */ public void setBatchConflationEnabled(Boolean batchConflationEnabled) { this.batchConflationEnabled = batchConflationEnabled; @@ -171,10 +181,11 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< } /** - * Set the Aysync Event Queue's (AEQ) interval between sending batches. + * Configures the {@link AsyncEventQueue} (AEQ) interval between sending batches. * - * @param batchTimeInterval an integer value indicating the maximum number of milliseconds that can elapse - * between sending batches. + * @param batchTimeInterval {@link Integer} specifying the maximum number of milliseconds + * that can elapse between sending batches. + * @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setBatchTimeInterval(int) */ public void setBatchTimeInterval(Integer batchTimeInterval) { this.batchTimeInterval = batchTimeInterval; @@ -185,36 +196,69 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< } /** - * Set the Async Event Queue (AEQ) disk write synchronization policy. + * Configures the {@link AsyncEventQueue} (AEQ) disk write synchronization policy. * - * @param diskSynchronous a boolean value indicating whether disk writes are synchronous. + * @param diskSynchronous boolean value indicating whether disk writes are synchronous. + * @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setDiskSynchronous(boolean) */ public void setDiskSynchronous(Boolean diskSynchronous) { this.diskSynchronous = diskSynchronous; } /** - * Set the number of dispatcher threads used to process Region Events from the associated Async Event Queue (AEQ). + * Configures the number of dispatcher threads used to process Region Events + * from the associated {@link AsyncEventQueue} (AEQ). * - * @param dispatcherThreads an Integer indicating the number of dispatcher threads used to process Region Events - * from the associated Queue. + * @param dispatcherThreads {@link Integer} specifying the number of dispatcher threads used + * to process {@link Region} events from the associated queue. + * @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setDispatcherThreads(int) */ public void setDispatcherThreads(Integer dispatcherThreads) { this.dispatcherThreads = dispatcherThreads; } + /** + * Forwards expiration (action-based) destroy events to the {@link AsyncEventQueue} (AEQ). + * + * By default, destroy events are not added to the AEQ. Setting this attribute to {@literal true} + * will add all expiration destroy events to the AEQ. + * + * @param forwardExpirationDestroy boolean value indicating whether to forward expiration destroy events. + * @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setForwardExpirationDestroy(boolean) + * @see org.apache.geode.cache.ExpirationAttributes#getAction() + * @see org.apache.geode.cache.ExpirationAction#DESTROY + */ + public void setForwardExpirationDestroy(Boolean forwardExpirationDestroy) { + this.forwardExpirationDestroy = forwardExpirationDestroy; + } + + public void setGatewayEventFilters(List eventFilters) { + this.gatewayEventFilters = eventFilters; + } + + public void setGatewayEventSubstitutionFilter(GatewayEventSubstitutionFilter eventSubstitutionFilter) { + this.gatewayEventSubstitutionFilter = eventSubstitutionFilter; + } + public void setMaximumQueueMemory(Integer maximumQueueMemory) { this.maximumQueueMemory = maximumQueueMemory; } /** - * Set the Async Event Queue (AEQ) ordering policy (e.g. KEY, PARTITION, THREAD). When dispatcher threads - * are greater than 1, the ordering policy configures the way in which multiple dispatcher threads - * process Region events from the queue. + * Configures the {@link AsyncEventQueue} (AEQ) ordering policy (e.g. {@literal KEY}, {@literal PARTITION}, + * {@literal THREAD}). * - * @param orderPolicy a String to indicate the AEQ order policy. + * When dispatcher threads are greater than one, the ordering policy configures the way in which + * multiple dispatcher threads process Region events from the queue. + * + * @param orderPolicy {@link String} specifying the name of the AEQ order policy. + * @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setOrderPolicy(GatewaySender.OrderPolicy) */ public void setOrderPolicy(String orderPolicy) { + setOrderPolicy(GatewaySender.OrderPolicy.valueOf(String.valueOf(orderPolicy).toUpperCase())); + } + + public void setOrderPolicy(GatewaySender.OrderPolicy orderPolicy) { this.orderPolicy = orderPolicy; } @@ -233,5 +277,4 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< public void setPersistent(Boolean persistent) { this.persistent = persistent; } - } diff --git a/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java b/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java index 44e69680..d05e1c4b 100644 --- a/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java +++ b/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java @@ -46,12 +46,6 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean eventFilters; - - private List transportFilters; - private Boolean diskSynchronous; private Boolean batchConflationEnabled; private Boolean parallel; @@ -59,6 +53,10 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean eventFilters; + + private List transportFilters; + private String diskStoreReference; - private String orderPolicy; /** * Constructs an instance of the {@link GatewaySenderFactoryBean} class initialized with a reference to @@ -86,8 +87,9 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean getObjectType() { - return (gatewaySender != null ? gatewaySender.getClass() : GatewaySender.class); + return this.gatewaySender != null ? this.gatewaySender.getClass() : GatewaySender.class; } public void setAlertThreshold(Integer alertThreshold) { @@ -235,6 +229,10 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean + + @@ -2977,6 +2979,7 @@ if an inner bean. ]]> + diff --git a/src/main/resources/org/springframework/data/gemfire/config/spring-geode-1.1.xsd b/src/main/resources/org/springframework/data/gemfire/config/spring-geode-1.1.xsd index 46656697..20801af1 100644 --- a/src/main/resources/org/springframework/data/gemfire/config/spring-geode-1.1.xsd +++ b/src/main/resources/org/springframework/data/gemfire/config/spring-geode-1.1.xsd @@ -2931,8 +2931,7 @@ An AsyncEventListener bean definition for this AsyncEventQueue. (requires Gemfir - + + + diff --git a/src/test/java/org/springframework/data/gemfire/config/xml/AsyncEventQueueNamespaceTest.java b/src/test/java/org/springframework/data/gemfire/config/xml/AsyncEventQueueNamespaceTest.java new file mode 100644 index 00000000..40e4ef59 --- /dev/null +++ b/src/test/java/org/springframework/data/gemfire/config/xml/AsyncEventQueueNamespaceTest.java @@ -0,0 +1,194 @@ +/* + * Copyright 2012-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.config.xml; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.is; +import static org.hamcrest.Matchers.notNullValue; + +import java.util.List; +import java.util.stream.Collectors; + +import javax.annotation.Resource; + +import org.apache.geode.cache.EntryEvent; +import org.apache.geode.cache.asyncqueue.AsyncEvent; +import org.apache.geode.cache.asyncqueue.AsyncEventListener; +import org.apache.geode.cache.asyncqueue.AsyncEventQueue; +import org.apache.geode.cache.wan.GatewayEventFilter; +import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter; +import org.apache.geode.cache.wan.GatewayQueueEvent; +import org.apache.geode.cache.wan.GatewaySender; +import org.junit.Assert; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit4.SpringRunner; + +/** + * The AsyncEventQueueNamespaceTest class is a test suite of test cases testing the contract and functionality + * of configuring a Pivotal GemFire or Apache Geode {@link AsyncEventQueue} using the SDG XML namespace. + * + * @author John Blum + * @see org.junit.Test + * @see org.junit.runner.RunWith + * @see org.springframework.data.gemfire.config.AsyncEventQueueParser + * @see org.springframework.data.gemfire.wan.AsyncEventQueueFactoryBean + * @see org.springframework.test.context.ContextConfiguration + * @see org.springframework.test.context.junit4.SpringJUnit4ClassRunner + * @see org.apache.geode.cache.asyncqueue.AsyncEventQueue + * @since 1.0.0 + */ +@RunWith(SpringRunner.class) +@ContextConfiguration +@SuppressWarnings("all") +public class AsyncEventQueueNamespaceTest { + + @Resource(name = "TestAsyncEventQueue") + private AsyncEventQueue asyncEventQueue; + + @Resource(name = "TestAsyncEventQueueWithFilters") + private AsyncEventQueue asyncEventQueueWithFilters; + + @Test + public void asyncEventQueueIsConfiguredProperly() { + + Assert.assertThat(asyncEventQueue, is(notNullValue(AsyncEventQueue.class))); + Assert.assertThat(asyncEventQueue.getId(), is(equalTo("TestAsyncEventQueue"))); + Assert.assertThat(asyncEventQueue.isBatchConflationEnabled(), is(true)); + Assert.assertThat(asyncEventQueue.getBatchSize(), is(equalTo(100))); + Assert.assertThat(asyncEventQueue.getBatchTimeInterval(), is(equalTo(30))); + Assert.assertThat(asyncEventQueue.getDiskStoreName(), is(equalTo("TestDiskStore"))); + Assert.assertThat(asyncEventQueue.isDiskSynchronous(), is(true)); + Assert.assertThat(asyncEventQueue.getDispatcherThreads(), is(equalTo(4))); + Assert.assertThat(asyncEventQueue.isForwardExpirationDestroy(), is(true)); + Assert.assertThat(asyncEventQueue.getMaximumQueueMemory(), is(equalTo(50))); + Assert.assertThat(asyncEventQueue.getOrderPolicy(), is(equalTo(GatewaySender.OrderPolicy.KEY))); + Assert.assertThat(asyncEventQueue.isParallel(), is(false)); + Assert.assertThat(asyncEventQueue.isPersistent(), is(true)); + } + + @Test + public void asyncEventQueueListenerEqualsExpected() { + + AsyncEventListener asyncEventListener = asyncEventQueue.getAsyncEventListener(); + + Assert.assertThat(asyncEventListener, is(notNullValue(AsyncEventListener.class))); + Assert.assertThat(asyncEventListener.toString(), is(equalTo("TestAeqListener"))); + } + + @Test + public void asyncEventQueueWithFiltersIsConfiguredProperly() { + + assertThat(asyncEventQueueWithFilters).isNotNull(); + assertThat(asyncEventQueueWithFilters.getId()).isEqualTo("TestAsyncEventQueueWithFilters"); + + AsyncEventListener listener = asyncEventQueueWithFilters.getAsyncEventListener(); + + assertThat(listener).isNotNull(); + assertThat(listener.toString()).isEqualTo("TestListenerOne"); + + List gatewayEventFilters = asyncEventQueueWithFilters.getGatewayEventFilters(); + + assertThat(gatewayEventFilters).isNotNull(); + assertThat(gatewayEventFilters).hasSize(2); + assertThat(gatewayEventFilters.stream().map(Object::toString).collect(Collectors.toList())) + .containsExactly("GatewayEventFilterOne", "GatewayEventFilterTwo"); + + GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter = + asyncEventQueueWithFilters.getGatewayEventSubstitutionFilter(); + + assertThat(gatewayEventSubstitutionFilter).isNotNull(); + assertThat(gatewayEventSubstitutionFilter.toString()).isEqualTo("GatewayEventSubstitutionFilterOne"); + } + + public static class TestAsyncEventListener implements AsyncEventListener { + + private final String name; + + public TestAsyncEventListener(String name) { + this.name = name; + } + + @Override + public boolean processEvents(List events) { + return false; + } + + @Override + public void close() { + } + + @Override + public String toString() { + return this.name; + } + } + + public static class TestGatewayEventFilter implements GatewayEventFilter { + + private final String name; + + public TestGatewayEventFilter(String name) { + this.name = name; + } + + @Override + public boolean beforeEnqueue(GatewayQueueEvent event) { + return false; + } + + @Override + public boolean beforeTransmit(GatewayQueueEvent event) { + return false; + } + + @Override + public void afterAcknowledgement(GatewayQueueEvent event) { } + + @Override + public void close() { } + + @Override + public String toString() { + return this.name; + } + } + + public static class TestGatewayEventSubstitutionFilter implements GatewayEventSubstitutionFilter { + + private final String name; + + public TestGatewayEventSubstitutionFilter(String name) { + this.name = name; + } + @Override + public Object getSubstituteValue(EntryEvent event) { + return null; + } + + @Override + public void close() { } + + @Override + public String toString() { + return this.name; + } + } +} diff --git a/src/test/java/org/springframework/data/gemfire/test/StubAsyncEventQueueFactory.java b/src/test/java/org/springframework/data/gemfire/test/StubAsyncEventQueueFactory.java index d3af1a21..1f1c6e7e 100644 --- a/src/test/java/org/springframework/data/gemfire/test/StubAsyncEventQueueFactory.java +++ b/src/test/java/org/springframework/data/gemfire/test/StubAsyncEventQueueFactory.java @@ -49,14 +49,15 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory { private GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter; - private OrderPolicy orderPolicy; - private List gatewayEventFilters = new ArrayList<>(); + private OrderPolicy orderPolicy; + private String diskStoreName; @Override - public AsyncEventQueue create(final String name, final AsyncEventListener listener) { + public AsyncEventQueue create(String name, AsyncEventListener listener) { + when(asyncEventQueue.getAsyncEventListener()).thenReturn(listener); when(asyncEventQueue.getBatchSize()).thenReturn(this.batchSize); when(asyncEventQueue.getDiskStoreName()).thenReturn(this.diskStoreName); @@ -76,47 +77,30 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory { return this.asyncEventQueue; } + //The following added in 7.0.1 + public AsyncEventQueueFactory setBatchConflationEnabled(boolean arg0) { + this.batchConflationEnabled = arg0; + return this; + } + @Override public AsyncEventQueueFactory setBatchSize(int batchSize) { this.batchSize = batchSize; return this; } + @Override + public AsyncEventQueueFactory setBatchTimeInterval(int interval) { + this.batchTimeInterval = interval; + return this; + } + @Override public AsyncEventQueueFactory setDiskStoreName(String diskStoreName) { this.diskStoreName = diskStoreName; return this; } - @Override - public AsyncEventQueueFactory setMaximumQueueMemory(int maxQueueMemory) { - this.maxQueueMemory = maxQueueMemory; - return this; - } - - @Override - public AsyncEventQueueFactory setPersistent(boolean persistent) { - this.persistent = persistent; - return this; - } - - @Override - public AsyncEventQueueFactory setParallel(boolean parallel) { - this.parallel = parallel; - return this; - } - - //The following added in 7.0.1 - public AsyncEventQueueFactory setBatchConflationEnabled(boolean arg0) { - this.batchConflationEnabled = arg0; - return this; - } - - public AsyncEventQueueFactory setBatchTimeInterval(int arg0) { - this.batchTimeInterval = arg0; - return this; - } - public AsyncEventQueueFactory setDiskSynchronous(boolean arg0) { this.diskSynchronous = arg0; return this; @@ -127,11 +111,39 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory { return this; } + @Override + public AsyncEventQueueFactory setForwardExpirationDestroy(boolean forward) { + this.forwardExpirationDestroy = forward; + return this; + } + + public AsyncEventQueueFactory setGatewayEventSubstitutionListener(final GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter) { + this.gatewayEventSubstitutionFilter = gatewayEventSubstitutionFilter; + return this; + } + + public AsyncEventQueueFactory setMaximumQueueMemory(int maxQueueMemory) { + this.maxQueueMemory = maxQueueMemory; + return this; + } + public AsyncEventQueueFactory setOrderPolicy(OrderPolicy arg0) { this.orderPolicy = arg0; return this; } + @Override + public AsyncEventQueueFactory setParallel(boolean parallel) { + this.parallel = parallel; + return this; + } + + @Override + public AsyncEventQueueFactory setPersistent(boolean persistent) { + this.persistent = persistent; + return this; + } + @Override public AsyncEventQueueFactory addGatewayEventFilter(final GatewayEventFilter gatewayEventFilter) { gatewayEventFilters.add(gatewayEventFilter); @@ -143,16 +155,4 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory { gatewayEventFilters.remove(gatewayEventFilter); return this; } - - @Override - public AsyncEventQueueFactory setGatewayEventSubstitutionListener(final GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter) { - this.gatewayEventSubstitutionFilter = gatewayEventSubstitutionFilter; - return this; - } - - @Override - public AsyncEventQueueFactory setForwardExpirationDestroy(boolean forward) { - this.forwardExpirationDestroy = forward; - return this; - } } diff --git a/src/test/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBeanTest.java b/src/test/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBeanTest.java index c12c417d..f2102512 100644 --- a/src/test/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBeanTest.java +++ b/src/test/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBeanTest.java @@ -16,21 +16,27 @@ package org.springframework.data.gemfire.wan; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertNotNull; -import static org.junit.Assert.assertNull; -import static org.junit.Assert.assertSame; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.isA; import static org.mockito.Matchers.eq; -import static org.mockito.Matchers.notNull; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import java.util.Arrays; + import org.apache.geode.cache.Cache; import org.apache.geode.cache.asyncqueue.AsyncEventListener; import org.apache.geode.cache.asyncqueue.AsyncEventQueue; import org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory; +import org.apache.geode.cache.wan.GatewayEventFilter; +import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter; import org.apache.geode.cache.wan.GatewaySender; import org.junit.Test; import org.springframework.data.gemfire.TestUtils; @@ -52,277 +58,409 @@ import org.springframework.data.gemfire.TestUtils; */ public class AsyncEventQueueFactoryBeanTest { - protected Cache createMockCacheWithAsyncEventQueueInfrastructure( - final AsyncEventQueueFactory mockAsynEventQueueFactory) { - Cache mockCache = mock(Cache.class); - when((mockCache.createAsyncEventQueueFactory())).thenReturn(mockAsynEventQueueFactory); + private Cache mockCache() { + return mock(Cache.class); + } + + private Cache mockCache(AsyncEventQueueFactory mockAsyncEventQueueFactory) { + + Cache mockCache = mockCache(); + + when((mockCache.createAsyncEventQueueFactory())).thenReturn(mockAsyncEventQueueFactory); + return mockCache; } - protected AsyncEventQueueFactory createMockAsyncEventQueueFactory(final String asyncEventQueueId) { - AsyncEventQueueFactory mockAsyncEventQueueFactory = mock(AsyncEventQueueFactory.class); - AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class); + private AsyncEventQueueFactory mockAsyncEventQueueFactory(String asyncEventQueueId) { - when(mockAsyncEventQueue.getId()).thenReturn(asyncEventQueueId); - when(mockAsyncEventQueueFactory.create(eq(asyncEventQueueId), notNull(AsyncEventListener.class))) + AsyncEventQueueFactory mockAsyncEventQueueFactory = mock(AsyncEventQueueFactory.class); + + AsyncEventQueue mockAsyncEventQueue = mockAsyncEventQueue(asyncEventQueueId); + + when(mockAsyncEventQueueFactory.create(eq(asyncEventQueueId), isA(AsyncEventListener.class))) .thenReturn(mockAsyncEventQueue); return mockAsyncEventQueueFactory; } - protected AsyncEventListener createMockAsyncEventListener() { + private AsyncEventQueue mockAsyncEventQueue(String asyncEventQueueId) { + + AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class); + + when(mockAsyncEventQueue.getId()).thenReturn(asyncEventQueueId); + + return mockAsyncEventQueue; + } + + private AsyncEventListener mockAsyncEventListener() { return mock(AsyncEventListener.class); } - protected void verifyExpectations(final AsyncEventQueueFactory mockAsyncEventQueueFactory, - final AsyncEventQueueFactoryBean factoryBean) throws Exception { - Boolean parallel = TestUtils.readField("parallel", factoryBean); - - verify(mockAsyncEventQueueFactory).setParallel(eq(Boolean.TRUE.equals(parallel))); - - String orderPolicy = TestUtils.readField("orderPolicy", factoryBean); - - if (orderPolicy != null) { - verify(mockAsyncEventQueueFactory).setOrderPolicy( - eq(GatewaySender.OrderPolicy.valueOf(orderPolicy.toUpperCase()))); - } - - Integer dispatcherThreads = TestUtils.readField("dispatcherThreads", factoryBean); - - if (dispatcherThreads != null) { - verify(mockAsyncEventQueueFactory).setDispatcherThreads(eq(dispatcherThreads)); - } - - String diskStoreReference = TestUtils.readField("diskStoreReference", factoryBean); - - if (diskStoreReference != null) { - verify(mockAsyncEventQueueFactory).setDiskStoreName(eq(diskStoreReference)); - } - - Boolean diskSynchronous = TestUtils.readField("diskSynchronous", factoryBean); - - if (diskSynchronous != null) { - verify(mockAsyncEventQueueFactory).setDiskSynchronous(eq(diskSynchronous)); - } - - Boolean persistent = TestUtils.readField("persistent", factoryBean); - - if (persistent != null) { - verify(mockAsyncEventQueueFactory).setPersistent(eq(persistent)); - } - else { - verify(mockAsyncEventQueueFactory, never()).setPersistent(true); - } - } - @Test - public void testSetAsyncEventListener() throws Exception { - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(createMockAsyncEventQueueFactory("testEventQueue"))); + public void setAndGetAsyncEventListener() throws Exception { - AsyncEventListener listenerOne = createMockAsyncEventListener(); + AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(mockCache()); + + AsyncEventListener listenerOne = mockAsyncEventListener(); factoryBean.setAsyncEventListener(listenerOne); - assertSame(listenerOne, TestUtils.readField("asyncEventListener", factoryBean)); + assertThat(TestUtils.readField("asyncEventListener", factoryBean)) + .isSameAs(listenerOne); - AsyncEventListener listenerTwo = createMockAsyncEventListener(); + AsyncEventListener listenerTwo = mockAsyncEventListener(); factoryBean.setAsyncEventListener(listenerTwo); - assertSame(listenerTwo, TestUtils.readField("asyncEventListener", factoryBean)); + assertThat(TestUtils.readField("asyncEventListener", factoryBean)) + .isSameAs(listenerTwo); } @Test(expected = IllegalStateException.class) - public void testSetAsyncEventListenerAfterAsyncEventQueueCreation() throws Exception { - String asyncEventQueueId = "testEventQueue"; + public void setAsyncEventListenerAfterAsyncEventQueueCreationThrowsIllegalStateException() throws Exception { - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(createMockAsyncEventQueueFactory(asyncEventQueueId))); + AsyncEventListener mockAsyncEventListener = mockAsyncEventListener(); - factoryBean.setName(asyncEventQueueId); + AsyncEventQueue mockAsyncEventQueue = mockAsyncEventQueue("testEventQueue"); - AsyncEventListener listenerOne = createMockAsyncEventListener(); + AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(mockCache(), mockAsyncEventListener); - factoryBean.setAsyncEventListener(listenerOne); - - assertSame(listenerOne, TestUtils.readField("asyncEventListener", factoryBean)); - - factoryBean.doInit(); - - assertNotNull(TestUtils.readField("asyncEventQueue", factoryBean)); + factoryBean.setAsyncEventQueue(mockAsyncEventQueue); try { - factoryBean.setAsyncEventListener(createMockAsyncEventListener()); + factoryBean.setAsyncEventListener(mockAsyncEventListener()); } catch (IllegalStateException expected) { - assertEquals("Setting an AsyncEventListener is not allowed once the AsyncEventQueue has been created.", - expected.getMessage()); - assertSame(listenerOne, TestUtils.readField("asyncEventListener", factoryBean)); + + assertThat(expected) + .hasMessage("Setting an AsyncEventListener is not allowed once the AsyncEventQueue has been created"); + + assertThat(expected).hasNoCause(); + throw expected; } - } - - @Test(expected = IllegalArgumentException.class) - public void testDoInitWhenAsyncEventListenerIsNull() throws Exception { - try { - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(createMockAsyncEventQueueFactory("testEventQueue"))); - - assertNull(TestUtils.readField("asyncEventListener", factoryBean)); - - factoryBean.doInit(); - } - catch (Exception e) { - assertEquals("The AsyncEventListener cannot be null.", e.getMessage()); - throw e; + finally { + assertThat(TestUtils.readField("asyncEventListener", factoryBean)) + .isSameAs(mockAsyncEventListener); } } @Test - public void testConcurrentParallelAsyncEventQueue() throws Exception { - AsyncEventQueueFactory mockAsyncEventQueueFactory = createMockAsyncEventQueueFactory("000"); + public void doInitConfiguresAsyncEventQueue() throws Exception { - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory)); + AsyncEventListener mockListener = mockAsyncEventListener(); - factoryBean.setName("000"); - factoryBean.setAsyncEventListener(createMockAsyncEventListener()); - factoryBean.setDispatcherThreads(8); - factoryBean.setParallel(true); - factoryBean.doInit(); + AsyncEventQueueFactory mockAsyncEventQueueFactory = mockAsyncEventQueueFactory("testQueue"); - verifyExpectations(mockAsyncEventQueueFactory, factoryBean); + AsyncEventQueueFactoryBean factoryBean = + new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory), mockListener); - AsyncEventQueue eventQueue = factoryBean.getObject(); + GatewayEventFilter mockGatewayEventFilterOne = mock(GatewayEventFilter.class); + GatewayEventFilter mockGatewayEventFilterTwo = mock(GatewayEventFilter.class); - assertNotNull(eventQueue); - assertEquals("000", eventQueue.getId()); - } + GatewayEventSubstitutionFilter mockGatewayEventSubstitutionFilter = mock(GatewayEventSubstitutionFilter.class); - @Test - public void testParallelAsyncEventQueue() throws Exception { - AsyncEventQueueFactory mockAsyncEventQueueFactory = createMockAsyncEventQueueFactory("123"); - - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory)); - - factoryBean.setName("123"); - factoryBean.setAsyncEventListener(createMockAsyncEventListener()); - factoryBean.setParallel(true); - factoryBean.doInit(); - - verifyExpectations(mockAsyncEventQueueFactory, factoryBean); - - AsyncEventQueue eventQueue = factoryBean.getObject(); - - assertNotNull(eventQueue); - assertEquals("123", eventQueue.getId()); - } - - @Test(expected = IllegalArgumentException.class) - public void testParallelAsyncEventQueueWithOrderPolicy() { - AsyncEventQueueFactory mockAsyncEventQueueFactory = createMockAsyncEventQueueFactory("456"); - - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory)); - - factoryBean.setName("456"); - factoryBean.setAsyncEventListener(createMockAsyncEventListener()); - factoryBean.setOrderPolicy("Key"); - factoryBean.setParallel(true); - - try { - factoryBean.doInit(); - } - catch (IllegalArgumentException expected) { - assertEquals("Order Policy cannot be used with a Parallel Event Queue.", - expected.getMessage()); - throw expected; - } - } - - @Test - public void testSerialAsyncEventQueueWithOrderPolicy() throws Exception { - AsyncEventQueueFactory mockAsyncEventQueueFatory = createMockAsyncEventQueueFactory("789"); - - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFatory)); - - factoryBean.setName("789"); - factoryBean.setAsyncEventListener(createMockAsyncEventListener()); - factoryBean.setOrderPolicy("THREAD"); - factoryBean.setParallel(false); - factoryBean.doInit(); - - verifyExpectations(mockAsyncEventQueueFatory, factoryBean); - - AsyncEventQueue eventQueue = factoryBean.getObject(); - - assertNotNull(eventQueue); - assertEquals("789", eventQueue.getId()); - } - - @Test - public void testAsyncEventQueueWithOrderPolicyAndDispatcherThreads() throws Exception { - AsyncEventQueueFactory mockAsyncEventQueueFatory = createMockAsyncEventQueueFactory("abc"); - - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFatory)); - - factoryBean.setName("abc"); - factoryBean.setAsyncEventListener(createMockAsyncEventListener()); + factoryBean.setBatchConflationEnabled(true); + factoryBean.setBatchSize(1024); + factoryBean.setBatchTimeInterval(600); + factoryBean.setDiskStoreRef("testDiskStore"); + factoryBean.setDiskSynchronous(false); factoryBean.setDispatcherThreads(2); - factoryBean.setOrderPolicy("THREAD"); - factoryBean.doInit(); - - verifyExpectations(mockAsyncEventQueueFatory, factoryBean); - - AsyncEventQueue eventQueue = factoryBean.getObject(); - - assertNotNull(eventQueue); - assertEquals("abc", eventQueue.getId()); - } - - @Test - public void testAsyncEventQueueWithOverflowDiskStoreNoPersistence() throws Exception { - AsyncEventQueueFactory mockAsyncEventQueueFactory = createMockAsyncEventQueueFactory("123abc"); - - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory)); - - factoryBean.setName("123abc"); - factoryBean.setAsyncEventListener(createMockAsyncEventListener()); - factoryBean.setDiskStoreRef("queueOverflowDiskStore"); + factoryBean.setForwardExpirationDestroy(true); + factoryBean.setGatewayEventFilters(Arrays.asList(mockGatewayEventFilterOne, mockGatewayEventFilterTwo)); + factoryBean.setGatewayEventSubstitutionFilter(mockGatewayEventSubstitutionFilter); + factoryBean.setMaximumQueueMemory(8192); + factoryBean.setName("testQueue"); + factoryBean.setOrderPolicy(GatewaySender.OrderPolicy.PARTITION); + factoryBean.setParallel(false); factoryBean.setPersistent(false); factoryBean.doInit(); - verifyExpectations(mockAsyncEventQueueFactory, factoryBean); + verify(mockAsyncEventQueueFactory, times(1)).setBatchConflationEnabled(eq(true)); + verify(mockAsyncEventQueueFactory, times(1)).setBatchSize(eq(1024)); + verify(mockAsyncEventQueueFactory, times(1)).setBatchTimeInterval(eq(600)); + verify(mockAsyncEventQueueFactory, times(1)).setDiskStoreName(eq("testDiskStore")); + verify(mockAsyncEventQueueFactory, times(1)).setDiskSynchronous(eq(false)); + verify(mockAsyncEventQueueFactory, times(1)).setDispatcherThreads(eq(2)); + verify(mockAsyncEventQueueFactory, times(1)).setForwardExpirationDestroy(eq(true)); + verify(mockAsyncEventQueueFactory, times(1)) + .setGatewayEventSubstitutionListener(eq(mockGatewayEventSubstitutionFilter)); + verify(mockAsyncEventQueueFactory, times(1)).setMaximumQueueMemory(eq(8192)); + verify(mockAsyncEventQueueFactory, times(1)).setOrderPolicy(eq(GatewaySender.OrderPolicy.PARTITION)); + verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false)); + verify(mockAsyncEventQueueFactory, times(1)).setPersistent(eq(false)); + verify(mockAsyncEventQueueFactory, times(1)).addGatewayEventFilter(eq(mockGatewayEventFilterOne)); + verify(mockAsyncEventQueueFactory, times(1)).addGatewayEventFilter(eq(mockGatewayEventFilterTwo)); - AsyncEventQueue evenQueue = factoryBean.getObject(); + AsyncEventQueue asyncEventQueue = factoryBean.getObject(); - assertNotNull(evenQueue); - assertEquals("123abc", evenQueue.getId()); + assertThat(asyncEventQueue).isNotNull(); + assertThat(asyncEventQueue.getId()).isEqualTo("testQueue"); } @Test - public void testAsyncEventQueueWithDiskSynchronousSetPersistenceUnset() throws Exception { - AsyncEventQueueFactory mockAsyncEventQueueFactory = createMockAsyncEventQueueFactory("12345"); + public void doInitConfiguresConcurrentParallelAsyncEventQueue() throws Exception { - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory)); + AsyncEventQueueFactory mockAsyncEventQueueFactory = + mockAsyncEventQueueFactory("concurrentParallelQueue"); - factoryBean.setName("12345"); - factoryBean.setAsyncEventListener(createMockAsyncEventListener()); - factoryBean.setDiskSynchronous(true); + AsyncEventQueueFactoryBean factoryBean = + new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory)); + + factoryBean.setAsyncEventListener(mockAsyncEventListener()); + factoryBean.setDispatcherThreads(8); + factoryBean.setName("concurrentParallelQueue"); + factoryBean.setParallel(true); factoryBean.doInit(); - verifyExpectations(mockAsyncEventQueueFactory, factoryBean); + verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString()); + verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean()); + verify(mockAsyncEventQueueFactory, times(1)).setDispatcherThreads(eq(8)); + verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any()); + verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setOrderPolicy(any()); + verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true)); + verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any()); - AsyncEventQueue evenQueue = factoryBean.getObject(); + AsyncEventQueue asyncEventQueue = factoryBean.getObject(); - assertNotNull(evenQueue); - assertEquals("12345", evenQueue.getId()); + assertThat(asyncEventQueue).isNotNull(); + assertThat(asyncEventQueue.getId()).isEqualTo("concurrentParallelQueue"); } + @Test + public void doInitConfiguresParallelAsyncEventQueue() throws Exception { + + AsyncEventQueueFactory mockAsyncEventQueueFactory = + mockAsyncEventQueueFactory("parallelQueue"); + + AsyncEventQueueFactoryBean factoryBean = + new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory)); + + factoryBean.setAsyncEventListener(mockAsyncEventListener()); + factoryBean.setName("parallelQueue"); + factoryBean.setParallel(true); + factoryBean.doInit(); + + verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString()); + verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setDispatcherThreads(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any()); + verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setOrderPolicy(any()); + verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true)); + verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any()); + + AsyncEventQueue asyncEventQueue = factoryBean.getObject(); + + assertThat(asyncEventQueue).isNotNull(); + assertThat(asyncEventQueue.getId()).isEqualTo("parallelQueue"); + } + + @Test + public void doInitConfiguresConcurrentSerialAsyncEventQueue() throws Exception { + + AsyncEventQueueFactory mockAsyncEventQueueFactory = + mockAsyncEventQueueFactory("concurrentSerialQueue"); + + AsyncEventQueueFactoryBean factoryBean = + new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory)); + + factoryBean.setAsyncEventListener(mockAsyncEventListener()); + factoryBean.setDispatcherThreads(16); + factoryBean.setName("concurrentSerialQueue"); + factoryBean.setParallel(false); + factoryBean.doInit(); + + verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString()); + verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean()); + verify(mockAsyncEventQueueFactory, times(1)).setDispatcherThreads(eq(16)); + verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any()); + verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setOrderPolicy(any()); + verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false)); + verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any()); + + AsyncEventQueue asyncEventQueue = factoryBean.getObject(); + + assertThat(asyncEventQueue).isNotNull(); + assertThat(asyncEventQueue.getId()).isEqualTo("concurrentSerialQueue"); + } + + @Test + public void doInitConfiguresSerialAsyncEventQueue() throws Exception { + + AsyncEventQueueFactory mockAsyncEventQueueFactory = mockAsyncEventQueueFactory("serialQueue"); + + AsyncEventQueueFactoryBean factoryBean = + new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory)); + + factoryBean.setAsyncEventListener(mockAsyncEventListener()); + factoryBean.setName("serialQueue"); + factoryBean.doInit(); + + verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString()); + verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setDispatcherThreads(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any()); + verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setOrderPolicy(any()); + verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false)); + verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any()); + + AsyncEventQueue asyncEventQueue = factoryBean.getObject(); + + assertThat(asyncEventQueue).isNotNull(); + assertThat(asyncEventQueue.getId()).isEqualTo("serialQueue"); + } + + @Test + public void doInitConfiguresSerialAsyncEventQueueWithOrderPolicy() throws Exception { + + AsyncEventQueueFactory mockAsyncEventQueueFactory = + mockAsyncEventQueueFactory("orderedSerialQueue"); + + AsyncEventQueueFactoryBean factoryBean = + new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory)); + + factoryBean.setAsyncEventListener(mockAsyncEventListener()); + factoryBean.setName("orderedSerialQueue"); + factoryBean.setOrderPolicy(GatewaySender.OrderPolicy.THREAD); + factoryBean.doInit(); + + verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString()); + verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setDispatcherThreads(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any()); + verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt()); + verify(mockAsyncEventQueueFactory, times(1)).setOrderPolicy(eq(GatewaySender.OrderPolicy.THREAD)); + verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false)); + verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any()); + + AsyncEventQueue asyncEventQueue = factoryBean.getObject(); + + assertThat(asyncEventQueue).isNotNull(); + assertThat(asyncEventQueue.getId()).isEqualTo("orderedSerialQueue"); + } + + @Test(expected = IllegalStateException.class) + public void doInitWithNullAsyncEventListenerThrowsIllegalStateException() throws Exception { + + try { + new AsyncEventQueueFactoryBean(mockCache(), null).doInit(); + } + catch (IllegalStateException expected) { + + assertThat(expected).hasMessage("AsyncEventListener must not be null"); + assertThat(expected).hasNoCause(); + + throw expected; + } + } + + @Test(expected = IllegalStateException.class) + public void doInitWithParallelAsyncEventQueueHavingAnOrderPolicyThrowsIllegalStateException() throws Exception { + + AsyncEventQueueFactory mockAsyncEventQueueFactory = + mockAsyncEventQueueFactory("parallelQueue"); + + AsyncEventQueueFactoryBean factoryBean = + new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory)); + + factoryBean.setAsyncEventListener(mockAsyncEventListener()); + factoryBean.setName("parallelQueue"); + factoryBean.setOrderPolicy(GatewaySender.OrderPolicy.KEY); + factoryBean.setParallel(true); + + try { + factoryBean.doInit(); + } + catch (IllegalStateException expected) { + + assertThat(expected).hasMessage("OrderPolicy cannot be used with a Parallel AsyncEventQueue"); + assertThat(expected).hasNoCause(); + + throw expected; + } + finally { + + assertThat(factoryBean.getObject()).isNull(); + + verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString()); + verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setDispatcherThreads(eq(8)); + verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any()); + verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setOrderPolicy(any()); + verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true)); + verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any()); + } + } + + @Test + public void doInitConfiguresAsyncEventQueueWithSynchronousOverflowDiskStoreNoPersistence() throws Exception { + + AsyncEventQueueFactory mockAsyncEventQueueFactory = + mockAsyncEventQueueFactory("nonPersistentSynchronousOverflowQueue"); + + AsyncEventQueueFactoryBean factoryBean = + new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory)); + + factoryBean.setAsyncEventListener(mockAsyncEventListener()); + factoryBean.setDiskStoreRef("queueOverflowDiskStore"); + factoryBean.setDiskSynchronous(true); + factoryBean.setName("nonPersistentSynchronousOverflowQueue"); + factoryBean.setOrderPolicy(GatewaySender.OrderPolicy.KEY); + factoryBean.setPersistent(false); + factoryBean.doInit(); + + verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt()); + verify(mockAsyncEventQueueFactory, times(1)).setDiskStoreName("queueOverflowDiskStore"); + verify(mockAsyncEventQueueFactory, times(1)).setDiskSynchronous(eq(true)); + verify(mockAsyncEventQueueFactory, never()).setDispatcherThreads(anyInt()); + verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean()); + verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any()); + verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt()); + verify(mockAsyncEventQueueFactory, times(1)).setOrderPolicy(eq(GatewaySender.OrderPolicy.KEY)); + verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false)); + verify(mockAsyncEventQueueFactory, times(1)).setPersistent(eq(false)); + verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any()); + + AsyncEventQueue asyncEventQueue = factoryBean.getObject(); + + assertThat(asyncEventQueue).isNotNull(); + assertThat(asyncEventQueue.getId()).isEqualTo("nonPersistentSynchronousOverflowQueue"); + } } diff --git a/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java b/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java index f46e1326..e28bc697 100644 --- a/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java +++ b/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java @@ -65,17 +65,17 @@ public class GatewaySenderFactoryBeanTest { return mockGatewaySenderFactory; } - protected void verifyExpectations(final GatewaySenderFactoryBean factoryBean, - final GatewaySenderFactory mockGatewaySenderFactory) throws Exception { + protected void verifyExpectations(GatewaySenderFactoryBean factoryBean, + GatewaySenderFactory mockGatewaySenderFactory) throws Exception { + Boolean parallel = TestUtils.readField("parallel", factoryBean); verify(mockGatewaySenderFactory).setParallel(eq(Boolean.TRUE.equals(parallel))); - String orderPolicy = TestUtils.readField("orderPolicy", factoryBean); + GatewaySender.OrderPolicy orderPolicy = TestUtils.readField("orderPolicy", factoryBean); if (orderPolicy != null) { - verify(mockGatewaySenderFactory).setOrderPolicy( - eq(GatewaySender.OrderPolicy.valueOf(orderPolicy.toUpperCase()))); + verify(mockGatewaySenderFactory).setOrderPolicy(eq(orderPolicy)); } Integer dispatcherThreads = TestUtils.readField("dispatcherThreads", factoryBean); @@ -165,7 +165,7 @@ public class GatewaySenderFactoryBeanTest { factoryBean.doInit(); } catch (IllegalArgumentException expected) { - assertEquals("Order Policy cannot be used with a Parallel Gateway Sender Queue.", expected.getMessage()); + assertEquals("OrderPolicy cannot be used with a Parallel GatewaySender", expected.getMessage()); throw expected; } } @@ -194,7 +194,9 @@ public class GatewaySenderFactoryBeanTest { @Test public void testGatewaySenderWithOrderPolicyAndDispatcherThreads() throws Exception { - GatewaySenderFactory mockGatewaySenderFactory = createMockGatewaySenderFactory("g5", 42); + + GatewaySenderFactory mockGatewaySenderFactory = + createMockGatewaySenderFactory("g5", 42); GatewaySenderFactoryBean factoryBean = new GatewaySenderFactoryBean( createMockCacheWithGatewayInfrastructure(mockGatewaySenderFactory)); @@ -256,5 +258,4 @@ public class GatewaySenderFactoryBeanTest { assertEquals("g7", gatewaySender.getId()); assertEquals(51, gatewaySender.getRemoteDSId()); } - } diff --git a/src/test/resources/org/springframework/data/gemfire/config/xml/AsyncEventQueueNamespaceTest-context.xml b/src/test/resources/org/springframework/data/gemfire/config/xml/AsyncEventQueueNamespaceTest-context.xml new file mode 100644 index 00000000..7129e989 --- /dev/null +++ b/src/test/resources/org/springframework/data/gemfire/config/xml/AsyncEventQueueNamespaceTest-context.xml @@ -0,0 +1,64 @@ + + + + + + + AsyncEventQueueNamespaceTest + 0 + error + + + + + + + + + + + + + + + + + + + + + + + + + + + + + +