Skip to content

fix: wait for ACK_RESPONSE when acknowledging messages - #51

Open
danielbackes wants to merge 3 commits into
ikilobyte:mainfrom
danielbackes:bug/validate-ack-response
Open

danielbackes wants to merge 3 commits into
ikilobyte:mainfrom
danielbackes:bug/validate-ack-response

Conversation

@danielbackes

@danielbackes danielbackes commented Sep 17, 2026 •

Copy link
Copy Markdown

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 ACK command 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

  • Adds a request ID to the ACK command.
  • Waits for a response from the broker.
  • Queues MESSAGE commands received while waiting for the ACK_RESPONSE, so they can be processed later.
  • Throws an appropriate exception when the broker responds with a CLOSE_CONSUMER command.
  • Ensures that the received command is an ACK_RESPONSE.
  • Ensures that the response request ID matches the request ID sent with the command.
  • Throws an exception if the connection is closed.

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_RESPONSE

  • When PING commands arrived before the ACK_RESPONSE, the client handled them and continued waiting.
  • When MESSAGE commands arrived before the ACK_RESPONSE, the client queued the messages and continued waiting.
  • When the ACK_RESPONSE arrived, 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.

  • The connection loss was simulated by stopping the broker gracefully (docker compose down). When the worker tried to acknowledge the message, a RuntimeException was thrown: The consumer was closed before the message acknowledgment was confirmed.
  • The connection loss was simulated by stopping the broker abruptly (docker compose kill). When the worker tried to acknowledge the message, an IOException was thrown: socket is closed.
  • The connection loss occurred when the broker closed the connection after the keepalive timeout. When the worker tried to acknowledge the message, an IOException was thrown: socket is closed.

Test timeout

When the client did not receive an ACK_RESPONSE within its configured ackTimeout, it threw a RuntimeException: Timed out waiting for ACK response.

Fixes #49

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
danielbackes marked this pull request as ready for review September 17, 2026 00:14
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Bug: Message acknowledgment succeeds without receiving ACK_RESPONSE

1 participant