Skip to content

Commit 20b8a96

Browse files
committed
[fix][broker] Preserve replication recovery across read owner transitions
Retain ACK demand without scheduling a redundant turn while a read is pending. Rewind rejected read completions before another admission even across disconnection, and contain scheduling or startup errors after releasing ownership so they cannot clear a newer owner. Capture failed-send diagnostic state under the task monitor. Replace the live-broker failed-send fixture with controlled executor coverage, retain the throttling fault-injection assertion, and add deterministic recovery and concurrent ownership regressions. Assisted-by: Codex
1 parent 90f1ecf commit 20b8a96

4 files changed

Lines changed: 136 additions & 52 deletions

File tree

‎pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java‎

Lines changed: 25 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -324,14 +324,19 @@ private void requestReadProcessing(boolean requestRead, boolean runAsync) {
324324
if (processingReads) {
325325
return;
326326
}
327+
if (runAsync && !cancelReadRequested && !rewindRequested && hasPendingRead()) {
328+
// The pending read's completion will claim the owner and consume this demand. Do not queue
329+
// a turn per ACK just to discover that the result is not available yet.
330+
return;
331+
}
327332
processingReads = true;
328333
}
329334
if (runAsync) {
330335
try {
331336
// ACKs may hold the geo producer monitor on its IO thread. Retain ownership while queued,
332337
// so concurrent callbacks only publish work rather than starting another drain.
333338
brokerService.executor().execute(this::processReads);
334-
} catch (RuntimeException e) {
339+
} catch (Throwable e) {
335340
log.error().exception(e).log("Failed to schedule replication read processing");
336341
// Retain ownership through termination: its cleanup request must not start an inline drain
337342
// on the ACK thread either. A late read callback will settle any still-pending result.
@@ -400,7 +405,7 @@ private void discardPendingReadResults() {
400405
}
401406
}
402407

403-
private void handleReadRetrySchedulingFailure(Exception exception) {
408+
private void handleReadRetrySchedulingFailure(Throwable exception) {
404409
// Ownership has already been released. Never clear a newer owner's state here.
405410
log.error().exception(exception).log("Failed to schedule replication read retry");
406411
// A failed retry submission has no wakeup left if there are no producer ACKs in flight.
@@ -466,7 +471,7 @@ private boolean processRead() {
466471
if (retryDelayMillis > 0) {
467472
try {
468473
scheduleReadRetry(retryDelayMillis);
469-
} catch (Exception e) {
474+
} catch (Throwable e) {
470475
// Ownership was already released: a new owner might be running now. Do not let
471476
// this failure reach the owner cleanup in processReads and clear its ownership.
472477
handleReadRetrySchedulingFailure(e);
@@ -475,7 +480,7 @@ private boolean processRead() {
475480
if (state == Disconnected) {
476481
try {
477482
startProducer();
478-
} catch (Exception e) {
483+
} catch (Throwable e) {
479484
// The read-retry timer was accepted; this is not a timer scheduling failure.
480485
log.error().exception(e).log("Failed to restart replication producer; retry remains scheduled");
481486
}
@@ -525,7 +530,7 @@ private void scheduleReadRetry(long delayMillis) {
525530
}
526531
readMoreEntries();
527532
}, delayMillis, TimeUnit.MILLISECONDS);
528-
} catch (RuntimeException e) {
533+
} catch (Throwable e) {
529534
synchronized (inFlightTasks) {
530535
readRetryScheduled = false;
531536
}
@@ -623,18 +628,20 @@ protected static final class ProducerSendCallback implements SendCallback {
623628
public void sendComplete(Throwable exception, OpSendMsgStats opSendMsgStats) {
624629
boolean failed = exception != null && !(exception instanceof PulsarClientException.InvalidMessageException);
625630
if (failed) {
626-
replicator.log.error()
627-
.attr("inFlightTasks", replicator.inFlightTasks)
628-
.attr("pendingQueueSize", replicator.producer.getPendingQueueSize())
629-
.exception(exception)
630-
.log("Error producing on remote broker");
631+
int inFlightTaskCount;
631632
synchronized (replicator.inFlightTasks) {
632633
// Unlike asynchronous schema lookup, this recovery has no outstanding stage to wait for.
633634
// Publish cancellation and rewind together, leaving cursor work to the owner.
634635
replicator.inFlightTasks.forEach(task -> task.skipReadResultDueToCursorRewind = true);
635636
replicator.cancelReadRequested = true;
636637
replicator.rewindRequested = true;
638+
inFlightTaskCount = replicator.inFlightTasks.size();
637639
}
640+
replicator.log.error()
641+
.attr("inFlightTaskCount", inFlightTaskCount)
642+
.attr("pendingQueueSize", replicator.producer.getPendingQueueSize())
643+
.exception(exception)
644+
.log("Error producing on remote broker");
638645
// The failed send has completed from the producer queue perspective. The cursor rewind
639646
// makes the entry readable again, so this in-flight task must release its permit.
640647
inFlightTask.incCompletedEntries();
@@ -743,6 +750,13 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) {
743750
}
744751

745752
private void handleReadFailure(ManagedLedgerException exception, InFlightTask task) {
753+
if (exception.getCause() instanceof RejectedExecutionException) {
754+
synchronized (inFlightTasks) {
755+
// Completion may be rejected after advancing the cursor but before transferring entries.
756+
// Only this owner can restore the position before it admits another read.
757+
rewindRequested = true;
758+
}
759+
}
746760
if (state != Started) {
747761
return;
748762
}
@@ -751,13 +765,6 @@ private void handleReadFailure(ManagedLedgerException exception, InFlightTask ta
751765
terminate();
752766
return;
753767
}
754-
if (exception.getCause() instanceof RejectedExecutionException) {
755-
synchronized (inFlightTasks) {
756-
// Completion may be rejected after advancing the cursor but before transferring entries.
757-
// Only this owner can restore the position before it admits another read.
758-
rewindRequested = true;
759-
}
760-
}
761768
readBatchSize = brokerService.pulsar().getConfiguration().getDispatcherMinReadBatchSize();
762769
long waitTimeMillis = delayReadRetry();
763770
if (!(exception instanceof TooManyRequestsException)) {
@@ -777,7 +784,7 @@ protected long delayReadRetry() {
777784
}
778785
try {
779786
scheduleReadRetry(waitTimeMillis);
780-
} catch (Exception e) {
787+
} catch (Throwable e) {
781788
// A failed timer must not interrupt the caller's unsent-entry cleanup or schema rewind.
782789
handleReadRetrySchedulingFailure(e);
783790
}

‎pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorTest.java‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -933,8 +933,10 @@ public void testProbBKErrorWhenReplicating() throws Exception {
933933
// Smoke-test progress while failures keep alternating with successful reads. The deterministic
934934
// owner-loop tests verify that ACK demand resumes reads before the fallback timer.
935935
AtomicInteger readAttempts = new AtomicInteger();
936+
AtomicInteger injectedFailures = new AtomicInteger();
936937
Supplier<ManagedLedgerException> bkErrorOrNot = () -> {
937938
if (readAttempts.incrementAndGet() % 2 == 1) {
939+
injectedFailures.incrementAndGet();
938940
return new ManagedLedgerException.TooManyRequestsException("mocked error");
939941
}
940942
return null;
@@ -965,6 +967,7 @@ public void testProbBKErrorWhenReplicating() throws Exception {
965967
}
966968
assertEquals(received.size(), msgPublished.size());
967969
assertEquals(received, msgPublished);
970+
assertTrue(injectedFailures.get() > 0, "At least one BookKeeper read failure must be injected");
968971

969972
// cleanup.
970973
producer1.close();

‎pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorInflightTaskTest.java‎

Lines changed: 0 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,6 @@
2626
import static org.mockito.ArgumentMatchers.eq;
2727
import static org.mockito.ArgumentMatchers.same;
2828
import static org.mockito.Mockito.doAnswer;
29-
import static org.mockito.Mockito.doNothing;
3029
import static org.mockito.Mockito.mock;
3130
import static org.mockito.Mockito.never;
3231
import static org.mockito.Mockito.spy;
@@ -206,35 +205,6 @@ public void testReadEntriesFailedCompletesInFlightTaskAfterReplicatorTerminated(
206205
}
207206
}
208207

209-
@Test
210-
public void testFailedPublishCompletesInFlightTask() throws Exception {
211-
PersistentReplicator replicator = spy(getReplicator(topicName));
212-
doNothing().when(replicator).beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding.Failed_Publishing);
213-
doNothing().when(replicator).doRewindCursor(false);
214-
doNothing().when(replicator).readMoreEntries();
215-
216-
LinkedList<InFlightTask> inFlightTasks = replicator.inFlightTasks;
217-
List<InFlightTask> originalTasks = new ArrayList<>(inFlightTasks);
218-
inFlightTasks.clear();
219-
220-
try {
221-
InFlightTask task = new InFlightTask(PositionFactory.create(1, 1), 1, replicator.getReplicatorId());
222-
task.setEntries(Collections.singletonList(mock(Entry.class)));
223-
task.setSubmissionComplete(true);
224-
inFlightTasks.add(task);
225-
assertEquals(replicator.getPermitsIfNoPendingRead(), 999);
226-
227-
ProducerSendCallback callback = ProducerSendCallback.create(replicator, mock(Entry.class), null, task);
228-
callback.sendComplete(new PulsarClientException.ProducerBlockedQuotaExceededException("mocked"), null);
229-
230-
assertTrue(task.isDone());
231-
assertEquals(replicator.getPermitsIfNoPendingRead(), 1000);
232-
} finally {
233-
inFlightTasks.clear();
234-
inFlightTasks.addAll(originalTasks);
235-
}
236-
}
237-
238208
/**
239209
* Reproduces a geo-replication stall on the cursor-rewind path.
240210
*

‎pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorReadProcessingTest.java‎

Lines changed: 108 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -390,12 +390,13 @@ public void testReadFailureWaitsForTimerButNewReadDemandCanResumeImmediately() t
390390
public Object[][] retrySchedulingFailures() {
391391
return new Object[][] {
392392
{new RejectedExecutionException("retry rejected")},
393-
{new IllegalStateException("retry scheduling failed")}
393+
{new IllegalStateException("retry scheduling failed")},
394+
{new AssertionError("retry scheduler error")}
394395
};
395396
}
396397

397398
@Test(dataProvider = "retrySchedulingFailures")
398-
public void testReadFailureTerminatesWhenRetrySchedulingFails(RuntimeException failure) throws Exception {
399+
public void testReadFailureTerminatesWhenRetrySchedulingFails(Throwable failure) throws Exception {
399400
TestReplicatorFixture fixture = newTestReplicatorFixture();
400401
TestPersistentReplicator replicator = fixture.replicator;
401402
doThrow(failure).when(fixture.executor).schedule(any(Runnable.class), anyLong(), any(TimeUnit.class));
@@ -440,7 +441,59 @@ public void testProducerRestartFailureKeepsScheduledRetry() throws Exception {
440441
}
441442

442443
@Test
443-
public void testRejectedReadCompletionRewindsBeforeRetriedRead() throws Exception {
444+
public void testProducerRestartErrorCannotReleaseAnotherReadOwner() throws Exception {
445+
TestReplicatorFixture fixture = newTestReplicatorFixture();
446+
TestPersistentReplicator replicator = fixture.replicator;
447+
CountDownLatch resultPublished = new CountDownLatch(1);
448+
CountDownLatch releaseReadOwner = new CountDownLatch(1);
449+
AtomicInteger reads = new AtomicInteger();
450+
Entry entry = entry(0);
451+
replicator.entryObserver = (next, task, entries) -> {
452+
replicator.submittedEntries.add(next.getEntryId());
453+
task.incCompletedEntries();
454+
next.release();
455+
};
456+
doAnswer(invocation -> {
457+
if (reads.incrementAndGet() == 1) {
458+
ReadEntriesCallback callback = invocation.getArgument(2);
459+
callback.readEntriesComplete(List.of(entry), invocation.getArgument(3));
460+
resultPublished.countDown();
461+
await(releaseReadOwner);
462+
}
463+
return null;
464+
}).when(fixture.cursor).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any());
465+
Thread newOwner = new Thread(replicator::readMoreEntries, "replication-read-owner");
466+
replicator.markDisconnected();
467+
replicator.producerStartAction = () -> {
468+
replicator.markStarted();
469+
newOwner.start();
470+
await(resultPublished);
471+
throw new AssertionError("producer restart failed after another owner started");
472+
};
473+
try {
474+
replicator.readMoreEntries();
475+
replicator.readMoreEntries();
476+
assertThat(reads).hasValue(1);
477+
assertThat(replicator.submittedEntries).isEmpty();
478+
verify(entry, never()).release();
479+
verify(fixture.cursor, never()).cancelPendingReadRequest();
480+
verify(fixture.cursor, never()).rewind();
481+
} finally {
482+
releaseReadOwner.countDown();
483+
newOwner.join(TimeUnit.SECONDS.toMillis(10));
484+
assertThat(newOwner.isAlive()).isFalse();
485+
}
486+
assertThat(replicator.submittedEntries).containsExactly(0L);
487+
verify(entry).release();
488+
}
489+
490+
@DataProvider
491+
public Object[][] rejectedCompletionStates() {
492+
return new Object[][] {{State.Started}, {State.Disconnected}};
493+
}
494+
495+
@Test(dataProvider = "rejectedCompletionStates")
496+
public void testRejectedReadCompletionRewindsBeforeRetriedRead(State stateAtCompletion) throws Exception {
444497
TestReplicatorFixture fixture = newTestReplicatorFixture();
445498
List<ReadRequest> requests = new ArrayList<>();
446499
doAnswer(invocation -> {
@@ -457,11 +510,20 @@ public void testRejectedReadCompletionRewindsBeforeRetriedRead() throws Exceptio
457510
return null;
458511
}).when(fixture.cursor).rewind();
459512

513+
if (stateAtCompletion == State.Disconnected) {
514+
fixture.replicator.markDisconnected();
515+
}
460516
ReadRequest read = requests.get(0);
461517
read.callback.readEntriesFailed(new ManagedLedgerException(new RejectedExecutionException("handoff rejected")),
462518
read.context);
463519
assertThat(requests).hasSize(1);
464-
fixture.takeScheduledWork().command.run();
520+
verify(fixture.cursor).rewind();
521+
if (stateAtCompletion == State.Disconnected) {
522+
fixture.replicator.markStarted();
523+
fixture.replicator.readMoreEntries();
524+
} else {
525+
fixture.takeScheduledWork().command.run();
526+
}
465527

466528
assertThat(requests).hasSize(2);
467529
assertThat(((InFlightTask) requests.get(1).context).getReadPos()).isEqualTo(start);
@@ -513,6 +575,8 @@ public void testProducerAckDuringReservedReadRetriesFailedReadBeforeTimerWithout
513575

514576
assertThat(acknowledgement[0]).isNotNull();
515577
acknowledgement[0].sendComplete(null, null);
578+
assertThat(fixture.queuedWork).isEmpty();
579+
verify(fixture.executor, never()).execute(any(Runnable.class));
516580
synchronized (requests) {
517581
// The actual producer callback records demand but cannot issue another read while one is reserved.
518582
assertThat(requests).hasSize(2);
@@ -635,6 +699,38 @@ public void testRejectedAckHandoffTerminatesWithoutDrainingOnAckThread() throws
635699
verify(producer).closeAsync();
636700
}
637701

702+
@Test
703+
public void testFailedPublishCompletesInFlightTaskBeforeQueuedRecovery() throws Exception {
704+
TestReplicatorFixture fixture = newTestReplicatorFixture();
705+
TestPersistentReplicator replicator = fixture.replicator;
706+
ProducerImpl<?> producer = mock(ProducerImpl.class);
707+
when(producer.isWritable()).thenReturn(true);
708+
replicator.setProducerForTest(producer);
709+
Entry entry = entry(0);
710+
InFlightTask task = new InFlightTask(PositionFactory.create(1, 1), 1, replicator.getReplicatorId());
711+
task.setEntries(List.of(entry));
712+
task.setSubmissionComplete(true);
713+
replicator.inFlightTasks.add(task);
714+
assertThat(replicator.getPermitsIfNoPendingRead()).isEqualTo(999);
715+
716+
PersistentReplicator.ProducerSendCallback callback =
717+
PersistentReplicator.ProducerSendCallback.create(replicator, entry, null, task);
718+
callback.sendComplete(new PulsarClientException.ProducerBlockedQuotaExceededException("test failure"), null);
719+
720+
assertThat(task.isDone()).isTrue();
721+
assertThat(task.getCompletedEntries()).isEqualTo(1);
722+
verify(entry).release();
723+
assertThat(replicator.getPermitsIfNoPendingRead()).isEqualTo(1000);
724+
assertThat(fixture.queuedWork).hasSize(1);
725+
verify(fixture.cursor, never()).rewind();
726+
verify(fixture.cursor, never()).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any());
727+
728+
fixture.runQueuedWork();
729+
InOrder cursorCalls = inOrder(fixture.cursor);
730+
cursorCalls.verify(fixture.cursor).rewind();
731+
cursorCalls.verify(fixture.cursor).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any());
732+
}
733+
638734
@Test
639735
public void testRewindFailureRetriesWithoutTerminatingOrRecursing() throws Exception {
640736
TestReplicatorFixture fixture = newTestReplicatorFixture();
@@ -869,6 +965,7 @@ private static final class TestPersistentReplicator extends PersistentReplicator
869965
private final AtomicBoolean cancelBeforeCursorInvocation = new AtomicBoolean();
870966
private EntryObserver entryObserver;
871967
private RuntimeException producerStartFailure;
968+
private Runnable producerStartAction;
872969

873970
private TestPersistentReplicator(PersistentTopic topic, ManagedCursor cursor, BrokerService brokerService,
874971
PulsarClientImpl replicationClient, PulsarAdmin replicationAdmin)
@@ -881,6 +978,9 @@ private TestPersistentReplicator(PersistentTopic topic, ManagedCursor cursor, Br
881978
@Override
882979
protected void startProducer() {
883980
// The test drives read scheduling directly.
981+
if (producerStartAction != null) {
982+
producerStartAction.run();
983+
}
884984
if (producerStartFailure != null) {
885985
throw producerStartFailure;
886986
}
@@ -916,6 +1016,10 @@ private void markDisconnected() {
9161016
state = State.Disconnected;
9171017
}
9181018

1019+
private void markStarted() {
1020+
state = State.Started;
1021+
}
1022+
9191023
private void setProducerForTest(ProducerImpl<?> producer) {
9201024
this.producer = producer;
9211025
}

0 commit comments

Comments
 (0)