reifydb_store_multi/tier/read/
range.rs1use std::{collections::BTreeMap, ops::Bound};
5
6use reifydb_codec::key::encoded::{EncodedKey, EncodedKeyRange};
7use reifydb_core::interface::store::EntryKind;
8use reifydb_store::row::page::{PageId, key_range_of, page_of};
9use tracing::instrument;
10
11use crate::{
12 MultiVersionScope,
13 tier::{
14 RangeBatch, RangeCursor, RawEntry,
15 read::{MultiReadBufferTier, PageEntry, ResidentPage, ServedChunk},
16 },
17};
18
19impl MultiReadBufferTier {
20 pub fn page_of_key(&self, key: &EncodedKey) -> PageId {
21 page_of(key, self.bucket_shift())
22 }
23
24 pub fn page_key_range(&self, page: PageId) -> Option<EncodedKeyRange> {
25 key_range_of(page, self.bucket_shift())
26 }
27
28 pub fn invalidate_page(&self, page: PageId) {
29 let mut shard = self.shard_for(&page).lock();
30 shard.pages.remove(&page);
31 }
32
33 pub fn populate_page(&self, page: PageId, entries: Vec<RawEntry>, complete: bool) {
34 let shift = self.bucket_shift();
35 let range_complete = complete && key_range_of(page, shift).is_some();
36 let mut shard = self.shard_for(&page).lock();
37 let next = shard.next_tick;
38 let resident = shard.pages.entry(page).or_insert_with(|| ResidentPage {
39 entries: BTreeMap::new(),
40 hot: false,
41 tick: next,
42 range_complete: false,
43 warm_blocked: false,
44 });
45 for entry in entries {
46 match resident.entries.get(&entry.key) {
47 Some(existing) if existing.version > entry.version => continue,
48 _ => {
49 resident.entries.insert(
50 entry.key,
51 PageEntry {
52 version: entry.version,
53 value: entry.value,
54 previous: None,
55 },
56 );
57 }
58 }
59 }
60 resident.range_complete = range_complete;
61 resident.tick = next;
62 shard.next_tick = next + 1;
63 shard.evict_to_capacity();
64 }
65
66 pub fn finish_warm(&self, page: PageId, entries: Vec<RawEntry>) -> bool {
67 let shift = self.bucket_shift();
68 let range_complete = key_range_of(page, shift).is_some();
69 let mut shard = self.shard_for(&page).lock();
70 let Some(dirty) = shard.warming.remove(&page) else {
71 return false;
72 };
73 if dirty || !range_complete {
74 return false;
75 }
76 let next = shard.next_tick;
77 let resident = shard.pages.entry(page).or_insert_with(|| ResidentPage {
78 entries: BTreeMap::new(),
79 hot: false,
80 tick: next,
81 range_complete: false,
82 warm_blocked: false,
83 });
84 for entry in entries {
85 match resident.entries.get(&entry.key) {
86 Some(existing) if existing.version > entry.version => continue,
87 _ => {
88 resident.entries.insert(
89 entry.key,
90 PageEntry {
91 version: entry.version,
92 value: entry.value,
93 previous: None,
94 },
95 );
96 }
97 }
98 }
99 resident.range_complete = true;
100 resident.tick = next;
101 shard.next_tick = next + 1;
102 shard.evict_to_capacity();
103 true
104 }
105
106 #[allow(clippy::too_many_arguments)]
107 #[instrument(name = "store::multi::read::serve", level = "trace", skip(self, cursor, start, end), fields(table = ?table, descending = descending))]
108 pub fn serve_persistent_chunk(
109 &self,
110 table: EntryKind,
111 cursor: &mut RangeCursor,
112 start: &[u8],
113 end: &[u8],
114 scope: MultiVersionScope,
115 batch_size: usize,
116 descending: bool,
117 ) -> ServedChunk {
118 match table {
119 EntryKind::Source(_) => {}
120 EntryKind::Operator(_) | EntryKind::OperatorInternal(_) => {
121 return self.serve_operator_chunk(cursor, start, end, scope, batch_size, descending);
122 }
123 _ => return ServedChunk::Gap,
124 }
125
126 let shift = self.bucket_shift();
127 let range_lo = EncodedKey::new(start.to_vec());
128 let range_hi = EncodedKey::new(end.to_vec());
129 if range_lo > range_hi {
130 cursor.exhausted = true;
131 return ServedChunk::Served(RangeBatch::empty());
132 }
133
134 let mut out: Vec<RawEntry> = Vec::new();
135 let mut first = true;
136 let mut page = match &cursor.last_key {
137 Some(last) => page_of(last, shift),
138 None if descending => page_of(&range_hi, shift),
139 None => page_of(&range_lo, shift),
140 };
141
142 loop {
143 let Some(page_range) = key_range_of(page, shift) else {
144 if out.is_empty() {
145 return ServedChunk::Gap;
146 }
147 return served_chunk(out, cursor, false);
148 };
149 let (page_start, page_end) = match (page_range.start, page_range.end) {
150 (Bound::Included(s), Bound::Included(e)) => (s, e),
151 _ => {
152 if out.is_empty() {
153 return ServedChunk::Gap;
154 }
155 return served_chunk(out, cursor, false);
156 }
157 };
158
159 if descending {
160 if page_end < range_lo {
161 return served_chunk(out, cursor, true);
162 }
163 } else if page_start > range_hi {
164 return served_chunk(out, cursor, true);
165 }
166
167 let mut shard = self.shard_for(&page).lock();
168 let complete = shard.pages.get(&page).map(|p| p.range_complete).unwrap_or(false);
169 if !complete {
170 drop(shard);
171 if out.is_empty() {
172 return ServedChunk::Gap;
173 }
174 return served_chunk(out, cursor, false);
175 }
176
177 let tick = shard.next_tick;
178 let page_ref = shard.pages.get_mut(&page).expect("complete page present under lock");
179
180 let lo_bound: Bound<EncodedKey> = if first {
181 match &cursor.last_key {
182 Some(last) if !descending && *last >= range_lo => Bound::Excluded(last.clone()),
183 _ => Bound::Included(page_start.clone().max(range_lo.clone())),
184 }
185 } else {
186 Bound::Included(page_start.clone().max(range_lo.clone()))
187 };
188 let hi_bound: Bound<EncodedKey> = if first {
189 match &cursor.last_key {
190 Some(last) if descending && *last <= range_hi => Bound::Excluded(last.clone()),
191 _ => Bound::Included(page_end.clone().min(range_hi.clone())),
192 }
193 } else {
194 Bound::Included(page_end.clone().min(range_hi.clone()))
195 };
196
197 let mut full = false;
198 if descending {
199 for (key, entry) in page_ref.entries.range((lo_bound, hi_bound)).rev() {
200 if out.len() >= batch_size {
201 full = true;
202 break;
203 }
204 if scope.contains(entry.version) {
205 out.push(RawEntry {
206 key: key.clone(),
207 version: entry.version,
208 value: entry.value.clone(),
209 });
210 }
211 }
212 } else {
213 for (key, entry) in page_ref.entries.range((lo_bound, hi_bound)) {
214 if out.len() >= batch_size {
215 full = true;
216 break;
217 }
218 if scope.contains(entry.version) {
219 out.push(RawEntry {
220 key: key.clone(),
221 version: entry.version,
222 value: entry.value.clone(),
223 });
224 }
225 }
226 }
227
228 page_ref.hot = true;
229 page_ref.tick = tick;
230 shard.next_tick = tick + 1;
231 drop(shard);
232
233 if full {
234 return served_chunk(out, cursor, false);
235 }
236
237 if descending {
238 if page_start <= range_lo {
239 return served_chunk(out, cursor, true);
240 }
241 page = PageId {
242 kind: page.kind,
243 bucket: page.bucket + 1,
244 };
245 } else {
246 if page_end >= range_hi {
247 return served_chunk(out, cursor, true);
248 }
249 if page.bucket == 0 {
250 return served_chunk(out, cursor, true);
251 }
252 page = PageId {
253 kind: page.kind,
254 bucket: page.bucket - 1,
255 };
256 }
257 first = false;
258 }
259 }
260
261 fn serve_operator_chunk(
262 &self,
263 cursor: &mut RangeCursor,
264 start: &[u8],
265 end: &[u8],
266 scope: MultiVersionScope,
267 batch_size: usize,
268 descending: bool,
269 ) -> ServedChunk {
270 let shift = self.bucket_shift();
271 let range_lo = EncodedKey::new(start.to_vec());
272 let range_hi = EncodedKey::new(end.to_vec());
273 if range_lo > range_hi {
274 cursor.exhausted = true;
275 return ServedChunk::Served(RangeBatch::empty());
276 }
277
278 let page = page_of(&range_lo, shift);
279 let Some(page_range) = key_range_of(page, shift) else {
280 return ServedChunk::Gap;
281 };
282 let (Bound::Included(page_start), Bound::Included(page_end)) = (page_range.start, page_range.end)
283 else {
284 return ServedChunk::Gap;
285 };
286 if range_lo < page_start || range_hi > page_end {
287 return ServedChunk::Gap;
288 }
289
290 let mut shard = self.shard_for(&page).lock();
291 let complete = shard.pages.get(&page).map(|p| p.range_complete).unwrap_or(false);
292 if !complete {
293 return ServedChunk::Gap;
294 }
295
296 let tick = shard.next_tick;
297 let page_ref = shard.pages.get_mut(&page).expect("complete page present under lock");
298
299 let lo_bound: Bound<EncodedKey> = match &cursor.last_key {
300 Some(last) if !descending && *last >= range_lo => Bound::Excluded(last.clone()),
301 _ => Bound::Included(range_lo.clone()),
302 };
303 let hi_bound: Bound<EncodedKey> = match &cursor.last_key {
304 Some(last) if descending && *last <= range_hi => Bound::Excluded(last.clone()),
305 _ => Bound::Included(range_hi.clone()),
306 };
307
308 let mut out: Vec<RawEntry> = Vec::new();
309 let mut full = false;
310 if descending {
311 for (key, entry) in page_ref.entries.range((lo_bound, hi_bound)).rev() {
312 if out.len() >= batch_size {
313 full = true;
314 break;
315 }
316 if entry.version > scope.read() {
317 return ServedChunk::Gap;
318 }
319 if scope.contains(entry.version) {
320 out.push(RawEntry {
321 key: key.clone(),
322 version: entry.version,
323 value: entry.value.clone(),
324 });
325 }
326 }
327 } else {
328 for (key, entry) in page_ref.entries.range((lo_bound, hi_bound)) {
329 if out.len() >= batch_size {
330 full = true;
331 break;
332 }
333 if entry.version > scope.read() {
334 return ServedChunk::Gap;
335 }
336 if scope.contains(entry.version) {
337 out.push(RawEntry {
338 key: key.clone(),
339 version: entry.version,
340 value: entry.value.clone(),
341 });
342 }
343 }
344 }
345
346 page_ref.hot = true;
347 page_ref.tick = tick;
348 shard.next_tick = tick + 1;
349 drop(shard);
350
351 served_chunk(out, cursor, !full)
352 }
353
354 pub fn page_is_complete(&self, page: PageId) -> bool {
355 let shard = self.shard_for(&page).lock();
356 shard.pages.get(&page).map(|p| p.range_complete).unwrap_or(false)
357 }
358}
359
360fn served_chunk(out: Vec<RawEntry>, cursor: &mut RangeCursor, exhausted: bool) -> ServedChunk {
361 if let Some(last) = out.last() {
362 cursor.last_key = Some(last.key.clone());
363 }
364 cursor.exhausted = exhausted;
365 ServedChunk::Served(RangeBatch {
366 entries: out,
367 has_more: !exhausted,
368 })
369}