Oregami
Repositories/oxedyne/fe2o3

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

19.0 KiB, 97 runs

created by r1870400018:19647, 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//! Reconciling two logs by their difference rather than by their size.
2//!
3//! An invertible Bloom lookup table is a fixed-size sketch of a set. Subtract
4//! one peer's sketch from another's and what remains encodes the symmetric
5//! difference; peel it, and the two halves fall out -- what only they hold, and
6//! what only we hold. The cost is the size of the sketch, which is chosen from
7//! the expected difference, and not the size of either history. Two large logs
8//! that differ by a handful of operations reconcile in a few hundred bytes.
9//!
10//! # The key
11//!
12//! A table needs a fixed-length key, and an [`OpId`] does not have one: its
13//! encoding is two varints, so a name spends between two and twenty bytes
14//! depending on how far the replica and the counter have got. The sketch key is
15//! therefore its own spelling: sixteen bytes, the replica in eight big-endian
16//! bytes and the counter in eight more. Big-endian so that the byte order of the
17//! keys is the order of the identifiers, which costs nothing and makes a dumped
18//! table legible.
19//!
20//! # The size
21//!
22//! Per the sizing rule of [`oxedyne_fe2o3_data::iblt`]: about one and a half
23//! cells per expected difference at three hashes. The estimate is the caller's,
24//! because only the caller knows how long it has been since the last exchange.
25//! Below [`MIN_CELLS`] the table is padded, because one and a half cells per
26//! entry is an asymptotic statement and a handful of keys in a handful of cells
27//! stalls far too often to be worth the round trip.
28//!
29//! # When the estimate was wrong
30//!
31//! The peeling decoder stalls, and that is reported as [`Diff::Undecodable`]
32//! with the reason, not as an error. It is not a failure of anything: the
33//! estimate was a guess, the guess was low, and the frontier walk is still there
34//! to answer with. What must never happen is a decode that stalled being treated
35//! as a decode that finished, since the partial difference it recovered is
36//! exactly the arbitrary subset that must never be sent.
37//!
38//! # A sketch is compared under the sender's shape
39//!
40//! Two tables can only be subtracted if they agree on cells, hashes, key length
41//! and seed. Rather than make the peers negotiate that, a receiver builds its own
42//! table under the shape the arriving one declares. Both sides do it, so two
43//! peers that estimated differently still reconcile -- each answering under the
44//! other's shape -- and there is no configuration to get wrong.
45//!
46//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
47//! Anthropic Claude
48
49use crate::id::{
50 OpId,
51 ReplicaId,
52};
53use crate::log::OpLog;
54
55use oxedyne_fe2o3_core::prelude::*;
56use oxedyne_fe2o3_data::iblt::{
57 DecodeOutcome,
58 Iblt,
59 IbltConfig,
60};
61
62
63/// The bytes an operation name takes as a sketch key, which is what the table
64/// is keyed on and what a peer's table must agree with.
65pub const KEY_LEN: usize = 16;
66
67pub const HASHES: usize = 3; // hashes per key, what the sizing rule assumes
68
69pub const MIN_CELLS: usize = 16; // fewest cells, whatever the estimate
70
71// The most cells a sketch may declare before a receiver answers with the walk
72// instead. A peer that genuinely expects a difference this large wants a bulk
73// transfer, which is what the walk is, so the cap costs nothing that was worth
74// having.
75pub const MAX_CELLS: usize = 1 << 20;
76
77/// The most cells a sketch grown to answer a stalled decode may declare.
78///
79/// A table crosses in one message and nothing cuts one up, so the bound is the
80/// smallest body a carrier is met with in ordinary use: one mebibyte, which is
81/// nginx's default and what Ore's own relay client falls back to when a relay
82/// publishes no limit. At twenty-eight bytes a cell -- the key, the fingerprint
83/// and the count -- this many cells come to 896 KiB, which leaves the framing
84/// room to spare, and they cover a difference of about twenty-one thousand
85/// operations. A difference larger than that is a bulk transfer, which is what
86/// [`crate::sync::walk`] is for, so growth stops here and the walk answers.
87pub const GROW_CELLS: usize = 1 << 15;
88
89/// Two peers must sketch under the same seed to subtract at all, and they do:
90/// the seed travels in the table's own serialised form, and a receiver adopts
91/// it. This one is for a caller with no reason to choose another.
92pub const SEED: u64 = 0x4f52_4553_594e_4331;
93
94
95/// Why a sketch could not be turned into a difference.
96#[derive(Clone, Copy, Debug, Eq, PartialEq)]
97pub enum Fallback {
98 // The peeling decoder stalled: the difference was larger than the table was
99 // sized for. What it recovered before stopping is discarded, since a part of
100 // a difference is not a difference.
101 Incomplete {
102 remaining: usize, // cells still holding state when peeling stopped
103 recovered: usize, // names recovered before it stopped
104 },
105 // The table declares more cells than MAX_CELLS.
106 Oversized {
107 cells: usize, // the cell count declared
108 },
109}
110
111impl Fallback {
112 /// For a caller that logs the reason.
113 pub fn why(&self) -> String {
114 match self {
115 Self::Incomplete { remaining, recovered } => fmt!(
116 "the peeling decoder stalled with {} cell{} left, having recovered {} \
117 name{}", remaining, if *remaining == 1 { "" } else { "s" },
118 recovered, if *recovered == 1 { "" } else { "s" },
119 ),
120 Self::Oversized { cells } => fmt!(
121 "the sketch declares {} cells, and this reader will not allocate past \
122 {}", cells, MAX_CELLS,
123 ),
124 }
125 }
126}
127
128
129/// What subtracting two sketches yielded.
130#[derive(Clone, Debug, Eq, PartialEq)]
131pub enum Diff {
132 // The difference, whole, both halves ascending.
133 Decoded {
134 remote_only: Vec<OpId>, // held only by the remote peer
135 local_only: Vec<OpId>, // held only here, which is what is owed
136 },
137 Undecodable(Fallback), // not recovered, and why
138}
139
140
141/// The replica in eight big-endian bytes, then the counter in eight more, so
142/// that the byte order of the keys is the order of the identifiers.
143pub fn key(id: &OpId) -> [u8; KEY_LEN] {
144 let mut out = [0u8; KEY_LEN];
145 out[..8].copy_from_slice(&id.replica.inner().to_be_bytes());
146 out[8..].copy_from_slice(&id.counter.to_be_bytes());
147 out
148}
149
150pub fn key_id(bytes: &[u8])
151 -> Outcome<OpId>
152{
153 if bytes.len() != KEY_LEN {
154 return Err(err!(
155 "A sketch key is {} bytes, and {} arrived.", KEY_LEN, bytes.len();
156 Decode, Input, Size, Mismatch));
157 }
158 let mut replica = [0u8; 8];
159 replica.copy_from_slice(&bytes[..8]);
160 let mut counter = [0u8; 8];
161 counter.copy_from_slice(&bytes[8..]);
162 Ok(OpId::new(
163 ReplicaId::new(u64::from_be_bytes(replica)),
164 u64::from_be_bytes(counter),
165 ))
166}
167
168/// One and a half cells per expected name, floored at [`MIN_CELLS`] and capped
169/// at [`MAX_CELLS`].
170pub fn cells_for(estimate: usize) -> usize {
171 // Three halves, rounded up, without leaving the integers.
172 let want = estimate.saturating_mul(3).saturating_add(1) / 2;
173 want.max(MIN_CELLS).min(MAX_CELLS)
174}
175
176/// The estimate a table of `cells` was sized from, which is [`cells_for`] read
177/// backwards.
178///
179/// Never narrower and never wider than it has to be, which are two requirements
180/// and not one. Narrower is unsafe: this is what a peer answering a table sizes
181/// its own from, and an answer a cell short of the table it answers puts back on
182/// the wire the stall it was sent to settle. Wider is merely wrong, and it costs
183/// a round trip rather than bytes: a peer that answers a table with one cell more
184/// than it holds is a peer the other end then answers again, because an arriving
185/// table wider than the last one sent is what a re-opening is read off.
186///
187/// So the division is rounded down and corrected upwards only where it fell
188/// short, which makes every reachable width a fixed point of the pair.
189pub fn estimate_for(cells: usize) -> usize {
190 let want = cells.saturating_mul(2) / 3;
191 match cells_for(want) < cells {
192 true => want.saturating_add(1),
193 false => want,
194 }
195}
196
197/// The estimate to sketch under when a table of `cells` could not be decoded:
198/// twice the cells, and `None` where that would pass [`GROW_CELLS`].
199///
200/// Doubling rather than guessing again. What a stall says is that the difference
201/// is larger than the table, and nothing in it says by how much -- the peeling
202/// decoder stops without knowing what it did not recover -- so the only honest
203/// answer is to try a size that cannot be reached by repeating the mistake.
204/// `None` is the answer the walk takes.
205pub fn grown(cells: usize) -> Option<usize> {
206 let want = cells.saturating_mul(2);
207 match want > GROW_CELLS {
208 true => None,
209 false => Some(estimate_for(want)),
210 }
211}
212
213/// How many cells a serialised table declares, without the table being built.
214pub fn cells_in(bytes: &[u8])
215 -> Outcome<usize>
216{
217 Ok(res!(Iblt::cells_in(bytes)))
218}
219
220pub fn config(estimate: usize, seed: u64) -> IbltConfig {
221 IbltConfig {
222 num_cells: cells_for(estimate),
223 num_hashes: HASHES,
224 key_len: KEY_LEN,
225 value_len: 0,
226 seed,
227 }
228}
229
230pub fn sketch(log: &OpLog, cfg: IbltConfig)
231 -> Outcome<Iblt>
232{
233 res!(check(&cfg));
234 let mut table = res!(Iblt::new(cfg));
235 for rec in log.iter() {
236 res!(table.insert(&key(&rec.id()), &[]));
237 }
238 Ok(table)
239}
240
241/// Sized from an estimate of the difference, not from the log.
242pub fn sketch_bytes(log: &OpLog, estimate: usize, seed: u64)
243 -> Outcome<Vec<u8>>
244{
245 Ok(res!(sketch(log, config(estimate, seed))).to_bytes())
246}
247
248/// The arriving table's shape is adopted, so peers that estimated differently
249/// still reconcile. A shape that is not a sketch of operation names at all is an
250/// error; a shape that is one but too large to work with is a fallback, since
251/// that is a judgement about effort rather than a malformed message.
252pub fn reconcile(log: &OpLog, remote: &[u8])
253 -> Outcome<Diff>
254{
255 let mut diff = res!(Iblt::from_bytes(remote));
256 let cfg = diff.config();
257 res!(check(&cfg));
258 if cfg.num_cells > MAX_CELLS {
259 return Ok(Diff::Undecodable(Fallback::Oversized { cells: cfg.num_cells }));
260 }
261 let mine = res!(sketch(log, cfg));
262 res!(diff.subtract(&mine));
263 match res!(diff.decode()) {
264 DecodeOutcome::Complete { inserted, deleted } => {
265 // Inserted is what the arriving table had and ours did not.
266 let mut remote_only = res!(names(&inserted));
267 let mut local_only = res!(names(&deleted));
268 remote_only.sort();
269 local_only.sort();
270 // A name we do not hold cannot be owed by us, whatever the table said.
271 local_only.retain(|id| log.contains(id));
272 Ok(Diff::Decoded { remote_only, local_only })
273 },
274 DecodeOutcome::Incomplete { inserted, deleted, remaining_cells } =>
275 Ok(Diff::Undecodable(Fallback::Incomplete {
276 remaining: remaining_cells,
277 recovered: inserted.len() + deleted.len(),
278 })),
279 }
280}
281
282
283/// Refuses a table that is not a sketch of operation names.
284fn check(cfg: &IbltConfig)
285 -> Outcome<()>
286{
287 if cfg.key_len != KEY_LEN {
288 return Err(err!(
289 "A sketch of operation names keys on {} bytes, and this table keys on {}.",
290 KEY_LEN, cfg.key_len;
291 Decode, Input, Size, Mismatch));
292 }
293 if cfg.value_len != 0 {
294 return Err(err!(
295 "A sketch of operation names carries no values, and this table carries {} \
296 bytes of them per cell.", cfg.value_len;
297 Decode, Input, Mismatch));
298 }
299 Ok(())
300}
301
302fn names(recovered: &[(Vec<u8>, Vec<u8>)])
303 -> Outcome<Vec<OpId>>
304{
305 let mut out = Vec::with_capacity(recovered.len());
306 for (k, _) in recovered {
307 out.push(res!(key_id(k)));
308 }
309 Ok(out)
310}
311
312
313#[cfg(test)]
314mod tests {
315 use super::*;
316
317 use crate::op::{
318 Header,
319 Op,
320 Record,
321 };
322
323 fn oid(replica: u64, counter: u64) -> OpId {
324 OpId::new(ReplicaId::new(replica), counter)
325 }
326
327 /// A log of `n` marks chained by one replica, starting at counter one.
328 fn chain(replica: u64, n: u64)
329 -> Outcome<OpLog>
330 {
331 let mut log = OpLog::new();
332 let r = ReplicaId::new(replica);
333 for i in 0..n {
334 res!(log.author(r, Op::Mark { name: fmt!("m{}", i), body: None, time: None }));
335 }
336 Ok(log)
337 }
338
339 #[test]
340 fn the_key_is_sixteen_fixed_bytes() -> Outcome<()> {
341 let id = oid(1, 7);
342 let k = key(&id);
343 assert_eq!(k.len(), KEY_LEN);
344 assert_eq!(&k[..8], &[0, 0, 0, 0, 0, 0, 0, 1]);
345 assert_eq!(&k[8..], &[0, 0, 0, 0, 0, 0, 0, 7]);
346 for id in [oid(0, 1), oid(1, 1), oid(u64::MAX, u64::MAX), oid(300, 70_000)] {
347 assert_eq!(res!(key_id(&key(&id))), id);
348 }
349 // Byte order is identifier order, which a varint spelling would not give.
350 assert!(key(&oid(1, 2)) < key(&oid(1, 3)));
351 assert!(key(&oid(1, 9)) < key(&oid(2, 1)));
352 // And a key of the wrong length is refused rather than padded.
353 assert!(key_id(&[0u8; 8]).is_err());
354 assert!(key_id(&[0u8; 17]).is_err());
355 Ok(())
356 }
357
358 /// With a floor under it and a ceiling over it.
359 #[test]
360 fn sizing_is_three_halves_of_the_estimate() -> Outcome<()> {
361 assert_eq!(cells_for(0), MIN_CELLS);
362 assert_eq!(cells_for(1), MIN_CELLS);
363 assert_eq!(cells_for(10), MIN_CELLS, "under the floor");
364 assert_eq!(cells_for(100), 150);
365 assert_eq!(cells_for(101), 152, "rounded up");
366 assert_eq!(cells_for(usize::MAX), MAX_CELLS, "and it never overflows");
367 let cfg = config(100, SEED);
368 assert_eq!(cfg.num_cells, 150);
369 assert_eq!(cfg.num_hashes, HASHES);
370 assert_eq!(cfg.key_len, KEY_LEN);
371 assert_eq!(cfg.value_len, 0);
372 Ok(())
373 }
374
375 /// Both halves of it, whichever side is asking.
376 #[test]
377 fn a_small_difference_decodes_whole() -> Outcome<()> {
378 // A shared prefix, then each side writes its own.
379 let mut a = res!(chain(1, 20));
380 let mut b = a.clone();
381 let r1 = ReplicaId::new(1);
382 let r2 = ReplicaId::new(2);
383 let mut only_a = Vec::new();
384 let mut only_b = Vec::new();
385 for i in 0..3 {
386 only_a.push(res!(a.author(r1, Op::Mark { name: fmt!("a{}", i), body: None, time: None })).id());
387 only_b.push(res!(b.author(r2, Op::Mark { name: fmt!("b{}", i), body: None, time: None })).id());
388 }
389 let from_b = res!(sketch_bytes(&b, 8, SEED));
390 match res!(reconcile(&a, &from_b)) {
391 Diff::Decoded { remote_only, local_only } => {
392 assert_eq!(remote_only, { let mut v = only_b.clone(); v.sort(); v });
393 assert_eq!(local_only, { let mut v = only_a.clone(); v.sort(); v });
394 },
395 Diff::Undecodable(f) => return Err(err!(
396 "A difference of six decoded incompletely: {}.", f.why(); Test)),
397 }
398 // And the other way round, which is the same computation mirrored.
399 let from_a = res!(sketch_bytes(&a, 8, SEED));
400 match res!(reconcile(&b, &from_a)) {
401 Diff::Decoded { remote_only, local_only } => {
402 assert_eq!(remote_only, { let mut v = only_a.clone(); v.sort(); v });
403 assert_eq!(local_only, { let mut v = only_b; v.sort(); v });
404 },
405 Diff::Undecodable(f) => return Err(err!(
406 "A difference of six decoded incompletely: {}.", f.why(); Test)),
407 }
408 Ok(())
409 }
410
411 #[test]
412 fn no_difference_decodes_to_nothing() -> Outcome<()> {
413 let a = res!(chain(1, 30));
414 let b = a.clone();
415 match res!(reconcile(&a, &res!(sketch_bytes(&b, 4, SEED)))) {
416 Diff::Decoded { remote_only, local_only } => {
417 assert!(remote_only.is_empty());
418 assert!(local_only.is_empty());
419 },
420 Diff::Undecodable(f) => return Err(err!(
421 "Two identical logs failed to decode: {}.", f.why(); Test)),
422 }
423 Ok(())
424 }
425
426 /// Rather than handing back the part of the difference it got to.
427 #[test]
428 fn an_undersized_sketch_stalls_and_says_so() -> Outcome<()> {
429 let a = res!(chain(1, 200));
430 let b = res!(chain(2, 200));
431 // Disjoint histories: four hundred names of difference, sixteen cells.
432 match res!(reconcile(&a, &res!(sketch_bytes(&b, 0, SEED)))) {
433 Diff::Decoded { remote_only, local_only } => return Err(err!(
434 "A sketch of {} cells decoded a difference of 400, as {} and {}.",
435 MIN_CELLS, remote_only.len(), local_only.len(); Test)),
436 Diff::Undecodable(Fallback::Incomplete { remaining, .. }) => {
437 assert!(remaining > 0);
438 },
439 Diff::Undecodable(other) => return Err(err!(
440 "Expected a stalled decode, got {}.", other.why(); Test)),
441 }
442 Ok(())
443 }
444
445 /// Which is what lets peers that estimated differently still reconcile.
446 #[test]
447 fn a_receiver_adopts_the_arriving_shape() -> Outcome<()> {
448 let mut a = res!(chain(1, 40));
449 let b = a.clone();
450 res!(a.author(ReplicaId::new(1), Op::Mark { name: fmt!("extra"), body: None, time: None }));
451 // B sketches generously; A would have sketched tightly.
452 let from_b = res!(sketch_bytes(&b, 500, SEED));
453 assert!(res!(Iblt::from_bytes(&from_b)).config().num_cells == 750);
454 match res!(reconcile(&a, &from_b)) {
455 Diff::Decoded { remote_only, local_only } => {
456 assert!(remote_only.is_empty());
457 assert_eq!(local_only.len(), 1);
458 },
459 Diff::Undecodable(f) => return Err(err!(
460 "A generous sketch failed to decode: {}.", f.why(); Test)),
461 }
462 // A seed of the sender's choosing is adopted along with the shape.
463 let odd = res!(sketch_bytes(&b, 8, 0x1234));
464 assert!(matches!(res!(reconcile(&a, &odd)), Diff::Decoded { .. }));
465 Ok(())
466 }
467
468 /// A table that is not a sketch of operation names is an error; one that is,
469 /// but is too big to work with, is a fallback.
470 #[test]
471 fn a_table_of_the_wrong_shape_is_refused() -> Outcome<()> {
472 let log = res!(chain(1, 5));
473 let wrong = IbltConfig {
474 num_cells: MIN_CELLS,
475 num_hashes: HASHES,
476 key_len: 8,
477 value_len: 0,
478 seed: SEED,
479 };
480 let table = res!(Iblt::new(wrong));
481 assert!(reconcile(&log, &table.to_bytes()).is_err(), "keyed on eight bytes");
482 let valued = IbltConfig { key_len: KEY_LEN, value_len: 4, ..wrong };
483 let table = res!(Iblt::new(valued));
484 assert!(reconcile(&log, &table.to_bytes()).is_err(), "carrying values");
485 assert!(sketch(&log, valued).is_err());
486 // Rubbish where a table should be.
487 assert!(reconcile(&log, b"not a sketch").is_err());
488 Ok(())
489 }
490
491 /// Which is the whole reason to send one.
492 #[test]
493 fn the_cost_follows_the_estimate_not_the_log() -> Outcome<()> {
494 let small = res!(chain(1, 10));
495 let large = res!(chain(1, 1000));
496 let a = res!(sketch_bytes(&small, 8, SEED)).len();
497 let b = res!(sketch_bytes(&large, 8, SEED)).len();
498 assert_eq!(a, b, "a hundredfold more history, the same sketch");
499 // And the per-cell cost is the key, the fingerprint and the count.
500 assert_eq!(a, 40 + MIN_CELLS * (KEY_LEN + 8 + 4));
501 Ok(())
502 }
503
504 /// Because it is the complement of what both hold.
505 #[test]
506 fn a_decoded_difference_is_already_closed() -> Outcome<()> {
507 use crate::sync::walk::close;
508 use std::collections::BTreeSet;
509
510 let mut a = OpLog::new();
511 let r1 = ReplicaId::new(1);
512 let r2 = ReplicaId::new(2);
513 for i in 0..6 {
514 res!(a.author(r1, Op::Mark { name: fmt!("shared{}", i), body: None, time: None }));
515 }
516 let mut b = a.clone();
517 // A merge on each side, so the divergence is not a straight line.
518 for i in 0..3 {
519 res!(a.author(r1, Op::Mark { name: fmt!("a{}", i), body: None, time: None }));
520 res!(b.author(r2, Op::Mark { name: fmt!("b{}", i), body: None, time: None }));
521 }
522 res!(a.append(Record::new(
523 res!(Header::new(oid(1, 20), a.frontier())),
524 Op::Mark { name: fmt!("merge"), body: None, time: None },
525 )));
526 let local_only = match res!(reconcile(&a, &res!(sketch_bytes(&b, 16, SEED)))) {
527 Diff::Decoded { local_only, .. } => local_only,
528 Diff::Undecodable(f) => return Err(err!(
529 "The difference failed to decode: {}.", f.why(); Test)),
530 };
531 let held: BTreeSet<OpId> = a.iter()
532 .map(|rec| rec.id())
533 .filter(|id| !local_only.contains(id))
534 .collect();
535 let closed = close(&a, &local_only, &held);
536 assert_eq!(
537 closed.into_iter().collect::<Vec<_>>(), local_only,
538 "closing a decoded difference added something",
539 );
540 Ok(())
541 }
542}