Add mock objects for o.a.g.cache.asyncqueue.AsyncEventQueueFactory and o.a.g.cache.asyncqueue.AsyncEventQueue.
Add mock objects for o.a.g.cache.wan.GatewayReceiverFactory and o.a.g.cache.wan.GatewayReceiver. Add mock objects for o.a.g.cache.wan.GatewaySenderFactory and o.a.g.cache.wan.GatewaySender.
This commit is contained in:
@@ -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 extends GemFireCache> 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.<String>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<String> diskStoreName = new AtomicReference<>(null);
|
||||
AtomicReference<GatewayEventSubstitutionFilter> gatewayEventSubstitutionFilter =
|
||||
new AtomicReference<>(null);
|
||||
AtomicReference<GatewaySender.OrderPolicy> orderPolicy = new AtomicReference<>(GatewaySender.OrderPolicy.KEY);
|
||||
|
||||
CopyOnWriteArrayList<GatewayEventFilter> 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.<GatewayEventFilter>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<String> bindAddress = new AtomicReference<>(GatewayReceiver.DEFAULT_BIND_ADDRESS);
|
||||
AtomicReference<String> hostnameForSenders = new AtomicReference<>(GatewayReceiver.DEFAULT_HOSTNAME_FOR_SENDERS);
|
||||
|
||||
CopyOnWriteArrayList<GatewayTransportFilter> 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.<GatewayTransportFilter>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<String> diskStoreName = new AtomicReference<>(null);
|
||||
AtomicReference<GatewayEventSubstitutionFilter> gatewayEventSubstitutionFilter =
|
||||
new AtomicReference<>(null);
|
||||
AtomicReference<GatewaySender.OrderPolicy> orderPolicy = new AtomicReference<>(GatewaySender.DEFAULT_ORDER_POLICY);
|
||||
|
||||
CopyOnWriteArrayList<GatewayEventFilter> gatewayEventFilters = new CopyOnWriteArrayList<>();
|
||||
CopyOnWriteArrayList<GatewayTransportFilter> 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.<GatewayEventFilter>getArgument(0));
|
||||
|
||||
return mockGatewaySenderFactory;
|
||||
});
|
||||
|
||||
when(mockGatewaySenderFactory.removeGatewayTransportFilter(any(GatewayTransportFilter.class))).thenAnswer(invocation -> {
|
||||
|
||||
gatewayTransportFilters.remove(invocation.<GatewayTransportFilter>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.<GatewayEventFilter>getArgument(0));
|
||||
|
||||
return null;
|
||||
|
||||
}).when(mockGatewaySender).removeGatewayEventFilter(any(GatewayEventFilter.class));
|
||||
|
||||
return mockGatewaySender;
|
||||
});
|
||||
|
||||
return mockGatewaySenderFactory;
|
||||
}
|
||||
|
||||
public static PoolFactory mockPoolFactory() {
|
||||
|
||||
PoolFactory mockPoolFactory = mock(PoolFactory.class);
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user