reifydb_store_multi/tier/read/
mod.rs1mod 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}