oxedyne/fe2o3/fe2o3_o3db_sync/src/sweep.rs
13.2 KiB, 9 runs
created by r1870400018:38131, 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 | //! In-place, online orphan sweep for an Ozone store. |
| 2 | //! |
| 3 | //! Garbage collection here is supersession-based: a record's bytes become reclaimable only when a |
| 4 | //! newer record is written at the same stored key. A chunked value's chunk-data records are keyed |
| 5 | //! by the value's geometry, so an overwrite that changes the geometry -- or, on a pre-fix build, any |
| 6 | //! overwrite at all, because the chunk set was keyed by a fresh random ticket -- leaves the old |
| 7 | //! chunk records under keys no live bunch key names. Nothing rewrites those keys, so nothing |
| 8 | //! flags them old, so the collector never reaches them: they are orphans, and on the gateway they |
| 9 | //! reached ~19 GB. [`sweep_orphans`] reclaims them in place, against a live store, without downtime |
| 10 | //! or a second store. |
| 11 | //! |
| 12 | //! The sweep is the in-place counterpart of [`crate::migrate`]: the migration copies the live set |
| 13 | //! into a fresh store and achieves zero residue (no orphans, no tombstones, no empty file states) |
| 14 | //! but needs the store stopped and twice the disk transiently; the sweep runs online and reclaims |
| 15 | //! the large chunk orphans, at the cost of leaving a small dead tombstone per orphan behind (the |
| 16 | //! residue §"What it does not do" names). Operators run the sweep routinely and a migration rarely. |
| 17 | //! |
| 18 | //! # How it stays correct against concurrent writers |
| 19 | //! |
| 20 | //! A record's `Meta.time` is stamped once per write and cloned identically onto the bunch key and |
| 21 | //! every chunk of a chunked value (`api::OzoneApi::prepare_write`), so a chunked value's bunch key |
| 22 | //! and all its chunks share one timestamp. The sweep stamps `T0` before it scans and retires a |
| 23 | //! chunk-data key only when it is both (a) absent from the live set reconstructed from every live |
| 24 | //! bunch key and (b) stamped strictly before `T0 - epoch_skew`. A value written after `T0` has |
| 25 | //! chunk timestamps at or after `T0`, so it is never retired even if its bunch key was written too |
| 26 | //! late for the live scan to see it; a value present and still live before `T0` has its chunks in |
| 27 | //! the live set and is kept. A conservative `epoch_skew` (a few seconds) only defers reclaiming a |
| 28 | //! genuine orphan to a later sweep; it never retires a live chunk. `epoch_skew` of zero degenerates |
| 29 | //! to the quiesced/offline case, correct only when nothing is writing. |
| 30 | //! |
| 31 | //! # What it does not do |
| 32 | //! |
| 33 | //! Retiring a large chunk orphan converts it into a small dead tombstone (a deleted-kind marker |
| 34 | //! record at the chunk key). The sweep therefore turns gigabytes of chunk orphans into a few |
| 35 | //! mebibytes of tombstones -- a >99.9% reclaim -- but does not remove the key. Clearing the |
| 36 | //! tombstones is left to the offline migration, which drops them by never copying them. |
| 37 | //! |
| 38 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 39 | //! Anthropic Claude |
| 40 | |
| 41 | use crate::{ |
| 42 | prelude::*, |
| 43 | comm::response::Wait, |
| 44 | }; |
| 45 | |
| 46 | use oxedyne_fe2o3_data::time::Timestamp; |
| 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::HashSet, |
| 58 | fs, |
| 59 | path::Path, |
| 60 | time::Duration, |
| 61 | }; |
| 62 | |
| 63 | |
| 64 | /// What a sweep found and retired. |
| 65 | /// |
| 66 | /// `bytes_before` is the store's data-file footprint at the instant the sweep started, measured over |
| 67 | /// the store's root directory, so it is meaningful only when the store keeps its data under that root |
| 68 | /// (no absolute zone override pointing elsewhere), which is the single-directory arrangement this |
| 69 | /// tool is built for. There is deliberately no `bytes_after`: collection is asynchronous and the |
| 70 | /// sweep returns as soon as every tombstone write is acknowledged, so the reclaimed footprint is only |
| 71 | /// visible once the collector has settled -- a caller that wants it measures the root itself after its |
| 72 | /// own settle (the `o3db_sweep` example does exactly this). |
| 73 | #[derive(Clone, Debug)] |
| 74 | pub struct SweepReport { |
| 75 | pub scanned_main_keys: usize, // main keys the live scan emitted (Complete and bunch keys) |
| 76 | pub live_chunk_keys: usize, // chunk keys referenced by a live bunch key |
| 77 | pub chunk_data_keys: usize, // all chunk-data keys (cind >= 1) the inverse scan emitted |
| 78 | pub orphans_found: usize, // not in the live set and older than the epoch threshold |
| 79 | pub orphans_retired: usize, // tombstones written (equals orphans_found on success) |
| 80 | pub skipped_recent: usize, // not in the live set but at/after the epoch threshold: kept |
| 81 | pub bytes_before: u64, // zone data-file bytes when the sweep started |
| 82 | } |
| 83 | |
| 84 | impl SweepReport { |
| 85 | /// A human-readable multi-line summary for a log or a console. |
| 86 | pub fn summary(&self) -> String { |
| 87 | let mut s = String::new(); |
| 88 | s.push_str("Orphan sweep:\n"); |
| 89 | s.push_str(&fmt!(" scanned main keys : {}\n", self.scanned_main_keys)); |
| 90 | s.push_str(&fmt!(" live chunk keys : {}\n", self.live_chunk_keys)); |
| 91 | s.push_str(&fmt!(" chunk-data keys seen: {}\n", self.chunk_data_keys)); |
| 92 | s.push_str(&fmt!(" orphans retired : {} of {} found\n", self.orphans_retired, self.orphans_found)); |
| 93 | s.push_str(&fmt!(" skipped (too recent): {}\n", self.skipped_recent)); |
| 94 | s.push_str(&fmt!(" data bytes at start : {}\n", self.bytes_before)); |
| 95 | s |
| 96 | } |
| 97 | } |
| 98 | |
| 99 | |
| 100 | /// Reclaims orphaned chunk-data records in place and returns a report. |
| 101 | /// |
| 102 | /// The store must be running with garbage collection ON (unlike the migration, which needs it off): |
| 103 | /// the sweep retires an orphan by writing a deleted-kind tombstone at the chunk key, and it is the |
| 104 | /// running collector that then reclaims the superseded chunk bytes. The sweep issues no write to |
| 105 | /// any value; it only tombstones chunk keys it has proven orphaned. |
| 106 | /// |
| 107 | /// Safety rests on three things, each grounded in the store's own mechanics: |
| 108 | /// - **Exact membership.** The live set is reconstructed from every live bunch key exactly as the |
| 109 | /// reader reconstructs chunk keys, and membership is exact `[u64; 5]` tuple equality, so a |
| 110 | /// geometry-change orphan (same set id, different geometry) and a random-ticket orphan (unrelated |
| 111 | /// set id) are both absent from it while a live value's chunks are all present. |
| 112 | /// - **Conservative on uncertainty.** A failure to read any scanned main key aborts the sweep |
| 113 | /// rather than narrowing the live set, so a transient read error can never widen the orphan set. |
| 114 | /// - **Epoch guard.** Only chunk records stamped before `T0 - epoch_skew` are retired, so a value |
| 115 | /// written during the sweep is excluded even if its bunch key was enumerated late. `epoch_skew` |
| 116 | /// of zero is correct only for a quiesced store. |
| 117 | /// |
| 118 | /// # Arguments |
| 119 | /// * `scan_wait` - how long each of the two store-wide scans may take. A scan walks every index |
| 120 | /// file in every zone, so a large store needs far longer than the shared user-request deadline; |
| 121 | /// pass a generous wait. |
| 122 | /// * `epoch_skew` - the clock margin below `T0`; a few seconds guards against a writer whose clock |
| 123 | /// is marginally behind the sweeper's. Pass `Duration::ZERO` only for a store nothing is writing. |
| 124 | pub fn sweep_orphans< |
| 125 | const UIDL: usize, |
| 126 | UID: NumIdDat<UIDL> + 'static, |
| 127 | ENC: Encrypter + 'static, |
| 128 | KH: Hasher + 'static, |
| 129 | PR: Hasher + 'static, |
| 130 | CS: Checksummer + 'static, |
| 131 | >( |
| 132 | api: &OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 133 | user: UID, |
| 134 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 135 | scan_wait: Wait, |
| 136 | epoch_skew: Duration, |
| 137 | ) |
| 138 | -> Outcome<SweepReport> |
| 139 | { |
| 140 | let enc = api.schemes().encrypter(); |
| 141 | let or_enc = schms2.map(|s| s.encrypter()); |
| 142 | |
| 143 | // 1. Stamp the epoch before touching the store, and derive the retirement threshold. A chunk |
| 144 | // record is old enough to retire only if its timestamp is strictly below this. |
| 145 | let t0 = res!(Timestamp::now()); |
| 146 | let t0_dur: Duration = *t0; |
| 147 | let threshold_dur = t0_dur.saturating_sub(epoch_skew); |
| 148 | let threshold = Timestamp::new(threshold_dur.as_secs(), threshold_dur.subsec_nanos()); |
| 149 | |
| 150 | let bytes_before = res!(zone_data_bytes(api.db_root())); |
| 151 | |
| 152 | // 2. Build the LIVE chunk-key set from every live bunch key. Scan for main keys, fetch each to |
| 153 | // its stored value, and for a chunked value (a Tup5u64 bunch value) reconstruct its chunk |
| 154 | // keys exactly as the reader does. A read error aborts: never narrow the live set on |
| 155 | // uncertainty. An absent value (a tombstone, or a value mid-write) contributes no live |
| 156 | // chunks, which is safe because the epoch guard keeps any recent value's chunks regardless. |
| 157 | let main_opts = ScanOpts::all(); |
| 158 | let main_entries = res!(api.scan_with_wait(&main_opts, schms2, dup_wait(&scan_wait))); |
| 159 | let scanned_main_keys = main_entries.len(); |
| 160 | |
| 161 | let mut live: HashSet<[u64; 5]> = HashSet::new(); |
| 162 | for (kdat, _empty, _meta) in &main_entries { |
| 163 | let resp = res!(api.fetch_using_schemes(kdat, schms2)); |
| 164 | match res!(resp.recv_daticle(enc, or_enc)) { |
| 165 | (Some((Dat::Tup5u64(tup), _)), _) => { |
| 166 | // A chunked value: its bunch value is the part key naming its chunks. |
| 167 | let set_id = tup[0]; |
| 168 | let data_len = tup[2]; |
| 169 | let num_parts = tup[3]; |
| 170 | let part_size = tup[4]; |
| 171 | for i in 1..(num_parts + 1) { |
| 172 | live.insert([set_id, i, data_len, num_parts, part_size]); |
| 173 | } |
| 174 | }, |
| 175 | // An unchunked value, or an absent one (tombstone / mid-write): no live chunks. |
| 176 | _ => (), |
| 177 | } |
| 178 | } |
| 179 | let live_chunk_keys = live.len(); |
| 180 | |
| 181 | // 3. Enumerate every chunk-data key (the inverse scan), classify against the live set and the |
| 182 | // epoch, and collect the orphans to retire. Each candidate decodes to its Tup5u64 chunk |
| 183 | // key; one that does not is left untouched rather than guessed at. |
| 184 | let chunk_opts = ScanOpts::all().chunk_data_only(true); |
| 185 | let chunk_entries = res!(api.scan_with_wait(&chunk_opts, schms2, dup_wait(&scan_wait))); |
| 186 | let chunk_data_keys = chunk_entries.len(); |
| 187 | |
| 188 | let mut orphans: Vec<Dat> = Vec::new(); |
| 189 | let mut skipped_recent = 0usize; |
| 190 | for (kdat, _empty, meta) in &chunk_entries { |
| 191 | let tup = match kdat { |
| 192 | Dat::Tup5u64(arr) => *arr, |
| 193 | // A chunk-data key that is not a Tup5u64 cannot be classified; leave it alone. |
| 194 | _ => continue, |
| 195 | }; |
| 196 | if live.contains(&tup) { |
| 197 | continue; // Referenced by a live bunch key: keep. |
| 198 | } |
| 199 | // Not referenced: an orphan by membership. Retire only if old enough; a recent one is a |
| 200 | // value that may still be settling, so it is kept and reclaimed by a later sweep. |
| 201 | if meta.time < threshold { |
| 202 | orphans.push(kdat.clone()); |
| 203 | } else { |
| 204 | skipped_recent += 1; |
| 205 | } |
| 206 | } |
| 207 | let orphans_found = orphans.len(); |
| 208 | |
| 209 | // 4. Retire each orphan through the ordinary supersession path: an unencrypted deleted-kind |
| 210 | // tombstone at the chunk key, which supersedes the chunk record at the same stored key so the |
| 211 | // running collector flags its bytes old and reclaims them. One shared responder collects one |
| 212 | // acknowledgement per tombstone, so the function returns only once every write has landed, |
| 213 | // and a writer that fails one says so rather than being counted as an answer. |
| 214 | let mut orphans_retired = 0usize; |
| 215 | if orphans_found > 0 { |
| 216 | let resp = api.responder(); |
| 217 | for ck in &orphans { |
| 218 | res!(api.tombstone_chunk_key(ck, user, schms2, resp.clone())); |
| 219 | orphans_retired += 1; |
| 220 | } |
| 221 | let liveness = scan_wait.max_wait; |
| 222 | let durability = std::cmp::max(liveness, constant::DURABILITY_TIMEOUT); |
| 223 | res!(resp.recv_write_acks(orphans_retired, liveness, durability)); |
| 224 | } |
| 225 | |
| 226 | // The function returns here, as soon as every tombstone is acknowledged. It does NOT wait for |
| 227 | // the collector: reclamation is asynchronous, so a caller that wants the reclaimed footprint |
| 228 | // settles and re-measures the root itself (see the `o3db_sweep` example). |
| 229 | Ok(SweepReport { |
| 230 | scanned_main_keys, |
| 231 | live_chunk_keys, |
| 232 | chunk_data_keys, |
| 233 | orphans_found, |
| 234 | orphans_retired, |
| 235 | skipped_recent, |
| 236 | bytes_before, |
| 237 | }) |
| 238 | } |
| 239 | |
| 240 | /// A fresh copy of a `Wait`, since `scan_with_wait` and `recv_number` take it by value and it is not |
| 241 | /// `Copy`; its fields are, so this is a plain field copy. |
| 242 | fn dup_wait(w: &Wait) -> Wait { |
| 243 | Wait { |
| 244 | max_wait: w.max_wait, |
| 245 | check_interval: w.check_interval, |
| 246 | } |
| 247 | } |
| 248 | |
| 249 | /// Total bytes of the store's data files (`.dat`) under a root, walked recursively. Index files and |
| 250 | /// anything else are not counted; the reclaim the sweep drives shows up as data-file bytes. |
| 251 | fn zone_data_bytes(root: &Path) -> Outcome<u64> { |
| 252 | let mut total = 0u64; |
| 253 | let mut stack = vec![root.to_path_buf()]; |
| 254 | while let Some(dir) = stack.pop() { |
| 255 | let rd = match fs::read_dir(&dir) { |
| 256 | Ok(rd) => rd, |
| 257 | Err(_) => continue, // A directory that is not there contributes nothing. |
| 258 | }; |
| 259 | for entry in rd { |
| 260 | let entry = res!(entry); |
| 261 | let path = entry.path(); |
| 262 | if path.is_dir() { |
| 263 | stack.push(path); |
| 264 | } else if path.extension().map(|e| e == constant::DATA_FILE_EXT).unwrap_or(false) { |
| 265 | let meta = res!(entry.metadata()); |
| 266 | total += meta.len(); |
| 267 | } |
| 268 | } |
| 269 | } |
| 270 | Ok(total) |
| 271 | } |