+
+ @Test
+ @SuppressWarnings("checkstyle:illegalcatch")
+ public void testApplyStateRace() throws Exception {
+ final String leaderId = factory.generateActorId("leader-");
+ final String followerId = factory.generateActorId("follower-");
+
+ DefaultConfigParamsImpl config = new DefaultConfigParamsImpl();
+ config.setIsolatedLeaderCheckInterval(new FiniteDuration(1, TimeUnit.DAYS));
+ config.setCustomRaftPolicyImplementationClass(DisableElectionsRaftPolicy.class.getName());
+
+ ActorRef mockFollowerActorRef = factory.createActor(MessageCollectorActor.props());
+
+ TestRaftActor.Builder builder = TestRaftActor.newBuilder()
+ .id(leaderId)
+ .peerAddresses(Map.of(followerId, mockFollowerActorRef.path().toString()))
+ .config(config)
+ .collectorActor(factory.createActor(
+ MessageCollectorActor.props(), factory.generateActorId(leaderId + "-collector")));
+
+ TestActorRef<MockRaftActor> leaderActorRef = factory.createTestActor(
+ builder.props(), leaderId);
+ MockRaftActor leaderActor = leaderActorRef.underlyingActor();
+ leaderActor.waitForInitializeBehaviorComplete();
+
+ leaderActor.getRaftActorContext().getTermInformation().update(1, leaderId);
+ Leader leader = new Leader(leaderActor.getRaftActorContext());
+ leaderActor.setCurrentBehavior(leader);
+
+ final ExecutorService executorService = Executors.newSingleThreadExecutor();
+
+ leaderActor.setPersistence(new PersistentDataProvider(leaderActor) {
+ @Override
+ public <T> void persistAsync(final T entry, final Procedure<T> procedure) {
+ // needs to be executed from another thread to simulate the persistence actor calling this callback
+ executorService.submit(() -> {
+ try {
+ procedure.apply(entry);
+ } catch (Exception e) {
+ TEST_LOG.info("Fail during async persist callback", e);
+ }
+ }, "persistence-callback");
+ }
+ });
+
+ leader.getFollower(followerId).setNextIndex(0);
+ leader.getFollower(followerId).setMatchIndex(-1);
+
+ // hitting this is flimsy so run multiple times to improve the chance of things
+ // blowing up while breaking actor containment
+ final TestPersist message =
+ new TestPersist(leaderActorRef, new MockIdentifier("1"), new MockPayload("1"));
+ for (int i = 0; i < 100; i++) {
+ leaderActorRef.tell(message, null);
+
+ AppendEntriesReply reply =
+ new AppendEntriesReply(followerId, 1, true, i, 1, (short) 5);
+ leaderActorRef.tell(reply, mockFollowerActorRef);
+ }
+
+ await("Persistence callback.").atMost(5, TimeUnit.SECONDS).until(() -> leaderActor.getState().size() == 100);
+ executorService.shutdown();
+ }