Oregami
Repositories/oxedyne/fe2o3

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

58.1 KiB, 256 runs

created by r1870400018:19639, 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 one peer says to another.
2//!
3//! A handful of messages, and both peers may send every one of them. A message
4//! has a daticle form, for a caller that keeps its wire in daticles, and a byte
5//! form that begins with a magic and a version, for a caller that keeps a wire
6//! of bytes. The version is there because these bytes cross between machines
7//! that were built at different times, and a reader that cannot tell an old
8//! spelling from a new one will eventually mistake one for the other.
9//!
10//! # Framing belongs to the transport
11//!
12//! [`Message::decode`] reads a message that occupies the whole of the buffer it
13//! is given. Where one message ends and the next begins is a question the
14//! carrier already answers -- a datagram has a length, a stream has whatever
15//! framing it was given -- and answering it twice is how the two answers come to
16//! disagree.
17//!
18//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
19//! Anthropic Claude
20
21use crate::id::OpId;
22use crate::segment::Entry;
23
24use oxedyne_fe2o3_core::prelude::*;
25use oxedyne_fe2o3_jdat::prelude::*;
26
27
28/// The bytes every sync message begins with.
29pub const MAGIC: [u8; 6] = *b"ORESYN";
30
31/// The newest format version this module writes.
32///
33/// Raised to 2 when the vocabulary gained [`Message::Part`], a piece of an
34/// operation too large for the carrier to take whole. The four older kinds did
35/// not move a byte and are still stamped [`VERSION_MIN`], because a message
36/// declares the version it needs rather than the version its writer was built
37/// at: a peer that speaks only version 1 still reads every message a version 1
38/// peer could have sent, and refuses the one it could not have.
39///
40/// That is the whole of why the version moved rather than a kind byte being
41/// added to [`crate::segment::Entry`] beside [`crate::segment::KIND_PACKED`].
42/// A kind was right there because a segment stamps one version over many
43/// records, so a bump would condemn every plain record beside the new one. A
44/// message stamps its own version in its own bytes, so a bump condemns nothing
45/// it does not have to, and it puts the refusal at the header -- before a
46/// decode, naming both versions -- rather than inside `Entry::from_dat`, after
47/// the handshake has already said the frame was compatible.
48///
49/// Raised to 3 for [`Message::Forgotten`], by the same rule and for the same
50/// reason: a carrier's copy of a history cannot drop what a forget took unless
51/// somebody hands it the stubs to write in their place, and a carrier that
52/// cannot read the repository could never build them itself. A part is still
53/// stamped 2 -- see [`version_for`] -- so a peer built before any of this reads
54/// every message it could have read before.
55///
56/// Raised to 4 for [`Message::Resume`], and that one is about a carrier that
57/// keeps nothing. A bounded reply cuts a frontier walk into a run of sessions,
58/// and a carrier with no memory of the last one works out the same owed set and
59/// sends the same prefix again; the cursor is how the end that does remember
60/// tells it where to carry on from. A peer too old to read one is not refused
61/// anything: it is simply never sent one, and answers from the frontier as it
62/// always did.
63pub const VERSION: u8 = 4;
64
65/// The oldest format version this module reads, and the version a message is
66/// stamped with when it needs nothing newer.
67///
68/// It stays at 1 because every version since is a strict superset: the framing
69/// is identical, the four original kinds are spelled exactly as they were, and
70/// the only thing a version 1 reader cannot do is read a kind that did not exist
71/// when it was built. [`highest_kind`] is what keeps that true of the bytes and
72/// not merely of the intention.
73pub const VERSION_MIN: u8 = 1;
74
75// The kind byte each message is tagged with on the wire.
76pub const KIND_HELLO: u8 = 1;
77pub const KIND_SKETCH: u8 = 2;
78pub const KIND_SEND: u8 = 3;
79pub const KIND_DONE: u8 = 4;
80pub const KIND_PART: u8 = 5;
81pub const KIND_FORGOTTEN: u8 = 6;
82pub const KIND_RESUME: u8 = 7;
83
84/// Most bytes an operation may come to when its pieces are put back together.
85///
86/// A declared piece count is an instruction to allocate and it arrives from
87/// wherever the message did, so nothing is ever sized from it: [`Parts`] grows
88/// by what has actually arrived and refuses at this. The number is the largest
89/// single operation whose arrival costs no more than opening the repository it
90/// belongs to already costs -- an operation of *n* bytes is held about three
91/// times over while it is taken in, and fe2o3's whole history peaks at 263 MB --
92/// and it is the same sixty-four mebibytes a packed run may inflate to.
93pub const PART_MAX: usize = 64 << 20;
94
95
96/// The highest message kind a peer speaking the given format version may send.
97///
98/// The rule the vocabulary grows by, and it is the same rule
99/// [`crate::segment::highest_code`] states for operations: a new kind goes
100/// strictly above the existing ones, [`VERSION`] rises by one, this gains a
101/// branch, and [`VERSION_MIN`] stays where it is. A message declaring version 1
102/// and carrying kind 5 is refused rather than read, because a peer that would
103/// read it is a peer that would also have to guess at what else the sender
104/// thought version 1 meant.
105pub const fn highest_kind(version: u8) -> u8 {
106 match version {
107 0 | 1 => KIND_DONE,
108 2 => KIND_PART,
109 3 => KIND_FORGOTTEN,
110 _ => KIND_RESUME,
111 }
112}
113
114/// The oldest format version whose vocabulary spells the given kind.
115///
116/// The inverse of [`highest_kind`], and worked out from it rather than written
117/// down beside it, so that a kind added to one cannot be forgotten in the other.
118/// This is what stamps a message with the version it needs: a part still says 2
119/// after [`VERSION`] moved to 3, because nothing about a part changed and a peer
120/// that could read one before still can.
121pub const fn version_for(kind: u8) -> u8 {
122 let mut version = VERSION_MIN;
123 while version < VERSION && highest_kind(version) < kind {
124 version += 1;
125 }
126 version
127}
128
129
130/// One thing a peer says.
131///
132/// The fields are public because a caller may legitimately want to work on a
133/// message before it goes out -- to seal the records of a [`Message::Send`] with
134/// a key this crate knows nothing about, most obviously. The encoding is
135/// canonical whatever a caller builds: frontiers are written ascending and
136/// without repetition, and a decoder refuses anything else.
137#[derive(Clone, Debug, Eq, PartialEq)]
138pub enum Message {
139 // Opens a frontier walk.
140 Hello {
141 heads: Vec<OpId>, // the speaker's frontier
142 },
143 // Opens a reconciliation. The frontier rides along so that a receiver whose
144 // decode stalls can answer with the walk in the same turn, rather than
145 // spending a round trip asking for what it was already told. The count is
146 // advisory: it lets a peer judge whether the estimate the table was sized
147 // from was sensible, and nothing is decided by it.
148 Sketch {
149 heads: Vec<OpId>, // the speaker's frontier
150 cells: Vec<u8>, // the serialised table
151 count: u64, // operations the speaker's log holds
152 },
153 // Operations the speaker owes, causally closed against what the receiver
154 // holds. Order carries no meaning -- OpLog::absorb places a batch however it
155 // is shuffled -- but the sender writes them in a causal order anyway, so that
156 // one pass places them all.
157 Send {
158 entries: Vec<Entry>, // bare or sealed
159 },
160 Done,
161 // One piece of an operation the carrier cannot take whole. A part is not an
162 // operation and never becomes one on its own: nothing places it, nothing
163 // signs it, no log holds it and no segment may carry it. What it carries is
164 // a slice of the bytes an [`Entry`] encodes to, and the entry those slices
165 // make is byte for byte the entry that was signed, so a signature crosses a
166 // carrier that had to cut it up exactly as it crosses one that did not.
167 //
168 // The pieces of one operation arrive in order and are put back together by
169 // [`Parts`], which is the transport's work and not a session's.
170 Part {
171 id: OpId, // the operation the pieces make
172 seq: u64, // which piece this is, counted from zero
173 total: u64, // how many pieces there are
174 bytes: Vec<u8>, // this piece of the entry's encoded form
175 },
176 // What to write in place of the operations a [`crate::op::Op::Forget`]
177 // names, so that a peer holding this history stops holding their bytes.
178 // Every identifier here is one the receiver already holds and nothing
179 // absorbs any of it: what changes is the form a record is held in, not the
180 // set of operations there are, so a session neither sends this nor answers
181 // it and a carrier acts on it after the session has finished.
182 //
183 // The SENDER builds the stubs, always. A carrier of a veiled repository can
184 // read neither the forget nor the record it names, so it could not build
185 // one if it wanted to; a carrier of a plain one could, and is sent them
186 // anyway, because two ways of arriving at the same record is one more than
187 // the act needs.
188 Forgotten {
189 entries: Vec<Entry>, // each a Forgotten record under the original's header
190 },
191 // How far into what the receiver last owed this speaker the speaker was
192 // carried, so that a receiver keeping nothing between sessions does not begin
193 // again. A bounded reply cuts a frontier walk into a run of sessions, and an
194 // owed set is worked out afresh in each of them: where the speaker's own
195 // frontier does not move -- which is what a pull-only mirror's never does --
196 // the same prefix is sent every time and the exchange makes no progress at
197 // all.
198 //
199 // It names an operation and not a count, because a count is a position in a
200 // set that is recomputed and an identifier is a position in the log, which is
201 // not. What it claims is that the speaker holds every operation of the
202 // receiver's log up to and including this one, and that is a claim the
203 // receiver's own earlier answers make true: everything before it was either
204 // sent to the speaker or already subtracted as held.
205 //
206 // It rides beside an opening rather than inside one, so that the four
207 // original kinds are still spelled exactly as they were and a peer that never
208 // learned this one is simply never sent it.
209 Resume {
210 at: OpId, // the last operation of the receiver's owed set that arrived
211 },
212}
213
214impl Message {
215
216 /// Puts the frontier in canonical order.
217 pub fn hello(heads: Vec<OpId>) -> Self {
218 Self::Hello { heads: canonical(heads) }
219 }
220
221 /// Puts the frontier in canonical order.
222 pub fn sketch(heads: Vec<OpId>, cells: Vec<u8>, count: u64) -> Self {
223 Self::Sketch { heads: canonical(heads), cells, count }
224 }
225
226 pub fn kind(&self) -> u8 {
227 match self {
228 Self::Hello { .. } => KIND_HELLO,
229 Self::Sketch { .. } => KIND_SKETCH,
230 Self::Send { .. } => KIND_SEND,
231 Self::Done => KIND_DONE,
232 Self::Part { .. } => KIND_PART,
233 Self::Forgotten { .. } => KIND_FORGOTTEN,
234 Self::Resume { .. } => KIND_RESUME,
235 }
236 }
237
238 /// The oldest format version that can express this message.
239 ///
240 /// A message is stamped with what it needs and not with what its writer was
241 /// built at, so every message a version 1 peer could have sent is still
242 /// spelled and stamped exactly as it was.
243 pub fn version(&self) -> u8 {
244 version_for(self.kind())
245 }
246
247 /// For messages about messages.
248 pub fn name(&self) -> &'static str {
249 match self {
250 Self::Hello { .. } => "hello",
251 Self::Sketch { .. } => "sketch",
252 Self::Send { .. } => "send",
253 Self::Done => "done",
254 Self::Part { .. } => "part",
255 Self::Forgotten { .. } => "forgotten",
256 Self::Resume { .. } => "resume",
257 }
258 }
259
260 /// Does the message open an exchange, which is what a peer answers?
261 pub fn is_opening(&self) -> bool {
262 matches!(self, Self::Hello { .. } | Self::Sketch { .. })
263 }
264
265 /// Empty for the messages that carry no frontier.
266 pub fn heads(&self) -> &[OpId] {
267 match self {
268 Self::Hello { heads } => heads,
269 Self::Sketch { heads, .. } => heads,
270 Self::Send { .. }
271 | Self::Done
272 | Self::Part { .. }
273 | Self::Forgotten { .. }
274 | Self::Resume { .. } => &[],
275 }
276 }
277
278 /// Where a caller that requires provenance does its checking, before the
279 /// message reaches a session: every [`Entry::Sealed`] can be put to
280 /// [`crate::envelope::Envelope::verify`] under whatever scheme the caller
281 /// holds. No scheme is chosen here and none is assumed.
282 ///
283 /// A [`Message::Forgotten`] answers here as well, so that what a carrier is
284 /// asked to write over its store is checked exactly as what it is asked to
285 /// absorb. What separates the two is [`Message::replacements`].
286 pub fn entries(&self) -> &[Entry] {
287 match self {
288 Self::Send { entries } => entries,
289 Self::Forgotten { entries } => entries,
290 Self::Hello { .. }
291 | Self::Sketch { .. }
292 | Self::Done
293 | Self::Part { .. }
294 | Self::Resume { .. } => &[],
295 }
296 }
297
298 /// The records a [`Message::Forgotten`] asks to be written in place of the
299 /// ones they name, and nothing for every other message.
300 ///
301 /// Asked separately from [`Message::entries`] because the two are acted on
302 /// differently and only the caller can tell them apart: an entry of a send is
303 /// an operation to absorb, and one of these is an operation already held,
304 /// arriving in the form it is to be kept in from now on. A caller that
305 /// absorbed one would place nothing; a caller that wrote a send's entries
306 /// over its store would destroy a history.
307 pub fn replacements(&self) -> &[Entry] {
308 match self {
309 Self::Forgotten { entries } => entries,
310 _ => &[],
311 }
312 }
313
314 /// Does the message hand operations over, whole or in pieces?
315 ///
316 /// A carrier bounding what it will send has to know the difference between a
317 /// message that carries work and one that only says something, because a
318 /// bound that stops before any work has gone is a turn that tells the far end
319 /// nothing it did not know -- and it comes back, and is told nothing again.
320 /// A part carries work and holds no entries, which is why asking
321 /// [`Message::entries`] is not the same question.
322 pub fn carries_operations(&self) -> bool {
323 match self {
324 Self::Send { entries } => !entries.is_empty(),
325 Self::Forgotten { entries } => !entries.is_empty(),
326 Self::Part { .. } => true,
327 _ => false,
328 }
329 }
330
331 /// The shape is `[kind, body]`. The table of a sketch is a [`Dat::BU64`]: it
332 /// readily exceeds the 255 bytes a [`Dat::BU8`] length field can express, and
333 /// a truncated length there would corrupt silently.
334 pub fn to_dat(&self) -> Dat {
335 let body = match self {
336 Self::Hello { heads } => Dat::List(
337 canonical(heads.clone()).iter().map(|h| h.to_dat()).collect(),
338 ),
339 Self::Sketch { heads, cells, count } => Dat::List(vec![
340 Dat::List(canonical(heads.clone()).iter().map(|h| h.to_dat()).collect()),
341 Dat::BU64(cells.clone()),
342 Dat::U64(*count),
343 ]),
344 Self::Send { entries } => Dat::List(
345 entries.iter().map(|e| e.to_dat()).collect(),
346 ),
347 Self::Done => Dat::List(Vec::new()),
348 Self::Part { id, seq, total, bytes } => Dat::List(vec![
349 id.to_dat(),
350 Dat::U64(*seq),
351 Dat::U64(*total),
352 Dat::BU64(bytes.clone()),
353 ]),
354 Self::Forgotten { entries } => Dat::List(
355 entries.iter().map(|e| e.to_dat()).collect(),
356 ),
357 Self::Resume { at } => at.to_dat(),
358 };
359 Dat::List(vec![Dat::U8(self.kind()), body])
360 }
361
362 pub fn from_dat(dat: &Dat)
363 -> Outcome<Self>
364 {
365 let v = match dat {
366 Dat::List(v) if v.len() == 2 => v,
367 _ => return Err(err!(
368 "A Message expects a 2-element Dat::List, got {:?}.", dat;
369 Decode, Input, Mismatch)),
370 };
371 let kind = match &v[0] {
372 Dat::U8(k) => *k,
373 other => return Err(err!(
374 "A Message kind expects Dat::U8, got {:?}.", other;
375 Decode, Input, Mismatch)),
376 };
377 match kind {
378 KIND_HELLO => Ok(Self::Hello {
379 heads: res!(heads_from_dat(&v[1], "hello")),
380 }),
381 KIND_SKETCH => {
382 let f = match &v[1] {
383 Dat::List(f) if f.len() == 3 => f,
384 other => return Err(err!(
385 "A sketch message expects a 3-element Dat::List, got {:?}.",
386 other;
387 Decode, Input, Mismatch)),
388 };
389 let cells = match &f[1] {
390 Dat::BU64(b) => b.clone(),
391 other => return Err(err!(
392 "A sketch message's table expects Dat::BU64, got {:?}.", other;
393 Decode, Input, Mismatch)),
394 };
395 let count = match &f[2] {
396 Dat::U64(n) => *n,
397 other => return Err(err!(
398 "A sketch message's count expects Dat::U64, got {:?}.", other;
399 Decode, Input, Mismatch)),
400 };
401 Ok(Self::Sketch {
402 heads: res!(heads_from_dat(&f[0], "sketch")),
403 cells,
404 count,
405 })
406 },
407 KIND_SEND => {
408 let listed = match &v[1] {
409 Dat::List(e) => e,
410 other => return Err(err!(
411 "A send message's operations expect Dat::List, got {:?}.", other;
412 Decode, Input, Mismatch)),
413 };
414 let mut entries = Vec::with_capacity(listed.len());
415 for item in listed {
416 entries.push(res!(Entry::from_dat(item)));
417 }
418 Ok(Self::Send { entries })
419 },
420 KIND_DONE => match &v[1] {
421 Dat::List(f) if f.is_empty() => Ok(Self::Done),
422 other => Err(err!(
423 "A done message expects an empty Dat::List, got {:?}.", other;
424 Decode, Input, Mismatch)),
425 },
426 KIND_PART => {
427 let f = match &v[1] {
428 Dat::List(f) if f.len() == 4 => f,
429 other => return Err(err!(
430 "A part message expects a 4-element Dat::List, got {:?}.", other;
431 Decode, Input, Mismatch)),
432 };
433 let id = res!(OpId::from_dat(&f[0]));
434 let seq = match &f[1] {
435 Dat::U64(n) => *n,
436 other => return Err(err!(
437 "A part message's position expects Dat::U64, got {:?}.", other;
438 Decode, Input, Mismatch)),
439 };
440 let total = match &f[2] {
441 Dat::U64(n) => *n,
442 other => return Err(err!(
443 "A part message's count expects Dat::U64, got {:?}.", other;
444 Decode, Input, Mismatch)),
445 };
446 // One width and not the narrower ones a shorter piece would also
447 // fit, for the reason a veiled body is written the same way: two
448 // spellings of one piece would put back together into the same
449 // operation and hash differently on the way.
450 let bytes = match &f[3] {
451 Dat::BU64(b) => b.clone(),
452 other => return Err(err!(
453 "The piece {} of {} is encoded {:?}; it is written under a \
454 64-bit length, whatever its size.", seq, id, other;
455 Decode, Input, Mismatch)),
456 };
457 if total == 0 {
458 return Err(err!(
459 "A part message says the operation {} was cut into no pieces \
460 at all.", id;
461 Decode, Input, Invalid));
462 }
463 if seq >= total {
464 return Err(err!(
465 "A part message is piece {} of the {} the operation {} was cut \
466 into, counting from zero.", seq, total, id;
467 Decode, Input, Range));
468 }
469 if bytes.is_empty() {
470 return Err(err!(
471 "The piece {} of {} carries no bytes, so it makes no progress \
472 towards the operation and a peer could send them for ever.",
473 seq, id;
474 Decode, Input, Invalid));
475 }
476 Ok(Self::Part { id, seq, total, bytes })
477 },
478 KIND_FORGOTTEN => {
479 let listed = match &v[1] {
480 Dat::List(e) => e,
481 other => return Err(err!(
482 "A forgotten message's records expect Dat::List, got {:?}.", other;
483 Decode, Input, Mismatch)),
484 };
485 if listed.is_empty() {
486 return Err(err!(
487 "A forgotten message names no record to write. It is an \
488 instruction to a carrier and not a statement about a session, \
489 so an empty one asks for nothing and says nothing.";
490 Decode, Input, Missing));
491 }
492 let mut entries = Vec::with_capacity(listed.len());
493 for item in listed {
494 entries.push(res!(Entry::from_dat(item)));
495 }
496 Ok(Self::Forgotten { entries })
497 },
498 KIND_RESUME => Ok(Self::Resume { at: res!(OpId::from_dat(&v[1])) }),
499 other => Err(err!(
500 "A Message is tagged {}, which names no message this version knows.",
501 other;
502 Decode, Input, Invalid)),
503 }
504 }
505
506 /// The magic, then the version, then the daticle form.
507 pub fn encode_into(&self, buf: &mut Vec<u8>)
508 -> Outcome<()>
509 {
510 buf.extend_from_slice(&MAGIC);
511 buf.push(self.version());
512 let body = res!(self.to_dat().to_bytes(Vec::new()));
513 buf.extend_from_slice(&body);
514 Ok(())
515 }
516
517 pub fn encode(&self)
518 -> Outcome<Vec<u8>>
519 {
520 let mut buf = Vec::new();
521 res!(self.encode_into(&mut buf));
522 Ok(buf)
523 }
524
525 /// What [`Message::encode`] comes to, without any of it being built.
526 ///
527 /// Both ends bound what they put in one body, and a bound is decided by
528 /// measuring messages that are then mostly not sent. Serialising to measure
529 /// makes that decision cost the payload it declines: a clone of fe2o3 measured
530 /// its whole history twice over to send it once.
531 pub fn encoded_len(&self)
532 -> Outcome<usize>
533 {
534 let body = res!(self.to_dat().byte_len().ok_or_else(|| err!(
535 "A {} message holds a daticle whose encoded length cannot be known \
536 without encoding it, which is a kind no message was built to carry.",
537 self.name();
538 Bug, Invalid)));
539 Ok(MAGIC.len() + 1 + body)
540 }
541
542
543 /// Cuts one entry into pieces, each of which encodes to at most `cap` bytes.
544 ///
545 /// This is the answer to an operation the carrier cannot take whole, and it
546 /// is the only one that does not cost something else. Raising the carrier's
547 /// limit fixes one operation and leaves the next; splitting the operation
548 /// itself would change what was signed; compressing it moves the threshold
549 /// rather than removing it. Cutting the *bytes* leaves the operation exactly
550 /// as it was, which is what [`Parts`] puts back.
551 ///
552 /// The pieces must be sent in the order they are handed back, and nothing may
553 /// come between them.
554 ///
555 /// # Arguments
556 ///
557 /// * `cap` - The most one piece may encode to, framing excluded; a carrier
558 /// that adds a length prefix subtracts it before asking.
559 pub fn part(entry: &Entry, cap: usize)
560 -> Outcome<Vec<Self>>
561 {
562 let id = res!(entry.id());
563 let whole = res!(entry.to_dat().to_bytes(Vec::new()));
564 if whole.len() > PART_MAX {
565 return Err(err!(
566 "The operation {} encodes to {} bytes, and a carrier will put back \
567 together at most {}.", id, whole.len(), PART_MAX;
568 Invalid, Data, Excessive));
569 }
570 // A piece's framing is fixed except for the two compact length prefixes a
571 // daticle list carries, which widen with what they measure. So the room is
572 // taken from an empty piece and then corrected against a real one, rather
573 // than guessed at with a margin that would be wrong at some size nobody
574 // tried.
575 let empty = Self::Part { id, seq: 0, total: 1, bytes: Vec::new() };
576 let mut room = cap.saturating_sub(res!(empty.encoded_len()));
577 loop {
578 if room == 0 {
579 return Err(err!(
580 "A carrier taking {} bytes cannot carry a piece of the operation \
581 {}: the framing of a piece comes to {} on its own.",
582 cap, id, res!(empty.encoded_len());
583 Invalid, Input, Range));
584 }
585 let sample = Self::Part {
586 id,
587 seq: 0,
588 total: 1,
589 bytes: whole[..std::cmp::min(room, whole.len())].to_vec(),
590 };
591 let sized = res!(sample.encoded_len());
592 if sized <= cap {
593 break;
594 }
595 room = room.saturating_sub(sized - cap);
596 }
597 let total = whole.len().div_ceil(room) as u64;
598 let mut out = Vec::with_capacity(total as usize);
599 for (seq, piece) in whole.chunks(room).enumerate() {
600 out.push(Self::Part {
601 id,
602 seq: seq as u64,
603 total,
604 bytes: piece.to_vec(),
605 });
606 }
607 Ok(out)
608 }
609
610 /// The message must occupy the whole of `buf`; where one ends is the
611 /// transport's question.
612 pub fn decode(buf: &[u8])
613 -> Outcome<Self>
614 {
615 let at = MAGIC.len() + 1;
616 if buf.len() < at {
617 return Err(err!(
618 "A sync message of {} byte{} is too short to carry even its header.",
619 buf.len(), if buf.len() == 1 { "" } else { "s" };
620 Decode, Input, Missing));
621 }
622 if buf[..MAGIC.len()] != MAGIC {
623 return Err(err!(
624 "A sync message begins {:02x?}, which is not the magic {:02x?}.",
625 &buf[..MAGIC.len()], MAGIC;
626 Decode, Input, Invalid));
627 }
628 let version = buf[MAGIC.len()];
629 if version < VERSION_MIN || version > VERSION {
630 return Err(err!(
631 "A sync message declares format version {}, and this reader knows \
632 versions {} to {}.", version, VERSION_MIN, VERSION;
633 Decode, Input, Version, Mismatch));
634 }
635 let (dat, used) = res!(Dat::from_bytes(&buf[at..]));
636 if used != buf.len() - at {
637 return Err(err!(
638 "A sync message body of {} bytes decoded from only {} of them.",
639 buf.len() - at, used;
640 Decode, Input, Mismatch));
641 }
642 let msg = res!(Self::from_dat(&dat));
643 // The version and the kind have to agree, or the superset promise
644 // VERSION_MIN rests on is a promise about intentions. A message stamped
645 // with a version that could not have expressed it is refused here rather
646 // than read, since a sender that spelled one of them wrong may have
647 // spelled others wrong too.
648 if msg.kind() > highest_kind(version) {
649 return Err(err!(
650 "A sync message declares format version {} and carries a {} message, \
651 which is kind {}; version {} carries up to kind {}.",
652 version, msg.name(), msg.kind(), version, highest_kind(version);
653 Decode, Input, Version, Mismatch));
654 }
655 Ok(msg)
656 }
657}
658
659
660
661/// One operation being put back together out of the pieces it crossed in.
662///
663/// A carrier holds one of these for each peer it is taking a push from, and a
664/// peer that stops halfway leaves one holding bytes that will never be
665/// completed. **Nothing is done with those bytes and nothing durable is written
666/// from them.** A run begins at piece zero and a piece zero discards whatever
667/// was held, so a push that died and was run again completes rather than
668/// doubling; and since the operation was never absorbed, the rerun offers it
669/// again of its own accord. That is Ore's resume-is-rerun, unchanged: the only
670/// state a carrier keeps is a buffer it is free to throw away, and throwing it
671/// away costs the sender a repetition and nothing else.
672///
673/// Every other message is handed straight back, so a caller folds one of these
674/// over an arriving run and gets the run it would have had if the carrier had
675/// been able to take the operation whole.
676#[derive(Clone, Debug, Default)]
677pub struct Parts {
678 held: Option<Holding>,
679}
680
681#[derive(Clone, Debug)]
682struct Holding {
683 id: OpId, // the operation the pieces make
684 total: u64, // how many pieces were declared
685 next: u64, // which piece is expected now
686 bytes: Vec<u8>, // what has arrived, in the order it arrived
687}
688
689impl Parts {
690
691 pub fn new() -> Self {
692 Self::default()
693 }
694
695 /// Is something half arrived?
696 pub fn pending(&self) -> bool {
697 self.held.is_some()
698 }
699
700 /// How many bytes are being held towards an operation that has not finished
701 /// arriving.
702 pub fn held(&self) -> usize {
703 match &self.held {
704 Some(h) => h.bytes.len(),
705 None => 0,
706 }
707 }
708
709 /// Throws away whatever is half arrived.
710 pub fn forget(&mut self) {
711 self.held = None;
712 }
713
714 /// Takes one message, and hands back the message the caller should act on.
715 ///
716 /// A piece that does not complete an operation yields nothing, because there
717 /// is nothing yet to act on. The piece that completes one yields the send the
718 /// operation would have crossed in had the carrier been able to take it
719 /// whole, so everything downstream sees exactly what it would have seen.
720 ///
721 /// A run that is interrupted by anything at all is refused rather than
722 /// patched over: the pieces of one operation are sent together, so a message
723 /// between them means the two ends disagree about what is being carried, and
724 /// a carrier that guessed would be putting together an operation nobody sent.
725 pub fn absorb(&mut self, msg: Message)
726 -> Outcome<Option<Message>>
727 {
728 let (id, seq, total, bytes) = match msg {
729 Message::Part { id, seq, total, bytes } => (id, seq, total, bytes),
730 other => {
731 if let Some(h) = self.held.take() {
732 return Err(err!(
733 "A {} message arrived between the pieces of {}, of which {} of \
734 {} had crossed. The pieces of one operation are sent together \
735 and nothing comes between them.",
736 other.name(), h.id, h.next, h.total;
737 Invalid, Input, Order));
738 }
739 return Ok(Some(other));
740 },
741 };
742 if seq == 0 {
743 // A run begins, and whatever was held for this peer is a run that
744 // stopped. Nothing was absorbed from it, so nothing is lost by letting
745 // it go; keeping it would be the carrier deciding which of two attempts
746 // the sender meant.
747 self.held = Some(Holding { id, total, next: 0, bytes: Vec::new() });
748 }
749 let held = match &mut self.held {
750 Some(h) => h,
751 None => return Err(err!(
752 "The piece {} of {} arrived with nothing before it. A run of pieces \
753 begins at zero, and a carrier that started in the middle would be \
754 putting together an operation it had only part of.", seq, id;
755 Invalid, Input, Order, Missing)),
756 };
757 if held.id != id || held.total != total || held.next != seq {
758 let (was_id, was_total, want) = (held.id, held.total, held.next);
759 self.held = None;
760 return Err(err!(
761 "The piece {} of {} of {} arrived where piece {} of {} of {} was \
762 expected. The pieces of one operation are sent in order and nothing \
763 comes between them.",
764 seq, total, id, want, was_total, was_id;
765 Invalid, Input, Order, Mismatch));
766 }
767 // Grown by what has arrived and never sized from what was declared: a
768 // count is an instruction to allocate and it comes from wherever the
769 // message did.
770 if held.bytes.len() + bytes.len() > PART_MAX {
771 let (was_id, sofar) = (held.id, held.bytes.len());
772 self.held = None;
773 return Err(err!(
774 "The pieces of {} come to more than the {} bytes a carrier will put \
775 back together: {} had arrived and another {} followed.",
776 was_id, PART_MAX, sofar, bytes.len();
777 Invalid, Data, Excessive));
778 }
779 held.bytes.extend_from_slice(&bytes);
780 held.next += 1;
781 if held.next < held.total {
782 return Ok(None);
783 }
784 let done = match self.held.take() {
785 Some(h) => h,
786 None => return Err(err!(
787 "The pieces of {} completed and there is nothing held.", id;
788 Bug, Unreachable)),
789 };
790 let (dat, used) = res!(Dat::from_bytes(&done.bytes));
791 if used != done.bytes.len() {
792 return Err(err!(
793 "The {} pieces of {} came to {} bytes and decoded from only {} of \
794 them.", done.total, done.id, done.bytes.len(), used;
795 Decode, Input, Mismatch));
796 }
797 let entry = res!(Entry::from_dat(&dat));
798 // The clear identifier and the one inside are compared for the reason a
799 // veiled entry's two headers are: a carrier that cut an operation up is a
800 // carrier that could have relabelled the pieces, and the record inside is
801 // the signed one.
802 let inside = res!(entry.id());
803 if inside != done.id {
804 return Err(err!(
805 "The pieces said they made the operation {} and they make {}. The \
806 record inside is the signed one and is what to believe.",
807 done.id, inside;
808 Invalid, Input, Security, Mismatch));
809 }
810 Ok(Some(Message::Send { entries: vec![entry] }))
811 }
812}
813
814
815/// The order the encoding spells a frontier in: ascending, without repetition.
816fn canonical(heads: Vec<OpId>) -> Vec<OpId> {
817 let mut heads = heads;
818 heads.sort();
819 heads.dedup();
820 heads
821}
822
823/// Refuses a frontier that is not in canonical order, rather than sorting it.
824fn heads_from_dat(dat: &Dat, what: &str)
825 -> Outcome<Vec<OpId>>
826{
827 let listed = match dat {
828 Dat::List(v) => v,
829 other => return Err(err!(
830 "A {} message's frontier expects Dat::List, got {:?}.", what, other;
831 Decode, Input, Mismatch)),
832 };
833 let mut heads: Vec<OpId> = Vec::with_capacity(listed.len());
834 for item in listed {
835 let id = res!(OpId::from_dat(item));
836 if let Some(last) = heads.last() {
837 if id <= *last {
838 return Err(err!(
839 "A {} message lists {} after {}; a frontier is encoded ascending \
840 and without repetition.", what, id, last;
841 Decode, Input, Order));
842 }
843 }
844 heads.push(id);
845 }
846 Ok(heads)
847}
848
849
850#[cfg(test)]
851mod tests {
852 use super::*;
853
854 use crate::envelope::Envelope;
855 use crate::id::{
856 Anchor,
857 ContentId,
858 ReplicaId,
859 };
860 use crate::op::{
861 Header,
862 Op,
863 Placing,
864 Record,
865 };
866 use crate::test_support::StubSigner;
867
868 fn oid(replica: u64, counter: u64) -> OpId {
869 OpId::new(ReplicaId::new(replica), counter)
870 }
871
872 /// One of each message, with both entry forms among them.
873 fn samples()
874 -> Outcome<Vec<Message>>
875 {
876 let rec = Record::new(
877 res!(Header::new(oid(2, 3), vec![oid(1, 1), oid(1, 2)])),
878 Op::FileCreate { path: b"notes.md".to_vec() },
879 );
880 let sealed = res!(Envelope::seal_record(
881 &StubSigner::with_seed(3),
882 &Record::root(oid(1, 1), Op::Mark { name: fmt!("start"), body: None, time: None }),
883 ));
884 Ok(vec![
885 Message::hello(vec![oid(3, 9), oid(1, 2)]),
886 Message::hello(Vec::new()),
887 Message::sketch(vec![oid(1, 2)], vec![0x11; 300], 42),
888 Message::Send { entries: vec![
889 Entry::Bare(rec),
890 Entry::Sealed(sealed),
891 ] },
892 Message::Send { entries: Vec::new() },
893 Message::Done,
894 Message::Part { id: oid(2, 3), seq: 0, total: 3, bytes: vec![0x5a; 40] },
895 Message::Part { id: oid(2, 3), seq: 2, total: 3, bytes: vec![0x01] },
896 Message::Forgotten { entries: res!(stubs()) },
897 Message::Resume { at: oid(2, 3) },
898 ])
899 }
900
901 /// Two records in a forgotten operation's place: one bare, one sealed, and
902 /// both shapes a splice can keep.
903 fn stubs()
904 -> Outcome<Vec<Entry>>
905 {
906 let void = Record::new(
907 res!(Header::new(oid(2, 4), vec![oid(2, 3)])),
908 Op::Forgotten { placing: Placing::Void },
909 );
910 let placed = Record::new(
911 res!(Header::new(oid(2, 5), vec![oid(2, 4)])),
912 Op::Forgotten { placing: Placing::Splice {
913 left: Some(Anchor::after(ContentId::new(oid(2, 3), 0))),
914 right: None,
915 remove: Vec::new(),
916 len: 17,
917 } },
918 );
919 Ok(vec![
920 Entry::Bare(void),
921 Entry::Sealed(res!(Envelope::seal_record(&StubSigner::with_seed(5), &placed))),
922 ])
923 }
924
925 #[test]
926 fn messages_round_trip() -> Outcome<()> {
927 for msg in res!(samples()) {
928 assert_eq!(res!(Message::from_dat(&msg.to_dat())), msg, "as a daticle");
929 let bytes = res!(msg.encode());
930 assert_eq!(res!(Message::decode(&bytes)), msg, "as bytes");
931 assert_eq!(&bytes[..MAGIC.len()], &MAGIC);
932 // The version a message is stamped with is the version it needs, so
933 // every message a version 1 peer could have sent still says 1.
934 assert_eq!(bytes[MAGIC.len()], msg.version(), "{} is stamped wrongly", msg.name());
935 }
936 Ok(())
937 }
938
939 /// A frontier is a set, so it is encoded ascending and without repetition
940 /// however it was handed over, and a decoder refuses any other spelling.
941 /// A run of pieces puts the entry back byte for byte, signature included.
942 ///
943 /// This is the whole claim chunking makes and the one thing it may not get
944 /// wrong. The bytes compared are the entry's own encoding, so a reassembly
945 /// that produced an equal-looking entry by re-encoding it would still fail;
946 /// and the envelope is verified afterwards, because a signature is over bytes
947 /// and equality of a struct is not equality of bytes.
948 #[test]
949 fn a_run_of_pieces_puts_the_signed_entry_back_byte_for_byte() -> Outcome<()> {
950 let signer = StubSigner::with_seed(11);
951 let rec = Record::new(
952 res!(Header::new(oid(4, 12), vec![oid(4, 11)])),
953 Op::Mark {
954 name: fmt!("a mark whose body is large enough to need cutting up"),
955 body: Some(vec![0xa7u8; 20_000]),
956 time: Some(1_755_000_000),
957 },
958 );
959 let entry = Entry::Sealed(res!(Envelope::seal_record(&signer, &rec)));
960 let was = res!(entry.to_dat().to_bytes(Vec::new()));
961 let pieces = res!(Message::part(&entry, 4_096));
962 assert!(pieces.len() > 4, "an entry of {} bytes made {} pieces", was.len(), pieces.len());
963
964 let mut parts = Parts::new();
965 let mut got = None;
966 for (at, piece) in pieces.iter().enumerate() {
967 // Every piece crosses the wire, because that is where a spelling that
968 // only round-trips in memory would go wrong.
969 let piece = res!(Message::decode(&res!(piece.encode())));
970 match res!(parts.absorb(piece)) {
971 Some(msg) => {
972 assert_eq!(at, pieces.len() - 1, "the run completed early");
973 got = Some(msg);
974 },
975 None => assert!(parts.pending(), "a piece was taken and nothing is held"),
976 }
977 }
978 assert!(!parts.pending(), "the run completed and something is still held");
979 let back = match got {
980 Some(Message::Send { entries }) => entries,
981 other => return Err(err!(
982 "A completed run yielded {:?} rather than a send.", other; Test, Mismatch)),
983 };
984 assert_eq!(back.len(), 1);
985 assert_eq!(res!(back[0].to_dat().to_bytes(Vec::new())), was,
986 "the entry that came back is not the entry that went, byte for byte");
987 // And the signature over those bytes still verifies.
988 match &back[0] {
989 Entry::Sealed(env) => assert!(res!(env.verify(&signer)),
990 "the signature did not survive being cut up"),
991 other => return Err(err!(
992 "A sealed entry came back as a {}.", other.name(); Test, Mismatch)),
993 }
994 Ok(())
995 }
996
997 /// Every piece fits the carrier that asked for it, at every cap, and the
998 /// pieces are the whole in order.
999 ///
1000 /// BOTH HALVES, because either alone is satisfied by something useless: one
1001 /// piece per byte respects any cap, and one piece carrying everything
1002 /// preserves any order. The first half is the defect this exists for -- a
1003 /// carrier told a limit and handed something over it is the failure that
1004 /// reads as the relay being down.
1005 #[test]
1006 fn every_piece_fits_the_cap_and_the_pieces_are_the_whole() -> Outcome<()> {
1007 let entry = Entry::Bare(Record::new(
1008 res!(Header::new(oid(9, 2), vec![oid(9, 1)])),
1009 Op::Mark { name: fmt!("wide"), body: Some(vec![0x7eu8; 30_000]), time: None },
1010 ));
1011 let was = res!(entry.to_dat().to_bytes(Vec::new()));
1012 for cap in [96usize, 200, 1_024, 4_096, 65_536, 1 << 20] {
1013 let pieces = res!(Message::part(&entry, cap));
1014 let mut seen: Vec<u8> = Vec::new();
1015 for (at, piece) in pieces.iter().enumerate() {
1016 let sized = res!(piece.encode()).len();
1017 assert!(sized <= cap,
1018 "piece {} of {} encodes to {} against a cap of {}",
1019 at, pieces.len(), sized, cap);
1020 match piece {
1021 Message::Part { id, seq, total, bytes } => {
1022 assert_eq!(*id, res!(entry.id()));
1023 assert_eq!(*seq, at as u64, "the pieces are out of order");
1024 assert_eq!(*total as usize, pieces.len(), "the count is wrong");
1025 seen.extend_from_slice(bytes);
1026 },
1027 other => return Err(err!(
1028 "Parting an entry produced a {} message.", other.name();
1029 Test, Mismatch)),
1030 }
1031 }
1032 assert_eq!(seen, was, "the pieces at a cap of {} are not the whole", cap);
1033 }
1034 // A cap under the framing of a piece is refused by name rather than
1035 // producing a piece that does not fit it.
1036 assert!(Message::part(&entry, 8).is_err(), "a cap of eight bytes was accepted");
1037 Ok(())
1038 }
1039
1040 /// A run that is broken in any way is refused, and nothing half assembled is
1041 /// kept.
1042 #[test]
1043 fn a_broken_run_is_refused_rather_than_patched() -> Outcome<()> {
1044 let entry = Entry::Bare(Record::new(
1045 res!(Header::new(oid(5, 2), vec![oid(5, 1)])),
1046 Op::Mark { name: fmt!("five"), body: Some(vec![0x33u8; 2_000]), time: None },
1047 ));
1048 let pieces = res!(Message::part(&entry, 512));
1049 assert!(pieces.len() >= 4);
1050
1051 // A run that begins in the middle.
1052 let mut parts = Parts::new();
1053 assert!(parts.absorb(pieces[1].clone()).is_err(), "a run beginning at piece one");
1054 assert!(!parts.pending());
1055
1056 // A piece skipped.
1057 let mut parts = Parts::new();
1058 res!(parts.absorb(pieces[0].clone()));
1059 assert!(parts.absorb(pieces[2].clone()).is_err(), "a piece skipped");
1060 assert!(!parts.pending(), "a refused run left something held");
1061
1062 // Another operation's piece in the middle of this one.
1063 let other = res!(Message::part(&Entry::Bare(Record::new(
1064 res!(Header::new(oid(6, 2), vec![oid(6, 1)])),
1065 Op::Mark { name: fmt!("six"), body: Some(vec![0x44u8; 2_000]), time: None },
1066 )), 512));
1067 let mut parts = Parts::new();
1068 res!(parts.absorb(pieces[0].clone()));
1069 assert!(parts.absorb(other[1].clone()).is_err(), "another operation's piece");
1070
1071 // Anything at all between the pieces.
1072 let mut parts = Parts::new();
1073 res!(parts.absorb(pieces[0].clone()));
1074 assert!(parts.absorb(Message::Done).is_err(), "a done between the pieces");
1075 assert!(!parts.pending());
1076
1077 // Pieces that make bytes which are not an entry.
1078 let mut parts = Parts::new();
1079 let rubbish: Vec<Message> = (0..2u64)
1080 .map(|seq| Message::Part {
1081 id: oid(5, 2),
1082 seq,
1083 total: 2,
1084 bytes: vec![0xffu8; 8],
1085 })
1086 .collect();
1087 res!(parts.absorb(rubbish[0].clone()));
1088 assert!(parts.absorb(rubbish[1].clone()).is_err(), "bytes that are not an entry");
1089
1090 // Pieces that make an entry which is not the one they claim.
1091 let whole = res!(entry.to_dat().to_bytes(Vec::new()));
1092 let mut parts = Parts::new();
1093 assert!(parts.absorb(Message::Part {
1094 id: oid(7, 7),
1095 seq: 0,
1096 total: 1,
1097 bytes: whole,
1098 }).is_err(), "pieces labelled as another operation");
1099
1100 // And the spellings a decoder refuses outright.
1101 for wrong in [
1102 Message::Part { id: oid(5, 2), seq: 0, total: 0, bytes: vec![1] },
1103 Message::Part { id: oid(5, 2), seq: 3, total: 3, bytes: vec![1] },
1104 Message::Part { id: oid(5, 2), seq: 0, total: 2, bytes: Vec::new() },
1105 ] {
1106 assert!(Message::from_dat(&wrong.to_dat()).is_err(),
1107 "a part {:?} was accepted", wrong);
1108 }
1109 Ok(())
1110 }
1111
1112 /// A push that died halfway and was run again completes rather than doubling.
1113 ///
1114 /// The carrier keeps nothing durable and the operation was never absorbed, so
1115 /// the rerun offers the whole of it again; piece zero is what says a run has
1116 /// begun, and it throws away what the attempt before it left.
1117 #[test]
1118 fn a_rerun_completes_rather_than_doubling() -> Outcome<()> {
1119 let entry = Entry::Bare(Record::new(
1120 res!(Header::new(oid(8, 2), vec![oid(8, 1)])),
1121 Op::Mark { name: fmt!("eight"), body: Some(vec![0x21u8; 4_000]), time: None },
1122 ));
1123 let was = res!(entry.to_dat().to_bytes(Vec::new()));
1124 let pieces = res!(Message::part(&entry, 512));
1125 assert!(pieces.len() >= 4);
1126
1127 let mut parts = Parts::new();
1128 // The push gets halfway and stops.
1129 for piece in &pieces[..2] {
1130 assert!(res!(parts.absorb(piece.clone())).is_none());
1131 }
1132 assert!(parts.pending(), "nothing was held when the push stopped");
1133 assert!(parts.held() > 0);
1134
1135 // It is run again from the beginning, against a carrier still holding the
1136 // remains of the first attempt.
1137 let mut got = None;
1138 for piece in &pieces {
1139 if let Some(msg) = res!(parts.absorb(piece.clone())) {
1140 got = Some(msg);
1141 }
1142 }
1143 let back = match got {
1144 Some(Message::Send { entries }) => entries,
1145 other => return Err(err!(
1146 "The rerun yielded {:?} rather than a send.", other; Test, Mismatch)),
1147 };
1148 assert_eq!(back.len(), 1, "the rerun produced {} entries", back.len());
1149 assert_eq!(res!(back[0].to_dat().to_bytes(Vec::new())), was,
1150 "the rerun did not put the operation back as it was");
1151 assert!(!parts.pending());
1152 Ok(())
1153 }
1154
1155 /// A peer that speaks only version 1 refuses a part by name, and reads every
1156 /// message a version 1 peer could have sent.
1157 ///
1158 /// The refusal is at the header, before a decode, which is the reason the
1159 /// version moved rather than a kind byte being added to an entry: a kind
1160 /// would have been refused inside `Entry::from_dat`, after the handshake had
1161 /// already said the frame was compatible.
1162 #[test]
1163 fn a_version_and_a_kind_have_to_agree() -> Outcome<()> {
1164 let part = Message::Part { id: oid(1, 1), seq: 0, total: 1, bytes: vec![0x01] };
1165 let bytes = res!(part.encode());
1166 assert_eq!(bytes[MAGIC.len()], version_for(KIND_PART),
1167 "a part is stamped with the version it needs");
1168
1169 // The same message stamped as version 1, which is what an old peer would
1170 // have to believe to read it.
1171 let mut lying = bytes.clone();
1172 lying[MAGIC.len()] = VERSION_MIN;
1173 let e = match Message::decode(&lying) {
1174 Ok(got) => return Err(err!(
1175 "A part stamped version {} decoded as a {}.", VERSION_MIN, got.name();
1176 Test, Mismatch)),
1177 Err(e) => e,
1178 };
1179 assert!(fmt!("{}", e).contains("kind"), "message was {}", e);
1180
1181 assert_eq!(highest_kind(VERSION_MIN), KIND_DONE);
1182 assert_eq!(highest_kind(VERSION), KIND_RESUME);
1183 // Every older message still says what it always said, so an old peer reads
1184 // it unchanged: the four original kinds say 1, a part still says 2 and a
1185 // replacement still says 3, although VERSION has moved past all of them.
1186 for msg in res!(samples()) {
1187 let stamped = res!(msg.encode())[MAGIC.len()];
1188 match msg {
1189 Message::Part { .. } => assert_eq!(stamped, 2,
1190 "a part stopped being a version 2 message"),
1191 Message::Forgotten { .. } => assert_eq!(stamped, 3,
1192 "a replacement stopped being a version 3 message"),
1193 Message::Resume { .. } => assert_eq!(stamped, VERSION),
1194 _ => assert_eq!(stamped, VERSION_MIN,
1195 "the {} message stopped being a version 1 message", msg.name()),
1196 }
1197 }
1198 // And the two halves of the rule are each other's inverse, at every kind
1199 // there is: a message is stamped with the oldest version that admits its
1200 // kind, so a reader following `highest_kind` never refuses one a writer
1201 // was entitled to send.
1202 for kind in [KIND_HELLO, KIND_SKETCH, KIND_SEND, KIND_DONE, KIND_PART,
1203 KIND_FORGOTTEN, KIND_RESUME]
1204 {
1205 let at = version_for(kind);
1206 assert!(highest_kind(at) >= kind,
1207 "kind {} is stamped version {}, which admits only up to {}",
1208 kind, at, highest_kind(at));
1209 assert!(at == VERSION_MIN || highest_kind(at - 1) < kind,
1210 "kind {} is stamped version {}, and version {} would have carried it",
1211 kind, at, at - 1);
1212 }
1213 Ok(())
1214 }
1215
1216 /// A peer built before the replacement message refuses one by name, at the
1217 /// header, and reads everything it could read before.
1218 ///
1219 /// The rule of rule six: nothing new is negotiated. The version byte already
1220 /// says what a reader must understand, and a reader that does not is told so
1221 /// before it decodes a byte of the body.
1222 #[test]
1223 fn an_older_peer_refuses_a_replacement_and_nothing_else() -> Outcome<()> {
1224 let msg = Message::Forgotten { entries: res!(stubs()) };
1225 let bytes = res!(msg.encode());
1226 assert_eq!(bytes[MAGIC.len()], version_for(KIND_FORGOTTEN),
1227 "a replacement is stamped with the version it needs");
1228 // What a version 2 peer would be doing if it read this one: believing a
1229 // stamp the sender could not honestly have written.
1230 let mut lying = bytes.clone();
1231 lying[MAGIC.len()] = 2;
1232 let e = match Message::decode(&lying) {
1233 Ok(got) => return Err(err!(
1234 "A replacement stamped version 2 decoded as a {}.", got.name();
1235 Test, Mismatch)),
1236 Err(e) => e,
1237 };
1238 let said = fmt!("{}", e);
1239 assert!(said.contains("kind"), "message was {}", said);
1240 // A version this reader does not know is refused at the header too, which
1241 // is what a NEWER peer's message meets here.
1242 let mut ahead = bytes.clone();
1243 ahead[MAGIC.len()] = VERSION + 1;
1244 let e = match Message::decode(&ahead) {
1245 Ok(got) => return Err(err!(
1246 "A message stamped version {} decoded as a {}.", VERSION + 1, got.name();
1247 Test, Mismatch)),
1248 Err(e) => e,
1249 };
1250 assert!(fmt!("{}", e).contains("format version"), "message was {}", e);
1251 Ok(())
1252 }
1253
1254 /// A replacement carries its records where provenance is checked, and says
1255 /// separately that they are replacements and not an offer.
1256 #[test]
1257 fn a_replacement_is_not_an_offer() -> Outcome<()> {
1258 let msg = Message::Forgotten { entries: res!(stubs()) };
1259 assert_eq!(msg.entries().len(), 2, "the records are where a verifier looks");
1260 assert_eq!(msg.replacements().len(), 2);
1261 assert!(msg.heads().is_empty());
1262 assert!(!msg.is_opening());
1263 assert!(msg.carries_operations(), "a body of nothing else still does work");
1264 // A send is the other way round: entries to absorb, and nothing to write
1265 // over anything.
1266 let send = Message::Send { entries: res!(stubs()) };
1267 assert!(send.replacements().is_empty(),
1268 "a send's entries were read as records to write over the store");
1269 // And an empty one is refused rather than carried, since it asks for
1270 // nothing.
1271 let empty = Message::Forgotten { entries: Vec::new() };
1272 assert!(Message::from_dat(&empty.to_dat()).is_err(),
1273 "an empty replacement was taken for a message");
1274 Ok(())
1275 }
1276
1277 /// The bytes of a piece, frozen.
1278 ///
1279 /// A part is the one message that crosses in pieces, so a change to its
1280 /// spelling stops two builds putting the same operation back together while
1281 /// each remains correct against itself -- the failure a golden exists for.
1282 /// The payload is the entry's own bytes and not a re-encoding of them, which
1283 /// is what keeps the signature over them intact, and freezing them here is
1284 /// what says so.
1285 #[test]
1286 fn the_part_bytes_are_frozen() -> Outcome<()> {
1287 let entry = Entry::Bare(Record::new(
1288 res!(Header::new(oid(2, 3), vec![oid(1, 7)])),
1289 Op::Mark { name: fmt!("v1"), body: None, time: None },
1290 ));
1291 let pieces = res!(Message::part(&entry, 96));
1292 assert_eq!(pieces.len(), 2, "the entry was not cut into two pieces");
1293 let want: &[u8] = &[
1294 // The magic, and the version a part needs.
1295 0x4f, 0x52, 0x45, 0x53, 0x59, 0x4e,
1296 0x02,
1297 // The message: a two-element list of the kind and the body, 86 bytes.
1298 0x33, 0x21, 0x56,
1299 // The kind: part.
1300 0x0a, 0x05,
1301 // The body: identifier, position, count, bytes; 81 bytes.
1302 0x33, 0x21, 0x51,
1303 // The operation the pieces make, r2:3.
1304 0x33, 0x21, 0x12,
1305 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02,
1306 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x03,
1307 // Piece zero.
1308 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
1309 // Of two.
1310 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02,
1311 // Thirty-three bytes, under a 64-bit length whatever its size.
1312 0x47, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x21,
1313 // The first thirty-three bytes of the entry, which are the
1314 // entry's own bytes and not a re-encoding of it: the kind, the
1315 // record, the header, the identifier r2:3.
1316 0x33, 0x21, 0x3f,
1317 0x0a, 0x01,
1318 0x33, 0x21, 0x3a,
1319 0x33, 0x21, 0x2d,
1320 0x33, 0x21, 0x12,
1321 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02,
1322 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x03,
1323 0x33,
1324 ];
1325 assert_eq!(res!(pieces[0].encode()), want, "the part message format has changed");
1326 // And the frozen bytes still read.
1327 assert_eq!(res!(Message::decode(want)), pieces[0]);
1328 // The payload is a prefix of the entry's own bytes, which is the claim.
1329 let whole = res!(entry.to_dat().to_bytes(Vec::new()));
1330 assert_eq!(&want[63..], &whole[..33],
1331 "a piece carries something other than the entry's own bytes");
1332 Ok(())
1333 }
1334
1335 /// The bytes of a cursor, frozen.
1336 ///
1337 /// It is one identifier, and the framing around it is the whole of what a
1338 /// peer too old to read this one meets: the version byte, refused at the
1339 /// header by name. What the test holds down is that a cursor carries nothing
1340 /// else -- no count, no set, no claim about what was sent -- because a carrier
1341 /// that keeps nothing between sessions must be told where to resume in a
1342 /// message small enough to ride beside every opening.
1343 #[test]
1344 fn the_resume_bytes_are_frozen() -> Outcome<()> {
1345 let msg = Message::Resume { at: oid(2, 3) };
1346 let want: &[u8] = &[
1347 // The magic, and the version a cursor needs.
1348 0x4f, 0x52, 0x45, 0x53, 0x59, 0x4e,
1349 0x04,
1350 // The message: a two-element list of the kind and the body, 23 bytes.
1351 0x33, 0x21, 0x17,
1352 // The kind: resume.
1353 0x0a, 0x07,
1354 // The body: the operation this end was carried as far as, r2:3.
1355 0x33, 0x21, 0x12,
1356 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02,
1357 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x03,
1358 ];
1359 assert_eq!(res!(msg.encode()), want, "the resume message format has changed");
1360 assert_eq!(res!(Message::decode(want)), msg);
1361 // And a peer built before it refuses it at the header, naming both
1362 // versions, rather than reading past a kind it does not know.
1363 let mut older = want.to_vec();
1364 older[6] = version_for(KIND_FORGOTTEN);
1365 assert!(Message::decode(&older).is_err(),
1366 "a cursor stamped as version 3 was read by a version 3 vocabulary");
1367 Ok(())
1368 }
1369
1370 /// The bytes of a replacement, frozen.
1371 ///
1372 /// Two builds that disagree about this are two builds that disagree about
1373 /// which record a carrier is to write over which, and the one that gets it
1374 /// wrong destroys a history rather than failing to read one. What it freezes
1375 /// besides the framing is that the record inside is an ordinary entry under
1376 /// the forgotten operation's own header -- the same spelling a segment
1377 /// carries, so what the carrier writes down is what it was handed.
1378 #[test]
1379 fn the_forgotten_bytes_are_frozen() -> Outcome<()> {
1380 let msg = Message::Forgotten {
1381 entries: vec![Entry::Bare(Record::new(
1382 res!(Header::new(oid(2, 3), vec![oid(1, 7)])),
1383 Op::Forgotten { placing: Placing::Void },
1384 ))],
1385 };
1386 let want: &[u8] = &[
1387 // The magic, and the version a replacement needs.
1388 0x4f, 0x52, 0x45, 0x53, 0x59, 0x4e,
1389 0x03,
1390 // The message: a two-element list of the kind and the body, 71 bytes.
1391 0x33, 0x21, 0x47,
1392 // The kind: forgotten.
1393 0x0a, 0x06,
1394 // The body: a list of one entry, 66 bytes.
1395 0x33, 0x21, 0x42,
1396 // The entry, 63 bytes: the kind, and the record.
1397 0x33, 0x21, 0x3f,
1398 // Bare.
1399 0x0a, 0x01,
1400 // The record, 58 bytes: the header, the operation.
1401 0x33, 0x21, 0x3a,
1402 // The header, 45 bytes, and it is the FORGOTTEN
1403 // operation's own: the identifier r2:3, then the parents.
1404 0x33, 0x21, 0x2d,
1405 0x33, 0x21, 0x12,
1406 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02,
1407 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x03,
1408 // One parent, r1:7.
1409 0x33, 0x21, 0x15,
1410 0x33, 0x21, 0x12,
1411 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x01,
1412 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x07,
1413 // The operation, 7 bytes: the Forgotten code, and the
1414 // shape it keeps, which here is a void.
1415 0x33, 0x21, 0x07,
1416 0x0a, 0x10,
1417 0x33, 0x21, 0x02,
1418 0x0a, 0x00,
1419 ];
1420 assert_eq!(res!(msg.encode()), want, "the replacement message format has changed");
1421 // And the frozen bytes still read.
1422 assert_eq!(res!(Message::decode(want)), msg);
1423 Ok(())
1424 }
1425
1426 #[test]
1427 fn a_frontier_is_encoded_canonically() -> Outcome<()> {
1428 let msg = Message::hello(vec![oid(3, 1), oid(1, 1), oid(3, 1), oid(2, 1)]);
1429 assert_eq!(msg.heads(), vec![oid(1, 1), oid(2, 1), oid(3, 1)]);
1430 // Built by hand, out of order, it still encodes canonically.
1431 let odd = Message::Hello { heads: vec![oid(3, 1), oid(1, 1)] };
1432 assert_eq!(
1433 res!(Message::from_dat(&odd.to_dat())),
1434 Message::hello(vec![oid(1, 1), oid(3, 1)]),
1435 );
1436 // And a daticle spelling it otherwise is refused rather than sorted.
1437 let wrong = Dat::List(vec![
1438 Dat::U8(KIND_HELLO),
1439 Dat::List(vec![oid(3, 1).to_dat(), oid(1, 1).to_dat()]),
1440 ]);
1441 let _ = msg;
1442 let e = match Message::from_dat(&wrong) {
1443 Ok(_) => return Err(err!("A frontier out of order was accepted."; Test)),
1444 Err(e) => e,
1445 };
1446 assert!(fmt!("{}", e).contains("ascending"), "message was {}", e);
1447 // Repetition likewise.
1448 let twice = Dat::List(vec![
1449 Dat::U8(KIND_HELLO),
1450 Dat::List(vec![oid(1, 1).to_dat(), oid(1, 1).to_dat()]),
1451 ]);
1452 assert!(Message::from_dat(&twice).is_err());
1453 Ok(())
1454 }
1455
1456 /// Truncating a message anywhere is a typed error, never a panic and never a
1457 /// half-read message.
1458 #[test]
1459 fn truncation_at_every_offset_is_clean() -> Outcome<()> {
1460 for msg in res!(samples()) {
1461 let bytes = res!(msg.encode());
1462 for cut in 0..bytes.len() {
1463 match Message::decode(&bytes[..cut]) {
1464 Ok(got) => return Err(err!(
1465 "A {} message cut at {} of {} decoded as a {}.",
1466 msg.name(), cut, bytes.len(), got.name();
1467 Test, Mismatch)),
1468 Err(_) => {},
1469 }
1470 }
1471 assert_eq!(res!(Message::decode(&bytes)), msg);
1472 // And trailing rubbish is refused, not ignored.
1473 let mut extra = bytes.clone();
1474 extra.push(0x00);
1475 assert!(Message::decode(&extra).is_err(), "{} with a trailing byte", msg.name());
1476 }
1477 Ok(())
1478 }
1479
1480 #[test]
1481 fn a_message_that_is_not_one_is_refused() -> Outcome<()> {
1482 assert!(Message::decode(b"").is_err());
1483 assert!(Message::decode(b"not a sync message").is_err());
1484 let mut wrong = MAGIC.to_vec();
1485 wrong.push(VERSION + 1);
1486 wrong.push(0);
1487 let e = match Message::decode(&wrong) {
1488 Ok(_) => return Err(err!("An unknown version was accepted."; Test)),
1489 Err(e) => e,
1490 };
1491 assert!(fmt!("{}", e).contains("version"), "message was {}", e);
1492 // A kind nobody knows.
1493 let odd = Dat::List(vec![Dat::U8(99), Dat::List(Vec::new())]);
1494 assert!(Message::from_dat(&odd).is_err());
1495 // A done message carrying something.
1496 let heavy = Dat::List(vec![
1497 Dat::U8(KIND_DONE),
1498 Dat::List(vec![Dat::U64(1)]),
1499 ]);
1500 assert!(Message::from_dat(&heavy).is_err());
1501 Ok(())
1502 }
1503
1504 #[test]
1505 fn accessors_report_what_is_there() -> Outcome<()> {
1506 let msgs = res!(samples());
1507 assert!(msgs[0].is_opening());
1508 assert!(msgs[2].is_opening());
1509 assert!(!msgs[3].is_opening());
1510 assert!(!Message::Done.is_opening());
1511 assert_eq!(msgs[2].heads(), vec![oid(1, 2)]);
1512 assert!(msgs[3].heads().is_empty());
1513 assert_eq!(msgs[3].entries().len(), 2);
1514 assert!(msgs[0].entries().is_empty());
1515 assert_eq!(Message::Done.name(), "done");
1516 assert_eq!(Message::Done.kind(), KIND_DONE);
1517 Ok(())
1518 }
1519
1520 /// A format that changes by accident leaves two versions of this crate unable
1521 /// to speak to each other, and every other test in this file would pass
1522 /// regardless: they all encode and decode with the same code. This one is the
1523 /// fixed point. If it fails and the change was deliberate, the version byte is
1524 /// the thing to raise.
1525 #[test]
1526 fn the_message_bytes_are_frozen() -> Outcome<()> {
1527 let msg = Message::Send {
1528 entries: vec![Entry::Bare(Record::new(
1529 res!(Header::new(oid(2, 3), vec![oid(1, 7)])),
1530 Op::Mark { name: fmt!("v1"), body: None, time: None },
1531 ))],
1532 };
1533 let want: &[u8] = &[
1534 // The magic and the version.
1535 0x4f, 0x52, 0x45, 0x53, 0x59, 0x4e,
1536 0x01,
1537 // The message: a two-element list of the kind and the body, 71 bytes.
1538 0x33, 0x21, 0x47,
1539 // The kind: send.
1540 0x0a, 0x03,
1541 // The body: a list of one entry, 66 bytes.
1542 0x33, 0x21, 0x42,
1543 // The entry, 63 bytes: the kind, and the record.
1544 0x33, 0x21, 0x3f,
1545 // Bare.
1546 0x0a, 0x01,
1547 // The record, 58 bytes: the header, the operation.
1548 0x33, 0x21, 0x3a,
1549 // The header, 45 bytes: the identifier, then the parents.
1550 0x33, 0x21, 0x2d,
1551 // The identifier r2:3, as two 64-bit integers.
1552 0x33, 0x21, 0x12,
1553 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02,
1554 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x03,
1555 // One parent, r1:7.
1556 0x33, 0x21, 0x15,
1557 0x33, 0x21, 0x12,
1558 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x01,
1559 0x0d, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x07,
1560 // The operation, 7 bytes: the Mark code, and the name "v1".
1561 0x33, 0x21, 0x07,
1562 0x0a, 0x04,
1563 0x29, 0x21, 0x02, 0x76, 0x31,
1564 ];
1565 assert_eq!(res!(msg.encode()), want, "the sync message format has changed");
1566 // And the frozen bytes still read.
1567 assert_eq!(res!(Message::decode(want)), msg);
1568 Ok(())
1569 }
1570}
1571