SGF-505 - Add support for forwardExpirationDestroy in the AsyncEventQueueFactoryBean API and XML namespace.

This commit is contained in:
John Blum
2016-06-20 15:08:52 -07:00
parent fbe1ef1080
commit a5d6777589
7 changed files with 81 additions and 72 deletions

View File

@@ -42,6 +42,7 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser {
@Override
protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) {
builder.setLazyInit(false);
parseAsyncEventListener(element, parserContext, builder);
@@ -54,12 +55,11 @@ 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, "ignore-eviction-and-expiration");
ParsingUtils.setPropertyValue(element, builder, "foward-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");
ParsingUtils.setPropertyValue(element, builder, NAME_ATTRIBUTE);
if (!StringUtils.hasText(element.getAttribute(NAME_ATTRIBUTE))) {
@@ -82,8 +82,7 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser {
}
}
private void parseAsyncEventListener(final Element element, final ParserContext parserContext,
final BeanDefinitionBuilder builder) {
private void parseAsyncEventListener(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) {
Element asyncEventListenerElement = DomUtils.getChildElementByTagName(element, "async-event-listener");
@@ -97,7 +96,7 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser {
}
}
private void parseCache(final Element element, final BeanDefinitionBuilder builder) {
private void parseCache(Element element, BeanDefinitionBuilder builder) {
String cacheRefAttribute = element.getAttribute("cache-ref");
String cacheName = (StringUtils.hasText(cacheRefAttribute) ? cacheRefAttribute
@@ -106,7 +105,7 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser {
builder.addConstructorArgReference(cacheName);
}
private void parseDiskStore(final Element element, final BeanDefinitionBuilder builder) {
private void parseDiskStore(Element element, BeanDefinitionBuilder builder) {
ParsingUtils.setPropertyValue(element, builder, "disk-store-ref");
String diskStoreRef = element.getAttribute("disk-store-ref");
@@ -115,5 +114,4 @@ class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser {
builder.addDependsOn(diskStoreRef);
}
}
}

View File

@@ -26,7 +26,7 @@ import com.gemstone.gemfire.cache.wan.GatewaySender;
/**
* FactoryBean for creating GemFire {@link AsyncEventQueue}s.
*
*
* @author David Turanski
* @author John Blum
*/
@@ -39,7 +39,7 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
private Boolean batchConflationEnabled;
private Boolean diskSynchronous;
private Boolean ignoreEvictionAndExpiration;
private Boolean forwardExpirationDestroy;
private Boolean parallel;
private Boolean persistent;
@@ -53,21 +53,21 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
/**
* Constructs an instance of the AsyncEventQueueFactoryBean for creating an GemFire AsyncEventQueue.
*
*
* @param cache the GemFire Cache reference.
* @see #AsyncEventQueueFactoryBean(com.gemstone.gemfire.cache.Cache, com.gemstone.gemfire.cache.asyncqueue.AsyncEventListener)
*/
public AsyncEventQueueFactoryBean(final Cache cache) {
public AsyncEventQueueFactoryBean(Cache cache) {
this(cache, null);
}
/**
* Constructs an instance of the AsyncEventQueueFactoryBean for creating an GemFire AsyncEventQueue.
*
*
* @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);
}
@@ -79,12 +79,13 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
@Override
public Class<?> getObjectType() {
return AsyncEventQueue.class;
return (asyncEventQueue != null ? asyncEventQueue.getClass() : AsyncEventQueue.class);
}
@Override
protected void doInit() {
Assert.notNull(this.asyncEventListener, "The AsyncEventListener cannot be null.");
Assert.notNull(this.asyncEventListener, "AsyncEventListener must not be null");
AsyncEventQueueFactory asyncEventQueueFactory = (this.factory != null ? (AsyncEventQueueFactory) factory
: cache.createAsyncEventQueueFactory());
@@ -113,8 +114,8 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
asyncEventQueueFactory.setDiskSynchronous(diskSynchronous);
}
if (ignoreEvictionAndExpiration != null) {
asyncEventQueueFactory.setIgnoreEvictionAndExpiration(ignoreEvictionAndExpiration);
if (forwardExpirationDestroy != null) {
asyncEventQueueFactory.setForwardExpirationDestroy(forwardExpirationDestroy);
}
if (maximumQueueMemory != null) {
@@ -124,10 +125,10 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
asyncEventQueueFactory.setParallel(isParallelEventQueue());
if (orderPolicy != null) {
Assert.isTrue(isSerialEventQueue(), "Order Policy cannot be used with a Parallel Event Queue.");
Assert.isTrue(isSerialEventQueue(), "Order Policy cannot be used with a Parallel Event Queue");
Assert.isTrue(VALID_ORDER_POLICIES.contains(orderPolicy.toUpperCase()), String.format(
"The value of Order Policy '$1%s' is invalid.", orderPolicy));
"The value of Order Policy '$1%s' is invalid", orderPolicy));
asyncEventQueueFactory.setOrderPolicy(GatewaySender.OrderPolicy.valueOf(orderPolicy.toUpperCase()));
}
@@ -152,7 +153,8 @@ 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;
}
@@ -209,9 +211,19 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
this.dispatcherThreads = dispatcherThreads;
}
/* (non-Javadoc) */
public void setIgnoreEvictionAndExpiration(Boolean ignoreEvictionAndExpiration) {
this.ignoreEvictionAndExpiration = ignoreEvictionAndExpiration;
/**
* 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 com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueueFactory#setForwardExpirationDestroy(boolean)
* @see com.gemstone.gemfire.cache.ExpirationAttributes#getAction()
* @see com.gemstone.gemfire.cache.ExpirationAction#DESTROY
*/
public void setForwardExpirationDestroy(Boolean forwardExpirationDestroy) {
this.forwardExpirationDestroy = forwardExpirationDestroy;
}
public void setMaximumQueueMemory(Integer maximumQueueMemory) {
@@ -233,14 +245,14 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
this.parallel = parallel;
}
public boolean isSerialEventQueue() {
return !isParallelEventQueue();
}
public boolean isParallelEventQueue() {
return Boolean.TRUE.equals(parallel);
}
public boolean isSerialEventQueue() {
return !isParallelEventQueue();
}
public void setPersistent(Boolean persistent) {
this.persistent = persistent;
}

View File

@@ -2952,7 +2952,7 @@ if an inner bean.
]]></xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="ignore-eviction-and-expiration" type="xsd:string" default="true"/>
<xsd:attribute name="forward-expiration-destroy" type="xsd:string" default="false" use="optional"/>
<xsd:attributeGroup ref="commonWANQueueAttributes" />
</xsd:complexType>
<!-- -->

View File

@@ -67,7 +67,7 @@ public class AsyncEventQueueNamespaceTest {
assertThat(asyncEventQueue.getDiskStoreName(), is(equalTo("TestDiskStore")));
assertThat(asyncEventQueue.isDiskSynchronous(), is(true));
assertThat(asyncEventQueue.getDispatcherThreads(), is(equalTo(4)));
assertThat(asyncEventQueue.isIgnoreEvictionAndExpiration(), is(false));
assertThat(asyncEventQueue.isForwardExpirationDestroy(), is(false));
assertThat(asyncEventQueue.getMaximumQueueMemory(), is(equalTo(50)));
assertThat(asyncEventQueue.getOrderPolicy(), is(equalTo(GatewaySender.OrderPolicy.KEY)));
assertThat(asyncEventQueue.isParallel(), is(false));

View File

@@ -1,11 +1,11 @@
/*
* Copyright 2002-2013 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.
@@ -38,7 +38,7 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
private boolean batchConflationEnabled;
private boolean diskSynchronous;
private boolean ignoreEvictionAndExpiration;
private boolean forwardExpirationDestroy;
private boolean parallel;
private boolean persistent;
@@ -67,7 +67,7 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
when(asyncEventQueue.getGatewayEventFilters()).thenReturn(Collections.unmodifiableList(gatewayEventFilters));
when(asyncEventQueue.getGatewayEventSubstitutionFilter()).thenReturn(this.gatewayEventSubstitutionFilter);
when(asyncEventQueue.getId()).thenReturn(name);
when(asyncEventQueue.isIgnoreEvictionAndExpiration()).thenReturn(this.ignoreEvictionAndExpiration);
when(asyncEventQueue.isForwardExpirationDestroy()).thenReturn(this.forwardExpirationDestroy);
when(asyncEventQueue.getMaximumQueueMemory()).thenReturn(this.maxQueueMemory);
when(asyncEventQueue.getOrderPolicy()).thenReturn(this.orderPolicy);
when(asyncEventQueue.isParallel()).thenReturn(this.parallel);
@@ -108,8 +108,8 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
return this;
}
public AsyncEventQueueFactory setIgnoreEvictionAndExpiration(boolean ignoreEvictionAndExpiration) {
this.ignoreEvictionAndExpiration = ignoreEvictionAndExpiration;
public AsyncEventQueueFactory setForwardExpirationDestroy(boolean forwardExpirationDestroy) {
this.forwardExpirationDestroy = forwardExpirationDestroy;
return this;
}

View File

@@ -53,14 +53,13 @@ import com.gemstone.gemfire.cache.wan.GatewaySender;
*/
public class AsyncEventQueueFactoryBeanTest {
protected Cache createMockCacheWithAsyncEventQueueInfrastructure(
final AsyncEventQueueFactory mockAsynEventQueueFactory) {
protected Cache createMockCacheWithAsyncEventQueueInfrastructure(AsyncEventQueueFactory mockAsyncEventQueueFactory) {
Cache mockCache = mock(Cache.class);
when((mockCache.createAsyncEventQueueFactory())).thenReturn(mockAsynEventQueueFactory);
when((mockCache.createAsyncEventQueueFactory())).thenReturn(mockAsyncEventQueueFactory);
return mockCache;
}
protected AsyncEventQueueFactory createMockAsyncEventQueueFactory(final String asyncEventQueueId) {
protected AsyncEventQueueFactory createMockAsyncEventQueueFactory(String asyncEventQueueId) {
AsyncEventQueueFactory mockAsyncEventQueueFactory = mock(AsyncEventQueueFactory.class);
AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class);
@@ -71,48 +70,49 @@ public class AsyncEventQueueFactoryBeanTest {
return mockAsyncEventQueueFactory;
}
protected AsyncEventListener createMockAsyncEventListener() {
protected AsyncEventListener mockAsyncEventListener() {
return mock(AsyncEventListener.class);
}
protected void verifyExpectations(final AsyncEventQueueFactory mockAsyncEventQueueFactory,
final AsyncEventQueueFactoryBean factoryBean) throws Exception {
Boolean parallel = TestUtils.readField("parallel", factoryBean);
protected void verifyExpectations(AsyncEventQueueFactory asyncEventQueueFactory,
AsyncEventQueueFactoryBean asyncEventQueueFactoryBean) throws Exception {
verify(mockAsyncEventQueueFactory).setParallel(eq(Boolean.TRUE.equals(parallel)));
Boolean parallel = TestUtils.readField("parallel", asyncEventQueueFactoryBean);
String orderPolicy = TestUtils.readField("orderPolicy", factoryBean);
verify(asyncEventQueueFactory).setParallel(eq(Boolean.TRUE.equals(parallel)));
String orderPolicy = TestUtils.readField("orderPolicy", asyncEventQueueFactoryBean);
if (orderPolicy != null) {
verify(mockAsyncEventQueueFactory).setOrderPolicy(eq(GatewaySender.OrderPolicy.valueOf(
verify(asyncEventQueueFactory).setOrderPolicy(eq(GatewaySender.OrderPolicy.valueOf(
orderPolicy.toUpperCase())));
}
Integer dispatcherThreads = TestUtils.readField("dispatcherThreads", factoryBean);
Integer dispatcherThreads = TestUtils.readField("dispatcherThreads", asyncEventQueueFactoryBean);
if (dispatcherThreads != null) {
verify(mockAsyncEventQueueFactory).setDispatcherThreads(eq(dispatcherThreads));
verify(asyncEventQueueFactory).setDispatcherThreads(eq(dispatcherThreads));
}
String diskStoreReference = TestUtils.readField("diskStoreReference", factoryBean);
String diskStoreReference = TestUtils.readField("diskStoreReference", asyncEventQueueFactoryBean);
if (diskStoreReference != null) {
verify(mockAsyncEventQueueFactory).setDiskStoreName(eq(diskStoreReference));
verify(asyncEventQueueFactory).setDiskStoreName(eq(diskStoreReference));
}
Boolean diskSynchronous = TestUtils.readField("diskSynchronous", factoryBean);
Boolean diskSynchronous = TestUtils.readField("diskSynchronous", asyncEventQueueFactoryBean);
if (diskSynchronous != null) {
verify(mockAsyncEventQueueFactory).setDiskSynchronous(eq(diskSynchronous));
verify(asyncEventQueueFactory).setDiskSynchronous(eq(diskSynchronous));
}
Boolean persistent = TestUtils.readField("persistent", factoryBean);
Boolean persistent = TestUtils.readField("persistent", asyncEventQueueFactoryBean);
if (persistent != null) {
verify(mockAsyncEventQueueFactory).setPersistent(eq(persistent));
verify(asyncEventQueueFactory).setPersistent(eq(persistent));
}
else {
verify(mockAsyncEventQueueFactory, never()).setPersistent(true);
verify(asyncEventQueueFactory, never()).setPersistent(true);
}
}
@@ -121,13 +121,13 @@ public class AsyncEventQueueFactoryBeanTest {
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
createMockCacheWithAsyncEventQueueInfrastructure(createMockAsyncEventQueueFactory("testEventQueue")));
AsyncEventListener listenerOne = createMockAsyncEventListener();
AsyncEventListener listenerOne = mockAsyncEventListener();
factoryBean.setAsyncEventListener(listenerOne);
assertSame(listenerOne, TestUtils.readField("asyncEventListener", factoryBean));
AsyncEventListener listenerTwo = createMockAsyncEventListener();
AsyncEventListener listenerTwo = mockAsyncEventListener();
factoryBean.setAsyncEventListener(listenerTwo);
@@ -143,7 +143,7 @@ public class AsyncEventQueueFactoryBeanTest {
factoryBean.setName(asyncEventQueueId);
AsyncEventListener listenerOne = createMockAsyncEventListener();
AsyncEventListener listenerOne = mockAsyncEventListener();
factoryBean.setAsyncEventListener(listenerOne);
@@ -154,10 +154,10 @@ public class AsyncEventQueueFactoryBeanTest {
assertNotNull(TestUtils.readField("asyncEventQueue", factoryBean));
try {
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
factoryBean.setAsyncEventListener(mockAsyncEventListener());
}
catch (IllegalStateException expected) {
assertEquals("Setting an AsyncEventListener is not allowed once the AsyncEventQueue has been created.",
assertEquals("Setting an AsyncEventListener is not allowed once the AsyncEventQueue has been created",
expected.getMessage());
assertSame(listenerOne, TestUtils.readField("asyncEventListener", factoryBean));
throw expected;
@@ -175,7 +175,7 @@ public class AsyncEventQueueFactoryBeanTest {
factoryBean.doInit();
}
catch (Exception e) {
assertEquals("The AsyncEventListener cannot be null.", e.getMessage());
assertEquals("AsyncEventListener must not be null", e.getMessage());
throw e;
}
}
@@ -188,7 +188,7 @@ public class AsyncEventQueueFactoryBeanTest {
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory));
factoryBean.setName("000");
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
factoryBean.setAsyncEventListener(mockAsyncEventListener());
factoryBean.setDispatcherThreads(8);
factoryBean.setParallel(true);
factoryBean.doInit();
@@ -209,7 +209,7 @@ public class AsyncEventQueueFactoryBeanTest {
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory));
factoryBean.setName("123");
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
factoryBean.setAsyncEventListener(mockAsyncEventListener());
factoryBean.setParallel(true);
factoryBean.doInit();
@@ -229,7 +229,7 @@ public class AsyncEventQueueFactoryBeanTest {
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory));
factoryBean.setName("456");
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
factoryBean.setAsyncEventListener(mockAsyncEventListener());
factoryBean.setOrderPolicy("Key");
factoryBean.setParallel(true);
@@ -237,7 +237,7 @@ public class AsyncEventQueueFactoryBeanTest {
factoryBean.doInit();
}
catch (IllegalArgumentException expected) {
assertEquals("Order Policy cannot be used with a Parallel Event Queue.",
assertEquals("Order Policy cannot be used with a Parallel Event Queue",
expected.getMessage());
throw expected;
}
@@ -251,7 +251,7 @@ public class AsyncEventQueueFactoryBeanTest {
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFatory));
factoryBean.setName("789");
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
factoryBean.setAsyncEventListener(mockAsyncEventListener());
factoryBean.setOrderPolicy("THREAD");
factoryBean.setParallel(false);
factoryBean.doInit();
@@ -272,7 +272,7 @@ public class AsyncEventQueueFactoryBeanTest {
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFatory));
factoryBean.setName("abc");
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
factoryBean.setAsyncEventListener(mockAsyncEventListener());
factoryBean.setDispatcherThreads(2);
factoryBean.setOrderPolicy("THREAD");
factoryBean.doInit();
@@ -293,7 +293,7 @@ public class AsyncEventQueueFactoryBeanTest {
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory));
factoryBean.setName("123abc");
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
factoryBean.setAsyncEventListener(mockAsyncEventListener());
factoryBean.setDiskStoreRef("queueOverflowDiskStore");
factoryBean.setPersistent(false);
factoryBean.doInit();
@@ -314,7 +314,7 @@ public class AsyncEventQueueFactoryBeanTest {
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFactory));
factoryBean.setName("12345");
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
factoryBean.setAsyncEventListener(mockAsyncEventListener());
factoryBean.setDiskSynchronous(true);
factoryBean.doInit();
@@ -325,5 +325,4 @@ public class AsyncEventQueueFactoryBeanTest {
assertNotNull(evenQueue);
assertEquals("12345", evenQueue.getId());
}
}

View File

@@ -31,7 +31,7 @@
disk-store-ref="TestDiskStore"
disk-synchronous="true"
dispatcher-threads="4"
ignore-eviction-and-expiration="false"
forward-expiration-destroy="true"
maximum-queue-memory="50"
order-policy="KEY"
parallel="false"