Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/tests/dist_anti_entropy.rs

13.6 KiB, 23 runs

created by r1870400018:11422, 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 IBLT anti-entropy reconciliation cycle.
3//!
4//! Tests cover:
5//!
6//! - Digest envelope construction and rejection on cohort / unknown tables.
7//! - Digest handling when the peers are already in sync (zero diff).
8//! - Digest handling with a small symmetric difference (both directions).
9//! - Reply -> Push follow-up, exchanging records in the direction the
10//! recipient of the reply lacks them.
11//! - Bulk-reply fallback when the sketch is overloaded.
12//! - End-to-end convergence across a three-message round.
13//!
14//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
15//! Anthropic Claude
16
17use oxedyne_fe2o3_core::prelude::*;
18use oxedyne_fe2o3_o3db_sync::kademlia::id::NodeId;
19use oxedyne_fe2o3_o3db_sync::oam::config::OamConfig;
20use oxedyne_fe2o3_o3db_sync::dist::{
21 config::{
22 Consistency,
23 DistOzoneConfig,
24 TableConfig,
25 },
26 engine::DistOzone,
27 record::{
28 Record,
29 RecordId,
30 },
31 storage::{
32 MemoryStorage,
33 Storage,
34 },
35 transport::{
36 Envelope,
37 MsgKind,
38 },
39};
40
41use std::time::Duration;
42
43
44fn node_id_from_u8(b: u8) -> NodeId {
45 let mut bytes = [0u8; 32];
46 bytes[31] = b;
47 NodeId::from_bytes(bytes)
48}
49
50fn record_id_from_u8(b: u8) -> RecordId {
51 let mut bytes = [0u8; 32];
52 bytes[0] = b;
53 RecordId::from_bytes(bytes)
54}
55
56/// Each considers itself the sole holder of everything in its tables:
57/// replication factor equals network size, so a placement re-check never drops
58/// a record.
59fn build_two_engines(
60 a: NodeId,
61 b: NodeId,
62 table_cells: usize,
63)
64 -> Outcome<(DistOzone<MemoryStorage>, DistOzone<MemoryStorage>)>
65{
66 let oam = res!(OamConfig::new(2, 2));
67 let table = res!(TableConfig::new(
68 "identity",
69 Consistency::Eventual,
70 Duration::from_secs(30),
71 table_cells,
72 ));
73 let cfg_a = res!(DistOzoneConfig::new(
74 a, vec![b], oam, vec![table.clone()],
75 ));
76 let cfg_b = res!(DistOzoneConfig::new(
77 b, vec![a], oam, vec![table],
78 ));
79 Ok((
80 res!(DistOzone::new(cfg_a, MemoryStorage::new())),
81 res!(DistOzone::new(cfg_b, MemoryStorage::new())),
82 ))
83}
84
85
86#[test]
87fn build_anti_entropy_request_rejects_unknown_table() -> Outcome<()> {
88 let (engine_a, _) = res!(build_two_engines(
89 node_id_from_u8(1), node_id_from_u8(2), 64,
90 ));
91 assert!(engine_a.build_anti_entropy_request(
92 "missing", node_id_from_u8(2),
93 ).is_err());
94 Ok(())
95}
96
97#[test]
98fn build_anti_entropy_request_rejects_cohort_table() -> Outcome<()> {
99 let a = node_id_from_u8(1);
100 let b = node_id_from_u8(2);
101 let oam = res!(OamConfig::new(2, 2));
102 let eventual = res!(TableConfig::eventual("identity"));
103 let cohort = res!(TableConfig::cohort_default("treasury"));
104 let cfg = res!(DistOzoneConfig::new(
105 a, vec![b], oam, vec![eventual, cohort],
106 ));
107 let engine = res!(DistOzone::new(cfg, MemoryStorage::new()));
108 assert!(engine.build_anti_entropy_request("treasury", b).is_err());
109 Ok(())
110}
111
112#[test]
113fn anti_entropy_with_identical_tables_decodes_empty() -> Outcome<()> {
114 // Both engines hold the same record; digest exchange must produce an
115 // empty reply.
116 let a = node_id_from_u8(1);
117 let b = node_id_from_u8(2);
118 let (engine_a, engine_b) = res!(build_two_engines(a, b, 64));
119
120 let rid = record_id_from_u8(7);
121 let record = Record::new(rid, "identity", b"same".to_vec());
122 res!(engine_a.storage().put(&record));
123 res!(engine_b.storage().put(&record));
124
125 let digest = res!(engine_a.build_anti_entropy_request("identity", b));
126 let out = res!(engine_b.handle_envelope(digest));
127 assert_eq!(out.outbound.len(), 1);
128 match &out.outbound[0].body {
129 MsgKind::AntiEntropyReply { records, requested_ids, bulk, .. } => {
130 assert!(records.is_empty(), "expected no records in reply");
131 assert!(requested_ids.is_empty(), "expected no requested ids");
132 assert!(!bulk);
133 },
134 other => panic!("expected AntiEntropyReply, got {:?}", other),
135 }
136 Ok(())
137}
138
139#[test]
140fn anti_entropy_sender_has_extra_record_receiver_requests_it() -> Outcome<()> {
141 // A has record X, B does not. A's digest is received by B; B's reply
142 // requests X.
143 let a = node_id_from_u8(1);
144 let b = node_id_from_u8(2);
145 let (engine_a, engine_b) = res!(build_two_engines(a, b, 64));
146
147 let rid = record_id_from_u8(3);
148 let record = Record::new(rid, "identity", b"a-only".to_vec());
149 res!(engine_a.storage().put(&record));
150
151 let digest = res!(engine_a.build_anti_entropy_request("identity", b));
152 let out = res!(engine_b.handle_envelope(digest));
153 assert_eq!(out.outbound.len(), 1);
154 match &out.outbound[0].body {
155 MsgKind::AntiEntropyReply { records, requested_ids, bulk, .. } => {
156 assert!(records.is_empty(),
157 "B shouldn't have anything to send to A");
158 assert_eq!(requested_ids.len(), 1);
159 assert_eq!(requested_ids[0], rid);
160 assert!(!bulk);
161 },
162 other => panic!("expected AntiEntropyReply, got {:?}", other),
163 }
164 Ok(())
165}
166
167#[test]
168fn anti_entropy_receiver_has_extra_record_replies_with_it() -> Outcome<()> {
169 // B has record Y, A does not. A's digest contains nothing new; B's
170 // reply contains Y as records-for-sender.
171 let a = node_id_from_u8(1);
172 let b = node_id_from_u8(2);
173 let (engine_a, engine_b) = res!(build_two_engines(a, b, 64));
174
175 let rid = record_id_from_u8(5);
176 let record = Record::new(rid, "identity", b"b-only".to_vec());
177 res!(engine_b.storage().put(&record));
178
179 let digest = res!(engine_a.build_anti_entropy_request("identity", b));
180 let out = res!(engine_b.handle_envelope(digest));
181 assert_eq!(out.outbound.len(), 1);
182 match &out.outbound[0].body {
183 MsgKind::AntiEntropyReply { records, requested_ids, bulk, .. } => {
184 assert_eq!(records.len(), 1);
185 assert_eq!(records[0], record);
186 assert!(requested_ids.is_empty());
187 assert!(!bulk);
188 },
189 other => panic!("expected AntiEntropyReply, got {:?}", other),
190 }
191 Ok(())
192}
193
194#[test]
195fn anti_entropy_reply_persists_received_records() -> Outcome<()> {
196 // A receives a reply from B containing a record A is missing. A
197 // persists it through handle_envelope.
198 let a = node_id_from_u8(1);
199 let b = node_id_from_u8(2);
200 let (engine_a, _) = res!(build_two_engines(a, b, 64));
201
202 let rid = record_id_from_u8(11);
203 let record = Record::new(rid, "identity", b"payload".to_vec());
204
205 let reply = Envelope::new(b, a, MsgKind::AntiEntropyReply {
206 table: "identity".to_string(),
207 records: vec![record.clone()],
208 requested_ids: Vec::new(),
209 bulk: false,
210 });
211 let out = res!(engine_a.handle_envelope(reply));
212 assert!(out.outbound.is_empty(),
213 "no requested ids => no push follow-up");
214 let stored = res!(engine_a.storage().get("identity", &rid));
215 assert_eq!(stored, Some(record));
216 Ok(())
217}
218
219#[test]
220fn anti_entropy_reply_triggers_push_for_requested_ids() -> Outcome<()> {
221 // A holds a record that B's reply requests; A's handle_envelope on
222 // that reply emits an AntiEntropyPush carrying the record.
223 let a = node_id_from_u8(1);
224 let b = node_id_from_u8(2);
225 let (engine_a, _) = res!(build_two_engines(a, b, 64));
226
227 let rid = record_id_from_u8(13);
228 let record = Record::new(rid, "identity", b"x".to_vec());
229 res!(engine_a.storage().put(&record));
230
231 let reply = Envelope::new(b, a, MsgKind::AntiEntropyReply {
232 table: "identity".to_string(),
233 records: Vec::new(),
234 requested_ids: vec![rid],
235 bulk: false,
236 });
237 let out = res!(engine_a.handle_envelope(reply));
238 assert_eq!(out.outbound.len(), 1);
239 match &out.outbound[0].body {
240 MsgKind::AntiEntropyPush { table, records } => {
241 assert_eq!(table, "identity");
242 assert_eq!(records.len(), 1);
243 assert_eq!(records[0], record);
244 },
245 other => panic!("expected AntiEntropyPush, got {:?}", other),
246 }
247 Ok(())
248}
249
250#[test]
251fn anti_entropy_push_persists_records() -> Outcome<()> {
252 // B receives a Push from A carrying a record B is missing.
253 let a = node_id_from_u8(1);
254 let b = node_id_from_u8(2);
255 let (_, engine_b) = res!(build_two_engines(a, b, 64));
256
257 let rid = record_id_from_u8(17);
258 let record = Record::new(rid, "identity", b"pushed".to_vec());
259 let push = Envelope::new(a, b, MsgKind::AntiEntropyPush {
260 table: "identity".to_string(),
261 records: vec![record.clone()],
262 });
263 let out = res!(engine_b.handle_envelope(push));
264 assert!(out.outbound.is_empty());
265 assert_eq!(
266 res!(engine_b.storage().get("identity", &rid)),
267 Some(record),
268 );
269 Ok(())
270}
271
272#[test]
273fn anti_entropy_overload_triggers_bulk_reply() -> Outcome<()> {
274 // With a tiny sketch (num_cells = 6) and many records the IBLT decode
275 // is overloaded. B must fall back to a bulk reply, including every
276 // record it holds for the table.
277 let a = node_id_from_u8(1);
278 let b = node_id_from_u8(2);
279 let (engine_a, engine_b) = res!(build_two_engines(a, b, 6));
280
281 // Populate B with 50 records that A has none of.
282 for i in 0..50u8 {
283 let rid = record_id_from_u8(i);
284 let record = Record::new(
285 rid, "identity", vec![i, i ^ 0xa5],
286 );
287 res!(engine_b.storage().put(&record));
288 }
289
290 let digest = res!(engine_a.build_anti_entropy_request("identity", b));
291 let out = res!(engine_b.handle_envelope(digest));
292 assert_eq!(out.outbound.len(), 1);
293 match &out.outbound[0].body {
294 MsgKind::AntiEntropyReply { records, requested_ids, bulk, .. } => {
295 assert!(*bulk, "expected bulk fallback on overloaded sketch");
296 assert_eq!(records.len(), 50);
297 assert!(requested_ids.is_empty());
298 },
299 other => panic!("expected AntiEntropyReply, got {:?}", other),
300 }
301 Ok(())
302}
303
304#[test]
305fn anti_entropy_full_round_converges_two_peers() -> Outcome<()> {
306 // End-to-end: A has records {1, 2, 3}; B has records {2, 3, 4}. One
307 // anti-entropy round leaves both peers holding {1, 2, 3, 4}.
308 let a = node_id_from_u8(1);
309 let b = node_id_from_u8(2);
310 let (engine_a, engine_b) = res!(build_two_engines(a, b, 64));
311
312 let rec = |i: u8| -> Record {
313 Record::new(
314 record_id_from_u8(i),
315 "identity",
316 vec![i],
317 )
318 };
319 for i in [1u8, 2, 3] {
320 res!(engine_a.storage().put(&rec(i)));
321 }
322 for i in [2u8, 3, 4] {
323 res!(engine_b.storage().put(&rec(i)));
324 }
325
326 // Round 1: A -> B digest. B replies.
327 let digest = res!(engine_a.build_anti_entropy_request("identity", b));
328 let reply_out = res!(engine_b.handle_envelope(digest));
329 assert_eq!(reply_out.outbound.len(), 1);
330
331 // A handles reply; should persist rec(4) and push rec(1).
332 let reply_env = reply_out.outbound[0].clone();
333 let push_out = res!(engine_a.handle_envelope(reply_env));
334 assert_eq!(push_out.outbound.len(), 1, "expected a push for rec(1)");
335
336 // B handles the push; should persist rec(1).
337 let push_env = push_out.outbound[0].clone();
338 let tail_out = res!(engine_b.handle_envelope(push_env));
339 assert!(tail_out.outbound.is_empty());
340
341 // Final state: both hold {1, 2, 3, 4}.
342 for i in 1u8..=4 {
343 let rid = record_id_from_u8(i);
344 assert_eq!(
345 res!(engine_a.storage().get("identity", &rid)),
346 Some(rec(i)),
347 "engine_a missing record {}", i,
348 );
349 assert_eq!(
350 res!(engine_b.storage().get("identity", &rid)),
351 Some(rec(i)),
352 "engine_b missing record {}", i,
353 );
354 }
355 Ok(())
356}
357
358#[test]
359fn anti_entropy_digest_rejects_cohort_table_on_receive() -> Outcome<()> {
360 // Malicious or misconfigured sender directs a digest at a cohort-
361 // backed table; handler rejects.
362 let a = node_id_from_u8(1);
363 let b = node_id_from_u8(2);
364 let oam = res!(OamConfig::new(2, 2));
365 let eventual = res!(TableConfig::eventual("identity"));
366 let cohort = res!(TableConfig::cohort_default("treasury"));
367 let cfg_b = res!(DistOzoneConfig::new(
368 b, vec![a], oam, vec![eventual, cohort],
369 ));
370 let engine_b = res!(DistOzone::new(cfg_b, MemoryStorage::new()));
371
372 // Hand-built envelope with an empty sketch -- we won't get that far.
373 let fake_digest = Envelope::new(a, b, MsgKind::AntiEntropyDigest {
374 table: "treasury".to_string(),
375 sketch: vec![0u8; 100],
376 });
377 assert!(engine_b.handle_envelope(fake_digest).is_err());
378 Ok(())
379}
380
381#[test]
382fn anti_entropy_digest_rejects_unknown_table() -> Outcome<()> {
383 let a = node_id_from_u8(1);
384 let b = node_id_from_u8(2);
385 let (_, engine_b) = res!(build_two_engines(a, b, 64));
386 let fake = Envelope::new(a, b, MsgKind::AntiEntropyDigest {
387 table: "missing".to_string(),
388 sketch: vec![0u8; 100],
389 });
390 assert!(engine_b.handle_envelope(fake).is_err());
391 Ok(())
392}
393
394#[test]
395fn anti_entropy_digest_rejects_mismatched_sketch_shape() -> Outcome<()> {
396 // Wire up two engines with differing iblt_cells for the same table.
397 // A builds a digest; B's handler detects the config mismatch and
398 // errors rather than silently decoding nonsense.
399 let a = node_id_from_u8(1);
400 let b = node_id_from_u8(2);
401 let oam = res!(OamConfig::new(2, 2));
402
403 let table_a = res!(TableConfig::new(
404 "identity",
405 Consistency::Eventual,
406 Duration::from_secs(30),
407 64,
408 ));
409 let table_b = res!(TableConfig::new(
410 "identity",
411 Consistency::Eventual,
412 Duration::from_secs(30),
413 128,
414 ));
415 let cfg_a = res!(DistOzoneConfig::new(
416 a, vec![b], oam, vec![table_a],
417 ));
418 let cfg_b = res!(DistOzoneConfig::new(
419 b, vec![a], oam, vec![table_b],
420 ));
421 let engine_a = res!(DistOzone::new(cfg_a, MemoryStorage::new()));
422 let engine_b = res!(DistOzone::new(cfg_b, MemoryStorage::new()));
423
424 let digest = res!(engine_a.build_anti_entropy_request("identity", b));
425 assert!(engine_b.handle_envelope(digest).is_err());
426 Ok(())
427}
428
429#[test]
430fn anti_entropy_reply_placement_rechecks_incoming() -> Outcome<()> {
431 // Receiver-side placement re-check: a reply delivering a record to a
432 // peer that does not consider itself a holder must drop the record
433 // rather than persisting it.
434 let a = node_id_from_u8(1);
435 let b = node_id_from_u8(2);
436 let oam = res!(OamConfig::new(0, 2)); // n=0: no peer is ever a holder.
437 let table = res!(TableConfig::eventual("identity"));
438 let cfg_b = res!(DistOzoneConfig::new(
439 b, vec![a], oam, vec![table],
440 ));
441 let engine_b = res!(DistOzone::new(cfg_b, MemoryStorage::new()));
442
443 let record = Record::new(
444 record_id_from_u8(99),
445 "identity",
446 b"nope".to_vec(),
447 );
448 let reply = Envelope::new(a, b, MsgKind::AntiEntropyReply {
449 table: "identity".to_string(),
450 records: vec![record.clone()],
451 requested_ids: Vec::new(),
452 bulk: false,
453 });
454 let _ = res!(engine_b.handle_envelope(reply));
455 assert_eq!(res!(engine_b.storage().len()), 0);
456 Ok(())
457}