1use std::collections::HashMap;
4use std::sync::atomic::{AtomicU64, Ordering};
5use std::sync::{Mutex, MutexGuard, PoisonError};
6use std::time::Duration;
7
8use sva_formula::{Hash, Lanes};
9use sva_samples::{FilterTrace, Label};
10
11use super::entry_bytes::{self, RawF64, SampleCodec, Unread};
12use super::evict;
13use super::{Cache, Entry, Expected, Payload, PayloadKind};
14
15pub const FORMAT: u32 = 1;
17
18const MAGIC: u32 = u32::from_le_bytes(*b"SVP1");
19const HEADER: usize = 4 + 4 + 16 + 8;
21const CHECKSUM_ROTATE: u32 = 29;
22
23pub trait Medium: Sync {
26 fn size(&self) -> u64;
27
28 fn read_at(&self, off: u64, buf: &mut [u8]) -> bool;
29
30 fn write_at(&self, off: u64, bytes: &[u8]) -> bool;
31
32 fn truncate(&self, len: u64) -> bool;
33
34 fn flush(&self) -> bool;
35}
36
37#[derive(Default)]
39pub struct VecMedium(Mutex<Vec<u8>>);
40
41impl VecMedium {
42 pub fn holding(bytes: Vec<u8>) -> VecMedium {
43 VecMedium(Mutex::new(bytes))
44 }
45
46 pub fn bytes(&self) -> Vec<u8> {
47 self.held().clone()
48 }
49
50 fn held(&self) -> MutexGuard<'_, Vec<u8>> {
51 self.0.lock().unwrap_or_else(PoisonError::into_inner)
52 }
53}
54
55impl Medium for VecMedium {
56 fn size(&self) -> u64 {
57 self.held().len() as u64
58 }
59
60 fn read_at(&self, off: u64, buf: &mut [u8]) -> bool {
61 let held = self.held();
62 let Some(from) = held.get(off as usize..off as usize + buf.len()) else {
63 return false;
64 };
65 buf.copy_from_slice(from);
66 true
67 }
68
69 fn write_at(&self, off: u64, bytes: &[u8]) -> bool {
70 let mut held = self.held();
71 let end = off as usize + bytes.len();
72 if held.len() < end {
73 held.resize(end, 0);
74 }
75 held[off as usize..end].copy_from_slice(bytes);
76 true
77 }
78
79 fn truncate(&self, len: u64) -> bool {
80 self.held().resize(len as usize, 0);
81 true
82 }
83
84 fn flush(&self) -> bool {
85 true
86 }
87}
88
89#[derive(Clone, Copy)]
90struct Slot {
91 off: u64,
92 body: u32,
93 read: u64,
94}
95
96impl Slot {
97 fn bytes(self) -> u64 {
98 HEADER as u64 + u64::from(self.body)
99 }
100}
101
102struct State {
103 index: HashMap<Hash, Slot>,
104 end: u64,
105 clock: u64,
106 unflushed: bool,
108}
109
110pub struct Pack<M: Medium> {
115 medium: M,
116 state: Mutex<State>,
117 codec: Box<dyn SampleCodec>,
118 max_bytes: u64,
119 evicted: AtomicU64,
120 faults: AtomicU64,
121}
122
123impl<M: Medium> Pack<M> {
124 pub fn open(medium: M, max_bytes: u64) -> Pack<M> {
125 let mut state = State {
126 index: HashMap::new(),
127 end: 0,
128 clock: 0,
129 unflushed: false,
130 };
131 let len = medium.size();
132 while let Some((key, slot)) = record_at(&medium, state.end, len, state.clock) {
133 state.index.insert(key, slot);
134 state.end += slot.bytes();
135 state.clock += 1;
136 }
137 if state.end < len {
138 medium.truncate(state.end);
139 state.unflushed = true;
140 }
141 Pack {
142 medium,
143 state: Mutex::new(state),
144 codec: Box::new(RawF64),
145 max_bytes,
146 evicted: AtomicU64::new(0),
147 faults: AtomicU64::new(0),
148 }
149 }
150
151 pub fn coded(mut self, codec: Box<dyn SampleCodec>) -> Pack<M> {
152 self.codec = codec;
153 self
154 }
155
156 pub fn medium(&self) -> &M {
157 &self.medium
158 }
159
160 fn locked(&self) -> MutexGuard<'_, State> {
161 self.state.lock().unwrap_or_else(PoisonError::into_inner)
162 }
163
164 fn fault(&self) {
165 self.faults.fetch_add(1, Ordering::Relaxed);
166 }
167
168 fn compact(&self, state: &mut State) {
171 let mut order: Vec<(Hash, Slot)> = state.index.iter().map(|(k, s)| (*k, *s)).collect();
172 order.sort_by_key(|(_, slot)| slot.off);
173 let mut cursor = 0;
174 for (key, slot) in order {
175 if slot.off != cursor {
176 let mut record = vec![0; slot.bytes() as usize];
177 if !self.medium.read_at(slot.off, &mut record)
178 || !self.medium.write_at(cursor, &record)
179 {
180 self.fault();
181 state.index.remove(&key);
182 continue;
183 }
184 }
185 state.index.insert(
186 key,
187 Slot {
188 off: cursor,
189 ..slot
190 },
191 );
192 cursor += slot.bytes();
193 }
194 if !self.medium.truncate(cursor) {
195 self.fault();
196 }
197 state.end = cursor;
198 state.unflushed = true;
199 }
200}
201
202fn checksum(key: Hash, body: &[u8]) -> u64 {
203 let mut lanes = Lanes::<CHECKSUM_ROTATE>::default();
204 lanes.word(key.0);
205 lanes.word(key.1);
206 lanes.word(body.len() as u64);
207 for chunk in body.chunks(8) {
208 let mut word = [0; 8];
209 word[..chunk.len()].copy_from_slice(chunk);
210 lanes.word(u64::from_le_bytes(word));
211 }
212 let Hash(a, b) = lanes.finish();
213 a ^ b
214}
215
216fn header(key: Hash, body: &[u8]) -> [u8; HEADER] {
217 let mut out = [0; HEADER];
218 out[..4].copy_from_slice(&MAGIC.to_le_bytes());
219 out[4..8].copy_from_slice(&(body.len() as u32).to_le_bytes());
220 out[8..16].copy_from_slice(&key.0.to_le_bytes());
221 out[16..24].copy_from_slice(&key.1.to_le_bytes());
222 out[24..].copy_from_slice(&checksum(key, body).to_le_bytes());
223 out
224}
225
226fn word<const N: usize>(bytes: &[u8], at: usize) -> [u8; N] {
227 bytes[at..at + N].try_into().expect("a header field")
228}
229
230fn record_at(medium: &impl Medium, off: u64, len: u64, read: u64) -> Option<(Hash, Slot)> {
232 let (key, body) = read_record(medium, off, len)?;
233 Some((
234 key,
235 Slot {
236 off,
237 body: body.len() as u32,
238 read,
239 },
240 ))
241}
242
243fn read_record(medium: &impl Medium, off: u64, len: u64) -> Option<(Hash, Vec<u8>)> {
244 let mut head = [0; HEADER];
245 if off + HEADER as u64 > len || !medium.read_at(off, &mut head) {
246 return None;
247 }
248 if u32::from_le_bytes(word(&head, 0)) != MAGIC {
249 return None;
250 }
251 let body_len = u32::from_le_bytes(word(&head, 4));
252 if off + HEADER as u64 + u64::from(body_len) > len {
253 return None;
254 }
255 let key = Hash(
256 u64::from_le_bytes(word(&head, 8)),
257 u64::from_le_bytes(word(&head, 16)),
258 );
259 let mut body = vec![0; body_len as usize];
260 if !medium.read_at(off + HEADER as u64, &mut body) {
261 return None;
262 }
263 (checksum(key, &body) == u64::from_le_bytes(word(&head, 24))).then_some((key, body))
264}
265
266impl<M: Medium> Cache for Pack<M> {
267 fn max_bytes(&self) -> u64 {
268 self.max_bytes
269 }
270
271 fn held_bytes(&self) -> u64 {
273 self.locked().end
274 }
275
276 fn evicted_bytes(&self) -> u64 {
277 self.evicted.load(Ordering::Relaxed)
278 }
279
280 fn faults(&self) -> u64 {
281 self.faults.load(Ordering::Relaxed)
282 }
283
284 fn holds(&self, key: Hash) -> bool {
285 self.locked().index.contains_key(&key)
286 }
287
288 fn load(&self, key: Hash, node: &str, expected: Expected) -> Option<Entry> {
289 let Expected::Samples {
290 rate,
291 width,
292 samples,
293 } = expected
294 else {
295 return None;
296 };
297 let mut state = self.locked();
298 let slot = *state.index.get(&key)?;
299 let body = match read_record(&self.medium, slot.off, state.end) {
300 Some((found, body)) if found == key => body,
301 _ => {
302 self.fault();
303 state.index.remove(&key);
304 return None;
305 }
306 };
307 match entry_bytes::decode(&body, node, (rate, width, samples), self.codec.as_ref()) {
308 Ok(entry) => {
309 let read = state.clock;
310 state.clock += 1;
311 state.index.insert(key, Slot { read, ..slot });
312 Some(entry)
313 }
314 Err(Unread::OtherCodec) => None,
315 Err(Unread::Damaged) => {
316 self.fault();
317 state.index.remove(&key);
318 None
319 }
320 }
321 }
322
323 fn worth_storing(&self, _cost: Duration, _bytes: usize, kind: PayloadKind) -> bool {
325 kind == PayloadKind::Samples
326 }
327
328 fn store(&self, key: Hash, payload: &Payload, traces: &[FilterTrace], label: Option<&Label>) {
329 let Payload::Samples(buffer) = payload else {
330 return;
331 };
332 let body = entry_bytes::encode(buffer, traces, label, self.codec.as_ref());
333 let Ok(body_len) = u32::try_from(body.len()) else {
334 self.fault();
335 return;
336 };
337 let mut record = header(key, &body).to_vec();
338 record.extend_from_slice(&body);
339 let mut state = self.locked();
340 let off = state.end;
341 state.unflushed = true;
342 if !self.medium.write_at(off, &record) {
343 self.fault();
344 self.medium.truncate(off);
345 return;
346 }
347 let read = state.clock;
348 state.clock += 1;
349 state.end += record.len() as u64;
350 state.index.insert(
351 key,
352 Slot {
353 off,
354 body: body_len,
355 read,
356 },
357 );
358 }
359
360 fn sweep(&self) {
361 let mut state = self.locked();
362 let mut evicted = 0;
363 if state.end > self.max_bytes {
364 let order = state
365 .index
366 .iter()
367 .map(|(key, slot)| (slot.read, slot.bytes(), *key))
368 .collect();
369 let index = &mut state.index;
370 evicted =
371 evict::to_cap(order, self.max_bytes, |key| index.remove(key).is_some()).evicted;
372 self.compact(&mut state);
373 }
374 self.evicted.store(evicted, Ordering::Relaxed);
375 if std::mem::take(&mut state.unflushed) && !self.medium.flush() {
376 self.fault();
377 }
378 }
379}