/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 | | } |