Skip to content

[fix][client] Return flow-control permit when client seek by messageID - #26726

Open
programmerahul wants to merge 3 commits into
apache:masterfrom
programmerahul:fix/chunked-message-seek-startmessageid-permit-leak
Open

programmerahul wants to merge 3 commits into
apache:masterfrom
programmerahul:fix/chunked-message-seek-startmessageid-permit-leak

Conversation

@programmerahul

@programmerahul programmerahul commented Sep 27, 2026 •

Copy link
Copy Markdown
Contributor

Fixes #26727

Motivation

When a consumer seeks to (or is created with) a specific startMessageId, the broker re-dispatches the boundary message, which the client filters out in messageReceived() via the isSameEntry(msgId) && isPriorEntryIndex(...) block. That block released the payload and returned without calling increaseAvailablePermits(), so the dropped message's outstanding flow-control permit was never repaid (it is not delivered, so messageProcessed() never runs for it).

Each seek-to-messageId therefore leaks one permit. It is masked at normal receiver queue sizes (seek re-establishes flow control), but with receiverQueueSize=1 the single leaked permit exhausts the whole budget and the consumer stalls immediately after the seek -- the message after the seek target is never delivered. Not chunking-specific; applies to plain messages too.

Modifications

return the permit in the drop block. Adds a test (plain messages, receiverQueueSize=1) that stalls without the fix and passes with it.

Verifying this change

  • Make sure that the change passes the CI checks.

(Please pick either of the following options)

This change added tests and can be verified as follows:

  • *Added test-case so that no permit leak while seek operation *

Does this pull request potentially affect one of the following parts:

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

…geId boundary message

When a consumer seeks to (or is created with) a specific startMessageId, the broker
re-dispatches the boundary message, which the client filters out in messageReceived()
via the `isSameEntry(msgId) && isPriorEntryIndex(...)` block. That block released the
payload and returned without calling increaseAvailablePermits(), so the dropped
message's outstanding flow-control permit was never repaid (it is not delivered, so
messageProcessed() never runs for it).

Each seek-to-messageId therefore leaks one permit. It is masked at normal receiver
queue sizes (seek re-establishes flow control), but with receiverQueueSize=1 the single
leaked permit exhausts the whole budget and the consumer stalls immediately after the
seek -- the message after the seek target is never delivered. Not chunking-specific;
applies to plain messages too.

Fix: return the permit in the drop block. Adds a test (plain messages,
receiverQueueSize=1) that stalls without the fix and passes with it.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>

@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.

Thanks for tracking down this subtle flow-control leak and for adding a regression test with receiverQueueSize=1 that makes it visible.

I checked the premise and it holds: after a seek the consumer reconnects, ConsumerImpl.java:1071 resets the client permits and ConsumerImpl.java:975 sends a fresh full Flow, and only then does the broker re-dispatch the boundary entry. That entry consumes one broker-side permit and is dropped at ConsumerImpl.java:1541 before messageProcessed() can repay it, so the reconnect reset does not cover it. I also ran the new test with the added increaseAvailablePermits(cnx) removed and it fails (receive returns null); with the fix it passes. I found no double-count with chunked messages (non-last chunks are credited in processMessageChunk, the last chunk is repaid by the new line), zero-queue consumers, the payload-processor path or MultiTopicsConsumerImpl, and partially dropped batches are already refunded through the skipped-messages accounting. Two small suggestions inline.

// it here to avoid leaking a permit for the boundary message that a seek/startMessageId
// caused to be re-dispatched. (For a chunked message the non-last chunks were already
// credited at the top of this method; this repays the single remaining permit.)
increaseAvailablePermits(cnx);

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.

[QUALITY] refund numMessages rather than a fixed 1 so an undecryptable batch is repaid exactly

This block is also reached for an undecryptable batch (isMessageUndecryptable with the CONSUME crypto failure action), where the broker charged numMessages permits for the entry but only one is repaid here. Using increaseAvailablePermits(cnx, numMessages) (as the duplicate-ack drop at ConsumerImpl.java:1479 does) repays exactly what was consumed, and is identical for the plain and chunked cases where numMessages is 1. Not a regression, just a cheap way to close the remaining gap.

}

/**
* Seeking to a specific messageId re-dispatches the boundary message, which the consumer filters

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.

[QUALITY] the seek test is not chunking-specific and could live with the other seek tests

The test uses plain, non-chunked messages and the PR description says the bug is not chunking-specific, so it would be easier to find in a seek-focused class (for example the existing seek tests in SimpleProducerConsumerTest). The Javadoc that quotes the production code block will also go stale as ConsumerImpl evolves; a short description of the symptom would age better.

It would also be worth adding a chunked variant (a chunked message as the seek boundary with receiverQueueSize=1), since the code comment in ConsumerImpl claims that case is covered.

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] Return flow-control permit when client seek by messageID

2 participants