Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/tests/scan_under_gc.rs

13.7 KiB, 36 runs

created by r1870400018:21762, 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//! Ozone scan-under-garbage-collection integration test.
2//!
3//! A scan is a foreground request with a person waiting on it; garbage
4//! collection is background work with no deadline at all. This test
5//! holds the two apart.
6//!
7//! It builds a store whose files are all past the collection trigger
8//! while collection is switched off, so no work has been done yet and
9//! none is queued. A quiet scan then measures what the walk costs with
10//! nothing in its way. Collection is switched on and a short burst of
11//! writes dispatches the whole backlog at once, and the same scan is
12//! issued again. Two properties are checked:
13//!
14//! 1. **A concurrent scan sees every key, and does not fail.**
15//! Collection rewrites a file's index. A scanner walking that index
16//! at the wrong moment can see a file with no index at all and drop
17//! every key whose current value lives in it, or read a half-written
18//! record and fail outright. Entry count and outcome are checked on
19//! every scan. This is the sharp half of the test: before the index
20//! rebuild was corrected it failed here every run.
21//!
22//! 2. **Latency does not track the collector's backlog.** The scan
23//! issued while the backlog is being worked through costs about what
24//! the quiet one did. This half is a tripwire rather than a proof.
25//! A one-shot backlog can only ever cost about as much as one full
26//! scan -- the collector's work over every file and the scan's walk
27//! over every file are both a fixed cost per record, so the whole
28//! backlog is roughly one scan's worth -- and the allowance below is
29//! set well clear of ordinary noise. What it catches is a change
30//! that puts the scan back on a queue behind unbounded background
31//! work; what it cannot show, at any size a test can afford, is the
32//! multi-second wait a store under sustained churn produces.
33//!
34//! A burst is used rather than a sustained write storm because it makes
35//! the backlog a known quantity. Under a storm the collector is idle
36//! most of the time -- work arrives for it in proportion to the write
37//! rate, and it serves that work about as fast as several writers can
38//! produce it -- so a queue forms only by luck, and the test would then
39//! measure the machine rather than the database.
40//!
41//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
42//! Anthropic Claude
43
44use oxedyne_fe2o3_core::{
45 prelude::*,
46 alt::Override,
47};
48use oxedyne_fe2o3_crypto::enc::EncryptionScheme;
49use oxedyne_fe2o3_hash::{
50 csum::ChecksumScheme,
51 hash::HashScheme,
52};
53use oxedyne_fe2o3_iop_db::api::{
54 Database,
55 RestSchemesOverride,
56 ScanOpts,
57};
58use oxedyne_fe2o3_jdat::prelude::*;
59use oxedyne_fe2o3_o3db_sync::{
60 data::core::RestSchemesInput,
61 test::setup,
62};
63
64use std::{
65 path::Path,
66 thread,
67 time::{
68 Duration,
69 Instant,
70 },
71};
72
73/// The dimensions of one run.
74struct Shape {
75 dir: &'static str, // directory name for this run's store
76 keys: usize, // distinct keys held in the store
77 // Value payload size in bytes, before encoding. Large, so that a
78 // collection moves far more than a scan reads.
79 val_bytes: usize,
80 file_bytes: u64, // and so the size of one collection's work
81}
82
83// The store the crate's own test suite can afford.
84const QUICK: Shape = Shape {
85 dir: "./test_db_scan_under_gc",
86 keys: 6_000,
87 val_bytes: 20_000,
88 file_bytes: 8_000_000,
89};
90
91// A store several times larger, where the backlog is deep enough to cost more
92// than the user request deadline outright rather than only running late. Writes
93// some hundreds of megabytes.
94const DEEP: Shape = Shape {
95 dir: "./test_db_scan_stall",
96 keys: 30_000,
97 val_bytes: 20_000,
98 file_bytes: 8_000_000,
99};
100
101// Fraction of the keys superseded before collection is switched on. It has to
102// clear OLD_DATA_PERCENT_GC_TRIGGER, which is measured against the maximum data
103// file size rather than the file's own size.
104const OLD_FRAC: f64 = 0.5;
105
106const NUM_BASELINE: usize = 3; // quiet scans establishing the latency baseline
107const NUM_LOADED: usize = 5; // scans issued while the backlog is worked off
108
109/// The spacing of the keys superseded to set the backlog going: about two per
110/// file, which reaches every file holding original
111/// records while keeping the burst far shorter than the collection work
112/// it dispatches. A finer stride lets the collector work the queue down
113/// as fast as it is filled, and the test then measures nothing.
114fn trigger_stride(shape: &Shape) -> usize {
115 let per_file = (shape.file_bytes as usize) / shape.val_bytes;
116 std::cmp::max(1, per_file / 2)
117}
118
119/// Runs the scan-under-collection test as its own integration binary.
120///
121/// It is kept out of `tests/main.rs` deliberately: it is the slowest
122/// test in the crate and the only one that wants the disk to itself, so
123/// `cargo test -p oxedyne_fe2o3_o3db_sync --test scan_under_gc` can run
124/// it alone.
125#[test]
126fn main() -> Outcome<()> {
127
128 // Logging every stored key costs several times what the writes do,
129 // and this test writes a great many of them. `test` is the lowest
130 // level that still carries the test's own reporting.
131 log_set_level!("test");
132
133 let outcome = test_scan_under_gc(&QUICK);
134
135 log_finish_wait!();
136
137 outcome
138}
139
140/// The same test against a store deep enough to reproduce the
141/// production symptom rather than only the mechanism behind it. It
142/// writes for some minutes, so it is ignored by default:
143///
144/// ```ignore
145/// cargo test -p oxedyne_fe2o3_o3db_sync --test scan_under_gc -- --ignored --nocapture
146/// ```
147#[test]
148#[ignore]
149fn deep() -> Outcome<()> {
150
151 log_set_level!("test");
152
153 let outcome = test_scan_under_gc(&DEEP);
154
155 log_finish_wait!();
156
157 outcome
158}
159
160pub fn test_scan_under_gc(shape: &Shape) -> Outcome<()> {
161
162 let db_root = res!(Path::new(shape.dir).canonicalize().or_else(|_| {
163 ok!(std::fs::create_dir_all(shape.dir));
164 Path::new(shape.dir).canonicalize()
165 }));
166
167 // Fixed key so the test is deterministic.
168 let enckey = [0x37u8; 32];
169 let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..]));
170 let crc32 = ChecksumScheme::new_crc32();
171 let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> =
172 RestSchemesOverride::default()
173 .set_encrypter(Override::Default(aes_gcm.clone()));
174 let schms2 = Some(&schms2);
175 let user = setup::Uid::default();
176
177 let schms_input = RestSchemesInput::new(
178 Some(aes_gcm.clone()),
179 None::<HashScheme>,
180 None::<HashScheme>,
181 Some(crc32.clone()),
182 );
183
184 let mut cfg = res!(setup::default_cfg());
185 // One zone, and one bot of each kind in it, so there is exactly one
186 // queue per role and no ambiguity about what a scan is waiting for.
187 cfg.num_zones = 1;
188 cfg.num_cbots_per_zone = 1;
189 cfg.num_fbots_per_zone = 1;
190 cfg.num_wbots_per_zone = 1;
191 // One collector is the sharpest form of the defect: a scan sent to
192 // the collector's pool is then guaranteed to sit behind whatever is
193 // already queued there. It is a legal production setting, and with
194 // more collectors the defect merely becomes intermittent.
195 cfg.num_igbots_per_zone = 1;
196 cfg.data_file_max_bytes = shape.file_bytes;
197 // Values are well under the file size, so chunking stays out of the
198 // way and each file holds whole records only.
199 cfg.rest_chunk_threshold = shape.file_bytes / 8;
200 cfg.rest_chunk_bytes = shape.file_bytes / 40;
201 // Small enough that the cache jettisons values, so a scan cannot be
202 // served from a fully resident cache by accident.
203 cfg.cache_size_limit_bytes = 40_000_000;
204 cfg.zone_overrides = mapdat!{
205 1u16 => mapdat!{ "dir" => "", "max_size" => 8_000_000_000u64 },
206 }.get_map().unwrap();
207
208 test!(sync_log::stream(), "+---------------------------------------------+");
209 test!(sync_log::stream(), "| SCAN UNDER GARBAGE COLLECTION TEST |");
210 test!(sync_log::stream(), "+---------------------------------------------+");
211 test!(sync_log::stream(), "{} keys of {} bytes in {} byte files.",
212 shape.keys, shape.val_bytes, shape.file_bytes);
213
214 // Collection stays off while the store is built, so that nothing is
215 // collected until the test chooses the moment.
216 let mut db = res!(setup::start_db(
217 db_root.clone(),
218 Some(cfg.clone()),
219 schms_input.clone(),
220 None,
221 false, // gc off
222 true, // wipe
223 ));
224 thread::sleep(Duration::from_secs(1));
225
226 let payload = |i: usize, round: usize| -> Dat {
227 let seed = fmt!("k{:06}r{:02}:", i, round);
228 let mut s = String::with_capacity(shape.val_bytes);
229 while s.len() < shape.val_bytes {
230 s.push_str(&seed);
231 }
232 s.truncate(shape.val_bytes);
233 Dat::Str(s)
234 };
235
236 // 1. Populate.
237 test!(sync_log::stream(), "Populating {} keys.", shape.keys);
238 let t = Instant::now();
239 for i in 0..shape.keys {
240 res!(db.insert(dat!(fmt!("rec:{:06}", i)), payload(i, 0), user, schms2));
241 }
242 test!(sync_log::stream(), "Populated in {:?}.", t.elapsed());
243
244 // 2. Supersede every second key, with collection still off. Half of
245 // every file is now garbage -- well past the trigger -- and the
246 // other half is untouched, so a later write to any of those keys
247 // still lands on the file it was written to and can set its
248 // collection going.
249 let num_old = ((shape.keys as f64) * OLD_FRAC) as usize;
250 test!(sync_log::stream(), "Superseding {} keys with collection off.", num_old);
251 let t = Instant::now();
252 let mut i = 0;
253 while i < shape.keys {
254 res!(db.insert(dat!(fmt!("rec:{:06}", i)), payload(i, 1), user, schms2));
255 i += 2;
256 }
257 test!(sync_log::stream(), "Superseded in {:?}.", t.elapsed());
258 thread::sleep(Duration::from_secs(2));
259
260 // 3. Quiet baseline. The collector has done nothing and has nothing
261 // queued, so this is the cost of the walk alone.
262 let mut worst_quiet = Duration::ZERO;
263 for _ in 0..NUM_BASELINE {
264 let t = Instant::now();
265 let entries = res!(db.scan(&ScanOpts::all(), schms2));
266 let dt = t.elapsed();
267 if dt > worst_quiet {
268 worst_quiet = dt;
269 }
270 if entries.len() != shape.keys {
271 return Err(err!(
272 "Quiet scan returned {} entries, expected {}.",
273 entries.len(), shape.keys;
274 Test, Mismatch));
275 }
276 }
277 test!(sync_log::stream(), "Worst quiet scan latency {:?}.", worst_quiet);
278
279 // 4. Switch collection on and dispatch the backlog. Each of these
280 // writes supersedes a record in a different file, and the file
281 // bot answers each by handing one collection to the zone's
282 // collector. They return as soon as the message is sent, so by
283 // the end of the burst the collector's queue holds the lot.
284 res!(ok!(db.updated_api()).activate_gc(true));
285 test!(sync_log::stream(), "Collection enabled; dispatching the backlog.");
286 let stride = trigger_stride(shape);
287 let t = Instant::now();
288 let mut triggers = 0;
289 let mut i = 1;
290 while i < shape.keys {
291 res!(db.insert(dat!(fmt!("rec:{:06}", i)), payload(i, 2), user, schms2));
292 triggers += 1;
293 i += stride;
294 }
295 test!(sync_log::stream(), "{} triggers dispatched in {:?}.", triggers, t.elapsed());
296
297 // 5. Scan while the collector works through it. This is an operator
298 // refreshing a view.
299 let mut worst_loaded = Duration::ZERO;
300 let mut failures: Vec<String> = Vec::new();
301 for _ in 0..NUM_LOADED {
302 let t = Instant::now();
303 match db.scan(&ScanOpts::all(), schms2) {
304 Err(e) => {
305 let dt = t.elapsed();
306 if dt > worst_loaded {
307 worst_loaded = dt;
308 }
309 failures.push(fmt!("scan failed after {:?}: {}", dt, e));
310 },
311 Ok(entries) => {
312 let dt = t.elapsed();
313 if dt > worst_loaded {
314 worst_loaded = dt;
315 }
316 if entries.len() != shape.keys {
317 failures.push(fmt!(
318 "scan during collection returned {} entries, expected {}",
319 entries.len(), shape.keys));
320 }
321 },
322 }
323 }
324 test!(sync_log::stream(), "Worst scan latency during collection {:?}, {} misbehaved.",
325 worst_loaded, failures.len());
326 for f in failures.iter().take(5) {
327 test!(sync_log::stream(), " {}", f);
328 }
329
330 thread::sleep(Duration::from_secs(5));
331 res!(db.shutdown());
332
333 // 6. Judgement.
334 if !failures.is_empty() {
335 return Err(err!(
336 "{} scan(s) issued during garbage collection did not behave: first was '{}'.",
337 failures.len(), failures[0];
338 Test, Mismatch));
339 }
340 // A scan served independently of the collector costs what the walk
341 // costs. The allowance is loose enough for the disk contention a
342 // running collection genuinely causes, and far tighter than a
343 // backlog of collections.
344 let allowance = worst_quiet * 3 + Duration::from_millis(200);
345 if worst_loaded > allowance {
346 return Err(err!(
347 "Worst scan latency during collection was {:?}, more than the allowance \
348 of {:?} derived from the worst quiet latency of {:?}. The scan is \
349 waiting on the collector, and on a larger store that wait passes the \
350 user request deadline.",
351 worst_loaded, allowance, worst_quiet;
352 Test, Excessive));
353 }
354
355 test!(sync_log::stream(),
356 "Scan under garbage collection test passed: worst during collection {:?}, \
357 worst quiet {:?}.", worst_loaded, worst_quiet);
358 Ok(())
359}