From cccadf68aa9fa80af968266c11768c052c25a11d Mon Sep 17 00:00:00 2001 From: John Blum Date: Tue, 1 Dec 2020 16:09:43 -0800 Subject: [PATCH] Add metadata and state to track the RepositoryAsyncEventListener.processEvents(:List) method invocations. Specifically, the listener state will track: * The number of processEvents(:List) method invocations. * Whether the processEvents(..) method has ever been invoked. * And, whether the processEvents(..) method has been invoked since the last check. Introduces a protected doProcessEvents(:List) method and changes processEvents(:List) method to final. --- .../cache/RepositoryAsyncEventListener.java | 55 ++++++++++++++++++- ...RepositoryAsyncEventListenerUnitTests.java | 35 ++++++++++++ 2 files changed, 88 insertions(+), 2 deletions(-) diff --git a/spring-geode/src/main/java/org/springframework/geode/cache/RepositoryAsyncEventListener.java b/spring-geode/src/main/java/org/springframework/geode/cache/RepositoryAsyncEventListener.java index 230fd8ac..a865f08c 100644 --- a/spring-geode/src/main/java/org/springframework/geode/cache/RepositoryAsyncEventListener.java +++ b/spring-geode/src/main/java/org/springframework/geode/cache/RepositoryAsyncEventListener.java @@ -21,6 +21,7 @@ import java.util.Objects; import java.util.Optional; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; import java.util.function.Function; import org.apache.geode.cache.Operation; @@ -51,6 +52,10 @@ public class RepositoryAsyncEventListener implements AsyncEventListener { private AsyncEventErrorHandler asyncEventErrorHandler = DEFAULT_ASYNC_EVENT_ERROR_HANDLER; + private final AtomicBoolean hasFired = new AtomicBoolean(false); + + private final AtomicLong firedCount = new AtomicLong(0L); + private final CrudRepository repository; private final List> repositoryFunctions = new CopyOnWriteArrayList<>(); @@ -76,6 +81,41 @@ public class RepositoryAsyncEventListener implements AsyncEventListener { )); } + /** + * Determines whether this listener has (ever) been fired (triggered) by the GemFire/Geode AEQ system. + * + * @return a boolean value indicating whether this listener has been fired (triggered). + * @see #hasFiredSinceLastCheck() + */ + @SuppressWarnings("unused") + public boolean hasFired() { + return getFiredCount() > 0; + } + + /** + * Determines whether this listener has been fired (triggered) by the GemFire/Geode AEQ system + * since the last check. + * + * A call to this method clears the flag. + * + * @return a boolean value indicating whether this listener has been fired (triggered) since the last check. + * @see #hasFired() + */ + @SuppressWarnings("unused") + public boolean hasFiredSinceLastCheck() { + return this.hasFired.compareAndSet(true, false); + } + + /** + * Determines how many times this listener has been fired (triggered) by the GemFire/Geode AEQ system. + * + * @return a {@link Long} value indicating how many times this listener has been fired (triggered). + */ + @SuppressWarnings("unused") + public long getFiredCount() { + return this.firedCount.get(); + } + /** * Configures an {@link AsyncEventErrorHandler} to handle errors that may occur when this listener is invoked with * a batch of {@link AsyncEvent AsyncEvents}. @@ -147,8 +187,19 @@ public class RepositoryAsyncEventListener implements AsyncEventListener { * @see java.util.List */ @Override - @SuppressWarnings("unchecked") - public boolean processEvents(List events) { + public final boolean processEvents(List events) { + + this.firedCount.incrementAndGet(); + this.hasFired.set(true); + + return doProcessEvents(events); + } + + /** + * @see #processEvents(List) + */ + @SuppressWarnings({ "rawtypes", "unchecked" }) + protected boolean doProcessEvents(List events) { AtomicBoolean result = new AtomicBoolean(true); diff --git a/spring-geode/src/test/java/org/springframework/geode/cache/RepositoryAsyncEventListenerUnitTests.java b/spring-geode/src/test/java/org/springframework/geode/cache/RepositoryAsyncEventListenerUnitTests.java index 2ccf164c..8e5c132b 100644 --- a/spring-geode/src/test/java/org/springframework/geode/cache/RepositoryAsyncEventListenerUnitTests.java +++ b/spring-geode/src/test/java/org/springframework/geode/cache/RepositoryAsyncEventListenerUnitTests.java @@ -385,6 +385,41 @@ public class RepositoryAsyncEventListenerUnitTests { verifyNoInteractions(mockRepository); } + @Test + public void processEventsCountsInvocations() { + + CrudRepository mockRepository = mock(CrudRepository.class); + + RepositoryAsyncEventListener listener = spy(new RepositoryAsyncEventListener<>(mockRepository)); + + doReturn(true).when(listener).doProcessEvents(any()); + + assertThat(listener).isNotNull(); + assertThat(listener.getRepository()).isEqualTo(mockRepository); + assertThat(listener.getFiredCount()).isZero(); + assertThat(listener.hasFired()).isFalse(); + assertThat(listener.hasFiredSinceLastCheck()).isFalse(); + + listener.processEvents(Collections.emptyList()); + + assertThat(listener.getFiredCount()).isOne(); + assertThat(listener.hasFired()).isTrue(); + assertThat(listener.hasFiredSinceLastCheck()).isTrue(); + assertThat(listener.getFiredCount()).isOne(); + assertThat(listener.hasFired()).isTrue(); + assertThat(listener.hasFiredSinceLastCheck()).isFalse(); + + listener.processEvents(Collections.singletonList(mock(AsyncEvent.class))); + listener.processEvents(Collections.emptyList()); + + assertThat(listener.getFiredCount()).isEqualTo(3); + assertThat(listener.hasFired()).isTrue(); + assertThat(listener.hasFiredSinceLastCheck()).isTrue(); + assertThat(listener.getFiredCount()).isEqualTo(3); + assertThat(listener.hasFired()).isTrue(); + assertThat(listener.hasFiredSinceLastCheck()).isFalse(); + } + @Test public void constructAsyncEventError() {