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/oracle.rs
Line
Count
Source
1
use arc_swap::ArcSwap;
2
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
3
use std::sync::Arc;
4
#[cfg(not(target_arch = "wasm32"))]
5
use std::sync::Mutex;
6
#[cfg(not(target_arch = "wasm32"))]
7
use std::thread::JoinHandle;
8
use web_time::{Duration, Instant, SystemTime, UNIX_EPOCH};
9
10
/// A timestamp oracle for monotonically increasing time
11
pub(crate) struct Oracle {
12
  // The inner strcuture of an Oracle
13
  pub(crate) inner: Arc<Inner>,
14
}
15
16
impl Drop for Oracle {
17
26.2k
  fn drop(&mut self) {
18
26.2k
    self.shutdown();
19
26.2k
  }
20
}
21
22
/// The inner structure of the timestamp oracle
23
pub(crate) struct Inner {
24
  /// The latest monotonic counter for this oracle
25
  pub(crate) timestamp: AtomicU64,
26
  /// The reference time when this Oracle was synced
27
  pub(crate) reference: ArcSwap<(u64, Instant)>,
28
  /// Specifies whether timestamp syncing is enabled in the background
29
  pub(crate) resync_enabled: AtomicBool,
30
  /// Stores a handle to the current timestamp syncing background thread
31
  #[cfg(not(target_arch = "wasm32"))]
32
  pub(crate) resync_handle: Mutex<Option<JoinHandle<()>>>,
33
  /// Interval at which the oracle resyncs with the system clock
34
  #[cfg_attr(target_arch = "wasm32", allow(dead_code))]
35
  pub(crate) resync_interval: Duration,
36
}
37
38
impl Oracle {
39
  /// Creates a new timestamp oracle with the specified resync interval
40
26.2k
  pub fn new(resync_interval: Duration) -> Arc<Self> {
41
    // Get the current unix time in nanoseconds
42
26.2k
    let reference_unix = Self::current_unix_ns();
43
    // Get a new monotonically increasing clock
44
26.2k
    let reference_time = Instant::now();
45
    // Return the current timestamp oracle
46
26.2k
    let oracle = Self {
47
26.2k
      inner: Arc::new(Inner {
48
26.2k
        timestamp: AtomicU64::new(reference_unix),
49
26.2k
        reference: ArcSwap::new(Arc::new((reference_unix, reference_time))),
50
26.2k
        resync_enabled: AtomicBool::new(true),
51
26.2k
        #[cfg(not(target_arch = "wasm32"))]
52
26.2k
        resync_handle: Mutex::new(None),
53
26.2k
        resync_interval,
54
26.2k
      }),
55
26.2k
    };
56
    // Start up the resyncing thread
57
    #[cfg(not(target_arch = "wasm32"))]
58
26.2k
    oracle.worker_resync();
59
    // Return the oracle
60
26.2k
    Arc::new(oracle)
61
26.2k
  }
62
63
  /// Returns the current timestamp for this oracle
64
  #[cfg(test)]
65
  #[inline]
66
  pub fn current_timestamp(&self) -> u64 {
67
    self.inner.timestamp.load(Ordering::Acquire)
68
  }
69
70
  /// Gets the current system time in nanoseconds since the Unix epoch
71
  #[inline]
72
52.6k
  pub(crate) fn current_unix_ns() -> u64 {
73
    // Get the current system time
74
52.6k
    let timestamp = SystemTime::now().duration_since(UNIX_EPOCH);
75
    // Count the nanoseconds since the Unix epoch
76
52.6k
    timestamp.unwrap_or_default().as_nanos() as u64
77
52.6k
  }
78
79
  /// Gets the current estimated time in nanoseconds since the Unix epoch
80
  #[inline]
81
90.0k
  pub(crate) fn current_time_ns(&self) -> u64 {
82
    // Get the current reference time
83
90.0k
    let reference = self.inner.reference.load();
84
    // Calculate the nanoseconds since the Unix epoch
85
90.0k
    reference.0 + reference.1.elapsed().as_nanos() as u64
86
90.0k
  }
87
88
  /// Shutdown the oracle resync, waiting for background threads to exit
89
26.2k
  fn shutdown(&self) {
90
    // Disable timestamp resyncing
91
26.2k
    self.inner.resync_enabled.store(false, Ordering::Release);
92
    // Wait for the timestamp resyncing thread to exit
93
    #[cfg(not(target_arch = "wasm32"))]
94
26.2k
    if let Some(handle) = self.inner.resync_handle.lock().unwrap().take() {
95
26.2k
      handle.thread().unpark();
96
26.2k
      handle.join().unwrap();
97
26.2k
    }
98
26.2k
  }
99
100
  /// Start the resyncing thread after creating the oracle
101
  #[cfg(not(target_arch = "wasm32"))]
102
26.2k
  fn worker_resync(&self) {
103
    // Clone the underlying oracle inner
104
26.2k
    let oracle = self.inner.clone();
105
    // Store the resync interval for the thread
106
26.2k
    let interval = oracle.resync_interval;
107
    // Spawn a new thread to handle timestamp resyncing
108
26.2k
    let handle = std::thread::spawn(move || {
109
      // Check whether the timestamp resync process is enabled
110
52.6k
      while oracle.resync_enabled.load(Ordering::Acquire) {
111
26.3k
        // Wait for a specified time interval
112
26.3k
        std::thread::park_timeout(interval);
113
26.3k
        // Get the current unix time in nanoseconds
114
26.3k
        let reference_unix = Self::current_unix_ns();
115
26.3k
        // Get a new monotonically increasing clock
116
26.3k
        let reference_time = Instant::now();
117
26.3k
        // Store the timestamp and monotonic instant
118
26.3k
        oracle.reference.store(Arc::new((reference_unix, reference_time)));
119
26.3k
      }
120
26.2k
    });
121
    // Store and track the thread handle
122
26.2k
    *self.inner.resync_handle.lock().unwrap() = Some(handle);
123
26.2k
  }
124
}