Repository navigation
[fix][ml] Fail queued adds when a managed ledger is terminated during a ledger rollover - #26680
Merged
lhotari merged 2 commits intoSep 25, 2026
Merged
Conversation
… a ledger rollover asyncTerminate() sets the Terminated state and closes the current ledger regardless of a rollover in progress, and the rollover callbacks did not check for it: createComplete() only special-cased the Closed state, then either reset the state on a failed creation or went on to updateLedgersIdsComplete(), which set LedgerOpened and re-submitted the queued adds to the new ledger. A managed ledger terminated while a new ledger was being created became writable again, and the queued adds were persisted after the terminated position stored in the metadata, where no reader ever goes. Check for the Terminated state in createComplete() and around the ledgers-list update that follows the creation (before each attempt, and in its success and failure callbacks): keep the state, fail the queued adds with ManagedLedgerTerminatedException and discard the new ledger.
lhotari
reviewed
Sep 22, 2026
lhotari
left a comment
Member
There was a problem hiding this comment.
Thanks for working on this rollover termination race. The queued-add cleanup paths look good, but the in-flight metadata update still has a version-conflict ordering that can fence the ledger or leave the termination unstored.
…st updates The terminate stored its position without taking the metadata mutex, so it could race the ledgers list update of a rollover it had overtaken, both with the same expected version. Whichever write lost failed with BadVersionException and fenced the managed ledger: either the terminate failed without storing its position, or the in-memory state flipped from Terminated to Fenced after the terminate had already succeeded. Store the terminated position under the metadata mutex, deferring the update while another one is in flight, as the other updates of the ledgers list do.
lhotari
approved these changes
Sep 25, 2026
lhotari
left a comment
Member
There was a problem hiding this comment.
LGTM. The termination write now goes through the same metadataMutex as the ledger rollover, so the version race raised in the previous round is closed on both orderings.
1 of 11 tasks
This was referenced Sep 26, 2026
ascentstream-bot
pushed a commit
to ascentstream/pulsar
that referenced
this pull request
Oct 1, 2026
… a ledger rollover (apache#26680) (cherry picked from commit 5eca6a9)
ascentstream-bot
pushed a commit
to ascentstream/pulsar
that referenced
this pull request
Oct 2, 2026
… a ledger rollover (apache#26680) (cherry picked from commit 5eca6a9)
Radiancebobo
pushed a commit
to Radiancebobo/pulsar
that referenced
this pull request
Oct 8, 2026
… a ledger rollover (apache#26680)
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.
Motivation
ManagedLedgerImpl.asyncTerminate()sets theTerminatedstate and closes the current ledger regardless of a ledger rollover that may be in progress (CreatingLedger, or the update of the ledgers list that follows the creation). The rollover callbacks did not check for the terminated state either:createComplete()only special-cased theClosedstate, and then either reset the state toClosedLedger/WriteFailedwhen the creation failed, or went on toupdateLedgersIdsComplete(), which unconditionally setLedgerOpenedand re-submitted every queuedOpAddEntryto the new ledger.So a managed ledger terminated while a new ledger was being created became writable again, and the adds queued during the rollover were persisted after the terminated position stored in the metadata. The reproducer sees entries written at
4:0after the ledger was terminated at3:0, and the ledger keeps rolling over from there. Readers stop at the terminated position, so those entries are acknowledged to the producer but never delivered.Scalable topics seal a segment with
PersistentTopic.terminate()under producer load, where a rollover in progress is a common state to hit. This is the companion of #26678, which covers the adds that were already in flight on the ledger closed by the terminate (ledgerClosed()): this PR covers the adds that were queued waiting for the next ledger. The two changes touch different code paths.Modifications
ManagedLedgerImpl.createComplete(): if the managed ledger was terminated while the ledger was being created (whether the creation succeeded, failed or timed out), abort the rollover instead of reopening the managed ledger for writes.metadataMutexis held by another operation, which leaves room for the terminate to complete in between) and in both its success and failure callbacks. Without the check before the attempt, the new ledger would be persisted next to the terminated position, and a terminated managed ledger could later be recovered with a last ledger that no longer exists.abortRolloverAfterTerminate()) keeps theTerminatedstate, fails the queued adds withManagedLedgerTerminatedException(asinternalAsyncAddEntry()does for adds arriving after the terminate), and closes and deletes the ledger that was just created so that it is not leaked.asyncTerminate()stored the terminated position without takingmetadataMutex, so it could race the ledgers-list update of the rollover it had overtaken, both with the same expected version: whichever write lost failed withBadVersionExceptionand fenced the managed ledger, either failing the terminate without storing its position, or flipping the in-memory state fromTerminatedtoFencedafter the terminate had already succeeded. The terminated position is now stored undermetadataMutex, deferring the write while another update of the ledgers list is in flight, as all the other updates of the ledgers list do.Verifying this change
This change added tests and can be verified as follows.
Eight new tests in
ManagedLedgerTerminationTestterminate the managed ledger at each point of a rollover in progress, driven by gates on the mock BookKeeper client and on the metadata store rather than by timing:metadataMutexis held;Each of them asserts that the queued adds fail with
ManagedLedgerTerminatedException, that the state staysTerminated, that nothing was written past the terminated position (in memory, in the storedManagedLedgerInfo, and in BookKeeper, where the created ledger is deleted), that a later add is still rejected, and, for the first one, that the terminated state is what gets recovered on reopen. All of them fail without the fix, and each of the added checks is covered by a test that fails when only that check is removed.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes