Add protected API calls to construct the AEQ, the factory and the listener.
Add builder methods to configure Function post-processors applied to the AEQ, factory and listener. Add API to configure an AsyncEventErrorHandler. Edit Javadoc. Resolves gh-58.
This commit is contained in:
@@ -20,11 +20,13 @@ import java.util.List;
|
|||||||
import java.util.Objects;
|
import java.util.Objects;
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
import java.util.UUID;
|
import java.util.UUID;
|
||||||
|
import java.util.function.Function;
|
||||||
import java.util.function.Predicate;
|
import java.util.function.Predicate;
|
||||||
|
|
||||||
import org.apache.geode.cache.Cache;
|
import org.apache.geode.cache.Cache;
|
||||||
import org.apache.geode.cache.DiskStore;
|
import org.apache.geode.cache.DiskStore;
|
||||||
import org.apache.geode.cache.Region;
|
import org.apache.geode.cache.Region;
|
||||||
|
import org.apache.geode.cache.asyncqueue.AsyncEvent;
|
||||||
import org.apache.geode.cache.asyncqueue.AsyncEventListener;
|
import org.apache.geode.cache.asyncqueue.AsyncEventListener;
|
||||||
import org.apache.geode.cache.asyncqueue.AsyncEventQueue;
|
import org.apache.geode.cache.asyncqueue.AsyncEventQueue;
|
||||||
import org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory;
|
import org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory;
|
||||||
@@ -38,6 +40,7 @@ import org.springframework.data.gemfire.config.annotation.RegionConfigurer;
|
|||||||
import org.springframework.data.gemfire.util.ArrayUtils;
|
import org.springframework.data.gemfire.util.ArrayUtils;
|
||||||
import org.springframework.data.gemfire.util.CollectionUtils;
|
import org.springframework.data.gemfire.util.CollectionUtils;
|
||||||
import org.springframework.data.repository.CrudRepository;
|
import org.springframework.data.repository.CrudRepository;
|
||||||
|
import org.springframework.geode.cache.RepositoryAsyncEventListener.AsyncEventErrorHandler;
|
||||||
import org.springframework.lang.NonNull;
|
import org.springframework.lang.NonNull;
|
||||||
import org.springframework.lang.Nullable;
|
import org.springframework.lang.Nullable;
|
||||||
import org.springframework.util.Assert;
|
import org.springframework.util.Assert;
|
||||||
@@ -49,6 +52,7 @@ import org.springframework.util.StringUtils;
|
|||||||
* abstraction.
|
* abstraction.
|
||||||
*
|
*
|
||||||
* @author John Blum
|
* @author John Blum
|
||||||
|
* @see java.util.function.Function
|
||||||
* @see java.util.function.Predicate
|
* @see java.util.function.Predicate
|
||||||
* @see org.apache.geode.cache.Cache
|
* @see org.apache.geode.cache.Cache
|
||||||
* @see org.apache.geode.cache.Region
|
* @see org.apache.geode.cache.Region
|
||||||
@@ -61,10 +65,13 @@ import org.springframework.util.StringUtils;
|
|||||||
* @see org.springframework.data.gemfire.PeerRegionFactoryBean
|
* @see org.springframework.data.gemfire.PeerRegionFactoryBean
|
||||||
* @see org.springframework.data.gemfire.config.annotation.RegionConfigurer
|
* @see org.springframework.data.gemfire.config.annotation.RegionConfigurer
|
||||||
* @see org.springframework.data.repository.CrudRepository
|
* @see org.springframework.data.repository.CrudRepository
|
||||||
|
* @see org.springframework.geode.cache.RepositoryAsyncEventListener.AsyncEventErrorHandler
|
||||||
* @since 1.4.0
|
* @since 1.4.0
|
||||||
*/
|
*/
|
||||||
public class AsyncInlineCachingRegionConfigurer<T, ID> implements RegionConfigurer {
|
public class AsyncInlineCachingRegionConfigurer<T, ID> implements RegionConfigurer {
|
||||||
|
|
||||||
|
protected static final Predicate<String> DEFAULT_REGION_BEAN_NAME_PREDICATE = beanName -> false;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Factory method used to construct a new instance of {@link AsyncInlineCachingRegionConfigurer} initialized with
|
* Factory method used to construct a new instance of {@link AsyncInlineCachingRegionConfigurer} initialized with
|
||||||
* the given Spring Data {@link CrudRepository} and {@link Predicate} identifying the target {@link Region}
|
* the given Spring Data {@link CrudRepository} and {@link Predicate} identifying the target {@link Region}
|
||||||
@@ -101,16 +108,18 @@ public class AsyncInlineCachingRegionConfigurer<T, ID> implements RegionConfigur
|
|||||||
* {@literal Asynchronous Inline Caching} will be configured.
|
* {@literal Asynchronous Inline Caching} will be configured.
|
||||||
* @return a new {@link AsyncInlineCachingRegionConfigurer}.
|
* @return a new {@link AsyncInlineCachingRegionConfigurer}.
|
||||||
* @throws IllegalArgumentException if {@link CrudRepository} is {@literal null}.
|
* @throws IllegalArgumentException if {@link CrudRepository} is {@literal null}.
|
||||||
* @see #AsyncInlineCachingRegionConfigurer(CrudRepository, Predicate)
|
|
||||||
* @see org.springframework.data.repository.CrudRepository
|
* @see org.springframework.data.repository.CrudRepository
|
||||||
|
* @see #create(CrudRepository, Predicate)
|
||||||
* @see java.lang.String
|
* @see java.lang.String
|
||||||
*/
|
*/
|
||||||
public static <T, ID> AsyncInlineCachingRegionConfigurer<T, ID> create(@NonNull CrudRepository<T, ID> repository,
|
public static <T, ID> AsyncInlineCachingRegionConfigurer<T, ID> create(@NonNull CrudRepository<T, ID> repository,
|
||||||
@Nullable String regionBeanName) {
|
@Nullable String regionBeanName) {
|
||||||
|
|
||||||
return new AsyncInlineCachingRegionConfigurer<>(repository, Predicate.isEqual(regionBeanName));
|
return create(repository, Predicate.isEqual(regionBeanName));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private AsyncEventErrorHandler asyncEventErrorHandler;
|
||||||
|
|
||||||
private Boolean batchConflationEnabled;
|
private Boolean batchConflationEnabled;
|
||||||
private Boolean diskSynchronous;
|
private Boolean diskSynchronous;
|
||||||
private Boolean forwardExpirationDestroy;
|
private Boolean forwardExpirationDestroy;
|
||||||
@@ -120,6 +129,12 @@ public class AsyncInlineCachingRegionConfigurer<T, ID> implements RegionConfigur
|
|||||||
|
|
||||||
private final CrudRepository<T, ID> repository;
|
private final CrudRepository<T, ID> repository;
|
||||||
|
|
||||||
|
private Function<AsyncEventListener, AsyncEventListener> asyncEventListenerPostProcessor;
|
||||||
|
|
||||||
|
private Function<AsyncEventQueue, AsyncEventQueue> asyncEventQueuePostProcessor;
|
||||||
|
|
||||||
|
private Function<AsyncEventQueueFactory, AsyncEventQueueFactory> asyncEventQueueFactoryPostProcessor;
|
||||||
|
|
||||||
private Integer batchSize;
|
private Integer batchSize;
|
||||||
private Integer batchTimeInterval;
|
private Integer batchTimeInterval;
|
||||||
private Integer dispatcherThreads;
|
private Integer dispatcherThreads;
|
||||||
@@ -155,7 +170,7 @@ public class AsyncInlineCachingRegionConfigurer<T, ID> implements RegionConfigur
|
|||||||
Assert.notNull(repository, "CrudRepository must not be null");
|
Assert.notNull(repository, "CrudRepository must not be null");
|
||||||
|
|
||||||
this.repository = repository;
|
this.repository = repository;
|
||||||
this.regionBeanName = regionBeanName != null ? regionBeanName : beanName -> false;
|
this.regionBeanName = regionBeanName != null ? regionBeanName : DEFAULT_REGION_BEAN_NAME_PREDICATE;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -195,6 +210,7 @@ public class AsyncInlineCachingRegionConfigurer<T, ID> implements RegionConfigur
|
|||||||
* @param bean {@link PeerRegionFactoryBean} containing the configuration of the target {@link Region}
|
* @param bean {@link PeerRegionFactoryBean} containing the configuration of the target {@link Region}
|
||||||
* in the Spring container.
|
* in the Spring container.
|
||||||
* @see org.springframework.data.gemfire.PeerRegionFactoryBean
|
* @see org.springframework.data.gemfire.PeerRegionFactoryBean
|
||||||
|
* @see #newAsyncEventQueue(Cache, String)
|
||||||
*/
|
*/
|
||||||
@Override
|
@Override
|
||||||
public void configure(String beanName, PeerRegionFactoryBean<?, ?> bean) {
|
public void configure(String beanName, PeerRegionFactoryBean<?, ?> bean) {
|
||||||
@@ -230,12 +246,17 @@ public class AsyncInlineCachingRegionConfigurer<T, ID> implements RegionConfigur
|
|||||||
* @return a new {@link AsyncEventQueue}.
|
* @return a new {@link AsyncEventQueue}.
|
||||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueue
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueue
|
||||||
* @see org.apache.geode.cache.Cache
|
* @see org.apache.geode.cache.Cache
|
||||||
* @see #newRepositoryAsyncEventListener()
|
|
||||||
* @see #generateId(String)
|
* @see #generateId(String)
|
||||||
|
* @see #newAsyncEventQueueFactory(Cache)
|
||||||
|
* @see #newAsyncEventQueue(AsyncEventQueueFactory, String, AsyncEventListener)
|
||||||
|
* @see #newRepositoryAsyncEventListener()
|
||||||
|
* @see #postProcess(AsyncEventListener)
|
||||||
|
* @see #postProcess(AsyncEventQueue)
|
||||||
|
* @see #postProcess(AsyncEventQueueFactory)
|
||||||
*/
|
*/
|
||||||
protected AsyncEventQueue newAsyncEventQueue(@NonNull Cache peerCache, @NonNull String regionBeanName) {
|
protected AsyncEventQueue newAsyncEventQueue(@NonNull Cache peerCache, @NonNull String regionBeanName) {
|
||||||
|
|
||||||
AsyncEventQueueFactory asyncEventQueueFactory = peerCache.createAsyncEventQueueFactory();
|
AsyncEventQueueFactory asyncEventQueueFactory = newAsyncEventQueueFactory(peerCache);
|
||||||
|
|
||||||
Optional.ofNullable(this.batchConflationEnabled).ifPresent(asyncEventQueueFactory::setBatchConflationEnabled);
|
Optional.ofNullable(this.batchConflationEnabled).ifPresent(asyncEventQueueFactory::setBatchConflationEnabled);
|
||||||
Optional.ofNullable(this.batchSize).ifPresent(asyncEventQueueFactory::setBatchSize);
|
Optional.ofNullable(this.batchSize).ifPresent(asyncEventQueueFactory::setBatchSize);
|
||||||
@@ -258,20 +279,236 @@ public class AsyncInlineCachingRegionConfigurer<T, ID> implements RegionConfigur
|
|||||||
asyncEventQueueFactory.pauseEventDispatching();
|
asyncEventQueueFactory.pauseEventDispatching();
|
||||||
}
|
}
|
||||||
|
|
||||||
return asyncEventQueueFactory.create(generateId(regionBeanName), newRepositoryAsyncEventListener());
|
String asyncEventQueueId = generateId(regionBeanName);
|
||||||
|
|
||||||
|
AsyncEventListener asyncEventListener = newRepositoryAsyncEventListener();
|
||||||
|
|
||||||
|
asyncEventListener = postProcess(asyncEventListener);
|
||||||
|
asyncEventQueueFactory = postProcess(asyncEventQueueFactory);
|
||||||
|
|
||||||
|
AsyncEventQueue asyncEventQueue =
|
||||||
|
newAsyncEventQueue(asyncEventQueueFactory, asyncEventQueueId, asyncEventListener);
|
||||||
|
|
||||||
|
asyncEventQueue = postProcess(asyncEventQueue);
|
||||||
|
|
||||||
|
return asyncEventQueue;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Constructs (creates) a new instance of {@link AsyncEventQueue} using the given {@link AsyncEventQueueFactory}
|
||||||
|
* with the given {@link String AEQ ID} and {@link AsyncEventListener}.
|
||||||
|
*
|
||||||
|
* @param factory {@link AsyncEventQueueFactory} used to create the {@link AsyncEventQueue};
|
||||||
|
* must not be {@literal null}.
|
||||||
|
* @param asyncEventQueueId {@link String} containing the {@literal ID} for the {@link AsyncEventQueue};
|
||||||
|
* must not be {@literal null}.
|
||||||
|
* @param listener {@link AsyncEventListener} registered with the {@link AsyncEventQueue} to process cache events
|
||||||
|
* from the {@link AsyncEventQueue} attached to the {@link Region}.
|
||||||
|
* @return a new {@link AsyncEventQueue} with the {@link String ID} and registered {@link AsyncEventListener}.
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueue
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventListener
|
||||||
|
*/
|
||||||
|
protected @NonNull AsyncEventQueue newAsyncEventQueue(@NonNull AsyncEventQueueFactory factory,
|
||||||
|
@NonNull String asyncEventQueueId, @NonNull AsyncEventListener listener) {
|
||||||
|
|
||||||
|
return factory.create(asyncEventQueueId, listener);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Constructs (creates) a new instance of the {@link AsyncEventQueueFactory} from the given {@literal peer}
|
||||||
|
* {@link Cache}.
|
||||||
|
*
|
||||||
|
* @param peerCache {@literal Peer} {@link Cache} instance used to create an instance of
|
||||||
|
* the {@link AsyncEventQueueFactory}; must not be {@literal null}.
|
||||||
|
* @return a new instance of {@link AsyncEventQueueFactory} to create a {@link AsyncEventQueue}.
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory
|
||||||
|
* @see org.apache.geode.cache.Cache
|
||||||
|
*/
|
||||||
|
protected @NonNull AsyncEventQueueFactory newAsyncEventQueueFactory(@NonNull Cache peerCache) {
|
||||||
|
return peerCache.createAsyncEventQueueFactory();
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Constructs a new Apache Geode {@link AsyncEventListener} to register on an {@link AsyncEventQueue} attached to
|
* Constructs a new Apache Geode {@link AsyncEventListener} to register on an {@link AsyncEventQueue} attached to
|
||||||
* the target {@link Region}, which uses the {@link CrudRepository} to perform data access operations on an external
|
* the target {@link Region}, which uses the {@link CrudRepository} to perform data access operations on an external
|
||||||
* data source asynchronously when cache events and operations occur on the target {@link Region}.
|
* backend data source asynchronously when cache events and operations occur on the target {@link Region}.
|
||||||
*
|
*
|
||||||
* @return a new {@link RepositoryAsyncEventListener}.
|
* @return a new {@link RepositoryAsyncEventListener}.
|
||||||
* @see org.springframework.geode.cache.RepositoryAsyncEventListener
|
* @see org.springframework.geode.cache.RepositoryAsyncEventListener
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventListener
|
||||||
|
* @see #newRepositoryAsyncEventListener(CrudRepository)
|
||||||
* @see #getRepository()
|
* @see #getRepository()
|
||||||
*/
|
*/
|
||||||
protected @NonNull AsyncEventListener newRepositoryAsyncEventListener() {
|
protected @NonNull AsyncEventListener newRepositoryAsyncEventListener() {
|
||||||
return new RepositoryAsyncEventListener<>(getRepository());
|
return newRepositoryAsyncEventListener(getRepository());
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Constructs a new Apache Geode {@link AsyncEventListener} to register on an {@link AsyncEventQueue} attached to
|
||||||
|
* the target {@link Region}, which uses the given {@link CrudRepository} to perform data access operations on an
|
||||||
|
* external, backend data source asynchronously when cache events and operations occur on the target {@link Region}.
|
||||||
|
*
|
||||||
|
* @param repository Spring Data {@link CrudRepository} used to perform data access operations on the external,
|
||||||
|
* backend data source; must not be {@literal null}.
|
||||||
|
* @return a new {@link RepositoryAsyncEventListener}.
|
||||||
|
* @see org.springframework.geode.cache.RepositoryAsyncEventListener
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventListener
|
||||||
|
* @see #getRepository()
|
||||||
|
*/
|
||||||
|
protected @NonNull AsyncEventListener newRepositoryAsyncEventListener(@NonNull CrudRepository<T, ID> repository) {
|
||||||
|
return new RepositoryAsyncEventListener<>(repository);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Applies the user-defined {@link Function} to the framework constructed/provided {@link AsyncEventListener}
|
||||||
|
* for post processing.
|
||||||
|
*
|
||||||
|
* @param asyncEventListener {@link AsyncEventListener} constructed by the framework and post processed by
|
||||||
|
* end-user code encapsulated in the {@link #applyToListener(Function) configured} {@link Function}.
|
||||||
|
* @return the post-processed {@link AsyncEventListener}.
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventListener
|
||||||
|
* @see #applyToListener(Function)
|
||||||
|
*/
|
||||||
|
protected @NonNull AsyncEventListener postProcess(@NonNull AsyncEventListener asyncEventListener) {
|
||||||
|
return resolveAsyncEventListenerPostProcessor().apply(asyncEventListener);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Applies the user-defined {@link Function} to the framework constructed/provided {@link AsyncEventQueue}
|
||||||
|
* for post processing.
|
||||||
|
*
|
||||||
|
* @param asyncEventQueue {@link AsyncEventQueue} constructed by the framework and post processed by
|
||||||
|
* end-user code encapsulated in the {@link #applyToQueue(Function) configured} {@link Function}.
|
||||||
|
* @return the post-processed {@link AsyncEventQueue}.
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueue
|
||||||
|
* @see #applyToQueue(Function)
|
||||||
|
*/
|
||||||
|
protected @NonNull AsyncEventQueue postProcess(@NonNull AsyncEventQueue asyncEventQueue) {
|
||||||
|
|
||||||
|
Function<AsyncEventQueue, AsyncEventQueue> asyncEventQueuePostProcessor =
|
||||||
|
this.asyncEventQueuePostProcessor;
|
||||||
|
|
||||||
|
return asyncEventQueuePostProcessor != null
|
||||||
|
? asyncEventQueuePostProcessor.apply(asyncEventQueue)
|
||||||
|
: asyncEventQueue;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Applies the user-defined {@link Function} to the framework constructed/provided {@link AsyncEventQueueFactory}
|
||||||
|
* for post processing.
|
||||||
|
*
|
||||||
|
* @param asyncEventQueueFactory {@link AsyncEventQueueFactory} constructed by the framework and post processed by
|
||||||
|
* end-user code encapsulated in the {@link #applyToQueueFactory(Function) configured} {@link Function}.
|
||||||
|
* @return the post-processed {@link AsyncEventQueueFactory}.
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory
|
||||||
|
* @see #applyToQueueFactory(Function)
|
||||||
|
*/
|
||||||
|
protected @NonNull AsyncEventQueueFactory postProcess(@NonNull AsyncEventQueueFactory asyncEventQueueFactory) {
|
||||||
|
|
||||||
|
Function<AsyncEventQueueFactory, AsyncEventQueueFactory> asyncEventQueueFactoryPostProcessor =
|
||||||
|
this.asyncEventQueueFactoryPostProcessor;
|
||||||
|
|
||||||
|
return asyncEventQueueFactoryPostProcessor != null
|
||||||
|
? asyncEventQueueFactoryPostProcessor.apply(asyncEventQueueFactory)
|
||||||
|
: asyncEventQueueFactory;
|
||||||
|
}
|
||||||
|
|
||||||
|
@SuppressWarnings("unchecked")
|
||||||
|
private @NonNull Function<AsyncEventListener, AsyncEventListener> resolveAsyncEventListenerPostProcessor() {
|
||||||
|
|
||||||
|
AsyncEventErrorHandler asyncEventErrorHandler = this.asyncEventErrorHandler;
|
||||||
|
|
||||||
|
Function<AsyncEventListener, AsyncEventListener> resolvedListenerPostProcessor = asyncEventErrorHandler != null
|
||||||
|
? listener -> {
|
||||||
|
|
||||||
|
if (listener instanceof RepositoryAsyncEventListener) {
|
||||||
|
((RepositoryAsyncEventListener<T, ID>) listener).setAsyncEventErrorHandler(asyncEventErrorHandler);
|
||||||
|
}
|
||||||
|
|
||||||
|
return listener;
|
||||||
|
}
|
||||||
|
: Function.identity();
|
||||||
|
|
||||||
|
Function<AsyncEventListener, AsyncEventListener> asyncEventListenerPostProcessor =
|
||||||
|
this.asyncEventListenerPostProcessor;
|
||||||
|
|
||||||
|
if (asyncEventListenerPostProcessor != null) {
|
||||||
|
resolvedListenerPostProcessor = resolvedListenerPostProcessor.andThen(asyncEventListenerPostProcessor);
|
||||||
|
}
|
||||||
|
|
||||||
|
return resolvedListenerPostProcessor;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Builder method used to configure the given user-defined {@link Function} applied to the framework constructed
|
||||||
|
* and provided {@link AsyncEventListener} for post processing.
|
||||||
|
*
|
||||||
|
* @param asyncEventListenerPostProcessor user-defined {@link Function} encapsulating the logic applied to
|
||||||
|
* the framework constructed/provided {@link AsyncEventListener} for post-processing.
|
||||||
|
* @return this {@link AsyncInlineCachingRegionConfigurer}.
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventListener
|
||||||
|
* @see java.util.function.Function
|
||||||
|
*/
|
||||||
|
public AsyncInlineCachingRegionConfigurer<T, ID> applyToListener(
|
||||||
|
@Nullable Function<AsyncEventListener, AsyncEventListener> asyncEventListenerPostProcessor) {
|
||||||
|
|
||||||
|
this.asyncEventListenerPostProcessor = asyncEventListenerPostProcessor;
|
||||||
|
|
||||||
|
return this;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Builder method used to configure the given user-defined {@link Function} applied to the framework constructed
|
||||||
|
* and provided {@link AsyncEventQueue} for post processing.
|
||||||
|
*
|
||||||
|
* @param asyncEventQueuePostProcessor user-defined {@link Function} encapsulating the logic applied to
|
||||||
|
* the framework constructed {@link AsyncEventQueue} for post-processing.
|
||||||
|
* @return this {@link AsyncInlineCachingRegionConfigurer}.
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueue
|
||||||
|
* @see java.util.function.Function
|
||||||
|
*/
|
||||||
|
public AsyncInlineCachingRegionConfigurer<T, ID> applyToQueue(
|
||||||
|
@Nullable Function<AsyncEventQueue, AsyncEventQueue> asyncEventQueuePostProcessor) {
|
||||||
|
|
||||||
|
this.asyncEventQueuePostProcessor = asyncEventQueuePostProcessor;
|
||||||
|
|
||||||
|
return this;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Builder method used to configure the given user-defined {@link Function} applied to the framework constructed
|
||||||
|
* and provided {@link AsyncEventQueueFactory} for post processing.
|
||||||
|
*
|
||||||
|
* @param asyncEventQueueFactoryPostProcessor user-defined {@link Function} encapsulating the logic applied to
|
||||||
|
* the framework constructed {@link AsyncEventQueueFactory} for post-processing.
|
||||||
|
* @return this {@link AsyncInlineCachingRegionConfigurer}.
|
||||||
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory
|
||||||
|
* @see java.util.function.Function
|
||||||
|
*/
|
||||||
|
public AsyncInlineCachingRegionConfigurer<T, ID> applyToQueueFactory(
|
||||||
|
@Nullable Function<AsyncEventQueueFactory, AsyncEventQueueFactory> asyncEventQueueFactoryPostProcessor) {
|
||||||
|
|
||||||
|
this.asyncEventQueueFactoryPostProcessor = asyncEventQueueFactoryPostProcessor;
|
||||||
|
|
||||||
|
return this;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Builder method used to configure a {@link AsyncEventErrorHandler} to handle errors thrown while processing
|
||||||
|
* {@link AsyncEvent AsyncEvents} in the {@link AsyncEventListener}.
|
||||||
|
*
|
||||||
|
* @param errorHandler {@link AsyncEventErrorHandler} used to handle errors thrown while processing
|
||||||
|
* {@link AsyncEvent AsyncEvents} in the {@link AsyncEventListener}.
|
||||||
|
* @return this {@link AsyncInlineCachingRegionConfigurer}.
|
||||||
|
* @see org.springframework.geode.cache.RepositoryAsyncEventListener.AsyncEventErrorHandler
|
||||||
|
*/
|
||||||
|
public AsyncInlineCachingRegionConfigurer<T, ID> withAsyncEventErrorHandler(
|
||||||
|
@Nullable AsyncEventErrorHandler errorHandler) {
|
||||||
|
|
||||||
|
this.asyncEventErrorHandler = errorHandler;
|
||||||
|
|
||||||
|
return this;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -18,6 +18,7 @@ package org.springframework.geode.cache;
|
|||||||
import static org.assertj.core.api.Assertions.assertThat;
|
import static org.assertj.core.api.Assertions.assertThat;
|
||||||
import static org.mockito.ArgumentMatchers.eq;
|
import static org.mockito.ArgumentMatchers.eq;
|
||||||
import static org.mockito.Mockito.doReturn;
|
import static org.mockito.Mockito.doReturn;
|
||||||
|
import static org.mockito.Mockito.inOrder;
|
||||||
import static org.mockito.Mockito.mock;
|
import static org.mockito.Mockito.mock;
|
||||||
import static org.mockito.Mockito.spy;
|
import static org.mockito.Mockito.spy;
|
||||||
import static org.mockito.Mockito.times;
|
import static org.mockito.Mockito.times;
|
||||||
@@ -27,9 +28,11 @@ import static org.mockito.Mockito.verifyNoMoreInteractions;
|
|||||||
|
|
||||||
import java.time.Duration;
|
import java.time.Duration;
|
||||||
import java.util.Arrays;
|
import java.util.Arrays;
|
||||||
|
import java.util.function.Function;
|
||||||
import java.util.function.Predicate;
|
import java.util.function.Predicate;
|
||||||
|
|
||||||
import org.junit.Test;
|
import org.junit.Test;
|
||||||
|
import org.mockito.InOrder;
|
||||||
|
|
||||||
import org.apache.geode.cache.Cache;
|
import org.apache.geode.cache.Cache;
|
||||||
import org.apache.geode.cache.asyncqueue.AsyncEventListener;
|
import org.apache.geode.cache.asyncqueue.AsyncEventListener;
|
||||||
@@ -42,6 +45,7 @@ import org.apache.geode.cache.wan.GatewaySender;
|
|||||||
import org.springframework.data.gemfire.PeerRegionFactoryBean;
|
import org.springframework.data.gemfire.PeerRegionFactoryBean;
|
||||||
import org.springframework.data.gemfire.util.ArrayUtils;
|
import org.springframework.data.gemfire.util.ArrayUtils;
|
||||||
import org.springframework.data.repository.CrudRepository;
|
import org.springframework.data.repository.CrudRepository;
|
||||||
|
import org.springframework.geode.cache.RepositoryAsyncEventListener.AsyncEventErrorHandler;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Unit Tests for {@link AsyncInlineCachingRegionConfigurer}.
|
* Unit Tests for {@link AsyncInlineCachingRegionConfigurer}.
|
||||||
@@ -54,9 +58,12 @@ import org.springframework.data.repository.CrudRepository;
|
|||||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventListener
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventListener
|
||||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueue
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueue
|
||||||
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory
|
* @see org.apache.geode.cache.asyncqueue.AsyncEventQueueFactory
|
||||||
|
* @see org.apache.geode.cache.wan.GatewayEventFilter
|
||||||
|
* @see org.apache.geode.cache.wan.GatewayEventSubstitutionFilter
|
||||||
* @see org.springframework.data.gemfire.PeerRegionFactoryBean
|
* @see org.springframework.data.gemfire.PeerRegionFactoryBean
|
||||||
* @see org.springframework.data.repository.CrudRepository
|
* @see org.springframework.data.repository.CrudRepository
|
||||||
* @see org.springframework.geode.cache.AsyncInlineCachingRegionConfigurer
|
* @see org.springframework.geode.cache.AsyncInlineCachingRegionConfigurer
|
||||||
|
* @see org.springframework.geode.cache.RepositoryAsyncEventListener.AsyncEventErrorHandler
|
||||||
* @since 1.4.0
|
* @since 1.4.0
|
||||||
*/
|
*/
|
||||||
public class AsyncInlineCachingRegionConfigurerUnitTests {
|
public class AsyncInlineCachingRegionConfigurerUnitTests {
|
||||||
@@ -66,13 +73,13 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
|
|
||||||
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
||||||
|
|
||||||
Predicate<String> regionBeanName = Predicate.isEqual("TestRegion");
|
Predicate<String> testRegionBeanName = Predicate.isEqual("TestRegion");
|
||||||
|
|
||||||
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
||||||
new AsyncInlineCachingRegionConfigurer<>(mockRepository, regionBeanName);
|
new AsyncInlineCachingRegionConfigurer<>(mockRepository, testRegionBeanName);
|
||||||
|
|
||||||
assertThat(regionConfigurer).isNotNull();
|
assertThat(regionConfigurer).isNotNull();
|
||||||
assertThat(regionConfigurer.getRegionBeanName()).isEqualTo(regionBeanName);
|
assertThat(regionConfigurer.getRegionBeanName()).isEqualTo(testRegionBeanName);
|
||||||
assertThat(regionConfigurer.getRepository()).isEqualTo(mockRepository);
|
assertThat(regionConfigurer.getRepository()).isEqualTo(mockRepository);
|
||||||
|
|
||||||
verifyNoInteractions(mockRepository);
|
verifyNoInteractions(mockRepository);
|
||||||
@@ -87,7 +94,8 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
new AsyncInlineCachingRegionConfigurer<>(mockRepository, null);
|
new AsyncInlineCachingRegionConfigurer<>(mockRepository, null);
|
||||||
|
|
||||||
assertThat(regionConfigurer).isNotNull();
|
assertThat(regionConfigurer).isNotNull();
|
||||||
assertThat(regionConfigurer.getRegionBeanName()).isNotNull();
|
assertThat(regionConfigurer.getRegionBeanName())
|
||||||
|
.isEqualTo(AsyncInlineCachingRegionConfigurer.DEFAULT_REGION_BEAN_NAME_PREDICATE);
|
||||||
assertThat(regionConfigurer.getRepository()).isEqualTo(mockRepository);
|
assertThat(regionConfigurer.getRepository()).isEqualTo(mockRepository);
|
||||||
|
|
||||||
verifyNoInteractions(mockRepository);
|
verifyNoInteractions(mockRepository);
|
||||||
@@ -113,14 +121,30 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
|
|
||||||
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
||||||
|
|
||||||
Predicate<String> regionBeanName = Predicate.isEqual("TestRegion");
|
Predicate<String> testRegionBeanName = Predicate.isEqual("TestRegion");
|
||||||
|
|
||||||
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
||||||
AsyncInlineCachingRegionConfigurer.create(mockRepository, regionBeanName);
|
AsyncInlineCachingRegionConfigurer.create(mockRepository, testRegionBeanName);
|
||||||
|
|
||||||
|
assertThat(regionConfigurer).isNotNull();
|
||||||
|
assertThat(regionConfigurer.getRegionBeanName()).isEqualTo(testRegionBeanName);
|
||||||
|
assertThat(regionConfigurer.getRepository()).isEqualTo(mockRepository);
|
||||||
|
|
||||||
|
verifyNoInteractions(mockRepository);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void createAsyncInlineCachingRegionConfigurerFromCrudRepositoryAndNullPredicate() {
|
||||||
|
|
||||||
|
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
||||||
|
|
||||||
|
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
||||||
|
AsyncInlineCachingRegionConfigurer.create(mockRepository, (Predicate<String>) null);
|
||||||
|
|
||||||
assertThat(regionConfigurer).isNotNull();
|
assertThat(regionConfigurer).isNotNull();
|
||||||
assertThat(regionConfigurer.getRepository()).isEqualTo(mockRepository);
|
assertThat(regionConfigurer.getRepository()).isEqualTo(mockRepository);
|
||||||
assertThat(regionConfigurer.getRegionBeanName()).isEqualTo(regionBeanName);
|
assertThat(regionConfigurer.getRegionBeanName())
|
||||||
|
.isEqualTo(AsyncInlineCachingRegionConfigurer.DEFAULT_REGION_BEAN_NAME_PREDICATE);
|
||||||
|
|
||||||
verifyNoInteractions(mockRepository);
|
verifyNoInteractions(mockRepository);
|
||||||
}
|
}
|
||||||
@@ -134,10 +158,30 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
AsyncInlineCachingRegionConfigurer.create(mockRepository, "MockRegion");
|
AsyncInlineCachingRegionConfigurer.create(mockRepository, "MockRegion");
|
||||||
|
|
||||||
assertThat(regionConfigurer).isNotNull();
|
assertThat(regionConfigurer).isNotNull();
|
||||||
assertThat(regionConfigurer.getRepository()).isEqualTo(mockRepository);
|
|
||||||
assertThat(regionConfigurer.getRegionBeanName()).isNotNull();
|
assertThat(regionConfigurer.getRegionBeanName()).isNotNull();
|
||||||
assertThat(regionConfigurer.getRegionBeanName().test("MockRegion")).isTrue();
|
assertThat(regionConfigurer.getRegionBeanName().test("MockRegion")).isTrue();
|
||||||
|
assertThat(regionConfigurer.getRegionBeanName().test("MOCKREGION")).isFalse();
|
||||||
|
assertThat(regionConfigurer.getRegionBeanName().test("mockregion")).isFalse();
|
||||||
assertThat(regionConfigurer.getRegionBeanName().test("TestRegion")).isFalse();
|
assertThat(regionConfigurer.getRegionBeanName().test("TestRegion")).isFalse();
|
||||||
|
assertThat(regionConfigurer.getRepository()).isEqualTo(mockRepository);
|
||||||
|
|
||||||
|
verifyNoInteractions(mockRepository);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void createAsyncInlineCachingRegionConfigurerFromCrudRepositoryAndNullString() {
|
||||||
|
|
||||||
|
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
||||||
|
|
||||||
|
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
||||||
|
AsyncInlineCachingRegionConfigurer.create(mockRepository, (String) null);
|
||||||
|
|
||||||
|
assertThat(regionConfigurer).isNotNull();
|
||||||
|
assertThat(regionConfigurer.getRegionBeanName()).isNotNull();
|
||||||
|
assertThat(regionConfigurer.getRegionBeanName().test(null)).isTrue();
|
||||||
|
assertThat(regionConfigurer.getRegionBeanName().test("MockRegion")).isFalse();
|
||||||
|
assertThat(regionConfigurer.getRegionBeanName().test("TestRegion")).isFalse();
|
||||||
|
assertThat(regionConfigurer.getRepository()).isEqualTo(mockRepository);
|
||||||
|
|
||||||
verifyNoInteractions(mockRepository);
|
verifyNoInteractions(mockRepository);
|
||||||
}
|
}
|
||||||
@@ -181,11 +225,11 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
|
|
||||||
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
||||||
|
|
||||||
|
PeerRegionFactoryBean<?, ?> peerRegionFactoryBean = mock(PeerRegionFactoryBean.class);
|
||||||
|
|
||||||
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
||||||
spy(new AsyncInlineCachingRegionConfigurer<>(mockRepository, Predicate.isEqual("TestRegion")));
|
spy(new AsyncInlineCachingRegionConfigurer<>(mockRepository, Predicate.isEqual("TestRegion")));
|
||||||
|
|
||||||
PeerRegionFactoryBean<?, ?> peerRegionFactoryBean = mock(PeerRegionFactoryBean.class);
|
|
||||||
|
|
||||||
doReturn(mockCache).when(peerRegionFactoryBean).getCache();
|
doReturn(mockCache).when(peerRegionFactoryBean).getCache();
|
||||||
doReturn(mockAsyncEventQueue).when(regionConfigurer).newAsyncEventQueue(eq(mockCache), eq("TestRegion"));
|
doReturn(mockAsyncEventQueue).when(regionConfigurer).newAsyncEventQueue(eq(mockCache), eq("TestRegion"));
|
||||||
|
|
||||||
@@ -194,9 +238,9 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
verify(regionConfigurer, times(1))
|
verify(regionConfigurer, times(1))
|
||||||
.configure(eq("TestRegion"), eq(peerRegionFactoryBean));
|
.configure(eq("TestRegion"), eq(peerRegionFactoryBean));
|
||||||
verify(regionConfigurer, times(1)).getRegionBeanName();
|
verify(regionConfigurer, times(1)).getRegionBeanName();
|
||||||
|
verify(peerRegionFactoryBean, times(1)).getCache();
|
||||||
verify(regionConfigurer, times(1))
|
verify(regionConfigurer, times(1))
|
||||||
.newAsyncEventQueue(eq(mockCache), eq("TestRegion"));
|
.newAsyncEventQueue(eq(mockCache), eq("TestRegion"));
|
||||||
verify(peerRegionFactoryBean, times(1)).getCache();
|
|
||||||
verify(peerRegionFactoryBean, times(1))
|
verify(peerRegionFactoryBean, times(1))
|
||||||
.setAsyncEventQueues(eq(ArrayUtils.asArray(mockAsyncEventQueue)));
|
.setAsyncEventQueues(eq(ArrayUtils.asArray(mockAsyncEventQueue)));
|
||||||
|
|
||||||
@@ -210,17 +254,20 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
||||||
|
|
||||||
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
||||||
spy(new AsyncInlineCachingRegionConfigurer<>(mockRepository, Predicate.isEqual("TestRegion")));
|
spy(new AsyncInlineCachingRegionConfigurer<>(mockRepository, Predicate.isEqual("NoRegion")));
|
||||||
|
|
||||||
PeerRegionFactoryBean<?, ?> peerRegionFactoryBean = mock(PeerRegionFactoryBean.class);
|
PeerRegionFactoryBean<?, ?> peerRegionFactoryBean = mock(PeerRegionFactoryBean.class);
|
||||||
|
|
||||||
regionConfigurer.configure("MockRegion", peerRegionFactoryBean);
|
regionConfigurer.configure("MockRegion", peerRegionFactoryBean);
|
||||||
|
regionConfigurer.configure("TestRegion", peerRegionFactoryBean);
|
||||||
|
|
||||||
verify(regionConfigurer, times(1))
|
verify(regionConfigurer, times(1))
|
||||||
.configure(eq("MockRegion"), eq(peerRegionFactoryBean));
|
.configure(eq("MockRegion"), eq(peerRegionFactoryBean));
|
||||||
verify(regionConfigurer, times(1)).getRegionBeanName();
|
verify(regionConfigurer, times(1))
|
||||||
|
.configure(eq("TestRegion"), eq(peerRegionFactoryBean));
|
||||||
|
verify(regionConfigurer, times(2)).getRegionBeanName();
|
||||||
verifyNoMoreInteractions(regionConfigurer);
|
verifyNoMoreInteractions(regionConfigurer);
|
||||||
verifyNoInteractions(peerRegionFactoryBean);
|
verifyNoInteractions(mockRepository, peerRegionFactoryBean);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
@@ -275,9 +322,9 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
@SuppressWarnings({ "rawtypes", "unchecked" })
|
@SuppressWarnings({ "rawtypes", "unchecked" })
|
||||||
public void newAsyncEventQueueCreatesAsyncEventQueueFromCacheInitializedWithRegionConfigurer() {
|
public void newAsyncEventQueueCreatesQueueFromCacheInitializedWithRegionConfigurer() {
|
||||||
|
|
||||||
AsyncEventListener mockAsyncEventListener = mock(AsyncEventListener.class, "Mock AEQ Listener");
|
AsyncEventErrorHandler mockAsyncEventErrorHandler = mock(AsyncEventErrorHandler.class);
|
||||||
|
|
||||||
AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class, "Mock AEQ");
|
AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class, "Mock AEQ");
|
||||||
|
|
||||||
@@ -289,11 +336,18 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
|
|
||||||
Duration batchTimeInterval = Duration.ofSeconds(15);
|
Duration batchTimeInterval = Duration.ofSeconds(15);
|
||||||
|
|
||||||
|
Function<AsyncEventListener, AsyncEventListener> mockAsyncEventListenerFunction = mock(Function.class);
|
||||||
|
Function<AsyncEventQueue, AsyncEventQueue> mockAsyncEventQueueFunction = mock(Function.class);
|
||||||
|
Function<AsyncEventQueueFactory, AsyncEventQueueFactory> mockAsyncEventQueueFactoryFunction = mock(Function.class);
|
||||||
|
|
||||||
GatewayEventFilter mockEventFilterOne = mock(GatewayEventFilter.class);
|
GatewayEventFilter mockEventFilterOne = mock(GatewayEventFilter.class);
|
||||||
GatewayEventFilter mockEventFilterTwo = mock(GatewayEventFilter.class);
|
GatewayEventFilter mockEventFilterTwo = mock(GatewayEventFilter.class);
|
||||||
|
|
||||||
GatewayEventSubstitutionFilter mockEventSubstitutionFilter = mock(GatewayEventSubstitutionFilter.class);
|
GatewayEventSubstitutionFilter mockEventSubstitutionFilter = mock(GatewayEventSubstitutionFilter.class);
|
||||||
|
|
||||||
|
RepositoryAsyncEventListener mockAsyncEventListener =
|
||||||
|
mock(RepositoryAsyncEventListener.class, "Mock AEQ Listener");
|
||||||
|
|
||||||
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
||||||
spy(new AsyncInlineCachingRegionConfigurer<>(mockRepository, Predicate.isEqual("TestRegion")));
|
spy(new AsyncInlineCachingRegionConfigurer<>(mockRepository, Predicate.isEqual("TestRegion")));
|
||||||
|
|
||||||
@@ -301,7 +355,14 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
doReturn(mockAsyncEventQueue).when(mockAsyncEventQueueFactory).create(eq("123"), eq(mockAsyncEventListener));
|
doReturn(mockAsyncEventQueue).when(mockAsyncEventQueueFactory).create(eq("123"), eq(mockAsyncEventListener));
|
||||||
doReturn("123").when(regionConfigurer).generateId(eq("TestRegion"));
|
doReturn("123").when(regionConfigurer).generateId(eq("TestRegion"));
|
||||||
doReturn(mockAsyncEventListener).when(regionConfigurer).newRepositoryAsyncEventListener();
|
doReturn(mockAsyncEventListener).when(regionConfigurer).newRepositoryAsyncEventListener();
|
||||||
|
doReturn(mockAsyncEventListener).when(mockAsyncEventListenerFunction).apply(eq(mockAsyncEventListener));
|
||||||
|
doReturn(mockAsyncEventQueue).when(mockAsyncEventQueueFunction).apply(eq(mockAsyncEventQueue));
|
||||||
|
doReturn(mockAsyncEventQueueFactory).when(mockAsyncEventQueueFactoryFunction).apply(eq(mockAsyncEventQueueFactory));
|
||||||
|
|
||||||
|
assertThat(regionConfigurer.applyToListener(mockAsyncEventListenerFunction)).isSameAs(regionConfigurer);
|
||||||
|
assertThat(regionConfigurer.applyToQueue(mockAsyncEventQueueFunction)).isSameAs(regionConfigurer);
|
||||||
|
assertThat(regionConfigurer.applyToQueueFactory(mockAsyncEventQueueFactoryFunction)).isSameAs(regionConfigurer);
|
||||||
|
assertThat(regionConfigurer.withAsyncEventErrorHandler(mockAsyncEventErrorHandler)).isSameAs(regionConfigurer);
|
||||||
assertThat(regionConfigurer.withParallelQueue()).isSameAs(regionConfigurer);
|
assertThat(regionConfigurer.withParallelQueue()).isSameAs(regionConfigurer);
|
||||||
assertThat(regionConfigurer.withPersistentQueue()).isSameAs(regionConfigurer);
|
assertThat(regionConfigurer.withPersistentQueue()).isSameAs(regionConfigurer);
|
||||||
assertThat(regionConfigurer.withQueueBatchConflationEnabled()).isSameAs(regionConfigurer);
|
assertThat(regionConfigurer.withQueueBatchConflationEnabled()).isSameAs(regionConfigurer);
|
||||||
@@ -323,45 +384,103 @@ public class AsyncInlineCachingRegionConfigurerUnitTests {
|
|||||||
assertThat(regionConfigurer.newAsyncEventQueue(mockCache, "TestRegion"))
|
assertThat(regionConfigurer.newAsyncEventQueue(mockCache, "TestRegion"))
|
||||||
.isEqualTo(mockAsyncEventQueue);
|
.isEqualTo(mockAsyncEventQueue);
|
||||||
|
|
||||||
verify(mockCache, times(1)).createAsyncEventQueueFactory();
|
InOrder order = inOrder(regionConfigurer, mockAsyncEventListener, mockAsyncEventQueueFactory, mockCache,
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setBatchConflationEnabled(eq(true));
|
mockAsyncEventListenerFunction, mockAsyncEventQueueFunction, mockAsyncEventQueueFactoryFunction);
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setBatchSize(eq(224));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setBatchTimeInterval(eq((int) batchTimeInterval.toMillis()));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setDiskStoreName(eq("TestDiskStore"));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setDiskSynchronous(eq(true));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setDispatcherThreads(eq(8));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setForwardExpirationDestroy(eq(true));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setGatewayEventSubstitutionListener(eq(mockEventSubstitutionFilter));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setMaximumQueueMemory(eq(51));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setOrderPolicy(eq(GatewaySender.OrderPolicy.THREAD));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).setPersistent(eq(true));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).addGatewayEventFilter(eq(mockEventFilterOne));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).addGatewayEventFilter(eq(mockEventFilterTwo));
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).pauseEventDispatching();
|
|
||||||
verify(mockAsyncEventQueueFactory, times(1)).create(eq("123"), eq(mockAsyncEventListener));
|
|
||||||
verify(regionConfigurer, times(1)).generateId(eq("TestRegion"));
|
|
||||||
verify(regionConfigurer, times(1)).newRepositoryAsyncEventListener();
|
|
||||||
|
|
||||||
verifyNoMoreInteractions(mockCache, mockAsyncEventQueueFactory);
|
order.verify(regionConfigurer, times(1)).newAsyncEventQueueFactory(eq(mockCache));
|
||||||
|
order.verify(mockCache, times(1)).createAsyncEventQueueFactory();
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setBatchConflationEnabled(eq(true));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setBatchSize(eq(224));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setBatchTimeInterval(eq((int) batchTimeInterval.toMillis()));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setDiskStoreName(eq("TestDiskStore"));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setDiskSynchronous(eq(true));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setDispatcherThreads(eq(8));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setForwardExpirationDestroy(eq(true));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setGatewayEventSubstitutionListener(eq(mockEventSubstitutionFilter));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setMaximumQueueMemory(eq(51));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setOrderPolicy(eq(GatewaySender.OrderPolicy.THREAD));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).setPersistent(eq(true));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).addGatewayEventFilter(eq(mockEventFilterOne));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).addGatewayEventFilter(eq(mockEventFilterTwo));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).pauseEventDispatching();
|
||||||
|
order.verify(regionConfigurer, times(1)).generateId(eq("TestRegion"));
|
||||||
|
order.verify(regionConfigurer, times(1)).newRepositoryAsyncEventListener();
|
||||||
|
order.verify(regionConfigurer, times(1)).postProcess(eq(mockAsyncEventListener));
|
||||||
|
order.verify(mockAsyncEventListener, times(1)).setAsyncEventErrorHandler(eq(mockAsyncEventErrorHandler));
|
||||||
|
order.verify(mockAsyncEventListenerFunction, times(1)).apply(eq(mockAsyncEventListener));
|
||||||
|
order.verify(regionConfigurer, times(1)).postProcess(eq(mockAsyncEventQueueFactory));
|
||||||
|
order.verify(mockAsyncEventQueueFactoryFunction, times(1)).apply(eq(mockAsyncEventQueueFactory));
|
||||||
|
order.verify(mockAsyncEventQueueFactory, times(1)).create(eq("123"), eq(mockAsyncEventListener));
|
||||||
|
order.verify(regionConfigurer, times(1)).postProcess(eq(mockAsyncEventQueue));
|
||||||
|
order.verify(mockAsyncEventQueueFunction, times(1)).apply(eq(mockAsyncEventQueue));
|
||||||
|
|
||||||
verifyNoInteractions(mockAsyncEventListener, mockAsyncEventQueue, mockRepository, mockEventFilterOne,
|
verifyNoMoreInteractions(mockCache, mockAsyncEventListener, mockAsyncEventListenerFunction,
|
||||||
mockEventFilterTwo, mockEventSubstitutionFilter);
|
mockAsyncEventQueueFactory, mockAsyncEventQueueFunction, mockAsyncEventQueueFactoryFunction);
|
||||||
|
|
||||||
|
verifyNoInteractions(mockAsyncEventErrorHandler, mockAsyncEventQueue, mockRepository,
|
||||||
|
mockEventFilterOne, mockEventFilterTwo, mockEventSubstitutionFilter);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void newRepositoryAsyncEventListener() {
|
public void newAsyncEventQueueCreatesQueueFromFactoryWithIdAndListener() {
|
||||||
|
|
||||||
|
AsyncEventListener mockAsyncEventListener = mock(AsyncEventListener.class, "Mock AEQ Listener");
|
||||||
|
|
||||||
|
AsyncEventQueue mockAsyncEventQueue = mock(AsyncEventQueue.class, "Mock AEQ");
|
||||||
|
|
||||||
|
AsyncEventQueueFactory mockAsyncEventQueueFactory = mock(AsyncEventQueueFactory.class, "Mock AEQ Factory");
|
||||||
|
|
||||||
|
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
||||||
|
|
||||||
|
String asyncEventQueueId = "abc123";
|
||||||
|
|
||||||
|
doReturn(mockAsyncEventQueue).when(mockAsyncEventQueueFactory)
|
||||||
|
.create(eq(asyncEventQueueId), eq(mockAsyncEventListener));
|
||||||
|
|
||||||
|
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
||||||
|
new AsyncInlineCachingRegionConfigurer<>(mockRepository, Predicate.isEqual("TestRegion"));
|
||||||
|
|
||||||
|
assertThat(regionConfigurer.newAsyncEventQueue(mockAsyncEventQueueFactory, asyncEventQueueId, mockAsyncEventListener))
|
||||||
|
.isEqualTo(mockAsyncEventQueue);
|
||||||
|
|
||||||
|
verify(mockAsyncEventQueueFactory, times(1))
|
||||||
|
.create(eq(asyncEventQueueId), eq(mockAsyncEventListener));
|
||||||
|
verifyNoMoreInteractions(mockAsyncEventQueueFactory);
|
||||||
|
verifyNoInteractions(mockAsyncEventListener, mockAsyncEventQueue, mockRepository);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void newAsyncEventQueueFactoryCallsCacheCreateAsyncEventQueueFactory() {
|
||||||
|
|
||||||
|
Cache mockCache = mock(Cache.class, "Mock Peer Cache");
|
||||||
|
|
||||||
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
|
||||||
|
|
||||||
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
AsyncInlineCachingRegionConfigurer<?, ?> regionConfigurer =
|
||||||
spy(new AsyncInlineCachingRegionConfigurer<>(mockRepository, Predicate.isEqual("TestRegion")));
|
new AsyncInlineCachingRegionConfigurer<>(mockRepository, Predicate.isEqual("TestRegion"));
|
||||||
|
|
||||||
|
regionConfigurer.newAsyncEventQueueFactory(mockCache);
|
||||||
|
|
||||||
|
verify(mockCache, times(1)).createAsyncEventQueueFactory();
|
||||||
|
verifyNoMoreInteractions(mockCache);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
@SuppressWarnings({ "rawtypes", "unchecked" })
|
||||||
|
public void newRepositoryAsyncEventListenerReturnsNewListener() {
|
||||||
|
|
||||||
|
CrudRepository mockRepository = mock(CrudRepository.class);
|
||||||
|
|
||||||
|
AsyncInlineCachingRegionConfigurer regionConfigurer =
|
||||||
|
spy(new AsyncInlineCachingRegionConfigurer(mockRepository, Predicate.isEqual("TestRegion")));
|
||||||
|
|
||||||
AsyncEventListener listener = regionConfigurer.newRepositoryAsyncEventListener();
|
AsyncEventListener listener = regionConfigurer.newRepositoryAsyncEventListener();
|
||||||
|
|
||||||
assertThat(listener).isInstanceOf(RepositoryAsyncEventListener.class);
|
assertThat(listener).isInstanceOf(RepositoryAsyncEventListener.class);
|
||||||
assertThat(((RepositoryAsyncEventListener<?, ?>) listener).getRepository()).isEqualTo(mockRepository);
|
assertThat(((RepositoryAsyncEventListener<?, ?>) listener).getRepository()).isEqualTo(mockRepository);
|
||||||
|
|
||||||
|
verify(regionConfigurer, times(1)).newRepositoryAsyncEventListener(eq(mockRepository));
|
||||||
verify(regionConfigurer, times(1)).getRepository();
|
verify(regionConfigurer, times(1)).getRepository();
|
||||||
verifyNoInteractions(mockRepository);
|
verifyNoInteractions(mockRepository);
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user