Skip to main content

kernel/
pool.rs

1//! A fixed array of frames, allocated once and never grown.
2//!
3//! Eviction is CLOCK: one reference bit per frame, a hand that sweeps. LRU's
4//! per-entry links are pure overhead at this budget.
5//!
6//! SACRIFICE (Law 4): CLOCK approximates LRU, so a pathological access order can
7//! evict a page that LRU would have kept. Bought: one bit per frame instead of
8//! two pointers, and no list surgery on every hit.
9//!
10//! Each frame carries a pin count plus a `writer` flag. The count keeps a
11//! frame out of `victim`'s reach while any guard references it; the flag
12//! additionally enforces reader/writer exclusivity between `get` and
13//! `get_mut` on the same page, the same way `RefCell` enforces it between
14//! `borrow` and `borrow_mut` — a conflicting request panics rather than
15//! handing out a second `&mut [u8]` over memory another guard already
16//! points into. Without that flag, `get` and `get_mut` on the same page_no
17//! could alias a live `&[u8]`/`&mut [u8]` pair, which the `unsafe fn`s on
18//! `AlignedRegion` exist specifically to rule out.
19
20use crate::budget::{Class, MemoryBudget, Reservation};
21use crate::io::{AlignedRegion, Barrier, FileIo};
22use crate::page::PAGE_SIZE;
23use crate::{Error, Result};
24use std::cell::RefCell;
25use std::collections::HashMap;
26use std::sync::Arc;
27
28const FREE_MAGIC: [u8; 8] = *b"SEKFREE\0";
29const FREE_VERSION: u16 = 2;
30const FREE_HEADER_LEN: usize = 24;
31
32#[derive(Debug, Default, Clone, Copy)]
33pub struct PoolStats {
34    pub hits: u64, pub misses: u64, pub evictions: u64,
35    /// Total frames the pool reserved at open -- its CAPACITY, not its
36    /// residency. The region is allocated once and never grown, so this is
37    /// the anonymous memory the pool holds from open onward whether or not
38    /// the frames are populated.
39    ///
40    /// It was called `frames_used`, and that name cost a reviewer a false
41    /// finding: 170 MB "used" against a 5 MB database reads as an accounting
42    /// bug rather than as a preallocated arena working exactly as designed. A
43    /// number that invites a wrong inference is a reporting defect even when
44    /// the arithmetic is right. See also `a_full_scan_pins_one_leaf_at_a_time`
45    /// in btree.rs, which independently found this field cannot move and so
46    /// cannot be asserted on for any specific operation -- `peak_pins` is the
47    /// field for that.
48    pub frames_total: usize,
49    pub peak_pins: u32,
50    /// Barriers actually placed by `flush_all`, counted at the point they are
51    /// issued rather than where they are chosen. Asserting on what a caller
52    /// like `Store::checkpoint` *passed* would only prove it passed
53    /// something -- these live here because this is where the syscall
54    /// actually happens, so a caller that silently ignored its own argument
55    /// could not still make the test pass.
56    pub sync_data_calls: u64,
57    pub sync_full_calls: u64,
58    /// Dirty frames written by `flush_all` (monotonic since pool construction).
59    pub dirty_pages_flushed: u64,
60}
61
62struct Frame {
63    page_no: u32,
64    present: bool,
65    dirty: bool,
66    referenced: bool,
67    /// Number of live PinnedRead/PinnedWrite guards over this frame.
68    pins: u32,
69    /// True while one of those live guards is a PinnedWrite. Invariant:
70    /// `writer` implies `pins > 0`.
71    writer: bool,
72    /// True once the slot directory passed full validation during THIS
73    /// residency (2f-A1). Cleared on every load from disk; our own PageMut
74    /// writes maintain the slot invariants, so they do not clear it.
75    validated: bool,
76}
77
78struct Inner {
79    sweep_steps: u64,
80    frames: Vec<Frame>,
81    table: HashMap<u32, usize>,
82    hand: usize,
83    stats: PoolStats,
84    next_page: u32,
85    limits: Option<crate::limits::ResourceLimits>,
86    free_count: usize,
87    epoch_allocated_pages: u64,
88    /// Generation stamped into every page sealed to disk (2n): the epoch
89    /// being built, i.e. last published generation + 1.
90    stamp_gen: u64,
91    /// Pages superseded by CoW shadows, keyed by the epoch that freed them
92    /// (2n). A page freed during epoch g still serves the g-1 tree until
93    /// epoch g+1 publishes (the dual-slot fallback), so it may be reused
94    /// only once `reuse_limit >= g` -- the store advances the limit after
95    /// each checkpoint, already min'd against the oldest live reader.
96    free: std::collections::BTreeMap<u64, Vec<u32>>,
97    /// Highest freed-generation currently safe to recycle (0 = nothing).
98    reuse_limit: u64,
99    /// Verified birth evidence for retired tree and overflow pages.
100    /// Zero/unknown birth evidence keeps the conservative generation horizon.
101    free_birth: HashMap<u32, u64>,
102    promotion_readers: Vec<u64>,
103    promotion_checked_through: u64,
104    /// Pages recycled during the CURRENT epoch: exempt from the frozen
105    /// check (their number is below the boundary but their content is
106    /// this epoch's). Cleared when the boundary advances.
107    thawed: std::collections::HashSet<u32>,
108    live_pins: u32,
109    /// See `set_frozen_boundary`.
110    frozen_boundary: u32,
111}
112
113pub struct BufferPool {
114    file: Arc<dyn FileIo>,
115    region: AlignedRegion,
116    inner: RefCell<Inner>,
117    /// Which leaf-cell encoding new cells are written in for THIS database.
118    ///
119    /// The encoding is a property of the stored file, declared in its own
120    /// header, not of the build that opens it (Law 8: a release must read AND
121    /// write every database an earlier release wrote, whatever cargo features
122    /// that build had). Every build decodes both families unconditionally;
123    /// this flag only decides what a WRITE produces, and the owner installs it
124    /// from the database's declared features right after opening the pool.
125    /// The `compact-cells` cargo feature survives only as the default for
126    /// databases this build CREATES, which is what seeds the initial value.
127    compact_cells: std::cell::Cell<bool>,
128    _res: Reservation,
129}
130
131impl BufferPool {
132    pub fn new(file: Arc<dyn FileIo>, budget: Arc<MemoryBudget>, frames: usize) -> Result<Self> {
133        let file_len = file.len()?;
134        if file_len % PAGE_SIZE as u64 != 0 {
135            return Err(Error::Corrupt { page_no: 0, why: "data file is not page aligned" });
136        }
137        let pages = file_len / PAGE_SIZE as u64;
138        if pages > u32::MAX as u64 {
139            return Err(Error::Corrupt { page_no: 0, why: "data file has too many pages" });
140        }
141        let bytes = frames.checked_mul(PAGE_SIZE).ok_or(Error::OutOfBudget)?;
142        let res = budget.reserve(Class::Pool, bytes)?;
143        let region = AlignedRegion::new(bytes)?;
144        let next_page = pages as u32;
145        Ok(BufferPool {
146            file,
147            region,
148            inner: RefCell::new(Inner {
149                sweep_steps: 0,
150                frozen_boundary: 0,
151                frames: (0..frames).map(|_| Frame {
152                    page_no: 0, present: false, dirty: false, referenced: false,
153                    pins: 0, writer: false, validated: false,
154                }).collect(),
155                table: HashMap::with_capacity(frames * 2),
156                hand: 0,
157                stats: PoolStats { frames_total: frames, ..Default::default() },
158                next_page,
159                limits: None,
160                free_count: 0,
161                epoch_allocated_pages: 0,
162                stamp_gen: 1,
163                free: std::collections::BTreeMap::new(),
164                reuse_limit: 0,
165                free_birth: HashMap::new(),
166                promotion_readers: Vec::new(),
167                promotion_checked_through: 0,
168                thawed: std::collections::HashSet::new(),
169                live_pins: 0,
170            }),
171            compact_cells: std::cell::Cell::new(cfg!(feature = "compact-cells")),
172            _res: res,
173        })
174    }
175
176    /// The cell encoding new cells on this pool are written in.
177    pub fn compact_cells(&self) -> bool { self.compact_cells.get() }
178    /// Install the encoding the DATABASE declares. Called once per open, by
179    /// the owner that read the header; never derived from build flags after
180    /// creation, and never changed while a tree is being written.
181    pub fn set_compact_cells(&self, on: bool) { self.compact_cells.set(on); }
182
183    pub fn resource_limits(&self) -> Option<crate::limits::ResourceLimits> { self.inner.borrow().limits }
184    pub(crate) fn set_resource_limits(&self, limits: crate::limits::ResourceLimits) -> Result<()> {
185        let limits = limits.validate()?;
186        let mut inner = self.inner.borrow_mut();
187        if inner.next_page as u64 * PAGE_SIZE as u64 > limits.data_bytes {
188            return Err(Error::ResourceLimit("existing data exceeds data allowance"));
189        }
190        inner.limits = Some(limits);
191        Ok(())
192    }
193    pub fn tracked_pages(&self) -> (usize, usize) {
194        let i = self.inner.borrow(); (i.free_count, i.thawed.len())
195    }
196
197    pub fn stats(&self) -> PoolStats { self.inner.borrow().stats }
198
199    /// Pages allocated in this file so far. The scan uses it as a ceiling on how
200    /// many leaves a sibling chain can legitimately visit.
201    pub fn page_count(&self) -> u32 { self.inner.borrow().next_page }
202    /// Allocations (including reuse) since the last root publication. This
203    /// counts touched shadow/overflow pages even after dirty-frame eviction.
204    pub fn epoch_allocated_bytes(&self) -> u64 {
205        self.inner.borrow().epoch_allocated_pages.saturating_mul(PAGE_SIZE as u64)
206    }
207
208    /// Set the generation stamped into pages sealed from now on (2n): the
209    /// store calls this at open and after every checkpoint flip.
210    pub fn set_stamp_gen(&self, gen: u64) { self.inner.borrow_mut().stamp_gen = gen; }
211    /// Birth bound for a writable page retired before its first publication.
212    pub(crate) fn write_generation(&self) -> u64 { self.inner.borrow().stamp_gen }
213
214    /// Record `page_no` as superseded (2n): its replacement was just
215    /// shadowed into a fresh page. Keyed by the CURRENT epoch; see the
216    /// `free` field for when it becomes reusable. Never pages 0/1 (meta).
217    pub fn free_page(&self, page_no: u32) -> Result<()> {
218        if page_no < 2 { return Ok(()); }
219        if self.file.manages_free_pages() {
220            { let mut inner = self.inner.borrow_mut();
221              if let Some(fi) = inner.table.remove(&page_no) {
222                  assert_eq!(inner.frames[fi].pins, 0, "freeing a pinned page");
223                  inner.frames[fi].present=false;inner.frames[fi].dirty=false;
224              }
225            }
226            return self.file.push_free_page(page_no);
227        }
228        let mut inner = self.inner.borrow_mut();
229        if inner.limits.is_some_and(|l| inner.free_count >= l.tracked_pages as usize) {
230            return Err(Error::ResourceLimit("retired-page bookkeeping full"));
231        }
232        let g = inner.stamp_gen;
233        inner.free.entry(g).or_default().push(page_no);
234        inner.free_count += 1;
235        Ok(())
236    }
237
238    /// `birth` came from an independently verified immutable source page.
239    /// The original generation is read before making its shadow, so a cache
240    /// eviction cannot erase this evidence. No extra disk read is required.
241    pub(crate) fn free_shadow_page(&self, page_no: u32, birth: u64) -> Result<()> {
242        if self.file.manages_free_pages() { return self.free_page(page_no); }
243        self.free_page(page_no)?;
244        let mut inner = self.inner.borrow_mut();
245        if birth > 0 && birth <= inner.stamp_gen {
246            inner.free_birth.insert(page_no, birth);
247        }
248        Ok(())
249    }
250
251    /// Promote only retired pages absent from BOTH metadata roots and every
252    /// currently registered snapshot. A page exists in [birth, retirement).
253    /// New readers can only select the current/fallback roots; generation-zero
254    /// reservations and every ambiguous reader disable this optimization.
255    /// Persist promotion in the ordinary freelist (generation 1 = certified
256    /// free), while v2 birth fields preserve still-pending lifetimes across writer reopen.
257    pub(crate) fn refresh_reuse(&self, published: u64, readers: Option<&[u64]>) {
258        let mut inner = self.inner.borrow_mut();
259        let Some(readers) = readers else { inner.reuse_limit = 0; return; };
260        let oldest = readers.first().copied().unwrap_or(u64::MAX);
261        let fallback_limit = published.saturating_sub(1);
262        inner.reuse_limit = fallback_limit.min(oldest);
263        let lower = if inner.promotion_readers == readers {
264            oldest.max(inner.promotion_checked_through)
265        } else {
266            inner.promotion_readers = readers.to_vec();
267            oldest
268        };
269        inner.promotion_checked_through = fallback_limit;
270        if lower >= fallback_limit { return; }
271        let mut promoted = Vec::new();
272        let mut birth = std::mem::take(&mut inner.free_birth);
273        for (&retired, pages) in inner.free.range_mut((std::ops::Bound::Excluded(lower), std::ops::Bound::Included(fallback_limit))) {
274            pages.retain(|p| {
275                let Some(&born) = birth.get(p) else { return true; };
276                let first = readers.partition_point(|g| *g < born);
277                let pinned = readers.get(first).is_some_and(|g| *g < retired);
278                if !pinned { promoted.push(*p); birth.remove(p); }
279                pinned
280            });
281        }
282        inner.free.retain(|_, pages| !pages.is_empty());
283        inner.free_birth = birth;
284        if !promoted.is_empty() { inner.free.entry(1).or_default().extend(promoted); }
285    }
286
287    /// Advance the recycling horizon (2n): pages freed at generations
288    /// <= `limit` may be handed out again. The store computes the limit
289    /// (published - 1, min'd with the oldest live snapshot reader).
290    pub fn set_reuse_limit(&self, limit: u64) { self.inner.borrow_mut().reuse_limit = limit; }
291
292    /// Serialize the freelist for the checkpoint's sidecar file (2n step D):
293    /// [magic 8][version u16][reserved 6][published generation u64], then
294    /// [freed-at-gen u64][n u32][(page u32, birth u64) * n]... and a crc32c trailer.
295    /// The publication generation is what makes an old, otherwise valid
296    /// sidecar unusable beside a newer data file. Loss, mismatch or corruption
297    /// can therefore only LEAK pages, never recycle live ones. Version 2
298    /// adds eight bytes per retired page to avoid losing lifetime evidence on
299    /// reopen. Version 1 is rejected as derived state, never guessed/migrated.
300    pub fn export_free(&self, published_generation: u64) -> Vec<u8> {
301        let inner = self.inner.borrow();
302        let mut v = Vec::with_capacity(FREE_HEADER_LEN + 4 + 12 * inner.free.len() + 12 * inner.free_count);
303        v.extend_from_slice(&Self::empty_free(published_generation)[..FREE_HEADER_LEN]);
304        for (g, pages) in &inner.free {
305            v.extend_from_slice(&g.to_le_bytes());
306            let n = u32::try_from(pages.len()).expect("one file cannot contain more than u32 pages");
307            v.extend_from_slice(&n.to_le_bytes());
308            for p in pages {
309                v.extend_from_slice(&p.to_le_bytes());
310                v.extend_from_slice(&inner.free_birth.get(p).copied().unwrap_or(0).to_le_bytes());
311            }
312        }
313        let c = crc32c::crc32c(&v);
314        v.extend_from_slice(&c.to_le_bytes());
315        v
316    }
317
318    /// The canonical empty sidecar used by recovery after it renumbers every
319    /// page. Kept here so ordinary checkpoints and recovery cannot drift into
320    /// two encoders for the same durability handshake.
321    pub(crate) fn empty_free(published_generation: u64) -> Vec<u8> {
322        let mut v = Vec::with_capacity(FREE_HEADER_LEN + 4);
323        v.extend_from_slice(&FREE_MAGIC);
324        v.extend_from_slice(&FREE_VERSION.to_le_bytes());
325        v.extend_from_slice(&[0; 6]);
326        v.extend_from_slice(&published_generation.to_le_bytes());
327        let c = crc32c::crc32c(&v);
328        v.extend_from_slice(&c.to_le_bytes());
329        v
330    }
331
332    /// Decode without mutating the pool. Checkpoint/recovery use this on a
333    /// freshly reopened candidate before renaming it over the standing
334    /// sidecar: build, independently verify, then publish.
335    fn walk_free<F>(
336        bytes: &[u8],
337        expected_generation: u64,
338        page_count: u32,
339        mut visit: F,
340    ) -> Option<()>
341    where
342        F: FnMut(u64, u32, u64),
343    {
344        if bytes.len() < FREE_HEADER_LEN + 4 { return None; }
345        let (body, tail) = bytes.split_at(bytes.len() - 4);
346        if crc32c::crc32c(body) != u32::from_le_bytes(tail.try_into().ok()?)
347            || body.get(..8)? != FREE_MAGIC
348            || u16::from_le_bytes(body.get(8..10)?.try_into().ok()?) != FREE_VERSION
349            || body.get(10..16)? != [0; 6]
350            || u64::from_le_bytes(body.get(16..24)?.try_into().ok()?) != expected_generation
351        {
352            return None;
353        }
354        let mut pos = FREE_HEADER_LEN;
355        let mut previous_generation = None;
356        let mut seen = std::collections::HashSet::new();
357        while pos < body.len() {
358            let header_end = pos.checked_add(12)?;
359            if header_end > body.len() { return None; }
360            let g = u64::from_le_bytes(body[pos..pos + 8].try_into().unwrap());
361            let n = u32::from_le_bytes(body[pos + 8..pos + 12].try_into().unwrap()) as usize;
362            if g == 0
363                || g > expected_generation
364                || n == 0
365                || previous_generation.is_some_and(|previous| g <= previous)
366            {
367                return None;
368            }
369            previous_generation = Some(g);
370            pos = header_end;
371            let pages_bytes = n.checked_mul(12)?;
372            let pages_end = pos.checked_add(pages_bytes)?;
373            if pages_end > body.len() { return None; }
374            for i in 0..n {
375                let at = pos + i * 12;
376                let p = u32::from_le_bytes(body[at..at + 4].try_into().unwrap());
377                let birth = u64::from_le_bytes(body[at + 4..at + 12].try_into().unwrap());
378                if p < 2 || p >= page_count || birth > g || !seen.insert(p) { return None; }
379                visit(g, p, birth);
380            }
381            pos = pages_end;
382        }
383        Some(())
384    }
385
386    pub(crate) fn verify_free(
387        bytes: &[u8],
388        expected_generation: u64,
389        page_count: u32,
390    ) -> bool {
391        Self::walk_free(bytes, expected_generation, page_count, |_, _, _| {}).is_some()
392    }
393
394    /// Load a persisted freelist (reopen). Anything malformed or stamped for
395    /// another data generation becomes an empty list (the leak-only posture).
396    /// Every page number is rejected unless it is inside this exact file.
397    pub fn import_free(&self, bytes: &[u8], expected_generation: u64) -> bool {
398        let page_count = self.inner.borrow().next_page;
399        if self.resource_limits().is_some_and(|l| bytes.len() as u64 > l.freelist_bytes()) { return false; }
400        let mut count = 0usize;
401        let mut free = std::collections::BTreeMap::new();
402        let mut births = HashMap::new();
403        let valid = Self::walk_free(bytes, expected_generation, page_count, |g, p, birth| {
404            count += 1;
405            free.entry(g).or_insert_with(Vec::new).push(p);
406            if birth != 0 { births.insert(p, birth); }
407        });
408        if valid.is_none() || self.resource_limits().is_some_and(|l| count > l.tracked_pages as usize) {
409            return false;
410        }
411        let mut inner = self.inner.borrow_mut();
412        inner.free = free;
413        inner.free_count = count;
414        inner.free_birth = births;
415        inner.promotion_checked_through = 0;
416        inner.promotion_readers.clear();
417        true
418    }
419
420    pub(crate) fn sync_dir(&self) -> Result<()> { self.file.sync_dir() }
421    pub(crate) fn file_ref(&self) -> &dyn FileIo { &*self.file }
422
423    /// (eligible-now, waiting-on-horizon) freelist depths (tests/probes).
424    pub fn free_pages_split(&self) -> (usize, usize) {
425        let inner = self.inner.borrow();
426        let lim = inner.reuse_limit;
427        let el: usize = inner.free.range(..=lim).map(|(_, v)| v.len()).sum();
428        let tot: usize = inner.free.values().map(|v| v.len()).sum();
429        (el, tot - el)
430    }
431    /// Sum of pages currently waiting on the freelist (tests/probes).
432    pub fn free_pages_pending(&self) -> usize {
433        self.inner.borrow().free.values().map(|v| v.len()).sum()
434    }
435
436    /// Pop one recyclable page number, if any (2n).
437    fn pop_free(inner: &mut Inner) -> Option<u32> {
438        let limit = inner.reuse_limit;
439        let g = *inner.free.range(..=limit).next()?.0;
440        let v = inner.free.get_mut(&g)?;
441        let p = v.pop()?;
442        inner.free_birth.remove(&p);
443        inner.free_count -= 1;
444        if v.is_empty() { inner.free.remove(&g); }
445        Some(p)
446    }
447
448    /// 2f: page numbers below this boundary belong to the last PUBLISHED
449    /// root (the checkpoint a snapshot reader may be standing on) and are
450    /// immutable -- writers shadow them to fresh numbers instead of editing
451    /// in place. `get_mut` asserts it. Meta slots (pages 0 and 1) are the
452    /// publication mechanism itself and are exempt. Boundary 0 = no epoch
453    /// published yet, everything mutable (fresh store before first publish).
454    pub fn set_frozen_boundary(&self) {
455        let mut inner = self.inner.borrow_mut();
456        inner.frozen_boundary = inner.next_page;
457        // recycled pages just published with this epoch: frozen again
458        inner.thawed.clear();
459        inner.epoch_allocated_pages = 0;
460    }
461    pub fn frozen_boundary(&self) -> u32 { self.inner.borrow().frozen_boundary }
462    /// Experimental stable-page pager owns snapshot versions outside the pool.
463    pub fn finish_stable_page_epoch(&self) {
464        let mut inner = self.inner.borrow_mut();
465        assert_eq!(inner.frozen_boundary, 0);
466        inner.thawed.clear();
467        inner.epoch_allocated_pages = 0;
468    }
469    pub fn is_frozen(&self, page_no: u32) -> bool {
470        let inner = self.inner.borrow();
471        page_no >= 2 && page_no < inner.frozen_boundary && !inner.thawed.contains(&page_no)
472    }
473    pub fn io_stats(&self) -> Option<&crate::io::IoStats> { self.file.stats() }
474
475    /// Re-baseline the high-water mark to the pins currently held, so a caller
476    /// can measure the peak over one specific operation.
477    pub fn reset_peak_pins(&self) {
478        let mut inner = self.inner.borrow_mut();
479        inner.stats.peak_pins = inner.live_pins;
480    }
481
482    /// Diagnostic: total clock-sweep steps across all victim() calls. A CLOCK
483    /// pathology shows here as steps >> misses -- every hit re-arms a ref bit,
484    /// so a miss arriving after many hits must strip a large stretch of the
485    /// table before it finds a victim.
486    pub fn sweep_steps(&self) -> u64 { self.inner.borrow().sweep_steps }
487
488    fn victim(&self, inner: &mut Inner) -> Result<usize> {
489        let n = inner.frames.len();
490        for _ in 0..(n * 4) {
491            inner.sweep_steps += 1;
492            let i = inner.hand;
493            inner.hand = (inner.hand + 1) % n;
494            if inner.frames[i].pins > 0 { continue; }
495            if !inner.frames[i].present { return Ok(i); }
496            if inner.frames[i].referenced { inner.frames[i].referenced = false; continue; }
497            if inner.frames[i].dirty {
498                let no = inner.frames[i].page_no;
499                // SAFETY: frame i has pins == 0 here. Every guard that could
500                // reference frame i (PinnedRead or PinnedWrite) decrements
501                // pins on drop before releasing its borrow, so pins == 0
502                // means no `&[u8]`/`&mut [u8]` into this frame is live.
503                crate::page::seal(unsafe { self.region.page_mut(i) }, inner.stamp_gen);
504                self.file.write_at(unsafe { self.region.page(i) }, no as u64 * PAGE_SIZE as u64)?;
505                crate::write_stats::add(
506                    crate::write_stats::Phase::FinalPages,
507                    PAGE_SIZE as u64,
508                );
509                inner.frames[i].dirty = false;
510            }
511            let old = inner.frames[i].page_no;
512            inner.table.remove(&old);
513            inner.frames[i].present = false;
514            inner.stats.evictions += 1;
515            return Ok(i);
516        }
517        Err(Error::OutOfBudget) // every frame is pinned
518    }
519
520    /// Resolve `page_no` to a resident frame, pinning it. `write` selects
521    /// which exclusivity rule applies: a write pin requires the frame to
522    /// currently have no other live pin at all; a read pin merely requires
523    /// no live writer.
524    fn load(&self, page_no: u32, write: bool) -> Result<usize> {
525        let mut inner = self.inner.borrow_mut();
526        if let Some(&i) = inner.table.get(&page_no) {
527            let f = &inner.frames[i];
528            assert!(
529                !(write && f.pins > 0),
530                "page {page_no} is already pinned; get_mut requires exclusive access"
531            );
532            assert!(
533                write || !f.writer,
534                "page {page_no} is already pinned by a writer"
535            );
536            inner.stats.hits += 1;
537            inner.frames[i].referenced = true;
538            inner.frames[i].pins += 1;
539            inner.live_pins += 1;
540            if inner.live_pins > inner.stats.peak_pins { inner.stats.peak_pins = inner.live_pins; }
541            if write { inner.frames[i].writer = true; }
542            return Ok(i);
543        }
544        inner.stats.misses += 1;
545        let i = self.victim(&mut inner)?;
546        // SAFETY: `victim` returned frame i with pins == 0 and removed it
547        // from the table, so no reference to it is live.
548        self.file.read_at(unsafe { self.region.page_mut(i) }, page_no as u64 * PAGE_SIZE as u64)?;
549        // The medium boundary: the one place a checksum needs verifying.
550        // Failure leaves the pool untouched (frame not yet in the table).
551        crate::page::PageRef::open(unsafe { self.region.page(i) }, page_no)?;
552        inner.frames[i] = Frame {
553            page_no, present: true, dirty: false, referenced: true, pins: 1, writer: write, validated: false,
554        };
555        inner.table.insert(page_no, i);
556        inner.live_pins += 1;
557        if inner.live_pins > inner.stats.peak_pins { inner.stats.peak_pins = inner.live_pins; }
558        Ok(i)
559    }
560
561    /// Read `n` consecutive pages in ONE pread into `buf`, WITHOUT caching
562    /// them (2e ablation A1, for overflow chains). Returns Ok(false) if any
563    /// page in the run is resident -- a resident frame may be dirtier than
564    /// disk, so the caller must take the per-page pool path instead. Each
565    /// page's checksum is verified at this medium boundary exactly as `load`
566    /// verifies it; a failure refuses the whole run and caches nothing.
567    ///
568    /// SACRIFICE (Law 4): pages read this way do not warm the cache -- an
569    /// immediate re-read pays the pread again. Bought: one syscall per chain
570    /// instead of one per page, and single-use payload pages never evict a
571    /// hot tree page.
572    pub fn read_run_uncached(&self, start: u32, n: u32, buf: &mut Vec<u8>) -> Result<bool> {
573        {
574            let inner = self.inner.borrow();
575            if n == 0 || start.checked_add(n).map_or(true, |e| e > inner.next_page) {
576                return Err(Error::Corrupt { page_no: start, why: "overflow run out of bounds" });
577            }
578            for p in start..start + n {
579                if inner.table.contains_key(&p) { return Ok(false); }
580            }
581        }
582        buf.clear();
583        buf.resize(n as usize * PAGE_SIZE, 0);
584        self.file.read_at(buf, start as u64 * PAGE_SIZE as u64)?;
585        for i in 0..n as usize {
586            crate::page::PageRef::open(&buf[i * PAGE_SIZE..(i + 1) * PAGE_SIZE], start + i as u32)?;
587        }
588        Ok(true)
589    }
590
591    pub fn get(&self, page_no: u32) -> Result<PinnedRead<'_>> {
592        let i = self.load(page_no, false)?;
593        Ok(PinnedRead { pool: self, frame: i })
594    }
595
596    pub fn get_mut(&self, page_no: u32) -> Result<PinnedWrite<'_>> {
597        // The 2f no-overwrite invariant, enforced at the single chokepoint
598        // every in-place write passes through: a frozen page is part of a
599        // published root a snapshot reader may hold; editing it would tear
600        // that reader's view. Loud, not a Result -- reaching here means a
601        // write path forgot to shadow, which is our bug, never the caller's.
602        assert!(
603            !self.is_frozen(page_no),
604            "page {page_no} is frozen (published at the last checkpoint); write paths must shadow it"
605        );
606        let i = self.load(page_no, true)?;
607        self.inner.borrow_mut().frames[i].dirty = true;
608        Ok(PinnedWrite { pool: self, frame: i, page_no })
609    }
610
611    /// Reserve a brand-new page at the end of the file.
612    ///
613    /// `next_page` only ever increases and a page number handed out here is
614    /// never reclaimed or reused -- `recover.rs`'s duplicate-key resolution
615    /// depends on that: among two surviving leaves claiming the same key, it
616    /// trusts the higher page number as the one written later. A change that
617    /// reclaims/reuses page numbers would silently break that reasoning and
618    /// must revisit `recover.rs`'s `tag`/`untag`/dedup logic.
619    fn admit_allocation(i: &Inner) -> Result<()> {
620        if let Some(l) = i.limits {
621            let reuse = i.free.range(..=i.reuse_limit).next().is_some();
622            if reuse && i.thawed.len() >= l.tracked_pages as usize {
623                return Err(Error::ResourceLimit("recycled-page bookkeeping full"));
624            }
625            if !reuse && i.next_page as u64 >= l.data_bytes / PAGE_SIZE as u64 {
626                return Err(Error::ResourceLimit("data extent full; snapshots or fallback roots may pin pages"));
627            }
628        }
629        Ok(())
630    }
631
632    pub fn allocate(&self) -> Result<PinnedWrite<'_>> {
633        let managed = self.file.manages_free_pages();
634        let external_free = if managed { self.file.pop_free_page()? } else { None };
635        let page_no = { let mut inner = self.inner.borrow_mut();
636            Self::admit_allocation(&inner)?;
637            inner.epoch_allocated_pages += 1;
638            match if managed { external_free } else { Self::pop_free(&mut inner) } {
639                Some(p) => {
640                    // recycled: its number sits below the frozen boundary but
641                    // its content belongs to THIS epoch -- thaw it, and drop
642                    // any stale cached frame for the old content.
643                    inner.thawed.insert(p);
644                    if let Some(&fi) = inner.table.get(&p) {
645                        debug_assert_eq!(inner.frames[fi].pins, 0,
646                            "recycling page {p} while a guard holds its stale frame");
647                        inner.table.remove(&p);
648                        inner.frames[fi].present = false;
649                        inner.frames[fi].dirty = false;
650                    }
651                    p
652                }
653                None => { let p = inner.next_page; inner.next_page = p.checked_add(1).ok_or(Error::TooLarge)?; p }
654            } };
655        let mut inner = self.inner.borrow_mut();
656        let i = self.victim(&mut inner)?;
657        // SAFETY: freshly evicted frame, pins == 0, not yet in the table, so
658        // no reference to it is live.
659        //
660        // Left as a VALID empty Free page, not zeroes: a caller that only wants
661        // the page number (bulk.rs reserves one by allocate-and-drop, three
662        // times) leaves the frame dirty, and an eviction before anything writes
663        // it would publish ZEROES to disk -- bad magic, undetectable until read.
664        // No path through the pool may publish a page a reader must refuse.
665        {
666            let b = unsafe { self.region.page_mut(i) };
667            b.fill(0);
668            crate::page::PageMut::init(b, crate::page::PageKind::Free, 0, page_no).finalise(0);
669        }
670        inner.frames[i] = Frame {
671            page_no, present: true, dirty: true, referenced: true, pins: 1, writer: true, validated: false,
672        };
673        inner.table.insert(page_no, i);
674        inner.live_pins += 1;
675        if inner.live_pins > inner.stats.peak_pins { inner.stats.peak_pins = inner.live_pins; }
676        drop(inner);
677        Ok(PinnedWrite { pool: self, frame: i, page_no })
678    }
679
680    /// Reserve a page number for a bulk packer that will write the complete,
681    /// checksummed page directly. Packed pages are unreachable until graft
682    /// publication, so caching every one only forces the fixed pool to evict
683    /// and issue a 4 KiB write per page. The direct writer batches them while
684    /// preserving the same allocator/freelist invariants.
685    pub(crate) fn allocate_unpooled(&self) -> Result<u32> {
686        let mut inner = self.inner.borrow_mut();
687        Self::admit_allocation(&inner)?;
688        inner.epoch_allocated_pages += 1;
689        Ok(match Self::pop_free(&mut inner) {
690            Some(page_no) => {
691                inner.thawed.insert(page_no);
692                if let Some(&frame) = inner.table.get(&page_no) {
693                    assert_eq!(inner.frames[frame].pins, 0,
694                        "recycling page {page_no} while a guard holds its stale frame");
695                    inner.table.remove(&page_no);
696                    inner.frames[frame].present = false;
697                    inner.frames[frame].dirty = false;
698                }
699                page_no
700            }
701            None => {
702                let page_no = inner.next_page;
703                inner.next_page = page_no.checked_add(1).ok_or(Error::TooLarge)?;
704                page_no
705            }
706        })
707    }
708
709    pub(crate) fn stamp_generation(&self) -> u64 { self.inner.borrow().stamp_gen }
710
711    pub(crate) fn write_unpooled_run(&self, first_page: u32, bytes: &[u8]) -> Result<()> {
712        if bytes.is_empty() || bytes.len() % PAGE_SIZE != 0 {
713            return Err(Error::Corrupt { page_no: first_page, why: "bulk page run is not aligned" });
714        }
715        let pages = u32::try_from(bytes.len() / PAGE_SIZE).map_err(|_| Error::TooLarge)?;
716        if first_page.checked_add(pages).is_none_or(|end| end > self.page_count()) {
717            return Err(Error::Corrupt { page_no: first_page, why: "bulk page run exceeds allocation" });
718        }
719        self.file.write_at(bytes, first_page as u64 * PAGE_SIZE as u64)?;
720        crate::write_stats::add(
721            crate::write_stats::Phase::CandidatePages,
722            bytes.len() as u64,
723        );
724        Ok(())
725    }
726
727    pub fn flush_all(&self, barrier: Barrier) -> Result<()> {
728        let mut inner = self.inner.borrow_mut();
729        for i in 0..inner.frames.len() {
730            if inner.frames[i].present && inner.frames[i].dirty {
731                // A dirty frame that is still pinned means a guard is live
732                // across a checkpoint. Skipping it silently would let
733                // flush_all return Ok(()) having left dirty data unwritten —
734                // a checkpoint that quietly declines to write is exactly how
735                // a row goes missing later, and Law 3 exists because a
736                // process that reports success while dropping state is the
737                // shape that loses data. `get`/`get_mut` already treat a
738                // conflicting access as a contract violation; this is the
739                // same violation and gets the same answer: loud, not a
740                // Result, because the single writer never holds a guard
741                // across a checkpoint, so this firing means a bug in our own
742                // code, not a condition a caller could act on.
743                assert_eq!(
744                    inner.frames[i].pins, 0,
745                    "flush_all with frame {i} (page {}) pinned and dirty",
746                    inner.frames[i].page_no
747                );
748                let no = inner.frames[i].page_no;
749                // SAFETY: pins == 0 was just asserted, so no guard to this
750                // frame is live.
751                crate::page::seal(unsafe { self.region.page_mut(i) }, inner.stamp_gen);
752                self.file.write_at(unsafe { self.region.page(i) }, no as u64 * PAGE_SIZE as u64)?;
753                crate::write_stats::add(
754                    crate::write_stats::Phase::FinalPages,
755                    PAGE_SIZE as u64,
756                );
757                inner.stats.dirty_pages_flushed += 1;
758                inner.frames[i].dirty = false;
759            }
760        }
761        // Count, then drop the borrow, then issue the syscall: a RefCell
762        // guard held across a blocking call (F_FULLFSYNC is ~65x an ordinary
763        // fsync on macOS) is a hazard worth not creating even though nothing
764        // else can reach `inner` mid-syscall in this single-writer engine.
765        match barrier {
766            Barrier::Full => { inner.stats.sync_full_calls += 1; drop(inner); self.file.sync_full() }
767            Barrier::Data => { inner.stats.sync_data_calls += 1; drop(inner); self.file.sync_data() }
768            Barrier::None => { drop(inner); Ok(()) }
769        }
770    }
771
772    fn unpin(&self, frame: usize, writer: bool) {
773        let mut inner = self.inner.borrow_mut();
774        inner.frames[frame].pins -= 1;
775        inner.live_pins -= 1;
776        if writer { inner.frames[frame].writer = false; }
777    }
778}
779
780pub struct PinnedRead<'a> { pool: &'a BufferPool, frame: usize }
781impl std::ops::Deref for PinnedRead<'_> {
782    type Target = [u8];
783    // SAFETY: this guard holds a read pin on the frame. `load` refuses to
784    // hand out a write pin (PinnedWrite) for the same frame while any pin is
785    // live, so no `&mut [u8]` can coexist with this borrow; `victim` refuses
786    // to touch a pinned frame, so it cannot be evicted or overwritten either.
787    fn deref(&self) -> &[u8] { unsafe { self.pool.region.page(self.frame) } }
788}
789impl PinnedRead<'_> {
790    /// 2f-A1: has this frame's slot directory been fully validated during
791    /// its current residency? See `Frame::validated`.
792    pub fn validated(&self) -> bool {
793        self.pool.inner.borrow().frames[self.frame].validated
794    }
795    pub fn set_validated(&self) {
796        self.pool.inner.borrow_mut().frames[self.frame].validated = true;
797    }
798}
799impl Drop for PinnedRead<'_> { fn drop(&mut self) { self.pool.unpin(self.frame, false) } }
800
801pub struct PinnedWrite<'a> { pool: &'a BufferPool, frame: usize, page_no: u32 }
802impl PinnedWrite<'_> {
803    pub fn page_no(&self) -> u32 { self.page_no }
804    // SAFETY: `&mut self` plus a write pin means this guard is the only live
805    // reference to the frame: `load` refuses to create a write pin unless
806    // pins was 0, and refuses any further pin (read or write) while this
807    // frame's writer flag is set, so nothing else can alias it.
808    pub fn bytes_mut(&mut self) -> &mut [u8] { unsafe { self.pool.region.page_mut(self.frame) } }
809    // SAFETY: `&self` plus the write pin; no other `&mut` can coexist with
810    // this borrow for the same reason as `bytes_mut`, and Rust's own borrow
811    // checker prevents this call from overlapping a `bytes_mut` call on the
812    // same guard.
813    pub fn bytes(&self) -> &[u8] { unsafe { self.pool.region.page(self.frame) } }
814}
815impl Drop for PinnedWrite<'_> { fn drop(&mut self) { self.pool.unpin(self.frame, true) } }
816
817#[cfg(test)]
818mod tests {
819    use super::*;
820    use crate::budget::MemoryBudget;
821    use crate::io::{open_file, IoMode};
822    use crate::page::{PageKind, PageMut, PageRef, PAGE_SIZE};
823    use std::sync::Arc;
824
825    fn pool_with(frames: usize) -> (BufferPool, tempfile::TempDir) {
826        let dir = tempfile::tempdir().unwrap();
827        let (f, _) = open_file(&dir.path().join("t.db"), IoMode::Buffered).unwrap();
828        let budget = Arc::new(MemoryBudget::new(64 * 1024 * 1024));
829        (BufferPool::new(f.into(), budget, frames).unwrap(), dir)
830    }
831
832    /// The whole claim in one test: a pool far smaller than the working set
833    /// still answers every page correctly.
834    #[test]
835    fn a_pool_smaller_than_the_working_set_still_serves_every_page() {
836        let (pool, _d) = pool_with(8);
837        let n = 400u32;
838        for _ in 0..n {
839            let mut w = pool.allocate().unwrap();
840            let no = w.page_no();
841            let mut p = PageMut::init(w.bytes_mut(), PageKind::Leaf, 1, no);
842            p.insert_slot(0, &no.to_le_bytes()).unwrap();
843            p.finalise(0);
844            drop(w);
845        }
846        pool.flush_all(Barrier::Data).unwrap();
847
848        for i in 0..n {
849            let r = pool.get(i).unwrap();
850            let pr = PageRef::open(&r, i).unwrap();
851            assert_eq!(pr.slot(0), &i.to_le_bytes());
852        }
853        assert!(pool.stats().evictions > 0, "8 frames over 400 pages must evict");
854        assert_eq!(pool.stats().frames_total, 8, "the pool must never grow");
855    }
856
857    #[test]
858    fn a_pinned_page_is_never_evicted() {
859        let (pool, _d) = pool_with(4);
860        for _ in 0..4 { let mut w = pool.allocate().unwrap(); let no = w.page_no();
861            PageMut::init(w.bytes_mut(), PageKind::Leaf, 1, no).finalise(0); }
862        pool.flush_all(Barrier::Data).unwrap();
863
864        let held = pool.get(0).unwrap();          // pin page 0 for the whole test
865        for i in 1..4 { let _ = pool.get(i).unwrap(); }
866        // Touching more pages than there are frames must not have stolen page 0.
867        let pr = PageRef::open(&held, 0).unwrap();
868        assert_eq!(pr.page_no(), 0);
869    }
870
871    #[test]
872    fn a_dirty_page_survives_eviction_because_it_was_written_out() {
873        let (pool, _d) = pool_with(2);
874        let mut w = pool.allocate().unwrap();
875        let no = w.page_no();
876        let mut p = PageMut::init(w.bytes_mut(), PageKind::Leaf, 9, no);
877        p.insert_slot(0, b"survives").unwrap();
878        p.finalise(7);
879        drop(w);
880        // Force eviction by touching more pages than frames.
881        for _ in 0..4 { let mut x = pool.allocate().unwrap(); let n2 = x.page_no();
882            PageMut::init(x.bytes_mut(), PageKind::Leaf, 9, n2).finalise(0); }
883        let r = pool.get(no).unwrap();
884        assert_eq!(PageRef::open(&r, no).unwrap().slot(0), b"survives");
885    }
886
887    /// The exclusivity rules are the whole reason `AlignedRegion::page_mut` is
888    /// an `unsafe fn`. They were added because a pin COUNT cannot distinguish a
889    /// reader from a writer, so `get` followed by `get_mut` on the same page
890    /// handed out a `&mut [u8]` aliasing a live `&[u8]` — reachable with two
891    /// ordinary safe calls. These three tests are what stop that returning.
892    #[test]
893    #[should_panic(expected = "already pinned")]
894    fn taking_a_writer_on_a_page_a_reader_holds_is_refused() {
895        let (pool, _d) = pool_with(4);
896        { let mut w = pool.allocate().unwrap(); let no = w.page_no();
897          PageMut::init(w.bytes_mut(), PageKind::Leaf, 1, no).finalise(0); }
898        pool.flush_all(Barrier::Data).unwrap();
899        let _reader = pool.get(0).unwrap();
900        let _writer = pool.get_mut(0).unwrap();   // must panic, not alias
901    }
902
903    #[test]
904    #[should_panic(expected = "writer")]
905    fn taking_a_reader_on_a_page_a_writer_holds_is_refused() {
906        let (pool, _d) = pool_with(4);
907        { let mut w = pool.allocate().unwrap(); let no = w.page_no();
908          PageMut::init(w.bytes_mut(), PageKind::Leaf, 1, no).finalise(0); }
909        pool.flush_all(Barrier::Data).unwrap();
910        let _writer = pool.get_mut(0).unwrap();
911        let _reader = pool.get(0).unwrap();       // must panic, not alias
912    }
913
914    /// A checkpoint that returns Ok() having left a dirty page unwritten is a
915    /// silent data-loss shape. It must be loud.
916    #[test]
917    #[should_panic(expected = "pinned and dirty")]
918    fn flushing_while_a_dirty_page_is_pinned_is_refused() {
919        let (pool, _d) = pool_with(4);
920        let mut w = pool.allocate().unwrap();
921        let no = w.page_no();
922        PageMut::init(w.bytes_mut(), PageKind::Leaf, 1, no).finalise(0);
923        pool.flush_all(Barrier::Data).unwrap();                // w is still alive: must panic
924    }
925
926    #[test]
927    fn the_budget_refuses_a_pool_it_cannot_fund() {
928        let dir = tempfile::tempdir().unwrap();
929        let (f, _) = open_file(&dir.path().join("t.db"), IoMode::Buffered).unwrap();
930        let budget = Arc::new(MemoryBudget::new(PAGE_SIZE * 4));
931        assert!(BufferPool::new(f.into(), budget, 1000).is_err());
932    }
933
934    #[test]
935    fn a_file_with_a_partial_page_is_refused() {
936        let dir = tempfile::tempdir().unwrap();
937        let path = dir.path().join("partial.db");
938        std::fs::write(&path, b"one bad trailing byte").unwrap();
939        let (f, _) = open_file(&path, IoMode::Buffered).unwrap();
940        let budget = Arc::new(MemoryBudget::new(4 * PAGE_SIZE));
941
942        assert!(matches!(
943            BufferPool::new(f.into(), budget, 4),
944            Err(Error::Corrupt { page_no: 0, .. })
945        ));
946    }
947}