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 a865f08c..035002f7 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 @@ -189,10 +189,13 @@ public class RepositoryAsyncEventListener implements AsyncEventListener { @Override public final boolean processEvents(List events) { - this.firedCount.incrementAndGet(); - this.hasFired.set(true); - - return doProcessEvents(events); + try { + return doProcessEvents(events); + } + finally { + this.firedCount.incrementAndGet(); + this.hasFired.set(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 8e5c132b..056d133b 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 @@ -420,6 +420,41 @@ public class RepositoryAsyncEventListenerUnitTests { assertThat(listener.hasFiredSinceLastCheck()).isFalse(); } + @Test(expected = IllegalStateException.class) + public void processEventsCountsInvocationsEvenWhenAnExceptionIsThrown() { + + CrudRepository mockRepository = mock(CrudRepository.class); + + RepositoryAsyncEventListener listener = spy(new RepositoryAsyncEventListener<>(mockRepository)); + + doThrow(new IllegalStateException("TEST")).when(listener).doProcessEvents(any()); + + assertThat(listener).isNotNull(); + assertThat(listener.getFiredCount()).isZero(); + assertThat(listener.hasFired()).isFalse(); + assertThat(listener.hasFiredSinceLastCheck()).isFalse(); + + try { + listener.processEvents(Collections.singletonList(mock(AsyncEvent.class))); + } + catch (IllegalStateException expected) { + + assertThat(expected).hasMessage("TEST"); + assertThat(expected).hasNoCause(); + + throw expected; + } + finally { + + assertThat(listener.getFiredCount()).isOne(); + assertThat(listener.hasFired()).isTrue(); + assertThat(listener.hasFiredSinceLastCheck()).isTrue(); + assertThat(listener.getFiredCount()).isOne(); + assertThat(listener.hasFired()).isTrue(); + assertThat(listener.hasFiredSinceLastCheck()).isFalse(); + } + } + @Test public void constructAsyncEventError() {