reifydb_store_multi/tier/read/
point.rs1use std::{collections::BTreeMap, sync::atomic::Ordering};
5
6use reifydb_codec::key::encoded::EncodedKey;
7use reifydb_core::{
8 common::CommitVersion,
9 interface::store::{EntryKind, classify_key},
10};
11use reifydb_store::row::page::{PageId, page_of};
12use reifydb_value::util::cowvec::CowVec;
13use tracing::instrument;
14
15use crate::tier::{
16 VersionedGetResult,
17 read::{MultiReadBufferTier, PageEntry, ResidentPage},
18};
19
20impl MultiReadBufferTier {
21 pub fn get(&self, key: &EncodedKey, version: CommitVersion) -> VersionedGetResult {
22 match classify_key(key) {
23 EntryKind::Operator(_) | EntryKind::OperatorInternal(_) => self.get_operator(key, version),
24 EntryKind::Source(_) => self.get_source(key, version),
25 _ => self.get_multi(key, version),
26 }
27 }
28
29 #[instrument(name = "store::multi::read::get::operator", level = "trace", skip(self, key), fields(version = version.0))]
30 fn get_operator(&self, key: &EncodedKey, version: CommitVersion) -> VersionedGetResult {
31 self.get_impl(key, version)
32 }
33
34 #[instrument(name = "store::multi::read::get::source", level = "trace", skip(self, key), fields(version = version.0))]
35 fn get_source(&self, key: &EncodedKey, version: CommitVersion) -> VersionedGetResult {
36 self.get_impl(key, version)
37 }
38
39 #[instrument(name = "store::multi::read::get::multi", level = "trace", skip(self, key), fields(version = version.0))]
40 fn get_multi(&self, key: &EncodedKey, version: CommitVersion) -> VersionedGetResult {
41 self.get_impl(key, version)
42 }
43
44 fn get_impl(&self, key: &EncodedKey, version: CommitVersion) -> VersionedGetResult {
45 let page_id = page_of(key, self.bucket_shift());
46 let mut shard = self.shard_for(&page_id).lock();
47 let next = shard.next_tick;
48 let result = {
49 let Some(page) = shard.pages.get_mut(&page_id) else {
50 return VersionedGetResult::NotFound;
51 };
52 let Some(entry) = page.entries.get(key) else {
53 if page.range_complete {
54 page.hot = true;
55 page.tick = next;
56 return VersionedGetResult::Tombstone;
57 }
58 return VersionedGetResult::NotFound;
59 };
60 let served = if entry.version <= version {
61 Some((entry.version, entry.value.clone()))
62 } else {
63 match &entry.previous {
64 Some((prev_version, prev_value)) if *prev_version <= version => {
65 Some((*prev_version, prev_value.clone()))
66 }
67 _ => None,
68 }
69 };
70 let Some((served_version, served_value)) = served else {
71 return VersionedGetResult::NotFound;
72 };
73 let result = match served_value {
74 Some(value) => VersionedGetResult::Value {
75 value,
76 version: served_version,
77 },
78 None => VersionedGetResult::Tombstone,
79 };
80 page.hot = true;
81 page.tick = next;
82 result
83 };
84 shard.next_tick = next + 1;
85 result
86 }
87
88 pub fn insert(&self, key: EncodedKey, version: CommitVersion, value: Option<CowVec<u8>>) {
89 let page_id = page_of(&key, self.bucket_shift());
90 let mut shard = self.shard_for(&page_id).lock();
91 let next = shard.next_tick;
92 match shard.pages.get_mut(&page_id) {
93 Some(page) => {
94 match page.entries.get_mut(&key) {
95 Some(existing) if existing.version > version => return,
96 Some(existing) if existing.version == version => {
97 existing.value = value;
98 existing.previous = None;
99 }
100 Some(existing) => {
101 existing.previous = Some((existing.version, existing.value.take()));
102 existing.version = version;
103 existing.value = value;
104 }
105 None => {
106 page.entries.insert(
107 key,
108 PageEntry {
109 version,
110 value,
111 previous: None,
112 },
113 );
114 }
115 }
116 page.hot = true;
117 page.tick = next;
118 }
119 None => {
120 let mut entries = BTreeMap::new();
121 entries.insert(
122 key,
123 PageEntry {
124 version,
125 value,
126 previous: None,
127 },
128 );
129 shard.pages.insert(
130 page_id,
131 ResidentPage {
132 entries,
133 hot: false,
134 tick: next,
135 range_complete: false,
136 warm_blocked: false,
137 },
138 );
139 }
140 }
141 shard.next_tick = next + 1;
142 shard.evict_to_capacity();
143 }
144
145 pub fn invalidate(&self, key: &EncodedKey) {
146 let page_id = page_of(key, self.bucket_shift());
147 let mut shard = self.shard_for(&page_id).lock();
148 if let Some(dirty) = shard.warming.get_mut(&page_id) {
149 *dirty = true;
150 }
151 let now_empty = match shard.pages.get_mut(&page_id) {
152 Some(page) => {
153 page.entries.remove(key);
154 page.range_complete = false;
155 page.entries.is_empty()
156 }
157 None => false,
158 };
159 if now_empty {
160 shard.pages.remove(&page_id);
161 }
162 }
163
164 pub fn remove_dropped(&self, key: &EncodedKey) {
165 let page_id = page_of(key, self.bucket_shift());
166 let mut shard = self.shard_for(&page_id).lock();
167 if let Some(dirty) = shard.warming.get_mut(&page_id) {
168 *dirty = true;
169 }
170 let now_empty_incomplete = match shard.pages.get_mut(&page_id) {
171 Some(page) => {
172 page.entries.remove(key);
173 page.entries.is_empty() && !page.range_complete
174 }
175 None => false,
176 };
177 if now_empty_incomplete {
178 shard.pages.remove(&page_id);
179 }
180 }
181
182 pub fn remove_dropped_through(&self, key: &EncodedKey, through: CommitVersion) {
183 let page_id = page_of(key, self.bucket_shift());
184 let mut shard = self.shard_for(&page_id).lock();
185 if let Some(dirty) = shard.warming.get_mut(&page_id) {
186 *dirty = true;
187 }
188 let now_empty_incomplete = match shard.pages.get_mut(&page_id) {
189 Some(page) => {
190 if let Some(entry) = page.entries.get_mut(key) {
191 if entry.version <= through {
192 page.entries.remove(key);
193 } else if entry.previous.as_ref().is_some_and(|(v, _)| *v <= through) {
194 entry.previous = None;
195 }
196 }
197 page.entries.is_empty() && !page.range_complete
198 }
199 None => false,
200 };
201 if now_empty_incomplete {
202 shard.pages.remove(&page_id);
203 }
204 }
205
206 pub fn page_is_warm_candidate(&self, page: PageId) -> bool {
207 let shard = self.shard_for(&page).lock();
208 match shard.pages.get(&page) {
209 Some(p) => !p.range_complete && !p.warm_blocked,
210 None => true,
211 }
212 }
213
214 pub fn set_warm_blocked(&self, page: PageId) {
215 let mut shard = self.shard_for(&page).lock();
216 let next = shard.next_tick;
217 shard.pages
218 .entry(page)
219 .or_insert_with(|| ResidentPage {
220 entries: BTreeMap::new(),
221 hot: false,
222 tick: next,
223 range_complete: false,
224 warm_blocked: false,
225 })
226 .warm_blocked = true;
227 }
228
229 pub fn begin_warm(&self, page: PageId) -> bool {
230 let mut shard = self.shard_for(&page).lock();
231 if shard.warming.contains_key(&page) {
232 return false;
233 }
234 shard.warming.insert(page, false);
235 true
236 }
237
238 pub fn abort_warm(&self, page: PageId) {
239 let mut shard = self.shard_for(&page).lock();
240 shard.warming.remove(&page);
241 }
242
243 pub fn clear(&self) {
244 for shard in self.inner.shards.iter() {
245 let mut shard = shard.lock();
246 shard.pages.clear();
247 shard.warming.clear();
248 shard.next_tick = 0;
249 }
250 }
251
252 pub fn set_capacity(&self, resident_pages: usize) {
253 let page_cap = (resident_pages / self.inner.shards.len()).max(1);
254 for shard in self.inner.shards.iter() {
255 let mut shard = shard.lock();
256 shard.page_cap = page_cap;
257 shard.evict_to_capacity();
258 }
259 }
260
261 pub fn reconfigure(&self, resident_pages: usize, page_size_rows: u64) {
262 let bucket_shift = page_size_rows.max(1).trailing_zeros() as u8;
263 let page_cap = (resident_pages / self.inner.shards.len()).max(1);
264 self.inner.bucket_shift.store(bucket_shift, Ordering::Relaxed);
265 for shard in self.inner.shards.iter() {
266 let mut shard = shard.lock();
267 shard.page_cap = page_cap;
268 shard.pages.clear();
269 shard.warming.clear();
270 shard.next_tick = 0;
271 }
272 }
273}