feoxdb/core/store/
internal.rs1use ahash::RandomState;
2use bytes::Bytes;
3use scc::HashMap;
4use std::sync::atomic::Ordering;
5use std::sync::Arc;
6
7use crate::constants::Operation;
8use crate::core::record::{Record, TreeSlot};
9use crate::error::{FeoxError, Result};
10use crate::storage::write_buffer::WriteBuffer;
11
12use super::FeoxStore;
13
14impl FeoxStore {
15 pub(super) fn update_record_with_ttl(
16 &self,
17 old_record: &Record,
18 value: &[u8],
19 timestamp: u64,
20 explicit_timestamp: bool,
21 ttl_expiry: u64,
22 ) -> Result<bool> {
23 let new_record = if ttl_expiry > 0 && self.enable_ttl {
24 Arc::new(Record::new_with_timestamp_ttl(
25 old_record.key.clone(),
26 value.to_vec(),
27 timestamp,
28 ttl_expiry,
29 ))
30 } else {
31 Arc::new(Record::new(
32 old_record.key.clone(),
33 value.to_vec(),
34 timestamp,
35 ))
36 };
37
38 let new_size = self.calculate_record_size(old_record.key.len(), value.len());
39
40 let key_vec = new_record.key.clone();
41
42 let old_record_arc = match self.hash_table.entry(key_vec.clone()) {
43 scc::hash_map::Entry::Occupied(mut entry) => {
44 let old_record_arc = Arc::clone(entry.get());
45 if !std::ptr::eq(old_record, old_record_arc.as_ref())
46 && timestamp <= old_record.retirement_timestamp()
47 {
48 return Err(FeoxError::OlderTimestamp);
49 }
50 if timestamp <= old_record_arc.timestamp {
51 return Err(FeoxError::OlderTimestamp);
52 }
53 let old_size = old_record_arc.calculate_size();
54 let reservation = self.reserve_memory(new_size.saturating_sub(old_size))?;
55 old_record_arc.link_successor(&new_record);
56 old_record_arc.refcount.store(0, Ordering::Release);
57 entry.insert(Arc::clone(&new_record));
58 self.publish_to_tree(&key_vec, Arc::clone(&new_record));
59 self.observe_published_timestamp(&key_vec, timestamp, explicit_timestamp);
60 reservation.commit();
61 if old_size > new_size {
62 self.release_memory(old_size - new_size);
63 }
64 self.note_ttl_transition(
65 old_record_arc.ttl_expiry.load(Ordering::Acquire),
66 new_record.ttl_expiry.load(Ordering::Acquire),
67 );
68 old_record_arc
69 }
70 scc::hash_map::Entry::Vacant(_entry) => {
71 if timestamp <= old_record.retirement_timestamp() {
72 return Err(FeoxError::OlderTimestamp);
73 }
74 return Err(FeoxError::KeyNotFound);
75 }
76 };
77
78 if !self.memory_only {
79 if self.enable_caching {
80 if let Some(ref cache) = self.cache {
81 cache.remove_for_record(&key_vec, &old_record_arc);
82 }
83 }
84
85 if let Some(ref wb) = self.write_buffer {
86 wb.add_replacement(new_record, old_record_arc)?;
87 }
88 }
89
90 Ok(false)
91 }
92
93 pub(super) fn update_record_with_ttl_bytes(
95 &self,
96 old_record: &Record,
97 value: Bytes,
98 timestamp: u64,
99 explicit_timestamp: bool,
100 ttl_expiry: u64,
101 ) -> Result<bool> {
102 let new_record = if ttl_expiry > 0 && self.enable_ttl {
103 Arc::new(Record::new_from_bytes_with_ttl(
104 old_record.key.clone(),
105 value,
106 timestamp,
107 ttl_expiry,
108 ))
109 } else {
110 Arc::new(Record::new_from_bytes(
111 old_record.key.clone(),
112 value,
113 timestamp,
114 ))
115 };
116
117 let new_size = new_record.calculate_size();
118
119 let key_vec = new_record.key.clone();
120
121 let old_record_arc = match self.hash_table.entry(key_vec.clone()) {
122 scc::hash_map::Entry::Occupied(mut entry) => {
123 let old_record_arc = Arc::clone(entry.get());
124 if !std::ptr::eq(old_record, old_record_arc.as_ref())
125 && timestamp <= old_record.retirement_timestamp()
126 {
127 return Err(FeoxError::OlderTimestamp);
128 }
129 if timestamp <= old_record_arc.timestamp {
130 return Err(FeoxError::OlderTimestamp);
131 }
132 let old_size = old_record_arc.calculate_size();
133 let reservation = self.reserve_memory(new_size.saturating_sub(old_size))?;
134 old_record_arc.link_successor(&new_record);
135 old_record_arc.refcount.store(0, Ordering::Release);
136 entry.insert(Arc::clone(&new_record));
137 self.publish_to_tree(&key_vec, Arc::clone(&new_record));
138 self.observe_published_timestamp(&key_vec, timestamp, explicit_timestamp);
139 reservation.commit();
140 if old_size > new_size {
141 self.release_memory(old_size - new_size);
142 }
143 self.note_ttl_transition(
144 old_record_arc.ttl_expiry.load(Ordering::Acquire),
145 new_record.ttl_expiry.load(Ordering::Acquire),
146 );
147 old_record_arc
148 }
149 scc::hash_map::Entry::Vacant(_entry) => {
150 if timestamp <= old_record.retirement_timestamp() {
151 return Err(FeoxError::OlderTimestamp);
152 }
153 return Err(FeoxError::KeyNotFound);
154 }
155 };
156
157 if !self.memory_only {
158 if self.enable_caching {
159 if let Some(ref cache) = self.cache {
160 cache.remove_for_record(&key_vec, &old_record_arc);
161 }
162 }
163
164 if let Some(ref wb) = self.write_buffer {
165 wb.add_replacement(new_record, old_record_arc)?;
166 }
167 }
168
169 Ok(false)
170 }
171
172 pub(crate) fn get_hash_table(&self) -> &HashMap<Vec<u8>, Arc<Record>, RandomState> {
174 &self.hash_table
175 }
176
177 #[inline]
178 pub(super) fn publish_to_tree(&self, key: &[u8], record: Arc<Record>) {
179 self.tree
180 .get(key)
181 .expect("missing ordered index entry")
182 .value()
183 .store(record);
184 }
185
186 #[inline]
187 pub(super) fn insert_into_tree(&self, key: Vec<u8>, record: Arc<Record>) {
188 self.tree.insert(key, TreeSlot::new(record));
189 }
190
191 pub(crate) fn remove_from_tree(&self, key: &[u8]) {
193 self.tree.remove(key);
194 }
195
196 pub(crate) fn get_write_buffer(&self) -> Option<&Arc<WriteBuffer>> {
198 self.write_buffer.as_ref()
199 }
200
201 pub(crate) fn remove_cached(&self, key: &[u8], record: &Arc<Record>) {
202 if !self.memory_only && self.enable_caching {
203 if let Some(ref cache) = self.cache {
204 cache.remove_for_record(key, record);
205 }
206 }
207 }
208
209 pub(crate) fn note_expired_record(&self, record_size: usize) {
210 self.stats.record_count.fetch_sub(1, Ordering::Relaxed);
211 self.stats
212 .memory_usage
213 .fetch_sub(record_size, Ordering::Relaxed);
214 let _ =
215 self.stats
216 .keys_with_ttl
217 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |count| {
218 Some(count.saturating_sub(1))
219 });
220 self.stats
221 .ttl_expired_active
222 .fetch_add(1, Ordering::Relaxed);
223 }
224
225 pub(super) fn retire_expired_if_current(
226 &self,
227 key: &[u8],
228 expected: &Arc<Record>,
229 now: u64,
230 ) -> Result<bool> {
231 let retired = match self.hash_table.entry(key.to_vec()) {
232 scc::hash_map::Entry::Occupied(entry) => {
233 let record = entry.get();
234 let expiry = record.ttl_expiry.load(Ordering::Acquire);
235 if !Arc::ptr_eq(record, expected) || expiry == 0 || expiry >= now {
236 return Ok(false);
237 }
238
239 let record = Arc::clone(record);
240 let record_size = record.calculate_size();
241 self.version_clock.observe(key, now);
242 record.retired_at.store(now, Ordering::Release);
243 record.refcount.store(0, Ordering::Release);
244 self.tree.remove(key);
245 self.note_expired_record(record_size);
246 let _ = entry.remove();
247 record
248 }
249 scc::hash_map::Entry::Vacant(_) => return Ok(false),
250 };
251
252 let record = retired;
253 self.remove_cached(key, &record);
254 if let Some(write_buffer) = self.write_buffer.as_ref() {
255 write_buffer.add_write(Operation::Delete, record, 0)?;
256 }
257 Ok(true)
258 }
259}