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 2def72b5..e7b4a085 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 @@ -62,19 +62,37 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser { ParsingUtils.setPropertyValue(element, builder, "batch-time-interval"); ParsingUtils.setPropertyValue(element, builder, "disk-synchronous"); ParsingUtils.setPropertyValue(element, builder, "dispatcher-threads"); - ParsingUtils.setPropertyValue(element, builder, "foward-expiration-destroy"); + 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"); + + 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; @@ -91,11 +109,11 @@ 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); @@ -106,15 +124,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 eab60861..c603c31c 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 acd1e8cd..baae7113 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 { @@ -47,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. @@ -67,81 +87,59 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< * @param asyncEventListener required {@link 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 != null ? asyncEventQueue.getClass() : AsyncEventQueue.class); + return this.asyncEventQueue != null ? this.asyncEventQueue.getClass() : AsyncEventQueue.class; } @Override protected void doInit() { - Assert.notNull(this.asyncEventListener, "AsyncEventListener must not be null"); + Assert.state(this.asyncEventListener != null, "AsyncEventListener must not be null"); - AsyncEventQueueFactory asyncEventQueueFactory = (this.factory != null ? (AsyncEventQueueFactory) factory - : cache.createAsyncEventQueueFactory()); + AsyncEventQueueFactory asyncEventQueueFactory = + this.factory != null ? (AsyncEventQueueFactory) this.factory : this.cache.createAsyncEventQueueFactory(); - if (batchSize != null) { - asyncEventQueueFactory.setBatchSize(batchSize); - } - - 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 (forwardExpirationDestroy != null) { - asyncEventQueueFactory.setForwardExpirationDestroy(forwardExpirationDestroy); - } - - 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(); } @@ -151,6 +149,7 @@ 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"); @@ -158,16 +157,20 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< } /** - * @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; @@ -178,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; @@ -192,19 +196,22 @@ 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; @@ -213,8 +220,8 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< /** * 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. + * 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) @@ -225,18 +232,33 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< 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; } diff --git a/src/main/resources/org/springframework/data/gemfire/config/spring-geode-2.1.xsd b/src/main/resources/org/springframework/data/gemfire/config/spring-geode-2.1.xsd index 86f2929f..7864a135 100644 --- a/src/main/resources/org/springframework/data/gemfire/config/spring-geode-2.1.xsd +++ b/src/main/resources/org/springframework/data/gemfire/config/spring-geode-2.1.xsd @@ -2955,8 +2955,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 index e06b1e59..40e4ef59 100644 --- a/src/test/java/org/springframework/data/gemfire/config/xml/AsyncEventQueueNamespaceTest.java +++ b/src/test/java/org/springframework/data/gemfire/config/xml/AsyncEventQueueNamespaceTest.java @@ -17,22 +17,29 @@ 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 static org.junit.Assert.assertThat; 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.beans.factory.annotation.Autowired; import org.springframework.test.context.ContextConfiguration; -import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.test.context.junit4.SpringRunner; /** * The AsyncEventQueueNamespaceTest class is a test suite of test cases testing the contract and functionality @@ -48,37 +55,67 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; * @see org.apache.geode.cache.asyncqueue.AsyncEventQueue * @since 1.0.0 */ -@RunWith(SpringJUnit4ClassRunner.class) +@RunWith(SpringRunner.class) @ContextConfiguration @SuppressWarnings("all") public class AsyncEventQueueNamespaceTest { - @Autowired + @Resource(name = "TestAsyncEventQueue") private AsyncEventQueue asyncEventQueue; + @Resource(name = "TestAsyncEventQueueWithFilters") + private AsyncEventQueue asyncEventQueueWithFilters; + @Test public void asyncEventQueueIsConfiguredProperly() { - assertThat(asyncEventQueue, is(notNullValue(AsyncEventQueue.class))); - assertThat(asyncEventQueue.getId(), is(equalTo("TestAsyncEventQueue"))); - assertThat(asyncEventQueue.isBatchConflationEnabled(), is(true)); - assertThat(asyncEventQueue.getBatchSize(), is(equalTo(100))); - assertThat(asyncEventQueue.getBatchTimeInterval(), is(equalTo(30))); - assertThat(asyncEventQueue.getDiskStoreName(), is(equalTo("TestDiskStore"))); - assertThat(asyncEventQueue.isDiskSynchronous(), is(true)); - assertThat(asyncEventQueue.getDispatcherThreads(), is(equalTo(4))); - assertThat(asyncEventQueue.isForwardExpirationDestroy(), is(false)); - assertThat(asyncEventQueue.getMaximumQueueMemory(), is(equalTo(50))); - assertThat(asyncEventQueue.getOrderPolicy(), is(equalTo(GatewaySender.OrderPolicy.KEY))); - assertThat(asyncEventQueue.isParallel(), is(false)); - assertThat(asyncEventQueue.isPersistent(), is(true)); + + 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(); - assertThat(asyncEventListener, is(notNullValue(AsyncEventListener.class))); - assertThat(asyncEventListener.toString(), is(equalTo("TestAeqListener"))); + 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 { @@ -103,4 +140,55 @@ public class AsyncEventQueueNamespaceTest { 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 a4f12d51..0009d02f 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 List gatewayEventFilters = new ArrayList(); + private OrderPolicy orderPolicy; private String diskStoreName; @Override public AsyncEventQueue create(String name, AsyncEventListener listener) { + when(asyncEventQueue.getAsyncEventListener()).thenReturn(listener); when(asyncEventQueue.isBatchConflationEnabled()).thenReturn(this.batchConflationEnabled); when(asyncEventQueue.getBatchSize()).thenReturn(this.batchSize); @@ -113,6 +114,11 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory { return this; } + public AsyncEventQueueFactory setGatewayEventSubstitutionListener(final GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter) { + this.gatewayEventSubstitutionFilter = gatewayEventSubstitutionFilter; + return this; + } + public AsyncEventQueueFactory setMaximumQueueMemory(int maxQueueMemory) { this.maxQueueMemory = maxQueueMemory; return this; @@ -142,9 +148,4 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory { gatewayEventFilters.remove(gatewayEventFilter); return this; } - - public AsyncEventQueueFactory setGatewayEventSubstitutionListener(final GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter) { - this.gatewayEventSubstitutionFilter = gatewayEventSubstitutionFilter; - 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 fa83a1e6..32ad6726 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,276 +58,409 @@ import org.springframework.data.gemfire.TestUtils; */ public class AsyncEventQueueFactoryBeanTest { - protected Cache createMockCacheWithAsyncEventQueueInfrastructure(AsyncEventQueueFactory mockAsyncEventQueueFactory) { - Cache mockCache = mock(Cache.class); + 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(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 mockAsyncEventListener() { + 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(AsyncEventQueueFactory asyncEventQueueFactory, - AsyncEventQueueFactoryBean asyncEventQueueFactoryBean) throws Exception { - - Boolean parallel = TestUtils.readField("parallel", asyncEventQueueFactoryBean); - - verify(asyncEventQueueFactory).setParallel(eq(Boolean.TRUE.equals(parallel))); - - String orderPolicy = TestUtils.readField("orderPolicy", asyncEventQueueFactoryBean); - - if (orderPolicy != null) { - verify(asyncEventQueueFactory).setOrderPolicy(eq(GatewaySender.OrderPolicy.valueOf( - orderPolicy.toUpperCase()))); - } - - Integer dispatcherThreads = TestUtils.readField("dispatcherThreads", asyncEventQueueFactoryBean); - - if (dispatcherThreads != null) { - verify(asyncEventQueueFactory).setDispatcherThreads(eq(dispatcherThreads)); - } - - String diskStoreReference = TestUtils.readField("diskStoreReference", asyncEventQueueFactoryBean); - - if (diskStoreReference != null) { - verify(asyncEventQueueFactory).setDiskStoreName(eq(diskStoreReference)); - } - - Boolean diskSynchronous = TestUtils.readField("diskSynchronous", asyncEventQueueFactoryBean); - - if (diskSynchronous != null) { - verify(asyncEventQueueFactory).setDiskSynchronous(eq(diskSynchronous)); - } - - Boolean persistent = TestUtils.readField("persistent", asyncEventQueueFactoryBean); - - if (persistent != null) { - verify(asyncEventQueueFactory).setPersistent(eq(persistent)); - } - else { - verify(asyncEventQueueFactory, never()).setPersistent(true); - } - } - @Test - public void testSetAsyncEventListener() throws Exception { - AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean( - createMockCacheWithAsyncEventQueueInfrastructure(createMockAsyncEventQueueFactory("testEventQueue"))); + public void setAndGetAsyncEventListener() throws Exception { + + 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 = 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 = mockAsyncEventListener(); + 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(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("AsyncEventListener must not 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(mockAsyncEventListener()); - 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(mockAsyncEventListener()); - 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(mockAsyncEventListener()); - 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(mockAsyncEventListener()); - 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(mockAsyncEventListener()); + 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(mockAsyncEventListener()); - 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"); + + AsyncEventQueueFactoryBean factoryBean = + new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory)); - factoryBean.setName("12345"); factoryBean.setAsyncEventListener(mockAsyncEventListener()); - factoryBean.setDiskSynchronous(true); + 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/resources/org/springframework/data/gemfire/config/xml/AsyncEventQueueNamespaceTest-context.xml b/src/test/resources/org/springframework/data/gemfire/config/xml/AsyncEventQueueNamespaceTest-context.xml index db919340..9419d5e4 100644 --- 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 @@ -17,7 +17,7 @@ AsyncEventQueueNamespaceTest 0 - warning + error @@ -41,9 +41,24 @@ parallel="false" persistent="true"> - + + + + + + + + + + + + + +