From 72c9405e7ef0a27ff2bbbb55557d3455d5172745 Mon Sep 17 00:00:00 2001 From: zjncs <18910855655@163.com> Date: Wed, 9 Sep 2026 18:04:56 +0800 Subject: [PATCH] [ISSUE #B4] Compare the producer group null-safely in EndTransactionProcessor checkPrepareMessage dereferenced the PROPERTY_PRODUCER_GROUP value read from the prepared message without a null check, so a half message that lacks the property (written by an old client or restored from a store without it) failed the whole commit/rollback request with an NPE instead of the intended 'producer group wrong' rejection. Use Objects.equals, which treats a missing property as a mismatch. Signed-off-by: zjncs <18910855655@163.com> --- .../broker/processor/EndTransactionProcessor.java | 5 ++++- .../processor/EndTransactionProcessorTest.java | 15 +++++++++++++++ 2 files changed, 19 insertions(+), 1 deletion(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/EndTransactionProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/EndTransactionProcessor.java index c568f1e4bf2..dff9d46d8cc 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/EndTransactionProcessor.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/EndTransactionProcessor.java @@ -245,7 +245,10 @@ private RemotingCommand checkPrepareMessage(MessageExt msgExt, EndTransactionReq return response; } final String pgroupRead = msgExt.getProperty(MessageConst.PROPERTY_PRODUCER_GROUP); - if (!pgroupRead.equals(requestHeader.getProducerGroup())) { + // The prepared message may lack the producer group property (messages + // written by old clients or restored stores); compare null-safely and + // reject instead of failing the request with an NPE. + if (!Objects.equals(pgroupRead, requestHeader.getProducerGroup())) { response.setCode(ResponseCode.SYSTEM_ERROR); response.setRemark("The producer group wrong"); return response; diff --git a/broker/src/test/java/org/apache/rocketmq/broker/processor/EndTransactionProcessorTest.java b/broker/src/test/java/org/apache/rocketmq/broker/processor/EndTransactionProcessorTest.java index 7b2d4bdf26a..e9ed1d45588 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/processor/EndTransactionProcessorTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/processor/EndTransactionProcessorTest.java @@ -202,6 +202,21 @@ public void testProcessRequestAllowsMissingTopicForCompatibility() throws Remoti assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS); } + @Test + public void testProcessRequestRejectsPreparedMessageWithoutProducerGroup() throws RemotingCommandException { + OperationResult prepareResult = createResponse(ResponseCode.SUCCESS); + MessageAccessor.clearProperty(prepareResult.getPrepareMessage(), MessageConst.PROPERTY_PRODUCER_GROUP); + when(transactionMsgService.commitMessage(any(EndTransactionRequestHeader.class))).thenReturn(prepareResult); + + RemotingCommand request = createEndTransactionMsgCommand(MessageSysFlag.TRANSACTION_COMMIT_TYPE, false); + RemotingCommand response = endTransactionProcessor.processRequest(handlerContext, request); + + assertThat(response.getCode()).isEqualTo(ResponseCode.SYSTEM_ERROR); + assertThat(response.getRemark()).contains("producer group wrong"); + verify(messageStore, never()).putMessage(any(MessageExtBrokerInner.class)); + verify(transactionMsgService, never()).deletePrepareMessage(any(MessageExt.class)); + } + private MessageExt createDefaultMessageExt() { MessageExt messageExt = new MessageExt(); messageExt.setMsgId("12345678");