1use 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
12pub 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
51pub const FETCH_READS: usize = 4;
53
54#[derive(Default)]
56pub(crate) struct Fetched {
57 pub(crate) handed: Vec<(Hash, Vec<Arc<Buffer>>)>,
58 pub(crate) left: Vec<(Hash, Extent)>,
59}
60
61pub struct Tier<B = Nothing> {
64 memory: Memory,
65 disk: Option<Store<B>>,
66 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 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 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 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 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
259pub(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}