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 | |
| 16 | use crate::id::{ |
| 17 | OpId, |
| 18 | ReplicaId, |
| 19 | }; |
| 20 | use crate::log::OpLog; |
| 21 | use crate::op::{ |
| 22 | Header, |
| 23 | Op, |
| 24 | Record, |
| 25 | }; |
| 26 | use crate::sync::msg::Message; |
| 27 | use crate::sync::session::{ |
| 28 | Growth, |
| 29 | Mode, |
| 30 | Session, |
| 31 | Step, |
| 32 | FANOUT, |
| 33 | }; |
| 34 | use crate::sync::sketch::{ |
| 35 | Fallback, |
| 36 | MIN_CELLS, |
| 37 | }; |
| 38 | |
| 39 | use oxedyne_fe2o3_core::prelude::*; |
| 40 | |
| 41 | use std::collections::BTreeSet; |
| 42 | |
| 43 | |
| 44 | /// A small linear congruential generator, so a failure can be reproduced. |
| 45 | struct Rng(u64); |
| 46 | |
| 47 | impl 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)] |
| 67 | struct 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. |
| 78 | fn 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. |
| 134 | fn 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 | |
| 165 | fn 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. |
| 177 | fn 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. |
| 185 | fn 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] |
| 211 | fn 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] |
| 226 | fn 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] |
| 246 | fn 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. |
| 283 | fn 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] |
| 338 | fn 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] |
| 359 | fn 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] |
| 375 | fn 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] |
| 390 | fn 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] |
| 404 | fn 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] |
| 436 | fn 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] |
| 494 | fn 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] |
| 538 | fn 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] |
| 570 | fn 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] |
| 606 | fn 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] |
| 639 | fn 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] |
| 697 | fn 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] |
| 721 | fn 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] |
| 754 | fn 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] |
| 780 | fn 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] |
| 810 | fn 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] |
| 868 | fn 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] |
| 894 | fn 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] |
| 925 | fn 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] |
| 1000 | fn 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. |
| 1043 | fn 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] |
| 1085 | fn 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 | } |