/rust/registry/src/index.crates.io-1949cf8c6b5b557f/surrealmx-0.22.0/src/inner.rs
Line | Count | Source |
1 | | // Copyright © SurrealDB Ltd |
2 | | // |
3 | | // Licensed under the Apache License, Version 2.0 (the "License"); |
4 | | // you may not use this file except in compliance with the License. |
5 | | // You may obtain a copy of the License at |
6 | | // |
7 | | // http://www.apache.org/licenses/LICENSE-2.0 |
8 | | // |
9 | | // Unless required by applicable law or agreed to in writing, software |
10 | | // distributed under the License is distributed on an "AS IS" BASIS, |
11 | | // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
12 | | // See the License for the specific language governing permissions and |
13 | | // limitations under the License. |
14 | | |
15 | | //! This module stores the inner in-memory database type. |
16 | | |
17 | | use crate::oracle::Oracle; |
18 | | #[cfg(not(target_arch = "wasm32"))] |
19 | | use crate::persistence::Persistence; |
20 | | use crate::queue::{Commit, Merge}; |
21 | | use crate::versions::Versions; |
22 | | use crate::DatabaseOptions; |
23 | | use bytes::Bytes; |
24 | | use crossbeam_queue::SegQueue; |
25 | | use crossbeam_skiplist::SkipMap; |
26 | | use parking_lot::RwLock; |
27 | | use std::collections::HashSet; |
28 | | use std::sync::atomic::{fence, AtomicBool, AtomicU64, Ordering}; |
29 | | use std::sync::Arc; |
30 | | #[cfg(not(target_arch = "wasm32"))] |
31 | | use std::thread::JoinHandle; |
32 | | use std::time::Duration; |
33 | | |
34 | | /// Sentinel value stored in a counter entry while its owning [`Inner`] |
35 | | /// SkipMap entry is being removed by [`crate::tx::Transaction::drop`]. A |
36 | | /// concurrent `register_counter` will observe this marker and retry with a |
37 | | /// fresh entry rather than incrementing a detached counter. |
38 | | pub(crate) const COUNTER_TOMBSTONE: u64 = u64::MAX; |
39 | | |
40 | | /// The inner structure of the transactional in-memory database |
41 | | pub struct Inner { |
42 | | /// The timestamp version oracle |
43 | | pub(crate) oracle: Arc<Oracle>, |
44 | | /// The underlying lock-free skip-list datastructure |
45 | | pub(crate) datastore: SkipMap<Bytes, RwLock<Versions>>, |
46 | | /// A count of total transactions grouped by oracle version |
47 | | pub(crate) counter_by_oracle: SkipMap<u64, Arc<AtomicU64>>, |
48 | | /// A count of total transactions grouped by commit id |
49 | | pub(crate) counter_by_commit: SkipMap<u64, Arc<AtomicU64>>, |
50 | | /// The transaction commit queue attempt sequence number |
51 | | pub(crate) transaction_queue_id: AtomicU64, |
52 | | /// The transaction commit queue success sequence number |
53 | | pub(crate) transaction_commit_id: AtomicU64, |
54 | | /// The transaction merge queue attempt sequence number |
55 | | pub(crate) transaction_merge_id: AtomicU64, |
56 | | /// The transaction commit queue list of modifications |
57 | | pub(crate) transaction_commit_queue: SkipMap<u64, Arc<Commit>>, |
58 | | /// Transaction updates which are committed but not yet applied |
59 | | pub(crate) transaction_merge_queue: SkipMap<u64, Arc<Merge>>, |
60 | | /// The epoch duration to determine how long to store versioned data |
61 | | pub(crate) garbage_collection_epoch: RwLock<Option<Duration>>, |
62 | | /// Optional persistence handler |
63 | | #[cfg(not(target_arch = "wasm32"))] |
64 | | pub(crate) persistence: RwLock<Option<Arc<Persistence>>>, |
65 | | /// Specifies whether background worker threads are enabled |
66 | | pub(crate) background_threads_enabled: AtomicBool, |
67 | | /// Stores a handle to the current transaction cleanup background thread |
68 | | #[cfg(not(target_arch = "wasm32"))] |
69 | | pub(crate) transaction_cleanup_handle: RwLock<Option<JoinHandle<()>>>, |
70 | | /// Stores a handle to the current garbage collection background thread |
71 | | #[cfg(not(target_arch = "wasm32"))] |
72 | | pub(crate) garbage_collection_handle: RwLock<Option<JoinHandle<()>>>, |
73 | | /// Keys with stale versions pending incremental garbage collection |
74 | | pub(crate) gc_dirty_keys: SegQueue<Bytes>, |
75 | | /// Watermark below which versions are about to be (or have been) |
76 | | /// reclaimed by the garbage collector. The GC sweeper publishes its |
77 | | /// intended `cleanup_ts` here with a SeqCst `fetch_max` *before* |
78 | | /// actually reclaiming any versions. `register_counter` validates |
79 | | /// `v_r >= gc_floor` after publishing its counter; if a reader |
80 | | /// finds `gc_floor > v_r`, it rolls back and reloads the oracle, |
81 | | /// landing on a fresh snapshot above the floor. This closes the |
82 | | /// scan/gc race in the BG sweeper: a sweeper that misses an |
83 | | /// in-flight reader on its first scan publishes `gc_floor`, fences, |
84 | | /// then re-scans — and either the reader saw the new floor and |
85 | | /// retried, or its publish is now visible to the re-scan so the |
86 | | /// sweeper's final cleanup_ts is bounded by it. |
87 | | pub(crate) gc_floor: AtomicU64, |
88 | | /// Threshold after which transaction state is reset |
89 | | pub(crate) reset_threshold: usize, |
90 | | } |
91 | | |
92 | | impl Inner { |
93 | | /// Create a new [`Inner`] structure with the given oracle resync interval. |
94 | 26.2k | pub fn new(opts: &DatabaseOptions) -> Self { |
95 | 26.2k | Self { |
96 | 26.2k | oracle: Oracle::new(opts.resync_interval), |
97 | 26.2k | datastore: SkipMap::new(), |
98 | 26.2k | counter_by_oracle: SkipMap::new(), |
99 | 26.2k | counter_by_commit: SkipMap::new(), |
100 | 26.2k | transaction_queue_id: AtomicU64::new(0), |
101 | 26.2k | transaction_commit_id: AtomicU64::new(0), |
102 | 26.2k | transaction_merge_id: AtomicU64::new(0), |
103 | 26.2k | transaction_commit_queue: SkipMap::new(), |
104 | 26.2k | transaction_merge_queue: SkipMap::new(), |
105 | 26.2k | garbage_collection_epoch: RwLock::new(None), |
106 | 26.2k | #[cfg(not(target_arch = "wasm32"))] |
107 | 26.2k | persistence: RwLock::new(None), |
108 | 26.2k | background_threads_enabled: AtomicBool::new(true), |
109 | 26.2k | #[cfg(not(target_arch = "wasm32"))] |
110 | 26.2k | transaction_cleanup_handle: RwLock::new(None), |
111 | 26.2k | #[cfg(not(target_arch = "wasm32"))] |
112 | 26.2k | garbage_collection_handle: RwLock::new(None), |
113 | 26.2k | gc_dirty_keys: SegQueue::new(), |
114 | 26.2k | gc_floor: AtomicU64::new(0), |
115 | 26.2k | reset_threshold: opts.reset_threshold, |
116 | 26.2k | } |
117 | 26.2k | } |
118 | | } |
119 | | |
120 | | impl Inner { |
121 | | /// Returns the earliest active reader's snapshot version, or `fallback` |
122 | | /// if no reader is currently registered. Pairs with the SeqCst |
123 | | /// load-and-fence in `register_counter` to give writers a watermark |
124 | | /// that observes every reader whose registration is totally ordered |
125 | | /// before the fence below. |
126 | | #[inline] |
127 | 0 | pub(crate) fn earliest_active_version(&self, fallback: u64) -> u64 { |
128 | 0 | earliest_active(&self.counter_by_oracle, fallback) |
129 | 0 | } |
130 | | |
131 | | /// Returns the earliest active reader's start commit id, or `fallback` |
132 | | /// if no reader is currently registered. See `earliest_active_version`. |
133 | | #[inline] |
134 | 551 | pub(crate) fn earliest_active_commit(&self, fallback: u64) -> u64 { |
135 | 551 | earliest_active(&self.counter_by_commit, fallback) |
136 | 551 | } |
137 | | |
138 | | /// Compute the next `cleanup_ts`, publishing it into `gc_floor` so |
139 | | /// concurrent `register_counter` retries any reader whose snapshot |
140 | | /// is below it, then re-scan to bound by any newly-arrived reader. |
141 | | /// |
142 | | /// The proposed value is capped at the current oracle timestamp. |
143 | | /// Without that cap, an idle database (oracle frozen while wall |
144 | | /// clock advances) would push `gc_floor` above any value a future |
145 | | /// reader could load — causing `register_counter` to spin forever |
146 | | /// retrying. |
147 | 0 | pub(crate) fn compute_cleanup_ts(&self) -> u64 { |
148 | 0 | let now = self.oracle.current_time_ns(); |
149 | 0 | let history = self.garbage_collection_epoch.read().unwrap_or_default().as_nanos(); |
150 | 0 | let history_cutoff = now.saturating_sub(history as u64); |
151 | 0 | let earliest_tx = self.earliest_active_version(now); |
152 | 0 | let oracle_now = self.oracle.inner.timestamp.load(Ordering::SeqCst); |
153 | | // `gc_floor` must stay <= oracle so any future reader's load |
154 | | // (which returns >= current oracle) satisfies the floor check. |
155 | 0 | let proposed = history_cutoff.min(earliest_tx).min(oracle_now); |
156 | | // Publish proposed cleanup_ts via `gc_floor` BEFORE reclaiming. |
157 | | // A `register_counter` that fences after CAS-publish will see |
158 | | // either the old floor (and its publish becomes visible to our |
159 | | // re-scan below via fence-fence SC ordering) or the new floor |
160 | | // (and will retry to a fresh snapshot above it). |
161 | 0 | self.gc_floor.fetch_max(proposed, Ordering::SeqCst); |
162 | 0 | fence(Ordering::SeqCst); |
163 | 0 | let earliest_after = self.earliest_active_version(now); |
164 | 0 | proposed.min(earliest_after) |
165 | 0 | } |
166 | | |
167 | | /// Drain the dirty-key queue, reclaiming stale versions on each key and |
168 | | /// unlinking any whose version chain becomes empty. |
169 | | /// |
170 | | /// Keys are de-duplicated within a pass: a hot key committed many times |
171 | | /// between gc cycles is enqueued once per commit, but reclaiming it more |
172 | | /// than once under a fixed `cleanup_ts` is wasted lock traffic — every |
173 | | /// pass after the first is a no-op. A version added by a commit that |
174 | | /// re-dirties the key mid-drain is newer than `cleanup_ts` and so not |
175 | | /// reclaimable this pass anyway; it is caught on the next commit or the |
176 | | /// periodic full scan. |
177 | 0 | pub(crate) fn run_gc_dirty_inner(&self, cleanup_ts: u64) { |
178 | 0 | let mut seen = HashSet::new(); |
179 | | // Drain all keys from the dirty queue |
180 | 0 | while let Some(key) = self.gc_dirty_keys.pop() { |
181 | | // Skip keys already reclaimed in this pass |
182 | 0 | if !seen.insert(key.clone()) { |
183 | 0 | continue; |
184 | 0 | } |
185 | | // Reclaim stale versions on this key |
186 | 0 | self.gc_key(&key, cleanup_ts); |
187 | | } |
188 | 0 | } |
189 | | |
190 | | /// Scan the entire datastore, reclaiming stale versions on every key. |
191 | 0 | pub(crate) fn run_gc_full(&self, cleanup_ts: u64) { |
192 | | // Iterate over the entire datastore |
193 | 0 | for entry in self.datastore.iter() { |
194 | | // Get a mutable reference to the versions list |
195 | 0 | let mut versions = entry.value().write(); |
196 | | // Clean up unnecessary older versions |
197 | 0 | if versions.gc_older_versions(cleanup_ts) == 0 { |
198 | 0 | // Remove under the version write lock (see `gc_key`). |
199 | 0 | entry.remove(); |
200 | 0 | } |
201 | | } |
202 | 0 | } |
203 | | |
204 | | /// Reclaim stale versions on a single key, unlinking the entry if its |
205 | | /// version chain becomes empty. |
206 | 0 | fn gc_key(&self, key: &Bytes, cleanup_ts: u64) { |
207 | | // Look the key up in the datastore |
208 | 0 | if let Some(entry) = self.datastore.get(key) { |
209 | | // Get a mutable reference to the versions list |
210 | 0 | let mut versions = entry.value().write(); |
211 | | // Clean up unnecessary older versions |
212 | 0 | if versions.gc_older_versions(cleanup_ts) == 0 { |
213 | 0 | // Remove the entry while still holding the version write lock, |
214 | 0 | // so a committer blocked on that lock observes `is_removed()` |
215 | 0 | // and re-inserts rather than writing into a node we are about |
216 | 0 | // to unlink. `Entry::remove` also unlinks at the cursor with |
217 | 0 | // no second key lookup. |
218 | 0 | entry.remove(); |
219 | 0 | } |
220 | 0 | } |
221 | 0 | } |
222 | | } |
223 | | |
224 | | #[inline] |
225 | 551 | fn earliest_active(map: &SkipMap<u64, Arc<AtomicU64>>, fallback: u64) -> u64 { |
226 | 551 | fence(Ordering::SeqCst); |
227 | 551 | for entry in map.iter() { |
228 | 10 | let c = entry.value().load(Ordering::Acquire); |
229 | 10 | if c != 0 && c != COUNTER_TOMBSTONE { |
230 | 10 | return *entry.key(); |
231 | 0 | } |
232 | | } |
233 | 541 | fallback |
234 | 551 | } |
235 | | |
236 | | impl Default for Inner { |
237 | 0 | fn default() -> Self { |
238 | 0 | Self::new(&DatabaseOptions::default()) |
239 | 0 | } |
240 | | } |