Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/tests/dist_engine.rs

20.9 KiB, 26 runs

created by r1870400018:11396, 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 distributed-Ozone engine.
3//!
4//! Tests cover the pure state-machine behaviour of the engine against the
5//! in-memory [`MemoryStorage`] adapter: configuration validation, placement
6//! decisions, write-path outbound construction, read-path local/remote
7//! branching, inbound handling, response correlation, and peer-set / OAM
8//! mutation.
9//!
10//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
11//! Anthropic Claude
12
13use oxedyne_fe2o3_core::prelude::*;
14use oxedyne_fe2o3_o3db_sync::kademlia::id::NodeId;
15use oxedyne_fe2o3_o3db_sync::oam::{
16 config::OamConfig,
17 threshold::Threshold,
18};
19use oxedyne_fe2o3_o3db_sync::dist::{
20 config::{
21 Consistency,
22 DistOzoneConfig,
23 TableConfig,
24 },
25 engine::{
26 DistOzone,
27 GetOutcome,
28 InboundOutcome,
29 PollOutcome,
30 },
31 peer_set::PeerSet,
32 placement::Placement,
33 record::{
34 Record,
35 RecordId,
36 },
37 storage::{
38 MemoryStorage,
39 Storage,
40 },
41 transport::{
42 Envelope,
43 MsgKind,
44 },
45};
46
47
48/// Deterministic splitmix64 for reproducible tests.
49struct Rng { state: u64 }
50
51impl Rng {
52 fn new(seed: u64) -> Self { Self { state: seed } }
53
54 fn next_u64(&mut self) -> u64 {
55 self.state = self.state.wrapping_add(0x9E3779B97F4A7C15);
56 let mut z = self.state;
57 z = (z ^ (z >> 30)).wrapping_mul(0xBF58476D1CE4E5B9);
58 z = (z ^ (z >> 27)).wrapping_mul(0x94D049BB133111EB);
59 z ^ (z >> 31)
60 }
61
62 fn next_id(&mut self) -> NodeId {
63 let mut bytes = [0u8; 32];
64 for i in 0..4 {
65 let word = self.next_u64().to_le_bytes();
66 bytes[i * 8..(i + 1) * 8].copy_from_slice(&word);
67 }
68 NodeId::from_bytes(bytes)
69 }
70
71 fn next_record_id(&mut self) -> RecordId {
72 RecordId(*self.next_id().as_bytes())
73 }
74}
75
76
77fn node_id_from_u8(b: u8) -> NodeId {
78 let mut bytes = [0u8; 32];
79 bytes[31] = b;
80 NodeId::from_bytes(bytes)
81}
82
83
84// ---------------------------------------------------------------------------
85// Config tests.
86// ---------------------------------------------------------------------------
87
88#[test]
89fn config_rejects_empty_table_list() -> Outcome<()> {
90 let oam = res!(OamConfig::new(20, 500));
91 let err = DistOzoneConfig::new(
92 node_id_from_u8(1),
93 Vec::new(),
94 oam,
95 Vec::new(),
96 );
97 assert!(err.is_err());
98 Ok(())
99}
100
101#[test]
102fn config_rejects_duplicate_table_names() -> Outcome<()> {
103 let oam = res!(OamConfig::new(20, 500));
104 let tables = vec![
105 res!(TableConfig::eventual("identity")),
106 res!(TableConfig::eventual("identity")),
107 ];
108 let err = DistOzoneConfig::new(node_id_from_u8(1), Vec::new(), oam, tables);
109 assert!(err.is_err());
110 Ok(())
111}
112
113#[test]
114fn config_rejects_invalid_cohort_lambda() -> Outcome<()> {
115 assert!(TableConfig::new(
116 "treasury",
117 Consistency::Cohort { lambda: 4 },
118 TableConfig::DEFAULT_AE,
119 TableConfig::DEFAULT_IBLT_CELLS,
120 ).is_err());
121 assert!(TableConfig::new(
122 "treasury",
123 Consistency::Cohort { lambda: 6 },
124 TableConfig::DEFAULT_AE,
125 TableConfig::DEFAULT_IBLT_CELLS,
126 ).is_err());
127 assert!(TableConfig::new(
128 "treasury",
129 Consistency::Cohort { lambda: 7 },
130 TableConfig::DEFAULT_AE,
131 TableConfig::DEFAULT_IBLT_CELLS,
132 ).is_ok());
133 // Zero cells rejected.
134 assert!(TableConfig::new(
135 "identity",
136 Consistency::Eventual,
137 TableConfig::DEFAULT_AE,
138 0,
139 ).is_err());
140 Ok(())
141}
142
143#[test]
144fn config_rejects_empty_table_name() -> Outcome<()> {
145 assert!(TableConfig::eventual("").is_err());
146 Ok(())
147}
148
149
150// ---------------------------------------------------------------------------
151// Peer set tests.
152// ---------------------------------------------------------------------------
153
154#[test]
155fn peer_set_filters_self_on_bootstrap() -> Outcome<()> {
156 let me = node_id_from_u8(1);
157 let a = node_id_from_u8(2);
158 let b = node_id_from_u8(3);
159 let set = PeerSet::from_bootstrap(&me, vec![a, me, b, a]);
160 // Self removed, duplicates removed, sorted.
161 assert_eq!(set.len(), 2);
162 assert!(!set.contains(&me));
163 assert!(set.contains(&a));
164 assert!(set.contains(&b));
165 Ok(())
166}
167
168#[test]
169fn peer_set_insert_is_idempotent() -> Outcome<()> {
170 let mut set = PeerSet::new();
171 let a = node_id_from_u8(7);
172 assert!(set.insert(a));
173 assert!(!set.insert(a));
174 assert_eq!(set.len(), 1);
175 Ok(())
176}
177
178#[test]
179fn peer_set_is_sorted() -> Outcome<()> {
180 let mut rng = Rng::new(0x5a5a);
181 let mut set = PeerSet::new();
182 for _ in 0..20 {
183 set.insert(rng.next_id());
184 }
185 for pair in set.as_slice().windows(2) {
186 assert!(pair[0].as_bytes() < pair[1].as_bytes(),
187 "peer set not sorted");
188 }
189 Ok(())
190}
191
192
193// ---------------------------------------------------------------------------
194// Placement service tests.
195// ---------------------------------------------------------------------------
196
197#[test]
198fn placement_local_holder_branches_on_threshold() -> Outcome<()> {
199 // Saturated threshold: every peer is a holder.
200 let oam = res!(OamConfig::new(500, 500));
201 let local = node_id_from_u8(1);
202 let p = Placement::new(local, oam);
203 assert!(matches!(p.threshold(), Threshold::All));
204 assert!(p.i_am_holder(&RecordId::from_bytes([0; 32])));
205
206 // Zero threshold: no peer is a holder.
207 let oam = res!(OamConfig::new(0, 500));
208 let p = Placement::new(local, oam);
209 assert!(matches!(p.threshold(), Threshold::None));
210 assert!(!p.i_am_holder(&RecordId::from_bytes([0; 32])));
211 Ok(())
212}
213
214#[test]
215fn placement_update_oam_refreshes_threshold() -> Outcome<()> {
216 let local = node_id_from_u8(1);
217 let mut p = Placement::new(local, res!(OamConfig::new(1, 1_000_000)));
218 let first = *res!(p.threshold().as_bytes().ok_or_else(|| err!(
219 "expected Bounded threshold"; Bug)));
220 p.update_oam(res!(OamConfig::new(1, 2)));
221 let second = *res!(p.threshold().as_bytes().ok_or_else(|| err!(
222 "expected Bounded threshold"; Bug)));
223 assert_ne!(first, second);
224 Ok(())
225}
226
227#[test]
228fn placement_holder_count_matches_decision() -> Outcome<()> {
229 let oam = res!(OamConfig::new(100, 200));
230 let me = node_id_from_u8(1);
231 let p = Placement::new(me, oam);
232 let mut set = PeerSet::new();
233 let mut rng = Rng::new(0xdead);
234 for _ in 0..50 {
235 set.insert(rng.next_id());
236 }
237 let rid = rng.next_record_id();
238 let decision = p.decide(&rid, &set);
239 assert_eq!(
240 decision.holder_count(),
241 decision.remote_holders.len() + usize::from(decision.local_is_holder),
242 );
243 Ok(())
244}
245
246
247// ---------------------------------------------------------------------------
248// DistOzone write-path tests.
249// ---------------------------------------------------------------------------
250
251fn build_engine(
252 local: NodeId,
253 peers: Vec<NodeId>,
254 replication: u64,
255 network_size: u64,
256)
257 -> Outcome<DistOzone<MemoryStorage>>
258{
259 let oam = res!(OamConfig::new(replication, network_size));
260 let tables = vec![
261 res!(TableConfig::eventual("identity")),
262 res!(TableConfig::eventual("escrow")),
263 ];
264 let cfg = res!(DistOzoneConfig::new(local, peers, oam, tables));
265 DistOzone::new(cfg, MemoryStorage::new())
266}
267
268
269#[test]
270fn put_rejects_unknown_table() -> Outcome<()> {
271 let me = node_id_from_u8(1);
272 let engine = res!(build_engine(me, Vec::new(), 20, 500));
273 let rid = RecordId::from_bytes([0; 32]);
274 let record = Record::new(rid, "missing", b"v".to_vec());
275 assert!(engine.put(record).is_err());
276 Ok(())
277}
278
279#[test]
280fn put_on_cohort_table_enters_consensus() -> Outcome<()> {
281 // On a cohort-backed table with no peers, the sole member is the
282 // local node -- which is therefore the leader. `put` opens a HotStuff
283 // round and reports consensus_pending. The record is *not* persisted
284 // yet (that happens on Decide).
285 let me = node_id_from_u8(1);
286 let oam = res!(OamConfig::new(20, 500));
287 let tables = vec![res!(TableConfig::cohort_default("treasury"))];
288 let cfg = res!(DistOzoneConfig::new(me, Vec::new(), oam, tables));
289 let engine = res!(DistOzone::new(cfg, MemoryStorage::new()));
290 let rid = RecordId::from_bytes([0; 32]);
291 let record = Record::new(rid, "treasury", b"v".to_vec());
292 let outcome = res!(engine.put(record.clone()));
293 assert!(!outcome.local_persisted);
294 assert_eq!(outcome.consensus_pending, Some((record.table.clone(), rid)));
295 Ok(())
296}
297
298#[test]
299fn put_with_saturated_threshold_reaches_everyone() -> Outcome<()> {
300 // Replication == network size => everyone holds everything.
301 let me = node_id_from_u8(1);
302 let peers: Vec<NodeId> = (2..=5).map(node_id_from_u8).collect();
303 let engine = res!(build_engine(me, peers.clone(), 5, 5));
304 let rid = RecordId::from_bytes([7; 32]);
305 let record = Record::new(rid, "identity", b"hello".to_vec());
306 let outcome = res!(engine.put(record.clone()));
307 assert!(outcome.local_persisted);
308 assert_eq!(outcome.outbound.len(), peers.len());
309 // Every outbound goes to a distinct remote peer.
310 let mut seen: Vec<NodeId> = outcome.outbound.iter().map(|e| e.to).collect();
311 seen.sort_by(|a, b| a.as_bytes().cmp(b.as_bytes()));
312 seen.dedup();
313 assert_eq!(seen.len(), peers.len());
314 // Every outbound is a ReplicatePut with the right record.
315 for env in &outcome.outbound {
316 assert_eq!(env.from, me);
317 match &env.body {
318 MsgKind::ReplicatePut { record: r } => assert_eq!(r, &record),
319 other => panic!("unexpected outbound: {:?}", other),
320 }
321 }
322 // Local store got the record.
323 let stored = res!(engine.storage().get("identity", &rid));
324 assert_eq!(stored, Some(record));
325 Ok(())
326}
327
328#[test]
329fn put_with_empty_threshold_writes_nowhere() -> Outcome<()> {
330 let me = node_id_from_u8(1);
331 let peers: Vec<NodeId> = (2..=10).map(node_id_from_u8).collect();
332 let engine = res!(build_engine(me, peers, 0, 10)); // n=0: no holders.
333 let record = Record::new(
334 RecordId::from_bytes([9; 32]), "identity", b"x".to_vec(),
335 );
336 let outcome = res!(engine.put(record));
337 assert!(!outcome.local_persisted);
338 assert!(outcome.outbound.is_empty());
339 Ok(())
340}
341
342
343// ---------------------------------------------------------------------------
344// DistOzone read-path tests.
345// ---------------------------------------------------------------------------
346
347#[test]
348fn get_local_returns_record_when_holder() -> Outcome<()> {
349 let me = node_id_from_u8(1);
350 let engine = res!(build_engine(me, Vec::new(), 100, 100));
351 let rid = RecordId::from_bytes([4; 32]);
352 let record = Record::new(rid, "identity", b"v".to_vec());
353 let _ = res!(engine.put(record.clone()));
354 let outcome = res!(engine.get("identity", &rid));
355 match outcome {
356 GetOutcome::Local(r) => assert_eq!(r, record),
357 other => panic!("expected Local, got {:?}", other),
358 }
359 Ok(())
360}
361
362#[test]
363fn get_local_miss_when_holder_but_no_record() -> Outcome<()> {
364 let me = node_id_from_u8(1);
365 let engine = res!(build_engine(me, Vec::new(), 100, 100));
366 let outcome = res!(engine.get("identity", &RecordId::from_bytes([9; 32])));
367 assert!(matches!(outcome, GetOutcome::LocalMiss));
368 Ok(())
369}
370
371#[test]
372fn get_with_zero_threshold_and_peers_routes_remote() -> Outcome<()> {
373 // n = 0 means the local peer is not a holder. With peers known, a remote
374 // read is scheduled.
375 let me = node_id_from_u8(1);
376 let peers: Vec<NodeId> = (2..=10).map(node_id_from_u8).collect();
377 let engine = res!(build_engine(me, peers.clone(), 0, 10));
378 let rid = RecordId::from_bytes([5; 32]);
379 let outcome = res!(engine.get("identity", &rid));
380 match outcome {
381 GetOutcome::Remote { request_id, outbound } => {
382 assert!(request_id > 0);
383 assert!(!outbound.is_empty());
384 assert!(outbound.len() <= peers.len());
385 for env in &outbound {
386 assert_eq!(env.from, me);
387 match &env.body {
388 MsgKind::GetRequest { request_id: rid_msg, table, id } => {
389 assert_eq!(*rid_msg, request_id);
390 assert_eq!(table, "identity");
391 assert_eq!(id, &rid);
392 },
393 other => panic!("unexpected outbound: {:?}", other),
394 }
395 }
396 },
397 other => panic!("expected Remote, got {:?}", other),
398 }
399 Ok(())
400}
401
402#[test]
403fn get_no_peers_and_not_holder_returns_notargets() -> Outcome<()> {
404 let me = node_id_from_u8(1);
405 let engine = res!(build_engine(me, Vec::new(), 0, 10));
406 let rid = RecordId::from_bytes([0; 32]);
407 let outcome = res!(engine.get("identity", &rid));
408 assert!(matches!(outcome, GetOutcome::NoTargets));
409 Ok(())
410}
411
412
413// ---------------------------------------------------------------------------
414// Inbound handling tests.
415// ---------------------------------------------------------------------------
416
417#[test]
418fn handle_replicate_put_persists_when_holder() -> Outcome<()> {
419 let me = node_id_from_u8(1);
420 let engine = res!(build_engine(me, Vec::new(), 100, 100));
421 let sender = node_id_from_u8(2);
422 let record = Record::new(
423 RecordId::from_bytes([3; 32]),
424 "identity",
425 b"payload".to_vec(),
426 );
427 let env = Envelope::new(sender, me, MsgKind::ReplicatePut {
428 record: record.clone(),
429 });
430 let out = res!(engine.handle_envelope(env));
431 assert!(out.outbound.is_empty());
432 let stored = res!(engine.storage().get("identity", &record.id));
433 assert_eq!(stored, Some(record));
434 Ok(())
435}
436
437#[test]
438fn handle_replicate_put_drops_when_not_holder() -> Outcome<()> {
439 // n = 0: not a holder of anything. Incoming put must not be persisted.
440 let me = node_id_from_u8(1);
441 let engine = res!(build_engine(me, Vec::new(), 0, 10));
442 let sender = node_id_from_u8(2);
443 let record = Record::new(
444 RecordId::from_bytes([3; 32]),
445 "identity",
446 b"payload".to_vec(),
447 );
448 let env = Envelope::new(sender, me, MsgKind::ReplicatePut {
449 record: record.clone(),
450 });
451 let out = res!(engine.handle_envelope(env));
452 assert!(out.outbound.is_empty());
453 assert_eq!(res!(engine.storage().len()), 0);
454 Ok(())
455}
456
457#[test]
458fn handle_replicate_put_rejects_unknown_table() -> Outcome<()> {
459 let me = node_id_from_u8(1);
460 let engine = res!(build_engine(me, Vec::new(), 100, 100));
461 let sender = node_id_from_u8(2);
462 let record = Record::new(
463 RecordId::from_bytes([3; 32]),
464 "missing",
465 b"payload".to_vec(),
466 );
467 let env = Envelope::new(sender, me, MsgKind::ReplicatePut { record });
468 assert!(engine.handle_envelope(env).is_err());
469 Ok(())
470}
471
472#[test]
473fn handle_get_request_responds_with_record() -> Outcome<()> {
474 let me = node_id_from_u8(1);
475 let engine = res!(build_engine(me, Vec::new(), 100, 100));
476 let record = Record::new(
477 RecordId::from_bytes([5; 32]),
478 "identity",
479 b"aaa".to_vec(),
480 );
481 let _ = res!(engine.put(record.clone()));
482 let requester = node_id_from_u8(2);
483 let env = Envelope::new(requester, me, MsgKind::GetRequest {
484 request_id: 42,
485 table: "identity".to_string(),
486 id: record.id,
487 });
488 let out = res!(engine.handle_envelope(env));
489 assert_eq!(out.outbound.len(), 1);
490 let reply = &out.outbound[0];
491 assert_eq!(reply.to, requester);
492 match &reply.body {
493 MsgKind::GetResponse { request_id, record: r } => {
494 assert_eq!(*request_id, 42);
495 assert_eq!(r.as_ref(), Some(&record));
496 },
497 other => panic!("expected GetResponse, got {:?}", other),
498 }
499 Ok(())
500}
501
502#[test]
503fn handle_get_response_completes_pending_read() -> Outcome<()> {
504 let me = node_id_from_u8(1);
505 let peers: Vec<NodeId> = (2..=5).map(node_id_from_u8).collect();
506 let mut engine = res!(build_engine(me, peers, 0, 10));
507 engine.set_read_fanout(1);
508
509 let rid = RecordId::from_bytes([7; 32]);
510 let outcome = res!(engine.get("identity", &rid));
511 let (request_id, outbound) = match outcome {
512 GetOutcome::Remote { request_id, outbound } => (request_id, outbound),
513 other => panic!("expected Remote, got {:?}", other),
514 };
515 assert_eq!(outbound.len(), 1);
516
517 // Still pending.
518 assert!(matches!(res!(engine.poll_get(request_id)), PollOutcome::Pending));
519
520 // Target replies with the record.
521 let target = outbound[0].to;
522 let record = Record::new(rid, "identity", b"found".to_vec());
523 let response = Envelope::new(target, me, MsgKind::GetResponse {
524 request_id,
525 record: Some(record.clone()),
526 });
527 let in_out = res!(engine.handle_envelope(response));
528 assert_eq!(in_out.completed_get, Some(request_id));
529
530 match res!(engine.poll_get(request_id)) {
531 PollOutcome::Record(r) => assert_eq!(r, record),
532 other => panic!("expected Record, got {:?}", other),
533 }
534 Ok(())
535}
536
537#[test]
538fn handle_get_response_resolves_notfound_after_all_miss() -> Outcome<()> {
539 let me = node_id_from_u8(1);
540 let peers: Vec<NodeId> = (2..=4).map(node_id_from_u8).collect();
541 let mut engine = res!(build_engine(me, peers, 0, 10));
542 engine.set_read_fanout(3);
543
544 let rid = RecordId::from_bytes([2; 32]);
545 let (request_id, outbound) = match res!(engine.get("identity", &rid)) {
546 GetOutcome::Remote { request_id, outbound } => (request_id, outbound),
547 other => panic!("expected Remote, got {:?}", other),
548 };
549 assert_eq!(outbound.len(), 3);
550
551 // All three targets reply empty.
552 let mut last: Option<InboundOutcome> = None;
553 for env in outbound {
554 let resp = Envelope::new(env.to, me, MsgKind::GetResponse {
555 request_id,
556 record: None,
557 });
558 last = Some(res!(engine.handle_envelope(resp)));
559 }
560 let last = res!(last.ok_or_else(|| err!("no outbound"; Bug)));
561 assert_eq!(last.completed_get, Some(request_id));
562 assert!(matches!(res!(engine.poll_get(request_id)), PollOutcome::NotFound));
563 Ok(())
564}
565
566#[test]
567fn handle_get_response_unknown_request_id_is_ignored() -> Outcome<()> {
568 let me = node_id_from_u8(1);
569 let engine = res!(build_engine(me, Vec::new(), 100, 100));
570 let sender = node_id_from_u8(2);
571 let env = Envelope::new(sender, me, MsgKind::GetResponse {
572 request_id: 9999,
573 record: None,
574 });
575 let out = res!(engine.handle_envelope(env));
576 assert!(out.outbound.is_empty());
577 assert_eq!(out.completed_get, None);
578 Ok(())
579}
580
581#[test]
582fn cancel_get_drops_pending_state() -> Outcome<()> {
583 let me = node_id_from_u8(1);
584 let peers: Vec<NodeId> = (2..=5).map(node_id_from_u8).collect();
585 let engine = res!(build_engine(me, peers, 0, 10));
586 let rid = RecordId::from_bytes([0xaa; 32]);
587 let (request_id, _) = match res!(engine.get("identity", &rid)) {
588 GetOutcome::Remote { request_id, outbound } => (request_id, outbound),
589 other => panic!("expected Remote, got {:?}", other),
590 };
591 res!(engine.cancel_get(request_id));
592 assert!(matches!(res!(engine.poll_get(request_id)), PollOutcome::Unknown));
593 Ok(())
594}
595
596#[test]
597fn handle_envelope_rejects_misaddressed() -> Outcome<()> {
598 let me = node_id_from_u8(1);
599 let engine = res!(build_engine(me, Vec::new(), 100, 100));
600 let sender = node_id_from_u8(2);
601 let other = node_id_from_u8(3);
602 let env = Envelope::new(sender, other, MsgKind::ReplicatePut {
603 record: Record::new(
604 RecordId::from_bytes([0; 32]), "identity", b"v".to_vec(),
605 ),
606 });
607 let out = res!(engine.handle_envelope(env));
608 assert!(out.outbound.is_empty());
609 assert_eq!(out.completed_get, None);
610 // Nothing got stored.
611 assert_eq!(res!(engine.storage().len()), 0);
612 Ok(())
613}
614
615
616// ---------------------------------------------------------------------------
617// Peer-set / OAM mutation.
618// ---------------------------------------------------------------------------
619
620#[test]
621fn insert_peer_rejects_self() -> Outcome<()> {
622 let me = node_id_from_u8(1);
623 let mut engine = res!(build_engine(me, Vec::new(), 20, 500));
624 assert!(!engine.insert_peer(me));
625 assert_eq!(engine.peer_set().len(), 0);
626 Ok(())
627}
628
629#[test]
630fn insert_and_remove_peer_round_trip() -> Outcome<()> {
631 let me = node_id_from_u8(1);
632 let a = node_id_from_u8(2);
633 let mut engine = res!(build_engine(me, Vec::new(), 20, 500));
634 assert!(engine.insert_peer(a));
635 assert!(!engine.insert_peer(a));
636 assert_eq!(engine.peer_set().len(), 1);
637 assert!(engine.remove_peer(&a));
638 assert!(!engine.remove_peer(&a));
639 assert_eq!(engine.peer_set().len(), 0);
640 Ok(())
641}
642
643#[test]
644fn update_network_size_refreshes_threshold() -> Outcome<()> {
645 let me = node_id_from_u8(1);
646 let mut engine = res!(build_engine(me, Vec::new(), 1, 1_000_000));
647 let first = *res!(engine.placement().threshold().as_bytes()
648 .ok_or_else(|| err!("expected Bounded threshold"; Bug)));
649 res!(engine.update_network_size(2));
650 let second = *res!(engine.placement().threshold().as_bytes()
651 .ok_or_else(|| err!("expected Bounded threshold"; Bug)));
652 assert_ne!(first, second);
653 // n=1, N=2: top bit set.
654 assert_eq!(second[0], 0x80);
655 Ok(())
656}
657
658
659// ---------------------------------------------------------------------------
660// End-to-end simulation across two engines.
661// ---------------------------------------------------------------------------
662
663#[test]
664fn two_peer_replicate_and_read_back() -> Outcome<()> {
665 // Peer A (me) writes a record; peer B's engine receives the replicate and
666 // persists it. Then B serves a get request from a third peer (simulated
667 // by us sending a GetRequest to B).
668 let a = node_id_from_u8(1);
669 let b = node_id_from_u8(2);
670
671 let engine_a = res!(build_engine(a, vec![b], 5, 5));
672 let engine_b = res!(build_engine(b, vec![a], 5, 5));
673
674 let rid = RecordId::from_bytes([0x1a; 32]);
675 let record = Record::new(rid, "identity", b"hello".to_vec());
676
677 // A puts.
678 let put_outcome = res!(engine_a.put(record.clone()));
679 assert!(put_outcome.local_persisted);
680 assert_eq!(put_outcome.outbound.len(), 1);
681 let env_to_b = &put_outcome.outbound[0];
682 assert_eq!(env_to_b.to, b);
683
684 // B handles A's replicate put.
685 let in_out = res!(engine_b.handle_envelope(env_to_b.clone()));
686 assert!(in_out.outbound.is_empty());
687 let stored = res!(engine_b.storage().get("identity", &rid));
688 assert_eq!(stored, Some(record.clone()));
689
690 // Simulate a third peer asking B for the record.
691 let requester = node_id_from_u8(99);
692 let ask = Envelope::new(requester, b, MsgKind::GetRequest {
693 request_id: 7,
694 table: "identity".to_string(),
695 id: rid,
696 });
697 let out = res!(engine_b.handle_envelope(ask));
698 assert_eq!(out.outbound.len(), 1);
699 match &out.outbound[0].body {
700 MsgKind::GetResponse { request_id, record: r } => {
701 assert_eq!(*request_id, 7);
702 assert_eq!(r.as_ref(), Some(&record));
703 },
704 other => panic!("expected GetResponse, got {:?}", other),
705 }
706 Ok(())
707}