Skip to main content

reifydb_store_multi/tier/read/
range.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}