DATAGEODE-91 - Impossible to define event filter for AsyncEventQueue.

This commit is contained in:
John Blum
2018-03-22 17:58:19 -07:00
parent 781813448f
commit 867ad2062f
9 changed files with 620 additions and 341 deletions

View File

@@ -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");

View File

@@ -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<T> implements BeanNameAware, FactoryBean<T>,
InitializingBean, DisposableBean {
protected static final List<String> VALID_ORDER_POLICIES = Arrays.asList("KEY", "PARTITION", "THREAD");
protected Log log = LogFactory.getLog(getClass());
public abstract class AbstractWANComponentFactoryBean<T>
implements BeanNameAware, DisposableBean, FactoryBean<T>, 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<T> 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<T> 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 { }
}

View File

@@ -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<AsyncEventQueue> {
@@ -47,8 +62,13 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
private Integer dispatcherThreads;
private Integer maximumQueueMemory;
private GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter;
private GatewaySender.OrderPolicy orderPolicy;
private List<GatewayEventFilter> 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<GatewayEventFilter> 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;
}

View File

@@ -2931,8 +2931,7 @@ An AsyncEventListener bean definition for this AsyncEventQueue. (requires Gemfir
</xsd:annotation>
<xsd:complexType>
<xsd:sequence>
<xsd:any namespace="##other" processContents="skip"
minOccurs="0" maxOccurs="unbounded">
<xsd:any namespace="##other" processContents="skip" minOccurs="0" maxOccurs="unbounded">
<xsd:annotation>
<xsd:documentation><![CDATA[
Inner bean definition of the async event listener
@@ -2950,6 +2949,8 @@ use inner bean declarations.
</xsd:attribute>
</xsd:complexType>
</xsd:element>
<xsd:element name="event-filter" type="gatewayEventFilterType" minOccurs="0" maxOccurs="1"/>
<xsd:element name="event-substitution-filter" type="gatewayEventSubstitutionFilterType" minOccurs="0" maxOccurs="1"/>
</xsd:sequence>
<xsd:attribute name="name" type="xsd:string" use="optional">
<xsd:annotation>

View File

@@ -2971,6 +2971,8 @@ use inner bean declarations.
</xsd:attribute>
</xsd:complexType>
</xsd:element>
<xsd:element name="event-filter" type="gatewayEventFilterType" minOccurs="0" maxOccurs="1"/>
<xsd:element name="event-substitution-filter" type="gatewayEventSubstitutionFilterType" minOccurs="0" maxOccurs="1"/>
</xsd:sequence>
<xsd:attribute name="name" type="xsd:string" use="optional">
<xsd:annotation>

View File

@@ -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<GatewayEventFilter> 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<Object, Object> {
private final String name;
public TestGatewayEventSubstitutionFilter(String name) {
this.name = name;
}
@Override
public Object getSubstituteValue(EntryEvent<Object, Object> event) {
return null;
}
@Override
public void close() { }
@Override
public String toString() {
return this.name;
}
}
}

View File

@@ -49,14 +49,15 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
private GatewayEventSubstitutionFilter<?, ?> gatewayEventSubstitutionFilter;
private OrderPolicy orderPolicy;
private List<GatewayEventFilter> gatewayEventFilters = new ArrayList<>();
private List<GatewayEventFilter> gatewayEventFilters = new ArrayList<GatewayEventFilter>();
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;
}
}

View File

@@ -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.<AsyncEventListener>readField("asyncEventListener", factoryBean))
.isSameAs(listenerOne);
AsyncEventListener listenerTwo = mockAsyncEventListener();
factoryBean.setAsyncEventListener(listenerTwo);
assertSame(listenerTwo, TestUtils.readField("asyncEventListener", factoryBean));
assertThat(TestUtils.<AsyncEventListener>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.<AsyncEventListener>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");
}
}

View File

@@ -17,7 +17,7 @@
<util:properties id="gemfireProperties">
<prop key="name">AsyncEventQueueNamespaceTest</prop>
<prop key="mcast-port">0</prop>
<prop key="log-level">warning</prop>
<prop key="log-level">error</prop>
</util:properties>
<context:property-placeholder/>
@@ -41,9 +41,24 @@
parallel="false"
persistent="true">
<gfe:async-event-listener>
<bean class="org.springframework.data.gemfire.config.xml.AsyncEventQueueNamespaceTest.TestAsyncEventListener"
c:name="TestAeqListener"/>
<bean c:name="TestAeqListener"
class="org.springframework.data.gemfire.config.xml.AsyncEventQueueNamespaceTest.TestAsyncEventListener"/>
</gfe:async-event-listener>
</gfe:async-event-queue>
<bean id="testListenerOne" class="org.springframework.data.gemfire.config.xml.AsyncEventQueueNamespaceTest.TestAsyncEventListener"
c:name="TestListenerOne"/>
<gfe:async-event-queue id="TestAsyncEventQueueWithFilters">
<gfe:async-event-listener ref="testListenerOne"/>
<gfe:event-filter>
<bean class="org.springframework.data.gemfire.config.xml.AsyncEventQueueNamespaceTest.TestGatewayEventFilter" c:name="GatewayEventFilterOne"/>
<bean class="org.springframework.data.gemfire.config.xml.AsyncEventQueueNamespaceTest.TestGatewayEventFilter" c:name="GatewayEventFilterTwo"/>
</gfe:event-filter>
<gfe:event-substitution-filter>
<bean class="org.springframework.data.gemfire.config.xml.AsyncEventQueueNamespaceTest.TestGatewayEventSubstitutionFilter"
c:name="GatewayEventSubstitutionFilterOne"/>
</gfe:event-substitution-filter>
</gfe:async-event-queue>
</beans>