Fixes JIRA issue SGF-232 - user is unable to specify Order Policy with a Serial Async Event Queue. Plus, made additional changes to the GatewaySenderFactoryBean class and associated tests/code for SGF-231.
This commit is contained in:
@@ -15,6 +15,9 @@
|
||||
*/
|
||||
package org.springframework.data.gemfire.wan;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.List;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.springframework.beans.factory.BeanNameAware;
|
||||
@@ -34,6 +37,8 @@ import com.gemstone.gemfire.cache.Cache;
|
||||
public abstract class AbstractWANComponentFactoryBean<T> implements FactoryBean<T>, InitializingBean, BeanNameAware,
|
||||
DisposableBean {
|
||||
|
||||
protected static final List<String> VALID_ORDER_POLICIES = Arrays.asList("KEY", "PARTITION", "THREAD");
|
||||
|
||||
protected Log log = LogFactory.getLog(getClass());
|
||||
|
||||
protected final Cache cache;
|
||||
|
||||
@@ -35,8 +35,6 @@ import com.gemstone.gemfire.cache.util.Gateway;
|
||||
*/
|
||||
public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<AsyncEventQueue> {
|
||||
|
||||
private static List<String> validOrderPolicyValues = Arrays.asList("KEY", "PARTITION", "THREAD");
|
||||
|
||||
private AsyncEventListener asyncEventListener;
|
||||
|
||||
private AsyncEventQueue asyncEventQueue;
|
||||
@@ -60,7 +58,7 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
* @param cache the GemFire Cache reference.
|
||||
* @see #AsyncEventQueueFactoryBean(com.gemstone.gemfire.cache.Cache, com.gemstone.gemfire.cache.asyncqueue.AsyncEventListener)
|
||||
*/
|
||||
public AsyncEventQueueFactoryBean(Cache cache) {
|
||||
public AsyncEventQueueFactoryBean(final Cache cache) {
|
||||
this(cache, null);
|
||||
}
|
||||
|
||||
@@ -70,7 +68,7 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
* @param cache the GemFire Cache reference.
|
||||
* @param asyncEventListener required {@link AsyncEventListener}
|
||||
*/
|
||||
public AsyncEventQueueFactoryBean(Cache cache, AsyncEventListener asyncEventListener) {
|
||||
public AsyncEventQueueFactoryBean(final Cache cache, final AsyncEventListener asyncEventListener) {
|
||||
super(cache);
|
||||
setAsyncEventListener(asyncEventListener);
|
||||
}
|
||||
@@ -93,13 +91,13 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
: cache.createAsyncEventQueueFactory());
|
||||
|
||||
if (diskStoreRef != null) {
|
||||
persistent = (persistent == null) ? Boolean.TRUE : persistent;
|
||||
persistent = (persistent == null || persistent);
|
||||
Assert.isTrue(persistent, "Specifying a 'disk store' requires the persistent property to be true.");
|
||||
asyncEventQueueFactory.setDiskStoreName(diskStoreRef);
|
||||
}
|
||||
|
||||
if (diskSynchronous != null) {
|
||||
persistent = (persistent == null) ? Boolean.TRUE : persistent;
|
||||
persistent = (persistent == null || persistent);
|
||||
Assert.isTrue(persistent, "Specifying 'disk synchronous' requires the persistent property to be true.");
|
||||
asyncEventQueueFactory.setDiskSynchronous(diskSynchronous);
|
||||
}
|
||||
@@ -121,6 +119,7 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
}
|
||||
|
||||
if (dispatcherThreads != null) {
|
||||
Assert.isTrue(isSerialEventQueue(), "The number of Dispatcher Threads cannot be specified with a Parallel Event Queue.");
|
||||
asyncEventQueueFactory.setDispatcherThreads(dispatcherThreads);
|
||||
}
|
||||
|
||||
@@ -128,15 +127,13 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
asyncEventQueueFactory.setMaximumQueueMemory(maximumQueueMemory);
|
||||
}
|
||||
|
||||
if (parallel != null) {
|
||||
asyncEventQueueFactory.setParallel(parallel);
|
||||
}
|
||||
|
||||
if (orderPolicy != null) {
|
||||
Assert.isTrue(parallel, "specifying an order policy requires the parallel property to be true");
|
||||
asyncEventQueueFactory.setParallel(isParallelEventQueue());
|
||||
|
||||
Assert.isTrue(validOrderPolicyValues.contains(orderPolicy.toUpperCase()), String.format(
|
||||
"The value of order policy:'$1%s'' is invalid.", orderPolicy));
|
||||
if (orderPolicy != null) {
|
||||
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));
|
||||
|
||||
asyncEventQueueFactory.setOrderPolicy(Gateway.OrderPolicy.valueOf(orderPolicy.toUpperCase()));
|
||||
}
|
||||
@@ -181,11 +178,12 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
this.parallel = parallel;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param validOrderPolicyValues the validOrderPolicyValues to set
|
||||
*/
|
||||
public static void setValidOrderPolicyValues(List<String> validOrderPolicyValues) {
|
||||
AsyncEventQueueFactoryBean.validOrderPolicyValues = validOrderPolicyValues;
|
||||
public boolean isSerialEventQueue() {
|
||||
return !isParallelEventQueue();
|
||||
}
|
||||
|
||||
public boolean isParallelEventQueue() {
|
||||
return Boolean.TRUE.equals(parallel);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -38,8 +38,6 @@ import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
|
||||
public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<GatewaySender>
|
||||
implements SmartLifecycle {
|
||||
|
||||
private static final List<String> VALID_ORDER_POLICIES = Arrays.asList("KEY", "PARTITION", "THREAD");
|
||||
|
||||
private boolean manualStart = false;
|
||||
|
||||
private int remoteDistributedSystemId;
|
||||
|
||||
@@ -16,6 +16,7 @@
|
||||
package org.springframework.data.gemfire.config;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertFalse;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
@@ -95,7 +96,7 @@ public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest {
|
||||
AsyncEventQueue aseq = ctx.getBean("async-event-queue", AsyncEventQueue.class);
|
||||
assertEquals(10, aseq.getBatchSize());
|
||||
assertTrue(aseq.isPersistent());
|
||||
assertTrue(aseq.isParallel());
|
||||
assertFalse(aseq.isParallel());
|
||||
assertEquals("diskstore", aseq.getDiskStoreName());
|
||||
assertEquals(50, aseq.getMaximumQueueMemory());
|
||||
assertEquals(3, aseq.getBatchTimeInterval());
|
||||
@@ -160,7 +161,7 @@ public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest {
|
||||
assertEquals(50, gws.getMaximumQueueMemory());
|
||||
assertEquals(OrderPolicy.THREAD, gws.getOrderPolicy());
|
||||
assertTrue(gws.isPersistenceEnabled());
|
||||
assertTrue(gws.isParallel());
|
||||
assertFalse(gws.isParallel());
|
||||
assertEquals(16536, gws.getSocketBufferSize());
|
||||
assertEquals(3000, gws.getSocketReadTimeout());
|
||||
}
|
||||
|
||||
@@ -23,6 +23,7 @@ import static org.junit.Assert.assertSame;
|
||||
import static org.mockito.Matchers.eq;
|
||||
import static org.mockito.Matchers.notNull;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
import org.junit.Test;
|
||||
@@ -32,6 +33,7 @@ import com.gemstone.gemfire.cache.Cache;
|
||||
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventListener;
|
||||
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueue;
|
||||
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueueFactory;
|
||||
import com.gemstone.gemfire.cache.util.Gateway;
|
||||
|
||||
/**
|
||||
* The AsyncEventQueueFactoryBeanTest class is a test suite of test cases testing the contract and functionality
|
||||
@@ -40,32 +42,60 @@ import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueueFactory;
|
||||
* @author John Blum
|
||||
* @see org.junit.Test
|
||||
* @see org.mockito.Mockito
|
||||
* @see org.springframework.data.gemfire.TestUtils
|
||||
* @see org.springframework.data.gemfire.wan.AsyncEventQueueFactoryBean
|
||||
* @see com.gemstone.gemfire.cache.Cache
|
||||
* @see com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueue
|
||||
* @see com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueueFactory
|
||||
* @since 1.3.3
|
||||
*/
|
||||
public class AsyncEventQueueFactoryBeanTest {
|
||||
|
||||
protected Cache createMockCacheWithAsyncInfrastructure(final String asyncEventQueueId) {
|
||||
protected Cache createMockCacheWithAsyncEventQueueInfrastructure(
|
||||
final AsyncEventQueueFactory mockAsynEventQueueFactory) {
|
||||
Cache mockCache = mock(Cache.class);
|
||||
AsyncEventQueueFactory mockAsynEventQueueFactory = mock(AsyncEventQueueFactory.class);
|
||||
when((mockCache.createAsyncEventQueueFactory())).thenReturn(mockAsynEventQueueFactory);
|
||||
return mockCache;
|
||||
}
|
||||
|
||||
protected AsyncEventQueueFactory createMockAsyncEventQueueFactory(final String asyncEventQueueId) {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFactory = mock(AsyncEventQueueFactory.class);
|
||||
AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class);
|
||||
|
||||
when((mockCache.createAsyncEventQueueFactory())).thenReturn(mockAsynEventQueueFactory);
|
||||
when(mockAsynEventQueueFactory.create(eq(asyncEventQueueId), notNull(AsyncEventListener.class)))
|
||||
.thenReturn(mockAsyncEventQueue);
|
||||
when(mockAsyncEventQueue.getId()).thenReturn(asyncEventQueueId);
|
||||
when(mockAsyncEventQueueFactory.create(eq(asyncEventQueueId), notNull(AsyncEventListener.class)))
|
||||
.thenReturn(mockAsyncEventQueue);
|
||||
|
||||
return mockCache;
|
||||
return mockAsyncEventQueueFactory;
|
||||
}
|
||||
|
||||
protected AsyncEventListener createMockAsyncEventListener() {
|
||||
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(Gateway.OrderPolicy.valueOf(orderPolicy.toUpperCase())));
|
||||
}
|
||||
|
||||
Integer dispatcherThreads = TestUtils.readField("dispatcherThreads", factoryBean);
|
||||
|
||||
if (dispatcherThreads != null) {
|
||||
verify(mockAsyncEventQueueFactory).setDispatcherThreads(eq(dispatcherThreads));
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSetAsyncEventListener() throws Exception {
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncInfrastructure("testEventQueue"));
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(createMockAsyncEventQueueFactory("testEventQueue")));
|
||||
|
||||
AsyncEventListener listenerOne = createMockAsyncEventListener();
|
||||
|
||||
@@ -85,7 +115,7 @@ public class AsyncEventQueueFactoryBeanTest {
|
||||
String asyncEventQueueId = "testEventQueue";
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncInfrastructure(asyncEventQueueId));
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(createMockAsyncEventQueueFactory(asyncEventQueueId)));
|
||||
|
||||
factoryBean.setName(asyncEventQueueId);
|
||||
|
||||
@@ -114,7 +144,7 @@ public class AsyncEventQueueFactoryBeanTest {
|
||||
public void testDoInitWhenAsyncEventListenerIsNull() throws Exception {
|
||||
try {
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncInfrastructure("testEventQueue"));
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(createMockAsyncEventQueueFactory("testEventQueue")));
|
||||
|
||||
assertNull(TestUtils.readField("asyncEventListener", factoryBean));
|
||||
|
||||
@@ -126,4 +156,110 @@ public class AsyncEventQueueFactoryBeanTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testParallelAsyncEventQueue() throws Exception {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFatory = createMockAsyncEventQueueFactory("123");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFatory));
|
||||
|
||||
factoryBean.setName("123");
|
||||
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
|
||||
factoryBean.setParallel(true);
|
||||
factoryBean.doInit();
|
||||
|
||||
verifyExpectations(mockAsyncEventQueueFatory, factoryBean);
|
||||
|
||||
AsyncEventQueue eventQueue = factoryBean.getObject();
|
||||
|
||||
assertNotNull(eventQueue);
|
||||
assertEquals("123", eventQueue.getId());
|
||||
}
|
||||
|
||||
@Test(expected = IllegalArgumentException.class)
|
||||
public void testParallelAsyncEventQueueWithDispatcherThreads() {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFatory = createMockAsyncEventQueueFactory("456");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFatory));
|
||||
|
||||
factoryBean.setName("456");
|
||||
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
|
||||
factoryBean.setDispatcherThreads(1);
|
||||
factoryBean.setParallel(true);
|
||||
|
||||
try {
|
||||
factoryBean.doInit();
|
||||
}
|
||||
catch (IllegalArgumentException expected) {
|
||||
assertEquals("The number of Dispatcher Threads cannot be specified with a Parallel Event Queue.",
|
||||
expected.getMessage());
|
||||
throw expected;
|
||||
}
|
||||
}
|
||||
|
||||
@Test(expected = IllegalArgumentException.class)
|
||||
public void testParallelAsyncEventQueueWithOrderPolicy() {
|
||||
AsyncEventQueueFactory mockAsyncEventQueueFatory = createMockAsyncEventQueueFactory("456");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFatory));
|
||||
|
||||
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("789");
|
||||
|
||||
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(
|
||||
createMockCacheWithAsyncEventQueueInfrastructure(mockAsyncEventQueueFatory));
|
||||
|
||||
factoryBean.setName("789");
|
||||
factoryBean.setAsyncEventListener(createMockAsyncEventListener());
|
||||
factoryBean.setDispatcherThreads(2);
|
||||
factoryBean.setOrderPolicy("THREAD");
|
||||
factoryBean.doInit();
|
||||
|
||||
verifyExpectations(mockAsyncEventQueueFatory, factoryBean);
|
||||
|
||||
AsyncEventQueue eventQueue = factoryBean.getObject();
|
||||
|
||||
assertNotNull(eventQueue);
|
||||
assertEquals("789", eventQueue.getId());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -150,7 +150,7 @@ public class GatewaySenderFactoryBeanTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSerialGatewaySenderWithOrderPolicy() throws Exception {
|
||||
public void testSerialGatewaySenderWithDispatcherThreads() throws Exception {
|
||||
GatewaySenderFactory mockGatewaySenderFactory = createMockGatewaySenderFactory("g4", 21);
|
||||
|
||||
GatewaySenderFactoryBean factoryBean = new GatewaySenderFactoryBean(
|
||||
|
||||
@@ -23,7 +23,7 @@
|
||||
maximum-queue-memory="50"
|
||||
order-policy="THREAD"
|
||||
persistent="true"
|
||||
parallel="true"
|
||||
parallel="false"
|
||||
socket-buffer-size="16536"
|
||||
socket-read-timeout="3000">
|
||||
<gfe:event-filter>
|
||||
@@ -41,7 +41,7 @@
|
||||
persistent="true"
|
||||
disk-store-ref="diskstore"
|
||||
maximum-queue-memory="50"
|
||||
parallel="true"
|
||||
parallel="false"
|
||||
batch-conflation-enabled="true"
|
||||
batch-time-interval="3"
|
||||
dispatcher-threads="4"
|
||||
@@ -81,4 +81,4 @@
|
||||
<bean id="event-filter" class="org.springframework.data.gemfire.config.GemfireV7GatewayNamespaceTest.TestEventFilter"/>
|
||||
<bean id="transport-filter" class="org.springframework.data.gemfire.config.GemfireV7GatewayNamespaceTest.TestTransportFilter"/>
|
||||
|
||||
</beans>
|
||||
</beans>
|
||||
|
||||
Reference in New Issue
Block a user