Skip to content

[fix][broker] Prevent stale read completions from stranding Failover subscriptions - #26174

Merged
lhotari merged 3 commits into
apache:masterfrom
alexandrebrg:fix-failover-stuck-read-26164
Jul 23, 2026
Merged

[fix][broker] Prevent stale read completions from stranding Failover subscriptions#26174
lhotari merged 3 commits into
apache:masterfrom
alexandrebrg:fix-failover-stuck-read-26164

Conversation

@alexandrebrg

@alexandrebrg alexandrebrg commented Jul 10, 2026

Copy link
Copy Markdown
Contributor

Motivation

On a single-active-consumer (Failover / Exclusive) subscription, the dispatcher tracks its one
outstanding read with a plain volatile boolean havePendingRead, with no notion of which read
the flag refers to — while the cursor state it mirrors is generation-guarded
((op, opReadId) in ManagedCursorImpl).

internalRedeliverUnacknowledgedMessages clears havePendingRead, rewinds, and re-arms a fresh
tail-wait read without cancelling a read that is already in flight
(cursor.cancelPendingReadRequest() can only cancel a waiting op). When the disowned read
later completes, readEntriesComplete / readEntriesFailed run havePendingRead = false as
their first statement — clearing a flag that now describes the newer armed read. Result: the
cursor holds an armed waitingReadOp while the dispatcher believes no read is pending.

The desync is benign while the cursor remains in waitingCursors. It becomes permanent on the
next last-consumer disconnect: cancelPendingRead() short-circuits on the false flag
(if (havePendingRead && ...)) so the armed op survives, and
PersistentSubscription#removeConsumer then removes the cursor from waitingCursors. From
there the subscription is stuck forever:

  • every subsequent arm CAS-fails on the leftover op (ConcurrentWaitCallbackException), which
    readEntriesFailed returns on without rescheduling;
  • addWaitingCursor is only reachable after a successful arm, so the cursor can never re-enter
    waitingCursors;
  • every publish's notifyCursors() polls a queue the cursor is not in.

The consumer stays connected with permits, the backlog grows, msgOutCounter stays 0, and only
a topic unload recovers — exactly the production signature reported in #26164
(waitingReadOp=true, pendingReadOps=0, waitingCursorsCount=0, cursor state=Open). On
Failover, every redeliver form (ack-timeout, negative ack, explicit redeliver, reconnect epoch
bump) funnels into internalRedeliverUnacknowledgedMessages, so any of them racing an in-flight
read completion can mint the desync.

Modifications

  • Add a monotonic readOpEpoch (plain long, guarded by the dispatcher monitor) to
    PersistentDispatcherSingleActiveConsumer, minted alongside the only
    havePendingRead = true write in readMoreEntries.
  • Capture the epoch into the read-completion continuation; readEntriesComplete and
    readEntriesFailed check it first and ignore stale completions (entries released, no
    dispatcher state mutated, DEBUG log carrying both the stale and the current epoch), so a
    superseded read can no longer clear the newer read's flag. This restores the invariant
    armed ⇒ havePendingRead == true, which in turn makes the disconnect-path short-circuit
    safe: whenever an op is armed, cancelPendingRead() now actually cancels it.
  • Add PersistentDispatcherSingleActiveConsumerStuckReadTest: a deterministic reproduction
    using the real ManagedLedgerImpl + ManagedCursorImpl + PersistentSubscription +
    PersistentDispatcherSingleActiveConsumer, with the dispatcher's ordered executor replaced
    by a manual task board so the test chooses between the two production-possible arrival
    orders of the racing tasks (the stale completion is posted by the BK completion chain, the
    redeliver by the client-command thread — both orders occur in production). The consumer
    lifecycle runs through the real PersistentSubscription#addConsumer /
    removeConsumer(Consumer, boolean), so dispatcher removal, deactivateCursor() and
    removeWaitingCursor() happen in the production order; the spied dispatcher is injected
    through the existing production seam PersistentSubscription#reuseOrCreateDispatcher.
    Both tests assert the externally observable contract of [Bug][broker] Failover subscription stops dispatching forever after active-consumer handover #26164: after the last consumer
    disconnects and a fresh consumer reconnects, a message published afterwards must actually
    reach that consumer's sendMessages and advance the cursor — the internal-state tuple is
    kept only as a secondary diagnostic. The repro method fails on current master and passes
    with the fix; a negative control runs the same scenario in the benign order and passes on
    both. Fidelity notes (the one remaining filterEntriesForConsumer stub, needed because the
    test publishes raw bytes rather than serialized Pulsar messages; the telescoped ack;
    newEntriesCheckDelayInMillis=0 as a determinism pin) are documented in the class Javadoc.

Notes for reviewers:

  • The fix is completion-side. One narrow corner intentionally keeps master's behavior: if the
    redeliver's re-arm bails (e.g. no permits), the disowned read's completion still counts as
    current and dispatches — a benign at-least-once duplicate, unchanged from today.
  • PersistentDispatcherMultipleConsumers / ...Classic keep their own havePendingRead;
    their redeliver model is structurally different (replay queue, no rewind + immediate re-arm),
    so they are deliberately out of scope here.
  • readEntriesFailed is @VisibleForTesting; its two test call sites were updated for the
    added parameter. No public API is touched.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • PersistentDispatcherSingleActiveConsumerStuckReadTest#testFailoverConsumerStuckWhenRedeliverRacesInFlightReadCompletion
    — deterministic reproduction of the strand: fails on master, passes with this fix (verified
    over repeated forced re-executions, plus the inverse: re-fails when the fix is reverted).
  • PersistentDispatcherSingleActiveConsumerStuckReadTest#testRedeliverAfterReadCompletionDoesNotStrandCursor
    — negative control (same staging, benign task order): passes with and without the fix,
    showing the arrival order alone is the trigger.
  • Run:
    ./gradlew :pulsar-broker:test --tests 'org.apache.pulsar.broker.service.persistent.PersistentDispatcherSingleActiveConsumerStuckReadTest' -PtestFailFast=false -PtestRetryCount=0
  • Existing coverage exercised and green: PersistentDispatcherSingleActiveConsumerTest,
    PersistentTopicTest, FailoverSubscriptionTest (incl. the waitingCursors invariant tests
    from [fix][broker] Fix Broker OOM due to too many waiting cursors and reuse a recycled OpReadEntry incorrectly #24551), MessageRedeliveryTest.

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

Fixes #26164

…subscriptions

On a single-active-consumer (Failover/Exclusive) subscription,
internalRedeliverUnacknowledgedMessages clears havePendingRead and re-arms a
fresh tail-wait read without cancelling a read that is already in flight (an
in-flight read cannot be cancelled). When the disowned read later completes,
readEntriesComplete/readEntriesFailed unconditionally clear havePendingRead
again while the fresh op stays armed on the cursor: the dispatcher now believes
no read is pending while the cursor holds an armed waiting read.

The desync is benign while the cursor remains in waitingCursors, but on the
next last-consumer disconnect cancelPendingRead() short-circuits on the false
flag (the armed op survives) while PersistentSubscription#removeConsumer
removes the cursor from waitingCursors. From then on the subscription is
permanently stuck: every re-arm CAS-fails on the leftover op
(ConcurrentWaitCallbackException is returned without rescheduling) and every
publish's notifyCursors() misses the cursor. Only a topic unload recovers it.

The cursor side is already generation-guarded ((op, opReadId)); the
dispatcher's boolean mirror is not — that asymmetry is the defect.

Fixes apache#26164

Signed-off-by: Alexandre Burgoni <alexandre.burgoni@clever.cloud>
@alexandrebrg
alexandrebrg force-pushed the fix-failover-stuck-read-26164 branch from f12a486 to 8f682f0 Compare July 13, 2026 08:44
@alexandrebrg
alexandrebrg marked this pull request as ready for review July 13, 2026 09:46

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

please address the review comments and update the stale PR description

@lhotari lhotari added the triage/lhotari/important lhotari's triaging label for important issues or PRs label Jul 22, 2026
- Assert the externally observable contract of apache#26164: both tests now
  reconnect a fresh consumer after the last-consumer disconnect and
  require the next published message to actually reach
  Consumer#sendMessages and advance the cursor read position. The
  internal state tuple is kept only as a secondary diagnostic, so the
  regression test is no longer coupled to one fix strategy.
- Drive the consumer lifecycle through the real
  PersistentSubscription#addConsumer and
  removeConsumer(Consumer, boolean) instead of replaying two selected
  steps by hand, so dispatcher removal, deactivateCursor() and
  removeWaitingCursor() run in the production order. The spied
  dispatcher is injected through the existing
  PersistentSubscription#reuseOrCreateDispatcher test seam; no
  production code was added for the test.
- Replace the hard-coded line-number references in the test Javadoc
  with durable method-name anchors, which cannot silently rot as the
  referenced files change.
- Record both the stale and the current readOpEpoch on the two
  stale-read debug log paths, so a suspected recurrence in production
  shows which epochs disagreed.
@lhotari
lhotari dismissed Denovo1998’s stale review July 23, 2026 19:25

review comments have been addressed

@lhotari
lhotari merged commit 0679b0d into apache:master Jul 23, 2026
43 checks passed
lhotari pushed a commit that referenced this pull request Jul 25, 2026
…subscriptions (#26174)

Signed-off-by: Alexandre Burgoni <alexandre.burgoni@clever.cloud>
(cherry picked from commit 0679b0d)
lhotari pushed a commit that referenced this pull request Jul 25, 2026
…subscriptions (#26174)

Signed-off-by: Alexandre Burgoni <alexandre.burgoni@clever.cloud>
(cherry picked from commit 0679b0d)
(cherry picked from commit 8cbbb6e)
sandeep-ctds pushed a commit to datastax/pulsar that referenced this pull request Jul 31, 2026
…subscriptions (apache#26174)

Signed-off-by: Alexandre Burgoni <alexandre.burgoni@clever.cloud>
(cherry picked from commit 0679b0d)
(cherry picked from commit 8cbbb6e)
sandeep-ctds pushed a commit to datastax/pulsar that referenced this pull request Jul 31, 2026
…subscriptions (apache#26174)

Signed-off-by: Alexandre Burgoni <alexandre.burgoni@clever.cloud>
(cherry picked from commit 0679b0d)
(cherry picked from commit 8cbbb6e)
sandeep-ctds pushed a commit to datastax/pulsar that referenced this pull request Jul 31, 2026
…subscriptions (apache#26174)

Signed-off-by: Alexandre Burgoni <alexandre.burgoni@clever.cloud>
(cherry picked from commit 0679b0d)
(cherry picked from commit 8cbbb6e)
nodece pushed a commit to ascentstream/pulsar that referenced this pull request Aug 28, 2026
…subscriptions (apache#26174)

Signed-off-by: Alexandre Burgoni <alexandre.burgoni@clever.cloud>
(cherry picked from commit 0679b0d)
(cherry picked from commit 8cbbb6e)
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][broker] Failover subscription stops dispatching forever after active-consumer handover

3 participants