@@ -386,11 +386,19 @@ public void testReadFailureWaitsForTimerButNewReadDemandCanResumeImmediately() t
386386 }
387387 }
388388
389- @ Test
390- public void testReadFailureTerminatesWhenRetrySchedulingIsRejectedWhileExecutorIsRunning () throws Exception {
389+ @ DataProvider
390+ public Object [][] retrySchedulingFailures () {
391+ return new Object [][] {
392+ {new RejectedExecutionException ("retry rejected" )},
393+ {new IllegalStateException ("retry scheduling failed" )}
394+ };
395+ }
396+
397+ @ Test (dataProvider = "retrySchedulingFailures" )
398+ public void testReadFailureTerminatesWhenRetrySchedulingFails (RuntimeException failure ) throws Exception {
391399 TestReplicatorFixture fixture = newTestReplicatorFixture ();
392400 TestPersistentReplicator replicator = fixture .replicator ;
393- fixture .rejectScheduledWork . set ( true );
401+ doThrow ( failure ). when ( fixture .executor ). schedule ( any ( Runnable . class ), anyLong (), any ( TimeUnit . class ) );
394402 when (fixture .executor .isShuttingDown ()).thenReturn (false );
395403 List <ReadRequest > requests = new ArrayList <>();
396404 doAnswer (invocation -> {
@@ -417,6 +425,52 @@ public void testReadFailureTerminatesWhenRetrySchedulingIsRejectedWhileExecutorI
417425 verify (fixture .cursor ).setInactive ();
418426 }
419427
428+ @ Test
429+ public void testProducerRestartFailureKeepsScheduledRetry () throws Exception {
430+ TestReplicatorFixture fixture = newTestReplicatorFixture ();
431+ fixture .replicator .markDisconnected ();
432+ fixture .replicator .producerStartFailure = new IllegalStateException ("producer startup failed" );
433+
434+ fixture .replicator .readMoreEntries ();
435+ assertThat (fixture .replicator .getState ()).isEqualTo (State .Disconnected );
436+ fixture .takeScheduledWork ().command .run ();
437+ assertThat (fixture .replicator .getState ()).isEqualTo (State .Disconnected );
438+ assertThat (fixture .scheduledWork ).hasSize (1 );
439+ verify (fixture .cursor , never ()).asyncReadEntriesOrWait (anyInt (), anyLong (), any (), any (), any ());
440+ }
441+
442+ @ Test
443+ public void testRejectedReadCompletionRewindsBeforeRetriedRead () throws Exception {
444+ TestReplicatorFixture fixture = newTestReplicatorFixture ();
445+ List <ReadRequest > requests = new ArrayList <>();
446+ doAnswer (invocation -> {
447+ requests .add (new ReadRequest (invocation .getArgument (2 ), invocation .getArgument (3 )));
448+ return null ;
449+ }).when (fixture .cursor ).asyncReadEntriesOrWait (anyInt (), anyLong (), any (), any (), any ());
450+ Position start = PositionFactory .create (1 , 1 );
451+ Position advanced = PositionFactory .create (1 , 3 );
452+ when (fixture .cursor .getReadPosition ()).thenReturn (start );
453+ fixture .replicator .readMoreEntries ();
454+ when (fixture .cursor .getReadPosition ()).thenReturn (advanced );
455+ doAnswer (invocation -> {
456+ when (fixture .cursor .getReadPosition ()).thenReturn (start );
457+ return null ;
458+ }).when (fixture .cursor ).rewind ();
459+
460+ ReadRequest read = requests .get (0 );
461+ read .callback .readEntriesFailed (new ManagedLedgerException (new RejectedExecutionException ("handoff rejected" )),
462+ read .context );
463+ assertThat (requests ).hasSize (1 );
464+ fixture .takeScheduledWork ().command .run ();
465+
466+ assertThat (requests ).hasSize (2 );
467+ assertThat (((InFlightTask ) requests .get (1 ).context ).getReadPos ()).isEqualTo (start );
468+ InOrder cursorCalls = inOrder (fixture .cursor );
469+ cursorCalls .verify (fixture .cursor ).asyncReadEntriesOrWait (anyInt (), anyLong (), any (), any (), any ());
470+ cursorCalls .verify (fixture .cursor ).rewind ();
471+ cursorCalls .verify (fixture .cursor ).asyncReadEntriesOrWait (anyInt (), anyLong (), any (), any (), any ());
472+ }
473+
420474 @ Test
421475 public void testProducerAckDuringReservedReadRetriesFailedReadBeforeTimerWithoutOverlappingReads ()
422476 throws Exception {
@@ -725,15 +779,11 @@ private static TestReplicatorFixture newTestReplicatorFixture(ServiceConfigurati
725779 EventLoopGroup executor = mock (EventLoopGroup .class );
726780 Queue <Runnable > queuedWork = new ConcurrentLinkedQueue <>();
727781 Queue <ScheduledWork > scheduledWork = new ConcurrentLinkedQueue <>();
728- AtomicBoolean rejectScheduledWork = new AtomicBoolean ();
729782 doAnswer (invocation -> {
730783 queuedWork .add (invocation .getArgument (0 ));
731784 return null ;
732785 }).when (executor ).execute (any (Runnable .class ));
733786 doAnswer (invocation -> {
734- if (rejectScheduledWork .get ()) {
735- throw new RejectedExecutionException ("test retry scheduler rejection" );
736- }
737787 scheduledWork .add (new ScheduledWork (invocation .getArgument (0 ), invocation .getArgument (1 ),
738788 invocation .getArgument (2 )));
739789 return null ;
@@ -765,7 +815,7 @@ private static TestReplicatorFixture newTestReplicatorFixture(ServiceConfigurati
765815
766816 TestPersistentReplicator replicator = new TestPersistentReplicator (topic , cursor , brokerService ,
767817 replicationClient , mock (PulsarAdmin .class ));
768- return new TestReplicatorFixture (replicator , cursor , executor , queuedWork , scheduledWork , rejectScheduledWork );
818+ return new TestReplicatorFixture (replicator , cursor , executor , queuedWork , scheduledWork );
769819 }
770820
771821 private record ReadRequest (ReadEntriesCallback callback , Object context ) {
@@ -780,18 +830,15 @@ private static final class TestReplicatorFixture {
780830 private final EventLoopGroup executor ;
781831 private final Queue <Runnable > queuedWork ;
782832 private final Queue <ScheduledWork > scheduledWork ;
783- private final AtomicBoolean rejectScheduledWork ;
784833
785834 private TestReplicatorFixture (TestPersistentReplicator replicator , ManagedCursor cursor ,
786835 EventLoopGroup executor ,
787- Queue <Runnable > queuedWork , Queue <ScheduledWork > scheduledWork ,
788- AtomicBoolean rejectScheduledWork ) {
836+ Queue <Runnable > queuedWork , Queue <ScheduledWork > scheduledWork ) {
789837 this .replicator = replicator ;
790838 this .cursor = cursor ;
791839 this .executor = executor ;
792840 this .queuedWork = queuedWork ;
793841 this .scheduledWork = scheduledWork ;
794- this .rejectScheduledWork = rejectScheduledWork ;
795842 }
796843
797844 private void runQueuedWork () {
@@ -821,6 +868,7 @@ private static final class TestPersistentReplicator extends PersistentReplicator
821868 private final AtomicBoolean rewindBeforeCursorInvocation = new AtomicBoolean ();
822869 private final AtomicBoolean cancelBeforeCursorInvocation = new AtomicBoolean ();
823870 private EntryObserver entryObserver ;
871+ private RuntimeException producerStartFailure ;
824872
825873 private TestPersistentReplicator (PersistentTopic topic , ManagedCursor cursor , BrokerService brokerService ,
826874 PulsarClientImpl replicationClient , PulsarAdmin replicationAdmin )
@@ -833,6 +881,9 @@ private TestPersistentReplicator(PersistentTopic topic, ManagedCursor cursor, Br
833881 @ Override
834882 protected void startProducer () {
835883 // The test drives read scheduling directly.
884+ if (producerStartFailure != null ) {
885+ throw producerStartFailure ;
886+ }
836887 }
837888
838889 @ Override
@@ -861,6 +912,10 @@ private void markTerminated() {
861912 state = State .Terminated ;
862913 }
863914
915+ private void markDisconnected () {
916+ state = State .Disconnected ;
917+ }
918+
864919 private void setProducerForTest (ProducerImpl <?> producer ) {
865920 this .producer = producer ;
866921 }
0 commit comments