[fix][Client] permit leak in chunked message discard path - #26661
Conversation
|
@lhotari please review |
|
@poorbarcode @void-ptr974 @merlimat @codelipenghui @BewareMyPower @Demogorgon314 can any please review this PR. This is a reproducible bug, and cause the consumption to stop |
|
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. |
|
NITS pointed by claude:
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.
▎ "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.
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.
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. |
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 |
|
lhotari
left a comment
There was a problem hiding this comment.
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.
Co-authored-by: Rahul Prasad <rahul.prasad3@flipkart.com> (cherry picked from commit 4c245b6)
Co-authored-by: Rahul Prasad <rahul.prasad3@flipkart.com> (cherry picked from commit 4c245b6)
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
This change added tests and can be verified as follows:
If the box was checked, please highlight the changes