This is an automated email from the ASF dual-hosted git repository. lhotari pushed a commit to branch branch-4.2 in repository https://gitbox.apache.org/repos/asf/pulsar.git
commit 964a723c7e07a3a07c93e01c819e08326ba89ac4 Author: 道君- Tao Jiuming <[email protected]> AuthorDate: Wed Jun 24 18:02:51 2026 +0800 [fix][broker] Guard BucketDelayedDeliveryTracker.nextDeliveryTime against empty queues (#26080) (cherry picked from commit 51e2b9a6bc9f739d3fc33b24821df15c588d4055) --- .../bucket/BucketDelayedDeliveryTracker.java | 5 +++ .../bucket/BucketDelayedDeliveryTrackerTest.java | 46 ++++++++++++++++++++++ 2 files changed, 51 insertions(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java index 7a1b95a3a4a..4b6121fd3e4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java @@ -593,6 +593,11 @@ public class BucketDelayedDeliveryTracker extends AbstractDelayedDeliveryTracker return sharedBucketPriorityQueue.peekN1(); } else if (sharedBucketPriorityQueue.isEmpty() && !lastMutableBucket.isEmpty()) { return lastMutableBucket.nextDeliveryTime(); + } else if (lastMutableBucket.isEmpty() && sharedBucketPriorityQueue.isEmpty()) { + // numberDelayedMessages can be > 0 while both queues are empty (e.g. remaining + // messages live in not-yet-loaded snapshot segments). Returning Long.MAX_VALUE + // signals "no imminent delivery" without throwing on the empty queues. + return Long.MAX_VALUE; } long timestamp = lastMutableBucket.nextDeliveryTime(); long bucketTimestamp = sharedBucketPriorityQueue.peekN1(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java index ee0ec0d9294..4becc00ab3f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java @@ -676,4 +676,50 @@ public class BucketDelayedDeliveryTrackerTest extends AbstractDeliveryTrackerTes ts.close(); } + + @Test + public void testGetScheduledMessagesWhenAllOrphaned() throws Exception { + // Reproduces IAE in nextDeliveryTime: when every delayed message lies below the + // mark-delete position, the filter in getScheduledMessages pops the in-memory + // messages without returning them. If the immutable bucket has additional messages + // still in storage (later snapshot segments), numberDelayedMessages stays > 0 + // while both the mutable bucket and the shared priority queue are empty. + // The trailing updateTimer -> nextDeliveryTime must not throw. + long firstLedgerId = 50L; + TrackerWithStorage ts = createTrackerWithMockLedger(firstLedgerId, 50); + + // Five delayed messages on the same orphaned ledger (ledgerId < firstLedgerId). + // They share a mutable bucket because seal requires a strictly greater ledgerId. + // Timestamps are 100ms apart so each lands in its own snapshot segment + // (timeStep=10ms); only the first segment is loaded into the shared queue at seal. + for (int i = 1; i <= 5; i++) { + ts.tracker.addMessage(1, i, i * 100); + } + // A new orphaned ledgerId triggers the seal, producing immutable bucket [1..1] + // with 5 messages across 5 segments; shared queue holds just the first segment. + ts.tracker.addMessage(2, 1, 600); + + Awaitility.await().untilAsserted(() -> + Assert.assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().values().stream() + .noneMatch(x -> x.merging))); + + // In strict deliver-at mode getCutoffTime() is just clock.millis(), so advancing the + // clock past the trigger message's deliverAt (600) is enough for + // moveScheduledMessageToSharedQueue to flush the mutable bucket into the shared queue. + ts.clockTime.set(700); + + // Both queues end up empty (filter pops the two in-memory messages), but + // numberDelayedMessages is still 4 (segments 2..5 remain in storage). + NavigableSet<Position> scheduledMessages = ts.tracker.getScheduledMessages(10); + assertTrue(scheduledMessages.isEmpty(), + "Orphaned messages should be filtered out, not returned"); + assertTrue(ts.tracker.getNumberOfDelayedMessages() > 0, + "Remaining storage-only messages should keep the counter > 0"); + + // hasMessageAvailable calls nextDeliveryTime while numberDelayedMessages > 0; + // it must not throw IAE. + assertFalse(ts.tracker.hasMessageAvailable()); + + ts.close(); + } }
