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