Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_ore/src/sync/tests.rs

39.4 KiB, 91 runs

created by r1870400018:19659, 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//! Convergence, which is the only thing a sync protocol is for.
2//!
3//! Every test here ends the same way: two logs that had diverged hold the same
4//! operations and the same frontier. The divergences differ -- a clone, two
5//! histories with nothing in common, one side far ahead, both sides a little
6//! ahead -- and so do the modes, but the assertion does not.
7//!
8//! The pipe carries bytes and not values. Every message is encoded where it is
9//! sent and decoded where it arrives, so the codec is exercised by every test
10//! and the byte counts the mode comparison rests on are the bytes that would
11//! cross a wire.
12//!
13//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
14//! Anthropic Claude
15
16use crate::id::{
17 OpId,
18 ReplicaId,
19};
20use crate::log::OpLog;
21use crate::op::{
22 Header,
23 Op,
24 Record,
25};
26use crate::sync::msg::Message;
27use crate::sync::session::{
28 Growth,
29 Mode,
30 Session,
31 Step,
32 FANOUT,
33};
34use crate::sync::sketch::{
35 Fallback,
36 MIN_CELLS,
37};
38
39use oxedyne_fe2o3_core::prelude::*;
40
41use std::collections::BTreeSet;
42
43
44/// A small linear congruential generator, so a failure can be reproduced.
45struct Rng(u64);
46
47impl Rng {
48 fn new(seed: u64) -> Self {
49 Self(seed ^ 0x9e37_79b9_7f4a_7c15)
50 }
51
52 fn next(&mut self) -> usize {
53 self.0 = self.0
54 .wrapping_mul(6_364_136_223_846_793_005)
55 .wrapping_add(1_442_695_040_888_963_407);
56 (self.0 >> 33) as usize
57 }
58
59 fn below(&mut self, n: usize) -> usize {
60 if n == 0 { 0 } else { self.next() % n }
61 }
62}
63
64
65/// What one exchange cost and how it went.
66#[derive(Debug, Default)]
67struct Tally {
68 bytes: usize, // on the wire, both directions
69 messages: usize, // on the wire, both directions
70 ops: usize, // handed over, both directions
71 fell_back: Option<Fallback>, // either side falling back to the walk
72}
73
74
75/// Both peers open at once, which is the harder case: neither has heard anything
76/// when it works out what to say. A caller whose transport has a caller and a
77/// callee simply does not call `open` on the callee's side.
78fn exchange(a: &mut OpLog, b: &mut OpLog, mode: Mode)
79 -> Outcome<Tally>
80{
81 let mut sa = Session::new(mode);
82 let mut sb = Session::new(mode);
83 let mut tally = Tally::default();
84 // Each queue holds what is in flight towards that peer, as bytes.
85 let mut to_a: Vec<Vec<u8>> = Vec::new();
86 let mut to_b: Vec<Vec<u8>> = Vec::new();
87 let put = |q: &mut Vec<Vec<u8>>, t: &mut Tally, msg: &Message| -> Outcome<()> {
88 let bytes = res!(msg.encode());
89 t.bytes += bytes.len();
90 t.messages += 1;
91 t.ops += msg.entries().len();
92 q.push(bytes);
93 Ok(())
94 };
95 res!(put(&mut to_b, &mut tally, &res!(sa.open(a))));
96 res!(put(&mut to_a, &mut tally, &res!(sb.open(b))));
97 let mut guard = 0usize;
98 while !(sa.is_converged() && sb.is_converged()) {
99 guard += 1;
100 if guard > 64 {
101 return Err(err!(
102 "Two sessions exchanged {} times without converging.", guard;
103 Test, Excessive));
104 }
105 if to_a.is_empty() && to_b.is_empty() {
106 return Err(err!(
107 "The pipe emptied with neither side converged: a is {}, b is {}.",
108 sa.is_converged(), sb.is_converged();
109 Test, Missing));
110 }
111 for bytes in std::mem::take(&mut to_a) {
112 let turn = res!(sa.receive(a, res!(Message::decode(&bytes))));
113 if let Step::FellBack(reason) = turn.step {
114 tally.fell_back = Some(reason);
115 }
116 for msg in &turn.send {
117 res!(put(&mut to_b, &mut tally, msg));
118 }
119 }
120 for bytes in std::mem::take(&mut to_b) {
121 let turn = res!(sb.receive(b, res!(Message::decode(&bytes))));
122 if let Step::FellBack(reason) = turn.step {
123 tally.fell_back = Some(reason);
124 }
125 for msg in &turn.send {
126 res!(put(&mut to_a, &mut tally, msg));
127 }
128 }
129 }
130 Ok(tally)
131}
132
133/// Asserts that two logs hold the same operations and the same frontier.
134fn agree(a: &OpLog, b: &OpLog)
135 -> Outcome<()>
136{
137 let mut ids_a: Vec<OpId> = a.iter().map(|rec| rec.id()).collect();
138 let mut ids_b: Vec<OpId> = b.iter().map(|rec| rec.id()).collect();
139 ids_a.sort();
140 ids_b.sort();
141 if ids_a != ids_b {
142 let missing: Vec<OpId> = ids_a.iter().filter(|id| !ids_b.contains(id)).copied().collect();
143 let extra: Vec<OpId> = ids_b.iter().filter(|id| !ids_a.contains(id)).copied().collect();
144 return Err(err!(
145 "The logs hold {} and {} operations; {:?} is absent from the second and \
146 {:?} from the first.", ids_a.len(), ids_b.len(), missing, extra;
147 Test, Mismatch));
148 }
149 if a.frontier() != b.frontier() {
150 return Err(err!(
151 "The logs agree on their operations and not on their frontiers: {:?} \
152 against {:?}.", a.frontier(), b.frontier();
153 Test, Mismatch));
154 }
155 // And every record is the same record, not merely the same name.
156 for id in &ids_a {
157 if a.get(id) != b.get(id) {
158 return Err(err!(
159 "The logs hold different records for {}.", id; Test, Mismatch));
160 }
161 }
162 Ok(())
163}
164
165fn write(log: &mut OpLog, replica: u64, n: usize, tag: &str)
166 -> Outcome<()>
167{
168 let r = ReplicaId::new(replica);
169 for i in 0..n {
170 res!(log.author(r, Op::Mark { name: fmt!("{}{}", tag, i), body: None, time: None }));
171 }
172 Ok(())
173}
174
175/// An operation naming the whole frontier, which is a merge wherever the frontier
176/// is wider than one.
177fn merge(log: &mut OpLog, replica: u64, tag: &str)
178 -> Outcome<()>
179{
180 res!(log.author(ReplicaId::new(replica), Op::Mark { name: fmt!("{}", tag), body: None, time: None }));
181 Ok(())
182}
183
184/// A shared prefix, then each side writes and merges: the everyday divergence.
185fn diverged(prefix: usize, left: usize, right: usize)
186 -> Outcome<(OpLog, OpLog)>
187{
188 let mut a = OpLog::new();
189 res!(write(&mut a, 1, prefix, "shared"));
190 let mut b = a.clone();
191 res!(write(&mut a, 1, left, "a"));
192 res!(write(&mut b, 2, right, "b"));
193 // A merge on each side, so neither history is a straight line.
194 if left > 1 {
195 res!(a.append(Record::new(
196 res!(Header::new(a.next_id(ReplicaId::new(3)), a.frontier())),
197 Op::Mark { name: fmt!("merge-a"), body: None, time: None },
198 )));
199 }
200 if right > 1 {
201 res!(b.append(Record::new(
202 res!(Header::new(b.next_id(ReplicaId::new(4)), b.frontier())),
203 Op::Mark { name: fmt!("merge-b"), body: None, time: None },
204 )));
205 }
206 Ok((a, b))
207}
208
209
210#[test]
211fn a_fresh_peer_takes_the_whole_history() -> Outcome<()> {
212 for mode in [Mode::Walk, Mode::sketch(64)] {
213 let mut a = OpLog::new();
214 res!(write(&mut a, 1, 12, "x"));
215 res!(merge(&mut a, 2, "join"));
216 let mut b = OpLog::new();
217 let tally = res!(exchange(&mut a, &mut b, mode));
218 res!(agree(&a, &b));
219 assert_eq!(b.len(), 13, "{:?}", mode);
220 assert_eq!(tally.ops, 13, "exactly the history, once, under {:?}", mode);
221 }
222 Ok(())
223}
224
225#[test]
226fn disjoint_histories_join() -> Outcome<()> {
227 for mode in [Mode::Walk, Mode::sketch(64)] {
228 let mut a = OpLog::new();
229 res!(write(&mut a, 1, 7, "a"));
230 let mut b = OpLog::new();
231 res!(write(&mut b, 2, 9, "b"));
232 res!(exchange(&mut a, &mut b, mode));
233 res!(agree(&a, &b));
234 assert_eq!(a.len(), 16, "{:?}", mode);
235 // Two roots and no merge between them, so the frontier is wide.
236 assert_eq!(a.frontier().len(), 2);
237 }
238 Ok(())
239}
240
241/// The peer that is behind is the loose direction: it cannot subtract a head it
242/// has never seen, so it offers its whole log back, all five operations of which
243/// are dropped on arrival. That is the walk's cost, and the reason for the other
244/// mode.
245#[test]
246fn one_sided_divergence_converges() -> Outcome<()> {
247 let mut a = OpLog::new();
248 res!(write(&mut a, 1, 5, "shared"));
249 let mut b = a.clone();
250 res!(write(&mut a, 1, 60, "ahead"));
251 res!(merge(&mut a, 3, "tip"));
252 let before = b.len();
253 let mut sa = Session::new(Mode::Walk);
254 let mut sb = Session::new(Mode::Walk);
255 let mut to_a = vec![res!(sb.open(&b))];
256 let mut to_b = vec![res!(sa.open(&a))];
257 let mut guard = 0usize;
258 while !(sa.is_converged() && sb.is_converged()) {
259 guard += 1;
260 assert!(guard < 16);
261 for msg in std::mem::take(&mut to_a) {
262 to_b.extend(res!(sa.receive(&mut a, msg)).send);
263 }
264 for msg in std::mem::take(&mut to_b) {
265 to_a.extend(res!(sb.receive(&mut b, msg)).send);
266 }
267 }
268 res!(agree(&a, &b));
269 assert_eq!(b.len(), before + 61);
270 assert_eq!(sa.ops_sent(), 61, "the news and nothing else");
271 assert_eq!(sb.ops_sent(), 5, "and the shared prefix, loosely, the other way");
272 assert_eq!(sa.ops_absorbed(), 0, "every one of which was already held");
273 Ok(())
274}
275
276/// One clone, cut into sessions by a carrier that bounds its reply, with `carry`
277/// saying whether one session hands the next what it learned.
278///
279/// The bound is the relay's: what does not fit is not sent and not remembered,
280/// and `Done` claims only that this end will send no more this turn, so the
281/// puller opens again. An append-order prefix of an owed set is causally closed
282/// by construction, which is what makes cutting one safe.
283fn clone_in_sessions(a: &mut OpLog, cut: usize, carry: bool)
284 -> Outcome<(usize, usize, OpLog)>
285{
286 let mut b = OpLog::new();
287 let mut known: BTreeSet<OpId> = BTreeSet::new();
288 let mut offered = 0usize;
289 let mut sessions = 0usize;
290 while b.len() < a.len() {
291 sessions += 1;
292 if sessions > 32 {
293 return Err(err!(
294 "The sessions stopped making progress, at {} of {}.", b.len(), a.len();
295 Test, Excessive));
296 }
297 let mut sa = Session::new(Mode::Walk);
298 let mut sb = if carry {
299 Session::knowing(Mode::Walk, known.clone())
300 } else {
301 Session::new(Mode::Walk)
302 };
303 let hello = res!(sb.open(&b));
304 let mut room = cut;
305 for msg in res!(sa.receive(a, hello)).send {
306 let msg = match msg {
307 Message::Send { entries } => {
308 let take = std::cmp::min(room, entries.len());
309 room -= take;
310 Message::Send { entries: entries.into_iter().take(take).collect() }
311 },
312 other => other,
313 };
314 for out in &res!(sb.receive(&mut b, msg)).send {
315 offered += out.entries().len();
316 }
317 }
318 if !sb.is_converged() {
319 return Err(err!(
320 "A session ended unconverged after {} of them.", sessions; Test, Missing));
321 }
322 known = sb.known().clone();
323 }
324 Ok((sessions, offered, b))
325}
326
327/// **The defect this exists for.** A relay bounds its reply, so a clone larger
328/// than the bound is a run of sessions, and every session that began again from
329/// the frontier alone offered back the whole prefix the sessions before it had
330/// just delivered. Nothing in a frontier says "you handed me this": a head the
331/// puller has never seen subtracts nothing, so its own log looked entirely owed.
332///
333/// Measured on fe2o3 on 12026-08-22, 35,314 operations over sixteen sessions:
334/// 166,224 operations offered, every one of them already at the far end, 717 MB
335/// up against 87 MB down. Here the same shape at four sessions, with the loose
336/// run kept beside the tight one so the cost has a number rather than a claim.
337#[test]
338fn a_bounded_clone_offers_back_nothing_it_was_given() -> Outcome<()> {
339 let mut a = OpLog::new();
340 res!(write(&mut a, 1, 40, "x"));
341
342 let (sessions, offered, b) = res!(clone_in_sessions(&mut a, 10, true));
343 res!(agree(&a, &b));
344 assert!(sessions > 1, "the bound was never reached, so nothing here was exercised");
345 assert_eq!(offered, 0, "a clone owes nothing and offered {} operations", offered);
346
347 // The same clone, with each session starting from the frontier alone.
348 let (loose_sessions, loose, b) = res!(clone_in_sessions(&mut a, 10, false));
349 res!(agree(&a, &b));
350 assert_eq!(loose_sessions, sessions, "remembering changed how many sessions it took");
351 // Ten, then twenty, then thirty: the triangular number the fe2o3 clone paid
352 // 717 MB of.
353 assert_eq!(loose, 60, "the loose run's cost is not what it was");
354 Ok(())
355}
356
357/// The result does not depend on the mode.
358#[test]
359fn a_symmetric_divergence_converges_either_way() -> Outcome<()> {
360 let (mut wa, mut wb) = res!(diverged(20, 4, 5));
361 res!(exchange(&mut wa, &mut wb, Mode::Walk));
362 res!(agree(&wa, &wb));
363 let (mut sa, mut sb) = res!(diverged(20, 4, 5));
364 let tally = res!(exchange(&mut sa, &mut sb, Mode::sketch(16)));
365 res!(agree(&sa, &sb));
366 assert!(tally.fell_back.is_none(), "an estimate of 16 for a difference of 11");
367 // The two modes reach the same place.
368 res!(agree(&wa, &sa));
369 assert_eq!(wa.len(), 31);
370 Ok(())
371}
372
373/// One exchange, and no operations.
374#[test]
375fn logs_that_already_agree_send_nothing() -> Outcome<()> {
376 for mode in [Mode::Walk, Mode::sketch(4)] {
377 let mut a = OpLog::new();
378 res!(write(&mut a, 1, 10, "x"));
379 let mut b = a.clone();
380 let tally = res!(exchange(&mut a, &mut b, mode));
381 res!(agree(&a, &b));
382 assert_eq!(tally.ops, 0, "{:?}", mode);
383 assert_eq!(tally.messages, 4, "an opening and a done each, under {:?}", mode);
384 }
385 Ok(())
386}
387
388/// The degenerate case must not be a special one.
389#[test]
390fn two_empty_logs_converge() -> Outcome<()> {
391 for mode in [Mode::Walk, Mode::sketch(0)] {
392 let mut a = OpLog::new();
393 let mut b = OpLog::new();
394 let tally = res!(exchange(&mut a, &mut b, mode));
395 res!(agree(&a, &b));
396 assert!(a.is_empty());
397 assert_eq!(tally.ops, 0);
398 }
399 Ok(())
400}
401
402/// The walk answers in the same turn.
403#[test]
404fn an_undersized_sketch_falls_back_and_still_converges() -> Outcome<()> {
405 // Two hundred apiece with nothing in common: four hundred of difference,
406 // sketched as though there were none.
407 let mut a = OpLog::new();
408 res!(write(&mut a, 1, 200, "a"));
409 let mut b = OpLog::new();
410 res!(write(&mut b, 2, 200, "b"));
411 let tally = res!(exchange(&mut a, &mut b, Mode::sketch(0)));
412 res!(agree(&a, &b));
413 assert_eq!(a.len(), 400);
414 match tally.fell_back {
415 Some(Fallback::Incomplete { remaining, .. }) => assert!(remaining > 0),
416 Some(other) => return Err(err!(
417 "Expected a stalled decode, got {}.", other.why(); Test, Mismatch)),
418 None => return Err(err!(
419 "A sketch of sixteen cells decoded a difference of four hundred.";
420 Test, Mismatch)),
421 }
422 Ok(())
423}
424
425/// And a stall grows before it falls back, and then gives the walk the turn.
426///
427/// The order is the point. A table that stalls says the estimate was low and
428/// says nothing about by how much, so doubling it is tried before a frontier is
429/// walked -- against a peer whose head this end does not hold, the walk's owed
430/// set is the whole log. What stops the climb is not a count of growths: every
431/// growth at least doubles, and it ends at whichever of `GROW_CELLS` and the size
432/// [`Mode::between`] would not have chosen -- a table sized for as much as the
433/// smaller log -- comes first. Two logs that share nothing reach the second of
434/// those in a handful of turns, which is this test.
435#[test]
436fn a_fallback_is_reported_and_remembered() -> Outcome<()> {
437 let mut a = OpLog::new();
438 res!(write(&mut a, 1, 120, "a"));
439 let mut b = OpLog::new();
440 res!(write(&mut b, 2, 120, "b"));
441 let mut sa = Session::new(Mode::sketch(0));
442 let mut sb = Session::new(Mode::sketch(0));
443 // B is called upon: it never opens on its own account, and answers what it
444 // is sent.
445 let opening = res!(sa.open(&a));
446 let turn = res!(sb.receive(&mut b, res!(Message::decode(&res!(opening.encode())))));
447 // Sixteen cells against a difference of two hundred and forty: the first
448 // answer is a bigger table and not the walk, and nothing is handed over with
449 // it.
450 match turn.step {
451 Step::Grew { cells } => assert!(cells > MIN_CELLS,
452 "it grew to {} cells, which is no larger than the table it answered",
453 cells),
454 other => return Err(err!("The first answer was {:?}.", other; Test, Mismatch)),
455 }
456 assert_eq!(turn.send.len(), 1, "a turn that grew handed something over as well");
457 assert!(sb.fell_back().is_none(), "growing was recorded as falling back");
458 assert!(!sb.is_converged());
459 // The rest of the exchange, by hand, so that the sticky flag can be read at
460 // the end. Both ends grow once, both stall again, and the walk answers.
461 let mut queue = turn.send;
462 let mut guard = 0usize;
463 while !(sa.is_converged() && sb.is_converged()) {
464 guard += 1;
465 if guard > 16 {
466 return Err(err!(
467 "The exchange took {} turns to converge.", guard; Test, Excessive));
468 }
469 let mut back = Vec::new();
470 for msg in std::mem::take(&mut queue) {
471 for out in res!(sa.receive(&mut a, msg)).send {
472 back.push(out);
473 }
474 }
475 for msg in back {
476 for out in res!(sb.receive(&mut b, msg)).send {
477 queue.push(out);
478 }
479 }
480 }
481 assert!(sa.is_converged());
482 assert!(sb.is_converged());
483 assert!(sb.fell_back().is_some(), "the fallback is remembered past convergence");
484 assert!(sa.fell_back().is_some(), "and both sides made one");
485 res!(agree(&a, &b));
486 Ok(())
487}
488
489/// Everything that arrives closes causally against what the receiver already
490/// holds. This is the property the whole protocol exists to preserve, so it is
491/// asserted against the receiver's log at the moment of arrival rather than
492/// inferred from the outcome.
493#[test]
494fn every_batch_closes_on_arrival() -> Outcome<()> {
495 use crate::sync::walk::arrival_gap;
496
497 for mode in [Mode::Walk, Mode::sketch(8)] {
498 let (mut a, mut b) = res!(diverged(15, 6, 4));
499 let mut sa = Session::new(mode);
500 let mut sb = Session::new(mode);
501 let mut to_a = vec![res!(sb.open(&b))];
502 let mut to_b = vec![res!(sa.open(&a))];
503 let mut checked = 0usize;
504 let mut guard = 0usize;
505 while !(sa.is_converged() && sb.is_converged()) {
506 guard += 1;
507 assert!(guard < 64, "no convergence under {:?}", mode);
508 for msg in std::mem::take(&mut to_a) {
509 if let Message::Send { entries } = &msg {
510 assert!(!entries.is_empty());
511 assert_eq!(
512 res!(arrival_gap(&a, entries)), None,
513 "a batch arriving at a under {:?} has a hole", mode,
514 );
515 checked += 1;
516 }
517 to_b.extend(res!(sa.receive(&mut a, msg)).send);
518 }
519 for msg in std::mem::take(&mut to_b) {
520 if let Message::Send { entries } = &msg {
521 assert_eq!(
522 res!(arrival_gap(&b, entries)), None,
523 "a batch arriving at b under {:?} has a hole", mode,
524 );
525 checked += 1;
526 }
527 to_a.extend(res!(sb.receive(&mut b, msg)).send);
528 }
529 }
530 assert_eq!(checked, 2, "one batch each way under {:?}", mode);
531 res!(agree(&a, &b));
532 }
533 Ok(())
534}
535
536/// Refused whole: the log is untouched.
537#[test]
538fn a_batch_with_a_hole_is_refused() -> Outcome<()> {
539 let mut a = OpLog::new();
540 res!(write(&mut a, 1, 6, "a"));
541 let mut b = OpLog::new();
542 let mut sa = Session::new(Mode::Walk);
543 let mut sb = Session::new(Mode::Walk);
544 let opening = res!(sb.open(&b));
545 let turn = res!(sa.receive(&mut a, opening));
546 // The batch a would have sent, with its first operation taken out.
547 let mut holed = None;
548 for msg in turn.send {
549 if let Message::Send { entries } = msg {
550 holed = Some(Message::Send { entries: entries[1..].to_vec() });
551 }
552 }
553 let holed = match holed {
554 Some(m) => m,
555 None => return Err(err!("No batch was sent to hole."; Test, Missing)),
556 };
557 let e = match sb.receive(&mut b, holed) {
558 Ok(_) => return Err(err!("A batch with a hole was absorbed."; Test, Mismatch)),
559 Err(e) => e,
560 };
561 let text = fmt!("{}", e);
562 assert!(text.contains("closures and not subsets"), "message was {}", text);
563 assert!(b.is_empty(), "nothing was absorbed");
564 Ok(())
565}
566
567/// Dropped rather than refused, which is what makes a loose owed set cost bytes
568/// and not correctness.
569#[test]
570fn repeated_operations_are_dropped() -> Outcome<()> {
571 let mut a = OpLog::new();
572 res!(write(&mut a, 1, 5, "a"));
573 let mut b = a.clone();
574 let entries: Vec<crate::segment::Entry> = a.iter()
575 .map(|rec| crate::segment::Entry::Bare(rec.clone()))
576 .collect();
577 let mut s = Session::new(Mode::Walk);
578 // The same batch twice over, to catch a repetition within one message as
579 // well as one the log already holds.
580 let mut twice = entries.clone();
581 twice.extend(entries);
582 let turn = res!(s.receive(&mut b, Message::Send { entries: twice }));
583 assert_eq!(turn.step, Step::NeedMore);
584 assert_eq!(s.ops_absorbed(), 0);
585 assert_eq!(b.len(), 5, "the log is as it was");
586 // And into a log that holds none of it, a repeated batch places each once.
587 let mut fresh = OpLog::new();
588 let mut once: Vec<crate::segment::Entry> = a.iter()
589 .map(|rec| crate::segment::Entry::Bare(rec.clone()))
590 .collect();
591 once.extend(once.clone());
592 let mut s = Session::new(Mode::Walk);
593 res!(s.receive(&mut fresh, Message::Send { entries: once }));
594 assert_eq!(fresh.len(), 5);
595 assert_eq!(s.ops_absorbed(), 5);
596 Ok(())
597}
598
599/// That ratio is the whole reason the sketch mode exists.
600///
601/// Two logs of 204 operations differing by eight, both peers speaking: the walk
602/// spends 29,608 bytes and moves 408 operations, the sketch spends 2,142 and
603/// moves 8. Tripling the shared history leaves the sketch at 2,142 exactly and
604/// would triple the walk.
605#[test]
606fn the_sketch_costs_the_difference_and_the_walk_costs_the_log() -> Outcome<()> {
607 // Two hundred shared operations, three new on each side.
608 let (mut wa, mut wb) = res!(diverged(200, 3, 3));
609 let walk = res!(exchange(&mut wa, &mut wb, Mode::Walk));
610 res!(agree(&wa, &wb));
611 let (mut sa, mut sb) = res!(diverged(200, 3, 3));
612 let sketch = res!(exchange(&mut sa, &mut sb, Mode::sketch(16)));
613 res!(agree(&sa, &sb));
614 assert!(sketch.fell_back.is_none());
615 // The walk cannot subtract a head it has never seen, so each side sends its
616 // whole log; the sketch sends the difference.
617 assert_eq!(walk.ops, 2 * 204, "each side sent everything it holds");
618 assert_eq!(sketch.ops, 8, "four operations each way");
619 assert!(
620 sketch.bytes * 8 < walk.bytes,
621 "the sketch cost {} bytes and the walk {}: not the ratio the mode is for",
622 sketch.bytes, walk.bytes,
623 );
624 // And the sketch's own cost does not grow with the history: the same
625 // divergence over a longer shared prefix costs the same.
626 let (mut la, mut lb) = res!(diverged(600, 3, 3));
627 let longer = res!(exchange(&mut la, &mut lb, Mode::sketch(16)));
628 res!(agree(&la, &lb));
629 assert_eq!(longer.ops, sketch.ops);
630 assert!(
631 longer.bytes < sketch.bytes + 64,
632 "a history three times longer cost {} bytes against {}",
633 longer.bytes, sketch.bytes,
634 );
635 Ok(())
636}
637
638#[test]
639fn a_soak_of_random_divergences_converges() -> Outcome<()> {
640 let mut rng = Rng::new(0x5e_ed_50_ac);
641 for trial in 0..60 {
642 let prefix = rng.below(25);
643 let left = rng.below(6);
644 let right = rng.below(6);
645 let mut a = OpLog::new();
646 res!(write(&mut a, 1, prefix, "s"));
647 let mut b = a.clone();
648 // Each side writes its own, sometimes merging what it finds.
649 for i in 0..left {
650 if rng.below(4) == 0 && !a.is_empty() {
651 res!(a.append(Record::new(
652 res!(Header::new(a.next_id(ReplicaId::new(5)), a.frontier())),
653 Op::Mark { name: fmt!("ma{}", i), body: None, time: None },
654 )));
655 } else {
656 res!(write(&mut a, 1 + rng.below(2) as u64, 1, &fmt!("a{}", i)));
657 }
658 }
659 for i in 0..right {
660 if rng.below(4) == 0 && !b.is_empty() {
661 res!(b.append(Record::new(
662 res!(Header::new(b.next_id(ReplicaId::new(6)), b.frontier())),
663 Op::Mark { name: fmt!("mb{}", i), body: None, time: None },
664 )));
665 } else {
666 res!(write(&mut b, 3 + rng.below(2) as u64, 1, &fmt!("b{}", i)));
667 }
668 }
669 let want = a.len() + b.len() - prefix;
670 let mode = match trial % 3 {
671 0 => Mode::Walk,
672 1 => Mode::sketch(16),
673 // Deliberately too small some of the time, so the fallback is soaked
674 // as well as the happy path.
675 _ => Mode::sketch(rng.below(3)),
676 };
677 let mut wa = a.clone();
678 let mut wb = b.clone();
679 match exchange(&mut wa, &mut wb, mode) {
680 Ok(_) => {},
681 Err(e) => return Err(err!(e,
682 "Trial {} failed under {:?}: a prefix of {}, then {} and {}.",
683 trial, mode, prefix, left, right; Test)),
684 }
685 match agree(&wa, &wb) {
686 Ok(()) => {},
687 Err(e) => return Err(err!(e,
688 "Trial {} did not converge under {:?}.", trial, mode; Test)),
689 }
690 assert_eq!(wa.len(), want, "trial {} lost or invented an operation", trial);
691 }
692 Ok(())
693}
694
695/// There is no client and no server.
696#[test]
697fn either_side_may_open() -> Outcome<()> {
698 // Only A opens; B answers what it was sent and opens in the same turn.
699 let (mut a, mut b) = res!(diverged(10, 3, 2));
700 let mut sa = Session::new(Mode::Walk);
701 let mut sb = Session::new(Mode::Walk);
702 let mut to_b = vec![res!(sa.open(&a))];
703 let mut to_a: Vec<Message> = Vec::new();
704 let mut guard = 0usize;
705 while !(sa.is_converged() && sb.is_converged()) {
706 guard += 1;
707 assert!(guard < 16, "no convergence with a single opener");
708 for msg in std::mem::take(&mut to_b) {
709 to_a.extend(res!(sb.receive(&mut b, msg)).send);
710 }
711 for msg in std::mem::take(&mut to_a) {
712 to_b.extend(res!(sa.receive(&mut a, msg)).send);
713 }
714 }
715 res!(agree(&a, &b));
716 assert!(sa.ops_sent() > 0 && sb.ops_sent() > 0, "both sides had news");
717 Ok(())
718}
719
720#[test]
721fn a_session_counts_what_it_moved() -> Outcome<()> {
722 let mut a = OpLog::new();
723 res!(write(&mut a, 1, 9, "a"));
724 let mut b = OpLog::new();
725 res!(write(&mut b, 2, 4, "b"));
726 let mut sa = Session::new(Mode::Walk);
727 let mut sb = Session::new(Mode::Walk);
728 let mut to_a = vec![res!(sb.open(&b))];
729 let mut to_b = vec![res!(sa.open(&a))];
730 let mut guard = 0usize;
731 while !(sa.is_converged() && sb.is_converged()) {
732 guard += 1;
733 assert!(guard < 16);
734 for msg in std::mem::take(&mut to_a) {
735 to_b.extend(res!(sa.receive(&mut a, msg)).send);
736 }
737 for msg in std::mem::take(&mut to_b) {
738 to_a.extend(res!(sb.receive(&mut b, msg)).send);
739 }
740 }
741 res!(agree(&a, &b));
742 assert_eq!(sa.ops_sent(), 9);
743 assert_eq!(sa.ops_absorbed(), 4);
744 assert_eq!(sb.ops_sent(), 4);
745 assert_eq!(sb.ops_absorbed(), 9);
746 assert_eq!(sa.mode(), Mode::Walk);
747 Ok(())
748}
749
750/// The rule is judged on shapes alone and cannot see overlap, so two logs of a
751/// size that in fact share nothing are still offered a sketch. That is exactly
752/// the bad estimate the fallback exists for, and it costs no round trip.
753#[test]
754fn the_mode_is_chosen_from_the_two_shapes() -> Outcome<()> {
755 // A clone: nothing here, a history there.
756 assert_eq!(Mode::between(0, 0, 400, 1), Mode::Walk);
757 // A small divergence over a large shared history.
758 match Mode::between(400, 1, 402, 1) {
759 Mode::Sketch { estimate, .. } => assert_eq!(estimate, 2 + FANOUT * 2),
760 other => return Err(err!("A small divergence chose {:?}.", other; Test, Mismatch)),
761 }
762 // Two logs of a size, both wide open: the guess reaches the smaller log, so
763 // the walk is the cheaper answer.
764 assert_eq!(Mode::between(20, 4, 20, 4), Mode::Walk);
765 // One head each over the same length is a sketch, whatever the two logs turn
766 // out to share.
767 assert!(matches!(Mode::between(30, 1, 30, 1), Mode::Sketch { .. }));
768 // And two empty logs do not divide by nothing.
769 assert_eq!(Mode::between(0, 0, 0, 0), Mode::Walk);
770 Ok(())
771}
772
773/// A session refuses a piece of an operation rather than trying to place it.
774///
775/// Putting a run of pieces back together is the carrier's work, and a session
776/// that took one would be placing something nobody signed. The refusal names
777/// where the work belongs, because a caller meeting it has reached for the wrong
778/// layer rather than made a mistake.
779#[test]
780fn a_session_refuses_a_piece_of_an_operation() -> Outcome<()> {
781 let mut log = OpLog::default();
782 let mut session = Session::new(Mode::Walk);
783 let e = match session.receive(&mut log, Message::Part {
784 id: OpId::new(ReplicaId::new(1), 1),
785 seq: 0,
786 total: 4,
787 bytes: vec![0x01],
788 }) {
789 Ok(_) => return Err(err!("A session placed a piece of an operation."; Test)),
790 Err(e) => e,
791 };
792 assert!(fmt!("{}", e).contains("Parts"), "the refusal does not say where the work belongs: {}", e);
793 assert_eq!(log.len(), 0, "a refused piece changed the log");
794 Ok(())
795}
796
797/// What every message kind comes to, measured and then encoded, and the two
798/// numbers compared.
799///
800/// [`Message::encoded_len`] is what both ends decide a body's contents by, so a
801/// number one short of the truth is a body over a bound that was published and a
802/// proxy closing the connection. The corpus is every kind, at the shapes whose
803/// lengths are decided by something other than the message itself: a frontier of
804/// none and of many, a send of none, one and many, and the pieces an operation
805/// too large for the carrier is cut into.
806///
807/// Proved red by adding one to the answer, and again by leaving the magic and the
808/// version out of it.
809#[test]
810fn encoded_len_is_what_the_message_encodes_to() -> Outcome<()> {
811 let head = res!(Header::new(
812 OpId::new(ReplicaId::new(5), 7),
813 vec![OpId::new(ReplicaId::new(1), 1)],
814 ));
815 let entry = |len: usize| crate::segment::Entry::Bare(Record::new(head.clone(), Op::Proposal {
816 title: fmt!("of {} bytes", len),
817 body: vec![0x5a; len],
818 voice: fmt!("wren"),
819 time: 1_755_400_000,
820 }));
821 let heads: Vec<OpId> = (1..=40).map(|i| OpId::new(ReplicaId::new(i), i)).collect();
822 // One piece over a mebibyte, so the pieces are many and the last one short.
823 let big = crate::segment::Entry::Bare(Record::new(head.clone(), Op::Proposal {
824 title: fmt!("a large one"),
825 body: vec![0xa5; 3_000_000],
826 voice: fmt!("wren"),
827 time: 1_755_400_000,
828 }));
829 let mut corpus = vec![
830 (fmt!("an empty hello"), Message::hello(Vec::new())),
831 (fmt!("a wide hello"), Message::hello(heads.clone())),
832 (fmt!("an empty sketch"), Message::sketch(Vec::new(), Vec::new(), 0)),
833 (fmt!("a sketch"), Message::sketch(heads, vec![0x11; 4_000], u64::MAX)),
834 (fmt!("an empty send"), Message::Send { entries: Vec::new() }),
835 (fmt!("a send of one"), Message::Send { entries: vec![entry(0)] }),
836 (fmt!("a send of many"), Message::Send {
837 entries: (0..64).map(|i| entry(i * 37)).collect(),
838 }),
839 (fmt!("a done"), Message::Done),
840 ];
841 let pieces = res!(Message::part(&big, 1 << 20));
842 assert!(pieces.len() > 2, "the operation went in {} pieces", pieces.len());
843 for (i, piece) in pieces.into_iter().enumerate() {
844 corpus.push((fmt!("piece {}", i), piece));
845 }
846 let mut kinds = BTreeSet::new();
847 for (name, msg) in corpus {
848 kinds.insert(msg.kind());
849 let said = res!(msg.encoded_len());
850 let wrote = res!(msg.encode()).len();
851 assert_eq!(said, wrote,
852 "{} is measured at {} bytes and encodes to {}", name, said, wrote);
853 }
854 assert_eq!(kinds.len(), 5, "the corpus covers {} of the five message kinds", kinds.len());
855 Ok(())
856}
857
858
859/// Two logs that already agree cost a handful of hundred bytes, whatever the
860/// history behind them.
861///
862/// The number nothing else in this file pins. A sync that carries nothing is the
863/// ordinary outcome of a repository somebody syncs often, and what it costs is
864/// the whole argument for sketching: the exchange is proportional to the
865/// difference and not to the history, so a thousand-fold larger history has to
866/// cost the same nothing.
867#[test]
868fn a_sync_that_carries_nothing_costs_almost_nothing() -> Outcome<()> {
869 let mut a = OpLog::new();
870 res!(write(&mut a, 1, 1_400, "x"));
871 res!(merge(&mut a, 2, "join"));
872 let mut b = a.clone();
873 let mode = Mode::between(a.len(), a.frontier().len(), b.len(), 0);
874 let tally = res!(exchange(&mut a, &mut b, mode));
875 res!(agree(&a, &b));
876 assert_eq!(tally.ops, 0, "an exchange between equals handed something over");
877 assert!(tally.bytes < 4 << 10,
878 "a no-op sync of {} operations cost {} bytes over {} messages",
879 a.len(), tally.bytes, tally.messages);
880 assert!(tally.fell_back.is_none(), "it fell back: {:?}", tally.fell_back);
881 Ok(())
882}
883
884/// A two-sided divergence reconciles by sketch, at every size worth trying, and
885/// never reaches the walk.
886///
887/// The case the estimate cannot see. Each side holds a head the other does not,
888/// so the two logs' lengths say nothing about how far apart they are -- k against
889/// k + 2 is two, and the truth is 2k + 2. Where the first table is too small the
890/// answer is a larger table, so what this asserts is that the sketch path
891/// finishes the job: the logs agree, nothing fell back to the walk, and the whole
892/// exchange stays far under the history it reconciles.
893#[test]
894fn a_two_sided_divergence_decodes() -> Outcome<()> {
895 for k in [1usize, 2, 5, 9, 17, 33, 64] {
896 let (mut a, mut b) = res!(diverged(300, k, k + 2));
897 let whole = a.len();
898 // What each end knows before it opens: its own shape, the other's length,
899 // and that the other's frontier is news to it.
900 let mode = Mode::between(a.len(), a.frontier().len(), b.len(), b.frontier().len());
901 let tally = res!(exchange(&mut a, &mut b, mode));
902 res!(agree(&a, &b));
903 assert!(tally.fell_back.is_none(),
904 "at k = {} the exchange fell back to the walk: {:?}", k, tally.fell_back);
905 // The walk would have offered the whole log from each side, since neither
906 // can subtract the other's tip.
907 assert!(tally.ops < whole,
908 "at k = {} the exchange handed over {} operations of a {} operation log",
909 k, tally.ops, whole);
910 }
911 Ok(())
912}
913
914/// A cursor carries a bounded walk forward, against a peer whose head the log
915/// does not hold.
916///
917/// The fault in one place. A peer that has written anything of its own presents
918/// a frontier this log cannot subtract, so the owed set is the whole log however
919/// much of it that peer already holds -- and a carrier with a bounded reply sends
920/// the same prefix every session, for ever. The cursor is what a session with no
921/// memory is told instead, and it is one identifier: everything at or before it
922/// in the append order is held, so the next owed set begins where the last one
923/// stopped.
924#[test]
925fn a_cursor_moves_a_bounded_walk_along() -> Outcome<()> {
926 let mut here = OpLog::new();
927 res!(write(&mut here, 1, 60, "x"));
928 // A peer holding the whole of it and one operation of its own, which is the
929 // head this log has never seen.
930 let mut there = here.clone();
931 res!(write(&mut there, 2, 1, "mine"));
932 let heads = there.frontier();
933 assert_eq!(heads.len(), 1);
934 assert!(!here.contains(&heads[0]), "the peer's head is one this log holds");
935
936 // Without a cursor, every session owes the same whole log.
937 let mut first = Session::new(Mode::Walk);
938 let turn = res!(first.receive(&mut here.clone(), Message::hello(heads.clone())));
939 let offered = turn.send.iter().map(|m| m.entries().len()).sum::<usize>();
940 assert_eq!(offered, 60, "the walk offered {} of a 60 operation log", offered);
941
942 // With one, the owed set begins after the operation it names. Twenty at a
943 // time, which is what a bounded reply leaves behind.
944 let mut at = 0usize;
945 let mut sessions = 0usize;
946 while at < 60 {
947 sessions += 1;
948 if sessions > 8 {
949 return Err(err!(
950 "{} sessions carried the walk to {} of 60.", sessions, at;
951 Test, Excessive));
952 }
953 let mut session = Session::new(Mode::Walk);
954 if at > 0 {
955 let cursor = match here.at(at - 1) {
956 Some(rec) => rec.id(),
957 None => return Err(err!("The log lost operation {}.", at; Test, Missing)),
958 };
959 res!(session.receive(&mut here.clone(), Message::Resume { at: cursor }));
960 }
961 let turn = res!(session.receive(&mut here.clone(), Message::hello(heads.clone())));
962 let mut sent: Vec<OpId> = Vec::new();
963 for msg in &turn.send {
964 for entry in msg.entries() {
965 sent.push(res!(entry.id()));
966 }
967 }
968 assert_eq!(sent.len(), 60 - at,
969 "at cursor {} the session owed {} operations", at, sent.len());
970 match here.at(at) {
971 Some(rec) => assert_eq!(sent[0], rec.id(),
972 "the session began at {} rather than at the cursor", sent[0]),
973 None => return Err(err!("The log lost operation {}.", at; Test, Missing)),
974 }
975 // What a bounded reply would have carried of it.
976 at += 20;
977 }
978 assert_eq!(sessions, 3, "the walk took {} sessions at twenty a turn", sessions);
979 Ok(())
980}
981
982
983/// An end that grows against a peer that cannot answer a grown table strands it,
984/// so it does not grow.
985///
986/// **The compatibility failure of the whole idea, and it is not symmetric.** A
987/// grown table is a question, and the answer to it is the peer opening again at
988/// the new shape. A peer built before growth existed opens once: fed a grown
989/// table it decodes it, hands over what it owes, and the end that grew is left
990/// holding a table nobody will answer -- it said only the table, so it never
991/// worked out what it owed, and a carrier that keeps nothing between requests has
992/// no opening to answer on the visit after.
993///
994/// Both halves are asserted here, because the second is the reason the first is a
995/// knob rather than a rule. Told that its peer opens once, the end that would
996/// have grown walks instead and the two logs agree; told nothing, it grows and
997/// the pipe empties with it unconverged, which is the state a carrier reports as
998/// unfinished.
999#[test]
1000fn growth_against_a_peer_that_opens_once_is_refused() -> Outcome<()> {
1001 // A divergence the first table cannot hold, which is what makes the question
1002 // arise at all.
1003 let spread = |a: &OpLog, b: &OpLog| Mode::between(
1004 a.len(), a.frontier().len(), b.len(), b.frontier().len());
1005
1006 // A opens once and never again, which is every build before the cursor. B
1007 // answers it, and is told what A is.
1008 let (mut a, mut b) = res!(diverged(300, 30, 32));
1009 let mode = spread(&a, &b);
1010 let mut sa = Session::new(mode).with_growth(Growth::Refused);
1011 let mut sb = Session::new(mode).with_growth(Growth::Refused);
1012 assert!(res!(called_upon(&mut sa, &mut a, &mut sb, &mut b)),
1013 "an exchange with a peer that opens once did not finish");
1014 res!(agree(&a, &b));
1015 assert!(sb.fell_back().is_some(), "B grew against a peer that opens once");
1016
1017 // The same exchange with B told nothing: it grows, and A -- which opens once
1018 // -- answers the grown table and then has nothing more to say. B never made
1019 // its own opening count, so it is left unconverged with the pipe empty.
1020 let (mut a, mut b) = res!(diverged(300, 30, 32));
1021 let mut sa = Session::new(mode).with_growth(Growth::Refused);
1022 let mut sb = Session::new(mode);
1023 assert!(!res!(called_upon(&mut sa, &mut a, &mut sb, &mut b)),
1024 "a peer that opens once answered a grown table, so growth costs nothing \
1025 against it and this knob is not needed");
1026 assert!(!sb.is_converged(), "B is the end left holding a table nobody answered");
1027
1028 // And between two ends that both re-open the same divergence settles by
1029 // sketch, so what is being refused above is a saving and not the exchange.
1030 let (mut c, mut d) = res!(diverged(300, 30, 32));
1031 let tally = res!(exchange(&mut c, &mut d, mode));
1032 res!(agree(&c, &d));
1033 assert!(tally.fell_back.is_none(),
1034 "two ends that both re-open fell back: {:?}", tally.fell_back);
1035 Ok(())
1036}
1037
1038/// Runs an exchange where `sa` speaks first and `sb` answers, and says whether
1039/// both ends converged before the pipe emptied.
1040///
1041/// Unlike [`exchange`] it is not an error for the pipe to empty with one end
1042/// unfinished, because that is the outcome one of its callers is asserting.
1043fn called_upon(
1044 sa: &mut Session,
1045 a: &mut OpLog,
1046 sb: &mut Session,
1047 b: &mut OpLog,
1048)
1049 -> Outcome<bool>
1050{
1051 let mut queue = vec![res!(sa.open(a))];
1052 let mut guard = 0usize;
1053 while !(sa.is_converged() && sb.is_converged()) {
1054 guard += 1;
1055 if guard > 16 {
1056 return Err(err!(
1057 "An exchange took {} turns without converging.", guard; Test, Excessive));
1058 }
1059 if queue.is_empty() {
1060 return Ok(false);
1061 }
1062 let mut back = Vec::new();
1063 for msg in std::mem::take(&mut queue) {
1064 for out in res!(sb.receive(b, msg)).send {
1065 back.push(out);
1066 }
1067 }
1068 for msg in back {
1069 for out in res!(sa.receive(a, msg)).send {
1070 queue.push(out);
1071 }
1072 }
1073 }
1074 Ok(true)
1075}
1076
1077/// Reading the sizing rule backwards lands exactly where it started.
1078///
1079/// Both directions cost something and they are not the same something. A cell
1080/// short of the table being answered is a stall put straight back on the wire; a
1081/// cell over is a round trip, because an arriving table wider than the last one
1082/// sent is exactly what a re-opening is read off, so a peer that answers with one
1083/// cell more than it was given is answered again.
1084#[test]
1085fn a_width_is_a_fixed_point_of_the_sizing_rule() -> Outcome<()> {
1086 use crate::sync::sketch::{
1087 cells_for,
1088 estimate_for,
1089 MAX_CELLS,
1090 MIN_CELLS,
1091 };
1092
1093 // Every width a sketch can actually declare, which is the image of the sizing
1094 // rule and not every number between its ends.
1095 let mut estimate = 0usize;
1096 let mut seen = 0usize;
1097 while estimate <= (2 * MAX_CELLS) / 3 + 4 {
1098 let cells = cells_for(estimate);
1099 assert!(cells >= MIN_CELLS && cells <= MAX_CELLS);
1100 assert_eq!(cells_for(estimate_for(cells)), cells,
1101 "a table of {} cells is read back as an estimate of {}, which sizes {}",
1102 cells, estimate_for(cells), cells_for(estimate_for(cells)));
1103 seen += 1;
1104 estimate += 1;
1105 }
1106 assert!(seen > 600_000, "only {} widths were tried", seen);
1107 // And a width the rule cannot produce is rounded up rather than down, since
1108 // narrower is the direction that stalls.
1109 for cells in MIN_CELLS..4096 {
1110 assert!(cells_for(estimate_for(cells)) >= cells,
1111 "a table of {} cells is answered with {}", cells,
1112 cells_for(estimate_for(cells)));
1113 }
1114 Ok(())
1115}