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 | |
| 21 | use crate::id::OpId; |
| 22 | use crate::segment::Entry; |
| 23 | |
| 24 | use oxedyne_fe2o3_core::prelude::*; |
| 25 | use oxedyne_fe2o3_jdat::prelude::*; |
| 26 | |
| 27 | |
| 28 | /// The bytes every sync message begins with. |
| 29 | pub 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. |
| 63 | pub 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. |
| 73 | pub const VERSION_MIN: u8 = 1; |
| 74 | |
| 75 | // The kind byte each message is tagged with on the wire. |
| 76 | pub const KIND_HELLO: u8 = 1; |
| 77 | pub const KIND_SKETCH: u8 = 2; |
| 78 | pub const KIND_SEND: u8 = 3; |
| 79 | pub const KIND_DONE: u8 = 4; |
| 80 | pub const KIND_PART: u8 = 5; |
| 81 | pub const KIND_FORGOTTEN: u8 = 6; |
| 82 | pub 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. |
| 93 | pub 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. |
| 105 | pub 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. |
| 121 | pub 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)] |
| 138 | pub 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 | |
| 214 | impl 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)] |
| 677 | pub struct Parts { |
| 678 | held: Option<Holding>, |
| 679 | } |
| 680 | |
| 681 | #[derive(Clone, Debug)] |
| 682 | struct 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 | |
| 689 | impl 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. |
| 816 | fn 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. |
| 824 | fn 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)] |
| 851 | mod 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 |