oxedyne/ore/relay/src/serve.rs
59.5 KiB, 124 runs
created by r2848102244:151, 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 | //! Answering a request, and listening for one. |
| 2 | //! |
| 3 | //! [`respond`] is the whole of the relay and touches no socket: a method, a |
| 4 | //! path, a credential and a body go in, and a status and some bytes come out. |
| 5 | //! That is deliberate. The transport underneath is a detail -- a plain listener |
| 6 | //! here, a vhost's API handler on a server that already terminates TLS -- and a |
| 7 | //! relay whose behaviour lives inside its listener could not be mounted twice. |
| 8 | //! |
| 9 | //! # The exchange, in turns rather than requests |
| 10 | //! |
| 11 | //! The client posts its opening. The relay runs a `Session` over its own log and |
| 12 | //! answers with its own opening, as much of what the client is owed as fits under |
| 13 | //! [`Host::reply_bytes`], and `Done` -- which claims only that this end will send |
| 14 | //! no more this turn. The client absorbs, and posts what the relay is owed and |
| 15 | //! `Done`; the relay absorbs, and answers with nothing. |
| 16 | //! |
| 17 | //! This heading said "in two round trips" until 2026-08-20, and for a history |
| 18 | //! small enough to cross whole it still is two. Larger, it is not, in both |
| 19 | //! directions at once: a client whose log does not then cover the frontier the |
| 20 | //! relay opened with opens a fresh session, and one turn may leave the client as |
| 21 | //! several posts, since a request is bounded as well (`proto::POST_BYTES`). |
| 22 | //! Measured on a 58 MB clone: fourteen sessions, twenty-eight requests. |
| 23 | //! |
| 24 | //! No transfer state is kept on either side, which is what makes resume a rerun: |
| 25 | //! an exchange interrupted anywhere leaves both logs valid, everything absorbed is |
| 26 | //! causally closed and durable, and the next attempt opens a fresh session over |
| 27 | //! frontiers that already reflect what crossed. A session that stopped short is |
| 28 | //! that same case reached on purpose rather than by failure. |
| 29 | //! |
| 30 | //! Each request is a session of its own at this end, because the relay keeps |
| 31 | //! nothing between them. That is sound rather than a shortcut: the client's |
| 32 | //! second message is closed against the frontier the relay stated in the first, |
| 33 | //! and the relay's log only grows, so a closure that held then holds now. |
| 34 | //! |
| 35 | //! # One request at a time |
| 36 | //! |
| 37 | //! The listener answers connections in turn, and a request that writes holds the |
| 38 | //! repository's lock while it does. A relay for one team wants no more, and the |
| 39 | //! alternative -- concurrent writers into one segment -- is the one thing the |
| 40 | //! lock exists to prevent. |
| 41 | //! |
| 42 | //! # A repository it cannot read |
| 43 | //! |
| 44 | //! A veiled repository arrives as entries whose identifiers and parents are in |
| 45 | //! clear and whose bodies are encrypted under a key no relay holds. Nothing in |
| 46 | //! the exchange changes: the same session runs, over a log whose records are the |
| 47 | //! stand-ins of `ore_store::veil`, and every question it asks of that log is |
| 48 | //! answered from a header. What arrived veiled is set aside on the way in, put |
| 49 | //! back on the way out, and written to the segments as it stands, so the relay |
| 50 | //! stores and serves ciphertext and converges two replicas that never meet |
| 51 | //! without being able to say what either of them wrote. |
| 52 | |
| 53 | use crate::acl::Role; |
| 54 | use crate::host::{ |
| 55 | Host, |
| 56 | Hosted, |
| 57 | }; |
| 58 | use crate::proto::{ |
| 59 | self, |
| 60 | Presented, |
| 61 | }; |
| 62 | use crate::reading::{ |
| 63 | Carried, |
| 64 | Took, |
| 65 | Whence, |
| 66 | }; |
| 67 | |
| 68 | use ore_store::keys::Binding; |
| 69 | use ore_store::veilkey::{ |
| 70 | VeilBinding, |
| 71 | Wrap, |
| 72 | }; |
| 73 | use ore_store::store::{ |
| 74 | arrived, |
| 75 | Consumed, |
| 76 | Keep, |
| 77 | Replayed, |
| 78 | seal_one, |
| 79 | Verify, |
| 80 | }; |
| 81 | use ore_store::veil::{ |
| 82 | placehold, |
| 83 | restore_one, |
| 84 | }; |
| 85 | |
| 86 | use oxedyne_fe2o3_core::prelude::*; |
| 87 | use oxedyne_fe2o3_jdat::prelude::*; |
| 88 | use oxedyne_fe2o3_ore::id::OpId; |
| 89 | use oxedyne_fe2o3_ore::log::OpLog; |
| 90 | use oxedyne_fe2o3_ore::segment::{ |
| 91 | self, |
| 92 | Entry, |
| 93 | Veiled, |
| 94 | }; |
| 95 | use oxedyne_fe2o3_ore::sync::{ |
| 96 | msg, |
| 97 | Message, |
| 98 | Mode, |
| 99 | Session, |
| 100 | }; |
| 101 | |
| 102 | use std::collections::{ |
| 103 | BTreeMap, |
| 104 | BTreeSet, |
| 105 | }; |
| 106 | use std::net::SocketAddr; |
| 107 | use std::pin::Pin; |
| 108 | use std::sync::Arc; |
| 109 | |
| 110 | use oxedyne_fe2o3_net::constant; |
| 111 | use oxedyne_fe2o3_net::http::{ |
| 112 | fields::{ |
| 113 | HeaderFieldCategory, |
| 114 | HeaderFieldValue, |
| 115 | HeaderName, |
| 116 | }, |
| 117 | header::{ |
| 118 | HttpHeadline, |
| 119 | HttpMethod, |
| 120 | }, |
| 121 | msg::{ |
| 122 | HttpMessage, |
| 123 | ReadLimits, |
| 124 | }, |
| 125 | status::HttpStatus, |
| 126 | }; |
| 127 | |
| 128 | use tokio::net::{ |
| 129 | TcpListener, |
| 130 | TcpStream, |
| 131 | }; |
| 132 | |
| 133 | |
| 134 | /// The content type frames are carried under. |
| 135 | pub const FRAMES_TYPE: &str = "application/octet-stream"; |
| 136 | /// The content type a daticle answer is carried under. |
| 137 | pub const TEXT_TYPE: &str = "text/plain; charset=utf-8"; |
| 138 | |
| 139 | /// The header a sync answer counts its absorption in. |
| 140 | /// |
| 141 | /// A client cannot work out what the relay took from what it offered: the walk |
| 142 | /// is loose, so a peer sends more than it needs to and the receiver drops what it |
| 143 | /// holds. The count is a courtesy for the line the command prints, and nothing |
| 144 | /// rests on it. |
| 145 | pub const HEADER_ABSORBED: &str = "x-ore-absorbed"; |
| 146 | |
| 147 | |
| 148 | /// What a request asked for. |
| 149 | pub struct Request { |
| 150 | /// The method, as the transport spelled it. |
| 151 | pub method: String, |
| 152 | /// The path, without any query string. |
| 153 | pub path: String, |
| 154 | /// What the request said about who is making it, unchecked. |
| 155 | pub cred: Option<Presented>, |
| 156 | /// The body. |
| 157 | pub body: Vec<u8>, |
| 158 | } |
| 159 | |
| 160 | |
| 161 | /// What the relay answers. |
| 162 | pub struct Reply { |
| 163 | /// The status code. |
| 164 | pub status: u16, |
| 165 | /// The content type. |
| 166 | pub kind: String, |
| 167 | /// The body. |
| 168 | pub body: Vec<u8>, |
| 169 | /// How many operations this answer absorbed, where it absorbed any. |
| 170 | pub absorbed: Option<usize>, |
| 171 | } |
| 172 | |
| 173 | impl Reply { |
| 174 | |
| 175 | /// A reply of words, which is what every refusal is. |
| 176 | pub fn text(status: u16, said: String) -> Self { |
| 177 | Self { |
| 178 | status, |
| 179 | kind: fmt!("{}", TEXT_TYPE), |
| 180 | body: said.into_bytes(), |
| 181 | absorbed: None, |
| 182 | } |
| 183 | } |
| 184 | |
| 185 | /// A reply of frames. |
| 186 | pub fn frames(body: Vec<u8>, absorbed: usize) -> Self { |
| 187 | Self { |
| 188 | status: 200, |
| 189 | kind: fmt!("{}", FRAMES_TYPE), |
| 190 | body, |
| 191 | absorbed: Some(absorbed), |
| 192 | } |
| 193 | } |
| 194 | |
| 195 | /// A reply of a daticle, written as text. |
| 196 | pub fn dat(dat: &Dat) |
| 197 | -> Outcome<Self> |
| 198 | { |
| 199 | Ok(Self::text(200, fmt!("{}\n", res!(dat.jdat_to_lines(" "))))) |
| 200 | } |
| 201 | } |
| 202 | |
| 203 | |
| 204 | /// Answers one request. |
| 205 | /// |
| 206 | /// Nothing here fails: a request that cannot be served becomes a status and a |
| 207 | /// sentence, because a caller on the far side of a socket has no error type to |
| 208 | /// receive. What the sentence says is what a person reading a failed `ore sync` |
| 209 | /// will see. |
| 210 | pub fn respond(host: &Host, req: &Request) -> Reply { |
| 211 | match answer(host, req) { |
| 212 | Ok(reply) => reply, |
| 213 | Err(e) => Reply::text(500, fmt!( |
| 214 | "The relay could not answer {} {}: {}", req.method, req.path, e.plain(), |
| 215 | )), |
| 216 | } |
| 217 | } |
| 218 | |
| 219 | /// Answers one request, or fails in a way [`respond`] turns into a status. |
| 220 | fn answer(host: &Host, req: &Request) |
| 221 | -> Outcome<Reply> |
| 222 | { |
| 223 | // The version endpoint answers before anything is authenticated, so an old |
| 224 | // client fails with a sentence naming both sides rather than a decode error |
| 225 | // part way through an exchange. |
| 226 | if req.path == proto::PREFIX { |
| 227 | return versions(host); |
| 228 | } |
| 229 | // What this relay said it would take, refused where it was not taken. A proxy |
| 230 | // that caps a body closes the connection part way through it, so the client |
| 231 | // sees a broken pipe and reads it as the relay being down; a relay that |
| 232 | // enforces its own published number answers instead, and says both. |
| 233 | if req.body.len() > host.post_bytes { |
| 234 | return Ok(Reply::text(413, fmt!( |
| 235 | "A request body of {} bytes reaches this relay, and it takes {}. An \ |
| 236 | operation larger than that crosses in pieces -- {} says which ORESYN \ |
| 237 | versions this relay speaks, and a client that speaks version {} or above \ |
| 238 | cuts one up rather than sending it whole.", |
| 239 | req.body.len(), host.post_bytes, proto::PREFIX, msg::VERSION, |
| 240 | ))); |
| 241 | } |
| 242 | let (account, name, verb) = match route(&req.path) { |
| 243 | Some(parts) => parts, |
| 244 | None => return Ok(Reply::text(404, fmt!( |
| 245 | "{:?} names nothing this relay serves. A repository is at {}/{}/<account>/\ |
| 246 | <name>/sync, and {} says which versions this relay speaks.", |
| 247 | req.path, proto::PREFIX, proto::VERSION, proto::PREFIX, |
| 248 | ))), |
| 249 | }; |
| 250 | // The signature is checked before a repository is touched, and says nothing |
| 251 | // about whether this key may do what it is asking. |
| 252 | let mut key: Option<Vec<u8>> = None; |
| 253 | if let Some(cred) = &req.cred { |
| 254 | let now = res!(proto::now()); |
| 255 | match cred.check(&req.method, &req.path, &req.body, now) { |
| 256 | Ok(()) => key = Some(cred.public.clone()), |
| 257 | Err(e) => return Ok(Reply::text(401, fmt!("{}", e.plain()))), |
| 258 | } |
| 259 | } |
| 260 | let hosted = match res!(host.open(&account, &name)) { |
| 261 | Some(h) => h, |
| 262 | None => return Ok(Reply::text(404, fmt!( |
| 263 | "This relay does not hold {}/{}. A repository is created by an \ |
| 264 | administrator with `ore-relay create`, not by pushing to it.", |
| 265 | account, name, |
| 266 | ))), |
| 267 | }; |
| 268 | match (req.method.as_str(), verb.as_str()) { |
| 269 | ("POST", "sync") => sync(host, &hosted, req, key.as_deref()), |
| 270 | ("GET", "keys") => keys(host, &hosted, req, key.as_deref(), false), |
| 271 | ("POST", "keys") => keys(host, &hosted, req, key.as_deref(), true), |
| 272 | ("GET", "wraps") => wraps(&hosted, req, key.as_deref(), false), |
| 273 | ("POST", "wraps") => wraps(&hosted, req, key.as_deref(), true), |
| 274 | (method, verb) => Ok(Reply::text(405, fmt!( |
| 275 | "{} {} is not something this relay does. A sync is posted to \ |
| 276 | .../{}/sync, bindings are read from and deposited at .../{}/keys, and \ |
| 277 | veil key bindings and wraps at .../{}.", |
| 278 | method, verb, "sync", "keys", "wraps", |
| 279 | ))), |
| 280 | } |
| 281 | } |
| 282 | |
| 283 | /// Splits a path into the account, the repository and what is being asked of it. |
| 284 | /// |
| 285 | /// `None` for anything that is not this transport's shape, which includes a path |
| 286 | /// of another version: a relay that guessed at one would be answering a client |
| 287 | /// it does not understand. |
| 288 | fn route(path: &str) -> Option<(String, String, String)> { |
| 289 | let want = fmt!("{}/{}/", proto::PREFIX, proto::VERSION); |
| 290 | let rest = match path.strip_prefix(&want) { |
| 291 | Some(r) => r, |
| 292 | None => return None, |
| 293 | }; |
| 294 | let parts: Vec<&str> = rest.split('/').collect(); |
| 295 | if parts.len() != 3 || parts.iter().any(|p| p.is_empty()) { |
| 296 | return None; |
| 297 | } |
| 298 | Some((fmt!("{}", parts[0]), fmt!("{}", parts[1]), fmt!("{}", parts[2]))) |
| 299 | } |
| 300 | |
| 301 | /// Answers with the versions this relay speaks, of the transport and of the two |
| 302 | /// formats it carries. |
| 303 | fn versions(host: &Host) |
| 304 | -> Outcome<Reply> |
| 305 | { |
| 306 | let mut map = DaticleMap::new(); |
| 307 | map.insert(Dat::Str(fmt!("relay")), Dat::Str(fmt!("ore"))); |
| 308 | map.insert(Dat::Str(fmt!("transport")), Dat::List(vec![Dat::Str(fmt!("{}", proto::VERSION))])); |
| 309 | // Both ends of the range, exactly as `oreseg` carries both. A relay that |
| 310 | // speaks up to version 2 can take an operation larger than one request body, |
| 311 | // because that is the version the piece message arrived in; one that speaks |
| 312 | // only version 1 cannot, and a client reads that here rather than finding out |
| 313 | // when its push is refused. |
| 314 | map.insert(Dat::Str(fmt!("oresyn")), Dat::List(vec![ |
| 315 | Dat::U8(msg::VERSION_MIN), |
| 316 | Dat::U8(msg::VERSION), |
| 317 | ])); |
| 318 | map.insert(Dat::Str(fmt!("oreseg")), Dat::List(vec![ |
| 319 | Dat::U8(segment::VERSION_MIN), |
| 320 | Dat::U8(segment::VERSION), |
| 321 | ])); |
| 322 | // The entry forms, which the segment version does not speak for: a kind is a |
| 323 | // separate axis, so a relay that can carry a veiled repository says so here |
| 324 | // rather than leaving a client to find out at the first push. |
| 325 | map.insert(Dat::Str(fmt!("kinds")), Dat::List(vec![ |
| 326 | Dat::U8(segment::KIND_BARE), |
| 327 | Dat::U8(segment::KIND_SEALED), |
| 328 | Dat::U8(segment::KIND_VEILED), |
| 329 | ])); |
| 330 | // What this relay will take in one request body, so that a client sizes to |
| 331 | // the serving end rather than to a number compiled into its own build. It is |
| 332 | // It is the operator's number since 2026-08-22 -- `--post-bytes` on |
| 333 | // `ore-relay serve` -- and this relay refuses a body over it rather than |
| 334 | // letting a proxy close the connection and leave the client guessing. |
| 335 | map.insert(Dat::Str(fmt!("post")), Dat::U64(host.post_bytes as u64)); |
| 336 | Reply::dat(&Dat::Map(map)) |
| 337 | } |
| 338 | |
| 339 | /// Reads the bindings a repository carries, and deposits any that arrived. |
| 340 | fn keys(host: &Host, hosted: &Hosted, req: &Request, key: Option<&[u8]>, deposit: bool) |
| 341 | -> Outcome<Reply> |
| 342 | { |
| 343 | let want = if deposit { Role::Push } else { Role::Pull }; |
| 344 | if !hosted.acl.may(key, want) { |
| 345 | return Ok(refused(hosted, want)); |
| 346 | } |
| 347 | if deposit && !req.body.is_empty() { |
| 348 | let text = match String::from_utf8(req.body.clone()) { |
| 349 | Ok(t) => t, |
| 350 | Err(e) => return Ok(Reply::text(400, fmt!( |
| 351 | "A deposit of bindings is JDAT text, and this body is not: {}", e, |
| 352 | ))), |
| 353 | }; |
| 354 | let dat = match Dat::decode_string(text) { |
| 355 | Ok(d) => d, |
| 356 | Err(e) => return Ok(Reply::text(400, fmt!( |
| 357 | "A deposit of bindings is not readable JDAT: {}", e.plain(), |
| 358 | ))), |
| 359 | }; |
| 360 | let listed = match &dat { |
| 361 | Dat::List(l) => l, |
| 362 | other => return Ok(Reply::text(400, fmt!( |
| 363 | "A deposit of bindings expects a list, got {:?}.", other, |
| 364 | ))), |
| 365 | }; |
| 366 | let mut offered = Vec::new(); |
| 367 | for item in listed { |
| 368 | match Binding::from_dat(item) { |
| 369 | Ok(b) => offered.push(b), |
| 370 | Err(e) => return Ok(Reply::text(400, fmt!( |
| 371 | "A deposited binding could not be read: {}", e.plain(), |
| 372 | ))), |
| 373 | } |
| 374 | } |
| 375 | let _lock = match hosted.store.lock() { |
| 376 | Ok(l) => l, |
| 377 | Err(e) => return Ok(Reply::text(503, fmt!("{}", e.plain()))), |
| 378 | }; |
| 379 | res!(hosted.learn(&offered)); |
| 380 | } |
| 381 | // The size of the log rides along so that a client can size a sketch without |
| 382 | // spending a round trip asking. |
| 383 | // |
| 384 | // Only the two counts are wanted, and this used to say so by reading under |
| 385 | // `Keep::Nothing` -- which meant a reading of its own, and a whole one, on the |
| 386 | // one endpoint every `ore sync` calls before it syncs. Measured on the fe2o3 |
| 387 | // copy: 89.1 MB of the 90.6 MB a warm clone cost the relay was this line. It |
| 388 | // takes the reading `sync` holds instead, which costs a caller asking only for |
| 389 | // bindings the envelopes it will not read, and costs the sync that follows |
| 390 | // nothing at all. |
| 391 | let replayed = res!(hold(host, hosted)); |
| 392 | let held: Vec<Dat> = res!(hosted.bindings()).iter().map(|b| b.to_dat()).collect(); |
| 393 | let mut map = DaticleMap::new(); |
| 394 | map.insert(Dat::Str(fmt!("keys")), Dat::List(held)); |
| 395 | map.insert(Dat::Str(fmt!("ops")), Dat::U64(replayed.log.len() as u64)); |
| 396 | map.insert(Dat::Str(fmt!("heads")), Dat::U64(replayed.log.frontier().len() as u64)); |
| 397 | // Carried here as well as on `GET /ore` for the reason the counts above are: |
| 398 | // a client that is about to push needs it before it pushes, and asking |
| 399 | // separately would be a round trip spent on a number. The message versions |
| 400 | // ride along for the same reason and are wanted by the same caller: what a |
| 401 | // client does with an operation larger than `post` depends on whether this |
| 402 | // relay can take it in pieces. |
| 403 | map.insert(Dat::Str(fmt!("post")), Dat::U64(host.post_bytes as u64)); |
| 404 | map.insert(Dat::Str(fmt!("oresyn")), Dat::List(vec![ |
| 405 | Dat::U8(msg::VERSION_MIN), |
| 406 | Dat::U8(msg::VERSION), |
| 407 | ])); |
| 408 | Reply::dat(&Dat::Map(map)) |
| 409 | } |
| 410 | |
| 411 | /// Reads the veil key bindings and wraps a repository carries, and deposits any |
| 412 | /// that arrived. |
| 413 | /// |
| 414 | /// **A wrap is served to anybody who may pull, and that is not a leak.** It is |
| 415 | /// one repository's content key encrypted to one replica's veil key, and the |
| 416 | /// secret that opens it never leaves the machine that minted it, so the relay |
| 417 | /// carrying it learns nothing and neither does anybody else who fetches it. |
| 418 | /// Handing them out is how a second replica comes to read a veiled repository at |
| 419 | /// all; a reviewer who narrows this to the addressee has removed the mechanism, |
| 420 | /// because the addressee is exactly who the relay cannot identify. |
| 421 | fn wraps(hosted: &Hosted, req: &Request, key: Option<&[u8]>, deposit: bool) |
| 422 | -> Outcome<Reply> |
| 423 | { |
| 424 | let want = if deposit { Role::Push } else { Role::Pull }; |
| 425 | if !hosted.acl.may(key, want) { |
| 426 | return Ok(refused(hosted, want)); |
| 427 | } |
| 428 | let mut took = 0usize; |
| 429 | let mut kept = 0usize; |
| 430 | if deposit && !req.body.is_empty() { |
| 431 | let text = match String::from_utf8(req.body.clone()) { |
| 432 | Ok(t) => t, |
| 433 | Err(e) => return Ok(Reply::text(400, fmt!( |
| 434 | "A deposit of veil keys and wraps is JDAT text, and this body is not: {}", e, |
| 435 | ))), |
| 436 | }; |
| 437 | let dat = match Dat::decode_string(text) { |
| 438 | Ok(d) => d, |
| 439 | Err(e) => return Ok(Reply::text(400, fmt!( |
| 440 | "A deposit of veil keys and wraps is not readable JDAT: {}", e.plain(), |
| 441 | ))), |
| 442 | }; |
| 443 | let map = match &dat { |
| 444 | Dat::Map(m) => m, |
| 445 | other => return Ok(Reply::text(400, fmt!( |
| 446 | "A deposit of veil keys and wraps expects a map of \"veils\" and \ |
| 447 | \"wraps\", got {:?}.", other, |
| 448 | ))), |
| 449 | }; |
| 450 | let listed = |name: &str| -> Result<Vec<Dat>, String> { |
| 451 | match map.get(&Dat::Str(fmt!("{}", name))) { |
| 452 | Some(Dat::List(l)) => Ok(l.clone()), |
| 453 | None => Ok(Vec::new()), |
| 454 | Some(other) => Err(fmt!( |
| 455 | "A deposit's {:?} expects a list, got {:?}.", name, other)), |
| 456 | } |
| 457 | }; |
| 458 | let veils = match listed("veils") { |
| 459 | Ok(l) => l, |
| 460 | Err(said) => return Ok(Reply::text(400, said)), |
| 461 | }; |
| 462 | let carried = match listed("wraps") { |
| 463 | Ok(l) => l, |
| 464 | Err(said) => return Ok(Reply::text(400, said)), |
| 465 | }; |
| 466 | let mut offered = Vec::new(); |
| 467 | for item in &veils { |
| 468 | match VeilBinding::from_dat(item) { |
| 469 | Ok(b) => offered.push(b), |
| 470 | Err(e) => return Ok(Reply::text(400, fmt!( |
| 471 | "A deposited veil key binding could not be read: {}", e.plain(), |
| 472 | ))), |
| 473 | } |
| 474 | } |
| 475 | let mut arriving = Vec::new(); |
| 476 | for item in &carried { |
| 477 | match Wrap::from_dat(item) { |
| 478 | Ok(w) => arriving.push(w), |
| 479 | Err(e) => return Ok(Reply::text(400, fmt!( |
| 480 | "A deposited wrap could not be read: {}", e.plain(), |
| 481 | ))), |
| 482 | } |
| 483 | } |
| 484 | let _lock = match hosted.store.lock() { |
| 485 | Ok(l) => l, |
| 486 | Err(e) => return Ok(Reply::text(503, fmt!("{}", e.plain()))), |
| 487 | }; |
| 488 | took = res!(hosted.learn_veils(&offered)); |
| 489 | kept = res!(hosted.keep_wraps(&arriving)); |
| 490 | } |
| 491 | let held: Vec<Dat> = res!(hosted.veil_bindings()).iter().map(|b| b.to_dat()).collect(); |
| 492 | let carried: Vec<Dat> = res!(hosted.wraps()).iter().map(|w| w.to_dat()).collect(); |
| 493 | let mut map = DaticleMap::new(); |
| 494 | map.insert(Dat::Str(fmt!("veils")), Dat::List(held)); |
| 495 | map.insert(Dat::Str(fmt!("wraps")), Dat::List(carried)); |
| 496 | map.insert(Dat::Str(fmt!("learnt")), Dat::U64(took as u64)); |
| 497 | map.insert(Dat::Str(fmt!("kept")), Dat::U64(kept as u64)); |
| 498 | Reply::dat(&Dat::Map(map)) |
| 499 | } |
| 500 | |
| 501 | /// Holds a reading of the hosted log and hands it over. |
| 502 | /// |
| 503 | /// The reading is taken out of the cache, extended over whatever has been |
| 504 | /// appended since it was last read, and put straight back. This is the whole of |
| 505 | /// what a request that does not write needs. |
| 506 | /// |
| 507 | /// A request that *does* write cannot use it: it has to put the reading back |
| 508 | /// itself, over a cursor naming the bytes it wrote, and it has to hold the |
| 509 | /// repository's lock across the two. See [`sync`], which does that by hand. |
| 510 | fn hold(host: &Host, hosted: &Hosted) |
| 511 | -> Outcome<Arc<Replayed>> |
| 512 | { |
| 513 | let carried = res!(host.readings().take(&hosted.dir)); |
| 514 | let (read, cursor, took) = match read_on(hosted, carried) { |
| 515 | Ok(got) => got, |
| 516 | Err(e) => { |
| 517 | // The reading was taken out above and the read that was to replace it |
| 518 | // failed, so this relay is holding none. Said here, so that the next |
| 519 | // request against this repository is a first read and not a lost one. |
| 520 | res!(host.readings().unsee(&hosted.dir)); |
| 521 | return Err(e); |
| 522 | }, |
| 523 | }; |
| 524 | res!(host.readings().keep( |
| 525 | &hosted.dir, res!(hosted.store.bytes()), cursor, took, Arc::clone(&read))); |
| 526 | Ok(read) |
| 527 | } |
| 528 | |
| 529 | /// Runs one turn of a session against the hosted log. |
| 530 | fn sync(host: &Host, hosted: &Hosted, req: &Request, key: Option<&[u8]>) |
| 531 | -> Outcome<Reply> |
| 532 | { |
| 533 | let cap = host.reply_bytes; |
| 534 | let arriving = match proto::unframe(&req.body) { |
| 535 | Ok(m) => m, |
| 536 | Err(e) => return Ok(Reply::text(400, fmt!("{}", e.plain()))), |
| 537 | }; |
| 538 | // An operation larger than one request body arrives in pieces across several |
| 539 | // requests, so the pieces are put back together before anything else looks at |
| 540 | // them. Everything below this line sees what it would have seen had the |
| 541 | // operation been able to cross whole. |
| 542 | // |
| 543 | // The replica is the one the request signed itself with, which is what makes |
| 544 | // two pushes at once two runs and not one interleaved mess. A caller with no |
| 545 | // credential shares one buffer, which is all an unsigned push can expect. |
| 546 | let replica = match &req.cred { |
| 547 | Some(cred) => cred.replica, |
| 548 | None => 0, |
| 549 | }; |
| 550 | let arriving = match host.rejoin(&hosted.label, replica, arriving) { |
| 551 | Ok(m) => m, |
| 552 | Err(e) => return Ok(Reply::text(409, fmt!("{}", e.plain()))), |
| 553 | }; |
| 554 | // Whatever arrived veiled is set aside here and a stand-in put in its place, |
| 555 | // so that everything below this line is working with headers whether or not |
| 556 | // there was ever an operation to read. What is set aside is what goes back on |
| 557 | // the wire and to the disk; see `ore_store::veil`. |
| 558 | let mut veils: BTreeMap<OpId, Veiled> = BTreeMap::new(); |
| 559 | let mut msgs = Vec::with_capacity(arriving.len()); |
| 560 | for msg in arriving { |
| 561 | msgs.push(res!(placehold(msg, &mut veils))); |
| 562 | } |
| 563 | // Reading is the least any request can be, and it is asked for first so that |
| 564 | // nothing below this touches the repository on behalf of a caller the list |
| 565 | // answers with nothing. |
| 566 | if !hosted.acl.may(key, Role::Pull) { |
| 567 | return Ok(refused(hosted, Role::Pull)); |
| 568 | } |
| 569 | let offers = msgs.iter().any(|m| !m.entries().is_empty()); |
| 570 | let may_push = hosted.acl.may(key, Role::Push); |
| 571 | let _lock = match hosted.store.lock() { |
| 572 | Ok(l) => l, |
| 573 | Err(e) => return Ok(Reply::text(503, fmt!("{}", e.plain()))), |
| 574 | }; |
| 575 | // The reading this relay is holding of this repository, taken OUT of the cache |
| 576 | // and not borrowed: what happens to it below is that a session absorbs into it, |
| 577 | // and a reading left in the cache while that happened would be one this relay |
| 578 | // was handing to somebody else mid-absorption. Every way out of this function |
| 579 | // either puts it back with what it learned or says that it did not. |
| 580 | let carried = res!(host.readings().take(&hosted.dir)); |
| 581 | let (mut read, mut cursor, took) = match read_on(hosted, carried) { |
| 582 | Ok(got) => got, |
| 583 | Err(e) => { |
| 584 | // The reading was taken out above and the read that was to replace it |
| 585 | // failed, so this relay is holding none. Said here, so that the next |
| 586 | // request against this repository is a first read and not a lost one. |
| 587 | res!(host.readings().unsee(&hosted.dir)); |
| 588 | return Err(e); |
| 589 | }, |
| 590 | }; |
| 591 | let before: BTreeSet<OpId> = read.log.iter().map(|rec| rec.id()).collect(); |
| 592 | // Is this request a write? Not for what it carries, but for what the relay |
| 593 | // lacks of it. The walk is loose: a peer that cannot subtract the other's tip |
| 594 | // offers its whole log, so a replica that cloned once offers the clone straight |
| 595 | // back the moment the author has moved on. Read off the offer alone, a `pull` |
| 596 | // grant is good for one command and every one after it is refused for pushing, |
| 597 | // which is not what the caller was doing. |
| 598 | // |
| 599 | // It is compared against the log this request holds, and holding that log is |
| 600 | // what serving any read costs, so the answer is worked out for a caller that is |
| 601 | // entitled to make the relay do the work regardless. |
| 602 | let mut refuse = false; |
| 603 | if offers && !may_push { |
| 604 | 'walk: for msg in &msgs { |
| 605 | for entry in msg.entries() { |
| 606 | if !before.contains(&res!(entry.id())) { |
| 607 | refuse = true; |
| 608 | break 'walk; |
| 609 | } |
| 610 | } |
| 611 | } |
| 612 | } |
| 613 | if refuse { |
| 614 | // Nothing has been absorbed and nothing merged, so what is put back is |
| 615 | // exactly what the read above produced. |
| 616 | res!(host.readings().keep( |
| 617 | &hosted.dir, res!(hosted.store.bytes()), cursor, took, read)); |
| 618 | return Ok(refused(hosted, Role::Push)); |
| 619 | } |
| 620 | let mut session = Session::new(mode_for(&read.log, &msgs)); |
| 621 | let mut out: Vec<Message> = Vec::new(); |
| 622 | // Where the session stopped, and whether it had absorbed anything by then. The |
| 623 | // pair is carried out of the block below rather than acted on inside it, |
| 624 | // because putting the reading back needs the whole `Arc` and the block is |
| 625 | // holding the only mutable borrow of it. |
| 626 | let mut stopped: Option<String> = None; |
| 627 | let absorbed; |
| 628 | { |
| 629 | // Copied only if somebody is still reading the last one. On a relay that |
| 630 | // answers one request at a time nobody is: the cache let go of what it held |
| 631 | // a few lines above, and this request is the only holder. |
| 632 | let base = Arc::make_mut(&mut read); |
| 633 | // What arrives keeps the form it arrived in, so the seal on it survives the |
| 634 | // hop. An operation the log already holds keeps the form already on disk, |
| 635 | // since upgrading it would mean rewriting a segment. |
| 636 | for msg in &msgs { |
| 637 | for entry in msg.entries() { |
| 638 | if let Entry::Sealed(env) = entry { |
| 639 | let id = res!(entry.id()); |
| 640 | if !before.contains(&id) { |
| 641 | base.envelopes.insert(id, env.clone()); |
| 642 | } |
| 643 | } |
| 644 | } |
| 645 | } |
| 646 | // The veils already on disk win over the ones that arrived a moment ago, for |
| 647 | // the reason the envelopes do, which is why an arriving one is only put where |
| 648 | // the reading holds none. |
| 649 | for (id, veiled) in veils { |
| 650 | base.veils.entry(id).or_insert(veiled); |
| 651 | } |
| 652 | let held = base.log.len(); |
| 653 | for msg in msgs { |
| 654 | match session.receive(&mut base.log, msg) { |
| 655 | Ok(turn) => out.extend(turn.send), |
| 656 | // A batch with a hole in it, or one whose frames will not read: the |
| 657 | // client is told which operation and why, and nothing is absorbed. |
| 658 | Err(e) => { |
| 659 | stopped = Some(e.plain()); |
| 660 | break; |
| 661 | }, |
| 662 | } |
| 663 | } |
| 664 | absorbed = base.log.len() != held; |
| 665 | } |
| 666 | if let Some(said) = stopped { |
| 667 | // Nothing reached the segments, so a log an earlier message of this batch had |
| 668 | // already absorbed into is a log the disk does not answer to, and the reading |
| 669 | // goes rather than being filed against bytes that do not hold it. Where |
| 670 | // nothing was absorbed it is what the read produced and it stays. |
| 671 | match absorbed { |
| 672 | false => res!(host.readings().keep( |
| 673 | &hosted.dir, res!(hosted.store.bytes()), cursor, took, read)), |
| 674 | true => res!(host.readings().unsee(&hosted.dir)), |
| 675 | } |
| 676 | return Ok(Reply::text(409, said)); |
| 677 | } |
| 678 | // Everything that arrived reaches the segments before a word is said about |
| 679 | // it, so a reply the client acts on is a reply the relay has already kept. |
| 680 | let fresh = arrived(&read.log, &before, &read.envelopes, &read.veils); |
| 681 | if let Err(e) = hosted.store.append(&fresh, None) { |
| 682 | // The log holds operations the segments may or may not, and nothing can say |
| 683 | // which, so the reading goes and the next request reads the store as it |
| 684 | // stands. |
| 685 | res!(host.readings().unsee(&hosted.dir)); |
| 686 | return Err(e); |
| 687 | } |
| 688 | // A reading is a note about bytes and this request has just written some, so |
| 689 | // the few that were written are read back and the cursor carried over them. |
| 690 | // Without this a push would leave a reading nothing could be taken up from, and |
| 691 | // the request after every push would read the whole history -- which is the |
| 692 | // exact shape of the fault that shipped in the forge on 2026-08-22. |
| 693 | let mut hold = true; |
| 694 | if !fresh.is_empty() { |
| 695 | match carry_over(hosted, &cursor, Arc::make_mut(&mut read), fresh.len()) { |
| 696 | Ok(moved) => cursor = moved, |
| 697 | Err(e) => { |
| 698 | fault!("{}: {} operation{} reached the segments and the reading this \ |
| 699 | relay holds could not be carried over the bytes they were written \ |
| 700 | to, so it is being let go and the next request will read the whole \ |
| 701 | history: {}", hosted.label, fresh.len(), |
| 702 | if fresh.len() == 1 { "" } else { "s" }, e.plain()); |
| 703 | hold = false; |
| 704 | }, |
| 705 | } |
| 706 | } |
| 707 | // Put back before the reply is built rather than after, so that a reply this |
| 708 | // relay cannot frame does not also cost the next caller the whole history. |
| 709 | match hold { |
| 710 | true => res!(host.readings().keep( |
| 711 | &hosted.dir, res!(hosted.store.bytes()), cursor, took, Arc::clone(&read))), |
| 712 | false => res!(host.readings().unsee(&hosted.dir)), |
| 713 | } |
| 714 | // The reply is bounded like the request, and for a failure one end further on: |
| 715 | // framed whole, a large clone arrives as one body, on a receiving end where |
| 716 | // there is no proxy to blame and a phone or a small VPS has nothing to raise. |
| 717 | // |
| 718 | // This comment used to add "and is held entire while it is taken in, six times |
| 719 | // over". The bound is right; that reason was not. Measured 2026-08-20, the |
| 720 | // same clone across fourteen bounded replies peaked at 346,612 kB against |
| 721 | // 347,532 kB unbounded -- under one percent. The six times is the engine's |
| 722 | // per-operation cost of holding a history, which `ore log` pays with no |
| 723 | // network at all. What the bound buys is that no single body has to be |
| 724 | // materialised whole, which is a real failure and a different one. |
| 725 | // |
| 726 | // What does not fit is not sent and not remembered -- `Done` says only that |
| 727 | // this end will send no more this turn, and the client, which was told this |
| 728 | // relay's frontier in the opening above, sees that its log does not cover it |
| 729 | // and opens a fresh session. Resume is rerun, which is what §4.3 already |
| 730 | // promised for an interrupted exchange. |
| 731 | // |
| 732 | // And what does not fit is no longer built either. `outgoing` stops at the |
| 733 | // bound rather than substituting and serialising the whole owed set for |
| 734 | // `proto::upto` to throw away, which is where two thirds of this relay's |
| 735 | // processor went until 2026-08-23. |
| 736 | let mut reply = res!(outgoing(out, &read.envelopes, &read.veils, cap)); |
| 737 | let (fits, held_back) = res!(proto::upto(&reply, cap)); |
| 738 | if held_back { |
| 739 | reply.truncate(fits); |
| 740 | reply.push(Message::Done); |
| 741 | } |
| 742 | Ok(Reply::frames( |
| 743 | res!(proto::frame(&reply)), |
| 744 | session.ops_absorbed(), |
| 745 | )) |
| 746 | } |
| 747 | |
| 748 | /// Reads the hosted log, taking up from where the last request of it stopped. |
| 749 | /// |
| 750 | /// **Only the read is incremental, and on a relay that is the whole of the cost.** |
| 751 | /// A history is appended to and never rewritten, so the operations already read |
| 752 | /// are the operations still there, and reading them again is re-deriving a |
| 753 | /// constant. A relay renders nothing and verifies nothing, so nothing follows the |
| 754 | /// read that has to be done afresh: what a request pays after this is the session |
| 755 | /// walk and the bytes it sends. |
| 756 | /// |
| 757 | /// A store that will not be taken up -- because a segment already read is not the |
| 758 | /// segment it was, or because a reading placed an operation before a parent that |
| 759 | /// came later in the file -- is read whole instead. The refusal is |
| 760 | /// [`ore_store::store::Consumed`]'s and it is deliberate: the answer to it is to |
| 761 | /// read the store as it stands, which is what a puller is owed, and never to carry |
| 762 | /// a log across a change nothing has looked at. |
| 763 | /// |
| 764 | /// **Every whole read of a repository this relay has already read says so, in the |
| 765 | /// log, as it happens.** A fallback that is only slower is a fallback nobody |
| 766 | /// finds: exactly this change shipped in the forge on 2026-08-22 taking the whole |
| 767 | /// path on every push, was measured at no gain by the lane that deployed it, and |
| 768 | /// left no line anywhere saying which path it had taken. [`Whence`] names the five |
| 769 | /// ways a read can start and [`Took`] counts what each one cost, so the question |
| 770 | /// is answered by a number rather than by a stopwatch. |
| 771 | fn read_on(hosted: &Hosted, carried: Carried) |
| 772 | -> Outcome<(Arc<Replayed>, Consumed, Took)> |
| 773 | { |
| 774 | let (mut read, from, mut whence) = match carried { |
| 775 | Carried::From(held, from) if from.resumable() => (held, from, Whence::On), |
| 776 | Carried::From(_, from) => { |
| 777 | fault!("{}: the reading this relay holds of this repository placed an \ |
| 778 | operation before a parent that came later in its file, {} records into \ |
| 779 | {} bytes, so it cannot be taken up from and the store is being read \ |
| 780 | whole instead.", hosted.label, from.records(), from.bytes()); |
| 781 | (Arc::new(Replayed::new()), Consumed::new(), Whence::Unresumable) |
| 782 | }, |
| 783 | Carried::Fresh => (Arc::new(Replayed::new()), Consumed::new(), Whence::First), |
| 784 | Carried::Lost => { |
| 785 | fault!("{}: this relay has read this repository before and is holding no \ |
| 786 | reading of it, so the store is being read whole instead of extended. A \ |
| 787 | reading let go on purpose is let go where this relay can see it, so this \ |
| 788 | is a defect in the reading cache and it costs every request after it the \ |
| 789 | whole history again.", hosted.label); |
| 790 | (Arc::new(Replayed::new()), Consumed::new(), Whence::Lost) |
| 791 | }, |
| 792 | }; |
| 793 | let base = Arc::make_mut(&mut read); |
| 794 | // Nothing is verified on the way in. A signature the relay could check is one |
| 795 | // the puller must check anyway, and a relay that held a trust set would be |
| 796 | // the beginning of an authority. |
| 797 | let taken = match hosted.store.replay_since(&from, base, Verify::Nothing, Keep::Envelopes) { |
| 798 | Ok(taken) => taken, |
| 799 | Err(e) => { |
| 800 | // A first read that fails is the store failing, and the caller is told so. |
| 801 | if from.is_empty() { |
| 802 | return Err(e); |
| 803 | } |
| 804 | fault!("{}: the store would not be taken up from where the last request of \ |
| 805 | it stopped, so it is being read whole instead: {}", hosted.label, e); |
| 806 | whence = Whence::Refused; |
| 807 | *base = Replayed::new(); |
| 808 | res!(hosted.store.replay_since( |
| 809 | &Consumed::new(), base, Verify::Nothing, Keep::Envelopes)) |
| 810 | }, |
| 811 | }; |
| 812 | let skipped = match whence { |
| 813 | Whence::On => from.bytes(), |
| 814 | _ => 0, |
| 815 | }; |
| 816 | let took = Took { |
| 817 | whence, |
| 818 | skipped, |
| 819 | read: taken.bytes().saturating_sub(skipped), |
| 820 | }; |
| 821 | Ok((read, taken, took)) |
| 822 | } |
| 823 | |
| 824 | /// Reads back what this request appended, so that the reading names those bytes |
| 825 | /// too. |
| 826 | /// |
| 827 | /// **The records are reconciled, not discarded.** The log already holds them -- |
| 828 | /// the session absorbed them a moment ago -- so what comes off the disk is put to |
| 829 | /// the two questions that say whether the read and the writing agree: are there |
| 830 | /// exactly as many operations in those bytes as this request wrote, and is every |
| 831 | /// one of them an operation the log holds? A cursor kept over bytes that said |
| 832 | /// something else would be a reading filed against a store it was not taken from, |
| 833 | /// which is the whole fault [`ore_store::store::Consumed`] exists to prevent, |
| 834 | /// arrived at by a shorter road. |
| 835 | /// |
| 836 | /// Where they disagree the caller lets the reading go, and the next request reads |
| 837 | /// the store as it stands. |
| 838 | fn carry_over(hosted: &Hosted, from: &Consumed, into: &mut Replayed, appended: usize) |
| 839 | -> Outcome<Consumed> |
| 840 | { |
| 841 | let got = res!(hosted.store.read_since(from, Verify::Nothing, Keep::Envelopes)); |
| 842 | if got.records.len() != appended { |
| 843 | return Err(err!( |
| 844 | "This request appended {} operation{} to {} and reading the segments back \ |
| 845 | found {}.", appended, if appended == 1 { "" } else { "s" }, |
| 846 | hosted.label, got.records.len(); |
| 847 | Invalid, Data, Mismatch)); |
| 848 | } |
| 849 | for rec in &got.records { |
| 850 | if !into.log.contains(&rec.id()) { |
| 851 | return Err(err!( |
| 852 | "The segments of {} gave back the operation {}, which this request did \ |
| 853 | not absorb and the log does not hold.", hosted.label, rec.id(); |
| 854 | Invalid, Data, Mismatch)); |
| 855 | } |
| 856 | } |
| 857 | // What is on the disk wins, for the reason it does on the way in: an operation |
| 858 | // the log already holds keeps the form already written, since upgrading it |
| 859 | // would mean rewriting a segment. |
| 860 | for (id, env) in got.envelopes { |
| 861 | into.envelopes.insert(id, env); |
| 862 | } |
| 863 | for (id, veiled) in got.veils { |
| 864 | into.veils.insert(id, veiled); |
| 865 | } |
| 866 | Ok(got.cursor) |
| 867 | } |
| 868 | |
| 869 | /// Puts the provenance and the veils back on what is going out, splits a send |
| 870 | /// that is too large into several that are not, and stops once it has the reply |
| 871 | /// bound's worth. |
| 872 | /// |
| 873 | /// A session builds a send set out of the log, and the log holds records, so what |
| 874 | /// it produces is bare. Two substitutions put back what the relay was given: an |
| 875 | /// envelope where it holds one, and a veiled entry where the record in the log is |
| 876 | /// only the stand-in for one. The second is what keeps a stand-in inside the |
| 877 | /// relay: it is the only path by which an entry the session chose reaches the |
| 878 | /// wire. |
| 879 | /// |
| 880 | /// # Why it stops |
| 881 | /// |
| 882 | /// A session owes what the far end's frontier does not cover, which on a clone is |
| 883 | /// the whole history, and a reply carries [`Host::reply_bytes`] of it. This built |
| 884 | /// the owed turn entire -- copying every envelope, serialising every entry to |
| 885 | /// measure it, and serialising every message again to find where the bound fell |
| 886 | /// -- and then threw all but the first six mebibytes away. Measured on the fe2o3 |
| 887 | /// copy at the deployed bound: thirty-two requests, of which sixteen carry |
| 888 | /// operations, 1.95 s of the relay's 2.48 s spent here, and the first request of |
| 889 | /// the clone alone spent 346 ms building eighty-nine megabytes to send five. |
| 890 | /// |
| 891 | /// So entries are substituted, measured and batched one at a time, and the walk |
| 892 | /// stops at the first entry that would take the running total past `cap`. |
| 893 | /// **Nothing that could have been kept is dropped.** The frames of a unit come to |
| 894 | /// more than the entries in it, so an entry that takes the entries past `cap` is |
| 895 | /// in a unit that takes the frames past it too, and the one unit |
| 896 | /// [`proto::upto`] lets through over the bound -- the first carrying operations |
| 897 | /// -- is a batch bounded by [`proto::BATCH_BYTES`] and the caller's own minimum, |
| 898 | /// both under `cap`. The single exception is an operation larger than the whole |
| 899 | /// reply, and that is why the first entry is taken whatever it comes to. |
| 900 | /// |
| 901 | /// The exact cut is still [`proto::upto`]'s, over a list that is now a reply long |
| 902 | /// instead of a history long, so the bound is decided in one place and the two |
| 903 | /// cannot disagree. |
| 904 | /// |
| 905 | /// No state is kept for this and none is needed. What is not sent is not |
| 906 | /// remembered: the client sees that its log does not cover the frontier this |
| 907 | /// relay opened with and comes back, and the session that answers it works the |
| 908 | /// owed set out afresh over a frontier that has moved. |
| 909 | fn outgoing( |
| 910 | msgs: Vec<Message>, |
| 911 | envelopes: &BTreeMap<OpId, oxedyne_fe2o3_ore::envelope::Envelope>, |
| 912 | veils: &BTreeMap<OpId, Veiled>, |
| 913 | cap: usize, |
| 914 | ) |
| 915 | -> Outcome<Vec<Message>> |
| 916 | { |
| 917 | let mut out = Vec::new(); |
| 918 | // What the entries taken so far come to, across every send in the turn. A |
| 919 | // session sends one, and a turn of several is bounded as a whole rather than |
| 920 | // once each. |
| 921 | let mut taken = 0usize; |
| 922 | for msg in msgs { |
| 923 | let entries = match msg { |
| 924 | Message::Send { entries } => entries, |
| 925 | other => { |
| 926 | out.push(other); |
| 927 | continue; |
| 928 | }, |
| 929 | }; |
| 930 | // Never larger than the reply that carries them: an operator who lowers the |
| 931 | // reply bound is asking for smaller answers, and batches that ignored it |
| 932 | // would make every reply a single oversized message the bound then has to |
| 933 | // let through anyway. An operation larger than the bound on its own leaves |
| 934 | // as a run of pieces, which `proto::upto` keeps whole -- a reply is |
| 935 | // answered by an end that keeps nothing between requests, so half a run |
| 936 | // sent is half a run the client throws away and the next session sends |
| 937 | // again. |
| 938 | let mut batching = proto::Batching::upto( |
| 939 | std::cmp::min(proto::BATCH_BYTES, cap), cap.saturating_sub(taken)); |
| 940 | for entry in entries { |
| 941 | let entry = res!(restore_one(res!(seal_one(entry, envelopes)), veils)); |
| 942 | match res!(batching.take(entry)) { |
| 943 | proto::Fit::Took(msgs) => out.extend(msgs), |
| 944 | // Enough. The entry handed back is dropped and so is everything after |
| 945 | // it, because neither could have been kept: an entry that takes the |
| 946 | // entries past `cap` is in a unit whose frames come to more than `cap`, |
| 947 | // and the only unit `proto::upto` lets through over the bound is the |
| 948 | // first one carrying operations, which is the batch this entry did not |
| 949 | // fit into. |
| 950 | proto::Fit::Full(_) => break, |
| 951 | } |
| 952 | } |
| 953 | taken += batching.taken(); |
| 954 | out.extend(batching.rest()); |
| 955 | } |
| 956 | Ok(out) |
| 957 | } |
| 958 | |
| 959 | /// Returns the mode this end should open in, judged from the opening it was |
| 960 | /// sent. |
| 961 | /// |
| 962 | /// A client that opened with a sketch is answered with one, so the saving runs |
| 963 | /// both ways; a client that walked is walked back to. The estimate is the |
| 964 | /// engine's own rule over the two shapes, the client's shape being what its |
| 965 | /// sketch message states. |
| 966 | fn mode_for(log: &OpLog, msgs: &[Message]) -> Mode { |
| 967 | for msg in msgs { |
| 968 | match msg { |
| 969 | Message::Hello { .. } => return Mode::Walk, |
| 970 | Message::Sketch { heads, count, .. } => return Mode::between( |
| 971 | log.len(), |
| 972 | log.frontier().len(), |
| 973 | *count as usize, |
| 974 | heads.len(), |
| 975 | ), |
| 976 | _ => (), |
| 977 | } |
| 978 | } |
| 979 | Mode::Walk |
| 980 | } |
| 981 | |
| 982 | /// Says no, and says what would have been needed. |
| 983 | fn refused(hosted: &Hosted, want: Role) -> Reply { |
| 984 | Reply::text(403, fmt!( |
| 985 | "That key holds no {} on {}. Access to a hosted repository is granted by \ |
| 986 | an administrator with `ore-relay grant`, and is a fact about this relay's \ |
| 987 | disk rather than about the history.", |
| 988 | want.name(), hosted.label, |
| 989 | )) |
| 990 | } |
| 991 | |
| 992 | |
| 993 | /// Listens on an address, answering one connection at a time, for ever. |
| 994 | pub async fn listen(host: Host, addr: SocketAddr) |
| 995 | -> Outcome<()> |
| 996 | { |
| 997 | let listener = match TcpListener::bind(addr).await { |
| 998 | Ok(l) => l, |
| 999 | Err(e) => return Err(err!(e, |
| 1000 | "The relay could not listen on {}.", addr; |
| 1001 | IO, Network, Init)), |
| 1002 | }; |
| 1003 | let bound = match listener.local_addr() { |
| 1004 | Ok(a) => a, |
| 1005 | Err(e) => return Err(err!(e, |
| 1006 | "The relay listened on {} and cannot say where.", addr; |
| 1007 | IO, Network, Init)), |
| 1008 | }; |
| 1009 | println!("relay listening on {}, holding {}", bound, host.dir.display()); |
| 1010 | loop { |
| 1011 | let (stream, peer) = match listener.accept().await { |
| 1012 | Ok(pair) => pair, |
| 1013 | Err(e) => { |
| 1014 | // One connection failing to arrive is not the relay failing. |
| 1015 | eprintln!("relay: a connection could not be accepted: {}", e); |
| 1016 | continue; |
| 1017 | }, |
| 1018 | }; |
| 1019 | if let Err(e) = carry(&host, stream).await { |
| 1020 | eprintln!("relay: {} was answered with nothing: {}", peer, e.plain()); |
| 1021 | } |
| 1022 | } |
| 1023 | } |
| 1024 | |
| 1025 | /// Reads one request off a connection, answers it, and closes. |
| 1026 | /// |
| 1027 | /// One request per connection, as the client's own transport does, which is a |
| 1028 | /// connection state machine neither end has to have. |
| 1029 | /// |
| 1030 | /// The reason given here used to be that keep-alive "would save a handshake on an |
| 1031 | /// exchange that is two round trips long". An exchange is no longer two: bounding |
| 1032 | /// both directions made a 58 MB clone twenty-eight requests, so what is being |
| 1033 | /// declined is twenty-seven handshakes and, over TLS, twenty-seven of the client's |
| 1034 | /// per-request `letsencrypt_client_config`. Still declined, and now knowingly. |
| 1035 | async fn carry(host: &Host, mut stream: TcpStream) |
| 1036 | -> Outcome<()> |
| 1037 | { |
| 1038 | let limits = ReadLimits { |
| 1039 | max_header_bytes: Some(64 << 10), |
| 1040 | max_body_bytes: Some(proto::FRAME_LIMIT), |
| 1041 | header_read_timeout: Some(std::time::Duration::from_secs(30)), |
| 1042 | }; |
| 1043 | let (msg, _) = res!(HttpMessage::read::< |
| 1044 | { constant::HTTP_DEFAULT_HEADER_CHUNK_SIZE }, |
| 1045 | { constant::HTTP_DEFAULT_BODY_CHUNK_SIZE }, |
| 1046 | _, |
| 1047 | >( |
| 1048 | Pin::new(&mut stream), |
| 1049 | &Vec::new(), |
| 1050 | Some(true), |
| 1051 | Some(&limits), |
| 1052 | ).await); |
| 1053 | let msg = match msg { |
| 1054 | Some(m) => m, |
| 1055 | None => return Ok(()), |
| 1056 | }; |
| 1057 | let (method, path) = match &msg.header.headline { |
| 1058 | HttpHeadline::Request { method, loc } => ( |
| 1059 | fmt!("{}", method_name(method)), |
| 1060 | fmt!("{}", loc.path.as_str()), |
| 1061 | ), |
| 1062 | HttpHeadline::Response { .. } => return Err(err!( |
| 1063 | "A response arrived where a request was expected."; |
| 1064 | Invalid, Input, Mismatch)), |
| 1065 | }; |
| 1066 | let cred = presented(&msg); |
| 1067 | let reply = respond(host, &Request { |
| 1068 | method, |
| 1069 | path, |
| 1070 | cred, |
| 1071 | body: msg.body, |
| 1072 | }); |
| 1073 | let status = match HttpStatus::from_repr(reply.status) { |
| 1074 | Some(s) => s, |
| 1075 | None => HttpStatus::InternalServerError, |
| 1076 | }; |
| 1077 | let mut out = HttpMessage::new_response(status); |
| 1078 | out.insert( |
| 1079 | HeaderName::ContentType, |
| 1080 | HeaderFieldValue::Generic(reply.kind), |
| 1081 | Some(HeaderFieldCategory::Entity as u16), |
| 1082 | ); |
| 1083 | if let Some(n) = reply.absorbed { |
| 1084 | out.insert( |
| 1085 | HeaderName::from(HEADER_ABSORBED), |
| 1086 | HeaderFieldValue::Generic(fmt!("{}", n)), |
| 1087 | Some(HeaderFieldCategory::Other as u16), |
| 1088 | ); |
| 1089 | } |
| 1090 | out.set_connection_close(true); |
| 1091 | out.body = reply.body; |
| 1092 | res!(out.write_all(&mut stream).await); |
| 1093 | Ok(()) |
| 1094 | } |
| 1095 | |
| 1096 | /// Returns what a request said about who is making it, or nothing where it said |
| 1097 | /// nothing. |
| 1098 | /// |
| 1099 | /// A header that is there and malformed is treated as no credential at all: the |
| 1100 | /// access list then refuses everything a public pull does not cover, which is the |
| 1101 | /// same answer a wrong signature gets and needs no separate path. |
| 1102 | fn presented(msg: &HttpMessage) -> Option<Presented> { |
| 1103 | let field = |name: &str| -> Option<String> { |
| 1104 | msg.header.fields |
| 1105 | .get_one(&HeaderName::from(name)) |
| 1106 | .map(|v| fmt!("{}", v)) |
| 1107 | }; |
| 1108 | let replica = match field(proto::HEADER_REPLICA) { Some(v) => v, None => return None }; |
| 1109 | let key = match field(proto::HEADER_KEY) { Some(v) => v, None => return None }; |
| 1110 | let stamp = match field(proto::HEADER_TIME) { Some(v) => v, None => return None }; |
| 1111 | let sig = match field(proto::HEADER_SIG) { Some(v) => v, None => return None }; |
| 1112 | Presented::read(&replica, &key, &stamp, &sig).ok() |
| 1113 | } |
| 1114 | |
| 1115 | /// Returns the method's name, which is what the signature covers. |
| 1116 | fn method_name(method: &HttpMethod) -> &'static str { |
| 1117 | match method { |
| 1118 | HttpMethod::CONNECT => "CONNECT", |
| 1119 | HttpMethod::DELETE => "DELETE", |
| 1120 | HttpMethod::GET => "GET", |
| 1121 | HttpMethod::HEAD => "HEAD", |
| 1122 | HttpMethod::OPTIONS => "OPTIONS", |
| 1123 | HttpMethod::PATCH => "PATCH", |
| 1124 | HttpMethod::POST => "POST", |
| 1125 | HttpMethod::PUT => "PUT", |
| 1126 | HttpMethod::TRACE => "TRACE", |
| 1127 | } |
| 1128 | } |
| 1129 | |
| 1130 | |
| 1131 | #[cfg(test)] |
| 1132 | mod tests { |
| 1133 | use super::*; |
| 1134 | |
| 1135 | use crate::acl::Acl; |
| 1136 | |
| 1137 | use oxedyne_fe2o3_ore::id::ReplicaId; |
| 1138 | use oxedyne_fe2o3_ore::op::{ |
| 1139 | Header, |
| 1140 | Op, |
| 1141 | Record, |
| 1142 | }; |
| 1143 | |
| 1144 | use std::fs; |
| 1145 | use std::path::PathBuf; |
| 1146 | use std::time::{ |
| 1147 | SystemTime, |
| 1148 | UNIX_EPOCH, |
| 1149 | }; |
| 1150 | |
| 1151 | /// A directory that removes itself however the test ends. |
| 1152 | struct Scratch { |
| 1153 | path: PathBuf, |
| 1154 | } |
| 1155 | |
| 1156 | impl Scratch { |
| 1157 | fn new(what: &str) |
| 1158 | -> Outcome<Self> |
| 1159 | { |
| 1160 | let stamp = res!(SystemTime::now().duration_since(UNIX_EPOCH)); |
| 1161 | let path = std::env::temp_dir().join(fmt!( |
| 1162 | "ore_relay_{}_{}_{}", what, std::process::id(), stamp.as_nanos(), |
| 1163 | )); |
| 1164 | res!(fs::create_dir_all(&path)); |
| 1165 | Ok(Self { path }) |
| 1166 | } |
| 1167 | } |
| 1168 | |
| 1169 | impl Drop for Scratch { |
| 1170 | fn drop(&mut self) { |
| 1171 | let _ = fs::remove_dir_all(&self.path); |
| 1172 | } |
| 1173 | } |
| 1174 | |
| 1175 | /// A run of operations, each the child of the one before it. |
| 1176 | fn history(replica: u64, from: u64, how_many: u64) |
| 1177 | -> Outcome<Vec<Entry>> |
| 1178 | { |
| 1179 | let mut out = Vec::new(); |
| 1180 | for n in from..from + how_many { |
| 1181 | // Counters start at one, and the first operation of a replica has no |
| 1182 | // parent. |
| 1183 | let parents = match n { |
| 1184 | 1 => Vec::new(), |
| 1185 | _ => vec![OpId::new(ReplicaId::new(replica), n - 1)], |
| 1186 | }; |
| 1187 | out.push(Entry::Bare(Record::new( |
| 1188 | res!(Header::new(OpId::new(ReplicaId::new(replica), n), parents)), |
| 1189 | Op::FileCreate { path: fmt!("file{}.txt", n).into_bytes() }, |
| 1190 | ))); |
| 1191 | } |
| 1192 | Ok(out) |
| 1193 | } |
| 1194 | |
| 1195 | /// What a client with an empty log posts to open a session. |
| 1196 | fn opening(account: &str, name: &str) |
| 1197 | -> Outcome<Request> |
| 1198 | { |
| 1199 | let mine = OpLog::new(); |
| 1200 | let mut session = Session::new(Mode::Walk); |
| 1201 | let said = vec![res!(session.open(&mine)), Message::Done]; |
| 1202 | Ok(Request { |
| 1203 | method: fmt!("POST"), |
| 1204 | path: fmt!("{}/{}/{}/{}/sync", proto::PREFIX, proto::VERSION, account, name), |
| 1205 | cred: None, |
| 1206 | body: res!(proto::frame(&said)), |
| 1207 | }) |
| 1208 | } |
| 1209 | |
| 1210 | /// What a request carries to say who is making it. |
| 1211 | fn credential(signer: &ore_store::keys::Signing, method: &str, path: &str, body: &[u8]) |
| 1212 | -> Outcome<Presented> |
| 1213 | { |
| 1214 | let headers = res!(Presented::sign(signer, method, path, body)); |
| 1215 | let value = |name: &str| -> Outcome<String> { |
| 1216 | match headers.iter().find(|(n, _)| n == name) { |
| 1217 | Some((_, v)) => Ok(v.clone()), |
| 1218 | None => Err(err!( |
| 1219 | "A signed request carries no {} header.", name; Test, Missing)), |
| 1220 | } |
| 1221 | }; |
| 1222 | Presented::read( |
| 1223 | &res!(value(proto::HEADER_REPLICA)), |
| 1224 | &res!(value(proto::HEADER_KEY)), |
| 1225 | &res!(value(proto::HEADER_TIME)), |
| 1226 | &res!(value(proto::HEADER_SIG)), |
| 1227 | ) |
| 1228 | } |
| 1229 | |
| 1230 | /// One request, answered, with a failure that says what the relay said. |
| 1231 | fn served(host: &Host, req: &Request) |
| 1232 | -> Outcome<Reply> |
| 1233 | { |
| 1234 | let reply = respond(host, req); |
| 1235 | if reply.status != 200 { |
| 1236 | return Err(err!("The relay answered {}: {}", |
| 1237 | reply.status, String::from_utf8_lossy(&reply.body); Test, Invalid)); |
| 1238 | } |
| 1239 | Ok(reply) |
| 1240 | } |
| 1241 | |
| 1242 | /// What the reading cache says the last read of a repository had to do. |
| 1243 | fn took(host: &Host, hosted: &Hosted) |
| 1244 | -> Outcome<Took> |
| 1245 | { |
| 1246 | Ok(res!(res!(host.readings().took(&hosted.dir)).ok_or_else(|| err!( |
| 1247 | "This relay is holding no reading of {}, so it cannot say what the last \ |
| 1248 | read of it did.", hosted.label; |
| 1249 | Test, Missing)))) |
| 1250 | } |
| 1251 | |
| 1252 | /// A request reads the bytes appended since the last one and no others, and a |
| 1253 | /// reading this relay drops is reported rather than paid for in silence. |
| 1254 | /// |
| 1255 | /// Wall time cannot tell a resumed read from a whole one on a busy host, which |
| 1256 | /// is why this asks [`Took`] instead. The same change shipped in the forge on |
| 1257 | /// 2026-08-22 taking the whole path every time, was reported as met, and was |
| 1258 | /// caught only by a lane that measured four arms and found no difference |
| 1259 | /// between them. |
| 1260 | /// |
| 1261 | /// Proved red three ways: having `sync` pass `Carried::Fresh` to `read_on` |
| 1262 | /// whatever the cache held, which makes every read `First`; dropping the |
| 1263 | /// `carry_over` call after the append, which makes the read after a write |
| 1264 | /// `Refused` and whole; and having `Readings::take` answer `Carried::Fresh` |
| 1265 | /// where it holds nothing, which makes the dropped reading read as an ordinary |
| 1266 | /// first request. |
| 1267 | #[test] |
| 1268 | fn a_request_reads_what_was_appended_since_the_last_one() -> Outcome<()> { |
| 1269 | let scratch = res!(Scratch::new("readings")); |
| 1270 | let host = Host::at(&scratch.path); |
| 1271 | let hosted = res!(host.create("someone", "notes", &Acl::new(b"owner".to_vec(), true))); |
| 1272 | res!(hosted.store.append(&res!(history(7, 1, 40)), None)); |
| 1273 | let stored = res!(hosted.store.bytes()); |
| 1274 | assert!(stored > 0, "the fixture wrote no segments"); |
| 1275 | |
| 1276 | // The first request, which was always going to read the whole store. |
| 1277 | res!(served(&host, &res!(opening("someone", "notes")))); |
| 1278 | let first = res!(took(&host, &hosted)); |
| 1279 | assert_eq!(first.whence, Whence::First, "the first read was not a first read"); |
| 1280 | assert_eq!(first.skipped, 0, "the first read skipped bytes nothing had read"); |
| 1281 | assert_eq!(first.read, stored, "the first read did not read the store"); |
| 1282 | |
| 1283 | // A second, over a store nothing has touched since. |
| 1284 | res!(served(&host, &res!(opening("someone", "notes")))); |
| 1285 | let again = res!(took(&host, &hosted)); |
| 1286 | assert_eq!(again.whence, Whence::On, |
| 1287 | "the second request read the store whole, which is the fault this exists \ |
| 1288 | to catch"); |
| 1289 | assert_eq!(again.read, 0, |
| 1290 | "the second request read {} bytes of a store nothing had appended to", |
| 1291 | again.read); |
| 1292 | assert_eq!(again.skipped, stored, "the second request did not skip the store"); |
| 1293 | |
| 1294 | // And one over a store that has grown, which reads the growth and no more. |
| 1295 | res!(hosted.store.append(&res!(history(7, 41, 10)), None)); |
| 1296 | let grown = res!(hosted.store.bytes()); |
| 1297 | assert!(grown > stored, "the fixture appended nothing"); |
| 1298 | res!(served(&host, &res!(opening("someone", "notes")))); |
| 1299 | let after = res!(took(&host, &hosted)); |
| 1300 | assert_eq!(after.whence, Whence::On, "the read over an appended store was whole"); |
| 1301 | assert_eq!(after.skipped, stored, "the read over an appended store skipped nothing"); |
| 1302 | assert_eq!(after.read, grown - stored, |
| 1303 | "the read over an appended store read {} bytes for {} of growth", |
| 1304 | after.read, grown - stored); |
| 1305 | |
| 1306 | // A reading this relay took and did not put back. Nothing ordinary does it, |
| 1307 | // so the next read says so and starts at the first byte. |
| 1308 | let _dropped = res!(host.readings().take(&hosted.dir)); |
| 1309 | res!(served(&host, &res!(opening("someone", "notes")))); |
| 1310 | let lost = res!(took(&host, &hosted)); |
| 1311 | assert_eq!(lost.whence, Whence::Lost, |
| 1312 | "a reading this relay dropped read as an ordinary first request, which is \ |
| 1313 | a fallback nobody would ever find"); |
| 1314 | assert_eq!(lost.read, grown, "the read after a lost reading did not read the store"); |
| 1315 | Ok(()) |
| 1316 | } |
| 1317 | |
| 1318 | /// A reading survives the request that writes to the store, which is the one |
| 1319 | /// moment an incremental read exists for. |
| 1320 | /// |
| 1321 | /// A push is the write a relay actually sees, and it is where the forge's |
| 1322 | /// version of this shipped inert: the cache was emptied on every push, so the |
| 1323 | /// read that was to take up where the last one stopped was handed nothing on |
| 1324 | /// the only path that mattered. Here the reading is carried over the bytes the |
| 1325 | /// push wrote -- see [`carry_over`] -- so the request after a push skips |
| 1326 | /// everything the push did not write. |
| 1327 | /// |
| 1328 | /// Proved red by dropping the `carry_over` call, which leaves the cursor naming |
| 1329 | /// bytes the store has grown past; the next read is then `Whence::Refused` and |
| 1330 | /// whole, and says so. |
| 1331 | #[test] |
| 1332 | fn a_push_leaves_a_reading_the_next_request_takes_up_from() -> Outcome<()> { |
| 1333 | let scratch = res!(Scratch::new("pushed")); |
| 1334 | let host = Host::at(&scratch.path); |
| 1335 | // A push is signed, always: an unsigned caller may pull from a public |
| 1336 | // repository and may never write to one. So the key that will do the pushing |
| 1337 | // owns it. |
| 1338 | let signer = res!(ore_store::keys::Signing::mint(ReplicaId::new(9))); |
| 1339 | let path = fmt!("{}/{}/someone/notes/sync", proto::PREFIX, proto::VERSION); |
| 1340 | let body = res!(proto::frame(&vec![ |
| 1341 | Message::Send { entries: res!(history(9, 1, 1)) }, |
| 1342 | Message::Done, |
| 1343 | ])); |
| 1344 | let cred = res!(credential(&signer, "POST", &path, &body)); |
| 1345 | let hosted = res!(host.create( |
| 1346 | "someone", "notes", &Acl::new(cred.public.clone(), true))); |
| 1347 | res!(hosted.store.append(&res!(history(7, 1, 40)), None)); |
| 1348 | |
| 1349 | // A first request, so that there is a reading for the push to keep. |
| 1350 | res!(served(&host, &res!(opening("someone", "notes")))); |
| 1351 | let before = res!(hosted.store.bytes()); |
| 1352 | |
| 1353 | // A push of one operation the relay does not hold. |
| 1354 | let reply = res!(served(&host, &Request { |
| 1355 | method: fmt!("POST"), |
| 1356 | path, |
| 1357 | cred: Some(cred), |
| 1358 | body, |
| 1359 | })); |
| 1360 | assert_eq!(reply.absorbed, Some(1), "the push absorbed nothing to carry over"); |
| 1361 | let grown = res!(hosted.store.bytes()); |
| 1362 | assert!(grown > before, "the push wrote no segments"); |
| 1363 | let pushed = res!(took(&host, &hosted)); |
| 1364 | assert_eq!(pushed.whence, Whence::On, "the push read the store whole"); |
| 1365 | |
| 1366 | // And the request after it takes up where the push left off, having read the |
| 1367 | // bytes the push wrote as part of carrying the reading over them. |
| 1368 | res!(served(&host, &res!(opening("someone", "notes")))); |
| 1369 | let after = res!(took(&host, &hosted)); |
| 1370 | assert_eq!(after.whence, Whence::On, |
| 1371 | "the request after a push read the store whole, which is the fault the \ |
| 1372 | forge shipped on 2026-08-22"); |
| 1373 | assert_eq!(after.skipped, grown, "the request after a push skipped {} of {}", |
| 1374 | after.skipped, grown); |
| 1375 | assert_eq!(after.read, 0, "the request after a push read {} bytes", after.read); |
| 1376 | Ok(()) |
| 1377 | } |
| 1378 | |
| 1379 | /// The entries carried by a run of messages, in the order they appear. |
| 1380 | fn carried(msgs: &[Message]) |
| 1381 | -> Outcome<Vec<OpId>> |
| 1382 | { |
| 1383 | let mut out = Vec::new(); |
| 1384 | for msg in msgs { |
| 1385 | if let Message::Send { entries } = msg { |
| 1386 | for entry in entries { |
| 1387 | out.push(res!(entry.id())); |
| 1388 | } |
| 1389 | } |
| 1390 | } |
| 1391 | Ok(out) |
| 1392 | } |
| 1393 | |
| 1394 | /// [`outgoing`] builds a reply and not a history: it stops at the bound rather |
| 1395 | /// than substituting, measuring and batching everything the session owes and |
| 1396 | /// letting [`proto::upto`] throw the rest away. |
| 1397 | /// |
| 1398 | /// The cost this exists to stop is not the reply, which is bounded either way. |
| 1399 | /// It is the work: every envelope copied, every entry serialised to be |
| 1400 | /// measured, and every message serialised again to find where the bound fell. |
| 1401 | /// On a clone of fe2o3's history at the deployed bound that was 1.95 s of the |
| 1402 | /// relay's 2.48 s, and eighty-nine megabytes built to send six, thirty-two |
| 1403 | /// times over. |
| 1404 | /// |
| 1405 | /// So this asserts what [`outgoing`] hands back, not what the reply carries. |
| 1406 | /// Truncating after leaves the reply right and this test red, which is the |
| 1407 | /// reason for asking here. |
| 1408 | /// |
| 1409 | /// Proved red by building the batching with `Batching::to` instead of |
| 1410 | /// `Batching::upto`, which takes all two hundred entries and fails the first |
| 1411 | /// assertion with the whole history in hand. |
| 1412 | /// |
| 1413 | /// **It cannot see the `break`.** A batching that is full hands every further |
| 1414 | /// entry straight back, so dropping the `break` leaves what is built exactly as |
| 1415 | /// it is and costs only the measuring of the entries beyond the bound. That is |
| 1416 | /// the difference between this being cheap and this being right, and no |
| 1417 | /// assertion here reaches it; the measurement in `~/.cache/ore-trials/lane-g` |
| 1418 | /// is what does. |
| 1419 | #[test] |
| 1420 | fn outgoing_builds_a_reply_and_not_a_history() -> Outcome<()> { |
| 1421 | let entries = res!(history(7, 1, 200)); |
| 1422 | let ids = res!(carried(&[Message::Send { entries: entries.clone() }])); |
| 1423 | let whole = res!(proto::frame(&vec![Message::Send { entries: entries.clone() }])).len(); |
| 1424 | let envelopes = BTreeMap::new(); |
| 1425 | let veils = BTreeMap::new(); |
| 1426 | |
| 1427 | // A bound a fifth of the history, which is the shape of a real clone. |
| 1428 | let cap = whole / 5; |
| 1429 | let turn = vec![Message::Send { entries: entries.clone() }, Message::Done]; |
| 1430 | let built = res!(outgoing(turn, &envelopes, &veils, cap)); |
| 1431 | let got = res!(carried(&built)); |
| 1432 | assert!(got.len() < ids.len(), |
| 1433 | "outgoing built all {} entries against a bound a fifth of them, which is \ |
| 1434 | the whole history serialised to send a fifth of it", ids.len()); |
| 1435 | assert!(!got.is_empty(), "outgoing built nothing, so the reply carries no progress"); |
| 1436 | assert_eq!(got, ids[..got.len()], |
| 1437 | "outgoing built entries that are not the first {} in append order, and a \ |
| 1438 | prefix of the owed set is the only thing a peer can absorb", got.len()); |
| 1439 | |
| 1440 | // It stops at the bound and not short of it: everything it built is kept, so |
| 1441 | // the reply is a full one and nothing was built to be thrown away. |
| 1442 | let (fits, _) = res!(proto::upto(&built, cap)); |
| 1443 | let sent = res!(carried(&built[..fits])); |
| 1444 | assert_eq!(sent, got, |
| 1445 | "outgoing built {} entries and the reply carried {}", got.len(), sent.len()); |
| 1446 | |
| 1447 | // And a bound nothing reaches leaves the turn whole, or a clone that fits in |
| 1448 | // one reply would still take two. |
| 1449 | let all = res!(outgoing( |
| 1450 | vec![Message::Send { entries }, Message::Done], &envelopes, &veils, whole * 2)); |
| 1451 | assert_eq!(res!(carried(&all)), ids, "a turn under the bound was cut short"); |
| 1452 | Ok(()) |
| 1453 | } |
| 1454 | |
| 1455 | /// A clone crosses whole under a bound smaller than the history, and no reply |
| 1456 | /// carries more than the bound. |
| 1457 | /// |
| 1458 | /// What [`outgoing`] stopping early must not cost. The relay keeps nothing |
| 1459 | /// between requests, so a turn that stops short is a turn the next session |
| 1460 | /// works out afresh over a frontier that has moved; this walks that loop until |
| 1461 | /// the client holds everything, and fails rather than spins where a request |
| 1462 | /// moves it no further. |
| 1463 | /// |
| 1464 | /// Proved red by having [`outgoing`] break before taking any entry, which |
| 1465 | /// leaves every reply carrying nothing and stops the loop on the first request. |
| 1466 | #[test] |
| 1467 | fn a_bounded_clone_crosses_whole() -> Outcome<()> { |
| 1468 | let scratch = res!(Scratch::new("bounded")); |
| 1469 | let want = 400usize; |
| 1470 | let bound = 4 << 10; |
| 1471 | let host = Host::at(&scratch.path).with_reply_bytes(bound); |
| 1472 | let hosted = res!(host.create("someone", "notes", &Acl::new(b"owner".to_vec(), true))); |
| 1473 | res!(hosted.store.append(&res!(history(7, 1, want as u64)), None)); |
| 1474 | |
| 1475 | let mut mine = OpLog::new(); |
| 1476 | let mut requests = 0usize; |
| 1477 | let mut largest = 0usize; |
| 1478 | while mine.len() < want { |
| 1479 | // A session each time, because a bounded reply leaves this log short of |
| 1480 | // the frontier it was told and the next visit opens afresh. |
| 1481 | let mut session = Session::new(Mode::Walk); |
| 1482 | let said = vec![res!(session.open(&mine)), Message::Done]; |
| 1483 | let reply = res!(served(&host, &Request { |
| 1484 | method: fmt!("POST"), |
| 1485 | path: fmt!("{}/{}/someone/notes/sync", proto::PREFIX, proto::VERSION), |
| 1486 | cred: None, |
| 1487 | body: res!(proto::frame(&said)), |
| 1488 | })); |
| 1489 | requests += 1; |
| 1490 | largest = largest.max(reply.body.len()); |
| 1491 | let before = mine.len(); |
| 1492 | for msg in res!(proto::unframe(&reply.body)) { |
| 1493 | res!(session.receive(&mut mine, msg)); |
| 1494 | } |
| 1495 | if mine.len() == before { |
| 1496 | return Err(err!( |
| 1497 | "Request {} of a bounded clone moved it no further than {} of {} \ |
| 1498 | operations, so the exchange never ends.", requests, mine.len(), want; |
| 1499 | Test, Invalid)); |
| 1500 | } |
| 1501 | } |
| 1502 | assert_eq!(mine.len(), want, "the clone holds {} of {} operations", mine.len(), want); |
| 1503 | assert!(requests > 1, |
| 1504 | "the whole history crossed in one reply, so the bound never bit and this \ |
| 1505 | proves nothing"); |
| 1506 | // The opening the relay answers with rides in front of the batch, and |
| 1507 | // `proto::upto` lets the first unit carrying operations through whatever its |
| 1508 | // size so that a reply always carries progress. So a reply may come to the |
| 1509 | // bound and an opening, and never to a second batch over it. In service the |
| 1510 | // question does not arise: `proto::BATCH_BYTES` is two thirds of |
| 1511 | // `proto::REPLY_BYTES`, and it is only a bound below one batch that lets a |
| 1512 | // batch fill the whole of it. |
| 1513 | let opening = res!(proto::frame(&vec![ |
| 1514 | Message::hello(vec![OpId::new(ReplicaId::new(7), want as u64)]), |
| 1515 | ])).len(); |
| 1516 | assert!(largest <= bound + opening, |
| 1517 | "a reply carried {} bytes against a bound of {} and an opening of {}", |
| 1518 | largest, bound, opening); |
| 1519 | Ok(()) |
| 1520 | } |
| 1521 | |
| 1522 | /// A path is split into an account, a repository and a verb, and anything |
| 1523 | /// else is nothing. |
| 1524 | #[test] |
| 1525 | fn a_path_is_read_or_it_is_not() -> Outcome<()> { |
| 1526 | let got = match route("/ore/v1/oxedyne/fe2o3/sync") { |
| 1527 | Some(parts) => parts, |
| 1528 | None => return Err(err!("The ordinary path was not read."; Test, Missing)), |
| 1529 | }; |
| 1530 | assert_eq!(got, (fmt!("oxedyne"), fmt!("fe2o3"), fmt!("sync"))); |
| 1531 | for bad in [ |
| 1532 | "/ore/v1/oxedyne/fe2o3", // no verb |
| 1533 | "/ore/v1/oxedyne/fe2o3/sync/more", // too many parts |
| 1534 | "/ore/v2/oxedyne/fe2o3/sync", // a version this relay does not speak |
| 1535 | "/ore/v1//fe2o3/sync", // an empty account |
| 1536 | "/oxedyne/fe2o3/sync", // no prefix |
| 1537 | "/", |
| 1538 | ] { |
| 1539 | if route(bad).is_some() { |
| 1540 | return Err(err!("The path {:?} was read as a route.", bad; Test, Invalid)); |
| 1541 | } |
| 1542 | } |
| 1543 | Ok(()) |
| 1544 | } |
| 1545 | } |