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