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
| 1 | use 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 | |
| 17 | use oxedyne_fe2o3_core::byte::FromBytes; |
| 18 | use oxedyne_fe2o3_iop_db::api::{ |
| 19 | Meta, |
| 20 | ScanOpts, |
| 21 | }; |
| 22 | use oxedyne_fe2o3_jdat::{ |
| 23 | Dat, |
| 24 | id::NumIdDat, |
| 25 | }; |
| 26 | |
| 27 | use std::{ |
| 28 | collections::{ |
| 29 | BTreeMap, |
| 30 | HashMap, |
| 31 | }, |
| 32 | fs, |
| 33 | io::BufReader, |
| 34 | sync::Arc, |
| 35 | thread, |
| 36 | time::Duration, |
| 37 | }; |
| 38 | |
| 39 | const SCAN_COVERAGE_ATTEMPTS: usize = 3; |
| 40 | |
| 41 | const SCAN_COVERAGE_SETTLE: Duration = Duration::from_millis(20); |
| 42 | |
| 43 | #[derive(Clone, Copy, Debug, Default)] |
| 44 | struct ZoneFilePresence { |
| 45 | dat: bool, |
| 46 | ind: bool, |
| 47 | } |
| 48 | |
| 49 | #[derive(Clone, Copy, Debug)] |
| 50 | struct Shortfall { |
| 51 | fnum: FileNum, |
| 52 | dat_len: u64, |
| 53 | covered: u64, |
| 54 | } |
| 55 | |
| 56 | pub 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 | |
| 81 | impl< |
| 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 | |
| 94 | impl< |
| 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 | |
| 107 | impl< |
| 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 | |
| 159 | impl< |
| 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 | |
| 439 | fn 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 | } |