Oregami
Repositories/oxedyne/ore

oxedyne/ore/relay/src/reading.rs

20.4 KiB, 1 run

created by r2848102244:1274, 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 readings this relay holds between requests, so that serving a read costs
2//! the bytes since the last one rather than the whole log.
3//!
4//! # Why this exists, in numbers
5//!
6//! A relay replayed its whole store on every request. Measured against a copy of
7//! fe2o3's history -- 35,446 operations in 89,000,226 bytes of segments, over 44
8//! files -- a clone under the deployed [`crate::proto::REPLY_BYTES`] is thirty-two
9//! requests, and every one of them read the store whole. So did the `keys`
10//! endpoint every `ore sync` calls first. Thirty-three whole reads at 89,000,555
11//! bytes of `rchar` each: 2,937 MB of reading for 89 MB of history, the same
12//! figure to the byte across five runs, and 16.3 s of the relay's processor
13//! against a clone that took 19.4 s of wall clock over the loopback.
14//!
15//! The same store, read once and held: 1.7 MB of `rchar` and 2.4 s. What is left
16//! is the growing segment's prefix, read again on each request and hashed to prove
17//! the bytes already read have not moved -- 50,636 bytes here, and never more than
18//! [`ore_store::store::SEGMENT_LIMIT`].
19//!
20//! A history is appended to and never rewritten, so the operations a request read
21//! are the operations the next request would read. [`ore_store::store::Store::replay_since`]
22//! is the read that takes up where the last one stopped; this is what holds the
23//! reading it takes up from.
24//!
25//! # What is held, and what stops it eating the host
26//!
27//! One [`Replayed`] per repository -- the log and the envelopes, which is what
28//! [`ore_store::store::Keep::Envelopes`] asks for and what a session sends from.
29//! Nothing is rendered here and nothing is verified here: a relay checks no
30//! signature, so what it holds is smaller than a forge's render of the same
31//! history and is the whole of what serving a sync needs.
32//!
33//! The readings held are bounded by [`READING_BUDGET`] against an estimate of
34//! what each costs. When one more takes the total past the bound the least
35//! recently used go until it fits, and **a reading that alone would not fit is
36//! used and not kept** -- so a repository too large for the bound leaves the relay
37//! as it found it rather than emptying the cache on the way past.
38//!
39//! # A lost reading is reported, never paid for in silence
40//!
41//! A `Readings` holds the entries and the paths they are filed under in two
42//! records rather than one, so that losing a reading can be seen. A request
43//! handed nothing for a repository this relay has read is [`Whence::Lost`], and it
44//! writes a line as it happens.
45//!
46//! That is written because the same change shipped inert once. On 2026-08-22 the
47//! forge's resumable read was deployed with its cache emptied on every push, the
48//! fallback silent, and reported as met; it was caught by a lane that measured
49//! four arms and found no difference between them. Wall time on a busy host
50//! cannot tell a resumed read from a whole one, so [`Took`] counts the bytes each
51//! read skipped and the bytes it decoded, and the question is answered by a
52//! number.
53
54use ore_store::store::{
55 Consumed,
56 Replayed,
57};
58
59use oxedyne_fe2o3_core::prelude::*;
60
61use std::collections::BTreeSet;
62use std::path::{
63 Path,
64 PathBuf,
65};
66use std::sync::{
67 Arc,
68 RwLock,
69};
70
71
72/// What a reading is estimated to cost in memory, per byte of segment.
73///
74/// Measured on the fe2o3 copy above: a relay that has read its 89,000,226 bytes
75/// of segments rests 195 MB above where it started, which is 2.2 times them. A
76/// log holds every operation decoded, and [`ore_store::store::Keep::Envelopes`]
77/// holds the encoded record a second time so that provenance can be handed on
78/// rather than stopping here.
79///
80/// **Rounded up on purpose.** An estimate that under-counts turns a bound into a
81/// number that does not bound, and the measurement is one format and one history:
82/// a repository of many small operations carries more per byte of segment than
83/// one of few large ones. The figure has to be right about the ratio between two
84/// repositories rather than about either of them.
85const READING_WEIGHT: u64 = 3;
86
87/// What the readings held at once may come to, estimated.
88///
89/// The arithmetic, on the host the deployed relay serves from: 1,962 MB, of which
90/// the forge in the same process budgets 400 MB of renders and serving one clone
91/// of a history the size of fe2o3 peaks near 400 MB. 320 MB of readings holds one
92/// history of that size -- the one that costs 15.6 s of relay processor to serve
93/// without it -- and leaves the box a third of itself with everything warm.
94///
95/// It costs less than it looks. A relay that read the store on every request was
96/// resident at 317 MB after one clone against 286 MB for one that holds the
97/// reading: the whole read allocated and freed a log on every request and the
98/// allocator kept the pages, so what this makes resident is largely what was
99/// resident anyway, and it is a reading rather than fragmentation.
100///
101/// The bound bites when a relay is given a second repository that size. What it
102/// does then is described above: the least recently used goes, and whoever asks
103/// for it next pays the whole read once, which is what every request paid before
104/// this existed.
105///
106/// A constant with a setter rather than a configuration field, because the number
107/// that matters is a property of the host and the operator already chooses the
108/// host. See [`crate::host::Host::with_reading_bytes`], which is what
109/// `ore-relay serve --reading-bytes` and the tests reach.
110pub const READING_BUDGET: u64 = 320 << 20;
111
112
113/// What the relay was holding for a repository when a request began.
114///
115/// A request is handed one of these rather than an `Option`, because "nothing"
116/// has two meanings and only one of them is ordinary. [`Carried::Fresh`] is the
117/// first request against a repository, which was always going to read the whole
118/// log. [`Carried::Lost`] is a reading this relay took and then let go without
119/// meaning to, which costs every request after it the whole log again.
120pub enum Carried {
121 From(Arc<Replayed>, Consumed), // the reading the last request left, and how far it got
122 Fresh, // nothing, this relay never having read the repository
123 Lost, // nothing, although it has read it: a reading was dropped
124}
125
126
127/// Where a read of a repository started, and why it started there.
128///
129/// [`Whence::On`] and [`Whence::First`] are the two ordinary answers. Each of the
130/// other three is a whole read of a store this relay had already read, which on a
131/// history of any size is most of what a puller waits for, and each writes a line
132/// as it happens.
133#[derive(Clone, Copy, Debug, PartialEq, Eq)]
134pub enum Whence {
135 On, // where the last read of this repository stopped
136 First, // at the first byte, this relay never having read it
137 Lost, // at the first byte, a reading having been taken and then dropped
138 Unresumable, // at the first byte, the reading held not being one to take up from
139 Refused, // at the first byte, the store refusing to be taken up over
140}
141
142impl Whence {
143
144 /// Did this read start at the first byte of the store?
145 pub fn is_whole(&self) -> bool {
146 !matches!(self, Self::On)
147 }
148}
149
150
151/// What one read of a repository had to do, in the bytes of its segments.
152///
153/// The counter a measurement asks instead of a stopwatch. A resumed read of the
154/// fe2o3 copy skips eighty-nine million bytes and decodes none; a whole one
155/// decodes all eighty-nine million and skips none. Wall time on a busy host
156/// cannot tell those apart and this can.
157///
158/// It counts what was *decoded*, not what the kernel handed over: the growing
159/// segment's prefix is read again on every resumed read, to hash it and prove the
160/// bytes have not moved, and those bytes appear here as skipped. The physical
161/// figure is the process's own `rchar`, which is where a measurement of this
162/// belongs anyway.
163#[derive(Clone, Copy, Debug, PartialEq, Eq)]
164pub struct Took {
165 pub whence: Whence, // where the read started
166 pub skipped: u64, // bytes an earlier read had taken, that this one did not decode
167 pub read: u64, // bytes this one decoded
168}
169
170
171/// One repository's reading, and how far into its segments it got.
172struct Held {
173 at: PathBuf, // where the repository lives, which is what a reading is of
174 weight: u64, // what the reading is estimated to cost in memory
175 cursor: Consumed, // how far into the segments the reading got
176 took: Took, // what the read that left it here had to do
177 read: Arc<Replayed>, // the log and the envelopes themselves
178}
179
180
181/// What holding one more reading means for the ones already held.
182enum Room {
183 Keep(usize), // hold it, once this many of the least recently used have gone
184 Pass, // use it and hold nothing
185}
186
187/// Says whether one more reading can be held, and what has to go for it.
188///
189/// `held` is what is held already, least recently used first, and `adding` is the
190/// estimate of the one being offered. A reading that would not fit **on its own**
191/// is passed over rather than admitted: admitting it would empty the cache and
192/// then leave the relay holding the one thing it cannot afford, which is worse
193/// than holding nothing and is how a single oversized repository takes a host
194/// down.
195fn room(held: &[u64], adding: u64, budget: u64) -> Room {
196 if adding > budget {
197 return Room::Pass;
198 }
199 let mut total: u64 = held.iter().sum::<u64>().saturating_add(adding);
200 let mut go = 0;
201 while total > budget && go < held.len() {
202 total -= held[go];
203 go += 1;
204 }
205 Room::Keep(go)
206}
207
208
209/// The readings this relay is holding.
210///
211/// # Two records of one thing, on purpose
212///
213/// `held` carries the readings and `seen` carries only the paths they are filed
214/// under, and the second exists so that losing one can be *seen*. A path is added
215/// to `seen` when an entry is kept and removed by every path that does not keep
216/// one, so the two never disagree unless something is wrong. A request handed no
217/// reading for a repository `seen` names is one this relay took and dropped, and
218/// it is reported as [`Whence::Lost`] rather than paid for in silence.
219///
220/// # Held in a list because the order is the policy
221///
222/// Least recently used first, which is what eviction takes. A relay holding
223/// enough repositories for the scan to matter is a relay whose readings stopped
224/// fitting long before.
225pub struct Readings {
226 held: RwLock<Vec<Held>>, // the readings, least recently used first
227 seen: RwLock<BTreeSet<PathBuf>>, // what `held` has an entry for, and nothing else
228 budget: u64, // what they may come to, estimated
229}
230
231impl Readings {
232
233 /// A cache holding [`READING_BUDGET`] worth of readings.
234 pub fn new() -> Self {
235 Self::to(READING_BUDGET)
236 }
237
238 /// A cache holding the given estimate worth of readings.
239 ///
240 /// The bound is a parameter for the tests, which have to be able to ask what
241 /// happens when a reading will not fit without building a repository the size
242 /// of fe2o3 to ask it with.
243 pub fn to(budget: u64) -> Self {
244 Self {
245 held: RwLock::new(Vec::new()),
246 seen: RwLock::new(BTreeSet::new()),
247 budget,
248 }
249 }
250
251 /// Takes what is held for a repository out of the cache and hands it over, so
252 /// that the request about to read extends it in place rather than copying it.
253 ///
254 /// **It is taken out and not borrowed.** What the request does to a log is
255 /// absorb into it, and what it then writes to the segments is the difference;
256 /// a reading left in the cache while that happened would be one this relay was
257 /// handing to somebody else mid-absorption. The request puts it back with
258 /// [`Readings::keep`] when it has a cursor over the bytes it wrote, and every
259 /// other way out of a request calls [`Readings::unsee`] instead.
260 ///
261 /// Where nothing is held, this says which nothing it is: a repository this
262 /// relay has never read, or one it has read and is no longer holding a reading
263 /// of. The second is a defect and the request is told so -- see [`Carried`].
264 pub fn take(&self, at: &Path)
265 -> Outcome<Carried>
266 {
267 let got = {
268 let mut held = lock_write!(self.held,
269 "The lock over the readings this relay holds is poisoned, which means an \
270 earlier request stopped part way through one.");
271 held.iter().position(|one| one.at == at).map(|n| held.remove(n))
272 };
273 if let Some(one) = got {
274 return Ok(Carried::From(one.read, one.cursor));
275 }
276 let seen = lock_read!(self.seen,
277 "The lock over the repositories this relay holds a reading of is poisoned, \
278 which means an earlier request stopped part way through one.");
279 match seen.contains(at) {
280 true => Ok(Carried::Lost),
281 false => Ok(Carried::Fresh),
282 }
283 }
284
285 /// Holds a reading, dropping the least recently used until the total fits.
286 ///
287 /// `stored` is what the segments came to when the reading was taken, which is
288 /// what the estimate is made from.
289 pub fn keep(
290 &self,
291 at: &Path,
292 stored: u64,
293 cursor: Consumed,
294 took: Took,
295 read: Arc<Replayed>,
296 )
297 -> Outcome<()>
298 {
299 let weight = stored.saturating_mul(READING_WEIGHT);
300 let mut held = lock_write!(self.held,
301 "The lock over the readings this relay holds is poisoned, which means an \
302 earlier request stopped part way through one.");
303 // Nothing should be here: the request putting a reading back is the one that
304 // took the last one out. Where something is, it was read from segments this
305 // one has since read past, and holding both would be holding two readings of
306 // one repository under one bound.
307 held.retain(|one| one.at != at);
308 let weights: Vec<u64> = held.iter().map(|one| one.weight).collect();
309 let mut gone: Vec<PathBuf> = Vec::new();
310 let mut kept = false;
311 match room(&weights, weight, self.budget) {
312 Room::Pass => gone.push(at.to_path_buf()),
313 Room::Keep(go) => {
314 gone.extend(held.drain(..go).map(|one| one.at));
315 held.push(Held {
316 at: at.to_path_buf(),
317 weight,
318 cursor,
319 took,
320 read,
321 });
322 kept = true;
323 },
324 }
325 drop(held);
326 // The two must agree: a repository with an entry is one a reading is held
327 // for, and one without is not. Evicting a reading is a decision this cache
328 // took rather than a reading it mislaid, so nothing is reported when the
329 // next request for it starts afresh.
330 let mut seen = lock_write!(self.seen,
331 "The lock over the repositories this relay holds a reading of is poisoned, \
332 which means an earlier request stopped part way through one.");
333 for one in gone {
334 seen.remove(&one);
335 }
336 if kept {
337 seen.insert(at.to_path_buf());
338 }
339 Ok(())
340 }
341
342 /// Forgets that a reading exists for a repository, because none does.
343 ///
344 /// Every way out of a request that does not put a reading back calls this, so
345 /// that the next request against the repository is a first read and not a lost
346 /// one. A request that absorbed operations it could not then account for on the
347 /// disk is the case this exists for: what it holds is a log the segments do not
348 /// answer to, and the only safe thing to do with it is let it go and say so.
349 pub fn unsee(&self, at: &Path)
350 -> Outcome<()>
351 {
352 lock_write!(self.seen,
353 "The lock over the repositories this relay holds a reading of is poisoned, \
354 which means an earlier request stopped part way through one.")
355 .remove(at);
356 Ok(())
357 }
358
359 /// What the last read of a repository had to do, where a reading is held.
360 ///
361 /// For the measurements and for the tests. A relay that published this would be
362 /// telling a stranger how much of the host's work their request had just made.
363 pub fn took(&self, at: &Path)
364 -> Outcome<Option<Took>>
365 {
366 let held = lock_read!(self.held,
367 "The lock over the readings this relay holds is poisoned, which means an \
368 earlier request stopped part way through one.");
369 Ok(held.iter().find(|one| one.at == at).map(|one| one.took))
370 }
371
372 /// How many readings are held, and what they are estimated to come to.
373 pub fn holding(&self)
374 -> Outcome<(usize, u64)>
375 {
376 let held = lock_read!(self.held,
377 "The lock over the readings this relay holds is poisoned, which means an \
378 earlier request stopped part way through one.");
379 Ok((held.len(), held.iter().map(|one| one.weight).sum()))
380 }
381}
382
383impl Default for Readings {
384 fn default() -> Self {
385 Self::new()
386 }
387}
388
389impl std::fmt::Debug for Readings {
390 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
391 match self.held.read() {
392 Ok(held) => write!(f, "Readings {{ held: {}, weight: {} }}",
393 held.len(), held.iter().map(|one| one.weight).sum::<u64>()),
394 Err(_) => write!(f, "Readings {{ poisoned }}"),
395 }
396 }
397}
398
399
400#[cfg(test)]
401mod tests {
402 use super::*;
403
404 /// A reading of nothing, which every test here files under a path of its own.
405 fn nothing() -> Arc<Replayed> {
406 Arc::new(Replayed::new())
407 }
408
409 fn plainly() -> Took {
410 Took { whence: Whence::First, skipped: 0, read: 0 }
411 }
412
413 /// A reading put back is handed to the next request, and taking it out leaves
414 /// the cache holding none.
415 #[test]
416 fn a_reading_put_back_is_handed_to_the_next_request() -> Outcome<()> {
417 let held = Readings::new();
418 let at = Path::new("/nowhere/one");
419 match res!(held.take(at)) {
420 Carried::Fresh => (),
421 _ => return Err(err!(
422 "A repository this cache has never seen was not fresh."; Test, Invalid)),
423 }
424 res!(held.keep(at, 1_000, Consumed::new(), plainly(), nothing()));
425 assert_eq!(res!(held.holding()), (1, 3_000), "the estimate is three times the segments");
426 match res!(held.take(at)) {
427 Carried::From(..) => (),
428 _ => return Err(err!(
429 "The reading that was put back was not handed over."; Test, Missing)),
430 }
431 assert_eq!(res!(held.holding()), (0, 0), "taking a reading left it in the cache");
432 Ok(())
433 }
434
435 /// A reading this cache took and did not put back is reported as lost, and one
436 /// it let go of on purpose is not.
437 ///
438 /// The whole of the fault that shipped on 2026-08-22: the fallback to a whole
439 /// read was silent, so nothing anywhere said that the incremental path was
440 /// never being taken. Proved red by having [`Readings::take`] answer
441 /// `Carried::Fresh` wherever it holds nothing, which makes both halves of this
442 /// pass as a first read.
443 #[test]
444 fn a_reading_that_went_missing_is_not_a_first_read() -> Outcome<()> {
445 let held = Readings::new();
446 let at = Path::new("/nowhere/two");
447 res!(held.keep(at, 1_000, Consumed::new(), plainly(), nothing()));
448 // Taken out and not put back, which is what every failed request does.
449 let _ = res!(held.take(at));
450 match res!(held.take(at)) {
451 Carried::Lost => (),
452 Carried::Fresh => return Err(err!(
453 "A reading this cache dropped was reported as a first read, which is \
454 the fault this distinction exists to catch."; Test, Invalid)),
455 Carried::From(..) => return Err(err!(
456 "A reading that was never put back was handed over."; Test, Invalid)),
457 }
458 // And once the cache has said it is holding none, the next is a first read.
459 res!(held.unsee(at));
460 match res!(held.take(at)) {
461 Carried::Fresh => (),
462 _ => return Err(err!(
463 "A reading this cache let go of on purpose was still reported lost.";
464 Test, Invalid)),
465 }
466 Ok(())
467 }
468
469 /// The least recently used go when a further reading will not fit, and one
470 /// that would not fit on its own is not held at all.
471 ///
472 /// The second is what stops one oversized repository emptying the cache and
473 /// then taking the host down with the one thing it cannot afford. Proved red
474 /// by dropping the `adding > budget` arm from [`room`], which then evicts
475 /// everything and holds the oversized reading.
476 #[test]
477 fn a_reading_that_will_not_fit_alone_is_not_held() -> Outcome<()> {
478 let held = Readings::to(9_000);
479 let one = Path::new("/nowhere/a");
480 let two = Path::new("/nowhere/b");
481 let big = Path::new("/nowhere/big");
482 res!(held.keep(one, 1_000, Consumed::new(), plainly(), nothing()));
483 res!(held.keep(two, 2_000, Consumed::new(), plainly(), nothing()));
484 assert_eq!(res!(held.holding()), (2, 9_000));
485 // One that fits only once the oldest has gone.
486 let three = Path::new("/nowhere/c");
487 res!(held.keep(three, 1_000, Consumed::new(), plainly(), nothing()));
488 assert_eq!(res!(held.holding()), (2, 9_000), "the least recently used did not go");
489 match res!(held.take(one)) {
490 Carried::Fresh => (),
491 _ => return Err(err!(
492 "The evicted reading was still held, or was reported as lost when the \
493 cache let it go on purpose."; Test, Invalid)),
494 }
495 // And one that will not fit however much goes.
496 res!(held.keep(big, 100_000, Consumed::new(), plainly(), nothing()));
497 let (count, weight) = res!(held.holding());
498 assert!(weight <= 9_000, "the cache holds {} against a bound of 9,000", weight);
499 assert!(count > 0, "the cache emptied itself for a reading it then did not hold");
500 match res!(held.take(big)) {
501 Carried::Fresh => (),
502 _ => return Err(err!(
503 "A reading too large to hold was held, or was reported lost."; Test, Invalid)),
504 }
505 Ok(())
506 }
507}