Skip to main content

sva_engine/cache/
pack.rs

1// Concern: keeps rendered buffers as records appended to one file-like medium, across sessions | Non-concern: what the medium is, one entry's bytes (entry_bytes.rs) | IO: (Hash) -> a buffer + traces
2
3use 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
15/// Bumped with the record layout; a caller naming its file by it never opens an older one.
16pub const FORMAT: u32 = 1;
17
18const MAGIC: u32 = u32::from_le_bytes(*b"SVP1");
19/// magic, body length, key, checksum.
20const HEADER: usize = 4 + 4 + 16 + 8;
21const CHECKSUM_ROTATE: u32 = 29;
22
23/// Bytes at offsets, as one file offers them. A medium that fails answers `false`, and the pack
24/// reads that as a miss or a write that did not happen.
25pub 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/// A medium in this process's heap, for a caller that wants the pack's bytes themselves.
38#[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    /// A flush is dear on some media, so a sweep after nothing but hits touches none.
107    unflushed: bool,
108}
109
110/// Records `[magic u32][len u32][key 16B][checksum u64][entry]`, appended and never rewritten
111/// in place except by a sweep. Only samples are kept, as on disk. Opening reads every record,
112/// the later of two under one key winning, and cuts the medium at the first that does not
113/// check out: a torn tail is a miss, never a misread.
114pub 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    /// Survivors move toward the start in offset order, so each lands at or before where it
169    /// was and none is overwritten before it is read. Cut short, the tail simply fails to scan.
170    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
230/// The record at `off`, if it is whole and its checksum holds.
231fn 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    /// The medium's whole length, records a later one replaced included, until a sweep.
272    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    /// No clock enters it: a value a browser renders has no measured cost to weigh.
324    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}