From a71b1bf0296f0301ee75d156109359e036b135e3 Mon Sep 17 00:00:00 2001 From: Gerard <213731635+GerardGao@users.noreply.github.com> Date: Wed, 26 Aug 2026 13:38:44 +0800 Subject: [PATCH 1/2] fix(broker): reject timer delays that overflow the absolute delivery time transformTimerMessage converted relative timer delays with unchecked long multiplication and addition, so a very large delay overflowed into a past timestamp, bypassed the future-delay validation, and took the immediate-message path. Use Math.multiplyExact/Math.addExact so overflow returns WHEEL_TIMER_MSG_ILLEGAL. Fixes #10872 --- .../rocketmq/broker/util/HookUtils.java | 11 ++- .../util/HookUtilsTimerOverflowTest.java | 77 +++++++++++++++++++ 2 files changed, 84 insertions(+), 4 deletions(-) create mode 100644 broker/src/test/java/org/apache/rocketmq/broker/util/HookUtilsTimerOverflowTest.java diff --git a/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java b/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java index 94be46ea405..60ec502b135 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java @@ -206,19 +206,22 @@ private static PutMessageResult transformTimerMessage(BrokerController brokerCon //do transform int delayLevel = msg.getDelayTimeLevel(); long deliverMs; + long now = System.currentTimeMillis(); try { if (msg.getProperty(MessageConst.PROPERTY_TIMER_DELAY_SEC) != null) { - deliverMs = System.currentTimeMillis() + Long.parseLong(msg.getProperty(MessageConst.PROPERTY_TIMER_DELAY_SEC)) * 1000; + long delaySec = Long.parseLong(msg.getProperty(MessageConst.PROPERTY_TIMER_DELAY_SEC)); + deliverMs = Math.addExact(now, Math.multiplyExact(delaySec, 1000L)); } else if (msg.getProperty(MessageConst.PROPERTY_TIMER_DELAY_MS) != null) { - deliverMs = System.currentTimeMillis() + Long.parseLong(msg.getProperty(MessageConst.PROPERTY_TIMER_DELAY_MS)); + long delayMs = Long.parseLong(msg.getProperty(MessageConst.PROPERTY_TIMER_DELAY_MS)); + deliverMs = Math.addExact(now, delayMs); } else { deliverMs = Long.parseLong(msg.getProperty(MessageConst.PROPERTY_TIMER_DELIVER_MS)); } } catch (Exception e) { return new PutMessageResult(PutMessageStatus.WHEEL_TIMER_MSG_ILLEGAL, null); } - if (deliverMs > System.currentTimeMillis()) { - if (delayLevel <= 0 && deliverMs - System.currentTimeMillis() > brokerController.getMessageStoreConfig().getTimerMaxDelaySec() * 1000L) { + if (deliverMs > now) { + if (delayLevel <= 0 && deliverMs - now > brokerController.getMessageStoreConfig().getTimerMaxDelaySec() * 1000L) { return new PutMessageResult(PutMessageStatus.WHEEL_TIMER_MSG_ILLEGAL, null); } diff --git a/broker/src/test/java/org/apache/rocketmq/broker/util/HookUtilsTimerOverflowTest.java b/broker/src/test/java/org/apache/rocketmq/broker/util/HookUtilsTimerOverflowTest.java new file mode 100644 index 00000000000..3e470a63f47 --- /dev/null +++ b/broker/src/test/java/org/apache/rocketmq/broker/util/HookUtilsTimerOverflowTest.java @@ -0,0 +1,77 @@ +/* + * 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.broker.util; + +import org.apache.rocketmq.broker.BrokerController; +import org.apache.rocketmq.common.message.MessageAccessor; +import org.apache.rocketmq.common.message.MessageConst; +import org.apache.rocketmq.common.message.MessageExtBrokerInner; +import org.apache.rocketmq.store.PutMessageResult; +import org.apache.rocketmq.store.PutMessageStatus; +import org.apache.rocketmq.store.config.MessageStoreConfig; +import org.junit.Assert; +import org.junit.Test; +import org.mockito.Mockito; + +public class HookUtilsTimerOverflowTest { + + private BrokerController newTimerEnabledController() { + BrokerController brokerController = Mockito.mock(BrokerController.class); + MessageStoreConfig messageStoreConfig = new MessageStoreConfig(); + messageStoreConfig.setTimerWheelEnable(true); + Mockito.when(brokerController.getMessageStoreConfig()).thenReturn(messageStoreConfig); + return brokerController; + } + + private MessageExtBrokerInner newTimerMessage(String key, String value) { + MessageExtBrokerInner msg = new MessageExtBrokerInner(); + msg.setTopic("OverflowTestTopic"); + MessageAccessor.putProperty(msg, key, value); + return msg; + } + + // Regression for https://github.com/apache/rocketmq/issues/10872: an + // overflowing relative delay used to wrap around into a past timestamp, + // bypass the future-delay validation, and take the immediate-message path. + @Test + public void testTimerDelaySecOverflowRejected() { + PutMessageResult result = HookUtils.handleScheduleMessage(newTimerEnabledController(), + newTimerMessage(MessageConst.PROPERTY_TIMER_DELAY_SEC, String.valueOf(Long.MAX_VALUE))); + Assert.assertNotNull(result); + Assert.assertEquals(PutMessageStatus.WHEEL_TIMER_MSG_ILLEGAL, result.getPutMessageStatus()); + } + + @Test + public void testTimerDelayMsOverflowRejected() { + PutMessageResult result = HookUtils.handleScheduleMessage(newTimerEnabledController(), + newTimerMessage(MessageConst.PROPERTY_TIMER_DELAY_MS, String.valueOf(Long.MAX_VALUE))); + Assert.assertNotNull(result); + Assert.assertEquals(PutMessageStatus.WHEEL_TIMER_MSG_ILLEGAL, result.getPutMessageStatus()); + } + + // Covers the future-delivery validation branch with a delay that does not + // overflow but exceeds the configured maximum. + @Test + public void testTimerDelaySecExceedsMaxRejected() { + BrokerController brokerController = newTimerEnabledController(); + brokerController.getMessageStoreConfig().setTimerMaxDelaySec(1); + PutMessageResult result = HookUtils.handleScheduleMessage(brokerController, + newTimerMessage(MessageConst.PROPERTY_TIMER_DELAY_SEC, "10")); + Assert.assertNotNull(result); + Assert.assertEquals(PutMessageStatus.WHEEL_TIMER_MSG_ILLEGAL, result.getPutMessageStatus()); + } +} From f2427de60f6906b74adac9f2addabbab7997ae7a Mon Sep 17 00:00:00 2001 From: gerardgao Date: Thu, 27 Aug 2026 11:03:33 +0800 Subject: [PATCH 2/2] test(broker): cover timer delay overflow and success paths Cover the remaining transformTimerMessage branches so the exact-math overflow checks and the millisecond success path are exercised: - Long.MAX_VALUE / 1000 overflows addExact after multiplyExact succeeds. - A valid PROPERTY_TIMER_DELAY_MS rewrites the message to the wheel-timer topic with PROPERTY_TIMER_OUT_MS set. Split the nested exact-math expression in HookUtils so each overflow source is independently testable and reported. --- .../rocketmq/broker/util/HookUtils.java | 3 +- .../util/HookUtilsTimerOverflowTest.java | 28 +++++++++++++++++++ 2 files changed, 30 insertions(+), 1 deletion(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java b/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java index 60ec502b135..06eb90b6eaf 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java @@ -210,7 +210,8 @@ private static PutMessageResult transformTimerMessage(BrokerController brokerCon try { if (msg.getProperty(MessageConst.PROPERTY_TIMER_DELAY_SEC) != null) { long delaySec = Long.parseLong(msg.getProperty(MessageConst.PROPERTY_TIMER_DELAY_SEC)); - deliverMs = Math.addExact(now, Math.multiplyExact(delaySec, 1000L)); + deliverMs = Math.multiplyExact(delaySec, 1000L); + deliverMs = Math.addExact(now, deliverMs); } else if (msg.getProperty(MessageConst.PROPERTY_TIMER_DELAY_MS) != null) { long delayMs = Long.parseLong(msg.getProperty(MessageConst.PROPERTY_TIMER_DELAY_MS)); deliverMs = Math.addExact(now, delayMs); diff --git a/broker/src/test/java/org/apache/rocketmq/broker/util/HookUtilsTimerOverflowTest.java b/broker/src/test/java/org/apache/rocketmq/broker/util/HookUtilsTimerOverflowTest.java index 3e470a63f47..9c107b9b7aa 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/util/HookUtilsTimerOverflowTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/util/HookUtilsTimerOverflowTest.java @@ -23,6 +23,7 @@ import org.apache.rocketmq.store.PutMessageResult; import org.apache.rocketmq.store.PutMessageStatus; import org.apache.rocketmq.store.config.MessageStoreConfig; +import org.apache.rocketmq.store.timer.TimerMessageStore; import org.junit.Assert; import org.junit.Test; import org.mockito.Mockito; @@ -63,6 +64,33 @@ public void testTimerDelayMsOverflowRejected() { Assert.assertEquals(PutMessageStatus.WHEEL_TIMER_MSG_ILLEGAL, result.getPutMessageStatus()); } + // Long.MAX_VALUE / 1000 keeps multiplyExact within range but overflows the + // subsequent addExact against the current time, exercising its throw path. + @Test + public void testTimerDelaySecAddExactOverflowRejected() { + PutMessageResult result = HookUtils.handleScheduleMessage(newTimerEnabledController(), + newTimerMessage(MessageConst.PROPERTY_TIMER_DELAY_SEC, String.valueOf(Long.MAX_VALUE / 1000))); + Assert.assertNotNull(result); + Assert.assertEquals(PutMessageStatus.WHEEL_TIMER_MSG_ILLEGAL, result.getPutMessageStatus()); + } + + // Exercises the successful timer transformation for a millisecond delay: + // the message is rewritten to the wheel-timer topic with the normalized + // delivery time stored in PROPERTY_TIMER_OUT_MS. + @Test + public void testTimerDelayMsSuccessPath() { + BrokerController brokerController = newTimerEnabledController(); + TimerMessageStore timerMessageStore = Mockito.mock(TimerMessageStore.class); + Mockito.when(brokerController.getTimerMessageStore()).thenReturn(timerMessageStore); + Mockito.when(timerMessageStore.isReject(Mockito.anyLong())).thenReturn(false); + + MessageExtBrokerInner msg = newTimerMessage(MessageConst.PROPERTY_TIMER_DELAY_MS, "1000"); + PutMessageResult result = HookUtils.handleScheduleMessage(brokerController, msg); + Assert.assertNull(result); + Assert.assertNotNull(msg.getProperty(MessageConst.PROPERTY_TIMER_OUT_MS)); + Assert.assertEquals(TimerMessageStore.TIMER_TOPIC, msg.getTopic()); + } + // Covers the future-delivery validation branch with a delay that does not // overflow but exceeds the configured maximum. @Test