Skip to main content

Module consensus_based_total_order_broadcast

Module consensus_based_total_order_broadcast 

Source
Expand description

Consensus-based total-order broadcast.

Status: transcription. Space: unbounded — unordered, delivered and the family of consensus instances all grow with the number of entries handled. That is the page, and docs/bounded-space.md is explicit that inheriting the book’s omissions is correct of a transcription and disqualifying of an implementation. Bounding any of them weakens a guarantee to a scope and belongs to a change with a proposal.

Cachin, Guerraoui & Rodrigues, Module 6.1 (TotalOrderBroadcast) and Algorithm 6.1, quoted from the book:

Algorithm 6.1: Consensus-Based Total-Order Broadcast
Implements: TotalOrderBroadcast, instance tob.
Uses:
    ReliableBroadcast, instance rb;
    Consensus (multiple instances).

upon event ⟨ tob, Init ⟩ do
    unordered := ∅;
    delivered := ∅;
    round := 1;
    wait := FALSE;

upon event ⟨ tob, Broadcast | m ⟩ do
    trigger ⟨ rb, Broadcast | m ⟩;

upon event ⟨ rb, Deliver | p, m ⟩ do
    if m ∉ delivered then
        unordered := unordered ∪ {(p, m)};

upon unordered ≠ ∅ ∧ wait = FALSE do
    wait := TRUE;
    Initialize a new instance c.round of consensus;
    trigger ⟨ c.round, Propose | unordered ⟩;

upon event ⟨ c.r, Decide | decided ⟩ such that r = round do
    forall (s, m) ∈ sort(decided) do     // by the order in the resulting sorted list
        trigger ⟨ tob, Deliver | s, m ⟩;
    delivered := delivered ∪ decided;
    unordered := unordered \ decided;
    round := round + 1;
    wait := FALSE;

The shape is one round at a time: everything reliable broadcast has delivered and this process has not yet ordered is proposed as a set, consensus agrees on a set, and every process turns that same set into the same sequence by sorting it. Ordering is therefore agreed without anyone communicating about order at all — the sort does that work, which is why it must be deterministic.

§Departures from the page

  • A read. The port this satisfies offers crate::total_order_log::TotalOrderLog::read, which the book’s abstraction does not: its clients observe deliveries. The algorithm already maintains delivered, so the read exposes what the page keeps and does not offer, served locally. See the port’s own documentation.

  • Consensus instances are held explicitly, keyed by round. The page writes c.round and ⟨ c.r, Decide ⟩, so instances are a family addressed by round; the book’s runtime routes to them and this one does not. They are created on demand — including for a round this process has not reached, which is what lets a peer that is ahead make progress — and never pruned, as the page has them. Creation runs the instance’s ⟨ Init ⟩ before the event that provoked it: “Initialize a new instance c.round” is an event the book’s runtime delivers, and skipping it leaves the instance’s failure detector without its timers — which no fault-free run notices, because deciding under the initial epoch never consults the detector. A crash is then never detected and the survivors stall, which is what the suite’s crash property caught.

  • The conditional event handler is discharged here. such that r = round is not a guard that discards. The book states its meaning: “An algorithm that uses conditional event handlers relies on the run-time system to buffer external events until the condition on internal variables becomes satisfied.” Cx has no such facility, so a decision for a round this process has not reached is held in decisions and acted on when round catches up. crate::leader_driven_consensus’s pending is the same pattern, for the same reason.

  • A consensus message carries its round. The page addresses instances; nothing on this wire does. Unlike crate::epoch_consensus, whose instance stamps its own messages because the epoch is its identity, the round is this layer’s concept and not the consensus’s — so this layer stamps, and the stamp is the one thing it cannot delegate.

  • unordered is deduplicated on the pair (p, m), not on m alone. The page’s if m ∉ delivered reads against a delivered that holds pairs, so one of the two is loose; the pair is what makes unordered \ decided well defined, and it is what is used here. The consequence, which the page shares: one process appending the same value twice contributes one entry.

  • The consensus beneath is not a type parameter, and its link is fixed. Every other composing layer here takes its child as a parameter with a default. This one cannot for the consensus: instances are created at run time, one per round, so the layer would need a link factory rather than a link — a runtime indirection the static composition model exists to avoid. The reliable broadcast keeps its parameter, because there is exactly one of it and the caller supplies it once. Revisit if a stack ever wants rounds agreed over something else.

  • The sort is a BTreeSet’s own order. Proposing an ordered set means sort(decided) is iteration, and every process computes the same sequence because Ord is a function of the values rather than of anything local. The guard forbidding hash-keyed maps in these three crates exists for the converse reason: an iteration order that varies per process is exactly what would break agreement here.

Structs§

ConsensusBasedTotalOrderBroadcast
A totally ordered log, agreed by one consensus instance per round.

Enums§

Cmd
Requests from the layer above.
Ind
Indications to the layer above.
Wire
This layer’s messages: the broadcast’s, and a consensus instance’s stamped with its round.

Type Aliases§

Batch
What a round proposes and decides: the set of entries not yet ordered.
Carried
What reliable broadcast carries for this layer.
Consensus
The consensus one round runs. Not a type parameter — see the module’s departures.
ConsensusCarried
What the consensus beneath carries for this layer.
Slot
One entry with the process that appended it — the page’s (s, m).