Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/worker/bot_scan.rs

14.5 KiB, 60 runs

created by r1870400018:21736, 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

1use crate::{
2 prelude::*,
3 bots::{
4 base::bot_deps::*,
5 worker::worker_deps::*,
6 },
7 file::{
8 core::FileAccess,
9 floc::FileNum,
10 stored::{
11 StoredIndex,
12 StoredKey,
13 },
14 },
15};
16
17use oxedyne_fe2o3_core::byte::FromBytes;
18use oxedyne_fe2o3_iop_db::api::{
19 Meta,
20 ScanOpts,
21};
22use oxedyne_fe2o3_jdat::{
23 Dat,
24 id::NumIdDat,
25};
26
27use std::{
28 collections::{
29 BTreeMap,
30 HashMap,
31 },
32 fs,
33 io::BufReader,
34 sync::Arc,
35 thread,
36 time::Duration,
37};
38
39const SCAN_COVERAGE_ATTEMPTS: usize = 3;
40
41const SCAN_COVERAGE_SETTLE: Duration = Duration::from_millis(20);
42
43#[derive(Clone, Copy, Debug, Default)]
44struct ZoneFilePresence {
45 dat: bool,
46 ind: bool,
47}
48
49#[derive(Clone, Copy, Debug)]
50struct Shortfall {
51 fnum: FileNum,
52 dat_len: u64,
53 covered: u64,
54}
55
56pub struct ScanBot<
57 const UIDL: usize,
58 UID: NumIdDat<UIDL>,
59 ENC: Encrypter,
60 KH: Hasher,
61 PR: Hasher,
62 CS: Checksummer,
63>{
64 // Identity
65 wind: WorkerInd,
66 wtyp: WorkerType,
67 // Bot
68 sem: Semaphore,
69 errc: Arc<Mutex<usize>>,
70 log_stream_id: String,
71 // Config
72 zdir: ZoneDir,
73 // Comms
74 chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>,
75 // API
76 api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>,
77 // State
78 inited: bool,
79}
80
81impl<
82 const UIDL: usize,
83 UID: NumIdDat<UIDL> + 'static,
84 ENC: Encrypter + 'static,
85 KH: Hasher + 'static,
86 PR: Hasher,
87 CS: Checksummer,
88>
89 WorkerBot<UIDL, UID, ENC, KH, PR, CS> for ScanBot<UIDL, UID, ENC, KH, PR, CS>
90{
91 workerbot_methods!();
92}
93
94impl<
95 const UIDL: usize,
96 UID: NumIdDat<UIDL> + 'static,
97 ENC: Encrypter + 'static,
98 KH: Hasher + 'static,
99 PR: Hasher,
100 CS: Checksummer,
101>
102 OzoneBot<UIDL, UID, ENC, KH, PR, CS> for ScanBot<UIDL, UID, ENC, KH, PR, CS>
103{
104 ozonebot_methods!();
105}
106
107impl<
108 const UIDL: usize,
109 UID: NumIdDat<UIDL> + 'static,
110 ENC: Encrypter + 'static,
111 KH: Hasher + 'static,
112 PR: Hasher,
113 CS: Checksummer,
114>
115 Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for ScanBot<UIDL, UID, ENC, KH, PR, CS>
116{
117 bot_methods!();
118
119 fn go(&mut self) {
120
121 sync_log::set_stream(self.log_stream_id());
122
123 if self.no_init() { return; }
124 self.now_listening();
125 loop {
126 if self.listen().must_end() { break; }
127 }
128 }
129
130 fn listen(&mut self) -> LoopBreak {
131 match self.chan_in().recv() {
132 Err(e) => self.err_cannot_receive(err!(e,
133 "{}: Waiting for message.", self.ozid();
134 IO, Channel)),
135 Ok(msg) => {
136 if let Some(msg) = self.listen_worker(msg) {
137 match msg {
138 OzoneMsg::ScanRequest {
139 opts,
140 schms2: _,
141 resp,
142 } => {
143 let result = self.scan_zone(&opts);
144 let msg = match result {
145 Ok(entries) => OzoneMsg::ScanEntries(entries),
146 Err(e) => OzoneMsg::Error(e),
147 };
148 self.respond(Ok(msg), &resp);
149 },
150 _ => return self.listen_more(msg),
151 }
152 }
153 },
154 }
155 LoopBreak(false)
156 }
157}
158
159impl<
160 const UIDL: usize,
161 UID: NumIdDat<UIDL> + 'static,
162 ENC: Encrypter + 'static,
163 KH: Hasher + 'static,
164 PR: Hasher,
165 CS: Checksummer,
166>
167 ScanBot<UIDL, UID, ENC, KH, PR, CS>
168{
169 pub fn new(
170 args: ZoneWorkerInitArgs<UIDL, UID, ENC, KH, PR, CS>,
171 )
172 -> Self
173 {
174 Self {
175 // Identity
176 wind: args.wind,
177 wtyp: args.wtyp,
178 // Bot
179 sem: args.sem,
180 errc: Arc::new(Mutex::new(0)),
181 log_stream_id: args.log_stream_id,
182 // Config
183 zdir: ZoneDir::default(),
184 // Comms
185 chan_in: args.chan_in,
186 // API
187 api: args.api,
188 // State
189 inited: false,
190 }
191 }
192
193 fn scan_zone(
194 &mut self,
195 opts: &ScanOpts,
196 )
197 -> Outcome<Vec<(Dat, Dat, Meta<UIDL, UID>)>>
198 {
199 // The chunk index of each record is carried alongside its key and meta so the
200 // main-key / chunk-data classification is made on the newest record per key, not on
201 // whichever happened to be walked last: a chunk-data record later superseded by a
202 // Complete tombstone must classify as the tombstone, or the inverse scan would emit a
203 // chunk key a delete has already begun reclaiming.
204 let mut live: HashMap<Vec<u8>, (Dat, Meta<UIDL, UID>, Option<usize>)>
205 = HashMap::new();
206 let mut short: Vec<Shortfall> = Vec::new();
207
208 // A pass that comes up short is repeated from scratch rather
209 // than patched, because the deduplication depends on files
210 // being visited in ascending order: re-walking one file after
211 // the others would let an older record overwrite a newer one.
212 for attempt in 0..SCAN_COVERAGE_ATTEMPTS {
213 live.clear();
214 short = res!(self.scan_pass(&mut live));
215 if short.is_empty() {
216 break;
217 }
218 if attempt + 1 < SCAN_COVERAGE_ATTEMPTS {
219 thread::sleep(SCAN_COVERAGE_SETTLE);
220 }
221 }
222
223 if !short.is_empty() {
224 let mut detail = String::new();
225 for s in &short {
226 detail.push_str(&fmt!(
227 " file {} holds {} bytes of records and its index accounts for {};",
228 s.fnum, s.dat_len, s.covered));
229 }
230 return Err(err!(
231 "{}: Scan of this zone is short: {} index file(s) do not account \
232 for what their data files hold, and there is no way to say how \
233 many records are missing, so nothing is reported rather than an \
234 answer that would look complete.{} The records themselves are \
235 intact and readable by key throughout; an index file that does \
236 not account for its data file is rebuilt from that data file on \
237 the next start.",
238 self.ozid(), short.len(), detail;
239 Data, Mismatch, Missing));
240 }
241
242 let mut out: Vec<(Dat, Dat, Meta<UIDL, UID>)> =
243 Vec::with_capacity(live.len());
244 for (_kbyts, (kdat, meta, cind)) in live.into_iter() {
245 // A chunk-data record has a chunk index >= 1. The default scan keeps the main user
246 // keys (Complete, and the bunch key at index 0) and elides those; `chunk_data_only`
247 // inverts it, keeping only the chunk-data keys.
248 let is_chunk_data = matches!(cind, Some(c) if c >= 1);
249 if opts.chunk_data_only != is_chunk_data {
250 continue;
251 }
252 if !scan_matches_prefix(&kdat, opts.prefix.as_ref()) {
253 continue;
254 }
255 out.push((kdat, Dat::Empty, meta));
256 if let Some(lim) = opts.limit {
257 if out.len() >= lim {
258 break;
259 }
260 }
261 }
262 if opts.include_values {
263 warn!(sync_log::stream(),
264 "{}: scan called with include_values=true; scan v1 \
265 returns Dat::Empty for every value. Fetch individual \
266 values via get() once a key is selected.",
267 self.ozid());
268 }
269 Ok(out)
270 }
271
272 fn scan_pass(
273 &mut self,
274 live: &mut HashMap<Vec<u8>, (Dat, Meta<UIDL, UID>, Option<usize>)>,
275 )
276 -> Outcome<Vec<Shortfall>>
277 {
278 let files = res!(self.list_zone_files());
279 let mut short = Vec::new();
280 for (fnum, present) in files {
281 let before = self.data_file_len(fnum);
282 // Every index file present is walked, including one whose data
283 // file has gone: a collection that deletes a wholly superseded
284 // file removes the data file first, and every key that file held
285 // has a newer copy in a higher-numbered file which overwrites it
286 // here anyway.
287 let covered = if present.ind {
288 res!(self.scan_walk_ind_file(fnum, live))
289 } else {
290 0
291 };
292 let after = self.data_file_len(fnum);
293 // A data file absent on either side has been collected away, or
294 // was never there: there is nothing for the index to be short of.
295 let dat_len = match (present.dat, before, after) {
296 (true, Some(a), Some(b)) => std::cmp::min(a, b),
297 _ => continue,
298 };
299 if covered < dat_len {
300 short.push(Shortfall { fnum, dat_len, covered });
301 }
302 }
303 Ok(short)
304 }
305
306 fn data_file_len(&self, fnum: FileNum) -> Option<u64> {
307 let mut path = self.zdir().dir.clone();
308 path.push(ZoneDir::relative_file_path(&FileType::Data, fnum));
309 match fs::metadata(&path) {
310 Ok(m) => Some(m.len()),
311 Err(_) => None,
312 }
313 }
314
315 fn list_zone_files(&self) -> Outcome<BTreeMap<FileNum, ZoneFilePresence>> {
316 let mut files: BTreeMap<FileNum, ZoneFilePresence> = BTreeMap::new();
317 for entry in res!(fs::read_dir(&self.zdir().dir)) {
318 let entry = res!(entry);
319 let path = entry.path();
320 if !path.is_file() {
321 continue;
322 }
323 if ZoneDir::is_gc_temp_file(&path) {
324 continue;
325 }
326 let (fnum, ftyp) = match ZoneDir::ozone_file_number_and_type(&path) {
327 Ok(t) => t,
328 Err(_) => continue,
329 };
330 let present = files.entry(fnum).or_insert_with(ZoneFilePresence::default);
331 match ftyp {
332 FileType::Data => present.dat = true,
333 FileType::Index => present.ind = true,
334 }
335 }
336 Ok(files)
337 }
338
339 fn scan_walk_ind_file(
340 &mut self,
341 fnum: FileNum,
342 live: &mut HashMap<Vec<u8>, (Dat, Meta<UIDL, UID>, Option<usize>)>,
343 )
344 -> Outcome<u64>
345 {
346 let file = match self.zdir().open_ozone_file(
347 fnum,
348 &FileType::Index,
349 &FileAccess::Reading,
350 ) {
351 Ok((_, file)) => file,
352 Err(_) => {
353 trace!(sync_log::stream(),
354 "{}: Index file {} went away between listing and opening, \
355 which is what collecting an entirely superseded file looks \
356 like; skipping it.", self.ozid(), fnum);
357 return Ok(0);
358 },
359 };
360 let mut reader = BufReader::new(file);
361 let typ = FileType::Index;
362 let mut pos = 0usize;
363 // Bytes of key-value data this index accounts for.
364 let mut covered = 0u64;
365
366 loop {
367 // 1. Load the StoredKey from the file.
368 let (key, meta) = match StoredKey::load(
369 &mut reader,
370 self.api().schemes().checksummer().clone(),
371 ) {
372 Err(e) => return Err(err!(e,
373 "{}: While scanning {:?} file {} at position {}.",
374 self.ozid(), typ, fnum, pos;
375 IO, File, Read)),
376 Ok(None) => break,
377 Ok(Some((skey, _, n))) => {
378 pos += n;
379 let meta = skey.meta().clone();
380 (skey.into_key(), meta)
381 },
382 };
383 // 2. Skip the matching StoredIndex. We do not need the
384 // location -- we are not reading values in v1.
385 match StoredIndex::read(
386 &mut reader,
387 fnum,
388 self.api().schemes().checksummer().clone(),
389 ) {
390 Err(e) => return Err(err!(e,
391 "{}: While scanning stored index in {:?} file {} \
392 at position {}.",
393 self.ozid(), typ, fnum, pos;
394 IO, File, Read)),
395 Ok((None, _)) => return Err(err!(
396 "{}: Missing StoredIndex at end of {:?} file {}.",
397 self.ozid(), typ, fnum;
398 Missing)),
399 Ok((Some(sindex), n)) => {
400 pos += n;
401 // Counted before the chunk entries are elided below:
402 // the data file holds those records too, so leaving
403 // them out here would make every chunked value look
404 // like an under-count.
405 covered += sindex.keyval_len();
406 },
407 }
408
409 // 3. Record the chunk index so the caller can classify on the newest record.
410 // `Complete` keys have no index, a bunch key is index 0, and a chunk-data
411 // record is index >= 1. The main-key / chunk-data split is applied in
412 // `scan_zone` after the newest-wins dedup below, not here, so a chunk-data
413 // record superseded by a later Complete tombstone classifies as the tombstone.
414 let cind = key.index();
415
416 // 4. Decode the raw key bytes to a Dat.
417 let kbyts = key.into_bytes();
418 let (kdat, _n_decoded) = match Dat::from_bytes(&kbyts) {
419 Ok(pair) => pair,
420 Err(e) => {
421 warn!(sync_log::stream(),
422 "{}: Could not decode scanned key bytes in file {} \
423 at position {}: {}. Skipping entry.",
424 self.ozid(), fnum, pos, e);
425 continue;
426 },
427 };
428
429 // 5. Insert into the live map. Later occurrences of the
430 // same raw key bytes (from higher fnum or later in
431 // the same file) overwrite, which is exactly the
432 // stale-filtering behaviour we want.
433 live.insert(kbyts, (kdat, meta, cind));
434 }
435 Ok(covered)
436 }
437}
438
439fn scan_matches_prefix(kdat: &Dat, prefix: Option<&Dat>) -> bool {
440 match prefix {
441 None => true,
442 Some(Dat::Str(p)) => match kdat {
443 Dat::Str(s) => s.starts_with(p.as_str()),
444 _ => false,
445 },
446 Some(other) => kdat == other,
447 }
448}