Skip to main content

sva_engine/cache/
memory.rs

1// Concern: the memory tier: each resident value and its node under one cap, what it evicts and writes back, its misses | Non-concern: the disk | IO: (key) -> held, answered; (key, samples, node) -> kept
2
3#[cfg(test)]
4mod chains;
5
6use std::cell::RefCell;
7use std::collections::{HashMap, HashSet};
8use std::sync::{Arc, Mutex, MutexGuard};
9
10use sva_formula::Hash;
11use sva_samples::{Buffer, Extent, Label};
12
13use super::stats::{Outcome, Recording};
14use super::stored::{Header, Samples};
15use super::{Entry, Expected, Payload, PayloadKind, Run, Stored, joined};
16
17pub const DEFAULT_CACHE_BYTES: u64 = 2 << 30;
18
19/// Samples between two states a run keeps, so a reader resumes from one at most this far back.
20pub const DEFAULT_MARK_EVERY: usize = 16_384;
21
22/// What passed between memory and the disk beneath it, and what memory let go.
23#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
24pub struct Counters {
25    pub disk_lookups: u64,
26    pub disk_reads: u64,
27    pub disk_read_bytes: u64,
28    /// Disk answers made resident: a header looked up, or samples read.
29    pub promotions: u64,
30    pub writebacks: u64,
31    /// Every entry's hits, summed: each answer of a node or a value memory held.
32    pub hits: u64,
33    pub probation_evictions: u64,
34    pub protected_evictions: u64,
35}
36
37impl Counters {
38    pub fn since(self, then: Counters) -> Counters {
39        Counters {
40            disk_lookups: self.disk_lookups - then.disk_lookups,
41            disk_reads: self.disk_reads - then.disk_reads,
42            disk_read_bytes: self.disk_read_bytes - then.disk_read_bytes,
43            promotions: self.promotions - then.promotions,
44            writebacks: self.writebacks - then.writebacks,
45            hits: self.hits - then.hits,
46            probation_evictions: self.probation_evictions - then.probation_evictions,
47            protected_evictions: self.protected_evictions - then.protected_evictions,
48        }
49    }
50
51    pub fn evictions(self) -> u64 {
52        self.probation_evictions + self.protected_evictions
53    }
54}
55
56/// What a keep states of the node its samples answer: whether it is the target, whether two
57/// or more values read it, and whether a stateful run wrote it.
58#[derive(Clone, Copy, Debug)]
59pub(crate) struct Facts {
60    pub(crate) target: bool,
61    pub(crate) shared: bool,
62    pub(crate) stateful: bool,
63}
64
65/// What one keep sends memory: the samples a value computed, and the node they answer.
66pub(crate) struct Keep<'a> {
67    pub(crate) samples: Option<Payload>,
68    pub(crate) label: Option<&'a Label>,
69    pub(crate) slot: Option<Hash>,
70    pub(crate) node: Option<(Stored, Offered, Facts)>,
71}
72
73/// The protected segment holds at most this share of the cap, so probation always has a
74/// quarter for a new entry to earn its first hit in.
75const PROTECTED: (u64, u64) = (3, 4);
76
77#[derive(Clone, Copy, Debug, PartialEq, Eq)]
78pub(crate) enum Kept {
79    Held,
80    Replaced,
81    Refused,
82}
83
84/// What memory answers a node's key with: a miss holds for the round it was met in and later
85/// ones begun before it, so another holder's commit is found once a new round begins.
86pub(crate) enum Known {
87    Hit(Arc<Stored>),
88    Miss,
89    Unknown,
90}
91
92/// Where a kept node's sample `n` is.
93#[derive(Clone, Debug, PartialEq)]
94pub(crate) enum Offered {
95    /// The value held under the node's own key: a run's through each segment before it.
96    Own,
97    Moves {
98        of: Hash,
99        by: i64,
100    },
101    /// Samples no value holds, the node's own.
102    Held(Vec<Arc<Buffer>>),
103}
104
105pub(crate) struct Writeback {
106    pub(crate) key: Hash,
107    pub(crate) head: Header,
108    pub(crate) parts: Vec<Arc<Buffer>>,
109}
110
111/// A value's samples, and the node they were kept or read off the disk as.
112struct Item {
113    payload: Option<Payload>,
114    label: Option<Label>,
115    slot: Option<Hash>,
116    node: Option<Node>,
117}
118
119struct Node {
120    source: Source,
121    /// Holds what the disk lacks: each keep of more samples sets it again.
122    dirty: bool,
123    bound: bool,
124}
125
126enum Source {
127    Disk {
128        head: Box<Header>,
129        chunks: Vec<Arc<Buffer>>,
130    },
131    Offered {
132        stored: Box<Stored>,
133        offered: Offered,
134    },
135}
136
137impl Node {
138    fn stored(&self) -> &Stored {
139        match &self.source {
140            Source::Disk { head, .. } => head.stored(),
141            Source::Offered { stored, .. } => stored,
142        }
143    }
144
145    fn offered(&self) -> Option<&Offered> {
146        match &self.source {
147            Source::Offered { offered, .. } => Some(offered),
148            Source::Disk { .. } => None,
149        }
150    }
151
152    fn head(&self) -> Option<&Header> {
153        match &self.source {
154            Source::Disk { head, .. } => Some(head),
155            Source::Offered { .. } => None,
156        }
157    }
158
159    fn chunks(&self) -> &[Arc<Buffer>] {
160        match &self.source {
161            Source::Disk { chunks, .. }
162            | Source::Offered {
163                offered: Offered::Held(chunks),
164                ..
165            } => chunks,
166            Source::Offered { .. } => &[],
167        }
168    }
169}
170
171struct Held {
172    item: Item,
173    read: u64,
174    since: u64,
175    hit_round: u64,
176    protected: bool,
177}
178
179impl Held {
180    fn admitted(item: Item, at: u64) -> Held {
181        Held {
182            item,
183            read: at,
184            since: at,
185            hit_round: 0,
186            protected: false,
187        }
188    }
189
190    fn bytes(&self) -> u64 {
191        let value = self.item.payload.as_ref().map_or(0, |p| p.bytes() as u64);
192        let node = self.item.node.as_ref().map_or(&[][..], Node::chunks);
193        value + node.iter().map(|c| planes(c)).sum::<u64>()
194    }
195}
196
197fn planes(b: &Buffer) -> u64 {
198    (b.len() * b.width() * size_of::<f64>()) as u64
199}
200
201#[derive(Default)]
202struct State {
203    entries: HashMap<Hash, Held>,
204    slots: HashMap<Hash, Hash>,
205    misses: HashMap<Hash, u64>,
206    pending: Vec<Writeback>,
207    bytes: u64,
208    max_bytes: u64,
209    clock: u64,
210    round: u64,
211    mark_every: usize,
212    counters: Counters,
213    disk: bool,
214    /// Writes that failed, and why the latest did: none is tried again until a persist.
215    failed: (u64, Option<String>),
216    links: RefCell<Links>,
217    #[cfg(test)]
218    walked: std::cell::Cell<u64>,
219}
220
221/// Where a chain of moved nodes ends and the shift there; a chain that loops ends nowhere.
222#[derive(Clone, Copy)]
223struct Link {
224    end: Hash,
225    by: i64,
226    unbound: bool,
227    looped: bool,
228}
229
230/// Each moved node's link, a cache every change to `entries` evicts what it found through.
231#[derive(Default)]
232struct Links {
233    /// Each link, with the key it was found through.
234    found: HashMap<Hash, (Link, Hash)>,
235    through: HashMap<Hash, HashSet<Hash>>,
236}
237
238impl Links {
239    fn keep(&mut self, at: Hash, next: Hash, link: Link) {
240        self.found.insert(at, (link, next));
241        self.through.entry(next).or_default().insert(at);
242    }
243
244    fn evict(&mut self, key: Hash) {
245        let mut open = vec![key];
246        while let Some(at) = open.pop() {
247            if let Some((_, next)) = self.found.remove(&at)
248                && let Some(readers) = self.through.get_mut(&next)
249            {
250                readers.remove(&at);
251                if readers.is_empty() {
252                    self.through.remove(&next);
253                }
254            }
255            open.extend(self.through.remove(&at).into_iter().flatten());
256        }
257    }
258}
259
260impl State {
261    fn tick(&mut self) -> u64 {
262        self.clock += 1;
263        self.clock
264    }
265
266    fn node(&self, key: Hash) -> Option<&Node> {
267        self.entries.get(&key)?.item.node.as_ref()
268    }
269
270    fn stands(&self, key: Hash) -> bool {
271        let bound = |at: &Hash| self.node(*at).is_some_and(|node| node.bound);
272        self.doomed(key).iter().any(bound)
273    }
274
275    fn on_disk(&self, key: Hash, payload: &Payload) -> bool {
276        let Some(head) = self.node(key).and_then(Node::head) else {
277            return false;
278        };
279        let stored = head.stored();
280        match payload {
281            Payload::Segments(parts) => parts.iter().all(|part| stored.holds(part.extent())),
282            Payload::Run(run) => stored.holds(run.samples.extent()),
283            Payload::Frames(_) => false,
284        }
285    }
286
287    /// A run's segment, then each before it, ending where the one after it starts.
288    fn segments(&self, key: Hash) -> Vec<(Hash, Extent)> {
289        let (mut out, mut at, mut end) = (Vec::new(), Some(key), i64::MAX);
290        while let Some(k) = at {
291            out.push((k, Extent::new(i64::MIN, end)));
292            let held = self
293                .entries
294                .get(&k)
295                .and_then(|held| held.item.payload.as_ref());
296            let Some(Payload::Run(run)) = held else {
297                break;
298            };
299            (at, end) = (run.parent, run.samples.start);
300        }
301        out
302    }
303
304    fn reads(&self, key: Hash) -> Vec<Hash> {
305        match self.node(key).and_then(Node::offered) {
306            None | Some(Offered::Held(_)) => Vec::new(),
307            Some(Offered::Own) => self
308                .segments(key)
309                .into_iter()
310                .skip(1)
311                .map(|(k, _)| k)
312                .collect(),
313            Some(Offered::Moves { of, .. }) => vec![*of],
314        }
315    }
316
317    fn doomed(&self, key: Hash) -> Vec<Hash> {
318        let mut readers: HashMap<Hash, Vec<Hash>> = HashMap::new();
319        for at in self.entries.keys() {
320            for read in self.reads(*at) {
321                readers.entry(read).or_default().push(*at);
322            }
323        }
324        let mut out = vec![key];
325        let mut k = 0;
326        while k < out.len() {
327            for at in readers.get(&out[k]).into_iter().flatten() {
328                if !out.contains(at) {
329                    out.push(*at);
330                }
331            }
332            k += 1;
333        }
334        out
335    }
336
337    /// `key` and every node reading its samples gone, each dirty node flushed first: a later
338    /// keep brings back what it keeps on computing.
339    fn remove(&mut self, key: Hash) -> bool {
340        if !self.entries.contains_key(&key) {
341            return false;
342        }
343        let doomed = self.doomed(key);
344        for at in &doomed {
345            self.flush(*at);
346        }
347        for at in doomed {
348            self.discard(at);
349        }
350        true
351    }
352
353    fn discard(&mut self, key: Hash) {
354        let Some(gone) = self.entries.remove(&key) else {
355            return;
356        };
357        self.links.get_mut().evict(key);
358        self.bytes -= gone.bytes();
359        if let Some(slot) = gone.item.slot
360            && self.slots.get(&slot) == Some(&key)
361        {
362            self.slots.remove(&slot);
363        }
364    }
365
366    fn unnoded(&mut self, key: Hash) {
367        let Some(held) = self.entries.get_mut(&key) else {
368            return;
369        };
370        if held.item.payload.is_none() {
371            return self.discard(key);
372        }
373        let before = held.bytes();
374        held.item.node = None;
375        self.bytes = self.bytes - before + held.bytes();
376        self.links.get_mut().evict(key);
377    }
378
379    /// A dirty node's header and what memory holds of it, on their way to the disk.
380    fn flush(&mut self, key: Hash) {
381        let Some(Node {
382            dirty: dirty @ true,
383            ..
384        }) = self
385            .entries
386            .get_mut(&key)
387            .and_then(|held| held.item.node.as_mut())
388        else {
389            return;
390        };
391        *dirty = false;
392        if let Some(back) = self.written(key) {
393            self.pending.push(back);
394        }
395    }
396
397    fn written(&self, key: Hash) -> Option<Writeback> {
398        let Some(Node {
399            source: Source::Offered { stored, offered },
400            ..
401        }) = self.node(key)
402        else {
403            return None;
404        };
405        let head = |samples| Header::new((**stored).clone(), samples);
406        let referred = match offered {
407            Offered::Moves { of, by } => match self.referred(*of, *by) {
408                Some(found) => Some(found),
409                None if self.unbound(*of) => None,
410                None => return None,
411            },
412            _ => None,
413        };
414        Some(match referred {
415            Some((of, by)) => Writeback {
416                key,
417                head: head(Samples::Of { key: of, by }),
418                parts: Vec::new(),
419            },
420            None => Writeback {
421                key,
422                head: head(Samples::None),
423                parts: self
424                    .offered_parts(key, offered, Extent::EVERYWHERE)
425                    .unwrap_or_default(),
426            },
427        })
428    }
429
430    fn moving(&self, key: Hash) -> Option<(Hash, i64, bool)> {
431        let node = self.node(key)?;
432        match node.offered()? {
433            Offered::Moves { of, by } => Some((*of, *by, node.bound)),
434            _ => None,
435        }
436    }
437
438    /// Walked from `key` to the first link already found, each link walked kept.
439    fn link(&self, key: Hash) -> Link {
440        let (mut walked, mut seen, mut at) = (Vec::new(), HashSet::new(), key);
441        let mut link = loop {
442            if let Some((held, _)) = self.links.borrow().found.get(&at) {
443                break *held;
444            }
445            let looped = !seen.insert(at);
446            match self.moving(at) {
447                Some((of, by, bound)) if !looped => {
448                    #[cfg(test)]
449                    self.walked.set(self.walked.get() + 1);
450                    walked.push((at, of, by, !bound));
451                    at = of;
452                }
453                moving => {
454                    break Link {
455                        end: at,
456                        by: 0,
457                        unbound: false,
458                        looped: moving.is_some(),
459                    };
460                }
461            }
462        };
463        for (at, of, by, unbound) in walked.into_iter().rev() {
464            link.by += by;
465            link.unbound |= unbound;
466            if !link.looped {
467                self.links.borrow_mut().keep(at, of, link);
468            }
469        }
470        link
471    }
472
473    /// The node `stored` names, as `offered` says, in place of the last one under its key, and,
474    /// `sole`, with no samples to keep, under its slot; whether memory writes it to the disk. A
475    /// node the disk holds stays as the disk holds it.
476    fn noded(
477        &mut self,
478        stored: Stored,
479        offered: Offered,
480        (slot, sole, facts): (Option<Hash>, bool, Facts),
481    ) -> bool {
482        let key = stored.key;
483        if self.node(key).is_some_and(|node| node.head().is_some()) {
484            return true;
485        }
486        let bound = self.disk && writes(&stored, facts);
487        let read = self.tick();
488        let same = self.node(key).is_some_and(|node| match &node.source {
489            Source::Offered {
490                stored: had,
491                offered: was,
492            } => **had == stored && *was == offered,
493            Source::Disk { .. } => false,
494        });
495        let node = Node {
496            source: Source::Offered {
497                stored: Box::new(stored),
498                offered,
499            },
500            dirty: bound,
501            bound,
502        };
503        match self.entries.get_mut(&key) {
504            Some(held) if same => {
505                let node = held.item.node.as_mut().expect("the node held");
506                node.dirty |= node.bound;
507                held.read = read;
508            }
509            Some(held) => {
510                let before = held.bytes();
511                held.item.node = Some(node);
512                held.item.slot = slot;
513                held.read = read;
514                let after = held.bytes();
515                self.bytes = self.bytes - before + after;
516                self.links.get_mut().evict(key);
517            }
518            None => {
519                let item = Item {
520                    payload: None,
521                    label: None,
522                    slot,
523                    node: Some(node),
524                };
525                let held = Held::admitted(item, read);
526                self.bytes += held.bytes();
527                self.admit(key, held);
528            }
529        }
530        self.misses.remove(&key);
531        if let Some(last) = slot
532            .filter(|_| sole)
533            .and_then(|slot| self.slots.insert(slot, key))
534            && last != key
535        {
536            self.remove(last);
537        }
538        bound
539    }
540
541    fn admit(&mut self, key: Hash, held: Held) -> Option<Held> {
542        self.links.get_mut().evict(key);
543        self.entries.insert(key, held)
544    }
545
546    fn foot(&self, key: Hash) -> Option<(Hash, i64)> {
547        let link = self.link(key);
548        (!link.looped).then_some((link.end, link.by))
549    }
550
551    /// A link or its end is no node the disk is written, or is only a value.
552    fn unbound(&self, key: Hash) -> bool {
553        let link = self.link(key);
554        let end = self.entries.get(&link.end).map(|held| &held.item.node);
555        link.unbound || matches!(end, Some(None | Some(Node { bound: false, .. })))
556    }
557
558    fn referred(&self, of: Hash, by: i64) -> Option<(Hash, i64)> {
559        let link = self.link(of);
560        if link.looped || link.unbound {
561            return None;
562        }
563        let more = link.by;
564        let node = self.node(link.end)?;
565        match (node.head().map(Header::samples), node.bound) {
566            (Some(Samples::Entry { file, shift, .. } | Samples::Staged { file, shift, .. }), _) => {
567                Some((*file, by + more - shift))
568            }
569            (Some(_), _) | (None, false) => None,
570            (None, true) => Some((link.end, by + more)),
571        }
572    }
573
574    fn parts(&self, key: Hash) -> Option<Vec<(Arc<Buffer>, Extent)>> {
575        self.entries.get(&key)?;
576        let mut out = Vec::new();
577        for (k, within) in self.segments(key) {
578            let parts = match self
579                .entries
580                .get(&k)
581                .and_then(|held| held.item.payload.as_ref())
582            {
583                Some(Payload::Segments(parts)) => parts.clone(),
584                Some(Payload::Run(run)) => vec![Arc::clone(&run.samples)],
585                _ => Vec::new(),
586            };
587            let met = |part: Arc<Buffer>| {
588                let met = part.extent().intersect(within);
589                (!met.is_empty()).then_some((part, met))
590            };
591            out.extend(parts.into_iter().filter_map(met));
592        }
593        Some(out)
594    }
595
596    fn own(&self, key: Hash, over: Extent) -> Option<Vec<Arc<Buffer>>> {
597        let met = self.parts(key)?.into_iter();
598        let met = met.filter(|(_, held)| !held.intersect(over).is_empty());
599        Some(met.map(|(part, held)| clipped(&part, held)).collect())
600    }
601
602    fn offered_parts(
603        &self,
604        key: Hash,
605        offered: &Offered,
606        over: Extent,
607    ) -> Option<Vec<Arc<Buffer>>> {
608        let (parts, by): (Vec<Arc<Buffer>>, i64) = match offered {
609            Offered::Own => (self.own(key, over)?, 0),
610            Offered::Moves { of, by } => {
611                let (foot, more) = self.foot(*of)?;
612                (self.resident(foot, over.shifted(by + more))?.0, by + more)
613            }
614            Offered::Held(parts) => (parts.clone(), 0),
615        };
616        let over = over.shifted(by);
617        let meets = |b: &&Arc<Buffer>| !b.extent().intersect(over).is_empty();
618        Some(parts.iter().filter(meets).map(|p| moved(p, -by)).collect())
619    }
620
621    /// What memory holds of `key` meeting `over`, and for a node off the disk, what of `over`
622    /// the disk holds that memory lacks, with its layout.
623    fn resident(&self, key: Hash, over: Extent) -> Option<Resident> {
624        let source = match &self.entries.get(&key)?.item.node {
625            None => return Some((self.own(key, over)?, None)),
626            Some(node) => &node.source,
627        };
628        match source {
629            Source::Offered { offered, .. } => {
630                Some((self.offered_parts(key, offered, over)?, None))
631            }
632            Source::Disk { head, chunks } => {
633                let meets = |b: &&Arc<Buffer>| !b.extent().intersect(over).is_empty();
634                let parts: Vec<Arc<Buffer>> = chunks.iter().filter(meets).cloned().collect();
635                let mut lacks = Vec::new();
636                for held in head.stored().extents() {
637                    let mut asked = vec![held.intersect(over)];
638                    for part in chunks {
639                        asked = asked
640                            .into_iter()
641                            .flat_map(|a| minus(a, part.extent()))
642                            .collect();
643                    }
644                    lacks.extend(asked.into_iter().filter(|a| !a.is_empty()));
645                }
646                let hull = lacks.iter().fold(Extent::NOWHERE, |h, e| h.hull(*e));
647                Some((parts, (!hull.is_empty()).then(|| ((**head).clone(), hull))))
648            }
649        }
650    }
651
652    fn coverage(&self, key: Hash) -> Option<Vec<Extent>> {
653        let (foot, by) = self.foot(key)?;
654        let node = self.entries.get(&foot)?.item.node.as_ref();
655        let held = match (node.and_then(Node::head), node.and_then(Node::offered)) {
656            (Some(head), _) => head.stored().extents().to_vec(),
657            (None, None | Some(Offered::Own)) => {
658                let parts = self.parts(foot)?.into_iter();
659                parts.map(|(_, held)| held).collect()
660            }
661            (None, Some(Offered::Held(parts))) => parts.iter().map(|p| p.extent()).collect(),
662            (None, Some(Offered::Moves { .. })) => return None,
663        };
664        Some(held.into_iter().map(|e| e.shifted(-by)).collect())
665    }
666
667    /// A node off the disk keeps its header.
668    fn evict(&mut self, key: Hash) {
669        let Some(held) = self.entries.get_mut(&key) else {
670            return;
671        };
672        let protected = held.protected;
673        let before = held.bytes();
674        let gone = match held.item.node.as_mut().map(|node| &mut node.source) {
675            Some(Source::Disk { chunks, .. }) => {
676                chunks.clear();
677                held.item.payload = None;
678                self.bytes -= before;
679                true
680            }
681            _ => self.remove(key),
682        };
683        match (gone, protected) {
684            (false, _) => {}
685            (true, false) => self.counters.probation_evictions += 1,
686            (true, true) => self.counters.protected_evictions += 1,
687        }
688    }
689
690    /// A hit on `key` and on every entry whose samples it answers with: each moves to the
691    /// protected segment, whose least recent move back to probation past its share.
692    fn hit(&mut self, key: Hash, round: Option<u64>) {
693        let tick = self.tick();
694        let mut promoted = false;
695        for at in self.under(key) {
696            let Some(held) = self.entries.get_mut(&at) else {
697                continue;
698            };
699            held.read = tick;
700            if round.is_some_and(|round| held.hit_round == round) {
701                continue;
702            }
703            held.hit_round = round.unwrap_or(0);
704            promoted |= !std::mem::replace(&mut held.protected, true);
705            self.counters.hits += 1;
706        }
707        if promoted {
708            self.shared();
709        }
710    }
711
712    fn under(&self, key: Hash) -> Vec<Hash> {
713        let mut out = vec![key];
714        let mut k = 0;
715        while k < out.len() {
716            for at in self.reads(out[k]) {
717                if !out.contains(&at) {
718                    out.push(at);
719                }
720            }
721            k += 1;
722        }
723        out
724    }
725
726    /// Protected within its share: its least recent move back to probation, as its newest.
727    fn shared(&mut self) {
728        let most =
729            (u128::from(self.max_bytes) * u128::from(PROTECTED.0) / u128::from(PROTECTED.1)) as u64;
730        let mut protected: Vec<(u64, Hash, u64)> = self
731            .entries
732            .iter()
733            .filter(|(_, held)| held.protected)
734            .map(|(key, held)| (held.read, *key, held.bytes()))
735            .collect();
736        let mut bytes: u64 = protected.iter().map(|(_, _, b)| b).sum();
737        protected.sort_unstable();
738        for (_, key, held) in protected {
739            if bytes <= most {
740                break;
741            }
742            let tick = self.tick();
743            if let Some(entry) = self.entries.get_mut(&key) {
744                entry.protected = false;
745                entry.since = tick;
746                bytes -= held;
747            }
748        }
749    }
750
751    /// Under the cap: an entry too large for it goes first, then probation oldest first,
752    /// then protected least recently read.
753    fn bounded(&mut self) {
754        if self.bytes <= self.max_bytes {
755            return;
756        }
757        self.shared();
758        let max = self.max_bytes;
759        let mut order: Vec<(u8, u64, Hash)> = self
760            .entries
761            .iter()
762            .filter(|(_, held)| held.bytes() > 0)
763            .map(|(key, held)| match (held.bytes() > max, held.protected) {
764                (true, _) => (0, held.since, *key),
765                (false, false) => (1, held.since, *key),
766                (false, true) => (2, held.read, *key),
767            })
768            .collect();
769        order.sort_unstable();
770        for (_, _, key) in order {
771            if self.bytes <= self.max_bytes {
772                break;
773            }
774            self.evict(key);
775        }
776    }
777}
778
779type Resident = (Vec<Arc<Buffer>>, Option<(Header, Extent)>);
780
781fn minus(e: Extent, cut: Extent) -> Vec<Extent> {
782    let met = e.intersect(cut);
783    if met.is_empty() {
784        return vec![e];
785    }
786    vec![Extent::new(e.start, met.start), Extent::new(met.end, e.end)]
787}
788
789fn clipped(part: &Arc<Buffer>, met: Extent) -> Arc<Buffer> {
790    let held = part.extent();
791    match met == held {
792        true => Arc::clone(part),
793        false => Arc::new(part.over(met, held)),
794    }
795}
796
797fn moved(part: &Arc<Buffer>, by: i64) -> Arc<Buffer> {
798    if by == 0 {
799        return Arc::clone(part);
800    }
801    let mut out = (**part).clone();
802    out.start += by;
803    Arc::new(out)
804}
805
806/// The memory tier: the one owner of every value and node held resident, and of what memory
807/// knows a disk beneath it holds or lacks. A clone is a handle on the same memory.
808#[derive(Clone)]
809pub(crate) struct Memory {
810    state: Arc<Mutex<State>>,
811}
812
813impl Default for Memory {
814    fn default() -> Memory {
815        Memory::holding(DEFAULT_CACHE_BYTES)
816    }
817}
818
819impl Memory {
820    pub(crate) fn holding(max_bytes: u64) -> Memory {
821        Memory {
822            state: Arc::new(Mutex::new(State {
823                max_bytes,
824                mark_every: DEFAULT_MARK_EVERY,
825                ..State::default()
826            })),
827        }
828    }
829
830    /// Over a disk: each node a value keeps is written back before it goes.
831    pub(crate) fn over_disk(max_bytes: u64) -> Memory {
832        let memory = Memory::holding(max_bytes);
833        memory.locked().disk = true;
834        memory
835    }
836
837    fn locked(&self) -> MutexGuard<'_, State> {
838        self.state.lock().unwrap_or_else(|poisoned| {
839            let mut state = poisoned.into_inner();
840            state.entries.clear();
841            *state.links.get_mut() = Links::default();
842            state.slots.clear();
843            state.bytes = 0;
844            self.state.clear_poison();
845            state
846        })
847    }
848
849    pub(crate) fn max_bytes(&self) -> u64 {
850        self.locked().max_bytes
851    }
852
853    pub(crate) fn set_max_bytes(&self, max_bytes: u64) {
854        let mut state = self.locked();
855        state.max_bytes = max_bytes;
856        state.bounded();
857    }
858
859    pub(crate) fn bytes(&self) -> u64 {
860        self.locked().bytes
861    }
862
863    pub(crate) fn entries(&self) -> usize {
864        self.locked().entries.len()
865    }
866
867    pub(crate) fn holds(&self, key: Hash) -> bool {
868        self.locked().entries.contains_key(&key)
869    }
870
871    pub(crate) fn counters(&self) -> Counters {
872        self.locked().counters
873    }
874
875    pub(crate) fn count(&self, by: impl FnOnce(&mut Counters)) {
876        by(&mut self.locked().counters);
877    }
878
879    pub(crate) fn mark_every(&self) -> usize {
880        self.locked().mark_every
881    }
882
883    pub(crate) fn set_mark_every(&self, samples: usize) {
884        self.locked().mark_every = samples.max(1);
885    }
886
887    /// Whether memory takes a value computed: a memory of no bytes takes none.
888    pub(crate) fn keeps(&self) -> bool {
889        self.locked().max_bytes > 0
890    }
891
892    /// A new round of lookups: what any earlier round missed is asked again.
893    pub(crate) fn begin(&self) -> u64 {
894        let mut state = self.locked();
895        state.round += 1;
896        let round = state.round;
897        state.misses.retain(|_, met| *met + 1 >= round);
898        round
899    }
900
901    pub(crate) fn answer(&self, key: Hash, round: u64) -> Known {
902        let mut state = self.locked();
903        if state.node(key).is_some() {
904            let covered = state.coverage(key);
905            if let Some(held) = covered.clone().filter(|held| !held.is_empty()) {
906                state.hit(key, Some(round));
907                let node = state.node(key).expect("a node held");
908                return Known::Hit(Arc::new(node.stored().holding(held)));
909            }
910            match covered {
911                None => {
912                    state.remove(key);
913                }
914                Some(_) => return Known::Miss,
915            }
916        }
917        match state.misses.get(&key) {
918            Some(met) if *met >= round => Known::Miss,
919            _ => Known::Unknown,
920        }
921    }
922
923    /// `answer`, told to `seen` once a hit or a miss.
924    pub(crate) fn answered(
925        &self,
926        (key, round): (Hash, u64),
927        (node, seen): (&str, &mut Recording),
928    ) -> Known {
929        let known = self.answer(key, round);
930        let outcome = match &known {
931            Known::Hit(_) => Outcome::Hit,
932            Known::Miss => Outcome::ComputedNotStored,
933            Known::Unknown => return known,
934        };
935        let state = self.locked();
936        let held = state
937            .entries
938            .get(&key)
939            .and_then(|held| held.item.payload.as_ref());
940        let kind = held.map_or(PayloadKind::Segments, Payload::kind);
941        drop(state);
942        seen.answered(key, node, kind, outcome);
943        known
944    }
945
946    pub(crate) fn miss(&self, key: Hash, round: u64) {
947        let mut state = self.locked();
948        let met = state.misses.entry(key).or_insert(round);
949        *met = (*met).max(round);
950    }
951
952    /// A header the disk answered, resident from now on, unless memory took the node meanwhile.
953    pub(crate) fn promote(&self, head: Header) {
954        let mut state = self.locked();
955        let (key, read) = (head.stored().key, state.tick());
956        if state.node(key).is_some() {
957            return;
958        }
959        state.misses.remove(&key);
960        state.counters.promotions += 1;
961        let node = Node {
962            source: Source::Disk {
963                head: Box::new(head),
964                chunks: Vec::new(),
965            },
966            dirty: false,
967            bound: true,
968        };
969        match state.entries.get_mut(&key) {
970            Some(held) => held.item.node = Some(node),
971            None => {
972                let item = Item {
973                    payload: None,
974                    label: None,
975                    slot: None,
976                    node: Some(node),
977                };
978                state.admit(key, Held::admitted(item, read));
979            }
980        }
981        state.links.get_mut().evict(key);
982    }
983
984    pub(crate) fn promote_samples(&self, key: Hash, read: Vec<Buffer>) -> Vec<Arc<Buffer>> {
985        let read: Vec<Arc<Buffer>> = read.into_iter().map(Arc::new).collect();
986        let mut state = self.locked();
987        let tick = state.tick();
988        let Some(held) = state.entries.get_mut(&key) else {
989            return read;
990        };
991        let Some(Node {
992            source: Source::Disk { chunks, .. },
993            ..
994        }) = &mut held.item.node
995        else {
996            return read;
997        };
998        let mut added = 0;
999        for part in &read {
1000            let covered = chunks
1001                .iter()
1002                .any(|c| c.extent().intersect(part.extent()) == part.extent());
1003            if !covered {
1004                added += planes(part);
1005                chunks.push(Arc::clone(part));
1006            }
1007        }
1008        chunks.sort_by_key(|c| c.start);
1009        held.read = tick;
1010        state.bytes += added;
1011        state.counters.promotions += 1;
1012        state.bounded();
1013        read
1014    }
1015
1016    /// What memory holds of `key` over `over`, and what it lacks there that the disk holds.
1017    pub(crate) fn resident(&self, key: Hash, over: Extent) -> Resident {
1018        let mut state = self.locked();
1019        let tick = state.tick();
1020        if let Some(held) = state.entries.get_mut(&key) {
1021            held.read = tick;
1022        }
1023        state.resident(key, over).unwrap_or_default()
1024    }
1025
1026    pub(crate) fn forget(&self, key: Hash) {
1027        self.locked().remove(key);
1028    }
1029
1030    /// A value's samples, joined to what memory holds under `key`, and the node they answer,
1031    /// in place of the last one under its key or slot: the node first, so samples a node over
1032    /// the disk stands on are kept to be written back. A node holding none of its samples that
1033    /// memory will write nowhere goes. What became of the samples is told to `seen`.
1034    pub(crate) fn keep(&self, key: Hash, keep: Keep, seen: &mut Recording) {
1035        let Keep {
1036            samples,
1037            label,
1038            slot,
1039            node,
1040        } = keep;
1041        let sole = samples.is_none();
1042        let noded = node.map(|(stored, offered, facts)| {
1043            let key = stored.key;
1044            let bound = self.locked().noded(stored, offered, (slot, sole, facts));
1045            (key, bound)
1046        });
1047        let kept = samples.map(|samples| self.merge(key, samples, label, slot));
1048        if let Some(kept) = kept {
1049            seen.kept(key, kept);
1050        }
1051        if let Some((key, bound)) = noded {
1052            let mut state = self.locked();
1053            let covered = state.coverage(key).is_some_and(|held| !held.is_empty());
1054            if !covered && !bound {
1055                state.unnoded(key);
1056            }
1057            state.bounded();
1058        }
1059    }
1060
1061    /// The nodes read after `since`, least recent first, and the clock now.
1062    pub(crate) fn read_since(&self, since: u64) -> (Vec<Hash>, u64) {
1063        let state = self.locked();
1064        let mut read: Vec<(u64, Hash)> = state
1065            .entries
1066            .iter()
1067            .filter(|(_, held)| held.read > since && held.item.node.is_some())
1068            .map(|(key, held)| (held.read, *key))
1069            .collect();
1070        read.sort_unstable();
1071        (read.into_iter().map(|(_, key)| key).collect(), state.clock)
1072    }
1073
1074    /// Every node not yet on the disk, on its way there.
1075    pub(crate) fn flush(&self) {
1076        let mut state = self.locked();
1077        let keys: Vec<Hash> = state.entries.keys().copied().collect();
1078        for key in keys {
1079            state.flush(key);
1080        }
1081        state.failed.1 = None;
1082    }
1083
1084    /// The nodes on their way to the disk; none while a write that failed waits for a persist.
1085    pub(crate) fn pending(&self) -> Vec<Writeback> {
1086        let mut state = self.locked();
1087        match state.failed.1 {
1088            Some(_) => Vec::new(),
1089            None => std::mem::take(&mut state.pending),
1090        }
1091    }
1092
1093    /// `left` still on its way, after a write failed for `why`.
1094    pub(crate) fn failed(&self, why: String, left: Vec<Writeback>) {
1095        let mut state = self.locked();
1096        state.failed.0 += 1;
1097        state.failed.1 = Some(why);
1098        let mut more = std::mem::take(&mut state.pending);
1099        state.pending = left;
1100        state.pending.append(&mut more);
1101    }
1102
1103    pub(crate) fn written(&self) {
1104        self.locked().counters.writebacks += 1;
1105    }
1106
1107    pub(crate) fn failures(&self) -> (u64, Option<String>) {
1108        self.locked().failed.clone()
1109    }
1110
1111    pub(crate) fn blocked(&self) -> bool {
1112        self.locked().failed.1.is_some()
1113    }
1114
1115    /// The disk committed what was staged: a header that named staged samples names nothing
1116    /// now.
1117    pub(crate) fn committed(&self) {
1118        let mut state = self.locked();
1119        let staged: Vec<Hash> = state
1120            .entries
1121            .iter()
1122            .filter(|(_, held)| {
1123                let head = held.item.node.as_ref().and_then(Node::head);
1124                head.is_some_and(|head| matches!(head.samples(), Samples::Staged { .. }))
1125            })
1126            .map(|(key, _)| *key)
1127            .collect();
1128        for key in staged {
1129            state.unnoded(key);
1130        }
1131    }
1132
1133    /// What `key` holds, shared, never copied, told to `seen`.
1134    pub(crate) fn load(
1135        &self,
1136        key: Hash,
1137        expected: Expected,
1138        (node, seen): (&str, &mut Recording),
1139    ) -> Option<Entry> {
1140        let entry = self.entry(key, expected);
1141        let outcome = match entry {
1142            Some(_) => Outcome::Hit,
1143            None => Outcome::ComputedNotStored,
1144        };
1145        seen.answered(key, node, expected.kind(), outcome);
1146        entry
1147    }
1148
1149    /// A run's segments held under `keys` in turn, each handed to `take` until it takes none or
1150    /// asks no more, told to `seen` as one lookup of the last: a prefix where it took fewer.
1151    pub(crate) fn runs(
1152        &self,
1153        keys: &[Hash],
1154        expected: Expected,
1155        (node, seen): (&str, &mut Recording),
1156        mut take: impl FnMut(usize, Arc<Run>) -> (bool, bool),
1157    ) {
1158        let mut taken = 0;
1159        for (k, key) in keys.iter().enumerate() {
1160            let Some(run) = self.entry(*key, expected).and_then(|e| e.payload.run()) else {
1161                break;
1162            };
1163            let (took, more) = take(k, run);
1164            taken += usize::from(took);
1165            if !(took && more) {
1166                break;
1167            }
1168        }
1169        let outcome = match taken {
1170            0 => Outcome::ComputedNotStored,
1171            n if n == keys.len() => Outcome::Hit,
1172            _ => Outcome::Prefix,
1173        };
1174        let last = *keys.last().expect("a run of a segment or more");
1175        seen.answered(last, node, PayloadKind::Run, outcome);
1176    }
1177
1178    fn entry(&self, key: Hash, expected: Expected) -> Option<Entry> {
1179        let mut state = self.locked();
1180        let item = &state.entries.get(&key)?.item;
1181        let entry = Entry {
1182            payload: item.payload.clone()?,
1183            label: item.label.clone(),
1184        };
1185        if !entry.payload.answers(expected) {
1186            state.remove(key);
1187            return None;
1188        }
1189        state.hit(key, None);
1190        Some(entry)
1191    }
1192
1193    /// A value's segments join those held under `key`, and a run continuing the one held
1194    /// there extends it, each in place; anything else replaces what `key` held.
1195    fn merge(
1196        &self,
1197        key: Hash,
1198        payload: Payload,
1199        label: Option<&Label>,
1200        slot: Option<Hash>,
1201    ) -> Kept {
1202        let mut state = self.locked();
1203        let tick = state.tick();
1204        let joined = match state.entries.get_mut(&key) {
1205            Some(held) if held.item.slot == slot && held.item.payload.is_some() => {
1206                let before = held.bytes();
1207                let had = held.item.payload.as_mut().expect("a value held");
1208                let payload = match (had, payload) {
1209                    (Payload::Segments(parts), Payload::Segments(more)) => {
1210                        joined(parts, more);
1211                        None
1212                    }
1213                    (Payload::Run(run), Payload::Run(more)) if overlaps(run, &more) => {
1214                        let from = (run.end() - more.samples.start).max(0) as usize;
1215                        let run = Arc::make_mut(run);
1216                        let samples = Arc::make_mut(&mut run.samples);
1217                        for (held, more) in samples.planes.iter_mut().zip(&more.samples.planes) {
1218                            held.extend_from_slice(&more[from.min(more.len())..]);
1219                        }
1220                        run.marks
1221                            .extend(more.marks.iter().map(|(at, m)| (*at, m.clone())));
1222                        None
1223                    }
1224                    (_, payload) => Some(payload),
1225                };
1226                match payload {
1227                    None => {
1228                        held.read = tick;
1229                        let after = held.bytes();
1230                        Ok((before, after))
1231                    }
1232                    Some(payload) => Err(payload),
1233                }
1234            }
1235            _ => Err(payload),
1236        };
1237        match joined {
1238            Ok((before, after)) => {
1239                state.bytes = state.bytes - before + after;
1240                state.bounded();
1241                Kept::Held
1242            }
1243            Err(payload) => {
1244                drop(state);
1245                self.store(key, payload, label, slot)
1246            }
1247        }
1248    }
1249
1250    /// A value the disk holds under `key` is held there alone. One too large to stay is
1251    /// refused, unless a node over the disk stands on it: then it is kept only to be evicted,
1252    /// and so written back.
1253    fn store(
1254        &self,
1255        key: Hash,
1256        payload: Payload,
1257        label: Option<&Label>,
1258        slot: Option<Hash>,
1259    ) -> Kept {
1260        let mut state = self.locked();
1261        if state.on_disk(key, &payload) {
1262            return Kept::Held;
1263        }
1264        let bytes = payload.bytes() as u64;
1265        if bytes > state.max_bytes && !state.stands(key) {
1266            return Kept::Refused;
1267        }
1268        let read = state.tick();
1269        let replaced = match slot.and_then(|slot| state.slots.insert(slot, key)) {
1270            Some(last) if last != key => state.remove(last),
1271            _ => false,
1272        };
1273        let old = state.entries.remove(&key);
1274        if let Some(old) = &old {
1275            state.bytes -= old.bytes();
1276        }
1277        let (node, earned) = match old {
1278            Some(old) => (
1279                old.item.node,
1280                Some((old.read, old.since, old.hit_round, old.protected)),
1281            ),
1282            None => (None, None),
1283        };
1284        let item = Item {
1285            payload: Some(payload),
1286            label: label.cloned(),
1287            slot,
1288            node,
1289        };
1290        let held = match earned {
1291            Some((read, since, hit_round, protected)) if item.node.is_some() => Held {
1292                item,
1293                read,
1294                since,
1295                hit_round,
1296                protected,
1297            },
1298            _ => Held::admitted(item, read),
1299        };
1300        state.bytes += held.bytes();
1301        state.admit(key, held);
1302        state.bounded();
1303        match replaced {
1304            true => Kept::Replaced,
1305            false => Kept::Held,
1306        }
1307    }
1308}
1309
1310/// The target, and a node a later render is answered by that two or more values read or a
1311/// stateful run wrote: what one reader alone recomputes from its own reads is never written.
1312fn writes(stored: &Stored, facts: Facts) -> bool {
1313    facts.target || (stored.readable && (facts.shared || facts.stateful))
1314}
1315
1316/// A run that starts inside or at the end of the one held continues it: what it holds past
1317/// that one's end is laid on, the samples both hold being the same.
1318fn overlaps(held: &super::Run, more: &super::Run) -> bool {
1319    let (a, b) = (held.samples.start, held.end());
1320    a <= more.samples.start && more.samples.start <= b
1321}
1322
1323#[cfg(test)]
1324mod tests {
1325    use super::*;
1326
1327    /// Only a colliding key reaches this, so no render can: the entry is a miss, and goes.
1328    #[test]
1329    fn an_entry_that_does_not_answer_what_was_asked_is_a_miss_and_goes() {
1330        let memory = Memory::default();
1331        let key = Hash(7, 11);
1332        let four = Payload::Segments(vec![Arc::new(Buffer::mono(8_000, vec![0.25; 4]))]);
1333        for (rate, width) in [(48_000, 1), (8_000, 2)] {
1334            memory.store(key, four.clone(), None, None);
1335            let asked = Expected::Segments { rate, width };
1336            assert!(memory.entry(key, asked).is_none());
1337            assert!(!memory.holds(key));
1338            assert_eq!(memory.bytes(), 0);
1339        }
1340    }
1341
1342    #[test]
1343    fn a_load_shares_the_samples_it_holds() {
1344        let memory = Memory::default();
1345        let key = Hash(3, 5);
1346        let part = Arc::new(Buffer::mono(8_000, vec![0.5; 64]));
1347        memory.store(key, Payload::Segments(vec![Arc::clone(&part)]), None, None);
1348        let asked = Expected::Segments {
1349            rate: 8_000,
1350            width: 1,
1351        };
1352        for _ in 0..2 {
1353            let loaded = memory.entry(key, asked).expect("a hit");
1354            let Payload::Segments(parts) = loaded.payload else {
1355                panic!("segments were stored");
1356            };
1357            assert!(Arc::ptr_eq(&parts[0], &part), "the stored part itself");
1358        }
1359    }
1360}