diff --git a/src/main/java/org/springframework/data/gemfire/tests/mock/MockGemFireObjectsSupport.java b/src/main/java/org/springframework/data/gemfire/tests/mock/MockGemFireObjectsSupport.java index 6d0deea..dba6332 100644 --- a/src/main/java/org/springframework/data/gemfire/tests/mock/MockGemFireObjectsSupport.java +++ b/src/main/java/org/springframework/data/gemfire/tests/mock/MockGemFireObjectsSupport.java @@ -52,6 +52,7 @@ import java.util.Properties; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.CopyOnWriteArraySet; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -82,6 +83,9 @@ import org.apache.geode.cache.RegionService; import org.apache.geode.cache.RegionShortcut; import org.apache.geode.cache.Scope; import org.apache.geode.cache.SubscriptionAttributes; +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.client.ClientCache; import org.apache.geode.cache.client.ClientCacheFactory; import org.apache.geode.cache.client.ClientRegionFactory; @@ -99,6 +103,13 @@ import org.apache.geode.cache.query.QueryService; import org.apache.geode.cache.query.QueryStatistics; import org.apache.geode.cache.server.CacheServer; import org.apache.geode.cache.server.ClientSubscriptionConfig; +import org.apache.geode.cache.wan.GatewayEventFilter; +import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter; +import org.apache.geode.cache.wan.GatewayReceiver; +import org.apache.geode.cache.wan.GatewayReceiverFactory; +import org.apache.geode.cache.wan.GatewaySender; +import org.apache.geode.cache.wan.GatewaySenderFactory; +import org.apache.geode.cache.wan.GatewayTransportFilter; import org.apache.geode.compression.Compressor; import org.apache.geode.distributed.DistributedMember; import org.apache.geode.distributed.DistributedSystem; @@ -282,7 +293,7 @@ public abstract class MockGemFireObjectsSupport extends MockObjectsSupport { * is a root {@link Region}. */ private static boolean isRootRegion(String regionPath) { - return (regionPath.lastIndexOf(Region.SEPARATOR) <= 0); + return regionPath.lastIndexOf(Region.SEPARATOR) <= 0; } /** @@ -298,7 +309,8 @@ public abstract class MockGemFireObjectsSupport extends MockObjectsSupport { regionPath = regionPath.replaceAll(REPEATING_REGION_SEPARATOR, Region.SEPARATOR); regionPath = regionPath.endsWith(Region.SEPARATOR) - ? regionPath.substring(0, regionPath.length() - 1) : regionPath; + ? regionPath.substring(0, regionPath.length() - 1) + : regionPath; return regionPath; } @@ -314,7 +326,7 @@ public abstract class MockGemFireObjectsSupport extends MockObjectsSupport { * @see org.apache.geode.cache.GemFireCache */ private static T rememberMockedGemFireCache(T mockedGemFireCache, - boolean useSingletonCache) { + boolean useSingletonCache) { return Optional.ofNullable(mockedGemFireCache) .map(it -> { @@ -455,8 +467,8 @@ public abstract class MockGemFireObjectsSupport extends MockObjectsSupport { when(mockGemFireCache.listRegionAttributes()).thenReturn(Collections.unmodifiableMap(regionAttributes)); - doThrow(newUnsupportedOperationException(NOT_SUPPORTED)).when(mockGemFireCache) - .loadCacheXml(any(InputStream.class)); + doThrow(newUnsupportedOperationException(NOT_SUPPORTED)) + .when(mockGemFireCache).loadCacheXml(any(InputStream.class)); return mockRegionServiceApi(mockGemFireCache); } @@ -564,7 +576,124 @@ public abstract class MockGemFireObjectsSupport extends MockObjectsSupport { when(mockCache.createRegionFactory(anyString())).thenAnswer(invocation -> mockRegionFactory(mockCache, invocation.getArgument(0))); - return mockQueryService(mockCacheApi(mockCache)); + return mockQueryService( + mockGatewaySenderFactory( + mockGatewayReceiverFactory( + mockAsyncEventQueueFactory( + mockCacheApi(mockCache))))); + } + + public static Cache mockAsyncEventQueueFactory(Cache mockCache) { + + when(mockCache.createAsyncEventQueueFactory()).thenAnswer(invocation -> mockAsyncEventQueueFactory()); + + return mockCache; + } + + public static AsyncEventQueueFactory mockAsyncEventQueueFactory() { + + AsyncEventQueueFactory mockAsyncEventQueueFactory = mock(AsyncEventQueueFactory.class); + + AtomicBoolean batchConflationEnabled = new AtomicBoolean(false); + AtomicBoolean diskSynchronous = new AtomicBoolean(true); + AtomicBoolean forwardExpirationDestroy = new AtomicBoolean(false); + AtomicBoolean parallel = new AtomicBoolean(false); + AtomicBoolean persistent = new AtomicBoolean(false); + + AtomicInteger batchSize = new AtomicInteger(100); + AtomicInteger batchTimeInterval = new AtomicInteger(5); + AtomicInteger dispatcherThreads = new AtomicInteger(5); + AtomicInteger maximumQueueMemory = new AtomicInteger(100); + + AtomicReference diskStoreName = new AtomicReference<>(null); + AtomicReference gatewayEventSubstitutionFilter = + new AtomicReference<>(null); + AtomicReference orderPolicy = new AtomicReference<>(GatewaySender.OrderPolicy.KEY); + + CopyOnWriteArrayList gatewayEventFilters = new CopyOnWriteArrayList<>(); + + when(mockAsyncEventQueueFactory.setBatchConflationEnabled(anyBoolean())) + .thenAnswer(newSetter(batchConflationEnabled, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setBatchSize(anyInt())) + .thenAnswer(newSetter(batchSize, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setBatchTimeInterval(anyInt())) + .thenAnswer(newSetter(batchTimeInterval, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setDiskStoreName(anyString())) + .thenAnswer(newSetter(diskStoreName, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setDiskSynchronous(anyBoolean())) + .thenAnswer(newSetter(diskSynchronous, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setDispatcherThreads(anyInt())) + .thenAnswer(newSetter(dispatcherThreads, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setForwardExpirationDestroy(anyBoolean())) + .thenAnswer(newSetter(forwardExpirationDestroy, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setGatewayEventSubstitutionListener(any(GatewayEventSubstitutionFilter.class))) + .thenAnswer(newSetter(gatewayEventSubstitutionFilter, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setMaximumQueueMemory(anyInt())) + .thenAnswer(newSetter(maximumQueueMemory, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setOrderPolicy(any(GatewaySender.OrderPolicy.class))) + .thenAnswer(newSetter(orderPolicy, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setParallel(anyBoolean())) + .thenAnswer(newSetter(parallel, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.setPersistent(anyBoolean())) + .thenAnswer(newSetter(persistent, mockAsyncEventQueueFactory)); + + when(mockAsyncEventQueueFactory.addGatewayEventFilter(any(GatewayEventFilter.class))).thenAnswer(invocation -> { + + gatewayEventFilters.add(invocation.getArgument(0)); + + return mockAsyncEventQueueFactory; + }); + + when(mockAsyncEventQueueFactory.removeGatewayEventFilter(any(GatewayEventFilter.class))).thenAnswer(invocation -> { + + gatewayEventFilters.remove(invocation.getArgument(0)); + + return mockAsyncEventQueueFactory; + }); + + when(mockAsyncEventQueueFactory.create(anyString(), any(AsyncEventListener.class))).thenAnswer(invocation -> { + + String asyncEventQueueId = invocation.getArgument(0); + + AsyncEventListener listener = invocation.getArgument(1); + + AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class); + + when(mockAsyncEventQueue.isBatchConflationEnabled()).thenAnswer(newGetter(batchConflationEnabled)); + when(mockAsyncEventQueue.isDiskSynchronous()).thenAnswer(newGetter(diskSynchronous)); + when(mockAsyncEventQueue.isForwardExpirationDestroy()).thenAnswer(newGetter(forwardExpirationDestroy)); + when(mockAsyncEventQueue.isParallel()).thenAnswer(newGetter(parallel)); + when(mockAsyncEventQueue.isPersistent()).thenAnswer(newGetter(persistent)); + when(mockAsyncEventQueue.isPrimary()).thenReturn(false); + + when(mockAsyncEventQueue.getAsyncEventListener()).thenReturn(listener); + when(mockAsyncEventQueue.getBatchSize()).thenAnswer(newGetter(batchSize)); + when(mockAsyncEventQueue.getBatchTimeInterval()).thenAnswer(newGetter(batchTimeInterval)); + when(mockAsyncEventQueue.getDiskStoreName()).thenAnswer(newGetter(diskStoreName)); + when(mockAsyncEventQueue.getDispatcherThreads()).thenAnswer(newGetter(dispatcherThreads)); + when(mockAsyncEventQueue.getGatewayEventFilters()).thenReturn(gatewayEventFilters); + when(mockAsyncEventQueue.getGatewayEventSubstitutionFilter()).thenAnswer(newGetter(gatewayEventSubstitutionFilter)); + when(mockAsyncEventQueue.getId()).thenReturn(asyncEventQueueId); + when(mockAsyncEventQueue.getMaximumQueueMemory()).thenAnswer(newGetter(maximumQueueMemory)); + when(mockAsyncEventQueue.getOrderPolicy()).thenAnswer(newGetter(orderPolicy)); + + when(mockAsyncEventQueue.size()).thenReturn(0); + + return mockAsyncEventQueue; + }); + + return mockAsyncEventQueueFactory; } public static CacheServer mockCacheServer() { @@ -985,6 +1114,259 @@ public abstract class MockGemFireObjectsSupport extends MockObjectsSupport { return mockDistributedSystem; } + public static Cache mockGatewayReceiverFactory(Cache mockCache) { + + when(mockCache.createGatewayReceiverFactory()).thenAnswer(invocation -> mockGatewayReceiverFactory()); + + return mockCache; + } + + public static GatewayReceiverFactory mockGatewayReceiverFactory() { + + GatewayReceiverFactory mockGatewayReceiverFactory = mock(GatewayReceiverFactory.class); + + AtomicBoolean manualStart = new AtomicBoolean(GatewayReceiver.DEFAULT_MANUAL_START); + + AtomicInteger endPort = new AtomicInteger(GatewayReceiver.DEFAULT_END_PORT); + AtomicInteger maximumTimeBetweenPings = new AtomicInteger(GatewayReceiver.DEFAULT_MAXIMUM_TIME_BETWEEN_PINGS); + AtomicInteger socketBufferSize = new AtomicInteger(GatewayReceiver.DEFAULT_SOCKET_BUFFER_SIZE); + AtomicInteger startPort = new AtomicInteger(GatewayReceiver.DEFAULT_START_PORT); + + AtomicReference bindAddress = new AtomicReference<>(GatewayReceiver.DEFAULT_BIND_ADDRESS); + AtomicReference hostnameForSenders = new AtomicReference<>(GatewayReceiver.DEFAULT_HOSTNAME_FOR_SENDERS); + + CopyOnWriteArrayList gatewayTransportFilters = new CopyOnWriteArrayList<>(); + + when(mockGatewayReceiverFactory.setBindAddress(anyString())) + .thenAnswer(newSetter(bindAddress, mockGatewayReceiverFactory)); + + when(mockGatewayReceiverFactory.setEndPort(anyInt())) + .thenAnswer(newSetter(endPort, mockGatewayReceiverFactory)); + + when(mockGatewayReceiverFactory.setHostnameForSenders(anyString())) + .thenAnswer(newSetter(hostnameForSenders, mockGatewayReceiverFactory)); + + when(mockGatewayReceiverFactory.setManualStart(anyBoolean())) + .thenAnswer(newSetter(manualStart, mockGatewayReceiverFactory)); + + when(mockGatewayReceiverFactory.setMaximumTimeBetweenPings(anyInt())) + .thenAnswer(newSetter(maximumTimeBetweenPings, mockGatewayReceiverFactory)); + + when(mockGatewayReceiverFactory.setSocketBufferSize(anyInt())) + .thenAnswer(newSetter(socketBufferSize, mockGatewayReceiverFactory)); + + when(mockGatewayReceiverFactory.setStartPort(anyInt())) + .thenAnswer(newSetter(startPort, mockGatewayReceiverFactory)); + + when(mockGatewayReceiverFactory.addGatewayTransportFilter(any(GatewayTransportFilter.class))).thenAnswer(invocation -> { + + gatewayTransportFilters.add(invocation.getArgument(0)); + + return mockGatewayReceiverFactory; + }); + + when(mockGatewayReceiverFactory.removeGatewayTransportFilter(any(GatewayTransportFilter.class))).thenAnswer(invocation -> { + + gatewayTransportFilters.remove(invocation.getArgument(0)); + + return mockGatewayReceiverFactory; + }); + + when(mockGatewayReceiverFactory.create()).thenAnswer(invocation -> { + + AtomicBoolean running = new AtomicBoolean(false); + + GatewayReceiver mockGatewayReceiver = mock(GatewayReceiver.class); + + when(mockGatewayReceiver.isManualStart()).thenAnswer(newGetter(manualStart)); + when(mockGatewayReceiver.isRunning()).thenAnswer(newGetter(running)); + + when(mockGatewayReceiver.getBindAddress()).thenAnswer(newGetter(bindAddress)); + when(mockGatewayReceiver.getEndPort()).thenAnswer(newGetter(endPort)); + when(mockGatewayReceiver.getGatewayTransportFilters()).thenReturn(gatewayTransportFilters); + when(mockGatewayReceiver.getHost()).thenAnswer(newGetter(hostnameForSenders)); + //when(mockGatewayReceiver.getHostnameForSenders()).thenAnswer(newGetter(hostnameForSenders)); + when(mockGatewayReceiver.getMaximumTimeBetweenPings()).thenAnswer(newGetter(maximumTimeBetweenPings)); + when(mockGatewayReceiver.getPort()).thenReturn(0); + when(mockGatewayReceiver.getSocketBufferSize()).thenAnswer(newGetter(socketBufferSize)); + when(mockGatewayReceiver.getStartPort()).thenAnswer(newGetter(startPort)); + + doAnswer(newSetter(running, true, null)).when(mockGatewayReceiver).start(); + doAnswer(newSetter(running, false, null)).when(mockGatewayReceiver).stop(); + + return mockGatewayReceiver; + }); + + return mockGatewayReceiverFactory; + } + + public static Cache mockGatewaySenderFactory(Cache mockCache) { + + when(mockCache.createGatewaySenderFactory()).thenAnswer(invocation -> mockGatewaySenderFactory()); + + return mockCache; + } + + public static GatewaySenderFactory mockGatewaySenderFactory() { + + GatewaySenderFactory mockGatewaySenderFactory = mock(GatewaySenderFactory.class); + + AtomicBoolean batchConflationEnabled = new AtomicBoolean(GatewaySender.DEFAULT_BATCH_CONFLATION); + AtomicBoolean diskSynchronous = new AtomicBoolean(GatewaySender.DEFAULT_DISK_SYNCHRONOUS); + AtomicBoolean parallel = new AtomicBoolean(GatewaySender.DEFAULT_IS_PARALLEL); + AtomicBoolean persistenceEnabled = new AtomicBoolean(GatewaySender.DEFAULT_PERSISTENCE_ENABLED); + + AtomicInteger alertThreshold = new AtomicInteger(GatewaySender.DEFAULT_ALERT_THRESHOLD); + AtomicInteger batchSize = new AtomicInteger(GatewaySender.DEFAULT_BATCH_SIZE); + AtomicInteger batchTimeInterval = new AtomicInteger(GatewaySender.DEFAULT_BATCH_TIME_INTERVAL); + AtomicInteger dispatcherThreads = new AtomicInteger(GatewaySender.DEFAULT_DISPATCHER_THREADS); + AtomicInteger maximumQueueMemory = new AtomicInteger(GatewaySender.DEFAULT_MAXIMUM_QUEUE_MEMORY); + AtomicInteger parallelFactorForReplicatedRegion = + new AtomicInteger(GatewaySender.DEFAULT_PARALLELISM_REPLICATED_REGION); + AtomicInteger socketBufferSize = new AtomicInteger(GatewaySender.DEFAULT_SOCKET_BUFFER_SIZE); + AtomicInteger socketReadTimeout = new AtomicInteger(GatewaySender.DEFAULT_SOCKET_READ_TIMEOUT); + + AtomicReference diskStoreName = new AtomicReference<>(null); + AtomicReference gatewayEventSubstitutionFilter = + new AtomicReference<>(null); + AtomicReference orderPolicy = new AtomicReference<>(GatewaySender.DEFAULT_ORDER_POLICY); + + CopyOnWriteArrayList gatewayEventFilters = new CopyOnWriteArrayList<>(); + CopyOnWriteArrayList gatewayTransportFilters = new CopyOnWriteArrayList<>(); + + when(mockGatewaySenderFactory.setAlertThreshold(anyInt())) + .thenAnswer(newSetter(alertThreshold, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setBatchConflationEnabled(anyBoolean())) + .thenAnswer(newSetter(batchConflationEnabled, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setBatchSize(anyInt())).thenAnswer(newSetter(batchSize, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setBatchTimeInterval(anyInt())) + .thenAnswer(newSetter(batchTimeInterval, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setDiskStoreName(anyString())) + .thenAnswer(newSetter(diskStoreName, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setDiskSynchronous(anyBoolean())) + .thenAnswer(newSetter(diskSynchronous, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setDispatcherThreads(anyInt())) + .thenAnswer(newSetter(dispatcherThreads, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setGatewayEventSubstitutionFilter(any(GatewayEventSubstitutionFilter.class))) + .thenAnswer(newSetter(gatewayEventSubstitutionFilter, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setMaximumQueueMemory(anyInt())) + .thenAnswer(newSetter(maximumQueueMemory, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setOrderPolicy(any(GatewaySender.OrderPolicy.class))) + .thenAnswer(newSetter(orderPolicy, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setParallel(anyBoolean())).thenAnswer(newSetter(parallel, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setParallelFactorForReplicatedRegion(anyInt())) + .thenAnswer(newSetter(parallelFactorForReplicatedRegion, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setPersistenceEnabled(anyBoolean())) + .thenAnswer(newSetter(persistenceEnabled, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setSocketBufferSize(anyInt())) + .thenAnswer(newSetter(socketBufferSize, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.setSocketReadTimeout(anyInt())) + .thenAnswer(newSetter(socketReadTimeout, mockGatewaySenderFactory)); + + when(mockGatewaySenderFactory.addGatewayEventFilter(any(GatewayEventFilter.class))).thenAnswer(invocation -> { + + gatewayEventFilters.add(invocation.getArgument(0)); + + return mockGatewaySenderFactory; + }); + + when(mockGatewaySenderFactory.addGatewayTransportFilter(any(GatewayTransportFilter.class))).thenAnswer(invocation -> { + + gatewayTransportFilters.add(invocation.getArgument(0)); + + return mockGatewaySenderFactory; + }); + + when(mockGatewaySenderFactory.removeGatewayEventFilter(any(GatewayEventFilter.class))).thenAnswer(invocation -> { + + gatewayEventFilters.remove(invocation.getArgument(0)); + + return mockGatewaySenderFactory; + }); + + when(mockGatewaySenderFactory.removeGatewayTransportFilter(any(GatewayTransportFilter.class))).thenAnswer(invocation -> { + + gatewayTransportFilters.remove(invocation.getArgument(0)); + + return mockGatewaySenderFactory; + }); + + when(mockGatewaySenderFactory.create(anyString(), anyInt())).thenAnswer(invocation -> { + + GatewaySender mockGatewaySender = mock(GatewaySender.class); + + AtomicBoolean destroyed = new AtomicBoolean(false); + AtomicBoolean running = new AtomicBoolean(false); + + Integer remoteDistributedSystemId = invocation.getArgument(1); + + String gatewaySenderId = invocation.getArgument(0); + + when(mockGatewaySender.isBatchConflationEnabled()).thenAnswer(newGetter(batchConflationEnabled)); + when(mockGatewaySender.isDiskSynchronous()).thenAnswer(newGetter(diskSynchronous)); + when(mockGatewaySender.isParallel()).thenAnswer(newGetter(parallel)); + when(mockGatewaySender.isPaused()).thenAnswer(newGetter(running)); + when(mockGatewaySender.isPersistenceEnabled()).thenAnswer(newGetter(persistenceEnabled)); + when(mockGatewaySender.isRunning()).thenAnswer(newGetter(running)); + + when(mockGatewaySender.getAlertThreshold()).thenAnswer(newGetter(alertThreshold)); + when(mockGatewaySender.getBatchSize()).thenAnswer(newGetter(batchSize)); + when(mockGatewaySender.getBatchTimeInterval()).thenAnswer(newGetter(batchTimeInterval)); + when(mockGatewaySender.getDiskStoreName()).thenAnswer(newGetter(diskStoreName)); + when(mockGatewaySender.getDispatcherThreads()).thenAnswer(newGetter(dispatcherThreads)); + when(mockGatewaySender.getGatewayEventFilters()).thenReturn(gatewayEventFilters); + when(mockGatewaySender.getGatewayEventSubstitutionFilter()).thenAnswer(newGetter(gatewayEventSubstitutionFilter)); + when(mockGatewaySender.getGatewayTransportFilters()).thenReturn(gatewayTransportFilters); + when(mockGatewaySender.getId()).thenReturn(gatewaySenderId); + when(mockGatewaySender.getMaximumQueueMemory()).thenAnswer(newGetter(maximumQueueMemory)); + when(mockGatewaySender.getMaxParallelismForReplicatedRegion()).thenAnswer(newGetter(parallelFactorForReplicatedRegion)); + when(mockGatewaySender.getOrderPolicy()).thenAnswer(newGetter(orderPolicy)); + when(mockGatewaySender.getRemoteDSId()).thenReturn(remoteDistributedSystemId); + when(mockGatewaySender.getSocketBufferSize()).thenAnswer(newGetter(socketBufferSize)); + when(mockGatewaySender.getSocketReadTimeout()).thenAnswer(newGetter(socketReadTimeout)); + + doAnswer(it -> { + + gatewayEventFilters.add(it.getArgument(0)); + + return null; + + }).when(mockGatewaySender).addGatewayEventFilter(any(GatewayEventFilter.class)); + + doAnswer(newSetter(destroyed, null)).when(mockGatewaySender).destroy(); + doAnswer(newSetter(running, false, null)).when(mockGatewaySender).pause(); + doAnswer(newSetter(running, true, null)).when(mockGatewaySender).resume(); + doAnswer(newSetter(running, true, null)).when(mockGatewaySender).start(); + doAnswer(newSetter(running, false, null)).when(mockGatewaySender).stop(); + + doAnswer(it -> { + + gatewayEventFilters.remove(it.getArgument(0)); + + return null; + + }).when(mockGatewaySender).removeGatewayEventFilter(any(GatewayEventFilter.class)); + + return mockGatewaySender; + }); + + return mockGatewaySenderFactory; + } + public static PoolFactory mockPoolFactory() { PoolFactory mockPoolFactory = mock(PoolFactory.class); diff --git a/src/main/java/org/springframework/data/gemfire/tests/mock/config/MockGemFireObjectsBeanPostProcessor.java b/src/main/java/org/springframework/data/gemfire/tests/mock/config/MockGemFireObjectsBeanPostProcessor.java index 32a12cd..007b1f9 100644 --- a/src/main/java/org/springframework/data/gemfire/tests/mock/config/MockGemFireObjectsBeanPostProcessor.java +++ b/src/main/java/org/springframework/data/gemfire/tests/mock/config/MockGemFireObjectsBeanPostProcessor.java @@ -46,7 +46,7 @@ import org.springframework.lang.Nullable; * @see org.springframework.data.gemfire.CacheFactoryBean * @see org.springframework.data.gemfire.client.ClientCacheFactoryBean * @see org.springframework.data.gemfire.client.PoolFactoryBean - * @see org.springframework.data.gemfire.test.mock.MockGemFireObjectsSupport + * @see org.springframework.data.gemfire.tests.mock.MockGemFireObjectsSupport * @since 2.0.0 */ public class MockGemFireObjectsBeanPostProcessor implements BeanPostProcessor {