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