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: