oxedyne/ore/store/src/repack.rs
42.7 KiB, 25 runs
created by r2848102244:749, 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 | //! Rewriting the container, without touching a byte of what was signed. |
| 2 | //! |
| 3 | //! # The invariant this rests on |
| 4 | //! |
| 5 | //! An Ed25519 signature covers the **plaintext encoding of the record** -- the |
| 6 | //! operation, the identifier that names it and the parents that place it. That |
| 7 | //! is all of what [`Envelope::seal_record`] presents to the signer, and it says |
| 8 | //! nothing whatever about how the record is stored. So the container may be |
| 9 | //! rewritten freely: a segment may be split, joined, renumbered, given a |
| 10 | //! different header, or written under a different layout entirely, and every |
| 11 | //! signature in it still verifies. Nothing here re-signs anything, and nothing |
| 12 | //! here may: an entry goes out exactly as it came in. |
| 13 | //! |
| 14 | //! [`Packed::shape`] is how that is enforced rather than merely intended. It is |
| 15 | //! a digest over every record digest in the log, in order and **across** segment |
| 16 | //! boundaries, and each record digest covers `[kind] || body`, which for a |
| 17 | //! sealed entry is the whole envelope -- signer, signature and payload. Two logs |
| 18 | //! with the same shape hold the same operations, sealed the same way, in the |
| 19 | //! same order. The rewritten log is read back and its shape compared with the |
| 20 | //! original's before anything is swapped, and a difference refuses the repack |
| 21 | //! with the old log untouched. |
| 22 | //! |
| 23 | //! # Why it exists |
| 24 | //! |
| 25 | //! Ore removes nothing, ever, and until now it also rewrote nothing, ever. That |
| 26 | //! second absence is what turned ordinary storage choices into one-way doors: a |
| 27 | //! compression scheme, a key table or a batch signature had to be decided before |
| 28 | //! an import because there would be no way to apply them afterwards. It also |
| 29 | //! left one shape stuck: [`Store::append`] rolls to a new segment on size, but |
| 30 | //! only ever at the tail, so an import -- which writes its whole history in one |
| 31 | //! call -- leaves a single segment of the whole repository, 55 MB on the trial |
| 32 | //! history and 87.8 MB on fe2o3's, that nothing will ever divide again. |
| 33 | //! |
| 34 | //! # What it does not do |
| 35 | //! |
| 36 | //! No compression, no key table, no batch or Merkle signature. Two of those |
| 37 | //! three change what is signed, and all three are separate decisions. This is |
| 38 | //! the tool that makes them applyable later. |
| 39 | //! |
| 40 | //! # Interruption |
| 41 | //! |
| 42 | //! The new log is written aside as `log.new`, read back and checked, and only |
| 43 | //! then swapped in by renaming `log` to `log.old` and `log.new` to `log`. A |
| 44 | //! store may be the only copy of somebody's history, so the rule is that at |
| 45 | //! every instant either `log` is the complete old history, or `log` is the |
| 46 | //! complete new one, or `log` is absent and `log.old` is the complete old one. |
| 47 | //! Only the third needs a person, it is one rename wide, and |
| 48 | //! [`Store::segments`] says so by name when it meets it. |
| 49 | |
| 50 | use crate::keys; |
| 51 | use crate::store::{ |
| 52 | self, |
| 53 | Begin, |
| 54 | Store, |
| 55 | Verify, |
| 56 | Keep, |
| 57 | LOG_DIR, |
| 58 | SEGMENT_LIMIT, |
| 59 | SEGMENT_SUFFIX, |
| 60 | hasher, |
| 61 | }; |
| 62 | use crate::verdict::{ |
| 63 | self, |
| 64 | Fold, |
| 65 | Verdicts, |
| 66 | }; |
| 67 | |
| 68 | use oxedyne_fe2o3_core::prelude::*; |
| 69 | use oxedyne_fe2o3_hash::hash::HashScheme; |
| 70 | use oxedyne_fe2o3_hash::sha256::Sha256; |
| 71 | use oxedyne_fe2o3_ore::id::ReplicaId; |
| 72 | use oxedyne_fe2o3_ore::segment::{ |
| 73 | self, |
| 74 | Entry, |
| 75 | Head, |
| 76 | Integrity, |
| 77 | }; |
| 78 | |
| 79 | use std::fs; |
| 80 | use std::path::{ |
| 81 | Path, |
| 82 | PathBuf, |
| 83 | }; |
| 84 | |
| 85 | |
| 86 | // The two directories a repack works through |
| 87 | pub const NEW_DIR: &str = "log.new"; |
| 88 | pub const OLD_DIR: &str = "log.old"; |
| 89 | |
| 90 | const BATCH: usize = 2_000; // entries verified together, as `Store::replay` does |
| 91 | |
| 92 | |
| 93 | /// A digest over every record digest in a log, in order and across segment |
| 94 | /// boundaries. |
| 95 | pub type Shape = [u8; 32]; |
| 96 | |
| 97 | |
| 98 | /// Whether the rewritten segments carry their records compressed. |
| 99 | /// |
| 100 | /// Both directions are a repack, which is the whole point of it being here: a |
| 101 | /// store packed today is unpacked by running it again the other way, and the |
| 102 | /// operations and the signatures are the same either way. Nothing else in Ore |
| 103 | /// can undo a storage decision. |
| 104 | #[derive(Clone, Copy, Debug, Eq, PartialEq)] |
| 105 | pub enum Packing { |
| 106 | Packed, // each segment's records deflated together as one run |
| 107 | Plain, // each record framed on its own, which is what an append writes |
| 108 | } |
| 109 | |
| 110 | /// How the rewritten container is laid out. |
| 111 | #[derive(Clone, Copy, Debug)] |
| 112 | pub struct Plan { |
| 113 | pub segment_bytes: u64, // PLAIN records a segment is filled to before the next |
| 114 | pub packing: Packing, |
| 115 | } |
| 116 | |
| 117 | impl Default for Plan { |
| 118 | /// [`SEGMENT_LIMIT`] of plain records, and no packing, so that a repacked |
| 119 | /// store is shaped like one that grew. |
| 120 | /// |
| 121 | /// The size is measured in the records' own framing rather than in what lands |
| 122 | /// on the disk, so a segment holds the same operations whether it is packed |
| 123 | /// or not. That is what makes a packed store and a plain one comparable, and |
| 124 | /// it is why a packed segment file is smaller than this rather than equal to |
| 125 | /// it. |
| 126 | fn default() -> Self { |
| 127 | Self { segment_bytes: SEGMENT_LIMIT, packing: Packing::Plain } |
| 128 | } |
| 129 | } |
| 130 | |
| 131 | |
| 132 | /// What a repack did. |
| 133 | #[derive(Clone, Copy, Debug)] |
| 134 | pub struct Packed { |
| 135 | pub operations: usize, |
| 136 | pub segments_before: usize, |
| 137 | pub segments_after: usize, |
| 138 | pub bytes_before: u64, |
| 139 | pub bytes_after: u64, |
| 140 | pub sealed: usize, // entries carrying a signature |
| 141 | pub bare: usize, // entries carrying none |
| 142 | pub veiled: usize, // entries nothing here can read |
| 143 | pub shape: Shape, // the same before and after, or it did not happen |
| 144 | } |
| 145 | |
| 146 | |
| 147 | /// What one pass over a log observed. |
| 148 | struct Seen { |
| 149 | shape: Sha256, // running, over every record digest in order |
| 150 | folds: Vec<(PathBuf, Fold)>, |
| 151 | operations: usize, |
| 152 | sealed: usize, |
| 153 | bare: usize, |
| 154 | veiled: usize, |
| 155 | } |
| 156 | |
| 157 | impl Seen { |
| 158 | fn new() -> Self { |
| 159 | Self { |
| 160 | shape: Sha256::new(), |
| 161 | folds: Vec::new(), |
| 162 | operations: 0, |
| 163 | sealed: 0, |
| 164 | bare: 0, |
| 165 | veiled: 0, |
| 166 | } |
| 167 | } |
| 168 | |
| 169 | fn count(&mut self, entry: &Entry) { |
| 170 | self.operations += 1; |
| 171 | match entry { |
| 172 | Entry::Sealed(_) => self.sealed += 1, |
| 173 | Entry::Bare(_) => self.bare += 1, |
| 174 | Entry::Veiled(_) => self.veiled += 1, |
| 175 | } |
| 176 | } |
| 177 | } |
| 178 | |
| 179 | |
| 180 | /// A segment being filled, and what its header will have to say. |
| 181 | struct Filling { |
| 182 | writer: segment::Writer<HashScheme, 8>, |
| 183 | top: u8, // highest wire code among the records so far |
| 184 | hint: Option<ReplicaId>, // the one replica writing, while there is only one |
| 185 | mixed: bool, // whether a second replica has written into it |
| 186 | // The entries themselves, kept only where they are going to be packed. The |
| 187 | // writer above is then the size oracle and nothing else: it frames each entry |
| 188 | // so that the segment is filled to the same number of records whether it ends |
| 189 | // up packed or plain, and its bytes are thrown away. |
| 190 | held: Vec<Entry>, |
| 191 | } |
| 192 | |
| 193 | |
| 194 | /// Rewrites the store's container, keeping every operation and every signature. |
| 195 | /// |
| 196 | /// `how` is the same choice a replay makes: a working copy checks every |
| 197 | /// signature on the way through and again on the way back, a relay checks |
| 198 | /// neither because it holds no keys and may be carrying veiled entries it could |
| 199 | /// not check if it wanted to. |
| 200 | /// |
| 201 | /// The old log is not touched until the new one has been read back and found to |
| 202 | /// hold the same operations, sealed the same way, in the same order. |
| 203 | pub fn repack(dir: &Path, how: Verify, plan: &Plan) |
| 204 | -> Outcome<Packed> |
| 205 | { |
| 206 | let store = Store::at(dir); |
| 207 | res!(settled(dir)); |
| 208 | // Held across the whole of it, because a repack is a read-modify-write of |
| 209 | // every segment there is. |
| 210 | let _lock = res!(store.lock()); |
| 211 | // Asked again under the lock: taking it is what makes the answer stay true. |
| 212 | res!(settled(dir)); |
| 213 | |
| 214 | let before = res!(store.segments()); |
| 215 | let bytes_before = res!(store.bytes()); |
| 216 | let new_dir = dir.join(NEW_DIR); |
| 217 | // Rubbish from a repack that was killed. It only ever becomes the log by |
| 218 | // being renamed, so a `log.new` on disk is a repack that did not get that |
| 219 | // far, and there is nothing else it could be. |
| 220 | if new_dir.exists() { |
| 221 | res!(fs::remove_dir_all(&new_dir)); |
| 222 | } |
| 223 | res!(fs::create_dir_all(&new_dir)); |
| 224 | |
| 225 | let outcome = write_through(&store, &before, &new_dir, how, plan); |
| 226 | let (source, written) = match outcome { |
| 227 | Ok(pair) => pair, |
| 228 | Err(e) => { |
| 229 | let _ = fs::remove_dir_all(&new_dir); |
| 230 | return Err(e); |
| 231 | }, |
| 232 | }; |
| 233 | |
| 234 | // Read back what was written, checking it as a stranger would. Its shape is |
| 235 | // the claim: the same records, in the same order, whatever segment they |
| 236 | // landed in. |
| 237 | let back = match read_through(&written, how) { |
| 238 | Ok(seen) => seen, |
| 239 | Err(e) => { |
| 240 | let _ = fs::remove_dir_all(&new_dir); |
| 241 | return Err(err!(e, |
| 242 | "The repacked log at {:?} could not be read back, so it was removed \ |
| 243 | and the store was left exactly as it was.", new_dir; |
| 244 | Invalid, Data, File)); |
| 245 | }, |
| 246 | }; |
| 247 | let was = source.shape.finish(); |
| 248 | let now = back.shape.clone().finish(); |
| 249 | if was != now || source.operations != back.operations { |
| 250 | let _ = fs::remove_dir_all(&new_dir); |
| 251 | return Err(err!( |
| 252 | "The repacked log holds {} operations shaped {:02x?} and the log it was \ |
| 253 | made from holds {} shaped {:02x?}. A rewrite of the container must not \ |
| 254 | change either, so the repack was abandoned and the store left exactly as \ |
| 255 | it was.", |
| 256 | back.operations, &now[..8], source.operations, &was[..8]; |
| 257 | Bug, Invalid, Data, Mismatch)); |
| 258 | } |
| 259 | |
| 260 | let bytes_after = res!(sized(&new_dir)); |
| 261 | let segments_after = back.folds.len(); |
| 262 | res!(swap(dir)); |
| 263 | // The verdicts of the read-back, moved onto the names the segments now have. |
| 264 | // They are earned: those bytes were read and every signature in them was |
| 265 | // checked, a moment ago, by the pass above. |
| 266 | res!(settle_verdicts(dir, how, &back)); |
| 267 | |
| 268 | let old_dir = dir.join(OLD_DIR); |
| 269 | match fs::remove_dir_all(&old_dir) { |
| 270 | Ok(()) => (), |
| 271 | Err(e) => return Err(err!(e, |
| 272 | "The repack finished and the log it replaced could not be removed from \ |
| 273 | {:?}. The store is the repacked one and is complete; remove that \ |
| 274 | directory by hand.", old_dir; |
| 275 | IO, File, Write)), |
| 276 | } |
| 277 | |
| 278 | Ok(Packed { |
| 279 | operations: source.operations, |
| 280 | segments_before: before.len(), |
| 281 | segments_after, |
| 282 | bytes_before, |
| 283 | bytes_after, |
| 284 | sealed: source.sealed, |
| 285 | bare: source.bare, |
| 286 | veiled: source.veiled, |
| 287 | shape: was, |
| 288 | }) |
| 289 | } |
| 290 | |
| 291 | |
| 292 | /// Refuses a store part way through a repack, rather than guessing which half of |
| 293 | /// it to believe. |
| 294 | fn settled(dir: &Path) |
| 295 | -> Outcome<()> |
| 296 | { |
| 297 | let log = dir.join(LOG_DIR); |
| 298 | let old = dir.join(OLD_DIR); |
| 299 | if !log.is_dir() { |
| 300 | if old.is_dir() { |
| 301 | return Err(err!( |
| 302 | "There is no log at {:?} and there is one at {:?}. A repack was \ |
| 303 | interrupted between the two renames that swap its result in, which is \ |
| 304 | the one moment at which that can happen. The history is whole and it \ |
| 305 | is in {:?}: move it back with `mv {} {}` and nothing has been lost.", |
| 306 | log, old, old, old.display(), log.display(); |
| 307 | Invalid, Data, File, Missing)); |
| 308 | } |
| 309 | return Err(err!( |
| 310 | "There is no log at {:?}, so there is no store here to repack.", log; |
| 311 | Invalid, Input, Missing)); |
| 312 | } |
| 313 | if old.is_dir() { |
| 314 | return Err(err!( |
| 315 | "There is a log at {:?} and also one at {:?}. A repack swapped its result \ |
| 316 | in and was interrupted before it removed what it replaced. The log in \ |
| 317 | place is the repacked one, and it was read back and checked before the \ |
| 318 | swap; {:?} is the history as it stood before, holding the same operations. \ |
| 319 | Remove it once you are satisfied, and run this again.", |
| 320 | log, old, old; |
| 321 | Invalid, Data, File, Exists)); |
| 322 | } |
| 323 | Ok(()) |
| 324 | } |
| 325 | |
| 326 | |
| 327 | /// Reads every segment of `store`, checking what `how` asks, and writes the |
| 328 | /// entries out again into `into` under `plan`. |
| 329 | fn write_through( |
| 330 | store: &Store, |
| 331 | from: &[(u64, PathBuf)], |
| 332 | into: &Path, |
| 333 | how: Verify, |
| 334 | plan: &Plan, |
| 335 | ) |
| 336 | -> Outcome<(Seen, Vec<PathBuf>)> |
| 337 | { |
| 338 | let _ = store; |
| 339 | let mut seen = Seen::new(); |
| 340 | let mut filling: Option<Filling> = None; |
| 341 | let mut written: Vec<PathBuf> = Vec::new(); |
| 342 | for (_, path) in from { |
| 343 | let whence = fmt!("the segment {}", path.display()); |
| 344 | res!(Store::stream_batches(path, Begin::Start, Integrity::Checked, None, BATCH, |
| 345 | |entries, digests| { |
| 346 | seen.shape.update(digests); |
| 347 | // Checked here so that a repack never rewrites a history it could not |
| 348 | // stand behind. The entries themselves are passed on untouched; what |
| 349 | // this establishes is that they were sound when they were passed. |
| 350 | let told = match how { |
| 351 | Verify::Signatures(trust) => Some(res!( |
| 352 | keys::check_all(entries, trust, &whence, Keep::Nothing))), |
| 353 | Verify::Nothing => None, |
| 354 | }; |
| 355 | for (i, entry) in entries.iter().enumerate() { |
| 356 | seen.count(entry); |
| 357 | let (code, replica) = match &told { |
| 358 | // Already decoded by the check, and in the same order: a veiled |
| 359 | // entry never reaches here, because checking one is refused. |
| 360 | Some(checked) => { |
| 361 | let one = res!(checked.get(i).ok_or_else(|| err!( |
| 362 | "Entry {} of a batch from {} has no verdict beside it.", |
| 363 | i, whence; Bug, Missing))); |
| 364 | (Some(one.rec.op.code()), one.rec.id().replica) |
| 365 | }, |
| 366 | None => res!(about(entry)), |
| 367 | }; |
| 368 | res!(place(&mut filling, entry, code, replica, into, plan, &mut written)); |
| 369 | } |
| 370 | Ok(()) |
| 371 | })); |
| 372 | } |
| 373 | if let Some(full) = filling.take() { |
| 374 | res!(close(full, into, &mut written, plan)); |
| 375 | } |
| 376 | Ok((seen, written)) |
| 377 | } |
| 378 | |
| 379 | |
| 380 | /// What a segment header has to be able to say about an entry. |
| 381 | /// |
| 382 | /// A veiled entry's operation is ciphertext, so nothing here can name its wire |
| 383 | /// code. It is reported as unknown, and the caller declares the highest version |
| 384 | /// there is rather than guessing at a lower one that might not carry it. |
| 385 | fn about(entry: &Entry) |
| 386 | -> Outcome<(Option<u8>, ReplicaId)> |
| 387 | { |
| 388 | match entry { |
| 389 | Entry::Veiled(v) => Ok((None, v.head.id().replica)), |
| 390 | other => { |
| 391 | let rec = res!(other.peek()); |
| 392 | Ok((Some(rec.op.code()), rec.id().replica)) |
| 393 | }, |
| 394 | } |
| 395 | } |
| 396 | |
| 397 | |
| 398 | /// Puts one entry into the segment being filled, starting one where there is |
| 399 | /// none and closing one that has reached its size. |
| 400 | fn place( |
| 401 | filling: &mut Option<Filling>, |
| 402 | entry: &Entry, |
| 403 | code: Option<u8>, |
| 404 | replica: ReplicaId, |
| 405 | into: &Path, |
| 406 | plan: &Plan, |
| 407 | written: &mut Vec<PathBuf>, |
| 408 | ) |
| 409 | -> Outcome<()> |
| 410 | { |
| 411 | if filling.is_none() { |
| 412 | *filling = Some(Filling { |
| 413 | // Resumed over a bare header at the current version, so that nothing is |
| 414 | // refused for its wire code while the segment is being filled. The |
| 415 | // header that is finally written is decided from the records that went |
| 416 | // in, and a resumed writer hands back the records alone, so the two are |
| 417 | // independent -- which is the whole reason this is done this way round. |
| 418 | writer: res!(segment::Writer::resume( |
| 419 | &Head::new(None).encode(), hasher(), store::SALT)), |
| 420 | top: 0, |
| 421 | hint: Some(replica), |
| 422 | mixed: false, |
| 423 | held: Vec::new(), |
| 424 | }); |
| 425 | } |
| 426 | let full = res!(filling.as_mut().ok_or_else(|| err!( |
| 427 | "A segment was to be filled and there was none."; Bug, Missing))); |
| 428 | res!(full.writer.push(entry)); |
| 429 | if plan.packing == Packing::Packed { |
| 430 | full.held.push(entry.clone()); |
| 431 | } |
| 432 | match code { |
| 433 | Some(c) if c > full.top => full.top = c, |
| 434 | Some(_) => (), |
| 435 | None => full.top = segment::highest_code(segment::VERSION), |
| 436 | } |
| 437 | if full.hint != Some(replica) { |
| 438 | full.mixed = true; |
| 439 | } |
| 440 | if full.writer.bytes().len() as u64 >= plan.segment_bytes { |
| 441 | let done = res!(filling.take().ok_or_else(|| err!( |
| 442 | "A full segment went missing between being filled and being closed."; |
| 443 | Bug, Missing))); |
| 444 | res!(close(done, into, written, plan)); |
| 445 | } |
| 446 | Ok(()) |
| 447 | } |
| 448 | |
| 449 | |
| 450 | /// Writes a filled segment out under the header its records call for. |
| 451 | /// |
| 452 | /// The version is the lowest that spells every wire code inside, so a repack |
| 453 | /// does not bump a segment of old operations to a version an older reader would |
| 454 | /// refuse. The hint is the one replica that wrote it, and nothing where more |
| 455 | /// than one did, which is what a hint is for. |
| 456 | fn close(full: Filling, into: &Path, written: &mut Vec<PathBuf>, plan: &Plan) |
| 457 | -> Outcome<()> |
| 458 | { |
| 459 | let head = Head { |
| 460 | version: version_for(full.top), |
| 461 | replica: match full.mixed { |
| 462 | true => None, |
| 463 | false => full.hint, |
| 464 | }, |
| 465 | }; |
| 466 | let mut bytes = head.encode(); |
| 467 | match plan.packing { |
| 468 | // The records go under the compressor together, at the version this |
| 469 | // segment turned out to need: a run is refused a record the version has no |
| 470 | // code for exactly as a plain push is. |
| 471 | Packing::Packed => { |
| 472 | let mut writer = segment::Writer::new(&head, hasher(), store::SALT); |
| 473 | res!(writer.push_packed(&full.held)); |
| 474 | let packed = writer.finish(); |
| 475 | // `Writer::new` emitted the header, so what it hands back is the whole |
| 476 | // segment rather than the records alone. |
| 477 | bytes = packed; |
| 478 | }, |
| 479 | Packing::Plain => bytes.extend_from_slice(&full.writer.finish()), |
| 480 | } |
| 481 | let path = into.join(fmt!("{:06}{}", written.len(), SEGMENT_SUFFIX)); |
| 482 | match fs::write(&path, &bytes) { |
| 483 | Ok(()) => (), |
| 484 | Err(e) => return Err(err!(e, |
| 485 | "The repacked segment {:?} of {} bytes could not be written.", |
| 486 | path, bytes.len(); |
| 487 | IO, File, Write)), |
| 488 | } |
| 489 | written.push(path); |
| 490 | Ok(()) |
| 491 | } |
| 492 | |
| 493 | |
| 494 | /// The lowest format version whose vocabulary spells `top`. |
| 495 | fn version_for(top: u8) -> u8 { |
| 496 | let mut version = segment::VERSION_MIN; |
| 497 | while version < segment::VERSION && segment::highest_code(version) < top { |
| 498 | version += 1; |
| 499 | } |
| 500 | version |
| 501 | } |
| 502 | |
| 503 | |
| 504 | /// Reads a written log back, checking what `how` asks and folding what it finds. |
| 505 | fn read_through(paths: &[PathBuf], how: Verify) |
| 506 | -> Outcome<Seen> |
| 507 | { |
| 508 | let mut seen = Seen::new(); |
| 509 | for path in paths { |
| 510 | let whence = fmt!("the segment {}", path.display()); |
| 511 | let mut sum = Sha256::new(); |
| 512 | let head = res!(Store::stream_batches(path, Begin::Start, Integrity::Checked, None, BATCH, |
| 513 | |entries, digests| { |
| 514 | seen.shape.update(digests); |
| 515 | sum.update(digests); |
| 516 | if let Verify::Signatures(trust) = how { |
| 517 | res!(keys::check_all(entries, trust, &whence, Keep::Nothing)); |
| 518 | } |
| 519 | for entry in entries.iter() { |
| 520 | seen.count(entry); |
| 521 | } |
| 522 | Ok(()) |
| 523 | })); |
| 524 | seen.folds.push((path.clone(), Store::fold_from(sum, &head))); |
| 525 | } |
| 526 | Ok(seen) |
| 527 | } |
| 528 | |
| 529 | |
| 530 | /// The total size of the segments in a directory. |
| 531 | fn sized(dir: &Path) |
| 532 | -> Outcome<u64> |
| 533 | { |
| 534 | let mut total = 0u64; |
| 535 | for entry in res!(fs::read_dir(dir)) { |
| 536 | let entry = res!(entry); |
| 537 | total += res!(entry.metadata()).len(); |
| 538 | } |
| 539 | Ok(total) |
| 540 | } |
| 541 | |
| 542 | |
| 543 | /// Puts the repacked log in place of the one it was made from. |
| 544 | /// |
| 545 | /// Two renames, and between them is the only instant at which a store does not |
| 546 | /// have a log. Both directories are complete on disk throughout, so an |
| 547 | /// interruption there loses nothing; [`settled`] is what turns the state it |
| 548 | /// leaves into a sentence rather than a puzzle. |
| 549 | fn swap(dir: &Path) |
| 550 | -> Outcome<()> |
| 551 | { |
| 552 | let log = dir.join(LOG_DIR); |
| 553 | let old = dir.join(OLD_DIR); |
| 554 | let new = dir.join(NEW_DIR); |
| 555 | match fs::rename(&log, &old) { |
| 556 | Ok(()) => (), |
| 557 | Err(e) => return Err(err!(e, |
| 558 | "The log {:?} could not be moved aside to {:?}, so nothing was changed.", |
| 559 | log, old; |
| 560 | IO, File, Write)), |
| 561 | } |
| 562 | match fs::rename(&new, &log) { |
| 563 | Ok(()) => Ok(()), |
| 564 | Err(e) => { |
| 565 | // Put it back rather than leave a store with no log at all. |
| 566 | match fs::rename(&old, &log) { |
| 567 | Ok(()) => Err(err!(e, |
| 568 | "The repacked log {:?} could not be moved into place at {:?}. The \ |
| 569 | log it was to replace has been put back and the store is as it was.", |
| 570 | new, log; |
| 571 | IO, File, Write)), |
| 572 | Err(back) => Err(err!(back, |
| 573 | "The repacked log {:?} could not be moved into place at {:?}, and \ |
| 574 | the log it was to replace could not be moved back from {:?}. The \ |
| 575 | history is whole and it is in {:?}: move it back by hand.", |
| 576 | new, log, old, old; |
| 577 | IO, File, Write)), |
| 578 | } |
| 579 | }, |
| 580 | } |
| 581 | } |
| 582 | |
| 583 | |
| 584 | /// Records the verdicts of the read-back against the names the segments now |
| 585 | /// have, or removes the file where nothing was checked. |
| 586 | fn settle_verdicts(dir: &Path, how: Verify, back: &Seen) |
| 587 | -> Outcome<()> |
| 588 | { |
| 589 | let path = Verdicts::path_of(dir); |
| 590 | if let Verify::Nothing = how { |
| 591 | // Nothing was checked, so there is nothing to say. What was there was |
| 592 | // about segments that no longer exist. |
| 593 | let _ = fs::remove_file(&path); |
| 594 | return Ok(()); |
| 595 | } |
| 596 | let log = dir.join(LOG_DIR); |
| 597 | let mut said = Verdicts::new(); |
| 598 | // One reading of the clock for every note the repack leaves, for the reason |
| 599 | // [`Store::replay_since`] takes one. |
| 600 | let at = verdict::now(); |
| 601 | let last = back.folds.len().saturating_sub(1); |
| 602 | for (i, (was, fold)) in back.folds.iter().enumerate() { |
| 603 | let name = match was.file_name() { |
| 604 | Some(n) => n.to_string_lossy().into_owned(), |
| 605 | None => continue, |
| 606 | }; |
| 607 | let now = log.join(&name); |
| 608 | let meta = match fs::metadata(&now) { |
| 609 | Ok(m) => m, |
| 610 | Err(_) => continue, |
| 611 | }; |
| 612 | // The same rule the replay uses: the tail is the one the next append |
| 613 | // continues, and a fold over it names a state it may already have left. |
| 614 | if i == last && meta.len() < SEGMENT_LIMIT { |
| 615 | continue; |
| 616 | } |
| 617 | let modified = match meta.modified() { |
| 618 | Ok(t) => match t.duration_since(std::time::UNIX_EPOCH) { |
| 619 | Ok(d) => d.as_nanos(), |
| 620 | Err(_) => continue, |
| 621 | }, |
| 622 | Err(_) => continue, |
| 623 | }; |
| 624 | said.record(&name, meta.len(), modified, *fold, at); |
| 625 | } |
| 626 | // Derived state, so a store that cannot be written to is one that verifies |
| 627 | // next time. Not the repack's failure. |
| 628 | let _ = said.write(dir); |
| 629 | Ok(()) |
| 630 | } |
| 631 | |
| 632 | |
| 633 | /// The shape of a log as it stands, for a caller that wants to compare two |
| 634 | /// stores without repacking either. |
| 635 | pub fn shape_of(store: &Store, how: Verify) |
| 636 | -> Outcome<Shape> |
| 637 | { |
| 638 | let paths: Vec<PathBuf> = res!(store.segments()).into_iter().map(|(_, p)| p).collect(); |
| 639 | let seen = res!(read_through(&paths, how)); |
| 640 | Ok(seen.shape.finish()) |
| 641 | } |
| 642 | |
| 643 | |
| 644 | #[cfg(test)] |
| 645 | mod tests { |
| 646 | use super::*; |
| 647 | |
| 648 | use crate::keys::{ |
| 649 | Signing, |
| 650 | Trust, |
| 651 | }; |
| 652 | use crate::store::{ |
| 653 | Keep, |
| 654 | SALT, |
| 655 | }; |
| 656 | |
| 657 | use oxedyne_fe2o3_ore::id::OpId; |
| 658 | use oxedyne_fe2o3_ore::op::{ |
| 659 | Header, |
| 660 | Op, |
| 661 | Record, |
| 662 | }; |
| 663 | |
| 664 | use std::time::{ |
| 665 | SystemTime, |
| 666 | UNIX_EPOCH, |
| 667 | }; |
| 668 | |
| 669 | const MARKS: u64 = 3_000; // padded marks, enough for several segments |
| 670 | |
| 671 | /// A directory that removes itself however the test ends. |
| 672 | struct Scratch { |
| 673 | path: PathBuf, |
| 674 | } |
| 675 | |
| 676 | impl Scratch { |
| 677 | fn new(what: &str) |
| 678 | -> Outcome<Self> |
| 679 | { |
| 680 | let stamp = res!(SystemTime::now().duration_since(UNIX_EPOCH)); |
| 681 | let path = std::env::temp_dir().join(fmt!( |
| 682 | "ore_repack_{}_{}_{}", what, std::process::id(), stamp.as_nanos(), |
| 683 | )); |
| 684 | res!(fs::create_dir_all(&path)); |
| 685 | Ok(Self { path }) |
| 686 | } |
| 687 | } |
| 688 | |
| 689 | impl Drop for Scratch { |
| 690 | fn drop(&mut self) { |
| 691 | let _ = fs::remove_dir_all(&self.path); |
| 692 | } |
| 693 | } |
| 694 | |
| 695 | /// A store holding its whole history in one segment, which is the shape an |
| 696 | /// import leaves and the shape nothing but a repack can change. |
| 697 | fn one_big_segment(dir: &Path, key: &Signing) |
| 698 | -> Outcome<Store> |
| 699 | { |
| 700 | let replica = ReplicaId::new(1); |
| 701 | let store = res!(Store::create(dir, Some(replica))); |
| 702 | let pad = "x".repeat(500); |
| 703 | let mut entries = Vec::new(); |
| 704 | let mut parents = Vec::new(); |
| 705 | for i in 1..=MARKS { |
| 706 | let rec = Record::new( |
| 707 | res!(Header::new(OpId::new(replica, i), parents.clone())), |
| 708 | Op::Mark { name: fmt!("mark {} {}", i, pad), body: None, time: None }, |
| 709 | ); |
| 710 | parents = vec![OpId::new(replica, i)]; |
| 711 | entries.push(Entry::Sealed(res!(key.seal(&rec)))); |
| 712 | } |
| 713 | res!(store.append(&entries, Some(replica))); |
| 714 | assert_eq!(res!(store.segments()).len(), 1, |
| 715 | "an append writes what it is given into one segment, however much it is"); |
| 716 | Ok(store) |
| 717 | } |
| 718 | |
| 719 | /// Every operation, its parents and its seal survive, and the container does |
| 720 | /// not. |
| 721 | #[test] |
| 722 | fn the_operations_survive_and_the_container_does_not() -> Outcome<()> { |
| 723 | let scratch = res!(Scratch::new("through")); |
| 724 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 725 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 726 | let trust = Trust::of(&[key.binding()]); |
| 727 | |
| 728 | let was = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes)); |
| 729 | let bytes_before = res!(store.bytes()); |
| 730 | |
| 731 | let done = res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 732 | segment_bytes: 200_000, |
| 733 | packing: Packing::Plain, |
| 734 | })); |
| 735 | assert_eq!(done.operations, MARKS as usize); |
| 736 | assert_eq!(done.segments_before, 1); |
| 737 | assert!(done.segments_after > 4, |
| 738 | "one segment became several, which is the thing nothing else can do: {}", |
| 739 | done.segments_after); |
| 740 | assert_eq!(res!(store.segments()).len(), done.segments_after); |
| 741 | assert!(bytes_before > 0); |
| 742 | |
| 743 | let now = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes)); |
| 744 | assert_eq!(now.log.len(), was.log.len()); |
| 745 | let before: Vec<(OpId, Vec<OpId>)> = was.log.iter() |
| 746 | .map(|r| (r.id(), r.parents().to_vec())).collect(); |
| 747 | let after: Vec<(OpId, Vec<OpId>)> = now.log.iter() |
| 748 | .map(|r| (r.id(), r.parents().to_vec())).collect(); |
| 749 | assert_eq!(before, after, "every identifier and every parent, in order"); |
| 750 | assert_eq!(was.envelopes, now.envelopes, "and every signature, byte for byte"); |
| 751 | assert_eq!(was.prov, now.prov, "attributed exactly as before"); |
| 752 | Ok(()) |
| 753 | } |
| 754 | |
| 755 | /// However a repack ends, the store holds the operations it held. |
| 756 | /// |
| 757 | /// This is the invariant, and it is asserted over the outcome rather than |
| 758 | /// beside it: a repack that succeeds must have kept them, and a repack that |
| 759 | /// refuses must have left them. Breaking the writer so that it drops an |
| 760 | /// operation turns this into a refusal and the assertion still holds; |
| 761 | /// breaking the shape comparison as well lets the loss through, and it does |
| 762 | /// not. That pair is what the shape comparison is worth. |
| 763 | #[test] |
| 764 | fn a_repack_never_leaves_a_store_holding_different_operations() -> Outcome<()> { |
| 765 | let scratch = res!(Scratch::new("invariant")); |
| 766 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 767 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 768 | let trust = Trust::of(&[key.binding()]); |
| 769 | let held = |store: &Store| -> Outcome<Vec<(OpId, Vec<OpId>, Option<Vec<u8>>)>> { |
| 770 | let done = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes)); |
| 771 | Ok(done.log.iter().map(|r| ( |
| 772 | r.id(), |
| 773 | r.parents().to_vec(), |
| 774 | done.envelopes.get(&r.id()).map(|e| e.signature().to_vec()), |
| 775 | )).collect()) |
| 776 | }; |
| 777 | let before = res!(held(&store)); |
| 778 | assert_eq!(before.len(), MARKS as usize); |
| 779 | |
| 780 | let outcome = repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 781 | segment_bytes: 200_000, |
| 782 | packing: Packing::Plain, |
| 783 | }); |
| 784 | let after = res!(held(&store)); |
| 785 | assert_eq!(before, after, |
| 786 | "every operation, every parent and every signature, whether the repack \ |
| 787 | succeeded or refused"); |
| 788 | if outcome.is_err() { |
| 789 | assert!(!scratch.path.join(NEW_DIR).exists(), "and nothing left lying about"); |
| 790 | assert!(!scratch.path.join(OLD_DIR).exists()); |
| 791 | } |
| 792 | Ok(()) |
| 793 | } |
| 794 | |
| 795 | /// A log missing one operation is a different shape, which is the whole of |
| 796 | /// what the comparison in [`repack`] is asking. |
| 797 | #[test] |
| 798 | fn a_shape_is_the_operations_in_order() -> Outcome<()> { |
| 799 | let scratch = res!(Scratch::new("shape")); |
| 800 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 801 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 802 | let trust = Trust::of(&[key.binding()]); |
| 803 | let paths: Vec<PathBuf> = res!(store.segments()).into_iter() |
| 804 | .map(|(_, p)| p).collect(); |
| 805 | let whole = res!(read_through(&paths, Verify::Signatures(&trust))); |
| 806 | |
| 807 | let mut short = Seen::new(); |
| 808 | let mut skipped = false; |
| 809 | res!(Store::stream_batches(&paths[0], Begin::Start, Integrity::Checked, None, BATCH, |
| 810 | |entries, digests| { |
| 811 | match skipped { |
| 812 | false => { |
| 813 | skipped = true; |
| 814 | short.shape.update(&digests[32..]); // every digest but the first |
| 815 | for entry in entries.iter().skip(1) { |
| 816 | short.count(entry); |
| 817 | } |
| 818 | }, |
| 819 | true => { |
| 820 | short.shape.update(digests); |
| 821 | for entry in entries.iter() { |
| 822 | short.count(entry); |
| 823 | } |
| 824 | }, |
| 825 | } |
| 826 | Ok(()) |
| 827 | })); |
| 828 | assert_ne!(whole.shape.clone().finish(), short.shape.finish()); |
| 829 | assert_eq!(whole.operations, short.operations + 1); |
| 830 | |
| 831 | // And the same operations read through a different segmentation are the |
| 832 | // same shape, which is the other half of it. |
| 833 | res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 834 | segment_bytes: 200_000, |
| 835 | packing: Packing::Plain, |
| 836 | })); |
| 837 | let cut: Vec<PathBuf> = res!(store.segments()).into_iter() |
| 838 | .map(|(_, p)| p).collect(); |
| 839 | assert!(cut.len() > 4); |
| 840 | let again = res!(read_through(&cut, Verify::Signatures(&trust))); |
| 841 | assert_eq!(whole.shape.finish(), again.shape.finish(), |
| 842 | "a shape says nothing about which segment a record landed in"); |
| 843 | Ok(()) |
| 844 | } |
| 845 | |
| 846 | /// A repack that checks nothing says nothing about what it checked. |
| 847 | #[test] |
| 848 | fn an_unverified_repack_leaves_no_verdict() -> Outcome<()> { |
| 849 | let scratch = res!(Scratch::new("unverified")); |
| 850 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 851 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 852 | let trust = Trust::of(&[key.binding()]); |
| 853 | res!(store.replay(Verify::Signatures(&trust), Keep::Nothing)); |
| 854 | assert!(Verdicts::path_of(&scratch.path).is_file(), "there is one to lose"); |
| 855 | |
| 856 | res!(repack(&scratch.path, Verify::Nothing, &Plan { segment_bytes: 200_000, packing: Packing::Plain })); |
| 857 | assert!(!Verdicts::path_of(&scratch.path).is_file(), |
| 858 | "and it went with the segments it was about"); |
| 859 | Ok(()) |
| 860 | } |
| 861 | |
| 862 | /// A store part way through a repack is refused, by name, with the move that |
| 863 | /// mends it. |
| 864 | #[test] |
| 865 | fn an_interrupted_swap_is_named_and_not_guessed_at() -> Outcome<()> { |
| 866 | let scratch = res!(Scratch::new("interrupted")); |
| 867 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 868 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 869 | let trust = Trust::of(&[key.binding()]); |
| 870 | res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 871 | segment_bytes: 200_000, |
| 872 | packing: Packing::Plain, |
| 873 | })); |
| 874 | |
| 875 | // The state a kill between the two renames leaves: no log, and the whole |
| 876 | // history in `log.old`. |
| 877 | let log = scratch.path.join(LOG_DIR); |
| 878 | let old = scratch.path.join(OLD_DIR); |
| 879 | res!(fs::rename(&log, &old)); |
| 880 | |
| 881 | let said = match store.replay(Verify::Signatures(&trust), Keep::Nothing) { |
| 882 | Ok(_) => return Err(err!( |
| 883 | "A store with no log replayed."; Test, Invalid)), |
| 884 | Err(e) => fmt!("{}", e.plain()), |
| 885 | }; |
| 886 | assert!(said.contains("interrupted between the two renames"), |
| 887 | "and says what happened: {}", said); |
| 888 | assert!(said.contains("Move it back"), "and what to do: {}", said); |
| 889 | |
| 890 | let refused = match repack(&scratch.path, Verify::Signatures(&trust), &Plan::default()) { |
| 891 | Ok(_) => return Err(err!( |
| 892 | "A repack ran over a store part way through one."; Test, Invalid)), |
| 893 | Err(e) => fmt!("{}", e.plain()), |
| 894 | }; |
| 895 | assert!(refused.contains("interrupted"), "and repack says so too: {}", refused); |
| 896 | |
| 897 | // The remedy works, and the history is all there. |
| 898 | res!(fs::rename(&old, &log)); |
| 899 | let done = res!(store.replay(Verify::Signatures(&trust), Keep::Nothing)); |
| 900 | assert_eq!(done.log.len(), MARKS as usize, "nothing was lost"); |
| 901 | Ok(()) |
| 902 | } |
| 903 | |
| 904 | /// Rubbish a killed repack left behind is cleared, and the history it was |
| 905 | /// made from is not. |
| 906 | #[test] |
| 907 | fn a_half_written_repack_is_rubbish_and_is_treated_as_such() -> Outcome<()> { |
| 908 | let scratch = res!(Scratch::new("halfway")); |
| 909 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 910 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 911 | let trust = Trust::of(&[key.binding()]); |
| 912 | let before = res!(Store::fold_of(&store.segment_path(0))); |
| 913 | |
| 914 | // A `log.new` holding a truncated segment: what a kill part way through |
| 915 | // the writing leaves. It only ever becomes the log by being renamed, so |
| 916 | // there is nothing else it could be. |
| 917 | let new = scratch.path.join(NEW_DIR); |
| 918 | res!(fs::create_dir_all(&new)); |
| 919 | res!(fs::write(new.join("000000.seg"), b"ORESEG rubbish")); |
| 920 | // Numbered past anything this repack will write, so that clearing the |
| 921 | // directory is the only thing that can remove it. |
| 922 | res!(fs::write(new.join("000099.seg"), b"ORESEG more rubbish")); |
| 923 | |
| 924 | let done = res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 925 | segment_bytes: 200_000, |
| 926 | packing: Packing::Plain, |
| 927 | })); |
| 928 | assert_eq!(done.operations, MARKS as usize); |
| 929 | assert_ne!(before, res!(Store::fold_of(&store.segment_path(0))), |
| 930 | "the container was rewritten"); |
| 931 | let now = res!(store.replay(Verify::Signatures(&trust), Keep::Nothing)); |
| 932 | assert_eq!(now.log.len(), MARKS as usize); |
| 933 | assert_eq!(res!(store.segments()).len(), done.segments_after, |
| 934 | "and nothing of the killed repack came through with it"); |
| 935 | Ok(()) |
| 936 | } |
| 937 | |
| 938 | /// A repacked segment's verdict is about the bytes that were read back, so a |
| 939 | /// tamper afterwards is still caught. |
| 940 | #[test] |
| 941 | fn the_verdicts_a_repack_leaves_are_about_what_it_read() -> Outcome<()> { |
| 942 | let scratch = res!(Scratch::new("verdicts")); |
| 943 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 944 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 945 | let trust = Trust::of(&[key.binding()]); |
| 946 | res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 947 | segment_bytes: 200_000, |
| 948 | packing: Packing::Plain, |
| 949 | })); |
| 950 | |
| 951 | let said = Verdicts::read(&scratch.path); |
| 952 | assert!(!said.is_empty(), "the read-back earned them"); |
| 953 | // Every fold in the file is the fold of the segment it names, taken now, |
| 954 | // from the disk, rather than from what the repack said about itself. |
| 955 | for (n, path) in res!(store.segments()) { |
| 956 | let meta = res!(fs::metadata(&path)); |
| 957 | if n == res!(store.segments()).len() as u64 - 1 && meta.len() < SEGMENT_LIMIT { |
| 958 | continue; |
| 959 | } |
| 960 | let name = fmt!("{:06}{}", n, SEGMENT_SUFFIX); |
| 961 | let modified = res!(res!(meta.modified()).duration_since( |
| 962 | std::time::UNIX_EPOCH)).as_nanos(); |
| 963 | let expected = res!(said.expected(&name, meta.len(), modified, verdict::now(), 0) |
| 964 | .ok_or_else(|| err!( |
| 965 | "No verdict was left for {}.", name; Test, Missing))); |
| 966 | assert_eq!(expected, res!(Store::fold_of(&path)), |
| 967 | "the verdict for {} names the bytes on the disk", name); |
| 968 | } |
| 969 | Ok(()) |
| 970 | } |
| 971 | |
| 972 | /// **Packing is applied and revoked, and the history is the same all three |
| 973 | /// times.** |
| 974 | /// |
| 975 | /// This is what "revisable" has to mean to be worth anything: a store packed |
| 976 | /// today can be unpacked tomorrow, and neither direction is a migration. |
| 977 | /// Everything is compared against the state before either -- the shape, the |
| 978 | /// operations, their parents and their signatures -- so a round trip that |
| 979 | /// merely agreed with itself would not pass. |
| 980 | #[test] |
| 981 | fn packing_is_applied_and_revoked_and_the_history_does_not_move() -> Outcome<()> { |
| 982 | let scratch = res!(Scratch::new("round_trip")); |
| 983 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 984 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 985 | let trust = Trust::of(&[key.binding()]); |
| 986 | |
| 987 | let held = |store: &Store| -> Outcome<Vec<(OpId, Vec<OpId>, Option<Vec<u8>>)>> { |
| 988 | let done = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes)); |
| 989 | Ok(done.log.iter().map(|r| ( |
| 990 | r.id(), |
| 991 | r.parents().to_vec(), |
| 992 | done.envelopes.get(&r.id()).map(|e| e.signature().to_vec()), |
| 993 | )).collect()) |
| 994 | }; |
| 995 | // Repacked plain FIRST, so that the baseline differs from what follows in |
| 996 | // the packing alone. Against the store as it arrived the comparison would |
| 997 | // also carry the re-segmentation -- eleven segment headers where there had |
| 998 | // been one -- and would fail for a reason that has nothing to do with |
| 999 | // compression, which is how this test found its own footing. |
| 1000 | res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 1001 | segment_bytes: 200_000, |
| 1002 | packing: Packing::Plain, |
| 1003 | })); |
| 1004 | let was = res!(held(&store)); |
| 1005 | let shape = res!(shape_of(&store, Verify::Signatures(&trust))); |
| 1006 | let plain_bytes = res!(store.bytes()); |
| 1007 | let plain_segments = res!(store.segments()).len(); |
| 1008 | assert_eq!(was.len(), MARKS as usize); |
| 1009 | assert!(plain_segments > 4, "the baseline is several segments: {}", plain_segments); |
| 1010 | |
| 1011 | // PACK. |
| 1012 | let packed = res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 1013 | segment_bytes: 200_000, |
| 1014 | packing: Packing::Packed, |
| 1015 | })); |
| 1016 | assert_eq!(packed.shape, shape, "the shape is what it was"); |
| 1017 | assert_eq!(packed.operations, MARKS as usize); |
| 1018 | let packed_bytes = res!(store.bytes()); |
| 1019 | assert!(packed_bytes < plain_bytes, |
| 1020 | "packing made it smaller: {} against {}", packed_bytes, plain_bytes); |
| 1021 | assert_eq!(res!(held(&store)), was, "and every operation, parent and signature"); |
| 1022 | assert_eq!(res!(shape_of(&store, Verify::Signatures(&trust))), shape); |
| 1023 | // The bytes on disk really are packed, and not merely smaller. |
| 1024 | let mut found = false; |
| 1025 | for (_, path) in res!(store.segments()) { |
| 1026 | let bytes = res!(fs::read(&path)); |
| 1027 | let (head, used) = res!(res!(segment::Head::decode(&bytes)).ok_or_else(|| err!( |
| 1028 | "no header"; Test, Missing))); |
| 1029 | assert_eq!(head.version, segment::VERSION_MIN, "the version did not move"); |
| 1030 | if bytes[used] == segment::KIND_PACKED { |
| 1031 | found = true; |
| 1032 | } |
| 1033 | } |
| 1034 | assert!(found, "at least one segment says it carries a run"); |
| 1035 | |
| 1036 | // REVOKE. |
| 1037 | let plain = res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 1038 | segment_bytes: 200_000, |
| 1039 | packing: Packing::Plain, |
| 1040 | })); |
| 1041 | assert_eq!(plain.shape, shape, "and unpacking did not move it either"); |
| 1042 | assert_eq!(res!(held(&store)), was, "nor any operation, parent or signature"); |
| 1043 | assert_eq!(res!(store.bytes()), plain_bytes, |
| 1044 | "and the store is the size the plain one was, byte for byte"); |
| 1045 | assert_eq!(res!(store.segments()).len(), plain_segments, |
| 1046 | "in the same number of segments"); |
| 1047 | for (_, path) in res!(store.segments()) { |
| 1048 | let bytes = res!(fs::read(&path)); |
| 1049 | let (_, used) = res!(res!(segment::Head::decode(&bytes)).ok_or_else(|| err!( |
| 1050 | "no header"; Test, Missing))); |
| 1051 | assert_ne!(bytes[used], segment::KIND_PACKED, "and nothing is packed any more"); |
| 1052 | } |
| 1053 | Ok(()) |
| 1054 | } |
| 1055 | |
| 1056 | /// A packed store and a plain one holding one history have one shape. |
| 1057 | /// |
| 1058 | /// Said on its own because it is the property the round trip rests on rather |
| 1059 | /// than a consequence of it: two stores, never repacked into each other, and |
| 1060 | /// the same answer. |
| 1061 | #[test] |
| 1062 | fn a_packed_store_and_a_plain_one_have_the_same_shape() -> Outcome<()> { |
| 1063 | let first = res!(Scratch::new("shape_packed")); |
| 1064 | let second = res!(Scratch::new("shape_plain")); |
| 1065 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 1066 | let one = res!(one_big_segment(&first.path, &key)); |
| 1067 | let two = res!(one_big_segment(&second.path, &key)); |
| 1068 | let trust = Trust::of(&[key.binding()]); |
| 1069 | assert_eq!( |
| 1070 | res!(shape_of(&one, Verify::Signatures(&trust))), |
| 1071 | res!(shape_of(&two, Verify::Signatures(&trust))), |
| 1072 | "the two fixtures start alike, or nothing below means anything"); |
| 1073 | |
| 1074 | res!(repack(&first.path, Verify::Signatures(&trust), &Plan { |
| 1075 | segment_bytes: 200_000, |
| 1076 | packing: Packing::Packed, |
| 1077 | })); |
| 1078 | res!(repack(&second.path, Verify::Signatures(&trust), &Plan { |
| 1079 | segment_bytes: 200_000, |
| 1080 | packing: Packing::Plain, |
| 1081 | })); |
| 1082 | assert_ne!(res!(one.bytes()), res!(two.bytes()), "and they are not the same bytes"); |
| 1083 | assert_eq!( |
| 1084 | res!(shape_of(&one, Verify::Signatures(&trust))), |
| 1085 | res!(shape_of(&two, Verify::Signatures(&trust))), |
| 1086 | "a shape says nothing about whether the records were compressed"); |
| 1087 | Ok(()) |
| 1088 | } |
| 1089 | |
| 1090 | /// Packing changes the disk and does not change the wire. |
| 1091 | /// |
| 1092 | /// A sync carries [`Entry`] values inside an `ORESYN` message and never a |
| 1093 | /// byte of segment framing, so what a peer receives is the same whether the |
| 1094 | /// store it came from was packed or not. The disk saving is a disk saving and |
| 1095 | /// must not be reported as a network one, which is why this is asserted |
| 1096 | /// rather than argued: the entries are compared byte for byte, in the |
| 1097 | /// encoding a message would put them in. |
| 1098 | #[test] |
| 1099 | fn packing_changes_the_disk_and_not_the_wire() -> Outcome<()> { |
| 1100 | let scratch = res!(Scratch::new("wire")); |
| 1101 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 1102 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 1103 | let trust = Trust::of(&[key.binding()]); |
| 1104 | |
| 1105 | // The send set itself, which is what a message carries: an `Entry` compares |
| 1106 | // as the bytes it holds -- the payload, the public key and the signature -- |
| 1107 | // so equality here is equality on the wire. |
| 1108 | let sending = |store: &Store| -> Outcome<Vec<Entry>> { |
| 1109 | let done = res!(store.replay(Verify::Signatures(&trust), Keep::Envelopes)); |
| 1110 | let before = std::collections::BTreeSet::new(); |
| 1111 | let veils = std::collections::BTreeMap::new(); |
| 1112 | Ok(crate::store::arrived(&done.log, &before, &done.envelopes, &veils)) |
| 1113 | }; |
| 1114 | let plain = res!(sending(&store)); |
| 1115 | assert_eq!(plain.len(), MARKS as usize, "there is a payload to compare"); |
| 1116 | assert!(plain.iter().all(|e| matches!(e, Entry::Sealed(_))), |
| 1117 | "and it is the sealed form, which is what a peer has to check"); |
| 1118 | |
| 1119 | res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 1120 | segment_bytes: 200_000, |
| 1121 | packing: Packing::Packed, |
| 1122 | })); |
| 1123 | assert!(res!(store.bytes()) < 1_800_000, "the disk really did shrink"); |
| 1124 | assert_eq!(res!(sending(&store)), plain, |
| 1125 | "and the bytes a sync would hand over are the same ones"); |
| 1126 | Ok(()) |
| 1127 | } |
| 1128 | |
| 1129 | /// A repack takes the lock, and leaves nobody holding it. |
| 1130 | /// |
| 1131 | /// It is a read-modify-write of every segment there is, so a command |
| 1132 | /// appending to one while it ran would have its operation written into a log |
| 1133 | /// that was about to be replaced by one made without it. |
| 1134 | #[test] |
| 1135 | fn a_repack_takes_the_lock_and_gives_it_back() -> Outcome<()> { |
| 1136 | let scratch = res!(Scratch::new("lock")); |
| 1137 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 1138 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 1139 | let trust = Trust::of(&[key.binding()]); |
| 1140 | |
| 1141 | let held = res!(store.lock()); |
| 1142 | let refused = match repack(&scratch.path, Verify::Signatures(&trust), &Plan::default()) { |
| 1143 | Ok(_) => return Err(err!( |
| 1144 | "A repack ran while another command held the lock."; Test, Invalid)), |
| 1145 | Err(e) => fmt!("{}", e.plain()), |
| 1146 | }; |
| 1147 | assert!(refused.contains("Another `ore` command"), "and says so: {}", refused); |
| 1148 | drop(held); |
| 1149 | |
| 1150 | res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 1151 | segment_bytes: 200_000, |
| 1152 | packing: Packing::Plain, |
| 1153 | })); |
| 1154 | // The file stays, because that is where the lock is kept rather than being |
| 1155 | // the lock. What went with the repack is the note, so nothing reads a |
| 1156 | // finished command's name as a running one's. |
| 1157 | let where_it_lives = crate::store::Lock::path_of(&scratch.path); |
| 1158 | assert!(where_it_lives.is_file(), "and the lock file is still where it lives"); |
| 1159 | assert!(res!(fs::read_to_string(&where_it_lives)).trim().is_empty(), |
| 1160 | "the repack left its name in the lock"); |
| 1161 | // A lock somebody is genuinely holding is still refused, and named. |
| 1162 | let stranded = res!(store.lock()); |
| 1163 | let said = match store.lock() { |
| 1164 | Ok(_) => return Err(err!("A held lock was taken twice."; Test, Invalid)), |
| 1165 | Err(e) => fmt!("{}", e.plain()), |
| 1166 | }; |
| 1167 | assert!(said.contains("Another `ore` command"), "and says why: {}", said); |
| 1168 | assert!(said.contains(&fmt!("process {} ", std::process::id())), |
| 1169 | "and names what holds it: {}", said); |
| 1170 | drop(stranded); |
| 1171 | let done = res!(store.replay(Verify::Signatures(&trust), Keep::Nothing)); |
| 1172 | assert_eq!(done.log.len(), MARKS as usize, "while the history reads regardless"); |
| 1173 | Ok(()) |
| 1174 | } |
| 1175 | |
| 1176 | /// A segment of old operations is not bumped to a version an older reader |
| 1177 | /// would refuse. |
| 1178 | #[test] |
| 1179 | fn a_segment_declares_the_lowest_version_that_carries_it() -> Outcome<()> { |
| 1180 | assert_eq!(version_for(oxedyne_fe2o3_ore::op::CODE_MARK), segment::VERSION_MIN); |
| 1181 | assert_eq!(version_for(oxedyne_fe2o3_ore::op::CODE_NOTE), segment::VERSION_MIN); |
| 1182 | assert_eq!(version_for(oxedyne_fe2o3_ore::op::CODE_FILE_MODE), 3); |
| 1183 | assert_eq!(version_for(oxedyne_fe2o3_ore::op::CODE_REVERTS), segment::VERSION); |
| 1184 | assert_eq!(version_for(segment::highest_code(segment::VERSION)), segment::VERSION); |
| 1185 | |
| 1186 | let scratch = res!(Scratch::new("version")); |
| 1187 | let key = res!(Signing::mint(ReplicaId::new(1))); |
| 1188 | let store = res!(one_big_segment(&scratch.path, &key)); |
| 1189 | let trust = Trust::of(&[key.binding()]); |
| 1190 | res!(repack(&scratch.path, Verify::Signatures(&trust), &Plan { |
| 1191 | segment_bytes: 200_000, |
| 1192 | packing: Packing::Plain, |
| 1193 | })); |
| 1194 | // Marks alone, so every repacked segment declares the oldest version there |
| 1195 | // is and every reader that ever read this repository still can. |
| 1196 | for (_, path) in res!(store.segments()) { |
| 1197 | let bytes = res!(fs::read(&path)); |
| 1198 | let (head, _) = res!(res!(Head::decode(&bytes)).ok_or_else(|| err!( |
| 1199 | "The repacked segment {:?} has no header.", path; Test, Missing))); |
| 1200 | assert_eq!(head.version, segment::VERSION_MIN, |
| 1201 | "a segment of marks needs nothing newer"); |
| 1202 | assert_eq!(head.replica, Some(ReplicaId::new(1)), |
| 1203 | "and one replica wrote all of it"); |
| 1204 | } |
| 1205 | let _ = SALT; |
| 1206 | Ok(()) |
| 1207 | } |
| 1208 | } |