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