oxedyne/fe2o3/fe2o3_o3db_sync/src/migrate.rs
14.8 KiB, 11 runs
created by r1870400018:38112, 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 | //! Live-set migration and compaction for an Ozone store. |
| 2 | //! |
| 3 | //! [`migrate_live_set`] copies exactly the live key set of a source store into |
| 4 | //! a fresh target store and, in doing so, drops the orphaned chunk-data records |
| 5 | //! that a pre-fix build leaked (chunk keys left behind when a chunked value was |
| 6 | //! overwritten or deleted under a fresh random ticket). Nothing about the copy |
| 7 | //! is chunk-aware: it rests on three properties of the store this crate builds, |
| 8 | //! each verified against the code rather than assumed. |
| 9 | //! |
| 10 | //! - **`scan` answers from the live index and emits only main user keys.** A |
| 11 | //! scan walks every index file in every zone and returns each live |
| 12 | //! `Key::Complete` and each live bunch key `Key::Chunk(_, 0)`, eliding the |
| 13 | //! chunk-data keys (index `>= 1`) that carry a chunked value's bytes |
| 14 | //! (`bots/worker/bot_scan.rs`, the `c >= 1` continue). An orphaned chunk key |
| 15 | //! has no live bunch key naming it, so a scan never emits it and the copy |
| 16 | //! never carries it: that, and only that, is what reclaims the leak. |
| 17 | //! - **A scan enumerates every live `Complete` key independent of references.** |
| 18 | //! The walk is over the index files themselves, not a reachability walk from |
| 19 | //! any root, so a value stored under its own `Complete` key that nothing else |
| 20 | //! points at -- a chat body keyed `sbody:<account>:<hex>`, referenced only by |
| 21 | //! the hashes inside a separate `sync:` record -- is enumerated and copied |
| 22 | //! like any other key. A reachability-based copy would drop it; this does |
| 23 | //! not. |
| 24 | //! - **`get` on a bunch key returns the whole current value.** `get_wait` |
| 25 | //! hands a `Dat::Tup5u64` bunch key to `fetch_chunks`, which reconstructs the |
| 26 | //! current chunk keys from the stored `PartKey` and rejoins, decrypts and |
| 27 | //! decodes them, so a copied value is the complete, current value of whatever |
| 28 | //! kind it was stored as. |
| 29 | //! |
| 30 | //! Both stores are opened under the same schemes, including the same at-rest |
| 31 | //! encryption key: `get_wait` on the source returns plaintext and `store` on |
| 32 | //! the target re-encrypts it, so the target holds the same values, freshly |
| 33 | //! encrypted. The source is only ever read; every write goes to the target. |
| 34 | //! |
| 35 | //! A key whose newest record is a deletion tombstone reads back as absent |
| 36 | //! (`Responder::recv_daticle` maps the deleted-kind marker to `None`), so a |
| 37 | //! scan emits it but `get_wait` yields nothing. The migration copies nothing |
| 38 | //! for such a key, which is the correct compaction: the target simply lacks the |
| 39 | //! key and reads it as absent, exactly as the source did. Tombstones are dead |
| 40 | //! weight that the copy leaves behind. |
| 41 | |
| 42 | use crate::{ |
| 43 | prelude::*, |
| 44 | comm::response::Wait, |
| 45 | }; |
| 46 | |
| 47 | use oxedyne_fe2o3_jdat::{ |
| 48 | prelude::*, |
| 49 | id::NumIdDat, |
| 50 | }; |
| 51 | use oxedyne_fe2o3_iop_db::api::{ |
| 52 | RestSchemesOverride, |
| 53 | ScanOpts, |
| 54 | }; |
| 55 | |
| 56 | use std::{ |
| 57 | collections::BTreeMap, |
| 58 | fs, |
| 59 | path::Path, |
| 60 | }; |
| 61 | |
| 62 | |
| 63 | /// What a migration copied, and whether it verified. |
| 64 | /// |
| 65 | /// The byte and file totals are measured over each store's root directory, so |
| 66 | /// they are meaningful only when the store keeps its data under that root (no |
| 67 | /// absolute zone override pointing elsewhere), which is the single-directory |
| 68 | /// arrangement this tool is built for. |
| 69 | #[derive(Clone, Debug)] |
| 70 | pub struct MigrationReport { |
| 71 | pub scanned: usize, // keys the source scan emitted |
| 72 | pub copied: usize, // keys with a value present, copied to the target |
| 73 | pub tombstones: usize, // scanned keys whose value was absent (deleted) |
| 74 | pub verified: usize, // copied keys whose target bytes matched the source |
| 75 | pub per_prefix: BTreeMap<String, usize>,// copied keys grouped by key prefix (before first ':') |
| 76 | pub source_bytes: u64, |
| 77 | pub target_bytes: u64, |
| 78 | pub source_files: usize, |
| 79 | pub target_files: usize, |
| 80 | pub verified_ok: bool, // every copied key read back byte-identical, counts agreed |
| 81 | } |
| 82 | |
| 83 | impl MigrationReport { |
| 84 | /// A human-readable multi-line summary for a log or a console. |
| 85 | pub fn summary(&self) -> String { |
| 86 | let mut s = String::new(); |
| 87 | s.push_str(&fmt!( |
| 88 | "Migration {}:\n", |
| 89 | if self.verified_ok { "VERIFIED" } else { "NOT verified" })); |
| 90 | s.push_str(&fmt!(" scanned live keys : {}\n", self.scanned)); |
| 91 | s.push_str(&fmt!(" copied (present) : {}\n", self.copied)); |
| 92 | s.push_str(&fmt!(" tombstones (drop) : {}\n", self.tombstones)); |
| 93 | s.push_str(&fmt!(" verified round-trip: {}\n", self.verified)); |
| 94 | s.push_str(&fmt!(" source bytes : {} in {} files\n", self.source_bytes, self.source_files)); |
| 95 | s.push_str(&fmt!(" target bytes : {} in {} files\n", self.target_bytes, self.target_files)); |
| 96 | if self.source_bytes > 0 { |
| 97 | let pct = (self.target_bytes as f64 / self.source_bytes as f64) * 100.0; |
| 98 | s.push_str(&fmt!(" target is {:.2}% of source by bytes\n", pct)); |
| 99 | } |
| 100 | s.push_str(" copied keys by prefix:\n"); |
| 101 | for (prefix, n) in &self.per_prefix { |
| 102 | s.push_str(&fmt!(" {:<24} {}\n", prefix, n)); |
| 103 | } |
| 104 | s |
| 105 | } |
| 106 | } |
| 107 | |
| 108 | |
| 109 | /// Copies the live key set of `source` into `target`, verifies the copy, and |
| 110 | /// returns a report. |
| 111 | /// |
| 112 | /// The function is fail-closed: it returns an error if any copied key does not |
| 113 | /// read back from the target byte-identical to the source, or if the target's |
| 114 | /// live-key count does not equal the number of present source keys. A returned |
| 115 | /// `Ok` report always has `verified_ok == true`. Deciding whether the target |
| 116 | /// is *small enough* to trust (the orphan-shrink gate) is left to the caller, |
| 117 | /// which knows how much shrink to expect; the report carries the byte totals it |
| 118 | /// needs. |
| 119 | /// |
| 120 | /// `target` must be freshly created with the *same* configuration as `source` |
| 121 | /// (so chunk geometry and encryption match) and opened under the same schemes. |
| 122 | /// `source` must be opened with garbage collection off and must not be written |
| 123 | /// to by anyone else for the duration; this function issues no write to it. |
| 124 | /// |
| 125 | /// # Arguments |
| 126 | /// * `scan_wait` - how long each zone's scan may take. A scan walks every index |
| 127 | /// file in the store, so a large store needs far longer than the shared |
| 128 | /// user-request deadline; pass a generous wait. |
| 129 | pub fn migrate_live_set< |
| 130 | const UIDL: usize, |
| 131 | UID: NumIdDat<UIDL> + 'static, |
| 132 | ENC: Encrypter + 'static, |
| 133 | KH: Hasher + 'static, |
| 134 | PR: Hasher + 'static, |
| 135 | CS: Checksummer + 'static, |
| 136 | >( |
| 137 | source: &OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 138 | target: &OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 139 | user: UID, |
| 140 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 141 | scan_wait: Wait, |
| 142 | ) |
| 143 | -> Outcome<MigrationReport> |
| 144 | { |
| 145 | // 1. Enumerate every live main key in the source, across all zones. No |
| 146 | // prefix, no limit, values excluded: scan v1 returns Dat::Empty for |
| 147 | // values by design, and the value is fetched per key below. |
| 148 | let opts = ScanOpts::all(); |
| 149 | let entries = res!(source.scan_with_wait(&opts, schms2, dup_wait(&scan_wait))); |
| 150 | let scanned = entries.len(); |
| 151 | |
| 152 | let mut copied = 0usize; |
| 153 | let mut tombstones = 0usize; |
| 154 | let mut per_prefix: BTreeMap<String, usize> = BTreeMap::new(); |
| 155 | |
| 156 | // 2. Copy each present value into the target. A key whose value is absent |
| 157 | // is a live tombstone: it is dropped, not copied. |
| 158 | for (kdat, _empty, _meta) in &entries { |
| 159 | match res!(source.get_wait(kdat, schms2)) { |
| 160 | None => { |
| 161 | tombstones += 1; |
| 162 | }, |
| 163 | Some((val, _meta)) => { |
| 164 | res!(store_and_wait(target, kdat.clone(), val, user, schms2)); |
| 165 | copied += 1; |
| 166 | let prefix = key_prefix(kdat); |
| 167 | *per_prefix.entry(prefix).or_insert(0) += 1; |
| 168 | }, |
| 169 | } |
| 170 | } |
| 171 | |
| 172 | // 3. Verify before reporting success. Every present source key must read |
| 173 | // back from the target byte-identical, and the target must hold exactly |
| 174 | // the present source keys and no more. |
| 175 | let verified = res!(verify_migration(source, target, user, schms2, scan_wait, copied)); |
| 176 | |
| 177 | let (source_bytes, source_files) = res!(dir_bytes_and_files(source.db_root())); |
| 178 | let (target_bytes, target_files) = res!(dir_bytes_and_files(target.db_root())); |
| 179 | |
| 180 | Ok(MigrationReport { |
| 181 | scanned, |
| 182 | copied, |
| 183 | tombstones, |
| 184 | verified, |
| 185 | per_prefix, |
| 186 | source_bytes, |
| 187 | target_bytes, |
| 188 | source_files, |
| 189 | target_files, |
| 190 | verified_ok: true, |
| 191 | }) |
| 192 | } |
| 193 | |
| 194 | /// Confirms the target holds exactly the present source keys, each byte-identical. |
| 195 | /// |
| 196 | /// Returns the number of keys checked (the present source keys). Errors on the |
| 197 | /// first mismatch or count disagreement, so a caller that gets `Ok` has a |
| 198 | /// verified copy. `expected_present` is the count of keys the copy pass wrote, |
| 199 | /// used to cross-check against an independent re-scan of both stores. |
| 200 | pub fn verify_migration< |
| 201 | const UIDL: usize, |
| 202 | UID: NumIdDat<UIDL> + 'static, |
| 203 | ENC: Encrypter + 'static, |
| 204 | KH: Hasher + 'static, |
| 205 | PR: Hasher + 'static, |
| 206 | CS: Checksummer + 'static, |
| 207 | >( |
| 208 | source: &OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 209 | target: &OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 210 | _user: UID, |
| 211 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 212 | scan_wait: Wait, |
| 213 | expected_present: usize, |
| 214 | ) |
| 215 | -> Outcome<usize> |
| 216 | { |
| 217 | let opts = ScanOpts::all(); |
| 218 | |
| 219 | // A. Re-scan the source independently and compare each key's value against |
| 220 | // the target, byte for byte. A tombstone (absent value) must be absent |
| 221 | // in the target too. |
| 222 | let src_entries = res!(source.scan_with_wait(&opts, schms2, dup_wait(&scan_wait))); |
| 223 | let mut present = 0usize; |
| 224 | for (kdat, _empty, _meta) in &src_entries { |
| 225 | let sv = res!(source.get_wait(kdat, schms2)); |
| 226 | let tv = res!(target.get_wait(kdat, schms2)); |
| 227 | match (sv, tv) { |
| 228 | (None, None) => { |
| 229 | // A tombstoned source key: correctly absent from the target. |
| 230 | }, |
| 231 | (None, Some(_)) => return Err(err!( |
| 232 | "Verify: source key {:?} is a tombstone (absent) but the target \ |
| 233 | holds a value for it; the copy invented a key.", kdat; |
| 234 | Data, Mismatch)), |
| 235 | (Some(_), None) => return Err(err!( |
| 236 | "Verify: source key {:?} has a value but the target does not; the \ |
| 237 | copy dropped a live key.", kdat; |
| 238 | Data, Missing)), |
| 239 | (Some((svd, _)), Some((tvd, _))) => { |
| 240 | let sb = res!(svd.as_bytes()); |
| 241 | let tb = res!(tvd.as_bytes()); |
| 242 | if sb != tb { |
| 243 | return Err(err!( |
| 244 | "Verify: source key {:?} reads back {} value bytes from the \ |
| 245 | target but {} from the source; the copied value differs.", |
| 246 | kdat, tb.len(), sb.len(); |
| 247 | Data, Mismatch)); |
| 248 | } |
| 249 | present += 1; |
| 250 | }, |
| 251 | } |
| 252 | } |
| 253 | |
| 254 | if present != expected_present { |
| 255 | return Err(err!( |
| 256 | "Verify: the copy pass wrote {} present keys, but an independent \ |
| 257 | re-scan of the source finds {} present keys; the source changed \ |
| 258 | underneath the migration, or a key was miscounted.", |
| 259 | expected_present, present; |
| 260 | Data, Mismatch)); |
| 261 | } |
| 262 | |
| 263 | // B. Re-scan the target: it must hold exactly the present source keys, none |
| 264 | // of them a tombstone, and no extras. |
| 265 | let tgt_entries = res!(target.scan_with_wait(&opts, schms2, dup_wait(&scan_wait))); |
| 266 | let mut tgt_present = 0usize; |
| 267 | for (kdat, _empty, _meta) in &tgt_entries { |
| 268 | match res!(target.get_wait(kdat, schms2)) { |
| 269 | None => return Err(err!( |
| 270 | "Verify: the target scan emitted key {:?} but it reads back absent; \ |
| 271 | the target holds a tombstone, which a fresh copy should never write.", |
| 272 | kdat; |
| 273 | Data, Mismatch)), |
| 274 | Some(_) => tgt_present += 1, |
| 275 | } |
| 276 | } |
| 277 | if tgt_present != present { |
| 278 | return Err(err!( |
| 279 | "Verify: the target holds {} live keys but the source has {} present \ |
| 280 | keys; the live-key counts do not agree.", |
| 281 | tgt_present, present; |
| 282 | Data, Mismatch)); |
| 283 | } |
| 284 | |
| 285 | Ok(present) |
| 286 | } |
| 287 | |
| 288 | /// Stores one key-value pair into `api` and waits for every write part to be |
| 289 | /// acknowledged, so the value is in the target index before anything reads it. |
| 290 | /// |
| 291 | /// A store dispatches `n` write parts -- one `Complete` record, or one bunch key |
| 292 | /// plus one record per chunk -- and reports `n` first as an `OzoneMsg::Chunks`, |
| 293 | /// then answers each part twice, written and then durable. Waiting for all `n` |
| 294 | /// final answers, rather than on ordering between writers, is what makes the |
| 295 | /// copy safe to verify immediately. |
| 296 | fn store_and_wait< |
| 297 | const UIDL: usize, |
| 298 | UID: NumIdDat<UIDL> + 'static, |
| 299 | ENC: Encrypter + 'static, |
| 300 | KH: Hasher + 'static, |
| 301 | PR: Hasher + 'static, |
| 302 | CS: Checksummer + 'static, |
| 303 | >( |
| 304 | api: &OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 305 | k: Dat, |
| 306 | v: Dat, |
| 307 | user: UID, |
| 308 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 309 | ) |
| 310 | -> Outcome<()> |
| 311 | { |
| 312 | let resp = res!(api.store_using_schemes(k, v, user, schms2)); |
| 313 | res!(resp.recv_store_ack()); |
| 314 | Ok(()) |
| 315 | } |
| 316 | |
| 317 | /// A fresh copy of a `Wait`, since `scan_with_wait` takes it by value and it is |
| 318 | /// not `Copy`; its fields are, so this is a plain field copy. |
| 319 | fn dup_wait(w: &Wait) -> Wait { |
| 320 | Wait { |
| 321 | max_wait: w.max_wait, |
| 322 | check_interval: w.check_interval, |
| 323 | } |
| 324 | } |
| 325 | |
| 326 | /// The part of a string key before its first ':', or the whole key if it has |
| 327 | /// none; a non-string key groups under `<non-str>`. Used only to group the |
| 328 | /// report so an operator can eyeball that every expected keyspace is present. |
| 329 | fn key_prefix(k: &Dat) -> String { |
| 330 | match k { |
| 331 | Dat::Str(s) => match s.find(':') { |
| 332 | Some(i) => s[..i].to_string(), |
| 333 | None => s.clone(), |
| 334 | }, |
| 335 | _ => "<non-str>".to_string(), |
| 336 | } |
| 337 | } |
| 338 | |
| 339 | /// Total bytes and file count under a directory, walked recursively. Symlinks |
| 340 | /// are not followed; a store keeps regular files only. |
| 341 | fn dir_bytes_and_files(root: &Path) -> Outcome<(u64, usize)> { |
| 342 | let mut bytes = 0u64; |
| 343 | let mut files = 0usize; |
| 344 | let mut stack = vec![root.to_path_buf()]; |
| 345 | while let Some(dir) = stack.pop() { |
| 346 | let rd = match fs::read_dir(&dir) { |
| 347 | Ok(rd) => rd, |
| 348 | Err(_) => continue, // A directory that is not there contributes nothing. |
| 349 | }; |
| 350 | for entry in rd { |
| 351 | let entry = res!(entry); |
| 352 | let path = entry.path(); |
| 353 | let meta = res!(entry.metadata()); |
| 354 | if meta.is_dir() { |
| 355 | stack.push(path); |
| 356 | } else if meta.is_file() { |
| 357 | bytes += meta.len(); |
| 358 | files += 1; |
| 359 | } |
| 360 | } |
| 361 | } |
| 362 | Ok((bytes, files)) |
| 363 | } |