Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/dist/hotstuff/replica.rs

17.1 KiB, 172 runs

created by r1870400018:11226, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! The per-replica HotStuff state machine.
2//!
3//! [`Replica`] is a pure, deterministic state machine. It consumes a leader-
4//! side [`Replica::propose`], [`Replica::on_proposal`], [`Replica::on_vote`],
5//! [`Replica::on_new_view`] and [`Replica::on_timeout`] and emits a list of
6//! [`Command`]s the caller turns into wire sends or a terminal decide.
7//!
8//! The implementation here is Basic HotStuff with:
9//!
10//! - Three-phase happy path: `Prepare -> PreCommit -> Commit -> Decide`.
11//! - Round-robin leader rotation: `leader_for(view) = (view - 1) mod cohort_size`.
12//! - View change via [`Replica::on_timeout`] + [`Replica::on_new_view`].
13//! - The classical safety predicate `safeBlock`: on a Prepare proposal in
14//! view `v > 1` the replica accepts iff the proposal's justify QC either
15//! endorses the replica's locked block or is from a view strictly newer
16//! than the locked QC's view.
17//! - Opaque ride-along signatures (not verified by this primitive).
18//!
19//! Not yet in scope: checkpointing of multiple decisions, Byzantine-fault
20//! simulation tests, signature aggregation (we pass individual signatures
21//! through the QC).
22//!
23//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
24//! Anthropic Claude
25
26use super::types::{
27 BLOCK_HASH_LEN,
28 BlockHash,
29 NewView,
30 Phase,
31 Proposal,
32 Qc,
33 ReplicaId,
34 ViewId,
35 Vote,
36};
37
38use oxedyne_fe2o3_core::prelude::*;
39
40use std::collections::{
41 BTreeMap,
42 HashMap,
43};
44
45
46/// Per-replica configuration.
47#[derive(Clone, Copy, Debug)]
48pub struct Config {
49 pub cohort_size: usize, // λ in the spec
50 pub f: usize, // z = floor((λ - 1) / 3)
51 pub self_id: ReplicaId,
52}
53
54impl Config {
55 /// The quorum threshold `λ - z`.
56 pub fn quorum(&self) -> usize {
57 self.cohort_size - self.f
58 }
59
60 /// Round-robin rotation: view 1's leader is replica 0, view 2's is replica
61 /// 1, and so on.
62 pub fn leader_for(&self, view: ViewId) -> ReplicaId {
63 let idx = (view.saturating_sub(1) as usize) % self.cohort_size;
64 idx as ReplicaId
65 }
66
67 pub fn is_leader_for(&self, view: ViewId) -> bool {
68 self.self_id == self.leader_for(view)
69 }
70
71 pub fn validate(&self) -> Outcome<()> {
72 if self.cohort_size == 0 {
73 return Err(err!(
74 "HotStuff cohort_size must be > 0.";
75 Invalid, Input));
76 }
77 if self.cohort_size < 3 * self.f + 1 {
78 return Err(err!(
79 "HotStuff cohort_size ({}) must be at least 3f+1 = {}.",
80 self.cohort_size, 3 * self.f + 1;
81 Invalid, Input));
82 }
83 if (self.self_id as usize) >= self.cohort_size {
84 return Err(err!(
85 "HotStuff self_id {} out of range (cohort_size = {}).",
86 self.self_id, self.cohort_size;
87 Invalid, Input));
88 }
89 Ok(())
90 }
91}
92
93
94/// A command the state machine asks the caller to perform.
95#[derive(Clone, Debug)]
96pub enum Command {
97 // Goes to every replica including the leader, whose own vote is recorded
98 // through Replica::on_proposal.
99 BroadcastProposal(Proposal),
100 // Normally to the current leader.
101 SendVote {
102 to: ReplicaId,
103 vote: Vote,
104 },
105 // Normally to the incoming leader.
106 SendNewView {
107 to: ReplicaId,
108 new_view: NewView,
109 },
110 // The caller hands the block back to the application; the state machine
111 // is inert afterwards.
112 Decide {
113 view: ViewId,
114 block: Vec<u8>,
115 },
116}
117
118
119/// A pure HotStuff replica state machine (three-phase, with view change).
120pub struct Replica {
121 cfg: Config,
122 view: ViewId, // starts at 1, advances on on_timeout
123 blocks: HashMap<BlockHash, Vec<u8>>,
124 // Leader-side aggregation for the current view. The BTreeMap keeps voters
125 // in ascending-unique order, which is what the QC requires.
126 gathered: HashMap<(Phase, BlockHash), BTreeMap<ReplicaId, Vote>>,
127 new_views: BTreeMap<ReplicaId, NewView>,
128 last_voted: Option<Phase>,
129 // Carried forward across views and included in outgoing NewView messages.
130 prepare_qc: Option<Qc>,
131 // The highest pre-commit QC seen; gates safeBlock on Prepare proposals.
132 locked_qc: Option<Qc>,
133 decided: bool, // the state machine is inert once set
134 // Stops a later NewView quorum re-opening a Prepare in the same view.
135 prepared_this_view: bool,
136}
137
138impl Replica {
139 pub fn new(cfg: Config) -> Outcome<Self> {
140 res!(cfg.validate());
141 Ok(Self {
142 cfg,
143 view: 1,
144 blocks: HashMap::new(),
145 gathered: HashMap::new(),
146 new_views: BTreeMap::new(),
147 last_voted: None,
148 prepare_qc: None,
149 locked_qc: None,
150 decided: false,
151 prepared_this_view: false,
152 })
153 }
154
155 pub fn config(&self) -> Config {
156 self.cfg
157 }
158
159 pub fn view(&self) -> ViewId {
160 self.view
161 }
162
163 pub fn has_decided(&self) -> bool {
164 self.decided
165 }
166
167 pub fn locked_qc(&self) -> Option<&Qc> {
168 self.locked_qc.as_ref()
169 }
170
171 pub fn prepare_qc(&self) -> Option<&Qc> {
172 self.prepare_qc.as_ref()
173 }
174
175 /// Leader-only. Valid in view 1, which has no previous state, or after a
176 /// NewView quorum in view > 1 where no replica reported a prepare QC --
177 /// otherwise the leader must propose the block the highest prepare QC
178 /// already pinned.
179 pub fn propose(&mut self, block: Vec<u8>, block_hash: BlockHash) -> Outcome<Vec<Command>> {
180 if !self.cfg.is_leader_for(self.view) {
181 return Err(err!(
182 "Only the leader of view {} may call propose (self_id = {}).",
183 self.view, self.cfg.self_id;
184 Invalid, Order));
185 }
186 if self.decided {
187 return Err(err!(
188 "propose called after Decide.";
189 Invalid, Order));
190 }
191 if self.prepared_this_view {
192 return Err(err!(
193 "propose called twice in view {}.", self.view;
194 Invalid, Order));
195 }
196 self.blocks.insert(block_hash, block.clone());
197 self.prepared_this_view = true;
198 let proposal = Proposal {
199 view: self.view,
200 phase: Phase::Prepare,
201 block_hash,
202 block: Some(block),
203 justify: None,
204 };
205 Ok(vec![Command::BroadcastProposal(proposal)])
206 }
207
208 pub fn on_proposal(&mut self, proposal: Proposal) -> Outcome<Vec<Command>> {
209 if self.decided {
210 return Ok(Vec::new());
211 }
212 if proposal.view != self.view {
213 return Err(err!(
214 "Proposal view {} does not match current view {}.",
215 proposal.view, self.view;
216 Invalid, Input, Mismatch));
217 }
218
219 match proposal.phase {
220 Phase::Prepare => {
221 let block = match proposal.block {
222 Some(b) => b,
223 None => return Err(err!(
224 "Prepare proposal must carry a block payload.";
225 Invalid, Input, Missing)),
226 };
227 // View 1: no justify is legal.
228 // View > 1: a justify is required if this replica holds a
229 // locked QC (otherwise the safeBlock rule cannot be checked
230 // and a locked block could be overridden by a fresh proposal).
231 // A justify-less Prepare in view > 1 is legal iff no lock is
232 // held -- the new leader learned via NewView that no replica
233 // was locked and is starting fresh.
234 if self.view == 1 {
235 if proposal.justify.is_some() {
236 return Err(err!(
237 "Prepare proposal in view 1 must not carry a justify QC.";
238 Invalid, Input));
239 }
240 } else {
241 match &proposal.justify {
242 Some(justify) => {
243 if justify.phase != Phase::Prepare {
244 return Err(err!(
245 "Prepare proposal's justify must be a Prepare QC, got {:?}.",
246 justify.phase;
247 Invalid, Input));
248 }
249 res!(justify.validate(
250 justify.view,
251 Phase::Prepare,
252 &justify.block_hash,
253 self.cfg.quorum(),
254 self.cfg.cohort_size,
255 ));
256 if !self.safe_block(&proposal.block_hash, justify) {
257 return Err(err!(
258 "Prepare proposal fails safeBlock: justify view {} vs locked view {:?}.",
259 justify.view,
260 self.locked_qc.as_ref().map(|q| q.view);
261 Security, Invalid));
262 }
263 },
264 None => {
265 if self.locked_qc.is_some() {
266 return Err(err!(
267 "Prepare proposal in view {} lacks a justify QC \
268 but this replica is locked on view {:?} -- cannot \
269 accept a fresh proposal that would abandon the lock.",
270 self.view,
271 self.locked_qc.as_ref().map(|q| q.view);
272 Security, Invalid, Missing));
273 }
274 },
275 }
276 }
277 self.blocks.insert(proposal.block_hash, block);
278 },
279 Phase::PreCommit => {
280 let justify = match &proposal.justify {
281 Some(q) => q,
282 None => return Err(err!(
283 "PreCommit proposal requires a Prepare QC.";
284 Invalid, Input, Missing)),
285 };
286 res!(justify.validate(
287 self.view,
288 Phase::Prepare,
289 &proposal.block_hash,
290 self.cfg.quorum(),
291 self.cfg.cohort_size,
292 ));
293 // Record as our prepare_qc if newer.
294 if self.prepare_qc.as_ref().map(|q| q.view).unwrap_or(0) < justify.view {
295 self.prepare_qc = Some(justify.clone());
296 }
297 if proposal.block.is_some() {
298 return Err(err!(
299 "PreCommit proposal must not re-send the block payload.";
300 Invalid, Input));
301 }
302 },
303 Phase::Commit => {
304 let justify = match &proposal.justify {
305 Some(q) => q,
306 None => return Err(err!(
307 "Commit proposal requires a PreCommit QC.";
308 Invalid, Input, Missing)),
309 };
310 res!(justify.validate(
311 self.view,
312 Phase::PreCommit,
313 &proposal.block_hash,
314 self.cfg.quorum(),
315 self.cfg.cohort_size,
316 ));
317 // Locking: record the pre-commit QC as locked_qc.
318 if self.locked_qc.as_ref().map(|q| q.view).unwrap_or(0) < justify.view {
319 self.locked_qc = Some(justify.clone());
320 }
321 if proposal.block.is_some() {
322 return Err(err!(
323 "Commit proposal must not re-send the block payload.";
324 Invalid, Input));
325 }
326 },
327 Phase::Decide => {
328 let justify = match &proposal.justify {
329 Some(q) => q,
330 None => return Err(err!(
331 "Decide proposal requires a Commit QC.";
332 Invalid, Input, Missing)),
333 };
334 res!(justify.validate(
335 self.view,
336 Phase::Commit,
337 &proposal.block_hash,
338 self.cfg.quorum(),
339 self.cfg.cohort_size,
340 ));
341 let block = match self.blocks.get(&proposal.block_hash) {
342 Some(b) => b.clone(),
343 None => return Err(err!(
344 "Decide proposal references unknown block hash.";
345 Invalid, Input, Missing)),
346 };
347 self.decided = true;
348 return Ok(vec![Command::Decide { view: self.view, block }]);
349 },
350 }
351
352 // Cast vote unless we've already voted in this or a later phase.
353 if let Some(p) = self.last_voted {
354 if self.phase_rank(p) >= self.phase_rank(proposal.phase) {
355 return Ok(Vec::new());
356 }
357 }
358 self.last_voted = Some(proposal.phase);
359 let vote = Vote {
360 view: self.view,
361 phase: proposal.phase,
362 block_hash: proposal.block_hash,
363 voter: self.cfg.self_id,
364 signature: self.self_signature(proposal.phase, &proposal.block_hash),
365 };
366 let leader = self.cfg.leader_for(self.view);
367 Ok(vec![Command::SendVote { to: leader, vote }])
368 }
369
370 /// Leader-only; a non-leader ignores votes.
371 pub fn on_vote(&mut self, vote: Vote) -> Outcome<Vec<Command>> {
372 if !self.cfg.is_leader_for(self.view) {
373 return Ok(Vec::new());
374 }
375 if self.decided {
376 return Ok(Vec::new());
377 }
378 if vote.view != self.view {
379 return Err(err!(
380 "Vote view {} does not match current view {}.",
381 vote.view, self.view;
382 Invalid, Input, Mismatch));
383 }
384 if (vote.voter as usize) >= self.cfg.cohort_size {
385 return Err(err!(
386 "Vote voter id {} out of range (cohort_size = {}).",
387 vote.voter, self.cfg.cohort_size;
388 Invalid, Input));
389 }
390 if vote.phase == Phase::Decide {
391 return Err(err!(
392 "Decide votes are not part of the protocol.";
393 Invalid, Input));
394 }
395 let key = (vote.phase, vote.block_hash);
396 let entry = self.gathered.entry(key).or_insert_with(BTreeMap::new);
397 entry.insert(vote.voter, vote);
398 if entry.len() < self.cfg.quorum() {
399 return Ok(Vec::new());
400 }
401 let (phase, block_hash) = key;
402 let sigs: Vec<(ReplicaId, Vec<u8>)> = entry.iter()
403 .map(|(id, v)| (*id, v.signature.clone()))
404 .collect();
405 let qc = Qc { view: self.view, phase, block_hash, signatures: sigs };
406 self.gathered.remove(&key);
407 // Update leader's own prepare_qc / locked_qc bookkeeping as it forms
408 // each QC, so subsequent views behave consistently if this replica
409 // remains in the cohort after a leader handover.
410 match phase {
411 Phase::Prepare => {
412 if self.prepare_qc.as_ref().map(|q| q.view).unwrap_or(0) < qc.view {
413 self.prepare_qc = Some(qc.clone());
414 }
415 },
416 Phase::PreCommit => {
417 if self.locked_qc.as_ref().map(|q| q.view).unwrap_or(0) < qc.view {
418 self.locked_qc = Some(qc.clone());
419 }
420 },
421 _ => {},
422 }
423 let next_phase = match phase.next() {
424 Some(p) => p,
425 None => return Err(err!(
426 "Quorum reached on terminal phase Decide -- invalid state.";
427 Invalid, Order, Bug)),
428 };
429 let proposal = Proposal {
430 view: self.view,
431 phase: next_phase,
432 block_hash,
433 block: None,
434 justify: Some(qc),
435 };
436 Ok(vec![Command::BroadcastProposal(proposal)])
437 }
438
439 /// Advances the replica to the next view and emits a [`Command::SendNewView`]
440 /// targeting the new leader.
441 ///
442 /// The caller is expected to invoke this when its local timer fires
443 /// without observing progress in the current view. The state machine
444 /// has no notion of time -- timeout policy is the caller's to set.
445 pub fn on_timeout(&mut self) -> Outcome<Vec<Command>> {
446 if self.decided {
447 return Ok(Vec::new());
448 }
449 self.view = self.view.wrapping_add(1);
450 self.gathered.clear();
451 self.new_views.clear();
452 self.last_voted = None;
453 self.prepared_this_view = false;
454 let new_view = NewView {
455 view: self.view,
456 sender: self.cfg.self_id,
457 prepare_qc: self.prepare_qc.clone(),
458 };
459 let to = self.cfg.leader_for(self.view);
460 Ok(vec![Command::SendNewView { to, new_view }])
461 }
462
463 /// Leader-only: consumes an incoming [`NewView`] and, on accumulating a
464 /// quorum, opens the next view's Prepare proposal if a prior prepare QC
465 /// pinned a block.
466 ///
467 /// If no participating replica reported a prepare QC the primitive does
468 /// *not* speculate a block -- the leader must then call
469 /// [`Replica::propose`] with a fresh block of its choice.
470 pub fn on_new_view(&mut self, nv: NewView) -> Outcome<Vec<Command>> {
471 if !self.cfg.is_leader_for(self.view) {
472 return Ok(Vec::new());
473 }
474 if self.decided {
475 return Ok(Vec::new());
476 }
477 if nv.view != self.view {
478 return Err(err!(
479 "NewView view {} does not match current view {}.",
480 nv.view, self.view;
481 Invalid, Input, Mismatch));
482 }
483 if (nv.sender as usize) >= self.cfg.cohort_size {
484 return Err(err!(
485 "NewView sender {} out of range (cohort_size = {}).",
486 nv.sender, self.cfg.cohort_size;
487 Invalid, Input));
488 }
489 // Validate the carried prepare_qc, if any.
490 if let Some(qc) = &nv.prepare_qc {
491 res!(qc.validate(
492 qc.view,
493 Phase::Prepare,
494 &qc.block_hash,
495 self.cfg.quorum(),
496 self.cfg.cohort_size,
497 ));
498 // Update the leader's own prepare_qc if this one is higher --
499 // helps leaders that were offline for earlier views catch up.
500 if self.prepare_qc.as_ref().map(|q| q.view).unwrap_or(0) < qc.view {
501 self.prepare_qc = Some(qc.clone());
502 }
503 }
504 self.new_views.insert(nv.sender, nv);
505 if self.new_views.len() < self.cfg.quorum() {
506 return Ok(Vec::new());
507 }
508 if self.prepared_this_view {
509 // Already opened the view; further NewViews are dropped.
510 return Ok(Vec::new());
511 }
512 // Pick the highest prepare_qc among gathered NewViews.
513 let highest: Option<Qc> = self.new_views.values()
514 .filter_map(|nv| nv.prepare_qc.clone())
515 .max_by_key(|qc| qc.view);
516 match highest {
517 Some(qc) => {
518 // The block hash is pinned by qc; we need the block bytes.
519 let block_bytes = match self.blocks.get(&qc.block_hash) {
520 Some(b) => b.clone(),
521 None => return Err(err!(
522 "Leader for view {} lacks block payload for pinned hash. \
523 The caller must arrange payload retrieval (e.g. out-of-band \
524 fetch) before the new view's Prepare can be opened.",
525 self.view;
526 Missing, Data)),
527 };
528 self.prepared_this_view = true;
529 let proposal = Proposal {
530 view: self.view,
531 phase: Phase::Prepare,
532 block_hash: qc.block_hash,
533 block: Some(block_bytes),
534 justify: Some(qc),
535 };
536 Ok(vec![Command::BroadcastProposal(proposal)])
537 },
538 None => {
539 // No replica had a prepare QC. Leader must call propose()
540 // with a fresh block of its choice. Nothing to emit here.
541 Ok(Vec::new())
542 },
543 }
544 }
545
546 /// Is this replica the leader, holding a NewView quorum in which no prior
547 /// prepare QC pinned a block? If so it must supply a fresh block through
548 /// [`Replica::propose`].
549 pub fn awaiting_fresh_block(&self) -> bool {
550 self.cfg.is_leader_for(self.view)
551 && !self.decided
552 && !self.prepared_this_view
553 && self.new_views.len() >= self.cfg.quorum()
554 && self.new_views.values().all(|nv| nv.prepare_qc.is_none())
555 }
556
557 fn safe_block(&self, _block_hash: &BlockHash, justify: &Qc) -> bool {
558 match &self.locked_qc {
559 None => true,
560 Some(lqc) => {
561 // Liveness: a later view overrides.
562 if justify.view > lqc.view {
563 return true;
564 }
565 // Safety: the justified block must match our locked block.
566 justify.block_hash == lqc.block_hash
567 },
568 }
569 }
570
571 fn phase_rank(&self, p: Phase) -> u8 {
572 match p {
573 Phase::Prepare => 1,
574 Phase::PreCommit => 2,
575 Phase::Commit => 3,
576 Phase::Decide => 4,
577 }
578 }
579
580 fn self_signature(&self, phase: Phase, block_hash: &BlockHash) -> Vec<u8> {
581 let mut out = Vec::with_capacity(2 + 1 + BLOCK_HASH_LEN);
582 out.extend_from_slice(&self.cfg.self_id.to_le_bytes());
583 out.push(phase_tag(phase));
584 out.extend_from_slice(block_hash);
585 out
586 }
587}
588
589fn phase_tag(p: Phase) -> u8 {
590 match p {
591 Phase::Prepare => 1,
592 Phase::PreCommit => 2,
593 Phase::Commit => 3,
594 Phase::Decide => 4,
595 }
596}