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}