Skip to main content

sva_engine/cache/
memory.rs

1// Concern: keeps rendered buffers in this process's own heap under their content hash | Non-concern: what a store is for (cache.rs), outliving the process (disk.rs) | IO: (Hash) -> a buffer + traces
2
3use 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
15/// What a process can hold, not what a filesystem can.
16pub 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    /// A holder that panicked mid-mutation may have left an entry half-written, and every
64    /// entry is derived data: the store empties, charged as evicted, and keeps serving.
65    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    /// A hit is a memcpy, so anything computed is kept until the budget says otherwise.
123    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    /// The read time is the tick `load` stamped the entry with.
141    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}