Oregami
Repositories/oxedyne/fe2o3

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
41use crate::{
42 prelude::*,
43 comm::response::Wait,
44};
45
46use oxedyne_fe2o3_data::time::Timestamp;
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::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)]
74pub 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
84impl 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.
124pub 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.
242fn 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.
251fn 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}