+ @Test
+ public void testElectionScheduledWhenAnyRaftRPCReceived() {
+ MockRaftActorContext context = createActorContext();
+ follower = createBehavior(context);
+ follower.handleMessage(leaderActor, new RaftRPC() {
+ private static final long serialVersionUID = 1L;
+
+ @Override
+ public long getTerm() {
+ return 100;
+ }
+ });
+ verify(follower).scheduleElection(any(FiniteDuration.class));
+ }
+
+ @Test
+ public void testElectionNotScheduledWhenNonRaftRPCMessageReceived() {
+ MockRaftActorContext context = createActorContext();
+ follower = createBehavior(context);
+ follower.handleMessage(leaderActor, "non-raft-rpc");
+ verify(follower, never()).scheduleElection(any(FiniteDuration.class));
+ }
+
+ @Test
+ public void testCaptureSnapshotOnLastEntryInAppendEntries() {
+ String id = "testCaptureSnapshotOnLastEntryInAppendEntries";
+ logStart(id);
+
+ InMemoryJournal.addEntry(id, 1, new UpdateElectionTerm(1, null));
+
+ DefaultConfigParamsImpl config = new DefaultConfigParamsImpl();
+ config.setSnapshotBatchCount(2);
+ config.setCustomRaftPolicyImplementationClass(DisableElectionsRaftPolicy.class.getName());
+
+ final AtomicReference<MockRaftActor> followerRaftActor = new AtomicReference<>();
+ RaftActorSnapshotCohort snapshotCohort = newRaftActorSnapshotCohort(followerRaftActor);
+ Builder builder = MockRaftActor.builder().persistent(Optional.of(true)).id(id)
+ .peerAddresses(ImmutableMap.of("leader", "")).config(config).snapshotCohort(snapshotCohort);
+ TestActorRef<MockRaftActor> followerActorRef = actorFactory.createTestActor(builder.props()
+ .withDispatcher(Dispatchers.DefaultDispatcherId()), id);
+ followerRaftActor.set(followerActorRef.underlyingActor());
+ followerRaftActor.get().waitForInitializeBehaviorComplete();
+
+ InMemorySnapshotStore.addSnapshotSavedLatch(id);
+ InMemoryJournal.addDeleteMessagesCompleteLatch(id);
+ InMemoryJournal.addWriteMessagesCompleteLatch(id, 1, ApplyJournalEntries.class);
+
+ List<ReplicatedLogEntry> entries = Arrays.asList(
+ newReplicatedLogEntry(1, 0, "one"), newReplicatedLogEntry(1, 1, "two"));
+
+ AppendEntries appendEntries = new AppendEntries(1, "leader", -1, -1, entries, 1, -1, (short)0);
+
+ followerActorRef.tell(appendEntries, leaderActor);
+
+ AppendEntriesReply reply = MessageCollectorActor.expectFirstMatching(leaderActor, AppendEntriesReply.class);
+ assertEquals("isSuccess", true, reply.isSuccess());
+
+ final Snapshot snapshot = InMemorySnapshotStore.waitForSavedSnapshot(id, Snapshot.class);
+
+ InMemoryJournal.waitForDeleteMessagesComplete(id);
+ InMemoryJournal.waitForWriteMessagesComplete(id);
+ // We expect the ApplyJournalEntries for index 1 to remain in the persisted log b/c it's still queued for
+ // persistence by the time we initiate capture so the last persisted journal sequence number doesn't include it.
+ // This is OK - on recovery it will be a no-op since index 1 has already been applied.
+ List<Object> journalEntries = InMemoryJournal.get(id, Object.class);
+ assertEquals("Persisted journal entries size: " + journalEntries, 1, journalEntries.size());
+ assertEquals("Persisted journal entry type", ApplyJournalEntries.class, journalEntries.get(0).getClass());
+ assertEquals("ApplyJournalEntries index", 1, ((ApplyJournalEntries)journalEntries.get(0)).getToIndex());
+
+ assertEquals("Snapshot unapplied size", 0, snapshot.getUnAppliedEntries().size());
+ assertEquals("Snapshot getLastAppliedTerm", 1, snapshot.getLastAppliedTerm());
+ assertEquals("Snapshot getLastAppliedIndex", 1, snapshot.getLastAppliedIndex());
+ assertEquals("Snapshot getLastTerm", 1, snapshot.getLastTerm());
+ assertEquals("Snapshot getLastIndex", 1, snapshot.getLastIndex());
+ assertEquals("Snapshot state", ImmutableList.of(entries.get(0).getData(), entries.get(1).getData()),
+ MockRaftActor.fromState(snapshot.getState()));
+ }
+
+ @Test
+ public void testCaptureSnapshotOnMiddleEntryInAppendEntries() {
+ String id = "testCaptureSnapshotOnMiddleEntryInAppendEntries";
+ logStart(id);
+
+ InMemoryJournal.addEntry(id, 1, new UpdateElectionTerm(1, null));
+
+ DefaultConfigParamsImpl config = new DefaultConfigParamsImpl();
+ config.setSnapshotBatchCount(2);
+ config.setCustomRaftPolicyImplementationClass(DisableElectionsRaftPolicy.class.getName());
+
+ final AtomicReference<MockRaftActor> followerRaftActor = new AtomicReference<>();
+ RaftActorSnapshotCohort snapshotCohort = newRaftActorSnapshotCohort(followerRaftActor);
+ Builder builder = MockRaftActor.builder().persistent(Optional.of(true)).id(id)
+ .peerAddresses(ImmutableMap.of("leader", "")).config(config).snapshotCohort(snapshotCohort);
+ TestActorRef<MockRaftActor> followerActorRef = actorFactory.createTestActor(builder.props()
+ .withDispatcher(Dispatchers.DefaultDispatcherId()), id);
+ followerRaftActor.set(followerActorRef.underlyingActor());
+ followerRaftActor.get().waitForInitializeBehaviorComplete();
+
+ InMemorySnapshotStore.addSnapshotSavedLatch(id);
+ InMemoryJournal.addDeleteMessagesCompleteLatch(id);
+ InMemoryJournal.addWriteMessagesCompleteLatch(id, 1, ApplyJournalEntries.class);
+
+ List<ReplicatedLogEntry> entries = Arrays.asList(
+ newReplicatedLogEntry(1, 0, "one"), newReplicatedLogEntry(1, 1, "two"),
+ newReplicatedLogEntry(1, 2, "three"));
+
+ AppendEntries appendEntries = new AppendEntries(1, "leader", -1, -1, entries, 2, -1, (short)0);
+
+ followerActorRef.tell(appendEntries, leaderActor);
+
+ AppendEntriesReply reply = MessageCollectorActor.expectFirstMatching(leaderActor, AppendEntriesReply.class);
+ assertEquals("isSuccess", true, reply.isSuccess());
+
+ final Snapshot snapshot = InMemorySnapshotStore.waitForSavedSnapshot(id, Snapshot.class);
+
+ InMemoryJournal.waitForDeleteMessagesComplete(id);
+ InMemoryJournal.waitForWriteMessagesComplete(id);
+ // We expect the ApplyJournalEntries for index 2 to remain in the persisted log b/c it's still queued for
+ // persistence by the time we initiate capture so the last persisted journal sequence number doesn't include it.
+ // This is OK - on recovery it will be a no-op since index 2 has already been applied.
+ List<Object> journalEntries = InMemoryJournal.get(id, Object.class);
+ assertEquals("Persisted journal entries size: " + journalEntries, 1, journalEntries.size());
+ assertEquals("Persisted journal entry type", ApplyJournalEntries.class, journalEntries.get(0).getClass());
+ assertEquals("ApplyJournalEntries index", 2, ((ApplyJournalEntries)journalEntries.get(0)).getToIndex());
+
+ assertEquals("Snapshot unapplied size", 0, snapshot.getUnAppliedEntries().size());
+ assertEquals("Snapshot getLastAppliedTerm", 1, snapshot.getLastAppliedTerm());
+ assertEquals("Snapshot getLastAppliedIndex", 2, snapshot.getLastAppliedIndex());
+ assertEquals("Snapshot getLastTerm", 1, snapshot.getLastTerm());
+ assertEquals("Snapshot getLastIndex", 2, snapshot.getLastIndex());
+ assertEquals("Snapshot state", ImmutableList.of(entries.get(0).getData(), entries.get(1).getData(),
+ entries.get(2).getData()), MockRaftActor.fromState(snapshot.getState()));
+
+ assertEquals("Journal size", 0, followerRaftActor.get().getReplicatedLog().size());
+ assertEquals("Snapshot index", 2, followerRaftActor.get().getReplicatedLog().getSnapshotIndex());
+
+ // Reinstate the actor from persistence
+
+ actorFactory.killActor(followerActorRef, new TestKit(getSystem()));
+
+ followerActorRef = actorFactory.createTestActor(builder.props()
+ .withDispatcher(Dispatchers.DefaultDispatcherId()), id);
+ followerRaftActor.set(followerActorRef.underlyingActor());
+ followerRaftActor.get().waitForInitializeBehaviorComplete();
+
+ assertEquals("Journal size", 0, followerRaftActor.get().getReplicatedLog().size());
+ assertEquals("Last index", 2, followerRaftActor.get().getReplicatedLog().lastIndex());
+ assertEquals("Last applied index", 2, followerRaftActor.get().getRaftActorContext().getLastApplied());
+ assertEquals("Commit index", 2, followerRaftActor.get().getRaftActorContext().getCommitIndex());
+ assertEquals("State", ImmutableList.of(entries.get(0).getData(), entries.get(1).getData(),
+ entries.get(2).getData()), followerRaftActor.get().getState());
+ }
+
+ @Test
+ public void testCaptureSnapshotOnAppendEntriesWithUnapplied() {
+ String id = "testCaptureSnapshotOnAppendEntriesWithUnapplied";
+ logStart(id);
+
+ InMemoryJournal.addEntry(id, 1, new UpdateElectionTerm(1, null));
+
+ DefaultConfigParamsImpl config = new DefaultConfigParamsImpl();
+ config.setSnapshotBatchCount(1);
+ config.setCustomRaftPolicyImplementationClass(DisableElectionsRaftPolicy.class.getName());
+
+ final AtomicReference<MockRaftActor> followerRaftActor = new AtomicReference<>();
+ RaftActorSnapshotCohort snapshotCohort = newRaftActorSnapshotCohort(followerRaftActor);
+ Builder builder = MockRaftActor.builder().persistent(Optional.of(true)).id(id)
+ .peerAddresses(ImmutableMap.of("leader", "")).config(config).snapshotCohort(snapshotCohort);
+ TestActorRef<MockRaftActor> followerActorRef = actorFactory.createTestActor(builder.props()
+ .withDispatcher(Dispatchers.DefaultDispatcherId()), id);
+ followerRaftActor.set(followerActorRef.underlyingActor());
+ followerRaftActor.get().waitForInitializeBehaviorComplete();
+
+ InMemorySnapshotStore.addSnapshotSavedLatch(id);
+ InMemoryJournal.addDeleteMessagesCompleteLatch(id);
+ InMemoryJournal.addWriteMessagesCompleteLatch(id, 1, ApplyJournalEntries.class);
+
+ List<ReplicatedLogEntry> entries = Arrays.asList(
+ newReplicatedLogEntry(1, 0, "one"), newReplicatedLogEntry(1, 1, "two"),
+ newReplicatedLogEntry(1, 2, "three"));
+
+ AppendEntries appendEntries = new AppendEntries(1, "leader", -1, -1, entries, 0, -1, (short)0);
+
+ followerActorRef.tell(appendEntries, leaderActor);
+
+ AppendEntriesReply reply = MessageCollectorActor.expectFirstMatching(leaderActor, AppendEntriesReply.class);
+ assertEquals("isSuccess", true, reply.isSuccess());
+
+ final Snapshot snapshot = InMemorySnapshotStore.waitForSavedSnapshot(id, Snapshot.class);
+
+ InMemoryJournal.waitForDeleteMessagesComplete(id);
+ InMemoryJournal.waitForWriteMessagesComplete(id);
+ // We expect the ApplyJournalEntries for index 0 to remain in the persisted log b/c it's still queued for
+ // persistence by the time we initiate capture so the last persisted journal sequence number doesn't include it.
+ // This is OK - on recovery it will be a no-op since index 0 has already been applied.
+ List<Object> journalEntries = InMemoryJournal.get(id, Object.class);
+ assertEquals("Persisted journal entries size: " + journalEntries, 1, journalEntries.size());
+ assertEquals("Persisted journal entry type", ApplyJournalEntries.class, journalEntries.get(0).getClass());
+ assertEquals("ApplyJournalEntries index", 0, ((ApplyJournalEntries)journalEntries.get(0)).getToIndex());
+
+ assertEquals("Snapshot unapplied size", 2, snapshot.getUnAppliedEntries().size());
+ assertEquals("Snapshot unapplied entry index", 1, snapshot.getUnAppliedEntries().get(0).getIndex());
+ assertEquals("Snapshot unapplied entry index", 2, snapshot.getUnAppliedEntries().get(1).getIndex());
+ assertEquals("Snapshot getLastAppliedTerm", 1, snapshot.getLastAppliedTerm());
+ assertEquals("Snapshot getLastAppliedIndex", 0, snapshot.getLastAppliedIndex());
+ assertEquals("Snapshot getLastTerm", 1, snapshot.getLastTerm());
+ assertEquals("Snapshot getLastIndex", 2, snapshot.getLastIndex());
+ assertEquals("Snapshot state", ImmutableList.of(entries.get(0).getData()),
+ MockRaftActor.fromState(snapshot.getState()));
+ }
+
+ @Test
+ public void testNeedsLeaderAddress() {
+ logStart("testNeedsLeaderAddress");
+
+ MockRaftActorContext context = createActorContext();
+ context.setReplicatedLog(new MockRaftActorContext.SimpleReplicatedLog());
+ context.addToPeers("leader", null, VotingState.VOTING);
+ ((DefaultConfigParamsImpl)context.getConfigParams()).setPeerAddressResolver(NoopPeerAddressResolver.INSTANCE);
+
+ follower = createBehavior(context);
+
+ follower.handleMessage(leaderActor,
+ new AppendEntries(1, "leader", -1, -1, Collections.emptyList(), -1, -1, (short)0));
+
+ AppendEntriesReply reply = MessageCollectorActor.expectFirstMatching(leaderActor, AppendEntriesReply.class);
+ assertTrue(reply.isNeedsLeaderAddress());
+ MessageCollectorActor.clearMessages(leaderActor);
+
+ PeerAddressResolver mockResolver = mock(PeerAddressResolver.class);
+ ((DefaultConfigParamsImpl)context.getConfigParams()).setPeerAddressResolver(mockResolver);
+
+ follower.handleMessage(leaderActor, new AppendEntries(1, "leader", -1, -1, Collections.emptyList(), -1, -1,
+ (short)0, RaftVersions.CURRENT_VERSION, leaderActor.path().toString()));
+
+ reply = MessageCollectorActor.expectFirstMatching(leaderActor, AppendEntriesReply.class);
+ assertFalse(reply.isNeedsLeaderAddress());
+
+ verify(mockResolver).setResolved("leader", leaderActor.path().toString());
+ }
+
+ @SuppressWarnings("checkstyle:IllegalCatch")
+ private static RaftActorSnapshotCohort newRaftActorSnapshotCohort(
+ final AtomicReference<MockRaftActor> followerRaftActor) {
+ RaftActorSnapshotCohort snapshotCohort = new RaftActorSnapshotCohort() {
+ @Override
+ public void createSnapshot(final ActorRef actorRef, final Optional<OutputStream> installSnapshotStream) {
+ try {
+ actorRef.tell(new CaptureSnapshotReply(new MockSnapshotState(followerRaftActor.get().getState()),
+ installSnapshotStream), actorRef);
+ } catch (RuntimeException e) {
+ throw e;
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ @Override
+ public void applySnapshot(final State snapshotState) {
+ }
+
+ @Override
+ public State deserializeSnapshot(final ByteSource snapshotBytes) {
+ throw new UnsupportedOperationException();
+ }
+ };
+ return snapshotCohort;
+ }
+
+ public byte[] getNextChunk(final ByteString bs, final int offset, final int chunkSize) {