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() {