Skip to content

[fix][broker] Cancel queued transaction snapshot recovery on topic close - #26335

Merged
Demogorgon314 merged 13 commits into
apache:masterfrom
Demogorgon314:Demogorgon314/Cancel-queued-transaction-snapshot-recovery-on-topic-close
Sep 9, 2026
Merged

Demogorgon314 merged 13 commits into
apache:masterfrom
Demogorgon314:Demogorgon314/Cancel-queued-transaction-snapshot-recovery-on-topic-close

Conversation

@Demogorgon314

Copy link
Copy Markdown
Member

Motivation

Transaction snapshot recovery tasks can remain queued after a topic is closed.
The queued task retains the snapshot processor, topic, and managed ledger,
causing closed topics to accumulate when topic loading is repeatedly retried.

Modifications

  • Cancel queued snapshot recovery tasks when the transaction buffer closes.
  • Wait for running recovery to stop before releasing snapshot resources.
  • Prevent recovery from succeeding after close has started.
  • Stop segmented recovery between I/O operations after close.
  • Handle recovery task submission rejection asynchronously.
  • Add tests for queued, running, completed, retried, and rejected recovery.

Topic close may now wait for an in-flight recovery I/O operation to finish before releasing snapshot resources.

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

@Demogorgon314 Demogorgon314 self-assigned this Aug 15, 2026

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

Thanks for taking this on — the underlying problem is real, and the shape of the fix (a per-attempt state machine plus a recoveryStoppedFuture barrier) is the right one. The three points from the earlier review look genuinely addressed: the stopped-future barrier does cover synchronous continuations, recoveryIndexUpdateFuture is safely published (written on the recovery thread before stoppedFuture.complete, read only in a dependent of it), and the isClosed() check is now at readSegmentEntries entry ahead of openReadOnlyManagedLedger. The new tests are deterministic — latch-driven, no sleeps, no reflection — and SnapshotSegmentAbortedTxnProcessorCloseTest genuinely pins the fix rather than passing vacuously.

My concern is not the barrier itself but what ended up behind it. Because recoverFromSnapshot() hands back future.copy() and TopicTransactionBufferRecover attaches a non-async thenAccept, the continuation that runs inside future.complete(...) is not a small callback — it is the entire transaction-buffer replay: open a non-durable cursor and read every entry through to LAC. Waiting for that before releasing resources is correct for safety, but it puts the whole replay inside topic close, and the two most common close types gate the managed-ledger close on it. Two consequences follow that I think need addressing before merge (the first two inline comments); a third is a residual gap in the barrier itself.

One unrelated, pre-existing bug I noticed while reading PersistentWorker, mentioned only so it does not stay buried — it is not something I think you should fix here. In the Clear branch, taskQueue.forEach(pair -> pair.getRight().getRight().get().completeExceptionally(...)): pair.getRight().getRight() is the Supplier, so .get() starts every queued WriteSegment/DeleteSegment instead of cancelling it, and then completes the newly started task's future rather than the stored taskExecutedResult (which stays incomplete forever). Compare executeTask(), which correctly uses pair.getRight().getKey() for the result future. The effect is that topic deletion launches the writes it means to cancel and then deletes the segments concurrently. Worth a separate issue/PR.

Cancel queued snapshot recovery tasks when the transaction buffer closes and wait for running recovery before releasing snapshot resources.

This prevents the recovery executor queue from retaining closed topics and managed ledgers.
… after failed deletion

Run blocking replay on the snapshot recovery executor and cancel recovery only when the buffer is closed. Gate cursor setup before starting recovery and cover transaction executor responsiveness and failed topic deletion.
Remove same-instance recovery retry coverage and consolidate post-recovery close assertions. Test late recovery completion with a standard future instead of intercepting continuation registration.

Assisted-by: Codex
@Demogorgon314
Demogorgon314 force-pushed the Demogorgon314/Cancel-queued-transaction-snapshot-recovery-on-topic-close branch from 7259653 to 29d2141 Compare September 7, 2026 03:09

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

All six points from my last review are addressed, and the way they are addressed is better than what I asked for. Rather than widening the close barrier you narrowed it to processor-owned work and made the replay itself a separately cancellable task that observes close.

One new problem, and it is the only thing I would hold on. Taking the transaction buffer's monitor in closeAsync() introduces a lock-ordering inversion against a monitor the buffer already acquires in the other direction, so a replay racing a fenced publish can deadlock two shared broker threads. Detail in the inline comment.

Two smaller notes, neither blocking:

  • PersistentWorker is unchanged, so the residual from my earlier comment is still open: putAbortedTxnAndPosition from the replay can enqueue a WriteSegment right up to the moment close takes the buffer monitor, and closeResources() then releases the writers without draining the queue. Pre-existing, still not something I am asking you to widen the scope for — just recording that the new design does not close it either.
  • recoverFromSnapshot() and closeAsync() are the state machine, but neither is final (AbstractSnapshotAbortedTxnProcessor.java:57 and :143). Marking them final would stop a subclass from half-overriding the protocol; isClosed() already is.

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

The lock-order fix landed in the shape I suggested and it holds up on inspection: removeTxnAndUpdateMaxReadPosition no longer touches the topic, all three of its callers invalidate lastDispatchablePosition after leaving the buffer monitor, and testTxnCompletionUpdatesTopicOutsideBufferLock pins that with a Thread.holdsLock assertion across all four replay/commit/abort combinations. The final nit is done as well, and moving the abort path's complete(null) out of the monitor is a bonus improvement.

One edge of the same cycle is still open, and it is one I missed last round rather than anything new here: recoverComplete() and noNeedToRecover() complete transactionBufferFuture while holding the buffer monitor, checkIfTBRecoverCompletely() hands that same future to handleGetLastMessageId, and that continuation can take the topic monitor inline. Detail and a suggested fix are on TopicTransactionBuffer.java:159 — it is the discipline closeAsync() already applies to itself.

Nothing else in the delta needs attention; it is exactly the requested edits plus the test.

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

Both open points are closed, and I checked the code rather than taking the reply for it.

recoverComplete() now decides the state transition under the monitor and completes transactionBufferFuture after leaving it, on the failure branch as well as the success one:

if (recoveryFailure != null) {
getTransactionBufferFuture().completeExceptionally(recoveryFailure);
} else {
getTransactionBufferFuture().complete(null);
}

noNeedToRecover() does the same at TopicTransactionBuffer.java:188-199. The other two completions were already outside a monitor — recoverExceptionally at TopicTransactionBuffer.java:241-244 and closeAsync() at TopicTransactionBuffer.java:756 — and nothing outside this class completes that future, so no path is left that completes it while the buffer monitor is held. Putting the rule on the field itself at TopicTransactionBuffer.java:97-98 is the right place for it.

testRecoveryNotifiesLastPositionQueryOutsideBufferLock pins it rather than leaving it to argument, and both data-provider values earn their place: null drives noNeedToRecover(), a position drives recoverComplete().

One non-blocking flag on the second commit, detail inline: moving handleLowWaterMark past completableFuture.complete(null) is not only a lock move — it changes how far a single low-water-mark trigger drains the backlog. It looks right and it does not reopen the cycle, but nothing pins it and the commit path was left as it was.

The PersistentWorker residual on the segment processor is unchanged and still deliberately out of scope, exactly as I said in my last review.

Otherwise this looks good to me.

@Demogorgon314
Demogorgon314 merged commit d2bc207 into apache:master Sep 9, 2026
43 checks passed
lhotari pushed a commit that referenced this pull request Sep 9, 2026
lhotari pushed a commit that referenced this pull request Sep 9, 2026
lhotari added a commit that referenced this pull request Sep 9, 2026
lhotari pushed a commit that referenced this pull request Sep 9, 2026
lhotari added a commit that referenced this pull request Sep 12, 2026
… release assertion

Fix the replay fixture introduced by
fb093d2
([fix][broker] Cancel queued transaction snapshot recovery on topic
close, #26335), cherry-picked from
d2bc207.

Upstream and branch-4.2 mock Entry.getMessageMetadata(), so the fixture
needs neither serialized metadata nor a real reference-counted entry.
Branch-4.0 lacks that API and parses metadata from the entry buffer.
The backport therefore replaced the mock with a serialized EntryImpl,
but omitted required producerName, sequenceId and publishTime fields
and retained Mockito.verify(marker).release() on the real entry.
It also retained the buffer's original reference after EntryImpl.create.

Build valid commit/abort markers with the production marker helpers,
release the original buffer reference, and assert that recovery releases
the remaining reference instead of verifying a non-mock object.

Validation: all 8 TopicTransactionBufferCloseTest cases pass;
Checkstyle passes.
@lhotari lhotari added this to the 5.0.0-M2 milestone Sep 12, 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