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

lizhimins 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 e58f9fd35a [ISSUE #10959] Fix RocksDBConsumeQueue.iterateFrom should 
reject offset below min offset (#10960)
e58f9fd35a is described below

commit e58f9fd35ad45b4bd11d749d2a368a3f6b3b9d38
Author: ymwneu <[email protected]>
AuthorDate: Wed Aug 19 10:56:51 2026 +0800

    [ISSUE #10959] Fix RocksDBConsumeQueue.iterateFrom should reject offset 
below min offset (#10960)
    
    startIndex < minOffset was previously accepted as long as it stayed
    within [0, maxOffset), returning a non-null LargeRocksDBConsumeQueueIterator
    whose next() can yield a null CqUnit once the underlying data has been
    purged. Callers like ScheduleMessageService rely on iterateFrom returning
    null to detect and correct an out-of-range offset (matching ConsumeQueue's
    getMinLogicOffset() check); without it they NPE dereferencing the null
    CqUnit instead.
    
    Add the same startIndex >= getMinOffsetInQueue() bound used by the
    file-based ConsumeQueue to both iterateFrom overloads.
    
    Co-authored-by: maowei.ymw <[email protected]>
---
 .../rocketmq/store/queue/RocksDBConsumeQueue.java  |  6 ++++--
 .../store/queue/RocksDBConsumeQueueTest.java       | 22 ++++++++++++++++++++++
 2 files changed, 26 insertions(+), 2 deletions(-)

diff --git 
a/store/src/main/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueue.java 
b/store/src/main/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueue.java
index 86b4d3ef8b..19672d1b6b 100644
--- 
a/store/src/main/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueue.java
+++ 
b/store/src/main/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueue.java
@@ -298,8 +298,9 @@ public class RocksDBConsumeQueue implements 
ConsumeQueueInterface {
 
     @Override
     public ReferredIterator<CqUnit> iterateFrom(final long startIndex) {
+        long minCqOffset = getMinOffsetInQueue();
         long maxCqOffset = getMaxOffsetInQueue();
-        if (startIndex < maxCqOffset && startIndex >= 0) {
+        if (startIndex >= minCqOffset && startIndex < maxCqOffset && 
startIndex >= 0) {
             int num = pullNum(startIndex, maxCqOffset);
             return new LargeRocksDBConsumeQueueIterator(startIndex, num);
         }
@@ -308,8 +309,9 @@ public class RocksDBConsumeQueue implements 
ConsumeQueueInterface {
 
     @Override
     public ReferredIterator<CqUnit> iterateFrom(long startIndex, int count) 
throws RocksDBException {
+        long minCqOffset = getMinOffsetInQueue();
         long maxCqOffset = getMaxOffsetInQueue();
-        if (startIndex < maxCqOffset) {
+        if (startIndex >= minCqOffset && startIndex < maxCqOffset) {
             int num = Math.min((int)(maxCqOffset - startIndex), count);
             return iterateFrom0(startIndex, num, maxCqOffset);
         }
diff --git 
a/store/src/test/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueueTest.java
 
b/store/src/test/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueueTest.java
index 702d91fb07..f41d350437 100644
--- 
a/store/src/test/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueueTest.java
+++ 
b/store/src/test/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueueTest.java
@@ -120,6 +120,28 @@ public class RocksDBConsumeQueueTest extends QueueTestBase 
{
         assertFalse(it.hasNext());
     }
 
+    @Test
+    public void testIterateFrom_startIndexBelowMinOffset_returnsNull() throws 
Exception {
+        if (MixAll.isMac()) {
+            return;
+        }
+        DefaultMessageStore messageStore = mock(DefaultMessageStore.class);
+        RocksDBConsumeQueueStore rocksDBConsumeQueueStore = 
mock(RocksDBConsumeQueueStore.class);
+        
when(messageStore.getQueueStore()).thenReturn(rocksDBConsumeQueueStore);
+        when(rocksDBConsumeQueueStore.getMinOffsetInQueue(anyString(), 
anyInt())).thenReturn(9000L);
+        when(rocksDBConsumeQueueStore.getMaxOffsetInQueue(anyString(), 
anyInt())).thenReturn(10000L);
+
+        RocksDBConsumeQueue consumeQueue = new 
RocksDBConsumeQueue(messageStore.getMessageStoreConfig(), 
rocksDBConsumeQueueStore, "topic", 0);
+
+        // startIndex has already been purged (below the queue's current min 
offset); must not
+        // return a non-null iterator that later yields a null CqUnit and NPEs 
its caller.
+        assertNull(consumeQueue.iterateFrom(8000));
+        assertNull(consumeQueue.iterateFrom(8000, 10));
+
+        // startIndex within [min, max) still works as before.
+        assertNotNull(consumeQueue.iterateFrom(9000));
+    }
+
     @Test
     public void testLmqCounter_running() throws ConsumeQueueException {
         messageStore.getMessageStoreConfig().setEnableMultiDispatch(true);

Reply via email to