Skip to main content

sva_engine/cache/
persist.rs

1// Concern: node values on a backend beneath memory, staged then committed: version, budget, eviction | Non-concern: the bytes' medium (a Backend) | IO: (key) -> Header; samples -> staged; persist()
2
3use std::collections::BTreeMap;
4use std::sync::{Arc, Mutex, MutexGuard};
5
6use sva_formula::{FOURIER_DUAL_RULES_VERSION, Hash};
7use sva_samples::{Buffer, Extent};
8
9use super::codec::{self, STORE_FORMAT};
10use super::index::{Index, key_of, name_of};
11use super::joined;
12use super::stored::{Header, Samples};
13
14pub const DEFAULT_STORE_BYTES: u64 = 2 << 30;
15
16fn version() -> String {
17    format!("sva store format {STORE_FORMAT} table {FOURIER_DUAL_RULES_VERSION}")
18}
19
20pub const INDEX_NAME: &str = "index";
21
22const META: &str = "meta";
23
24/// Named bytes. A `put` is whole or absent: no reader sees half of one.
25pub trait Backend: Sized {
26    type Lock;
27    /// The one exclusive lock over these names, shared by every store over them.
28    fn lock(&self) -> impl Future<Output = Result<Self::Lock, String>>;
29    fn get(&self, name: &str) -> impl Future<Output = Result<Option<Vec<u8>>, String>>;
30    /// At most `len` bytes from `from` on, fewer where the name ends sooner.
31    fn get_range(
32        &self,
33        name: &str,
34        from: u64,
35        len: u64,
36    ) -> impl Future<Output = Result<Option<Vec<u8>>, String>>;
37    fn put(&self, name: &str, bytes: &[u8]) -> impl Future<Output = Result<(), String>>;
38    /// False where another holder has `name` open, so it stays.
39    fn delete(&self, name: &str) -> impl Future<Output = Result<bool, String>>;
40    /// Each name, with its size in bytes.
41    fn list(&self) -> impl Future<Output = Result<Vec<(String, u64)>, String>>;
42    /// An area beside these names that no other store writes and `list` never names.
43    fn staging(&self) -> impl Future<Output = Result<Self, String>>;
44    /// `name` moved from here into `to` in one step, over whatever `to` held under it; as
45    /// `delete`, false where either is open.
46    fn rename(&self, name: &str, to: &Self) -> impl Future<Output = Result<bool, String>>;
47}
48
49#[derive(Clone, Default)]
50struct Staged {
51    chunks: Vec<(String, Extent)>,
52    meta: bool,
53}
54
55/// Node values under their keys, beneath the memory tier alone: it stages, and only `persist`
56/// commits each value whole by rename, evicting the least recently used past the budget. Every
57/// change to the backend's names runs under its lock; lookups and reads never wait on it.
58pub struct Store<B> {
59    backend: B,
60    staging: B,
61    max_bytes: u64,
62    index: Mutex<Index>,
63    staged: Mutex<BTreeMap<Hash, Staged>>,
64    written: Mutex<u64>,
65}
66
67enum Prepared {
68    Ready { bytes: u64 },
69    Dropped,
70    Refused,
71}
72
73#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
74pub struct Persisted {
75    pub written: usize,
76    pub evicted: usize,
77    /// Values larger than the whole budget.
78    pub refused: usize,
79}
80
81fn locked<T>(held: &Mutex<T>) -> MutexGuard<'_, T> {
82    held.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
83}
84
85impl<B: Backend> Store<B> {
86    /// Over another format's index, a fresh one: what that format left is neither read nor
87    /// removed, nor counted against the budget.
88    pub async fn open(backend: B, max_bytes: u64) -> Result<Store<B>, String> {
89        let version = version();
90        let held = backend.lock().await?;
91        let found = backend.get(INDEX_NAME).await?;
92        let index = match found.and_then(|text| Index::read(&text, &version)) {
93            Some(index) => index,
94            None => {
95                let empty = Index::default();
96                backend
97                    .put(INDEX_NAME, empty.text(&version).as_bytes())
98                    .await?;
99                empty
100            }
101        };
102        let staging = backend.staging().await?;
103        drop(held);
104        Ok(Store {
105            backend,
106            staging,
107            max_bytes,
108            index: Mutex::new(index),
109            staged: Mutex::new(BTreeMap::new()),
110            written: Mutex::new(0),
111        })
112    }
113
114    pub fn max_bytes(&self) -> u64 {
115        self.max_bytes
116    }
117
118    pub fn bytes(&self) -> u64 {
119        locked(&self.index).bytes
120    }
121
122    pub fn holds(&self, key: Hash) -> bool {
123        locked(&self.index).held.contains_key(&key)
124    }
125
126    pub(crate) fn used(&self, key: Hash) {
127        locked(&self.index).used(key);
128    }
129
130    /// Reads headers, never a sample; one standing for another's samples reads that one's.
131    pub(crate) async fn lookup(&self, key: Hash) -> Option<Header> {
132        let found = self.written(key).await?;
133        let &Samples::Of { key: of, by } = found.samples() else {
134            return Some(found);
135        };
136        let samples = match self.written(of).await?.into_parts().1 {
137            Samples::Entry { file, runs, .. } => Samples::Entry {
138                file,
139                runs,
140                shift: -by,
141            },
142            Samples::Staged { file, chunks, .. } => Samples::Staged {
143                file,
144                chunks,
145                shift: -by,
146            },
147            _ => return None,
148        };
149        Some(Header::new(found.into_parts().0, samples))
150    }
151
152    /// Staged first; past the index, the backend, for another store's commits.
153    async fn written(&self, key: Hash) -> Option<Header> {
154        if let Some(staged) = self.staged_value(key).await {
155            return Some(staged);
156        }
157        let name = name_of(key);
158        let Some(mut bytes) = self.backend.get_range(&name, 0, 8).await.ok()? else {
159            locked(&self.index).forget(key);
160            return None;
161        };
162        let rest = codec::head_len(&bytes)? - 8;
163        bytes.extend(self.backend.get_range(&name, 8, rest as u64).await.ok()??);
164        let (found, len) =
165            codec::read_head(&bytes, key).filter(|(found, _)| found.stored().key == key)?;
166        let mut index = locked(&self.index);
167        match index.held.contains_key(&key) {
168            true => index.used(key),
169            false => index.touch(key, len),
170        }
171        Some(found)
172    }
173
174    async fn staged_value(&self, key: Hash) -> Option<Header> {
175        let chunks = {
176            let staged = locked(&self.staged);
177            let held = staged.get(&key).filter(|staged| staged.meta)?;
178            held.chunks.clone()
179        };
180        let meta = self.staging.get(&staged_name(key, META)).await.ok()??;
181        let (found, _) = codec::read_head(&meta, key)?;
182        Some(match found.refers() {
183            true => found,
184            false => {
185                let samples = Samples::Staged {
186                    file: key,
187                    chunks,
188                    shift: 0,
189                };
190                Header::new(found.into_parts().0, samples)
191            }
192        })
193    }
194
195    /// Whole chunks holding `over`; `None` where they are gone or corrupt.
196    pub(crate) async fn read(&self, head: &Header, over: Extent) -> Option<Vec<Buffer>> {
197        let mut out = Vec::new();
198        match head.samples() {
199            Samples::None | Samples::Of { .. } => {}
200            Samples::Entry { file, runs, shift } => {
201                for run in runs {
202                    let met = run.extent().intersect(over.shifted(-shift));
203                    if met.is_empty() {
204                        continue;
205                    }
206                    let (from, to) = codec::chunks_over(run, met);
207                    let (at, len) = codec::span_of(run, from, to);
208                    let bytes = self.backend.get_range(&name_of(*file), at, len).await;
209                    let mut samples = codec::read_chunks(&bytes.ok()??, run, from, to)?;
210                    samples.start += shift;
211                    out.push(samples);
212                }
213            }
214            Samples::Staged { chunks, shift, .. } => {
215                for (name, e) in chunks {
216                    if !e.intersect(over.shifted(-shift)).is_empty() {
217                        let bytes = self.staging.get(name).await.ok()??;
218                        let mut samples = codec::read_chunk(&bytes)?;
219                        samples.start += shift;
220                        out.push(samples);
221                    }
222                }
223            }
224        }
225        Some(out)
226    }
227
228    pub(crate) async fn stage(&self, key: Hash, samples: &Buffer) -> Result<(), String> {
229        let n = {
230            let mut written = locked(&self.written);
231            *written += 1;
232            *written
233        };
234        let name = staged_name(key, &n.to_string());
235        self.staging.put(&name, &codec::chunk(samples)).await?;
236        locked(&self.staged)
237            .entry(key)
238            .or_default()
239            .chunks
240            .push((name, samples.extent()));
241        Ok(())
242    }
243
244    pub(crate) async fn stage_meta(&self, key: Hash, head: &Header) -> Result<(), String> {
245        let name = staged_name(key, META);
246        self.staging.put(&name, &codec::entry(head, &[])).await?;
247        locked(&self.staged).entry(key).or_default().meta = true;
248        Ok(())
249    }
250
251    /// Free while nothing is staged or used since and the budget holds. The index file is
252    /// every store's account and names each value before it moves: one cut short leaves a
253    /// gone entry named, never one unnamed. A value staged without its meta is dropped; one
254    /// failed or over an open entry stays staged.
255    pub(crate) async fn persist(&self) -> Result<Persisted, String> {
256        let idle = {
257            let index = locked(&self.index);
258            index.bytes <= self.max_bytes && index.unsynced().is_empty()
259        };
260        if idle && locked(&self.staged).is_empty() {
261            return Ok(Persisted::default());
262        }
263        let _held = self.backend.lock().await?;
264        let (mut index, seen) = self.current().await?;
265        let mut done = Persisted::default();
266        let mut moving = Vec::new();
267        let keys: Vec<Hash> = locked(&self.staged).keys().copied().collect();
268        for key in keys {
269            let Some(held) = locked(&self.staged).get(&key).cloned() else {
270                continue;
271            };
272            match self.prepared(key, &held).await? {
273                Prepared::Ready { bytes } => {
274                    moving.push((key, held, bytes));
275                    continue;
276                }
277                Prepared::Dropped => {}
278                Prepared::Refused => done.refused += 1,
279            }
280            self.unstaged(key, &held).await?;
281        }
282        if !moving.is_empty() {
283            let mut naming = index.clone();
284            for (key, _, bytes) in &moving {
285                naming.touch(*key, *bytes);
286            }
287            let text = naming.text(&version());
288            self.backend.put(INDEX_NAME, text.as_bytes()).await?;
289        }
290        for (key, held, bytes) in moving {
291            if !self.staging.rename(&name_of(key), &self.backend).await? {
292                continue;
293            }
294            index.touch(key, bytes);
295            done.written += 1;
296            self.unstaged(key, &held).await?;
297        }
298        for key in index.oldest_first() {
299            if index.bytes <= self.max_bytes {
300                break;
301            }
302            if self.backend.delete(&name_of(key)).await? {
303                index.forget(key);
304                done.evicted += 1;
305            }
306        }
307        self.backend
308            .put(INDEX_NAME, index.text(&version()).as_bytes())
309            .await?;
310        index.synced();
311        let mut local = locked(&self.index);
312        for key in local.used_since(seen) {
313            index.used(key);
314        }
315        *local = index;
316        Ok(done)
317    }
318
319    /// The index file, or else the directory listed, with this store's uses since; its clock.
320    async fn current(&self) -> Result<(Index, u64), String> {
321        let found = self.backend.get(INDEX_NAME).await?;
322        let mut index = match found.and_then(|text| Index::read(&text, &version())) {
323            Some(index) => index,
324            None => {
325                let listed = self.backend.list().await?;
326                let listed = listed
327                    .into_iter()
328                    .filter_map(|(name, bytes)| Some((key_of(&name)?, bytes)));
329                Index::listed(listed.collect(), &locked(&self.index))
330            }
331        };
332        let local = locked(&self.index);
333        for key in local.unsynced() {
334            index.used(key);
335        }
336        Ok((index, local.clock()))
337    }
338
339    async fn unstaged(&self, key: Hash, held: &Staged) -> Result<(), String> {
340        for (chunk, _) in &held.chunks {
341            self.staging.delete(chunk).await?;
342        }
343        self.staging.delete(&staged_name(key, META)).await?;
344        locked(&self.staged).remove(&key);
345        Ok(())
346    }
347
348    /// `key`'s staged value joined with its entry's, whole in staging under the entry's name.
349    async fn prepared(&self, key: Hash, held: &Staged) -> Result<Prepared, String> {
350        let whole = match held.meta {
351            true => self.staged_value_of(key, held).await?,
352            false => None,
353        };
354        let whole = whole.filter(|(head, runs)| !runs.is_empty() || head.refers());
355        let Some((head, mut runs)) = whole else {
356            return Ok(Prepared::Dropped);
357        };
358        let name = name_of(key);
359        if !head.refers() {
360            let held = self.backend.get(&name).await?;
361            if let Some(held) = held.and_then(|bytes| codec::read_runs(&bytes, key)) {
362                runs = joined_runs(runs, held);
363            }
364        }
365        let bytes = codec::entry(&head, &runs);
366        if bytes.len() as u64 > self.max_bytes {
367            return Ok(Prepared::Refused);
368        }
369        self.staging.put(&name, &bytes).await?;
370        Ok(Prepared::Ready {
371            bytes: bytes.len() as u64,
372        })
373    }
374
375    async fn staged_value_of(
376        &self,
377        key: Hash,
378        held: &Staged,
379    ) -> Result<Option<(Header, Vec<Buffer>)>, String> {
380        let Some(meta) = self.staging.get(&staged_name(key, META)).await? else {
381            return Ok(None);
382        };
383        let Some((head, _)) = codec::read_head(&meta, key) else {
384            return Ok(None);
385        };
386        let mut runs = Vec::new();
387        for (chunk, _) in &held.chunks {
388            let bytes = self.staging.get(chunk).await?;
389            match bytes.as_deref().and_then(codec::read_chunk) {
390                Some(samples) => joined(&mut runs, vec![Arc::new(samples)]),
391                None => return Ok(None),
392            }
393        }
394        let runs = runs.into_iter().map(Arc::unwrap_or_clone).collect();
395        Ok(Some((head, runs)))
396    }
397}
398
399fn joined_runs(staged: Vec<Buffer>, held: Vec<Buffer>) -> Vec<Buffer> {
400    let mut runs: Vec<Arc<Buffer>> = held.into_iter().map(Arc::new).collect();
401    joined(&mut runs, staged.into_iter().map(Arc::new).collect());
402    runs.into_iter().map(Arc::unwrap_or_clone).collect()
403}
404
405fn staged_name(key: Hash, part: &str) -> String {
406    format!("{}.{part}", name_of(key))
407}