diff --git a/spring-data-geode-test/src/main/java/org/springframework/data/gemfire/tests/mock/GemFireMockObjectsSupport.java b/spring-data-geode-test/src/main/java/org/springframework/data/gemfire/tests/mock/GemFireMockObjectsSupport.java index 535e65b..3362daf 100644 --- a/spring-data-geode-test/src/main/java/org/springframework/data/gemfire/tests/mock/GemFireMockObjectsSupport.java +++ b/spring-data-geode-test/src/main/java/org/springframework/data/gemfire/tests/mock/GemFireMockObjectsSupport.java @@ -246,12 +246,18 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport { private static final List cachedGemFireObjects = Collections.synchronizedList(new ArrayList<>()); + private static final Map asyncEventQueues = new ConcurrentHashMap<>(); + private static final Map diskStores = new ConcurrentHashMap<>(); + //private static final Map gatewaySenders = new ConcurrentHashMap<>(); + private static final Map> regions = new ConcurrentHashMap<>(); private static final Map> regionAttributes = new ConcurrentHashMap<>(); + //private static final Set gatewayReceivers = new ConcurrentSkipListSet<>(); + private static final Set registeredPoolNames = new ConcurrentSkipListSet<>(); private static final String CACHE_FACTORY_DS_PROPS_FIELD_NAME = "dsProps"; @@ -277,7 +283,10 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport { singletonCache.set(null); gemfireProperties.set(new Properties()); + asyncEventQueues.clear(); diskStores.clear(); + //gatewayReceivers.clear(); + //gatewaySenders.clear(); regions.clear(); regionAttributes.clear(); @@ -758,7 +767,12 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport { doAnswer(newSetter(searchTimeout, null)).when(mockCache).setSearchTimeout(anyInt()); when(mockCache.isServer()).thenReturn(true); + when(mockCache.getAsyncEventQueues()) + .thenAnswer(invocation -> Collections.unmodifiableSet(new HashSet<>(asyncEventQueues.values()))); when(mockCache.getCacheServers()).thenAnswer(invocation -> Collections.unmodifiableList(cacheServers)); + //when(mockCache.getGatewayReceivers()).thenAnswer(invocation -> Collections.unmodifiableSet(gatewayReceivers)); + //when(mockCache.getGatewaySenders()) + // .thenAnswer(invocation -> Collections.unmodifiableSet(new HashSet<>(gatewaySenders.values()))); when(mockCache.getLockLease()).thenAnswer(newGetter(lockLease)); when(mockCache.getLockTimeout()).thenAnswer(newGetter(lockTimeout)); when(mockCache.getMessageSyncInterval()).thenAnswer(newGetter(messageSyncInterval)); @@ -776,6 +790,14 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport { when(mockCache.createRegionFactory(anyString())).thenAnswer(invocation -> mockRegionFactory(mockCache, invocation.getArgument(0))); + doAnswer(invocation -> { + + String asyncEventQueueId = invocation.getArgument(0); + + return asyncEventQueues.get(asyncEventQueueId); + + }).when(mockCache).getAsyncEventQueue(anyString()); + return mockQueryService( mockGatewaySenderFactory( mockGatewayReceiverFactory( @@ -898,6 +920,8 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport { when(mockAsyncEventQueue.size()).thenReturn(0); + asyncEventQueues.put(asyncEventQueueId, mockAsyncEventQueue); + return mockAsyncEventQueue; }); @@ -1424,6 +1448,8 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport { doAnswer(newSetter(running, true, null)).when(mockGatewayReceiver).start(); doAnswer(newSetter(running, false, null)).when(mockGatewayReceiver).stop(); + //gatewayReceivers.add(mockGatewayReceiver); + return mockGatewayReceiver; }); @@ -1591,6 +1617,8 @@ public abstract class GemFireMockObjectsSupport extends MockObjectsSupport { }).when(mockGatewaySender).removeGatewayEventFilter(any(GatewayEventFilter.class)); + //gatewaySenders.put(gatewaySenderId, mockGatewaySender); + return mockGatewaySender; });