Oregami
Repositories/oxedyne/ore

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
53use crate::acl::Role;
54use crate::host::{
55 Host,
56 Hosted,
57};
58use crate::proto::{
59 self,
60 Presented,
61};
62use crate::reading::{
63 Carried,
64 Took,
65 Whence,
66};
67
68use ore_store::keys::Binding;
69use ore_store::veilkey::{
70 VeilBinding,
71 Wrap,
72};
73use ore_store::store::{
74 arrived,
75 Consumed,
76 Keep,
77 Replayed,
78 seal_one,
79 Verify,
80};
81use ore_store::veil::{
82 placehold,
83 restore_one,
84};
85
86use oxedyne_fe2o3_core::prelude::*;
87use oxedyne_fe2o3_jdat::prelude::*;
88use oxedyne_fe2o3_ore::id::OpId;
89use oxedyne_fe2o3_ore::log::OpLog;
90use oxedyne_fe2o3_ore::segment::{
91 self,
92 Entry,
93 Veiled,
94};
95use oxedyne_fe2o3_ore::sync::{
96 msg,
97 Message,
98 Mode,
99 Session,
100};
101
102use std::collections::{
103 BTreeMap,
104 BTreeSet,
105};
106use std::net::SocketAddr;
107use std::pin::Pin;
108use std::sync::Arc;
109
110use oxedyne_fe2o3_net::constant;
111use 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
128use tokio::net::{
129 TcpListener,
130 TcpStream,
131};
132
133
134/// The content type frames are carried under.
135pub const FRAMES_TYPE: &str = "application/octet-stream";
136/// The content type a daticle answer is carried under.
137pub 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.
145pub const HEADER_ABSORBED: &str = "x-ore-absorbed";
146
147
148/// What a request asked for.
149pub 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.
162pub 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
173impl 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.
210pub 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.
220fn 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.
288fn 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.
303fn 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.
340fn 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.
421fn 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.
510fn 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.
530fn 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.
771fn 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.
838fn 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.
909fn 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.
966fn 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.
983fn 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.
994pub 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.
1035async 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.
1102fn 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.
1116fn 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)]
1132mod 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}