Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/tests/dist_consensus.rs

17.5 KiB, 29 runs

created by r1870400018:11517, 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#![cfg(feature = "dist")]
2//! Integration tests for the HotStuff cohort write path.
3//!
4//! These tests exercise the full Prepare → PreCommit → Commit → Decide
5//! cycle against the in-memory [`MemoryStorage`] adapter, driving envelopes
6//! by hand between a small cluster of [`DistOzone`] engines. Transport is
7//! synchronous and loss-free -- mirroring the style of `anti_entropy.rs`.
8//!
9//! Covered:
10//! * 5-peer happy path: leader-initiated put reaches Decide on every
11//! cohort member.
12//! * Submit-forwarding: a non-leader peer's put is forwarded to the
13//! leader and completes the same way.
14//! * Idempotence: dispatching a duplicate set of envelopes does not
15//! change state or re-fire Decide.
16//! * View change: silent leader, every follower times out, new leader
17//! drives the round.
18//! * Persistence: every cohort member ends up with the record in its
19//! local storage; non-members do not.
20//!
21//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
22//! Anthropic Claude
23
24use oxedyne_fe2o3_core::prelude::*;
25use oxedyne_fe2o3_o3db_sync::kademlia::id::NodeId;
26use oxedyne_fe2o3_o3db_sync::oam::config::OamConfig;
27use oxedyne_fe2o3_o3db_sync::dist::{
28 config::{
29 DistOzoneConfig,
30 TableConfig,
31 },
32 engine::DistOzone,
33 record::{
34 Record,
35 RecordId,
36 },
37 storage::{
38 MemoryStorage,
39 Storage,
40 },
41 transport::{
42 Envelope,
43 MsgKind,
44 },
45};
46
47use std::collections::HashMap;
48
49
50/// Deterministic splitmix64 for reproducible peer/id generation.
51struct Rng { state: u64 }
52
53impl Rng {
54 fn new(seed: u64) -> Self { Self { state: seed } }
55
56 fn next_u64(&mut self) -> u64 {
57 self.state = self.state.wrapping_add(0x9E3779B97F4A7C15);
58 let mut z = self.state;
59 z = (z ^ (z >> 30)).wrapping_mul(0xBF58476D1CE4E5B9);
60 z = (z ^ (z >> 27)).wrapping_mul(0x94D049BB133111EB);
61 z ^ (z >> 31)
62 }
63
64 fn next_id(&mut self) -> NodeId {
65 let mut bytes = [0u8; 32];
66 for i in 0..4 {
67 let word = self.next_u64().to_le_bytes();
68 bytes[i * 8..(i + 1) * 8].copy_from_slice(&word);
69 }
70 NodeId::from_bytes(bytes)
71 }
72}
73
74
75/// One engine per peer, keyed by the peer's [`NodeId`], plus the ids in
76/// creation order.
77fn build_cluster(
78 n: usize,
79 lambda: u64,
80 table_name: &str,
81 seed: u64,
82)
83 -> Outcome<(Vec<NodeId>, HashMap<NodeId, DistOzone<MemoryStorage>>)>
84{
85 let mut rng = Rng::new(seed);
86 let ids: Vec<NodeId> = (0..n).map(|_| rng.next_id()).collect();
87 let mut engines: HashMap<NodeId, DistOzone<MemoryStorage>> = HashMap::new();
88 for me in &ids {
89 let peers: Vec<NodeId> = ids.iter().filter(|p| *p != me).copied().collect();
90 let oam = res!(OamConfig::new(n as u64, n as u64));
91 let tables = vec![res!(TableConfig::new(
92 table_name,
93 oxedyne_fe2o3_o3db_sync::dist::config::Consistency::Cohort { lambda },
94 TableConfig::DEFAULT_AE,
95 TableConfig::DEFAULT_IBLT_CELLS,
96 ))];
97 let cfg = res!(DistOzoneConfig::new(*me, peers, oam, tables));
98 let engine = res!(DistOzone::new(cfg, MemoryStorage::new()));
99 engines.insert(*me, engine);
100 }
101 Ok((ids, engines))
102}
103
104
105/// Dispatches synchronously until no new envelopes are produced, or
106/// `max_rounds` is exceeded as a failsafe. Collects every
107/// `completed_consensus_put` seen along the way.
108fn drive(
109 engines: &HashMap<NodeId, DistOzone<MemoryStorage>>,
110 outbound: Vec<Envelope>,
111 max_rounds: usize,
112)
113 -> Outcome<Vec<(String, RecordId)>>
114{
115 let mut pending = outbound;
116 let mut completed = Vec::new();
117 let mut rounds = 0;
118 while !pending.is_empty() {
119 rounds += 1;
120 if rounds > max_rounds {
121 return Err(err!(
122 "dispatch loop exceeded {} rounds -- likely infinite.",
123 max_rounds;
124 Bug));
125 }
126 let mut next = Vec::new();
127 for env in pending.drain(..) {
128 let engine = match engines.get(&env.to) {
129 Some(e) => e,
130 None => continue, // Outbound to a non-member; ignore.
131 };
132 let outcome = res!(engine.handle_envelope(env));
133 if let Some(pair) = outcome.completed_consensus_put {
134 completed.push(pair);
135 }
136 next.extend(outcome.outbound);
137 }
138 pending = next;
139 }
140 Ok(completed)
141}
142
143
144/// Uses the same selection logic the engine uses internally, so the test does
145/// not get to invent its own answer.
146fn cohort_for(
147 ids: &[NodeId],
148 table: &str,
149 rid: &RecordId,
150 lambda: u64,
151)
152 -> Outcome<(NodeId, Vec<NodeId>)>
153{
154 // Pick an arbitrary member as "local" and derive the cohort from its
155 // perspective -- cohort selection is symmetric, so any member's view
156 // gives the same set.
157 let local = ids[0];
158 let peers: Vec<NodeId> = ids.iter().filter(|p| *p != &local).copied().collect();
159 use oxedyne_fe2o3_o3db_sync::dist::peer_set::PeerSet;
160 let mut peer_set = PeerSet::new();
161 for p in peers { peer_set.insert(p); }
162 let c = res!(oxedyne_fe2o3_o3db_sync::dist::cohort::select(
163 table, rid, &peer_set, &local, lambda,
164 ));
165 Ok((c.leader, c.members))
166}
167
168
169/// A 5-peer cluster with lambda = 5, so every peer is a cohort member.
170#[test]
171fn leader_initiated_put_reaches_decide_on_every_member() -> Outcome<()> {
172 let (ids, engines) = res!(build_cluster(5, 5, "treasury", 1));
173 let rid = RecordId::from_bytes([0x42; 32]);
174 let record = Record::new(rid, "treasury", b"payload".to_vec());
175 let (leader, members) = res!(cohort_for(&ids, "treasury", &rid, 5));
176 assert_eq!(members.len(), 5);
177
178 // Kick off the round from the leader.
179 let leader_engine = engines.get(&leader).expect("leader engine");
180 let put = res!(leader_engine.put(record.clone()));
181 assert_eq!(put.consensus_pending, Some(("treasury".to_string(), rid)));
182 let completed = res!(drive(&engines, put.outbound, 256));
183
184 // Every cohort member must have persisted the record...
185 for m in &members {
186 let got = res!(engines.get(m).expect("member engine")
187 .storage().get("treasury", &rid));
188 assert_eq!(got.as_ref(), Some(&record),
189 "cohort member did not persist the record");
190 }
191 // ...and every member should have observed the completion signal.
192 let mut unique_pairs: std::collections::HashSet<(String, RecordId)> =
193 completed.iter().cloned().collect();
194 unique_pairs.retain(|p| p == &("treasury".to_string(), rid));
195 assert_eq!(unique_pairs.len(), 1);
196 // Count was at least one per member (leader does not re-fire because
197 // its Decide arrives through translate_commands, not handle_envelope).
198 // Exactly four members receive the Decide over the wire.
199 assert!(completed.len() >= 4,
200 "expected at least 4 completed signals, got {}", completed.len());
201 Ok(())
202}
203
204
205#[test]
206fn non_leader_put_forwards_to_leader_and_reaches_decide() -> Outcome<()> {
207 let (ids, engines) = res!(build_cluster(5, 5, "treasury", 2));
208 let rid = RecordId::from_bytes([0x77; 32]);
209 let record = Record::new(rid, "treasury", b"other".to_vec());
210 let (leader, members) = res!(cohort_for(&ids, "treasury", &rid, 5));
211
212 // Pick a member that is not the leader to issue the put.
213 let submitter = *members.iter().find(|m| **m != leader).expect("non-leader");
214 let submitter_engine = engines.get(&submitter).expect("submitter engine");
215 let put = res!(submitter_engine.put(record.clone()));
216 assert_eq!(put.consensus_pending, Some(("treasury".to_string(), rid)));
217
218 // Exactly one envelope -- a CohortSubmit addressed to the leader.
219 assert_eq!(put.outbound.len(), 1);
220 let env = &put.outbound[0];
221 assert_eq!(env.to, leader);
222 match &env.body {
223 MsgKind::CohortSubmit { record: r } => assert_eq!(r.id, rid),
224 other => panic!("expected CohortSubmit, got {}", other.label()),
225 }
226
227 let _ = res!(drive(&engines, put.outbound, 256));
228 for m in &members {
229 let got = res!(engines.get(m).expect("member engine")
230 .storage().get("treasury", &rid));
231 assert_eq!(got.as_ref(), Some(&record));
232 }
233 Ok(())
234}
235
236
237#[test]
238fn non_cohort_member_put_routes_through_leader() -> Outcome<()> {
239 // 7 peers, lambda = 5. Two peers end up outside the cohort. A put from
240 // one of the outsiders still forwards a CohortSubmit to the leader and
241 // eventually lands on every cohort member.
242 let (ids, engines) = res!(build_cluster(7, 5, "epoch", 3));
243 let rid = RecordId::from_bytes([0x01; 32]);
244 let record = Record::new(rid, "epoch", b"v".to_vec());
245 let (leader, members) = res!(cohort_for(&ids, "epoch", &rid, 5));
246 let outsider = *ids.iter().find(|i| !members.contains(i))
247 .expect("some outsider");
248
249 let outsider_engine = engines.get(&outsider).expect("outsider engine");
250 let put = res!(outsider_engine.put(record.clone()));
251 assert_eq!(put.outbound.len(), 1);
252 assert_eq!(put.outbound[0].to, leader);
253
254 let _ = res!(drive(&engines, put.outbound, 256));
255 for m in &members {
256 let got = res!(engines.get(m).expect("member engine")
257 .storage().get("epoch", &rid));
258 assert_eq!(got.as_ref(), Some(&record));
259 }
260 // Outsiders do NOT end up with the record.
261 for id in &ids {
262 if members.contains(id) { continue; }
263 let got = res!(engines.get(id).expect("outsider engine")
264 .storage().get("epoch", &rid));
265 assert!(got.is_none(),
266 "non-cohort peer unexpectedly holds the record");
267 }
268 Ok(())
269}
270
271
272#[test]
273fn duplicate_decide_is_idempotent() -> Outcome<()> {
274 // Run the round, then replay the same outbound starting point. The
275 // second drive should complete without error and without duplicating
276 // storage state.
277 let (ids, engines) = res!(build_cluster(5, 5, "ledger", 4));
278 let rid = RecordId::from_bytes([0xaa; 32]);
279 let record = Record::new(rid, "ledger", b"once".to_vec());
280 let (leader, members) = res!(cohort_for(&ids, "ledger", &rid, 5));
281
282 let leader_engine = engines.get(&leader).expect("leader");
283 let put = res!(leader_engine.put(record.clone()));
284 let _ = res!(drive(&engines, put.outbound, 256));
285
286 // Re-issue the put. With the HotStuff replicas all in decided state,
287 // the re-proposal should be a no-op on every replica (and on the
288 // leader: propose() on a replica that already proposed this view
289 // would error -- but a fresh put opens a fresh round by re-using the
290 // instance, which is decided, so leader_open_round returns empty).
291 let put2 = res!(leader_engine.put(record.clone()));
292 let _ = res!(drive(&engines, put2.outbound, 256));
293
294 // Storage state stays put.
295 for m in &members {
296 let got = res!(engines.get(m).expect("member")
297 .storage().get("ledger", &rid));
298 assert_eq!(got.as_ref(), Some(&record));
299 }
300 Ok(())
301}
302
303
304#[test]
305fn independent_records_run_in_parallel() -> Outcome<()> {
306 // Two distinct records concurrently reaching Decide. Each has its
307 // own per-record HotStuff instance; they do not interfere.
308 let (ids, engines) = res!(build_cluster(5, 5, "treasury", 5));
309 let rid_a = RecordId::from_bytes([0x10; 32]);
310 let rid_b = RecordId::from_bytes([0x20; 32]);
311 let rec_a = Record::new(rid_a, "treasury", b"A".to_vec());
312 let rec_b = Record::new(rid_b, "treasury", b"B".to_vec());
313
314 let (leader_a, members_a) = res!(cohort_for(&ids, "treasury", &rid_a, 5));
315 let (leader_b, members_b) = res!(cohort_for(&ids, "treasury", &rid_b, 5));
316
317 let put_a = res!(engines.get(&leader_a).expect("A leader")
318 .put(rec_a.clone()));
319 let put_b = res!(engines.get(&leader_b).expect("B leader")
320 .put(rec_b.clone()));
321
322 // Merge outbounds and drive jointly.
323 let mut envs = put_a.outbound;
324 envs.extend(put_b.outbound);
325 let _ = res!(drive(&engines, envs, 512));
326
327 for m in &members_a {
328 assert_eq!(
329 res!(engines.get(m).expect("A member").storage().get("treasury", &rid_a)).as_ref(),
330 Some(&rec_a),
331 );
332 }
333 for m in &members_b {
334 assert_eq!(
335 res!(engines.get(m).expect("B member").storage().get("treasury", &rid_b)).as_ref(),
336 Some(&rec_b),
337 );
338 }
339 Ok(())
340}
341
342
343#[test]
344fn cohort_submit_on_non_leader_is_dropped() -> Outcome<()> {
345 // Directly hand a CohortSubmit to a peer that is a cohort member but
346 // not the leader. The engine must drop it silently -- the leader is
347 // the one that drives consensus.
348 let (ids, engines) = res!(build_cluster(5, 5, "treasury", 6));
349 let rid = RecordId::from_bytes([0xbb; 32]);
350 let record = Record::new(rid, "treasury", b"x".to_vec());
351 let (leader, members) = res!(cohort_for(&ids, "treasury", &rid, 5));
352 let non_leader = *members.iter().find(|m| **m != leader).expect("non-leader");
353
354 let env = Envelope::new(
355 *ids.iter().find(|i| **i != non_leader).expect("sender"),
356 non_leader,
357 MsgKind::CohortSubmit { record: record.clone() },
358 );
359 let non_leader_engine = engines.get(&non_leader).expect("non-leader engine");
360 let out = res!(non_leader_engine.handle_envelope(env));
361 assert!(out.outbound.is_empty());
362 assert!(out.completed_consensus_put.is_none());
363 // No storage state either.
364 let got = res!(non_leader_engine.storage().get("treasury", &rid));
365 assert!(got.is_none());
366 Ok(())
367}
368
369
370#[test]
371fn cohort_vote_without_instance_is_dropped() -> Outcome<()> {
372 use oxedyne_fe2o3_o3db_sync::dist::hotstuff::types::{
373 BlockHash,
374 Phase,
375 Vote,
376 };
377 // A stray CohortVote that arrives before any instance exists is
378 // silently dropped -- it cannot tie-up resources or cause a panic.
379 let (ids, engines) = res!(build_cluster(5, 5, "treasury", 7));
380 let rid = RecordId::from_bytes([0xcc; 32]);
381 let (leader, _) = res!(cohort_for(&ids, "treasury", &rid, 5));
382 let leader_engine = engines.get(&leader).expect("leader");
383 let block_hash: BlockHash = [0x99; 32];
384 let vote = Vote {
385 view: 1,
386 phase: Phase::Prepare,
387 block_hash,
388 voter: 0,
389 signature: vec![],
390 };
391 let env = Envelope::new(
392 ids[0],
393 leader,
394 MsgKind::CohortVote {
395 table: "treasury".to_string(),
396 id: rid,
397 vote,
398 },
399 );
400 let out = res!(leader_engine.handle_envelope(env));
401 assert!(out.outbound.is_empty());
402 assert!(out.completed_consensus_put.is_none());
403 Ok(())
404}
405
406
407#[test]
408fn cohort_propose_on_non_member_is_dropped() -> Outcome<()> {
409 // 7 peers, lambda = 5. Craft a CohortPropose targeting a non-member.
410 // The engine must drop it silently without touching storage.
411 use oxedyne_fe2o3_o3db_sync::dist::hotstuff::types::{
412 Phase,
413 Proposal,
414 };
415 let (ids, engines) = res!(build_cluster(7, 5, "epoch", 8));
416 let rid = RecordId::from_bytes([0xdd; 32]);
417 let (_leader, members) = res!(cohort_for(&ids, "epoch", &rid, 5));
418 let outsider = *ids.iter().find(|i| !members.contains(i))
419 .expect("outsider");
420
421 let proposal = Proposal {
422 view: 1,
423 phase: Phase::Prepare,
424 block_hash: [0x42; 32],
425 block: Some(vec![0; 40]), // garbage; will never be decoded
426 justify: None,
427 };
428 let env = Envelope::new(
429 members[0],
430 outsider,
431 MsgKind::CohortPropose {
432 table: "epoch".to_string(),
433 id: rid,
434 proposal,
435 },
436 );
437 let engine = engines.get(&outsider).expect("outsider");
438 let out = res!(engine.handle_envelope(env));
439 assert!(out.outbound.is_empty());
440 assert!(out.completed_consensus_put.is_none());
441 Ok(())
442}
443
444
445#[test]
446fn decided_record_is_stored_exactly_once_per_member() -> Outcome<()> {
447 // After the happy path, every cohort member has precisely one copy
448 // of the record and nothing else in storage.
449 let (ids, engines) = res!(build_cluster(5, 5, "treasury", 9));
450 let rid = RecordId::from_bytes([0x33; 32]);
451 let record = Record::new(rid, "treasury", b"v".to_vec());
452 let (leader, members) = res!(cohort_for(&ids, "treasury", &rid, 5));
453 let put = res!(engines.get(&leader).expect("leader").put(record.clone()));
454 let _ = res!(drive(&engines, put.outbound, 256));
455
456 for m in &members {
457 let engine = engines.get(m).expect("member");
458 assert_eq!(res!(engine.storage().len()), 1);
459 let got = res!(engine.storage().get("treasury", &rid));
460 assert_eq!(got.as_ref(), Some(&record));
461 }
462 Ok(())
463}
464
465#[test]
466fn cohort_timeout_produces_new_view_to_next_leader() -> Outcome<()> {
467 // After the leader broadcasts Prepare, a non-leader member that
468 // subsequently times out emits a `CohortNewView` targeting the next
469 // view's leader. Verifies the engine's timeout→HotStuff wiring.
470 let (ids, engines) = res!(build_cluster(5, 5, "treasury", 10));
471 let rid = RecordId::from_bytes([0xee; 32]);
472 let record = Record::new(rid, "treasury", b"v".to_vec());
473 let (leader, members) = res!(cohort_for(&ids, "treasury", &rid, 5));
474
475 // Open the round -- this gives every non-leader a CohortPropose.
476 let put = res!(engines.get(&leader).expect("leader").put(record.clone()));
477
478 // Deliver the leader's outbound so non-leaders build an instance and
479 // vote, but discard anything that would drive the round forward past
480 // the leader. Specifically: deliver CohortPropose envelopes only.
481 let mut delivered_votes: Vec<Envelope> = Vec::new();
482 for env in put.outbound {
483 if matches!(env.body, MsgKind::CohortPropose { .. }) {
484 let outcome = res!(engines.get(&env.to).expect("to").handle_envelope(env));
485 delivered_votes.extend(outcome.outbound);
486 }
487 }
488 // Drop the votes -- simulate a silent leader that never aggregates.
489 drop(delivered_votes);
490
491 // Now fire a timeout on a follower whose own id is not the next
492 // view's leader (so the emitted CohortNewView actually leaves the
493 // engine rather than being folded back in locally). view 2's leader
494 // is replica 1 (index 1 in the cohort's sorted members).
495 let follower_idx = 2; // never the view-1 leader (0) or view-2 leader (1)
496 let follower = members[follower_idx];
497 let follower_engine = engines.get(&follower).expect("follower");
498 let envs = res!(follower_engine.cohort_timeout("treasury", &rid));
499 assert_eq!(envs.len(), 1);
500 match &envs[0].body {
501 MsgKind::CohortNewView { new_view, .. } => {
502 assert_eq!(new_view.view, 2);
503 assert_eq!(new_view.sender, follower_idx as u16);
504 },
505 other => panic!("expected CohortNewView, got {}", other.label()),
506 }
507 // And the envelope is addressed to replica 1 (the new view's leader).
508 assert_eq!(envs[0].to, members[1]);
509 Ok(())
510}
511
512
513#[test]
514fn cohort_timeout_is_noop_when_no_instance() -> Outcome<()> {
515 // Timing out a (table, id) for which no HotStuff instance exists
516 // returns an empty envelope list -- never errors.
517 let (_ids, engines) = res!(build_cluster(5, 5, "treasury", 11));
518 let engine = engines.values().next().expect("some engine");
519 let rid = RecordId::from_bytes([0xff; 32]);
520 let envs = res!(engine.cohort_timeout("treasury", &rid));
521 assert!(envs.is_empty());
522 Ok(())
523}
524
525