From 6434cbd52f58d651edf3b31d939747e5a8c3f495 Mon Sep 17 00:00:00 2001 From: liuhy Date: Thu, 2 Jul 2026 23:20:02 +0800 Subject: [PATCH 1/2] Prevent POP revive checkpoint ordering overflow POP revive checkpoints are ordered by reviveOffset before mergeAndRevive advances through the pending checkpoints. The previous comparator subtracted two long offsets and cast the result to int, so large offset gaps could reverse the intended ordering. Use Comparator.comparingLong to preserve the full long ordering and add a regression test covering offsets separated by more than Integer.MAX_VALUE. Constraint: reviveOffset is a long queue offset and can exceed int comparison range Rejected: Keep subtraction with a wider cast | Comparator APIs express the ordering directly and avoid overflow Confidence: high Scope-risk: narrow Directive: Do not compare long queue offsets by subtraction in POP revive ordering Tested: mvn -pl broker -Dtest=PopReviveServiceTest test -Dspotbugs.skip=true -Dcheckstyle.skip=true Tested: git diff --check --- .../broker/processor/PopReviveService.java | 3 ++- .../broker/processor/PopReviveServiceTest.java | 14 ++++++++++++++ 2 files changed, 16 insertions(+), 1 deletion(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java index 07f16e98965..e5eec54f5b9 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java @@ -53,6 +53,7 @@ import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Collections; +import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -725,7 +726,7 @@ ArrayList genSortList() { return sortList; } sortList = new ArrayList<>(map.values()); - sortList.sort((o1, o2) -> (int) (o1.getReviveOffset() - o2.getReviveOffset())); + sortList.sort(Comparator.comparingLong(PopCheckPoint::getReviveOffset)); return sortList; } } diff --git a/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java b/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java index fa7e9982e1f..7b9f9e7fddc 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java @@ -442,6 +442,20 @@ public void testReviveMsgFromBatchAck() throws Throwable { assertEquals(1, commitOffsetCaptor.getValue().longValue()); } + @Test + public void testGenSortListShouldSortLargeReviveOffsets() { + PopReviveService.ConsumeReviveObj consumeReviveObj = new PopReviveService.ConsumeReviveObj(); + PopCheckPoint lowOffsetCk = buildPopCheckPoint(0, 0, 1); + PopCheckPoint highOffsetCk = buildPopCheckPoint(1, 0, (long) Integer.MAX_VALUE + 2); + consumeReviveObj.map.put("high", highOffsetCk); + consumeReviveObj.map.put("low", lowOffsetCk); + + List sortList = consumeReviveObj.genSortList(); + + assertEquals(lowOffsetCk, sortList.get(0)); + assertEquals(highOffsetCk, sortList.get(1)); + } + public static MessageExtBrokerInner buildBatchAckMsg(BatchAckMsg batchAckMsg, long deliverMs, long reviveOffset, long deliverTime) { MessageExtBrokerInner result = buildBatchAckInnerMessage(REVIVE_TOPIC, batchAckMsg, REVIVE_QUEUE_ID, STORE_HOST, deliverMs, PopMessageProcessor.genAckUniqueId(batchAckMsg)); result.setQueueOffset(reviveOffset); From 5da7a034991e253955bf8c6416462b58aab9051c Mon Sep 17 00:00:00 2001 From: liuhy Date: Fri, 10 Jul 2026 20:30:44 +0800 Subject: [PATCH 2/2] Prevent POP checkpoint compareTo overflow on the revive path PopReviveService keeps inflight checkpoints in a TreeMap (inflightReviveRequestMap), which orders keys via PopCheckPoint.compareTo. That comparator subtracted two long startOffsets and cast the result to int, the same overflow pattern already fixed for the reviveOffset sort in the prior commit. With large offset gaps the TreeMap could misorder keys and corrupt the revive ordering. Use Long.compare to preserve the full long ordering and add a regression test covering startOffsets separated by more than Integer.MAX_VALUE, including TreeMap ordering which exercises the actual PopReviveService usage. Constraint: startOffset is a long queue offset and can exceed int comparison range Rejected: Keep subtraction with a wider cast | Long.compare expresses the ordering directly and avoids overflow Confidence: high Scope-risk: narrow Directive: Do not compare long queue offsets by subtraction in POP checkpoint ordering Tested: mvn -pl store -Dtest=PopCheckPointTest test Tested: git diff --check Co-Authored-By: Claude --- .../rocketmq/store/pop/PopCheckPoint.java | 2 +- .../rocketmq/store/pop/PopCheckPointTest.java | 58 +++++++++++++++++++ 2 files changed, 59 insertions(+), 1 deletion(-) create mode 100644 store/src/test/java/org/apache/rocketmq/store/pop/PopCheckPointTest.java diff --git a/store/src/main/java/org/apache/rocketmq/store/pop/PopCheckPoint.java b/store/src/main/java/org/apache/rocketmq/store/pop/PopCheckPoint.java index 803ebc68957..bba3fd33141 100644 --- a/store/src/main/java/org/apache/rocketmq/store/pop/PopCheckPoint.java +++ b/store/src/main/java/org/apache/rocketmq/store/pop/PopCheckPoint.java @@ -212,6 +212,6 @@ public String toString() { @Override public int compareTo(PopCheckPoint o) { - return (int) (this.getStartOffset() - o.getStartOffset()); + return Long.compare(this.getStartOffset(), o.getStartOffset()); } } diff --git a/store/src/test/java/org/apache/rocketmq/store/pop/PopCheckPointTest.java b/store/src/test/java/org/apache/rocketmq/store/pop/PopCheckPointTest.java new file mode 100644 index 00000000000..af004761ebe --- /dev/null +++ b/store/src/test/java/org/apache/rocketmq/store/pop/PopCheckPointTest.java @@ -0,0 +1,58 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.store.pop; + +import java.util.TreeMap; +import org.junit.Test; + +import static org.assertj.core.api.Assertions.assertThat; + +public class PopCheckPointTest { + + private static PopCheckPoint build(long startOffset) { + PopCheckPoint ck = new PopCheckPoint(); + ck.setStartOffset(startOffset); + return ck; + } + + @Test + public void testCompareToWithoutOverflow() { + PopCheckPoint low = build(0); + PopCheckPoint high = build((long) Integer.MAX_VALUE + 2L); + + // The legacy (int)(a - b) comparator overflowed and reported the high + // offset as the smaller one. Long.compare must keep them ordered. + assertThat(low.compareTo(high)).isNegative(); + assertThat(high.compareTo(low)).isPositive(); + assertThat(low.compareTo(low)).isZero(); + } + + @Test + public void testTreeMapOrdersLargeStartOffsets() { + // PopReviveService keeps checkpoints in a TreeMap, + // so a broken compareTo corrupts ordering of large offsets. + TreeMap map = new TreeMap<>(); + PopCheckPoint high = build((long) Integer.MAX_VALUE + 2L); + PopCheckPoint low = build(0); + map.put(high, true); + map.put(low, true); + + assertThat(map.firstKey().getStartOffset()).isEqualTo(0L); + assertThat(map.lastKey().getStartOffset()).isEqualTo((long) Integer.MAX_VALUE + 2L); + } +}