Skip to content

[fix][broker] Fix dispatcher permit accounting race with recycled batch index acks - #26809

Merged
merlimat merged 1 commit into
apache:masterfrom
merlimat:mmerli/dispatcher-acked-index-permit-race
Oct 2, 2026
Merged

merlimat merged 1 commit into
apache:masterfrom
merlimat:mmerli/dispatcher-acked-index-permit-race

Conversation

@merlimat

@merlimat merlimat commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

Motivation

A Shared or Key_Shared dispatcher keeps totalAvailablePermits, which should always equal the sum of its
consumers' permits. It uses this total to decide how much to read. On every send, the dispatcher and the consumer
each subtract the number of messages delivered. Batch-index acks are on by default. When a partially acknowledged
batch is redelivered, only its unacked messages are delivered, so both subtract messages in batch - acked indexes.
The acked-index count comes from a pooled EntryBatchIndexesAcks object.

The dispatcher read that count after Consumer.sendMessages() returned. By then sendMessages() has handed the
object to the consumer connection's event loop. That task writes the messages and then recycles the object, which
clears every slot. The dispatcher runs on a different thread, so nothing orders the two:

dispatcher thread                                  consumer connection event loop
-----------------                                  ------------------------------
fill batchIndexesAcks (4 of 10 indexes acked)
consumer.sendMessages(...)
  consumer permits -= 10 - 4 = 6   (correct)
  post the write task  ──────────────────────────▶ write the 6 messages
                                                   batchIndexesAcks.recycle()  (clears all slots)
batchIndexesAcks.getTotalAckedIndexCount() → 0
totalAvailablePermits -= 10 - 0    (should be 6)

When the event loop gets there first, the dispatcher subtracts every message in the batch. Its total then falls
below the sum of the consumers' permits, by 4 in this example. The error adds up with each occurrence and is only
reset when all consumers disconnect. Until then, reads are smaller than the consumers can take. The dispatcher is
also reading, without synchronization, an object that another thread is clearing.

The same pattern is in three places:

  • PersistentDispatcherMultipleConsumers.trySendMessagesToConsumers
  • PersistentDispatcherMultipleConsumers.sendChunkedMessagesToConsumers, used for any read that contains a chunk
  • PersistentStickyKeyDispatcherMultipleConsumers.trySendMessagesToConsumers

Consumer-level permits are correct, because Consumer.sendMessages() reads the count before the handoff.

Modifications

  • In the three dispatch paths, read batchIndexesAcks.getTotalAckedIndexCount() into a local variable before
    calling sendMessages(). Use it for the permit update and the debug log, so the dispatcher no longer touches the
    object after the handoff. Without a race the result is unchanged, because sendMessages() never modifies the
    object.
  • Add SharedDispatcherPermitAccountingTest.testRedeliveredPartiallyAckedBatchDoesNotLosePermits. The test's
    connection stub recycles batchIndexesAcks before sendMessages() returns. That is the losing interleaving, so
    the race reproduces every time. The test dispatches a redelivered batch of 10 messages with 4 acked indexes
    through each of the three paths. It then checks that the dispatcher total equals the sum of the consumers'
    permits. Without the fix, the total is 4 permits short on every path.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • SharedDispatcherPermitAccountingTest.testRedeliveredPartiallyAckedBatchDoesNotLosePermits covers the Shared,
    Shared chunked and Key_Shared paths. It fails without the fix (for example, expected: 194 but was: 190) and
    passes with it.
  • The existing dispatcher, Key_Shared, batch-index-ack and Shared chunking suites pass, including
    PersistentDispatcherMultipleConsumersTest, PersistentStickyKeyDispatcherMultipleConsumersTest,
    KeySharedSubscriptionTest, BatchMessageWithBatchIndexLevelTest, BatchMessageIndexAckTest and
    MessageChunkingSharedTest.

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

…ch index acks

The Shared and Key_Shared dispatchers read
EntryBatchIndexesAcks.getTotalAckedIndexCount() after Consumer.sendMessages()
returned, to update the dispatcher-level TOTAL_AVAILABLE_PERMITS. By then
sendMessages has handed the object to the consumer channel's event loop,
which recycles it (nulling every slot) once the entries are written. When
the dispatcher runs on another thread, the read can see 0 or a partial
count, so the dispatcher subtracts all the messages of a partially acked
batch instead of only the unacked ones, and its permit total drifts below
the sum of the consumers' permits.

Read the count once before calling sendMessages in the three dispatch
paths (Shared, Shared chunked and Key_Shared), and use that value for the
permit update and the debug log. Consumer.sendMessages already computes
its own count before the handoff, so consumer-level permits were not
affected.

@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

@merlimat
merlimat merged commit 1046481 into apache:master Oct 2, 2026
44 checks passed
lhotari pushed a commit that referenced this pull request Oct 2, 2026
…ch index acks (#26809)

(cherry picked from commit 1046481)

[branch-4.x] Also applied to the classic Shared and Key_Shared dispatchers
(PersistentDispatcherMultipleConsumersClassic, PersistentStickyKeyDispatcherMultipleConsumersClassic),
which were removed on master by #26687. The test covers both implementations.
lhotari pushed a commit that referenced this pull request Oct 2, 2026
…ch index acks (#26809)

(cherry picked from commit 1046481)

[branch-4.x] Also applied to the classic Shared and Key_Shared dispatchers
(PersistentDispatcherMultipleConsumersClassic, PersistentStickyKeyDispatcherMultipleConsumersClassic),
which were removed on master by #26687. The test covers both implementations.
dao-jun pushed a commit to ascentstream/pulsar that referenced this pull request Oct 5, 2026
lhotari pushed a commit that referenced this pull request Oct 5, 2026
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.

3 participants