1use 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 pub frames_total: usize,
49 pub peak_pins: u32,
50 pub sync_data_calls: u64,
57 pub sync_full_calls: u64,
58 pub dirty_pages_flushed: u64,
60}
61
62struct Frame {
63 page_no: u32,
64 present: bool,
65 dirty: bool,
66 referenced: bool,
67 pins: u32,
69 writer: bool,
72 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 stamp_gen: u64,
91 free: std::collections::BTreeMap<u64, Vec<u32>>,
97 reuse_limit: u64,
99 free_birth: HashMap<u32, u64>,
102 promotion_readers: Vec<u64>,
103 promotion_checked_through: u64,
104 thawed: std::collections::HashSet<u32>,
108 live_pins: u32,
109 frozen_boundary: u32,
111}
112
113pub struct BufferPool {
114 file: Arc<dyn FileIo>,
115 region: AlignedRegion,
116 inner: RefCell<Inner>,
117 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 pub fn compact_cells(&self) -> bool { self.compact_cells.get() }
178 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 pub fn page_count(&self) -> u32 { self.inner.borrow().next_page }
202 pub fn epoch_allocated_bytes(&self) -> u64 {
205 self.inner.borrow().epoch_allocated_pages.saturating_mul(PAGE_SIZE as u64)
206 }
207
208 pub fn set_stamp_gen(&self, gen: u64) { self.inner.borrow_mut().stamp_gen = gen; }
211 pub(crate) fn write_generation(&self) -> u64 { self.inner.borrow().stamp_gen }
213
214 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 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 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 pub fn set_reuse_limit(&self, limit: u64) { self.inner.borrow_mut().reuse_limit = limit; }
291
292 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 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 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 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 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 pub fn free_pages_pending(&self) -> usize {
433 self.inner.borrow().free.values().map(|v| v.len()).sum()
434 }
435
436 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 pub fn set_frozen_boundary(&self) {
455 let mut inner = self.inner.borrow_mut();
456 inner.frozen_boundary = inner.next_page;
457 inner.thawed.clear();
459 inner.epoch_allocated_pages = 0;
460 }
461 pub fn frozen_boundary(&self) -> u32 { self.inner.borrow().frozen_boundary }
462 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 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 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 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) }
519
520 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 self.file.read_at(unsafe { self.region.page_mut(i) }, page_no as u64 * PAGE_SIZE as u64)?;
549 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 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 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 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 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 {
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 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 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 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 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 fn deref(&self) -> &[u8] { unsafe { self.pool.region.page(self.frame) } }
788}
789impl PinnedRead<'_> {
790 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 pub fn bytes_mut(&mut self) -> &mut [u8] { unsafe { self.pool.region.page_mut(self.frame) } }
809 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 #[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(); for i in 1..4 { let _ = pool.get(i).unwrap(); }
866 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 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 #[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(); }
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(); }
913
914 #[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(); }
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}