Oregami
Repositories/oxedyne/fe2o3

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
62use crate::id::OpId;
63use crate::log::OpLog;
64use crate::segment::Entry;
65
66use oxedyne_fe2o3_core::prelude::*;
67
68use 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.
79pub 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.
108pub 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.
128pub 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.
150pub 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.
174pub 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
194pub 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)]
202mod 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}