This is an automated email from the ASF dual-hosted git repository.

RongtongJin pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/rocketmq.git


The following commit(s) were added to refs/heads/develop by this push:
     new 9ea2ccdc2b [ISSUE #11230] Fix self-eviction on repeated exclusive 
LiteTopic subscription (#11231)
9ea2ccdc2b is described below

commit 9ea2ccdc2b41445aee5e177e14e501d606f5a1ed
Author: Xiao Yang <[email protected]>
AuthorDate: Thu Oct 8 17:24:04 2026 +0800

    [ISSUE #11230] Fix self-eviction on repeated exclusive LiteTopic 
subscription (#11231)
---
 .../broker/lite/LiteSubscriptionRegistryImpl.java  |  5 ++--
 .../lite/LiteSubscriptionRegistryImplTest.java     | 29 ++++++++++++++++++++++
 2 files changed, 32 insertions(+), 2 deletions(-)

diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java
 
b/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java
index 8571d664c4..b818ff29cb 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java
@@ -333,7 +333,8 @@ public class LiteSubscriptionRegistryImpl extends 
ServiceThread implements LiteS
             return;
         }
         List<ClientGroup> toRemove = clientSet.stream()
-            .filter(clientGroup -> Objects.equals(group, clientGroup.group))
+            .filter(clientGroup -> Objects.equals(group, clientGroup.group)
+                && !Objects.equals(newClientId, clientGroup.clientId))
             .collect(Collectors.toList());
 
         toRemove.forEach(clientGroup -> {
@@ -511,4 +512,4 @@ public class LiteSubscriptionRegistryImpl extends 
ServiceThread implements LiteS
         return exclusiveEvictionTombstones.contains(clientId, lmqName);
     }
 
-}
\ No newline at end of file
+}
diff --git 
a/broker/src/test/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImplTest.java
 
b/broker/src/test/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImplTest.java
index 505613508a..4fc94e150f 100644
--- 
a/broker/src/test/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImplTest.java
+++ 
b/broker/src/test/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImplTest.java
@@ -47,9 +47,11 @@ import static org.junit.Assert.assertNotNull;
 import static org.junit.Assert.assertNull;
 import static org.junit.Assert.assertThrows;
 import static org.junit.Assert.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -237,6 +239,33 @@ public class LiteSubscriptionRegistryImplTest {
         verify(mockListener).onRegister(clientId2, group, "lmq1");
     }
 
+    @Test
+    public void 
testAddPartialSubscription_ExclusiveModeDoesNotEvictSameClient() {
+        String clientId = "testClient";
+        String group = "testGroup";
+        String topic = "testTopic";
+        String lmqName = LiteUtil.toLmqName(topic, "liteTopic");
+        Channel channel = mock(Channel.class);
+
+        SubscriptionGroupConfig groupConfig = new SubscriptionGroupConfig();
+        groupConfig.setLiteSubExclusive(true);
+        
when(mockSubscriptionGroupManager.findSubscriptionGroupConfig(group)).thenReturn(groupConfig);
+        when(mockLifecycleManager.isSubscriptionActive(topic, 
lmqName)).thenReturn(true);
+
+        registry.updateClientChannel(clientId, channel);
+        registry.addPartialSubscription(clientId, group, topic, 
Collections.singleton(lmqName), null);
+        registry.addPartialSubscription(clientId, group, topic, 
Collections.singleton(lmqName), null);
+
+        assertNotNull(registry.getLiteSubscription(clientId));
+        
assertTrue(registry.getLiteSubscription(clientId).getLmqSet().contains(lmqName));
+        assertEquals(1, registry.getActiveSubscriptionNum());
+        assertFalse(registry.hasExclusiveEvictionTombstone(clientId, lmqName));
+        verify(mockBroker2Client, never()).notifyUnsubscribeLite(
+            eq(channel), any(NotifyUnsubscribeLiteRequestHeader.class));
+        verify(mockListener).onRegister(clientId, group, lmqName);
+        verify(mockListener, never()).onUnregister(clientId, group, lmqName);
+    }
+
     /**
      * Test removePartialSubscription removes partial subscription correctly
      */

Reply via email to