Skip to main content

reifydb_store_multi/tier/read/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4//! Read buffer tier of the multi-version store, caching keys the commit buffer has evicted so a repeated
5//! point read need not fall through to persistent every time. The previous slot is only ever filled by an
6//! in-place supersede, so it stays version-adjacent to the current slot. Range scans consult this tier only
7//! for `range_complete` buckets, and the always-scanned commit buffer still wins on version, so the cache
8//! can never mask a newer value nor resurrect a deleted one.
9
10mod point;
11mod pool;
12mod range;
13#[cfg(test)]
14mod tests;
15
16use std::{
17	collections::{BTreeMap, HashMap},
18	mem::size_of,
19	sync::Arc,
20};
21
22use reifydb_codec::key::encoded::EncodedKey;
23use reifydb_core::{common::CommitVersion, util::budget::MemoryBudget};
24use reifydb_runtime::sync::mutex::Mutex;
25use reifydb_store::row::page::{DEFAULT_BUCKET_SHIFT, PageId};
26use reifydb_value::{byte_size::ByteSize, util::cowvec::CowVec};
27
28use crate::tier::RangeBatch;
29
30#[derive(Clone, Copy, Debug)]
31pub struct ReadBufferConfig {
32	pub resident_pages: usize,
33	pub resident_bytes: Option<ByteSize>,
34	pub shards: usize,
35	pub bucket_shift: u8,
36}
37
38impl Default for ReadBufferConfig {
39	fn default() -> Self {
40		Self {
41			resident_pages: 1024,
42			resident_bytes: Some(ByteSize::from_gib(2)),
43			shards: 16,
44			bucket_shift: DEFAULT_BUCKET_SHIFT,
45		}
46	}
47}
48
49#[derive(Clone)]
50struct PageEntry {
51	version: CommitVersion,
52	value: Option<CowVec<u8>>,
53	previous: Option<(CommitVersion, Option<CowVec<u8>>)>,
54}
55
56struct ResidentPage {
57	entries: BTreeMap<EncodedKey, PageEntry>,
58	bytes: usize,
59	payload: usize,
60	hot: bool,
61	tick: u64,
62	range_complete: bool,
63	warm_blocked: bool,
64}
65
66const NODE_FILL_DIVISOR: usize = 2;
67
68const ENTRY_OVERHEAD: usize = NODE_FILL_DIVISOR * (size_of::<EncodedKey>() + size_of::<PageEntry>());
69
70fn value_len(value: &Option<CowVec<u8>>) -> usize {
71	value.as_ref().map_or(0, |bytes| bytes.len())
72}
73
74#[derive(Clone, Copy, Default)]
75struct EntryFootprint {
76	resident: usize,
77	payload: usize,
78}
79
80fn entry_footprint(key: &EncodedKey, entry: &PageEntry) -> EntryFootprint {
81	let version_payload = key.len() + size_of::<CommitVersion>();
82	EntryFootprint {
83		resident: ENTRY_OVERHEAD
84			+ key.heap_bytes() + value_len(&entry.value)
85			+ entry.previous.as_ref().map_or(0, |(_, value)| value_len(value)),
86		payload: version_payload
87			+ value_len(&entry.value)
88			+ entry.previous.as_ref().map_or(0, |(_, value)| version_payload + value_len(value)),
89	}
90}
91
92fn account(bytes: &mut usize, payload: &mut usize, budget: &MemoryBudget, old: EntryFootprint, new: EntryFootprint) {
93	if new.resident >= old.resident {
94		let delta = new.resident - old.resident;
95		*bytes += delta;
96		budget.charge(ByteSize::from_bytes(delta as u64));
97	} else {
98		let delta = old.resident - new.resident;
99		*bytes -= delta;
100		budget.release(ByteSize::from_bytes(delta as u64));
101	}
102	if new.payload >= old.payload {
103		*payload += new.payload - old.payload;
104	} else {
105		*payload -= old.payload - new.payload;
106	}
107}
108
109pub enum ServedChunk {
110	Served(RangeBatch),
111	Gap,
112}
113
114#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
115pub struct ReadBufferWarmMetrics {
116	pub warms_started: u64,
117	pub warms_completed: u64,
118	pub warms_dirty_aborted: u64,
119	pub warms_aborted: u64,
120	pub pages_warm_blocked: u64,
121	pub pages_evicted: u64,
122	pub complete_pages_invalidated: u64,
123}
124
125#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
126pub struct ReadBufferReadMetrics {
127	pub point_hits: u64,
128	pub previous_hits: u64,
129	pub point_misses: u64,
130	pub range_served: u64,
131	pub range_gaps: u64,
132}
133
134#[derive(Clone, Copy, Debug)]
135pub struct ReadBufferStateMetrics {
136	pub used: ByteSize,
137	pub limit: ByteSize,
138	pub pages: usize,
139	pub page_cap: usize,
140	pub payload: ByteSize,
141	pub entries: usize,
142	pub hot_pages: usize,
143	pub complete_pages: usize,
144	pub blocked_pages: usize,
145	pub warming: usize,
146}
147
148#[derive(Clone, Copy, Debug)]
149pub struct ReadBufferShardMetrics {
150	pub shard: usize,
151	pub state: ReadBufferStateMetrics,
152	pub warms: ReadBufferWarmMetrics,
153	pub reads: ReadBufferReadMetrics,
154}
155
156struct Shard {
157	pages: HashMap<PageId, ResidentPage>,
158	warming: HashMap<PageId, bool>,
159	next_tick: u64,
160	page_cap: usize,
161	budget: MemoryBudget,
162	warm_metrics: ReadBufferWarmMetrics,
163	read_metrics: ReadBufferReadMetrics,
164}
165
166struct PoolInner {
167	shards: Box<[Mutex<Shard>]>,
168	bucket_shift: u8,
169}
170
171#[derive(Clone)]
172pub struct MultiReadBufferTier {
173	inner: Arc<PoolInner>,
174}