fix: wait for ACK_RESPONSE when acknowledging messages - #51
Open
danielbackes wants to merge 3 commits into
Open
danielbackes wants to merge 3 commits into
danielbackes wants to merge 3 commits into
Conversation
Update the ack method to wait for the ackResponse before proceeding. Any messages decoded while waiting are safely added to the queue to prevent data loss. Refs ikilobyte#49
ack() previously reused the pollFrame()/sendAckRequest() machinery to transparently reconnect and resend the CommandAck if the connection dropped mid-wait. That's the wrong recovery model for this specific call: a dropped connection during ack() can mean the broker already decided this consumer is gone and redelivered the message to another consumer on the same subscription (Shared/Key_Shared). Silently reconnecting and resending can't undo that redelivery, and it hides from the caller the fact that their handler may now be running twice. ack() now throws IOException immediately on a dropped connection, regardless of ConsumerOptions::getReconnectPolicy(), so the caller gets an honest, timely signal instead of the library quietly retrying underneath it. receive() is unaffected and keeps its own inline reconnect handling for the ordinary receive-loop case. Removes the now-unused pollFrame()/sendAckRequest() helpers. Refs ikilobyte#49
danielbackes
marked this pull request as ready for review
September 17, 2026 00:14
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
This PR ensures that an acknowledgment is considered successful only after the Pulsar broker returns the corresponding
ACK_RESPONSE.Previously, the consumer only wrote the
ACKcommand to the socket. A successful write confirmed that the local socket accepted the command, but it did not confirm that the broker received and processed the acknowledgment.Changes
ACKcommand.MESSAGEcommands received while waiting for theACK_RESPONSE, so they can be processed later.CLOSE_CONSUMERcommand.ACK_RESPONSE.Motivation
If the broker becomes unavailable after delivering a message but before processing its acknowledgment, the client could previously return successfully from the acknowledgment operation.
The application could then consider the message acknowledged even though Pulsar still regarded it as unacknowledged, causing the message to be redelivered.
Tests Performed
Project tests
Manual tests
Test whether the client can wait for an
ACK_RESPONSEPINGcommands arrived before theACK_RESPONSE, the client handled them and continued waiting.MESSAGEcommands arrived before theACK_RESPONSE, the client queued the messages and continued waiting.ACK_RESPONSEarrived, the client stopped waiting and considered the acknowledgment successful.Test whether the client raises an exception when the connection is lost
In each scenario, the message was delivered successfully. A connection loss was then simulated while the handler was running.
docker compose down). When the worker tried to acknowledge the message, aRuntimeExceptionwas thrown:The consumer was closed before the message acknowledgment was confirmed.docker compose kill). When the worker tried to acknowledge the message, anIOExceptionwas thrown:socket is closed.IOExceptionwas thrown:socket is closed.Test timeout
When the client did not receive an
ACK_RESPONSEwithin its configuredackTimeout, it threw aRuntimeException:Timed out waiting for ACK response.Fixes #49