Skip to content

[fix][broker] Check deliverAt before containsMessage in bucket addMessage - #26230

Merged
dao-jun merged 1 commit into
apache:masterfrom
nodece:fix-bucket-addmessage-expired-check-order
Jul 23, 2026
Merged

dao-jun merged 1 commit into
apache:masterfrom
nodece:fix-bucket-addmessage-expired-check-order

Conversation

@nodece

@nodece nodece commented Jul 23, 2026

Copy link
Copy Markdown
Member

Motivation

BucketDelayedDeliveryTracker uses lazy loading for sealed bucket segments. During segment loading, expired messages are filtered based on their deliverAt timestamp and are not added into sharedBucketPriorityQueue.

delayedIndexBitMap is persisted with the bucket metadata and restored from BookKeeper. It is used for position deduplication and does not indicate whether a message is currently loaded into sharedBucketPriorityQueue.

Therefore, during lazy loading, a message can temporarily exist in delayedIndexBitMap but not in sharedBucketPriorityQueue.

The current BucketDelayedDeliveryTracker.addMessage implementation checks containsMessage before checking whether the message has already expired (deliverAt <= cutoffTime).

This causes an issue when an expired message exists only in the delayed index:

  1. The dispatcher reads the message from the managed ledger.
  2. containsMessage returns true because the position exists in delayedIndexBitMap.
  3. addMessage returns true, and the dispatcher skips the message.
  4. The message is not in sharedBucketPriorityQueue, so delayed delivery cannot return it either.

The message is not lost from storage, but it can remain skip:

ManagedLedger:
    message exists

delayedIndexBitMap:
    position exists (dedup index)

sharedBucketPriorityQueue:
    position missing

Modifications

  • Change the check order in BucketDelayedDeliveryTracker.addMessage.

    • Check deliverAt <= cutoffTime before containsMessage.
    • Expired messages return false so they can be delivered immediately by the dispatcher.
    • Deduplication is only applied to messages that are still delayed.
  • Add testExpiredIndexedMessageReturnsFalse.

    • Add a message with future deliverAt.
    • Advance the clock beyond deliverAt.
    • Re-add the same position.
    • Assert that addMessage returns false.
  • Add testRecoverThenExpireAddMessage.

    • Seal a bucket.
    • Recover the tracker.
    • Advance the clock beyond the message delivery time.
    • Re-add a recovered position that exists in delayedIndexBitMap but has not been loaded into sharedBucketPriorityQueue.
    • Assert that addMessage returns false.

In addition, the dispatcher side will add an explicit expiration check to make the delivery path more robust. The tracker-side fix ensures that an expired message is not incorrectly blocked by the delayed index when it should bypass delayed delivery tracking.

@dao-jun dao-jun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@dao-jun
dao-jun merged commit 094f270 into apache:master Jul 23, 2026
44 checks passed
@lhotari

lhotari commented Jul 23, 2026

Copy link
Copy Markdown
Member

@nodece Thanks for the fix — the reordering itself looks right (it matches the InMemoryDelayedDeliveryTracker semantics, where expiry is decided before dedup). One follow-up observation from a local review (Claude Fable 5 + Codex gpt-5.6-sol — both flagged this independently):

When addMessage now returns false for an expired message that is still tracked, the entry is not removed from the tracker. If the position is still in sharedBucketPriorityQueue (e.g. the dispatcher re-reads an already-due tracked message during replay/rewind), the dispatcher delivers it immediately, and getScheduledMessages() later returns the same position again when it pops it from the queue — a duplicate-delivery window. numberDelayedMessages also transiently overcounts until the stale entry pops.

Making the false path remove the entry isn't a one-liner, though: TripleLongPriorityQueue has no positional removal (only head pop), and clearing just the bitmap bit wouldn't be enough because the pop path adds the position and decrements the counter unconditionally:

positions.add(PositionFactory.create(ledgerId, entryId));
sharedBucketPriorityQueue.pop();
removeIndexBit(ledgerId, entryId);
--n;
numberDelayedMessages.decrementAndGet();

A complete fix would clear the index bit in addMessage before returning false and make the pop path skip positions whose bit is already gone, mirroring the existing firstLiveLedgerId branch:

if (firstLiveLedgerId != null && ledgerId < firstLiveLedgerId) {
sharedBucketPriorityQueue.pop();
if (removeIndexBit(ledgerId, entryId)) {
numberDelayedMessages.decrementAndGet();
}
continue;
}

Given Pulsar's at-least-once semantics this may well be acceptable as-is (an occasional redelivery beats the pre-fix indefinite skip), but it seemed worth recording — and possibly worth handling together with the dispatcher-side expiration check mentioned as follow-up work in the PR description.

@nodece
nodece deleted the fix-bucket-addmessage-expired-check-order branch July 24, 2026 03:29
@nodece

nodece commented Jul 24, 2026

Copy link
Copy Markdown
Member Author

@lhotari Good catch. I only considered the case where the expired message could be skipped indefinitely, but I missed the potential duplicate-delivery window caused by leaving the stale entry in the priority queue.

The current Bucket implementation has become quite complicated over time — different bucket types, index maintenance, storage lifecycle, and message counters are mixed together, which makes these edge cases harder to reason about.

I'm currently refactoring this part and will try to address this issue as part of the redesign. The direction is to make the index state and queue state more consistent, and avoid relying on implicit cleanup during pop operations.

lhotari pushed a commit that referenced this pull request Jul 25, 2026
lhotari pushed a commit that referenced this pull request Jul 25, 2026
sandeep-ctds pushed a commit to datastax/pulsar that referenced this pull request Jul 31, 2026
sandeep-ctds pushed a commit to datastax/pulsar that referenced this pull request Jul 31, 2026
sandeep-ctds pushed a commit to datastax/pulsar that referenced this pull request Jul 31, 2026
nodece added a commit to ascentstream/pulsar that referenced this pull request Aug 28, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants