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);