Oregami
Repositories/oxedyne/ore

oxedyne/ore/relay/src/proto.rs

44.8 KiB, 107 runs

created by r2848102244:149, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! The transport both ends speak: where a request goes, how the bytes are
2//! framed, and how a request says who is making it.
3//!
4//! Two version surfaces, deliberately separate. An `ORESYN` message carries its
5//! own magic and version byte and refuses an unknown version by name; that
6//! governs message *content* and is the engine's. This module's version lives in
7//! the path -- `/ore/v1/...` -- and governs framing, authentication and the
8//! binding vocabulary, so a transport change never masquerades as a protocol
9//! change. An unauthenticated `GET /ore` says which of each the relay speaks, so
10//! an old client fails with a sentence naming both sides rather than a decode
11//! error part way through an exchange.
12//!
13//! # Framing belongs here
14//!
15//! `ORESYN` messages own their bytes but not their boundaries. A request or
16//! response body is a sequence of frames, each a four-byte big-endian length
17//! followed by one encoded message. Where a whole exchange fits in two bodies it
18//! is two round trips: the client posts its opening, the response carries the
19//! relay's opening, what the client is owed and `Done`; the client posts what the
20//! relay is owed and `Done`, and the response acknowledges.
21//!
22//! It is only two where it fits. Since 2026-08-20 both directions are bounded --
23//! [`POST_BYTES`] on the way in, [`REPLY_BYTES`] on the way out -- so a turn's
24//! frames go out across as many requests as [`groups`] makes of them, and a reply
25//! that was cut leaves the client to open a fresh session. A 58 MB clone measured
26//! twenty-eight requests over fourteen sessions.
27//!
28//! # A request signs itself
29//!
30//! Replica keys already exist and are the right credential; no second identity
31//! system is introduced. A request carries the replica identifier, the public
32//! key, a timestamp and a signature over [`statement`] -- the method, the path,
33//! the replica, the timestamp and a digest of the body. The relay checks
34//! freshness within [`SKEW`] to blunt replay, and verifies the signature before
35//! it touches a repository.
36//!
37//! The replica identifier is inside what is signed, so a request that verifies
38//! proves both control of the key and which replica claims it. That is enough
39//! for the relay to decide what this caller may do, and not enough to hand on:
40//! a third party cannot check a signature over a request it never saw. Bindings
41//! that travel are [`ore_store::keys::Binding`] values, each signed by the key it
42//! binds, deposited over the keys route -- including the ones belonging to
43//! replicas that never speak to this relay, which an imported git history is full
44//! of. So the relay carries bindings and cannot mint them.
45//!
46//! TLS is the deployment's, not this module's. The request signature is not a
47//! substitute for it but an authentication inside it, so the relay never handles
48//! a password and stores no client secret.
49
50use ore_store::keys::{
51 algorithm,
52 bytes_of,
53 text_of,
54 Signing,
55};
56
57use oxedyne_fe2o3_core::prelude::*;
58use oxedyne_fe2o3_hash::sha256::Sha256;
59use oxedyne_fe2o3_jdat::prelude::*;
60use oxedyne_fe2o3_iop_crypto::keys::KeyManager;
61use oxedyne_fe2o3_iop_crypto::sign::Signer;
62use oxedyne_fe2o3_ore::segment::Entry;
63use oxedyne_fe2o3_ore::sync::Message;
64
65use std::time::{
66 SystemTime,
67 UNIX_EPOCH,
68};
69
70
71/// Version of this transport, which appears in every path.
72pub const VERSION: &str = "v1";
73
74/// What every path of this transport begins with.
75pub const PREFIX: &str = "/ore";
76
77/// What a request signs, so that a change of shape cannot be read as an old one.
78pub const REQUEST_TAG: &str = "ORE-REQ-1";
79
80/// The header naming the replica making a request.
81pub const HEADER_REPLICA: &str = "x-ore-replica";
82/// The header carrying the public key the request is signed with.
83pub const HEADER_KEY: &str = "x-ore-key";
84/// The header carrying the request's timestamp, in seconds since the epoch.
85pub const HEADER_TIME: &str = "x-ore-time";
86/// The header carrying the signature over [`statement`].
87pub const HEADER_SIG: &str = "x-ore-sig";
88
89/// How far a request's timestamp may stand from the relay's clock, in seconds.
90///
91/// Wide enough that an unsynchronised laptop still works, narrow enough that a
92/// captured request is not replayable for long. It is a blunting and not a
93/// defence: a replayed push is a push of operations the log already holds, which
94/// the log drops.
95pub const SKEW: u64 = 300;
96
97/// How many bytes of entries one `Send` message may carry before the transport
98/// starts another.
99///
100/// A session answers with the whole owed set in a single message, so a clone of
101/// a large history would be one message and an interruption at ninety percent
102/// would lose ninety percent. Splitting is legal without touching the engine: a
103/// receiver handles any number of `Send` messages before `Done`, each is closure
104/// checked independently, and any append-order prefix of an owed set is causally
105/// closed by construction.
106pub const BATCH_BYTES: usize = 4 << 20;
107
108/// How large a single frame may be, which bounds what a peer can make the reader
109/// allocate.
110pub const FRAME_LIMIT: usize = 64 << 20;
111
112/// What a length prefix costs in front of every framed message.
113pub const FRAME_PREFIX: usize = 4;
114
115/// What a `Send` message costs around the entries it carries.
116///
117/// The magic, the version, the kind, and the two compact list prefixes. It is
118/// used to decide whether one entry can travel as a message at all, so it is
119/// rounded up: an entry that fits with room to spare travels whole, and one that
120/// is close enough to the cap for this to matter is cut up instead, which costs
121/// a piece and is always safe.
122const SEND_FRAMING: usize = 32;
123
124/// How much one HTTP request body may carry, whatever it is carrying.
125///
126/// [`BATCH_BYTES`] bounds a single message; this bounds a REQUEST, and until
127/// 2026-08-20 nothing did. `frame` concatenates every message it is handed, so
128/// a session with a thousand messages to send built one body of a thousand
129/// batches and posted it. That works against a relay reached directly and fails
130/// against every ordinary deployment, because a reverse proxy in front of the
131/// relay caps request bodies: Steel's own default is 8 MiB (`http_max_body_bytes`
132/// in `fe2o3_steel`), nginx's is 1 MiB, and neither is unusual.
133///
134/// The failure it produced is not a clean refusal either. The proxy closes the
135/// connection part way through the body, so the client sees `Broken pipe` while
136/// writing and the relay never sees a request at all -- which reads as the relay
137/// being down. Measured on the first repository of any size to try it: an 84 MB
138/// push of 44,182 operations, against a proxy that would have taken 8.
139///
140/// Six mebibytes, which is under Steel's default with room for the headers and
141/// above [`BATCH_BYTES`] so that a chunk carrying one maximal message still fits.
142/// A body is never sized to a particular proxy's limit; it is sized so that the
143/// ordinary ones do not have to be reconfigured to accept it.
144///
145/// This is what a relay *publishes*, and no longer only what a client assumes.
146/// Six mebibytes was itself a guess at somebody else's proxy -- the same fault one
147/// layer up from the one it fixed -- and it dies behind a stock nginx exactly as
148/// eighty-four megabytes died behind Steel. A relay states its own limit on
149/// `GET /ore` and beside its bindings, and a client that is told nothing assumes
150/// [`POST_FALLBACK`].
151///
152/// Half a fix, and the half that is missing should be said plainly. This is a
153/// `const`, published as it was compiled: there is no `Host` field for it as
154/// there is for [`REPLY_BYTES`] and no `--post-bytes` on `ore-relay serve`, so an
155/// operator whose proxy takes less than six mebibytes cannot lower what this relay
156/// advertises without a rebuild. A client also holds this number as a ceiling of
157/// its own and posts the smaller of the two, so a relay that publishes more than
158/// its client was built for gains nothing by saying so.
159pub const POST_BYTES: usize = 6 << 20;
160
161/// What a client posts to a relay that does not say what it will accept.
162///
163/// One mebibyte, which is nginx's default and the smallest limit in ordinary use.
164/// It is deliberately far below [`POST_BYTES`]: a relay too old to publish a
165/// limit is also a relay whose proxy is unknown, and that is the case with no
166/// evidence to debug from. Guessing low costs round trips, which the walk's own
167/// posture calls a cost rather than a fault; guessing high costs a connection
168/// closed part way through a body, which reads as the relay being down.
169pub const POST_FALLBACK: usize = 1 << 20;
170
171/// How much one HTTP response body may carry.
172///
173/// The mirror of [`POST_BYTES`], and it exists for a failure nobody had met,
174/// because proxies cap requests and not responses. Measured 2026-08-20: a 58 MB
175/// clone arriving as one response peaked at 347,532 kB resident in the receiving
176/// process, about six times the payload. That dies on the receiving end, where
177/// there is no proxy to blame and a phone or a small VPS has nothing to raise.
178///
179/// This doc used to name the cause: that `serve.rs` framed the whole reply as one
180/// body and the client writes its segments once at the end of an exchange. It is
181/// not the cause. The same clone across fourteen bounded replies peaked at
182/// 346,612 kB, under one percent lower, because the six times is the engine's
183/// per-operation cost of holding a history -- `ore log` pays it with no network
184/// at all. The bound stays because it is still worth having: it means no single
185/// response body must be materialised whole, whatever the history's size.
186///
187/// A relay that reaches this stops adding to the reply and says [`Message::Done`],
188/// which is a claim about what this end will send and not about what the two logs
189/// hold. The client notices that its log does not cover the frontier the relay
190/// opened with, and comes back. Nothing is remembered between the two visits by
191/// either end.
192pub const REPLY_BYTES: usize = 6 << 20;
193
194
195/// Returns how many leading messages frame to no more than `cap`, and whether
196/// anything was left behind.
197///
198/// Always keeps at least one message, so a single message larger than the cap
199/// still goes rather than an exchange stalling on it; [`BATCH_BYTES`] is what
200/// bounds that case, and it is below every cap here.
201///
202/// # Arguments
203///
204/// * `cap` - The most the kept messages may frame to; see [`REPLY_BYTES`].
205pub fn upto(msgs: &[Message], cap: usize)
206 -> Outcome<(usize, bool)>
207{
208 let mut running = 0usize;
209 let mut carried = false;
210 let mut at = 0usize;
211 while at < msgs.len() {
212 // A run of pieces is one unit here, because cutting inside one would send
213 // half an operation that the far end must throw away -- and the next
214 // session would recompute the same owed set and send the same half again,
215 // for ever. A reply is the one direction where that trap exists: a request
216 // is cut across several bodies on purpose and the RELAY holds the pieces
217 // between them, whereas a reply is answered by an end that keeps nothing,
218 // so what it does not send whole it can never continue.
219 let end = run_end(msgs, at);
220 let mut size = 0usize;
221 for msg in &msgs[at..end] {
222 // Four bytes of length prefix, exactly as `frame` writes it.
223 size += res!(msg.encode()).len() + FRAME_PREFIX;
224 }
225 // The cap yields to progress. A reply that fits the bound by carrying only
226 // the opening tells the client nothing it did not know, so the client comes
227 // back and is told nothing again -- an exchange that never ends and never
228 // fails. At least one unit carrying operations goes, whatever it costs.
229 //
230 // In ordinary use this never arises for a whole message, because `split`
231 // bounds one at `BATCH_BYTES` and the callers keep that at or below the
232 // cap. It arises for a run whenever an operation is larger than the reply
233 // bound, which is the case this was written for.
234 let works = msgs[at..end].iter().any(|m| m.carries_operations());
235 let must = !carried && works;
236 if at > 0 && !must && running + size > cap {
237 return Ok((at, true));
238 }
239 running += size;
240 carried = carried || works;
241 at = end;
242 }
243 Ok((msgs.len(), false))
244}
245
246
247/// Returns where the run of pieces beginning at `at` ends, which is one past
248/// `at` for anything that is not a piece.
249///
250/// A run is the pieces of one operation, in order and with nothing between them,
251/// which is what [`split`] emits and what `Message::Parts` insists on. A run that
252/// is malformed is not repaired here: the end is taken at the first message that
253/// does not continue it, and the far end refuses the run by name.
254fn run_end(msgs: &[Message], at: usize) -> usize {
255 let (id, total) = match msgs.get(at) {
256 Some(Message::Part { id, seq: 0, total, .. }) => (*id, *total),
257 _ => return at + 1,
258 };
259 let mut end = at + 1;
260 while end < msgs.len() {
261 match &msgs[end] {
262 Message::Part { id: got, seq, total: n, .. }
263 if *got == id && *n == total && *seq == (end - at) as u64 => end += 1,
264 _ => break,
265 }
266 }
267 end
268}
269
270
271/// Returns the path a sync exchange is posted to.
272pub fn sync_path(account: &str, name: &str) -> String {
273 fmt!("{}/{}/{}/{}/sync", PREFIX, VERSION, account, name)
274}
275
276/// Returns the path key bindings are fetched from and deposited at.
277pub fn keys_path(account: &str, name: &str) -> String {
278 fmt!("{}/{}/{}/{}/keys", PREFIX, VERSION, account, name)
279}
280
281/// Returns the path veil key bindings and wraps are fetched from and deposited
282/// at.
283///
284/// A route of its own rather than a second list on [`keys_path`]. The bindings
285/// there populate the trust set that decides provenance, and a veil binding
286/// decides nothing about provenance; and leaving that route's bytes exactly as
287/// they were keeps a relay built before this able to answer a client built
288/// after it, which is worth more than a saved round trip.
289pub fn wraps_path(account: &str, name: &str) -> String {
290 fmt!("{}/{}/{}/{}/wraps", PREFIX, VERSION, account, name)
291}
292
293/// Returns the seconds since the epoch, for a request's timestamp.
294pub fn now()
295 -> Outcome<u64>
296{
297 Ok(res!(SystemTime::now().duration_since(UNIX_EPOCH)).as_secs())
298}
299
300
301/// Returns the bytes a request is signed over.
302///
303/// The body is covered by its digest rather than its length, so a relay that
304/// checks the signature has checked the bytes it is about to read. Every field
305/// ends in a line feed and none may contain one, so no two different requests
306/// can produce one statement.
307pub fn statement(method: &str, path: &str, replica: u64, stamp: u64, body: &[u8]) -> Vec<u8> {
308 let mut sha = Sha256::new();
309 sha.update(body);
310 fmt!(
311 "{}\n{}\n{}\n{}\n{}\n{}\n",
312 REQUEST_TAG, method, path, replica, stamp, text_of(&sha.finish()),
313 ).into_bytes()
314}
315
316
317/// What a request said about who is making it, before any of it is believed.
318#[derive(Clone, Debug)]
319pub struct Presented {
320 /// The replica the key claims to belong to.
321 pub replica: u64,
322 /// The public key the signature is to be checked against.
323 pub public: Vec<u8>,
324 /// When the request says it was made, in seconds since the epoch.
325 pub stamp: u64,
326 /// The signature over [`statement`].
327 pub sig: Vec<u8>,
328}
329
330impl Presented {
331
332 /// Reads the four header values, or says which is missing or malformed.
333 pub fn read(replica: &str, key: &str, stamp: &str, sig: &str)
334 -> Outcome<Self>
335 {
336 let replica = match replica.trim().parse::<u64>() {
337 Ok(n) => n,
338 Err(e) => return Err(err!(e,
339 "The {} header {:?} is not a replica identifier.", HEADER_REPLICA, replica;
340 Invalid, Input, Decode)),
341 };
342 let stamp = match stamp.trim().parse::<u64>() {
343 Ok(n) => n,
344 Err(e) => return Err(err!(e,
345 "The {} header {:?} is not a number of seconds.", HEADER_TIME, stamp;
346 Invalid, Input, Decode)),
347 };
348 Ok(Self {
349 replica,
350 public: res!(bytes_of(key.trim())),
351 stamp,
352 sig: res!(bytes_of(sig.trim())),
353 })
354 }
355
356 /// Returns the four headers a client sends, signed with its key.
357 pub fn sign(signer: &Signing, method: &str, path: &str, body: &[u8])
358 -> Outcome<Vec<(String, String)>>
359 {
360 let replica = signer.replica.inner();
361 let stamp = res!(now());
362 let sig = res!(signer.sign(&statement(method, path, replica, stamp, body)));
363 Ok(vec![
364 (fmt!("{}", HEADER_REPLICA), fmt!("{}", replica)),
365 (fmt!("{}", HEADER_KEY), signer.public_text()),
366 (fmt!("{}", HEADER_TIME), fmt!("{}", stamp)),
367 (fmt!("{}", HEADER_SIG), text_of(&sig)),
368 ])
369 }
370
371 /// Checks the signature and the freshness, and says so plainly if either
372 /// fails.
373 ///
374 /// Nothing about the repository is consulted here. Whether this key may do
375 /// what it is asking to do is [`crate::acl`]'s question, and it is asked
376 /// afterwards.
377 pub fn check(&self, method: &str, path: &str, body: &[u8], now: u64)
378 -> Outcome<()>
379 {
380 let drift = now.abs_diff(self.stamp);
381 if drift > SKEW {
382 return Err(err!(
383 "A request is stamped {} seconds from this relay's clock, and {} is as \
384 far as it will take. Either the clock at one end is wrong, or this \
385 request has been kept and sent again.", drift, SKEW;
386 Invalid, Input, Security, Timeout));
387 }
388 let scheme = match algorithm().clone_with_keys(Some(&self.public), None) {
389 Ok(s) => s,
390 Err(e) => return Err(err!(e,
391 "A request carries a {} byte public key, which the signature scheme does \
392 not accept.", self.public.len();
393 Invalid, Input, Key)),
394 };
395 let want = statement(method, path, self.replica, self.stamp, body);
396 let ok = match scheme.verify(&want, &self.sig) {
397 Ok(v) => v,
398 Err(e) => return Err(err!(e,
399 "The signature on a request could not be checked: its {} byte public key \
400 and {} byte signature are not a pair.", self.public.len(), self.sig.len();
401 Invalid, Input, Security, Mismatch)),
402 };
403 if !ok {
404 return Err(err!(
405 "The signature on a request does not verify against the public key it \
406 carries. Either the request was altered on the way, or it was never \
407 signed by the holder of that key.";
408 Invalid, Input, Security, Mismatch));
409 }
410 Ok(())
411 }
412}
413
414
415/// Cuts a run of messages into groups, each of which frames to at most `cap`
416/// bytes, so that one HTTP request body carries one group.
417///
418/// A message that is larger than `cap` on its own still goes, alone, rather than
419/// being refused: `split` has already bounded what a single message can be, and a
420/// group of one is the smallest request that can carry it. Refusing here would
421/// turn a payload the protocol had already made legal into a sync that cannot
422/// complete.
423///
424/// Returns index ranges rather than copies. The bodies are built one at a time by
425/// the caller, so a session with a thousand batches to send never holds a
426/// thousand batches of framed bytes at once -- which was the other half of what
427/// the single-body version cost.
428///
429/// # Arguments
430/// * `msgs` - The messages to divide, in the order they must be sent.
431/// * `cap` - The most any one group may frame to; see [`POST_BYTES`].
432pub fn groups(msgs: &[Message], cap: usize)
433 -> Outcome<Vec<std::ops::Range<usize>>>
434{
435 let mut out = Vec::new();
436 let mut at = 0usize;
437 let mut start = 0usize;
438 let mut running = 0usize;
439 while at < msgs.len() {
440 // Four bytes of length prefix, exactly as `frame` writes it.
441 let size = res!(msgs[at].encode()).len() + 4;
442 if running > 0 && running + size > cap {
443 out.push(start..at);
444 start = at;
445 running = 0;
446 }
447 running += size;
448 at += 1;
449 }
450 if start < msgs.len() {
451 out.push(start..msgs.len());
452 }
453 Ok(out)
454}
455
456/// Writes messages as frames: a four-byte big-endian length, then the message.
457pub fn frame(msgs: &[Message])
458 -> Outcome<Vec<u8>>
459{
460 let mut out = Vec::new();
461 for msg in msgs {
462 let bytes = res!(msg.encode());
463 if bytes.len() > FRAME_LIMIT {
464 return Err(err!(
465 "A {} message comes to {} bytes, and a frame carries at most {}.",
466 msg.name(), bytes.len(), FRAME_LIMIT;
467 Invalid, Data, Excessive));
468 }
469 out.extend_from_slice(&(bytes.len() as u32).to_be_bytes());
470 out.extend_from_slice(&bytes);
471 }
472 Ok(out)
473}
474
475/// Reads the frames of a body back into messages.
476///
477/// A body that ends part way through a frame is refused rather than truncated:
478/// half a message is not a message, and a peer that absorbed the frames before it
479/// would be absorbing an arbitrary subset.
480pub fn unframe(bytes: &[u8])
481 -> Outcome<Vec<Message>>
482{
483 let mut out = Vec::new();
484 let mut at = 0usize;
485 while at < bytes.len() {
486 if bytes.len() - at < 4 {
487 return Err(err!(
488 "A body of {} bytes ends with {} of a frame's four byte length.",
489 bytes.len(), bytes.len() - at;
490 Decode, Input, Missing));
491 }
492 let mut len = [0u8; 4];
493 len.copy_from_slice(&bytes[at..at + 4]);
494 let len = u32::from_be_bytes(len) as usize;
495 at += 4;
496 if len > FRAME_LIMIT {
497 return Err(err!(
498 "A frame declares {} bytes, and a reader takes at most {}.",
499 len, FRAME_LIMIT;
500 Decode, Input, Excessive));
501 }
502 if bytes.len() - at < len {
503 return Err(err!(
504 "A frame declares {} bytes and only {} follow it.",
505 len, bytes.len() - at;
506 Decode, Input, Missing));
507 }
508 out.push(res!(Message::decode(&bytes[at..at + len])));
509 at += len;
510 }
511 Ok(out)
512}
513
514
515/// Splits a message that carries more than `cap` bytes of entries into several
516/// that do not.
517///
518/// Only a `Send` is split, and only along its entries, which arrive in the log's
519/// append order. Append order is a linear extension of the causal order, so every
520/// prefix of an owed set is causally closed and each piece is absorbed on its own
521/// merits. Anything else is handed back whole.
522///
523/// # One entry larger than the cap
524///
525/// Until 2026-08-22 such an entry went alone and over the cap, because an entry
526/// was the smallest thing the protocol had. That made the cap a wish: fe2o3's
527/// history holds a 22,153,680 byte compiled binary, added and later deleted, and
528/// the import reads git history, so it is one operation and it became one
529/// 22,153,955 byte request body against a proxy that would take eight mebibytes.
530/// The client had sized itself correctly to what the relay published and still
531/// had to send it, because a published limit only helps a sender that can
532/// subdivide below it.
533///
534/// So an entry that will not fit a message under `cap` is cut into
535/// `Message::Part` pieces instead, and every message this returns frames to at
536/// most `cap`. The pieces of one entry are contiguous and must stay in the order
537/// they are given.
538pub fn split(msg: Message, cap: usize)
539 -> Outcome<Vec<Message>>
540{
541 let entries = match msg {
542 Message::Send { entries } => entries,
543 other => return Ok(vec![other]),
544 };
545 let mut batching = Batching::to(cap);
546 let mut out: Vec<Message> = Vec::new();
547 for entry in entries {
548 match res!(batching.take(entry)) {
549 Fit::Took(msgs) => out.extend(msgs),
550 // A batching with no bound on what it takes never fills.
551 Fit::Full(entry) => return Err(err!(
552 "A batching gathering the whole of a send handed {} back rather than \
553 taking it.", res!(entry.id());
554 Bug, Unreachable)),
555 }
556 }
557 out.extend(batching.rest());
558 Ok(out)
559}
560
561
562/// Gathers entries into `Send` messages no larger than a cap, one entry at a
563/// time.
564///
565/// [`split`] is this driven over a whole message, and is what a caller that means
566/// to send everything wants. This is for the caller that does not: a relay bounds
567/// its reply, and an entry offered here is an entry that has been measured and
568/// copied, so a relay that offered the whole owed set would pay for the history
569/// it is about to throw away. Ask [`Batching::taken`] between entries and stop.
570///
571/// The batching is [`split`]'s, unchanged and in one place, so the two cannot
572/// drift: a batch is flushed when the next entry would take it past the cap, and
573/// an entry larger than a message becomes a run of [`Message::Part`] pieces which
574/// leaves contiguously and in order.
575pub struct Batching {
576 cap: usize, // what one message may frame to
577 room: usize, // what its entries may come to, the framing taken off
578 budget: usize, // what every entry taken may come to
579 batch: Vec<Entry>, // the entries gathered for the next message
580 held: usize, // what those come to
581 taken: usize, // what every entry taken so far came to
582 emitted: bool, // a message has been handed back
583}
584
585
586/// What became of an entry a [`Batching`] was offered.
587pub enum Fit {
588 Took(Vec<Message>), // taken; these messages are complete because of it
589 Full(Entry), // not taken, the budget being met; handed straight back
590}
591
592impl Batching {
593
594 /// Gathering into messages of at most `cap` framed bytes, all of them.
595 pub fn to(cap: usize) -> Self {
596 Self::upto(cap, usize::MAX)
597 }
598
599 /// The same, stopping once `budget` bytes of entries have been taken.
600 ///
601 /// The first entry is always taken, whatever it comes to, so that a budget
602 /// smaller than one operation still makes progress rather than sending nothing
603 /// for ever. That is [`upto`]'s rule at the other end of the same reply, and
604 /// the two must agree or an exchange stalls.
605 pub fn upto(cap: usize, budget: usize) -> Self {
606 Self {
607 cap,
608 room: cap.saturating_sub(SEND_FRAMING + FRAME_PREFIX),
609 budget,
610 batch: Vec::new(),
611 held: 0,
612 taken: 0,
613 emitted: false,
614 }
615 }
616
617 /// What every entry taken so far came to.
618 ///
619 /// Entries and not frames, so it is under what the messages will come to and
620 /// never over.
621 pub fn taken(&self) -> usize {
622 self.taken
623 }
624
625 /// Takes one entry, and hands back the messages that are complete because of
626 /// it, which is usually none.
627 ///
628 /// An entry that would take the total past the budget is handed straight back
629 /// instead, and nothing after it should be offered. It is measured before that
630 /// is known -- an entry's size is what serialising it says -- but it is not
631 /// batched, not cut into pieces, and not copied into a message.
632 pub fn take(&mut self, entry: Entry)
633 -> Outcome<Fit>
634 {
635 // The entry as it will be sent, not the record inside it. A sealed entry
636 // carries a key and a signature beside its record and a veiled one carries
637 // ciphertext instead of it, so measuring the record would undercount the
638 // first and be unable to measure the second at all.
639 let size = res!(entry.to_dat().to_bytes(Vec::new())).len();
640 if self.taken > 0 && self.taken.saturating_add(size) > self.budget {
641 return Ok(Fit::Full(entry));
642 }
643 self.taken += size;
644 let mut out: Vec<Message> = Vec::new();
645 if size > self.room {
646 // What is already batched goes first, because the pieces of one entry
647 // must be contiguous and a batch flushed after them would sit inside the
648 // run.
649 if !self.batch.is_empty() {
650 out.push(Message::Send { entries: std::mem::take(&mut self.batch) });
651 self.held = 0;
652 }
653 out.extend(res!(Message::part(&entry, self.cap.saturating_sub(FRAME_PREFIX))));
654 self.emitted = true;
655 return Ok(Fit::Took(out));
656 }
657 if !self.batch.is_empty() && self.held + size > self.room {
658 out.push(Message::Send { entries: std::mem::take(&mut self.batch) });
659 self.held = 0;
660 self.emitted = true;
661 }
662 self.held += size;
663 self.batch.push(entry);
664 Ok(Fit::Took(out))
665 }
666
667 /// Hands back what is still gathered, which ends the run of messages.
668 pub fn rest(self) -> Vec<Message> {
669 // An empty send is still a send: a caller that built one meant to say it.
670 match self.batch.is_empty() && self.emitted {
671 true => Vec::new(),
672 false => vec![Message::Send { entries: self.batch }],
673 }
674 }
675}
676
677
678#[cfg(test)]
679mod tests {
680 use super::*;
681
682 use oxedyne_fe2o3_ore::id::{
683 OpId,
684 ReplicaId,
685 };
686 use oxedyne_fe2o3_ore::op::{
687 Header,
688 Op,
689 Record,
690 };
691 use oxedyne_fe2o3_ore::sync::Parts;
692
693 /// An operation identifier.
694 fn oid(replica: u64, counter: u64) -> OpId {
695 OpId::new(ReplicaId::new(replica), counter)
696 }
697
698 /// A bare entry of roughly known size.
699 fn entry(counter: u64, text: &str)
700 -> Outcome<Entry>
701 {
702 Ok(Entry::Bare(Record::new(
703 res!(Header::new(oid(1, counter), if counter == 1 {
704 Vec::new()
705 } else {
706 vec![oid(1, counter - 1)]
707 })),
708 Op::Mark { name: fmt!("{}", text), body: None, time: None },
709 )))
710 }
711
712 /// Frames survive the round trip, and a body cut short is either a whole
713 /// prefix of it or nothing at all -- never a message the sender did not send.
714 ///
715 /// A cut that lands on a frame boundary is a shorter body and reads as one;
716 /// that is the framing working rather than failing, and the receiver's own
717 /// closure check is what decides whether a prefix is enough to absorb.
718 #[test]
719 fn frames_round_trip_and_a_cut_body_is_never_misread() -> Outcome<()> {
720 let msgs = vec![
721 Message::hello(vec![oid(1, 2), oid(3, 1)]),
722 Message::Send { entries: vec![res!(entry(1, "one")), res!(entry(2, "two"))] },
723 Message::Done,
724 ];
725 let bytes = res!(frame(&msgs));
726 assert_eq!(res!(unframe(&bytes)), msgs);
727 let mut refused = 0usize;
728 for cut in 0..bytes.len() {
729 match unframe(&bytes[..cut]) {
730 Err(_) => refused += 1,
731 Ok(got) => {
732 if got.len() >= msgs.len() || got[..] != msgs[..got.len()] {
733 return Err(err!(
734 "A body cut at {} of {} read as {} messages that are not its \
735 prefix.", cut, bytes.len(), got.len();
736 Test, Mismatch));
737 }
738 },
739 }
740 }
741 assert!(refused > 0, "a cut inside a frame is refused rather than half read");
742 // An empty body is no frames, which is what an acknowledgement is.
743 assert!(res!(unframe(&[])).is_empty());
744 Ok(())
745 }
746
747 /// A send larger than the cap becomes several sends, each an append-order
748 /// prefix of what is left, and nothing is lost or repeated.
749 #[test]
750 fn a_large_send_is_split_in_append_order() -> Outcome<()> {
751 let mut entries = Vec::new();
752 for i in 1..=20u64 {
753 entries.push(res!(entry(i, &fmt!("mark{}", i))));
754 }
755 let whole = Message::Send { entries: entries.clone() };
756 // Above the framing of a send and well below twenty entries, so the run is
757 // divided and every entry still travels whole.
758 let batches = res!(split(whole, 400));
759 assert!(batches.len() > 1, "a cap of 400 bytes splits twenty operations");
760 let mut seen = Vec::new();
761 for batch in &batches {
762 assert!(!batch.entries().is_empty(), "no batch is empty");
763 assert!(res!(batch.encode()).len() + FRAME_PREFIX <= 400,
764 "a batch frames to more than the cap it was bounded by");
765 seen.extend(batch.entries().iter().cloned());
766 }
767 assert_eq!(seen, entries, "the batches are the whole, in order");
768 // Anything else is handed back as it stands.
769 assert_eq!(res!(split(Message::Done, 100)), vec![Message::Done]);
770 Ok(())
771 }
772
773 /// A wide entry of roughly known size.
774 fn wide(counter: u64, bytes: usize)
775 -> Outcome<Entry>
776 {
777 Ok(Entry::Bare(Record::new(
778 res!(Header::new(oid(1, counter), vec![oid(1, counter - 1)])),
779 Op::Mark {
780 name: fmt!("wide{}", counter),
781 body: Some(vec![0xc3u8; bytes]),
782 time: Some(1_755_000_000),
783 },
784 )))
785 }
786
787 /// An entry larger than one message is cut into pieces, and EVERY message the
788 /// split produces fits the cap.
789 ///
790 /// This is the defect. Until 2026-08-22 an entry that would not fit went alone
791 /// and over the cap, because an entry was the smallest thing the protocol had:
792 /// fe2o3's history holds a 22,153,680 byte compiled binary as one operation,
793 /// and it became a 22,153,955 byte request body against a proxy that would
794 /// take eight mebibytes. The client had sized itself correctly to the number
795 /// the relay published and still had to send it.
796 ///
797 /// BOTH HALVES, because either alone is satisfied by something useless: a
798 /// split that drops the entry respects any cap, and the one that shipped
799 /// preserved the entry and ignored the cap.
800 #[test]
801 fn an_entry_larger_than_a_message_is_cut_into_pieces() -> Outcome<()> {
802 let big = res!(wide(2, 20_000));
803 let cap = 4_096usize;
804 let msgs = res!(split(Message::Send { entries: vec![big.clone()] }, cap));
805 assert!(msgs.len() > 4, "a 20 kB entry made {} messages at a cap of {}", msgs.len(), cap);
806 for m in &msgs {
807 assert!(res!(m.encode()).len() + FRAME_PREFIX <= cap,
808 "a {} message frames to {} against a cap of {}",
809 m.name(), res!(m.encode()).len() + FRAME_PREFIX, cap);
810 assert!(matches!(m, Message::Part { .. }),
811 "an oversized entry produced a {} message", m.name());
812 }
813 // And the far end gets the entry back, byte for byte.
814 let mut held = Parts::new();
815 let mut got = None;
816 for m in msgs {
817 // Through the wire, since that is where a spelling that only
818 // round-trips in memory would go wrong.
819 if let Some(whole) = res!(held.absorb(res!(Message::decode(&res!(m.encode()))))) {
820 got = Some(whole);
821 }
822 }
823 assert!(!held.pending());
824 match got {
825 Some(Message::Send { entries }) => {
826 assert_eq!(entries.len(), 1);
827 assert_eq!(
828 res!(entries[0].to_dat().to_bytes(Vec::new())),
829 res!(big.to_dat().to_bytes(Vec::new())),
830 "the entry that came back is not the entry that went",
831 );
832 },
833 other => return Err(err!(
834 "The pieces made {:?} rather than a send.", other; Test, Mismatch)),
835 }
836 // A group is one request body, and none of them is over the cap either.
837 let msgs = res!(split(Message::Send { entries: vec![res!(wide(2, 20_000))] }, cap));
838 for span in res!(groups(&msgs, cap)) {
839 assert!(res!(frame(&msgs[span.clone()])).len() <= cap,
840 "a group of {} framed over the cap", span.end - span.start);
841 }
842 Ok(())
843 }
844
845 /// The pieces of one entry stay contiguous, with whatever was already batched
846 /// sent before them.
847 ///
848 /// A message between the pieces is what `Parts` refuses, so a split that
849 /// flushed a batch into the middle of a run would produce a push no relay
850 /// would take.
851 #[test]
852 fn a_run_of_pieces_is_never_interrupted_by_a_batch() -> Outcome<()> {
853 let entries = vec![
854 res!(entry(1, "one")),
855 res!(wide(2, 20_000)),
856 res!(entry(3, "three")),
857 ];
858 let msgs = res!(split(Message::Send { entries }, 4_096));
859 let mut seen_run = false;
860 let mut ended = false;
861 for m in &msgs {
862 match m {
863 Message::Part { .. } => {
864 assert!(!ended, "a run of pieces started again after it had ended");
865 seen_run = true;
866 },
867 _ => {
868 if seen_run {
869 ended = true;
870 }
871 },
872 }
873 }
874 assert!(seen_run, "the wide entry was not cut up at all");
875 // The first message carries what was batched before the run, and the last
876 // what came after it.
877 assert_eq!(msgs[0].entries().len(), 1, "the earlier entry did not go first");
878 match msgs.last() {
879 Some(Message::Send { entries }) => assert_eq!(entries.len(), 1,
880 "the later entry did not go last"),
881 other => return Err(err!(
882 "The split ended with {:?}.", other; Test, Mismatch)),
883 }
884 Ok(())
885 }
886
887 /// A reply is never cut inside a run of pieces.
888 ///
889 /// A reply is answered by an end that keeps nothing between requests, so half
890 /// a run sent is half a run the far end throws away -- and the next session
891 /// computes the same owed set and sends the same half again, for ever. The
892 /// bound therefore yields to a whole run, exactly as it already yields to one
893 /// message that carries operations.
894 #[test]
895 fn a_reply_is_never_cut_inside_a_run_of_pieces() -> Outcome<()> {
896 let cap = 4_096usize;
897 let mut msgs = vec![Message::hello(vec![oid(1, 1)])];
898 msgs.extend(res!(split(Message::Send { entries: vec![res!(wide(2, 20_000))] }, cap)));
899 msgs.push(Message::Done);
900 let run = msgs.len() - 2;
901 assert!(run > 4, "the fixture has only {} pieces", run);
902
903 // A bound far under the run: it goes anyway and goes WHOLE, because nothing
904 // carrying operations has gone yet and a reply that carries none is a turn
905 // that tells the client nothing it did not know.
906 let (fits, _) = res!(upto(&msgs, cap));
907 assert_eq!(fits, run + 1, "the run of {} pieces was cut at {}", run, fits);
908 assert!(matches!(msgs[fits], Message::Done),
909 "the cut landed at {}, which is not the end of the run", fits);
910
911 // And with something carrying operations already in the reply, the run is
912 // held back WHOLE rather than begun.
913 let mut ahead = vec![
914 Message::hello(vec![oid(1, 1)]),
915 Message::Send { entries: vec![res!(entry(1, "one"))] },
916 ];
917 let head = res!(frame(&ahead)).len();
918 ahead.extend(res!(split(Message::Send { entries: vec![res!(wide(2, 20_000))] }, cap)));
919 let (fits, held_back) = res!(upto(&ahead, head + cap));
920 assert!(held_back, "the run fitted a bound one piece wide");
921 assert_eq!(fits, 2, "the reply was cut at {} rather than before the run", fits);
922 assert!(!matches!(ahead[fits], Message::Send { .. }),
923 "the cut landed inside a batch rather than before the run");
924 Ok(())
925 }
926
927 /// A signed request verifies, and every alteration of what was signed makes
928 /// it fail.
929 #[test]
930 fn a_request_signature_covers_what_it_says_it_covers() -> Outcome<()> {
931 let key = res!(Signing::mint(ReplicaId::new(7)));
932 let body = b"the frames of an exchange".to_vec();
933 let path = sync_path("oxedyne", "ore");
934 let headers = res!(Presented::sign(&key, "POST", &path, &body));
935 let value = |name: &str| -> Outcome<String> {
936 match headers.iter().find(|(n, _)| n == name) {
937 Some((_, v)) => Ok(v.clone()),
938 None => Err(err!(
939 "A signed request carries no {} header.", name; Test, Missing)),
940 }
941 };
942 let cred = res!(Presented::read(
943 &res!(value(HEADER_REPLICA)),
944 &res!(value(HEADER_KEY)),
945 &res!(value(HEADER_TIME)),
946 &res!(value(HEADER_SIG)),
947 ));
948 assert_eq!(cred.replica, 7);
949 let now = res!(now());
950 res!(cred.check("POST", &path, &body, now));
951 // The method, the path, the body and the replica are each inside it.
952 assert!(cred.check("GET", &path, &body, now).is_err(), "the method");
953 assert!(cred.check("POST", &keys_path("oxedyne", "ore"), &body, now).is_err(), "the path");
954 assert!(cred.check("POST", &path, b"other frames", now).is_err(), "the body");
955 let mut lying = cred.clone();
956 lying.replica = 8;
957 assert!(lying.check("POST", &path, &body, now).is_err(), "the replica");
958 // And a stale request is refused before its signature is even looked at.
959 assert!(cred.check("POST", &path, &body, now + SKEW + 1).is_err(), "the age");
960 Ok(())
961 }
962
963 /// Every group frames to at most the cap, and together they are the whole run
964 /// in its original order.
965 ///
966 /// BOTH HALVES, because either alone is satisfied by something useless: a
967 /// grouping that returns nothing respects any cap, and one that returns a
968 /// single group covering everything preserves any order. It was the second
969 /// that shipped -- `frame` was handed the entire run and built one body of it,
970 /// 84 MB on the first repository large enough to notice, against a proxy that
971 /// would take 8 MiB.
972 #[test]
973 fn test_groups_respect_the_cap_and_lose_nothing() -> Outcome<()> {
974 let msgs: Vec<Message> = (1..=40u64)
975 .map(|i| Ok(Message::Send { entries: vec![res!(entry(i, "some content to size it"))] }))
976 .collect::<Outcome<Vec<_>>>()?;
977 let whole = res!(frame(&msgs)).len();
978 // A cap that must divide this run, so the test is not measuring a single
979 // group calling itself a success.
980 let cap = whole / 5;
981 let spans = res!(groups(&msgs, cap));
982 assert!(spans.len() > 1, "the run was not divided at all: {} group(s)", spans.len());
983 let mut seen = 0usize;
984 for span in &spans {
985 assert_eq!(span.start, seen, "the groups are not contiguous, so order is lost");
986 seen = span.end;
987 let body = res!(frame(&msgs[span.clone()]));
988 // One message larger than the cap is allowed through alone; anything
989 // else must fit, or the request this becomes will not.
990 assert!(body.len() <= cap || span.end - span.start == 1,
991 "a group of {} framed to {} against a cap of {}",
992 span.end - span.start, body.len(), cap);
993 }
994 assert_eq!(seen, msgs.len(), "the groups do not cover the run");
995 Ok(())
996 }
997
998 /// A single message larger than the cap still goes, alone.
999 ///
1000 /// The alternative is a sync that cannot complete: `split` has already decided
1001 /// what a message may be, and refusing to carry one here would make a payload
1002 /// the protocol calls legal impossible to send.
1003 #[test]
1004 fn test_a_message_over_the_cap_goes_alone_rather_than_not_at_all() -> Outcome<()> {
1005 let msgs = vec![
1006 Message::Send { entries: vec![res!(entry(1, &"x".repeat(4096)))] },
1007 Message::Send { entries: vec![res!(entry(2, &"y".repeat(4096)))] },
1008 ];
1009 let spans = res!(groups(&msgs, 16));
1010 assert_eq!(2, spans.len(), "two oversized messages did not get a group each");
1011 for span in &spans {
1012 assert_eq!(1, span.end - span.start, "an oversized message was grouped with another");
1013 }
1014 Ok(())
1015 }
1016
1017 /// An empty run is an empty grouping, and not one empty group.
1018 ///
1019 /// A group of nothing would become a POST of nothing, which is a round trip
1020 /// spent saying nothing at all.
1021 #[test]
1022 fn test_nothing_to_send_is_no_requests() -> Outcome<()> {
1023 assert!(res!(groups(&[], POST_BYTES)).is_empty(), "an empty run produced a request");
1024 Ok(())
1025 }
1026
1027 /// What fits is kept, what does not is reported as held back, and the kept
1028 /// part frames to no more than the cap.
1029 ///
1030 /// BOTH HALVES, because either alone is satisfied by something useless: a
1031 /// bound that keeps nothing respects any cap, and one that keeps everything
1032 /// and says nothing was held back is what shipped -- `serve.rs` framed the
1033 /// whole reply, so a 58 MB clone had to be materialised entire in the
1034 /// receiving process. That whole-body materialisation is the failure, and not
1035 /// the six-times peak this once blamed on it; see [`REPLY_BYTES`].
1036 #[test]
1037 fn test_upto_keeps_what_fits_and_admits_the_rest() -> Outcome<()> {
1038 let msgs: Vec<Message> = (1..=40u64)
1039 .map(|i| Ok(Message::Send { entries: vec![res!(entry(i, "some content to size it"))] }))
1040 .collect::<Outcome<Vec<_>>>()?;
1041 let whole = res!(frame(&msgs)).len();
1042
1043 // A cap that must bite.
1044 let cap = whole / 5;
1045 let (fits, held_back) = res!(upto(&msgs, cap));
1046 assert!(held_back, "a run of {} bytes fitted a cap of {}", whole, cap);
1047 assert!(fits > 0, "nothing was kept, which is a reply that carries no progress");
1048 assert!(fits < msgs.len(), "everything was kept against a cap a fifth its size");
1049 assert!(res!(frame(&msgs[..fits])).len() <= cap,
1050 "the kept part frames to more than the cap it was bounded by");
1051
1052 // A run that fits is left whole and says so, or a client would come back
1053 // for ever against a relay that had already sent everything.
1054 let (all, more) = res!(upto(&msgs, whole));
1055 assert_eq!(all, msgs.len(), "a run that fits was cut");
1056 assert!(!more, "a run that fits reported something held back");
1057
1058 // One message over the cap still goes, alone: `split` has already decided
1059 // what a message may be, and a reply that refuses to carry one would
1060 // stall the exchange rather than bound it.
1061 let big = vec![res!(entry(99, &"x".repeat(4096)))];
1062 let one = vec![Message::Send { entries: big }];
1063 let (kept, over) = res!(upto(&one, 8));
1064 assert_eq!(kept, 1, "a single oversized message was dropped, so nothing can be sent");
1065 assert!(!over, "a single message alone reported a remainder there is no room for");
1066 Ok(())
1067 }
1068
1069 /// A batching under a budget hands back the entry that would cross it, having
1070 /// taken every one that fits, and it takes the first entry whatever it comes
1071 /// to.
1072 ///
1073 /// The second half is not a nicety. [`upto`] lets the first unit carrying
1074 /// operations through over the cap for exactly this reason: a bound that can
1075 /// refuse everything is a bound under which a clone never finishes, and a peer
1076 /// that comes back to be told nothing again is worse than one told too much
1077 /// once. The two ends of the same reply have to agree about it.
1078 ///
1079 /// Proved red both ways. Dropping the `self.taken > 0` guard from
1080 /// [`Batching::take`] makes an entry larger than the budget refuse itself and
1081 /// the last assertion fails with nothing taken; dropping the budget test
1082 /// altogether takes all ten and the first fails.
1083 #[test]
1084 fn a_batching_under_a_budget_stops_at_it() -> Outcome<()> {
1085 let all: Vec<Entry> = (1..=10u64)
1086 .map(|n| entry(n, "some content to give it a size"))
1087 .collect::<Outcome<Vec<_>>>()?;
1088 let mut sizes = Vec::new();
1089 for one in &all {
1090 sizes.push(res!(one.to_dat().to_bytes(Vec::new())).len());
1091 }
1092 // Room for three and half of the fourth, which is room for three. Measured
1093 // rather than multiplied: the first operation of a replica names no parent
1094 // and is the smaller for it.
1095 let budget = sizes[0] + sizes[1] + sizes[2] + sizes[3] / 2;
1096
1097 let mut batching = Batching::upto(BATCH_BYTES, budget);
1098 let mut took = 0usize;
1099 let mut refused = 0usize;
1100 for entry in all.clone() {
1101 match res!(batching.take(entry)) {
1102 Fit::Took(_) => took += 1,
1103 Fit::Full(_) => {
1104 refused += 1;
1105 break;
1106 },
1107 }
1108 }
1109 assert_eq!(took, 3, "a budget of three and a half entries took {}", took);
1110 assert_eq!(refused, 1, "the entry that would cross the budget was not handed back");
1111 assert!(batching.taken() <= budget,
1112 "the entries taken come to {} against a budget of {}", batching.taken(), budget);
1113
1114 // What was gathered still leaves as a message, and carries what was taken.
1115 let rest = batching.rest();
1116 let carried: usize = rest.iter().map(|m| match m {
1117 Message::Send { entries } => entries.len(),
1118 _ => 0,
1119 }).sum();
1120 assert_eq!(carried, 3, "the messages carry {} of the 3 entries taken", carried);
1121
1122 // A budget under one entry still takes one, or nothing ever crosses.
1123 let mut narrow = Batching::upto(BATCH_BYTES, 1);
1124 match res!(narrow.take(res!(entry(1, "wider than the budget it is offered to")))) {
1125 Fit::Took(_) => (),
1126 Fit::Full(_) => return Err(err!(
1127 "A budget smaller than one operation refused the first one, so a \
1128 repository holding it can never be cloned."; Test, Invalid)),
1129 }
1130 match res!(narrow.take(res!(entry(2, "and the one after it")))) {
1131 Fit::Full(_) => (),
1132 Fit::Took(_) => return Err(err!(
1133 "A batching over its budget took a second entry."; Test, Invalid)),
1134 }
1135 Ok(())
1136 }
1137}