1use std::collections::{BTreeMap, HashMap};
4use std::sync::{Mutex, MutexGuard};
5
6use sva_formula::{Codomain, Hash};
7use sva_samples::{Buffer, Extent, Grid, Label};
8
9use super::{codec, joined};
10
11#[derive(Clone, Debug, PartialEq)]
13pub struct Stored {
14 pub key: Hash,
15 pub segments: Vec<Buffer>,
16 pub label: Label,
17 pub width: u8,
18 pub codomain: Codomain,
19 pub rate: Option<u32>,
20 pub grid: Grid,
21 pub support: Extent,
22 pub priced: u128,
24 pub moved: f64,
26 pub readable: bool,
28}
29
30impl Stored {
31 pub(crate) fn holds(&self, over: Extent) -> bool {
32 let mut from = over.start;
33 let mut parts: Vec<Extent> = self.segments.iter().map(Buffer::extent).collect();
34 parts.sort_by_key(|e| e.start);
35 for part in parts {
36 if from >= over.end || part.start > from {
37 break;
38 }
39 from = from.max(part.end);
40 }
41 from >= over.end
42 }
43}
44
45pub(crate) fn node_key(identity: Hash, rate: u32, profile: &sva_samples::Profile) -> Hash {
47 super::mixed(
48 identity,
49 &[
50 u64::from(rate),
51 profile.precision_bits as u64,
52 profile.ceiling_hz.to_bits(),
53 0x6e_6f_64_65_00_00_00_01,
54 ],
55 )
56}
57
58pub const DEFAULT_STORE_BYTES: u64 = 2 << 30;
59
60pub const STORE_VERSION: &str = concat!(
61 "sva-engine ",
62 env!("CARGO_PKG_VERSION"),
63 " build ",
64 env!("SVA_ENGINE_BUILD"),
65 " format 2"
66);
67
68pub const VERSION_NAME: &str = "version";
69
70const RECENCY_NAME: &str = "recency";
71
72const META: &str = "meta";
73
74pub trait Backend: Sized {
76 fn get(&self, name: &str) -> impl Future<Output = Result<Option<Vec<u8>>, String>>;
77 fn put(&self, name: &str, bytes: &[u8]) -> impl Future<Output = Result<(), String>>;
78 fn delete(&self, name: &str) -> impl Future<Output = Result<(), String>>;
79 fn list(&self) -> impl Future<Output = Result<Vec<(String, u64)>, String>>;
81 fn staging(&self) -> impl Future<Output = Result<Self, String>>;
84 fn rename(&self, name: &str, to: &Self) -> impl Future<Output = Result<(), String>>;
86}
87
88pub trait Through {
90 fn lookup(&self, key: Hash) -> impl Future<Output = Option<Stored>>;
91}
92
93impl<B: Backend> Through for Store<B> {
94 fn lookup(&self, key: Hash) -> impl Future<Output = Option<Stored>> {
95 Store::lookup(self, key)
96 }
97}
98
99pub struct NoStore;
100
101impl Through for NoStore {
102 async fn lookup(&self, _: Hash) -> Option<Stored> {
103 None
104 }
105}
106
107#[derive(Default)]
108struct Index {
109 held: HashMap<Hash, (u64, u64)>,
110 clock: u64,
111 bytes: u64,
112}
113
114impl Index {
115 fn touch(&mut self, key: Hash, bytes: u64) {
116 self.clock += 1;
117 let at = self.clock;
118 if let Some((old, _)) = self.held.insert(key, (bytes, at)) {
119 self.bytes -= old;
120 }
121 self.bytes += bytes;
122 }
123
124 fn read(&mut self, key: Hash) {
125 self.clock += 1;
126 let at = self.clock;
127 if let Some((_, held)) = self.held.get_mut(&key) {
128 *held = at;
129 }
130 }
131
132 fn sync(&mut self, listed: Vec<(Hash, u64)>) {
134 let held = std::mem::take(&mut self.held);
135 self.held = listed
136 .into_iter()
137 .map(|(key, bytes)| (key, (bytes, held.get(&key).map_or(0, |(_, at)| *at))))
138 .collect();
139 self.bytes = self.held.values().map(|(bytes, _)| bytes).sum();
140 }
141
142 fn forget(&mut self, key: Hash) {
143 if let Some((bytes, _)) = self.held.remove(&key) {
144 self.bytes -= bytes;
145 }
146 }
147
148 fn oldest_first(&self) -> Vec<Hash> {
149 let mut keys: Vec<(u64, Hash)> = self.held.iter().map(|(k, (_, at))| (*at, *k)).collect();
150 keys.sort_unstable();
151 keys.into_iter().map(|(_, k)| k).collect()
152 }
153}
154
155#[derive(Clone, Default)]
156struct Staged {
157 chunks: Vec<String>,
158 meta: bool,
159}
160
161pub struct Store<B> {
165 backend: B,
166 staging: B,
167 max_bytes: u64,
168 index: Mutex<Index>,
169 staged: Mutex<BTreeMap<Hash, Staged>>,
170 written: Mutex<u64>,
171}
172
173#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
174pub struct Persisted {
175 pub written: usize,
176 pub evicted: usize,
177}
178
179fn name_of(key: Hash) -> String {
180 format!("{:016x}{:016x}", key.0, key.1)
181}
182
183fn key_of(name: &str) -> Option<Hash> {
184 let hex = |s: &str| u64::from_str_radix(s, 16).ok();
185 let valid = name.len() == 32 && name.bytes().all(|b| b.is_ascii_hexdigit());
186 valid.then(|| Some(Hash(hex(&name[..16])?, hex(&name[16..])?)))?
187}
188
189fn locked<T>(held: &Mutex<T>) -> MutexGuard<'_, T> {
190 held.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
191}
192
193impl<B: Backend> Store<B> {
194 pub async fn open(backend: B, max_bytes: u64) -> Result<Store<B>, String> {
196 let version = backend.get(VERSION_NAME).await?;
197 if version.as_deref() != Some(STORE_VERSION.as_bytes()) {
198 for (name, _) in backend.list().await? {
199 if key_of(&name).is_some() || name == RECENCY_NAME || name == VERSION_NAME {
200 backend.delete(&name).await?;
201 }
202 }
203 backend.put(VERSION_NAME, STORE_VERSION.as_bytes()).await?;
204 }
205 let order = backend.get(RECENCY_NAME).await?.unwrap_or_default();
206 let order = String::from_utf8_lossy(&order);
207 let rank: HashMap<&str, usize> = order.lines().enumerate().map(|(k, n)| (n, k)).collect();
208 let mut held: Vec<(usize, Hash, u64)> = backend
209 .list()
210 .await?
211 .into_iter()
212 .filter_map(|(name, bytes)| {
213 let at = rank.get(name.as_str()).map_or(0, |k| k + 1);
214 Some((at, key_of(&name)?, bytes))
215 })
216 .collect();
217 held.sort_unstable();
218 let mut index = Index::default();
219 for (_, key, bytes) in held {
220 index.touch(key, bytes);
221 }
222 let staging = backend.staging().await?;
223 Ok(Store {
224 backend,
225 staging,
226 max_bytes,
227 index: Mutex::new(index),
228 staged: Mutex::new(BTreeMap::new()),
229 written: Mutex::new(0),
230 })
231 }
232
233 pub fn max_bytes(&self) -> u64 {
234 self.max_bytes
235 }
236
237 pub fn bytes(&self) -> u64 {
238 locked(&self.index).bytes
239 }
240
241 pub fn holds(&self, key: Hash) -> bool {
242 locked(&self.index).held.contains_key(&key)
243 }
244
245 pub(crate) async fn lookup(&self, key: Hash) -> Option<Stored> {
247 if let Some(staged) = self.staged_value(key).await {
248 return Some(staged);
249 }
250 let Some(bytes) = self.backend.get(&name_of(key)).await.ok()? else {
251 locked(&self.index).forget(key);
252 return None;
253 };
254 let found = codec::read_entry(&bytes).filter(|found| found.key == key)?;
255 let mut index = locked(&self.index);
256 match index.held.contains_key(&key) {
257 true => index.read(key),
258 false => index.touch(key, bytes.len() as u64),
259 }
260 Some(found)
261 }
262
263 async fn staged_value(&self, key: Hash) -> Option<Stored> {
264 let chunks = {
265 let staged = locked(&self.staged);
266 let held = staged.get(&key).filter(|staged| staged.meta)?;
267 held.chunks.clone()
268 };
269 let meta = self.staging.get(&staged_name(key, META)).await.ok()??;
270 let mut stored = codec::read_entry(&meta)?;
271 for chunk in chunks {
272 let bytes = self.staging.get(&chunk).await.ok()??;
273 joined(&mut stored.segments, vec![codec::read_chunk(&bytes)?]);
274 }
275 Some(stored)
276 }
277
278 pub(crate) async fn stage(&self, key: Hash, samples: &Buffer) -> Result<(), String> {
280 let n = {
281 let mut written = locked(&self.written);
282 *written += 1;
283 *written
284 };
285 let name = staged_name(key, &n.to_string());
286 self.staging.put(&name, &codec::chunk(samples)).await?;
287 locked(&self.staged)
288 .entry(key)
289 .or_default()
290 .chunks
291 .push(name);
292 Ok(())
293 }
294
295 pub(crate) async fn stage_meta(&self, key: Hash, stored: &Stored) -> Result<(), String> {
297 let name = staged_name(key, META);
298 self.staging.put(&name, &codec::entry(stored)).await?;
299 locked(&self.staged).entry(key).or_default().meta = true;
300 Ok(())
301 }
302
303 pub async fn persist(&self) -> Result<Persisted, String> {
306 let mut done = Persisted::default();
307 let keys: Vec<Hash> = locked(&self.staged).keys().copied().collect();
308 for key in keys {
309 let Some(held) = locked(&self.staged).get(&key).cloned() else {
310 continue;
311 };
312 if let Some(bytes) = self.committed(key, &held).await? {
313 locked(&self.index).touch(key, bytes);
314 done.written += 1;
315 }
316 for chunk in &held.chunks {
317 self.staging.delete(chunk).await?;
318 }
319 self.staging.delete(&staged_name(key, META)).await?;
320 locked(&self.staged).remove(&key);
321 }
322 let listed = self.backend.list().await?;
323 let listed = listed
324 .into_iter()
325 .filter_map(|(name, bytes)| Some((key_of(&name)?, bytes)));
326 locked(&self.index).sync(listed.collect());
327 let oldest = locked(&self.index).oldest_first();
328 for key in oldest {
329 if self.bytes() <= self.max_bytes {
330 break;
331 }
332 self.backend.delete(&name_of(key)).await?;
333 locked(&self.index).forget(key);
334 done.evicted += 1;
335 }
336 let order: String = locked(&self.index)
337 .oldest_first()
338 .into_iter()
339 .map(|key| name_of(key) + "\n")
340 .collect();
341 self.backend.put(RECENCY_NAME, order.as_bytes()).await?;
342 Ok(done)
343 }
344
345 async fn committed(&self, key: Hash, held: &Staged) -> Result<Option<u64>, String> {
347 let whole = match held.meta {
348 true => self.staged_value_of(key, held).await?,
349 false => None,
350 };
351 let Some(stored) = whole.filter(|stored| !stored.segments.is_empty()) else {
352 return Ok(None);
353 };
354 let bytes = codec::entry(&stored);
355 if bytes.len() as u64 > self.max_bytes {
356 return Ok(None);
357 }
358 let name = name_of(key);
359 self.staging.put(&name, &bytes).await?;
360 self.staging.rename(&name, &self.backend).await?;
361 Ok(Some(bytes.len() as u64))
362 }
363
364 async fn staged_value_of(&self, key: Hash, held: &Staged) -> Result<Option<Stored>, String> {
365 let Some(meta) = self.staging.get(&staged_name(key, META)).await? else {
366 return Ok(None);
367 };
368 let Some(mut stored) = codec::read_entry(&meta) else {
369 return Ok(None);
370 };
371 for chunk in &held.chunks {
372 let bytes = self.staging.get(chunk).await?;
373 match bytes.as_deref().and_then(codec::read_chunk) {
374 Some(samples) => joined(&mut stored.segments, vec![samples]),
375 None => return Ok(None),
376 }
377 }
378 Ok(Some(stored))
379 }
380}
381
382fn staged_name(key: Hash, part: &str) -> String {
383 format!("{}.{part}", name_of(key))
384}