Skip to main content

Module multi_paxos_replica

Module multi_paxos_replica 

Source
Expand description

The Multi-Paxos replica: slots become positions in a log.

Status: implementation. Space: the bookkeeping is bounded — decisions and performed by a retention window, requests and proposals by WINDOW — and the ordered sequence is deliberately not, because it is the data rather than the bookkeeping.

§4.2 is applied, and it is worth being clear which half. The section separates state that is unnecessary — the leader’s and acceptor’s, collected below a watermark once f + 1 replicas have applied past it, which is crate::multi_paxos_synod’s — from state that is unavoidable, which is this module’s:

because commands may be decided in multiple slots, replicas each maintain a set of all decisions to filter out such duplicates. In practice, it is often sufficient if such information is only kept for a certain amount of time, making the probability of duplicate execution negligible.

So the filter is bounded by a retention window rather than by the watermark beneath it, and the reason is worth stating because using the watermark is the obvious wrong move: that watermark says f + 1 replicas have applied up to a slot, which says nothing whatever about whether a command decided below it may be decided again above it. A duplicate filter has to outlive every slot a duplicate could span, and no watermark bounds that.

What that weakens, stated rather than discovered. “A command decided in two slots takes one position” holds within the retention window. A duplicate arriving more than RETAIN slots after the first would take a second position. As things stand that is a weakening of something unreachable — a command is minted at one replica and occupies one slot at a time, so a duplicate needs the client-retry deployment this port does not have — but it is the honest bound and it is in the specification.

The ordered sequence is exempt and is not collected. A log grows with what is appended to it; that is the data, and a log that discarded its entries would not be one. What would bound it is a snapshot, which is outside the paper and is also the point past which a lagging replica can no longer be caught up — see the catch-up below.

§Source

van Renesse, R. and Altinbuken, D. (2015) ‘Paxos Made Moderately Complex’, ACM Computing Surveys, 47(3), pp. 1–36 — §2.1 and Figure 1, “Pseudocode for a replica”, page 42:6.

Not van Renesse, R. (2011), the Cornell technical report of the same title. The two differ here more than anywhere else in the algorithm: the report’s replica keeps a single slot num and its propose(p) searches for the lowest slot not already proposed or decided, with no window at all. The survey splits that into slot_in and slot_out and adds WINDOW, which is the version transcribed below. A reader checking this code against a page needs to know which page.

process Replica(leaders, initial_state)
  var state := initial_state, slot_in := 1, slot_out := 1;
  var requests := ∅, proposals := ∅, decisions := ∅;

  function propose()
    while slot_in < slot_out + WINDOW ∧ ∃c : c ∈ requests do
      if ∃op : ⟨slot_in − WINDOW, ⟨·, ·, op⟩⟩ ∈ decisions ∧ isreconfig(op) then
        leaders := op.leaders;
      end if
      if ∄c' : ⟨slot_in, c'⟩ ∈ decisions then
        requests := requests \ {c};
        proposals := proposals ∪ {⟨slot_in, c⟩};
        ∀λ ∈ leaders : send(λ, ⟨propose, slot_in, c⟩);
      end if
      slot_in := slot_in + 1;
    end while
  end function

  function perform(⟨κ, cid, op⟩)
    if (∃s : s < slot_out ∧ ⟨s, ⟨κ, cid, op⟩⟩ ∈ decisions) ∨ isreconfig(op) then
      slot_out := slot_out + 1;
    else
      ⟨next, result⟩ := op(state);
      atomic
        state := next; slot_out := slot_out + 1;
      end atomic
      send(κ, ⟨response, cid, result⟩);
    end if
  end function

  for ever
    switch receive()
      case ⟨request, c⟩ :
        requests := requests ∪ {c};
      end case
      case ⟨decision, s, c⟩ :
        decisions := decisions ∪ {⟨s, c⟩};
        while ∃c' : ⟨slot_out, c'⟩ ∈ decisions do
          if ∃c'' : ⟨slot_out, c''⟩ ∈ proposals then
            proposals := proposals \ {⟨slot_out, c''⟩};
            if c'' ≠ c' then
              requests := requests ∪ {c''};
            end if
          end if
          perform(c');
        end while
      end case
    end switch
    propose();
  end for
end process

§The invariants, quoted

  • R1: “There are no two different commands decided for the same slot: ∀s, ρ1, ρ2, c1, c2 : ⟨s, c1⟩ ∈ ρ1.decisions ∧ ⟨s, c2⟩ ∈ ρ2.decisions ⇒ c1 = c2.” Held by the child, not here: it is the Synod protocol’s S1, and this layer relies on it rather than enforcing it. What this layer does is not break it — decisions never replaces an entry.

  • R2: “All commands up to slot out are in the set of decisions: ∀ρ, s : 1 ≤ s < ρ.slot out ⇒ ∃c : ⟨s, c⟩ ∈ ρ.decisions.” Held because slot_out advances only inside perform, which the drain loop calls only for a slot that has a decision.

  • R3: “For all replicas ρ, ρ.state is the result of applying the commands ⟨s, cs⟩ ∈ ρ.decisions to initial state for all s up to slot out, in order of slot number.” Here state is the ordered sequence — see the departures — so R3 is the statement that the sequence is exactly the decided commands below slot_out, in slot order, minus the ones perform skipped as already applied.

  • R4: “For each ρ, the variable ρ.slot out cannot decrease over time.” Structural: the only assignment is += 1.

  • R5: “A replica proposes commands only for slots for which it knows the configuration: ∀ρ : ρ.slot in < ρ.slot out + WINDOW.” The while guard in propose, kept although the reason for it is not — see below.

    The formula and the sentence beside it do not say quite the same thing, and the code follows the sentence. propose’s loop tests slot_in < slot_out + WINDOW at the top and increments slot_in at the bottom, so an exit with requests still queued leaves slot_in = slot_out + WINDOW exactly — the strict inequality does not hold of the variable between the loop ending and slot_out next advancing. What does hold, always, is the English: every slot ever proposed for is below slot_out + WINDOW at the moment of proposing, since the guard is checked before the slot is used. The suite asserts that form over MultiPaxosReplica::proposed_slots and the loop’s own ≤ over the variable, rather than asserting a formula the page’s own pseudocode breaks.

§Where each part of Figure 1 went

Figure 1Here
var state := initial_statethe ordered sequence; there is no application state machine, because the port is a log
slot_in, slot_outfields, counting slots
requests, proposals, decisionsfields; decisions is append-only, as the page has it
function propose()transfer plus the send loop in pump
the isreconfig branch in propose()absent — reconfiguration is a later change; the WINDOW guard around it is kept
function perform()perform, minus op(state) and the client response
case ⟨request, c⟩Cmd::Append, the port’s own
case ⟨decision, s, c⟩the child’s crate::multi_paxos_synod::Ind::Decision
∀λ ∈ leaders : send(λ, ⟨propose, …⟩)a call into the child, per §4.4
send(κ, ⟨response, cid, result⟩)absent, with state

§Departures from the page

  • state is a sequence, and there is no op(state). The port this satisfies is crate::total_order_log::TotalOrderLog: a log, not a replicated state machine. perform therefore appends the command to the sequence where the page applies it, and the client response goes with the result it would have carried. What survives is the part R1–R4 are about — that every replica applies the same commands in the same order — and the part that goes is the application on top of it.

  • ∀λ ∈ leaders : send(λ, ⟨propose, s, c⟩) is a call into the child. §4.4: “each machine that runs a replica also runs a leader… the replica can send a proposal for a particular slot to its local leader”. So leaders is not a set here, it is one child, and the fan-out to remote leaders is the child’s forwarding rather than this layer’s broadcast. That is what lets this module hold no link, no broadcast and no wire of its own; the cost is in crate::multi_paxos_synod’s own documentation, and it is that proposal delivery now rests on the leader detector.

  • The already-decided check is indexed rather than scanned. The page writes ∃s : s < slot_out ∧ ⟨s, ⟨κ, cid, op⟩⟩ ∈ decisions, a scan of every decision below slot_out. performed is exactly that set — perform runs once per slot in slot order, so what it has seen is {decisions[s] : s < slot_out} — and membership answers the same question in log time. It grows the same unbounded way decisions does, so it changes nothing about the space statement above.

  • requests is a queue, where the page says “any command”. ∃c : c ∈ requests picks an arbitrary member. FIFO is a refinement rather than a departure in the strict sense, but it is worth naming because it is what stops a command being passed over for ever while later ones are proposed — the page leaves that to whoever implements the choice.

  • WINDOW is kept and its reason is deferred. In the source the window exists because a reconfiguration decided in slot s takes effect at s + WINDOW, so a replica may not propose past the last slot whose configuration it knows. Reconfiguration is not built here and the membership is fixed for the run, so the guard is a pipeline cap and nothing more. It is kept rather than dropped so that R5 reads against the page, and this note is here so that a reader does not take the cap for the whole reason.

  • §4.2’s periodic report, and the catch-up that pays for collecting. A replica tells the consensus beneath it how far it has applied, periodically and not as a consequence of doing work — a replica applying nothing is the one whose position others most need, and a report riding its own traffic would fall silent exactly then.

    Collecting at f + 1 means f correct replicas may be behind, and what would have helped them is what was collected. §4.2’s own remedy is that “replicas can learn decisions […] from one another”, so a replica whose proposal is refused as collected asks a peer, and the peer answers from its decisions — the one place that still holds them. A replica further behind than RETAIN cannot be caught up, because nobody holds those decisions any more.

  • The fourth liveness violation, and its fix. Liu, Y.A., Chand, S. and Stoller, S.D. (2019) ‘Moderately Complex Paxos Made Simple’, PPDP ’19, is the cross-check this module’s source is read against. Its fourth liveness violation is this layer’s: if no decision arrives for a slot, slot_out stops moving, WINDOW fills, slot_in stops advancing and the replica wedges — with nothing on the page to get it out, because Figure 1 acts only on messages that arrive. The fix has two halves, and the other one is the leader’s: a proposal outstanding longer than repropose_after is proposed again for the same slot, and a leader that has seen the slot decided answers with the decision rather than dropping the repeat. Re-proposing into a new slot instead would leave the old one unfilled for ever, which is the wedge itself.

    The replica does not try to tell why a decision has not come. The proposal may have been forwarded to a process that has crashed, the decision may have been lost at a session ending, or consensus for the slot may simply still be running. Every case is answered by the same message and none is made unsafe by asking: where the slot is open the leader’s ∄c' guard drops the repeat, where the proposal was lost the repeat is the first the leader hears of it, and where the slot is decided the leader answers. So the threshold can be generous and wrong without being unsafe. The one thing it must not be is absent.

§A position is not a slot

[Position] and Slot are both u64 and are not interchangeable. Slots count from 1; Position::START is 0. More importantly perform skips a command already applied at a lower slot, so a slot can pass without a position being taken — the two diverge by exactly the number of commands decided more than once. Position(slot) is right until the first duplicate and then silently wrong, which is why a_command_decided_in_two_slots_takes_one_position drives that case deliberately rather than waiting for a run to produce it.

§Scope

This layer bridges nothing. A session ending reaches it from the child and is propagated in its own indications: its redundancy is the child’s, and the child’s is the other processes rather than anything that outlives a session. See docs/conditional-guarantees.md.

Crash-stop, and one step stronger than the child’s. MultiPaxosSynod keeps nothing durably, so a returning process has forgotten which ballots it took up; this layer adds that a returning process has forgotten its sequence, and would answer a read with a shorter one than it had already served. A total order that shortens is not a total order, so a crashed process here is crashed for good. §4.3 of the source is what changes it, and it is not this change.

Structs§

Command
The source’s c = ⟨κ, cid, op⟩: who asked, which request of theirs it is, and what it says.
MultiPaxosReplica
A totally ordered log, built on Multi-Paxos: one consensus per slot, under a stable leader that keeps phase one across all of them.

Enums§

Cmd
Requests from the layer above.
Ind
Indications to the layer above.

Constants§

RETAIN
How many decided slots a replica keeps for the duplicate filter — §4.2’s retention window.
WINDOW
How far slot_in may run ahead of slot_out — the source’s WINDOW.

Type Aliases§

Carried
What the Synod protocol beneath carries for this layer.
Synod
The consensus this layer runs over. A type alias rather than a parameter: the replica is the half of Multi-Paxos that Figure 1 describes, and the other half is what it is.