oxedyne/fe2o3/fe2o3_o3db_sync/src/data/cache.rs
17.7 KiB, 60 runs
created by r1870400018:779, 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 | file::floc::{ |
| 4 | FileLocation, |
| 5 | FileNum, |
| 6 | }, |
| 7 | base::id::OzoneBotId, |
| 8 | data::core::Key, |
| 9 | file::stored::RecordDigest, |
| 10 | }; |
| 11 | |
| 12 | use oxedyne_fe2o3_data::time::Timestamp; |
| 13 | use oxedyne_fe2o3_iop_db::api::Meta; |
| 14 | use oxedyne_fe2o3_jdat::{ |
| 15 | daticle::Dat, |
| 16 | id::NumIdDat, |
| 17 | }; |
| 18 | |
| 19 | use std::{ |
| 20 | collections::BTreeMap, |
| 21 | fmt, |
| 22 | marker::PhantomData, |
| 23 | }; |
| 24 | |
| 25 | #[derive(Clone, Debug, Default, Eq, PartialEq)] |
| 26 | pub struct CacheId(pub u16); |
| 27 | |
| 28 | impl fmt::Display for CacheId { |
| 29 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 30 | write!(f, "CacheId({})", self.0) |
| 31 | } |
| 32 | } |
| 33 | |
| 34 | impl CacheId { |
| 35 | pub fn new(c: u16) -> Self { |
| 36 | Self(c) |
| 37 | } |
| 38 | } |
| 39 | |
| 40 | /// Contains the value itself, or its location. Used for cache retrieval. |
| 41 | #[derive(Clone, Debug)] |
| 42 | pub enum ValueOrLocation< |
| 43 | const UIDL: usize, |
| 44 | UID: NumIdDat<UIDL>, |
| 45 | > { |
| 46 | Deleted(Meta<UIDL, UID>), |
| 47 | Value(Vec<u8>, Meta<UIDL, UID>), |
| 48 | Location(MetaLocation<UIDL, UID>), |
| 49 | } |
| 50 | |
| 51 | /// Used for cache storage. |
| 52 | #[derive(Clone, Debug)] |
| 53 | pub struct MetaLocation< |
| 54 | const UIDL: usize, |
| 55 | UID: NumIdDat<UIDL>, |
| 56 | > { |
| 57 | meta: Meta<UIDL, UID>, |
| 58 | floc: FileLocation, |
| 59 | } |
| 60 | |
| 61 | impl< |
| 62 | const UIDL: usize, |
| 63 | UID: NumIdDat<UIDL>, |
| 64 | > |
| 65 | MetaLocation<UIDL, UID> |
| 66 | { |
| 67 | pub fn meta(&self) -> &Meta<UIDL, UID> { &self.meta } |
| 68 | pub fn meta_move(self) -> Meta<UIDL, UID> { self.meta } |
| 69 | pub fn file_location(&self) -> &FileLocation { &self.floc } |
| 70 | pub fn file_number(&self) -> FileNum { self.floc.file_number() } |
| 71 | |
| 72 | pub fn new_start_position(&mut self, new_start: u64) { |
| 73 | self.floc.start = new_start |
| 74 | } |
| 75 | } |
| 76 | |
| 77 | #[derive(Clone, Debug)] |
| 78 | struct CacheSizes< |
| 79 | const UIDL: usize, |
| 80 | UID: NumIdDat<UIDL>, |
| 81 | > { |
| 82 | byte: usize, |
| 83 | meta: usize, |
| 84 | floc: usize, |
| 85 | mloc: usize, |
| 86 | phantom: PhantomData<UID>, |
| 87 | } |
| 88 | |
| 89 | impl< |
| 90 | const UIDL: usize, |
| 91 | UID: NumIdDat<UIDL>, |
| 92 | > |
| 93 | Default for CacheSizes<UIDL, UID> |
| 94 | { |
| 95 | fn default() -> Self { |
| 96 | Self { |
| 97 | byte: std::mem::size_of::<u8>(), |
| 98 | meta: UIDL, |
| 99 | floc: std::mem::size_of::<FileLocation>(), |
| 100 | mloc: std::mem::size_of::<MetaLocation<UIDL, UID>>(), |
| 101 | phantom: PhantomData, |
| 102 | } |
| 103 | } |
| 104 | } |
| 105 | |
| 106 | #[derive(Clone, Debug)] |
| 107 | pub struct KeyVal< |
| 108 | const UIDL: usize, |
| 109 | UID: NumIdDat<UIDL>, |
| 110 | > { |
| 111 | pub key: Key, |
| 112 | pub val: Vec<u8>, |
| 113 | pub chash: alias::ChooseHash, |
| 114 | pub meta: Meta<UIDL, UID>, |
| 115 | pub cbpind: usize |
| 116 | } |
| 117 | |
| 118 | impl< |
| 119 | const UIDL: usize, |
| 120 | UID: NumIdDat<UIDL>, |
| 121 | > |
| 122 | KeyVal<UIDL, UID> |
| 123 | { |
| 124 | pub fn stamp_time_now(&mut self) -> Outcome<()> { |
| 125 | self.meta.stamp_time_now() |
| 126 | } |
| 127 | } |
| 128 | |
| 129 | #[derive(Clone, Debug)] |
| 130 | pub enum CacheEntry< |
| 131 | const UIDL: usize, |
| 132 | UID: NumIdDat<UIDL>, |
| 133 | > { |
| 134 | LocatedValue(MetaLocation<UIDL, UID>, Option<Vec<u8>>), |
| 135 | Deleted(Meta<UIDL, UID>), |
| 136 | } |
| 137 | |
| 138 | /// A central goal of Ozone is to hold as much data as possible in volatile memory in zone caches. |
| 139 | #[derive(Clone, Debug, Default)] |
| 140 | pub struct Cache< |
| 141 | const UIDL: usize, |
| 142 | UID: NumIdDat<UIDL>, |
| 143 | > { |
| 144 | ozid: Option<OzoneBotId>, // Creator |
| 145 | map: BTreeMap<Vec<u8>, CacheEntry<UIDL, UID>>, |
| 146 | size: usize, // estimate of bytes stored in encoded form |
| 147 | lim: usize, // limit on size in [MB] |
| 148 | cwt: CacheWriteTracker, |
| 149 | csizes: CacheSizes<UIDL, UID>, |
| 150 | } |
| 151 | |
| 152 | impl< |
| 153 | const UIDL: usize, |
| 154 | UID: NumIdDat<UIDL>, |
| 155 | > |
| 156 | Cache<UIDL, UID> |
| 157 | { |
| 158 | const MLOC_SIZE: usize = std::mem::size_of::<MetaLocation<UIDL, UID>>(); |
| 159 | |
| 160 | pub fn new(ozid: Option<&OzoneBotId>) -> Self { |
| 161 | Self { |
| 162 | ozid: ozid.map(|id| id.clone()), |
| 163 | ..Default::default() |
| 164 | } |
| 165 | } |
| 166 | |
| 167 | fn ozid(&self) -> &Option<OzoneBotId> { &self.ozid } |
| 168 | |
| 169 | /// Getter for cache size in bytes. |
| 170 | pub fn get_size(&self) -> usize { self.size } |
| 171 | /// Getter for ancillary data structures size in bytes. |
| 172 | pub fn get_ancillary_size(&self) -> usize { self.cwt.size } |
| 173 | /// Getter for cache size limit in bytes. |
| 174 | pub fn get_lim(&self) -> usize { self.lim } |
| 175 | /// Getter for a reference to the cache map. |
| 176 | pub fn map(&self) -> &BTreeMap<Vec<u8>, CacheEntry<UIDL, UID>> { &self.map } |
| 177 | pub fn mloc_size(&self) -> usize { |
| 178 | self.csizes.mloc |
| 179 | } |
| 180 | |
| 181 | /// The cache size in mebibytes, one of which is 1024^2 bytes. |
| 182 | pub fn size_mb(&self) -> f64 { (self.size as f64) / 1_048_576.0 } |
| 183 | |
| 184 | pub fn lim_size_mb(&self) -> f64 { (self.lim as f64) / 1_048_576.0 } |
| 185 | |
| 186 | /// Returns the fraction of the cache size compared to the size limit. |
| 187 | pub fn size_fraction(&self) -> f64 { self.size_mb() / (self.lim as f64) } |
| 188 | |
| 189 | pub fn set_lim(&mut self, lim: usize) { |
| 190 | self.lim = lim; |
| 191 | } |
| 192 | |
| 193 | /// Insert key, value and location into the cache. |
| 194 | /// |
| 195 | /// Returns the location of whichever copy of the key is now superseded, and the record it |
| 196 | /// holds, for the caller to schedule for garbage collection, or `None` when nothing was |
| 197 | /// superseded. That is usually the copy the cache held, but when the offered copy is the |
| 198 | /// older of the two the cache keeps what it has and the offered location comes back instead. |
| 199 | pub fn insert( |
| 200 | &mut self, |
| 201 | kbyts: Vec<u8>, |
| 202 | val: Option<Vec<u8>>, |
| 203 | floc: FileLocation, |
| 204 | meta: Meta<UIDL, UID>, |
| 205 | ) |
| 206 | -> Outcome<Option<(FileLocation, RecordDigest)>> |
| 207 | { |
| 208 | let klen = kbyts.len(); |
| 209 | // 1. Make space in the cache if we are going to exceed the size limit. We could just |
| 210 | // jettison enough in order to fit the new value into the limit. But subsequently, this |
| 211 | // expensive jettison process would be activated more frequently as the cache continues |
| 212 | // to bump up against its limit. A compromise is to chop cache value storage back by a |
| 213 | // solid amount (say 20%), which would improve performance at the expense of older |
| 214 | // values. A bit like the difference between mowing the lawn every day, or just |
| 215 | // every few weeks. |
| 216 | if let Some(v) = &val { |
| 217 | let vlen = v.len(); |
| 218 | if self.size + vlen > self.lim { |
| 219 | let start_size = self.size; |
| 220 | let desired_cache_size = |
| 221 | ((1.0 - constant::CACHE_JETTISON_FRAC_OF_LIM) * (self.lim as f64)) as usize; |
| 222 | let keys = res!(self.cwt.jettison(self.size + vlen - desired_cache_size)); |
| 223 | let mut saved = 0; |
| 224 | for key in &keys { |
| 225 | match self.map.get_mut(key) { |
| 226 | Some(CacheEntry::LocatedValue(_, val2_opt)) => { |
| 227 | if let Some(val2) = val2_opt { |
| 228 | saved += val2.len(); |
| 229 | self.size = try_sub!(&self.size, val2.len()); |
| 230 | *val2_opt = None; |
| 231 | } |
| 232 | }, |
| 233 | _ => (), |
| 234 | } |
| 235 | } |
| 236 | trace!(sync_log::stream(), |
| 237 | "{:?}: Automatically jettisoned the oldest {} (~{:.1}%) cache values \ |
| 238 | to reduce size from {} to {} bytes.", |
| 239 | self.ozid(), keys.len(), |
| 240 | constant::CACHE_JETTISON_FRAC_OF_LIM * 100.0, |
| 241 | start_size, start_size - saved, |
| 242 | ); |
| 243 | } |
| 244 | } |
| 245 | // 2. See if the key already exists. |
| 246 | match self.map.get_mut(&kbyts) { |
| 247 | Some(CacheEntry::LocatedValue(mloc, val2)) => { |
| 248 | // 2.1 Only insert if the given data is newer. The cache keeps the newer |
| 249 | // copy, which makes the copy just offered the superseded one, so its |
| 250 | // location is what goes back for flagging as old. Returning nothing |
| 251 | // here would leave that copy marked current in its file state forever, |
| 252 | // and its bytes would never be reclaimable. A copy at the location |
| 253 | // already cached is the same record arriving twice, not a supersession, |
| 254 | // and must be left alone. |
| 255 | if meta.time <= mloc.meta.time { |
| 256 | if floc == *mloc.file_location() { |
| 257 | return Ok(None); |
| 258 | } |
| 259 | trace!(sync_log::stream(), |
| 260 | "{:?}: The value offered for key = {:?} at {:?} is stamped {:?}, \ |
| 261 | no newer than the cached {:?}, so the offered copy is superseded.", |
| 262 | self.ozid.clone(), kbyts, floc, meta.time, mloc.meta.time, |
| 263 | ); |
| 264 | let rid = res!(RecordDigest::new(&kbyts, &meta)); |
| 265 | return Ok(Some((floc, rid))); |
| 266 | } |
| 267 | // 2.2 It does, insert the new info and return the old floc. |
| 268 | let new_mloc = MetaLocation { |
| 269 | meta: meta.clone(), |
| 270 | floc, |
| 271 | }; |
| 272 | let old = (mloc.file_location().clone(), res!(RecordDigest::new(&kbyts, mloc.meta()))); |
| 273 | *mloc = new_mloc; |
| 274 | match val { |
| 275 | Some(v) => { |
| 276 | let vlen = v.len(); |
| 277 | match &val2 { |
| 278 | Some(v2) => self.size = try_sub!(&self.size, res!(Self::valsize(v2.len()))), |
| 279 | None => (), |
| 280 | } |
| 281 | res!(self.cwt.insert(res!(Timestamp::now()), &kbyts, vlen)); |
| 282 | self.size = try_add!(&self.size, res!(Self::valsize(vlen))); |
| 283 | *val2 = Some(v); |
| 284 | }, |
| 285 | None => (), // leave any existing value untouched |
| 286 | } |
| 287 | Ok(Some(old)) |
| 288 | }, |
| 289 | None | |
| 290 | Some(CacheEntry::Deleted(_)) => { |
| 291 | // 2.3 It doesn't exist or was deleted, so create the entry and insert. |
| 292 | match &val { |
| 293 | Some(v) => { |
| 294 | let vlen = v.len(); |
| 295 | res!(self.cwt.insert(res!(Timestamp::now()), &kbyts, vlen)); |
| 296 | self.size = try_add!(&self.size, res!(Self::valsize(vlen))); |
| 297 | }, |
| 298 | None => (), |
| 299 | } |
| 300 | let mloc = MetaLocation { |
| 301 | meta, |
| 302 | floc, |
| 303 | }; |
| 304 | self.map.insert(kbyts, CacheEntry::LocatedValue(mloc, val)); |
| 305 | self.size = try_add!(&self.size, klen); |
| 306 | Ok(None) |
| 307 | }, |
| 308 | } |
| 309 | } |
| 310 | |
| 311 | fn valsize(len: usize) -> Outcome<usize> { |
| 312 | Ok(try_add!(&Self::MLOC_SIZE, len)) |
| 313 | } |
| 314 | |
| 315 | /// Update file location information for a key. |
| 316 | pub fn update( |
| 317 | &mut self, |
| 318 | k: &Vec<u8>, |
| 319 | floc: FileLocation, |
| 320 | meta: Meta<UIDL, UID>, |
| 321 | ) |
| 322 | -> Outcome<()> |
| 323 | { |
| 324 | match self.map.get_mut(k) { |
| 325 | Some(CacheEntry::LocatedValue(mloc, _)) => { |
| 326 | *mloc = MetaLocation { meta, floc }; |
| 327 | Ok(()) |
| 328 | }, |
| 329 | Some(CacheEntry::Deleted(_)) => Err(err!( |
| 330 | "Key starting with {:?} has been deleted from cache.", |
| 331 | if k.len() > 8 { &k[..8] } else { &k }; |
| 332 | Missing, Data)), |
| 333 | None => Err(err!( |
| 334 | "Key starting with {:?} not present in cache.", |
| 335 | if k.len() > 8 { &k[..8] } else { &k }; |
| 336 | Unknown, Data)), |
| 337 | } |
| 338 | } |
| 339 | |
| 340 | /// Moves the cached location of a record a collection has carried to its new start, if the |
| 341 | /// cache still names that record, and gives the location it named before. A record is its |
| 342 | /// key and its stamp: matched by file alone, a key with two records in the file, the older |
| 343 | /// carried because its supersession had yet to arrive, had its location moved to the older |
| 344 | /// record's and back, and the move entries were spent against the wrong offsets. |
| 345 | pub fn reanchor( |
| 346 | &mut self, |
| 347 | k: &Vec<u8>, |
| 348 | loc: &FileLocation, |
| 349 | meta: &Meta<UIDL, UID>, |
| 350 | ) |
| 351 | -> Option<FileLocation> |
| 352 | { |
| 353 | match self.map.get_mut(k) { |
| 354 | Some(CacheEntry::LocatedValue(mloc, _)) => { |
| 355 | if mloc.floc.fnum == loc.fnum && mloc.meta == *meta { |
| 356 | let old_floc = mloc.floc.clone(); |
| 357 | mloc.floc.start = loc.start; |
| 358 | return Some(old_floc); |
| 359 | } |
| 360 | None |
| 361 | }, |
| 362 | _ => None, |
| 363 | } |
| 364 | } |
| 365 | |
| 366 | /// Looks for the key in the cache map and if present, returns the value if it is present, or |
| 367 | /// the latest location. If the value is a `Dat::Box`, this method will (recursively) |
| 368 | /// obtain the final value, but note that Rust has a default recursion limit of 128. |
| 369 | pub fn get(&self, k: &[u8]) -> Outcome<Option<ValueOrLocation<UIDL, UID>>> { |
| 370 | match self.map.get(k) { |
| 371 | Some(CacheEntry::LocatedValue(mloc, val)) => { |
| 372 | match &val { |
| 373 | Some(val) => { // Use cache value. |
| 374 | if val.len() > 1 && val[0] == Dat::BOX_CODE { |
| 375 | // Automatic key referral allows multiple keys to |
| 376 | // point to the same value. |
| 377 | return self.get(&val[1..]); |
| 378 | } |
| 379 | return Ok(Some(ValueOrLocation::Value( |
| 380 | val.clone(), |
| 381 | mloc.meta().clone(), |
| 382 | ))); |
| 383 | }, |
| 384 | // Return data location. |
| 385 | None => return Ok(Some(ValueOrLocation::Location(mloc.clone()))), |
| 386 | } |
| 387 | }, |
| 388 | Some(CacheEntry::Deleted(meta)) => |
| 389 | return Ok(Some(ValueOrLocation::Deleted(meta.clone()))), |
| 390 | None => return Ok(None), |
| 391 | } |
| 392 | } |
| 393 | |
| 394 | pub fn clear_all_values(&mut self) { |
| 395 | for (_k, centry) in self.map.iter_mut() { |
| 396 | if let CacheEntry::LocatedValue(_, val) = centry { |
| 397 | *val = None; |
| 398 | } |
| 399 | } |
| 400 | } |
| 401 | |
| 402 | } |
| 403 | |
| 404 | #[derive(Clone, Debug)] |
| 405 | struct CacheWriteTrackerInfo { |
| 406 | key: Vec<u8>, |
| 407 | hash: u64, |
| 408 | vlen: usize, |
| 409 | } |
| 410 | |
| 411 | /// # Cache resource management |
| 412 | /// Maintaining an ordered (forward) map of `Timestamp`s to keys can facilitate a first-in, |
| 413 | /// first-out cache size limiting strategy. In other words, this allows us to dump the oldest |
| 414 | /// values from the cache first. A reverse map is maintained to allow us to identify when we can |
| 415 | /// delete an old timestamp from the forward map for a given key. |
| 416 | /// ```ignore |
| 417 | /// |
| 418 | /// Forward map: Reverse map: |
| 419 | /// t1 -> k1 k1 -> t1 |
| 420 | /// t2 -> k2 k2 -> t2 |
| 421 | /// t3 -> k1 |
| 422 | /// |
| 423 | /// Now, when t3 -> k1 is added to the forward map, the presence of k1 -> t1 in the reverse map |
| 424 | /// tells us that the t1 -> k1 entry in the forward map is redundant and can be deleted. At the |
| 425 | /// same time, the reverse map is also updated |
| 426 | /// |
| 427 | /// Forward map: Reverse map: |
| 428 | /// t2 -> k2 k1 -> t3 |
| 429 | /// t3 -> k1 k2 -> t2 |
| 430 | /// |
| 431 | /// ``` |
| 432 | #[derive(Clone, Debug)] |
| 433 | pub struct CacheWriteTracker { |
| 434 | fwd: BTreeMap<Timestamp, CacheWriteTrackerInfo>, |
| 435 | rev: BTreeMap<u64, Timestamp>, |
| 436 | bs: usize, // base size |
| 437 | size: usize, // track total size estimate for data structure |
| 438 | } |
| 439 | |
| 440 | impl Default for CacheWriteTracker { |
| 441 | fn default() -> Self { |
| 442 | Self { |
| 443 | fwd: BTreeMap::new(), |
| 444 | rev: BTreeMap::new(), |
| 445 | bs: ( // base size for an entry in both maps, not including CacheWriteTrackerInfo::Key |
| 446 | 2 * std::mem::size_of::<Timestamp>() + |
| 447 | 2 * std::mem::size_of::<u64>() + |
| 448 | std::mem::size_of::<usize>() |
| 449 | ), |
| 450 | size: 0, |
| 451 | } |
| 452 | } |
| 453 | } |
| 454 | |
| 455 | impl CacheWriteTracker { |
| 456 | /// This is only used for value insertions into the cache, not file locations. |
| 457 | fn insert( |
| 458 | &mut self, |
| 459 | t3: Timestamp, |
| 460 | k1: &Vec<u8>, |
| 461 | vlen: usize, |
| 462 | ) |
| 463 | -> Outcome<()> |
| 464 | { |
| 465 | let hash = seahash::hash(&k1); |
| 466 | let cwti = CacheWriteTrackerInfo { |
| 467 | key: k1.clone(), |
| 468 | hash: hash, |
| 469 | vlen: vlen, |
| 470 | }; |
| 471 | self.fwd.insert(t3.clone(), cwti); |
| 472 | match self.rev.insert(hash, t3) { |
| 473 | Some(t1) => { |
| 474 | self.fwd.remove(&t1); |
| 475 | // just an update, no size change |
| 476 | }, |
| 477 | None => { |
| 478 | self.size = try_add!(&self.size, self.bs + k1.len()); // new insertion |
| 479 | }, |
| 480 | } |
| 481 | Ok(()) |
| 482 | } |
| 483 | |
| 484 | /// Identifies the oldest cached values whose lengths sum to at least the given value length, |
| 485 | /// deleting their entries in the `CacheWriteTracker` while returning the list of associated |
| 486 | /// keys, allowing the caller to scrub values from the cache, and advising of the size |
| 487 | /// reduction of the tracker. If the given value length exceeds the length of all existing |
| 488 | /// cached values, the entire `CacheWriteTracker` contents will be deleted and the desired |
| 489 | /// cache size reduction will not be achieved. |
| 490 | fn jettison( |
| 491 | &mut self, |
| 492 | vlen: usize, |
| 493 | ) |
| 494 | -> Outcome<Vec<Vec<u8>>> |
| 495 | { |
| 496 | let mut vlensum = 0; |
| 497 | let mut jettison = Vec::new(); |
| 498 | for (t, cwti) in &self.fwd { |
| 499 | jettison.push(t.clone()); |
| 500 | // Account for the value in the cache and for the |
| 501 | // entries in the fwd and rev maps here. |
| 502 | vlensum += cwti.vlen + self.bs + cwti.key.len(); |
| 503 | if vlensum > vlen { |
| 504 | break; |
| 505 | } |
| 506 | } |
| 507 | |
| 508 | let mut cwt_size_reduction = 0; |
| 509 | let mut keys = Vec::new(); |
| 510 | for t in jettison { |
| 511 | match self.fwd.remove(&t) { |
| 512 | Some(cwti) => { |
| 513 | self.rev.remove(&cwti.hash); |
| 514 | cwt_size_reduction += self.bs + cwti.key.len(); |
| 515 | keys.push(cwti.key); |
| 516 | }, |
| 517 | None => (), // unreachable |
| 518 | } |
| 519 | } |
| 520 | |
| 521 | self.size = try_sub!(&self.size, cwt_size_reduction); |
| 522 | |
| 523 | Ok(keys) |
| 524 | } |
| 525 | } |