From 7de0a95f078b11d5acd8f190c6ee7ad7aad29bde Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Thu, 21 Jul 2022 11:28:49 +0200 Subject: [PATCH] Unsubscribe only once on `Subscription.close()`. We now unsubscribe from Redis only once when closing a Subscription object to avoid duplicate traffic and unexpected Redis responses. Jedis does not have a mechanism to consume the additional unsubscribe response which leaves protocol frames on the InputStream leading to a corrupt state. Closes #2355 --- .../redis/connection/jedis/JedisSubscription.java | 1 + .../connection/lettuce/LettuceSubscription.java | 4 ---- .../connection/util/AbstractSubscription.java | 15 +++++++++++++-- .../jedis/JedisSubscriptionUnitTests.java | 12 ++++++++++++ 4 files changed, 26 insertions(+), 6 deletions(-) diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisSubscription.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisSubscription.java index 92d8f1117..502924627 100644 --- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisSubscription.java +++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisSubscription.java @@ -42,6 +42,7 @@ class JedisSubscription extends AbstractSubscription { */ @Override protected void doClose() { + if (!getChannels().isEmpty()) { jedisPubSub.unsubscribe(); } diff --git a/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceSubscription.java b/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceSubscription.java index f3d56342c..5295f49fd 100644 --- a/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceSubscription.java +++ b/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceSubscription.java @@ -80,10 +80,6 @@ public class LettuceSubscription extends AbstractSubscription { @Override protected void doClose() { - if (!isAlive()) { - return; - } - List> futures = new ArrayList<>(); if (!getChannels().isEmpty()) { diff --git a/src/main/java/org/springframework/data/redis/connection/util/AbstractSubscription.java b/src/main/java/org/springframework/data/redis/connection/util/AbstractSubscription.java index 1ffc6c226..877151e6b 100644 --- a/src/main/java/org/springframework/data/redis/connection/util/AbstractSubscription.java +++ b/src/main/java/org/springframework/data/redis/connection/util/AbstractSubscription.java @@ -103,8 +103,19 @@ public abstract class AbstractSubscription implements Subscription { */ @Override public void close() { - doClose(); - alive.set(false); + + if (alive.compareAndSet(true, false)) { + + doClose(); + + synchronized (channels) { + channels.clear(); + } + + synchronized (patterns) { + patterns.clear(); + } + } } /** diff --git a/src/test/java/org/springframework/data/redis/connection/jedis/JedisSubscriptionUnitTests.java b/src/test/java/org/springframework/data/redis/connection/jedis/JedisSubscriptionUnitTests.java index bb8fd99c8..c6a5bf88c 100644 --- a/src/test/java/org/springframework/data/redis/connection/jedis/JedisSubscriptionUnitTests.java +++ b/src/test/java/org/springframework/data/redis/connection/jedis/JedisSubscriptionUnitTests.java @@ -35,6 +35,7 @@ import org.springframework.data.redis.connection.RedisInvalidSubscriptionExcepti * Unit test of {@link JedisSubscription} * * @author Jennifer Hickey + * @author Mark Paluch */ @ExtendWith(MockitoExtension.class) class JedisSubscriptionUnitTests { @@ -305,4 +306,15 @@ class JedisSubscriptionUnitTests { verify(jedisPubSub, times(1)).punsubscribe(); } + @Test // GH-2355 + void closeTwiceShouldUnsubscribeOnce() { + + subscription.subscribe(new byte[][] { "a".getBytes() }); + + subscription.close(); + subscription.close(); + + verify(jedisPubSub, times(1)).unsubscribe(); + } + }