Oregami
Repositories/oxedyne/ore

oxedyne/ore/store/src/store.rs

116 KiB, 91 runs

created by r2848102244:159, 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//! A directory of segments, read into a log and appended to.
2//!
3//! ```text
4//! <dir>/log/000000.seg ORESEG segments, appended to until one passes the size
5//! <dir>/log/000001.seg threshold, whereupon the next append starts the next
6//! file. Numbering is ascending and replay order is
7//! numeric order.
8//! <dir>/lock Where the advisory lock any command that may append
9//! holds is kept. Its presence is not the lock; its
10//! contents name the process holding it, or are empty.
11//! <dir>/verified Which sealed segments have had their signatures
12//! checked here already. Derived, replica-local and safe
13//! to delete; see [`crate::verdict`].
14//! ```
15//!
16//! A working copy puts this at `<root>/.ore` and keeps its own configuration,
17//! key and bookkeeping beside it. A relay puts it at the repository it hosts and
18//! keeps nothing beside it but an access list, because a relay has no working
19//! copy and authors nothing.
20//!
21//! # One writer at a time
22//!
23//! Appending is read-modify-write: [`Store::append`] reads the segment on disk,
24//! resumes the writer over it, and appends what the resume produced. Two writers
25//! doing that at once would each chain their records onto a segment the other was
26//! also extending, and the loser's digests would not match the bytes that ended
27//! up in front of them. So every writer takes [`Lock`] first, and one that finds
28//! the lock held says so and stops rather than waiting.
29//!
30//! # The integrity hasher
31//!
32//! Segment records carry a digest, and the engine takes the hash function from
33//! its caller. Every Ore repository uses SHA-256 under the fixed eight byte salt
34//! [`SALT`], which is `b"oreseg01"`. Both are part of the on-disk format: a
35//! reader that brings a different pair will find every record's digest wrong.
36//!
37//! # Verifying is the reader's choice, and the relay does not
38//!
39//! A working copy checks every signature it meets on the way in and refuses a
40//! history holding one that does not hold. A relay does not: it is not an
41//! authority, it holds no trust set, and a signature it could check is one the
42//! puller must check anyway. So it stores what arrives and serves it back
43//! unaltered, and a client refusing a forgery is the check that counts. See
44//! [`Verify`].
45
46use crate::keys::{
47 self,
48 Prov,
49 Trust,
50};
51use crate::repack;
52use crate::veil::{
53 is_standin,
54 standin,
55};
56use crate::verdict::{
57 self,
58 Fold,
59 Verdicts,
60 VERDICT_LIFE,
61};
62
63use oxedyne_fe2o3_core::prelude::*;
64use oxedyne_fe2o3_hash::hash::HashScheme;
65use oxedyne_fe2o3_hash::sha256::{
66 DIGEST_LEN,
67 Sha256,
68};
69use oxedyne_fe2o3_ore::envelope::Envelope;
70use oxedyne_fe2o3_ore::id::{
71 OpId,
72 ReplicaId,
73};
74use oxedyne_fe2o3_ore::log::OpLog;
75use oxedyne_fe2o3_ore::op::Record;
76use oxedyne_fe2o3_ore::segment::{
77 self,
78 Entry,
79 Head,
80 Integrity,
81 Veiled,
82};
83
84use std::collections::{
85 BTreeMap,
86 BTreeSet,
87};
88use std::fs;
89use std::io::{
90 Seek,
91 Write,
92};
93use std::path::{
94 Path,
95 PathBuf,
96};
97use std::time::{
98 SystemTime,
99 UNIX_EPOCH,
100};
101
102
103/// Name of the segment directory, within the store directory.
104pub const LOG_DIR: &str = "log";
105/// Name of the lock file, within the store directory.
106pub const LOCK_FILE: &str = "lock";
107
108/// The salt every segment digest is computed under. Part of the format.
109
110pub const SALT: [u8; 8] = *b"oreseg01";
111
112/// Size a segment file must reach before the next append starts a new one.
113pub const SEGMENT_LIMIT: u64 = 1 << 20;
114
115/// Suffix every segment file carries.
116pub const SEGMENT_SUFFIX: &str = ".seg";
117
118// The words a cursor's own lines begin with. See [`Consumed::owns`].
119const SEG_LINE: &str = "seg";
120const DEFERRED_LINE: &str = "deferred";
121const SEALED_WORD: &str = "sealed";
122
123
124/// Returns the integrity hasher every segment of every Ore repository is
125/// written and checked with.
126pub fn hasher() -> HashScheme {
127 HashScheme::new_sha256()
128}
129
130
131/// Returns the entries the log holds that were not among `before`, in append
132/// order, each in the form its provenance deserves.
133///
134/// An operation that arrived sealed is written to the segments sealed, so the
135/// signature survives the hop rather than being checked once and thrown away. An
136/// operation that arrived veiled is written veiled, for the stronger reason that
137/// nothing here could write it any other way: the log holds a stand-in for it,
138/// and `veils` is where the entry itself was kept.
139pub fn arrived(
140 log: &OpLog,
141 before: &BTreeSet<OpId>,
142 envelopes: &BTreeMap<OpId, Envelope>,
143 veils: &BTreeMap<OpId, Veiled>,
144)
145 -> Vec<Entry>
146{
147 log.iter()
148 .filter(|rec| !before.contains(&rec.id()))
149 .map(|rec| match veils.get(&rec.id()) {
150 Some(v) => Entry::Veiled(v.clone()),
151 None => match envelopes.get(&rec.id()) {
152 Some(env) => Entry::Sealed(env.clone()),
153 None => Entry::Bare(rec.clone()),
154 },
155 })
156 .collect()
157}
158
159/// Replaces bare entries with the envelopes held for them.
160///
161/// A session builds a send set out of a log, and a log holds records rather than
162/// envelopes, so what it produces is bare. This is the substitution the engine's
163/// sync module leaves to its caller: provenance crosses as provenance, and an
164/// operation that was signed when it was written is signed when it is sent on.
165pub fn sealed(entries: Vec<Entry>, envelopes: &BTreeMap<OpId, Envelope>)
166 -> Outcome<Vec<Entry>>
167{
168 let mut out = Vec::with_capacity(entries.len());
169 for entry in entries {
170 out.push(res!(seal_one(entry, envelopes)));
171 }
172 Ok(out)
173}
174
175/// The same substitution, one entry at a time.
176///
177/// For a caller that is bounding what it sends and so must not touch the entries
178/// past the bound. Substituting the whole owed set and truncating after is what
179/// the relay did until 2026-08-23, and on a clone of fe2o3's history it copied
180/// eighty-nine megabytes of envelopes on every one of thirty-two requests to send
181/// six.
182pub fn seal_one(entry: Entry, envelopes: &BTreeMap<OpId, Envelope>)
183 -> Outcome<Entry>
184{
185 let id = res!(entry.id());
186 Ok(match (&entry, envelopes.get(&id)) {
187 (Entry::Bare(_), Some(env)) => Entry::Sealed(env.clone()),
188 _ => entry,
189 })
190}
191
192
193/// Returns the highest operation wire code among the entries, which is what a
194/// segment's declared version has to spell before it can carry them.
195///
196/// A sealed entry is opened to be asked, without its signature being checked:
197/// what is wanted is the shape of the operation and not a claim about who wrote
198/// it, and [`Store::append`] is writing the entry down either way. See
199/// [`segment::highest_code`].
200/// A veiled entry is skipped, because it carries no code to be asked for: its
201/// operation is ciphertext and the version bounds only what a reader of the
202/// segment could be handed. See [`segment::KIND_VEILED`].
203fn top_code(entries: &[Entry])
204 -> Outcome<u8>
205{
206 let mut top = 0u8;
207 for entry in entries {
208 if entry.is_veiled() {
209 continue;
210 }
211 let code = res!(entry.peek()).op.code();
212 if code > top {
213 top = code;
214 }
215 }
216 Ok(top)
217}
218
219
220/// The right to append to one store, held for as long as the value lives.
221///
222/// # The lock is the kernel's, and the file is only where it is kept
223///
224/// `<dir>/lock` is opened, created where it is not there, and an exclusive
225/// advisory lock is taken on the open handle. That lock is what excludes a second
226/// caller, and the kernel releases it when the last handle on it closes -- on a
227/// clean drop, on a panic, on `SIGTERM`, and on `SIGKILL`, because closing every
228/// descriptor is something a dying process does not get a say in.
229///
230/// **The file is left behind on purpose.** Its presence used to be the lock, so a
231/// command that was killed left a file that refused every later command in that
232/// tree until somebody deleted it by hand. A git hook was written without a
233/// `timeout` around `ore` for exactly that reason, on 2026-08-22, choosing a hook
234/// that might hang over a store that might be bricked. Nothing now reads the file's
235/// presence as anything, so being killed costs the mark and nothing else.
236///
237/// Two designs were rejected in getting here, and both are unsound rather than
238/// merely inelegant.
239///
240/// **Asking whether the recorded process is alive.** The note says a PID, so a
241/// later caller could look. A PID says nothing on its own: it is reused, so the
242/// question "is a process with this number running" is not the question "is the
243/// process that took this lock running", and a wrong answer in the permissive
244/// direction is two writers in one store. Answering it soundly means recording the
245/// boot identifier and the process start time beside the PID and refusing wherever
246/// any of the three cannot be read -- a machine-specific answer to a question the
247/// kernel is already answering here for nothing.
248///
249/// **Removing the file on the way out.** Keeping the advisory lock and still
250/// unlinking on drop reopens a race the lock closes: a caller that opened this
251/// inode a moment before the unlink takes a lock on a file nobody can reach any
252/// more, while a third caller creates a new file at the path and locks that. Two
253/// writers again, and closing it needs a re-stat of the path against the handle and
254/// a retry loop. Not unlinking needs none of it.
255///
256/// What is written in the file is a note, and only a note: which process holds it
257/// and since when, so the refusal a second caller gets names something. It is
258/// cleared on the way out, so a note that is there is a note that is current.
259pub struct Lock {
260 file: fs::File, // the advisory lock lives on this handle and dies with it
261}
262
263impl Lock {
264
265 /// Returns the path of the lock file of the store directory `dir`.
266 pub fn path_of(dir: &Path) -> PathBuf {
267 dir.join(LOCK_FILE)
268 }
269
270 /// Takes the lock, or says what holds it.
271 pub fn take(dir: &Path)
272 -> Outcome<Self>
273 {
274 let path = Self::path_of(dir);
275 let mut file = match fs::OpenOptions::new()
276 .create(true).read(true).write(true).open(&path)
277 {
278 Ok(f) => f,
279 Err(e) => return Err(err!(e,
280 "The lock {:?} could not be opened.", path;
281 IO, File, Write)),
282 };
283 match file.try_lock() {
284 Ok(()) => (),
285 Err(fs::TryLockError::WouldBlock) => {
286 // Read only now, and only because the kernel has just said somebody
287 // holds this. A blank note is not a puzzle: the holder clears the
288 // note as it leaves and writes its own the instant it arrives, so
289 // nothing but an arrival that has not finished can leave one blank.
290 let held = match fs::read_to_string(&path) {
291 Ok(t) if !t.trim().is_empty() => fmt!("{}", t.trim()),
292 Ok(_) => fmt!("it has only just begun"),
293 Err(_) => fmt!("its note cannot be read"),
294 };
295 return Err(err!(
296 "Another `ore` command is working in {:?} ({}). Wait for it to \
297 finish.", dir, held;
298 Invalid, Data, Conflict, Exists));
299 },
300 Err(fs::TryLockError::Error(e)) => return Err(err!(e,
301 "The lock {:?} could not be taken.", path;
302 IO, File, Write)),
303 }
304 // Whatever a killed holder left is not this one's, and goes before this
305 // one's name does. The handle was opened without truncating and has not been
306 // written to, so its cursor is at the start already.
307 match file.set_len(0) {
308 Ok(()) => (),
309 Err(e) => return Err(err!(e,
310 "The lock {:?} could not be cleared.", path;
311 IO, File, Write)),
312 }
313 let stamp = res!(SystemTime::now().duration_since(UNIX_EPOCH));
314 let said = fmt!("process {} since {}\n", std::process::id(), stamp.as_secs());
315 match file.write_all(said.as_bytes()) {
316 Ok(()) => (),
317 Err(e) => return Err(err!(e,
318 "The lock {:?} could not be written.", path;
319 IO, File, Write)),
320 }
321 Ok(Self { file })
322 }
323}
324
325impl Drop for Lock {
326 fn drop(&mut self) {
327 // The note goes, so that one left in the file is one whose process is still
328 // in it. The file stays and the handle closing is what releases the lock.
329 let _ = self.file.set_len(0);
330 let _ = self.file.unlock();
331 }
332}
333
334
335/// How much a reader checks of what it finds.
336#[derive(Clone, Copy, Debug)]
337pub enum Verify<'a> {
338 /// Every sealed entry's signature is checked, and one that does not hold
339 /// refuses the whole replay by name. The trust set decides only whether a
340 /// verified signature is attributable, never whether it is accepted.
341 Signatures(&'a Trust),
342 /// Entries are taken as they are found. What a relay does: it holds no keys,
343 /// makes no claim about provenance, and hands on exactly what it was given
344 /// for the puller to check.
345 Nothing,
346}
347
348
349/// Whether the envelopes are wanted once the replay is over.
350///
351/// An envelope holds the whole encoded record a second time, so keeping every
352/// one costs a second copy of the log: 47.8 MB on a 44,629 operation history,
353/// measured, against a 55 MB log. Only a caller that will hand operations on to
354/// another replica has any use for them, and every verb that just reads was
355/// paying for them.
356#[derive(Clone, Copy, Debug, Eq, PartialEq)]
357pub enum Keep {
358 Envelopes, // a caller that will hand these operations to somebody else
359 Nothing, // a caller that only reads
360}
361
362
363/// Where a read of a segment file starts.
364///
365/// A segment grows only at its end and nothing already written is ever revisited,
366/// so a reader that kept where it stopped reads only what has arrived since. What
367/// it has to be handed back is what the bytes it will see do not carry: the header
368/// the file declared, and how many records stand in front of the byte it starts
369/// at. See [`Consumed`], which is where both are kept.
370#[derive(Clone, Copy, Debug)]
371pub(crate) enum Begin {
372 /// At the first byte, with the header still to be read.
373 Start,
374 /// At a byte a previous read of this file stopped at, and therefore knows to
375 /// be the start of a record.
376 Since {
377 at: u64,
378 head: Head,
379 ordinal: usize,
380 },
381}
382
383
384/// How far into a segment file a read got, and whether the file may grow past it.
385#[derive(Clone, Copy, Debug, Eq, PartialEq)]
386enum Reached {
387 /// The whole of a segment [`Store::append`] will never open again: one that
388 /// is not the newest, or one that has passed [`SEGMENT_LIMIT`]. Its bytes are
389 /// finished, so a later read that finds the same length and modification time
390 /// reads none of them.
391 Sealed,
392 /// As much of the newest segment as had been written, and what those bytes
393 /// hash to. The hash is the whole of what makes a resume safe over a file
394 /// that legitimately changes, and it is affordable because the file it covers
395 /// is under [`SEGMENT_LIMIT`] by the definition that put it here.
396 Growing([u8; DIGEST_LEN]),
397}
398
399
400/// One segment file, as far as a read got through it.
401#[derive(Clone, Debug, Eq, PartialEq)]
402struct Seen {
403 name: String, // the file's name within `log/`
404 len: u64, // how many of its bytes were read
405 modified: u128, // when the filesystem last said it was written
406 head: Head, // what it declared before its first record
407 records: usize, // how many records those bytes held
408 reached: Reached,
409}
410
411
412/// Where a read of a store stopped, so that the next may take up from there.
413///
414/// # What a resume point has to be
415///
416/// A byte offset, and one that is known to begin a record. A segment carries no
417/// index and a record is found only by reading the one before it, so **the only
418/// party that knows where a record begins is the read that stopped there**. That
419/// is why this is handed back by the read that produced it rather than built from
420/// the file, and why nothing here can be constructed by a caller with a number in
421/// its hand.
422///
423/// # What it refuses
424///
425/// A history is appended to and never rewritten, so a read may take up from where
426/// the last one stopped. A store where that is not true is a different store, and
427/// carrying a log across the difference would show a reader a repository that
428/// never existed. So [`Store::replay_since`] checks, before it reads a byte, that
429/// every segment this names is still there under the same name, that every sealed
430/// one still has the length and modification time it had, and that the one
431/// segment that may legitimately have grown still begins with the bytes it began
432/// with -- which is a hash over those bytes and not a guess from a timestamp. Any
433/// other state is refused by name, and the caller replays the store whole.
434///
435/// **The one thing it cannot see is a sealed segment rewritten with its length and
436/// modification time both put back.** Nothing reads those bytes again, so nothing
437/// can. That is the same evidence [`crate::verdict`] acts on and it is worth what
438/// the directory is worth: whoever can do it can rewrite `key` and the binary too.
439/// A caller that wants the whole log read is asking for [`Store::replay`], and
440/// that is the answer to it.
441///
442/// # It is this replica's own note and it never crosses a wire
443///
444/// Like a verdict, and for the same reason: it says what this process read, which
445/// a peer would have to take on trust.
446#[derive(Clone, Debug, Default, Eq, PartialEq)]
447pub struct Consumed {
448 seen: Vec<Seen>, // one per segment read, in replay order
449 deferred: usize, // records that had to wait for a parent later in the file
450}
451
452impl Consumed {
453
454 /// A cursor over nothing, which is what a read of a whole store starts from.
455 pub fn new() -> Self {
456 Self::default()
457 }
458
459 /// Has nothing been read yet?
460 pub fn is_empty(&self) -> bool {
461 self.seen.is_empty()
462 }
463
464 /// How many segment files the read got to.
465 pub fn segments(&self) -> usize {
466 self.seen.len()
467 }
468
469 /// How many bytes of them it read.
470 pub fn bytes(&self) -> u64 {
471 self.seen.iter().map(|one| one.len).sum()
472 }
473
474 /// How many operations those bytes held.
475 pub fn records(&self) -> usize {
476 self.seen.iter().map(|one| one.records).sum()
477 }
478
479 /// May a read take up from here?
480 ///
481 /// False where the read that produced this met an operation before its
482 /// parents. Nothing that writes an Ore segment does: what goes into one is
483 /// taken from a log in log order, and a log holds no operation before its
484 /// parents. Where it happens anyway, those records went into the log in an
485 /// order a whole replay would not have produced, and resuming from here would
486 /// carry that difference forward rather than show it -- so it is refused, and
487 /// [`Store::replay`] is what reads such a store.
488 pub fn resumable(&self) -> bool {
489 self.deferred == 0
490 }
491
492 /// Does a line beginning with this word belong to a cursor?
493 ///
494 /// A caller keeping a cursor inside a file of its own -- see
495 /// [`crate::verdict`] and the working copy index for the shape -- reads its
496 /// lines and this says which of them to hand to [`Consumed::from_lines`]. The
497 /// words are the cursor's business and not the caller's, which is why they
498 /// are asked about here rather than spelled out there.
499 pub fn owns(word: &str) -> bool {
500 word == SEG_LINE || word == DEFERRED_LINE
501 }
502
503 /// The cursor as lines of text, for a caller keeping one on disk.
504 ///
505 /// One line per segment, in replay order, and one saying how many records had
506 /// to wait. A cursor over nothing is no lines at all, which reads back as a
507 /// cursor over nothing.
508 pub fn to_lines(&self) -> Vec<String> {
509 let mut out = Vec::with_capacity(self.seen.len() + 1);
510 for one in &self.seen {
511 let replica = match one.head.replica {
512 Some(r) => fmt!("{}", r.inner()),
513 None => fmt!("-"),
514 };
515 let reached = match one.reached {
516 Reached::Sealed => fmt!("{}", SEALED_WORD),
517 Reached::Growing(hash) => keys::text_of(&hash),
518 };
519 out.push(fmt!("{} {} {} {} {} {} {} {}",
520 SEG_LINE, one.name, one.len, one.modified, one.records,
521 one.head.version, replica, reached));
522 }
523 out.push(fmt!("{} {}", DEFERRED_LINE, self.deferred));
524 out
525 }
526
527 /// Reads back what [`Consumed::to_lines`] wrote, or nothing where the lines
528 /// are not exactly that.
529 ///
530 /// **Nothing here is an error.** A cursor is a note this replica left itself
531 /// and the only cost of not being able to read one is reading the store
532 /// whole, which is what a caller holding no cursor does anyway. So a line
533 /// this build does not understand, a number that will not parse and a hash of
534 /// the wrong length all give the same answer: nothing is known.
535 pub fn from_lines<'a, I>(lines: I)
536 -> Option<Self>
537 where
538 I: IntoIterator<Item = &'a str>,
539 {
540 let mut seen: Vec<Seen> = Vec::new();
541 let mut deferred: Option<usize> = None;
542 for line in lines {
543 let field: Vec<&str> = line.split_whitespace().collect();
544 match field.first() {
545 Some(&DEFERRED_LINE) => {
546 if field.len() != 2 || deferred.is_some() {
547 return None;
548 }
549 deferred = match field[1].parse::<usize>() {
550 Ok(n) => Some(n),
551 Err(_) => return None,
552 };
553 },
554 Some(&SEG_LINE) => {
555 if field.len() != 8 {
556 return None;
557 }
558 let len = match field[2].parse::<u64>() {
559 Ok(n) => n,
560 Err(_) => return None,
561 };
562 let modified = match field[3].parse::<u128>() {
563 Ok(n) => n,
564 Err(_) => return None,
565 };
566 let records = match field[4].parse::<usize>() {
567 Ok(n) => n,
568 Err(_) => return None,
569 };
570 let version = match field[5].parse::<u8>() {
571 Ok(n) => n,
572 Err(_) => return None,
573 };
574 let replica = match field[6] {
575 "-" => None,
576 r => match r.parse::<u64>() {
577 Ok(n) => Some(ReplicaId::new(n)),
578 Err(_) => return None,
579 },
580 };
581 let reached = match field[7] {
582 SEALED_WORD => Reached::Sealed,
583 text => {
584 let bytes = match keys::bytes_of(text) {
585 Ok(b) => b,
586 Err(_) => return None,
587 };
588 if bytes.len() != DIGEST_LEN {
589 return None;
590 }
591 let mut hash = [0u8; DIGEST_LEN];
592 hash.copy_from_slice(&bytes);
593 Reached::Growing(hash)
594 },
595 };
596 seen.push(Seen {
597 name: fmt!("{}", field[1]),
598 len,
599 modified,
600 head: Head { version, replica },
601 records,
602 reached,
603 });
604 },
605 _ => return None,
606 }
607 }
608 deferred.map(|deferred| Self { seen, deferred })
609 }
610}
611
612
613/// Whether a store still holds exactly the bytes a cursor read.
614///
615/// See [`Store::settled_at`], which is the only thing that produces one.
616pub enum Settled {
617 /// Nothing has been appended and nothing has moved, so a reading taken under
618 /// the cursor still describes this store.
619 Same,
620 /// It has not, and this says how. The sentence names the segment and what
621 /// about it differs, and it is written to be printed.
622 Moved(String),
623}
624
625
626/// What a replay produced.
627///
628/// Cloneable because a reader may be holding one while the next read extends it:
629/// see [`Store::replay_since`], and [`std::sync::Arc::make_mut`], which is what a
630/// caller sharing one uses to pay for the copy only when somebody is still
631/// reading the old one.
632#[derive(Clone)]
633pub struct Replayed {
634 /// Every operation, in causal order.
635 pub log: OpLog,
636 /// The envelope every sealed operation was found in, kept so that provenance
637 /// can be handed on rather than stopping here.
638 pub envelopes: BTreeMap<OpId, Envelope>,
639 /// What is known about who wrote each operation. Empty under
640 /// [`Verify::Nothing`], which asks nothing and therefore knows nothing.
641 pub prov: BTreeMap<OpId, Prov>,
642 /// Every veiled entry found, by the operation its clear header names. The log
643 /// holds a stand-in for each; this is what the stand-in stands for, and what
644 /// goes back out in its place. Empty except at a carrier, since a working copy
645 /// refuses a veil it cannot read rather than holding one.
646 pub veils: BTreeMap<OpId, Veiled>,
647}
648
649
650impl Replayed {
651
652 /// What a replay of nothing produced, which is what a read fills.
653 pub fn new() -> Self {
654 Self {
655 log: OpLog::new(),
656 envelopes: BTreeMap::new(),
657 prov: BTreeMap::new(),
658 veils: BTreeMap::new(),
659 }
660 }
661}
662
663impl Default for Replayed {
664 fn default() -> Self {
665 Self::new()
666 }
667}
668
669
670/// What a read took out of the segments, with none of it placed in a log.
671///
672/// [`Store::replay_since`] fills a [`Replayed`] from one of these and is the
673/// only caller that should have to think about the difference. The other caller
674/// is one that already holds these operations -- it authored them a moment ago
675/// -- and wants the cursor over the bytes they were written to.
676#[derive(Debug, Default)]
677pub struct Read {
678 /// Every operation the read found, in log order, placed nowhere.
679 pub records: Vec<Record>,
680 pub envelopes: Vec<(OpId, Envelope)>,
681 pub prov: Vec<(OpId, Prov)>,
682 pub veils: Vec<(OpId, Veiled)>,
683 /// Where the read stopped, which is what anything derived from these
684 /// operations is derived under.
685 pub cursor: Consumed,
686}
687
688
689/// What every segment of one replay is asked for, whichever segment it is.
690#[derive(Clone, Copy)]
691struct Asked<'a> {
692 how: Verify<'a>,
693 keep: Keep,
694 fold: bool, // only a reader that checks signatures has a use for one
695}
696
697
698/// What a read has to do about a segment a cursor already names.
699enum Taken {
700 /// Nothing. The last read finished it and nothing has touched it since, so
701 /// what it holds is already in the caller's [`Replayed`].
702 Whole(Seen),
703 /// Take it up at the byte the last read stopped at, carrying the hasher that
704 /// has already absorbed every byte in front of that one.
705 Part(Seen, Sha256),
706}
707
708
709/// What one segment produced, held apart from the replay until the fold over it
710/// says whether it may be believed.
711struct Gathered {
712 records: Vec<Record>,
713 prov: Vec<(OpId, Prov)>,
714 envelopes: Vec<(OpId, Envelope)>,
715 veils: Vec<(OpId, Veiled)>,
716 fold: Option<Fold>, // none where nothing asked for one
717 raw: Option<[u8; DIGEST_LEN]>, // none where the segment cannot grow again
718 head: Head, // what the file declared, or was said to
719}
720
721
722/// A directory of segments.
723///
724/// The value is a path and nothing else: every method reads what is on disk when
725/// it is called, so two of these over one directory are the same store and there
726/// is no cache to fall out of date.
727#[derive(Clone, Debug)]
728pub struct Store {
729 /// The directory holding `log/`.
730 pub dir: PathBuf,
731}
732
733impl Store {
734
735 /// Names the store at a directory, without touching it.
736 pub fn at(dir: &Path) -> Self {
737 Self { dir: dir.to_path_buf() }
738 }
739
740 /// Creates the store's directories and its first, empty segment.
741 ///
742 /// An empty log is a segment holding a header and no records, so that a
743 /// reader never has to treat "no file" as a case of its own. `hint` is what
744 /// that header names as the replica whose work the segment mostly carries,
745 /// which is a hint and nothing else.
746 pub fn create(dir: &Path, hint: Option<ReplicaId>)
747 -> Outcome<Self>
748 {
749 let store = Self::at(dir);
750 res!(fs::create_dir_all(store.dir.join(LOG_DIR)));
751 let head = Head::new(hint);
752 let path = store.segment_path(0);
753 match fs::write(&path, head.encode()) {
754 Ok(()) => (),
755 Err(e) => return Err(err!(e,
756 "The first segment {:?} could not be written.", path;
757 IO, File, Write)),
758 }
759 Ok(store)
760 }
761
762 /// Reports whether there is a store here.
763 pub fn exists(&self) -> bool {
764 self.dir.join(LOG_DIR).is_dir()
765 }
766
767 /// Takes the lock on this store.
768 pub fn lock(&self)
769 -> Outcome<Lock>
770 {
771 Lock::take(&self.dir)
772 }
773
774 /// Returns the path of the segment with the given number.
775 pub fn segment_path(&self, n: u64) -> PathBuf {
776 self.dir.join(LOG_DIR).join(fmt!("{:06}{}", n, SEGMENT_SUFFIX))
777 }
778
779 /// Returns every segment, in replay order.
780 pub fn segments(&self)
781 -> Outcome<Vec<(u64, PathBuf)>>
782 {
783 let dir = self.dir.join(LOG_DIR);
784 let entries = match fs::read_dir(&dir) {
785 Ok(e) => e,
786 Err(e) => {
787 // A repack swaps its result in with two renames, and between them is
788 // the one instant at which a store has no log at all. Saying so, and
789 // naming the move that undoes it, is what turns that instant from a
790 // puzzle into a sentence. See [`crate::repack`].
791 let aside = self.dir.join(repack::OLD_DIR);
792 if aside.is_dir() {
793 return Err(err!(e,
794 "There is no log at {:?} and there is one at {:?}. A repack was \
795 interrupted between the two renames that put its result in place. \
796 Nothing has been lost: the history is whole and it is in {:?}. \
797 Move it back with `mv {} {}`.",
798 dir, aside, aside, aside.display(), dir.display();
799 IO, File, Read, Missing));
800 }
801 return Err(err!(e,
802 "The log directory {:?} could not be read.", dir;
803 IO, File, Read));
804 },
805 };
806 let mut out: Vec<(u64, PathBuf)> = Vec::new();
807 for entry in entries {
808 let entry = res!(entry);
809 let path = entry.path();
810 let name = path.file_name().map(|n| n.to_string_lossy().into_owned());
811 let name = match name {
812 Some(n) => n,
813 None => continue,
814 };
815 let stem = match name.strip_suffix(SEGMENT_SUFFIX) {
816 Some(s) => s,
817 None => continue,
818 };
819 let n = match stem.parse::<u64>() {
820 Ok(n) => n,
821 Err(e) => return Err(err!(e,
822 "The segment {:?} is not named by a number.", path;
823 Invalid, Input, Name)),
824 };
825 out.push((n, path));
826 }
827 out.sort();
828 Ok(out)
829 }
830
831 /// Returns the total size of the segments, which is the size of the log.
832 pub fn bytes(&self)
833 -> Outcome<u64>
834 {
835 let mut total = 0u64;
836 for (_, path) in res!(self.segments()) {
837 match fs::metadata(&path) {
838 Ok(m) => total += m.len(),
839 Err(e) => return Err(err!(e,
840 "The segment {:?} could not be measured.", path;
841 IO, File, Read)),
842 }
843 }
844 Ok(total)
845 }
846
847 /// Replays every segment into a log, checking what `how` asks it to check.
848 ///
849 /// A segment that will not decode names itself in the error, because a
850 /// truncated file is the failure a reader is most likely to meet and the
851 /// least useful to be told about anonymously.
852 ///
853 /// # A signature that does not hold stops the replay
854 ///
855 /// Under [`Verify::Signatures`] every sealed entry is put to
856 /// [`keys::check_all`], a segment at a time, which checks the segment's
857 /// signatures together where the scheme can and falls back to one at a time
858 /// where it must. One whose signature does not verify is refused by name:
859 /// the log is not loaded and the message says which operation and which file.
860 /// It is never quietly dropped, because an operation missing from a history
861 /// is a different history, and never quietly accepted, because that is what
862 /// the signature was for. A signature that verifies under a key the trust set
863 /// has not seen is not a failure; it is recorded as [`Prov::Unknown`].
864 /// Reads a segment in windows, handing entries to `take` in batches of at
865 /// most `batch`.
866 ///
867 /// The batch is what bounds the peak. Gathering a whole segment and then
868 /// verifying it holds the entries AND everything verification produces from
869 /// them at the same time -- measured at 66 MB of entries beside 102 MB of
870 /// records and envelopes, which is the largest single step in opening a
871 /// repository. Verifying a batch at a time lets each batch go as soon as it
872 /// has been turned into records.
873 ///
874 /// The batch is still large enough to keep what makes batch verification
875 /// worth having: one public key decompression per signer rather than one per
876 /// operation. A signer appears many times in two thousand operations.
877 /// Reads a segment in windows, handing `take` a batch of entries at a time
878 /// together with the digest each of them was checked against.
879 ///
880 /// The digests are the reader's own, computed to check each record and
881 /// otherwise dropped, so a caller that wants to name the exact bytes it read
882 /// pays a copy rather than a second hashing of the file: 7 ms over a 55 MB
883 /// segment against 205. They come a batch at a time, which is what stops them
884 /// growing to one per record in the segment. What is done with them is the
885 /// caller's -- [`Store::replay`] folds them per segment for
886 /// [`crate::verdict`], and [`crate::repack`] folds them across the whole log
887 /// as well.
888 ///
889 /// The header comes back at the end because a reader has it only once it has
890 /// read a record, and a fold that names a segment wants it.
891 ///
892 /// `begin` says where in the file to start, and `raw` is a hasher over the
893 /// bytes actually read, which is what a resume is checked against. Neither
894 /// is the reader's business: the reader is handed bytes and hands back
895 /// records, and where those bytes came from is settled here.
896 pub(crate) fn stream_batches<F>(
897 path: &Path,
898 begin: Begin,
899 check: Integrity,
900 raw: Option<&mut Sha256>,
901 batch: usize,
902 mut take: F,
903 )
904 -> Outcome<Head>
905 where
906 F: FnMut(&mut Vec<Entry>, &[u8]) -> Outcome<()>,
907 {
908 const WINDOW: usize = 64 * 1024;
909 let mut file = match fs::File::open(path) {
910 Ok(f) => f,
911 Err(e) => return Err(err!(e,
912 "The segment {:?} could not be opened.", path; IO, File, Read)),
913 };
914 let mut reader: segment::Reader<_, 8> = segment::Reader::tallying(hasher(), SALT)
915 .integrity(check);
916 if let Begin::Since { at, head, ordinal } = begin {
917 match file.seek(std::io::SeekFrom::Start(at)) {
918 Ok(_) => (),
919 Err(e) => return Err(err!(e,
920 "The segment {:?} could not be taken up at byte {}.", path, at;
921 IO, File, Read)),
922 }
923 // The header is at the front of the file and the read is starting past
924 // it, so it is handed back rather than read; the ordinal keeps a damaged
925 // record naming its place in the file. See
926 // [`oxedyne_fe2o3_ore::segment::Reader::take_up`].
927 reader.take_up(head, ordinal);
928 }
929 let mut raw = raw;
930 let mut entries: Vec<Entry> = Vec::new();
931 let mut window = vec![0u8; WINDOW];
932 loop {
933 let n = match std::io::Read::read(&mut file, &mut window) {
934 Ok(n) => n,
935 Err(e) => return Err(err!(e,
936 "The segment {:?} could not be read.", path; IO, File, Read)),
937 };
938 if n == 0 {
939 break;
940 }
941 if let Some(sum) = raw.as_mut() {
942 sum.update(&window[..n]);
943 }
944 reader.feed(&window[..n]);
945 // Pulled as they complete, which is what keeps the reader's own
946 // buffer from growing to the size of the file.
947 while let Some(entry) = match reader.next_entry() {
948 Ok(e) => e,
949 Err(e) => return Err(err!(e,
950 "The segment {:?} could not be decoded.", path; Decode, Input, File)),
951 } {
952 entries.push(entry);
953 if entries.len() >= batch {
954 let digests = reader.take_digests();
955 res!(take(&mut entries, &digests));
956 entries.clear();
957 }
958 }
959 }
960 reader.end();
961 while let Some(entry) = match reader.next_entry() {
962 Ok(e) => e,
963 Err(e) => return Err(err!(e,
964 "The segment {:?} could not be decoded.", path; Decode, Input, File)),
965 } {
966 entries.push(entry);
967 if entries.len() >= batch {
968 let digests = reader.take_digests();
969 res!(take(&mut entries, &digests));
970 entries.clear();
971 }
972 }
973 let digests = reader.take_digests();
974 if !entries.is_empty() || !digests.is_empty() {
975 res!(take(&mut entries, &digests));
976 }
977 let head = res!(reader.head().ok_or_else(|| err!(
978 "The segment {:?} was read to the end and never yielded a header.", path;
979 Bug, Missing)));
980 Ok(*head)
981 }
982
983 /// Folds a segment's record digests, and its header, into the one value a
984 /// verdict about it is filed under.
985 ///
986 /// The header goes in last because the fold is built as the records arrive
987 /// and the header is what the reader read before the first of them. The order
988 /// is fixed here and nowhere else.
989 pub(crate) fn fold_from(sum: Sha256, head: &Head) -> Fold {
990 let mut sum = sum;
991 sum.update(&head.encode());
992 sum.update(&SALT);
993 sum.finish()
994 }
995
996 /// Reads one segment, checking every signature in it unless `taken` says a
997 /// verdict over these bytes has already answered that.
998 ///
999 /// What comes back is held apart from the replay, because on the `taken` path
1000 /// nothing yet says the verdict was about the bytes that were actually read.
1001 /// The caller puts the fold to [`Verdicts::holds`] and either takes this or
1002 /// throws it away and reads the segment again.
1003 fn read_segment(
1004 path: &Path,
1005 begin: Begin,
1006 asked: Asked,
1007 taken: bool,
1008 raw: Option<Sha256>,
1009 whence: &str,
1010 )
1011 -> Outcome<Gathered>
1012 {
1013 // A BATCH at a time, not a segment. The entries of a batch are released
1014 // as soon as verification has turned them into records, so the two are
1015 // never both resident for a whole segment. Two thousand keeps a signer's
1016 // key decompressed across many operations, which is what batch
1017 // verification is for; a failure still names its own operation, because
1018 // `check_all` falls back to one at a time.
1019 const BATCH: usize = 2_000;
1020 let mut got = Gathered {
1021 records: Vec::new(),
1022 prov: Vec::new(),
1023 envelopes: Vec::new(),
1024 veils: Vec::new(),
1025 fold: None,
1026 raw: None,
1027 head: Head::default(),
1028 };
1029 let mut sum = Sha256::new();
1030 let mut over = raw;
1031 // The one place the digest gate is decided. `taken` is already the answer
1032 // to "has a verdict vouched for these exact bytes, whole and sealed", which
1033 // is what buys the skip of every signature just below; the framing digest
1034 // rests on the same warrant and no other.
1035 let check = match taken {
1036 true => Integrity::Vouched,
1037 false => Integrity::Checked,
1038 };
1039 let head = res!(Self::stream_batches(path, begin, check, over.as_mut(), BATCH, |entries, digests| {
1040 if asked.fold {
1041 sum.update(digests);
1042 }
1043 match asked.how {
1044 Verify::Signatures(trust) => {
1045 let checked = match taken {
1046 true => res!(keys::attribute_all(entries, trust, whence, asked.keep)),
1047 false => res!(keys::check_all(entries, trust, whence, asked.keep)),
1048 };
1049 for one in checked {
1050 let id = one.rec.id();
1051 if let (Some(env), Keep::Envelopes) = (one.env, asked.keep) {
1052 got.envelopes.push((id, env));
1053 }
1054 got.prov.push((id, one.prov));
1055 got.records.push(one.rec);
1056 }
1057 },
1058 Verify::Nothing => {
1059 for entry in entries.drain(..) {
1060 // A veiled entry has no record to put in the log, so what goes
1061 // in is the stand-in and what is kept beside it is the entry.
1062 if let Entry::Veiled(v) = entry {
1063 got.veils.push((v.head.id(), v.clone()));
1064 got.records.push(standin(v.head));
1065 continue;
1066 }
1067 let rec = res!(entry.peek());
1068 if let (Entry::Sealed(env), Keep::Envelopes) = (&entry, asked.keep) {
1069 got.envelopes.push((rec.id(), env.clone()));
1070 }
1071 got.records.push(rec);
1072 }
1073 },
1074 }
1075 Ok(())
1076 }));
1077 if asked.fold {
1078 got.fold = Some(Self::fold_from(sum, &head));
1079 }
1080 if let Some(sum) = over {
1081 got.raw = Some(sum.finish());
1082 }
1083 got.head = head;
1084 Ok(got)
1085 }
1086
1087 /// Is the store exactly the bytes a cursor read, with nothing appended since?
1088 ///
1089 /// A caller holding a reading of the log made under `from` -- on disk, across
1090 /// processes, wherever -- asks this before it acts on that reading. It is the
1091 /// same evidence [`Store::replay_since`] takes up on, put to the narrower
1092 /// question of whether there is anything to take up at all: every segment the
1093 /// cursor names still there under the same name, every sealed one with the
1094 /// length and modification time it was read at, the one segment that may
1095 /// legitimately grow still beginning with the bytes it began with -- read and
1096 /// hashed, not guessed from a timestamp -- **and** no bytes and no segments
1097 /// past what the cursor got to.
1098 ///
1099 /// What it cannot see is what [`Consumed`] says it cannot see, and no more.
1100 ///
1101 /// # It never fails for the wrong reason
1102 ///
1103 /// Everything that says "not the same bytes" comes back as
1104 /// [`Settled::Moved`], carrying the sentence that says which segment and how,
1105 /// so a caller can put it in front of somebody rather than fall back in
1106 /// silence. That includes a segment that cannot be opened: a caller told this
1107 /// reads the store whole, and the read is where such a failure belongs.
1108 pub fn settled_at(&self, from: &Consumed)
1109 -> Outcome<Settled>
1110 {
1111 if from.is_empty() {
1112 return Ok(Settled::Moved(fmt!(
1113 "the reading was taken from no segment at all")));
1114 }
1115 let segments = res!(self.segments());
1116 if segments.len() != from.seen.len() {
1117 return Ok(Settled::Moved(fmt!(
1118 "the reading was taken over {} segment{} and there {} {} in {:?} now",
1119 from.seen.len(), if from.seen.len() == 1 { "" } else { "s" },
1120 if segments.len() == 1 { "is" } else { "are" }, segments.len(), self.dir)));
1121 }
1122 let taking_up = match self.take_up(from, &segments) {
1123 Ok(t) => t,
1124 Err(e) => return Ok(Settled::Moved(e.plain())),
1125 };
1126 // The one segment an append may extend. `take_up` proves the bytes already
1127 // read are the bytes that were read; what is left is whether any have been
1128 // written after them, which is the whole difference between resuming a read
1129 // and having nothing to read.
1130 for (name, taken) in &taking_up {
1131 if let Taken::Part(was, _) = taken {
1132 let path = self.dir.join(LOG_DIR).join(name);
1133 let (len, _) = res!(Self::stamp(&path));
1134 if len != was.len {
1135 return Ok(Settled::Moved(fmt!(
1136 "{} bytes have been appended to the segment {} since the reading \
1137 was taken", len - was.len, name)));
1138 }
1139 }
1140 }
1141 Ok(Settled::Same)
1142 }
1143
1144 /// Replays the whole store into a fresh log.
1145 ///
1146 /// One line over [`Store::replay_since`] from a cursor over nothing, so that
1147 /// the whole read and the resumed read are not two implementations of one
1148 /// thing that could come to disagree. They are the same code, and a resumed
1149 /// replay produces what this produces because it *is* this, entered further
1150 /// along.
1151 pub fn replay(&self, how: Verify, keep: Keep)
1152 -> Outcome<Replayed>
1153 {
1154 let mut out = Replayed::new();
1155 res!(self.replay_since(&Consumed::new(), &mut out, how, keep));
1156 Ok(out)
1157 }
1158
1159 /// Reads what `from` has not already read, onto what it produced.
1160 ///
1161 /// `into` must be what the read that produced `from` filled, under the same
1162 /// [`Keep`]; it is extended in place, and what comes back is where this read
1163 /// stopped. A cursor over nothing and an empty [`Replayed`] read the whole
1164 /// store, which is what [`Store::replay`] is.
1165 ///
1166 /// **What this returns is what a whole replay returns**: the same operations
1167 /// in the same log order, the same provenance, the same envelopes and the same
1168 /// veils. Three things hold that up. The bytes of a segment already read are
1169 /// never revisited, because a history is appended to and never rewritten, and
1170 /// the cursor refuses to be resumed over a store where that stopped being true
1171 /// -- see [`Consumed`]. The records of a segment stand in log order, so
1172 /// placing them one file at a time places them exactly as placing them all at
1173 /// once would, and where a file is not in that order the cursor says so and
1174 /// this refuses. And the rest of what a replay produces is keyed by operation
1175 /// identity, so it cannot depend on how the reading was divided.
1176 ///
1177 /// # What is verified again, and what is trusted
1178 ///
1179 /// Nothing already read is read again, so nothing already read is verified
1180 /// again. What that rests on is not a timestamp: it is that an operation
1181 /// whose signature verified is an operation whose signature verifies, that
1182 /// nothing is ever removed from a log, and that the bytes those verifications
1183 /// were made over are checked to be the same bytes before a byte is read past
1184 /// them. New bytes are checked exactly as a whole replay checks them.
1185 ///
1186 /// This is the same argument [`crate::verdict`] makes and it stands on firmer
1187 /// ground: a verdict is a note left on disk by an earlier process, and a
1188 /// cursor is what this process read for itself. So a resumed read consults no
1189 /// verdict for a segment it skips, and it does not rewrite the verdicts of one
1190 /// either -- what it writes is what it already held plus what this read
1191 /// learned, never less. A segment that fills up between two reads gets its
1192 /// verdict from the next whole replay rather than from the read that finished
1193 /// it, because a fold over a segment is over the whole of it and a resumed
1194 /// read has only its tail.
1195 pub fn replay_since(
1196 &self,
1197 from: &Consumed,
1198 into: &mut Replayed,
1199 how: Verify,
1200 keep: Keep,
1201 )
1202 -> Outcome<Consumed>
1203 {
1204 let got = res!(self.read_since(from, how, keep));
1205 for (id, env) in got.envelopes {
1206 into.envelopes.insert(id, env);
1207 }
1208 for (id, prov) in got.prov {
1209 into.prov.insert(id, prov);
1210 }
1211 for (id, veiled) in got.veils {
1212 into.veils.insert(id, veiled);
1213 }
1214 let mut cursor = got.cursor;
1215 let records = got.records;
1216 let count = records.len();
1217 // Placed one at a time while each one's parents are already there, which is
1218 // what the whole batch's first pass through `OpLog::absorb` does, and then
1219 // absorbed for the rest. Splitting it this way changes nothing about the
1220 // order and says one thing more: how many records the file offered before
1221 // their parents, which is the one shape of history a later read cannot take
1222 // up from. See [`Consumed::resumable`].
1223 let mut waiting: Vec<Record> = Vec::new();
1224 for rec in records {
1225 match rec.parents().iter().all(|p| into.log.contains(p)) {
1226 true => res!(into.log.append(rec)),
1227 false => waiting.push(rec),
1228 }
1229 }
1230 cursor.deferred += waiting.len();
1231 let left = res!(into.log.absorb(waiting));
1232 if !left.is_empty() {
1233 return Err(err!(
1234 "{} of the {} operations in the log name parents the log does not \
1235 hold, starting with {}; the history is incomplete.",
1236 left.len(), count, left[0].id();
1237 Invalid, Input, Missing));
1238 }
1239 Ok(cursor)
1240 }
1241
1242 /// Reads what `from` has not already read, and places none of it.
1243 ///
1244 /// This is [`Store::replay_since`] without its second half. It exists for the
1245 /// one caller that already holds the operations and wants only the cursor: a
1246 /// command that has just appended has moved the store past the cursor it read
1247 /// under, and everything it derived from the log -- the history listing, the
1248 /// working copy index -- is then unattributable to any bytes and has to be
1249 /// thrown away. Reading the few bytes it wrote back gives it a cursor over
1250 /// them, and a mark stops costing the next command a whole read.
1251 ///
1252 /// **The records still come back and the caller must reconcile them.** They
1253 /// are what the disk says, and a caller that quietly dropped them would be
1254 /// keeping a cursor over bytes it had never looked at, which is the one thing
1255 /// [`Consumed`] exists to prevent. See [`crate::store::Read::records`].
1256 pub fn read_since(
1257 &self,
1258 from: &Consumed,
1259 how: Verify,
1260 keep: Keep,
1261 )
1262 -> Outcome<Read>
1263 {
1264 let segments = res!(self.segments());
1265 let newest = segments.last().map(|(n, _)| *n);
1266 // Every segment named by the cursor, measured and checked against what the
1267 // cursor says of it, before a byte is read past any of them.
1268 let taking_up = res!(self.take_up(from, &segments));
1269 // Only a reader that checks signatures has a verdict to look up or one to
1270 // leave behind. A relay checks nothing and therefore knows nothing, and
1271 // writing down what it did not do is exactly the mistake this must not
1272 // make.
1273 let checking = matches!(how, Verify::Signatures(_));
1274 let known = match checking {
1275 true => Verdicts::read(&self.dir),
1276 false => Verdicts::new(),
1277 };
1278 // A whole read starts from nothing, so a note about a segment no longer on
1279 // the disk falls out of the file. A resumed read reads only part of the
1280 // directory and must not prune what it did not look at, so it starts from
1281 // what is already recorded.
1282 let mut seen = match from.is_empty() {
1283 true => Verdicts::new(),
1284 false => known.clone(),
1285 };
1286 let asked = Asked { how, keep, fold: checking };
1287 // One reading of the clock for the whole replay, so that a note filed for
1288 // the last segment is not a moment younger than one filed for the first
1289 // and two runs of the same command cannot disagree about which of them
1290 // has gone off.
1291 let at = verdict::now();
1292 let mut out = Read::default();
1293 let mut records: Vec<Record> = Vec::new();
1294 let mut cursor = Consumed { seen: Vec::new(), deferred: from.deferred };
1295 for (n, path) in &segments {
1296 // STREAMED, not read whole.
1297 //
1298 // `fs::read` put the entire segment on the heap and `segment::decode`
1299 // then handed that slice to `Reader::feed`, which copies what it is
1300 // given -- so a 53 MB segment cost 106 MB of buffers before a single
1301 // record existed. The reader was built for exactly this and says so:
1302 // "memory is bounded by the largest single record rather than by the
1303 // segment", and "feeding a segment one byte at a time yields exactly
1304 // what feeding it all at once yields". Only the convenience wrapper
1305 // defeated it.
1306 let whence = fmt!("the segment {}", path.display());
1307 let (len, modified) = res!(Self::stamp(path));
1308 let name = match path.file_name() {
1309 Some(n) => n.to_string_lossy().into_owned(),
1310 None => continue,
1311 };
1312 // A segment nothing will ever append to. `Store::append` continues the
1313 // newest segment and only while it is under `SEGMENT_LIMIT`, so every
1314 // other segment in the directory is finished for good. The tail is left
1315 // out because a fold over it names a state it may already have left,
1316 // and an entry per append would be an entry per command.
1317 let sealed = Some(*n) != newest || len >= SEGMENT_LIMIT;
1318 let already = taking_up.get(&name);
1319 // A segment the last read finished, that nothing has touched since.
1320 // Nothing is read and nothing is placed; what it produced is already in
1321 // `into`, and its verdict is already in `seen`.
1322 if let Some(Taken::Whole(was)) = already {
1323 cursor.seen.push((*was).clone());
1324 continue;
1325 }
1326 let (begin, raw) = match already {
1327 Some(Taken::Part(was, prefix)) => (
1328 Begin::Since { at: was.len, head: was.head, ordinal: was.records },
1329 Some(prefix.clone()),
1330 ),
1331 // A segment about to be read from its first byte. Its bytes are
1332 // hashed only where it may still grow, since that hash is what the
1333 // next read checks the prefix against and a sealed segment has no
1334 // prefix to check -- and hashing every sealed segment would put the
1335 // cost of the whole log back into every read.
1336 _ => (Begin::Start, match sealed {
1337 true => None,
1338 false => Some(Sha256::new()),
1339 }),
1340 };
1341 // A GUESS, and treated as one. See `crate::verdict`: the name, the
1342 // length and the modification time decide only whether to read the
1343 // segment without checking it, and the fold decides whether that read
1344 // may be believed. Only a segment read whole has a fold to put to it.
1345 let whole = matches!(begin, Begin::Start);
1346 let guess = match checking && sealed && whole {
1347 true => known.expected(&name, len, modified, at, VERDICT_LIFE),
1348 false => None,
1349 };
1350 let asked = Asked { fold: checking && whole, ..asked };
1351 let mut got = res!(Self::read_segment(
1352 path, begin, asked, guess.is_some(), raw.clone(), &whence));
1353 if guess.is_some() {
1354 let fold = res!(got.fold.ok_or_else(|| err!(
1355 "The segment {:?} was read on a verdict and no fold was taken of \
1356 it, so nothing can say whether the verdict was about these bytes.",
1357 path;
1358 Bug, Missing)));
1359 if !known.holds(&fold) {
1360 // The file was not the file it looked like. Nothing checked what
1361 // it just produced, so all of it goes and the segment is read
1362 // again with every signature checked. That read is what refuses
1363 // a tampered segment; this is only what stops it being skipped.
1364 got = res!(Self::read_segment(
1365 path, begin, asked, false, raw, &whence));
1366 }
1367 }
1368 if let (true, true, Some(fold)) = (checking, sealed, got.fold) {
1369 seen.record(&name, len, modified, fold, at);
1370 }
1371 let reached = match sealed {
1372 true => Reached::Sealed,
1373 false => Reached::Growing(res!(got.raw.ok_or_else(|| err!(
1374 "The segment {:?} may still be appended to and no hash was taken \
1375 of the bytes read from it, so the next read could not tell that \
1376 they had not moved.", path;
1377 Bug, Missing)))),
1378 };
1379 let before = match already {
1380 Some(Taken::Part(was, _)) => was.records,
1381 _ => 0,
1382 };
1383 cursor.seen.push(Seen {
1384 name,
1385 len,
1386 modified,
1387 head: got.head,
1388 records: before + got.records.len(),
1389 reached,
1390 });
1391 out.envelopes.extend(got.envelopes);
1392 out.prov.extend(got.prov);
1393 out.veils.extend(got.veils);
1394 if records.is_empty() {
1395 records = got.records;
1396 } else {
1397 records.extend(got.records);
1398 }
1399 }
1400 if checking && seen != known {
1401 // Derived state, and a store that cannot be written to is a store that
1402 // verifies every time -- which is what every store did until this was
1403 // written. So a failure here is not the command's failure.
1404 let _ = seen.write(&self.dir);
1405 }
1406 out.records = records;
1407 out.cursor = cursor;
1408 Ok(out)
1409 }
1410
1411 /// Checks that every segment `from` names is still the file it was read from,
1412 /// and says how each is to be taken up.
1413 ///
1414 /// This is where a resume is refused, and it is refused before a byte is read
1415 /// past a segment rather than after. What it can see it names: a segment that
1416 /// has gone, one renamed, one whose length or modification time has moved when
1417 /// nothing should ever have opened it again, one that has shrunk, and -- for
1418 /// the one segment that may legitimately grow -- a prefix that does not hash
1419 /// to what the last read hashed it to. What it cannot see is set out in
1420 /// [`Consumed`].
1421 fn take_up(&self, from: &Consumed, segments: &[(u64, PathBuf)])
1422 -> Outcome<BTreeMap<String, Taken>>
1423 {
1424 let mut out = BTreeMap::new();
1425 if from.is_empty() {
1426 return Ok(out);
1427 }
1428 if !from.resumable() {
1429 return Err(err!(
1430 "The read this one would take up from placed {} operation{} only after \
1431 meeting a later one, because a segment offered them before their \
1432 parents. Their place in the log is therefore not the place a whole \
1433 replay would give them, and taking up from here would carry that \
1434 difference forward instead of showing it. Replay the store whole.",
1435 from.deferred, if from.deferred == 1 { "" } else { "s" };
1436 Invalid, Input, Order));
1437 }
1438 if from.seen.len() > segments.len() {
1439 return Err(err!(
1440 "The last read of {:?} got to {} segment{} and there are {} there now. \
1441 A log is only ever added to, so segments have been removed or the \
1442 directory has been replaced; the store cannot be taken up where it was \
1443 left.",
1444 self.dir, from.seen.len(),
1445 if from.seen.len() == 1 { "" } else { "s" }, segments.len();
1446 Invalid, Input, Mismatch));
1447 }
1448 for (was, (_, path)) in from.seen.iter().zip(segments.iter()) {
1449 let name = match path.file_name() {
1450 Some(n) => n.to_string_lossy().into_owned(),
1451 None => fmt!("{}", path.display()),
1452 };
1453 if name != was.name {
1454 return Err(err!(
1455 "The last read of {:?} took the segment {} where {} stands now. The \
1456 segments have been renumbered, which a log being added to never \
1457 does; the store cannot be taken up where it was left.",
1458 self.dir, was.name, name;
1459 Invalid, Input, Mismatch));
1460 }
1461 let (len, modified) = res!(Self::stamp(path));
1462 match was.reached {
1463 Reached::Sealed => {
1464 if len != was.len || modified != was.modified {
1465 return Err(err!(
1466 "The segment {:?} was read whole at {} bytes written {}, and \
1467 it is {} bytes written {} now. Nothing appends to a segment \
1468 that has been finished, so these are not the bytes that were \
1469 read; the store cannot be taken up where it was left.",
1470 path, was.len, was.modified, len, modified;
1471 Invalid, Data, Mismatch));
1472 }
1473 out.insert(name, Taken::Whole(was.clone()));
1474 },
1475 Reached::Growing(prefix) => {
1476 if len < was.len {
1477 return Err(err!(
1478 "The segment {:?} held {} bytes when it was read and holds {} \
1479 now. A log is only ever added to, so the store cannot be \
1480 taken up where it was left.",
1481 path, was.len, len;
1482 Invalid, Data, Mismatch));
1483 }
1484 // READ, not inferred from a timestamp. This is the one segment an
1485 // append may legitimately change, so a length and a modification
1486 // time say nothing here -- both move on every push. What the last
1487 // read left is a hash of the bytes it read, and the bytes are read
1488 // again and hashed again to meet it. They are bounded by
1489 // `SEGMENT_LIMIT`, because a segment past that is never the one
1490 // left growing.
1491 let sum = res!(Self::prefix_of(path, was.len));
1492 if sum.clone().finish() != prefix {
1493 return Err(err!(
1494 "The first {} bytes of the segment {:?} are not the bytes read \
1495 from it last time. A segment is only ever added to, so this \
1496 file has been rewritten rather than appended to, and it is a \
1497 different history; the store cannot be taken up where it was \
1498 left.",
1499 was.len, path;
1500 Invalid, Data, Mismatch));
1501 }
1502 out.insert(name, Taken::Part(was.clone(), sum));
1503 },
1504 }
1505 }
1506 Ok(out)
1507 }
1508
1509 /// Hashes the first `len` bytes of a file, and hands back the hasher so that
1510 /// the bytes after them can go into it too.
1511 fn prefix_of(path: &Path, len: u64)
1512 -> Outcome<Sha256>
1513 {
1514 const WINDOW: usize = 64 * 1024;
1515 let mut file = match fs::File::open(path) {
1516 Ok(f) => f,
1517 Err(e) => return Err(err!(e,
1518 "The segment {:?} could not be opened.", path; IO, File, Read)),
1519 };
1520 let mut sum = Sha256::new();
1521 let mut window = vec![0u8; WINDOW];
1522 let mut left = len;
1523 while left > 0 {
1524 let want = std::cmp::min(left as usize, WINDOW);
1525 let n = match std::io::Read::read(&mut file, &mut window[..want]) {
1526 Ok(n) => n,
1527 Err(e) => return Err(err!(e,
1528 "The segment {:?} could not be read.", path; IO, File, Read)),
1529 };
1530 if n == 0 {
1531 return Err(err!(
1532 "The segment {:?} ended {} bytes before the {} that were read from \
1533 it last time.", path, left, len;
1534 Invalid, Data, Mismatch));
1535 }
1536 sum.update(&window[..n]);
1537 left -= n as u64;
1538 }
1539 Ok(sum)
1540 }
1541
1542 /// The fold over a segment's records, which is what a verdict about it is
1543 /// filed under. See [`crate::verdict`].
1544 ///
1545 /// Reading a segment computes this on the way past, so nothing in a replay
1546 /// calls this; it is here for a caller that wants to say whether two segment
1547 /// files hold the same operations sealed the same way, which is the question
1548 /// a rewritten container has to answer.
1549 pub fn fold_of(path: &Path)
1550 -> Outcome<Fold>
1551 {
1552 let mut sum = Sha256::new();
1553 let head = res!(Self::stream_batches(path, Begin::Start, Integrity::Checked, None, 2_000,
1554 |entries, digests| {
1555 sum.update(digests);
1556 entries.clear();
1557 Ok(())
1558 }));
1559 Ok(Self::fold_from(sum, &head))
1560 }
1561
1562 /// Reads one segment the slow way -- every body hashed, every signature
1563 /// checked -- whatever any verdict says about it, and folds what it read.
1564 ///
1565 /// This is [`crate::sweep`]'s whole instrument, and it is deliberately the
1566 /// same code path a replay takes when it has no verdict to go on, entered
1567 /// with the vouching turned off. A sweep that read segments its own way would
1568 /// be checking something other than what the commands read.
1569 ///
1570 /// Nothing is kept. The records are dropped as they are made, so a sweep of a
1571 /// history far larger than memory costs one batch at a time.
1572 pub fn checked_fold(path: &Path, how: Verify, whence: &str)
1573 -> Outcome<Fold>
1574 {
1575 let asked = Asked { how, keep: Keep::Nothing, fold: true };
1576 let got = res!(Self::read_segment(path, Begin::Start, asked, false, None, whence));
1577 Ok(res!(got.fold.ok_or_else(|| err!(
1578 "The segment {:?} was read to be checked and no fold was taken of it.", path;
1579 Bug, Missing))))
1580 }
1581
1582 /// The size a segment file has and when it was last written, as the
1583 /// filesystem reports them.
1584 pub(crate) fn stamp(path: &Path)
1585 -> Outcome<(u64, u128)>
1586 {
1587 let meta = match fs::metadata(path) {
1588 Ok(m) => m,
1589 Err(e) => return Err(err!(e,
1590 "The segment {:?} could not be measured.", path;
1591 IO, File, Read)),
1592 };
1593 // A filesystem that will not say when a file was written is one where
1594 // nothing ever looks unchanged, so every segment is checked. Slow, and
1595 // right.
1596 let modified = match meta.modified() {
1597 Ok(t) => match t.duration_since(UNIX_EPOCH) {
1598 Ok(d) => d.as_nanos(),
1599 Err(_) => 0,
1600 },
1601 Err(_) => 0,
1602 };
1603 Ok((meta.len(), modified))
1604 }
1605
1606 /// Writes entries to the current segment, starting a new one once the
1607 /// current segment has passed [`SEGMENT_LIMIT`], or once what is being
1608 /// written is an operation that segment's format version has no code for.
1609 ///
1610 /// A segment already on disk is continued rather than rewritten: the engine
1611 /// reads what is there under the standard hasher and salt, and hands back the
1612 /// record bytes alone. A segment that could not be read to the end is
1613 /// therefore refused before anything is added to it.
1614 ///
1615 /// # Why size is not the only reason to start one
1616 ///
1617 /// The operation vocabulary grows upwards and a segment declares the version
1618 /// it was written at, so [`segment::Writer::push`] refuses a record whose wire
1619 /// code is above what that version spells -- which is what keeps an older
1620 /// version a genuine subset of a newer one rather than a promise the bytes
1621 /// break. The caller that refusal assumes is this one: a repository whose
1622 /// newest segment was written at version 3 meets its first
1623 /// [`oxedyne_fe2o3_ore::op::Op::Reverts`], and the record goes into a fresh
1624 /// segment at the current version instead of failing the command. Nothing
1625 /// already written is touched, and nothing about a log says its segments share
1626 /// a version.
1627 ///
1628 /// The hint is what a fresh segment's header carries, and it is a hint and
1629 /// nothing else: a reader may sort a directory of segments by it without
1630 /// opening them. Operations that arrived from elsewhere carry no one
1631 /// replica's name, so a sync passes `None` rather than putting one
1632 /// repository's name on somebody else's work.
1633 ///
1634 /// What form each entry takes is decided before it gets here: an operation
1635 /// that arrived sealed is written down sealed, so a store that passes history
1636 /// on passes the provenance with it.
1637 pub fn append(&self, entries: &[Entry], hint: Option<ReplicaId>)
1638 -> Outcome<()>
1639 {
1640 if entries.is_empty() {
1641 return Ok(());
1642 }
1643 // A stand-in is a carrier's note to itself that an operation is veiled, and
1644 // it is the one thing that must never be written down: the segment would
1645 // then hold the note instead of the operation, and the operation would be
1646 // gone. It cannot be recovered afterwards, so it is refused here, where the
1647 // veiled entry it displaced is still in the caller's hand.
1648 for entry in entries {
1649 if let Entry::Bare(rec) = entry {
1650 if is_standin(rec) {
1651 return Err(err!(
1652 "The stand-in for the veiled operation {} was about to be written \
1653 to the segments in place of the operation itself. A stand-in is \
1654 held in memory while a carrier places an operation it cannot read, \
1655 and what is written down is the veiled entry it stands for; \
1656 writing the stand-in would lose the operation for good.", rec.id();
1657 Bug, Invalid, Input, Data));
1658 }
1659 }
1660 }
1661 let segments = res!(self.segments());
1662 // The number a fresh segment takes, whether it is the first here or the one
1663 // after the newest.
1664 let next = match segments.last() {
1665 Some((n, _)) => n + 1,
1666 None => 0,
1667 };
1668 // The newest segment, where there is one under the size at which the next
1669 // append starts another.
1670 let carry_on = match segments.last() {
1671 Some((_, path)) => {
1672 let len = match fs::metadata(path) {
1673 Ok(m) => m.len(),
1674 Err(e) => return Err(err!(e,
1675 "The segment {:?} could not be measured.", path;
1676 IO, File, Read)),
1677 };
1678 match len >= SEGMENT_LIMIT {
1679 true => None,
1680 false => Some(path.clone()),
1681 }
1682 },
1683 None => None,
1684 };
1685 let mut resumed: Option<(PathBuf, segment::Writer<HashScheme, 8>)> = None;
1686 if let Some(path) = carry_on {
1687 let existing = match fs::read(&path) {
1688 Ok(b) => b,
1689 Err(e) => return Err(err!(e,
1690 "The segment {:?} could not be read.", path;
1691 IO, File, Read)),
1692 };
1693 let writer = match segment::Writer::resume(&existing, hasher(), SALT) {
1694 Ok(w) => w,
1695 Err(e) => return Err(err!(e,
1696 "The segment {:?} holds {} bytes and could not be continued.",
1697 path, existing.len();
1698 Invalid, Data, File)),
1699 };
1700 // A segment at the current version carries every code there is, so the
1701 // entries are not opened to ask. An older one is asked about, once, and
1702 // what it is asked is whether the highest code among them is one it
1703 // spells; a sealed entry has to be opened to answer that, which is why
1704 // the question is not put where the answer is already known.
1705 let known = writer.version() >= segment::VERSION
1706 || segment::highest_code(writer.version()) >= res!(top_code(entries));
1707 if known {
1708 resumed = Some((path, writer));
1709 }
1710 }
1711 let (path, mut writer) = match resumed {
1712 Some(pair) => pair,
1713 None => (
1714 self.segment_path(next),
1715 segment::Writer::new(&Head::new(hint), hasher(), SALT),
1716 ),
1717 };
1718 res!(writer.extend(entries));
1719 let bytes = writer.finish();
1720 let mut file = match fs::OpenOptions::new().create(true).append(true).open(&path) {
1721 Ok(f) => f,
1722 Err(e) => return Err(err!(e,
1723 "The segment {:?} could not be opened for appending.", path;
1724 IO, File, Write)),
1725 };
1726 match file.write_all(&bytes) {
1727 Ok(()) => (),
1728 Err(e) => return Err(err!(e,
1729 "{} bytes could not be written to the segment {:?}.",
1730 bytes.len(), path;
1731 IO, File, Write)),
1732 }
1733 match file.flush() {
1734 Ok(()) => Ok(()),
1735 Err(e) => Err(err!(e,
1736 "The segment {:?} could not be flushed.", path;
1737 IO, File, Write)),
1738 }
1739 }
1740}
1741
1742
1743#[cfg(test)]
1744mod tests {
1745 use super::*;
1746
1747 use crate::keys::Signing;
1748
1749 use oxedyne_fe2o3_ore::op::{
1750 Header,
1751 Op,
1752 };
1753
1754 use std::time::{
1755 SystemTime,
1756 UNIX_EPOCH,
1757 };
1758
1759 const MARKS: u64 = 1_600; // padded marks, enough to pass SEGMENT_LIMIT
1760
1761 /// A directory that removes itself however the test ends.
1762 struct Scratch {
1763 /// Where it is.
1764 path: PathBuf,
1765 }
1766
1767 impl Scratch {
1768 /// Makes a directory of its own.
1769 fn new(what: &str)
1770 -> Outcome<Self>
1771 {
1772 let stamp = res!(SystemTime::now().duration_since(UNIX_EPOCH));
1773 let path = std::env::temp_dir().join(fmt!(
1774 "ore_store_{}_{}_{}", what, std::process::id(), stamp.as_nanos(),
1775 ));
1776 res!(fs::create_dir_all(&path));
1777 Ok(Self { path })
1778 }
1779 }
1780
1781 impl Drop for Scratch {
1782 fn drop(&mut self) {
1783 let _ = fs::remove_dir_all(&self.path);
1784 }
1785 }
1786
1787 /// A caller that finds the lock held is always told what holds it.
1788 ///
1789 /// The file used to be created and then written to, so a second caller landing
1790 /// between the two read an empty file and was told the holder "left no note of
1791 /// itself" -- seen once in six concurrent marks on 2026-08-22. Nothing reads a
1792 /// note now except a caller the kernel has just refused, and the holder clears
1793 /// the note on the way out, so a blank one means one thing and says it.
1794 #[test]
1795 fn a_held_lock_always_says_what_holds_it() -> Outcome<()> {
1796 let scratch = res!(Scratch::new("lock_note"));
1797 let dir = std::sync::Arc::new(scratch.path.clone());
1798 let blank = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
1799 let refused = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
1800 let mut hands = Vec::new();
1801 for _ in 0..8 {
1802 let dir = std::sync::Arc::clone(&dir);
1803 let blank = std::sync::Arc::clone(&blank);
1804 let refused = std::sync::Arc::clone(&refused);
1805 hands.push(std::thread::spawn(move || {
1806 for _ in 0..400 {
1807 match Lock::take(&dir) {
1808 Ok(held) => drop(held),
1809 Err(e) => {
1810 let said = e.plain();
1811 refused.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1812 if said.contains("no note of itself") {
1813 blank.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1814 }
1815 },
1816 }
1817 }
1818 }));
1819 }
1820 for hand in hands {
1821 if hand.join().is_err() {
1822 return Err(err!("A thread taking the lock panicked."; Test, Invalid));
1823 }
1824 }
1825 let refused = refused.load(std::sync::atomic::Ordering::Relaxed);
1826 assert!(refused > 0,
1827 "nothing ever contended, so the fixture proves nothing about contention");
1828 assert_eq!(blank.load(std::sync::atomic::Ordering::Relaxed), 0,
1829 "{} of {} refusals could not say what held the lock",
1830 blank.load(std::sync::atomic::Ordering::Relaxed), refused);
1831 Ok(())
1832 }
1833
1834 /// A command that was killed does not brick the store: the next one takes the
1835 /// lock rather than telling somebody to delete a file by hand.
1836 ///
1837 /// What a killed process leaves behind is exactly what is fabricated here -- the
1838 /// file, a note naming a process that is gone, and no advisory lock, because the
1839 /// kernel drops that when the last handle closes however the process ended. A
1840 /// process is not killed in the test because the two states are one state and
1841 /// only one of them is deterministic.
1842 #[test]
1843 fn a_killed_command_does_not_brick_the_store() -> Outcome<()> {
1844 let scratch = res!(Scratch::new("lock_stale"));
1845 let path = Lock::path_of(&scratch.path);
1846 res!(fs::write(&path, b"process 4194304 since 1755000000\n"));
1847 let taken = match Lock::take(&scratch.path) {
1848 Ok(held) => held,
1849 Err(e) => return Err(err!(e,
1850 "A lock a killed command left behind still bricks the store."; Test)),
1851 };
1852 // And what was taken is a lock and not a shrug.
1853 match Lock::take(&scratch.path) {
1854 Ok(_) => return Err(err!("Two callers held one lock."; Test, Invalid)),
1855 Err(e) => assert!(e.plain().contains("Another `ore` command"),
1856 "a second caller was refused for the wrong reason: {}", e.plain()),
1857 }
1858 // The dead holder's name went with the lock, so nothing reads it as current.
1859 let now = res!(fs::read_to_string(&path));
1860 assert!(now.contains(&fmt!("process {} ", std::process::id())),
1861 "the lock does not name the process holding it: {:?}", now);
1862 drop(taken);
1863 let after = res!(fs::read_to_string(&path));
1864 assert!(after.trim().is_empty(),
1865 "a finished command left its name in the lock: {:?}", after);
1866 // The file staying is the point: it is where the lock is kept and not the
1867 // lock itself, so it is here and it stops nobody.
1868 assert!(path.is_file(), "the lock file was removed, which is the old race");
1869 drop(res!(Lock::take(&scratch.path)));
1870 Ok(())
1871 }
1872
1873 /// A store whose one segment is over [`SEGMENT_LIMIT`], so that nothing will
1874 /// ever append to it and a verdict may be recorded about it.
1875 ///
1876 /// The size is what makes it sealed and the size is the whole point: a
1877 /// segment under the limit is the tail, and the tail is never entered in the
1878 /// verdicts. Marks are used because a mark is the cheapest signed operation
1879 /// there is, and the names are padded so that the megabyte is reached by
1880 /// bytes rather than by signatures, which are what the fixture costs.
1881 fn a_sealed_store(dir: &Path, key: &Signing)
1882 -> Outcome<Store>
1883 {
1884 let replica = ReplicaId::new(1);
1885 let store = res!(Store::create(dir, Some(replica)));
1886 let mut entries = Vec::new();
1887 let mut parents = Vec::new();
1888 let pad = "x".repeat(500);
1889 for i in 1..=MARKS {
1890 let rec = Record::new(
1891 res!(Header::new(OpId::new(replica, i), parents.clone())),
1892 Op::Mark { name: fmt!("mark {} {}", i, pad), body: None, time: None },
1893 );
1894 parents = vec![OpId::new(replica, i)];
1895 entries.push(Entry::Sealed(res!(key.seal(&rec))));
1896 }
1897 res!(store.append(&entries, Some(replica)));
1898 let len = res!(fs::metadata(store.segment_path(0))).len();
1899 assert!(len >= SEGMENT_LIMIT,
1900 "the fixture must be over the segment limit or nothing is sealed: {} bytes", len);
1901 Ok(store)
1902 }
1903
1904 /// Rewrites the segment with one operation's payload altered, every record
1905 /// digest recomputed so the file is internally sound, and the length and the
1906 /// modification time put back as they were.
1907 ///
1908 /// This is what somebody who owns the disk can always do, and it is written
1909 /// to be as convincing as it can be: the file looks untouched from the
1910 /// outside and reads correctly from the inside. The only thing it cannot
1911 /// produce is a signature.
1912 fn tamper_in_place(path: &Path)
1913 -> Outcome<()>
1914 {
1915 let was = res!(fs::metadata(path));
1916 let when = res!(was.modified());
1917 let (head, entries) = res!(segment::decode(&res!(fs::read(path)), hasher(), SALT));
1918 let mut out = Vec::new();
1919 let mut done = false;
1920 for entry in &entries {
1921 match entry {
1922 Entry::Sealed(env) if !done => {
1923 done = true;
1924 let mut payload = env.payload().to_vec();
1925 let last = payload.len() - 1;
1926 payload[last] ^= 0x01;
1927 out.push(Entry::Sealed(Envelope::new(
1928 payload,
1929 env.signer().to_vec(),
1930 env.signature().to_vec(),
1931 )));
1932 },
1933 other => out.push(other.clone()),
1934 }
1935 }
1936 let mut writer = segment::Writer::new(&head, hasher(), SALT);
1937 res!(writer.extend(&out));
1938 res!(fs::write(path, writer.finish()));
1939 let now = res!(fs::metadata(path));
1940 assert_eq!(now.len(), was.len(),
1941 "the tamper must not change the length, or the length would catch it");
1942 let file = res!(fs::OpenOptions::new().write(true).open(path));
1943 res!(file.set_modified(when));
1944 Ok(())
1945 }
1946
1947 /// A cursor written down and read back is the cursor that was written.
1948 ///
1949 /// It is written down so that a reading kept between processes can say which
1950 /// bytes it was taken from, and a cursor that did not survive the round trip
1951 /// would say so about the wrong bytes. The store here has both kinds of
1952 /// segment in it -- one sealed and one still growing -- because the two are
1953 /// recorded differently and a form that carried only one of them would pass a
1954 /// test with only one.
1955 #[test]
1956 fn a_cursor_survives_being_written_down() -> Outcome<()> {
1957 let scratch = res!(Scratch::new("cursor_text"));
1958 let key = res!(Signing::mint(ReplicaId::new(1)));
1959 let store = res!(a_sealed_store(&scratch.path, &key));
1960 // One more operation, which starts the segment that is still growing.
1961 let replica = ReplicaId::new(1);
1962 let rec = Record::new(
1963 res!(Header::new(OpId::new(replica, MARKS + 1),
1964 vec![OpId::new(replica, MARKS)])),
1965 Op::Mark { name: fmt!("after"), body: None, time: None },
1966 );
1967 res!(store.append(&[Entry::Sealed(res!(key.seal(&rec)))], Some(replica)));
1968 let trust = Trust::of(&[key.binding()]);
1969
1970 let mut out = Replayed::new();
1971 let cursor = res!(store.replay_since(
1972 &Consumed::new(), &mut out, Verify::Signatures(&trust), Keep::Nothing));
1973 assert_eq!(cursor.segments(), 2, "the fixture must hold both kinds of segment");
1974 assert_eq!(cursor.records(), MARKS as usize + 1);
1975
1976 let lines = cursor.to_lines();
1977 for line in &lines {
1978 let word = match line.split_whitespace().next() {
1979 Some(w) => w,
1980 None => return Err(err!("A cursor wrote an empty line."; Test)),
1981 };
1982 assert!(Consumed::owns(word),
1983 "a cursor must own every line it writes, and it does not own {:?}", word);
1984 }
1985 let back = res!(Consumed::from_lines(lines.iter().map(|s| s.as_str()))
1986 .ok_or_else(|| err!("What a cursor wrote did not read back."; Test)));
1987 assert_eq!(back, cursor, "and it reads back as itself");
1988
1989 // And it is the same cursor to the store, not merely to `PartialEq`.
1990 assert!(matches!(res!(store.settled_at(&back)), Settled::Same),
1991 "a cursor read back from text still settles the store it was taken from");
1992 Ok(())
1993 }
1994
1995 /// Nothing a cursor cannot make sense of is an error, and none of it is a
1996 /// half-read cursor either.
1997 #[test]
1998 fn a_cursor_that_does_not_parse_is_no_cursor() -> Outcome<()> {
1999 for (what, lines) in [
2000 ("a word nothing here owns", vec!["stanza 000000.seg", "deferred 0"]),
2001 ("no count of what had to wait", vec!["seg 000000.seg 10 20 3 1 - sealed"]),
2002 ("two counts", vec!["deferred 0", "deferred 1"]),
2003 ("a length that is not a number", vec!["seg 000000.seg x 20 3 1 - sealed", "deferred 0"]),
2004 ("too few fields", vec!["seg 000000.seg 10 20 3 1 -", "deferred 0"]),
2005 ("a hash of the wrong length", vec!["seg 000000.seg 10 20 3 1 - AAAA=1", "deferred 0"]),
2006 ] {
2007 assert!(Consumed::from_lines(lines.iter().copied()).is_none(),
2008 "a cursor with {} is no cursor", what);
2009 }
2010 // And the shape those are variations of does read, so the test is about
2011 // the fault in each and not about the shape.
2012 assert!(Consumed::from_lines(
2013 ["seg 000000.seg 10 20 3 1 - sealed", "deferred 0"].iter().copied()).is_some(),
2014 "the shape they vary from must itself read");
2015 Ok(())
2016 }
2017
2018 /// A store nothing has touched settles at the cursor its own read produced,
2019 /// and one that has moved says how it moved rather than saying nothing.
2020 #[test]
2021 fn a_store_says_whether_it_is_where_a_reading_left_it() -> Outcome<()> {
2022 let scratch = res!(Scratch::new("settled"));
2023 let key = res!(Signing::mint(ReplicaId::new(1)));
2024 let store = res!(a_sealed_store(&scratch.path, &key));
2025 let trust = Trust::of(&[key.binding()]);
2026 let mut out = Replayed::new();
2027 let cursor = res!(store.replay_since(
2028 &Consumed::new(), &mut out, Verify::Signatures(&trust), Keep::Nothing));
2029 assert!(matches!(res!(store.settled_at(&cursor)), Settled::Same),
2030 "nothing has happened to the store since it was read");
2031
2032 // A cursor over nothing is not a reading of this store.
2033 match res!(store.settled_at(&Consumed::new())) {
2034 Settled::Moved(why) => assert!(why.contains("no segment"),
2035 "and it says which: {}", why),
2036 Settled::Same => return Err(err!(
2037 "A cursor over nothing settled a store holding {} operations.",
2038 MARKS; Test, Mismatch)),
2039 }
2040
2041 // One operation appended, which is what another process running a verb in
2042 // this store does.
2043 let replica = ReplicaId::new(1);
2044 let rec = Record::new(
2045 res!(Header::new(OpId::new(replica, MARKS + 1),
2046 vec![OpId::new(replica, MARKS)])),
2047 Op::Mark { name: fmt!("after"), body: None, time: None },
2048 );
2049 res!(store.append(&[Entry::Sealed(res!(key.seal(&rec)))], Some(replica)));
2050 match res!(store.settled_at(&cursor)) {
2051 Settled::Moved(why) => assert!(why.contains("segment"),
2052 "an append must be named for what it was: {}", why),
2053 Settled::Same => return Err(err!(
2054 "A store that has been appended to settled at a reading taken before \
2055 the append."; Test, Mismatch)),
2056 }
2057 Ok(())
2058 }
2059
2060 /// A verdict is about bytes, and a segment whose bytes have moved does not
2061 /// have one however untouched the file looks.
2062 ///
2063 /// **This is the test the cache rests on.** The tamper leaves the file the
2064 /// length it was and the modification time it was, so everything the cache
2065 /// uses to decide whether to skip the checking still says skip. What refuses
2066 /// it is the fold taken over the read: it is not a fold anything verified
2067 /// here, so the operations that read produced are thrown away and the segment
2068 /// is read again with every signature checked, which is where the forgery
2069 /// stops. Take the fold check out and this replay succeeds.
2070 #[test]
2071 fn a_verdict_is_about_the_bytes_and_not_the_file() -> Outcome<()> {
2072 let scratch = res!(Scratch::new("verdict_tamper"));
2073 let key = res!(Signing::mint(ReplicaId::new(1)));
2074 let store = res!(a_sealed_store(&scratch.path, &key));
2075 let trust = Trust::of(&[key.binding()]);
2076
2077 let done = res!(store.replay(Verify::Signatures(&trust), Keep::Nothing));
2078 assert_eq!(done.log.len(), MARKS as usize);
2079 let known = Verdicts::read(&store.dir);
2080 assert!(!known.is_empty(), "a checked replay of a sealed segment leaves a verdict");
2081
2082 let path = store.segment_path(0);
2083 let was = res!(fs::metadata(&path));
2084 res!(tamper_in_place(&path));
2085 let now = res!(fs::metadata(&path));
2086 assert_eq!(now.len(), was.len(), "the file is the length it was");
2087 assert_eq!(res!(now.modified()), res!(was.modified()), "and dated as it was");
2088 assert!(
2089 Verdicts::read(&store.dir).expected(
2090 "000000.seg", now.len(),
2091 res!(res!(now.modified()).duration_since(UNIX_EPOCH)).as_nanos(),
2092 verdict::now(), 0,
2093 ).is_some(),
2094 "so the cache still offers to skip the checking, which is the point",
2095 );
2096
2097 let refused = match store.replay(Verify::Signatures(&trust), Keep::Nothing) {
2098 Ok(_) => return Err(err!(
2099 "A tampered segment was replayed on a verdict taken over other bytes.";
2100 Test, Security)),
2101 Err(e) => fmt!("{}", e.plain()),
2102 };
2103 assert!(refused.contains("does not verify"), "and says why: {}", refused);
2104 Ok(())
2105 }
2106
2107 /// A verdict over exactly these bytes is what stops the checking, and it is
2108 /// worth exactly what the directory holding it is worth.
2109 ///
2110 /// The store is tampered with first and a verdict is then written by hand
2111 /// over the tampered bytes, which is what somebody who can write `.ore` can
2112 /// always do -- they can equally write `key` and `log/`. What it proves here
2113 /// is that the verdict is consulted at all: without it the same replay
2114 /// refuses the same store, which the test asserts first so that the two
2115 /// outcomes are told apart by the verdict and by nothing else.
2116 #[test]
2117 fn a_verdict_over_these_bytes_stops_the_checking() -> Outcome<()> {
2118 let scratch = res!(Scratch::new("verdict_skip"));
2119 let key = res!(Signing::mint(ReplicaId::new(1)));
2120 let store = res!(a_sealed_store(&scratch.path, &key));
2121 let trust = Trust::of(&[key.binding()]);
2122 res!(store.replay(Verify::Signatures(&trust), Keep::Nothing));
2123
2124 let path = store.segment_path(0);
2125 res!(tamper_in_place(&path));
2126 let _ = fs::remove_file(Verdicts::path_of(&store.dir));
2127 match store.replay(Verify::Signatures(&trust), Keep::Nothing) {
2128 Ok(_) => return Err(err!(
2129 "A tampered segment replayed with no verdict about it at all.";
2130 Test, Security)),
2131 Err(_) => (),
2132 }
2133
2134 let meta = res!(fs::metadata(&path));
2135 let modified = res!(res!(meta.modified()).duration_since(UNIX_EPOCH)).as_nanos();
2136 let mut said = Verdicts::new();
2137 said.record("000000.seg", meta.len(), modified, res!(Store::fold_of(&path)),
2138 verdict::now());
2139 res!(said.write(&store.dir));
2140
2141 let done = res!(store.replay(Verify::Signatures(&trust), Keep::Nothing));
2142 assert_eq!(done.log.len(), MARKS as usize,
2143 "the verdict was taken and nothing was checked");
2144 Ok(())
2145 }
2146
2147 /// The verdicts are beside the log and never in it, because everything in the
2148 /// log is something a peer is handed.
2149 #[test]
2150 fn the_verdicts_are_not_in_the_log() -> Outcome<()> {
2151 let scratch = res!(Scratch::new("verdict_place"));
2152 let key = res!(Signing::mint(ReplicaId::new(1)));
2153 let store = res!(a_sealed_store(&scratch.path, &key));
2154 let trust = Trust::of(&[key.binding()]);
2155 res!(store.replay(Verify::Signatures(&trust), Keep::Nothing));
2156
2157 assert!(Verdicts::path_of(&store.dir).is_file(), "it is written beside the log");
2158 // The directory itself is read, not `Store::segments`, which filters to
2159 // `.seg` and would say the same whatever else was in there. An earlier
2160 // version of this test asked `segments` and passed with the verdicts
2161 // written into `log/`.
2162 let mut inside: Vec<String> = Vec::new();
2163 for entry in res!(fs::read_dir(store.dir.join(LOG_DIR))) {
2164 inside.push(res!(entry).file_name().to_string_lossy().into_owned());
2165 }
2166 inside.sort();
2167 assert_eq!(inside, vec![fmt!("000000.seg")],
2168 "and the log holds segments and nothing else");
2169
2170 // Deleting it costs a verified replay and nothing else.
2171 res!(fs::remove_file(Verdicts::path_of(&store.dir)));
2172 let done = res!(store.replay(Verify::Signatures(&trust), Keep::Nothing));
2173 assert_eq!(done.log.len(), MARKS as usize);
2174 assert!(Verdicts::path_of(&store.dir).is_file(), "and it comes back");
2175 Ok(())
2176 }
2177
2178 /// The tail is not entered, because its bytes are still growing.
2179 ///
2180 /// A segment under [`SEGMENT_LIMIT`] is the one the next append continues, so
2181 /// a fold over it names a state it is about to leave and an entry would be
2182 /// written per command for ever. Nothing about it would be wrong -- the fold
2183 /// is over the bytes either way -- and it would be worth nothing.
2184 #[test]
2185 fn the_tail_is_left_out() -> Outcome<()> {
2186 let scratch = res!(Scratch::new("verdict_tail"));
2187 let key = res!(Signing::mint(ReplicaId::new(1)));
2188 let replica = ReplicaId::new(1);
2189 let store = res!(Store::create(&scratch.path, Some(replica)));
2190 let rec = Record::new(
2191 res!(Header::new(OpId::new(replica, 1), Vec::new())),
2192 Op::Mark { name: fmt!("start"), body: None, time: None },
2193 );
2194 res!(store.append(&[Entry::Sealed(res!(key.seal(&rec)))], Some(replica)));
2195 let len = res!(fs::metadata(store.segment_path(0))).len();
2196 assert!(len < SEGMENT_LIMIT, "the fixture must be under the limit: {} bytes", len);
2197
2198 let trust = Trust::of(&[key.binding()]);
2199 res!(store.replay(Verify::Signatures(&trust), Keep::Nothing));
2200 assert!(Verdicts::read(&store.dir).is_empty(),
2201 "nothing is recorded about a segment the next append will extend");
2202 Ok(())
2203 }
2204
2205 /// A relay checks nothing, so it has nothing to write down.
2206 #[test]
2207 fn a_replay_that_checks_nothing_leaves_no_verdict() -> Outcome<()> {
2208 let scratch = res!(Scratch::new("verdict_relay"));
2209 let key = res!(Signing::mint(ReplicaId::new(1)));
2210 let store = res!(a_sealed_store(&scratch.path, &key));
2211 res!(store.replay(Verify::Nothing, Keep::Nothing));
2212 assert!(!Verdicts::path_of(&store.dir).is_file(),
2213 "a reader that claims nothing records nothing");
2214 Ok(())
2215 }
2216
2217 /// A working copy refuses a forged operation and a relay serves it, which is
2218 /// the whole of "the relay is never an authority" in one assertion.
2219 ///
2220 /// The forgery is competent: the segment is rewritten whole so its digests
2221 /// match its records, which is what somebody who owns the disk can always do.
2222 /// The signature is the one thing they cannot produce, and it is what the
2223 /// working copy checks.
2224 #[test]
2225 fn a_forgery_stops_a_working_copy_and_not_a_relay() -> Outcome<()> {
2226 let scratch = res!(Scratch::new("verify"));
2227 let key = res!(Signing::mint(ReplicaId::new(1)));
2228 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2229 let rec = Record::new(
2230 res!(Header::new(OpId::new(ReplicaId::new(1), 1), Vec::new())),
2231 Op::Mark { name: fmt!("start"), body: None, time: None },
2232 );
2233 res!(store.append(&[Entry::Sealed(res!(key.seal(&rec)))], Some(ReplicaId::new(1))));
2234
2235 let trust = Trust::of(&[key.binding()]);
2236 let done = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes));
2237 assert_eq!(done.log.len(), 1);
2238 assert_eq!(done.prov.len(), 1, "a checked replay says what it found");
2239
2240 // Rewrite the one entry's payload and put the segment back together
2241 // properly, digests included.
2242 let (head, entries) = res!(segment::decode(
2243 &res!(fs::read(store.segment_path(0))), hasher(), SALT,
2244 ));
2245 let mut out = Vec::new();
2246 for entry in &entries {
2247 match entry {
2248 Entry::Sealed(env) => {
2249 let mut payload = env.payload().to_vec();
2250 let last = payload.len() - 1;
2251 payload[last] ^= 0x01;
2252 out.push(Entry::Sealed(Envelope::new(
2253 payload,
2254 env.signer().to_vec(),
2255 env.signature().to_vec(),
2256 )));
2257 },
2258 other => out.push(other.clone()),
2259 }
2260 }
2261 let mut writer = segment::Writer::new(&head, hasher(), SALT);
2262 res!(writer.extend(&out));
2263 res!(fs::write(store.segment_path(0), writer.finish()));
2264
2265 // The working copy refuses the whole log by name.
2266 let refused = match store.replay(Verify::Signatures(&trust), Keep::Envelopes) {
2267 Ok(_) => return Err(err!("A forged operation was replayed."; Test, Security)),
2268 Err(e) => fmt!("{}", e.plain()),
2269 };
2270 assert!(refused.contains("does not verify"), "and says why: {}", refused);
2271
2272 // The relay takes it as it stands, and claims nothing about it.
2273 let done = res!(store.replay(Verify::Nothing, Keep::Envelopes));
2274 assert_eq!(done.log.len(), 1, "a relay serves what it was given");
2275 assert!(done.prov.is_empty(), "and says nothing about who wrote it");
2276 assert_eq!(done.envelopes.len(), 1, "while keeping the seal to hand on");
2277 Ok(())
2278 }
2279
2280 /// A carrier's stand-in is refused at the segment rather than written down in
2281 /// place of the operation it stands for.
2282 ///
2283 /// This is the one way the veiled form could lose data rather than leak it. A
2284 /// relay holds a stand-in in memory while it places an operation it cannot
2285 /// read, and substitutes the veiled entry back for everything it writes or
2286 /// sends; if a path were ever added that missed the substitution, the segment
2287 /// would end up holding a mark saying "veiled operation r1:1" and the
2288 /// operation would be gone for good. So the store refuses it where the veiled
2289 /// entry it displaced is still in the caller's hand.
2290 #[test]
2291 fn a_standin_is_never_written_down() -> Outcome<()> {
2292 let scratch = res!(Scratch::new("standin"));
2293 let who = ReplicaId::new(1);
2294 let store = res!(Store::create(&scratch.path, Some(who)));
2295 let head = res!(Header::new(OpId::new(who, 1), Vec::new()));
2296 let refused = match store.append(&[Entry::Bare(crate::veil::standin(head))], Some(who)) {
2297 Ok(()) => return Err(err!(
2298 "A stand-in was written to the segments."; Test, Invalid)),
2299 Err(e) => fmt!("{}", e.plain()),
2300 };
2301 assert!(refused.contains("r1:1"),
2302 "the refusal names the operation that would have been lost: {}", refused);
2303 assert!(refused.contains("stand-in"), "and what it met: {}", refused);
2304 // An ordinary mark of its own is written without complaint, so what the
2305 // guard turns on is the stand-in and not the shape of a mark.
2306 let ordinary = Record::new(
2307 res!(Header::new(OpId::new(who, 1), Vec::new())),
2308 Op::Mark { name: fmt!("release 1.0"), body: None, time: None },
2309 );
2310 res!(store.append(&[Entry::Bare(ordinary)], Some(who)));
2311 assert_eq!(res!(store.replay(Verify::Nothing, Keep::Envelopes)).log.len(), 1);
2312 Ok(())
2313 }
2314
2315 /// Fills a store with sealed operations signed alternately by two keys.
2316 ///
2317 /// Two signers rather than one, because the batch decompresses a public key
2318 /// once per signer and keeps them in a cache: a segment signed by one key
2319 /// would not notice a cache that returned the wrong one.
2320 fn two_signed_run(dir: &Path, one: &Signing, two: &Signing, n: u64)
2321 -> Outcome<(Store, Vec<OpId>)>
2322 {
2323 let store = res!(Store::create(dir, Some(ReplicaId::new(1))));
2324 let mut entries = Vec::new();
2325 let mut ids = Vec::new();
2326 for i in 0..n {
2327 let id = OpId::new(ReplicaId::new(1), i + 1);
2328 ids.push(id);
2329 let rec = Record::new(
2330 res!(Header::new(id, Vec::new())),
2331 Op::Mark { name: fmt!("mark {}", i), body: None, time: None },
2332 );
2333 let who = if i % 2 == 0 { one } else { two };
2334 entries.push(Entry::Sealed(res!(who.seal(&rec))));
2335 }
2336 res!(store.append(&entries, Some(ReplicaId::new(1))));
2337 Ok((store, ids))
2338 }
2339
2340 /// One tampered signature among many sound ones is still refused by name.
2341 ///
2342 /// This is the property the batch is not allowed to cost. Checking a
2343 /// segment's signatures together can say only that something in it is wrong,
2344 /// so a refusal must fall back to checking them one at a time to find out
2345 /// what; the test spoils each position in turn and insists the refusal names
2346 /// that operation, names the file, and names no operation that is sound.
2347 #[test]
2348 fn a_tampered_signature_among_many_is_still_named() -> Outcome<()> {
2349 let total = 6u64;
2350 for spoiled in 0..total as usize {
2351 let scratch = res!(Scratch::new("batch"));
2352 let one = res!(Signing::mint(ReplicaId::new(1)));
2353 let two = res!(Signing::mint(ReplicaId::new(2)));
2354 let (store, ids) = res!(two_signed_run(&scratch.path, &one, &two, total));
2355 let trust = Trust::of(&[one.binding(), two.binding()]);
2356
2357 // Sound to begin with, and every operation attributed.
2358 let done = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes));
2359 assert_eq!(done.log.len(), total as usize);
2360 assert_eq!(done.prov.len(), total as usize);
2361 for id in &ids {
2362 assert_eq!(done.prov.get(id), Some(&Prov::Verified));
2363 }
2364
2365 // Rewrite the segment with one signature spoiled, digests and all,
2366 // which is what somebody who owns the disk can always do.
2367 let (head, found) = res!(segment::decode(
2368 &res!(fs::read(store.segment_path(0))), hasher(), SALT,
2369 ));
2370 let mut out = Vec::new();
2371 for (i, entry) in found.iter().enumerate() {
2372 match entry {
2373 Entry::Sealed(env) if i == spoiled => {
2374 // Within s rather than R, so the signature is still a
2375 // well formed one and the only thing wrong with it is
2376 // that it does not hold.
2377 let mut sig = env.signature().to_vec();
2378 sig[40] ^= 0x01;
2379 out.push(Entry::Sealed(Envelope::new(
2380 env.payload().to_vec(),
2381 env.signer().to_vec(),
2382 sig,
2383 )));
2384 },
2385 other => out.push(other.clone()),
2386 }
2387 }
2388 let mut writer = segment::Writer::new(&head, hasher(), SALT);
2389 res!(writer.extend(&out));
2390 res!(fs::write(store.segment_path(0), writer.finish()));
2391
2392 let refused = match store.replay(Verify::Signatures(&trust), Keep::Envelopes) {
2393 Ok(_) => return Err(err!(
2394 "The forged operation at position {} was replayed.", spoiled;
2395 Test, Security)),
2396 Err(e) => fmt!("{}", e.plain()),
2397 };
2398 assert!(refused.contains("does not verify"),
2399 "the refusal says why: {}", refused);
2400 let named = fmt!("{}", ids[spoiled]);
2401 assert!(refused.contains(&named),
2402 "the refusal names {} rather than saying only that something in the \
2403 segment failed: {}", named, refused);
2404 for (i, id) in ids.iter().enumerate() {
2405 if i != spoiled {
2406 assert!(!refused.contains(&fmt!("{}", id)),
2407 "the refusal names {}, whose signature is sound: {}", id, refused);
2408 }
2409 }
2410 assert!(refused.contains("000000.seg"),
2411 "the refusal names the file it was found in: {}", refused);
2412 }
2413 Ok(())
2414 }
2415
2416 /// A signature that holds under a key nobody here has seen is still not a
2417 /// failure, and is still marked as unattributed rather than verified.
2418 ///
2419 /// Checking a segment together answers one question -- is everything here
2420 /// sound? -- and it is not the question trust asks. The two must not be
2421 /// allowed to run into one another: a batch that held would be a poor reason
2422 /// to call an unknown key's work verified.
2423 #[test]
2424 fn an_unknown_key_still_verifies_and_is_still_unattributed() -> Outcome<()> {
2425 let scratch = res!(Scratch::new("unknown"));
2426 let one = res!(Signing::mint(ReplicaId::new(1)));
2427 let two = res!(Signing::mint(ReplicaId::new(2)));
2428 let (store, ids) = res!(two_signed_run(&scratch.path, &one, &two, 6));
2429 // Only the first key is known here.
2430 let trust = Trust::of(&[one.binding()]);
2431 let done = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes));
2432 assert_eq!(done.log.len(), 6);
2433 for (i, id) in ids.iter().enumerate() {
2434 let want = if i % 2 == 0 { Prov::Verified } else { Prov::Unknown };
2435 assert_eq!(done.prov.get(id), Some(&want),
2436 "the provenance of {} is not what the trust set says", id);
2437 }
2438 assert_eq!(done.envelopes.len(), 6, "the seals are kept to hand on");
2439 Ok(())
2440 }
2441
2442 /// A bare entry beside sealed ones is carried through the batch untouched.
2443 ///
2444 /// A run of entries is not all one kind, and only the sealed ones have a
2445 /// signature to check. The bare ones must come out in their place, in order,
2446 /// marked bare and carrying no envelope.
2447 #[test]
2448 fn bare_entries_pass_through_a_batch_in_their_place() -> Outcome<()> {
2449 let scratch = res!(Scratch::new("mixed"));
2450 let key = res!(Signing::mint(ReplicaId::new(1)));
2451 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2452 let mut entries = Vec::new();
2453 let mut ids = Vec::new();
2454 for i in 0..6u64 {
2455 let id = OpId::new(ReplicaId::new(1), i + 1);
2456 ids.push(id);
2457 let rec = Record::new(
2458 res!(Header::new(id, Vec::new())),
2459 Op::Mark { name: fmt!("mark {}", i), body: None, time: None },
2460 );
2461 entries.push(if i % 3 == 0 {
2462 Entry::Bare(rec)
2463 } else {
2464 Entry::Sealed(res!(key.seal(&rec)))
2465 });
2466 }
2467 res!(store.append(&entries, Some(ReplicaId::new(1))));
2468 let trust = Trust::of(&[key.binding()]);
2469 let done = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes));
2470 assert_eq!(done.log.len(), 6);
2471 for (i, id) in ids.iter().enumerate() {
2472 let want = if i % 3 == 0 { Prov::Bare } else { Prov::Verified };
2473 assert_eq!(done.prov.get(id), Some(&want));
2474 assert_eq!(done.envelopes.contains_key(id), i % 3 != 0,
2475 "only a sealed operation leaves an envelope behind");
2476 }
2477 // The order the log came out in is the order they went in.
2478 let order: Vec<OpId> = done.log.iter().map(|rec| rec.id()).collect();
2479 assert_eq!(order, ids);
2480 Ok(())
2481 }
2482
2483 /// An operation the newest segment's format version has no code for starts a
2484 /// fresh segment, and leaves the old one exactly as it was.
2485 ///
2486 /// This is the whole of what makes the vocabulary additive in practice rather
2487 /// than only on paper. Every repository written before a bump has a newest
2488 /// segment declaring the older version, and [`segment::Writer::push`] refuses
2489 /// a code that version does not spell; without a caller that answers the
2490 /// refusal by starting a segment, the first command using a new operation
2491 /// fails against every repository that already exists.
2492 ///
2493 /// The version 3 segment here is a real one and not a simulation: the header
2494 /// is written at version 3, the mark that goes into it is written through the
2495 /// ordinary append, and its bytes are compared before and after to say that
2496 /// rolling touched nothing that was already there.
2497 #[test]
2498 fn an_operation_an_old_segment_cannot_carry_starts_a_new_one() -> Outcome<()> {
2499 let scratch = res!(Scratch::new("roll"));
2500 let who = ReplicaId::new(1);
2501 let store = res!(Store::create(&scratch.path, Some(who)));
2502 // A repository as one written before the bump: the segment declares version
2503 // 3, whose vocabulary stops at `CODE_FILE_MODE`.
2504 res!(fs::write(
2505 store.segment_path(0),
2506 Head { version: 3, replica: Some(who) }.encode(),
2507 ));
2508
2509 // A mark carrying neither a body nor a time is wire code 4, which version 3
2510 // spells, so it continues the segment that is there.
2511 let plain = Record::new(
2512 res!(Header::new(OpId::new(who, 1), Vec::new())),
2513 Op::Mark { name: fmt!("v3"), body: None, time: None },
2514 );
2515 res!(store.append(&[Entry::Bare(plain)], Some(who)));
2516 assert_eq!(res!(store.segments()).len(), 1,
2517 "a code version 3 spells does not start a segment");
2518 let before = res!(fs::read(store.segment_path(0)));
2519 assert_eq!(before[6], 3, "the segment being appended to declares version 3");
2520
2521 // A mark carrying a time is wire code 9, which version 3 does not spell.
2522 let timed = Record::new(
2523 res!(Header::new(OpId::new(who, 2), vec![OpId::new(who, 1)])),
2524 Op::Mark { name: fmt!("v4"), body: None, time: Some(1_755_000_000) },
2525 );
2526 res!(store.append(&[Entry::Bare(timed)], Some(who)));
2527 let segments = res!(store.segments());
2528 assert_eq!(segments.len(), 2,
2529 "a code version 3 does not spell starts a segment: {:?}", segments);
2530 assert_eq!(res!(fs::read(store.segment_path(0))), before,
2531 "the version 3 segment is left byte for byte as it was");
2532 let after = res!(fs::read(store.segment_path(1)));
2533 assert_eq!(after[6], segment::VERSION,
2534 "and the fresh one declares the current version");
2535
2536 // Rolling is the answer to the version and not to every append: the fresh
2537 // segment carries the current version, so the next operation continues it.
2538 let more = Record::new(
2539 res!(Header::new(OpId::new(who, 3), vec![OpId::new(who, 2)])),
2540 Op::Mark { name: fmt!("after"), body: None, time: None },
2541 );
2542 res!(store.append(&[Entry::Bare(more)], Some(who)));
2543 assert_eq!(res!(store.segments()).len(), 2,
2544 "a segment at the current version is continued, not replaced");
2545
2546 // And the history reads back whole, across the two versions, in order.
2547 let done = res!(store.replay(Verify::Nothing, Keep::Envelopes));
2548 let names: Vec<String> = done.log.iter()
2549 .filter_map(|rec| match &rec.op {
2550 Op::Mark { name, .. } => Some(name.clone()),
2551 _ => None,
2552 })
2553 .collect();
2554 assert_eq!(names, vec![fmt!("v3"), fmt!("v4"), fmt!("after")],
2555 "the log spans both segments");
2556 let timed = res!(done.log.get(&OpId::new(who, 2))
2557 .ok_or_else(|| err!("The rolled operation is not in the log."; Test, Missing)));
2558 assert!(matches!(timed.op, Op::Mark { time: Some(1_755_000_000), .. }),
2559 "and the time it was written with came back: {:?}", timed.op);
2560 Ok(())
2561 }
2562 // -- Taking a read up where the last one stopped -------------------------
2563
2564 /// Appends `n` sealed marks, chained one to the next from `at`, and returns
2565 /// the counter to carry on from.
2566 ///
2567 /// One `append` a call, which is what a push is: the store is read, the
2568 /// writer is resumed over the newest segment and the records go on the end.
2569 fn push(store: &Store, key: &Signing, at: u64, n: u64, pad: usize)
2570 -> Outcome<u64>
2571 {
2572 let replica = ReplicaId::new(1);
2573 let filling = "y".repeat(pad);
2574 let mut entries = Vec::new();
2575 let mut parents = match at > 1 {
2576 true => vec![OpId::new(replica, at - 1)],
2577 false => Vec::new(),
2578 };
2579 for i in at..at + n {
2580 let rec = Record::new(
2581 res!(Header::new(OpId::new(replica, i), parents.clone())),
2582 Op::Mark { name: fmt!("mark {} {}", i, filling), body: None, time: None },
2583 );
2584 parents = vec![OpId::new(replica, i)];
2585 entries.push(Entry::Sealed(res!(key.seal(&rec))));
2586 }
2587 res!(store.append(&entries, Some(replica)));
2588 Ok(at + n)
2589 }
2590
2591 /// Fails unless two replays produced the same thing, in the same order.
2592 fn alike(whole: &Replayed, pieces: &Replayed, when: &str)
2593 -> Outcome<()>
2594 {
2595 let one: Vec<&Record> = whole.log.iter().collect();
2596 let two: Vec<&Record> = pieces.log.iter().collect();
2597 if one != two {
2598 let first = one.iter().zip(two.iter())
2599 .position(|(a, b)| a != b)
2600 .map(|at| fmt!("they first differ at position {}", at))
2601 .unwrap_or_else(|| fmt!("one is {} long and the other {}",
2602 one.len(), two.len()));
2603 return Err(err!(
2604 "{}: a replay in pieces did not produce the log a whole replay \
2605 produced; {}.", when, first;
2606 Test, Mismatch));
2607 }
2608 if whole.prov != pieces.prov {
2609 return Err(err!(
2610 "{}: a replay in pieces knows different provenance from a whole one.",
2611 when;
2612 Test, Mismatch));
2613 }
2614 if whole.envelopes != pieces.envelopes {
2615 return Err(err!(
2616 "{}: a replay in pieces kept different envelopes from a whole one.", when;
2617 Test, Mismatch));
2618 }
2619 if whole.veils != pieces.veils {
2620 return Err(err!(
2621 "{}: a replay in pieces kept different veils from a whole one.", when;
2622 Test, Mismatch));
2623 }
2624 Ok(())
2625 }
2626
2627 /// A store growing at its end is read where the last read stopped, and what
2628 /// comes out is what a whole replay produces.
2629 ///
2630 /// **The differential is the test.** Twelve pushes, read one at a time,
2631 /// against the whole store read afresh after each of them -- log order,
2632 /// provenance, envelopes and veils, every time. The sizes are chosen so that
2633 /// the run crosses [`SEGMENT_LIMIT`] twice: most pushes grow the newest
2634 /// segment in place, two of them take it past the limit, and the push after
2635 /// each of those starts a new file. So the mid-segment case, the rollover
2636 /// case and the two together are all in here, and the named tests below say
2637 /// which push was which.
2638 #[test]
2639 fn a_replay_in_pieces_produces_what_a_whole_replay_produces() -> Outcome<()> {
2640 let scratch = res!(Scratch::new("resume_differential"));
2641 let key = res!(Signing::mint(ReplicaId::new(1)));
2642 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2643 let trust = Trust::of(&[key.binding()]);
2644 let mut pieces = Replayed::new();
2645 let mut cursor = Consumed::new();
2646 let mut at = 1u64;
2647 let mut segments = 0usize;
2648 let mut rollovers = 0usize;
2649 let mut grew = 0usize;
2650 // Small, small, large enough to pass the limit, small again, and round.
2651 let run: [(u64, usize); 12] = [
2652 (3, 40), (5, 40), (130, 9_000), (2, 40), (4, 40), (7, 40),
2653 (130, 9_000), (3, 40), (2, 40), (9, 40), (1, 40), (4, 40),
2654 ];
2655 for (n, pad) in run {
2656 let before = res!(store.segments()).len();
2657 let held = cursor.clone();
2658 at = res!(push(&store, &key, at, n, pad));
2659 let after = res!(store.segments()).len();
2660 cursor = res!(store.replay_since(
2661 &cursor, &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2662 let whole = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes));
2663 res!(alike(&whole, &pieces, &fmt!("after {} operations", at - 1)));
2664 assert_eq!(cursor.records(), whole.log.len(),
2665 "the cursor counts operations the read did not produce");
2666 assert_eq!(cursor.segments(), after,
2667 "the cursor names segments the directory does not hold");
2668 match after > before {
2669 true => rollovers += 1,
2670 false => grew += 1,
2671 }
2672 if !held.is_empty() && held.segments() < after {
2673 segments += 1;
2674 }
2675 }
2676 assert!(rollovers >= 2, "the run must roll over, and it rolled over {} times",
2677 rollovers);
2678 assert!(grew >= 8, "the run must grow segments in place, and it grew {} times",
2679 grew);
2680 assert!(segments >= 2, "the run must resume across a new segment, and it did {}",
2681 segments);
2682 Ok(())
2683 }
2684
2685 /// A segment that has grown since it was read is taken up at the byte the
2686 /// last read stopped at, and the operations in front of that byte are not
2687 /// read again.
2688 #[test]
2689 fn a_read_takes_up_a_segment_that_grew_in_place() -> Outcome<()> {
2690 let scratch = res!(Scratch::new("resume_grew"));
2691 let key = res!(Signing::mint(ReplicaId::new(1)));
2692 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2693 let trust = Trust::of(&[key.binding()]);
2694 let mut pieces = Replayed::new();
2695
2696 let at = res!(push(&store, &key, 1, 40, 40));
2697 let first = res!(store.replay_since(
2698 &Consumed::new(), &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2699 assert_eq!(first.segments(), 1, "one push writes one segment");
2700 assert_eq!(first.records(), 40);
2701
2702 res!(push(&store, &key, at, 25, 40));
2703 assert_eq!(res!(store.segments()).len(), 1,
2704 "the second push must land in the same file, or this is the other test");
2705 let then = res!(store.replay_since(
2706 &first, &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2707 assert_eq!(then.segments(), 1);
2708 assert_eq!(then.records(), 65);
2709 assert!(then.bytes() > first.bytes(), "the file grew and the cursor did not");
2710 let whole = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes));
2711 res!(alike(&whole, &pieces, "after a segment grew in place"));
2712 Ok(())
2713 }
2714
2715 /// A store that has started a new segment since it was read reads the new
2716 /// file and leaves the finished one alone.
2717 #[test]
2718 fn a_read_takes_up_a_store_that_rolled_over() -> Outcome<()> {
2719 let scratch = res!(Scratch::new("resume_rolled"));
2720 let key = res!(Signing::mint(ReplicaId::new(1)));
2721 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2722 let trust = Trust::of(&[key.binding()]);
2723 let mut pieces = Replayed::new();
2724
2725 // Past the limit in one push, so that the first read finds the newest
2726 // segment already finished and the next push has to start another.
2727 let at = res!(push(&store, &key, 1, 130, 9_000));
2728 assert!(res!(fs::metadata(store.segment_path(0))).len() >= SEGMENT_LIMIT,
2729 "the first push must pass the segment limit or nothing rolls over");
2730 let first = res!(store.replay_since(
2731 &Consumed::new(), &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2732 assert_eq!(first.segments(), 1);
2733
2734 res!(push(&store, &key, at, 6, 40));
2735 assert_eq!(res!(store.segments()).len(), 2, "the push must start a second file");
2736 let then = res!(store.replay_since(
2737 &first, &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2738 assert_eq!(then.segments(), 2);
2739 assert_eq!(then.records(), 136);
2740 let whole = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes));
2741 res!(alike(&whole, &pieces, "after a store rolled over"));
2742 Ok(())
2743 }
2744
2745 /// A segment that grew past the limit and a new segment beside it, in one
2746 /// read: the finished file is taken up part way and the new one is read
2747 /// whole.
2748 #[test]
2749 fn a_read_takes_up_growth_and_a_rollover_together() -> Outcome<()> {
2750 let scratch = res!(Scratch::new("resume_both"));
2751 let key = res!(Signing::mint(ReplicaId::new(1)));
2752 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2753 let trust = Trust::of(&[key.binding()]);
2754 let mut pieces = Replayed::new();
2755
2756 let mut at = res!(push(&store, &key, 1, 12, 40));
2757 let first = res!(store.replay_since(
2758 &Consumed::new(), &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2759 assert_eq!(first.segments(), 1);
2760 assert!(res!(fs::metadata(store.segment_path(0))).len() < SEGMENT_LIMIT,
2761 "the first read must leave the newest segment still growing");
2762
2763 // One push takes the segment it is read up in past the limit, and the
2764 // next starts a file of its own -- so the read that follows has both to
2765 // do at once.
2766 at = res!(push(&store, &key, at, 130, 9_000));
2767 res!(push(&store, &key, at, 5, 40));
2768 assert_eq!(res!(store.segments()).len(), 2);
2769 assert!(res!(fs::metadata(store.segment_path(0))).len() >= SEGMENT_LIMIT);
2770
2771 let then = res!(store.replay_since(
2772 &first, &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2773 assert_eq!(then.segments(), 2);
2774 assert_eq!(then.records(), 147);
2775 let whole = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes));
2776 res!(alike(&whole, &pieces, "after growth and a rollover together"));
2777 Ok(())
2778 }
2779
2780 /// Nothing already read is read again, and the proof is that a read succeeds
2781 /// where the bytes already read cannot be opened at all.
2782 ///
2783 /// The finished segment is made unreadable, which is what a whole replay
2784 /// needs and a resumed one must not: the same store, at the same moment,
2785 /// refuses one and answers the other. Its length and its modification time
2786 /// are what the cursor checks, and a file with no read permission still has
2787 /// both.
2788 #[test]
2789 fn a_resumed_read_does_not_open_what_it_has_already_read() -> Outcome<()> {
2790 use std::os::unix::fs::PermissionsExt;
2791
2792 let scratch = res!(Scratch::new("resume_unopened"));
2793 let key = res!(Signing::mint(ReplicaId::new(1)));
2794 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2795 let trust = Trust::of(&[key.binding()]);
2796 let mut pieces = Replayed::new();
2797
2798 let at = res!(push(&store, &key, 1, 130, 9_000));
2799 let first = res!(store.replay_since(
2800 &Consumed::new(), &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2801 res!(push(&store, &key, at, 6, 40));
2802 let whole = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes));
2803
2804 let sealed = store.segment_path(0);
2805 let was = res!(fs::metadata(&sealed)).permissions();
2806 res!(fs::set_permissions(&sealed, fs::Permissions::from_mode(0o000)));
2807 let shut = store.replay(Verify::Signatures(&trust), Keep::Envelopes).is_err();
2808 let took = store.replay_since(
2809 &first, &mut pieces, Verify::Signatures(&trust), Keep::Envelopes);
2810 res!(fs::set_permissions(&sealed, was));
2811
2812 assert!(shut, "a whole replay must need the segment this one must not open");
2813 let then = res!(took);
2814 assert_eq!(then.records(), 136);
2815 res!(alike(&whole, &pieces, "reading only what arrived"));
2816 Ok(())
2817 }
2818
2819 /// A segment rewritten where it had already been read refuses to be taken up,
2820 /// however untouched the file looks.
2821 ///
2822 /// **This is the test the resume rests on.** The tamper puts the length and
2823 /// the modification time back, so everything a timestamp could say still says
2824 /// the file has not moved -- exactly as it does in
2825 /// [`a_verdict_is_about_the_bytes_and_not_the_file`]. What refuses it is the
2826 /// hash the last read took of the bytes it read, met by reading those bytes
2827 /// again. Take that comparison out and the read carries a log the store no
2828 /// longer holds.
2829 #[test]
2830 fn a_segment_rewritten_under_a_cursor_refuses_to_be_taken_up() -> Outcome<()> {
2831 let scratch = res!(Scratch::new("resume_rewritten"));
2832 let key = res!(Signing::mint(ReplicaId::new(1)));
2833 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2834 let trust = Trust::of(&[key.binding()]);
2835 let mut pieces = Replayed::new();
2836
2837 res!(push(&store, &key, 1, 40, 40));
2838 let first = res!(store.replay_since(
2839 &Consumed::new(), &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2840
2841 let path = store.segment_path(0);
2842 let was = res!(fs::metadata(&path));
2843 res!(tamper_in_place(&path));
2844 let now = res!(fs::metadata(&path));
2845 assert_eq!(now.len(), was.len(), "the file is the length it was");
2846 assert_eq!(res!(now.modified()), res!(was.modified()), "and dated as it was");
2847
2848 let refused = match store.replay_since(
2849 &first, &mut pieces, Verify::Signatures(&trust), Keep::Envelopes)
2850 {
2851 Ok(_) => return Err(err!(
2852 "A segment rewritten where it had been read was taken up as though it \
2853 had only grown.";
2854 Test, Security)),
2855 Err(e) => fmt!("{}", e.plain()),
2856 };
2857 assert!(refused.contains("rewritten rather than appended to"),
2858 "and says what it found: {}", refused);
2859 Ok(())
2860 }
2861
2862 /// A finished segment whose file has moved refuses to be taken up, and a
2863 /// shorter one does too.
2864 ///
2865 /// Nothing ever opens a finished segment again, so a length or a modification
2866 /// time that is not the one it was read at is not an append: it is a restored
2867 /// backup, a copy, a repack put back by hand, or somebody at the disk. Each
2868 /// is a different store and none may be read as a continuation of this one.
2869 #[test]
2870 fn a_finished_segment_that_moved_refuses_to_be_taken_up() -> Outcome<()> {
2871 let scratch = res!(Scratch::new("resume_moved"));
2872 let key = res!(Signing::mint(ReplicaId::new(1)));
2873 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2874 let trust = Trust::of(&[key.binding()]);
2875 let mut pieces = Replayed::new();
2876
2877 let at = res!(push(&store, &key, 1, 130, 9_000));
2878 res!(push(&store, &key, at, 6, 40));
2879 let first = res!(store.replay_since(
2880 &Consumed::new(), &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2881 assert_eq!(first.segments(), 2);
2882
2883 // The finished segment, dated a second later and otherwise untouched.
2884 let path = store.segment_path(0);
2885 let when = res!(res!(fs::metadata(&path)).modified());
2886 let file = res!(fs::OpenOptions::new().write(true).open(&path));
2887 res!(file.set_modified(when + std::time::Duration::from_secs(1)));
2888 let dated = match store.replay_since(
2889 &first, &mut pieces, Verify::Signatures(&trust), Keep::Envelopes)
2890 {
2891 Ok(_) => return Err(err!(
2892 "A finished segment written since it was read was taken up as though \
2893 nothing had happened to it.";
2894 Test, Security)),
2895 Err(e) => fmt!("{}", e.plain()),
2896 };
2897 assert!(dated.contains("Nothing appends to a segment that has been finished"),
2898 "and says why that cannot be an append: {}", dated);
2899 res!(file.set_modified(when));
2900
2901 // The one segment that may grow, shorter than it was.
2902 let tail = store.segment_path(1);
2903 let bytes = res!(fs::read(&tail));
2904 res!(fs::write(&tail, &bytes[..bytes.len() - 1]));
2905 let short = match store.replay_since(
2906 &first, &mut pieces, Verify::Signatures(&trust), Keep::Envelopes)
2907 {
2908 Ok(_) => return Err(err!(
2909 "A segment shorter than it was read at was taken up as though it had \
2910 grown.";
2911 Test, Security)),
2912 Err(e) => fmt!("{}", e.plain()),
2913 };
2914 assert!(short.contains("A log is only ever added to"),
2915 "and says why a shorter file is not one: {}", short);
2916 Ok(())
2917 }
2918
2919 /// A resumed read leaves the verdicts it was already holding, including the
2920 /// ones about segments it did not open.
2921 ///
2922 /// A verdict says every signature in a segment verified here, and a read that
2923 /// does not open the segment has learned nothing that could contradict it.
2924 /// What it must not do is drop it: a whole read rebuilds the file from
2925 /// nothing and a note about a segment that has gone falls out, which is right
2926 /// for a read that looked at every segment and wrong for one that looked at
2927 /// two.
2928 #[test]
2929 fn a_resumed_read_keeps_the_verdicts_it_did_not_look_at() -> Outcome<()> {
2930 let scratch = res!(Scratch::new("resume_verdicts"));
2931 let key = res!(Signing::mint(ReplicaId::new(1)));
2932 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2933 let trust = Trust::of(&[key.binding()]);
2934 let mut pieces = Replayed::new();
2935
2936 /// Is there a verdict recorded about the segment as it stands on disk?
2937 fn verdict_at(store: &Store, n: u64)
2938 -> Outcome<bool>
2939 {
2940 let path = store.segment_path(n);
2941 let (len, modified) = res!(Store::stamp(&path));
2942 let name = fmt!("{:06}{}", n, SEGMENT_SUFFIX);
2943 Ok(Verdicts::read(&store.dir).expected(&name, len, modified, verdict::now(), 0)
2944 .is_some())
2945 }
2946
2947 let mut at = res!(push(&store, &key, 1, 130, 9_000));
2948 let first = res!(store.replay_since(
2949 &Consumed::new(), &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2950 assert_eq!(first.segments(), 1);
2951 assert!(res!(verdict_at(&store, 0)),
2952 "the first read must leave a verdict about the finished segment");
2953
2954 // A whole segment and the start of another, so that the resumed read has a
2955 // verdict of its own to write and therefore rewrites the file. A segment
2956 // the resumed read finished gets no verdict of its own -- a fold is over
2957 // the whole of a file and a resumed read has only its tail -- so the write
2958 // has to be made to happen by a segment it read from the first byte.
2959 at = res!(push(&store, &key, at, 130, 9_000));
2960 res!(push(&store, &key, at, 4, 40));
2961 assert_eq!(res!(store.segments()).len(), 3);
2962 res!(store.replay_since(
2963 &first, &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
2964 assert!(res!(verdict_at(&store, 1)),
2965 "the resumed read learned a verdict about the segment it read whole");
2966 assert!(res!(verdict_at(&store, 0)),
2967 "and kept the one about the segment it never opened");
2968 Ok(())
2969 }
2970
2971 /// A cursor whose read met an operation before its parents refuses to be
2972 /// taken up, because the log it left is not in the order a whole replay
2973 /// would have put it in.
2974 #[test]
2975 fn a_read_that_met_an_operation_before_its_parents_refuses_to_be_taken_up()
2976 -> Outcome<()>
2977 {
2978 let scratch = res!(Scratch::new("resume_out_of_order"));
2979 let key = res!(Signing::mint(ReplicaId::new(1)));
2980 let store = res!(Store::create(&scratch.path, Some(ReplicaId::new(1))));
2981 let trust = Trust::of(&[key.binding()]);
2982 let replica = ReplicaId::new(2);
2983 // A child written in front of its parent, which nothing that writes an Ore
2984 // segment does and which the format does not forbid.
2985 let parent = Record::new(
2986 res!(Header::new(OpId::new(replica, 1), Vec::new())),
2987 Op::Mark { name: fmt!("the parent"), body: None, time: None },
2988 );
2989 let child = Record::new(
2990 res!(Header::new(OpId::new(ReplicaId::new(3), 2), vec![OpId::new(replica, 1)])),
2991 Op::Mark { name: fmt!("the child"), body: None, time: None },
2992 );
2993 res!(store.append(&[
2994 Entry::Sealed(res!(key.seal(&child))),
2995 Entry::Sealed(res!(key.seal(&parent))),
2996 ], None));
2997
2998 let mut pieces = Replayed::new();
2999 let first = res!(store.replay_since(
3000 &Consumed::new(), &mut pieces, Verify::Signatures(&trust), Keep::Envelopes));
3001 assert!(!first.resumable(), "the read met a child before its parent");
3002 let refused = match store.replay_since(
3003 &first, &mut pieces, Verify::Signatures(&trust), Keep::Envelopes)
3004 {
3005 Ok(_) => return Err(err!(
3006 "A cursor over a log placed out of order was taken up as though the \
3007 order were the one a whole replay gives.";
3008 Test, Invalid)),
3009 Err(e) => fmt!("{}", e.plain()),
3010 };
3011 assert!(refused.contains("Replay the store whole"),
3012 "and says what to do instead: {}", refused);
3013 Ok(())
3014 }
3015
3016}