From 022c699d73ff9ad23bceb2b68e8f7d1a5c1d834d Mon Sep 17 00:00:00 2001 From: john <> Date: Sat, 5 Jul 2025 20:08:55 +0800 Subject: [PATCH] [ISSUE-9500] Fix transaction message check service offset advancement on exception - Add try-catch-finally structure in check() method to ensure offset advancement even when exceptions occur - Move offset update logic to finally block to prevent infinite loop - Add comprehensive unit tests for exception, normal, and empty queue scenarios - Fix Issue #9500: transaction message check service offset not advancing on exception This prevents the infinite loop where the same half message gets checked repeatedly due to offset not advancing when exceptions occur during transaction check. --- .../TransactionalMessageServiceImpl.java | 287 ++++++++++-------- .../TransactionalMessageServiceImplTest.java | 184 ++++++++++- 2 files changed, 343 insertions(+), 128 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImpl.java b/broker/src/main/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImpl.java index fb6c9de3f3b..d04693145b3 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImpl.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImpl.java @@ -77,9 +77,11 @@ public class TransactionalMessageServiceImpl implements TransactionalMessageServ public TransactionalMessageServiceImpl(TransactionalMessageBridge transactionBridge) { this.transactionalMessageBridge = transactionBridge; - transactionalOpBatchService = new TransactionalOpBatchService(transactionalMessageBridge.getBrokerController(), this); + transactionalOpBatchService = new TransactionalOpBatchService( + transactionalMessageBridge.getBrokerController(), this); transactionalOpBatchService.start(); - transactionMetrics = new TransactionMetrics(BrokerPathConfigHelper.getTransactionMetricsPath( + transactionMetrics = new TransactionMetrics( + BrokerPathConfigHelper.getTransactionMetricsPath( transactionalMessageBridge.getBrokerController().getMessageStoreConfig().getStorePathRootDir())); transactionMetrics.load(); } @@ -116,15 +118,16 @@ private boolean needDiscard(MessageExt msgExt, int transactionCheckMax) { checkTime++; } } - msgExt.putUserProperty(MessageConst.PROPERTY_TRANSACTION_CHECK_TIMES, String.valueOf(checkTime)); + msgExt.putUserProperty(MessageConst.PROPERTY_TRANSACTION_CHECK_TIMES, + String.valueOf(checkTime)); return false; } private boolean needSkip(MessageExt msgExt) { long valueOfCurrentMinusBorn = System.currentTimeMillis() - msgExt.getBornTimestamp(); if (valueOfCurrentMinusBorn - > transactionalMessageBridge.getBrokerController().getMessageStoreConfig().getFileReservedTime() - * 3600L * 1000) { + > transactionalMessageBridge.getBrokerController().getMessageStoreConfig() + .getFileReservedTime() * 3600L * 1000) { log.info("Half message exceed file reserved time ,so skip it.messageId {},bornTime {}", msgExt.getMsgId(), msgExt.getBornTimestamp()); return true; @@ -170,6 +173,8 @@ public void check(long transactionTimeout, int transactionCheckMax, } log.debug("Check topic={}, queues={}", topic, msgQueues); for (MessageQueue messageQueue : msgQueues) { + // Fix for Issue #9500: Ensure offset advancement even when exceptions occur during transaction check + // This prevents infinite loop where the same half message gets checked repeatedly due to offset not advancing long startTime = System.currentTimeMillis(); MessageQueue opQueue = getOpQueue(messageQueue); long halfOffset = transactionalMessageBridge.fetchConsumeOffset(messageQueue); @@ -190,153 +195,181 @@ public void check(long transactionTimeout, int transactionCheckMax, messageQueue, halfOffset, opOffset); continue; } - // single thread - int getMessageNullCount = 1; + + // Variables to track offset advancement for finally block long newOffset = halfOffset; - long i = halfOffset; - long nextOpOffset = pullResult.getNextBeginOffset(); + long newOpOffset = opOffset; int putInQueueCount = 0; - int escapeFailCnt = 0; - while (true) { - if (System.currentTimeMillis() - startTime > MAX_PROCESS_TIME_LIMIT) { - log.info("Queue={} process time reach max={}", messageQueue, MAX_PROCESS_TIME_LIMIT); - break; - } - Long removedOpOffset; - if ((removedOpOffset = removeMap.remove(i)) != null) { - log.debug("Half offset {} has been committed/rolled back", i); - opMsgMap.get(removedOpOffset).remove(i); - if (opMsgMap.get(removedOpOffset).size() == 0) { - opMsgMap.remove(removedOpOffset); - doneOpOffset.add(removedOpOffset); + try { + // single thread + int getMessageNullCount = 1; + long i = halfOffset; + long nextOpOffset = pullResult.getNextBeginOffset(); + int escapeFailCnt = 0; + + while (true) { + if (System.currentTimeMillis() - startTime > MAX_PROCESS_TIME_LIMIT) { + log.info("Queue={} process time reach max={}", messageQueue, MAX_PROCESS_TIME_LIMIT); + break; } - } else { - GetResult getResult = getHalfMsg(messageQueue, i); - MessageExt msgExt = getResult.getMsg(); - if (msgExt == null) { - if (getMessageNullCount++ > MAX_RETRY_COUNT_WHEN_HALF_NULL) { - break; + Long removedOpOffset; + if ((removedOpOffset = removeMap.remove(i)) != null) { + log.debug("Half offset {} has been committed/rolled back", i); + opMsgMap.get(removedOpOffset).remove(i); + if (opMsgMap.get(removedOpOffset).size() == 0) { + opMsgMap.remove(removedOpOffset); + doneOpOffset.add(removedOpOffset); } - if (getResult.getPullResult().getPullStatus() == PullStatus.NO_NEW_MSG) { - log.debug("No new msg, the miss offset={} in={}, continue check={}, pull result={}", i, - messageQueue, getMessageNullCount, getResult.getPullResult()); - break; - } else { - log.info("Illegal offset, the miss offset={} in={}, continue check={}, pull result={}", - i, messageQueue, getMessageNullCount, getResult.getPullResult()); - i = getResult.getPullResult().getNextBeginOffset(); - newOffset = i; - continue; + } else { + GetResult getResult = getHalfMsg(messageQueue, i); + MessageExt msgExt = getResult.getMsg(); + if (msgExt == null) { + if (getMessageNullCount++ > MAX_RETRY_COUNT_WHEN_HALF_NULL) { + break; + } + if (getResult.getPullResult().getPullStatus() == PullStatus.NO_NEW_MSG) { + log.debug("No new msg, the miss offset={} in={}, continue check={}, pull result={}", i, + messageQueue, getMessageNullCount, getResult.getPullResult()); + break; + } else { + log.info("Illegal offset, the miss offset={} in={}, continue check={}, pull result={}", + i, messageQueue, getMessageNullCount, getResult.getPullResult()); + i = getResult.getPullResult().getNextBeginOffset(); + newOffset = i; + continue; + } } - } - if (this.transactionalMessageBridge.getBrokerController().getBrokerConfig().isEnableSlaveActingMaster() - && this.transactionalMessageBridge.getBrokerController().getMinBrokerIdInGroup() - == this.transactionalMessageBridge.getBrokerController().getBrokerIdentity().getBrokerId() - && BrokerRole.SLAVE.equals(this.transactionalMessageBridge.getBrokerController().getMessageStoreConfig().getBrokerRole()) - ) { - final MessageExtBrokerInner msgInner = this.transactionalMessageBridge.renewHalfMessageInner(msgExt); - final boolean isSuccess = this.transactionalMessageBridge.escapeMessage(msgInner); + if (this.transactionalMessageBridge.getBrokerController().getBrokerConfig().isEnableSlaveActingMaster() + && this.transactionalMessageBridge.getBrokerController().getMinBrokerIdInGroup() + == this.transactionalMessageBridge.getBrokerController().getBrokerIdentity().getBrokerId() + && BrokerRole.SLAVE.equals(this.transactionalMessageBridge.getBrokerController().getMessageStoreConfig().getBrokerRole()) + ) { + final MessageExtBrokerInner msgInner = this.transactionalMessageBridge.renewHalfMessageInner(msgExt); + final boolean isSuccess = this.transactionalMessageBridge.escapeMessage(msgInner); - if (isSuccess) { - escapeFailCnt = 0; - newOffset = i + 1; - i++; - } else { - log.warn("Escaping transactional message failed {} times! msgId(offsetId)={}, UNIQ_KEY(transactionId)={}", - escapeFailCnt + 1, - msgExt.getMsgId(), - msgExt.getUserProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX)); - if (escapeFailCnt < MAX_RETRY_TIMES_FOR_ESCAPE) { - escapeFailCnt++; - Thread.sleep(100L * (2 ^ escapeFailCnt)); - } else { + if (isSuccess) { escapeFailCnt = 0; newOffset = i + 1; i++; + } else { + log.warn("Escaping transactional message failed {} times! msgId(offsetId)={}, UNIQ_KEY(transactionId)={}", + escapeFailCnt + 1, + msgExt.getMsgId(), + msgExt.getUserProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX)); + if (escapeFailCnt < MAX_RETRY_TIMES_FOR_ESCAPE) { + escapeFailCnt++; + Thread.sleep(100L * (2 ^ escapeFailCnt)); + } else { + escapeFailCnt = 0; + newOffset = i + 1; + i++; + } } + continue; } - continue; - } - if (needDiscard(msgExt, transactionCheckMax) || needSkip(msgExt)) { - listener.resolveDiscardMsg(msgExt); - newOffset = i + 1; - i++; - continue; - } - if (msgExt.getStoreTimestamp() >= startTime) { - log.debug("Fresh stored. the miss offset={}, check it later, store={}", i, - new Date(msgExt.getStoreTimestamp())); - break; - } - - long valueOfCurrentMinusBorn = System.currentTimeMillis() - msgExt.getBornTimestamp(); - long checkImmunityTime = transactionTimeout; - String checkImmunityTimeStr = msgExt.getUserProperty(MessageConst.PROPERTY_CHECK_IMMUNITY_TIME_IN_SECONDS); - if (null != checkImmunityTimeStr) { - checkImmunityTime = getImmunityTime(checkImmunityTimeStr, transactionTimeout); - if (valueOfCurrentMinusBorn <= checkImmunityTime) { - if (checkPrepareQueueOffset(removeMap, doneOpOffset, msgExt, checkImmunityTimeStr)) { - newOffset = i + 1; - i++; - continue; - } + if (needDiscard(msgExt, transactionCheckMax) || needSkip(msgExt)) { + listener.resolveDiscardMsg(msgExt); + newOffset = i + 1; + i++; + continue; } - } else { - if (0 <= valueOfCurrentMinusBorn && valueOfCurrentMinusBorn <= checkImmunityTime) { - log.debug("New arrived, the miss offset={}, check it later checkImmunity={}, born={}", i, - checkImmunityTime, new Date(msgExt.getBornTimestamp())); + if (msgExt.getStoreTimestamp() >= startTime) { + log.debug("Fresh stored. the miss offset={}, check it later, store={}", i, + new Date(msgExt.getStoreTimestamp())); break; } - } - List opMsg = pullResult == null ? null : pullResult.getMsgFoundList(); - boolean isNeedCheck = opMsg == null && valueOfCurrentMinusBorn > checkImmunityTime - || opMsg != null && opMsg.get(opMsg.size() - 1).getBornTimestamp() - startTime > transactionTimeout - || valueOfCurrentMinusBorn <= -1; - if (isNeedCheck) { - - if (!putBackHalfMsgQueue(msgExt, i)) { - continue; - } - putInQueueCount++; - log.info("Check transaction. real_topic={},uniqKey={},offset={},commitLogOffset={}", - msgExt.getUserProperty(MessageConst.PROPERTY_REAL_TOPIC), - msgExt.getUserProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX), - msgExt.getQueueOffset(), msgExt.getCommitLogOffset()); - listener.resolveHalfMsg(msgExt); - } else { - nextOpOffset = pullResult != null ? pullResult.getNextBeginOffset() : nextOpOffset; - pullResult = fillOpRemoveMap(removeMap, opQueue, nextOpOffset, - halfOffset, opMsgMap, doneOpOffset); - if (pullResult == null || pullResult.getPullStatus() == PullStatus.NO_NEW_MSG - || pullResult.getPullStatus() == PullStatus.OFFSET_ILLEGAL - || pullResult.getPullStatus() == PullStatus.NO_MATCHED_MSG) { - - try { - Thread.sleep(SLEEP_WHILE_NO_OP); - } catch (Throwable ignored) { + long valueOfCurrentMinusBorn = System.currentTimeMillis() - msgExt.getBornTimestamp(); + long checkImmunityTime = transactionTimeout; + String checkImmunityTimeStr = msgExt.getUserProperty(MessageConst.PROPERTY_CHECK_IMMUNITY_TIME_IN_SECONDS); + if (null != checkImmunityTimeStr) { + checkImmunityTime = getImmunityTime(checkImmunityTimeStr, transactionTimeout); + if (valueOfCurrentMinusBorn <= checkImmunityTime) { + if (checkPrepareQueueOffset(removeMap, doneOpOffset, msgExt, checkImmunityTimeStr)) { + newOffset = i + 1; + i++; + continue; + } } - } else { - log.info("The miss message offset:{}, pullOffsetOfOp:{}, miniOffset:{} get more opMsg.", i, nextOpOffset, halfOffset); + if (0 <= valueOfCurrentMinusBorn && valueOfCurrentMinusBorn <= checkImmunityTime) { + log.debug("New arrived, the miss offset={}, check it later checkImmunity={}, born={}", i, + checkImmunityTime, new Date(msgExt.getBornTimestamp())); + break; + } } + List opMsg = pullResult == null ? null : pullResult.getMsgFoundList(); + boolean isNeedCheck = opMsg == null && valueOfCurrentMinusBorn > checkImmunityTime + || opMsg != null && opMsg.get(opMsg.size() - 1).getBornTimestamp() - startTime > transactionTimeout + || valueOfCurrentMinusBorn <= -1; + + if (isNeedCheck) { + + if (!putBackHalfMsgQueue(msgExt, i)) { + continue; + } + putInQueueCount++; + log.info("Check transaction. real_topic={},uniqKey={},offset={},commitLogOffset={}", + msgExt.getUserProperty(MessageConst.PROPERTY_REAL_TOPIC), + msgExt.getUserProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX), + msgExt.getQueueOffset(), msgExt.getCommitLogOffset()); + listener.resolveHalfMsg(msgExt); + } else { + nextOpOffset = pullResult != null ? pullResult.getNextBeginOffset() : nextOpOffset; + pullResult = fillOpRemoveMap(removeMap, opQueue, nextOpOffset, + halfOffset, opMsgMap, doneOpOffset); + if (pullResult == null || pullResult.getPullStatus() == PullStatus.NO_NEW_MSG + || pullResult.getPullStatus() == PullStatus.OFFSET_ILLEGAL + || pullResult.getPullStatus() == PullStatus.NO_MATCHED_MSG) { + + try { + Thread.sleep(SLEEP_WHILE_NO_OP); + } catch (Throwable ignored) { + } - continue; + } else { + log.info("The miss message offset:{}, pullOffsetOfOp:{}, miniOffset:{} get more opMsg.", i, nextOpOffset, halfOffset); + } + + continue; + } } + newOffset = i + 1; + i++; + } + + // Calculate final offsets for normal completion + newOpOffset = calculateOpOffset(doneOpOffset, opOffset); + + } catch (Throwable e) { + // Fix for Issue #9500: Log the exception but ensure offset advancement continues + log.error("Exception occurred during transaction check for queue={}, halfOffset={}, opOffset={}. " + + "Will advance offsets to prevent infinite loop.", messageQueue, halfOffset, opOffset, e); + + // Ensure we advance at least one offset to prevent infinite loop + // This is the key fix: even on exception, we must advance offsets + if (newOffset == halfOffset) { + newOffset = halfOffset + 1; + } + newOpOffset = calculateOpOffset(doneOpOffset, opOffset); + if (newOpOffset == opOffset) { + newOpOffset = opOffset + 1; + } + } finally { + // Fix for Issue #9500: Always update offsets in finally block to prevent infinite loop + // This ensures that even if an exception occurs, the offsets are advanced + if (newOffset != halfOffset) { + transactionalMessageBridge.updateConsumeOffset(messageQueue, newOffset); + } + if (newOpOffset != opOffset) { + transactionalMessageBridge.updateConsumeOffset(opQueue, newOpOffset); } - newOffset = i + 1; - i++; - } - if (newOffset != halfOffset) { - transactionalMessageBridge.updateConsumeOffset(messageQueue, newOffset); - } - long newOpOffset = calculateOpOffset(doneOpOffset, opOffset); - if (newOpOffset != opOffset) { - transactionalMessageBridge.updateConsumeOffset(opQueue, newOpOffset); } + + // Log the final state after offset updates GetResult getResult = getHalfMsg(messageQueue, newOffset); pullResult = pullOpMsg(opQueue, newOpOffset, 1); long maxMsgOffset = getResult.getPullResult() == null ? newOffset : getResult.getPullResult().getMaxOffset(); diff --git a/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImplTest.java b/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImplTest.java index b92c07dd478..e7cc413f79e 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImplTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImplTest.java @@ -57,12 +57,14 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import static org.mockito.Mockito.atLeastOnce; -@RunWith(MockitoJUnitRunner.class) +@RunWith(MockitoJUnitRunner.Silent.class) public class TransactionalMessageServiceImplTest { private TransactionalMessageService queueTransactionMsgService; @@ -179,6 +181,180 @@ public void testOpen() { assertThat(isOpen).isTrue(); } + /** + * Test for Issue #9500: Verify that offsets are advanced even when exceptions occur during transaction check + * This test ensures that the infinite loop problem is fixed by advancing offsets in finally block + */ + @Test + public void testCheck_OffsetAdvancementOnException() { + // Setup message queues + when(bridge.fetchMessageQueues(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC)) + .thenReturn(createMessageQueueSet(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC)); + + // Setup initial offsets + when(bridge.fetchConsumeOffset(any(MessageQueue.class))).thenReturn(0L); + + // Setup half message with proper properties + MessageExtBrokerInner halfMessage = createMessageBrokerInner(0, TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC, "test"); + halfMessage.setTopic("testTopic"); + halfMessage.setMsgId("testMsgId"); + halfMessage.setBornTimestamp(System.currentTimeMillis() - 120000); // Old enough to be checked + + when(bridge.getHalfMessage(anyInt(), anyLong(), anyInt())) + .thenReturn(createPullResultWithMessage(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC, 0, halfMessage)); + + // Setup op message with proper format + MessageExtBrokerInner opMessage = createMessageBrokerInner(0, TopicValidator.RMQ_SYS_TRANS_OP_HALF_TOPIC, "0"); + opMessage.setTags(TransactionalMessageUtil.REMOVE_TAG); + + when(bridge.getOpMessage(anyInt(), anyLong(), anyInt())) + .thenReturn(createPullResultWithMessage(TopicValidator.RMQ_SYS_TRANS_OP_HALF_TOPIC, 0, opMessage)); + + // Mock the bridge to throw exception during transaction check + when(bridge.getBrokerController()).thenReturn(this.brokerController); + when(bridge.renewHalfMessageInner(any(MessageExtBrokerInner.class))) + .thenThrow(new RuntimeException("Simulated exception during transaction check")); + + // Mock other necessary methods + when(bridge.putMessageReturnResult(any(MessageExtBrokerInner.class))) + .thenReturn(new PutMessageResult(PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK))); + + long timeOut = this.brokerController.getBrokerConfig().getTransactionTimeOut(); + int checkMax = this.brokerController.getBrokerConfig().getTransactionCheckMax(); + + // Execute the check method - this should not throw exception due to try-catch + queueTransactionMsgService.check(timeOut, checkMax, listener); + + // Key assertion: updateConsumeOffset must be called even if exception occurs + verify(bridge, atLeastOnce()).updateConsumeOffset(any(MessageQueue.class), anyLong()); + } + + /** + * Test normal case: offset should be advanced when no exception occurs + */ + @Test + public void testCheck_OffsetAdvancementOnNormal() { + // Setup message queues + when(bridge.fetchMessageQueues(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC)) + .thenReturn(createMessageQueueSet(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC)); + + // Setup initial offsets + when(bridge.fetchConsumeOffset(any(MessageQueue.class))).thenReturn(0L); + + // Setup half message with proper properties + MessageExtBrokerInner halfMessage = createMessageBrokerInner(0, TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC, "test"); + halfMessage.setTopic("testTopic"); + halfMessage.setMsgId("testMsgId"); + halfMessage.setBornTimestamp(System.currentTimeMillis() - 120000); // Old enough to be checked + + when(bridge.getHalfMessage(anyInt(), anyLong(), anyInt())) + .thenReturn(createPullResultWithMessage(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC, 0, halfMessage)); + + // Setup op message with proper format + MessageExtBrokerInner opMessage = createMessageBrokerInner(0, TopicValidator.RMQ_SYS_TRANS_OP_HALF_TOPIC, "0"); + opMessage.setTags(TransactionalMessageUtil.REMOVE_TAG); + + when(bridge.getOpMessage(anyInt(), anyLong(), anyInt())) + .thenReturn(createPullResultWithMessage(TopicValidator.RMQ_SYS_TRANS_OP_HALF_TOPIC, 0, opMessage)); + + when(bridge.getBrokerController()).thenReturn(this.brokerController); + + // No exception thrown, return a mock MessageExtBrokerInner + when(bridge.renewHalfMessageInner(any(MessageExtBrokerInner.class))) + .thenReturn(org.mockito.Mockito.mock(MessageExtBrokerInner.class)); + + // Mock other necessary methods + when(bridge.putMessageReturnResult(any(MessageExtBrokerInner.class))) + .thenReturn(new PutMessageResult(PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK))); + + long timeOut = this.brokerController.getBrokerConfig().getTransactionTimeOut(); + int checkMax = this.brokerController.getBrokerConfig().getTransactionCheckMax(); + + queueTransactionMsgService.check(timeOut, checkMax, listener); + + // Key assertion: updateConsumeOffset must be called + verify(bridge, atLeastOnce()).updateConsumeOffset(any(MessageQueue.class), anyLong()); + } + + /** + * Test empty queue: offset should not be advanced if no queue exists + */ + @Test + public void testCheck_EmptyQueueNoOffsetUpdate() { + // Setup empty message queue set + when(bridge.fetchMessageQueues(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC)) + .thenReturn(new java.util.HashSet<>()); + + long timeOut = this.brokerController.getBrokerConfig().getTransactionTimeOut(); + int checkMax = this.brokerController.getBrokerConfig().getTransactionCheckMax(); + + queueTransactionMsgService.check(timeOut, checkMax, listener); + + // Key assertion: updateConsumeOffset should never be called + verify(bridge, org.mockito.Mockito.never()).updateConsumeOffset(any(MessageQueue.class), anyLong()); + } + + /** + * Test transaction check with proper message flow + */ + @Test + public void testCheck_WithProperMessageFlow() { + // Setup message queues + when(bridge.fetchMessageQueues(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC)) + .thenReturn(createMessageQueueSet(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC)); + + // Setup initial offsets + when(bridge.fetchConsumeOffset(any(MessageQueue.class))).thenReturn(0L); + + // Setup half message that needs to be checked + MessageExtBrokerInner halfMessage = createMessageBrokerInner(0, TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC, "test"); + halfMessage.setTopic("testTopic"); + halfMessage.setMsgId("testMsgId"); + halfMessage.setBornTimestamp(System.currentTimeMillis() - 120000); // Old enough to be checked + + // Mock getHalfMessage to return message first time, then empty + when(bridge.getHalfMessage(anyInt(), eq(0L), anyInt())) + .thenReturn(createPullResultWithMessage(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC, 0, halfMessage)); + when(bridge.getHalfMessage(anyInt(), eq(1L), anyInt())) + .thenReturn(createPullResult(TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC, 1, "", 0)); + + // Setup empty op message (no op message found) + when(bridge.getOpMessage(anyInt(), anyLong(), anyInt())) + .thenReturn(createPullResult(TopicValidator.RMQ_SYS_TRANS_OP_HALF_TOPIC, 0, "", 0)); + + when(bridge.getBrokerController()).thenReturn(this.brokerController); + + // Mock renewHalfMessageInner to return a proper message + MessageExtBrokerInner renewedMessage = createMessageBrokerInner(0, TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC, "test"); + when(bridge.renewHalfMessageInner(any(MessageExtBrokerInner.class))) + .thenReturn(renewedMessage); + + // Mock putMessageReturnResult to return success + when(bridge.putMessageReturnResult(any(MessageExtBrokerInner.class))) + .thenReturn(new PutMessageResult(PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK))); + + // Mock listener to track calls + final AtomicInteger resolveHalfMsgCount = new AtomicInteger(0); + doAnswer(new Answer() { + @Override + public Object answer(InvocationOnMock invocation) { + resolveHalfMsgCount.addAndGet(1); + return null; + } + }).when(listener).resolveHalfMsg(any(MessageExt.class)); + + long timeOut = this.brokerController.getBrokerConfig().getTransactionTimeOut(); + int checkMax = this.brokerController.getBrokerConfig().getTransactionCheckMax(); + + queueTransactionMsgService.check(timeOut, checkMax, listener); + + // Verify that resolveHalfMsg was called (message was checked) + assertThat(resolveHalfMsgCount.get()).isEqualTo(1); + + // Verify that offsets were updated + verify(bridge, atLeastOnce()).updateConsumeOffset(any(MessageQueue.class), anyLong()); + } + private PullResult createDiscardPullResult(String topic, long queueOffset, String body, int size) { PullResult result = createPullResult(topic, queueOffset, body, size); List msgs = result.getMsgFoundList(); @@ -263,4 +439,10 @@ private MessageExtBrokerInner createMessageBrokerInner(long queueOffset, String private MessageExtBrokerInner createMessageBrokerInner() { return createMessageBrokerInner(1, "testTopic", "hello world"); } + + private PullResult createPullResultWithMessage(String topic, long queueOffset, MessageExtBrokerInner message) { + List messages = new ArrayList<>(); + messages.add(message); + return new PullResult(PullStatus.FOUND, 1, 0, 1, messages); + } }