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

Reply via email to