DATAGEODE-242 - Add support for AEQ pauseEventDispatching.

This commit is contained in:
John Blum
2019-10-27 23:27:05 -07:00
parent 7b17982158
commit 2b9a0979a3
11 changed files with 302 additions and 145 deletions

View File

@@ -14,19 +14,18 @@
* 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.junit.Test;
import org.junit.runner.RunWith;
import org.apache.geode.cache.EntryEvent;
import org.apache.geode.cache.asyncqueue.AsyncEvent;
import org.apache.geode.cache.asyncqueue.AsyncEventListener;
@@ -36,10 +35,6 @@ 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;
@@ -68,22 +63,26 @@ public class AsyncEventQueueNamespaceTest {
@Resource(name = "TestAsyncEventQueueWithFilters")
private AsyncEventQueue asyncEventQueueWithFilters;
@Resource(name = "TestPausedAsyncEventQueue")
private AsyncEventQueue pausedAsyncEventQueue;
@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));
assertThat(asyncEventQueue).isNotNull();
assertThat(asyncEventQueue.getId()).isEqualTo("TestAsyncEventQueue");
assertThat(asyncEventQueue.isBatchConflationEnabled()).isTrue();
assertThat(asyncEventQueue.getBatchSize()).isEqualTo(100);
assertThat(asyncEventQueue.getBatchTimeInterval()).isEqualTo(30);
assertThat(asyncEventQueue.getDiskStoreName()).isEqualTo("TestDiskStore");
assertThat(asyncEventQueue.isDiskSynchronous()).isTrue();
assertThat(asyncEventQueue.getDispatcherThreads()).isEqualTo(4);
assertThat(asyncEventQueue.isDispatchingPaused()).isFalse();
assertThat(asyncEventQueue.isForwardExpirationDestroy()).isTrue();
assertThat(asyncEventQueue.getMaximumQueueMemory()).isEqualTo(50);
assertThat(asyncEventQueue.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.KEY);
assertThat(asyncEventQueue.isParallel()).isFalse();
assertThat(asyncEventQueue.isPersistent()).isTrue();
}
@Test
@@ -91,8 +90,8 @@ public class AsyncEventQueueNamespaceTest {
AsyncEventListener asyncEventListener = asyncEventQueue.getAsyncEventListener();
Assert.assertThat(asyncEventListener, is(notNullValue(AsyncEventListener.class)));
Assert.assertThat(asyncEventListener.toString(), is(equalTo("TestAeqListener")));
assertThat(asyncEventListener).isNotNull();
assertThat(asyncEventListener.toString()).isEqualTo("TestAeqListener");
}
@Test
@@ -100,6 +99,7 @@ public class AsyncEventQueueNamespaceTest {
assertThat(asyncEventQueueWithFilters).isNotNull();
assertThat(asyncEventQueueWithFilters.getId()).isEqualTo("TestAsyncEventQueueWithFilters");
assertThat(asyncEventQueueWithFilters.isDispatchingPaused()).isFalse();
AsyncEventListener listener = asyncEventQueueWithFilters.getAsyncEventListener();
@@ -120,6 +120,14 @@ public class AsyncEventQueueNamespaceTest {
assertThat(gatewayEventSubstitutionFilter.toString()).isEqualTo("GatewayEventSubstitutionFilterOne");
}
@Test
public void pausedAsyncEventQueueIsConfiguredProperly() {
assertThat(pausedAsyncEventQueue).isNotNull();
assertThat(pausedAsyncEventQueue.getId()).isEqualTo("TestPausedAsyncEventQueue");
assertThat(pausedAsyncEventQueue.isDispatchingPaused()).isTrue();
}
public static class TestAsyncEventListener implements AsyncEventListener {
private final String name;
@@ -180,6 +188,7 @@ public class AsyncEventQueueNamespaceTest {
public TestGatewayEventSubstitutionFilter(String name) {
this.name = name;
}
@Override
public Object getSubstituteValue(EntryEvent<Object, Object> event) {
return null;

View File

@@ -12,6 +12,7 @@
*/
package org.springframework.data.gemfire.test;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@@ -40,6 +41,7 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
private boolean diskSynchronous;
private boolean forwardExpirationDestroy;
private boolean parallel;
private boolean pauseEventDispatching;
private boolean persistent;
private int batchSize;
@@ -65,6 +67,7 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
when(asyncEventQueue.getDiskStoreName()).thenReturn(this.diskStoreName);
when(asyncEventQueue.isDiskSynchronous()).thenReturn(this.diskSynchronous);
when(asyncEventQueue.getDispatcherThreads()).thenReturn(this.dispatcherThreads);
when(asyncEventQueue.isDispatchingPaused()).thenAnswer(invocation -> this.pauseEventDispatching);
when(asyncEventQueue.getGatewayEventFilters()).thenReturn(Collections.unmodifiableList(gatewayEventFilters));
when(asyncEventQueue.getGatewayEventSubstitutionFilter()).thenReturn(this.gatewayEventSubstitutionFilter);
when(asyncEventQueue.getId()).thenReturn(name);
@@ -74,6 +77,11 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
when(asyncEventQueue.isParallel()).thenReturn(this.parallel);
when(asyncEventQueue.isPersistent()).thenReturn(this.persistent);
doAnswer(invocation -> {
this.pauseEventDispatching = false;
return null;
}).when(asyncEventQueue).resumeEventDispatching();
return this.asyncEventQueue;
}
@@ -145,6 +153,7 @@ public class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
}
public AsyncEventQueueFactory pauseEventDispatching() {
this.pauseEventDispatching = true;
return this;
}

View File

@@ -63,6 +63,9 @@ import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
import java.util.stream.Collectors;
import org.mockito.ArgumentMatchers;
import org.mockito.stubbing.Answer;
import org.apache.geode.cache.AttributesMutator;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.CacheFactory;
@@ -128,8 +131,6 @@ import org.apache.geode.internal.concurrent.ConcurrentHashSet;
import org.apache.geode.pdx.PdxSerializer;
import org.apache.lucene.analysis.Analyzer;
import org.mockito.ArgumentMatchers;
import org.mockito.stubbing.Answer;
import org.springframework.data.gemfire.IndexType;
import org.springframework.data.gemfire.server.SubscriptionEvictionPolicy;
@@ -640,6 +641,7 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport {
AtomicBoolean diskSynchronous = new AtomicBoolean(true);
AtomicBoolean forwardExpirationDestroy = new AtomicBoolean(false);
AtomicBoolean parallel = new AtomicBoolean(false);
AtomicBoolean pauseEventDispatching = new AtomicBoolean(false);
AtomicBoolean persistent = new AtomicBoolean(false);
AtomicInteger batchSize = new AtomicInteger(100);
@@ -697,6 +699,9 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport {
return mockAsyncEventQueueFactory;
});
doAnswer(newSetter(pauseEventDispatching, true, mockAsyncEventQueueFactory))
.when(mockAsyncEventQueueFactory).pauseEventDispatching();
when(mockAsyncEventQueueFactory.removeGatewayEventFilter(any(GatewayEventFilter.class))).thenAnswer(invocation -> {
gatewayEventFilters.remove(invocation.<GatewayEventFilter>getArgument(0));
@@ -714,6 +719,7 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport {
when(mockAsyncEventQueue.isBatchConflationEnabled()).thenAnswer(newGetter(batchConflationEnabled));
when(mockAsyncEventQueue.isDiskSynchronous()).thenAnswer(newGetter(diskSynchronous));
when(mockAsyncEventQueue.isDispatchingPaused()).thenAnswer(newGetter(pauseEventDispatching));
when(mockAsyncEventQueue.isForwardExpirationDestroy()).thenAnswer(newGetter(forwardExpirationDestroy));
when(mockAsyncEventQueue.isParallel()).thenAnswer(newGetter(parallel));
when(mockAsyncEventQueue.isPersistent()).thenAnswer(newGetter(persistent));
@@ -730,6 +736,11 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport {
when(mockAsyncEventQueue.getMaximumQueueMemory()).thenAnswer(newGetter(maximumQueueMemory));
when(mockAsyncEventQueue.getOrderPolicy()).thenAnswer(newGetter(orderPolicy));
doAnswer(resumeEventDispatchingInvocation -> {
pauseEventDispatching.set(false);
return null;
}).when(mockAsyncEventQueue).resumeEventDispatching();
when(mockAsyncEventQueue.size()).thenReturn(0);
return mockAsyncEventQueue;

View File

@@ -14,7 +14,6 @@
* limitations under the License.
*
*/
package org.springframework.data.gemfire.util;
import static org.assertj.core.api.Assertions.assertThat;
@@ -33,6 +32,7 @@ import static org.springframework.data.gemfire.util.RuntimeExceptionFactory.newR
import java.util.Collections;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
import java.util.function.Supplier;
@@ -48,7 +48,7 @@ import org.springframework.beans.factory.config.RuntimeBeanReference;
import org.springframework.dao.InvalidDataAccessApiUsageException;
/**
* Unit tests for {@link SpringUtils}.
* Unit Tests for {@link SpringUtils}.
*
* @author John Blum
* @see org.junit.Test
@@ -68,38 +68,39 @@ public class SpringUtilsUnitTests {
@Test
public void addDependsOnToExistingDependencies() {
when(mockBeanDefinition.getDependsOn()).thenReturn(asArray("testBeanNameOne", "testBeanNameTwo"));
when(this.mockBeanDefinition.getDependsOn())
.thenReturn(asArray("testBeanNameOne", "testBeanNameTwo"));
assertThat(SpringUtils.addDependsOn(mockBeanDefinition, "testBeanNameThree"))
.isSameAs(mockBeanDefinition);
assertThat(SpringUtils.addDependsOn(this.mockBeanDefinition, "testBeanNameThree"))
.isSameAs(this.mockBeanDefinition);
verify(mockBeanDefinition, times(1)).getDependsOn();
verify(mockBeanDefinition, times(1))
verify(this.mockBeanDefinition, times(1)).getDependsOn();
verify(this.mockBeanDefinition, times(1))
.setDependsOn("testBeanNameOne", "testBeanNameTwo", "testBeanNameThree");
}
@Test
public void addDependsOnToNonExistingDependencies() {
when(mockBeanDefinition.getDependsOn()).thenReturn(null);
when(this.mockBeanDefinition.getDependsOn()).thenReturn(null);
assertThat(SpringUtils.addDependsOn(mockBeanDefinition, "testBeanName"))
.isSameAs(mockBeanDefinition);
assertThat(SpringUtils.addDependsOn(this.mockBeanDefinition, "testBeanName"))
.isSameAs(this.mockBeanDefinition);
verify(mockBeanDefinition, times(1)).getDependsOn();
verify(mockBeanDefinition, times(1)).setDependsOn("testBeanName");
verify(this.mockBeanDefinition, times(1)).getDependsOn();
verify(this.mockBeanDefinition, times(1)).setDependsOn("testBeanName");
}
@Test
public void addDependsOnWithMultipleDependenciesWithExistingDependencies() {
public void addDependsOnWithMultipleDependenciesToExistingDependencies() {
when(mockBeanDefinition.getDependsOn()).thenReturn(asArray("testBeanNameOne", "testBeanNameTwo"));
when(this.mockBeanDefinition.getDependsOn()).thenReturn(asArray("testBeanNameOne", "testBeanNameTwo"));
assertThat(SpringUtils.addDependsOn(mockBeanDefinition, "testBeanNameThree", "testBeanNameFour"))
.isSameAs(mockBeanDefinition);
assertThat(SpringUtils.addDependsOn(this.mockBeanDefinition, "testBeanNameThree", "testBeanNameFour"))
.isSameAs(this.mockBeanDefinition);
verify(mockBeanDefinition, times(1)).getDependsOn();
verify(mockBeanDefinition, times(1))
verify(this.mockBeanDefinition, times(1)).getDependsOn();
verify(this.mockBeanDefinition, times(1))
.setDependsOn("testBeanNameOne", "testBeanNameTwo", "testBeanNameThree", "testBeanNameFour");
}
@@ -109,12 +110,12 @@ public class SpringUtilsUnitTests {
MutablePropertyValues propertyValues =
new MutablePropertyValues(Collections.singletonMap("testProperty", "testValue"));
when(mockBeanDefinition.getPropertyValues()).thenReturn(propertyValues);
when(this.mockBeanDefinition.getPropertyValues()).thenReturn(propertyValues);
assertThat(SpringUtils.getPropertyValue(mockBeanDefinition, "testProperty").orElse(null))
assertThat(SpringUtils.getPropertyValue(this.mockBeanDefinition, "testProperty").orElse(null))
.isEqualTo("testValue");
verify(mockBeanDefinition, times(1)).getPropertyValues();
verify(this.mockBeanDefinition, times(1)).getPropertyValues();
}
@Test
@@ -124,13 +125,13 @@ public class SpringUtilsUnitTests {
PropertyValue testPropertyValue = spy(new PropertyValue("testProperty", null));
when(mockBeanDefinition.getPropertyValues()).thenReturn(testPropertyValues);
when(this.mockBeanDefinition.getPropertyValues()).thenReturn(testPropertyValues);
doReturn(testPropertyValue).when(testPropertyValues).getPropertyValue(anyString());
assertThat(SpringUtils.getPropertyValue(mockBeanDefinition, "testProperty").orElse(null))
assertThat(SpringUtils.getPropertyValue(this.mockBeanDefinition, "testProperty").orElse(null))
.isNull();
verify(mockBeanDefinition, times(1)).getPropertyValues();
verify(this.mockBeanDefinition, times(1)).getPropertyValues();
verify(testPropertyValues, times(1)).getPropertyValue(eq("testProperty"));
verify(testPropertyValue, times(1)).getValue();
}
@@ -140,24 +141,24 @@ public class SpringUtilsUnitTests {
MutablePropertyValues testPropertyValues = spy(new MutablePropertyValues());
when(mockBeanDefinition.getPropertyValues()).thenReturn(testPropertyValues);
when(this.mockBeanDefinition.getPropertyValues()).thenReturn(testPropertyValues);
assertThat(SpringUtils.getPropertyValue(mockBeanDefinition, "testProperty").orElse(null))
assertThat(SpringUtils.getPropertyValue(this.mockBeanDefinition, "testProperty").orElse(null))
.isNull();
verify(mockBeanDefinition, times(1)).getPropertyValues();
verify(this.mockBeanDefinition, times(1)).getPropertyValues();
verify(testPropertyValues, times(1)).getPropertyValue(eq("testProperty"));
}
@Test
public void getPropertyValueWithNullPropertyValuesReturnsNull() {
when(mockBeanDefinition.getPropertyValues()).thenReturn(null);
when(this.mockBeanDefinition.getPropertyValues()).thenReturn(null);
assertThat(SpringUtils.getPropertyValue(mockBeanDefinition, "testProperty").orElse(null))
assertThat(SpringUtils.getPropertyValue(this.mockBeanDefinition, "testProperty").orElse(null))
.isNull();
verify(mockBeanDefinition, times(1)).getPropertyValues();
verify(this.mockBeanDefinition, times(1)).getPropertyValues();
}
@Test
@@ -170,15 +171,13 @@ public class SpringUtilsUnitTests {
@SuppressWarnings("all")
public void setBeanDefinitionPropertyReference() {
BeanDefinition mockBeanDefinition = mock(BeanDefinition.class);
MutablePropertyValues mutablePropertyValues = new MutablePropertyValues();
when(mockBeanDefinition.getPropertyValues()).thenReturn(mutablePropertyValues);
when(this.mockBeanDefinition.getPropertyValues()).thenReturn(mutablePropertyValues);
assertThat(mutablePropertyValues.size()).isEqualTo(0);
SpringUtils.setPropertyReference(mockBeanDefinition, "testProperty", "testBean");
SpringUtils.setPropertyReference(this.mockBeanDefinition, "testProperty", "testBean");
assertThat(mutablePropertyValues.size()).isEqualTo(1);
assertThat(mutablePropertyValues.getPropertyValue("testProperty")).isNotNull();
@@ -192,15 +191,13 @@ public class SpringUtilsUnitTests {
@SuppressWarnings("all")
public void setBeanDefinitionPropertyValue() {
BeanDefinition mockBeanDefinition = mock(BeanDefinition.class);
MutablePropertyValues mutablePropertyValues = new MutablePropertyValues();
when(mockBeanDefinition.getPropertyValues()).thenReturn(mutablePropertyValues);
when(this.mockBeanDefinition.getPropertyValues()).thenReturn(mutablePropertyValues);
assertThat(mutablePropertyValues.size()).isEqualTo(0);
SpringUtils.setPropertyValue(mockBeanDefinition, "testProperty", "testValue");
SpringUtils.setPropertyValue(this.mockBeanDefinition, "testProperty", "testValue");
assertThat(mutablePropertyValues.size()).isEqualTo(1);
assertThat(mutablePropertyValues.getPropertyValue("testProperty")).isNotNull();
@@ -291,6 +288,7 @@ public class SpringUtilsUnitTests {
}
@Test
@SuppressWarnings("all")
public void equalsIgnoreNullIsFalse() {
assertThat(SpringUtils.equalsIgnoreNull(null, "null")).isFalse();
@@ -312,6 +310,7 @@ public class SpringUtilsUnitTests {
}
@Test
@SuppressWarnings("all")
public void nullOrEqualsWithNullIsTrue() {
assertThat(SpringUtils.nullOrEquals(null, "test")).isTrue();
}
@@ -327,6 +326,7 @@ public class SpringUtilsUnitTests {
}
@Test
@SuppressWarnings("all")
public void nullSafeEqualsWithNullObjectsIsFalse() {
assertThat(SpringUtils.nullSafeEquals(null, "test")).isFalse();
assertThat(SpringUtils.nullSafeEquals("test", null)).isFalse();
@@ -337,6 +337,20 @@ public class SpringUtilsUnitTests {
assertThat(SpringUtils.nullSafeEquals("test", "mock")).isFalse();
}
@Test
public void safeDoOperationWithNonThrowingOperation() {
AtomicReference<Object> operationValue = new AtomicReference<>();
assertThat(SpringUtils.safeDoOperation(() -> operationValue.set("TEST"))).isTrue();
assertThat(operationValue.get()).isEqualTo("TEST");
}
@Test
public void safeDoOperationWithThrowingOperation() {
assertThat(SpringUtils.safeDoOperation(() -> { throw new RuntimeException("TEST"); })).isFalse();
}
@Test
public void safeGetValueReturnsSuppliedValue() {
assertThat(SpringUtils.safeGetValue(() -> "test")).isEqualTo("test");

View File

@@ -30,6 +30,11 @@ import static org.mockito.Mockito.when;
import java.util.Arrays;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.Mock;
import org.mockito.junit.MockitoJUnitRunner;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.asyncqueue.AsyncEventListener;
import org.apache.geode.cache.asyncqueue.AsyncEventQueue;
@@ -38,38 +43,49 @@ 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;
/**
* The AsyncEventQueueFactoryBeanTest class is a test suite of test cases testing the contract and functionality
* of the AsyncEventQueueFactoryBean class.
* Unit Tests for {@link AsyncEventQueueFactoryBean}.
*
* @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 org.apache.geode.cache.Cache
* @see org.apache.geode.cache.asyncqueue.AsyncEventListener
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueue
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory
* @see org.apache.geode.cache.wan.GatewayEventFilter
* @see org.apache.geode.cache.wan.GatewayEventSubstitutionFilter
* @see org.apache.geode.cache.wan.GatewaySender
* @see org.springframework.data.gemfire.TestUtils
* @see org.springframework.data.gemfire.wan.AsyncEventQueueFactoryBean
* @since 1.3.3
*/
public class AsyncEventQueueFactoryBeanTest {
@RunWith(MockitoJUnitRunner.class)
public class AsyncEventQueueFactoryBeanUnitTests {
private Cache mockCache() {
return mock(Cache.class);
}
@Mock
private Cache mockCache;
private Cache mockCache(AsyncEventQueueFactory mockAsyncEventQueueFactory) {
Cache mockCache = mockCache();
when(this.mockCache.createAsyncEventQueueFactory()).thenReturn(mockAsyncEventQueueFactory);
when((mockCache.createAsyncEventQueueFactory())).thenReturn(mockAsyncEventQueueFactory);
return this.mockCache;
}
return mockCache;
private AsyncEventListener mockAsyncEventListener() {
return mock(AsyncEventListener.class);
}
private AsyncEventQueue mockAsyncEventQueue(String asyncEventQueueId) {
AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class);
when(mockAsyncEventQueue.getId()).thenReturn(asyncEventQueueId);
return mockAsyncEventQueue;
}
private AsyncEventQueueFactory mockAsyncEventQueueFactory(String asyncEventQueueId) {
@@ -84,23 +100,10 @@ public class AsyncEventQueueFactoryBeanTest {
return mockAsyncEventQueueFactory;
}
private AsyncEventQueue mockAsyncEventQueue(String asyncEventQueueId) {
AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class);
when(mockAsyncEventQueue.getId()).thenReturn(asyncEventQueueId);
return mockAsyncEventQueue;
}
private AsyncEventListener mockAsyncEventListener() {
return mock(AsyncEventListener.class);
}
@Test
public void setAndGetAsyncEventListener() throws Exception {
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(mockCache());
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(this.mockCache);
AsyncEventListener listenerOne = mockAsyncEventListener();
@@ -124,7 +127,7 @@ public class AsyncEventQueueFactoryBeanTest {
AsyncEventQueue mockAsyncEventQueue = mockAsyncEventQueue("testEventQueue");
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(mockCache(), mockAsyncEventListener);
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(this.mockCache, mockAsyncEventListener);
factoryBean.setAsyncEventQueue(mockAsyncEventQueue);
@@ -146,6 +149,36 @@ public class AsyncEventQueueFactoryBeanTest {
}
}
@Test
public void setAndGetAsyncEventQueue() {
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(this.mockCache);
assertThat(factoryBean.getAsyncEventQueue()).isNull();
AsyncEventQueue mockAsyncEventQueue = mockAsyncEventQueue("123");
factoryBean.setAsyncEventQueue(mockAsyncEventQueue);
assertThat(factoryBean.getAsyncEventQueue()).isEqualTo(mockAsyncEventQueue);
}
@Test
public void setAndGetPauseEventDispatching() {
AsyncEventQueueFactoryBean factoryBean = new AsyncEventQueueFactoryBean(this.mockCache);
assertThat(factoryBean.isPauseEventDispatching()).isFalse();
factoryBean.setPauseEventDispatching(true);
assertThat(factoryBean.isPauseEventDispatching()).isTrue();
factoryBean.setPauseEventDispatching(false);
assertThat(factoryBean.isPauseEventDispatching()).isFalse();
}
@Test
public void doInitConfiguresAsyncEventQueue() throws Exception {
@@ -184,14 +217,14 @@ public class AsyncEventQueueFactoryBeanTest {
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)).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));
verify(mockAsyncEventQueueFactory, never()).pauseEventDispatching();
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
@@ -227,6 +260,7 @@ public class AsyncEventQueueFactoryBeanTest {
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true));
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
verify(mockAsyncEventQueueFactory, never()).pauseEventDispatching();
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
@@ -261,6 +295,7 @@ public class AsyncEventQueueFactoryBeanTest {
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true));
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
verify(mockAsyncEventQueueFactory, never()).pauseEventDispatching();
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
@@ -296,6 +331,7 @@ public class AsyncEventQueueFactoryBeanTest {
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false));
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
verify(mockAsyncEventQueueFactory, never()).pauseEventDispatching();
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
@@ -328,6 +364,7 @@ public class AsyncEventQueueFactoryBeanTest {
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false));
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
verify(mockAsyncEventQueueFactory, never()).pauseEventDispatching();
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
@@ -362,6 +399,7 @@ public class AsyncEventQueueFactoryBeanTest {
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false));
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
verify(mockAsyncEventQueueFactory, never()).pauseEventDispatching();
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
@@ -370,10 +408,10 @@ public class AsyncEventQueueFactoryBeanTest {
}
@Test(expected = IllegalStateException.class)
public void doInitWithNullAsyncEventListenerThrowsIllegalStateException() throws Exception {
public void doInitWithNullAsyncEventListenerThrowsIllegalStateException() {
try {
new AsyncEventQueueFactoryBean(mockCache(), null).doInit();
new AsyncEventQueueFactoryBean(this.mockCache, null).doInit();
}
catch (IllegalStateException expected) {
@@ -425,14 +463,16 @@ public class AsyncEventQueueFactoryBeanTest {
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true));
verify(mockAsyncEventQueueFactory, never()).setPersistent(anyBoolean());
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
verify(mockAsyncEventQueueFactory, never()).pauseEventDispatching();
}
}
@Test
public void doInitConfiguresAsyncEventQueueWithSynchronousOverflowDiskStoreNoPersistence() throws Exception {
public void doInitConfiguresAsyncEventQueueWithSynchronousOverflowNonPersistentDiskStorePausingEventDispatching()
throws Exception {
AsyncEventQueueFactory mockAsyncEventQueueFactory =
mockAsyncEventQueueFactory("nonPersistentSynchronousOverflowQueue");
mockAsyncEventQueueFactory("SynchronousOverflowNonPersistentQueue");
AsyncEventQueueFactoryBean factoryBean =
new AsyncEventQueueFactoryBean(mockCache(mockAsyncEventQueueFactory));
@@ -440,9 +480,10 @@ public class AsyncEventQueueFactoryBeanTest {
factoryBean.setAsyncEventListener(mockAsyncEventListener());
factoryBean.setDiskStoreRef("queueOverflowDiskStore");
factoryBean.setDiskSynchronous(true);
factoryBean.setName("nonPersistentSynchronousOverflowQueue");
factoryBean.setName("SynchronousOverflowNonPersistentQueue");
factoryBean.setOrderPolicy(GatewaySender.OrderPolicy.KEY);
factoryBean.setPersistent(false);
factoryBean.setPauseEventDispatching(true);
factoryBean.doInit();
verify(mockAsyncEventQueueFactory, never()).setBatchConflationEnabled(anyBoolean());
@@ -458,10 +499,11 @@ public class AsyncEventQueueFactoryBeanTest {
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false));
verify(mockAsyncEventQueueFactory, times(1)).setPersistent(eq(false));
verify(mockAsyncEventQueueFactory, never()).addGatewayEventFilter(any());
verify(mockAsyncEventQueueFactory, times(1)).pauseEventDispatching();
AsyncEventQueue asyncEventQueue = factoryBean.getObject();
assertThat(asyncEventQueue).isNotNull();
assertThat(asyncEventQueue.getId()).isEqualTo("nonPersistentSynchronousOverflowQueue");
assertThat(asyncEventQueue.getId()).isEqualTo("SynchronousOverflowNonPersistentQueue");
}
}

View File

@@ -12,7 +12,7 @@
http://www.springframework.org/schema/util https://www.springframework.org/schema/util/spring-util.xsd
">
<bean class="org.springframework.data.gemfire.test.GemfireTestBeanPostProcessor"/>
<bean class="org.springframework.data.gemfire.test.mock.config.GemFireMockObjectsBeanPostProcessor"/>
<util:properties id="gemfireProperties">
<prop key="name">AsyncEventQueueNamespaceTest</prop>
@@ -38,10 +38,10 @@
maximum-queue-memory="50"
order-policy="KEY"
parallel="false"
pause-event-dispatching="false"
persistent="true">
<gfe:async-event-listener>
<bean c:name="TestAeqListener"
class="org.springframework.data.gemfire.config.xml.AsyncEventQueueNamespaceTest.TestAsyncEventListener"/>
<bean class="org.springframework.data.gemfire.config.xml.AsyncEventQueueNamespaceTest.TestAsyncEventListener" c:name="TestAeqListener"/>
</gfe:async-event-listener>
</gfe:async-event-queue>
@@ -60,4 +60,10 @@
</gfe:event-substitution-filter>
</gfe:async-event-queue>
<gfe:async-event-queue id="TestPausedAsyncEventQueue" pause-event-dispatching="true">
<gfe:async-event-listener>
<bean class="org.springframework.data.gemfire.config.xml.AsyncEventQueueNamespaceTest.TestAsyncEventListener" c:name="TestAeqListener"/>
</gfe:async-event-listener>
</gfe:async-event-queue>
</beans>