Oregami
Repositories/oxedyne/ore

oxedyne/ore/cli/src/relay.rs

46.2 KiB, 143 runs

created by r2848102244:126, 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//! `ore sync <url>` -- the same exchange, over a network.
2//!
3//! The engine's session is peer-symmetric and holds no connection, so what
4//! carries the bytes is the caller's choice. [`crate::sync`] carries them across
5//! a filesystem both ends can read. This module carries them to a machine that is
6//! always on, which is the one thing a filesystem transport cannot do: two
7//! replicas that are never awake at the same time still converge, each by talking
8//! to the relay once.
9//!
10//! # What is on the other end
11//!
12//! Not a service that merges, and not an authority. The relay holds a log and
13//! runs the same session over it that any peer would, which is why `ore sync
14//! <url>` is the same exchange as `ore sync <path>` and not a second protocol.
15//! It authors nothing, signs nothing into a history and verifies nothing on the
16//! way in; every property this end relies on is checked at this end, against
17//! material the relay cannot forge. A relay that fails, lies or withholds delays
18//! convergence exactly as a partition would.
19//!
20//! # Two things cross besides the operations
21//!
22//! Key bindings go first, over a channel of their own, because what arrives is
23//! verified against what is known and a key learned afterwards would be learned
24//! too late. A binding is signed by the very key it binds, so the relay carries
25//! bindings without being able to mint one: a fabricated binding fails its own
26//! signature, and the worst a relay can do is withhold one, which leaves the
27//! affected operations marked `?` rather than misattributed. That is the part the
28//! filesystem transport could not inherit, and this is where it is answered.
29//!
30//! The operations themselves cross sealed, and what arrives is verified before a
31//! session is allowed to absorb any of it.
32//!
33//! # When the relay is unreachable
34//!
35//! The capture happens first, which is wanted regardless; then the command fails
36//! with a message naming the relay and the error, changes nothing else, and exits
37//! non-zero. Nothing is queued, no daemon retries and no state marks the failure,
38//! because a resume is a rerun: the next attempt opens a fresh session over
39//! frontiers that already reflect whatever crossed. A relay that vanishes for a
40//! month costs a month of convergence and nothing else.
41
42use crate::arrived as delta;
43use crate::keys::Prov;
44use crate::repo::Repo;
45use crate::sync::Party;
46use crate::tree;
47use crate::verbs;
48
49use ore_relay::proto::{
50 self,
51 Presented,
52};
53use ore_relay::serve::HEADER_ABSORBED;
54
55use ore_store::keys::{
56 Binding,
57 Signing,
58};
59use ore_store::store::arrived;
60use ore_store::veilkey::{
61 VeilBinding,
62 Wrap,
63};
64
65use oxedyne_fe2o3_core::prelude::*;
66use oxedyne_fe2o3_jdat::prelude::*;
67use oxedyne_fe2o3_net::acme::trust::letsencrypt_client_config;
68use oxedyne_fe2o3_net::http::client::{
69 http_request,
70 https_request,
71};
72use oxedyne_fe2o3_net::http::fields::HeaderName;
73use oxedyne_fe2o3_net::http::header::HttpMethod;
74use oxedyne_fe2o3_net::http::msg::HttpMessage;
75use oxedyne_fe2o3_ore::id::OpId;
76use oxedyne_fe2o3_ore::log::OpLog;
77use oxedyne_fe2o3_ore::sync::{
78 msg,
79 Message,
80 Mode,
81 Parts,
82 Session,
83};
84
85use std::collections::{
86 BTreeMap,
87 BTreeSet,
88};
89
90use tokio::runtime::Runtime;
91
92
93/// How many turns one session may take before it is called a fault.
94///
95/// A session is two turns in the common case. The bound is far above that and is
96/// there so that a session which somehow never converges stops rather than
97/// looping. It does not bound the exchange -- `SESSION_LIMIT` does that -- and a
98/// turn is not a request either, since one turn may leave as several HTTP posts.
99/// The counter is reset per session at the top of the loop below, and the error
100/// it raises still calls these round trips, which is the name they had when the
101/// exchange was one session and a turn was one request.
102pub const TRIP_LIMIT: usize = 16;
103
104/// How many whole sessions one `ore sync` will run before it stops and says so.
105///
106/// A relay bounds its reply (`proto::REPLY_BYTES`), so a clone larger than that
107/// takes several sessions, and the count rises with the history: 58 MB over a
108/// six-mebibyte bound divides into about ten, and the clone of 44,541 operations
109/// measured on 2026-08-20 took fourteen. The measured figure is the one to size
110/// by, since a reply is cut at a whole message and the division assumes every one
111/// is filled to the bound. That history is growing by roughly 1,300 operations a
112/// day. Sixty-four leaves a margin over the largest thing anyone here has, and
113/// still ends rather than spinning if the two ends stop making progress.
114///
115/// Reaching it is not an error. What arrived is durable and causally closed;
116/// the exchange is simply unfinished, and running the command again continues
117/// it from where this one stopped.
118pub const SESSION_LIMIT: usize = 64;
119
120
121/// A repository on a relay, as a URL named it.
122#[derive(Clone, Debug)]
123pub struct Remote {
124 /// Whether the connection is wrapped in TLS.
125 pub tls: bool,
126 /// The host name, which is also the name the certificate must carry.
127 pub host: String,
128 /// The port.
129 pub port: u16,
130 /// The account the repository is under.
131 pub account: String,
132 /// The repository's name.
133 pub name: String,
134}
135
136impl Remote {
137
138 /// Reports whether an argument to `ore sync` names a relay rather than a
139 /// path.
140 pub fn is_url(arg: &str) -> bool {
141 arg.starts_with("http://") || arg.starts_with("https://")
142 }
143
144 /// Reads `https://<host>[:<port>]/<account>/<name>`.
145 pub fn parse(url: &str)
146 -> Outcome<Self>
147 {
148 let (tls, rest) = match url.strip_prefix("https://") {
149 Some(rest) => (true, rest),
150 None => match url.strip_prefix("http://") {
151 Some(rest) => (false, rest),
152 None => return Err(err!(
153 "{:?} is neither a path nor a relay. A relay is named \
154 https://<host>/<account>/<repository>.", url;
155 Invalid, Input)),
156 },
157 };
158 let rest = rest.trim_end_matches('/');
159 let parts: Vec<&str> = rest.split('/').collect();
160 if parts.len() != 3 || parts.iter().any(|p| p.is_empty()) {
161 return Err(err!(
162 "{:?} names no repository on a relay. The shape is \
163 https://<host>/<account>/<repository>, as \
164 https://oregami.example/oxedyne/ore.", url;
165 Invalid, Input, Missing));
166 }
167 let (host, port) = match parts[0].rsplit_once(':') {
168 Some((host, port)) => {
169 let n = match port.parse::<u16>() {
170 Ok(n) => n,
171 Err(e) => return Err(err!(e,
172 "The port in {:?} is not a number between 0 and 65535.", url;
173 Invalid, Input)),
174 };
175 (fmt!("{}", host), n)
176 },
177 None => (fmt!("{}", parts[0]), if tls { 443 } else { 80 }),
178 };
179 if host.is_empty() {
180 return Err(err!(
181 "{:?} names no host.", url;
182 Invalid, Input, Missing));
183 }
184 Ok(Self {
185 tls,
186 host,
187 port,
188 account: fmt!("{}", parts[1]),
189 name: fmt!("{}", parts[2]),
190 })
191 }
192
193 /// Returns how the relay is named in a message about it.
194 pub fn label(&self) -> String {
195 fmt!(
196 "{}://{}/{}/{}",
197 if self.tls { "https" } else { "http" }, self.authority(), self.account, self.name,
198 )
199 }
200
201 /// Returns the host and port, with the port left off where it is the usual
202 /// one.
203 fn authority(&self) -> String {
204 let usual = if self.tls { 443 } else { 80 };
205 if self.port == usual {
206 fmt!("{}", self.host)
207 } else {
208 fmt!("{}:{}", self.host, self.port)
209 }
210 }
211
212 /// Makes one signed request and returns what came back, whatever the relay
213 /// made of it.
214 ///
215 /// Only a failure to reach the relay at all is an error here. A status is an
216 /// answer, and a caller that has something to do with a refusal -- a puller
217 /// told it may not push, say -- needs to see it rather than be stopped by it.
218 fn attempt(
219 &self,
220 rt: &Runtime,
221 signer: &Signing,
222 method: HttpMethod,
223 path: &str,
224 kind: &str,
225 body: Vec<u8>,
226 )
227 -> Outcome<(u16, HttpMessage)>
228 {
229 let named = match method {
230 HttpMethod::GET => "GET",
231 HttpMethod::POST => "POST",
232 other => return Err(err!(
233 "This transport speaks GET and POST, not {}.", other;
234 Bug, Invalid, Input)),
235 };
236 let signed = res!(Presented::sign(signer, named, path, &body));
237 let mut headers: Vec<(&str, &str)> = signed
238 .iter()
239 .map(|(n, v)| (n.as_str(), v.as_str()))
240 .collect();
241 headers.push(("Content-Type", kind));
242 let reply = if self.tls {
243 let tls = res!(letsencrypt_client_config());
244 rt.block_on(https_request(
245 &self.host, self.port, method, path, &headers, &body, tls,
246 ))
247 } else {
248 rt.block_on(http_request(
249 &self.host, self.port, method, path, &headers, &body,
250 ))
251 };
252 let reply = match reply {
253 Ok(r) => r,
254 Err(e) => return Err(err!(e,
255 "The relay {} could not be reached. Nothing here has changed but the \
256 capture this command made before it tried, and running the command \
257 again is the whole of the retry.", self.label();
258 IO, Network)),
259 };
260 let status = res!(status_of(&reply));
261 Ok((status, reply))
262 }
263
264 /// Makes one signed request and insists on being served.
265 ///
266 /// A status the relay refused with becomes an error carrying the relay's own
267 /// sentence, because the relay is where the reason is known and repeating it
268 /// is more use than replacing it.
269 fn ask(
270 &self,
271 rt: &Runtime,
272 signer: &Signing,
273 method: HttpMethod,
274 path: &str,
275 kind: &str,
276 body: Vec<u8>,
277 )
278 -> Outcome<HttpMessage>
279 {
280 let named = if method == HttpMethod::GET { "GET" } else { "POST" };
281 let (status, reply) = res!(self.attempt(rt, signer, method, path, kind, body));
282 if status != 200 {
283 return Err(err!(
284 "The relay {} answered {} to {} {}: {}",
285 self.label(), status, named, path,
286 fmt!("{}", reply.body_as_string()).trim();
287 IO, Network, Invalid));
288 }
289 Ok(reply)
290 }
291}
292
293/// Returns the status code of a response.
294fn status_of(msg: &HttpMessage)
295 -> Outcome<u16>
296{
297 use oxedyne_fe2o3_net::http::header::HttpHeadline;
298 match &msg.header.headline {
299 HttpHeadline::Response { status } => Ok(*status as u16),
300 HttpHeadline::Request { .. } => Err(err!(
301 "A request arrived where a response was expected.";
302 IO, Network, Invalid, Mismatch)),
303 }
304}
305
306/// Returns a named response header, where there is one.
307fn header_of(msg: &HttpMessage, name: &str) -> Option<String> {
308 msg.header.fields.get_one(&HeaderName::from(name)).map(|v| fmt!("{}", v))
309}
310
311
312/// What the relay said about itself when the bindings crossed.
313struct Standing {
314 /// The bindings the relay carries, each certified by the key it binds.
315 bindings: Vec<Binding>,
316 /// How many operations the relay's log holds.
317 ops: usize,
318 /// How many heads its frontier carries.
319 heads: usize,
320 /// The most it will take in one request body.
321 post: usize,
322 /// The oldest and newest ORESYN versions it speaks.
323 speaks: (u8, u8),
324}
325
326/// Deposits this repository's certified bindings and takes the relay's.
327///
328/// One round trip does both, because a client that pushes is a client that has
329/// something to introduce and something to learn, and asking twice would be two
330/// round trips for one fact.
331fn bindings(rt: &Runtime, remote: &Remote, signer: &Signing, mine: &[Binding])
332 -> Outcome<Standing>
333{
334 let offered: Vec<Dat> = mine
335 .iter()
336 .filter(|b| b.is_certified())
337 .map(|b| b.to_dat())
338 .collect();
339 let body = fmt!("{}\n", res!(Dat::List(offered).jdat_to_lines(" "))).into_bytes();
340 let path = proto::keys_path(&remote.account, &remote.name);
341 // Depositing is a write, so a replica that may only pull is refused it. That
342 // is not a failure to sync: it takes the bindings it is allowed to take and
343 // introduces nobody, which is exactly what a reader does.
344 let (status, reply) = res!(remote.attempt(
345 rt, signer, HttpMethod::POST, &path, "text/plain", body,
346 ));
347 let reply = if status == 200 {
348 reply
349 } else if status == 403 {
350 res!(remote.ask(rt, signer, HttpMethod::GET, &path, "text/plain", Vec::new()))
351 } else {
352 return Err(err!(
353 "The relay {} answered {} to POST {}: {}",
354 remote.label(), status, path, fmt!("{}", reply.body_as_string()).trim();
355 IO, Network, Invalid));
356 };
357 let dat = match Dat::decode_string(fmt!("{}", reply.body_as_string())) {
358 Ok(d) => d,
359 Err(e) => return Err(err!(e,
360 "The relay {} answered the key channel with something that is not JDAT.",
361 remote.label();
362 Decode, Input)),
363 };
364 let map = match &dat {
365 Dat::Map(m) => m,
366 other => return Err(err!(
367 "The relay {} answered the key channel with {:?} rather than a map.",
368 remote.label(), other;
369 Decode, Input, Mismatch)),
370 };
371 let count = |key: &str| -> Outcome<usize> {
372 match map.get(&Dat::Str(fmt!("{}", key))) {
373 Some(Dat::U64(n)) => Ok(*n as usize),
374 Some(Dat::U32(n)) => Ok(*n as usize),
375 Some(Dat::U8(n)) => Ok(*n as usize),
376 other => Err(err!(
377 "The relay's {:?} expects a number, got {:?}.", key, other;
378 Decode, Input, Mismatch)),
379 }
380 };
381 let listed = match map.get(&Dat::Str(fmt!("keys"))) {
382 Some(Dat::List(l)) => l,
383 other => return Err(err!(
384 "The relay's bindings expect a list, got {:?}.", other;
385 Decode, Input, Mismatch)),
386 };
387 let mut held = Vec::new();
388 for item in listed {
389 let binding = res!(Binding::from_dat(item));
390 // A binding that does not certify itself is one the relay could have
391 // written, so it is dropped rather than believed.
392 if binding.is_certified() {
393 held.push(binding);
394 }
395 }
396 // A relay that does not say what it will accept is a relay whose proxy is
397 // unknown, so the assumption is the smallest limit in ordinary use rather
398 // than this build's own idea of a reasonable one. Guessing low costs
399 // requests; guessing high costs a connection closed part way through a body,
400 // which is the failure that reads as the relay being down.
401 let post = match map.get(&Dat::Str(fmt!("post"))) {
402 Some(_) => res!(count("post")),
403 None => proto::POST_FALLBACK,
404 };
405 // Which message versions the relay speaks, which decides what happens to an
406 // operation larger than `post`. A relay that says nothing is one built before
407 // the piece message existed, so it is taken at version 1 in both directions
408 // and told so rather than guessed at.
409 let speaks = match map.get(&Dat::Str(fmt!("oresyn"))) {
410 Some(Dat::List(v)) if !v.is_empty() => {
411 let byte = |d: &Dat| -> Outcome<u8> {
412 match d {
413 Dat::U8(n) => Ok(*n),
414 Dat::U64(n) => Ok(*n as u8),
415 other => Err(err!(
416 "The relay's ORESYN version {:?} is not a number.", other;
417 Decode, Input, Mismatch)),
418 }
419 };
420 (res!(byte(&v[0])), res!(byte(&v[v.len() - 1])))
421 },
422 _ => (msg::VERSION_MIN, msg::VERSION_MIN),
423 };
424 Ok(Standing {
425 bindings: held,
426 ops: res!(count("ops")),
427 heads: res!(count("heads")),
428 post,
429 speaks,
430 })
431}
432
433
434/// What the relay carried on the wraps route.
435struct Carried {
436 /// How many veil key bindings the relay holds.
437 veils: usize,
438 /// How many of them were new here.
439 learned: usize,
440 /// How many wraps the relay holds.
441 wraps: usize,
442 /// Whether a wrap addressed to this replica was opened, giving this
443 /// repository the content key it had not got.
444 opened: bool,
445 /// Whether the relay speaks this route at all.
446 spoken: bool,
447}
448
449/// Deposits this replica's veil key binding and whatever wraps it has made, takes
450/// the relay's, and installs a content key this repository has not got.
451///
452/// **The relay serves wraps to anybody who may pull, and that is not a leak.** A
453/// wrap is a content key encrypted to one veil key; the secret that opens it
454/// never leaves the machine that minted it, so the relay carrying wraps about is
455/// exactly how a second replica comes to read a veiled repository, and not a way
456/// for anybody else to.
457///
458/// A relay built before this route answers 405, which is not a failure to sync:
459/// the operations cross either way, and what is lost is only the enrolment.
460fn carry_wraps(
461 rt: &Runtime,
462 remote: &Remote,
463 signer: &Signing,
464 repo: &mut Repo,
465 dry: bool,
466)
467 -> Outcome<Carried>
468{
469 let mut out = Carried {
470 veils: 0, learned: 0, wraps: 0, opened: false, spoken: true,
471 };
472 // A repository with no veil key, no content key and no wrap of its own has
473 // nothing to publish, nothing to deposit and no use for a veil binding it
474 // learned, so it spends no round trip on the question. That is every
475 // repository that is not veiled, which is most of them.
476 if repo.veilkey.is_none() && repo.veil.is_none() && repo.cfg.wraps.is_empty() {
477 out.spoken = false;
478 return Ok(out);
479 }
480 let mine: Vec<Dat> = match &repo.veilkey {
481 Some(key) => vec![res!(key.binding(signer)).to_dat()],
482 None => Vec::new(),
483 };
484 let held: Vec<Dat> = repo.cfg.wraps.iter().map(|w| w.to_dat()).collect();
485 let mut map = DaticleMap::new();
486 map.insert(Dat::Str(fmt!("veils")), Dat::List(mine));
487 map.insert(Dat::Str(fmt!("wraps")), Dat::List(held));
488 let body = fmt!("{}\n", res!(Dat::Map(map).jdat_to_lines(" "))).into_bytes();
489 let path = proto::wraps_path(&remote.account, &remote.name);
490 let (status, reply) = res!(remote.attempt(
491 rt, signer, HttpMethod::POST, &path, "text/plain", body,
492 ));
493 let reply = match status {
494 200 => reply,
495 // Depositing is a write, so a replica that may only pull takes what it is
496 // allowed to take and introduces nobody. That is what a reader does, and
497 // a reader is exactly who is waiting for a wrap.
498 403 => res!(remote.ask(rt, signer, HttpMethod::GET, &path, "text/plain", Vec::new())),
499 404 | 405 => {
500 out.spoken = false;
501 return Ok(out);
502 },
503 other => return Err(err!(
504 "The relay {} answered {} to POST {}: {}",
505 remote.label(), other, path, fmt!("{}", reply.body_as_string()).trim();
506 IO, Network, Invalid)),
507 };
508 let dat = match Dat::decode_string(fmt!("{}", reply.body_as_string())) {
509 Ok(d) => d,
510 Err(e) => return Err(err!(e,
511 "The relay {} answered the wrap channel with something that is not JDAT.",
512 remote.label();
513 Decode, Input)),
514 };
515 let map = match &dat {
516 Dat::Map(m) => m,
517 other => return Err(err!(
518 "The relay {} answered the wrap channel with {:?} rather than a map.",
519 remote.label(), other;
520 Decode, Input, Mismatch)),
521 };
522 let listed = |name: &str| -> Outcome<Vec<Dat>> {
523 match map.get(&Dat::Str(fmt!("{}", name))) {
524 Some(Dat::List(l)) => Ok(l.clone()),
525 None => Ok(Vec::new()),
526 Some(other) => Err(err!(
527 "The relay\'s {:?} expect a list, got {:?}.", name, other;
528 Decode, Input, Mismatch)),
529 }
530 };
531 for item in res!(listed("veils")) {
532 let binding = res!(VeilBinding::from_dat(&item));
533 out.veils += 1;
534 // A binding whose chain does not hold is dropped rather than believed:
535 // the relay could have written it, and writing it down here would mean
536 // wrapping a content key to whoever the relay named.
537 if repo.cfg.learn_veil(binding) {
538 out.learned += 1;
539 }
540 }
541 let mut carried = Vec::new();
542 for item in res!(listed("wraps")) {
543 carried.push(res!(Wrap::from_dat(&item)));
544 }
545 out.wraps = carried.len();
546 // The content key, where this repository has not got one and a wrap here is
547 // addressed to it. This is the moment a second replica stops being enrolled
548 // and starts being able to read.
549 if !dry && repo.veil.is_none() {
550 if let Some(key) = &repo.veilkey {
551 if let Some(wrap) = carried.iter().find(|w| w.to == key.public()) {
552 let content = res!(wrap.open(key));
553 res!(repo.set_veil(Some(&content)));
554 out.opened = true;
555 }
556 }
557 }
558 Ok(out)
559}
560
561/// Empties a message of what it would hand over, keeping what it says.
562///
563/// A message carrying operations becomes the statement that nothing further is
564/// coming, which is true of a caller that has chosen to offer nothing: `Done`
565/// is a claim about what this end will send and not about what the two logs
566/// hold. The session has already counted the operations as told, so the
567/// exchange still closes; the far end simply never hears them.
568///
569/// This is what makes the request a read whatever the far end is short of. A
570/// relay reads a request for the operations it lacks of what is offered, so a
571/// replica holding nothing the relay has not got syncs plainly on a `pull` grant.
572/// One operation of its own is the thing it cannot hand over, and the relay
573/// refuses the whole request rather than its push half, so emptying the messages
574/// is how such a replica reads at all.
575fn withholding(msg: Message) -> Message {
576 match msg {
577 Message::Send { .. } => Message::Done,
578 other => other,
579 }
580}
581
582/// `ore sync <url>` -- the same exchange, carried to a relay, leaving this
583/// repository holding the union.
584///
585/// `taking` is `--pull-only`: nothing this end holds is handed over. The capture
586/// has already happened -- it is the first thing every verb does, and the one
587/// thing this command leaves behind if the network fails.
588pub fn sync(repo: &mut Repo, url: &str, taking: bool, dry: bool)
589 -> Outcome<()>
590{
591 let remote = res!(Remote::parse(url));
592 let signer = match &repo.signer {
593 Some(key) => key.clone(),
594 None => return Err(err!(
595 "A relay knows a replica by its key, and this repository holds none. Run \
596 `ore key` to mint one; syncing with another repository on this machine \
597 needs no key and still works.";
598 Invalid, Configuration, Missing, Key)),
599 };
600 let rt = match Runtime::new() {
601 Ok(r) => r,
602 Err(e) => return Err(err!(e,
603 "A runtime for the connection could not be started.";
604 IO, Init)),
605 };
606
607 // The keys cross before the operations do, because what arrives is verified
608 // against what is known and a key learned afterwards would be learned too
609 // late.
610 println!();
611 let standing = res!(bindings(&rt, &remote, &signer, &repo.cfg.keys));
612 let mut learned = 0usize;
613 for binding in &standing.bindings {
614 if repo.cfg.learn(binding.clone()) {
615 learned += 1;
616 }
617 }
618 // The wraps cross next, and before the operations do: a repository that takes
619 // its content key here is one that can read what arrives below, and one that
620 // took it a moment later could not.
621 let carried = res!(carry_wraps(&rt, &remote, &signer, repo, dry));
622 // A rehearsal writes nothing, the learned keys included: it is answering a
623 // question, and a question that changed the configuration would be a poor one.
624 if !dry {
625 res!(repo.save_config());
626 }
627 println!("relay {}", remote.label());
628 println!("keys {} known here ({} learned now), {} carried there",
629 repo.cfg.keys.len(), learned, standing.bindings.len());
630 if carried.spoken {
631 println!("veils {} known here ({} learned now), {} carried there, {} wrap{}",
632 repo.cfg.veils.len(), carried.learned, carried.veils, carried.wraps,
633 if carried.wraps == 1 { "" } else { "s" });
634 if carried.opened {
635 println!(" a wrap addressed to this replica was opened: this \
636 repository now holds the content key and can read what it takes");
637 }
638 }
639
640 // Held aside because the party below borrows the repository, and because the
641 // question it answers is asked once: is what leaves this machine readable?
642 let veil = repo.veil.clone();
643 let before: BTreeSet<OpId> = repo.log.iter().map(|rec| rec.id()).collect();
644 let was = repo.log.frontier();
645 let here_held = repo.log.len();
646 let mode = Mode::between(
647 repo.log.len(),
648 repo.log.frontier().len(),
649 standing.ops,
650 standing.heads,
651 );
652 let path = proto::sync_path(&remote.account, &remote.name);
653 let mut trips = 0usize;
654 let mut sessions = 0usize;
655 let mut requests = 0usize;
656 let mut unfinished = false;
657 let mut messages = 0usize;
658 let mut pieces = 0usize;
659 let mut up = 0usize; // request bodies, summed
660 let mut down = 0usize; // reply bodies, summed
661 let mut widest = 0usize; // the largest single request body
662 let mut absorbed_there = 0usize;
663 // A rehearsal absorbs into copies and then lets go of them, so the repository
664 // is left exactly as the capture above left it. It has to absorb into
665 // something: what would arrive is only knowable by taking it, verifying it
666 // under the keys this end knows and placing it, which is the same work the
667 // real thing does. The difference is what is kept, and nothing else.
668 let mut spare_log = if dry { repo.log.clone() } else { OpLog::default() };
669 let mut spare_env = if dry { repo.envelopes.clone() } else { BTreeMap::new() };
670 let mut spare_prov = if dry { repo.prov.clone() } else { BTreeMap::new() };
671 let (mut sent, received, fell_back, withheld) = {
672 let (log, envelopes, prov) = if dry {
673 (&mut spare_log, &mut spare_env, &mut spare_prov)
674 } else {
675 (&mut repo.log, &mut repo.envelopes, &mut repo.prov)
676 };
677 let mut here = Party {
678 name: fmt!("{}", repo.root.display()),
679 trust: repo.cfg.trust(),
680 require_signed: repo.cfg.require_signed,
681 log,
682 envelopes,
683 prov,
684 };
685 // One exchange is one session, and a session is not always enough.
686 //
687 // A relay whose reply would exceed `proto::REPLY_BYTES` sends what fits and
688 // says `Done`, which claims only that it will send no more this turn. So
689 // the loop below runs a whole session, asks whether this end now holds
690 // every head the relay opened with, and if it does not, opens another. The
691 // relay remembers nothing between them: the next opening carries the
692 // frontier this end has grown, and the owed set it computes is smaller by
693 // exactly what already arrived.
694 //
695 // It is the same property that makes an interrupted sync safe -- any
696 // append-order prefix of an owed set is causally closed by construction, so
697 // what arrived stands on its own -- applied on purpose rather than after a
698 // failure. Resume is rerun, and the rerun happens inside the one command.
699 let mut absorbed_here = 0usize;
700 // What the relay has SHOWN it holds, carried from one session to the next.
701 //
702 // The one thing a session boundary used to throw away. A session works out
703 // what it owes from the other end's frontier, and a frontier says nothing
704 // about what it handed over: the relay's heads are operations this end has
705 // not reached yet, so they subtract nothing, and a fresh session found its
706 // whole log owed. Measured on the fe2o3 clone of 12026-08-22, sixteen
707 // sessions over 35,314 operations: 166,224 operations offered, every one of
708 // them already at the relay, 717 MB up against 87 MB down.
709 //
710 // Two halves, and neither of them takes the relay's word for anything. The
711 // session writes down what arrived, because an end that hands an operation
712 // over holds it. The half below writes down what left here and was taken,
713 // because whether a request crossed is the carrier's question and not the
714 // session's -- and it is written down AFTER `withholding`, so a
715 // `--pull-only` exchange, which builds a send and never posts it, remembers
716 // nothing.
717 //
718 // It lives for the exchange and no longer. A failure anywhere below drops
719 // it with the command, so nothing survives to be believed on a run that
720 // never proved it.
721 let mut known: BTreeSet<OpId> = BTreeSet::new();
722 // The first session that fell back is what the summary names; a later one
723 // falling back for the same reason adds nothing to say.
724 let mut fell_back_any = None;
725 let mut withheld_all = 0usize;
726 let mut theirs: Vec<OpId> = Vec::new();
727 loop {
728 sessions += 1;
729 if sessions > SESSION_LIMIT {
730 // Not an error: what arrived is placed, durable and causally closed.
731 // The caller is told plainly that the exchange is unfinished, which is
732 // the one thing the summary must never round up to success.
733 unfinished = true;
734 break;
735 }
736 let mut session = Session::knowing(mode, known.clone());
737 // The bound below is on one session's turns, so it starts again with each.
738 trips = 0;
739 let mut out: Vec<Message> = vec![res!(session.open(here.log))];
740 // The relay's own limit, never this build's guess at it, and the smaller of
741 // the two where both are known: a client compiled with a larger idea than
742 // the deployment allows is the failure this whole change began with.
743 let cap = std::cmp::min(proto::POST_BYTES, standing.post);
744 while !out.is_empty() {
745 trips += 1;
746 if trips > TRIP_LIMIT {
747 return Err(err!(
748 "An exchange with {} took {} round trips without converging, which \
749 no sync needs; the sessions are not making progress.",
750 remote.label(), trips;
751 Bug, Excessive));
752 }
753 let mut going = Vec::new();
754 for msg in out.drain(..) {
755 let msg = if taking { withholding(msg) } else { msg };
756 // Read off the plain form, before the seal and the veil put it beyond
757 // reading, and after `withholding` has decided whether it goes at all.
758 // Written down before the post rather than after it because a post that
759 // is not answered well takes the whole command down with it, so nothing
760 // recorded here outlives a send that did not cross.
761 if let Message::Send { entries } = &msg {
762 for entry in entries {
763 known.insert(res!(entry.id()));
764 }
765 }
766 let msg = res!(here.seal_outgoing(msg));
767 // The veil goes on last, over the sealed form, so what the relay is
768 // handed is a header and ciphertext, and the signature is inside the
769 // ciphertext where only a reader holding the key ever meets it.
770 let msg = match &veil {
771 Some(v) => res!(v.veil(msg)),
772 None => msg,
773 };
774 // Cut to the SMALLER of a batch and what one request body may be,
775 // because both bound the same bytes and only the smaller of them
776 // binds. Until 2026-08-22 this was `BATCH_BYTES` alone, so against a
777 // relay publishing one mebibyte a four-mebibyte batch went out as a
778 // single oversized request that `groups` had to let through whole.
779 going.extend(res!(proto::split(msg, std::cmp::min(proto::BATCH_BYTES, cap))));
780 }
781 // An operation larger than one request body leaves in pieces, and a
782 // relay too old for the piece message would refuse them one at a time
783 // with nothing to say about why. Told here instead, before anything is
784 // posted, naming the operation and both versions.
785 if standing.speaks.1 < msg::VERSION {
786 if let Some(Message::Part { id, total, .. }) = going
787 .iter()
788 .find(|m| matches!(m, Message::Part { .. }))
789 {
790 return Err(err!(
791 "The operation {} is larger than the {} bytes {} takes in one \
792 request, so it has to cross in {} pieces, and that relay speaks \
793 ORESYN versions {} to {} where a piece is version {}. Nothing has \
794 been sent. The relay is the end that has to be brought forward; \
795 raising its request limit would carry this operation and not the \
796 next one.",
797 id, cap, remote.label(), total,
798 standing.speaks.0, standing.speaks.1, msg::VERSION;
799 Invalid, Network, Version, Mismatch));
800 }
801 }
802 messages += going.len();
803 pieces += going.iter().filter(|m| matches!(m, Message::Part { .. })).count();
804 // ONE REQUEST PER GROUP, not one request for everything there is to
805 // send. `frame` concatenates whatever it is handed, so this used to
806 // build a single body out of every batch a round had -- 84 MB on the
807 // first repository of any size -- and a reverse proxy in front of the
808 // relay closed the connection part way through writing it. See
809 // `proto::POST_BYTES`.
810 //
811 // The replies are read in order and their messages concatenated. The
812 // session is peer-symmetric and message-driven: it does not care
813 // whether what it is given arrived in one response or in twenty, only
814 // that the order held.
815 let mut arriving = Vec::new();
816 for span in res!(proto::groups(&going, cap)) {
817 requests += 1;
818 let body = res!(proto::frame(&going[span]));
819 up += body.len();
820 widest = widest.max(body.len());
821 let reply = res!(remote.ask(
822 &rt, &signer, HttpMethod::POST, &path, "application/octet-stream", body,
823 ));
824 if let Some(said) = header_of(&reply, HEADER_ABSORBED) {
825 if let Ok(n) = said.trim().parse::<usize>() {
826 absorbed_there += n;
827 }
828 }
829 down += reply.body.len();
830 arriving.extend(res!(proto::unframe(&reply.body)));
831 }
832 messages += arriving.len();
833 pieces += arriving.iter().filter(|m| matches!(m, Message::Part { .. })).count();
834 // An operation the relay could not put in one message arrives in pieces
835 // and is put back together here, before anything below this line sees
836 // it. `upto` keeps a run whole inside one reply, so a run that is
837 // unfinished at the end of a turn is a fault and not a bound being
838 // reached: the relay keeps nothing between requests, so what it did not
839 // send whole it can never continue.
840 let mut rejoined = Vec::with_capacity(arriving.len());
841 let mut held = Parts::new();
842 for msg in arriving {
843 if let Some(whole) = res!(held.absorb(msg)) {
844 rejoined.push(whole);
845 }
846 }
847 if held.pending() {
848 return Err(err!(
849 "{} stopped part way through an operation, having sent {} bytes of \
850 it. A relay sends the pieces of one operation together, because it \
851 keeps nothing between requests and could not continue one it had \
852 begun.", remote.label(), held.held();
853 IO, Network, Missing));
854 }
855 for msg in rejoined {
856 // What the relay says it holds. Both openings carry it, so knowing
857 // whether the exchange finished costs nothing on the wire: this end
858 // compares its own log against these heads once the session closes.
859 match &msg {
860 Message::Hello { heads } => theirs = heads.clone(),
861 Message::Sketch { heads, .. } => theirs = heads.clone(),
862 _ => {},
863 }
864 // The veil comes off first, because everything below this reads the
865 // operation. A repository holding no key for what arrived veiled is
866 // told so by name a line later rather than absorbing it unreadable.
867 let msg = match &veil {
868 Some(v) => res!(v.unveil(msg)),
869 None => msg,
870 };
871 // Verified at this end, under the keys this end knows, before the
872 // session is given the chance to absorb anything.
873 res!(here.check_incoming(&msg, &remote.label()));
874 let turn = res!(session.receive(&mut *here.log, msg));
875 out.extend(turn.send);
876 }
877 }
878 if !session.is_converged() {
879 return Err(err!(
880 "An exchange with {} ended with this end unfinished.", remote.label();
881 Bug, Missing));
882 }
883 absorbed_here += session.ops_absorbed();
884 // What this session watched arrive, added to what the loop above watched
885 // leave, and the next session opens knowing both.
886 known.extend(session.known().iter().copied());
887 fell_back_any = fell_back_any.or(session.fell_back());
888 // What was withheld is what the session counted as told and this end never
889 // put on the wire, so it is the difference between what the relay lacks
890 // and what it was given.
891 // The FIRST session's figure, never a sum of them.
892 //
893 // What this line answers is how much of this end's own work stayed here,
894 // which is a fact about the state the exchange started in. A later session
895 // asks it again of a repository that has meanwhile absorbed thousands of
896 // the relay's operations, and the walk is loose -- a peer that cannot
897 // subtract the other's tip offers its whole log -- so each session re-offers
898 // the growing prefix and a sum of them is a triangular number with no
899 // meaning at all. Measured on a fresh replica that held 44,541 operations
900 // and had held none ninety seconds earlier: 267,486 "kept back".
901 if taking && sessions == 1 {
902 withheld_all = session.ops_sent();
903 }
904 // Does this end now hold every head the relay opened with? A relay that
905 // held nothing back opened with a frontier this end has just finished
906 // absorbing, so the common case leaves after one pass.
907 if theirs.iter().all(|id| here.log.contains(id)) {
908 break;
909 }
910 }
911 (absorbed_there, absorbed_here, fell_back_any, withheld_all)
912 };
913
914 // Everything that arrived reaches the segments before the command says a word
915 // about what it did.
916 // A rehearsal keeps nothing. What arrived was placed in the copies, counted,
917 // and is about to be dropped with them.
918 if !dry {
919 let mine = arrived(&repo.log, &before, &repo.envelopes, &BTreeMap::new());
920 res!(repo.write_entries_from(&mine, None));
921 }
922 // What this repository can afterwards be asked about the exchange. The relay
923 // holds no working copy and asks nothing, so there is one record and not two.
924 if !dry {
925 res!(delta::record(&repo.root, &remote.label(), was.clone(), repo.log.frontier()));
926 }
927
928 // The mark naming where this visit left the two of them, authored here and
929 // pushed inside the same visit.
930 //
931 // A relay authors nothing, so the trick a path sync uses -- one mark written
932 // into both logs -- has to be spelled as a push here. It is one further
933 // request carrying one operation, and it is what closes the fixed point: left
934 // out, the relay is handed a frontier no mark names, the next replica to visit
935 // pulls that and names it itself, and the two replicas end one operation apart
936 // after every visit either of them will ever make.
937 //
938 // The batch closes causally, which is what [`Session::receive`] insists on:
939 // the mark's parents are this end's frontier, and after an exchange that
940 // withheld nothing the relay holds everything this end holds.
941 //
942 // It converges, and the second visit is what shows why. A visits and the relay
943 // ends holding A's operations and A's mark. B visits, receives both, and its
944 // frontier is two heads -- its own and A's mark -- so the mark it authors
945 // covers both and is pushed. A visits again and receives B's mark, which has
946 // A's own as an ancestor, so A's frontier is that single mark; `auto_mark`
947 // answers nothing on a frontier that is already one mark, and A pushes
948 // nothing. A fourth visit has nothing to exchange at all.
949 //
950 // `--pull-only` pushes nothing, this mark included: the whole point of the
951 // option is that the relay is not written to. `--dry-run` writes nothing
952 // anywhere.
953 let mut named: Option<OpId> = None;
954 if !dry && !taking && (sent > 0 || received > 0) {
955 if let Some(id) = res!(verbs::auto_mark(repo)) {
956 let mine = Message::Send { entries: vec![res!(repo.entry_of(&id))] };
957 let mine = match &repo.veil {
958 Some(v) => res!(v.veil(mine)),
959 None => mine,
960 };
961 let cap = std::cmp::min(proto::POST_BYTES, standing.post);
962 let mut going = res!(proto::split(mine, std::cmp::min(proto::BATCH_BYTES, cap)));
963 going.push(Message::Done);
964 messages += going.len();
965 // ONE REQUEST, in every case anybody will meet. This body carries one
966 // entry and a `Done`, and that entry is a mark -- a name, a time and a
967 // signature, measured at 229 bytes plus the name and independent of the
968 // history's size. It fits any limit a relay could sensibly publish.
969 //
970 // It is grouped anyway rather than framed whole, because "a mark is
971 // small" is a fact about the marks this tool writes and not about the
972 // operation code: `Op::Mark` carries an optional body, and a mark large
973 // enough to be cut up would otherwise leave here as a single oversized
974 // body -- exactly the failure the rest of this module exists to end. The
975 // absorbed count is summed over the group for the same reason, and the
976 // check below still insists on exactly ONE, which is what it always
977 // meant.
978 //
979 // (Until 2026-08-20 the reason given here was that `BATCH_BYTES` is under
980 // `POST_BYTES`, so a group would be one request anyway. That stopped being
981 // true when the cap became the relay's to state: `POST_FALLBACK` is 1 MiB
982 // and `BATCH_BYTES` is 4.)
983 let mut took = 0usize;
984 let mut answered = false;
985 for span in res!(proto::groups(&going, cap)) {
986 let body = res!(proto::frame(&going[span]));
987 up += body.len();
988 widest = widest.max(body.len());
989 requests += 1;
990 let got = res!(remote.ask(
991 &rt, &signer, HttpMethod::POST, &path, "application/octet-stream", body,
992 ));
993 down += got.body.len();
994 if let Some(n) = header_of(&got, HEADER_ABSORBED)
995 .and_then(|s| s.trim().parse::<usize>().ok())
996 {
997 took += n;
998 answered = true;
999 }
1000 }
1001 // The relay says how many it absorbed, and it is worth insisting on:
1002 // a mark that did not land leaves the relay unmarked, which is the
1003 // whole thing this request exists to prevent, and it would otherwise
1004 // fail silently and permanently.
1005 match if answered { Some(took) } else { None } {
1006 Some(1) => sent += 1,
1007 other => return Err(err!(
1008 "{} was sent the mark {} naming where this visit left the two of \
1009 them and says it absorbed {}, so the relay is left standing at a \
1010 point no mark names. The operations themselves are written at both \
1011 ends and nothing is lost; running the sync again names it.",
1012 remote.label(), id,
1013 match other {
1014 Some(n) => fmt!("{}", n),
1015 None => fmt!("nothing it can be asked about"),
1016 };
1017 Invalid, Data, Mismatch)),
1018 }
1019 named = Some(id);
1020 }
1021 }
1022
1023 println!();
1024 let opened = match mode {
1025 Mode::Walk => fmt!("frontier walk"),
1026 Mode::Sketch { estimate, .. } => fmt!(
1027 "sketch, sized for {} operation{} of difference over {} held",
1028 estimate, if estimate == 1 { "" } else { "s" }, here_held,
1029 ),
1030 };
1031 println!("mode {}", match fell_back {
1032 None => opened,
1033 Some(why) => fmt!(
1034 "{}, which fell back to the frontier walk: {}", opened, why.why()),
1035 });
1036 if dry {
1037 // Nothing crossed that is being kept, so the counts are said in the
1038 // conditional they belong in.
1039 println!("offered nothing: a rehearsal hands nothing over");
1040 println!("would take {} operation{} this repository does not hold",
1041 received, if received == 1 { "" } else { "s" });
1042 } else if taking {
1043 // Said whatever the counts are, because the whole point of the option is
1044 // that the two ends are not left agreeing and the caller has to know it.
1045 println!("offered nothing, by request");
1046 if withheld > 0 {
1047 println!("kept back {} operation{} the relay does not hold, which stay here \
1048 until a sync offers them", withheld, if withheld == 1 { "" } else { "s" });
1049 }
1050 println!("received {} operation{} this repository did not hold",
1051 received, if received == 1 { "" } else { "s" });
1052 } else if sent == 0 && received == 0 && !unfinished {
1053 println!("nothing to exchange: this repository and the relay already held the \
1054 same {} operation{}", here_held, if here_held == 1 { "" } else { "s" });
1055 } else {
1056 println!("sent {} operation{} the relay did not hold",
1057 sent, if sent == 1 { "" } else { "s" });
1058 println!("received {} operation{} this repository did not hold",
1059 received, if received == 1 { "" } else { "s" });
1060 }
1061 // Requests and sessions, not "round trips". Until now this line counted turns
1062 // of the session and called them round trips, so a push that went out as a
1063 // dozen HTTP requests reported two -- a true sentence about the state machine
1064 // that is false about the wire, which is the reading anybody debugging a proxy
1065 // would take from it.
1066 println!("traffic {} message{} in {} request{} over {} session{}, {} byte{}",
1067 messages, if messages == 1 { "" } else { "s" },
1068 requests, if requests == 1 { "" } else { "s" },
1069 sessions, if sessions == 1 { "" } else { "s" },
1070 up + down, if up + down == 1 { "" } else { "s" });
1071 // The two halves apart, because they are not alike and the difference is the
1072 // thing worth seeing. A loose frontier walk sends its whole log to a peer that
1073 // cannot subtract it, so a clone that takes 87 MB can offer several hundred on
1074 // the way; one number hides that and two name it. The largest body is beside
1075 // them because it is what a proxy in front of the relay refuses, and until now
1076 // the only way to read it was to instrument the proxy.
1077 println!(" {} byte{} up, {} down, largest request body {} against the {} \
1078 this relay publishes",
1079 up, if up == 1 { "" } else { "s" }, down, widest, standing.post);
1080 // Said only when it happened, because it is the answer to a question nobody
1081 // asks until an operation is too large to cross whole. An operation counted
1082 // here crossed in pieces and was put back together byte for byte at the far
1083 // end; nothing about it is stored differently for having been cut up.
1084 if pieces > 0 {
1085 println!("pieces {} of those messages carried an operation too large for one \
1086 request body", pieces);
1087 }
1088 // The exchange stopped short. Said here, in the summary, because a caller who
1089 // reads only the last few lines must not be able to take this for a finished
1090 // sync -- everything above it is true and none of it says the relay still
1091 // holds operations this repository has not seen.
1092 if unfinished {
1093 println!();
1094 println!("UNFINISHED: the relay still holds operations this repository does not.");
1095 println!(" What arrived is written and complete in itself; run `ore \
1096 sync {}` again to continue.", remote.label());
1097 }
1098 println!("this repository {} {} operation{}",
1099 if dry { "holds" } else { "now holds" },
1100 repo.log.len(), if repo.log.len() == 1 { "" } else { "s" });
1101 // Said out loud, because it is the one operation this command wrote that the
1102 // exchange did not ask for, and because it went to the relay as well.
1103 if let Some(id) = named {
1104 println!("named both ends at {}, so the next replica to visit finds a \
1105 point rather than a frontier", id);
1106 }
1107 println!("frontier {}", verbs::frontier_of(&repo.log.frontier()));
1108 let unknown = repo.prov.values().filter(|p| **p == Prov::Unknown).count();
1109 if unknown > 0 {
1110 println!("provenance {} operation{} signed by a key this repository does not \
1111 know, marked ?", unknown, if unknown == 1 { "" } else { "s" });
1112 }
1113
1114 // A rehearsal stops here. Nothing was written, so there is no working copy to
1115 // bring forward, and what would have arrived is described from the copies
1116 // before they are let go of.
1117 if dry {
1118 println!();
1119 println!("nothing was written: this was a rehearsal");
1120 res!(delta::describe(&spare_log, &repo.root, &spare_prov, repo.cfg.replica,
1121 &was, &spare_log.frontier(), "it would bring"));
1122 println!();
1123 println!("`ore sync {}` takes it", remote.label());
1124 return Ok(());
1125 }
1126
1127 // The relay holds no working copy, so there is no marker to leave anywhere:
1128 // this working copy is the only one the exchange touched.
1129 let tree = res!(tree::whole(&repo.log));
1130 let moved = res!(tree::materialise(&repo.root, &tree, tree::Surplus::Remove));
1131 println!();
1132 println!("the working copy here is now the merged state");
1133 verbs::report_moved_files(&moved);
1134 println!();
1135 verbs::summarise(repo, &tree)
1136}
1137
1138
1139#[cfg(test)]
1140mod tests {
1141 use super::*;
1142
1143 /// A URL names a host, a port, an account and a repository, and anything
1144 /// that does not is refused rather than guessed at.
1145 #[test]
1146 fn a_relay_url_is_read_or_it_is_not() -> Outcome<()> {
1147 let got = res!(Remote::parse("https://oregami.example/oxedyne/ore"));
1148 assert!(got.tls);
1149 assert_eq!(got.port, 443);
1150 assert_eq!(got.host, "oregami.example");
1151 assert_eq!(got.account, "oxedyne");
1152 assert_eq!(got.name, "ore");
1153 assert_eq!(got.label(), "https://oregami.example/oxedyne/ore");
1154
1155 let got = res!(Remote::parse("http://127.0.0.1:8420/a/b/"));
1156 assert!(!got.tls);
1157 assert_eq!(got.port, 8420);
1158 assert_eq!(got.label(), "http://127.0.0.1:8420/a/b");
1159
1160 for bad in [
1161 "https://oregami.example/oxedyne",
1162 "https://oregami.example",
1163 "https://oregami.example/a/b/c",
1164 "ftp://oregami.example/a/b",
1165 "/some/path",
1166 "https:///a/b",
1167 ] {
1168 if Remote::parse(bad).is_ok() {
1169 return Err(err!("The URL {:?} was accepted.", bad; Test, Invalid));
1170 }
1171 }
1172 assert!(Remote::is_url("http://x/a/b") && Remote::is_url("https://x/a/b"));
1173 assert!(!Remote::is_url("../elsewhere"));
1174 Ok(())
1175 }
1176}