Skip to main content

reifydb_store_multi/tier/read/
point.rs

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