Oregami
Repositories/oxedyne/ore

oxedyne/ore/relay/src/host.rs

42.3 KiB, 47 runs

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

download · who wrote it · its history

1//! The repositories a relay holds, and the key bindings it carries for them.
2//!
3//! ```text
4//! <data>/<account>/<name>/log/000000.seg the history, in the form it arrived
5//! <data>/<account>/<name>/acl who may read and write it
6//! <data>/<account>/<name>/keys self-certified bindings, carried
7//! <data>/<account>/<name>/veils veil key bindings, carried
8//! <data>/<account>/<name>/wraps content keys, wrapped per replica
9//! <data>/<account>/<name>/lock where the write lock is kept
10//! ```
11//!
12//! A hosted repository is addressed as `<account>/<name>` under the relay's host.
13//! The address is a relay-local label and not a property of the history: the same
14//! history hosted twice has two addresses, exactly as a git repository has two
15//! remotes. Deriving the name from the history was considered and rejected -- a
16//! log may hold several roots, an imported history's root is an accident of the
17//! import, and a rename would then be impossible. The label lives where the
18//! labelling authority is.
19//!
20//! # What is stored
21//!
22//! A bare repository: segments and the small records beside them, no working
23//! copy and no key. Entries are stored in the form they arrived, sealed
24//! preferred, so a signature survives the hop and is served back out intact.
25//! Dedup is the log's own behaviour -- an arriving operation the log holds is
26//! dropped -- so the union never stores twice, whichever replica carried it in.
27//!
28//! There is no history collection: the history is append-only and nothing being
29//! lost is the product's promise. What is collectable is relay bookkeeping, and
30//! the obliteration path the promise's stated exception will need is a design of
31//! its own.
32//!
33//! # The relay carries bindings and cannot mint them
34//!
35//! A binding is the statement "this key belongs to replica *n*", signed by the
36//! very key it binds. The relay stores only bindings that verify, so a fabricated
37//! one is refused at the door, and the worst the relay can do is withhold one --
38//! which marks the affected operations as signed by an unknown key rather than
39//! misattributing them.
40//!
41//! A veil binding cannot certify itself, because the key it names is an X25519
42//! key and X25519 does not sign. It is checked here by both of its links
43//! instead: the statement is signed by an Ed25519 key, and that key is the
44//! subject of a self-certified binding for the same replica. So the relay is no
45//! more able to introduce a reading key than it is a signing key.
46//!
47//! # The relay serves wraps to anybody who may pull, and this is not a leak
48//!
49//! A wrap is a repository's content key encrypted to one replica's veil key. It
50//! is carried by the very party it is keeping the repository from, which reads
51//! as a hole and is not one: without the matching secret, which never leaves the
52//! machine that minted it, a wrap is thirty-two bytes of nothing. This is the
53//! mechanism and not an oversight, and a reviewer who closes it has removed the
54//! only way a second replica ever reads a veiled repository.
55
56use crate::acl::Acl;
57use crate::reading::Readings;
58
59use ore_store::keys::Binding;
60use ore_store::store::Store;
61use ore_store::veilkey::{
62 VeilBinding,
63 Wrap,
64};
65
66use oxedyne_fe2o3_core::prelude::*;
67use oxedyne_fe2o3_jdat::prelude::*;
68use oxedyne_fe2o3_ore::sync::{
69 Message,
70 Parts,
71};
72
73use std::collections::BTreeMap;
74use std::fs;
75use std::path::{
76 Path,
77 PathBuf,
78};
79use std::sync::{
80 Arc,
81 Mutex,
82};
83use std::time::{
84 SystemTime,
85 UNIX_EPOCH,
86};
87
88
89/// Name of the file the carried bindings live in.
90pub const KEYS_FILE: &str = "keys";
91
92/// Name of the file the carried veil key bindings live in.
93pub const VEILS_FILE: &str = "veils";
94
95/// Name of the file the wrapped content keys live in.
96pub const WRAPS_FILE: &str = "wraps";
97
98/// How long a label may be.
99pub const LABEL_LIMIT: usize = 64;
100
101
102/// Checks that a label is one this relay will put in a path.
103///
104/// Letters, digits, hyphen, underscore and full stop, and nothing else. The
105/// restriction is not decoration: the label becomes a directory name, and a
106/// label carrying a separator or a pair of full stops is how a request reaches a
107/// directory nobody meant to host.
108pub fn check_label(what: &str, label: &str)
109 -> Outcome<()>
110{
111 if label.is_empty() || label.len() > LABEL_LIMIT {
112 return Err(err!(
113 "A repository's {} is {} characters, and one is between 1 and {}.",
114 what, label.len(), LABEL_LIMIT;
115 Invalid, Input, Range));
116 }
117 if label == "." || label == ".." {
118 return Err(err!(
119 "A repository's {} may not be {:?}.", what, label;
120 Invalid, Input));
121 }
122 for c in label.chars() {
123 if !(c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == '.') {
124 return Err(err!(
125 "A repository's {} is {:?}, which holds {:?}; a label is letters, digits, \
126 hyphen, underscore and full stop.", what, label, c;
127 Invalid, Input));
128 }
129 }
130 Ok(())
131}
132
133
134/// One hosted repository, opened.
135pub struct Hosted {
136 /// How it is addressed, as `<account>/<name>`.
137 pub label: String,
138 /// Where it lives.
139 pub dir: PathBuf,
140 /// Its segments.
141 pub store: Store,
142 /// Who may read and write it.
143 pub acl: Acl,
144}
145
146impl Hosted {
147
148 /// Returns the path of the carried bindings.
149 pub fn keys_path(&self) -> PathBuf {
150 self.dir.join(KEYS_FILE)
151 }
152
153 /// Reads the bindings the relay carries for this repository.
154 ///
155 /// Only bindings that certify themselves are handed back. One that does not
156 /// is dropped rather than raised: the file is the relay's own and a binding
157 /// that fails its own signature explains nothing to anybody.
158 pub fn bindings(&self)
159 -> Outcome<Vec<Binding>>
160 {
161 let path = self.keys_path();
162 if !path.is_file() {
163 return Ok(Vec::new());
164 }
165 let text = match fs::read_to_string(&path) {
166 Ok(t) => t,
167 Err(e) => return Err(err!(e,
168 "The carried bindings {:?} could not be read.", path;
169 IO, File, Read)),
170 };
171 let dat = match Dat::decode_string(text) {
172 Ok(d) => d,
173 Err(e) => return Err(err!(e,
174 "The carried bindings {:?} are not readable JDAT.", path;
175 Decode, Input)),
176 };
177 let listed = match &dat {
178 Dat::List(l) => l,
179 other => return Err(err!(
180 "The carried bindings {:?} expect a list, got {:?}.", path, other;
181 Decode, Input, Mismatch)),
182 };
183 let mut out = Vec::new();
184 for item in listed {
185 let binding = res!(Binding::from_dat(item));
186 if binding.is_certified() {
187 out.push(binding);
188 }
189 }
190 Ok(out)
191 }
192
193 /// Adds bindings the relay does not carry, and returns how many were new.
194 ///
195 /// A binding that does not verify against the key it binds is dropped, which
196 /// is the whole of the relay's inability to introduce a stranger.
197 pub fn learn(&self, offered: &[Binding])
198 -> Outcome<usize>
199 {
200 let mut held = res!(self.bindings());
201 let mut fresh = 0usize;
202 for binding in offered {
203 if !binding.is_certified() {
204 continue;
205 }
206 if held.iter().any(|b| b.public == binding.public) {
207 continue;
208 }
209 held.push(binding.clone());
210 fresh += 1;
211 }
212 if fresh == 0 {
213 return Ok(0);
214 }
215 held.sort();
216 let listed: Vec<Dat> = held.iter().map(|b| b.to_dat()).collect();
217 let text = res!(Dat::List(listed).jdat_to_lines(" "));
218 let path = self.keys_path();
219 match fs::write(&path, fmt!("{}\n", text)) {
220 Ok(()) => Ok(fresh),
221 Err(e) => Err(err!(e,
222 "The carried bindings {:?} could not be written.", path;
223 IO, File, Write)),
224 }
225 }
226
227 /// Returns the path of the carried veil key bindings.
228 pub fn veils_path(&self) -> PathBuf {
229 self.dir.join(VEILS_FILE)
230 }
231
232 /// Returns the path of the wraps.
233 pub fn wraps_path(&self) -> PathBuf {
234 self.dir.join(WRAPS_FILE)
235 }
236
237 /// Reads the veil key bindings the relay carries for this repository.
238 ///
239 /// Only the ones whose chain holds are handed back, and the chain is checked
240 /// against the signing bindings beside them: an Ed25519 key vouched for this
241 /// veil key, and that key is the subject of a binding that certifies itself
242 /// for the same replica. One that does not chain is dropped rather than
243 /// raised, exactly as an uncertified binding is.
244 pub fn veil_bindings(&self)
245 -> Outcome<Vec<VeilBinding>>
246 {
247 let listed = match res!(carried(&self.veils_path(), "veil key bindings")) {
248 Some(l) => l,
249 None => return Ok(Vec::new()),
250 };
251 let known = res!(self.bindings());
252 let mut out = Vec::new();
253 for item in &listed {
254 let binding = res!(VeilBinding::from_dat(item));
255 if binding.is_chained(&known) {
256 out.push(binding);
257 }
258 }
259 Ok(out)
260 }
261
262 /// Adds veil key bindings the relay does not carry, and returns how many were
263 /// new.
264 pub fn learn_veils(&self, offered: &[VeilBinding])
265 -> Outcome<usize>
266 {
267 let mut held = res!(self.veil_bindings());
268 let known = res!(self.bindings());
269 let mut fresh = 0usize;
270 for binding in offered {
271 if !binding.is_chained(&known) {
272 continue;
273 }
274 if held.iter().any(|b| b.public == binding.public) {
275 continue;
276 }
277 held.push(binding.clone());
278 fresh += 1;
279 }
280 if fresh == 0 {
281 return Ok(0);
282 }
283 held.sort();
284 res!(put(&self.veils_path(), "veil key bindings",
285 held.iter().map(|b| b.to_dat()).collect()));
286 Ok(fresh)
287 }
288
289 /// Reads the wraps the relay carries for this repository.
290 pub fn wraps(&self)
291 -> Outcome<Vec<Wrap>>
292 {
293 let listed = match res!(carried(&self.wraps_path(), "wraps")) {
294 Some(l) => l,
295 None => return Ok(Vec::new()),
296 };
297 let mut out = Vec::new();
298 for item in &listed {
299 out.push(res!(Wrap::from_dat(item)));
300 }
301 Ok(out)
302 }
303
304 /// Puts the wraps that arrived beside the ones already here, and returns how
305 /// many changed.
306 ///
307 /// One wrap per address, and a later wrap replaces the earlier: that is what
308 /// a fresh content key looks like from here, and refusing it would leave a
309 /// relay serving the key to a repository that has stopped using it. The
310 /// relay checks nothing about a wrap, because there is nothing about a wrap
311 /// it can check -- it cannot read one, and the replica it is addressed to
312 /// finds out at the moment it fails to open. What stands between a stranger
313 /// and this file is the push grant, which is where it belongs.
314 pub fn keep_wraps(&self, offered: &[Wrap])
315 -> Outcome<usize>
316 {
317 let mut held = res!(self.wraps());
318 let mut changed = 0usize;
319 for wrap in offered {
320 match held.iter_mut().find(|w| w.to == wrap.to) {
321 Some(held) => {
322 if *held == *wrap {
323 continue;
324 }
325 *held = wrap.clone();
326 },
327 None => held.push(wrap.clone()),
328 }
329 changed += 1;
330 }
331 if changed == 0 {
332 return Ok(0);
333 }
334 held.sort();
335 res!(put(&self.wraps_path(), "wraps",
336 held.iter().map(|w| w.to_dat()).collect()));
337 Ok(changed)
338 }
339}
340
341
342/// Reads one of the small carried lists beside a repository, where there is one.
343///
344/// Absence is `None` and not an empty list, so a caller can tell a repository
345/// nobody has deposited anything for from one whose file has been emptied.
346fn carried(path: &Path, what: &str)
347 -> Outcome<Option<Vec<Dat>>>
348{
349 if !path.is_file() {
350 return Ok(None);
351 }
352 let text = match fs::read_to_string(path) {
353 Ok(t) => t,
354 Err(e) => return Err(err!(e,
355 "The carried {} {:?} could not be read.", what, path;
356 IO, File, Read)),
357 };
358 let dat = match Dat::decode_string(text) {
359 Ok(d) => d,
360 Err(e) => return Err(err!(e,
361 "The carried {} {:?} are not readable JDAT.", what, path;
362 Decode, Input)),
363 };
364 match dat {
365 Dat::List(l) => Ok(Some(l)),
366 other => Err(err!(
367 "The carried {} {:?} expect a list, got {:?}.", what, path, other;
368 Decode, Input, Mismatch)),
369 }
370}
371
372/// Writes one of the small carried lists beside a repository.
373fn put(path: &Path, what: &str, listed: Vec<Dat>)
374 -> Outcome<()>
375{
376 let text = res!(Dat::List(listed).jdat_to_lines(" "));
377 match fs::write(path, fmt!("{}\n", text)) {
378 Ok(()) => Ok(()),
379 Err(e) => Err(err!(e,
380 "The carried {} {:?} could not be written.", what, path;
381 IO, File, Write)),
382 }
383}
384
385
386/// The reply bound a host takes when nothing says otherwise.
387fn ore_relay_reply_default() -> usize {
388 crate::proto::REPLY_BYTES
389}
390
391/// How long a half arrived operation is kept, in seconds.
392///
393/// An operation larger than one request body crosses in pieces, and a push that
394/// dies between two of them leaves the relay holding the pieces that did arrive.
395/// They are worth nothing on their own -- nothing was absorbed, nothing was
396/// written down -- so this is a bound on wasted memory and not a transfer state
397/// anybody resumes from. Five minutes is what a request's timestamp is allowed to
398/// drift by ([`crate::proto::SKEW`]), which is the same judgement about how long
399/// a client might reasonably take.
400pub const PART_KEEP: u64 = 300;
401
402/// How many peers may have a half arrived operation at once.
403///
404/// Each one costs what it has actually received, bounded by
405/// `sync::msg::PART_MAX`. The oldest is dropped when a further peer arrives,
406/// which costs that peer a repetition and costs the relay nothing.
407///
408/// It bounds the *number* of buffers and is kept alongside [`PART_BYTES`]
409/// because that is a different guarantee: a buffer costs its label, its map node
410/// and its bookkeeping whatever it holds, and a weight bound that measures only
411/// arrived bytes cannot see any of that. Ten thousand peers a single piece in
412/// would sit well inside a quarter of a gibibyte and still be ten thousand
413/// entries nobody bounded.
414pub const PART_PEERS: usize = 64;
415
416/// The most a relay will hold towards half arrived operations, across every peer
417/// at once.
418///
419/// The half that was missing until 2026-08-22. [`PART_PEERS`] said how many
420/// buffers there could be and said nothing about their weight, so the only
421/// ceiling on the weight was sixty-four multiplied by `sync::msg::PART_MAX` --
422/// four gibibytes, on a host with 1,962 MB. Reaching it needed thirty-two peers
423/// each pushing a large file and then falling quiet, which is not a forge's
424/// traffic; but a bound nobody can reach is not the same thing as a bound, and
425/// the count bound does not bound the thing that runs out.
426///
427/// A quarter of a gibibyte is what a relay can hold without the rest of the box
428/// noticing. Requests are taken one at a time, so the two figures that are
429/// genuinely concurrent are the resting spool and whatever the request in flight
430/// costs. Serving one clone of a history the size of fe2o3 peaks near 400 MB, so
431/// 400 and 268 against 1,962. Taking a push in costs a replay of the same
432/// history -- 263 MB -- and the operation being completed about three times over,
433/// the bytes and the daticle they decode to and the entry that makes, so two
434/// further `PART_MAX` beside the spool: 263 and 268 and 134, which lands in the
435/// same place. Either way about a third of the host, and the two thirds left are
436/// what everything else on it is living on.
437///
438/// Below `sync::msg::PART_MAX` the figure stops being the peak. The eviction
439/// below will not drop the buffer it is making room for, so one peer alone still
440/// carries a largest operation across and the spool passes the bound while it
441/// does. That is the right answer -- the alternative is an operation that can
442/// never cross at all -- but it means a figure under `PART_MAX` bounds the crowd
443/// and not the peak. Four times `PART_MAX` is four peers each halfway through
444/// the largest operation Ore will put back together, which is more at once than
445/// a forge sees.
446pub const PART_BYTES: usize = 256 << 20;
447
448
449/// One peer's half arrived operation.
450struct Waiting {
451 held: Parts, // what has been put together so far
452 at: u64, // when the last piece arrived, in seconds since the epoch
453}
454
455/// Every repository one relay holds.
456#[derive(Clone, Debug)]
457pub struct Host {
458 pub dir: PathBuf, // the directory the accounts live under
459 pub reply_bytes: usize, // the most one reply body may carry
460 pub post_bytes: usize, // the most one request body may carry
461 pub part_bytes: usize, // the most everything half arrived may come to
462 // What is half arrived, by repository and pushing replica. Nothing here
463 // reaches a disk and nothing here survives a restart, which is what keeps
464 // resume-is-rerun true: the worst a lost buffer costs is that the operation
465 // crosses again.
466 parts: Arc<Mutex<BTreeMap<(String, u64), Waiting>>>,
467 // The log each hosted repository was last read into, so that serving a read
468 // costs the bytes appended since rather than the whole history. Shared, so
469 // that a `Host` cloned per connection is still one relay holding one reading.
470 // Nothing here reaches a disk either, and a restart costs one whole read.
471 readings: Arc<Readings>,
472}
473
474impl std::fmt::Debug for Waiting {
475 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
476 write!(f, "Waiting {{ held: {} bytes, at: {} }}", self.held.held(), self.at)
477 }
478}
479
480impl Host {
481
482 /// Names the host at a directory, without touching it.
483 pub fn at(dir: &Path) -> Self {
484 Self {
485 dir: dir.to_path_buf(),
486 reply_bytes: ore_relay_reply_default(),
487 post_bytes: crate::proto::POST_BYTES,
488 part_bytes: PART_BYTES,
489 parts: Arc::new(Mutex::new(BTreeMap::new())),
490 readings: Arc::new(Readings::new()),
491 }
492 }
493
494 /// Bounds what the readings this relay holds may come to, for an operator with
495 /// less memory than [`crate::reading::READING_BUDGET`] assumes, and for a test
496 /// that would otherwise have to host a history the size of fe2o3 to reach the
497 /// boundary.
498 ///
499 /// Set to nothing it holds nothing, which is this relay as it stood before the
500 /// readings existed: every request reads the store whole. That is the arm to
501 /// reach for on a host where the memory matters more than the processor, and it
502 /// is the arm a measurement compares against.
503 pub fn with_reading_bytes(mut self, bytes: u64) -> Self {
504 self.readings = Arc::new(Readings::to(bytes));
505 self
506 }
507
508 /// The log each hosted repository was last read into.
509 pub fn readings(&self) -> &Readings {
510 &self.readings
511 }
512
513 /// Bounds what one reply may carry, for an operator whose callers have less
514 /// memory than the default assumes, and for a test that would otherwise have
515 /// to push six mebibytes to reach the boundary.
516 ///
517 /// It bounds the caller's body and not this relay's own peak. `serve::sync`
518 /// builds the whole owed turn through `outgoing` and truncates afterwards, so
519 /// lowering this spends the same memory here and hands out less.
520 pub fn with_reply_bytes(mut self, bytes: usize) -> Self {
521 self.reply_bytes = bytes;
522 self
523 }
524
525 /// Bounds what one request body may carry, and it is bound rather than
526 /// merely stated.
527 ///
528 /// The half that was missing until 2026-08-22. A relay published
529 /// [`crate::proto::POST_BYTES`] as a `const` compiled into it, so an operator
530 /// behind a proxy that takes less had no way to say so and no way to find out
531 /// except by watching a client's connection close part way through a body. A
532 /// relay now refuses what it said it would not take, with a sentence naming
533 /// both numbers, which is a better answer than the proxy's silence and is what
534 /// makes the published number worth reading.
535 pub fn with_post_bytes(mut self, bytes: usize) -> Self {
536 self.post_bytes = bytes;
537 self
538 }
539
540 /// Bounds what every half arrived operation may come to at once, for an
541 /// operator with less memory than [`PART_BYTES`] assumes, and for a test that
542 /// would otherwise have to hold a quarter of a gibibyte to reach the boundary.
543 ///
544 /// Set below `sync::msg::PART_MAX` it bounds the crowd rather than the peak:
545 /// the eviction will not drop the buffer it is making room for, so one peer
546 /// alone still carries a largest operation across and the spool passes the
547 /// figure while it does.
548 pub fn with_part_bytes(mut self, bytes: usize) -> Self {
549 self.part_bytes = bytes;
550 self
551 }
552
553 /// What every half arrived operation this relay holds comes to.
554 pub fn spooled(&self)
555 -> Outcome<usize>
556 {
557 let held = lock_mutex!(self.parts);
558 Ok(spooled(&held))
559 }
560
561 /// Folds an arriving run of pieces back into the operations they make.
562 ///
563 /// An operation larger than one request body crosses as a run of
564 /// `Message::Part`, and a run is longer than one request by construction --
565 /// that is the whole point of cutting it up -- so the relay holds what has
566 /// arrived between them. This is the only state a relay keeps across requests
567 /// and it is the smallest that will do the job: one buffer per pushing
568 /// replica, nothing written down, nothing surviving a restart.
569 ///
570 /// **A run that never finishes costs a repetition and nothing else.** The
571 /// operation was never absorbed, so the next `ore sync` offers it again of its
572 /// own accord and begins at piece zero, which throws away what the attempt
573 /// before it left. Resume is still rerun.
574 ///
575 /// **Two bounds hold the spool, and they are not the same guarantee.**
576 /// [`PART_PEERS`] stops one peer holding the other sixty-three out;
577 /// [`PART_BYTES`] stops sixty-four peers holding the host out. Either is
578 /// reached by dropping the oldest buffer, never the one this request is
579 /// filling, and the peer that lost one is told the relay dropped its run
580 /// rather than being left to read a fault as its own.
581 pub fn rejoin(&self, label: &str, replica: u64, arriving: Vec<Message>)
582 -> Outcome<Vec<Message>>
583 {
584 // Nothing to hold and nothing held: the ordinary request, which never
585 // touches the lock.
586 let key = (fmt!("{}", label), replica);
587 let mut held = lock_mutex!(self.parts);
588 let now = res!(now_secs());
589 held.retain(|_, w| now.saturating_sub(w.at) <= PART_KEEP);
590 if !arriving.iter().any(|m| matches!(m, Message::Part { .. })) && !held.contains_key(&key) {
591 return Ok(arriving);
592 }
593 // A run this relay holds nothing towards cannot be continued, and which end
594 // is at fault is not in doubt: a buffer goes when it falls quiet, when a
595 // further peer needs the room, and when the spool reaches its weight, all
596 // three of them this relay's own doing. Said here rather than left to
597 // `Parts`, which is handed a fresh buffer and cannot tell a run this relay
598 // threw away from one that genuinely began in the middle.
599 //
600 // Before any eviction, so that a peer repeating a run this relay has
601 // already dropped cannot cost the other peers theirs.
602 if !held.contains_key(&key) {
603 if let Some(Message::Part { id, seq, total, .. }) =
604 arriving.iter().find(|m| matches!(m, Message::Part { .. }))
605 {
606 if *seq != 0 {
607 return Err(err!(
608 "This relay is holding nothing towards the operation {}, so the \
609 run of pieces carrying it was dropped before this piece {} of {} \
610 arrived. A half arrived run goes when it has been quiet for {} \
611 seconds, when a further peer needs the room, or when everything \
612 half arrived reaches the {} bytes a relay holds at once. Nothing \
613 was absorbed from it and nothing was written down, so run the \
614 command again: it begins at piece zero and completes.",
615 id, seq, total, PART_KEEP, self.part_bytes;
616 Invalid, Input, Order, Missing));
617 }
618 }
619 }
620 if !held.contains_key(&key) && held.len() >= PART_PEERS {
621 // The oldest goes, which costs that peer a repetition. Refusing this
622 // one instead would let whoever got in first keep everybody else out.
623 if let Some(k) = oldest_but(&held, &key) {
624 held.remove(&k);
625 }
626 }
627 // And again by weight, since sixty-four buffers inside their count are
628 // still four gibibytes if nothing weighs them. What this request is about
629 // to add is counted before it is added, so the bound holds at the peak and
630 // not merely once the peak has passed.
631 let coming: usize = arriving.iter()
632 .map(|m| match m {
633 Message::Part { bytes, .. } => bytes.len(),
634 _ => 0,
635 })
636 .sum();
637 while spooled(&held) + coming > self.part_bytes {
638 match oldest_but(&held, &key) {
639 Some(k) => { held.remove(&k); },
640 // Nobody left but this peer's own run, which `sync::msg::PART_MAX`
641 // is what bounds. Dropping it here would be making room by throwing
642 // away the thing the room was for.
643 None => break,
644 }
645 }
646 let waiting = held.entry(key.clone()).or_insert_with(|| Waiting {
647 held: Parts::new(),
648 at: now,
649 });
650 waiting.at = now;
651 let mut out = Vec::with_capacity(arriving.len());
652 for msg in arriving {
653 match waiting.held.absorb(msg) {
654 Ok(Some(m)) => out.push(m),
655 Ok(None) => {},
656 // A run that broke is thrown away whole, so the next attempt starts
657 // from a relay holding nothing rather than from one holding half of
658 // something it cannot name.
659 Err(e) => {
660 held.remove(&key);
661 return Err(e);
662 },
663 }
664 }
665 if !waiting.held.pending() {
666 held.remove(&key);
667 }
668 Ok(out)
669 }
670
671 /// Returns where a repository lives, refusing a label that is not one.
672 pub fn path_of(&self, account: &str, name: &str)
673 -> Outcome<PathBuf>
674 {
675 res!(check_label("account", account));
676 res!(check_label("name", name));
677 Ok(self.dir.join(account).join(name))
678 }
679
680 /// Opens a repository, or says there is none of that name.
681 pub fn open(&self, account: &str, name: &str)
682 -> Outcome<Option<Hosted>>
683 {
684 let dir = res!(self.path_of(account, name));
685 let store = Store::at(&dir);
686 if !store.exists() {
687 return Ok(None);
688 }
689 Ok(Some(Hosted {
690 label: fmt!("{}/{}", account, name),
691 acl: res!(Acl::read(&dir)),
692 dir,
693 store,
694 }))
695 }
696
697 /// Returns every repository actually on this relay's disk, in address order.
698 ///
699 /// The only complete inventory there is. A forge's configuration names a subset
700 /// and its register names another, and a repository made at a shell before the
701 /// register existed is in neither -- while this relay goes on serving it, because
702 /// an address is resolved against the disk and an access list. A directory that
703 /// holds no store is passed over rather than refused: the data root is an
704 /// operator's directory and may hold something that is not a repository.
705 pub fn held(&self)
706 -> Outcome<Vec<(String, String)>>
707 {
708 let mut out = Vec::new();
709 if !self.dir.is_dir() {
710 // Nothing hosted yet, which is not a failure to list.
711 return Ok(out);
712 }
713 for account in res!(fs::read_dir(&self.dir)) {
714 let account = res!(account);
715 if !account.path().is_dir() {
716 continue;
717 }
718 // A name this relay could not have written is a name it does not serve,
719 // so it is not reported as something it holds.
720 let acct = match account.file_name().into_string() {
721 Ok(said) if check_label("account", &said).is_ok() => said,
722 _ => continue,
723 };
724 for repo in res!(fs::read_dir(account.path())) {
725 let repo = res!(repo);
726 let dir = repo.path();
727 if !dir.is_dir() {
728 continue;
729 }
730 let name = match repo.file_name().into_string() {
731 Ok(said) if check_label("name", &said).is_ok() => said,
732 _ => continue,
733 };
734 if Store::at(&dir).exists() {
735 out.push((acct.clone(), name));
736 }
737 }
738 }
739 out.sort();
740 Ok(out)
741 }
742
743 /// Creates a repository with an empty history and the given access list.
744 ///
745 /// Creation is an administrator's act at this rung rather than a self-service
746 /// one, so it happens here and not over the wire.
747 pub fn create(&self, account: &str, name: &str, acl: &Acl)
748 -> Outcome<Hosted>
749 {
750 let dir = res!(self.path_of(account, name));
751 if Store::at(&dir).exists() {
752 return Err(err!(
753 "This relay already holds {}/{}.", account, name;
754 Invalid, Input, Exists));
755 }
756 res!(fs::create_dir_all(&dir));
757 let store = res!(Store::create(&dir, None));
758 res!(acl.write(&dir));
759 Ok(Hosted {
760 label: fmt!("{}/{}", account, name),
761 dir,
762 store,
763 acl: acl.clone(),
764 })
765 }
766}
767
768
769/// What every half arrived operation in the spool comes to.
770fn spooled(held: &BTreeMap<(String, u64), Waiting>) -> usize {
771 held.values().map(|w| w.held.held()).sum()
772}
773
774/// The oldest buffer that is not the one this request is filling.
775///
776/// The exclusion is what stops a request evicting itself: a peer whose own run
777/// is the heaviest thing held would otherwise make room by dropping the run it
778/// was making room for, and would do it again on the next piece, for ever.
779fn oldest_but(held: &BTreeMap<(String, u64), Waiting>, key: &(String, u64))
780 -> Option<(String, u64)>
781{
782 held.iter()
783 .filter(|(k, _)| *k != key)
784 .min_by_key(|(_, w)| w.at)
785 .map(|(k, _)| k.clone())
786}
787
788
789/// Returns the seconds since the epoch.
790fn now_secs()
791 -> Outcome<u64>
792{
793 Ok(res!(SystemTime::now().duration_since(UNIX_EPOCH)).as_secs())
794}
795
796
797#[cfg(test)]
798mod tests {
799 use super::*;
800
801 use oxedyne_fe2o3_ore::id::{
802 OpId,
803 ReplicaId,
804 };
805 use oxedyne_fe2o3_ore::op::{
806 Header,
807 Op,
808 Record,
809 };
810 use oxedyne_fe2o3_ore::segment::Entry;
811
812 /// An entry too large to cross in one message.
813 fn wide(replica: u64, counter: u64, bytes: usize)
814 -> Outcome<Entry>
815 {
816 Ok(Entry::Bare(Record::new(
817 res!(Header::new(
818 OpId::new(ReplicaId::new(replica), counter),
819 vec![OpId::new(ReplicaId::new(replica), counter - 1)],
820 )),
821 Op::Mark {
822 name: fmt!("wide"),
823 body: Some(vec![0x5eu8; bytes]),
824 time: Some(1_755_000_000),
825 },
826 )))
827 }
828
829 /// The pieces of one operation, in the order they must be sent.
830 fn pieces(entry: &Entry, cap: usize)
831 -> Outcome<Vec<Message>>
832 {
833 Message::part(entry, cap)
834 }
835
836 /// A label that could reach outside the data directory is refused, and the
837 /// ordinary ones are not.
838 #[test]
839 fn a_label_that_is_a_path_is_refused() -> Outcome<()> {
840 for good in ["oxedyne", "fe2o3", "a", "some-repo_2.0"] {
841 res!(check_label("name", good));
842 }
843 for bad in ["", "..", ".", "a/b", "../etc", "a b", "a\\b", "a\0b", "café"] {
844 if check_label("name", bad).is_ok() {
845 return Err(err!("The label {:?} was accepted.", bad; Test, Invalid));
846 }
847 }
848 // And one past the limit.
849 let long = "a".repeat(LABEL_LIMIT + 1);
850 assert!(check_label("name", &long).is_err());
851 Ok(())
852 }
853
854 /// A directory that removes itself however the test ends.
855 struct Scratch {
856 path: PathBuf,
857 }
858
859 impl Scratch {
860 fn new(what: &str)
861 -> Outcome<Self>
862 {
863 let stamp = res!(SystemTime::now().duration_since(UNIX_EPOCH));
864 let path = std::env::temp_dir().join(fmt!(
865 "ore_relay_{}_{}_{}", what, std::process::id(), stamp.as_nanos(),
866 ));
867 res!(fs::create_dir_all(&path));
868 Ok(Self { path })
869 }
870 }
871
872 impl Drop for Scratch {
873 fn drop(&mut self) {
874 let _ = fs::remove_dir_all(&self.path);
875 }
876 }
877
878 /// Every repository on the disk is reported, and nothing that is not one is.
879 ///
880 /// The disk is the only complete inventory a relay or a forge has: a
881 /// configuration names a subset, a register names another, and a repository in
882 /// neither is served all the same.
883 #[test]
884 fn what_is_on_the_disk_is_what_is_held() -> Outcome<()> {
885 let scratch = res!(Scratch::new("held"));
886 let host = Host::at(&scratch.path);
887 assert!(res!(host.held()).is_empty(), "an empty data directory held something");
888 // A data root that was never written to at all is not a failure to list.
889 let never = Host::at(&scratch.path.join("never"));
890 assert!(res!(never.held()).is_empty(), "an absent data directory was an error");
891 for (account, name) in [
892 ("someone", "notes"),
893 ("oxedyne", "ore"),
894 ("oxedyne", "fe2o3"),
895 // A name this relay could not have written, so not one it serves.
896 ("oxedyne", "a b"),
897 ] {
898 let dir = scratch.path.join(account).join(name);
899 res!(fs::create_dir_all(&dir));
900 res!(Store::create(&dir, None));
901 }
902 // A directory holding no store, which an operator's data root may.
903 res!(fs::create_dir_all(scratch.path.join("oxedyne").join("empty")));
904 // A file beside the accounts, likewise.
905 res!(fs::write(scratch.path.join("README"), b"not an account"));
906 assert_eq!(res!(host.held()), vec![
907 (fmt!("oxedyne"), fmt!("fe2o3")),
908 (fmt!("oxedyne"), fmt!("ore")),
909 (fmt!("someone"), fmt!("notes")),
910 ], "the inventory is not what the disk holds");
911 Ok(())
912 }
913
914 /// A run of pieces spread over several requests is put back together, and
915 /// nothing downstream of the fold can tell it was ever cut up.
916 ///
917 /// This is the one place a relay keeps state between requests, and it keeps it
918 /// because a run longer than one request body is the whole point: an operation
919 /// of 22 MB against a body limit of 8 MiB cannot arrive in one request however
920 /// it is cut.
921 #[test]
922 fn a_run_spread_over_several_requests_is_put_back_together() -> Outcome<()> {
923 let host = Host::at(Path::new("/nowhere"));
924 let entry = res!(wide(7, 4, 20_000));
925 let run = res!(pieces(&entry, 4_096));
926 assert!(run.len() > 4, "the fixture is {} pieces", run.len());
927
928 // One request per piece, which is the shape a small body limit produces.
929 let mut got = Vec::new();
930 for (at, piece) in run.iter().enumerate() {
931 let out = res!(host.rejoin("acct/repo", 7, vec![piece.clone()]));
932 if at + 1 < run.len() {
933 assert!(out.is_empty(), "request {} answered before the run finished", at);
934 }
935 got.extend(out);
936 }
937 assert_eq!(got.len(), 1, "the run made {} messages", got.len());
938 match &got[0] {
939 Message::Send { entries } => {
940 assert_eq!(entries.len(), 1);
941 assert_eq!(
942 res!(entries[0].to_dat().to_bytes(Vec::new())),
943 res!(entry.to_dat().to_bytes(Vec::new())),
944 "the operation is not the one that was sent, byte for byte",
945 );
946 },
947 other => return Err(err!(
948 "The run made a {} message.", other.name(); Test, Mismatch)),
949 }
950 // A request carrying nothing to do with pieces is handed straight back.
951 let plain = vec![Message::hello(vec![OpId::new(ReplicaId::new(7), 4)]), Message::Done];
952 assert_eq!(res!(host.rejoin("acct/repo", 7, plain.clone())), plain);
953 Ok(())
954 }
955
956 /// A push that died halfway costs a repetition and nothing else.
957 ///
958 /// Resume is rerun, and this is what it means for an operation that crosses in
959 /// pieces: the relay holds what arrived, absorbed none of it and wrote none of
960 /// it down, so the next attempt begins at piece zero, which throws away what
961 /// the attempt before it left. The operation lands once.
962 #[test]
963 fn a_push_that_died_halfway_completes_on_a_rerun() -> Outcome<()> {
964 let host = Host::at(Path::new("/nowhere"));
965 let entry = res!(wide(3, 9, 20_000));
966 let run = res!(pieces(&entry, 4_096));
967
968 // Half of it arrives, and the pusher goes away.
969 for piece in &run[..2] {
970 assert!(res!(host.rejoin("acct/repo", 3, vec![piece.clone()])).is_empty());
971 }
972 // The whole of it arrives again, against a relay still holding the remains.
973 let mut got = Vec::new();
974 for piece in &run {
975 got.extend(res!(host.rejoin("acct/repo", 3, vec![piece.clone()])));
976 }
977 assert_eq!(got.len(), 1, "the rerun made {} operations, not one", got.len());
978 match &got[0] {
979 Message::Send { entries } => assert_eq!(
980 res!(entries[0].to_dat().to_bytes(Vec::new())),
981 res!(entry.to_dat().to_bytes(Vec::new())),
982 "the rerun did not put the operation back as it was",
983 ),
984 other => return Err(err!(
985 "The rerun made a {} message.", other.name(); Test, Mismatch)),
986 }
987 Ok(())
988 }
989
990 /// Two peers pushing at once do not interleave into one another's operation.
991 ///
992 /// The buffer is keyed by the replica the request signed itself with, so two
993 /// runs crossing at the same time are two runs. Keyed by the repository alone
994 /// they would be one, and each would refuse the other's pieces as out of
995 /// order -- neither push could ever complete while the other was running.
996 #[test]
997 fn two_peers_pushing_at_once_do_not_interleave() -> Outcome<()> {
998 let host = Host::at(Path::new("/nowhere"));
999 let one = res!(wide(1, 2, 12_000));
1000 let two = res!(wide(2, 2, 12_000));
1001 let run_one = res!(pieces(&one, 4_096));
1002 let run_two = res!(pieces(&two, 4_096));
1003 assert_eq!(run_one.len(), run_two.len());
1004
1005 let mut from_one = Vec::new();
1006 let mut from_two = Vec::new();
1007 for at in 0..run_one.len() {
1008 from_one.extend(res!(host.rejoin("acct/repo", 1, vec![run_one[at].clone()])));
1009 from_two.extend(res!(host.rejoin("acct/repo", 2, vec![run_two[at].clone()])));
1010 }
1011 assert_eq!(from_one.len(), 1, "the first peer's push did not complete");
1012 assert_eq!(from_two.len(), 1, "the second peer's push did not complete");
1013 assert_eq!(from_one[0].entries().len(), 1);
1014 assert_eq!(res!(from_one[0].entries()[0].id()), res!(one.id()));
1015 assert_eq!(res!(from_two[0].entries()[0].id()), res!(two.id()));
1016 Ok(())
1017 }
1018
1019 /// A run that breaks is thrown away whole, so the next attempt begins against
1020 /// a relay holding nothing.
1021 #[test]
1022 fn a_broken_run_leaves_the_relay_holding_nothing() -> Outcome<()> {
1023 let host = Host::at(Path::new("/nowhere"));
1024 let entry = res!(wide(5, 2, 12_000));
1025 let run = res!(pieces(&entry, 4_096));
1026
1027 res!(host.rejoin("acct/repo", 5, vec![run[0].clone()]));
1028 // A piece skipped.
1029 assert!(host.rejoin("acct/repo", 5, vec![run[2].clone()]).is_err());
1030 // And what was held went with it: a piece one now has nothing before it.
1031 assert!(host.rejoin("acct/repo", 5, vec![run[1].clone()]).is_err(),
1032 "the relay was still holding half a run");
1033 // Beginning again works.
1034 let mut got = Vec::new();
1035 for piece in &run {
1036 got.extend(res!(host.rejoin("acct/repo", 5, vec![piece.clone()])));
1037 }
1038 assert_eq!(got.len(), 1);
1039 Ok(())
1040 }
1041
1042 /// A relay holding half arrived operations for too many peers drops the
1043 /// oldest rather than refusing the newcomer.
1044 ///
1045 /// Refusing would let whoever got in first keep everybody else out, which is
1046 /// a denial of service a push grant should not buy.
1047 #[test]
1048 fn a_crowd_of_half_pushes_drops_the_oldest() -> Outcome<()> {
1049 let host = Host::at(Path::new("/nowhere"));
1050 let entry = res!(wide(1, 2, 12_000));
1051 let run = res!(pieces(&entry, 4_096));
1052 assert!(run.len() > 2);
1053
1054 for replica in 0..(PART_PEERS as u64 + 1) {
1055 res!(host.rejoin("acct/repo", replica, vec![run[0].clone()]));
1056 }
1057 // The first peer's buffer went, so its second piece has nothing before it.
1058 assert!(host.rejoin("acct/repo", 0, vec![run[1].clone()]).is_err(),
1059 "a crowded relay kept the oldest buffer and refused the newcomer");
1060 // The last one in is still there.
1061 res!(host.rejoin("acct/repo", PART_PEERS as u64, vec![run[1].clone()]));
1062 Ok(())
1063 }
1064
1065 /// How many bytes one piece of a run carries.
1066 fn weight(piece: &Message)
1067 -> Outcome<usize>
1068 {
1069 match piece {
1070 Message::Part { bytes, .. } => Ok(bytes.len()),
1071 other => Err(err!(
1072 "A run is made of pieces and this one is a {}.", other.name();
1073 Test, Mismatch)),
1074 }
1075 }
1076
1077 /// A spool that reaches its weight drops the oldest, however few peers are
1078 /// holding it.
1079 ///
1080 /// The peer bound cannot do this: three peers is three peers whether each is
1081 /// one piece in or halfway through the largest operation Ore will put back
1082 /// together, and it is the second of those that empties a relay's memory.
1083 #[test]
1084 fn a_heavy_spool_drops_the_oldest() -> Outcome<()> {
1085 let entry = res!(wide(1, 2, 20_000));
1086 let run = res!(pieces(&entry, 4_096));
1087 assert!(run.len() > 3, "the fixture is {} pieces", run.len());
1088 // Room for two pieces and half of a third, so a third peer one piece in
1089 // cannot be held beside the other two.
1090 let one = res!(weight(&run[0]));
1091 let host = Host::at(Path::new("/nowhere")).with_part_bytes(one * 2 + one / 2);
1092
1093 for replica in 1..=3u64 {
1094 res!(host.rejoin("acct/repo", replica, vec![run[0].clone()]));
1095 }
1096 assert!(res!(host.spooled()) <= one * 2 + one / 2,
1097 "the spool holds {} bytes against a bound of {}",
1098 res!(host.spooled()), one * 2 + one / 2);
1099 // The first peer's buffer went, three peers being well inside PART_PEERS.
1100 assert!(host.rejoin("acct/repo", 1, vec![run[1].clone()]).is_err(),
1101 "a spool at its weight kept the oldest buffer");
1102 // The last one in is still there, and carries on.
1103 res!(host.rejoin("acct/repo", 3, vec![run[1].clone()]));
1104 assert!(res!(host.spooled()) <= one * 2 + one / 2,
1105 "the spool holds {} bytes after the newcomer carried on",
1106 res!(host.spooled()));
1107 Ok(())
1108 }
1109
1110 /// One peer alone is never evicted to make room for itself, so an operation
1111 /// that fills the spool on its own still crosses.
1112 #[test]
1113 fn one_peer_filling_the_spool_still_completes() -> Outcome<()> {
1114 let entry = res!(wide(4, 6, 20_000));
1115 let run = res!(pieces(&entry, 4_096));
1116 // A bound one piece can meet and the whole run cannot, which is the shape
1117 // of a relay set below `sync::msg::PART_MAX`.
1118 let host = Host::at(Path::new("/nowhere"))
1119 .with_part_bytes(res!(weight(&run[0])));
1120
1121 let mut got = Vec::new();
1122 for piece in &run {
1123 got.extend(res!(host.rejoin("acct/repo", 4, vec![piece.clone()])));
1124 }
1125 assert_eq!(got.len(), 1, "the run made {} operations, not one", got.len());
1126 Ok(())
1127 }
1128
1129 /// A peer whose buffer the relay dropped is told the relay dropped it, and
1130 /// told to run the command again.
1131 ///
1132 /// It was told instead that its piece "arrived with nothing before it", which
1133 /// is true of the bytes and wrong about the reader: nothing the sender did was
1134 /// broken, the remedy is a rerun, and the sentence named neither.
1135 #[test]
1136 fn an_evicted_peer_is_told_to_run_the_command_again() -> Outcome<()> {
1137 let entry = res!(wide(1, 2, 20_000));
1138 let run = res!(pieces(&entry, 4_096));
1139 let one = res!(weight(&run[0]));
1140 let host = Host::at(Path::new("/nowhere")).with_part_bytes(one * 2 + one / 2);
1141
1142 for replica in 1..=3u64 {
1143 res!(host.rejoin("acct/repo", replica, vec![run[0].clone()]));
1144 }
1145 let said = match host.rejoin("acct/repo", 1, vec![run[1].clone()]) {
1146 Ok(_) => return Err(err!(
1147 "An evicted peer's next piece was taken."; Test, Invalid)),
1148 Err(e) => fmt!("{}", e.plain()),
1149 };
1150 assert!(said.contains("This relay is holding nothing"),
1151 "the peer was told: {}", said);
1152 assert!(said.contains("run the command again"),
1153 "the peer was told: {}", said);
1154 assert!(!said.contains("arrived with nothing before it"),
1155 "the peer was still told its own send was broken: {}", said);
1156 Ok(())
1157 }
1158
1159 /// A peer repeating a run this relay has already dropped does not cost the
1160 /// other peers theirs.
1161 ///
1162 /// The refusal comes before the eviction for exactly this: a peer that kept
1163 /// sending piece three of a run nobody holds would otherwise evict one buffer
1164 /// per request, and could empty the spool of everybody else's work without
1165 /// ever completing anything of its own.
1166 #[test]
1167 fn a_peer_repeating_a_dropped_run_evicts_nobody() -> Outcome<()> {
1168 let entry = res!(wide(1, 2, 20_000));
1169 let run = res!(pieces(&entry, 4_096));
1170 let one = res!(weight(&run[0]));
1171 let host = Host::at(Path::new("/nowhere")).with_part_bytes(one * 2 + one / 2);
1172
1173 // Two peers, each a piece in, filling the spool between them.
1174 for replica in 1..=2u64 {
1175 res!(host.rejoin("acct/repo", replica, vec![run[0].clone()]));
1176 }
1177 let before = res!(host.spooled());
1178 // A third says nothing but the middle of a run nobody holds, over and over.
1179 for _ in 0..8 {
1180 assert!(host.rejoin("acct/repo", 9, vec![run[2].clone()]).is_err());
1181 }
1182 assert_eq!(res!(host.spooled()), before,
1183 "a peer sending the middle of a dropped run took somebody else's buffer");
1184 // And the peer that was already halfway is still able to finish, which is
1185 // what the spool was being kept for.
1186 let mut got = Vec::new();
1187 for piece in &run[1..] {
1188 got.extend(res!(host.rejoin("acct/repo", 1, vec![piece.clone()])));
1189 }
1190 assert_eq!(got.len(), 1, "the peer that was already halfway could not finish");
1191 Ok(())
1192 }
1193}