sva_engine/cache/
memory.rs1use std::collections::BTreeMap;
4use std::sync::Mutex;
5use std::sync::atomic::{AtomicU64, Ordering};
6use std::time::Duration;
7
8use sva_samples::{FilterTrace, Label};
9
10use sva_formula::Hash;
11
12use super::evict;
13use super::{Cache, Entry, Expected, Payload, PayloadKind, Tier};
14
15pub const DEFAULT_MAX_BYTES: u64 = 2 << 30;
17
18struct Held {
19 payload: Payload,
20 traces: Vec<FilterTrace>,
21 label: Option<Label>,
22 read: u64,
23}
24
25impl Held {
26 fn bytes(&self) -> u64 {
27 self.payload.bytes() as u64
28 }
29
30 fn entry(&self, node: &str) -> Entry {
31 Entry {
32 payload: self.payload.clone(),
33 traces: super::renamed(&self.traces, node),
34 label: self.label.clone(),
35 tier: Tier::Memory,
36 }
37 }
38}
39
40pub struct MemoryCache {
41 entries: Mutex<BTreeMap<Hash, Held>>,
42 max_bytes: u64,
43 clock: AtomicU64,
44 held: AtomicU64,
45 evicted: AtomicU64,
46}
47
48impl MemoryCache {
49 pub fn new() -> MemoryCache {
50 MemoryCache::holding(DEFAULT_MAX_BYTES)
51 }
52
53 pub fn holding(max_bytes: u64) -> MemoryCache {
54 MemoryCache {
55 entries: Mutex::new(BTreeMap::new()),
56 max_bytes,
57 clock: AtomicU64::new(0),
58 held: AtomicU64::new(0),
59 evicted: AtomicU64::new(0),
60 }
61 }
62
63 fn locked(&self) -> std::sync::MutexGuard<'_, BTreeMap<Hash, Held>> {
66 match self.entries.lock() {
67 Ok(entries) => entries,
68 Err(poisoned) => {
69 let mut entries = poisoned.into_inner();
70 let dropped: u64 = entries.values().map(Held::bytes).sum();
71 entries.clear();
72 self.held.store(0, Ordering::Relaxed);
73 self.evicted.store(dropped, Ordering::Relaxed);
74 self.entries.clear_poison();
75 entries
76 }
77 }
78 }
79}
80
81impl Default for MemoryCache {
82 fn default() -> MemoryCache {
83 MemoryCache::new()
84 }
85}
86
87impl Cache for MemoryCache {
88 fn max_bytes(&self) -> u64 {
89 self.max_bytes
90 }
91
92 fn held_bytes(&self) -> u64 {
93 self.held.load(Ordering::Relaxed)
94 }
95
96 fn evicted_bytes(&self) -> u64 {
97 self.evicted.load(Ordering::Relaxed)
98 }
99
100 fn holds(&self, key: Hash) -> bool {
101 self.locked().contains_key(&key)
102 }
103
104 fn load(&self, key: Hash, node: &str, expected: Expected) -> Option<Entry> {
105 let tick = self.clock.fetch_add(1, Ordering::Relaxed);
106 let mut entries = self.locked();
107 let held = entries.get_mut(&key)?;
108 if !held.payload.answers(expected) {
109 entries.remove(&key);
110 return None;
111 }
112 held.read = tick;
113 Some(held.entry(node))
114 }
115
116 fn peek(&self, key: Hash, node: &str, expected: Expected) -> Option<Entry> {
117 let entries = self.locked();
118 let held = entries.get(&key)?;
119 held.payload.answers(expected).then(|| held.entry(node))
120 }
121
122 fn worth_storing(&self, _cost: Duration, _bytes: usize, _kind: PayloadKind) -> bool {
124 true
125 }
126
127 fn store(&self, key: Hash, payload: &Payload, traces: &[FilterTrace], label: Option<&Label>) {
128 let read = self.clock.fetch_add(1, Ordering::Relaxed);
129 self.locked().insert(
130 key,
131 Held {
132 payload: payload.clone(),
133 traces: traces.to_vec(),
134 label: label.cloned(),
135 read,
136 },
137 );
138 }
139
140 fn sweep(&self) {
142 let mut entries = self.locked();
143 let order: Vec<(u64, u64, Hash)> = entries
144 .iter()
145 .map(|(k, h)| (h.read, h.bytes(), *k))
146 .collect();
147 let swept = evict::to_cap(order, self.max_bytes, |key| entries.remove(key).is_some());
148 self.held.store(swept.held, Ordering::Relaxed);
149 self.evicted.store(swept.evicted, Ordering::Relaxed);
150 }
151}