1use std::sync::Arc;
4
5use sva_formula::Hash;
6use sva_samples::{Buffer, Extent};
7
8use super::memory::{CachePolicy, Counters, DEFAULT_CACHE_BYTES, Known, Memory, PrunePolicy};
9use super::persist::{Backend, Persisted, Store};
10
11pub 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
50pub const FETCH_READS: usize = 4;
52
53#[derive(Default)]
55pub(crate) struct Fetched {
56 pub(crate) handed: Vec<(Hash, Vec<Arc<Buffer>>)>,
57 pub(crate) left: Vec<(Hash, Extent)>,
58}
59
60pub 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 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 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 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 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 policy(&self) -> CachePolicy {
241 self.memory.policy()
242 }
243
244 pub fn set_policy(&self, policy: CachePolicy) {
245 self.memory.set_policy(policy);
246 }
247
248 pub fn prune(&self, policy: PrunePolicy) {
250 self.memory.prune(policy);
251 }
252
253 pub fn clear(&self) {
254 self.memory.clear();
255 }
256
257 pub fn mark_every(&self) -> usize {
258 self.memory.mark_every()
259 }
260
261 pub fn set_mark_every(&self, samples: usize) {
262 self.memory.set_mark_every(samples);
263 }
264}
265
266pub(crate) fn now<F: Future>(future: F) -> F::Output {
268 let mut future = std::pin::pin!(future);
269 let waker = std::task::Waker::noop();
270 match future
271 .as_mut()
272 .poll(&mut std::task::Context::from_waker(waker))
273 {
274 std::task::Poll::Ready(out) => out,
275 std::task::Poll::Pending => unreachable!("memory over nothing never waits"),
276 }
277}