Skip to content

[fix][Client] permit leak in chunked message discard path - #26661

Merged
lhotari merged 13 commits into
apache:masterfrom
programmerahul:fix/chunked-message-permit-leak
Oct 1, 2026
Merged

lhotari merged 13 commits into
apache:masterfrom
programmerahul:fix/chunked-message-permit-leak

Conversation

@programmerahul

@programmerahul programmerahul commented Sep 19, 2026 •

Copy link
Copy Markdown
Contributor

Fixes: #26675

Motivation

When a chunked message is discarded due to expiry, etc. then permit of last chunk is not returned. Which results in permit leak, and eventually after enough leaks( receiverQueue/2) the consumptions stops. This is because on broker the permit drops to 0, and client now cannot sent permit request, because it doesn't have enough messages.

this can be reproduced by setting the expireTimeOfIncompleteChunkedMessage to a very small value like 1ms, and on produce side, sending the chunked messages with chunks count of 10 or more. So that the chunks gets expired, and the last chunk of each message will leak 1 permit. and consumption will stop after (receiverQueue/2) leaks.

Modifications

if the last chunk of a chunked message is discarded, then increase the availablePermit on consumer

Verifying this change

  • Make sure that the change passes the CI checks.
    This change added tests and can be verified as follows:
    • Added a test-case for this

If the box was checked, please highlight the changes

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

@programmerahul

Copy link
Copy Markdown
Contributor Author

@lhotari please review

@programmerahul

Copy link
Copy Markdown
Contributor Author

@poorbarcode @void-ptr974 @merlimat @codelipenghui @BewareMyPower @Demogorgon314 can any please review this PR. This is a reproducible bug, and cause the consumption to stop

Comment thread pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java Outdated
@gauravAshok

gauravAshok commented Sep 23, 2026 •

Copy link
Copy Markdown

Claude's inputs:

There is a potentially similar issue at messageReceived:1541-1548 — after a chunked message is fully assembled, the isSameEntry(msgId) && isPriorEntryIndex(messageId.getEntryId()) branch drops it with uncompressedPayload.release(); return; and no credit. The last chunk's permit is never repaid. Compounding issue: by that point msgId has been rewritten to a ChunkMessageIdImpl (1531) while messageId.getEntryId() is still the raw last-chunk entry id, so the predicate mixes identities. Reachable on seek/startMessageId replay paths.

@gauravAshok

Copy link
Copy Markdown

NITS pointed by claude:

  1. Test asserts an exact permit count that depends on an unset config — fragile (main concern)

assertEquals(consumer.getAvailablePermits(), permitsAfterExpiry + numMessages).

increaseAvailablePermits (ConsumerImpl.java:1961-1973) resets the counter to 0 and sends a flow command once available >= getCurrentReceiverQueueSize() / 2. The test never sets receiverQueueSize, so it relies on the default (1000 → threshold 500). With 20 messages × 2 chunks the threshold isn't reached, so it passes — but the assertion is silently coupled to a default it doesn't state. Change the default, or bump numMessages, and it breaks for reasons unrelated to the bug.

Worse, the test is weaker than the bug it guards: the issue's failure mode is the reset-threshold never being reachable after receiverQueueSize/2 cumulative leaks. This test tears 20 messages with a 500 threshold — it never approaches the stall condition. It detects the missing credit arithmetically but would not catch the stall itself.

Minimal fix: pin .receiverQueueSize(N) explicitly so the assertion's premise is stated in the test. Stronger fix: set a small receiver queue, tear more than receiverQueueSize/2 messages, and assert that a subsequently-sent complete chunked message is still received — i.e. assert dispatch hasn't stalled. That's the actual user-visible symptom and it's robust to the reset behavior.

  1. Javadoc's closing sentence contradicts the test body

▎ "so the client's availablePermits must equal the total number of chunks delivered (2 * N)."

The test asserts permitsAfterExpiry + numMessages, not 2 * numMessages, and permitsAfterExpiry is read from the consumer rather than assumed. The 2 * N claim would also only hold if no threshold reset occurred. Drop that sentence or align it with what's asserted.

  1. autoAckOldestChunkedMessageOnQueueFull(true) and maxPendingChunkedMessage(100) are inert here

Only 20 pending contexts are created, so the queue-full eviction path never fires; expiry is what does the work. Harmless, but they suggest coverage the test doesn't have. Consider dropping them to keep the test's mechanism unambiguous.

  1. Test topic name is a non-@cleanup shared-namespace literal

persistent://my-property/my-ns/orphanChunkPermitLeak — fine and consistent with neighbors, just noting the class's other tests use the shared topicName field. No action needed.

@programmerahul

programmerahul commented Sep 27, 2026 •

Copy link
Copy Markdown
Contributor Author

Claude's inputs:

There is a potentially similar issue at messageReceived:1541-1548 — after a chunked message is fully assembled, the isSameEntry(msgId) && isPriorEntryIndex(messageId.getEntryId()) branch drops it with uncompressedPayload.release(); return; and no credit. The last chunk's permit is never repaid. Compounding issue: by that point msgId has been rewritten to a ChunkMessageIdImpl (1531) while messageId.getEntryId() is still the raw last-chunk entry id, so the predicate mixes identities. Reachable on seek/startMessageId replay paths.

This is applicable to normal messages also. It only leaks 1 permit when consumer uses consumer.seek(x); and isn't accumulated, since next consumer.seek(x) will reset permits. So I will raise a separate PR for this: #26726

@programmerahul

Copy link
Copy Markdown
Contributor Author

NITS pointed by claude:

  1. Test asserts an exact permit count that depends on an unset config — fragile (main concern)

assertEquals(consumer.getAvailablePermits(), permitsAfterExpiry + numMessages).

increaseAvailablePermits (ConsumerImpl.java:1961-1973) resets the counter to 0 and sends a flow command once available >= getCurrentReceiverQueueSize() / 2. The test never sets receiverQueueSize, so it relies on the default (1000 → threshold 500). With 20 messages × 2 chunks the threshold isn't reached, so it passes — but the assertion is silently coupled to a default it doesn't state. Change the default, or bump numMessages, and it breaks for reasons unrelated to the bug.

Worse, the test is weaker than the bug it guards: the issue's failure mode is the reset-threshold never being reachable after receiverQueueSize/2 cumulative leaks. This test tears 20 messages with a 500 threshold — it never approaches the stall condition. It detects the missing credit arithmetically but would not catch the stall itself.

Minimal fix: pin .receiverQueueSize(N) explicitly so the assertion's premise is stated in the test. Stronger fix: set a small receiver queue, tear more than receiverQueueSize/2 messages, and assert that a subsequently-sent complete chunked message is still received — i.e. assert dispatch hasn't stalled. That's the actual user-visible symptom and it's robust to the reset behavior.

  1. Javadoc's closing sentence contradicts the test body

▎ "so the client's availablePermits must equal the total number of chunks delivered (2 * N)."

The test asserts permitsAfterExpiry + numMessages, not 2 * numMessages, and permitsAfterExpiry is read from the consumer rather than assumed. The 2 * N claim would also only hold if no threshold reset occurred. Drop that sentence or align it with what's asserted.

  1. autoAckOldestChunkedMessageOnQueueFull(true) and maxPendingChunkedMessage(100) are inert here

Only 20 pending contexts are created, so the queue-full eviction path never fires; expiry is what does the work. Harmless, but they suggest coverage the test doesn't have. Consider dropping them to keep the test's mechanism unambiguous.

  1. Test topic name is a non-@cleanup shared-namespace literal

persistent://my-property/my-ns/orphanChunkPermitLeak — fine and consistent with neighbors, just noting the class's other tests use the shared topicName field. No action needed.

  1. Updated the testcase to simulate stall
  2. Removed those comments
  3. Remvoved these extra configs
  4. No change

@lhotari lhotari 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. Thanks for tracking this down and adding a regression test that reproduces the actual user-visible stall rather than just the internal counter.

I checked the specific claim — that a discarded last chunk leaks a flow-control permit — against the code and by mutation-testing the fix: with the increaseAvailablePermits(cnx) call at ConsumerImpl.java:1703-1705 reverted, testOrphanedLastChunkDoesNotStallDispatch fails exactly as described (dispatch stalls, assertNotNull fails); with the fix restored it passes. So the bug is real and this specific discard path is fixed.

I also checked the surrounding accounting: every non-last chunk is already credited on arrival (ConsumerImpl.java:1582-1584), so chunkedMessagesMap eviction (maxPendingChunkedMessage) and removeExpireIncompleteChunkedMessages don't need to credit anything themselves — they only ever hold non-last chunks, which were already credited. The new credit only fires in the orphan/out-of-order branch, not the duplicate-chunk branch (confirming the dead-code removal from earlier in this PR was safe), so there's no double-credit. processMessageChunk runs on the Netty I/O thread, the same thread that already calls increaseAvailablePermits for non-last chunks a few lines above, and the method itself uses the existing CAS-based accumulator — so the new call site doesn't add a new thread-safety concern. Batch-receive and multi-topics consumers go through the same per-partition ConsumerImpl/messageProcessed() path, so they're covered transparently.

Two earlier review comments look fully addressed at this head: the dead code in the duplicate-chunk branch is gone, and the test no longer asserts an exact permit count against an unstated receiver-queue default — it now asserts on the stall symptom itself, which is a stronger regression guard.

Two small, non-blocking notes below on adjacent code this PR already touches, plus a minor test-robustness suggestion.

@lhotari
lhotari merged commit 4c245b6 into apache:master Oct 1, 2026
43 checks passed
ascentstream-bot pushed a commit to ascentstream/pulsar that referenced this pull request Oct 1, 2026
Co-authored-by: Rahul Prasad <rahul.prasad3@flipkart.com>
(cherry picked from commit 4c245b6)
ascentstream-bot pushed a commit to ascentstream/pulsar that referenced this pull request Oct 2, 2026
Co-authored-by: Rahul Prasad <rahul.prasad3@flipkart.com>
(cherry picked from commit 4c245b6)
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.

[Bug] permit leak in chunked message discard path

3 participants