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

lollipopjin 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 cf27650066 [ISSUE #10953] Clamp shifted message count range to queue 
bounds (#10954)
cf27650066 is described below

commit cf27650066ac5a79e19a38d58c23ca321f70f7f9
Author: qianye <[email protected]>
AuthorDate: Tue Aug 18 12:02:27 2026 +0800

    [ISSUE #10953] Clamp shifted message count range to queue bounds (#10954)
---
 .../apache/rocketmq/store/DefaultMessageStore.java | 11 +++-
 .../rocketmq/store/DefaultMessageStoreTest.java    | 68 ++++++++++++++++++++++
 2 files changed, 77 insertions(+), 2 deletions(-)

diff --git 
a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java 
b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java
index 64ce41e47d..e2d95b963d 100644
--- a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java
+++ b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java
@@ -3184,12 +3184,19 @@ public class DefaultMessageStore implements 
MessageStore {
             return 0;
         }
 
-        // correct the "from" argument to min offset in queue if it is too 
small
+        // shift the range to min offset if it is too small, then clamp it to 
the current queue bounds
         long minOffset = consumeQueue.getMinOffsetInQueue();
+        long maxOffset = consumeQueue.getMaxOffsetInQueue();
         if (from < minOffset) {
             long diff = to - from;
             from = minOffset;
-            to = from + diff;
+            to = diff > maxOffset - from ? maxOffset : from + diff;
+        }
+
+        from = Math.min(from, maxOffset);
+        to = Math.min(to, maxOffset);
+        if (from >= to) {
+            return 0;
         }
 
         long msgCount = consumeQueue.estimateMessageCount(from, to, filter);
diff --git 
a/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java 
b/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java
index 39d837e7bc..05e186537b 100644
--- a/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java
+++ b/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java
@@ -20,10 +20,15 @@ package org.apache.rocketmq.store;
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertTrue;
+import static org.mockito.Mockito.anyLong;
 import static org.mockito.Mockito.atLeastOnce;
+import static org.mockito.Mockito.doReturn;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.same;
 import static org.mockito.Mockito.spy;
 import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
 
 import org.mockito.ArgumentCaptor;
 
@@ -346,6 +351,69 @@ public class DefaultMessageStoreTest {
         }
     }
 
+    @Test
+    public void testEstimateMessageCountShiftsRangeWithinQueueBounds() {
+        String topic = "FooBar";
+        int queueId = 0;
+        long minOffset = 100;
+        long to = 200;
+        long expectedCount = 50;
+        MessageFilter filter = mock(MessageFilter.class);
+        ConsumeQueueInterface consumeQueue = mock(ConsumeQueueInterface.class);
+        DefaultMessageStore store = spy(getDefaultMessageStore());
+        doReturn(consumeQueue).when(store).findConsumeQueue(topic, queueId);
+        when(consumeQueue.getMinOffsetInQueue()).thenReturn(minOffset);
+        when(consumeQueue.getMaxOffsetInQueue()).thenReturn(to);
+        when(consumeQueue.estimateMessageCount(minOffset, to, 
filter)).thenReturn(expectedCount);
+
+        long count = store.estimateMessageCount(topic, queueId, 0, to, filter);
+
+        assertThat(count).isEqualTo(expectedCount);
+        verify(consumeQueue).estimateMessageCount(minOffset, to, filter);
+    }
+
+    @Test
+    public void testEstimateMessageCountPreservesLengthWhenShiftedRangeFits() {
+        String topic = "FooBar";
+        int queueId = 0;
+        long minOffset = 100;
+        long maxOffset = 200;
+        long to = 50;
+        long shiftedTo = 150;
+        long expectedCount = 25;
+        MessageFilter filter = mock(MessageFilter.class);
+        ConsumeQueueInterface consumeQueue = mock(ConsumeQueueInterface.class);
+        DefaultMessageStore store = spy(getDefaultMessageStore());
+        doReturn(consumeQueue).when(store).findConsumeQueue(topic, queueId);
+        when(consumeQueue.getMinOffsetInQueue()).thenReturn(minOffset);
+        when(consumeQueue.getMaxOffsetInQueue()).thenReturn(maxOffset);
+        when(consumeQueue.estimateMessageCount(minOffset, shiftedTo, 
filter)).thenReturn(expectedCount);
+
+        long count = store.estimateMessageCount(topic, queueId, 0, to, filter);
+
+        assertThat(count).isEqualTo(expectedCount);
+        verify(consumeQueue).estimateMessageCount(minOffset, shiftedTo, 
filter);
+    }
+
+    @Test
+    public void testEstimateMessageCountReturnsZeroWhenRangeIsAfterMaxOffset() 
{
+        String topic = "FooBar";
+        int queueId = 0;
+        long minOffset = 100;
+        long maxOffset = 200;
+        MessageFilter filter = mock(MessageFilter.class);
+        ConsumeQueueInterface consumeQueue = mock(ConsumeQueueInterface.class);
+        DefaultMessageStore store = spy(getDefaultMessageStore());
+        doReturn(consumeQueue).when(store).findConsumeQueue(topic, queueId);
+        when(consumeQueue.getMinOffsetInQueue()).thenReturn(minOffset);
+        when(consumeQueue.getMaxOffsetInQueue()).thenReturn(maxOffset);
+
+        long count = store.estimateMessageCount(topic, queueId, 250, 300, 
filter);
+
+        assertThat(count).isZero();
+        verify(consumeQueue, never()).estimateMessageCount(anyLong(), 
anyLong(), same(filter));
+    }
+
     @Test
     public void testGetStoreTime_ParamIsNull() {
         long storeTime = getStoreTime(null);

Reply via email to