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
*/