SGF-726 - Impossible to define event filter for AsyncEventQueue.
This commit is contained in:
@@ -20,7 +20,6 @@ import org.springframework.beans.factory.config.RuntimeBeanReference;
|
||||
import org.springframework.beans.factory.support.BeanDefinitionBuilder;
|
||||
import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser;
|
||||
import org.springframework.beans.factory.xml.ParserContext;
|
||||
import org.springframework.data.gemfire.GemfireUtils;
|
||||
import org.springframework.data.gemfire.util.SpringUtils;
|
||||
import org.springframework.data.gemfire.wan.AsyncEventQueueFactoryBean;
|
||||
import org.springframework.util.StringUtils;
|
||||
@@ -50,34 +49,50 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser {
|
||||
*/
|
||||
@Override
|
||||
protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) {
|
||||
|
||||
builder.setLazyInit(false);
|
||||
|
||||
parseAsyncEventListener(element, parserContext, builder);
|
||||
parseCache(element, builder);
|
||||
parseDiskStore(element, builder);
|
||||
|
||||
ParsingUtils.setPropertyValue(element, builder, "enable-batch-conflation", "batchConflationEnabled");
|
||||
ParsingUtils.setPropertyValue(element, builder, "batch-conflation-enabled");
|
||||
ParsingUtils.setPropertyValue(element, builder, "batch-size");
|
||||
ParsingUtils.setPropertyValue(element, builder, "batch-time-interval");
|
||||
ParsingUtils.setPropertyValue(element, builder, "disk-synchronous");
|
||||
ParsingUtils.setPropertyValue(element, builder, "dispatcher-threads");
|
||||
ParsingUtils.setPropertyValue(element, builder, "forward-expiration-destroy");
|
||||
ParsingUtils.setPropertyValue(element, builder, "maximum-queue-memory");
|
||||
ParsingUtils.setPropertyValue(element, builder, "order-policy");
|
||||
ParsingUtils.setPropertyValue(element, builder, "parallel");
|
||||
ParsingUtils.setPropertyValue(element, builder, "persistent");
|
||||
|
||||
if (GemfireUtils.GEMFIRE_VERSION.compareTo("7.0.1") >= 0) {
|
||||
ParsingUtils.setPropertyValue(element, builder, "enable-batch-conflation", "batchConflationEnabled");
|
||||
ParsingUtils.setPropertyValue(element, builder, "batch-conflation-enabled");
|
||||
ParsingUtils.setPropertyValue(element, builder, "batch-time-interval");
|
||||
ParsingUtils.setPropertyValue(element, builder, "disk-synchronous");
|
||||
ParsingUtils.setPropertyValue(element, builder, "dispatcher-threads");
|
||||
ParsingUtils.setPropertyValue(element, builder, "order-policy");
|
||||
Element eventFilterElement = DomUtils.getChildElementByTagName(element, "event-filter");
|
||||
|
||||
if (eventFilterElement != null) {
|
||||
builder.addPropertyValue("gatewayEventFilters",
|
||||
ParsingUtils.parseRefOrNestedBeanDeclaration(eventFilterElement, parserContext, builder));
|
||||
}
|
||||
|
||||
Element eventSubstitutionFilterElement =
|
||||
DomUtils.getChildElementByTagName(element, "event-substitution-filter");
|
||||
|
||||
if (eventSubstitutionFilterElement != null) {
|
||||
builder.addPropertyValue("gatewayEventSubstitutionFilter",
|
||||
ParsingUtils.parseRefOrSingleNestedBeanDeclaration(eventSubstitutionFilterElement, parserContext, builder));
|
||||
}
|
||||
|
||||
ParsingUtils.setPropertyValue(element, builder, NAME_ATTRIBUTE);
|
||||
|
||||
if (!StringUtils.hasText(element.getAttribute(NAME_ATTRIBUTE))) {
|
||||
if (element.getParentNode().getNodeName().endsWith("region")) {
|
||||
|
||||
Element region = (Element) element.getParentNode();
|
||||
|
||||
String regionName = StringUtils.hasText(region.getAttribute(NAME_ATTRIBUTE))
|
||||
? region.getAttribute(NAME_ATTRIBUTE) : region.getAttribute(ID_ATTRIBUTE);
|
||||
? region.getAttribute(NAME_ATTRIBUTE)
|
||||
: region.getAttribute(ID_ATTRIBUTE);
|
||||
|
||||
int index = 0;
|
||||
|
||||
@@ -92,13 +107,12 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser {
|
||||
}
|
||||
}
|
||||
|
||||
/* (non-Javadoc) */
|
||||
private void parseAsyncEventListener(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) {
|
||||
|
||||
Element asyncEventListenerElement = DomUtils.getChildElementByTagName(element, "async-event-listener");
|
||||
|
||||
Object asyncEventListener = ParsingUtils.parseRefOrSingleNestedBeanDeclaration(asyncEventListenerElement,
|
||||
parserContext,
|
||||
builder);
|
||||
Object asyncEventListener =
|
||||
ParsingUtils.parseRefOrSingleNestedBeanDeclaration(asyncEventListenerElement, parserContext, builder);
|
||||
|
||||
builder.addPropertyValue("asyncEventListener", asyncEventListener);
|
||||
|
||||
@@ -109,15 +123,16 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser {
|
||||
|
||||
/* (non-Javadoc) */
|
||||
private void parseCache(Element element, BeanDefinitionBuilder builder) {
|
||||
|
||||
String cacheRefAttribute = element.getAttribute("cache-ref");
|
||||
String cacheName = SpringUtils.defaultIfEmpty(cacheRefAttribute,
|
||||
GemfireConstants.DEFAULT_GEMFIRE_CACHE_NAME);
|
||||
String cacheName = SpringUtils.defaultIfEmpty(cacheRefAttribute, GemfireConstants.DEFAULT_GEMFIRE_CACHE_NAME);
|
||||
|
||||
builder.addConstructorArgReference(cacheName);
|
||||
}
|
||||
|
||||
/* (non-Javadoc) */
|
||||
private void parseDiskStore(Element element, BeanDefinitionBuilder builder) {
|
||||
|
||||
ParsingUtils.setPropertyValue(element, builder, "disk-store-ref");
|
||||
|
||||
String diskStoreRef = element.getAttribute("disk-store-ref");
|
||||
|
||||
@@ -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 { }
|
||||
|
||||
}
|
||||
|
||||
@@ -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> {
|
||||
@@ -38,6 +53,7 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
|
||||
private Boolean batchConflationEnabled;
|
||||
private Boolean diskSynchronous;
|
||||
private Boolean forwardExpirationDestroy;
|
||||
private Boolean parallel;
|
||||
private Boolean persistent;
|
||||
|
||||
@@ -46,8 +62,13 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
private Integer dispatcherThreads;
|
||||
private Integer maximumQueueMemory;
|
||||
|
||||
private GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter;
|
||||
|
||||
private GatewaySender.OrderPolicy orderPolicy;
|
||||
|
||||
private List<GatewayEventFilter> gatewayEventFilters;
|
||||
|
||||
private String diskStoreReference;
|
||||
private String orderPolicy;
|
||||
|
||||
/**
|
||||
* Constructs an instance of the AsyncEventQueueFactoryBean for creating an GemFire AsyncEventQueue.
|
||||
@@ -65,77 +86,60 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
* @param cache the GemFire Cache reference.
|
||||
* @param asyncEventListener required {@link AsyncEventListener}
|
||||
*/
|
||||
public AsyncEventQueueFactoryBean(final Cache cache, final AsyncEventListener asyncEventListener) {
|
||||
public AsyncEventQueueFactoryBean(Cache cache, AsyncEventListener asyncEventListener) {
|
||||
|
||||
super(cache);
|
||||
|
||||
setAsyncEventListener(asyncEventListener);
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueue getObject() throws Exception {
|
||||
return asyncEventQueue;
|
||||
return this.asyncEventQueue;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Class<?> getObjectType() {
|
||||
return AsyncEventQueue.class;
|
||||
return this.asyncEventQueue != null ? this.asyncEventQueue.getClass() : AsyncEventQueue.class;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void doInit() {
|
||||
Assert.notNull(this.asyncEventListener, "The AsyncEventListener cannot be null.");
|
||||
|
||||
AsyncEventQueueFactory asyncEventQueueFactory = (this.factory != null ? (AsyncEventQueueFactory) factory
|
||||
: cache.createAsyncEventQueueFactory());
|
||||
Assert.state(this.asyncEventListener != null, "AsyncEventListener must not be null");
|
||||
|
||||
if (batchSize != null) {
|
||||
asyncEventQueueFactory.setBatchSize(batchSize);
|
||||
}
|
||||
AsyncEventQueueFactory asyncEventQueueFactory =
|
||||
this.factory != null ? (AsyncEventQueueFactory) this.factory : this.cache.createAsyncEventQueueFactory();
|
||||
|
||||
if (batchTimeInterval != null) {
|
||||
asyncEventQueueFactory.setBatchTimeInterval(batchTimeInterval);
|
||||
}
|
||||
|
||||
if (batchConflationEnabled != null) {
|
||||
asyncEventQueueFactory.setBatchConflationEnabled(batchConflationEnabled);
|
||||
}
|
||||
|
||||
if (dispatcherThreads != null) {
|
||||
asyncEventQueueFactory.setDispatcherThreads(dispatcherThreads);
|
||||
}
|
||||
|
||||
if (diskStoreReference != null) {
|
||||
asyncEventQueueFactory.setDiskStoreName(diskStoreReference);
|
||||
}
|
||||
|
||||
if (diskSynchronous != null) {
|
||||
asyncEventQueueFactory.setDiskSynchronous(diskSynchronous);
|
||||
}
|
||||
|
||||
if (maximumQueueMemory != null) {
|
||||
asyncEventQueueFactory.setMaximumQueueMemory(maximumQueueMemory);
|
||||
}
|
||||
Optional.ofNullable(this.batchConflationEnabled).ifPresent(asyncEventQueueFactory::setBatchConflationEnabled);
|
||||
Optional.ofNullable(this.batchSize).ifPresent(asyncEventQueueFactory::setBatchSize);
|
||||
Optional.ofNullable(this.batchTimeInterval).ifPresent(asyncEventQueueFactory::setBatchTimeInterval);
|
||||
Optional.ofNullable(this.diskStoreReference).ifPresent(asyncEventQueueFactory::setDiskStoreName);
|
||||
Optional.ofNullable(this.diskSynchronous).ifPresent(asyncEventQueueFactory::setDiskSynchronous);
|
||||
Optional.ofNullable(this.dispatcherThreads).ifPresent(asyncEventQueueFactory::setDispatcherThreads);
|
||||
Optional.ofNullable(this.forwardExpirationDestroy).ifPresent(asyncEventQueueFactory::setForwardExpirationDestroy);
|
||||
Optional.ofNullable(this.gatewayEventSubstitutionFilter).ifPresent(asyncEventQueueFactory::setGatewayEventSubstitutionListener);
|
||||
Optional.ofNullable(this.maximumQueueMemory).ifPresent(asyncEventQueueFactory::setMaximumQueueMemory);
|
||||
Optional.ofNullable(this.persistent).ifPresent(asyncEventQueueFactory::setPersistent);
|
||||
|
||||
asyncEventQueueFactory.setParallel(isParallelEventQueue());
|
||||
|
||||
if (orderPolicy != null) {
|
||||
Assert.isTrue(isSerialEventQueue(), "Order Policy cannot be used with a Parallel Event Queue.");
|
||||
nullSafeList(this.gatewayEventFilters).forEach(asyncEventQueueFactory::addGatewayEventFilter);
|
||||
|
||||
Assert.isTrue(VALID_ORDER_POLICIES.contains(orderPolicy.toUpperCase()), String.format(
|
||||
"The value of Order Policy '$1%s' is invalid.", orderPolicy));
|
||||
if (this.orderPolicy != null) {
|
||||
|
||||
asyncEventQueueFactory.setOrderPolicy(GatewaySender.OrderPolicy.valueOf(orderPolicy.toUpperCase()));
|
||||
Assert.state(isSerialEventQueue(), "OrderPolicy cannot be used with a Parallel AsyncEventQueue");
|
||||
|
||||
asyncEventQueueFactory.setOrderPolicy(this.orderPolicy);
|
||||
}
|
||||
|
||||
if (persistent != null) {
|
||||
asyncEventQueueFactory.setPersistent(persistent);
|
||||
}
|
||||
|
||||
asyncEventQueue = asyncEventQueueFactory.create(getName(), this.asyncEventListener);
|
||||
setAsyncEventQueue(asyncEventQueueFactory.create(getName(), this.asyncEventListener));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void destroy() throws Exception {
|
||||
if (!cache.isClosed()) {
|
||||
|
||||
if (!this.cache.isClosed()) {
|
||||
try {
|
||||
this.asyncEventListener.close();
|
||||
}
|
||||
@@ -145,22 +149,28 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
}
|
||||
|
||||
public final void setAsyncEventListener(AsyncEventListener listener) {
|
||||
|
||||
Assert.state(this.asyncEventQueue == null,
|
||||
"Setting an AsyncEventListener is not allowed once the AsyncEventQueue has been created.");
|
||||
"Setting an AsyncEventListener is not allowed once the AsyncEventQueue has been created");
|
||||
|
||||
this.asyncEventListener = listener;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param asyncEventQueue overrides Async Event Queue returned by this FactoryBean.
|
||||
* Configures the {@link AsyncEventQueue} returned by this {@link FactoryBean}.
|
||||
*
|
||||
* @param asyncEventQueue overrides {@link AsyncEventQueue} returned by this {@link FactoryBean}.
|
||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueue
|
||||
*/
|
||||
public void setAsyncEventQueue(AsyncEventQueue asyncEventQueue) {
|
||||
this.asyncEventQueue = asyncEventQueue;
|
||||
}
|
||||
|
||||
/**
|
||||
* Enable or disable the Async Event Queue's (AEQ) should conflate messages.
|
||||
* Enable or disable {@link AsyncEventQueue} (AEQ) message conflation.
|
||||
*
|
||||
* @param batchConflationEnabled a boolean value indicating whether to conflate queued events.
|
||||
* @param batchConflationEnabled {@link Boolean} indicating whether to conflate queued events.
|
||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setBatchConflationEnabled(boolean)
|
||||
*/
|
||||
public void setBatchConflationEnabled(Boolean batchConflationEnabled) {
|
||||
this.batchConflationEnabled = batchConflationEnabled;
|
||||
@@ -171,10 +181,11 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
}
|
||||
|
||||
/**
|
||||
* Set the Aysync Event Queue's (AEQ) interval between sending batches.
|
||||
* Configures the {@link AsyncEventQueue} (AEQ) interval between sending batches.
|
||||
*
|
||||
* @param batchTimeInterval an integer value indicating the maximum number of milliseconds that can elapse
|
||||
* between sending batches.
|
||||
* @param batchTimeInterval {@link Integer} specifying the maximum number of milliseconds
|
||||
* that can elapse between sending batches.
|
||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setBatchTimeInterval(int)
|
||||
*/
|
||||
public void setBatchTimeInterval(Integer batchTimeInterval) {
|
||||
this.batchTimeInterval = batchTimeInterval;
|
||||
@@ -185,36 +196,69 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
}
|
||||
|
||||
/**
|
||||
* Set the Async Event Queue (AEQ) disk write synchronization policy.
|
||||
* Configures the {@link AsyncEventQueue} (AEQ) disk write synchronization policy.
|
||||
*
|
||||
* @param diskSynchronous a boolean value indicating whether disk writes are synchronous.
|
||||
* @param diskSynchronous boolean value indicating whether disk writes are synchronous.
|
||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setDiskSynchronous(boolean)
|
||||
*/
|
||||
public void setDiskSynchronous(Boolean diskSynchronous) {
|
||||
this.diskSynchronous = diskSynchronous;
|
||||
}
|
||||
|
||||
/**
|
||||
* Set the number of dispatcher threads used to process Region Events from the associated Async Event Queue (AEQ).
|
||||
* Configures the number of dispatcher threads used to process Region Events
|
||||
* from the associated {@link AsyncEventQueue} (AEQ).
|
||||
*
|
||||
* @param dispatcherThreads an Integer indicating the number of dispatcher threads used to process Region Events
|
||||
* from the associated Queue.
|
||||
* @param dispatcherThreads {@link Integer} specifying the number of dispatcher threads used
|
||||
* to process {@link Region} events from the associated queue.
|
||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setDispatcherThreads(int)
|
||||
*/
|
||||
public void setDispatcherThreads(Integer dispatcherThreads) {
|
||||
this.dispatcherThreads = dispatcherThreads;
|
||||
}
|
||||
|
||||
/**
|
||||
* Forwards expiration (action-based) destroy events to the {@link AsyncEventQueue} (AEQ).
|
||||
*
|
||||
* By default, destroy events are not added to the AEQ. Setting this attribute to {@literal true}
|
||||
* will add all expiration destroy events to the AEQ.
|
||||
*
|
||||
* @param forwardExpirationDestroy boolean value indicating whether to forward expiration destroy events.
|
||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory#setForwardExpirationDestroy(boolean)
|
||||
* @see org.apache.geode.cache.ExpirationAttributes#getAction()
|
||||
* @see org.apache.geode.cache.ExpirationAction#DESTROY
|
||||
*/
|
||||
public void setForwardExpirationDestroy(Boolean forwardExpirationDestroy) {
|
||||
this.forwardExpirationDestroy = forwardExpirationDestroy;
|
||||
}
|
||||
|
||||
public void setGatewayEventFilters(List<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;
|
||||
}
|
||||
|
||||
@@ -233,5 +277,4 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
public void setPersistent(Boolean persistent) {
|
||||
this.persistent = persistent;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -46,12 +46,6 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
|
||||
private int remoteDistributedSystemId;
|
||||
|
||||
private GatewaySender gatewaySender;
|
||||
|
||||
private List<GatewayEventFilter> eventFilters;
|
||||
|
||||
private List<GatewayTransportFilter> transportFilters;
|
||||
|
||||
private Boolean diskSynchronous;
|
||||
private Boolean batchConflationEnabled;
|
||||
private Boolean parallel;
|
||||
@@ -59,6 +53,10 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
|
||||
private GatewayEventSubstitutionFilter eventSubstitutionFilter;
|
||||
|
||||
private GatewaySender gatewaySender;
|
||||
|
||||
private GatewaySender.OrderPolicy orderPolicy;
|
||||
|
||||
private Integer alertThreshold;
|
||||
private Integer batchSize;
|
||||
private Integer batchTimeInterval;
|
||||
@@ -67,8 +65,11 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
private Integer socketBufferSize;
|
||||
private Integer socketReadTimeout;
|
||||
|
||||
private List<GatewayEventFilter> eventFilters;
|
||||
|
||||
private List<GatewayTransportFilter> transportFilters;
|
||||
|
||||
private String diskStoreReference;
|
||||
private String orderPolicy;
|
||||
|
||||
/**
|
||||
* Constructs an instance of the {@link GatewaySenderFactoryBean} class initialized with a reference to
|
||||
@@ -86,8 +87,9 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
*/
|
||||
@Override
|
||||
protected void doInit() {
|
||||
GatewaySenderFactory gatewaySenderFactory = (this.factory != null ? (GatewaySenderFactory) factory
|
||||
: cache.createGatewaySenderFactory());
|
||||
|
||||
GatewaySenderFactory gatewaySenderFactory =
|
||||
this.factory != null ? (GatewaySenderFactory) this.factory : this.cache.createGatewaySenderFactory();
|
||||
|
||||
if (alertThreshold != null) {
|
||||
gatewaySenderFactory.setAlertThreshold(alertThreshold);
|
||||
@@ -132,12 +134,10 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
}
|
||||
|
||||
if (orderPolicy != null) {
|
||||
Assert.isTrue(isSerialGatewaySender(), "Order Policy cannot be used with a Parallel Gateway Sender Queue.");
|
||||
|
||||
Assert.isTrue(VALID_ORDER_POLICIES.contains(orderPolicy.toUpperCase()),
|
||||
String.format("The value for Order Policy '%s' is invalid.", orderPolicy));
|
||||
Assert.isTrue(isSerialGatewaySender(), "OrderPolicy cannot be used with a Parallel GatewaySender");
|
||||
|
||||
gatewaySenderFactory.setOrderPolicy(GatewaySender.OrderPolicy.valueOf(orderPolicy.toUpperCase()));
|
||||
gatewaySenderFactory.setOrderPolicy(this.orderPolicy);
|
||||
}
|
||||
|
||||
gatewaySenderFactory.setParallel(isParallelGatewaySender());
|
||||
@@ -162,20 +162,14 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
gatewaySender = wrapper;
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
public GatewaySender getObject() throws Exception {
|
||||
return gatewaySender;
|
||||
return this.gatewaySender;
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
public Class<?> getObjectType() {
|
||||
return (gatewaySender != null ? gatewaySender.getClass() : GatewaySender.class);
|
||||
return this.gatewaySender != null ? this.gatewaySender.getClass() : GatewaySender.class;
|
||||
}
|
||||
|
||||
public void setAlertThreshold(Integer alertThreshold) {
|
||||
@@ -235,6 +229,10 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
}
|
||||
|
||||
public void setOrderPolicy(String orderPolicy) {
|
||||
setOrderPolicy(GatewaySender.OrderPolicy.valueOf(String.valueOf(orderPolicy).toUpperCase()));
|
||||
}
|
||||
|
||||
public void setOrderPolicy(GatewaySender.OrderPolicy orderPolicy) {
|
||||
this.orderPolicy = orderPolicy;
|
||||
}
|
||||
|
||||
|
||||
@@ -2952,8 +2952,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
|
||||
@@ -2971,6 +2970,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>
|
||||
@@ -2980,6 +2981,7 @@ if an inner bean.
|
||||
]]></xsd:documentation>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
<xsd:attribute name="forward-expiration-destroy" type="xsd:string" default="false" use="optional"/>
|
||||
<xsd:attributeGroup ref="commonWANQueueAttributes" />
|
||||
</xsd:complexType>
|
||||
<!-- -->
|
||||
|
||||
@@ -0,0 +1,194 @@
|
||||
/*
|
||||
* Copyright 2012-2018 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*
|
||||
*/
|
||||
|
||||
package org.springframework.data.gemfire.config.xml;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.hamcrest.Matchers.equalTo;
|
||||
import static org.hamcrest.Matchers.is;
|
||||
import static org.hamcrest.Matchers.notNullValue;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
import javax.annotation.Resource;
|
||||
|
||||
import org.apache.geode.cache.EntryEvent;
|
||||
import org.apache.geode.cache.asyncqueue.AsyncEvent;
|
||||
import org.apache.geode.cache.asyncqueue.AsyncEventListener;
|
||||
import org.apache.geode.cache.asyncqueue.AsyncEventQueue;
|
||||
import org.apache.geode.cache.wan.GatewayEventFilter;
|
||||
import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter;
|
||||
import org.apache.geode.cache.wan.GatewayQueueEvent;
|
||||
import org.apache.geode.cache.wan.GatewaySender;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit4.SpringRunner;
|
||||
|
||||
/**
|
||||
* The AsyncEventQueueNamespaceTest class is a test suite of test cases testing the contract and functionality
|
||||
* of configuring a Pivotal GemFire or Apache Geode {@link AsyncEventQueue} using the SDG XML namespace.
|
||||
*
|
||||
* @author John Blum
|
||||
* @see org.junit.Test
|
||||
* @see org.junit.runner.RunWith
|
||||
* @see org.springframework.data.gemfire.config.AsyncEventQueueParser
|
||||
* @see org.springframework.data.gemfire.wan.AsyncEventQueueFactoryBean
|
||||
* @see org.springframework.test.context.ContextConfiguration
|
||||
* @see org.springframework.test.context.junit4.SpringJUnit4ClassRunner
|
||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueue
|
||||
* @since 1.0.0
|
||||
*/
|
||||
@RunWith(SpringRunner.class)
|
||||
@ContextConfiguration
|
||||
@SuppressWarnings("all")
|
||||
public class AsyncEventQueueNamespaceTest {
|
||||
|
||||
@Resource(name = "TestAsyncEventQueue")
|
||||
private AsyncEventQueue asyncEventQueue;
|
||||
|
||||
@Resource(name = "TestAsyncEventQueueWithFilters")
|
||||
private AsyncEventQueue asyncEventQueueWithFilters;
|
||||
|
||||
@Test
|
||||
public void asyncEventQueueIsConfiguredProperly() {
|
||||
|
||||
Assert.assertThat(asyncEventQueue, is(notNullValue(AsyncEventQueue.class)));
|
||||
Assert.assertThat(asyncEventQueue.getId(), is(equalTo("TestAsyncEventQueue")));
|
||||
Assert.assertThat(asyncEventQueue.isBatchConflationEnabled(), is(true));
|
||||
Assert.assertThat(asyncEventQueue.getBatchSize(), is(equalTo(100)));
|
||||
Assert.assertThat(asyncEventQueue.getBatchTimeInterval(), is(equalTo(30)));
|
||||
Assert.assertThat(asyncEventQueue.getDiskStoreName(), is(equalTo("TestDiskStore")));
|
||||
Assert.assertThat(asyncEventQueue.isDiskSynchronous(), is(true));
|
||||
Assert.assertThat(asyncEventQueue.getDispatcherThreads(), is(equalTo(4)));
|
||||
Assert.assertThat(asyncEventQueue.isForwardExpirationDestroy(), is(true));
|
||||
Assert.assertThat(asyncEventQueue.getMaximumQueueMemory(), is(equalTo(50)));
|
||||
Assert.assertThat(asyncEventQueue.getOrderPolicy(), is(equalTo(GatewaySender.OrderPolicy.KEY)));
|
||||
Assert.assertThat(asyncEventQueue.isParallel(), is(false));
|
||||
Assert.assertThat(asyncEventQueue.isPersistent(), is(true));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void asyncEventQueueListenerEqualsExpected() {
|
||||
|
||||
AsyncEventListener asyncEventListener = asyncEventQueue.getAsyncEventListener();
|
||||
|
||||
Assert.assertThat(asyncEventListener, is(notNullValue(AsyncEventListener.class)));
|
||||
Assert.assertThat(asyncEventListener.toString(), is(equalTo("TestAeqListener")));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void asyncEventQueueWithFiltersIsConfiguredProperly() {
|
||||
|
||||
assertThat(asyncEventQueueWithFilters).isNotNull();
|
||||
assertThat(asyncEventQueueWithFilters.getId()).isEqualTo("TestAsyncEventQueueWithFilters");
|
||||
|
||||
AsyncEventListener listener = asyncEventQueueWithFilters.getAsyncEventListener();
|
||||
|
||||
assertThat(listener).isNotNull();
|
||||
assertThat(listener.toString()).isEqualTo("TestListenerOne");
|
||||
|
||||
List<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 {
|
||||
|
||||
private final String name;
|
||||
|
||||
public TestAsyncEventListener(String name) {
|
||||
this.name = name;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean processEvents(List<AsyncEvent> events) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return this.name;
|
||||
}
|
||||
}
|
||||
|
||||
public static class TestGatewayEventFilter implements GatewayEventFilter {
|
||||
|
||||
private final String name;
|
||||
|
||||
public TestGatewayEventFilter(String name) {
|
||||
this.name = name;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean beforeEnqueue(GatewayQueueEvent event) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean beforeTransmit(GatewayQueueEvent event) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void afterAcknowledgement(GatewayQueueEvent event) { }
|
||||
|
||||
@Override
|
||||
public void close() { }
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return this.name;
|
||||
}
|
||||
}
|
||||
|
||||
public static class TestGatewayEventSubstitutionFilter implements GatewayEventSubstitutionFilter<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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -49,14 +49,15 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
|
||||
|
||||
private GatewayEventSubstitutionFilter<?, ?> gatewayEventSubstitutionFilter;
|
||||
|
||||
private OrderPolicy orderPolicy;
|
||||
|
||||
private List<GatewayEventFilter> gatewayEventFilters = new ArrayList<>();
|
||||
|
||||
private OrderPolicy orderPolicy;
|
||||
|
||||
private String diskStoreName;
|
||||
|
||||
@Override
|
||||
public AsyncEventQueue create(final String name, final AsyncEventListener listener) {
|
||||
public AsyncEventQueue create(String name, AsyncEventListener listener) {
|
||||
|
||||
when(asyncEventQueue.getAsyncEventListener()).thenReturn(listener);
|
||||
when(asyncEventQueue.getBatchSize()).thenReturn(this.batchSize);
|
||||
when(asyncEventQueue.getDiskStoreName()).thenReturn(this.diskStoreName);
|
||||
@@ -76,47 +77,30 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
|
||||
return this.asyncEventQueue;
|
||||
}
|
||||
|
||||
//The following added in 7.0.1
|
||||
public AsyncEventQueueFactory setBatchConflationEnabled(boolean arg0) {
|
||||
this.batchConflationEnabled = arg0;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setBatchSize(int batchSize) {
|
||||
this.batchSize = batchSize;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setBatchTimeInterval(int interval) {
|
||||
this.batchTimeInterval = interval;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setDiskStoreName(String diskStoreName) {
|
||||
this.diskStoreName = diskStoreName;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setMaximumQueueMemory(int maxQueueMemory) {
|
||||
this.maxQueueMemory = maxQueueMemory;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setPersistent(boolean persistent) {
|
||||
this.persistent = persistent;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setParallel(boolean parallel) {
|
||||
this.parallel = parallel;
|
||||
return this;
|
||||
}
|
||||
|
||||
//The following added in 7.0.1
|
||||
public AsyncEventQueueFactory setBatchConflationEnabled(boolean arg0) {
|
||||
this.batchConflationEnabled = arg0;
|
||||
return this;
|
||||
}
|
||||
|
||||
public AsyncEventQueueFactory setBatchTimeInterval(int arg0) {
|
||||
this.batchTimeInterval = arg0;
|
||||
return this;
|
||||
}
|
||||
|
||||
public AsyncEventQueueFactory setDiskSynchronous(boolean arg0) {
|
||||
this.diskSynchronous = arg0;
|
||||
return this;
|
||||
@@ -127,11 +111,39 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setForwardExpirationDestroy(boolean forward) {
|
||||
this.forwardExpirationDestroy = forward;
|
||||
return this;
|
||||
}
|
||||
|
||||
public AsyncEventQueueFactory setGatewayEventSubstitutionListener(final GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter) {
|
||||
this.gatewayEventSubstitutionFilter = gatewayEventSubstitutionFilter;
|
||||
return this;
|
||||
}
|
||||
|
||||
public AsyncEventQueueFactory setMaximumQueueMemory(int maxQueueMemory) {
|
||||
this.maxQueueMemory = maxQueueMemory;
|
||||
return this;
|
||||
}
|
||||
|
||||
public AsyncEventQueueFactory setOrderPolicy(OrderPolicy arg0) {
|
||||
this.orderPolicy = arg0;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setParallel(boolean parallel) {
|
||||
this.parallel = parallel;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setPersistent(boolean persistent) {
|
||||
this.persistent = persistent;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory addGatewayEventFilter(final GatewayEventFilter gatewayEventFilter) {
|
||||
gatewayEventFilters.add(gatewayEventFilter);
|
||||
@@ -143,16 +155,4 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
|
||||
gatewayEventFilters.remove(gatewayEventFilter);
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setGatewayEventSubstitutionListener(final GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter) {
|
||||
this.gatewayEventSubstitutionFilter = gatewayEventSubstitutionFilter;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncEventQueueFactory setForwardExpirationDestroy(boolean forward) {
|
||||
this.forwardExpirationDestroy = forward;
|
||||
return this;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -16,21 +16,27 @@
|
||||
|
||||
package org.springframework.data.gemfire.wan;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.junit.Assert.assertSame;
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.anyBoolean;
|
||||
import static org.mockito.ArgumentMatchers.anyInt;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
import static org.mockito.ArgumentMatchers.isA;
|
||||
import static org.mockito.Matchers.eq;
|
||||
import static org.mockito.Matchers.notNull;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.never;
|
||||
import static org.mockito.Mockito.times;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
import java.util.Arrays;
|
||||
|
||||
import org.apache.geode.cache.Cache;
|
||||
import org.apache.geode.cache.asyncqueue.AsyncEventListener;
|
||||
import org.apache.geode.cache.asyncqueue.AsyncEventQueue;
|
||||
import org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory;
|
||||
import org.apache.geode.cache.wan.GatewayEventFilter;
|
||||
import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter;
|
||||
import org.apache.geode.cache.wan.GatewaySender;
|
||||
import org.junit.Test;
|
||||
import org.springframework.data.gemfire.TestUtils;
|
||||
@@ -52,277 +58,409 @@ import org.springframework.data.gemfire.TestUtils;
|
||||
*/
|
||||
public class AsyncEventQueueFactoryBeanTest {
|
||||
|
||||
protected Cache createMockCacheWithAsyncEventQueueInfrastructure(
|
||||
final AsyncEventQueueFactory mockAsynEventQueueFactory) {
|
||||
Cache mockCache = mock(Cache.class);
|
||||
when((mockCache.createAsyncEventQueueFactory())).thenReturn(mockAsynEventQueueFactory);
|
||||
private Cache mockCache() {
|
||||
return mock(Cache.class);
|
||||
}
|
||||
|
||||
private Cache mockCache(AsyncEventQueueFactory mockAsyncEventQueueFactory) {
|
||||
|
||||
Cache mockCache = mockCache();
|
||||
|
||||
when((mockCache.createAsyncEventQueueFactory())).thenReturn(mockAsyncEventQueueFactory);
|
||||
|
||||
return mockCache;
|
||||
}
|
||||
|
||||
protected AsyncEventQueueFactory createMockAsyncEventQueueFactory(final String asyncEventQueueId) {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory = mock(AsyncEventQueueFactory.class);
|
||||
AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class);
|
||||
private AsyncEventQueueFactory mockAsyncEventQueueFactory(String asyncEventQueueId) {
|
||||
|
||||
when(mockAsyncEventQueue.getId()).thenReturn(asyncEventQueueId);
|
||||
when(mockAsyncEventQueueFactory.create(eq(asyncEventQueueId), notNull(AsyncEventListener.class)))
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory = mock(AsyncEventQueueFactory.class);
|
||||
|
||||
AsyncEventQueue mockAsyncEventQueue = mockAsyncEventQueue(asyncEventQueueId);
|
||||
|
||||
when(mockAsyncEventQueueFactory.create(eq(asyncEventQueueId), isA(AsyncEventListener.class)))
|
||||
.thenReturn(mockAsyncEventQueue);
|
||||
|
||||
return mockAsyncEventQueueFactory;
|
||||
}
|
||||
|
||||
protected AsyncEventListener createMockAsyncEventListener() {
|
||||
private AsyncEventQueue mockAsyncEventQueue(String asyncEventQueueId) {
|
||||
|
||||
AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class);
|
||||
|
||||
when(mockAsyncEventQueue.getId()).thenReturn(asyncEventQueueId);
|
||||
|
||||
return mockAsyncEventQueue;
|
||||
}
|
||||
|
||||
private AsyncEventListener mockAsyncEventListener() {
|
||||
return mock(AsyncEventListener.class);
|
||||
}
|
||||
|
||||
protected void verifyExpectations(final AsyncEventQueueFactory mockAsyncEventQueueFactory,
|
||||
final AsyncEventQueueFactoryBean factoryBean) throws Exception {
|
||||
Boolean parallel = TestUtils.readField("parallel", factoryBean);
|
||||
|
||||
verify(mockAsyncEventQueueFactory).setParallel(eq(Boolean.TRUE.equals(parallel)));
|
||||
|
||||
String orderPolicy = TestUtils.readField("orderPolicy", factoryBean);
|
||||
|
||||
if (orderPolicy != null) {
|
||||
verify(mockAsyncEventQueueFactory).setOrderPolicy(
|
||||
eq(GatewaySender.OrderPolicy.valueOf(orderPolicy.toUpperCase())));
|
||||
}
|
||||
|
||||
Integer dispatcherThreads = TestUtils.readField("dispatcherThreads", factoryBean);
|
||||
|
||||
if (dispatcherThreads != null) {
|
||||
verify(mockAsyncEventQueueFactory).setDispatcherThreads(eq(dispatcherThreads));
|
||||
}
|
||||
|
||||
String diskStoreReference = TestUtils.readField("diskStoreReference", factoryBean);
|
||||
|
||||
if (diskStoreReference != null) {
|
||||
verify(mockAsyncEventQueueFactory).setDiskStoreName(eq(diskStoreReference));
|
||||
}
|
||||
|
||||
Boolean diskSynchronous = TestUtils.readField("diskSynchronous", factoryBean);
|
||||
|
||||
if (diskSynchronous != null) {
|
||||
verify(mockAsyncEventQueueFactory).setDiskSynchronous(eq(diskSynchronous));
|
||||
}
|
||||
|
||||
Boolean persistent = TestUtils.readField("persistent", factoryBean);
|
||||
|
||||
if (persistent != null) {
|
||||
verify(mockAsyncEventQueueFactory).setPersistent(eq(persistent));
|
||||
}
|
||||
else {
|
||||
verify(mockAsyncEventQueueFactory, never()).setPersistent(true);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSetAsyncEventListener() throws Exception {
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(createMockAsyncEventQueueFactory("testEventQueue")));
|
||||
public void setAndGetAsyncEventListener() throws Exception {
|
||||
|
||||
AsyncEventListener listenerOne = createMockAsyncEventListener();
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(mockCache());
|
||||
|
||||
AsyncEventListener listenerOne = mockAsyncEventListener();
|
||||
|
||||
factoryBean.setAsyncEventListener(listenerOne);
|
||||
|
||||
assertSame(listenerOne, TestUtils.readField("asyncEventListener", factoryBean));
|
||||
assertThat(TestUtils.<AsyncEventListener>readField("asyncEventListener", factoryBean))
|
||||
.isSameAs(listenerOne);
|
||||
|
||||
AsyncEventListener listenerTwo = createMockAsyncEventListener();
|
||||
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 = createMockAsyncEventListener();
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(mockCache(), mockAsyncEventListener);
|
||||
|
||||
factoryBean.setAsyncEventListener(listenerOne);
|
||||
|
||||
assertSame(listenerOne, TestUtils.readField("asyncEventListener", factoryBean));
|
||||
|
||||
factoryBean.doInit();
|
||||
|
||||
assertNotNull(TestUtils.readField("asyncEventQueue", factoryBean));
|
||||
factoryBean.setAsyncEventQueue(mockAsyncEventQueue);
|
||||
|
||||
try {
|
||||
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
|
||||
factoryBean.setAsyncEventListener(mockAsyncEventListener());
|
||||
}
|
||||
catch (IllegalStateException expected) {
|
||||
assertEquals("Setting an AsyncEventListener is not allowed once the AsyncEventQueue has been created.",
|
||||
expected.getMessage());
|
||||
assertSame(listenerOne, TestUtils.readField("asyncEventListener", factoryBean));
|
||||
|
||||
assertThat(expected)
|
||||
.hasMessage("Setting an AsyncEventListener is not allowed once the AsyncEventQueue has been created");
|
||||
|
||||
assertThat(expected).hasNoCause();
|
||||
|
||||
throw expected;
|
||||
}
|
||||
}
|
||||
|
||||
@Test(expected = IllegalArgumentException.class)
|
||||
public void testDoInitWhenAsyncEventListenerIsNull() throws Exception {
|
||||
try {
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(createMockAsyncEventQueueFactory("testEventQueue")));
|
||||
|
||||
assertNull(TestUtils.readField("asyncEventListener", factoryBean));
|
||||
|
||||
factoryBean.doInit();
|
||||
}
|
||||
catch (Exception e) {
|
||||
assertEquals("The AsyncEventListener cannot be null.", e.getMessage());
|
||||
throw e;
|
||||
finally {
|
||||
assertThat(TestUtils.<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(createMockAsyncEventListener());
|
||||
factoryBean.setDispatcherThreads(8);
|
||||
factoryBean.setParallel(true);
|
||||
factoryBean.doInit();
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory = mockAsyncEventQueueFactory("testQueue");
|
||||
|
||||
verifyExpectations(mockAsyncEventQueueFactory, factoryBean);
|
||||
AsyncEventQueueFactoryBean factoryBean =
|
||||
new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory), mockListener);
|
||||
|
||||
AsyncEventQueue eventQueue = factoryBean.getObject();
|
||||
GatewayEventFilter mockGatewayEventFilterOne = mock(GatewayEventFilter.class);
|
||||
GatewayEventFilter mockGatewayEventFilterTwo = mock(GatewayEventFilter.class);
|
||||
|
||||
assertNotNull(eventQueue);
|
||||
assertEquals("000", eventQueue.getId());
|
||||
}
|
||||
GatewayEventSubstitutionFilter mockGatewayEventSubstitutionFilter = mock(GatewayEventSubstitutionFilter.class);
|
||||
|
||||
@Test
|
||||
public void testParallelAsyncEventQueue() throws Exception {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory = createMockAsyncEventQueueFactory("123");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory));
|
||||
|
||||
factoryBean.setName("123");
|
||||
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
|
||||
factoryBean.setParallel(true);
|
||||
factoryBean.doInit();
|
||||
|
||||
verifyExpectations(mockAsyncEventQueueFactory, factoryBean);
|
||||
|
||||
AsyncEventQueue eventQueue = factoryBean.getObject();
|
||||
|
||||
assertNotNull(eventQueue);
|
||||
assertEquals("123", eventQueue.getId());
|
||||
}
|
||||
|
||||
@Test(expected = IllegalArgumentException.class)
|
||||
public void testParallelAsyncEventQueueWithOrderPolicy() {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory = createMockAsyncEventQueueFactory("456");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory));
|
||||
|
||||
factoryBean.setName("456");
|
||||
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
|
||||
factoryBean.setOrderPolicy("Key");
|
||||
factoryBean.setParallel(true);
|
||||
|
||||
try {
|
||||
factoryBean.doInit();
|
||||
}
|
||||
catch (IllegalArgumentException expected) {
|
||||
assertEquals("Order Policy cannot be used with a Parallel Event Queue.",
|
||||
expected.getMessage());
|
||||
throw expected;
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSerialAsyncEventQueueWithOrderPolicy() throws Exception {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFatory = createMockAsyncEventQueueFactory("789");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFatory));
|
||||
|
||||
factoryBean.setName("789");
|
||||
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
|
||||
factoryBean.setOrderPolicy("THREAD");
|
||||
factoryBean.setParallel(false);
|
||||
factoryBean.doInit();
|
||||
|
||||
verifyExpectations(mockAsyncEventQueueFatory, factoryBean);
|
||||
|
||||
AsyncEventQueue eventQueue = factoryBean.getObject();
|
||||
|
||||
assertNotNull(eventQueue);
|
||||
assertEquals("789", eventQueue.getId());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testAsyncEventQueueWithOrderPolicyAndDispatcherThreads() throws Exception {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFatory = createMockAsyncEventQueueFactory("abc");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFatory));
|
||||
|
||||
factoryBean.setName("abc");
|
||||
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
|
||||
factoryBean.setBatchConflationEnabled(true);
|
||||
factoryBean.setBatchSize(1024);
|
||||
factoryBean.setBatchTimeInterval(600);
|
||||
factoryBean.setDiskStoreRef("testDiskStore");
|
||||
factoryBean.setDiskSynchronous(false);
|
||||
factoryBean.setDispatcherThreads(2);
|
||||
factoryBean.setOrderPolicy("THREAD");
|
||||
factoryBean.doInit();
|
||||
|
||||
verifyExpectations(mockAsyncEventQueueFatory, factoryBean);
|
||||
|
||||
AsyncEventQueue eventQueue = factoryBean.getObject();
|
||||
|
||||
assertNotNull(eventQueue);
|
||||
assertEquals("abc", eventQueue.getId());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testAsyncEventQueueWithOverflowDiskStoreNoPersistence() throws Exception {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory = createMockAsyncEventQueueFactory("123abc");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory));
|
||||
|
||||
factoryBean.setName("123abc");
|
||||
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
|
||||
factoryBean.setDiskStoreRef("queueOverflowDiskStore");
|
||||
factoryBean.setForwardExpirationDestroy(true);
|
||||
factoryBean.setGatewayEventFilters(Arrays.asList(mockGatewayEventFilterOne, mockGatewayEventFilterTwo));
|
||||
factoryBean.setGatewayEventSubstitutionFilter(mockGatewayEventSubstitutionFilter);
|
||||
factoryBean.setMaximumQueueMemory(8192);
|
||||
factoryBean.setName("testQueue");
|
||||
factoryBean.setOrderPolicy(GatewaySender.OrderPolicy.PARTITION);
|
||||
factoryBean.setParallel(false);
|
||||
factoryBean.setPersistent(false);
|
||||
factoryBean.doInit();
|
||||
|
||||
verifyExpectations(mockAsyncEventQueueFactory, factoryBean);
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setBatchConflationEnabled(eq(true));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setBatchSize(eq(1024));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setBatchTimeInterval(eq(600));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setDiskStoreName(eq("testDiskStore"));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setDiskSynchronous(eq(false));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setDispatcherThreads(eq(2));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setForwardExpirationDestroy(eq(true));
|
||||
verify(mockAsyncEventQueueFactory, times(1))
|
||||
.setGatewayEventSubstitutionListener(eq(mockGatewayEventSubstitutionFilter));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setMaximumQueueMemory(eq(8192));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setOrderPolicy(eq(GatewaySender.OrderPolicy.PARTITION));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setPersistent(eq(false));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).addGatewayEventFilter(eq(mockGatewayEventFilterOne));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).addGatewayEventFilter(eq(mockGatewayEventFilterTwo));
|
||||
|
||||
AsyncEventQueue evenQueue = factoryBean.getObject();
|
||||
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
|
||||
|
||||
assertNotNull(evenQueue);
|
||||
assertEquals("123abc", evenQueue.getId());
|
||||
assertThat(asyncEventQueue).isNotNull();
|
||||
assertThat(asyncEventQueue.getId()).isEqualTo("testQueue");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testAsyncEventQueueWithDiskSynchronousSetPersistenceUnset() throws Exception {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory = createMockAsyncEventQueueFactory("12345");
|
||||
public void doInitConfiguresConcurrentParallelAsyncEventQueue() throws Exception {
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory));
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory =
|
||||
mockAsyncEventQueueFactory("concurrentParallelQueue");
|
||||
|
||||
factoryBean.setName("12345");
|
||||
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
|
||||
factoryBean.setDiskSynchronous(true);
|
||||
AsyncEventQueueFactoryBean factoryBean =
|
||||
new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory));
|
||||
|
||||
factoryBean.setAsyncEventListener(mockAsyncEventListener());
|
||||
factoryBean.setDispatcherThreads(8);
|
||||
factoryBean.setName("concurrentParallelQueue");
|
||||
factoryBean.setParallel(true);
|
||||
factoryBean.doInit();
|
||||
|
||||
verifyExpectations(mockAsyncEventQueueFactory, factoryBean);
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setDispatcherThreads(eq(8));
|
||||
verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any());
|
||||
verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setOrderPolicy(any());
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true));
|
||||
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
|
||||
|
||||
AsyncEventQueue evenQueue = factoryBean.getObject();
|
||||
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
|
||||
|
||||
assertNotNull(evenQueue);
|
||||
assertEquals("12345", evenQueue.getId());
|
||||
assertThat(asyncEventQueue).isNotNull();
|
||||
assertThat(asyncEventQueue.getId()).isEqualTo("concurrentParallelQueue");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void doInitConfiguresParallelAsyncEventQueue() throws Exception {
|
||||
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory =
|
||||
mockAsyncEventQueueFactory("parallelQueue");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean =
|
||||
new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory));
|
||||
|
||||
factoryBean.setAsyncEventListener(mockAsyncEventListener());
|
||||
factoryBean.setName("parallelQueue");
|
||||
factoryBean.setParallel(true);
|
||||
factoryBean.doInit();
|
||||
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDispatcherThreads(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any());
|
||||
verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setOrderPolicy(any());
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true));
|
||||
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
|
||||
|
||||
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
|
||||
|
||||
assertThat(asyncEventQueue).isNotNull();
|
||||
assertThat(asyncEventQueue.getId()).isEqualTo("parallelQueue");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void doInitConfiguresConcurrentSerialAsyncEventQueue() throws Exception {
|
||||
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory =
|
||||
mockAsyncEventQueueFactory("concurrentSerialQueue");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean =
|
||||
new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory));
|
||||
|
||||
factoryBean.setAsyncEventListener(mockAsyncEventListener());
|
||||
factoryBean.setDispatcherThreads(16);
|
||||
factoryBean.setName("concurrentSerialQueue");
|
||||
factoryBean.setParallel(false);
|
||||
factoryBean.doInit();
|
||||
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setDispatcherThreads(eq(16));
|
||||
verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any());
|
||||
verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setOrderPolicy(any());
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false));
|
||||
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
|
||||
|
||||
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
|
||||
|
||||
assertThat(asyncEventQueue).isNotNull();
|
||||
assertThat(asyncEventQueue.getId()).isEqualTo("concurrentSerialQueue");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void doInitConfiguresSerialAsyncEventQueue() throws Exception {
|
||||
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory = mockAsyncEventQueueFactory("serialQueue");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean =
|
||||
new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory));
|
||||
|
||||
factoryBean.setAsyncEventListener(mockAsyncEventListener());
|
||||
factoryBean.setName("serialQueue");
|
||||
factoryBean.doInit();
|
||||
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDispatcherThreads(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any());
|
||||
verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setOrderPolicy(any());
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false));
|
||||
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
|
||||
|
||||
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
|
||||
|
||||
assertThat(asyncEventQueue).isNotNull();
|
||||
assertThat(asyncEventQueue.getId()).isEqualTo("serialQueue");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void doInitConfiguresSerialAsyncEventQueueWithOrderPolicy() throws Exception {
|
||||
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory =
|
||||
mockAsyncEventQueueFactory("orderedSerialQueue");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean =
|
||||
new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory));
|
||||
|
||||
factoryBean.setAsyncEventListener(mockAsyncEventListener());
|
||||
factoryBean.setName("orderedSerialQueue");
|
||||
factoryBean.setOrderPolicy(GatewaySender.OrderPolicy.THREAD);
|
||||
factoryBean.doInit();
|
||||
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDispatcherThreads(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any());
|
||||
verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setOrderPolicy(eq(GatewaySender.OrderPolicy.THREAD));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false));
|
||||
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
|
||||
|
||||
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
|
||||
|
||||
assertThat(asyncEventQueue).isNotNull();
|
||||
assertThat(asyncEventQueue.getId()).isEqualTo("orderedSerialQueue");
|
||||
}
|
||||
|
||||
@Test(expected = IllegalStateException.class)
|
||||
public void doInitWithNullAsyncEventListenerThrowsIllegalStateException() throws Exception {
|
||||
|
||||
try {
|
||||
new AsyncEventQueueFactoryBean(mockCache(), null).doInit();
|
||||
}
|
||||
catch (IllegalStateException expected) {
|
||||
|
||||
assertThat(expected).hasMessage("AsyncEventListener must not be null");
|
||||
assertThat(expected).hasNoCause();
|
||||
|
||||
throw expected;
|
||||
}
|
||||
}
|
||||
|
||||
@Test(expected = IllegalStateException.class)
|
||||
public void doInitWithParallelAsyncEventQueueHavingAnOrderPolicyThrowsIllegalStateException() throws Exception {
|
||||
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory =
|
||||
mockAsyncEventQueueFactory("parallelQueue");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean =
|
||||
new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory));
|
||||
|
||||
factoryBean.setAsyncEventListener(mockAsyncEventListener());
|
||||
factoryBean.setName("parallelQueue");
|
||||
factoryBean.setOrderPolicy(GatewaySender.OrderPolicy.KEY);
|
||||
factoryBean.setParallel(true);
|
||||
|
||||
try {
|
||||
factoryBean.doInit();
|
||||
}
|
||||
catch (IllegalStateException expected) {
|
||||
|
||||
assertThat(expected).hasMessage("OrderPolicy cannot be used with a Parallel AsyncEventQueue");
|
||||
assertThat(expected).hasNoCause();
|
||||
|
||||
throw expected;
|
||||
}
|
||||
finally {
|
||||
|
||||
assertThat(factoryBean.getObject()).isNull();
|
||||
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskStoreName(anyString());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDiskSynchronous(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setDispatcherThreads(eq(8));
|
||||
verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any());
|
||||
verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setOrderPolicy(any());
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true));
|
||||
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void doInitConfiguresAsyncEventQueueWithSynchronousOverflowDiskStoreNoPersistence() throws Exception {
|
||||
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory =
|
||||
mockAsyncEventQueueFactory("nonPersistentSynchronousOverflowQueue");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean =
|
||||
new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory));
|
||||
|
||||
factoryBean.setAsyncEventListener(mockAsyncEventListener());
|
||||
factoryBean.setDiskStoreRef("queueOverflowDiskStore");
|
||||
factoryBean.setDiskSynchronous(true);
|
||||
factoryBean.setName("nonPersistentSynchronousOverflowQueue");
|
||||
factoryBean.setOrderPolicy(GatewaySender.OrderPolicy.KEY);
|
||||
factoryBean.setPersistent(false);
|
||||
factoryBean.doInit();
|
||||
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchSize(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setBatchTimeInterval(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setDiskStoreName("queueOverflowDiskStore");
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setDiskSynchronous(eq(true));
|
||||
verify(mockAsyncEventQueueFactory, never()).setDispatcherThreads(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, never()).setForwardExpirationDestroy(anyBoolean());
|
||||
verify(mockAsyncEventQueueFactory, never()).setGatewayEventSubstitutionListener(any());
|
||||
verify(mockAsyncEventQueueFactory, never()).setMaximumQueueMemory(anyInt());
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setOrderPolicy(eq(GatewaySender.OrderPolicy.KEY));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false));
|
||||
verify(mockAsyncEventQueueFactory, times(1)).setPersistent(eq(false));
|
||||
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
|
||||
|
||||
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
|
||||
|
||||
assertThat(asyncEventQueue).isNotNull();
|
||||
assertThat(asyncEventQueue.getId()).isEqualTo("nonPersistentSynchronousOverflowQueue");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -65,17 +65,17 @@ public class GatewaySenderFactoryBeanTest {
|
||||
return mockGatewaySenderFactory;
|
||||
}
|
||||
|
||||
protected void verifyExpectations(final GatewaySenderFactoryBean factoryBean,
|
||||
final GatewaySenderFactory mockGatewaySenderFactory) throws Exception {
|
||||
protected void verifyExpectations(GatewaySenderFactoryBean factoryBean,
|
||||
GatewaySenderFactory mockGatewaySenderFactory) throws Exception {
|
||||
|
||||
Boolean parallel = TestUtils.readField("parallel", factoryBean);
|
||||
|
||||
verify(mockGatewaySenderFactory).setParallel(eq(Boolean.TRUE.equals(parallel)));
|
||||
|
||||
String orderPolicy = TestUtils.readField("orderPolicy", factoryBean);
|
||||
GatewaySender.OrderPolicy orderPolicy = TestUtils.readField("orderPolicy", factoryBean);
|
||||
|
||||
if (orderPolicy != null) {
|
||||
verify(mockGatewaySenderFactory).setOrderPolicy(
|
||||
eq(GatewaySender.OrderPolicy.valueOf(orderPolicy.toUpperCase())));
|
||||
verify(mockGatewaySenderFactory).setOrderPolicy(eq(orderPolicy));
|
||||
}
|
||||
|
||||
Integer dispatcherThreads = TestUtils.readField("dispatcherThreads", factoryBean);
|
||||
@@ -165,7 +165,7 @@ public class GatewaySenderFactoryBeanTest {
|
||||
factoryBean.doInit();
|
||||
}
|
||||
catch (IllegalArgumentException expected) {
|
||||
assertEquals("Order Policy cannot be used with a Parallel Gateway Sender Queue.", expected.getMessage());
|
||||
assertEquals("OrderPolicy cannot be used with a Parallel GatewaySender", expected.getMessage());
|
||||
throw expected;
|
||||
}
|
||||
}
|
||||
@@ -194,7 +194,9 @@ public class GatewaySenderFactoryBeanTest {
|
||||
|
||||
@Test
|
||||
public void testGatewaySenderWithOrderPolicyAndDispatcherThreads() throws Exception {
|
||||
GatewaySenderFactory mockGatewaySenderFactory = createMockGatewaySenderFactory("g5", 42);
|
||||
|
||||
GatewaySenderFactory mockGatewaySenderFactory =
|
||||
createMockGatewaySenderFactory("g5", 42);
|
||||
|
||||
GatewaySenderFactoryBean factoryBean = new GatewaySenderFactoryBean(
|
||||
createMockCacheWithGatewayInfrastructure(mockGatewaySenderFactory));
|
||||
@@ -256,5 +258,4 @@ public class GatewaySenderFactoryBeanTest {
|
||||
assertEquals("g7", gatewaySender.getId());
|
||||
assertEquals(51, gatewaySender.getRemoteDSId());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
<?xml version="1.0" encoding="utf-8"?>
|
||||
<beans xmlns="http://www.springframework.org/schema/beans"
|
||||
xmlns:c="http://www.springframework.org/schema/c"
|
||||
xmlns:context="http://www.springframework.org/schema/context"
|
||||
xmlns:gfe="http://www.springframework.org/schema/gemfire"
|
||||
xmlns:util="http://www.springframework.org/schema/util"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xsi:schemaLocation="
|
||||
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd
|
||||
http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context.xsd
|
||||
http://www.springframework.org/schema/gemfire http://www.springframework.org/schema/gemfire/spring-gemfire.xsd
|
||||
http://www.springframework.org/schema/util http://www.springframework.org/schema/util/spring-util.xsd
|
||||
">
|
||||
|
||||
<bean class="org.springframework.data.gemfire.test.GemfireTestBeanPostProcessor"/>
|
||||
|
||||
<util:properties id="gemfireProperties">
|
||||
<prop key="name">AsyncEventQueueNamespaceTest</prop>
|
||||
<prop key="mcast-port">0</prop>
|
||||
<prop key="log-level">error</prop>
|
||||
</util:properties>
|
||||
|
||||
<context:property-placeholder/>
|
||||
|
||||
<gfe:cache properties-ref="gemfireProperties"/>
|
||||
|
||||
<gfe:disk-store id="TestDiskStore">
|
||||
<gfe:disk-dir location="${java.io.tmpdir}" max-size="100"/>
|
||||
</gfe:disk-store>
|
||||
|
||||
<gfe:async-event-queue id="TestAsyncEventQueue"
|
||||
batch-conflation-enabled="true"
|
||||
batch-size="100"
|
||||
batch-time-interval="30"
|
||||
disk-store-ref="TestDiskStore"
|
||||
disk-synchronous="true"
|
||||
dispatcher-threads="4"
|
||||
forward-expiration-destroy="true"
|
||||
maximum-queue-memory="50"
|
||||
order-policy="KEY"
|
||||
parallel="false"
|
||||
persistent="true">
|
||||
<gfe:async-event-listener>
|
||||
<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>
|
||||
Reference in New Issue
Block a user