oxedyne/fe2o3/fe2o3_ore/src/sync/walk.rs
15.9 KiB, 92 runs
created by r1870400018:19641, 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 | //! What a peer at a given frontier is owed, and the closure that makes it safe |
| 2 | //! to send. |
| 3 | //! |
| 4 | //! # The owed set |
| 5 | //! |
| 6 | //! A peer says what its frontier is. Everything it holds is an ancestor of one |
| 7 | //! of those heads, because a log is causally closed and its frontier dominates |
| 8 | //! it. So of the operations we hold, the ones it demonstrably holds too are the |
| 9 | //! ancestors -- and the heads themselves -- of every head we can also see: |
| 10 | //! |
| 11 | //! ```text |
| 12 | //! roots = (their heads ∪ what they have handed us) ∩ our log |
| 13 | //! covered = ancestors-or-self, within our log, of roots |
| 14 | //! owed = our log \ covered |
| 15 | //! ``` |
| 16 | //! |
| 17 | //! That is sound: everything in `covered` is genuinely theirs, so nothing they |
| 18 | //! lack is ever left out. It is not tight. A head of theirs we have never seen |
| 19 | //! is news *for us*, and tells us nothing about what they hold, so the branch it |
| 20 | //! sits on cannot be subtracted; where both peers have written since they last |
| 21 | //! spoke, neither can subtract the other's tip and each sends its whole log. |
| 22 | //! What the receiver already holds it drops, so the cost of being loose is |
| 23 | //! bytes and never correctness. |
| 24 | //! |
| 25 | //! Two cases are exactly tight, and they are the common ones: |
| 26 | //! |
| 27 | //! - A peer that is behind us and has written nothing of its own -- a clone, a |
| 28 | //! fetch after someone else pushed -- has heads we hold, so `covered` is |
| 29 | //! precisely its log and `owed` is precisely the news. |
| 30 | //! - A peer that has everything we have gets an empty owed set and one message. |
| 31 | //! |
| 32 | //! Where both sides have written, [`crate::sync::sketch`] is what makes the |
| 33 | //! exchange proportional to the difference instead. |
| 34 | //! |
| 35 | //! # What they handed us |
| 36 | //! |
| 37 | //! The second root set is the one a frontier cannot report. A peer that hands an |
| 38 | //! operation over holds it -- a proof, not a claim -- and its heads say nothing |
| 39 | //! about it, since a head we have never seen subtracts nothing. Held in one |
| 40 | //! session that costs nothing, because a session sends what it owes once. It |
| 41 | //! costs a carrier that runs several: a bounded reply makes a large clone into a |
| 42 | //! run of sessions, and a session that started again from the frontier alone |
| 43 | //! offered back, every time, the whole prefix the sessions before it had just |
| 44 | //! delivered. |
| 45 | //! |
| 46 | //! Measured on the clone of fe2o3 of 12026-08-22, 35,314 operations over sixteen |
| 47 | //! sessions: 166,224 operations offered, every one of them already held at the |
| 48 | //! far end, 717 MB up against 87 MB down. [`crate::sync::Session::knowing`] is |
| 49 | //! how a carrier carries the answer across the boundary. |
| 50 | //! |
| 51 | //! # Closure at both ends |
| 52 | //! |
| 53 | //! The sender closes what it is about to send against what it believes the |
| 54 | //! receiver holds ([`close`]), and the receiver checks the property on arrival |
| 55 | //! rather than trusting it ([`arrival_gap`]). The first is a proof obligation |
| 56 | //! discharged where the information is; the second is what stops a peer that got |
| 57 | //! it wrong from leaving a hole in someone else's history. |
| 58 | //! |
| 59 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 60 | //! Anthropic Claude |
| 61 | |
| 62 | use crate::id::OpId; |
| 63 | use crate::log::OpLog; |
| 64 | use crate::segment::Entry; |
| 65 | |
| 66 | use oxedyne_fe2o3_core::prelude::*; |
| 67 | |
| 68 | use std::collections::BTreeSet; |
| 69 | |
| 70 | |
| 71 | /// The operations of `log` that a peer whose frontier is `heads` demonstrably |
| 72 | /// holds. Heads the log does not hold are skipped: they are operations the peer |
| 73 | /// has and we do not, and they say nothing about what we hold. |
| 74 | /// |
| 75 | /// `known` is the rest of what the peer has shown it holds, which is what it |
| 76 | /// handed over earlier in the same exchange. Those are roots of the same walk and |
| 77 | /// not a separate kind of thing: a log is causally closed, so a peer that holds an |
| 78 | /// operation holds every ancestor of it, exactly as it does of a head. |
| 79 | pub fn covered(log: &OpLog, heads: &[OpId], known: &BTreeSet<OpId>) -> BTreeSet<OpId> { |
| 80 | let mut seen: BTreeSet<OpId> = BTreeSet::new(); |
| 81 | let mut stack: Vec<OpId> = heads |
| 82 | .iter() |
| 83 | .chain(known.iter()) |
| 84 | .filter(|h| log.contains(h)) |
| 85 | .copied() |
| 86 | .collect(); |
| 87 | while let Some(id) = stack.pop() { |
| 88 | if !seen.insert(id) { |
| 89 | continue; |
| 90 | } |
| 91 | // A record in the log has every parent in the log, by the log's own |
| 92 | // append guard, so the walk never leaves it. |
| 93 | if let Some(rec) = log.get(&id) { |
| 94 | for p in rec.parents() { |
| 95 | if !seen.contains(p) { |
| 96 | stack.push(*p); |
| 97 | } |
| 98 | } |
| 99 | } |
| 100 | } |
| 101 | seen |
| 102 | } |
| 103 | |
| 104 | /// In the log's append order, which is a linear extension of the causal order, |
| 105 | /// so a receiver places the whole batch in one pass; nothing depends on it, |
| 106 | /// since [`OpLog::absorb`](crate::log::OpLog::absorb) takes a batch however it |
| 107 | /// is shuffled. |
| 108 | pub fn owed(log: &OpLog, heads: &[OpId], known: &BTreeSet<OpId>) -> Vec<OpId> { |
| 109 | let held = covered(log, heads, known); |
| 110 | log.iter() |
| 111 | .map(|rec| rec.id()) |
| 112 | .filter(|id| !held.contains(id)) |
| 113 | .collect() |
| 114 | } |
| 115 | |
| 116 | /// Extends a send set with every ancestor the receiver is not known to hold, so |
| 117 | /// that what arrives closes causally against what is already there. |
| 118 | /// |
| 119 | /// `held` is what the sender believes the receiver has. Ancestors outside both |
| 120 | /// sets are pulled in; an identifier the log does not hold is dropped, since |
| 121 | /// nothing can be said about an operation nobody has. |
| 122 | /// |
| 123 | /// For a send set worked out by [`owed`] this adds nothing, and the same is true |
| 124 | /// of a difference a sketch decoded in full: both are complements of an |
| 125 | /// ancestor-closed set, and the complement of an ancestor-closed set is closed |
| 126 | /// downwards by construction. The step is here because that is a property of the |
| 127 | /// *inputs*, and this function is what makes it a property of the output. |
| 128 | pub fn close(log: &OpLog, send: &[OpId], held: &BTreeSet<OpId>) -> BTreeSet<OpId> { |
| 129 | let mut out: BTreeSet<OpId> = BTreeSet::new(); |
| 130 | let mut stack: Vec<OpId> = Vec::new(); |
| 131 | for id in send { |
| 132 | if log.contains(id) && !held.contains(id) && out.insert(*id) { |
| 133 | stack.push(*id); |
| 134 | } |
| 135 | } |
| 136 | while let Some(id) = stack.pop() { |
| 137 | if let Some(rec) = log.get(&id) { |
| 138 | for p in rec.parents() { |
| 139 | if !held.contains(p) && out.insert(*p) { |
| 140 | stack.push(*p); |
| 141 | } |
| 142 | } |
| 143 | } |
| 144 | } |
| 145 | out |
| 146 | } |
| 147 | |
| 148 | /// In the log's append order. Fails where the log does not hold one of them, |
| 149 | /// which would mean the sender had worked out a set it cannot deliver. |
| 150 | pub fn entries_for(log: &OpLog, ids: &BTreeSet<OpId>) |
| 151 | -> Outcome<Vec<Entry>> |
| 152 | { |
| 153 | for id in ids { |
| 154 | if !log.contains(id) { |
| 155 | return Err(err!( |
| 156 | "The log does not hold {}, so it cannot be sent.", id; |
| 157 | Invalid, Input, Missing)); |
| 158 | } |
| 159 | } |
| 160 | Ok(log.iter() |
| 161 | .filter(|rec| ids.contains(&rec.id())) |
| 162 | .map(|rec| Entry::Bare(rec.clone())) |
| 163 | .collect()) |
| 164 | } |
| 165 | |
| 166 | /// Returns the first operation of a batch whose parent neither the batch nor the |
| 167 | /// log holds, together with that parent. |
| 168 | /// |
| 169 | /// `None` means the batch closes causally against the log: absorbing it leaves |
| 170 | /// no hole. This is the check a receiver makes before it absorbs anything, and |
| 171 | /// it is deliberately not [`crate::log::Causality::gap`], which sees only the set |
| 172 | /// it was built over and would have to be handed the whole log to answer the |
| 173 | /// question. |
| 174 | pub fn arrival_gap(log: &OpLog, entries: &[Entry]) |
| 175 | -> Outcome<Option<(OpId, OpId)>> |
| 176 | { |
| 177 | let mut arriving: BTreeSet<OpId> = BTreeSet::new(); |
| 178 | let mut records = Vec::with_capacity(entries.len()); |
| 179 | for entry in entries { |
| 180 | let rec = res!(entry.peek()); |
| 181 | arriving.insert(rec.id()); |
| 182 | records.push(rec); |
| 183 | } |
| 184 | for rec in &records { |
| 185 | for p in rec.parents() { |
| 186 | if !arriving.contains(p) && !log.contains(p) { |
| 187 | return Ok(Some((rec.id(), *p))); |
| 188 | } |
| 189 | } |
| 190 | } |
| 191 | Ok(None) |
| 192 | } |
| 193 | |
| 194 | pub fn closes(log: &OpLog, entries: &[Entry]) |
| 195 | -> Outcome<bool> |
| 196 | { |
| 197 | Ok(res!(arrival_gap(log, entries)).is_none()) |
| 198 | } |
| 199 | |
| 200 | |
| 201 | #[cfg(test)] |
| 202 | mod tests { |
| 203 | use super::*; |
| 204 | |
| 205 | use crate::id::ReplicaId; |
| 206 | use crate::op::{ |
| 207 | Header, |
| 208 | Op, |
| 209 | Record, |
| 210 | }; |
| 211 | |
| 212 | fn oid(replica: u64, counter: u64) -> OpId { |
| 213 | OpId::new(ReplicaId::new(replica), counter) |
| 214 | } |
| 215 | |
| 216 | /// Nothing remembered beyond the frontier, which is what a first session knows. |
| 217 | fn nil() -> BTreeSet<OpId> { |
| 218 | BTreeSet::new() |
| 219 | } |
| 220 | |
| 221 | fn rec(id: OpId, parents: Vec<OpId>, name: &str) |
| 222 | -> Outcome<Record> |
| 223 | { |
| 224 | Ok(Record::new( |
| 225 | res!(Header::new(id, parents)), |
| 226 | Op::Mark { name: fmt!("{}", name), body: None, time: None }, |
| 227 | )) |
| 228 | } |
| 229 | |
| 230 | /// A log of a chain, a fork, and a merge: |
| 231 | /// |
| 232 | /// ```text |
| 233 | /// a -- b -- d |
| 234 | /// \ / |
| 235 | /// ------c |
| 236 | /// ``` |
| 237 | fn forked() |
| 238 | -> Outcome<OpLog> |
| 239 | { |
| 240 | let mut log = OpLog::new(); |
| 241 | res!(log.append(Record::root(oid(1, 1), Op::Mark { name: fmt!("a"), body: None, time: None }))); |
| 242 | res!(log.append(res!(rec(oid(1, 2), vec![oid(1, 1)], "b")))); |
| 243 | res!(log.append(res!(rec(oid(2, 3), vec![oid(1, 1)], "c")))); |
| 244 | res!(log.append(res!(rec(oid(1, 4), vec![oid(1, 2), oid(2, 3)], "d")))); |
| 245 | Ok(log) |
| 246 | } |
| 247 | |
| 248 | /// The ancestors of a head are subtracted whichever branch they sit on. |
| 249 | #[test] |
| 250 | fn a_head_we_hold_covers_its_ancestors() -> Outcome<()> { |
| 251 | let log = res!(forked()); |
| 252 | assert_eq!( |
| 253 | covered(&log, &[oid(1, 2)], &nil()).into_iter().collect::<Vec<_>>(), |
| 254 | vec![oid(1, 1), oid(1, 2)], |
| 255 | ); |
| 256 | assert_eq!(owed(&log, &[oid(1, 2)], &nil()), vec![oid(2, 3), oid(1, 4)]); |
| 257 | // A merge covers both branches, so nothing is owed. |
| 258 | assert!(owed(&log, &[oid(1, 4)], &nil()).is_empty()); |
| 259 | // Two heads together cover their union. |
| 260 | assert!(owed(&log, &[oid(1, 2), oid(2, 3)], &nil()) == vec![oid(1, 4)]); |
| 261 | Ok(()) |
| 262 | } |
| 263 | |
| 264 | /// The clone case, and the answer is exactly the log. |
| 265 | #[test] |
| 266 | fn an_empty_peer_is_owed_everything() -> Outcome<()> { |
| 267 | let log = res!(forked()); |
| 268 | assert!(covered(&log, &[], &nil()).is_empty()); |
| 269 | assert_eq!(owed(&log, &[], &nil()), vec![oid(1, 1), oid(1, 2), oid(2, 3), oid(1, 4)]); |
| 270 | Ok(()) |
| 271 | } |
| 272 | |
| 273 | /// It is news for us, and says nothing about what the peer holds, so the owed |
| 274 | /// set is loose in that case and never wrong. |
| 275 | #[test] |
| 276 | fn a_head_we_do_not_hold_covers_nothing() -> Outcome<()> { |
| 277 | let log = res!(forked()); |
| 278 | assert!(covered(&log, &[oid(9, 9)], &nil()).is_empty()); |
| 279 | assert_eq!(owed(&log, &[oid(9, 9)], &nil()).len(), log.len(), "the whole log, loosely"); |
| 280 | // Mixed: the head we hold still does its work. |
| 281 | assert_eq!(owed(&log, &[oid(9, 9), oid(1, 2)], &nil()), vec![oid(2, 3), oid(1, 4)]); |
| 282 | Ok(()) |
| 283 | } |
| 284 | |
| 285 | /// The half a frontier cannot report. A peer's heads say what it holds *now*; |
| 286 | /// nothing in them says "you handed me this", and a head nobody here has seen |
| 287 | /// subtracts nothing at all. |
| 288 | #[test] |
| 289 | fn what_they_handed_us_is_not_offered_back() -> Outcome<()> { |
| 290 | let log = res!(forked()); |
| 291 | // Their tip is news to us, so on the frontier alone the whole log is owed. |
| 292 | let theirs = vec![oid(9, 9)]; |
| 293 | assert_eq!(owed(&log, &theirs, &nil()).len(), log.len()); |
| 294 | // Remember two they handed over, and it is exactly the rest. |
| 295 | let handed: BTreeSet<OpId> = [oid(1, 1), oid(1, 2)].into_iter().collect(); |
| 296 | assert_eq!( |
| 297 | covered(&log, &theirs, &handed).into_iter().collect::<Vec<_>>(), |
| 298 | vec![oid(1, 1), oid(1, 2)], |
| 299 | ); |
| 300 | assert_eq!(owed(&log, &theirs, &handed), vec![oid(2, 3), oid(1, 4)]); |
| 301 | // One that they handed over covers its ancestors as a head does, since a log |
| 302 | // is causally closed: the merge alone leaves nothing owed. |
| 303 | let handed: BTreeSet<OpId> = [oid(1, 4)].into_iter().collect(); |
| 304 | assert!(owed(&log, &theirs, &handed).is_empty()); |
| 305 | // And an operation nobody here holds says nothing, exactly as their tip does. |
| 306 | let handed: BTreeSet<OpId> = [oid(8, 8)].into_iter().collect(); |
| 307 | assert_eq!(owed(&log, &theirs, &handed).len(), log.len()); |
| 308 | Ok(()) |
| 309 | } |
| 310 | |
| 311 | /// Which is what makes it safe to subtract, and the reason a remembered root is |
| 312 | /// walked rather than merely removed: a peer that holds an operation holds every |
| 313 | /// ancestor of it, so what is subtracted has to be closed downwards or the |
| 314 | /// remainder arrives at the far end with a hole in it. |
| 315 | #[test] |
| 316 | fn what_is_subtracted_is_closed_downwards() -> Outcome<()> { |
| 317 | let log = res!(forked()); |
| 318 | for handed in [ |
| 319 | vec![], |
| 320 | vec![oid(1, 1)], |
| 321 | vec![oid(1, 2)], |
| 322 | vec![oid(2, 3)], |
| 323 | vec![oid(1, 2), oid(2, 3)], |
| 324 | // The merge, whose parents were never handed over by name. |
| 325 | vec![oid(1, 4)], |
| 326 | vec![oid(1, 4), oid(8, 8)], |
| 327 | ] { |
| 328 | let handed: BTreeSet<OpId> = handed.into_iter().collect(); |
| 329 | let held = covered(&log, &[oid(9, 9)], &handed); |
| 330 | for id in &held { |
| 331 | let rec = match log.get(id) { |
| 332 | Some(r) => r, |
| 333 | None => return Err(err!("The log lost {}.", id; Test, Missing)), |
| 334 | }; |
| 335 | for p in rec.parents() { |
| 336 | assert!(held.contains(p), |
| 337 | "{} is subtracted and its parent {} is not, remembering {:?}", |
| 338 | id, p, handed); |
| 339 | } |
| 340 | } |
| 341 | // So a peer holding exactly that much takes what is left of the log with |
| 342 | // no hole in it, which is the check the receiver makes for itself. |
| 343 | let mut peer = OpLog::new(); |
| 344 | let mut have = Vec::new(); |
| 345 | for rec in log.iter() { |
| 346 | if held.contains(&rec.id()) { |
| 347 | have.push(rec.clone()); |
| 348 | } |
| 349 | } |
| 350 | res!(peer.absorb(have)); |
| 351 | let ids: BTreeSet<OpId> = owed(&log, &[oid(9, 9)], &handed).into_iter().collect(); |
| 352 | let entries = res!(entries_for(&log, &ids)); |
| 353 | assert_eq!(res!(arrival_gap(&peer, &entries)), None, |
| 354 | "what is left over does not close, remembering {:?}", handed); |
| 355 | } |
| 356 | Ok(()) |
| 357 | } |
| 358 | |
| 359 | /// Which is what makes it safe to send in any order. |
| 360 | #[test] |
| 361 | fn the_owed_set_closes_against_the_peer() -> Outcome<()> { |
| 362 | let log = res!(forked()); |
| 363 | for heads in [ |
| 364 | vec![], |
| 365 | vec![oid(1, 1)], |
| 366 | vec![oid(1, 2)], |
| 367 | vec![oid(2, 3)], |
| 368 | vec![oid(1, 2), oid(2, 3)], |
| 369 | vec![oid(1, 4)], |
| 370 | vec![oid(9, 9)], |
| 371 | ] { |
| 372 | let held = covered(&log, &heads, &nil()); |
| 373 | let ids = owed(&log, &heads, &nil()); |
| 374 | // Every parent of everything sent is either sent or held. |
| 375 | for id in &ids { |
| 376 | let rec = match log.get(id) { |
| 377 | Some(r) => r, |
| 378 | None => return Err(err!("The log lost {}.", id; Test, Missing)), |
| 379 | }; |
| 380 | for p in rec.parents() { |
| 381 | assert!( |
| 382 | ids.contains(p) || held.contains(p), |
| 383 | "{} names {}, which is neither sent nor held at {:?}", id, p, heads, |
| 384 | ); |
| 385 | } |
| 386 | } |
| 387 | // Which is what closing adds nothing to. |
| 388 | let closed = close(&log, &ids, &held); |
| 389 | assert_eq!( |
| 390 | closed.into_iter().collect::<Vec<_>>(), |
| 391 | { let mut v = ids.clone(); v.sort(); v }, |
| 392 | "closing an owed set at {:?} added something", heads, |
| 393 | ); |
| 394 | } |
| 395 | Ok(()) |
| 396 | } |
| 397 | |
| 398 | /// Pulling in what it is missing, which is what the step is for. |
| 399 | #[test] |
| 400 | fn closing_repairs_a_hole() -> Outcome<()> { |
| 401 | let log = res!(forked()); |
| 402 | // The merge alone, with the peer holding nothing: its parents and their |
| 403 | // parent all have to go too. |
| 404 | let closed = close(&log, &[oid(1, 4)], &BTreeSet::new()); |
| 405 | // A set is a set, so it comes back in identifier order rather than in the |
| 406 | // log's append order. |
| 407 | assert_eq!( |
| 408 | closed.into_iter().collect::<Vec<_>>(), |
| 409 | vec![oid(1, 1), oid(1, 2), oid(1, 4), oid(2, 3)], |
| 410 | ); |
| 411 | // With the peer holding one branch, only the other is pulled in. |
| 412 | let held: BTreeSet<OpId> = [oid(1, 1), oid(1, 2)].into_iter().collect(); |
| 413 | let closed = close(&log, &[oid(1, 4)], &held); |
| 414 | assert_eq!(closed.into_iter().collect::<Vec<_>>(), vec![oid(1, 4), oid(2, 3)]); |
| 415 | // An identifier nobody holds is dropped rather than invented. |
| 416 | assert!(close(&log, &[oid(9, 9)], &BTreeSet::new()).is_empty()); |
| 417 | Ok(()) |
| 418 | } |
| 419 | |
| 420 | /// The gap names both operations, the one that arrived and the parent nobody |
| 421 | /// holds. |
| 422 | #[test] |
| 423 | fn arrival_names_the_hole() -> Outcome<()> { |
| 424 | let log = res!(forked()); |
| 425 | let mut fresh = OpLog::new(); |
| 426 | let all = res!(entries_for(&log, &owed(&log, &[], &nil()).into_iter().collect())); |
| 427 | assert!(res!(closes(&fresh, &all)), "a whole history closes against nothing"); |
| 428 | // Drop the root, and the batch no longer closes. |
| 429 | let short: Vec<Entry> = all[1..].to_vec(); |
| 430 | assert_eq!(res!(arrival_gap(&fresh, &short)), Some((oid(1, 2), oid(1, 1)))); |
| 431 | // Absorb the root, and it does. |
| 432 | res!(fresh.absorb(vec![res!(all[0].peek())])); |
| 433 | assert!(res!(closes(&fresh, &short))); |
| 434 | Ok(()) |
| 435 | } |
| 436 | |
| 437 | /// And a set naming what the log does not hold is refused rather than half |
| 438 | /// delivered. |
| 439 | #[test] |
| 440 | fn entries_follow_the_append_order() -> Outcome<()> { |
| 441 | let log = res!(forked()); |
| 442 | let ids: BTreeSet<OpId> = [oid(1, 4), oid(1, 1)].into_iter().collect(); |
| 443 | let got = res!(entries_for(&log, &ids)); |
| 444 | let mut names: Vec<OpId> = Vec::new(); |
| 445 | for entry in &got { |
| 446 | names.push(res!(entry.peek()).id()); |
| 447 | } |
| 448 | assert_eq!(names, vec![oid(1, 1), oid(1, 4)]); |
| 449 | let absent: BTreeSet<OpId> = [oid(9, 9)].into_iter().collect(); |
| 450 | assert!(entries_for(&log, &absent).is_err()); |
| 451 | Ok(()) |
| 452 | } |
| 453 | } |