Change AbstractAsyncEventOperationRepositoryFunction, CreateUpdateAsyncEventRepositoryFunction and RemoveAsyncEventRepositoryFunction to static member classes.

Add constructor requiring a reference to the associated RepositoryAsyncEventListener delegating AsyncEvent handling/processing to the Function.

Add alias methods from the RepositoryAsyncEventListener getAsyncEventErrorHandler() and getRepository() methods to the AsyncEventOperationRepositoryFunction classes.

Refatory type signatures and Generics usage.

Edit Javadoc.

Resolves gh-58.
This commit is contained in:
John Blum
2020-08-19 18:47:49 -07:00
parent ecffa18654
commit 6af91f5eb5
2 changed files with 204 additions and 74 deletions

View File

@@ -71,8 +71,8 @@ public class RepositoryAsyncEventListener<T, ID> implements AsyncEventListener {
this.repository = repository;
this.repositoryFunctions.addAll(Arrays.asList(
new CreateUpdateAsyncEventRepositoryFunction<>(),
new RemoveAsyncEventRepositoryFunction<>()
new CreateUpdateAsyncEventRepositoryFunction<>(this),
new RemoveAsyncEventRepositoryFunction<>(this)
));
}
@@ -324,9 +324,63 @@ public class RepositoryAsyncEventListener<T, ID> implements AsyncEventListener {
* @param <ID> {@link Class type} of the identifier of the entity.
* @see AsyncEventOperationRepositoryFunction
*/
protected abstract class AbstractAsyncEventOperationRepositoryFunction<T, ID>
public static abstract class AbstractAsyncEventOperationRepositoryFunction<T, ID>
implements AsyncEventOperationRepositoryFunction<T, ID> {
private final RepositoryAsyncEventListener<T, ID> listener;
/**
* Constructs an new instance of {@link AbstractAsyncEventOperationRepositoryFunction} initialized with
* the given, required {@link RepositoryAsyncEventListener} to which this function is associated.
*
* @param listener {@link RepositoryAsyncEventListener} processing {@link AsyncEvent AsyncEvents}
* by invoking this {@link Function} to handle them.
* @throws IllegalArgumentException if {@link RepositoryAsyncEventListener} is {@literal null}.
* @see RepositoryAsyncEventListener
*/
public AbstractAsyncEventOperationRepositoryFunction(@NonNull RepositoryAsyncEventListener<T, ID> listener) {
Assert.notNull(listener, "RepositoryAsyncEventListener must not be null");
this.listener = listener;
}
/**
* Alias to the {@link RepositoryAsyncEventListener#getAsyncEventErrorHandler() configured}
* {@link RepositoryAsyncEventListener} {@link AsyncEventErrorHandler}.
*
* @return the configured {@link AsyncEventErrorHandler}; never {@literal null}.
* @see RepositoryAsyncEventListener#getAsyncEventErrorHandler()
* @see AsyncEventErrorHandler
* @see #getListener()
*/
protected AsyncEventErrorHandler getErrorHandler() {
return getListener().getAsyncEventErrorHandler();
}
/**
* Returns a reference to the associated {@link RepositoryAsyncEventListener}.
*
* @return a reference to the associated {@link RepositoryAsyncEventListener}; never {@literal null}.
* @see RepositoryAsyncEventListener
*/
protected @NonNull RepositoryAsyncEventListener<T, ID> getListener() {
return this.listener;
}
/**
* Alias to the {@link RepositoryAsyncEventListener#getRepository() configured}
* {@link RepositoryAsyncEventListener} {@link CrudRepository}.
*
* @return the configured {@link CrudRepository}; never {@literal null}.
* @see org.springframework.data.repository.CrudRepository
* @see RepositoryAsyncEventListener#getRepository()
* @see #getListener()
*/
protected @NonNull CrudRepository<T, ID> getRepository() {
return getListener().getRepository();
}
/**
* Processes the given {@link AsyncEvent} by first determining whether the event can be processed by this
* {@link Function}, and then proceeds to extract the {@link AsyncEvent#getDeserializedValue() entity}
@@ -349,7 +403,7 @@ public class RepositoryAsyncEventListener<T, ID> implements AsyncEventListener {
* @see #canProcess(AsyncEvent)
* @see #doRepositoryOp(Object)
* @see #resolveEntity(AsyncEvent)
* @see #getAsyncEventErrorHandler()
* @see #getErrorHandler()
*/
@Override
public Boolean apply(@Nullable AsyncEvent<ID, T> event) {
@@ -367,7 +421,7 @@ public class RepositoryAsyncEventListener<T, ID> implements AsyncEventListener {
return false;
}
catch (Throwable cause) {
return getAsyncEventErrorHandler().apply(new AsyncEventError(event, cause));
return getErrorHandler().apply(new AsyncEventError(event, cause));
}
}
@@ -411,17 +465,29 @@ public class RepositoryAsyncEventListener<T, ID> implements AsyncEventListener {
*
* Invokes the {@link CrudRepository#save(Object)} data access operation.
*
* @param <S> {@link Class Subtype} of the entity tied to the event.
* @param <T> {@link Class type} of the entity tied to the event.
* @param <ID> {@link Class type} of the identifier of the entity.
*/
public class CreateUpdateAsyncEventRepositoryFunction<S extends T, ID>
extends AbstractAsyncEventOperationRepositoryFunction<S, ID> {
public static class CreateUpdateAsyncEventRepositoryFunction<T, ID>
extends AbstractAsyncEventOperationRepositoryFunction<T, ID> {
/**
* Constructs a new instance of {@link CreateUpdateAsyncEventRepositoryFunction} initialized with the given,
* required {@link RepositoryAsyncEventListener}.
*
* @param listener {@link RepositoryAsyncEventListener} forwarding {@link AsyncEvent AsyncEvents} for processing
* by this {@link Function}
* @see RepositoryAsyncEventListener
*/
public CreateUpdateAsyncEventRepositoryFunction(@NonNull RepositoryAsyncEventListener<T, ID> listener) {
super(listener);
}
/**
* @inheritDoc
*/
@Override
public boolean canProcess(@Nullable AsyncEvent<ID, S> event) {
public boolean canProcess(@Nullable AsyncEvent<ID, T> event) {
Operation operation = event != null ? event.getOperation() : null;
@@ -433,28 +499,39 @@ public class RepositoryAsyncEventListener<T, ID> implements AsyncEventListener {
*/
@Override
@SuppressWarnings("unchecked")
protected <R> R doRepositoryOp(S entity) {
protected <R> R doRepositoryOp(T entity) {
return (R) getRepository().save(entity);
}
}
/**
* An {@link AsyncEventOperationRepositoryFunction} capable of handling {@link Operation#REMOVE}
* {@link AsyncEvent AsyncEvents}.
* An {@link Function} implementation capable of handling {@link Operation#REMOVE} {@link AsyncEvent AsyncEvents}.
*
* Invokes the {@link CrudRepository#delete(Object)} data access operation.
*
* @param <S> {@link Class Subtype} of the entity tied to the event.
* @param <T> {@link Class type} of the entity tied to the event.
* @param <ID> {@link Class type} of the identifier of the entity.
*/
public class RemoveAsyncEventRepositoryFunction<S extends T, ID>
extends AbstractAsyncEventOperationRepositoryFunction<S, ID> {
public static class RemoveAsyncEventRepositoryFunction<T, ID>
extends AbstractAsyncEventOperationRepositoryFunction<T, ID> {
/**
* Constructs a new instance of {@link RemoveAsyncEventRepositoryFunction} initialized with the given, required
* {@link RepositoryAsyncEventListener}.
*
* @param listener {@link RepositoryAsyncEventListener} forwarding {@link AsyncEvent AsyncEvents} for processing
* by this {@link Function}
* @see RepositoryAsyncEventListener
*/
public RemoveAsyncEventRepositoryFunction(@NonNull RepositoryAsyncEventListener<T, ID> listener) {
super(listener);
}
/**
* @inheritDoc
*/
@Override
public boolean canProcess(@Nullable AsyncEvent<ID, S> event) {
public boolean canProcess(@Nullable AsyncEvent<ID, T> event) {
Operation operation = event != null ? event.getOperation() : null;
@@ -465,7 +542,7 @@ public class RepositoryAsyncEventListener<T, ID> implements AsyncEventListener {
* @inheritDoc
*/
@Override
protected <R> R doRepositoryOp(S entity) {
protected <R> R doRepositoryOp(T entity) {
getRepository().delete(entity);
return null;
}

View File

@@ -22,6 +22,7 @@ import static org.mockito.ArgumentMatchers.isA;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.doCallRealMethod;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.inOrder;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
@@ -42,6 +43,7 @@ import org.mockito.InOrder;
import org.apache.geode.cache.Operation;
import org.apache.geode.cache.asyncqueue.AsyncEvent;
import org.springframework.dao.QueryTimeoutException;
import org.springframework.data.repository.CrudRepository;
import org.springframework.geode.cache.RepositoryAsyncEventListener.AbstractAsyncEventOperationRepositoryFunction;
import org.springframework.geode.cache.RepositoryAsyncEventListener.AsyncEventError;
@@ -49,6 +51,7 @@ import org.springframework.geode.cache.RepositoryAsyncEventListener.AsyncEventEr
import org.springframework.geode.cache.RepositoryAsyncEventListener.AsyncEventOperationRepositoryFunction;
import org.springframework.geode.cache.RepositoryAsyncEventListener.CreateUpdateAsyncEventRepositoryFunction;
import org.springframework.geode.cache.RepositoryAsyncEventListener.RemoveAsyncEventRepositoryFunction;
import org.springframework.lang.NonNull;
/**
* Unit Tests for {@link RepositoryAsyncEventListener}.
@@ -99,11 +102,12 @@ public class RepositoryAsyncEventListenerUnitTests {
AsyncEventErrorHandler mockAsyncEventErrorHandler = mock(AsyncEventErrorHandler.class);
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
CrudRepository mockRepository = mock(CrudRepository.class);
RepositoryAsyncEventListener<?, ?> listener = new RepositoryAsyncEventListener<>(mockRepository);
RepositoryAsyncEventListener listener = new RepositoryAsyncEventListener<>(mockRepository);
assertThat(listener.getAsyncEventErrorHandler()).isEqualTo(RepositoryAsyncEventListener.DEFAULT_EVENT_ERROR_HANDLER);
assertThat(listener.getAsyncEventErrorHandler())
.isEqualTo(RepositoryAsyncEventListener.DEFAULT_EVENT_ERROR_HANDLER);
listener.setAsyncEventErrorHandler(mockAsyncEventErrorHandler);
@@ -111,7 +115,8 @@ public class RepositoryAsyncEventListenerUnitTests {
listener.setAsyncEventErrorHandler(null);
assertThat(listener.getAsyncEventErrorHandler()).isEqualTo(RepositoryAsyncEventListener.DEFAULT_EVENT_ERROR_HANDLER);
assertThat(listener.getAsyncEventErrorHandler())
.isEqualTo(RepositoryAsyncEventListener.DEFAULT_EVENT_ERROR_HANDLER);
verifyNoInteractions(mockAsyncEventErrorHandler, mockRepository);
}
@@ -119,9 +124,9 @@ public class RepositoryAsyncEventListenerUnitTests {
@Test
public void getRepositoryFunctionsHandlesCreateUpdateAndRemove() {
CrudRepository<Object, Object> mockRepository = mock(CrudRepository.class);
CrudRepository mockRepository = mock(CrudRepository.class);
RepositoryAsyncEventListener<Object, Object> listener = new RepositoryAsyncEventListener<>(mockRepository);
RepositoryAsyncEventListener listener = new RepositoryAsyncEventListener<>(mockRepository);
List<AsyncEventOperationRepositoryFunction<Object, Object>> repositoryFunctions =
listener.getRepositoryFunctions();
@@ -132,8 +137,10 @@ public class RepositoryAsyncEventListenerUnitTests {
.map(Object::getClass)
.collect(Collectors.toList());
assertThat(repositoryFunctionTypes).containsExactly(CreateUpdateAsyncEventRepositoryFunction.class,
RemoveAsyncEventRepositoryFunction.class);
assertThat(repositoryFunctionTypes).containsExactly(
CreateUpdateAsyncEventRepositoryFunction.class,
RemoveAsyncEventRepositoryFunction.class
);
verifyNoInteractions(mockRepository);
}
@@ -141,12 +148,11 @@ public class RepositoryAsyncEventListenerUnitTests {
@Test
public void registerAndUnregisterAsyncEventOperationRepositoryFunction() {
AsyncEventOperationRepositoryFunction<Object, Object> mockFunction =
mock(AsyncEventOperationRepositoryFunction.class);
AsyncEventOperationRepositoryFunction mockFunction = mock(AsyncEventOperationRepositoryFunction.class);
CrudRepository<Object, Object> mockRepository = mock(CrudRepository.class);
CrudRepository mockRepository = mock(CrudRepository.class);
RepositoryAsyncEventListener<Object, Object> listener = new RepositoryAsyncEventListener<>(mockRepository);
RepositoryAsyncEventListener listener = new RepositoryAsyncEventListener<>(mockRepository);
assertThat(listener.getRepositoryFunctions()).hasSize(2);
assertThat(listener.getRepositoryFunctions().get(0).getClass())
@@ -176,9 +182,9 @@ public class RepositoryAsyncEventListenerUnitTests {
@Test
public void processEventsIsNullSafe() {
CrudRepository<?, ?> mockRepository = mock(CrudRepository.class);
CrudRepository mockRepository = mock(CrudRepository.class);
RepositoryAsyncEventListener<?, ?> listener = new RepositoryAsyncEventListener<>(mockRepository);
RepositoryAsyncEventListener listener = new RepositoryAsyncEventListener<>(mockRepository);
assertThat(listener).isNotNull();
assertThat(listener.getRepository()).isEqualTo(mockRepository);
@@ -235,7 +241,7 @@ public class RepositoryAsyncEventListenerUnitTests {
List<AsyncEvent> mockEvents = Arrays.asList(mockEventOne, mockEventTwo, mockEventThree);
RepositoryAsyncEventListener listener = spy(new RepositoryAsyncEventListener<>(mockRepository));
RepositoryAsyncEventListener listener = spy(new RepositoryAsyncEventListener(mockRepository));
assertThat(listener).isNotNull();
assertThat(listener.getRepository()).isEqualTo(mockRepository);
@@ -302,7 +308,7 @@ public class RepositoryAsyncEventListenerUnitTests {
List<AsyncEvent> mockEvents = Arrays.asList(mockEventOne, mockEventTwo, mockEventThree);
RepositoryAsyncEventListener listener = spy(new RepositoryAsyncEventListener<>(mockRepository));
RepositoryAsyncEventListener listener = spy(new RepositoryAsyncEventListener(mockRepository));
assertThat(listener).isNotNull();
assertThat(listener.getRepository()).isEqualTo(mockRepository);
@@ -350,7 +356,7 @@ public class RepositoryAsyncEventListenerUnitTests {
CrudRepository mockRepository = mock(CrudRepository.class);
RepositoryAsyncEventListener listener = spy(new RepositoryAsyncEventListener<>(mockRepository));
RepositoryAsyncEventListener listener = spy(new RepositoryAsyncEventListener(mockRepository));
assertThat(listener).isNotNull();
assertThat(listener.getRepository()).isEqualTo(mockRepository);
@@ -429,7 +435,7 @@ public class RepositoryAsyncEventListenerUnitTests {
}
@Test
public void callAbstractAsyncEventOperationRepositoryFunctionApplyWhenFunctionCanProcessEvent() {
public void abstractAsyncEventOperationRepositoryFunctionApplyWhenFunctionCanProcessEvent() {
AbstractAsyncEventOperationRepositoryFunction repositoryFunction =
mock(AbstractAsyncEventOperationRepositoryFunction.class);
@@ -456,7 +462,7 @@ public class RepositoryAsyncEventListenerUnitTests {
}
@Test
public void callAbstractAsyncEventOperationRepositoryFunctionApplyWhenFunctionCannotProcessEvent() {
public void abstractAsyncEventOperationRepositoryFunctionApplyWhenFunctionCannotProcessEvent() {
AbstractAsyncEventOperationRepositoryFunction repositoryFunction =
mock(AbstractAsyncEventOperationRepositoryFunction.class);
@@ -476,14 +482,10 @@ public class RepositoryAsyncEventListenerUnitTests {
}
@Test
public void callAbstractAsyncEventOperationRepositoryFunctionApplyWhenRepositoryOperationThrowsException() {
CrudRepository mockRepository = mock(CrudRepository.class);
TestRepositoryAsyncEventListener listener = spy(new TestRepositoryAsyncEventListener(mockRepository));
public void abstractAsyncEventOperationRepositoryFunctionApplyWhenRepositoryOperationThrowsException() {
AbstractAsyncEventOperationRepositoryFunction repositoryFunction =
spy(listener.new TestAsyncEventOperationRepositoryFunction());
mock(AbstractAsyncEventOperationRepositoryFunction.class);
AsyncEvent mockEvent = mock(AsyncEvent.class);
@@ -491,17 +493,20 @@ public class RepositoryAsyncEventListenerUnitTests {
Object entity = "mock";
doReturn(mockEventErrorHandler).when(listener).getAsyncEventErrorHandler();
doCallRealMethod().when(repositoryFunction).apply(any());
doReturn(true).when(repositoryFunction).canProcess(eq(mockEvent));
doReturn(entity).when(repositoryFunction).resolveEntity(eq(mockEvent));
doThrow(new QueryTimeoutException("TEST")).when(repositoryFunction).doRepositoryOp(eq(entity));
doReturn(mockEventErrorHandler).when(repositoryFunction).getErrorHandler();
doAnswer(invocation -> {
AsyncEventError eventError = invocation.getArgument(0);
assertThat(eventError).isNotNull();
assertThat(eventError.getCause()).isInstanceOf(UnsupportedOperationException.class);
assertThat(eventError.getCause().getMessage()).isEqualTo("Not Implemented");
assertThat(eventError.getCause()).isInstanceOf(QueryTimeoutException.class);
assertThat(eventError.getCause().getMessage()).isEqualTo("TEST");
assertThat(eventError.getCause()).hasNoCause();
assertThat(eventError.getEvent()).isEqualTo(mockEvent);
return false;
@@ -510,17 +515,74 @@ public class RepositoryAsyncEventListenerUnitTests {
assertThat(repositoryFunction.apply(mockEvent)).isFalse();
InOrder order = inOrder(listener, mockEventErrorHandler, repositoryFunction);
InOrder order = inOrder(mockEventErrorHandler, repositoryFunction);
order.verify(repositoryFunction, times(1)).apply(eq(mockEvent));
order.verify(repositoryFunction, times(1)).canProcess(eq(mockEvent));
order.verify(repositoryFunction, times(1)).resolveEntity(eq(mockEvent));
order.verify(repositoryFunction, times(1)).doRepositoryOp(eq(entity));
order.verify(listener, times(1)).getAsyncEventErrorHandler();
order.verify(repositoryFunction, times(1)).getErrorHandler();
order.verify(mockEventErrorHandler, times(1)).apply(isA(AsyncEventError.class));
verifyNoMoreInteractions(repositoryFunction, listener, mockEventErrorHandler);
verifyNoInteractions(mockEvent, mockRepository);
verifyNoMoreInteractions(repositoryFunction, mockEventErrorHandler);
verifyNoInteractions(mockEvent);
}
@Test
public void asyncEventOperationRepositoryFunctionGetErrorHandlerCallsRepositoryAsyncEventListenerGetAsyncEventErrorHandler() {
AsyncEventErrorHandler mockErrorHandler = mock(AsyncEventErrorHandler.class);
CrudRepository mockRepository = mock(CrudRepository.class);
RepositoryAsyncEventListener listener = new RepositoryAsyncEventListener(mockRepository);
listener.setAsyncEventErrorHandler(mockErrorHandler);
listener = spy(listener);
AbstractAsyncEventOperationRepositoryFunction repositoryFunction =
new TestAsyncEventOperationRepositoryFunction(listener);
assertThat(repositoryFunction.getListener()).isSameAs(listener);
assertThat(repositoryFunction.getErrorHandler()).isEqualTo(mockErrorHandler);
verify(listener, times(1)).getAsyncEventErrorHandler();
verifyNoMoreInteractions(listener);
verifyNoInteractions(mockRepository);
}
@Test
public void asyncEventOperationRepositoryFunctionGetRepositoryCallsRepositoryAsyncEventListenerGetRepository() {
CrudRepository mockRepository = mock(CrudRepository.class);
RepositoryAsyncEventListener listener = spy(new RepositoryAsyncEventListener(mockRepository));
AbstractAsyncEventOperationRepositoryFunction repositoryFunction =
new TestAsyncEventOperationRepositoryFunction(listener);
assertThat(repositoryFunction.getListener()).isSameAs(listener);
assertThat(repositoryFunction.getRepository()).isEqualTo(mockRepository);
verify(listener, times(1)).getRepository();
verifyNoMoreInteractions(listener);
verifyNoInteractions(mockRepository);
}
@Test(expected = IllegalArgumentException.class)
public void constructAsyncEventOperationRepositoryFunctionWithNullRepositoryThrowsIllegalArgumentException() {
try {
new TestAsyncEventOperationRepositoryFunction(null);
}
catch (IllegalArgumentException expected) {
assertThat(expected).hasMessage("RepositoryAsyncEventListener must not be null");
assertThat(expected).hasNoCause();
throw expected;
}
}
@Test
@@ -669,13 +731,15 @@ public class RepositoryAsyncEventListenerUnitTests {
CrudRepository mockRepository = mock(CrudRepository.class);
TestRepositoryAsyncEventListener listener = new TestRepositoryAsyncEventListener(mockRepository);
RepositoryAsyncEventListener listener = new RepositoryAsyncEventListener(mockRepository);
CreateUpdateAsyncEventRepositoryFunction repositoryFunction =
spy(listener.new TestCreateUpdateAsyncEventOperationRepositoryFunction());
new CreateUpdateAsyncEventRepositoryFunction(listener);
doReturn("TEST").when(mockRepository).save(any());
assertThat(repositoryFunction.getListener()).isEqualTo(listener);
assertThat(repositoryFunction.getRepository()).isEqualTo(mockRepository);
assertThat(repositoryFunction.doRepositoryOp("MOCK")).isEqualTo("TEST");
verify(mockRepository, times(1)).save(eq("MOCK"));
@@ -742,40 +806,29 @@ public class RepositoryAsyncEventListenerUnitTests {
CrudRepository mockRepository = mock(CrudRepository.class);
TestRepositoryAsyncEventListener listener = spy(new TestRepositoryAsyncEventListener(mockRepository));
RepositoryAsyncEventListener listener = new RepositoryAsyncEventListener(mockRepository);
doReturn(mockRepository).when(listener).getRepository();
RemoveAsyncEventRepositoryFunction repositoryFunction =
spy(listener.new TestRemoveAsyncEventRepositoryFunction());
RemoveAsyncEventRepositoryFunction repositoryFunction = new RemoveAsyncEventRepositoryFunction(listener);
assertThat(repositoryFunction.getListener()).isEqualTo(listener);
assertThat(repositoryFunction.getRepository()).isEqualTo(mockRepository);
assertThat(repositoryFunction.doRepositoryOp("MOCK")).isNull();
verify(mockRepository, times(1)).delete(eq("MOCK"));
verifyNoMoreInteractions(mockRepository);
}
private static class TestRepositoryAsyncEventListener<T, ID> extends RepositoryAsyncEventListener<T, ID> {
private static final class TestAsyncEventOperationRepositoryFunction<T, ID>
extends AbstractAsyncEventOperationRepositoryFunction<T, ID> {
private TestRepositoryAsyncEventListener(CrudRepository<T, ID> repository) {
super(repository);
private TestAsyncEventOperationRepositoryFunction(RepositoryAsyncEventListener<T, ID> listener) {
super(listener);
}
private class TestAsyncEventOperationRepositoryFunction
extends AbstractAsyncEventOperationRepositoryFunction<T, ID> {
@Override
protected <R> R doRepositoryOp(T entity) {
throw new UnsupportedOperationException("Not Implemented");
}
@Override
protected <R> R doRepositoryOp(@NonNull T entity) {
throw new UnsupportedOperationException("Not Implemented");
}
private class TestCreateUpdateAsyncEventOperationRepositoryFunction
extends CreateUpdateAsyncEventRepositoryFunction<T, ID> {
}
private class TestRemoveAsyncEventRepositoryFunction extends RemoveAsyncEventRepositoryFunction<T, ID> { }
}
}