Coverage Report

Created: 2026-07-13 08:11

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}