Skip to main content

sva_engine/cache/
tier.rs

1// Concern: memory over a disk or over nothing, and the one passage between them | Non-concern: what memory keeps, the disk's medium | IO: (keys, needs) -> answers, samples; persist() -> commits
2
3use std::sync::Arc;
4
5use sva_formula::Hash;
6use sva_samples::{Buffer, Extent};
7
8use super::memory::{Counters, DEFAULT_CACHE_BYTES, Known, Memory};
9use super::persist::{Backend, Persisted, Store};
10
11/// No disk: memory over nothing never waits.
12pub enum Nothing {}
13
14impl Backend for Nothing {
15    type Lock = ();
16
17    async fn lock(&self) -> Result<(), String> {
18        match *self {}
19    }
20
21    async fn get(&self, _: &str) -> Result<Option<Vec<u8>>, String> {
22        match *self {}
23    }
24
25    async fn get_range(&self, _: &str, _: u64, _: u64) -> Result<Option<Vec<u8>>, String> {
26        match *self {}
27    }
28
29    async fn put(&self, _: &str, _: &[u8]) -> Result<(), String> {
30        match *self {}
31    }
32
33    async fn delete(&self, _: &str) -> Result<bool, String> {
34        match *self {}
35    }
36
37    async fn list(&self) -> Result<Vec<(String, u64)>, String> {
38        match *self {}
39    }
40
41    async fn staging(&self) -> Result<Nothing, String> {
42        match *self {}
43    }
44
45    async fn rename(&self, _: &str, _: &Nothing) -> Result<bool, String> {
46        match *self {}
47    }
48}
49
50/// The most disk reads one fetch makes.
51pub const FETCH_READS: usize = 4;
52
53/// What one fetch handed out, and the needs its reads did not reach.
54#[derive(Default)]
55pub(crate) struct Fetched {
56    pub(crate) handed: Vec<(Hash, Vec<Arc<Buffer>>)>,
57    pub(crate) left: Vec<(Hash, Extent)>,
58}
59
60/// Memory over a disk, or over nothing. Evaluation talks to memory alone; only memory decides
61/// whether a lookup, a read or a write passes through to the disk.
62pub struct Tier<B = Nothing> {
63    memory: Memory,
64    disk: Option<Store<B>>,
65}
66
67impl Default for Tier {
68    fn default() -> Tier {
69        Tier::new(DEFAULT_CACHE_BYTES)
70    }
71}
72
73impl Tier {
74    pub fn new(max_bytes: u64) -> Tier {
75        Tier::alone(max_bytes)
76    }
77}
78
79impl<B: Backend> Tier<B> {
80    pub fn alone(max_bytes: u64) -> Tier<B> {
81        Tier {
82            memory: Memory::holding(max_bytes),
83            disk: None,
84        }
85    }
86
87    pub fn over(disk: Store<B>, max_bytes: u64) -> Tier<B> {
88        Tier {
89            memory: Memory::over_disk(max_bytes),
90            disk: Some(disk),
91        }
92    }
93
94    pub fn disk(&self) -> Option<&Store<B>> {
95        self.disk.as_ref()
96    }
97
98    pub(crate) fn memory(&self) -> &Memory {
99        &self.memory
100    }
101
102    pub(crate) fn begin(&self) -> u64 {
103        self.memory.begin()
104    }
105
106    /// `key` answered: from memory, else off the disk and resident from then on.
107    pub(crate) async fn lookup(&self, key: Hash, round: u64) {
108        if !matches!(self.memory.answer(key, round), Known::Unknown) {
109            return;
110        }
111        let found = match &self.disk {
112            Some(disk) => {
113                self.memory.count(|c| c.disk_lookups += 1);
114                disk.lookup(key).await
115            }
116            None => None,
117        };
118        match found {
119            Some(head) => self.memory.promote(head),
120            None => self.memory.miss(key, round),
121        }
122    }
123
124    /// What each need asks of a node memory holds, and off the disk what memory lacks, at most
125    /// `FETCH_READS` reads a call, the rest left for the next. Memory first writes back.
126    pub(crate) async fn fetch(&self, needs: &[(Hash, Extent)]) -> Fetched {
127        self.write_back().await;
128        let mut out = Fetched::default();
129        let mut reads = 0;
130        for (key, over) in needs {
131            let (mut parts, lacks) = self.memory.resident(*key, *over);
132            if let (Some(disk), Some((head, gap))) = (&self.disk, lacks) {
133                if reads == FETCH_READS {
134                    out.left.push((*key, *over));
135                    continue;
136                }
137                reads += 1;
138                let read = disk.read(&head, gap).await;
139                let samples = read.iter().flatten().map(|b| b.len() * b.width);
140                let bytes = (samples.sum::<usize>() * size_of::<f64>()) as u64;
141                self.memory.count(|c| {
142                    c.disk_reads += 1;
143                    c.disk_read_bytes += bytes;
144                });
145                match read {
146                    Some(read) => parts.extend(self.memory.promote_samples(*key, read)),
147                    None => self.memory.forget(*key),
148                }
149            }
150            if !parts.is_empty() {
151                out.handed.push((*key, parts));
152            }
153        }
154        out
155    }
156
157    /// What memory let go of while not yet on the disk, written there; a failure stops every
158    /// write until the next persist.
159    async fn write_back(&self) {
160        if let Some(disk) = &self.disk {
161            let pending = self.memory.pending();
162            if let Err((why, left)) = self.written(disk, pending).await {
163                self.memory.failed(why, left);
164            }
165        }
166    }
167
168    async fn written(
169        &self,
170        disk: &Store<B>,
171        mut pending: Vec<super::memory::Writeback>,
172    ) -> Result<(), (String, Vec<super::memory::Writeback>)> {
173        pending.reverse();
174        while let Some(back) = pending.pop() {
175            let mut staged = Ok(());
176            for part in &back.parts {
177                staged = disk.stage(back.key, part).await;
178                if staged.is_err() {
179                    break;
180                }
181            }
182            let staged = match staged {
183                Ok(()) => disk.stage_meta(back.key, &back.head).await,
184                Err(why) => Err(why),
185            };
186            if let Err(why) = staged {
187                pending.push(back);
188                pending.reverse();
189                return Err((why, pending));
190            }
191            self.memory.written();
192        }
193        Ok(())
194    }
195
196    /// Every node memory holds that the disk lacks, written, then committed there.
197    pub async fn persist(&self) -> Result<Persisted, String> {
198        let Some(disk) = &self.disk else {
199            return Ok(Persisted::default());
200        };
201        self.memory.flush();
202        let pending = self.memory.pending();
203        if let Err((why, left)) = self.written(disk, pending).await {
204            self.memory.failed(why.clone(), left);
205            return Err(why);
206        }
207        let done = disk.persist().await?;
208        self.memory.committed();
209        Ok(done)
210    }
211
212    pub fn counters(&self) -> Counters {
213        self.memory.counters()
214    }
215
216    pub fn max_bytes(&self) -> u64 {
217        self.memory.max_bytes()
218    }
219
220    pub fn set_max_bytes(&self, max_bytes: u64) {
221        self.memory.set_max_bytes(max_bytes);
222    }
223
224    pub fn bytes(&self) -> u64 {
225        self.memory.bytes()
226    }
227
228    pub fn entries(&self) -> usize {
229        self.memory.entries()
230    }
231
232    pub fn evictions(&self) -> u64 {
233        self.memory.counters().evictions()
234    }
235
236    pub fn holds(&self, key: Hash) -> bool {
237        self.memory.holds(key)
238    }
239
240    pub fn mark_every(&self) -> usize {
241        self.memory.mark_every()
242    }
243
244    pub fn set_mark_every(&self, samples: usize) {
245        self.memory.set_mark_every(samples);
246    }
247}
248
249/// A future over memory alone, finished in its one poll.
250pub(crate) fn now<F: Future>(future: F) -> F::Output {
251    let mut future = std::pin::pin!(future);
252    let waker = std::task::Waker::noop();
253    match future
254        .as_mut()
255        .poll(&mut std::task::Context::from_waker(waker))
256    {
257        std::task::Poll::Ready(out) => out,
258        std::task::Poll::Pending => unreachable!("memory over nothing never waits"),
259    }
260}