This is an automated email from the ASF dual-hosted git repository.
RongtongJin 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 e533b663fb [ISSUE #10972] Encode timer message propertiesString after
internal properties are cleared (#10974)
e533b663fb is described below
commit e533b663fb7c9cc30eb15c182671a94b681b0ac6
Author: Jiahua Wang <[email protected]>
AuthorDate: Thu Sep 3 15:54:43 2026 +0800
[ISSUE #10972] Encode timer message propertiesString after internal
properties are cleared (#10974)
Co-authored-by: wangjiahua.wjh <[email protected]>
---
.../rocketmq/store/timer/TimerMessageStore.java | 4 +++-
.../store/timer/TimerMessageStoreTest.java | 25 ++++++++++++++++++++++
2 files changed, 28 insertions(+), 1 deletion(-)
diff --git
a/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java
b/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java
index 24d2faa9b6..9608af3213 100644
--- a/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java
+++ b/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java
@@ -1251,7 +1251,6 @@ public class TimerMessageStore {
long tagsCodeValue =
MessageExtBrokerInner.tagsString2tagsCode(topicFilterType,
msgInner.getTags());
msgInner.setTagsCode(tagsCodeValue);
-
msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgExt.getProperties()));
msgInner.setSysFlag(msgExt.getSysFlag());
msgInner.setBornTimestamp(msgExt.getBornTimestamp());
@@ -1270,6 +1269,9 @@ public class TimerMessageStore {
MessageAccessor.clearProperty(msgInner,
MessageConst.PROPERTY_REAL_TOPIC);
MessageAccessor.clearProperty(msgInner,
MessageConst.PROPERTY_REAL_QUEUE_ID);
}
+ // Encode after the properties are finalized so that the wire data
stays consistent
+ // with the property map, aligning with
TimerMessageRocksDBStore#convertMessage.
+
msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgInner.getProperties()));
return msgInner;
}
diff --git
a/store/src/test/java/org/apache/rocketmq/store/timer/TimerMessageStoreTest.java
b/store/src/test/java/org/apache/rocketmq/store/timer/TimerMessageStoreTest.java
index fe1a1177c6..e62075cf67 100644
---
a/store/src/test/java/org/apache/rocketmq/store/timer/TimerMessageStoreTest.java
+++
b/store/src/test/java/org/apache/rocketmq/store/timer/TimerMessageStoreTest.java
@@ -179,6 +179,31 @@ public class TimerMessageStoreTest {
return null;
}
+ @Test
+ public void testConvertMessagePropertiesStringMatchesProperties() throws
Exception {
+ final TimerMessageStore timerMessageStore =
createTimerMessageStore(null, true);
+
+ MessageExtBrokerInner msgExt = buildMessage(3000,
"TimerTest_testConvertMessage", false);
+ MessageAccessor.putProperty(msgExt, MessageConst.PROPERTY_REAL_TOPIC,
msgExt.getTopic());
+ MessageAccessor.putProperty(msgExt,
MessageConst.PROPERTY_REAL_QUEUE_ID, "0");
+
msgExt.setPropertiesString(MessageDecoder.messageProperties2String(msgExt.getProperties()));
+ msgExt.setTopic(TimerMessageStore.TIMER_TOPIC);
+
+ // delivered message: internal properties are cleared from both the
map and the wire data
+ MessageExtBrokerInner delivered =
timerMessageStore.convertMessage(msgExt, false);
+ assertEquals("TimerTest_testConvertMessage", delivered.getTopic());
+
assertFalse(delivered.getPropertiesString().contains(MessageConst.PROPERTY_REAL_TOPIC));
+
assertFalse(delivered.getPropertiesString().contains(MessageConst.PROPERTY_REAL_QUEUE_ID));
+
assertEquals(MessageDecoder.messageProperties2String(delivered.getProperties()),
delivered.getPropertiesString());
+
+ // rolled message: keeps REAL_TOPIC and stays consistent between the
map and the wire data
+ MessageExtBrokerInner rolled =
timerMessageStore.convertMessage(msgExt, true);
+ assertEquals(TimerMessageStore.TIMER_TOPIC, rolled.getTopic());
+
assertTrue(rolled.getPropertiesString().contains(MessageConst.PROPERTY_REAL_TOPIC));
+
assertTrue(rolled.getPropertiesString().contains(MessageConst.PROPERTY_REAL_QUEUE_ID));
+
assertEquals(MessageDecoder.messageProperties2String(rolled.getProperties()),
rolled.getPropertiesString());
+ }
+
@Test
public void testPutTimerMessage() throws Exception {
Assume.assumeFalse(MixAll.isWindows());