Skip to main content

sva_engine/cache/
persist.rs

1// Concern: node values on a backend, staged then committed: version, budget, eviction | Non-concern: the bytes' medium (a Backend) | IO: (key) -> Stored; samples -> staged; persist() -> commits
2
3use std::collections::{BTreeMap, HashMap};
4use std::sync::{Mutex, MutexGuard};
5
6use sva_formula::{Codomain, Hash};
7use sva_samples::{Buffer, Extent, Grid, Label};
8
9use super::{codec, joined};
10
11/// A node's samples, and what a render needs to read them in place of typing the node.
12#[derive(Clone, Debug, PartialEq)]
13pub struct Stored {
14    pub key: Hash,
15    pub segments: Vec<Buffer>,
16    pub label: Label,
17    pub width: u8,
18    pub codomain: Codomain,
19    pub rate: Option<u32>,
20    pub grid: Grid,
21    pub support: Extent,
22    /// Flops it and every value under it cost when computed.
23    pub priced: u128,
24    /// The most seconds a read under it moved to land on a sample.
25    pub moved: f64,
26    /// A reader may take these samples in place of computing the node.
27    pub readable: bool,
28}
29
30impl Stored {
31    pub(crate) fn holds(&self, over: Extent) -> bool {
32        let mut from = over.start;
33        let mut parts: Vec<Extent> = self.segments.iter().map(Buffer::extent).collect();
34        parts.sort_by_key(|e| e.start);
35        for part in parts {
36            if from >= over.end || part.start > from {
37                break;
38            }
39            from = from.max(part.end);
40        }
41        from >= over.end
42    }
43}
44
45/// One node's value at one rate and profile, named by its source.
46pub(crate) fn node_key(identity: Hash, rate: u32, profile: &sva_samples::Profile) -> Hash {
47    super::mixed(
48        identity,
49        &[
50            u64::from(rate),
51            profile.precision_bits as u64,
52            profile.ceiling_hz.to_bits(),
53            0x6e_6f_64_65_00_00_00_01,
54        ],
55    )
56}
57
58pub const DEFAULT_STORE_BYTES: u64 = 2 << 30;
59
60pub const STORE_VERSION: &str = concat!(
61    "sva-engine ",
62    env!("CARGO_PKG_VERSION"),
63    " build ",
64    env!("SVA_ENGINE_BUILD"),
65    " format 2"
66);
67
68pub const VERSION_NAME: &str = "version";
69
70const RECENCY_NAME: &str = "recency";
71
72const META: &str = "meta";
73
74/// Named byte storage. A `put` is whole or absent: no reader ever sees half of one.
75pub trait Backend: Sized {
76    fn get(&self, name: &str) -> impl Future<Output = Result<Option<Vec<u8>>, String>>;
77    fn put(&self, name: &str, bytes: &[u8]) -> impl Future<Output = Result<(), String>>;
78    fn delete(&self, name: &str) -> impl Future<Output = Result<(), String>>;
79    /// Every name held, with its size in bytes.
80    fn list(&self) -> impl Future<Output = Result<Vec<(String, u64)>, String>>;
81    /// An area of its own beside this one's names, which no other store writes and `list`
82    /// never names.
83    fn staging(&self) -> impl Future<Output = Result<Self, String>>;
84    /// `name` moved from here into `to` in one step, over whatever `to` held under it.
85    fn rename(&self, name: &str, to: &Self) -> impl Future<Output = Result<(), String>>;
86}
87
88/// What a stream looks a node up in before typing it: a store, or nothing at all.
89pub trait Through {
90    fn lookup(&self, key: Hash) -> impl Future<Output = Option<Stored>>;
91}
92
93impl<B: Backend> Through for Store<B> {
94    fn lookup(&self, key: Hash) -> impl Future<Output = Option<Stored>> {
95        Store::lookup(self, key)
96    }
97}
98
99pub struct NoStore;
100
101impl Through for NoStore {
102    async fn lookup(&self, _: Hash) -> Option<Stored> {
103        None
104    }
105}
106
107#[derive(Default)]
108struct Index {
109    held: HashMap<Hash, (u64, u64)>,
110    clock: u64,
111    bytes: u64,
112}
113
114impl Index {
115    fn touch(&mut self, key: Hash, bytes: u64) {
116        self.clock += 1;
117        let at = self.clock;
118        if let Some((old, _)) = self.held.insert(key, (bytes, at)) {
119            self.bytes -= old;
120        }
121        self.bytes += bytes;
122    }
123
124    fn read(&mut self, key: Hash) {
125        self.clock += 1;
126        let at = self.clock;
127        if let Some((_, held)) = self.held.get_mut(&key) {
128            *held = at;
129        }
130    }
131
132    /// The directory's entries at their sizes; one another store committed is least recent.
133    fn sync(&mut self, listed: Vec<(Hash, u64)>) {
134        let held = std::mem::take(&mut self.held);
135        self.held = listed
136            .into_iter()
137            .map(|(key, bytes)| (key, (bytes, held.get(&key).map_or(0, |(_, at)| *at))))
138            .collect();
139        self.bytes = self.held.values().map(|(bytes, _)| bytes).sum();
140    }
141
142    fn forget(&mut self, key: Hash) {
143        if let Some((bytes, _)) = self.held.remove(&key) {
144            self.bytes -= bytes;
145        }
146    }
147
148    fn oldest_first(&self) -> Vec<Hash> {
149        let mut keys: Vec<(u64, Hash)> = self.held.iter().map(|(k, (_, at))| (*at, *k)).collect();
150        keys.sort_unstable();
151        keys.into_iter().map(|(_, k)| k).collect()
152    }
153}
154
155#[derive(Clone, Default)]
156struct Staged {
157    chunks: Vec<String>,
158    meta: bool,
159}
160
161/// Node values under their keys. A render reads through it and spills what it computes to the
162/// staging area; only `persist` changes the store, committing each staged value whole by
163/// rename, least recently used going first past the budget.
164pub struct Store<B> {
165    backend: B,
166    staging: B,
167    max_bytes: u64,
168    index: Mutex<Index>,
169    staged: Mutex<BTreeMap<Hash, Staged>>,
170    written: Mutex<u64>,
171}
172
173#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
174pub struct Persisted {
175    pub written: usize,
176    pub evicted: usize,
177}
178
179fn name_of(key: Hash) -> String {
180    format!("{:016x}{:016x}", key.0, key.1)
181}
182
183fn key_of(name: &str) -> Option<Hash> {
184    let hex = |s: &str| u64::from_str_radix(s, 16).ok();
185    let valid = name.len() == 32 && name.bytes().all(|b| b.is_ascii_hexdigit());
186    valid.then(|| Some(Hash(hex(&name[..16])?, hex(&name[16..])?)))?
187}
188
189fn locked<T>(held: &Mutex<T>) -> MutexGuard<'_, T> {
190    held.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
191}
192
193impl<B: Backend> Store<B> {
194    /// A store another version wrote is emptied first.
195    pub async fn open(backend: B, max_bytes: u64) -> Result<Store<B>, String> {
196        let version = backend.get(VERSION_NAME).await?;
197        if version.as_deref() != Some(STORE_VERSION.as_bytes()) {
198            for (name, _) in backend.list().await? {
199                if key_of(&name).is_some() || name == RECENCY_NAME || name == VERSION_NAME {
200                    backend.delete(&name).await?;
201                }
202            }
203            backend.put(VERSION_NAME, STORE_VERSION.as_bytes()).await?;
204        }
205        let order = backend.get(RECENCY_NAME).await?.unwrap_or_default();
206        let order = String::from_utf8_lossy(&order);
207        let rank: HashMap<&str, usize> = order.lines().enumerate().map(|(k, n)| (n, k)).collect();
208        let mut held: Vec<(usize, Hash, u64)> = backend
209            .list()
210            .await?
211            .into_iter()
212            .filter_map(|(name, bytes)| {
213                let at = rank.get(name.as_str()).map_or(0, |k| k + 1);
214                Some((at, key_of(&name)?, bytes))
215            })
216            .collect();
217        held.sort_unstable();
218        let mut index = Index::default();
219        for (_, key, bytes) in held {
220            index.touch(key, bytes);
221        }
222        let staging = backend.staging().await?;
223        Ok(Store {
224            backend,
225            staging,
226            max_bytes,
227            index: Mutex::new(index),
228            staged: Mutex::new(BTreeMap::new()),
229            written: Mutex::new(0),
230        })
231    }
232
233    pub fn max_bytes(&self) -> u64 {
234        self.max_bytes
235    }
236
237    pub fn bytes(&self) -> u64 {
238        locked(&self.index).bytes
239    }
240
241    pub fn holds(&self, key: Hash) -> bool {
242        locked(&self.index).held.contains_key(&key)
243    }
244
245    /// Staged first; the backend asked past the index, which misses another store's commits.
246    pub(crate) async fn lookup(&self, key: Hash) -> Option<Stored> {
247        if let Some(staged) = self.staged_value(key).await {
248            return Some(staged);
249        }
250        let Some(bytes) = self.backend.get(&name_of(key)).await.ok()? else {
251            locked(&self.index).forget(key);
252            return None;
253        };
254        let found = codec::read_entry(&bytes).filter(|found| found.key == key)?;
255        let mut index = locked(&self.index);
256        match index.held.contains_key(&key) {
257            true => index.read(key),
258            false => index.touch(key, bytes.len() as u64),
259        }
260        Some(found)
261    }
262
263    async fn staged_value(&self, key: Hash) -> Option<Stored> {
264        let chunks = {
265            let staged = locked(&self.staged);
266            let held = staged.get(&key).filter(|staged| staged.meta)?;
267            held.chunks.clone()
268        };
269        let meta = self.staging.get(&staged_name(key, META)).await.ok()??;
270        let mut stored = codec::read_entry(&meta)?;
271        for chunk in chunks {
272            let bytes = self.staging.get(&chunk).await.ok()??;
273            joined(&mut stored.segments, vec![codec::read_chunk(&bytes)?]);
274        }
275        Some(stored)
276    }
277
278    /// One run of `key`'s samples, out of memory and into the staging area.
279    pub(crate) async fn stage(&self, key: Hash, samples: &Buffer) -> Result<(), String> {
280        let n = {
281            let mut written = locked(&self.written);
282            *written += 1;
283            *written
284        };
285        let name = staged_name(key, &n.to_string());
286        self.staging.put(&name, &codec::chunk(samples)).await?;
287        locked(&self.staged)
288            .entry(key)
289            .or_default()
290            .chunks
291            .push(name);
292        Ok(())
293    }
294
295    /// `stored` holds no samples.
296    pub(crate) async fn stage_meta(&self, key: Hash, stored: &Stored) -> Result<(), String> {
297        let name = staged_name(key, META);
298        self.staging.put(&name, &codec::entry(stored)).await?;
299        locked(&self.staged).entry(key).or_default().meta = true;
300        Ok(())
301    }
302
303    /// A staged value whose meta was never staged is dropped. Each key leaves the staging
304    /// record only once its files have, so a write that fails leaves every later one staged.
305    pub async fn persist(&self) -> Result<Persisted, String> {
306        let mut done = Persisted::default();
307        let keys: Vec<Hash> = locked(&self.staged).keys().copied().collect();
308        for key in keys {
309            let Some(held) = locked(&self.staged).get(&key).cloned() else {
310                continue;
311            };
312            if let Some(bytes) = self.committed(key, &held).await? {
313                locked(&self.index).touch(key, bytes);
314                done.written += 1;
315            }
316            for chunk in &held.chunks {
317                self.staging.delete(chunk).await?;
318            }
319            self.staging.delete(&staged_name(key, META)).await?;
320            locked(&self.staged).remove(&key);
321        }
322        let listed = self.backend.list().await?;
323        let listed = listed
324            .into_iter()
325            .filter_map(|(name, bytes)| Some((key_of(&name)?, bytes)));
326        locked(&self.index).sync(listed.collect());
327        let oldest = locked(&self.index).oldest_first();
328        for key in oldest {
329            if self.bytes() <= self.max_bytes {
330                break;
331            }
332            self.backend.delete(&name_of(key)).await?;
333            locked(&self.index).forget(key);
334            done.evicted += 1;
335        }
336        let order: String = locked(&self.index)
337            .oldest_first()
338            .into_iter()
339            .map(|key| name_of(key) + "\n")
340            .collect();
341        self.backend.put(RECENCY_NAME, order.as_bytes()).await?;
342        Ok(done)
343    }
344
345    /// `key`'s staged value moved whole into the store, and its size; none past the budget.
346    async fn committed(&self, key: Hash, held: &Staged) -> Result<Option<u64>, String> {
347        let whole = match held.meta {
348            true => self.staged_value_of(key, held).await?,
349            false => None,
350        };
351        let Some(stored) = whole.filter(|stored| !stored.segments.is_empty()) else {
352            return Ok(None);
353        };
354        let bytes = codec::entry(&stored);
355        if bytes.len() as u64 > self.max_bytes {
356            return Ok(None);
357        }
358        let name = name_of(key);
359        self.staging.put(&name, &bytes).await?;
360        self.staging.rename(&name, &self.backend).await?;
361        Ok(Some(bytes.len() as u64))
362    }
363
364    async fn staged_value_of(&self, key: Hash, held: &Staged) -> Result<Option<Stored>, String> {
365        let Some(meta) = self.staging.get(&staged_name(key, META)).await? else {
366            return Ok(None);
367        };
368        let Some(mut stored) = codec::read_entry(&meta) else {
369            return Ok(None);
370        };
371        for chunk in &held.chunks {
372            let bytes = self.staging.get(chunk).await?;
373            match bytes.as_deref().and_then(codec::read_chunk) {
374                Some(samples) => joined(&mut stored.segments, vec![samples]),
375                None => return Ok(None),
376            }
377        }
378        Ok(Some(stored))
379    }
380}
381
382fn staged_name(key: Hash, part: &str) -> String {
383    format!("{}.{part}", name_of(key))
384}