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); + } }