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