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 Committed {
71 Moved { bytes: u64 },
72 Busy,
73 Dropped,
74 Refused,
75}
76
77#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
78pub struct Persisted {
79 pub written: usize,
80 pub evicted: usize,
81 pub refused: usize,
83}
84
85fn locked<T>(held: &Mutex<T>) -> MutexGuard<'_, T> {
86 held.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
87}
88
89impl<B: Backend> Store<B> {
90 pub async fn open(backend: B, max_bytes: u64) -> Result<Store<B>, String> {
92 let version = version();
93 let held = backend.lock().await?;
94 let found = backend.get(INDEX_NAME).await?;
95 let index = match found.and_then(|text| Index::read(&text, &version)) {
96 Some(index) => index,
97 None => {
98 for (name, _) in backend.list().await? {
99 if key_of(&name).is_some() || RETIRED.contains(&name.as_str()) {
100 backend.delete(&name).await?;
101 }
102 }
103 let empty = Index::default();
104 backend
105 .put(INDEX_NAME, empty.text(&version).as_bytes())
106 .await?;
107 empty
108 }
109 };
110 let staging = backend.staging().await?;
111 drop(held);
112 Ok(Store {
113 backend,
114 staging,
115 max_bytes,
116 index: Mutex::new(index),
117 staged: Mutex::new(BTreeMap::new()),
118 written: Mutex::new(0),
119 })
120 }
121
122 pub fn max_bytes(&self) -> u64 {
123 self.max_bytes
124 }
125
126 pub fn bytes(&self) -> u64 {
127 locked(&self.index).bytes
128 }
129
130 pub fn holds(&self, key: Hash) -> bool {
131 locked(&self.index).held.contains_key(&key)
132 }
133
134 pub(crate) async fn lookup(&self, key: Hash) -> Option<Header> {
136 let found = self.written(key).await?;
137 let &Samples::Of { key: of, by } = found.samples() else {
138 return Some(found);
139 };
140 let samples = match self.written(of).await?.into_parts().1 {
141 Samples::Entry { file, runs, .. } => Samples::Entry {
142 file,
143 runs,
144 shift: -by,
145 },
146 Samples::Staged { chunks, .. } => Samples::Staged { chunks, shift: -by },
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 => Header::new(found.into_parts().0, Samples::Staged { chunks, shift: 0 }),
185 })
186 }
187
188 pub(crate) async fn read(&self, head: &Header, over: Extent) -> Option<Vec<Buffer>> {
190 let mut out = Vec::new();
191 match head.samples() {
192 Samples::None | Samples::Of { .. } => {}
193 Samples::Entry { file, runs, shift } => {
194 for run in runs {
195 let met = run.extent().intersect(over.shifted(-shift));
196 if met.is_empty() {
197 continue;
198 }
199 let chunk = codec::CHUNK as i64;
200 let from = ((met.start - run.start) / chunk) as usize;
201 let to = (met.end - run.start).div_euclid(chunk) as usize;
202 let to = to + usize::from((met.end - run.start) % chunk != 0);
203 let (at, len) = codec::span_of(run, from, to);
204 let bytes = self.backend.get_range(&name_of(*file), at, len).await;
205 let mut samples = codec::read_chunks(&bytes.ok()??, run, from, to)?;
206 samples.start += shift;
207 out.push(samples);
208 }
209 }
210 Samples::Staged { chunks, shift } => {
211 for (name, e) in chunks {
212 if !e.intersect(over.shifted(-shift)).is_empty() {
213 let bytes = self.staging.get(name).await.ok()??;
214 let mut samples = codec::read_chunk(&bytes)?;
215 samples.start += shift;
216 out.push(samples);
217 }
218 }
219 }
220 }
221 Some(out)
222 }
223
224 pub(crate) async fn stage(&self, key: Hash, samples: &Buffer) -> Result<(), String> {
225 let n = {
226 let mut written = locked(&self.written);
227 *written += 1;
228 *written
229 };
230 let name = staged_name(key, &n.to_string());
231 self.staging.put(&name, &codec::chunk(samples)).await?;
232 locked(&self.staged)
233 .entry(key)
234 .or_default()
235 .chunks
236 .push((name, samples.extent()));
237 Ok(())
238 }
239
240 pub(crate) async fn stage_meta(&self, key: Hash, head: &Header) -> Result<(), String> {
241 let name = staged_name(key, META);
242 self.staging.put(&name, &codec::entry(head, &[])).await?;
243 locked(&self.staged).entry(key).or_default().meta = true;
244 Ok(())
245 }
246
247 pub(crate) async fn persist(&self) -> Result<Persisted, String> {
251 let _held = self.backend.lock().await?;
252 let mut done = Persisted::default();
253 let keys: Vec<Hash> = locked(&self.staged).keys().copied().collect();
254 for key in keys {
255 let Some(held) = locked(&self.staged).get(&key).cloned() else {
256 continue;
257 };
258 match self.committed(key, &held).await? {
259 Committed::Busy => continue,
260 Committed::Moved { bytes } => {
261 locked(&self.index).touch(key, bytes);
262 done.written += 1;
263 }
264 Committed::Dropped => {}
265 Committed::Refused => done.refused += 1,
266 }
267 for (chunk, _) in &held.chunks {
268 self.staging.delete(chunk).await?;
269 }
270 self.staging.delete(&staged_name(key, META)).await?;
271 locked(&self.staged).remove(&key);
272 }
273 let listed = self.backend.list().await?;
274 let listed = listed
275 .into_iter()
276 .filter_map(|(name, bytes)| Some((key_of(&name)?, bytes)));
277 locked(&self.index).sync(listed.collect());
278 let oldest = locked(&self.index).oldest_first();
279 for key in oldest {
280 if self.bytes() <= self.max_bytes {
281 break;
282 }
283 if self.backend.delete(&name_of(key)).await? {
284 locked(&self.index).forget(key);
285 done.evicted += 1;
286 }
287 }
288 let text = locked(&self.index).text(&version());
289 self.backend.put(INDEX_NAME, text.as_bytes()).await?;
290 Ok(done)
291 }
292
293 async fn committed(&self, key: Hash, held: &Staged) -> Result<Committed, String> {
295 let whole = match held.meta {
296 true => self.staged_value_of(key, held).await?,
297 false => None,
298 };
299 let whole = whole.filter(|(head, runs)| !runs.is_empty() || head.refers());
300 let Some((head, mut runs)) = whole else {
301 return Ok(Committed::Dropped);
302 };
303 let name = name_of(key);
304 if !head.refers() {
305 let held = self.backend.get(&name).await?;
306 if let Some(held) = held.and_then(|bytes| codec::read_runs(&bytes, key)) {
307 runs = joined_runs(runs, held);
308 }
309 }
310 let bytes = codec::entry(&head, &runs);
311 if bytes.len() as u64 > self.max_bytes {
312 return Ok(Committed::Refused);
313 }
314 self.staging.put(&name, &bytes).await?;
315 Ok(match self.staging.rename(&name, &self.backend).await? {
316 true => Committed::Moved {
317 bytes: bytes.len() as u64,
318 },
319 false => Committed::Busy,
320 })
321 }
322
323 async fn staged_value_of(
324 &self,
325 key: Hash,
326 held: &Staged,
327 ) -> Result<Option<(Header, Vec<Buffer>)>, String> {
328 let Some(meta) = self.staging.get(&staged_name(key, META)).await? else {
329 return Ok(None);
330 };
331 let Some((head, _)) = codec::read_head(&meta, key) else {
332 return Ok(None);
333 };
334 let mut runs = Vec::new();
335 for (chunk, _) in &held.chunks {
336 let bytes = self.staging.get(chunk).await?;
337 match bytes.as_deref().and_then(codec::read_chunk) {
338 Some(samples) => joined(&mut runs, vec![Arc::new(samples)]),
339 None => return Ok(None),
340 }
341 }
342 let runs = runs.into_iter().map(Arc::unwrap_or_clone).collect();
343 Ok(Some((head, runs)))
344 }
345}
346
347fn joined_runs(staged: Vec<Buffer>, held: Vec<Buffer>) -> Vec<Buffer> {
348 let mut runs: Vec<Arc<Buffer>> = held.into_iter().map(Arc::new).collect();
349 joined(&mut runs, staged.into_iter().map(Arc::new).collect());
350 runs.into_iter().map(Arc::unwrap_or_clone).collect()
351}
352
353fn staged_name(key: Hash, part: &str) -> String {
354 format!("{}.{part}", name_of(key))
355}