+
+ long existingEntryTerm = getLogEntryTerm(matchEntry.getIndex());
+
+ log.debug("{}: matchEntry {} is present: existingEntryTerm: {}", logName(), matchEntry,
+ existingEntryTerm);
+
+ // existingEntryTerm == -1 means it's in the snapshot and not in the log. We don't know
+ // what the term was so we'll assume it matches.
+ if (existingEntryTerm == -1 || existingEntryTerm == matchEntry.getTerm()) {
+ continue;
+ }
+
+ if (!context.getRaftPolicy().applyModificationToStateBeforeConsensus()) {
+ log.info("{}: Removing entries from log starting at {}, commitIndex: {}, lastApplied: {}",
+ logName(), matchEntry.getIndex(), context.getCommitIndex(), context.getLastApplied());
+
+ // Entries do not match so remove all subsequent entries but only if the existing entries haven't
+ // been applied to the state yet.
+ if (matchEntry.getIndex() <= context.getLastApplied()
+ || !context.getReplicatedLog().removeFromAndPersist(matchEntry.getIndex())) {
+ // Could not remove the entries - this means the matchEntry index must be in the
+ // snapshot and not the log. In this case the prior entries are part of the state
+ // so we must send back a reply to force a snapshot to completely re-sync the
+ // follower's log and state.
+
+ log.info("{}: Could not remove entries - sending reply to force snapshot", logName());
+ sender.tell(new AppendEntriesReply(context.getId(), currentTerm(), false, lastIndex,
+ lastTerm(), context.getPayloadVersion(), true, needsLeaderAddress(),
+ appendEntries.getLeaderRaftVersion()), actor());
+ return false;
+ }
+
+ break;
+ } else {
+ sender.tell(new AppendEntriesReply(context.getId(), currentTerm(), false, lastIndex,
+ lastTerm(), context.getPayloadVersion(), true, needsLeaderAddress(),
+ appendEntries.getLeaderRaftVersion()), actor());
+ return false;
+ }
+ }
+ }
+
+ lastIndex = lastIndex();
+ log.debug("{}: After cleanup, lastIndex: {}, entries to be added from: {}", logName(), lastIndex,
+ addEntriesFrom);
+
+ // When persistence successfully completes for each new log entry appended, we need to determine if we
+ // should capture a snapshot to compact the persisted log. shouldCaptureSnapshot tracks whether or not
+ // one of the log entries has exceeded the log size threshold whereby a snapshot should be taken. However
+ // we don't initiate the snapshot at that log entry but rather after the last log entry has been persisted.
+ // This is done because subsequent log entries after the one that tripped the threshold may have been
+ // applied to the state already, as the persistence callback occurs async, and we want those entries
+ // purged from the persisted log as well.
+ final AtomicBoolean shouldCaptureSnapshot = new AtomicBoolean(false);
+ final Consumer<ReplicatedLogEntry> appendAndPersistCallback = logEntry -> {
+ final List<ReplicatedLogEntry> entries = appendEntries.getEntries();
+ final ReplicatedLogEntry lastEntryToAppend = entries.get(entries.size() - 1);
+ if (shouldCaptureSnapshot.get() && logEntry == lastEntryToAppend) {
+ context.getSnapshotManager().capture(context.getReplicatedLog().last(), getReplicatedToAllIndex());