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 | |
| 46 | use crate::keys::{ |
| 47 | self, |
| 48 | Prov, |
| 49 | Trust, |
| 50 | }; |
| 51 | use crate::repack; |
| 52 | use crate::veil::{ |
| 53 | is_standin, |
| 54 | standin, |
| 55 | }; |
| 56 | use crate::verdict::{ |
| 57 | self, |
| 58 | Fold, |
| 59 | Verdicts, |
| 60 | VERDICT_LIFE, |
| 61 | }; |
| 62 | |
| 63 | use oxedyne_fe2o3_core::prelude::*; |
| 64 | use oxedyne_fe2o3_hash::hash::HashScheme; |
| 65 | use oxedyne_fe2o3_hash::sha256::{ |
| 66 | DIGEST_LEN, |
| 67 | Sha256, |
| 68 | }; |
| 69 | use oxedyne_fe2o3_ore::envelope::Envelope; |
| 70 | use oxedyne_fe2o3_ore::id::{ |
| 71 | OpId, |
| 72 | ReplicaId, |
| 73 | }; |
| 74 | use oxedyne_fe2o3_ore::log::OpLog; |
| 75 | use oxedyne_fe2o3_ore::op::Record; |
| 76 | use oxedyne_fe2o3_ore::segment::{ |
| 77 | self, |
| 78 | Entry, |
| 79 | Head, |
| 80 | Integrity, |
| 81 | Veiled, |
| 82 | }; |
| 83 | |
| 84 | use std::collections::{ |
| 85 | BTreeMap, |
| 86 | BTreeSet, |
| 87 | }; |
| 88 | use std::fs; |
| 89 | use std::io::{ |
| 90 | Seek, |
| 91 | Write, |
| 92 | }; |
| 93 | use std::path::{ |
| 94 | Path, |
| 95 | PathBuf, |
| 96 | }; |
| 97 | use std::time::{ |
| 98 | SystemTime, |
| 99 | UNIX_EPOCH, |
| 100 | }; |
| 101 | |
| 102 | |
| 103 | /// Name of the segment directory, within the store directory. |
| 104 | pub const LOG_DIR: &str = "log"; |
| 105 | /// Name of the lock file, within the store directory. |
| 106 | pub const LOCK_FILE: &str = "lock"; |
| 107 | |
| 108 | /// The salt every segment digest is computed under. Part of the format. |
| 109 | |
| 110 | pub const SALT: [u8; 8] = *b"oreseg01"; |
| 111 | |
| 112 | /// Size a segment file must reach before the next append starts a new one. |
| 113 | pub const SEGMENT_LIMIT: u64 = 1 << 20; |
| 114 | |
| 115 | /// Suffix every segment file carries. |
| 116 | pub const SEGMENT_SUFFIX: &str = ".seg"; |
| 117 | |
| 118 | // The words a cursor's own lines begin with. See [`Consumed::owns`]. |
| 119 | const SEG_LINE: &str = "seg"; |
| 120 | const DEFERRED_LINE: &str = "deferred"; |
| 121 | const SEALED_WORD: &str = "sealed"; |
| 122 | |
| 123 | |
| 124 | /// Returns the integrity hasher every segment of every Ore repository is |
| 125 | /// written and checked with. |
| 126 | pub 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. |
| 139 | pub 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. |
| 165 | pub 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. |
| 182 | pub 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`]. |
| 203 | fn 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. |
| 259 | pub struct Lock { |
| 260 | file: fs::File, // the advisory lock lives on this handle and dies with it |
| 261 | } |
| 262 | |
| 263 | impl 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 | |
| 325 | impl 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)] |
| 337 | pub 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)] |
| 357 | pub 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)] |
| 371 | pub(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)] |
| 386 | enum 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)] |
| 402 | struct 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)] |
| 447 | pub 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 | |
| 452 | impl 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. |
| 616 | pub 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)] |
| 633 | pub 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 | |
| 650 | impl 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 | |
| 663 | impl 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)] |
| 677 | pub 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)] |
| 691 | struct 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. |
| 699 | enum 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. |
| 711 | struct 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)] |
| 728 | pub struct Store { |
| 729 | /// The directory holding `log/`. |
| 730 | pub dir: PathBuf, |
| 731 | } |
| 732 | |
| 733 | impl 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)] |
| 1744 | mod 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 | } |