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
88#[derive(Default)]
89struct Index {
90 held: HashMap<Hash, (u64, u64)>,
91 clock: u64,
92 bytes: u64,
93}
94
95impl Index {
96 fn touch(&mut self, key: Hash, bytes: u64) {
97 self.clock += 1;
98 let at = self.clock;
99 if let Some((old, _)) = self.held.insert(key, (bytes, at)) {
100 self.bytes -= old;
101 }
102 self.bytes += bytes;
103 }
104
105 fn read(&mut self, key: Hash) {
106 self.clock += 1;
107 let at = self.clock;
108 if let Some((_, held)) = self.held.get_mut(&key) {
109 *held = at;
110 }
111 }
112
113 fn forget(&mut self, key: Hash) {
114 if let Some((bytes, _)) = self.held.remove(&key) {
115 self.bytes -= bytes;
116 }
117 }
118
119 fn oldest_first(&self) -> Vec<Hash> {
120 let mut keys: Vec<(u64, Hash)> = self.held.iter().map(|(k, (_, at))| (*at, *k)).collect();
121 keys.sort_unstable();
122 keys.into_iter().map(|(_, k)| k).collect()
123 }
124}
125
126#[derive(Clone, Default)]
127struct Staged {
128 chunks: Vec<String>,
129 meta: bool,
130}
131
132pub struct Store<B> {
136 backend: B,
137 staging: B,
138 max_bytes: u64,
139 index: Mutex<Index>,
140 staged: Mutex<BTreeMap<Hash, Staged>>,
141 written: Mutex<u64>,
142}
143
144#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
145pub struct Persisted {
146 pub written: usize,
147 pub evicted: usize,
148}
149
150fn name_of(key: Hash) -> String {
151 format!("{:016x}{:016x}", key.0, key.1)
152}
153
154fn key_of(name: &str) -> Option<Hash> {
155 let hex = |s: &str| u64::from_str_radix(s, 16).ok();
156 let valid = name.len() == 32 && name.bytes().all(|b| b.is_ascii_hexdigit());
157 valid.then(|| Some(Hash(hex(&name[..16])?, hex(&name[16..])?)))?
158}
159
160fn locked<T>(held: &Mutex<T>) -> MutexGuard<'_, T> {
161 held.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
162}
163
164impl<B: Backend> Store<B> {
165 pub async fn open(backend: B, max_bytes: u64) -> Result<Store<B>, String> {
167 let version = backend.get(VERSION_NAME).await?;
168 if version.as_deref() != Some(STORE_VERSION.as_bytes()) {
169 for (name, _) in backend.list().await? {
170 if key_of(&name).is_some() || name == RECENCY_NAME || name == VERSION_NAME {
171 backend.delete(&name).await?;
172 }
173 }
174 backend.put(VERSION_NAME, STORE_VERSION.as_bytes()).await?;
175 }
176 let order = backend.get(RECENCY_NAME).await?.unwrap_or_default();
177 let order = String::from_utf8_lossy(&order);
178 let rank: HashMap<&str, usize> = order.lines().enumerate().map(|(k, n)| (n, k)).collect();
179 let mut held: Vec<(usize, Hash, u64)> = backend
180 .list()
181 .await?
182 .into_iter()
183 .filter_map(|(name, bytes)| {
184 let at = rank.get(name.as_str()).map_or(0, |k| k + 1);
185 Some((at, key_of(&name)?, bytes))
186 })
187 .collect();
188 held.sort_unstable();
189 let mut index = Index::default();
190 for (_, key, bytes) in held {
191 index.touch(key, bytes);
192 }
193 let staging = backend.staging().await?;
194 Ok(Store {
195 backend,
196 staging,
197 max_bytes,
198 index: Mutex::new(index),
199 staged: Mutex::new(BTreeMap::new()),
200 written: Mutex::new(0),
201 })
202 }
203
204 pub fn max_bytes(&self) -> u64 {
205 self.max_bytes
206 }
207
208 pub fn bytes(&self) -> u64 {
209 locked(&self.index).bytes
210 }
211
212 pub fn holds(&self, key: Hash) -> bool {
213 locked(&self.index).held.contains_key(&key)
214 }
215
216 pub(crate) async fn lookup(&self, key: Hash) -> Option<Stored> {
218 if let Some(staged) = self.staged_value(key).await {
219 return Some(staged);
220 }
221 if !self.holds(key) {
222 return None;
223 }
224 let bytes = self.backend.get(&name_of(key)).await.ok()??;
225 let found = codec::read_entry(&bytes).filter(|found| found.key == key)?;
226 locked(&self.index).read(key);
227 Some(found)
228 }
229
230 async fn staged_value(&self, key: Hash) -> Option<Stored> {
231 let chunks = {
232 let staged = locked(&self.staged);
233 let held = staged.get(&key).filter(|staged| staged.meta)?;
234 held.chunks.clone()
235 };
236 let meta = self.staging.get(&staged_name(key, META)).await.ok()??;
237 let mut stored = codec::read_entry(&meta)?;
238 for chunk in chunks {
239 let bytes = self.staging.get(&chunk).await.ok()??;
240 joined(&mut stored.segments, vec![codec::read_chunk(&bytes)?]);
241 }
242 Some(stored)
243 }
244
245 pub(crate) async fn stage(&self, key: Hash, samples: &Buffer) -> Result<(), String> {
247 let n = {
248 let mut written = locked(&self.written);
249 *written += 1;
250 *written
251 };
252 let name = staged_name(key, &n.to_string());
253 self.staging.put(&name, &codec::chunk(samples)).await?;
254 locked(&self.staged)
255 .entry(key)
256 .or_default()
257 .chunks
258 .push(name);
259 Ok(())
260 }
261
262 pub(crate) async fn stage_meta(&self, key: Hash, stored: &Stored) -> Result<(), String> {
264 let name = staged_name(key, META);
265 self.staging.put(&name, &codec::entry(stored)).await?;
266 locked(&self.staged).entry(key).or_default().meta = true;
267 Ok(())
268 }
269
270 pub async fn persist(&self) -> Result<Persisted, String> {
273 let mut done = Persisted::default();
274 let keys: Vec<Hash> = locked(&self.staged).keys().copied().collect();
275 for key in keys {
276 let Some(held) = locked(&self.staged).get(&key).cloned() else {
277 continue;
278 };
279 if let Some(bytes) = self.committed(key, &held).await? {
280 locked(&self.index).touch(key, bytes);
281 done.written += 1;
282 }
283 for chunk in &held.chunks {
284 self.staging.delete(chunk).await?;
285 }
286 self.staging.delete(&staged_name(key, META)).await?;
287 locked(&self.staged).remove(&key);
288 }
289 let oldest = locked(&self.index).oldest_first();
290 for key in oldest {
291 if self.bytes() <= self.max_bytes {
292 break;
293 }
294 self.backend.delete(&name_of(key)).await?;
295 locked(&self.index).forget(key);
296 done.evicted += 1;
297 }
298 let order: String = locked(&self.index)
299 .oldest_first()
300 .into_iter()
301 .map(|key| name_of(key) + "\n")
302 .collect();
303 self.backend.put(RECENCY_NAME, order.as_bytes()).await?;
304 Ok(done)
305 }
306
307 async fn committed(&self, key: Hash, held: &Staged) -> Result<Option<u64>, String> {
309 let whole = match held.meta {
310 true => self.staged_value_of(key, held).await?,
311 false => None,
312 };
313 let Some(stored) = whole.filter(|stored| !stored.segments.is_empty()) else {
314 return Ok(None);
315 };
316 let bytes = codec::entry(&stored);
317 if bytes.len() as u64 > self.max_bytes {
318 return Ok(None);
319 }
320 let name = name_of(key);
321 self.staging.put(&name, &bytes).await?;
322 self.staging.rename(&name, &self.backend).await?;
323 Ok(Some(bytes.len() as u64))
324 }
325
326 async fn staged_value_of(&self, key: Hash, held: &Staged) -> Result<Option<Stored>, String> {
327 let Some(meta) = self.staging.get(&staged_name(key, META)).await? else {
328 return Ok(None);
329 };
330 let Some(mut stored) = codec::read_entry(&meta) else {
331 return Ok(None);
332 };
333 for chunk in &held.chunks {
334 let bytes = self.staging.get(chunk).await?;
335 match bytes.as_deref().and_then(codec::read_chunk) {
336 Some(samples) => joined(&mut stored.segments, vec![samples]),
337 None => return Ok(None),
338 }
339 }
340 Ok(Some(stored))
341 }
342}
343
344fn staged_name(key: Hash, part: &str) -> String {
345 format!("{}.{part}", name_of(key))
346}