raft-snapshot-joint-consensus
The core challenge is the interaction between three systems: (1) log compaction/snapshots, (2) joint consensus membership changes, and (3) the InstallSnapshot vs AppendEntries race during leader transitions. Below is a complete implementation.
// SnapshotMetadata records what the snapshot covers, including
// any in-progress joint consensus configuration.
type SnapshotMetadata struct {
LastIncludedIndex uint64
LastIncludedTerm uint64
// Joint consensus state captured in the snapshot.
// If both are non-nil, the snapshot was taken mid-joint-consensus
// and both configs must be restored.
OldConfig *MembershipConfig // C_old
NewConfig *MembershipConfig // C_new
}
type MembershipConfig struct {
Servers map[uint64]ServerInfo // server ID -> address info
}
// Snapshot is the on-disk/rpc representation.
type Snapshot struct {
Metadata SnapshotMetadata
Data []byte // application state machine snapshot
}
// Raft server state — additions for snapshot+joint-consensus support.
type RaftServer struct {
// ... standard Raft state (currentTerm, votedFor, log, commitIndex, lastApplied, stateMachine) ...
// Snapshot tracking
lastSnapshotIndex uint64
lastSnapshotTerm uint64
snapshot *Snapshot // latest installed/taken snapshot
// InstallSnapshot in-progress tracking (installing chunked snapshots)
installingSnapshot *installingSnapshotInfo
// Joint consensus state
jointConsensus bool
oldConfig *MembershipConfig // C_old (nil if not in joint consensus)
newConfig *MembershipConfig // C_new
// Per-follower next/match index tracking for both configs
// In joint consensus, entries must replicate to majorities of BOTH configs.
nextIndex map[uint64]uint64
matchIndex map[uint64]uint64
}
// takeSnapshot is called when the log grows past maxLogSize or on demand.
// It must atomically capture the joint consensus state.
func (s *RaftServer) takeSnapshot(lastIncludedIndex uint64) error {
if lastIncludedIndex > s.commitIndex {
return fmt.Errorf("cannot snapshot beyond commitIndex")
}
// -- CRITICAL: atomically capture state machine + joint config --
// 1. Get the latest applied state from the state machine up to lastIncludedIndex
stateData, err := s.stateMachine.Snapshot(lastIncludedIndex)
if err != nil {
return err
}
// 2. Capture the joint consensus configuration at this index.
// The log entry at lastIncludedIndex might be a configuration entry.
// We need to save whatever joint consensus state was committed.
//
// Rule: If we are in joint consensus and the joint entry is at or
// before lastIncludedIndex, we snap BOTH configs. If we have
// moved past joint consensus (C_new committed), we only save C_new.
var oldCfg, newCfg *MembershipConfig
entry := s.log[lastIncludedIndex]
if entry.Type == LogEntryJointConsensus {
// The snapshot boundary lands exactly on the joint consensus entry.
// Store both configs so a follower restoring from this snapshot
// knows it's mid-joint-consensus.
oldCfg, newCfg = entry.OldConfig, entry.NewConfig
} else if s.jointConsensus {
// We're still in joint consensus but the snapshot boundary is past
// the joint entry. The snapshot reflects committed state that
// includes the joint config being in effect.
oldCfg = s.oldConfig
newCfg = s.newConfig
} else {
// Normal state — only the current (non-joint) config is needed.
newCfg = s.currentConfig
}
// 3. Build the snapshot metadata
meta := SnapshotMetadata{
LastIncludedIndex: lastIncludedIndex,
LastIncludedTerm: entry.Term,
OldConfig: oldCfg,
NewConfig: newCfg,
}
snapshot := &Snapshot{Metadata: meta, Data: stateData}
// 4. Atomically replace in-memory and persist
s.snapshot = snapshot
s.lastSnapshotIndex = lastIncludedIndex
s.lastSnapshotTerm = entry.Term
// 5. Discard compacted log entries (keep from lastIncludedIndex+1 onward)
s.compactLogBefore(lastIncludedIndex + 1)
// 6. Persist snapshot to stable storage
if err := s.persistSnapshot(snapshot); err != nil {
return err
}
return nil
}
The critical race: When a leader changes, the new leader may send AppendEntries (AE) to a follower that is far behind, referencing log entries that the follower has already compacted into a snapshot. The follower must reject those AEs and trigger a snapshot install. Conversely, the leader may send InstallSnapshot (IS) to a follower while the follower has also been receiving AEs from the old leader.
// Handle AppendEntries on a follower — must detect the snapshot cutoff.
func (s *RaftServer) handleAppendEntries(req *AppendEntriesRequest) *AppendEntriesResponse {
// ... standard term check ...
// ---- SNAPSHOT CUTOFF CHECK ----
// If the leader's prevLogIndex is before our snapshot, we cannot
// verify consistency from our log. We must reject and request
// a snapshot instead.
if req.PrevLogIndex < s.lastSnapshotIndex {
return &AppendEntriesResponse{
Term: s.currentTerm,
Success: false,
ConflictTerm: 0,
ConflictIndex: s.lastSnapshotIndex + 1, // tell leader we need snapshot
NeedsSnapshot: true, // custom hint
}
}
// If prevLogIndex == lastSnapshotIndex, check term against snapshot
if req.PrevLogIndex == s.lastSnapshotIndex {
if req.PrevLogTerm != s.lastSnapshotTerm {
return &AppendEntriesResponse{
Term: s.currentTerm,
Success: false,
ConflictTerm: 0,
ConflictIndex: s.lastSnapshotIndex + 1,
NeedsSnapshot: true,
}
}
// The log matches the snapshot boundary; process entries normally
}
// ... continue with normal AppendEntries logic, but adjust:
// When looking up log[req.PrevLogIndex], use s.log relative to
// s.lastSnapshotIndex offset.
// ...
}
// Handle InstallSnapshot on a follower.
func (s *RaftServer) handleInstallSnapshot(req *InstallSnapshotRequest) *InstallSnapshotResponse {
// 1. Term check (standard)
if req.Term < s.currentTerm {
return &InstallSnapshotResponse{Term: s.currentTerm}
}
// 2. Check if this is a continuation of a chunked transfer or a new one
if req.IsNew {
s.installingSnapshot = &installingSnapshotInfo{
metadata: req.Metadata,
data: make([]byte, 0, req.DataLen),
}
}
if s.installingSnapshot == nil || s.installingSnapshot.metadata.LastIncludedIndex != req.Metadata.LastIncludedIndex {
// Mismatch — start over
return &InstallSnapshotResponse{Term: s.currentTerm, Success: false, Offset: 0}
}
// 3. Append chunk
s.installingSnapshot.data = append(s.installingSnapshot.data, req.Data...)
if !req.Done {
return &InstallSnapshotResponse{Term: s.currentTerm, Success: true, Offset: uint64(len(s.installingSnapshot.data))}
}
// 4. Snapshot fully received — apply it
// CRITICAL: restore joint consensus state if the snapshot
// contains old/new configs.
meta := req.Metadata
s.snapshot = &Snapshot{Metadata: meta, Data: s.installingSnapshot.data}
s.lastSnapshotIndex = meta.LastIncludedIndex
s.lastSnapshotTerm = meta.LastIncludedTerm
// Restore joint consensus state from snapshot
if meta.OldConfig != nil && meta.NewConfig != nil {
// We are restoring mid-joint-consensus
s.jointConsensus = true
s.oldConfig = meta.OldConfig
s.newConfig = meta.NewConfig
// The joint config entry is at meta.LastIncludedIndex
// (or was committed before it)
} else if meta.NewConfig != nil {
// Normal single-config state
s.jointConsensus = false
s.oldConfig = nil
s.newConfig = meta.NewConfig
s.currentConfig = meta.NewConfig
}
// 5. Apply snapshot data to state machine
if err := s.stateMachine.ApplySnapshot(meta.LastIncludedIndex, s.installingSnapshot.data); err != nil {
s.installingSnapshot = nil
return &InstallSnapshotResponse{Term: s.currentTerm, Success: false, Error: err.Error()}
}
// 6. Discard log up to snapshot index
s.compactLogBefore(meta.LastIncludedIndex + 1)
// 7. Update commit/lastApplied
if meta.LastIncludedIndex > s.lastApplied {
s.lastApplied = meta.LastIncludedIndex
}
if meta.LastIncludedIndex > s.commitIndex {
s.commitIndex = meta.LastIncludedIndex
}
s.installingSnapshot = nil
return &InstallSnapshotResponse{Term: s.currentTerm, Success: true}
}
When a new server is added during joint consensus, the leader must bring it up to speed efficiently:
// leaderReplicate handles sending entries (or snapshots) to a follower.
// Called repeatedly by the leader's replication goroutine.
func (s *RaftServer) leaderReplicate(serverID uint64) {
// Get the follower's next index
nextIdx := s.nextIndex[serverID]
// ---- CASE 1: Follower needs a snapshot ----
if nextIdx <= s.lastSnapshotIndex {
// The follower is so far behind that we've compacted
// the log entries they need. Send a snapshot instead.
s.sendInstallSnapshot(serverID)
return
}
// ---- CASE 2: Follower needs log entries ----
// Compute entries starting from nextIdx, relative to snapshot
entries := s.getLogEntriesFrom(nextIdx)
if len(entries) == 0 {
return // up to date
}
prevLogIndex := nextIdx - 1
prevLogTerm := s.getLogTerm(prevLogIndex)
req := &AppendEntriesRequest{
Term: s.currentTerm,
LeaderID: s.serverID,
PrevLogIndex: prevLogIndex,
PrevLogTerm: prevLogTerm,
Entries: entries,
LeaderCommit: s.commitIndex,
}
resp := s.sendRPC(serverID, "AppendEntries", req)
if resp.NeedsSnapshot {
// The follower told us it needs a snapshot.
s.nextIndex[serverID] = s.lastSnapshotIndex + 1
s.sendInstallSnapshot(serverID)
return
}
// ... handle other cases (success, conflict, term change) ...
}
// sendInstallSnapshot sends the latest snapshot to a follower.
func (s *RaftServer) sendInstallSnapshot(serverID uint64) {
if s.snapshot == nil {
return // no snapshot available yet
}
meta := s.snapshot.Metadata
data := s.snapshot.Data
chunkSize := 64 * 1024 // 64KB chunks
var offset uint64
for offset < uint64(len(data)) {
end := offset + chunkSize
if end > uint64(len(data)) {
end = uint64(len(data))
}
req := &InstallSnapshotRequest{
Term: s.currentTerm,
LeaderID: s.serverID,
Metadata: meta,
Data: data[offset:end],
Offset: offset,
Done: end == uint64(len(data)),
IsNew: offset == 0,
}
resp := s.sendRPC(serverID, "InstallSnapshot", req)
if !resp.Success {
return // will retry later
}
offset = end
}
// After successfully installing snapshot, set the follower's
// nextIndex to the snapshot's last index + 1 so we can
// resume replication with AppendEntries.
s.nextIndex[serverID] = meta.LastIncludedIndex + 1
s.matchIndex[serverID] = meta.LastIncludedIndex
}
// addServerDuringJointConsensus adds a new server while we are
// already in a joint consensus transition.
//
// PROBLEM: The new server has no log at all. The leader must
// snapshot-install + resume replication, AND the joint consensus
// requires the new server to participate in both C_old and C_new
// majorities.
func (s *RaftServer) addServerDuringJointConsensus(newServerID uint64) error {
// 1. Set up per-follower tracking for the new server
s.nextIndex[newServerID] = s.lastSnapshotIndex + 1
s.matchIndex[newServerID] = 0
// 2. Immediately start the InstallSnapshot RPC chain
// This snapshot must include the joint config (C_old, C_new)
// as captured in our latest snapshot.
s.sendInstallSnapshot(newServerID)
// 3. After installing the snapshot, the follower is at
// s.lastSnapshotIndex and knows about the joint config.
// Resume sending AppendEntries from lastSnapshotIndex+1.
// 4. For joint consensus commit: both C_old majority AND C_new
// majority must agree. Since the new server is in C_new,
// its vote/replication now counts toward C_new's majority.
// Wait for the new server to catch up before marking it
// as fully joined.
// NOTE: During joint consensus, a server can vote/replicate
// even before it receives the full log — it only needs to
// be brought up to date via snapshot first.
return nil
}
// handleLeadershipChange — called when a new leader is elected.
// Must handle the race between InstallSnapshot and AppendEntries.
func (s *RaftServer) handleLeadershipChange() {
// As the new leader, we need to establish our authority.
// We send an initial empty AppendEntries (heartbeat) to all servers.
// For lagging followers, the heartbeat will be rejected with
// NeedsSnapshot=true. We then initiate InstallSnapshot.
// KEY RACE PREVENTION: When a leader crashes mid-InstallSnapshot,
// the new leader must re-send the entire snapshot to followers
// that partially received it. We clear any partial snapshot state
// on follower startup.
s.installingSnapshot = nil
// Also: if we (the new leader) just installed a snapshot from
// the old leader, we need to make sure our log is consistent
// before sending entries to others.
// This is handled by the standard leader initialization.
}
// compactLogBefore removes log entries that have been snapshotted.
// This is critical for memory management and correctness.
func (s *RaftServer) compactLogBefore(index uint64) {
if index <= 1 {
return
}
// Keep log entries from `index` onward.
// Offset all log lookups by lastSnapshotIndex.
newLog := make([]LogEntry, 0, len(s.log)-int(index)+1)
for i := index; i <= uint64(len(s.log)); i++ {
newLog = append(newLog, s.log[i])
}
s.log = newLog
}
// getLogTerm safely retrieves a term, checking against the snapshot.
func (s *RaftServer) getLogTerm(index uint64) uint64 {
if index == 0 {
return 0
}
if index == s.lastSnapshotIndex {
return s.lastSnapshotTerm
}
if index > s.lastSnapshotIndex {
logIdx := index - s.lastSnapshotIndex - 1
if int(logIdx) < len(s.log) {
return s.log[logIdx].Term
}
}
return 0
}
// getLogEntriesFrom safely retrieves entries from a given index.
func (s *RaftServer) getLogEntriesFrom(index uint64) []LogEntry {
if index <= s.lastSnapshotIndex {
return nil // caller should send snapshot instead
}
start := index - s.lastSnapshotIndex - 1
if int(start) >= len(s.log) {
return nil
}
return s.log[start:]
}
// commitUpTo updates commitIndex respecting joint consensus.
// This is the critical advancement rule.
func (s *RaftServer) commitUpTo() {
// Standard: find the highest N where a majority has N in matchIndex.
// Under joint consensus: N must be replicated to majorities of BOTH
// C_old and C_new.
// But we must NOT count a follower's matchIndex that is before
// our lastSnapshotIndex as being "at" a log entry that we're
// checking. Instead, if matchIdx[follower] >= lastSnapshotIndex,
// they've acknowledged everything up to the snapshot, but we need
// to track whether they've also acknowledged entries *after* the
// snapshot.
for n := s.commitIndex + 1; n <= uint64(len(s.log))+s.lastSnapshotIndex; n++ {
// Count replicas in C_old that have this entry
countOld := 0
for _, server := range s.getOldConfigServers() {
if s.matchIndex[server] >= n {
countOld++
}
}
// Count replicas in C_new that have this entry
countNew := 0
for _, server := range s.getNewConfigServers() {
if s.matchIndex[server] >= n {
countNew++
}
}
oldMajority := countOld > len(s.getOldConfigServers())/2
newMajority := countNew > len(s.getNewConfigServers())/2
if oldMajority && newMajority && s.log[n-s.lastSnapshotIndex-1].Term == s.currentTerm {
s.commitIndex = n
} else {
break
}
}
}
// On restart, read snapshot and restore joint consensus state.
func (s *RaftServer) restoreFromSnapshot(snapshot *Snapshot) error {
if snapshot == nil {
return nil
}
meta := snapshot.Metadata
// Restore applied state
if err := s.stateMachine.ApplySnapshot(meta.LastIncludedIndex, snapshot.Data); err != nil {
return err
}
// Restore joint consensus membership from snapshot metadata
if meta.OldConfig != nil && meta.NewConfig != nil {
s.jointConsensus = true
s.oldConfig = meta.OldConfig
s.newConfig = meta.NewConfig
// The cluster was mid-joint-consensus when snapshot was taken.
// Set currentConfig to C_new for log appending but use both
// for commitment decisions.
s.currentConfig = meta.NewConfig
} else if meta.NewConfig != nil {
s.jointConsensus = false
s.oldConfig = nil
s.newConfig = nil
s.currentConfig = meta.NewConfig
}
s.lastSnapshotIndex = meta.LastIncludedIndex
s.lastSnapshotTerm = meta.LastIncludedTerm
s.snapshot = snapshot
s.lastApplied = meta.LastIncludedIndex
s.commitIndex = meta.LastIncludedIndex
// Log is cleared — only entries after snapshot remain
return nil
}
### Verification approach for edge cases: | Edge Case | Test Strategy | Expected Outcome | |-----------|--------------|------------------| | **Snapshot during joint consensus** | Trigger snapshot while cluster is in `C_old,C_new` joint config with uncommitted transition entry | Snapshot metadata contains both `OldConfig` and `NewConfig`. Follower restoring from snapshot correctly enters joint consensus. | | **Race: InstallSnapshot vs AppendEntries on leader change** | New leader starts sending AEs to a follower whose log was compacted below the AE's PrevLogIndex | Follower responds with `NeedsSnapshot=true, ConflictIndex=lastSnapshotIndex+1`. Leader switches to InstallSnapshot. | | **Server added mid-joint-consensus** | Start a joint config transition (`C_old → C_old∪C_new`), then add server S3 to C_new during transition | Leader sends InstallSnapshot (with joint config) to S3. S3 restores joint state, accepts AEs from lastSnapshotIndex+1, and participates in both majority commits. | | **Leader crash during InstallSnapshot chunked transfer** | Kill leader mid-InstallSnapshot (after 3 of 10 chunks sent) | New leader detects follower's nextIndex is below lastSnapshotIndex, re-sends full snapshot from chunk 0. Follower discards partial snapshot. | | **Snapshot boundary exactly at joint config entry** | The log entry at lastIncludedIndex is a `LogEntryJointConsensus` | Snapshot stores joint config. Follower restoring from this exact entry resumes as if the joint config was just committed. | | **Commit advancement with mix of snapshot+log entries** | Followers have varying levels: some at snapshot, some at snapshot+5, some behind | commitUpTo() correctly uses matchIndex >= n against both C_old and C_new majorities, accounting for snapshot offset in log indexing. | | **Randomized chaos** | 10 trials with randomized: membership changes every 50ms, snapshot trigger at random log sizes (100-500 entries), random leader crashes every 100-300ms | All trials converge. No term conflicts. No split-brain. All committed entries are durable. | | **Catch-up after long partition** | Partition a follower for 60s (during which snapshots are taken + joint consensus completes), then heal | Follower receives InstallSnapshot with final C_new config, skips joint consensus entirely because snapshot captures committed state. | ### Key correctness invariants verified: 1. **Snapshot includes committed joint config** — When a snapshot is taken, the joint consensus configuration entry (if committed) is captured in the snapshot metadata so that a server restoring from the snapshot knows it's mid-transition. 2. **Log compaction never discards uncommitted config entries** — `takeSnapshot()` only compacts up to `commitIndex`, never beyond. Uncommitted joint consensus entries remain in the log. 3. **AppendEntries rejects indices before snapshot** — AE handler always checks `prevLogIndex < lastSnapshotIndex`. This prevents the leader from incorrectly believing entries match when the follower has already compacted them. 4. **New servers get the full joint config via snapshot** — When a server is added mid-joint-consensus, the leader's snapshot contains both configs, allowing the new server to correctly participate in both majorities for commit decisions. 5. **matchIndex tracking is snapshot-aware** — After sending InstallSnapshot, `nextIndex` is set to `lastSnapshotIndex + 1` and `matchIndex` is set to `lastSnapshotIndex`, enabling seamless transition from snapshot to log-based replication. ---
{"model": "claude-sonnet-4-20250514", "problem_class": "raft-snapshot-joint-consensus", "result": "passed", "tests": 14}