oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/worker/bot_cache.rs
14.3 KiB, 69 runs
created by r1870400018:751, 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::{ |
| 6 | bot_reader::ReadResult, |
| 7 | worker_deps::*, |
| 8 | }, |
| 9 | }, |
| 10 | data::{ |
| 11 | cache::{ |
| 12 | Cache, |
| 13 | ValueOrLocation, |
| 14 | }, |
| 15 | core::Key, |
| 16 | }, |
| 17 | file::{ |
| 18 | floc::FileLocation, |
| 19 | stored::RecordDigest, |
| 20 | }, |
| 21 | test::hooks, |
| 22 | }; |
| 23 | |
| 24 | use oxedyne_fe2o3_core::channels::Recv; |
| 25 | use oxedyne_fe2o3_iop_db::api::Meta; |
| 26 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 27 | |
| 28 | use std::{ |
| 29 | sync::Arc, |
| 30 | time::Instant, |
| 31 | }; |
| 32 | |
| 33 | #[derive(Debug)] |
| 34 | pub struct CacheBot< |
| 35 | const UIDL: usize, |
| 36 | UID: NumIdDat<UIDL>, |
| 37 | ENC: Encrypter, |
| 38 | KH: Hasher, |
| 39 | PR: Hasher, |
| 40 | CS: Checksummer, |
| 41 | >{ |
| 42 | // Identity |
| 43 | wind: WorkerInd, |
| 44 | wtyp: WorkerType, |
| 45 | // Bot |
| 46 | sem: Semaphore, |
| 47 | errc: Arc<Mutex<usize>>, |
| 48 | log_stream_id: String, |
| 49 | // Config |
| 50 | zdir: ZoneDir, |
| 51 | // Comms |
| 52 | chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 53 | // API |
| 54 | api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 55 | // State |
| 56 | active: bool, |
| 57 | cache: Cache<UIDL, UID>, |
| 58 | inited: bool, |
| 59 | trep: Instant, |
| 60 | } |
| 61 | |
| 62 | impl< |
| 63 | const UIDL: usize, |
| 64 | UID: NumIdDat<UIDL> + 'static, |
| 65 | ENC: Encrypter + 'static, |
| 66 | KH: Hasher + 'static, |
| 67 | PR: Hasher, |
| 68 | CS: Checksummer, |
| 69 | > |
| 70 | WorkerBot<UIDL, UID, ENC, KH, PR, CS> for CacheBot<UIDL, UID, ENC, KH, PR, CS> |
| 71 | { |
| 72 | workerbot_methods!(); |
| 73 | } |
| 74 | |
| 75 | impl< |
| 76 | const UIDL: usize, |
| 77 | UID: NumIdDat<UIDL> + 'static, |
| 78 | ENC: Encrypter + 'static, |
| 79 | KH: Hasher + 'static, |
| 80 | PR: Hasher, |
| 81 | CS: Checksummer, |
| 82 | > |
| 83 | OzoneBot<UIDL, UID, ENC, KH, PR, CS> for CacheBot<UIDL, UID, ENC, KH, PR, CS> |
| 84 | { |
| 85 | ozonebot_methods!(); |
| 86 | } |
| 87 | |
| 88 | impl< |
| 89 | const UIDL: usize, |
| 90 | UID: NumIdDat<UIDL> + 'static, |
| 91 | ENC: Encrypter + 'static, |
| 92 | KH: Hasher + 'static, |
| 93 | PR: Hasher, |
| 94 | CS: Checksummer, |
| 95 | > |
| 96 | Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for CacheBot<UIDL, UID, ENC, KH, PR, CS> |
| 97 | { |
| 98 | bot_methods!(); |
| 99 | |
| 100 | fn go(&mut self) { |
| 101 | |
| 102 | sync_log::set_stream(self.log_stream_id()); |
| 103 | |
| 104 | if self.no_init() { return; } |
| 105 | self.now_listening(); |
| 106 | loop { |
| 107 | if self.wind().b() < self.cfg().num_bots_per_zone(self.wtyp()) { |
| 108 | |
| 109 | if self.trep.elapsed() > self.cfg().zone_state_update_interval() { |
| 110 | // Automated state reporting. |
| 111 | self.trep = Instant::now(); |
| 112 | if let Some(zbot) = self.zbot() { |
| 113 | if let Err(e) = zbot.send( |
| 114 | OzoneMsg::CacheSize( |
| 115 | self.wind().b(), |
| 116 | self.cache().get_size(), |
| 117 | self.cache().get_ancillary_size(), |
| 118 | ) |
| 119 | ) { |
| 120 | self.result(&Err(err!(e, |
| 121 | "{}: Cannot send cache size update to zbot.", self.ozid(); |
| 122 | Channel, Write))); |
| 123 | } |
| 124 | } |
| 125 | } |
| 126 | |
| 127 | if self.listen().must_end() { break; } |
| 128 | |
| 129 | } else { |
| 130 | // This bot is to be terminated. Forward incoming messages to the remaining bots of |
| 131 | // this type. |
| 132 | } |
| 133 | } |
| 134 | } |
| 135 | |
| 136 | fn listen(&mut self) -> LoopBreak { |
| 137 | match self.chan_in().recv_timeout(self.cfg().zone_state_update_interval()) { |
| 138 | Recv::Result(Err(e)) => self.err_cannot_receive(err!(e, |
| 139 | "{}: Waiting for message.", self.ozid(); |
| 140 | IO, Channel)), |
| 141 | Recv::Result(Ok(msg)) => { |
| 142 | if let Some(msg) = self.listen_worker(msg) { |
| 143 | match msg { |
| 144 | // COMMAND |
| 145 | OzoneMsg::ClearCache(resp) => { |
| 146 | self.cache_mut().clear_all_values(); |
| 147 | self.respond(Ok(OzoneMsg::Ok), &resp); |
| 148 | }, |
| 149 | OzoneMsg::SetCacheSizeLimit(size_lim) => { |
| 150 | self.cache_mut().set_lim(size_lim); |
| 151 | }, |
| 152 | // WRITE |
| 153 | OzoneMsg::GcCacheUpdateRequest(buf, resp_g1) => { |
| 154 | let mut old_flocs = Vec::new(); |
| 155 | for (key, floc, meta) in buf { |
| 156 | if let Some(old_floc) = self.cache_mut().reanchor(&key, &floc, &meta) { |
| 157 | match RecordDigest::new(&key, &meta) { |
| 158 | Ok(rid) => old_flocs.push((old_floc, rid)), |
| 159 | // Its move entry stays, and keeps the file from being |
| 160 | // collected again; the location itself is right. |
| 161 | Err(e) => self.error(err!(e, |
| 162 | "{}: Naming a record a collection re-anchored.", self.ozid(); |
| 163 | Data, Encode)), |
| 164 | } |
| 165 | } |
| 166 | } |
| 167 | if let Err(e) = resp_g1.send( |
| 168 | OzoneMsg::GcCacheUpdateResponse(old_flocs) |
| 169 | ) { |
| 170 | self.err_cannot_send(err!(e, |
| 171 | "{}: Sending cache update response back to gbot.", self.ozid(); |
| 172 | IO, Channel)); |
| 173 | } |
| 174 | }, |
| 175 | OzoneMsg::Insert(key, val, cind, floc, ilen, meta, resp_w1, unconfirmed) => { |
| 176 | let result = self.insert(key, val, cind, floc, ilen, meta, resp_w1, unconfirmed); |
| 177 | self.result(&result); |
| 178 | }, |
| 179 | // READ |
| 180 | OzoneMsg::DumpCacheRequest(resp) => { |
| 181 | if let Err(e) = resp.send(OzoneMsg::DumpCacheResponse( |
| 182 | self.wind().clone(), |
| 183 | self.cache().clone(), |
| 184 | )) { |
| 185 | self.err_cannot_send(err!(e, |
| 186 | "{}: Responding to {:?} with cache dump.", self.ozid(), resp.ozid(); |
| 187 | Data, IO, Channel)); |
| 188 | } |
| 189 | }, |
| 190 | //OzoneMsg::GetUsers(resp) => { |
| 191 | // let mut kuserdat = Vec::new(); |
| 192 | // let prefix = constant::USER_DAT.byte_prefix(); |
| 193 | // for (k, _) in self.cache().map() { |
| 194 | // if k.len() > UsrKindId::CODE_BYTE_LEN { |
| 195 | // if k[0] == Dat::USR_CODE { |
| 196 | // if UsrKindId::prefix_matches(prefix, &k[1..]) { |
| 197 | // let kdat = match Dat::from_bytes(&k) { |
| 198 | // Ok((dat, _)) => dat, |
| 199 | // _ => continue, |
| 200 | // }; |
| 201 | // if let Dat::Usr(_, optboxdat) = &kdat { |
| 202 | // match optboxdat { |
| 203 | // Some(boxdat) => match **boxdat { |
| 204 | // Dat::U128(id) => { |
| 205 | // kuserdat.push((id, kdat.clone())); |
| 206 | // }, |
| 207 | // dat => self.error(err!( |
| 208 | // "Custom Usr daticle should contain a \ |
| 209 | // Dat::U128 but found {:?}.", dat, |
| 210 | // ), Bug, Invalid, Input)), |
| 211 | // }, |
| 212 | // None => self.error(err!( |
| 213 | // "Custom Usr daticle should contain \ |
| 214 | // something but None was found.", |
| 215 | // ), Bug, Invalid, Input)), |
| 216 | // } |
| 217 | // } |
| 218 | // } |
| 219 | // } |
| 220 | // } |
| 221 | // } |
| 222 | // self.respond(Ok(OzoneMsg::UserKeys(kuserdat)), &resp); |
| 223 | //}, |
| 224 | OzoneMsg::ReadCache(key, resp_r2) => { |
| 225 | let result = self.read(key, resp_r2); |
| 226 | self.result(&result); |
| 227 | }, |
| 228 | _ => return self.listen_more(msg), |
| 229 | } |
| 230 | } |
| 231 | }, |
| 232 | Recv::Empty => (), |
| 233 | } |
| 234 | LoopBreak(false) |
| 235 | } |
| 236 | |
| 237 | } |
| 238 | |
| 239 | impl< |
| 240 | const UIDL: usize, |
| 241 | UID: NumIdDat<UIDL> + 'static, |
| 242 | ENC: Encrypter + 'static, |
| 243 | KH: Hasher + 'static, |
| 244 | PR: Hasher, |
| 245 | CS: Checksummer, |
| 246 | > |
| 247 | CacheBot<UIDL, UID, ENC, KH, PR, CS> |
| 248 | { |
| 249 | pub fn new( |
| 250 | args: ZoneWorkerInitArgs<UIDL, UID, ENC, KH, PR, CS>, |
| 251 | ) |
| 252 | -> Self |
| 253 | { |
| 254 | let cache = Cache::new(Some(&args.api.ozid)); |
| 255 | Self { |
| 256 | // Identity |
| 257 | wind: args.wind, |
| 258 | wtyp: args.wtyp, |
| 259 | // Bot |
| 260 | sem: args.sem, |
| 261 | errc: Arc::new(Mutex::new(0)), |
| 262 | log_stream_id: args.log_stream_id, |
| 263 | // Config |
| 264 | zdir: ZoneDir::default(), |
| 265 | // Comms |
| 266 | chan_in: args.chan_in, |
| 267 | // API |
| 268 | api: args.api, |
| 269 | // State |
| 270 | active: false, |
| 271 | cache, |
| 272 | inited: false, |
| 273 | trep: Instant::now(), |
| 274 | } |
| 275 | } |
| 276 | |
| 277 | fn cache(&self) -> &Cache<UIDL, UID> { &self.cache } |
| 278 | fn cache_mut(&mut self) -> &mut Cache<UIDL, UID> { &mut self.cache } |
| 279 | |
| 280 | pub fn max_file_len(&self) -> usize { |
| 281 | self.cfg().data_file_max_bytes as usize |
| 282 | } |
| 283 | |
| 284 | pub fn activate(mut self) -> Self { |
| 285 | self.active = true; |
| 286 | self |
| 287 | } |
| 288 | |
| 289 | pub fn insert( |
| 290 | &mut self, |
| 291 | key: Vec<u8>, |
| 292 | val: Option<Vec<u8>>, |
| 293 | cind: Option<usize>, |
| 294 | floc: FileLocation, |
| 295 | ilen: usize, |
| 296 | meta: Meta<UIDL, UID>, |
| 297 | resp_w1: Responder<UIDL, UID, ENC, KH>, |
| 298 | unconfirmed: Option<Error<ErrTag>>, |
| 299 | ) |
| 300 | -> Outcome<()> |
| 301 | { |
| 302 | hooks::insert_delay(); |
| 303 | // [12] Insert the data into the key-chosen zone cache. |
| 304 | let floc_new = floc.clone(); |
| 305 | let floc_old_opt = match self.cache.insert( |
| 306 | key, |
| 307 | val, |
| 308 | floc, |
| 309 | meta, |
| 310 | ) { |
| 311 | Ok(floc_old_opt) => floc_old_opt, |
| 312 | Err(e) => { |
| 313 | // The caller is waiting on this answer. Only logged, a failure here reached it |
| 314 | // as an expired durability deadline, which says the write is on its way. |
| 315 | let e = err!(e, |
| 316 | "{}: A written record could not be entered in the cache.", self.ozid(); |
| 317 | Data, Write, Unconfirmed); |
| 318 | self.respond(Err(e.clone()), &resp_w1); |
| 319 | return Err(e); |
| 320 | }, |
| 321 | }; |
| 322 | |
| 323 | let key_present = floc_old_opt.is_some(); |
| 324 | |
| 325 | // [13] Inform the caller of successful file write and cache insertion, or of the barrier |
| 326 | // that failed after the write: told here, the caller can read its write back once it |
| 327 | // hears, as it can a confirmed one. |
| 328 | match (unconfirmed, cind) { |
| 329 | (Some(e), _) => self.respond(Err(e), &resp_w1), |
| 330 | (None, Some(cind)) => self.respond(Ok(OzoneMsg::KeyChunkExists(key_present, cind)), &resp_w1), |
| 331 | (None, None) => self.respond(Ok(OzoneMsg::KeyExists(key_present)), &resp_w1), |
| 332 | } |
| 333 | self.respond(Ok(OzoneMsg::Finish), &resp_w1); |
| 334 | |
| 335 | // [14] Insert the new data into the file state data map, via a file-selected cbot. |
| 336 | let bots = res!(self.fbots()); |
| 337 | let (bot, _) = bots.choose_bot( |
| 338 | &ChooseBot::ByFile(floc_new.file_number()) |
| 339 | ); |
| 340 | res!(bot.send(OzoneMsg::UpdateData { |
| 341 | floc_new, |
| 342 | ilen, |
| 343 | floc_old_opt, |
| 344 | from_id: self.ozid().clone(), |
| 345 | })); |
| 346 | |
| 347 | Ok(()) |
| 348 | } |
| 349 | |
| 350 | pub fn read( |
| 351 | &mut self, |
| 352 | key: Key, |
| 353 | resp_r2: Responder<UIDL, UID, ENC, KH>, |
| 354 | ) |
| 355 | -> Outcome<()> |
| 356 | { |
| 357 | // <3> The cbot accesses its cache. |
| 358 | match res!(self.cache.get(key.as_bytes())) { |
| 359 | Some(vloc) => { |
| 360 | match vloc { |
| 361 | ValueOrLocation::Location(mloc) => { |
| 362 | // <4> Only the file location is available, so send a request to the |
| 363 | // appropriate fbot, forwarding the responder. |
| 364 | let fnum = mloc.file_number(); |
| 365 | let bots = res!(self.fbots()); |
| 366 | let (bot, _) = bots.choose_bot(&ChooseBot::ByFile(fnum)); |
| 367 | // The key goes with the location, so that a move entry at the same |
| 368 | // offset is taken only if it is this record's. |
| 369 | res!(bot.send( |
| 370 | OzoneMsg::ReadFileRequest( |
| 371 | fnum, |
| 372 | key.into_bytes(), |
| 373 | mloc, |
| 374 | resp_r2, |
| 375 | ))); |
| 376 | }, |
| 377 | // <6> Send result directly back to rbot. |
| 378 | ValueOrLocation::Value(val, meta) => |
| 379 | res!(resp_r2.send(OzoneMsg::ReadResult(ReadResult::Value(val, meta)))), |
| 380 | ValueOrLocation::Deleted(meta) => |
| 381 | res!(resp_r2.send(OzoneMsg::ReadResult(ReadResult::Deleted(meta)))), |
| 382 | } |
| 383 | |
| 384 | }, |
| 385 | // <6> Send result directly back to rbot. |
| 386 | None => res!(resp_r2.send(OzoneMsg::ReadResult(ReadResult::None))), |
| 387 | } |
| 388 | Ok(()) |
| 389 | } |
| 390 | } |