From ed79d1435542254e365e6f2dc5ab64530594c53f Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Thu, 6 Dec 2018 11:46:55 +0100 Subject: [PATCH] DATAREDIS-905 - Fix race condition in StreamMessageListenerContainer subscription activation. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Stream subscriptions now report reliably their state reflecting regarding activation, cancellation and while being active. Previously, registering a subscription with immediate cancel could report an inactive subscription although StreamMessageListenerContainer could perform a stream read. When activating a subscription in StreamMessageListenerContainer, awaiting activation (awaitStart(…)), immediately cancelling the subscription, and reading active state via Subscription.isActive(), Subscription.isActive() could report in this case false. This was because we checked that the status was either running or the stream read has reached event looping. The state was CANCELLED and the task thread had not yet reached the event loop and so there was a gap. Original Pull Request: #377 --- .../springframework/data/redis/stream/StreamPollTask.java | 5 ++--- .../StreamMessageListenerContainerIntegrationTests.java | 2 +- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/src/main/java/org/springframework/data/redis/stream/StreamPollTask.java b/src/main/java/org/springframework/data/redis/stream/StreamPollTask.java index f1cd4abde..1e74cc59a 100644 --- a/src/main/java/org/springframework/data/redis/stream/StreamPollTask.java +++ b/src/main/java/org/springframework/data/redis/stream/StreamPollTask.java @@ -115,11 +115,11 @@ class StreamPollTask> implements Task { public void run() { pollState.starting(); - pollState.running(); try { isInEventLoop = true; + pollState.running(); doLoop(request.getStreamOffset().getKey()); } finally { isInEventLoop = false; @@ -144,8 +144,7 @@ class StreamPollTask> implements Task { } } catch (InterruptedException e) { - pollState.cancel(); - + cancel(); Thread.currentThread().interrupt(); } catch (RuntimeException e) { diff --git a/src/test/java/org/springframework/data/redis/stream/StreamMessageListenerContainerIntegrationTests.java b/src/test/java/org/springframework/data/redis/stream/StreamMessageListenerContainerIntegrationTests.java index 6155c5c2c..400024d86 100644 --- a/src/test/java/org/springframework/data/redis/stream/StreamMessageListenerContainerIntegrationTests.java +++ b/src/test/java/org/springframework/data/redis/stream/StreamMessageListenerContainerIntegrationTests.java @@ -53,7 +53,7 @@ import org.springframework.data.redis.stream.StreamMessageListenerContainer.Stre /** * Integration tests for {@link StreamMessageListenerContainer}. - * + * * @author Mark Paluch */ public class StreamMessageListenerContainerIntegrationTests {