Skip to main content

fs_core/
caching_device.rs

1//! Small LRU read-cache decorator. Caches only block-aligned, block-sized
2//! reads; everything else passes through. Writes invalidate any overlapping
3//! cached entries.
4
5use crate::block::{BlockDevice, BlockRead};
6use crate::error::Result;
7use std::collections::HashMap;
8use std::sync::{Arc, Condvar, Mutex};
9use std::thread::ThreadId;
10
11/// The largest block size a cache will accept.
12///
13/// `block()` allocates a whole block eagerly, and the `.min(size)` clamp
14/// that keeps a short device working bounds that allocation by the device
15/// rather than by anything sane — so a declared block size of 2^40 over a
16/// 1 TB image is a request for a terabyte. A ceiling is the only thing
17/// standing between a number off a disk and that allocation.
18///
19/// 64 MiB is far above anything a real filesystem declares — SquashFS tops
20/// out at 1 MiB, and the block sizes every other driver here uses are
21/// measured in kilobytes — and far below a size that could exhaust a host.
22pub const MAX_BLOCK_SIZE: u64 = 64 * 1024 * 1024;
23
24/// LRU read-cache wrapper.
25///
26/// # It caches a READ device, and writes through one only if it has one
27///
28/// This took an `Arc<dyn BlockDevice>` — the read *write* trait — and
29/// every driver in this family mounts a volume through an
30/// `Arc<dyn BlockRead>`. So a read-only mount could not wrap it at all,
31/// and four of the six drivers used no cache: not by choice, but
32/// because it was not expressible.
33///
34/// The read path never needed to write. It holds the read half now, and
35/// the writable half only when the caller had one to give:
36/// [`CachingDevice::new`] for a device that can be written,
37/// [`CachingDevice::read_only`] for one that cannot. A write to a cache
38/// built the second way is [`crate::Error::ReadOnly`], which is what the
39/// underlying device would have said.
40pub struct CachingDevice {
41    inner: Arc<dyn BlockRead>,
42    /// The same device again, present only when it can be written. Held
43    /// separately rather than as one handle so that "can this be
44    /// written" is a property of the type rather than a flag someone
45    /// has to remember to check.
46    writable: Option<Arc<dyn BlockDevice>>,
47    block_size: u64,
48    /// NOT IN `CacheState`, BECAUSE IT IS NOT STATE.
49    ///
50    /// It is a construction parameter, written once and never mutated,
51    /// and it used to live behind the mutex. Every read therefore took
52    /// the lock to compare `spanned` against it — including the reads
53    /// that were about to be passed straight through and never touch
54    /// the cache at all. On a hit the lock was taken twice, once here
55    /// and once inside `block()`, to serve bytes already in memory.
56    ///
57    /// A plain field beside `block_size` makes the bypass test lock-free
58    /// and says what the value is: fixed at construction, like the block
59    /// size next to it.
60    capacity: usize,
61    state: Mutex<CacheState>,
62    /// Signalled every time a block leaves [`CacheState::in_flight`],
63    /// whether its fetch succeeded or failed.
64    ///
65    /// One condition variable for all blocks rather than one per block.
66    /// A waiter wakes on any fetch completing and re-checks its own
67    /// block, so a wake it did not want costs it a lock and a scan of a
68    /// list bounded by the number of threads fetching. A map of
69    /// per-block condvars would avoid those wakes and has to be built,
70    /// looked up and torn down on every miss — the ordinary path — to
71    /// save work on the rare one.
72    fetched: Condvar,
73}
74
75/// One cached block, and its neighbours in recency order.
76///
77/// `newer`/`older` are indices into [`Lru::slots`], not pointers, so
78/// the list is intrusive without being unsafe. A slot's index is stable
79/// for as long as the node lives, which is what lets the index map
80/// point at it.
81struct Node {
82    block_start: u64,
83    data: Arc<Vec<u8>>,
84    /// Toward the head: more recently used. `None` at the head.
85    newer: Option<usize>,
86    /// Toward the tail: less recently used. `None` at the tail.
87    older: Option<usize>,
88}
89
90/// The cached blocks, in recency order, with O(1) lookup AND O(1)
91/// promotion.
92///
93/// # Why not a `VecDeque` any more
94///
95/// It was one, and a hit cost `iter().position(...)` -- a linear scan
96/// comparing offsets -- followed by `VecDeque::remove(pos)`, which
97/// shifts every element on the shorter side of `pos` to close the gap.
98/// Both halves are linear in the number of entries, so the cache got
99/// slower the larger it was asked to be. Measured on this crate under
100/// `--release`, working set equal to capacity so the steady state is
101/// ~100% hits, with a `CountingDevice` underneath proving no device
102/// traffic at all (`device_reads = 0` at every capacity, so the whole
103/// difference is the cache's own bookkeeping):
104///
105/// ```text
106/// capacity   us/read   vs capacity 8
107///        8    0.0157             1x
108///       64    0.0319           2.0x
109///      512    0.1924          12.3x
110///     4096    1.3306            85x
111/// ```
112///
113/// A cached read at capacity 4096 cost 85 times what the same cached
114/// read cost at capacity 8, and 8x more capacity from 512 to 4096 cost
115/// 6.9x more per read -- linear, which is what `capacity/2` comparisons
116/// plus `capacity/2` tuple moves predicts.
117///
118/// The capacities that matter are already past the knee, which is what
119/// ruled out the other candidate fix of documenting a ceiling:
120/// `am-fs-erofs` defaults to 512 metadata blocks, and one 3 MiB file
121/// read there touches 768 blocks -- about a millisecond of pure
122/// scanning to deliver bytes already in memory.
123///
124/// # Why not a `HashMap` alone
125///
126/// A map fixes the lookup and leaves the promotion: the entry still has
127/// to move to the head of a recency order, and in a `VecDeque` that is
128/// still `remove` plus `push_front`, still O(n), and it invalidates
129/// every index the map is holding. The recency order has to be a linked
130/// list for the promotion to be O(1), and then the map indexes into it.
131struct Lru {
132    /// Slab. `None` is a free slot, kept rather than compacted so live
133    /// indices stay valid.
134    slots: Vec<Option<Node>>,
135    /// Slots to reuse before growing `slots`.
136    free: Vec<usize>,
137    index: HashMap<u64, usize>,
138    /// Most recently used.
139    newest: Option<usize>,
140    /// Least recently used: the eviction end.
141    oldest: Option<usize>,
142}
143
144impl Lru {
145    fn with_capacity(capacity: usize) -> Self {
146        Lru {
147            slots: Vec::with_capacity(capacity),
148            free: Vec::new(),
149            index: HashMap::with_capacity(capacity),
150            newest: None,
151            oldest: None,
152        }
153    }
154
155    fn len(&self) -> usize {
156        self.index.len()
157    }
158
159    /// Take `i` out of the recency order, leaving the node in its slot.
160    fn unlink(&mut self, i: usize) {
161        let (newer, older) = {
162            let n = self.slots[i].as_ref().expect("unlink of a free slot");
163            (n.newer, n.older)
164        };
165        match newer {
166            Some(j) => self.slots[j].as_mut().expect("newer is live").older = older,
167            None => self.newest = older,
168        }
169        match older {
170            Some(j) => self.slots[j].as_mut().expect("older is live").newer = newer,
171            None => self.oldest = newer,
172        }
173        let n = self.slots[i].as_mut().expect("unlink of a free slot");
174        n.newer = None;
175        n.older = None;
176    }
177
178    /// Put `i` at the head of the recency order. It must not be linked.
179    fn link_newest(&mut self, i: usize) {
180        let old_head = self.newest;
181        {
182            let n = self.slots[i].as_mut().expect("link of a free slot");
183            n.newer = None;
184            n.older = old_head;
185        }
186        if let Some(j) = old_head {
187            self.slots[j].as_mut().expect("head is live").newer = Some(i);
188        } else {
189            self.oldest = Some(i);
190        }
191        self.newest = Some(i);
192    }
193
194    /// The block's data, promoted to most-recently-used. Two hash
195    /// lookups and a constant number of pointer writes, whatever the
196    /// capacity.
197    fn get(&mut self, block_start: u64) -> Option<Arc<Vec<u8>>> {
198        let i = *self.index.get(&block_start)?;
199        let data = self.slots[i]
200            .as_ref()
201            .expect("indexed slot is live")
202            .data
203            .clone();
204        if self.newest != Some(i) {
205            self.unlink(i);
206            self.link_newest(i);
207        }
208        Some(data)
209    }
210
211    /// Drop one block if it is held. Returns whether it was.
212    fn remove(&mut self, block_start: u64) -> bool {
213        let Some(i) = self.index.remove(&block_start) else {
214            return false;
215        };
216        self.unlink(i);
217        self.slots[i] = None;
218        self.free.push(i);
219        true
220    }
221
222    /// Insert at the head, evicting the least recently used first if
223    /// the cache is already at `capacity`.
224    ///
225    /// EVICT-THEN-INSERT UNCONDITIONALLY, which is the sequence the
226    /// `VecDeque` version used (`pop_back` under `len() >= capacity`,
227    /// then `push_front`). It matters at `capacity == 0`, where both
228    /// spellings leave exactly one entry held rather than none: a
229    /// behaviour worth preserving deliberately rather than changing
230    /// while moving house.
231    fn insert(&mut self, block_start: u64, data: Arc<Vec<u8>>, capacity: usize) {
232        // A RE-INSERT REPLACES RATHER THAN DUPLICATING. No caller does
233        // this today -- `block()` consults `get` under the same lock
234        // immediately before -- but a second slot for one block would
235        // leave the first linked in the recency list and unreachable
236        // through the index: a leak that also makes the list longer
237        // than the index, which is the kind of drift that shows up
238        // later as an eviction of something still held.
239        self.remove(block_start);
240        if self.len() >= capacity {
241            if let Some(oldest) = self.oldest {
242                let victim = self.slots[oldest]
243                    .as_ref()
244                    .expect("oldest is live")
245                    .block_start;
246                self.remove(victim);
247            }
248        }
249        let node = Node {
250            block_start,
251            data,
252            newer: None,
253            older: None,
254        };
255        let i = match self.free.pop() {
256            Some(i) => {
257                self.slots[i] = Some(node);
258                i
259            }
260            None => {
261                self.slots.push(Some(node));
262                self.slots.len() - 1
263            }
264        };
265        self.index.insert(block_start, i);
266        self.link_newest(i);
267    }
268
269    fn clear(&mut self) {
270        self.slots.clear();
271        self.free.clear();
272        self.index.clear();
273        self.newest = None;
274        self.oldest = None;
275    }
276
277    /// Drop every block for which `keep` is false.
278    ///
279    /// Linear in the number of entries, like the `retain` it replaces,
280    /// and deliberately so: this runs per WRITE, not per read, and the
281    /// blocks to drop have to be found by looking at all of them.
282    fn retain_blocks(&mut self, keep: impl Fn(u64) -> bool) {
283        let doomed: Vec<u64> = self.index.keys().copied().filter(|b| !keep(*b)).collect();
284        for b in doomed {
285            self.remove(b);
286        }
287    }
288
289    /// Most-recently-used first. Tests only: the recency ORDER is the
290    /// thing a linked list can get wrong in ways a hit rate cannot see.
291    #[cfg(test)]
292    fn recency_order(&self) -> Vec<u64> {
293        let mut out = Vec::with_capacity(self.len());
294        let mut cur = self.newest;
295        while let Some(i) = cur {
296            let n = self.slots[i].as_ref().expect("live");
297            out.push(n.block_start);
298            cur = n.older;
299        }
300        out
301    }
302
303    /// Least-recently-used first, walked the other way. Tests only, and
304    /// the reason it exists is that a singly-consistent list passes
305    /// every forward walk while being broken backwards -- which is the
306    /// half eviction uses.
307    #[cfg(test)]
308    fn recency_order_reversed(&self) -> Vec<u64> {
309        let mut out = Vec::with_capacity(self.len());
310        let mut cur = self.oldest;
311        while let Some(i) = cur {
312            let n = self.slots[i].as_ref().expect("live");
313            out.push(n.block_start);
314            cur = n.newer;
315        }
316        out
317    }
318}
319
320struct CacheState {
321    /// Fixed-capacity LRU; head is most-recently used. The capacity
322    /// itself is [`CachingDevice::capacity`] — it never changes, so it
323    /// is not kept under the lock.
324    ///
325    /// An index plus an intrusive recency list rather than a
326    /// `VecDeque`: see [`Lru`] for the measurement that decided it.
327    entries: Lru,
328    hits: u64,
329    misses: u64,
330    /// Bumped by every invalidation. A miss records it before it lets go
331    /// of the lock to read the device, and the insert on the way back in
332    /// is refused if it has moved — see `CachingDevice::block`.
333    generation: u64,
334    /// The blocks some thread is reading from the device right now.
335    ///
336    /// `generation` fences a miss against an *invalidation*; this fences
337    /// it against another *miss*. Nothing used to: two threads missing
338    /// the same block both read the device and both inserted, so a
339    /// capacity-4 cache could end up holding four copies of one block
340    /// having evicted the other three to make room. The more contention,
341    /// the worse it got, which is the opposite of what a cache is for.
342    ///
343    /// A `Vec` scanned linearly, like `entries` beside it. Its length is
344    /// the number of blocks being fetched concurrently, not the number
345    /// cached, so it is small in exactly the cases the scan would matter.
346    ///
347    /// EACH FETCH CARRIES THE THREAD DOING IT, so that a thread can ask
348    /// whether it is holding a fetch of its own before it waits for
349    /// anybody else's. Waiting for another thread is the point; waiting
350    /// while holding a fetch is how a cycle forms, whether the block
351    /// waited on is this thread's own or two links away round a ring of
352    /// re-entrant threads. See `block`.
353    in_flight: Vec<(u64, ThreadId)>,
354}
355
356impl CachingDevice {
357    /// Cache a device that can be written. Writes invalidate the
358    /// entries they overlap and go through to `inner`.
359    ///
360    /// The invalidation holds against concurrent readers as well as
361    /// sequential ones: once `write_at` has returned, no later read can
362    /// be served pre-write bytes from this cache, including a read that
363    /// was already in flight when the write began. What such a read
364    /// *returns* is still either side of the write — that is what racing
365    /// means — but it is not remembered.
366    ///
367    /// `block_size` must be non-zero and no larger than
368    /// [`MAX_BLOCK_SIZE`]; see [`CachingDevice::read_only`] for why that
369    /// is enforced at first use rather than here.
370    ///
371    /// `capacity` is documented on [`CachingDevice::read_only`]; it means
372    /// the same here, including that `0` still caches one block.
373    pub fn new(inner: Arc<dyn BlockDevice>, block_size: u64, capacity: usize) -> Arc<Self> {
374        Arc::new(Self {
375            inner: inner.clone(),
376            writable: Some(inner),
377            block_size,
378            capacity,
379            state: Mutex::new(CacheState {
380                entries: Lru::with_capacity(capacity),
381                hits: 0,
382                misses: 0,
383                generation: 0,
384                in_flight: Vec::new(),
385            }),
386            fetched: Condvar::new(),
387        })
388    }
389
390    /// Cache a device that is only ever read.
391    ///
392    /// The case every driver here actually has: a volume mounted for
393    /// reading, behind a `BlockRead` that was never a `BlockDevice`.
394    ///
395    /// `block_size` must be non-zero and no larger than
396    /// [`MAX_BLOCK_SIZE`]. Construction cannot refuse — it returns
397    /// `Arc<Self>`, not `Result` — so a block size outside that range is
398    /// refused by every read and every write instead, with an error
399    /// naming the offending size.
400    ///
401    /// `capacity` is a count of **blocks**, not bytes: at most that many
402    /// entries of up to `block_size` bytes each are held, so the memory
403    /// bound is their product and is the caller's to choose. There is no
404    /// ceiling; promotion and eviction are O(1) at any capacity.
405    ///
406    /// **`capacity = 0` does not disable the cache.** It holds one entry,
407    /// so a repeated single-block read is still served as a hit, and two
408    /// alternating blocks thrash it (#124). The behaviour is deliberate
409    /// and pinned, not an oversight. A caller that wants no cache at zero
410    /// should not construct one, and this constructor cannot do that for
411    /// it because it returns `Arc<Self>`:
412    ///
413    /// ```text
414    /// let dev: Arc<dyn BlockRead> = if blocks == 0 {
415    ///     dev
416    /// } else {
417    ///     CachingDevice::read_only(dev, block_size, blocks)
418    /// };
419    /// ```
420    ///
421    /// ```
422    /// # use std::sync::Arc;
423    /// # use fs_core::{BlockRead, CachingDevice, CountingDevice};
424    /// # struct Mem(Vec<u8>);
425    /// # impl BlockRead for Mem {
426    /// #     fn read_at(&self, offset: u64, buf: &mut [u8]) -> fs_core::Result<()> {
427    /// #         let o = offset as usize;
428    /// #         buf.copy_from_slice(&self.0[o..o + buf.len()]);
429    /// #         Ok(())
430    /// #     }
431    /// #     fn size_bytes(&self) -> u64 { self.0.len() as u64 }
432    /// # }
433    /// let device = Arc::new(CountingDevice::new(Arc::new(Mem(vec![7; 512 * 8]))));
434    /// let cache = CachingDevice::read_only(device.clone(), 512, 0);
435    /// let mut buf = [0u8; 16];
436    /// cache.read_at(0, &mut buf).unwrap();
437    /// cache.read_at(8, &mut buf).unwrap();
438    /// // A capacity of zero served the second read from the cache.
439    /// assert_eq!(device.reads(), 1);
440    /// assert_eq!(cache.stats(), (1, 1));
441    /// ```
442    pub fn read_only(inner: Arc<dyn BlockRead>, block_size: u64, capacity: usize) -> Arc<Self> {
443        Arc::new(Self {
444            inner,
445            writable: None,
446            block_size,
447            capacity,
448            state: Mutex::new(CacheState {
449                entries: Lru::with_capacity(capacity),
450                hits: 0,
451                misses: 0,
452                generation: 0,
453                in_flight: Vec::new(),
454            }),
455            fetched: Condvar::new(),
456        })
457    }
458
459    /// `(hits, misses)`, in that order.
460    ///
461    /// Both count **block lookups inside the cache**, not reads of this
462    /// device and not reads of `inner` (#123):
463    ///
464    /// - a hit is a block served from the cache, including one that was
465    ///   waiting on another thread's fetch of the same block;
466    /// - a miss is a block this cache fetched from `inner`, so `misses` is
467    ///   the number of fetches the cache itself made.
468    ///
469    /// A read the cache declines to serve moves **neither** counter. A
470    /// read reaching past the end of `inner`, and a read spanning enough
471    /// blocks that caching it would sweep the cache (more than one block,
472    /// and more than half of `capacity`), go straight to `inner`; an empty
473    /// read, or one refused for its block size, touches nothing. So `hits + misses` is not the number of `read_at` calls,
474    /// `misses` is a lower bound on the reads `inner` saw, and
475    /// `hits / (hits + misses)` is a rate over the reads the cache served,
476    /// which leaves out every read it chose not to. To count what reached
477    /// the device, put a [`CountingDevice`](crate::CountingDevice) under
478    /// the cache.
479    ///
480    /// ```
481    /// # use std::sync::Arc;
482    /// # use fs_core::{BlockRead, CachingDevice, CountingDevice};
483    /// # struct Mem(Vec<u8>);
484    /// # impl BlockRead for Mem {
485    /// #     fn read_at(&self, offset: u64, buf: &mut [u8]) -> fs_core::Result<()> {
486    /// #         let o = offset as usize;
487    /// #         buf.copy_from_slice(&self.0[o..o + buf.len()]);
488    /// #         Ok(())
489    /// #     }
490    /// #     fn size_bytes(&self) -> u64 { self.0.len() as u64 }
491    /// # }
492    /// let device = Arc::new(CountingDevice::new(Arc::new(Mem(vec![7; 512 * 32]))));
493    /// let cache = CachingDevice::read_only(device.clone(), 512, 8);
494    /// let mut one = [0u8; 16];
495    /// cache.read_at(0, &mut one).unwrap(); // miss: fetched
496    /// cache.read_at(8, &mut one).unwrap(); // hit
497    /// let mut big = vec![0u8; 512 * 6];
498    /// cache.read_at(0, &mut big).unwrap(); // 6 blocks > capacity / 2: bypassed
499    /// assert_eq!(cache.stats(), (1, 1));
500    /// assert_eq!(device.reads(), 2); // the bypassed read is not a miss
501    /// ```
502    pub fn stats(&self) -> (u64, u64) {
503        let s = self.state.lock().unwrap();
504        (s.hits, s.misses)
505    }
506
507    pub fn invalidate_all(&self) {
508        let mut s = self.state.lock().unwrap();
509        s.entries.clear();
510        s.generation = s.generation.wrapping_add(1);
511    }
512
513    fn invalidate_range(state: &mut CacheState, start: u64, end: u64, block_size: u64) {
514        state.entries.retain_blocks(|off| {
515            let block_end = off.saturating_add(block_size);
516            off >= end || block_end <= start
517        });
518        // BUMPED WHETHER OR NOT ANYTHING WAS DROPPED. The counter is not a
519        // record of what this sweep removed; it is a fence a concurrent
520        // miss can compare itself against, and a miss that is mid-flight
521        // over this range holds no entry for the sweep to find.
522        state.generation = state.generation.wrapping_add(1);
523    }
524
525    /// One invalidation sweep, taking and releasing the lock.
526    fn invalidate_for_write(&self, start: u64, end: u64) {
527        let mut s = self.state.lock().unwrap();
528        let bs = self.block_size;
529        Self::invalidate_range(&mut s, start, end, bs);
530    }
531
532    /// Refuse a block size the cache cannot work with.
533    ///
534    /// # Why this is checked here and not in the constructors
535    ///
536    /// It belongs in the constructors, and they cannot express it: both
537    /// return `Arc<Self>` rather than `Result`, and that signature is
538    /// published API in eleven sibling crates. Making them fallible to
539    /// catch a case no correct caller hits would be a breaking change to
540    /// all of them. So the refusal happens at first use instead, which
541    /// costs two comparisons against an immutable field per call and
542    /// turns both failures into an error the caller can handle.
543    ///
544    /// # Why a block size is worth checking at all
545    ///
546    /// Zero divides by zero on the first read, and integer division by
547    /// zero panics unconditionally — it is not governed by
548    /// `overflow-checks`, so a release build dies too. An absurd value
549    /// reaches `vec![0u8; len]` in `block()` with `len` bounded only by
550    /// the device. Both numbers come off a disk: a driver reads its block
551    /// size from a superblock and passes it through, so a truncated,
552    /// fuzzed or hostile image reaches this.
553    fn check_block_size(&self) -> Result<()> {
554        if self.block_size == 0 {
555            return Err(crate::error::Error::Custom(
556                "cache block size is zero".to_string(),
557            ));
558        }
559        if self.block_size > MAX_BLOCK_SIZE {
560            return Err(crate::error::Error::Custom(format!(
561                "cache block size {} exceeds the {MAX_BLOCK_SIZE}-byte ceiling",
562                self.block_size
563            )));
564        }
565        Ok(())
566    }
567}
568
569impl CachingDevice {
570    /// The cached block at `block_start`, fetching it if it is not held.
571    ///
572    /// # A MISS THAT OVERLAPPED AN INVALIDATION IS NOT CACHED
573    ///
574    /// The device read below runs with the lock released — holding a mutex
575    /// across I/O would serialise every reader, which is the whole reason
576    /// the lock is dropped. That leaves a window: a write can sweep the
577    /// cache and land on the device while this read is in flight, and the
578    /// bytes in hand are then the ones the device held *before* the write.
579    /// Inserting them puts a stale entry in a cache the sweep has already
580    /// gone past, and nothing would ever invalidate it again.
581    ///
582    /// So the miss records the generation counter before it lets go, and
583    /// declines to insert if any invalidation has happened since. The
584    /// caller still gets the bytes that were read — a read racing a write
585    /// may legitimately see either side of it — but the cache does not keep
586    /// them.
587    ///
588    /// The check is deliberately coarse: the counter is global rather than
589    /// per range, so a write to an unrelated block also costs this miss its
590    /// insert. Writes are far rarer than reads in these drivers, and the
591    /// price of being conservative is one extra device read.
592    fn block(&self, block_start: u64) -> Result<Arc<Vec<u8>>> {
593        // ONE FETCH PER BLOCK, NOT ONE PER MISSING THREAD.
594        //
595        // The lock has to be released for the device read -- holding a
596        // mutex across I/O would serialise every reader, which is what
597        // `perf: reads are positioned, and no longer take a lock`
598        // deliberately stopped doing. But releasing it used to mean N
599        // threads missing the same block all read the device and all
600        // inserted, so the cache held N copies of one block and had
601        // evicted up to N-1 others to store bytes it already had.
602        //
603        // So a thread that is going to fetch says so first, under the
604        // lock. A thread that finds someone else already fetching its
605        // block waits for them and then looks again, rather than making
606        // a second trip to the device for bytes already on their way.
607        //
608        // The loop is not a spin: `wait` blocks until a fetch finishes,
609        // and re-checking from the top afterwards is what makes it
610        // correct against a spurious wake, an invalidation that landed
611        // in the meantime, and a fetch that failed and left nothing.
612        //
613        // A THREAD THAT IS ALREADY FETCHING SOMETHING NEVER WAITS,
614        // WHOEVER OWNS THE BLOCK IT WANTS. WHICH IS WHY THE MARKER
615        // CARRIES AN OWNER.
616        //
617        // A device whose `read_at` reads back through the cache that
618        // wraps it re-enters this method while its own fetch is still
619        // outstanding. The narrow case is that it asks for the very
620        // block it is fetching, and waiting there is waiting for a
621        // fetch this thread is holding up: a deadlock. That shape is
622        // visible from the owner alone.
623        //
624        // THE GENERAL CASE IS NOT, and asking only "is this fetch
625        // mine" cannot see it. Thread A fetches X and re-enters for Y;
626        // thread B fetches Y and re-enters for X. Neither finds its own
627        // id against the block it wants, so both wait -- each for a
628        // fetch the other is holding up, on a condition variable only a
629        // completed fetch can signal. Nothing completes. It is the same
630        // deadlock one link longer, and there is no length at which
631        // comparing the owner to the caller starts to notice.
632        //
633        // So the question is not "is this fetch mine" but "am I holding
634        // one at all". A cycle needs every thread in it to be both
635        // holding a fetch and waiting for another one; a thread holding
636        // nothing cannot be waited on, so it cannot be in a cycle. A
637        // thread therefore waits only when it holds nothing, and that
638        // removes every cycle rather than the shortest one. It needs no
639        // wait-for graph and no new state to do it -- only a wider
640        // question asked of the list already being scanned.
641        //
642        // WHAT IT COSTS is a redundant device read in one narrow case:
643        // a re-entrant thread that wants a block another thread is
644        // fetching, where waiting would in fact have been safe. That is
645        // the same trade the self case already made, and only a device
646        // that re-enters can reach it. For every other caller nothing
647        // changes at all, because a thread holds a fetch here only
648        // while it is inside `read_at` on the device, and a device that
649        // cannot call back in leaves this false on every ordinary miss.
650        //
651        // THE GUARD IS HERE BECAUSE OF WHAT A HANG COSTS, not because
652        // re-entrancy is expected. It is not: no device in this crate
653        // can do it today, since none holds a handle to what wraps it,
654        // and stacked caches are unaffected because each has its own
655        // lock and its own list. But `CachingDevice` accepts any
656        // `BlockRead`, so the constraint is a property of the current
657        // devices rather than of this API -- a device with a slower
658        // tier behind it that consults a cache on miss would violate it
659        // the day it is written. And the drivers built on this crate
660        // run as filesystem extensions, where a deadlock inside a
661        // fetch is not a visible failure but a spinning cursor on a
662        // volume that cannot be unmounted, with no error and no log
663        // line to say which mount is stuck. Four lines and a thread id
664        // turn that into a redundant read.
665        let generation_at_miss = {
666            let mut s = self.state.lock().unwrap();
667            loop {
668                if let Some(data) = s.entries.get(block_start) {
669                    s.hits += 1;
670                    return Ok(data);
671                }
672                let mine = std::thread::current().id();
673                // Any fetch of this thread's, not just one of this
674                // block: holding any of them is what makes waiting
675                // unsafe. See the cycle argument above.
676                let holding_a_fetch = s.in_flight.iter().any(|(_, owner)| *owner == mine);
677                let being_fetched = s.in_flight.iter().any(|(o, _)| *o == block_start);
678                if being_fetched && !holding_a_fetch {
679                    // Counted as neither yet. It becomes a hit when the
680                    // fetch it is waiting for lands, and a miss if that
681                    // fetch fails and this thread has to do it instead
682                    // -- so `hits + misses` stays the number of calls,
683                    // and `misses` stays the number of device fetches.
684                    s = self.fetched.wait(s).unwrap();
685                    continue;
686                }
687                s.misses += 1;
688                s.in_flight.push((block_start, mine));
689                break s.generation;
690            }
691        };
692
693        // FROM HERE EVERY EXIT MUST CLEAR THE MARKER, including the `?`
694        // below and a panic inside the device. A fetch that vanished
695        // without clearing it would leave every later reader of that
696        // block waiting for a thread that is gone.
697        let _fetch = FetchGuard {
698            device: self,
699            block_start,
700            owner: std::thread::current().id(),
701        };
702
703        // THE LAST BLOCK OF A DEVICE IS OFTEN SHORT, and asking the
704        // device for a whole one past its end is an error rather than a
705        // short read. A SquashFS image is 4 KiB and its declared block
706        // size 128 KiB; without this clamp, caching such an image failed
707        // on the first read it ever made.
708        let size = self.inner.size_bytes();
709        let end = block_start.saturating_add(self.block_size).min(size);
710        let len = end.saturating_sub(block_start) as usize;
711        let mut block = vec![0u8; len];
712        self.inner.read_at(block_start, &mut block)?;
713        let data = Arc::new(block);
714
715        // Inserted BEFORE `_fetch` clears the marker and wakes the
716        // waiters, so a thread woken by this fetch finds the entry
717        // rather than an empty cache and a free marker.
718        {
719            let mut s = self.state.lock().unwrap();
720            // ALREADY HELD? Then return that copy and drop this one
721            // rather than holding the same block twice. Only a
722            // re-entrant fetch reaches this -- one fetch per block per
723            // thread is the rule for every other caller -- and without
724            // it a re-entrant read puts two entries in for one block,
725            // which is the defect this commit exists to remove, in a
726            // new place.
727            if let Some(held) = s.entries.get(block_start) {
728                return Ok(held);
729            }
730            if s.generation == generation_at_miss {
731                s.entries.insert(block_start, data.clone(), self.capacity);
732            }
733        }
734        Ok(data)
735    }
736}
737
738impl BlockRead for CachingDevice {
739    /// # A read is served from the blocks it falls in, whatever its size
740    ///
741    /// This used to serve a read only when it was **exactly one aligned
742    /// block**, and pass everything else through untouched — including
743    /// reads of bytes it was already holding.
744    ///
745    /// The drivers almost never read a whole block. Measured on
746    /// `am-fs-xfs` against a fixture with a 4096-byte block size, the
747    /// average read during a directory walk was **1040 bytes**: inodes
748    /// are read at inode size and group headers at sector size, so
749    /// roughly three quarters of reads missed by construction.
750    ///
751    /// # What it costs
752    ///
753    /// A 512-byte read of an uncached block now fetches 4096. That is a
754    /// trade of bytes for calls, and it is the right way round for these
755    /// drivers: the block being fetched is the one holding the inode,
756    /// and the next inode read is very often in it.
757    ///
758    /// # Where it still passes through
759    ///
760    /// A read larger than the cache's own capacity would evict
761    /// everything to hold one answer, so anything spanning more blocks
762    /// than a useful fraction of the cache goes straight to the device.
763    /// File data is read in large pieces and would otherwise push out
764    /// the metadata this exists to keep.
765    fn read_at(&self, offset: u64, buf: &mut [u8]) -> Result<()> {
766        self.check_block_size()?;
767        if buf.is_empty() {
768            return Ok(());
769        }
770
771        // THE END OF THE READ IS COMPUTED ONCE, CHECKED, AND BEFORE ANY
772        // DIVISION.
773        //
774        // This used to be `offset + buf.len()`, twice, unchecked, above
775        // the only bounds check the function has — which then made the
776        // same sum a third time with `saturating_add`, so the line that
777        // knew the sum could overflow sat below the two that did not.
778        //
779        // An offset here is computed by a driver from an on-disk field:
780        // an extent pointer, an inode block number, a directory offset.
781        // A wild one is an ordinary thing to be handed off a corrupt or
782        // hostile image rather than a mistake in the caller, which is
783        // the same argument `slice.rs` was fixed on.
784        //
785        // A sum that does not fit in a `u64` cannot name a byte on any
786        // device, so it is refused here rather than forwarded. `got: 0`
787        // because nothing was transferred — the convention the slice
788        // adapters already use for a read refused before it starts.
789        let Some(end) = offset.checked_add(buf.len() as u64) else {
790            return Err(crate::error::Error::ShortRead {
791                offset,
792                want: buf.len(),
793                got: 0,
794            });
795        };
796
797        // A READ RUNNING PAST THE END OF THE DEVICE IS THE DEVICE'S TO
798        // REFUSE. Serving it from clamped blocks would hand back a short
799        // answer with no error, which is worse than the failure the
800        // caller would otherwise have seen.
801        if end > self.inner.size_bytes() {
802            return self.inner.read_at(offset, buf);
803        }
804
805        let bs = self.block_size;
806        let first = offset / bs;
807        let last = (end - 1) / bs;
808        let spanned = (last - first + 1) as usize;
809
810        // A read big enough to sweep the cache is not worth caching.
811        //
812        // A SINGLE BLOCK IS NEVER "BIG ENOUGH", however small the cache.
813        // Without that clause a cache of one block bypasses every read
814        // it is ever given -- one block is more than half of one block --
815        // so the smallest cache anybody can ask for is the one that
816        // silently does nothing.
817        //
818        // AND THE TEST TAKES NO LOCK. `capacity` is a plain field on the
819        // device, so a read that is about to bypass the cache decides
820        // that without ever contending for the mutex -- see the field's
821        // own comment.
822        if spanned > 1 && spanned.saturating_mul(2) > self.capacity {
823            return self.inner.read_at(offset, buf);
824        }
825
826        let mut done = 0usize;
827        for index in first..=last {
828            let block_start = index * bs;
829            let block = self.block(block_start)?;
830            // Where this block overlaps what was asked for.
831            let from = (offset.max(block_start) - block_start) as usize;
832            // Bounded by what the BLOCK holds rather than by the block
833            // size, since the last one may be short.
834            let take = (block.len().saturating_sub(from)).min(buf.len() - done);
835            if take == 0 {
836                break;
837            }
838            buf[done..done + take].copy_from_slice(&block[from..from + take]);
839            done += take;
840        }
841
842        // A BUFFER THIS DID NOT FILL IS A FAILURE, NOT A SUCCESS.
843        //
844        // The loop can run out of bytes before the buffer is full: a
845        // cached block is clamped to the device's size when it is
846        // fetched, so a block held from when the device measured smaller
847        // is short, and the guard above — which asks the device its size
848        // a second time — lets the read through as in range. Reporting
849        // `Ok` there hands back exactly the short answer with no error
850        // that the guard's own comment refuses.
851        //
852        // It cannot fire while `size_bytes` is stable, which is the
853        // contract `BlockRead` now states: the guard bounds the read by
854        // the device, and the blocks between them cover everything up to
855        // that bound. So this is the device breaking its promise being
856        // caught rather than believed.
857        //
858        // `BlockDevice::set_len` IS THE ONE SANCTIONED WAY THAT NUMBER
859        // MOVES, and this is still the right answer when a read races
860        // one. `set_len` sweeps the blocks its new length makes short,
861        // either side of the device call, so a read that is not
862        // concurrent with one can never land here; a read that IS
863        // concurrent with one may, and an error naming what it got beats
864        // a short buffer reported as success. The defect this catches is
865        // a `set_len` that moved the length and swept nothing, which is
866        // exactly how a growth API re-creates #70 -- see this type's
867        // `set_len`.
868        if done != buf.len() {
869            return Err(crate::error::Error::ShortRead {
870                offset,
871                want: buf.len(),
872                got: done,
873            });
874        }
875        Ok(())
876    }
877
878    fn size_bytes(&self) -> u64 {
879        self.inner.size_bytes()
880    }
881}
882
883impl BlockDevice for CachingDevice {
884    /// # THE CACHE IS INVALIDATED EVEN IF THE WRITE THEN FAILS
885    ///
886    /// Deliberately: dropping entries the write would have made stale
887    /// costs a re-read, while keeping them past a write that half
888    /// succeeded serves bytes the device no longer holds.
889    ///
890    /// That applies to a write the device refused, not to one this type
891    /// refused on the device's behalf. A write rejected because there is
892    /// no writable half, or because the block size is unusable, never
893    /// reaches the device and so cannot have staled anything — those
894    /// return above both sweeps and leave the cache exactly as it was.
895    ///
896    /// # AND IT IS INVALIDATED TWICE, ONCE EITHER SIDE OF THE DEVICE
897    ///
898    /// One sweep before the write is not enough. Between it and the
899    /// device write landing, a concurrent [`CachingDevice::read_at`] can
900    /// miss, fetch pre-write bytes, and insert them behind the sweep —
901    /// an entry the sweep has already gone past and nothing else would
902    /// ever drop. The second sweep is what closes that window, together
903    /// with the generation check on the miss path that refuses such an
904    /// insert outright.
905    ///
906    /// Both are needed, and each covers what the other cannot. A read
907    /// that began BEFORE this write is caught by the counter, because it
908    /// recorded the generation before the first sweep bumped it. A read
909    /// that begins AFTER that sweep records the bumped value, so the
910    /// counter agrees with it and only the second sweep drops what it
911    /// inserted.
912    fn write_at(&self, offset: u64, buf: &[u8]) -> Result<()> {
913        // NO WRITABLE HALF IS ANSWERED FIRST, ABOVE EVERY OTHER REFUSAL.
914        //
915        // This used to sit below the block-size check and below the first
916        // sweep, and both positions were wrong for their own reason.
917        //
918        // Below the block-size check, a read-only cache whose block size
919        // came off a damaged disk answered a write with `Error::Custom`
920        // rather than the `Error::ReadOnly` this type's own documentation
921        // promises. That is not a cosmetic difference in the variant name:
922        // `stream.rs` maps `ReadOnly` to `PermissionDenied` and `Custom` to
923        // `io::Error::other`, so a caller branching on `PermissionDenied`
924        // to say "this volume is read-only" reported an uncategorised
925        // failure instead — on a read-only mount of a damaged image, which
926        // is exactly the configuration where a bad block size turns up.
927        //
928        // Below the first sweep, every refused write emptied the region it
929        // was refused for, so a read-only cache paid a re-read for a write
930        // that could never have staled anything.
931        //
932        // The block-size check does NOT move down to meet it. See below.
933        let Some(writable) = self.writable.as_ref() else {
934            return Err(crate::error::Error::ReadOnly);
935        };
936
937        // BEFORE EITHER SWEEP, AND THAT ORDERING IS LOAD-BEARING.
938        //
939        // A sweep cannot be done correctly with a block size of zero:
940        // `invalidate_range` computes `block_end = off + block_size`, so a
941        // zero block size makes `block_end == off` and turns the retain
942        // predicate into `*off >= end || *off <= start`, which KEEPS
943        // entries the write has made stale. Sweeping first and refusing
944        // afterwards would therefore do the one thing this type must never
945        // do, on the way to reporting an error.
946        //
947        // Refusing here is also a position that cannot break the two-sweep
948        // guarantee above. The invariant that guarantee rests on is "if the
949        // device was written, both sweeps ran" — and returning at this point
950        // means the device is never reached, so nothing was written and
951        // nothing needs sweeping. An early return anywhere BELOW would skip
952        // the second sweep after a write that may have landed, which is
953        // exactly the window that fix closed. The read-only refusal above
954        // is safe for the same reason and no other: it too returns before
955        // either sweep and touches neither the cache nor the device.
956        self.check_block_size()?;
957        let end = offset.saturating_add(buf.len() as u64);
958        self.invalidate_for_write(offset, end);
959        let result = writable.write_at(offset, buf);
960        // Unconditionally, for the same reason the first sweep is
961        // unconditional: a write that failed may still have landed.
962        self.invalidate_for_write(offset, end);
963        result
964    }
965
966    fn flush(&self) -> Result<()> {
967        match self.writable.as_ref() {
968            Some(writable) => writable.flush(),
969            // Nothing was written, so there is nothing to flush. An
970            // error here would make a caller that flushes defensively
971            // fail on a read-only volume.
972            None => Ok(()),
973        }
974    }
975
976    fn is_writable(&self) -> bool {
977        self.writable.as_ref().is_some_and(|w| w.is_writable())
978    }
979
980    /// # THE CACHE'S VIEW HAS TO MOVE WITH THE DEVICE, OR THIS IS #70
981    ///
982    /// `size_bytes` here forwards to the device, so the NUMBER follows a
983    /// grow for free. The entries do not. `block()` clamps every fetch to
984    /// `size_bytes()` at the moment it runs — "the last block of a device
985    /// is often short", as its own comment says — so an entry fetched
986    /// before the grow ends where the device used to. Leave it in place
987    /// and a later read across the old end is served from it, runs out of
988    /// bytes, and comes back `ShortRead` for a region the device now
989    /// holds perfectly well.
990    ///
991    /// Measured, with this method forwarding to the device and sweeping
992    /// nothing: a 6000-byte file behind a 4096-byte cache, the short
993    /// block warmed, `set_len(8192)`, then one read across the old end:
994    ///
995    /// ```text
996    /// ShortRead { offset: 5000, want: 3000, got: 1000 }
997    /// ```
998    ///
999    /// That is rust-fs-core#70's signature exactly — the grow succeeded,
1000    /// the file and the reported size agreed, and only a LATER CACHED
1001    /// READ found the hole. It is why the growth API is a change to this
1002    /// file as much as to `block.rs`.
1003    ///
1004    /// # FROM `min(old, new)` UPWARDS, WHICHEVER WAY THE LENGTH WENT
1005    ///
1006    /// A grow only makes the block STRADDLING the old end wrong, and that
1007    /// block starts below the old end — so the sweep has to begin at the
1008    /// old length, not at the first block boundary above it, and
1009    /// `invalidate_range` drops any block whose end passes `start`.
1010    ///
1011    /// A shrink makes everything from the new length up wrong instead.
1012    /// Taking the smaller of the two covers both without asking which
1013    /// happened, and the upper bound is `u64::MAX` because "the rest of
1014    /// the device" is what changed in either case.
1015    ///
1016    /// # `min` RATHER THAN `old`, AND NO TEST HERE CAN TELL THEM APART
1017    ///
1018    /// Said plainly because the alternative is a comment claiming a
1019    /// guarantee nobody measured. Sweeping from the OLD length alone
1020    /// leaves the blocks between the two lengths cached after a shrink,
1021    /// and that suite is EXIT=0 -- every arm in `tests/device_growth.rs`
1022    /// passes with `min` removed.
1023    ///
1024    /// It passes because nothing can read those entries. `read_at`
1025    /// forwards any read whose end passes `size_bytes()` straight to the
1026    /// device rather than serving it from blocks, so while the device is
1027    /// short they are unreachable; and a later grow sweeps from the
1028    /// smaller of ITS two lengths, which is the shrunk one, so they are
1029    /// dropped before they become reachable again.
1030    ///
1031    /// `min` ships anyway, and not for symmetry. The argument above rests
1032    /// on a bound in a DIFFERENT METHOD holding forever -- an entry that
1033    /// is stale but currently unreadable is one guard away from being
1034    /// stale and readable. This is the cheaper half of the invariant to
1035    /// state correctly, so it is stated correctly here rather than
1036    /// derived from somewhere else on every future read of this file.
1037    ///
1038    /// # AND IT IS SWEPT TWICE, ONCE EITHER SIDE, FOR `write_at`'S REASON
1039    ///
1040    /// One sweep before is not enough. Between it and the device call
1041    /// landing, a concurrent [`CachingDevice::read_at`] can miss, fetch a
1042    /// block clamped to the OLD length, and insert it behind the sweep.
1043    /// The second sweep closes that window, together with the generation
1044    /// check on the miss path. Each covers what the other cannot, in the
1045    /// same way and for the same reason `write_at` documents at length.
1046    ///
1047    /// Both run whether or not the device call succeeded, also for
1048    /// `write_at`'s reason: a `set_len` that failed may still have moved
1049    /// the file, and dropping entries needlessly costs a re-read while
1050    /// keeping stale ones serves bytes the device no longer has.
1051    ///
1052    /// # THE TWO REFUSALS ABOVE THE SWEEPS
1053    ///
1054    /// No writable half, and an unusable block size — the same pair
1055    /// `write_at` refuses on, in the same order, and above both sweeps
1056    /// for the same two reasons. `Error::ReadOnly` rather than
1057    /// `Error::Custom` when there is nothing to write, because
1058    /// [`crate::stream`] maps only the first to `PermissionDenied`; and
1059    /// the block-size check above the sweeps because `invalidate_range`
1060    /// cannot sweep correctly with a block size of zero, so sweeping
1061    /// first and refusing after would do the one thing this type must
1062    /// never do on the way to reporting an error. Neither refusal reaches
1063    /// the device, so neither can have staled anything.
1064    fn set_len(&self, new_len: u64) -> Result<()> {
1065        let Some(writable) = self.writable.as_ref() else {
1066            return Err(crate::error::Error::ReadOnly);
1067        };
1068        self.check_block_size()?;
1069
1070        // The lower of the two lengths: below it nothing changed, at or
1071        // above it everything may have.
1072        let from = self.inner.size_bytes().min(new_len);
1073        self.invalidate_for_write(from, u64::MAX);
1074        let result = writable.set_len(new_len);
1075        self.invalidate_for_write(from, u64::MAX);
1076        result
1077    }
1078
1079    /// The writable half's answer, or `false` when there is no writable
1080    /// half — a cache cannot grow a device it can only read.
1081    fn can_grow(&self) -> bool {
1082        self.writable.as_ref().is_some_and(|w| w.can_grow())
1083    }
1084}
1085
1086/// Clears one block from [`CacheState::in_flight`] and wakes whoever is
1087/// waiting for it.
1088///
1089/// A guard rather than a line at the end of `block()`, because the
1090/// paths out of that method are not all the happy one: the device read
1091/// uses `?`, and a device is free to panic. Either would step over a
1092/// manual cleanup and strand every future reader of that block on a
1093/// fetch that no longer exists.
1094struct FetchGuard<'a> {
1095    device: &'a CachingDevice,
1096    block_start: u64,
1097    owner: ThreadId,
1098}
1099
1100impl Drop for FetchGuard<'_> {
1101    fn drop(&mut self) {
1102        // REMOVES ONE ENTRY, NOT EVERY MATCH. `retain` would be wrong:
1103        // it reads as the same thing and would clear a second fetch of
1104        // this block that this guard does not own.
1105        //
1106        // A poisoned lock is left alone. Another thread panicked holding
1107        // it, so the cache is already unusable and every waiter will
1108        // surface the poison from its own `wait`; unwrapping here could
1109        // panic while a panic is already unwinding, which aborts.
1110        if let Ok(mut s) = self.device.state.lock() {
1111            if let Some(pos) = s
1112                .in_flight
1113                .iter()
1114                .position(|(o, owner)| *o == self.block_start && *owner == self.owner)
1115            {
1116                s.in_flight.remove(pos);
1117            }
1118        }
1119        // Outside the lock, and after it: waiters re-check the state, so
1120        // waking them before the marker was cleared would send them
1121        // straight back to sleep.
1122        self.device.fetched.notify_all();
1123    }
1124}
1125
1126#[cfg(test)]
1127mod tests {
1128    use super::*;
1129    use crate::test_device::Bytes;
1130
1131    const BS: u64 = 512;
1132
1133    fn backing() -> Arc<Bytes> {
1134        Arc::new(Bytes::new((0..4096u32).map(|i| i as u8).collect()))
1135    }
1136
1137    /// THE CASE THAT COULD NOT BE EXPRESSED BEFORE: a device that is
1138    /// only ever read, wrapped in a cache.
1139    ///
1140    /// Every driver in this family mounts through a `BlockRead`, so
1141    /// this is not an exotic configuration — it is the ordinary one,
1142    /// and requiring `BlockDevice` is why four of the six drivers used
1143    /// no cache at all.
1144    #[test]
1145    fn a_read_only_device_can_be_cached() {
1146        let inner = backing();
1147        let cache = CachingDevice::read_only(inner, BS, 4);
1148
1149        let mut first = vec![0u8; BS as usize];
1150        let mut again = vec![0u8; BS as usize];
1151        cache.read_at(0, &mut first).expect("first read");
1152        cache.read_at(0, &mut again).expect("second read");
1153
1154        assert_eq!(first, again, "the cache must serve what the device held");
1155        assert_eq!(cache.stats(), (1, 1), "one hit after one miss");
1156    }
1157
1158    /// THE CASE THE OLD HIT CONDITION MISSED: a read smaller than a
1159    /// block, of a block already held.
1160    ///
1161    /// Serving only exact aligned blocks meant the drivers' ordinary
1162    /// reads — an inode at inode size, a group header at sector size —
1163    /// went to the device every time, even when the block containing
1164    /// them was cached.
1165    #[test]
1166    fn a_read_smaller_than_a_block_is_served_from_it() {
1167        let cache = CachingDevice::read_only(backing(), BS, 8);
1168
1169        let mut whole = vec![0u8; BS as usize];
1170        cache.read_at(0, &mut whole).expect("warm the block");
1171        assert_eq!(cache.stats(), (0, 1), "one miss to fetch it");
1172
1173        // Four sub-block reads inside the block just fetched.
1174        for at in [0u64, 8, 100, 504] {
1175            let mut small = [0u8; 8];
1176            cache.read_at(at, &mut small).expect("sub-block read");
1177            assert_eq!(
1178                &small[..],
1179                &whole[at as usize..at as usize + 8],
1180                "the bytes must be the block's own, at the right offset"
1181            );
1182        }
1183        assert_eq!(cache.stats(), (4, 1), "four hits, and no further misses");
1184    }
1185
1186    /// A read crossing a block boundary is stitched from both blocks,
1187    /// and each is cached.
1188    #[test]
1189    fn a_read_spanning_two_blocks_is_stitched() {
1190        let inner = backing();
1191        let mut direct = vec![0u8; 16];
1192        inner.read_at(BS - 8, &mut direct).expect("read it plainly");
1193
1194        let cache = CachingDevice::read_only(backing(), BS, 8);
1195        let mut across = vec![0u8; 16];
1196        cache.read_at(BS - 8, &mut across).expect("spanning read");
1197
1198        assert_eq!(across, direct, "the same bytes the device would give");
1199        assert_eq!(cache.stats(), (0, 2), "one miss per block touched");
1200
1201        cache.read_at(BS - 8, &mut across).expect("again");
1202        assert_eq!(cache.stats(), (2, 2), "and both are held now");
1203    }
1204
1205    /// A read big enough to sweep the cache goes straight to the device.
1206    ///
1207    /// File data arrives in large pieces, and caching it would evict the
1208    /// metadata this exists to hold — the opposite of the point.
1209    #[test]
1210    fn a_read_that_would_sweep_the_cache_passes_through() {
1211        let cache = CachingDevice::read_only(backing(), BS, 4);
1212        let mut big = vec![0u8; (BS * 4) as usize];
1213        cache.read_at(0, &mut big).expect("a large read");
1214        assert_eq!(
1215            cache.stats(),
1216            (0, 0),
1217            "neither hit nor miss: it never consulted the cache"
1218        );
1219    }
1220
1221    /// A device smaller than one block still reads.
1222    ///
1223    /// THE CASE THAT BROKE. `am-fs-squashfs` declares a block size from
1224    /// the archive's superblock -- 128 KiB is the usual -- and a small
1225    /// image is a few kilobytes whole. Fetching "the block at zero"
1226    /// asked the device for 128 KiB it did not have, which is an error
1227    /// rather than a short read, so opening such an image with a cache
1228    /// failed on the very first read.
1229    #[test]
1230    fn a_device_shorter_than_a_block_still_reads() {
1231        let tiny: Arc<Bytes> = Arc::new(Bytes::new((0..100u32).map(|i| i as u8).collect()));
1232        let cache = CachingDevice::read_only(tiny, BS, 4);
1233
1234        let mut buf = vec![0u8; 40];
1235        cache
1236            .read_at(10, &mut buf)
1237            .expect("a read inside the device");
1238        assert_eq!(buf[0], 10, "the wrong bytes came back");
1239        assert_eq!(buf[39], 49);
1240
1241        // And the second one is a hit, so the short block was cached
1242        // rather than merely tolerated.
1243        cache.read_at(10, &mut buf).expect("again");
1244        assert_eq!(cache.stats(), (1, 1));
1245    }
1246
1247    /// A read running past the end of the device still fails.
1248    ///
1249    /// The clamp above must not turn "you asked for bytes that are not
1250    /// there" into a short answer with no error. That failure is
1251    /// invisible to the caller, which is the one kind this family of
1252    /// crates refuses to produce.
1253    #[test]
1254    fn a_read_past_the_end_is_still_an_error() {
1255        let tiny: Arc<Bytes> = Arc::new(Bytes::new(vec![0u8; 100]));
1256        let cache = CachingDevice::read_only(tiny, BS, 4);
1257
1258        let mut buf = vec![0u8; 40];
1259        assert!(
1260            cache.read_at(80, &mut buf).is_err(),
1261            "80 + 40 is past the end of a 100-byte device"
1262        );
1263    }
1264
1265    /// A READ THAT BYPASSES THE CACHE DOES NOT WAIT FOR THE CACHE LOCK.
1266    ///
1267    /// The state lock is held for the whole of the read, by a thread
1268    /// that is not doing the read. With the `capacity` comparison under
1269    /// that mutex the bypass cannot proceed and the receive times out;
1270    /// with it outside there is nothing to wait for. No timing
1271    /// threshold and no thread count, so nothing here is flaky on a
1272    /// loaded machine.
1273    #[test]
1274    fn a_bypassed_read_does_not_wait_for_the_cache_lock() {
1275        use std::sync::mpsc;
1276        use std::thread;
1277        use std::time::Duration;
1278
1279        // Capacity 4 and a read spanning 4 blocks: 4 * 2 > 4, so this
1280        // read bypasses -- the same read
1281        // `a_read_that_would_sweep_the_cache_passes_through` makes.
1282        let cache = CachingDevice::read_only(backing(), BS, 4);
1283        let held = cache.state.lock().expect("nothing else holds it yet");
1284
1285        let reader = Arc::clone(&cache);
1286        let (tx, rx) = mpsc::channel();
1287        thread::spawn(move || {
1288            let mut big = vec![0u8; (BS * 4) as usize];
1289            let outcome = reader.read_at(0, &mut big);
1290            // Sent whatever it is: the point is that the read RETURNED.
1291            let _ = tx.send(outcome);
1292        });
1293
1294        let outcome = rx.recv_timeout(Duration::from_secs(5)).expect(
1295            "a read that bypasses the cache must not block on the cache lock; \
1296             it timed out waiting for a mutex it has no reason to take",
1297        );
1298        outcome.expect("and the bypassed read itself must succeed");
1299
1300        // Released only now, and it has to be explicit: `stats()` takes
1301        // the same lock, so leaving the guard to the end of scope would
1302        // deadlock the assertion below.
1303        drop(held);
1304
1305        assert_eq!(
1306            cache.stats(),
1307            (0, 0),
1308            "it bypassed, so it neither hit nor missed"
1309        );
1310    }
1311
1312    /// A cache over a read-only device says so, and refuses a write
1313    /// with the answer the device underneath would have given.
1314    #[test]
1315    fn writing_through_a_read_only_cache_is_refused() {
1316        let cache = CachingDevice::read_only(backing(), BS, 4);
1317        assert!(!cache.is_writable());
1318        assert!(matches!(
1319            cache.write_at(0, &[1u8; 8]),
1320            Err(crate::error::Error::ReadOnly)
1321        ));
1322        // And flushing is not an error: a caller that flushes
1323        // defensively must not fail on a volume it never wrote.
1324        assert!(cache.flush().is_ok());
1325    }
1326}
1327
1328#[cfg(test)]
1329mod lru_tests {
1330    use super::*;
1331
1332    fn block(n: u8) -> Arc<Vec<u8>> {
1333        Arc::new(vec![n; 4])
1334    }
1335
1336    /// The list has to be walkable BOTH WAYS and agree with itself.
1337    ///
1338    /// A singly-consistent list passes every forward walk while being
1339    /// broken backwards -- and backwards is the half eviction uses, so
1340    /// the failure would surface as evicting the wrong block rather
1341    /// than as anything a hit rate could show.
1342    fn assert_consistent(lru: &Lru) {
1343        let forward = lru.recency_order();
1344        let mut backward = lru.recency_order_reversed();
1345        backward.reverse();
1346        assert_eq!(
1347            forward, backward,
1348            "the recency list disagrees with itself walked the other way"
1349        );
1350        assert_eq!(
1351            forward.len(),
1352            lru.len(),
1353            "the list holds {} nodes and the index {} -- they have drifted",
1354            forward.len(),
1355            lru.len()
1356        );
1357    }
1358
1359    #[test]
1360    fn a_hit_promotes_to_most_recently_used() {
1361        let mut lru = Lru::with_capacity(4);
1362        for i in 0..4u64 {
1363            lru.insert(i * 100, block(i as u8), 4);
1364        }
1365        assert_eq!(lru.recency_order(), vec![300, 200, 100, 0]);
1366        assert!(lru.get(100).is_some());
1367        assert_eq!(lru.recency_order(), vec![100, 300, 200, 0]);
1368        assert_consistent(&lru);
1369
1370        // Promoting what is already newest must not corrupt the ends.
1371        assert!(lru.get(100).is_some());
1372        assert_eq!(lru.recency_order(), vec![100, 300, 200, 0]);
1373        assert_consistent(&lru);
1374
1375        // Nor must promoting the oldest.
1376        assert!(lru.get(0).is_some());
1377        assert_eq!(lru.recency_order(), vec![0, 100, 300, 200]);
1378        assert_consistent(&lru);
1379    }
1380
1381    #[test]
1382    fn the_least_recently_used_is_what_gets_evicted() {
1383        let mut lru = Lru::with_capacity(3);
1384        for i in 0..3u64 {
1385            lru.insert(i * 100, block(i as u8), 3);
1386        }
1387        // Touch the oldest so it is no longer the victim.
1388        assert!(lru.get(0).is_some());
1389        lru.insert(999, block(9), 3);
1390        assert_eq!(lru.len(), 3);
1391        assert_eq!(
1392            lru.recency_order(),
1393            vec![999, 0, 200],
1394            "100 was least recently used and is the one that should be gone"
1395        );
1396        assert!(lru.get(100).is_none());
1397        assert_consistent(&lru);
1398    }
1399
1400    #[test]
1401    fn evicted_slots_are_reused_rather_than_growing_the_slab() {
1402        let mut lru = Lru::with_capacity(2);
1403        for i in 0..20u64 {
1404            lru.insert(i, block(i as u8), 2);
1405            assert_consistent(&lru);
1406        }
1407        assert_eq!(lru.len(), 2);
1408        assert!(
1409            lru.slots.len() <= 3,
1410            "twenty inserts at capacity 2 left {} slots: freed slots are not being \
1411             reused, so the slab grows without bound",
1412            lru.slots.len()
1413        );
1414    }
1415
1416    #[test]
1417    fn a_re_insert_replaces_and_does_not_leave_the_old_node_linked() {
1418        let mut lru = Lru::with_capacity(4);
1419        lru.insert(10, block(1), 4);
1420        lru.insert(20, block(2), 4);
1421        lru.insert(10, block(3), 4);
1422        assert_eq!(lru.len(), 1 + 1, "one entry per block, not one per insert");
1423        assert_eq!(lru.recency_order(), vec![10, 20]);
1424        assert_eq!(
1425            lru.get(10).as_deref().map(|v| v[0]),
1426            Some(3),
1427            "the newer data wins"
1428        );
1429        assert_consistent(&lru);
1430    }
1431
1432    #[test]
1433    fn removing_and_retaining_keep_the_list_consistent() {
1434        let mut lru = Lru::with_capacity(8);
1435        for i in 0..8u64 {
1436            lru.insert(i, block(i as u8), 8);
1437        }
1438        assert!(lru.remove(0), "the tail");
1439        assert_consistent(&lru);
1440        assert!(lru.remove(7), "the head");
1441        assert_consistent(&lru);
1442        assert!(lru.remove(4), "the middle");
1443        assert_consistent(&lru);
1444        assert!(!lru.remove(4), "already gone");
1445
1446        lru.retain_blocks(|b| b % 2 == 0);
1447        assert_consistent(&lru);
1448        assert_eq!(lru.recency_order(), vec![6, 2]);
1449
1450        lru.clear();
1451        assert_eq!(lru.len(), 0);
1452        assert!(lru.recency_order().is_empty());
1453        assert_consistent(&lru);
1454    }
1455
1456    /// `capacity == 0` holds exactly one entry, not none.
1457    ///
1458    /// That is what the `VecDeque` version did -- `pop_back` on an
1459    /// empty deque is a no-op, then `push_front` -- and it is
1460    /// preserved deliberately rather than changed while moving house.
1461    /// Pinned so the next person to touch `insert` finds out from a
1462    /// test rather than from a caller.
1463    #[test]
1464    fn capacity_zero_behaves_as_it_did_before() {
1465        let mut lru = Lru::with_capacity(0);
1466        lru.insert(1, block(1), 0);
1467        assert_eq!(lru.len(), 1);
1468        lru.insert(2, block(2), 0);
1469        assert_eq!(lru.len(), 1);
1470        assert_eq!(lru.recency_order(), vec![2]);
1471        assert_consistent(&lru);
1472    }
1473}