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