Oregami
Repositories/oxedyne/fe2o3

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
42use crate::{
43 prelude::*,
44 comm::response::Wait,
45};
46
47use oxedyne_fe2o3_jdat::{
48 prelude::*,
49 id::NumIdDat,
50};
51use oxedyne_fe2o3_iop_db::api::{
52 RestSchemesOverride,
53 ScanOpts,
54};
55
56use 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)]
70pub 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
83impl 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.
129pub 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.
200pub 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.
296fn 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.
319fn 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.
329fn 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.
341fn 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}