Skip to main content

Module uniform_reliable_broadcast

Module uniform_reliable_broadcast 

Source
Expand description

Uniform reliable broadcast.

Cachin, Guerraoui & Rodrigues, Module 3.3 and Algorithm 3.4 (“All-Ack Uniform Reliable Broadcast”).

Status: transcription. Space: unbounded. pending holds payloads, ack holds a process set per message, and delivered grows without limit. Deployable once collected — and this layer already computes the predicate it would need, since correct ⊆ ack[m] is a stability test: a message every correct process has seen can be dropped from pending and ack the moment it is delivered. See docs/bounded-space.md.

Reliable broadcast guarantees agreement only among correct processes. A process that delivers a message and then crashes may leave the survivors never delivering it — and if that delivery had any external effect, the divergence cannot be repaired from above. Uniform agreement quantifies over any process: delivered by anyone at all, correct or not, means eventually delivered by everyone correct.

It is bought by waiting. A message is delivered only once every process still believed correct has been seen to acknowledge it, so nobody can deliver something the others have not yet seen.

upon event ⟨ urb, Broadcast | m ⟩ do
    pending := pending ∪ {(self, m)};
    trigger ⟨ beb, Broadcast | [DATA, self, m] ⟩;

upon event ⟨ beb, Deliver | p, [DATA, s, m] ⟩ do
    ack[m] := ack[m] ∪ {p};
    if (s, m) ∉ pending then
        pending := pending ∪ {(s, m)};
        trigger ⟨ beb, Broadcast | [DATA, s, m] ⟩;

upon event ⟨ P, Crash | p ⟩ do
    correct := correct \ {p};

function candeliver(m) is  return correct ⊆ ack[m];

upon exists (s, m) ∈ pending such that candeliver(m) ∧ m ∉ delivered do
    delivered := delivered ∪ {m};
    trigger ⟨ urb, Deliver | s, m ⟩;

§This layer depends on a timing assumption, and cannot detect its failure

Uniform agreement holds only while the failure detector is accurate, which holds only while the network delivers within a known bound. A wrongly accused process is removed from correct, candeliver is satisfied too early, and a message can be delivered by some processes and not others.

That dependency is stated here rather than expressed with the scope annotation of docs/scope-annotated-modules.md, and deliberately: a scope must have a boundary the module can observe. This one has none. Synchrony failing arrives here as the detector reporting a crash, indistinguishable from the detector being right. An assumption a layer rests on but cannot detect is not a scope — tagging it would create an obligation no implementation could discharge and no test could exercise.

§Departures from the page

  • ack and delivered are keyed by an identifier carrying the originator and a per-sender sequence number, not by message content. The book’s ack[m] assumes messages are unique across senders; identical content broadcast twice must be delivered twice.
  • Two children both send, so this layer’s wire type is an enum distinguishing a broadcast payload from a heartbeat. It is the first multiplexing in the stack, and it is typed: a mis-wiring is a compile error rather than a silently undelivered message.
  • ⟨urb, Init⟩ is a separate event: new establishes the state, and [Protocol::on_init] starts the detector beneath. It was a Cmd::Start before the trait had an init event, which is why the only request now is Cmd::Broadcast.
  • Neither ack nor pending is garbage collected, as in the book. Long runs grow.

L is a parameter, so this one module is also what session_uniform_reliable_broadcast used to be. Algorithm 3.4 gains one clause and nothing else changes:

upon event ⟨ SessionEstablished | q ⟩ do
    forall (s, m) ∈ pending do
        trigger ⟨ beb, SendTo | q, [DATA, s, m] ⟩;

It reuses pending, which the algorithm already maintains, and SendTo, which is a narrowing of an existing action rather than a new communication step.

Two things about that clause are deliberate and neither is what one would write first. It is unconditional, not filtered by q ∉ ack[m]: the filtered version deadlocks, for the reason set out at resend_to, which is where a test found it. And it is directed at the peer whose scope returned rather than broadcast to everyone, because the ending was per peer and so is the repair. Nothing is attempted on the ending itself: the peer is unreachable at that moment and anything sent would be discarded.

§This layer keeps the perfect detector, and that is not an oversight

crate::eventually_perfect_failure_detector exists, and Ω moved onto it. This layer must not. Its agreement rests on strong accuracy by name — the book’s own proof invokes it — and one false suspicion breaks it, permanently and silently. A detector allowed to be wrong would turn a correct module into a broken one.

That is not a gap to be closed later: it is the whole subject of crate::majority_ack_uniform_reliable_broadcast, which asks has more than half relayed it? instead of has everyone still believed correct relayed it?, and drops the detector entirely. The pair is the demonstration of what strong accuracy is worth, and replacing the detector here would delete the demonstration rather than improve it.

URB1 [always]       Validity — conditional on the two mechanisms below
URB2 [incarnation]  No duplication — `delivered` is volatile, so a restart forgets it
URB3 [always]       No creation
URB4 [always]       Uniform agreement — conditional on the two mechanisms below

URB2 is [incarnation] by docs/scope-annotated-modules.md Corollary 7.2: the set that would have to survive is delivered, it is held in memory, and the boundary it cannot cross is this process’s own ⟨Init⟩.

URB1 and URB4 are [always] only because between two mechanisms no third outcome is left — the scope comes back and the resend repairs it, or the peer never returns and the detector’s timeout drops it from correct. Each carries a condition that is an assumption rather than a property of this code. The reconnection path needs the peer to be reachable again; the accusation path needs the detector’s synchrony assumption, and perfect_failure_detector is explicit that outside a synchronous system it accuses correct processes. Both failing at once is a permanent split: each side of a partition accuses the other, each has correct ⊆ ack[m] satisfied among itself, and both deliver — not uniform agreement failing on a technicality but two disjoint sets of processes proceeding as though the other did not exist. crate::majority_ack_uniform_reliable_broadcast cannot suffer that and blocks instead; that difference is what a quorum buys and a detector costs.

Read against reliable broadcast, which has neither mechanism and whose agreement is therefore scoped, this is the clearest statement of what a failure detector buys.

Structs§

BroadcastId
Names one broadcast uniquely: who originated it, and their sequence number for it.
Data
What this layer adds to a broadcast payload.
UniformReliableBroadcast
Broadcast with uniform agreement, over best-effort broadcast and a failure detector.

Enums§

Cmd
Requests from the layer above.
Ind
Indications to the layer above.
Wire
The wire type, multiplexing the two children.

Type Aliases§

BebMsg
What best-effort broadcast puts on the wire for this layer’s payloads.
Carried
What a link beneath this layer must carry.