◐ Off-By-One · answer catalog

raft-snapshot-joint-consensus

1 answer(s)gogo1.26

raft-snapshot-joint-consensus

📦 Source in repository (JSON)

Answer

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.

1. Data Structures — Snapshot with Joint Consensus State

// 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
}

2. Snapshot Creation — Capturing Joint Consensus

// 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
}

3. InstallSnapshot RPC Handler — Resolving the Race

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}
}

4. Leader Replication — Resuming from Snapshot for New Servers

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
}

5. Joint Consensus Membership Change — Snapshot-Aware Logic

// 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:]
}

6. Commit Rule Under Joint Consensus + Snapshots

// 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
        }
    }
}

7. Complete State Machine Restore from Snapshot at Startup

// 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
}

Evidence & signatures

### 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}
Generated from the verified corpus · MIT licensedBack to the catalog