|
27 | 27 | import static org.mockito.Mockito.doThrow; |
28 | 28 | import static org.mockito.Mockito.inOrder; |
29 | 29 | import static org.mockito.Mockito.mock; |
| 30 | +import static org.mockito.Mockito.never; |
30 | 31 | import static org.mockito.Mockito.times; |
31 | 32 | import static org.mockito.Mockito.verify; |
32 | 33 | import static org.mockito.Mockito.when; |
|
35 | 36 | import java.util.Collections; |
36 | 37 | import java.util.List; |
37 | 38 | import java.util.Queue; |
| 39 | +import java.util.concurrent.CompletableFuture; |
38 | 40 | import java.util.concurrent.ConcurrentLinkedQueue; |
39 | 41 | import java.util.concurrent.CountDownLatch; |
40 | 42 | import java.util.concurrent.RejectedExecutionException; |
|
55 | 57 | import org.apache.pulsar.broker.service.persistent.PersistentReplicator.InFlightTask; |
56 | 58 | import org.apache.pulsar.client.admin.PulsarAdmin; |
57 | 59 | import org.apache.pulsar.client.api.ProducerBuilder; |
| 60 | +import org.apache.pulsar.client.api.PulsarClientException; |
58 | 61 | import org.apache.pulsar.client.api.Schema; |
59 | 62 | import org.apache.pulsar.client.impl.ProducerImpl; |
60 | 63 | import org.apache.pulsar.client.impl.PulsarClientImpl; |
@@ -463,17 +466,119 @@ public void testProducerAckDuringReservedReadRetriesFailedReadBeforeTimerWithout |
463 | 466 |
|
464 | 467 | reservedRead.callback.readEntriesFailed(new ManagedLedgerException.TooManyRequestsException("read failed"), |
465 | 468 | reservedRead.context); |
| 469 | + fixture.runQueuedWork(); |
466 | 470 | ScheduledWork fallbackRetry = fixture.takeScheduledWork(); |
467 | 471 | synchronized (requests) { |
468 | 472 | // The ACK demand is consumed after the failed reservation settles, before its fallback timer runs. |
469 | 473 | assertThat(requests).hasSize(3); |
470 | 474 | } |
471 | 475 |
|
| 476 | + ReadRequest immediateRetry = requests.get(2); |
| 477 | + immediateRetry.callback.readEntriesFailed(new ManagedLedgerException.TooManyRequestsException("retry failed"), |
| 478 | + immediateRetry.context); |
| 479 | + fixture.runQueuedWork(); |
| 480 | + assertThat(requests).hasSize(3); |
| 481 | + assertThat(fixture.scheduledWork).isEmpty(); |
| 482 | + |
472 | 483 | fallbackRetry.command.run(); |
473 | 484 | synchronized (requests) { |
474 | | - // The fallback request cannot overlap the retry it found pending. |
475 | | - assertThat(requests).hasSize(3); |
| 485 | + // One ACK permits only one immediate retry. A second failure must wait for the existing timer. |
| 486 | + assertThat(requests).hasSize(4); |
| 487 | + } |
| 488 | + } |
| 489 | + |
| 490 | + @DataProvider |
| 491 | + public Object[][] sendCompletionOutcomes() { |
| 492 | + return new Object[][] {{false}, {true}}; |
| 493 | + } |
| 494 | + |
| 495 | + @Test(dataProvider = "sendCompletionOutcomes") |
| 496 | + public void testSendCompletionQueuesReadProcessingOutsideProducerMonitor(boolean failedSend) throws Exception { |
| 497 | + ServiceConfiguration configuration = new ServiceConfiguration(); |
| 498 | + configuration.setReplicationProducerQueueSize(2); |
| 499 | + TestReplicatorFixture fixture = newTestReplicatorFixture(configuration); |
| 500 | + TestPersistentReplicator replicator = fixture.replicator; |
| 501 | + ProducerImpl<?> producer = mock(ProducerImpl.class); |
| 502 | + when(producer.isWritable()).thenReturn(true); |
| 503 | + replicator.setProducerForTest(producer); |
| 504 | + List<ReadRequest> requests = new ArrayList<>(); |
| 505 | + doAnswer(invocation -> { |
| 506 | + assertThat(Thread.holdsLock(producer)).isFalse(); |
| 507 | + requests.add(new ReadRequest(invocation.getArgument(2), invocation.getArgument(3))); |
| 508 | + return null; |
| 509 | + }).when(fixture.cursor).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any()); |
| 510 | + doAnswer(invocation -> { |
| 511 | + assertThat(Thread.holdsLock(producer)).isFalse(); |
| 512 | + return null; |
| 513 | + }).when(fixture.cursor).rewind(); |
| 514 | + doAnswer(invocation -> { |
| 515 | + assertThat(Thread.holdsLock(producer)).isFalse(); |
| 516 | + return false; |
| 517 | + }).when(fixture.cursor).cancelPendingReadRequest(); |
| 518 | + List<PersistentReplicator.ProducerSendCallback> callbacks = new ArrayList<>(); |
| 519 | + replicator.entryObserver = (entry, task, entries) -> callbacks.add( |
| 520 | + PersistentReplicator.ProducerSendCallback.create(replicator, entry, null, task)); |
| 521 | + |
| 522 | + replicator.readMoreEntries(); |
| 523 | + ReadRequest read = requests.get(0); |
| 524 | + read.callback.readEntriesComplete(List.of(entry(0), entry(1)), read.context); |
| 525 | + assertThat(requests).hasSize(1); |
| 526 | + assertThat(callbacks).hasSize(2); |
| 527 | + synchronized (producer) { |
| 528 | + callbacks.get(0).sendComplete(failedSend ? new PulsarClientException("send failed") : null, null); |
| 529 | + callbacks.get(1).sendComplete(null, null); |
| 530 | + assertThat(requests).hasSize(1); |
| 531 | + verify(fixture.cursor, never()).rewind(); |
| 532 | + verify(fixture.cursor, never()).cancelPendingReadRequest(); |
| 533 | + assertThat(fixture.queuedWork).hasSize(1); |
| 534 | + } |
| 535 | + fixture.runQueuedWork(); |
| 536 | + assertThat(requests).hasSize(2); |
| 537 | + if (failedSend) { |
| 538 | + verify(fixture.cursor).rewind(); |
| 539 | + assertThat(replicator.waitForCursorRewindingRefCnf).isZero(); |
| 540 | + } |
| 541 | + } |
| 542 | + |
| 543 | + @Test |
| 544 | + public void testRejectedAckHandoffTerminatesWithoutDrainingOnAckThread() throws Exception { |
| 545 | + ServiceConfiguration configuration = new ServiceConfiguration(); |
| 546 | + configuration.setReplicationProducerQueueSize(2); |
| 547 | + TestReplicatorFixture fixture = newTestReplicatorFixture(configuration); |
| 548 | + TestPersistentReplicator replicator = fixture.replicator; |
| 549 | + ProducerImpl<?> producer = mock(ProducerImpl.class); |
| 550 | + when(producer.isWritable()).thenReturn(true); |
| 551 | + when(producer.closeAsync()).thenReturn(CompletableFuture.completedFuture(null)); |
| 552 | + replicator.setProducerForTest(producer); |
| 553 | + List<ReadRequest> requests = new ArrayList<>(); |
| 554 | + doAnswer(invocation -> { |
| 555 | + requests.add(new ReadRequest(invocation.getArgument(2), invocation.getArgument(3))); |
| 556 | + return null; |
| 557 | + }).when(fixture.cursor).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any()); |
| 558 | + List<PersistentReplicator.ProducerSendCallback> callbacks = new ArrayList<>(); |
| 559 | + replicator.entryObserver = (entry, task, entries) -> callbacks.add( |
| 560 | + PersistentReplicator.ProducerSendCallback.create(replicator, entry, null, task)); |
| 561 | + replicator.readMoreEntries(); |
| 562 | + Entry first = entry(0); |
| 563 | + Entry second = entry(1); |
| 564 | + ReadRequest read = requests.get(0); |
| 565 | + read.callback.readEntriesComplete(List.of(first, second), read.context); |
| 566 | + doThrow(new RejectedExecutionException("ACK handoff rejected")) |
| 567 | + .when(fixture.executor).execute(any(Runnable.class)); |
| 568 | + |
| 569 | + synchronized (producer) { |
| 570 | + callbacks.get(0).sendComplete(null, null); |
| 571 | + callbacks.get(1).sendComplete(null, null); |
476 | 572 | } |
| 573 | + |
| 574 | + assertThat(replicator.getState()).isEqualTo(State.Terminated); |
| 575 | + assertThat(requests).hasSize(1); |
| 576 | + assertThat(fixture.queuedWork).isEmpty(); |
| 577 | + verify(fixture.cursor, never()).rewind(); |
| 578 | + verify(fixture.cursor, never()).cancelPendingReadRequest(); |
| 579 | + verify(first).release(); |
| 580 | + verify(second).release(); |
| 581 | + verify(producer).closeAsync(); |
477 | 582 | } |
478 | 583 |
|
479 | 584 | @Test |
@@ -607,7 +712,6 @@ private static TestReplicatorFixture newTestReplicatorFixture() throws Exception |
607 | 712 | @SuppressWarnings("unchecked") |
608 | 713 | private static TestReplicatorFixture newTestReplicatorFixture(ServiceConfiguration configuration) throws Exception { |
609 | 714 | configuration.setClusterName("local"); |
610 | | - configuration.setReplicationProducerQueueSize(1000); |
611 | 715 | configuration.setDispatcherMaxReadBatchSize(1000); |
612 | 716 | configuration.setDispatcherMaxReadSizeBytes(1024 * 1024); |
613 | 717 |
|
|
0 commit comments