Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

crash consensus: preserving decisions through replacement and recovery (authored by agents unless marked 🧑)

research recommendation

  • start with crashes that overlap snapshot creation and membership changes
    • test whether acknowledged operations survive restart and machine replacement
    • test whether the system resumes serving requests under its stated failure limits
  • agent hypothesis: these boundaries offer a tractable project with practical value
    • novelty is unverified
    • the useful result may be a reproducible failure, a recovery contract or a negative result
  • second choice: quorum selection that explicitly prices recovery
    • normal-operation latency alone misses the cost of changing leaders
  • source ledger
    • verified quotations, reading depth and missing papers

what needs to stay true

  • a replicated state machine runs copies of one deterministic service
    • agreement chooses the same operations
    • execution applies them in a compatible order
    • recovery preserves decisions already exposed to clients
  • vocabulary
    • safety: executions never contradict the promised behavior
    • progress: operations eventually finish under the stated assumptions
    • quorum: a set of machines whose responses suffice for one protocol step
    • configuration: the machines and voting rules currently permitted
    • snapshot: saved service state replacing an older part of the operation log
    • linearizability: operations fit one legal sequential execution
      • each completed operation appears once between its request and response
      • an operation finishing before another starts appears first
  • foundational constraint
    • Fischer, Lynch and Paterson: “the possibility of nontermination, even with only one faulty process”
    • inference: recovery bounds require a timing or failure-detection assumption
  • persisted voting state is part of the protocol
    • Lamport: “some information can be remembered by an agent that has failed and restarted”
    • inference: restarting a process with intact storage differs from restarting it after losing storage
  • scope of this study
    • machines stop, restart or become slow
    • lost storage must be detected and repaired through an explicitly permitted protocol
    • arbitrary false messages require a malicious-failure protocol
      • see the sibling study

what existing work already addresses

  • understandability and membership
    • Raft requires “separate majorities from both the old and new configurations”
    • inference: replacing all voters by merely updating an address list can destroy the evidence for earlier decisions
    • research must compare the implemented transition rule
      • joint transitions and constrained one-member changes are different designs
  • application snapshots and damaged disks
    • Chandra, Griesemer and Redstone require snapshot and log to be “mutually consistent”
    • their damaged replica “participates in Paxos as a non-voting member”
    • inference: checksums alone do not restore lost promises
      • restoring voting rights needs additional evidence
  • cheaper voting during normal operation
    • Howard, Malkhi and Spiegelman: “intersection is required only across phases”
    • illustrative calculation for simple cardinality quorums
      • with five voters, a two-voter acceptance quorum can intersect every four-voter recovery quorum
      • two plus four exceeds five
      • normal writes can continue with two available voters
        • the surviving leader already completed recovery for the relevant instances
        • no higher ballot displaced that leader
      • starting a new leader still requires four
    • assumptions for this example
      • ordinary Paxos with correct persisted promises
      • all two-voter and four-voter subsets permitted
      • no change to membership during the illustrated round
  • faster membership changes
    • Whittaker et al: “decouple reconfiguration from the standard processing path”
    • inference: the registry of earlier configurations creates another recovery obligation
      • test its failure, replacement and cleanup separately
  • removing the leader as a throughput bottleneck
    • Whittaker et al: “scaling these components independently”
    • inference: measuring extra workers without equal resource budgets confounds a protocol comparison
  • allowing operations to choose nearby coordinators
    • Moraru, Andersen and Kaminsky claim “graceful performance degradation when replicas are slow or crash”
    • Whittaker et al caution that EPaxos “had several bugs go undiscovered for years”
    • inference: include dependency recovery and execution delay
      • counting only chosen operations can hide client-visible stalls
  • deciding when operations can execute
    • Enes et al: “executes it only after the timestamp becomes stable”
    • inference: a protocol can choose an operation quickly and still delay its result
      • measure the entire request-to-result interval
  • optimizing practical quorums
    • Whittaker et al consider “machine heterogeneity and workload skew”
    • inference: adaptive quorum research must exceed static workload optimization

recent work changes the comparison

  • Pineapple, NSDI 2025
    • Bantikyan et al evaluate in-memory systems “except for the etcd experiments”
    • shared registers handle reads and writes
      • Multi-Paxos handles operations that read and modify values
    • appendix A.2 says “reads need to back off and retry” during leader changes
    • agent implication: compare operation classes and persistence separately
      • test ballot-aware reads during leader replacement
  • Picsou, OSDI 2025
    • Frank et al: “allows both crash fault tolerant and Byzantine fault tolerant protocols to communicate”
    • agent implication: state transfer between groups needs a precise receipt rule
      • one machine receiving bytes differs from a replicated group retaining an accepted message
  • XLL, NSDI 2026
    • Shawger et al use “two-phase recovery”
    • restore the Raft log before replaying service operations
      • service data and its applied index flush atomically
      • correctness relies on a valid persisted log prefix
    • cleanup uses the “last persisted applied index”
      • §4.2.3
      • flush service state before deleting recoverable log entries
      • preserve live value references during relocation
    • §5.4 tests six failure points with ten runs each
      • read-back checks cover the write path
      • membership and snapshot overlap are not established by these tests
    • agent implication: project A must exceed this existing recovery and cleanup design
      • target interactions between member replacement, snapshots and shared-log relocation

project A: make recovery boundaries explicit

  • question
    • what must a snapshot carry before a restored machine may vote and serve reads?
  • proposed contribution
    • a small recovery interface with checkable obligations
    • a fault-testing corpus across that interface
  • candidate recovery-interface obligations
    • identifies the applied log position and its term or ballot
    • identifies the applicable configuration
    • preserves client request identifiers needed to avoid duplicate execution
    • preserves local durable voting promises separately from the service snapshot
      • installing another machine’s snapshot cannot overwrite newer local promises
      • the snapshot’s last included term does not replace the current term or vote
    • retains every byte still referenced by recovery metadata
  • first experiment
    • choose one maintained Raft implementation with real persisted state
    • add a small key-value service and a client history recorder
    • begin a snapshot while adding a replacement machine
    • crash a voter at each persistence boundary
    • repeat with a leader crash and delayed snapshot chunks
    • restart with intact storage
      • separately test explicitly detected disk loss
  • check
    • acknowledged writes remain observable
    • each recorded history has a valid linearizable explanation
    • duplicate client requests do not repeat application effects
    • surviving eligible voters eventually elect a leader after bounded delays return
  • baselines
    • unchanged implementation
    • existing snapshot and membership tests
    • its documented storage interface
  • output worth keeping
    • minimal reproducible execution with exact persisted files
    • invariant explaining the failure or why the interface excludes it
    • repair whose extra disk writes and recovery time are measured
  • stop or redirect
    • if existing tests cover the whole boundary, publish a clear coverage result and move to shared logging
    • if faults violate promised storage semantics, classify them as assumption failures
      • do not call them consensus bugs
  • novelty still to check
    • current implementation issues and fault tests
    • Grove and other verified storage interfaces
    • filesystem crash-consistency tools
    • XLL’s implementation of cleanup relocation and durable index updates

project B: quorum optimization that survives leader replacement

  • question
    • can normal-operation gains remain worthwhile once correlated failures and recovery time are included?
  • proposed contribution
    • choose quorum sets against an explicit service interruption budget
  • first experiment
    • use Flexible Paxos and a majority baseline
    • reproduce Quoracle’s static optimization before adding adaptation
    • emulate several regions with heterogeneous disk and network delays
    • fail the fastest region after a load-driven quorum choice
  • measure
    • normal-operation median and 99th-percentile request latency
    • time until the first successful operation after leader loss
    • fraction of requests finishing within a stated deadline
    • network bytes and equal total machine cost
  • hard requirement
    • every change preserves required intersections across rounds and configurations
    • use a proved transition rule rather than replacing quorum sets in memory
  • stop or redirect
    • if a static quorum choice matches adaptation, keep the simpler result
    • if gains require unjustified failure predictions, report the dependence explicitly
  • novelty still to check
    • weighted quorum adaptation and geography-aware Paxos
    • Matchmaker’s configuration selection policies
    • failure-aware quorum optimizers after Quoracle

project C: recovery cost of separating protocol roles

  • question
    • does independent scaling remain useful when recovery consumes the same network and disks as new requests?
  • proposed experiment
    • compare ordinary MultiPaxos with role-separated replication at equal resources
    • replay a realistic operation mix with large values and snapshots
    • fail a sequencer, broadcaster and execution replica separately
    • replace members while transferring snapshots
  • measure
    • request latency throughout recovery
    • time for replacement to become eligible to vote or execute
    • backlog size, network bytes and retained storage
  • research hypothesis
    • explicit recovery bandwidth limits can reduce client tail latency
    • excessively low limits can lengthen exposure to another failure
  • novelty still to check
    • existing catch-up throttling in databases
    • Compartmentalized Paxos recovery experiments
    • current disaggregated consensus and storage protocols

what would count as an advance

  • a demonstrated failure within a documented failure model
  • a smaller interface with an explicit proof or exhaustive bounded model check
  • a recovery policy with measured gains at equal cost
  • a reproducible negative result showing when a proposed optimization loses
  • unknown today
    • whether any project is new enough for a research paper
    • whether reviewed protocols reproduce their reported performance locally
    • whether the human would prefer implementation debugging or protocol design
  • recommendation under that uncertainty
    • run project A’s smallest experiment first
    • inspect recent artifacts and complete remaining full-paper reading before designing a new protocol

Last edited: