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
31pub struct MemoryCache {
32    entries: Mutex<BTreeMap<Hash, Held>>,
33    max_bytes: u64,
34    clock: AtomicU64,
35    held: AtomicU64,
36    evicted: AtomicU64,
37}
38
39impl MemoryCache {
40    pub fn new() -> MemoryCache {
41        MemoryCache::holding(DEFAULT_MAX_BYTES)
42    }
43
44    pub fn holding(max_bytes: u64) -> MemoryCache {
45        MemoryCache {
46            entries: Mutex::new(BTreeMap::new()),
47            max_bytes,
48            clock: AtomicU64::new(0),
49            held: AtomicU64::new(0),
50            evicted: AtomicU64::new(0),
51        }
52    }
53
54    /// A holder that panicked mid-mutation may have left an entry half-written, and every
55    /// entry is derived data: the store empties, charged as evicted, and keeps serving.
56    fn locked(&self) -> std::sync::MutexGuard<'_, BTreeMap<Hash, Held>> {
57        match self.entries.lock() {
58            Ok(entries) => entries,
59            Err(poisoned) => {
60                let mut entries = poisoned.into_inner();
61                let dropped: u64 = entries.values().map(Held::bytes).sum();
62                entries.clear();
63                self.held.store(0, Ordering::Relaxed);
64                self.evicted.store(dropped, Ordering::Relaxed);
65                self.entries.clear_poison();
66                entries
67            }
68        }
69    }
70}
71
72impl Default for MemoryCache {
73    fn default() -> MemoryCache {
74        MemoryCache::new()
75    }
76}
77
78impl Cache for MemoryCache {
79    fn max_bytes(&self) -> u64 {
80        self.max_bytes
81    }
82
83    fn held_bytes(&self) -> u64 {
84        self.held.load(Ordering::Relaxed)
85    }
86
87    fn evicted_bytes(&self) -> u64 {
88        self.evicted.load(Ordering::Relaxed)
89    }
90
91    fn holds(&self, key: Hash) -> bool {
92        self.locked().contains_key(&key)
93    }
94
95    fn load(&self, key: Hash, node: &str, expected: Expected) -> Option<Entry> {
96        let tick = self.clock.fetch_add(1, Ordering::Relaxed);
97        let mut entries = self.locked();
98        let held = entries.get_mut(&key)?;
99        if !held.payload.answers(expected) {
100            entries.remove(&key);
101            return None;
102        }
103        held.read = tick;
104        Some(Entry {
105            payload: held.payload.clone(),
106            traces: held
107                .traces
108                .iter()
109                .map(|t| FilterTrace {
110                    node: node.to_string(),
111                    ..t.clone()
112                })
113                .collect(),
114            label: held.label.clone(),
115            tier: Tier::Memory,
116        })
117    }
118
119    /// A hit is a memcpy, so anything computed is kept until the budget says otherwise.
120    fn worth_storing(&self, _cost: Duration, _bytes: usize, _kind: PayloadKind) -> bool {
121        true
122    }
123
124    fn store(&self, key: Hash, payload: &Payload, traces: &[FilterTrace], label: Option<&Label>) {
125        let read = self.clock.fetch_add(1, Ordering::Relaxed);
126        self.locked().insert(
127            key,
128            Held {
129                payload: payload.clone(),
130                traces: traces.to_vec(),
131                label: label.cloned(),
132                read,
133            },
134        );
135    }
136
137    /// The read time is the tick `load` stamped the entry with.
138    fn sweep(&self) {
139        let mut entries = self.locked();
140        let order: Vec<(u64, u64, Hash)> = entries
141            .iter()
142            .map(|(k, h)| (h.read, h.bytes(), *k))
143            .collect();
144        let swept = evict::to_cap(order, self.max_bytes, |key| entries.remove(key).is_some());
145        self.held.store(swept.held, Ordering::Relaxed);
146        self.evicted.store(swept.evicted, Ordering::Relaxed);
147    }
148}