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 | |
| 42 | use crate::arrived as delta; |
| 43 | use crate::keys::Prov; |
| 44 | use crate::repo::Repo; |
| 45 | use crate::sync::Party; |
| 46 | use crate::tree; |
| 47 | use crate::verbs; |
| 48 | |
| 49 | use ore_relay::proto::{ |
| 50 | self, |
| 51 | Presented, |
| 52 | }; |
| 53 | use ore_relay::serve::HEADER_ABSORBED; |
| 54 | |
| 55 | use ore_store::keys::{ |
| 56 | Binding, |
| 57 | Signing, |
| 58 | }; |
| 59 | use ore_store::store::arrived; |
| 60 | use ore_store::veilkey::{ |
| 61 | VeilBinding, |
| 62 | Wrap, |
| 63 | }; |
| 64 | |
| 65 | use oxedyne_fe2o3_core::prelude::*; |
| 66 | use oxedyne_fe2o3_jdat::prelude::*; |
| 67 | use oxedyne_fe2o3_net::acme::trust::letsencrypt_client_config; |
| 68 | use oxedyne_fe2o3_net::http::client::{ |
| 69 | http_request, |
| 70 | https_request, |
| 71 | }; |
| 72 | use oxedyne_fe2o3_net::http::fields::HeaderName; |
| 73 | use oxedyne_fe2o3_net::http::header::HttpMethod; |
| 74 | use oxedyne_fe2o3_net::http::msg::HttpMessage; |
| 75 | use oxedyne_fe2o3_ore::id::OpId; |
| 76 | use oxedyne_fe2o3_ore::log::OpLog; |
| 77 | use oxedyne_fe2o3_ore::sync::{ |
| 78 | msg, |
| 79 | Message, |
| 80 | Mode, |
| 81 | Parts, |
| 82 | Session, |
| 83 | }; |
| 84 | |
| 85 | use std::collections::{ |
| 86 | BTreeMap, |
| 87 | BTreeSet, |
| 88 | }; |
| 89 | |
| 90 | use 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. |
| 102 | pub 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. |
| 118 | pub const SESSION_LIMIT: usize = 64; |
| 119 | |
| 120 | |
| 121 | /// A repository on a relay, as a URL named it. |
| 122 | #[derive(Clone, Debug)] |
| 123 | pub 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 | |
| 136 | impl 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. |
| 294 | fn 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. |
| 307 | fn 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. |
| 313 | struct 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. |
| 331 | fn 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. |
| 435 | struct 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. |
| 460 | fn 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. |
| 575 | fn 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. |
| 588 | pub 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)] |
| 1140 | mod 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 | } |