Skip to main content

feoxdb/core/store/
internal.rs

1use 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    /// Update an existing record with a new value (Bytes version for zero-copy).
94    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    /// Get access to hash table (for TTL cleaner)
173    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    /// Remove from tree (for TTL cleaner)
192    pub(crate) fn remove_from_tree(&self, key: &[u8]) {
193        self.tree.remove(key);
194    }
195
196    /// Get write buffer (for TTL cleaner)
197    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}