Oregami
Repositories/oxedyne/fe2o3

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

1use 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
12use oxedyne_fe2o3_data::time::Timestamp;
13use oxedyne_fe2o3_iop_db::api::Meta;
14use oxedyne_fe2o3_jdat::{
15 daticle::Dat,
16 id::NumIdDat,
17};
18
19use std::{
20 collections::BTreeMap,
21 fmt,
22 marker::PhantomData,
23};
24
25#[derive(Clone, Debug, Default, Eq, PartialEq)]
26pub struct CacheId(pub u16);
27
28impl fmt::Display for CacheId {
29 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
30 write!(f, "CacheId({})", self.0)
31 }
32}
33
34impl 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)]
42pub 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)]
53pub struct MetaLocation<
54 const UIDL: usize,
55 UID: NumIdDat<UIDL>,
56> {
57 meta: Meta<UIDL, UID>,
58 floc: FileLocation,
59}
60
61impl<
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)]
78struct 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
89impl<
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)]
107pub 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
118impl<
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)]
130pub 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)]
140pub 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
152impl<
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)]
405struct 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)]
433pub 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
440impl 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
455impl 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}