Oregami
Repositories/oxedyne/fe2o3

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

23.7 KiB, 155 runs

created by r1870400018:19657, 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//! The driver: one message in, some messages out, and where things stand.
2//!
3//! A session holds no log and no connection. The caller passes the log it wants
4//! brought up to date on each call, reads the messages the session hands back,
5//! and puts them wherever messages go. Looping that over a pipe is the whole of
6//! a sync:
7//!
8//! ```text
9//! let mut s = Session::new(Mode::Walk);
10//! send(res!(s.open(&log)));
11//! while !s.is_converged() {
12//! let turn = res!(s.receive(&mut log, recv()));
13//! for msg in turn.send { send(msg); }
14//! }
15//! ```
16//!
17//! # Either peer may start
18//!
19//! A session that is handed an opening before it has made one answers with its
20//! own, so a peer that was called upon needs no separate path: it constructs a
21//! session and feeds it what arrived. Two peers that both open at once do not
22//! open twice. There is no client and no server here, only two sessions running
23//! the same code.
24//!
25//! # Convergence
26//!
27//! A session is converged when it has said everything it owes and heard the
28//! other side say the same. That is a statement about the conversation. It means
29//! the logs agree because the owed set was computed honestly at both ends, which
30//! the closure check on arrival is what enforces: a peer that sends a batch with
31//! a hole in it has its batch refused whole, and the session errs rather than
32//! absorbing part of it.
33//!
34//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
35//! Anthropic Claude
36
37use crate::id::OpId;
38use crate::log::OpLog;
39use crate::sync::msg::Message;
40use crate::sync::sketch::{
41 cells_for,
42 cells_in,
43 estimate_for,
44 grown,
45 reconcile,
46 sketch_bytes,
47 Diff,
48 Fallback,
49 GROW_CELLS,
50 SEED,
51};
52use crate::sync::walk::{
53 arrival_gap,
54 close,
55 covered,
56 entries_for,
57 owed,
58};
59
60use oxedyne_fe2o3_core::prelude::*;
61
62use std::collections::BTreeSet;
63
64
65/// How a session opens, and how it works out what it owes.
66#[derive(Clone, Copy, Debug, Eq, PartialEq)]
67pub enum Mode {
68 // Exchange frontiers, and send everything the other's frontier does not
69 // cover. Correct at any divergence, and loose where both sides have written.
70 Walk,
71 // Exchange sketches, and send the difference. Falls back to the walk, in the
72 // same turn, where the sketch turns out to have been too small: an estimate
73 // too low costs a fallback, one too high costs sketch bytes. A receiver
74 // adopts the seed it is sent, so the seed decides only what this peer opens
75 // with, and SEED is the answer where the caller has no reason to prefer
76 // another.
77 Sketch {
78 estimate: usize, // operations the two logs are guessed to differ by
79 seed: u64, // the seed both peers' tables are built under
80 },
81}
82
83// How many operations a head of a frontier is taken to stand for when the
84// difference between two logs is being guessed at. A log that has diverged
85// carries more than one head, and each head is a branch somebody wrote since the
86// two last spoke. How much they wrote is exactly what nobody knows, so this is a
87// guess with a fallback under it.
88pub const FANOUT: usize = 8;
89
90impl Mode {
91
92 /// Under the usual seed.
93 pub fn sketch(estimate: usize) -> Self {
94 Self::Sketch { estimate, seed: SEED }
95 }
96
97 /// The guess at the difference, and what it rests on is which of the other
98 /// peer's heads this log does not hold.
99 ///
100 /// **A head of theirs we hold is the whole answer, not a term in it.** A log
101 /// is causally closed and its frontier dominates it, so holding every head of
102 /// theirs means holding everything they hold: the difference is then exactly
103 /// the two logs' difference in length and nothing stands for branches that are
104 /// not there. That is the pull, the clone and the fetch after somebody else
105 /// pushed, and it is the common case.
106 ///
107 /// A head of theirs we do not hold says the divergence is two-sided, and a
108 /// length delta then measures nothing: two logs that have each written since
109 /// they last spoke can be the same length and share almost nothing. So each
110 /// frontier's heads are taken to stand for [`FANOUT`] operations apiece and the
111 /// delta rides on top of that, being evidence that one of the two tails is
112 /// longer than the other.
113 ///
114 /// It is still a guess, and a guess that is low costs a round trip rather than
115 /// a history: [`crate::sync::sketch::grown`] is what a stalled decode reaches
116 /// for before the walk.
117 ///
118 /// Where the guess comes to as much as the smaller log, sketching is pointless
119 /// -- a sketch sized for the whole history costs more than the history -- and
120 /// the walk is both cheaper and exact. That covers the clone case, where one
121 /// log is empty, and the case of two logs that share nothing.
122 ///
123 /// A guess of nothing means one log holds what the other does, which the walk
124 /// settles in one message; a sketch of nothing would be a table nobody needs.
125 ///
126 /// Both peers may call this and neither has to: the modes need not agree,
127 /// since a session answers whatever it is given.
128 ///
129 /// # Arguments
130 ///
131 /// * `unknown` - Heads of the other peer's frontier this log does not hold. A
132 /// caller who has their count and not their names has to pass that count,
133 /// which is the reading that assumes the worst.
134 pub fn between(
135 here_len: usize, // operations this peer holds
136 here_heads: usize, // heads of this peer's frontier
137 there_len: usize, // operations the other peer holds
138 unknown: usize, // heads of theirs this peer does not hold
139 )
140 -> Self
141 {
142 let delta = here_len.abs_diff(there_len);
143 let spread = match unknown {
144 0 => delta,
145 _ => delta + FANOUT * (here_heads + unknown),
146 };
147 if spread == 0 || spread >= here_len.min(there_len).max(1) {
148 Self::Walk
149 } else {
150 Self::sketch(spread)
151 }
152 }
153}
154
155
156/// Whether this end takes part in growing a sketch.
157///
158/// Three behaviours and one switch, because they are one mechanism: a stalled
159/// decode is answered with a larger table, an arriving table wider than the last
160/// one sent is answered by opening again at it, and the widest shape met is
161/// carried into the session after ([`Session::sizing`]). Refusing turns off all
162/// three, which is a peer built before any of it -- and that is the second thing
163/// this is for, since a carrier's older-peer behaviour can then be exercised
164/// against a real session rather than reasoned about.
165///
166/// **What it is set FROM is a fact about the other end: whether it re-opens.** A
167/// grown table is a question, and the answer to it is the peer opening again at
168/// the new shape; a peer that does not takes the grown table, answers what it
169/// owes, and is never given what it is owed, because the end that grew said
170/// nothing else and the exchange has nowhere to go. So an end that cannot tell
171/// refuses, and the walk answers instead, which every peer there has ever been
172/// understands.
173///
174/// The two ends need not agree, and in the ordinary case they do not: a client
175/// grows freely against a carrier that opens afresh on every request -- and so
176/// always answers an opening -- while that carrier grows only against a client
177/// that told it what it speaks.
178#[derive(Clone, Copy, Debug, Eq, PartialEq)]
179pub enum Growth {
180 // A stalled decode is answered with a larger table, a wider table with a
181 // wider answer, and the width reached is carried into the next session.
182 Allowed,
183 // None of the three. Both what a peer that opens once must be answered with,
184 // and what a build from before any of this does.
185 Refused,
186}
187
188
189/// Where a session stands after taking a message.
190#[derive(Clone, Copy, Debug, Eq, PartialEq)]
191pub enum Step {
192 Converged, // both sides have said everything they owe, so the logs agree
193 NeedMore, // more is expected from the other side
194 // A sketch could not be decoded, so the walk answered instead. The exchange
195 // carries on and will converge; the reason is worth a caller's attention only
196 // in that it says the estimate was low.
197 FellBack(Fallback),
198 // A sketch could not be decoded and a larger one was answered with, rather
199 // than the walk. Nothing was handed over on this turn: the answer is the
200 // bigger table alone, and what is owed crosses once a table decodes.
201 Grew {
202 cells: usize, // cells the table this end answered with declares
203 },
204}
205
206
207/// What a session hands back: what to send, and where things stand.
208#[derive(Clone, Debug, Eq, PartialEq)]
209pub struct Turn {
210 pub send: Vec<Message>, // to put on the wire, in order
211 pub step: Step, // where the session stands now
212}
213
214
215/// One side of an exchange.
216#[derive(Clone, Debug)]
217pub struct Session {
218 mode: Mode, // how this side opens and works out what it owes
219 opened: bool, // an opening has been made
220 told: bool, // everything owed has been sent
221 heard: bool, // the other side has said it is finished
222 fell_back: Option<Fallback>, // the first fallback, kept for the asking
223 known: BTreeSet<OpId>, // operations the far end has shown it holds
224 sent: usize, // operations handed over
225 absorbed: usize, // operations absorbed
226 offered: usize, // cells of the widest table this end has sent
227 widest: usize, // cells of the widest table met, sent or received
228 growth: Growth, // whether a stalled decode may be answered with a larger table
229}
230
231impl Session {
232
233 pub fn new(mode: Mode) -> Self {
234 Self::knowing(mode, BTreeSet::new())
235 }
236
237 /// A session that starts knowing the far end already holds `known`.
238 ///
239 /// The one thing a session boundary would otherwise throw away. No message says
240 /// "you handed me this": a frontier reports what its own end holds, and a head
241 /// nobody here has seen subtracts nothing, so a fresh session works out what it
242 /// owes from the frontier alone and offers back everything the sessions before
243 /// it took in. A carrier that runs several sessions -- which a bounded reply
244 /// makes it do -- hands one session's [`Session::known`] to the next, and the
245 /// re-offer stops.
246 ///
247 /// It is a proof and not a claim. Every identifier in it is one this end watched
248 /// arrive, or one the carrier watched leave and be taken; nothing in it rests on
249 /// what a peer says about itself, which is what makes it safe to subtract.
250 pub fn knowing(mode: Mode, known: BTreeSet<OpId>) -> Self {
251 Self {
252 mode,
253 opened: false,
254 told: false,
255 heard: false,
256 fell_back: None,
257 known,
258 sent: 0,
259 absorbed: 0,
260 offered: 0,
261 widest: 0,
262 growth: Growth::Allowed,
263 }
264 }
265
266 /// A session that takes no part in growing a sketch.
267 ///
268 /// For a caller that knows what the far end speaks, which over a carrier is
269 /// the carrier's business and not this module's. See [`Growth`].
270 pub fn with_growth(mut self, growth: Growth) -> Self {
271 self.growth = growth;
272 self
273 }
274
275 pub fn mode(&self) -> Mode {
276 self.mode
277 }
278
279 /// Have both sides finished, which is when the logs agree?
280 pub fn is_converged(&self) -> bool {
281 self.told && self.heard
282 }
283
284 /// The first fallback, and sticky, so a caller can ask once at the end rather
285 /// than watching every turn.
286 pub fn fell_back(&self) -> Option<Fallback> {
287 self.fell_back
288 }
289
290 /// The mode a session following this one should open in.
291 ///
292 /// The other thing a session boundary would otherwise throw away, and it is
293 /// [`Session::knowing`]'s companion. A table that had to be grown to before
294 /// anything decoded is the shape that works, and a carrier running several
295 /// sessions -- which a bounded reply makes it do -- would otherwise open each
296 /// of them at the estimate the first one was handed, stall, and grow again.
297 ///
298 /// A walk stays a walk. It was chosen because the difference was judged to be
299 /// as large as the history, and nothing a sketch did in passing changes that.
300 /// So does a session that refuses growth, which is one that behaves as a build
301 /// from before any of this: it opens every session at the estimate it was
302 /// handed.
303 pub fn sizing(&self) -> Mode {
304 if self.growth == Growth::Refused {
305 return self.mode;
306 }
307 match self.mode {
308 Mode::Walk => Mode::Walk,
309 Mode::Sketch { estimate, seed } => {
310 // At the same ceiling growth climbs to, and for the same reason: a
311 // table crosses in one message and nothing cuts one up, so a shape
312 // carried into the next session is a shape that has to fit the body
313 // this end posts. A peer may open wider than this end would have
314 // chosen -- MAX_CELLS is what a reader will allocate, not what a
315 // carrier will take -- and that is where an uncapped carry would put
316 // an oversized request on the wire.
317 let want = self.widest.min(GROW_CELLS);
318 match want > cells_for(estimate) {
319 true => Mode::Sketch { estimate: estimate_for(want), seed },
320 false => self.mode,
321 }
322 },
323 }
324 }
325
326 /// What the far end has shown it holds: what this session was told when it
327 /// opened, and every operation that has arrived in a [`Message::Send`] since.
328 ///
329 /// Not what it was sent. Whether a batch reached the far end is a question about
330 /// the carrier and not about the session, so a carrier that knows its request
331 /// was taken puts those identifiers in itself.
332 pub fn known(&self) -> &BTreeSet<OpId> {
333 &self.known
334 }
335
336 pub fn ops_sent(&self) -> usize {
337 self.sent
338 }
339
340 pub fn ops_absorbed(&self) -> usize {
341 self.absorbed
342 }
343
344 /// A session that is fed an opening before it has made one opens on its own
345 /// account, so calling this is a choice about who speaks first and not a
346 /// requirement.
347 pub fn open(&mut self, log: &OpLog)
348 -> Outcome<Message>
349 {
350 self.opened = true;
351 match self.mode {
352 Mode::Walk => Ok(Message::hello(log.frontier())),
353 Mode::Sketch { estimate, seed } => {
354 self.offered = self.offered.max(cells_for(estimate));
355 self.widest = self.widest.max(self.offered);
356 Ok(Message::sketch(
357 log.frontier(),
358 res!(sketch_bytes(log, estimate, seed)),
359 log.len() as u64,
360 ))
361 },
362 }
363 }
364
365 /// A [`Message::Send`] is checked for causal closure against the log before
366 /// anything is absorbed, and refused whole if it has a hole in it. Operations
367 /// the log already holds, and repetitions within one batch, are dropped
368 /// rather than refused: a peer that could not subtract one of our heads sends
369 /// more than it needs to, and that is a cost rather than a fault.
370 pub fn receive(&mut self, log: &mut OpLog, msg: Message)
371 -> Outcome<Turn>
372 {
373 match msg {
374 Message::Hello { heads } => self.answer(log, &heads, None, 0),
375 Message::Sketch { heads, cells, count } =>
376 self.answer(log, &heads, Some(&cells), count as usize),
377 Message::Send { entries } => {
378 if let Some((id, parent)) = res!(arrival_gap(log, &entries)) {
379 return Err(err!(
380 "An arriving batch of {} operation{} names the parent {} of {}, \
381 which neither the batch nor the log holds; a peer sends causal \
382 closures and not subsets.", entries.len(),
383 if entries.len() == 1 { "" } else { "s" }, parent, id;
384 Invalid, Input, Missing, Order));
385 }
386 let mut batch = Vec::with_capacity(entries.len());
387 let mut taken: BTreeSet<OpId> = BTreeSet::new();
388 for entry in &entries {
389 let rec = res!(entry.peek());
390 let id = rec.id();
391 // Written down whether or not it is news here. An operation a peer
392 // hands over is an operation that peer holds, and that is the whole of
393 // what a later session needs to stop offering it back.
394 self.known.insert(id);
395 if !log.contains(&id) && taken.insert(id) {
396 batch.push(rec);
397 }
398 }
399 let placed = batch.len();
400 let left = res!(log.absorb(batch));
401 self.absorbed += placed - left.len();
402 if !left.is_empty() {
403 return Err(err!(
404 "An arriving batch closed causally and yet left {} operation{} \
405 unplaced, starting at {}.", left.len(),
406 if left.len() == 1 { "" } else { "s" }, left[0].id();
407 Bug, Unreachable));
408 }
409 Ok(self.turn(Vec::new(), None))
410 },
411 Message::Done => {
412 self.heard = true;
413 Ok(self.turn(Vec::new(), None))
414 },
415 // A piece of an operation is not an operation, and a session that
416 // tried to place one would be placing something nobody signed. Putting
417 // the pieces back together is the carrier's work -- `Parts` in
418 // `sync::msg` -- and a session is handed the result.
419 Message::Part { id, seq, total, .. } => Err(err!(
420 "A session was handed piece {} of the {} the operation {} was cut \
421 into. A piece is a transport's business: `sync::msg::Parts` puts a \
422 run of them back into the send it was cut from, and that is what a \
423 session receives.", seq, total, id;
424 Invalid, Input, Mismatch)),
425 // How far the far end was carried into what this log last owed it.
426 // Written down as operations it holds, because that is what it is: the
427 // prefix is subtracted by the same walk that subtracts a head, and a
428 // carrier that keeps nothing between sessions is thereby told where to
429 // carry on from without being told what it already sent.
430 Message::Resume { at } => {
431 self.resume(log, &at);
432 Ok(self.turn(Vec::new(), None))
433 },
434 // A replacement names operations the receiver already holds and asks
435 // for the form they are held in to change, which is a fact about a
436 // store and not about a set of operations. A session reconciles sets,
437 // so it neither sends this nor answers it; a carrier acts on it once
438 // the session it arrived beside has finished.
439 Message::Forgotten { entries } => Err(err!(
440 "A session was handed {} record{} to write in place of operations it \
441 already holds. What a forget takes out of a store is the carrier's \
442 work -- `Message::replacements` is what hands them over -- and a \
443 session places operations and never rewrites one.",
444 entries.len(), if entries.len() == 1 { "" } else { "s" };
445 Invalid, Input, Mismatch)),
446 }
447 }
448
449 /// Answers an opening: what we owe, and our own opening if we have not made
450 /// one -- or if the one we made was narrower than the table we are answering.
451 fn answer(&mut self, log: &OpLog, heads: &[OpId], cells: Option<&[u8]>, theirs: usize)
452 -> Outcome<Turn>
453 {
454 // The shape the other end sketched under, read out of the table's header
455 // rather than out of the table.
456 let arriving = match cells {
457 Some(bytes) => res!(cells_in(bytes)),
458 None => 0,
459 };
460 self.widest = self.widest.max(arriving);
461 let taken = match cells {
462 Some(bytes) => Some(res!(reconcile(log, bytes))),
463 None => None,
464 };
465 // A stall is answered with a larger table before it is answered with the
466 // walk. What the walk would say here is the whole log: a peer whose head
467 // this end does not hold has nothing subtracted by it, and handing over a
468 // history to save a round trip is the trade the sketch exists to refuse.
469 // Nothing is owed on this turn and nothing is said to be finished; the
470 // bigger table goes on its own and what it decodes crosses next.
471 if let Some(Diff::Undecodable(Fallback::Incomplete { .. })) = &taken {
472 if let Some(turn) = res!(self.growing(log, theirs)) {
473 return Ok(turn);
474 }
475 }
476 // Never a narrower table than the one it answers. Both peers build under
477 // the arriving shape, so a shape this end could decode is one the other
478 // end can decode too, and answering it with less would stall over there
479 // and cost the round trip this turn just saved.
480 self.widen(arriving);
481 let mut out = Vec::new();
482 if !self.opened || (self.growth == Growth::Allowed && arriving > self.offered) {
483 out.push(res!(self.open(log)));
484 }
485 // What we owe, and what we believe they hold, which is what the send set
486 // is closed against.
487 let mut fallback: Option<Fallback> = None;
488 let (send, held) = match taken {
489 Some(Diff::Decoded { local_only, .. }) => {
490 // Everything we hold that is not ours alone, they hold too.
491 let owed_set: BTreeSet<OpId> = local_only.iter().copied().collect();
492 let held: BTreeSet<OpId> = log.iter()
493 .map(|rec| rec.id())
494 .filter(|id| !owed_set.contains(id))
495 .collect();
496 (local_only, held)
497 },
498 Some(Diff::Undecodable(reason)) => {
499 fallback = Some(reason);
500 if self.fell_back.is_none() {
501 self.fell_back = Some(reason);
502 }
503 (owed(log, heads, &self.known), covered(log, heads, &self.known))
504 },
505 None => (owed(log, heads, &self.known), covered(log, heads, &self.known)),
506 };
507 let ids = close(log, &send, &held);
508 if !ids.is_empty() {
509 let entries = res!(entries_for(log, &ids));
510 self.sent += entries.len();
511 out.push(Message::Send { entries });
512 }
513 out.push(Message::Done);
514 self.told = true;
515 Ok(self.turn(out, fallback.map(Step::FellBack)))
516 }
517
518 /// Opens again with a table of twice the cells, where there is one worth
519 /// having.
520 ///
521 /// `None` is every reason not to, and each of them leaves the walk to answer,
522 /// which it always can. The far end does not re-open, so a grown table would
523 /// be the last thing said -- see [`Growth`], and it is the reason this is a
524 /// knob at all. This end is walking anyway, so a sketch is not what it chose.
525 /// The double would pass [`crate::sync::sketch::GROW_CELLS`], which is what a
526 /// table has to cross in one message inside. Or the size it would climb to is
527 /// one [`Mode::between`] would not have chosen in the first place: a table
528 /// sized for as much as the smaller of the two logs costs more than that log,
529 /// and growth is what buys time against the walk rather than a substitute for
530 /// it.
531 ///
532 /// Every growth at least doubles, so the climb is geometric and each of those
533 /// three bounds is reached in a handful of turns. A session after this one
534 /// opens at [`Session::sizing`], so a bounded exchange keeps the size it
535 /// reached rather than starting again at the first estimate.
536 ///
537 /// # Arguments
538 ///
539 /// * `theirs` - Operations the other peer says it holds. It decides nothing
540 /// but how far this climbs, and a peer that overstates it is bounded by this
541 /// log's own length.
542 fn growing(&mut self, log: &OpLog, theirs: usize)
543 -> Outcome<Option<Turn>>
544 {
545 if self.growth == Growth::Refused {
546 return Ok(None);
547 }
548 let seed = match self.mode {
549 Mode::Sketch { seed, .. } => seed,
550 Mode::Walk => return Ok(None),
551 };
552 let estimate = match grown(self.widest) {
553 Some(e) => e,
554 None => return Ok(None),
555 };
556 if estimate >= log.len().min(theirs).max(1) {
557 return Ok(None);
558 }
559 self.mode = Mode::Sketch { estimate, seed };
560 let msg = res!(self.open(log));
561 Ok(Some(self.turn(vec![msg], Some(Step::Grew { cells: self.offered }))))
562 }
563
564 /// Raises this end's estimate so that the table it opens with is at least
565 /// `cells` wide, up to [`crate::sync::sketch::GROW_CELLS`], which is what a
566 /// carrier will take in one body.
567 ///
568 /// A walk is left alone: it was chosen because the difference was judged to be
569 /// as large as the history. So is a session that refuses growth, which is one
570 /// that does not open twice.
571 fn widen(&mut self, cells: usize) {
572 if self.growth == Growth::Refused {
573 return;
574 }
575 if let Mode::Sketch { estimate, seed } = self.mode {
576 let want = cells.min(GROW_CELLS);
577 if cells_for(estimate) < want {
578 self.mode = Mode::Sketch { estimate: estimate_for(want), seed };
579 }
580 }
581 }
582
583 /// Writes down that the far end holds every operation of this log up to and
584 /// including `at`, in the log's append order.
585 ///
586 /// Sound because of what the far end can only have come by it: everything
587 /// this end sent it went in append order, and everything before it that was
588 /// not sent was subtracted as already held. An append-order prefix is closed
589 /// downwards by the log's own append guard, so what is written down here is
590 /// closed downwards too, which is what makes it safe to subtract.
591 ///
592 /// An identifier this log does not hold says nothing and is dropped, exactly
593 /// as a head nobody here has seen is.
594 fn resume(&mut self, log: &OpLog, at: &OpId) {
595 let end = match log.position(at) {
596 Some(p) => p,
597 None => return,
598 };
599 for rec in log.iter().take(end + 1) {
600 self.known.insert(rec.id());
601 }
602 }
603
604 /// A step made on this turn is reported ahead of anything else.
605 fn turn(&self, send: Vec<Message>, step: Option<Step>) -> Turn {
606 let step = match step {
607 Some(said) => said,
608 None => if self.is_converged() {
609 Step::Converged
610 } else {
611 Step::NeedMore
612 },
613 };
614 Turn { send, step }
615 }
616}