Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down