Most replicated systems assume that a failed server stops. Raft, Multi-Paxos and ZooKeeper's Zab tolerate machines that crash, lose power or fall off the network, and they need 2f+1 replicas to survive f such failures. They do not tolerate a server that keeps running and lies: one that tells half the cluster it voted yes and the other half it voted no, forges a log entry, or reports a value it never stored. Byzantine fault tolerance (BFT) is the family of protocols that survive that behaviour, and it costs more replicas, more messages and more cryptography.

This article covers the fault model, the 3f+1 bound, the timing assumptions, PBFT step by step, HotStuff and Tendermint, a replica sketch and a trace of a lying leader, and then the practical question: when does BFT earn its cost, and when is a crash-fault protocol plus good engineering the better choice? For the crash-fault baseline, see Paxos and Raft consensus.

Advertisement

The Byzantine fault model

The name comes from Lamport, Shostak and Pease's 1982 Byzantine Generals paper. A Byzantine node may deviate from the protocol in any way: send conflicting messages to different peers (equivocation), stay silent, delay messages selectively, or collude with other faulty nodes. The cause does not matter: memory corruption, a misconfigured replica and an attacker holding a server's keys look the same.

Two limits make the model tractable: at most f of n nodes are faulty, and messages are authenticated with MACs or signatures, so a liar can lie in its own name but cannot impersonate an honest peer.

Why n must be at least 3f + 1

The core argument is about quorum intersection, the same property Quorum Systems Explained develops for crash faults. A protocol must make progress after hearing from n - f nodes, because f nodes may never answer. So a quorum has n - f members. Two quorums of that size overlap in at least n - 2f nodes. With crash faults any overlapping node is honest, so an overlap of one is enough, giving n >= 2f + 1. With Byzantine faults, up to f nodes in the overlap may be liars who tell each quorum what it wants to hear. For the overlap to contain at least one honest node that carries information between the quorums, it needs at least f + 1 members: n - 2f >= f + 1, which gives n >= 3f + 1.

With n = 3f + 1, a quorum is 2f + 1 nodes. The numbers grow quickly:

Faults tolerated (f)Crash-fault replicas (2f+1)BFT replicas (3f+1)BFT quorum (2f+1)
1343
2575
37107
10213121

Five or six replicas still tolerate only one Byzantine fault, so BFT clusters are sized 4, 7, 10 and so on.

Advertisement

Timing: safety always, liveness eventually

Fischer, Lynch and Paterson proved in 1985 that no deterministic protocol can guarantee consensus in a fully asynchronous network if even one node may crash. Practical BFT protocols sidestep this with the partial synchrony model of Dwork, Lynch and Stockmeyer: the network may behave arbitrarily for a while, but after some unknown global stabilisation time, messages arrive within a bound. The protocols are designed so that safety (no two correct nodes decide differently) holds regardless of timing, and liveness (requests eventually complete) holds only once the network behaves.

The practical consequence is timeouts. If a leader stops making progress, replicas time out and replace it. Too short and a slow network causes constant leader changes; too long and a faulty leader stalls the system. Implementations typically grow the timeout after each failed view change.

PBFT, step by step

Castro and Liskov's Practical Byzantine Fault Tolerance (OSDI 1999) was the first BFT state machine replication protocol fast enough for real services. Replicas move through numbered views; in view v, replica v mod n is the primary. The normal case has three phases:

  1. Request. A client sends a signed operation to the primary.
  2. Pre-prepare. The primary assigns a sequence number s and sends a pre-prepare message with the view, sequence number and request digest to every backup. This proposes an order.
  3. Prepare. Each backup that accepts the pre-prepare sends a prepare to all replicas. A replica is prepared for the request once it holds the pre-prepare plus 2f matching prepares from different replicas. Quorum intersection means no two different requests can be prepared at the same slot in one view.
  4. Commit. A prepared replica sends a commit to all. Once it holds 2f + 1 matching commits (its own included), the request is committed and survives a view change. The replica executes requests in sequence order.
  5. Reply. Each replica replies to the client, which accepts the result once f + 1 replicas return the same answer. At least one of those is honest, so the answer is correct.
PBFT normal case with n = 4, f = 1: request, pre-prepare, prepare, commit, replyClientPrimary R0Replica R1Replica R2Replica R3requestpre-prepare(primary assigns seq s)prepare (all-to-all)need 2f matchingcommit (all-to-all)need 2f+1 matchingreplyclient waits for f+1Two all-to-all rounds make the normal case quadratic in n. HotStuff replaces them with votessent to the leader, which aggregates them into a quorum certificate: linear per phase.
PBFT's normal case for four replicas tolerating one fault. The prepare and commit rounds are all-to-all, so message count grows with the square of n.

Two background mechanisms complete the protocol. Checkpoints: periodically replicas exchange state digests, and 2f + 1 matching digests form a stable checkpoint that lets them truncate the log. View change: if a backup's timer expires before a request commits, it stops accepting messages in the current view and broadcasts a view-change message with its latest stable checkpoint and its prepared certificates. The new primary collects 2f + 1 of these, re-proposes every request that was prepared in any of them, and starts the new view. That re-proposal keeps every committed request at its sequence number.

A replica, in code

The sketch below shows the normal-case logic of one replica. It omits checkpoints, view change and request batching, and it assumes an authenticated transport that drops messages with bad signatures.

from collections import defaultdict

class Replica:
    def __init__(self, rid, n, f, send_all, execute):
        self.rid, self.n, self.f = rid, n, f
        self.view = 0
        self.accepted = {}                  # (view, seq) -> digest from pre-prepare
        self.prepares = defaultdict(set)    # (view, seq, digest) -> sender ids
        self.commits = defaultdict(set)
        self.committed = {}                 # seq -> request
        self.requests = {}                  # digest -> request
        self.next_exec = 1
        self.send_all, self.execute = send_all, execute

    def primary(self):
        return self.view % self.n

    def on_pre_prepare(self, sender, view, seq, digest, request):
        if sender != self.primary() or view != self.view:
            return
        if self.accepted.get((view, seq), digest) != digest:
            return                          # primary equivocated: evidence for a view change
        if digest != hash_request(request):
            return
        self.accepted[(view, seq)] = digest
        self.requests[digest] = request
        self.send_all(("PREPARE", view, seq, digest, self.rid))

    def on_prepare(self, sender, view, seq, digest):
        if sender == self.primary():
            return                          # only backups send prepares
        key = (view, seq, digest)
        self.prepares[key].add(sender)
        if self.accepted.get((view, seq)) == digest and len(self.prepares[key]) >= 2 * self.f:
            if self.rid not in self.commits[key]:
                self.commits[key].add(self.rid)
                self.send_all(("COMMIT", view, seq, digest, self.rid))

    def on_commit(self, sender, view, seq, digest):
        key = (view, seq, digest)
        self.commits[key].add(sender)
        if len(self.commits[key]) >= 2 * self.f + 1 and seq not in self.committed:
            self.committed[seq] = self.requests[digest]
            while self.next_exec in self.committed:   # execute strictly in order
                self.execute(self.committed[self.next_exec])
                self.next_exec += 1

Note the ordering rule at the bottom: a replica may commit sequence 7 before 6 but must not execute it first, or honest replicas diverge. Real implementations also batch many requests per sequence number.

HotStuff and Tendermint

PBFT's all-to-all rounds cost O(n squared) messages per request, and its view change costs more. That is fine at four or seven replicas and painful at a hundred.

HotStuff (Yin, Malkhi, Reiter, Gueta and Abraham, PODC 2019) routes all votes through the leader. The leader aggregates 2f + 1 signed votes into a quorum certificate, a compact proof it can send to everyone, so each phase costs O(n) messages, including leader replacement. It adds a third voting phase compared with PBFT, and in exchange a new leader can proceed as soon as it hears from a quorum rather than waiting out a worst-case timeout, a property the authors call optimistic responsiveness. Chained HotStuff pipelines the phases, and derivatives such as DiemBFT were built on it.

Tendermint, now maintained as CometBFT, takes a rotating proposer and two voting rounds (prevote and precommit) with gossip-based message dissemination. Its liveness mechanism waits for a timeout at round changes, so it is not optimistically responsive, but the design is simple, and it has run public proof-of-stake networks for years.

All keep the 3f + 1 bound; they differ in communication cost, leader-replacement cost and latency.

Worked example: a lying primary

Take n = 4, f = 1, with primary R0 faulty. R0 wants R1 to order request A at sequence 5 and R2 and R3 to order request B at sequence 5, hoping to split the cluster.

  1. R0 sends pre-prepare(5, A) to R1 and pre-prepare(5, B) to R2 and R3.
  2. R1 broadcasts prepare(5, A). R2 and R3 broadcast prepare(5, B). As primary, R0 sends no prepares, and correct replicas ignore any it sends.
  3. R2 holds its pre-prepare for B plus matching prepares from itself and R3: 2f = 2, so R2 is prepared for B. R3 likewise. R1 holds only its own prepare for A, one short, so it cannot prepare A.
  4. R2 and R3 send commit(5, B). If R0 adds a commit, R2 and R3 each see 2f + 1 = 3 and execute B. If R0 withholds it, they stall at two commits. Either way R1 never commits A.
  5. Timers expire at replicas whose requests do not complete. R1, R2 and R3 send view-change messages, R1 becomes primary of view 1, and because R2 and R3 report a prepared certificate for B, the new primary re-proposes B at sequence 5. Every correct replica ends with B at position 5 and A is ordered later.

The faulty primary slowed the system but could not make two honest replicas execute different requests at the same position: safety against f liars, liveness once a correct leader runs on a timely network.

When BFT earns its cost

BFT pays for itself when replicas are run by parties who do not trust each other, or an adversary may control some of them. Blockchains and consortium ledgers are the clearest cases: validators belong to different organisations, and compromising one must not rewrite history. Safety-critical embedded systems use Byzantine-tolerant voting for a related reason: a corrupted sensor or processor can emit plausible wrong values rather than stopping.

It rarely pays inside one organisation's data centre. There the realistic faults are crashes, partitions, slow disks and gray failures, and a crash-fault protocol plus defences for the common forms of corruption covers them at lower cost: checksums on disk and network data, end-to-end verification, fencing tokens against stale leaders, and careful deployment. BFT's costs are concrete: an extra replica per fault, cryptographic checks on every message, heavier communication and more complex code.

There is also an independence problem. BFT tolerates f faulty replicas, but a bug in the shared software is not one fault: every replica running the same binary fails the same way at once. BFT protects against a bug only if replicas are diverse (different implementations, different teams) or if the fault is local, such as one corrupted disk or one compromised key.

Operating a BFT system

  • Key management is the security boundary. Anyone holding more than f validator keys can break safety. Keep keys in hardware security modules or remote signers, and plan rotations.
  • Place replicas for independence. Different clouds, regions, administrators and ideally implementations, so one incident cannot take out f + 1 nodes.
  • Tune timeouts to real latency. Measure round-trip times between replicas at the 99th percentile and set the base view-change timeout above it, with exponential backoff.
  • Watch for equivocation evidence. Conflicting signed messages from one replica are proof of misbehaviour. Log and alert on them; many systems penalise or eject the offender.
  • Monitor view changes. Frequent ones mean a faulty leader or a mis-tuned timeout.
  • Change membership deliberately. Membership changes alter f and the quorum size; use the protocol's reconfiguration path, never per-node config edits.

Failure modes

  • More than f faults. Safety is gone. With f + 1 colluding replicas in a 3f + 1 cluster, conflicting decisions become possible.
  • Correlated software bugs. One bug on all replicas is not tolerated; it behaves like n faults.
  • Timeout storms. Timeouts shorter than real latency during congestion cause repeated view changes and no progress, even with all replicas honest.
  • Client-side trust gaps. A client that accepts one replica's answer instead of f + 1 matching answers throws away the guarantee.
  • Non-deterministic execution. If executing a request depends on local time, randomness or map iteration order, honest replicas diverge after agreeing on the order.

Trade-offs

BFT buys resilience against lying nodes with 3f + 1 replicas instead of 2f + 1, cryptography on every message, higher latency and more code. PBFT gives low latency at small replica counts; HotStuff trades an extra phase for linear communication and cheap leader changes; Tendermint trades responsiveness for simplicity. A crash-fault protocol with strong integrity checks is cheaper whenever one trusted operator runs every replica. Choose BFT for the trust model you actually have.

What to do next

  1. Write down who operates each replica and who could compromise it. If one organisation controls every node, start from Raft or Paxos plus checksums.
  2. If you need BFT, size the cluster as 3f + 1 for the f you must tolerate and place replicas in independent failure and administrative domains.
  3. Prefer a mature implementation (for example CometBFT or a HotStuff-family library) over writing your own protocol.
  4. Audit request execution for determinism: no wall-clock time, randomness or unordered iteration in the state machine.
  5. Make clients wait for f + 1 matching replies or verify a quorum certificate.
  6. Set up alerts for view changes and equivocation evidence, and rehearse a validator key rotation before you need one.
Key takeaway: Byzantine fault tolerance lets replicas agree even when up to f of them lie, as long as n is at least 3f + 1 so that any two quorums of 2f + 1 share an honest node. PBFT achieves this with pre-prepare, prepare and commit phases plus view changes; HotStuff makes each phase linear with quorum certificates; Tendermint keeps the design simple at some cost in responsiveness. Safety holds under any timing, and liveness depends on timeouts once the network is timely. Use BFT when replicas are run by parties that do not trust each other; inside one trusted organisation, a crash-fault protocol with integrity checks is usually the better trade.