1use std::collections::BTreeMap;
4use std::sync::{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::Samples;
13pub use super::stored::Stored;
14
15pub const DEFAULT_STORE_BYTES: u64 = 2 << 30;
16
17fn version() -> String {
18 format!("sva store format {STORE_FORMAT}")
19}
20
21pub const INDEX_NAME: &str = "index";
22
23const RETIRED: [&str; 2] = ["version", "recency"];
25
26const META: &str = "meta";
27
28pub trait Backend: Sized {
30 type Lock;
31 fn lock(&self) -> impl Future<Output = Result<Self::Lock, String>>;
33 fn get(&self, name: &str) -> impl Future<Output = Result<Option<Vec<u8>>, String>>;
34 fn get_range(
36 &self,
37 name: &str,
38 from: u64,
39 len: u64,
40 ) -> impl Future<Output = Result<Option<Vec<u8>>, String>>;
41 fn put(&self, name: &str, bytes: &[u8]) -> impl Future<Output = Result<(), String>>;
42 fn delete(&self, name: &str) -> impl Future<Output = Result<bool, String>>;
44 fn list(&self) -> impl Future<Output = Result<Vec<(String, u64)>, String>>;
46 fn staging(&self) -> impl Future<Output = Result<Self, String>>;
48 fn rename(&self, name: &str, to: &Self) -> impl Future<Output = Result<bool, String>>;
51}
52
53pub trait Through {
55 fn lookup(&self, key: Hash) -> impl Future<Output = Option<Stored>>;
56 fn read(&self, stored: &Stored, over: Extent) -> impl Future<Output = Option<Vec<Buffer>>>;
58 fn epoch(&self) -> u64;
60}
61
62impl<B: Backend> Through for Store<B> {
63 fn lookup(&self, key: Hash) -> impl Future<Output = Option<Stored>> {
64 Store::lookup(self, key)
65 }
66
67 fn read(&self, stored: &Stored, over: Extent) -> impl Future<Output = Option<Vec<Buffer>>> {
68 Store::read(self, stored, over)
69 }
70
71 fn epoch(&self) -> u64 {
72 *locked(&self.epoch)
73 }
74}
75
76pub struct NoStore;
77
78impl Through for NoStore {
79 async fn lookup(&self, _: Hash) -> Option<Stored> {
80 None
81 }
82
83 async fn read(&self, _: &Stored, _: Extent) -> Option<Vec<Buffer>> {
84 None
85 }
86
87 fn epoch(&self) -> u64 {
88 0
89 }
90}
91
92#[derive(Clone, Default)]
93struct Staged {
94 chunks: Vec<(String, Extent)>,
95 meta: bool,
96}
97
98pub struct Store<B> {
102 backend: B,
103 staging: B,
104 max_bytes: u64,
105 index: Mutex<Index>,
106 staged: Mutex<BTreeMap<Hash, Staged>>,
107 written: Mutex<u64>,
108 epoch: Mutex<u64>,
109}
110
111enum Committed {
112 Moved { bytes: u64 },
113 Busy,
114 Dropped,
115}
116
117#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
118pub struct Persisted {
119 pub written: usize,
120 pub evicted: usize,
121}
122
123fn locked<T>(held: &Mutex<T>) -> MutexGuard<'_, T> {
124 held.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
125}
126
127impl<B: Backend> Store<B> {
128 pub async fn open(backend: B, max_bytes: u64) -> Result<Store<B>, String> {
130 let version = version();
131 let held = backend.lock().await?;
132 let found = backend.get(INDEX_NAME).await?;
133 let index = match found.and_then(|text| Index::read(&text, &version)) {
134 Some(index) => index,
135 None => {
136 for (name, _) in backend.list().await? {
137 if key_of(&name).is_some() || RETIRED.contains(&name.as_str()) {
138 backend.delete(&name).await?;
139 }
140 }
141 let empty = Index::default();
142 backend
143 .put(INDEX_NAME, empty.text(&version).as_bytes())
144 .await?;
145 empty
146 }
147 };
148 let staging = backend.staging().await?;
149 drop(held);
150 Ok(Store {
151 backend,
152 staging,
153 max_bytes,
154 index: Mutex::new(index),
155 staged: Mutex::new(BTreeMap::new()),
156 written: Mutex::new(0),
157 epoch: Mutex::new(0),
158 })
159 }
160
161 pub fn max_bytes(&self) -> u64 {
162 self.max_bytes
163 }
164
165 pub fn bytes(&self) -> u64 {
166 locked(&self.index).bytes
167 }
168
169 pub fn holds(&self, key: Hash) -> bool {
170 locked(&self.index).held.contains_key(&key)
171 }
172
173 pub(crate) async fn lookup(&self, key: Hash) -> Option<Stored> {
175 let found = self.written(key).await?;
176 let Samples::Of { key: of, by } = found.samples else {
177 return Some(found);
178 };
179 let samples = match self.written(of).await?.samples {
180 Samples::Entry { file, runs, .. } => Samples::Entry {
181 file,
182 runs,
183 shift: -by,
184 },
185 Samples::Staged { chunks, .. } => Samples::Staged { chunks, shift: -by },
186 _ => return None,
187 };
188 Some(Stored { samples, ..found })
189 }
190
191 async fn written(&self, key: Hash) -> Option<Stored> {
193 if let Some(staged) = self.staged_value(key).await {
194 return Some(staged);
195 }
196 let name = name_of(key);
197 let Some(mut bytes) = self.backend.get_range(&name, 0, 8).await.ok()? else {
198 locked(&self.index).forget(key);
199 return None;
200 };
201 let rest = codec::head_len(&bytes)? - 8;
202 bytes.extend(self.backend.get_range(&name, 8, rest as u64).await.ok()??);
203 let (found, len) = codec::read_head(&bytes, key).filter(|(found, _)| found.key == key)?;
204 let mut index = locked(&self.index);
205 match index.held.contains_key(&key) {
206 true => index.used(key),
207 false => index.touch(key, len),
208 }
209 Some(found)
210 }
211
212 async fn staged_value(&self, key: Hash) -> Option<Stored> {
213 let chunks = {
214 let staged = locked(&self.staged);
215 let held = staged.get(&key).filter(|staged| staged.meta)?;
216 held.chunks.clone()
217 };
218 let meta = self.staging.get(&staged_name(key, META)).await.ok()??;
219 let (mut stored, _) = codec::read_head(&meta, key)?;
220 if !stored.refers() {
221 stored.samples = Samples::Staged { chunks, shift: 0 };
222 }
223 Some(stored)
224 }
225
226 pub(crate) async fn read(&self, stored: &Stored, over: Extent) -> Option<Vec<Buffer>> {
227 let mut out = Vec::new();
228 match &stored.samples {
229 Samples::None | Samples::Of { .. } => {}
230 Samples::Entry { file, runs, shift } => {
231 for run in runs {
232 let met = run.extent().intersect(over.shifted(-shift));
233 if met.is_empty() {
234 continue;
235 }
236 let chunk = codec::CHUNK as i64;
237 let from = ((met.start - run.start) / chunk) as usize;
238 let to = (met.end - run.start).div_euclid(chunk) as usize;
239 let to = to + usize::from((met.end - run.start) % chunk != 0);
240 let (at, len) = codec::span_of(run, from, to);
241 let bytes = self.backend.get_range(&name_of(*file), at, len).await;
242 let mut samples = codec::read_chunks(&bytes.ok()??, run, from, to)?;
243 samples.start += shift;
244 out.push(samples);
245 }
246 }
247 Samples::Staged { chunks, shift } => {
248 for (name, e) in chunks {
249 if !e.intersect(over.shifted(-shift)).is_empty() {
250 let bytes = self.staging.get(name).await.ok()??;
251 let mut samples = codec::read_chunk(&bytes)?;
252 samples.start += shift;
253 out.push(samples);
254 }
255 }
256 }
257 }
258 Some(out)
259 }
260
261 pub(crate) async fn stage(&self, key: Hash, samples: &Buffer) -> Result<(), String> {
262 let n = {
263 let mut written = locked(&self.written);
264 *written += 1;
265 *written
266 };
267 let name = staged_name(key, &n.to_string());
268 self.staging.put(&name, &codec::chunk(samples)).await?;
269 locked(&self.staged)
270 .entry(key)
271 .or_default()
272 .chunks
273 .push((name, samples.extent()));
274 Ok(())
275 }
276
277 pub(crate) async fn stage_meta(&self, key: Hash, stored: &Stored) -> Result<(), String> {
279 let name = staged_name(key, META);
280 self.staging.put(&name, &codec::entry(stored, &[])).await?;
281 locked(&self.staged).entry(key).or_default().meta = true;
282 *locked(&self.epoch) += 1;
283 Ok(())
284 }
285
286 pub async fn persist(&self) -> Result<Persisted, String> {
290 let done = self.committed_all().await;
291 *locked(&self.epoch) += 1;
292 done
293 }
294
295 async fn committed_all(&self) -> Result<Persisted, String> {
296 let _held = self.backend.lock().await?;
297 let mut done = Persisted::default();
298 let keys: Vec<Hash> = locked(&self.staged).keys().copied().collect();
299 for key in keys {
300 let Some(held) = locked(&self.staged).get(&key).cloned() else {
301 continue;
302 };
303 match self.committed(key, &held).await? {
304 Committed::Busy => continue,
305 Committed::Moved { bytes } => {
306 locked(&self.index).touch(key, bytes);
307 done.written += 1;
308 }
309 Committed::Dropped => {}
310 }
311 for (chunk, _) in &held.chunks {
312 self.staging.delete(chunk).await?;
313 }
314 self.staging.delete(&staged_name(key, META)).await?;
315 locked(&self.staged).remove(&key);
316 }
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 locked(&self.index).sync(listed.collect());
322 let oldest = locked(&self.index).oldest_first();
323 for key in oldest {
324 if self.bytes() <= self.max_bytes {
325 break;
326 }
327 if self.backend.delete(&name_of(key)).await? {
328 locked(&self.index).forget(key);
329 done.evicted += 1;
330 }
331 }
332 let text = locked(&self.index).text(&version());
333 self.backend.put(INDEX_NAME, text.as_bytes()).await?;
334 Ok(done)
335 }
336
337 async fn committed(&self, key: Hash, held: &Staged) -> Result<Committed, String> {
339 let whole = match held.meta {
340 true => self.staged_value_of(key, held).await?,
341 false => None,
342 };
343 let whole = whole.filter(|(stored, runs)| !runs.is_empty() || stored.refers());
344 let Some((stored, runs)) = whole else {
345 return Ok(Committed::Dropped);
346 };
347 let bytes = codec::entry(&stored, &runs);
348 if bytes.len() as u64 > self.max_bytes {
349 return Ok(Committed::Dropped);
350 }
351 let name = name_of(key);
352 self.staging.put(&name, &bytes).await?;
353 Ok(match self.staging.rename(&name, &self.backend).await? {
354 true => Committed::Moved {
355 bytes: bytes.len() as u64,
356 },
357 false => Committed::Busy,
358 })
359 }
360
361 async fn staged_value_of(
362 &self,
363 key: Hash,
364 held: &Staged,
365 ) -> Result<Option<(Stored, Vec<Buffer>)>, String> {
366 let Some(meta) = self.staging.get(&staged_name(key, META)).await? else {
367 return Ok(None);
368 };
369 let Some((stored, _)) = codec::read_head(&meta, key) else {
370 return Ok(None);
371 };
372 let mut runs = Vec::new();
373 for (chunk, _) in &held.chunks {
374 let bytes = self.staging.get(chunk).await?;
375 match bytes.as_deref().and_then(codec::read_chunk) {
376 Some(samples) => joined(&mut runs, vec![samples]),
377 None => return Ok(None),
378 }
379 }
380 Ok(Some((stored, runs)))
381 }
382}
383
384fn staged_name(key: Hash, part: &str) -> String {
385 format!("{}.{part}", name_of(key))
386}