Skip to content

[fix][managed-ledger] Honor batch-index cap in isCursorDataFullyPersistable - #26553

Open
arimu1 wants to merge 4 commits into
apache:masterfrom
arimu1:fix/26499-isCursorDataFullyPersistable-batch-index
Open

arimu1 wants to merge 4 commits into
apache:masterfrom
arimu1:fix/26499-isCursorDataFullyPersistable-batch-index

Conversation

@arimu1

@arimu1 arimu1 commented Sep 12, 2026 •

Copy link
Copy Markdown

Fixes #26499

Motivation

ManagedCursorImpl.isCursorDataFullyPersistable() only compared the whole-entry
deletion range count to maxUnackedRangesToPersist. When partial batch ACK state
exceeded maxBatchDeletedIndexToPersist, the method still returned true even
though buildBatchEntryDeletionIndexInfoList() truncates on persist. That
mismatch can leave dispatchers running while some batch ACK records would not
survive a cursor-state write.

Modifications

  • Extend isCursorDataFullyPersistable() to return false when
    batchDeletedIndexes exceeds maxBatchDeletedIndexToPersist.
  • Add a managed-ledger test reproducing the issue scenario (three partial batch
    ACK records with limit 2), including boundary checks at the cap.
  • Document and test PIP-299 interaction: when
    dispatcherPauseOnAckStatePersistentEnabled is on, multi-consumer dispatchers
    also treat managedLedgerMaxBatchDeletedIndexToPersist as a pause threshold
    via isCursorDataFullyPersistable() (in addition to
    managedLedgerMaxUnackedRangesToPersist).

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • ./gradlew :managed-ledger:checkstyleTest :managed-ledger:test --tests org.apache.bookkeeper.mledger.impl.ManagedCursorBatchAckRecoveryTest.testIsCursorDataFullyPersistableReflectsBatchDeletedIndexLimit
  • ./gradlew :pulsar-broker:test --tests org.apache.pulsar.client.api.SubscriptionPauseOnAckStatPersistTest.testPauseOnBatchDeletedIndexLimitExceeded

Does this pull request potentially affect one of the following parts:

  • 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

@arimu1
arimu1 force-pushed the fix/26499-isCursorDataFullyPersistable-batch-index branch from 8d00f5b to 6da992c Compare September 12, 2026 13:24
@void-ptr974

Copy link
Copy Markdown
Contributor

One coordination note for future contributions: I opened #26499 and indicated that I was willing to submit a PR, although I had not opened one yet. In this situation, it would be helpful to leave a short comment before starting work, so contributors can avoid duplicating effort. No problem this time—just please coordinate on the issue first in similar cases.

One change is required: the new org.awaitility.Awaitility import is out of order. I verified that ./gradlew :managed-ledger:checkstyleTest fails with an ImportOrder violation. Please move it after the org.apache... imports.

I did not find another blocker in the current diff.

Comment on lines +417 to +418
return batchDeletedIndexes == null
|| batchDeletedIndexes.size() <= getConfig().getMaxBatchDeletedIndexToPersist();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we clarify the intended policy here? The method isCursorDataFullyPersistable() is used directly by multiple-consumer dispatchers when dispatcherPauseOnAckStatePersistentEnabled is enabled, so this change effectively makes managedLedgerMaxBatchDeletedIndexToPersist a dispatcher-pause threshold.

PIP-299 explicitly described pausing when managedLedgerMaxUnackedRangesToPersist is reached. While pausing on batch-index overflow also seems reasonable, this represents a real behavior change—not merely an improvement to the helper's accuracy.

If this is intended, could you update the PR description accordingly and add a dispatcher-level test that covers pause/resume behavior when the batch-index limit is exceeded?


@Test
public void testIsCursorDataFullyPersistableReflectsBatchDeletedIndexLimit() throws Exception {
ManagedLedgerConfig config = new ManagedLedgerConfig();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

After rebasing onto master, should we use rawEntryConfig() here?

This test writes raw byte arrays using ledger.addEntry(new byte[] {...}), and the existing tests in ManagedCursorBatchAckRecoveryTest have already migrated to rawEntryConfig() to prevent the cache/read path from parsing them as Pulsar message entries.

Keeping new ManagedLedgerConfig() here would make this new test inconsistent with the current setup and may cause issues after the rebase.

Comment on lines +71 to +76
long[][] ackSets = {{2L}, {Long.MIN_VALUE, 1L}, {5L, 0L, 3L}};
for (int i = 1; i < positions.size(); i++) {
Position position = positions.get(i);
cursor.delete(AckSetStateUtil.createPositionWithAckSet(
position.getLedgerId(), position.getEntryId(), ackSets[i - 1]));
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

buildBatchEntryDeletionIndexInfoList() can persist at most maxBatchDeletedIndexToPersist records. It should therefore be verified that, with a limit of two, two batch-deleted-index records are reported as fully persistable, and that adding a third switches the result to false.

This also protects the check against a future accidental regression of >= to > (or vice versa).

…stable

Also compare batchDeletedIndexes size to maxBatchDeletedIndexToPersist so
dispatchers see non-persistable ACK state when batch records exceed the limit.

Assisted-by: Cursor (Composer)
Align managed-ledger test with defaultConfig(), fix import order, and assert
persistability at the batch-index cap boundary. Add a broker test that pauses
and resumes dispatch when batch deleted index state exceeds the limit under
PIP-299 dispatcher pause policy.
@arimu1

arimu1 commented Sep 21, 2026

Copy link
Copy Markdown
Author

Thanks for the review — addressed in tip c7457d0 (rebased onto current master).

@void-ptr974 — Moved the Awaitility import after the org.apache.* imports; :managed-ledger:checkstyleTest passes.

@Denovo1998

  1. Dispatcher pause / PR description — Updated the PR body: with dispatcherPauseOnAckStatePersistentEnabled, multi-consumer dispatchers now also pause when batch deleted index state exceeds managedLedgerMaxBatchDeletedIndexToPersist because they call isCursorDataFullyPersistable(). Added SubscriptionPauseOnAckStatPersistTest#testPauseOnBatchDeletedIndexLimitExceeded for pause + resume on that path.
  2. rawEntryConfig() — After rebase, upstream renamed this helper to ManagedLedgerTestUtil.defaultConfig() ([improve][ml] Allow inline read completion and preserve replication order #26620); the managed-ledger test now uses defaultConfig() for raw-byte entries.
  3. Boundary (>= vs >) — testIsCursorDataFullyPersistableReflectsBatchDeletedIndexLimit now asserts true after the 1st and 2nd batch-index records (limit 2) and false after the 3rd.

Tests run locally (Temurin 21):

  • ./gradlew :managed-ledger:checkstyleTest :managed-ledger:test --tests …ManagedCursorBatchAckRecoveryTest.testIsCursorDataFullyPersistableReflectsBatchDeletedIndexLimit
  • ./gradlew :pulsar-broker:test --tests …SubscriptionPauseOnAckStatPersistTest.testPauseOnBatchDeletedIndexLimitExceeded

@arimu1
arimu1 force-pushed the fix/26499-isCursorDataFullyPersistable-batch-index branch from 6da992c to c7457d0 Compare September 21, 2026 11:04

Consumer<String> consumer = pulsarClient.newConsumer(Schema.STRING).topic(tpName)
.subscriptionName(subscription).subscriptionType(SubscriptionType.Shared)
.receiverQueueSize(1).enableBatchIndexAcknowledgment(true).isAckReceiptEnabled(true)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think there is still a test race condition around the consumer's initial Flow.

ConsumerImpl completes subscribeFuture before calling increaseAvailablePermits(cnx, getCurrentReceiverQueueSize()). After subscribe() returns, either the test's afterAckMessages() call or the connection thread sending the initial permit can win.

If the initial Flow wins, blockedDispatcherOnCursorDataCanNotFullyPersist remains false, and the dispatcher can already read one of the existing backlog messages. Since cancelPendingRead() cannot retract a message already delivered to the consumer queue, pausedReceive could become non-null even though the pause behavior is correct.

Could we make this deterministic by avoiding the initial permit? For example, using a zero-queue consumer with a pending receiveAsync() would allow us to set the blocked state first, then explicitly issue a single permit through the receive. The same future could be checked for timeout while paused and completed after the dispatcher resumes.

…ataFullyPersistable-batch-index

Signed-off-by: arimu1 <19286898+arimu1@users.noreply.github.com>
Use receiverQueueSize(0) and a pending receiveAsync so afterAckMessages
can set the blocked state before any permit is issued, addressing the
subscribe/Flow race noted in review.

Signed-off-by: arimu1 <19286898+arimu1@users.noreply.github.com>

@Denovo1998 Denovo1998 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update. I checked the zero-queue approach against the consumer Flow path, and this addresses the race I was concerned about.

With receiverQueueSize(0), no initial Flow is sent during subscribe, and the pending receiveAsync() explicitly issues the permit only after the dispatcher has entered the paused state. Reusing the same future after the batch-index state becomes persistable also covers the resume path cleanly.

The previous review comments look addressed now. LGTM.

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] isCursorDataFullyPersistable ignores the batch-index persistence limit

3 participants