Skip to main content

tuwunel_database/
engine.rs

1//! RocksDB engine: database-wide operations and shared resources.
2//!
3//! `Engine` owns the opened RocksDB instance together with the worker pool, the
4//! shared open-time context, and the flags fixed at open (read-only, secondary,
5//! checksums). Per-column-family reads and writes go through `Map`; the methods
6//! here act on the database as a whole: WAL flush and sync, memtable flush,
7//! manual compaction and primary catch-up, checkpoints, property queries, and
8//! the cork counter that coalesces WAL writes (see the `cork` module).
9
10mod backup;
11mod cf_opts;
12pub(crate) mod context;
13mod db_opts;
14pub(crate) mod descriptor;
15mod env;
16mod events;
17mod files;
18mod logger;
19mod memory_usage;
20mod open;
21mod repair;
22#[cfg(test)]
23mod tests;
24
25use std::{
26	collections::BTreeMap,
27	ffi::CStr,
28	path::Path,
29	sync::{
30		Arc, OnceLock, Weak,
31		atomic::{AtomicU32, Ordering},
32	},
33};
34
35use rocksdb::{
36	AsColumnFamilyRef, BoundColumnFamily, DBCommon, DBWithThreadMode, FlushOptions,
37	MultiThreaded, WaitForCompactOptions, WriteOptions, checkpoint::Checkpoint,
38};
39use tuwunel_core::{Err, Result, debug, implement, info, warn};
40
41use crate::{
42	Context, Map,
43	pool::Pool,
44	util::{map_err, result},
45};
46
47pub(crate) type CfIndex = BTreeMap<u32, Weak<Map>>;
48
49/// Handle to the opened RocksDB database and its shared resources.
50///
51/// One `Engine` exists per database, shared behind an `Arc` by every `Map`.
52pub struct Engine {
53	/// The opened RocksDB instance.
54	pub(crate) db: Db,
55
56	/// Thread pool offloading uncached, blocking database requests from the
57	/// tokio workers.
58	pub(crate) pool: Arc<Pool>,
59
60	/// Resources constructed before the database is opened and outliving it
61	/// (block caches, environment, column descriptors).
62	pub(crate) ctx: Arc<Context>,
63
64	/// Database was opened read-only; writes are rejected.
65	pub(super) read_only: bool,
66
67	/// Database was opened as a secondary follower of a primary instance.
68	pub(super) secondary: bool,
69
70	/// Verify block checksums on read.
71	pub(crate) checksums: bool,
72
73	/// Shared write options for atomic batch commits.
74	pub(crate) write_options: WriteOptions,
75
76	/// Resolves catalog column ids for post-commit watcher notification.
77	/// Runtime migration column families are intentionally absent.
78	cf_index: OnceLock<CfIndex>,
79
80	/// Live cork count; nonzero suppresses the per-write WAL flush.
81	corks: AtomicU32,
82}
83
84/// Backing RocksDB type: multi-threaded column-family access, no transactions.
85pub(crate) type Db = DBWithThreadMode<MultiThreaded>;
86
87impl Engine {
88	/// Block until outstanding background compactions finish.
89	///
90	/// Waits without a timeout and does not flush first; aborts the wait if
91	/// compaction has been paused.
92	#[tracing::instrument(
93		level = "info",
94		skip_all,
95		fields(
96			sequence = ?self.current_sequence(),
97		),
98	)]
99	pub fn wait_compactions_blocking(&self) -> Result {
100		let mut opts = WaitForCompactOptions::default();
101		opts.set_abort_on_pause(true);
102		opts.set_flush(false);
103		opts.set_timeout(0);
104
105		self.db.wait_for_compact(&opts).map_err(map_err)
106	}
107
108	/// Flush every column family's memtable to SST files (a RocksDB LSM-tree
109	/// flush).
110	///
111	/// Forces buffered writes out of memory into the on-disk LSM tree for all
112	/// opened families in one request, atomically when `rocksdb_atomic_flush`
113	/// is set. An LSM flush, not a libc `fflush(3)` or `fsync(2)`, and distinct
114	/// from the `flush` and `sync` methods here, which act on the write-ahead
115	/// log.
116	///
117	/// # Panics
118	///
119	/// Panics on a read-only or secondary database, which cannot flush.
120	#[tracing::instrument(
121		level = "info",
122		skip_all,
123		fields(
124			sequence = ?self.current_sequence(),
125		),
126	)]
127	pub fn sort(&self) -> Result {
128		assert!(!self.is_read_only(), "memtables cannot be flushed on a read-only database");
129
130		let cfs: Vec<_> = self
131			.db
132			.cf_names()
133			.iter()
134			.filter_map(|name| self.db.cf_handle(name))
135			.collect();
136
137		let cfs: Vec<&_> = cfs.iter().collect();
138		let opts = FlushOptions::default();
139
140		result(DBCommon::flush_cfs_opt(&self.db, &cfs, &opts))
141	}
142
143	/// Creates a physical checkpoint of the database at `path`.
144	///
145	/// `log_size` is forwarded to RocksDB as the write-ahead log size threshold
146	/// for flushing before checkpoint creation. A value of zero forces a flush.
147	#[tracing::instrument(level = "info", skip(self))]
148	pub fn checkpoint(&self, path: &Path, log_size: u64) -> Result {
149		let checkpoint = Checkpoint::new(&self.db).map_err(map_err)?;
150
151		checkpoint
152			.create_checkpoint_with_log_size(path, log_size)
153			.map_err(map_err)
154	}
155
156	/// Catch a secondary instance up to the primary's latest writes.
157	///
158	/// Replays the primary's newly appended WAL into this instance's view;
159	/// meaningful only when the database was opened as a secondary.
160	#[tracing::instrument(
161		level = "debug",
162		skip_all,
163		fields(
164			sequence = ?self.current_sequence(),
165		),
166	)]
167	pub fn update(&self) -> Result {
168		self.db
169			.try_catch_up_with_primary()
170			.map_err(map_err)
171	}
172
173	/// Flush the write-ahead log and fsync it to disk.
174	///
175	/// Once this returns the buffered writes survive power loss. Heavier than
176	/// `flush`, which stops at the OS page cache.
177	#[tracing::instrument(level = "info", skip_all)]
178	pub fn sync(&self) -> Result { result(DBCommon::flush_wal(&self.db, true)) }
179
180	/// Flush the buffered write-ahead log to the OS without an fsync.
181	///
182	/// Pushes WAL bytes to the page cache (durable against process crash, not
183	/// power loss). This is the per-write flush that corking suppresses.
184	#[tracing::instrument(level = "debug", skip_all)]
185	pub fn flush(&self) -> Result { result(DBCommon::flush_wal(&self.db, false)) }
186
187	/// Increment the cork count, suppressing the per-write WAL flush.
188	#[inline]
189	pub(crate) fn cork(&self) { self.corks.fetch_add(1, Ordering::Relaxed); }
190
191	/// Decrement the cork count; the per-write flush resumes at zero.
192	#[inline]
193	pub(crate) fn uncork(&self) { self.corks.fetch_sub(1, Ordering::Relaxed); }
194
195	/// Whether any cork is currently held.
196	///
197	/// When true, `Map` insert and remove skip their post-write WAL flush so
198	/// the records coalesce into one batch. Corking is purely a backend
199	/// write-buffering signal: it never changes application logic or any
200	/// observable database API behavior, because a write lands in the memtable
201	/// synchronously and reads back regardless of WAL flush state. See the
202	/// `cork` module.
203	#[inline]
204	pub fn corked(&self) -> bool { self.corks.load(Ordering::Relaxed) > 0 }
205
206	/// Query for database property by null-terminated name which is expected to
207	/// have a result with an integer representation. This is intended for
208	/// low-overhead programmatic use.
209	pub(crate) fn property_integer(
210		&self,
211		cf: &impl AsColumnFamilyRef,
212		name: &CStr,
213	) -> Result<u64> {
214		result(self.db.property_int_value_cf(cf, name))
215			.and_then(|val| val.map_or_else(|| Err!("Property {name:?} not found."), Ok))
216	}
217
218	/// Query for database property by name receiving the result in a string.
219	pub(crate) fn property(&self, cf: &impl AsColumnFamilyRef, name: &str) -> Result<String> {
220		result(self.db.property_value_cf(cf, name))
221			.and_then(|val| val.map_or_else(|| Err!("Property {name:?} not found."), Ok))
222	}
223
224	/// Look up a column-family handle by name.
225	///
226	/// The handle refers to a family opened with this database and remains tied
227	/// to the engine's lifetime.
228	///
229	/// # Panics
230	///
231	/// Panics if the family was not described before the database was opened.
232	pub(crate) fn cf(&self, name: &str) -> Arc<BoundColumnFamily<'_>> {
233		self.db
234			.cf_handle(name)
235			.expect("column must be described prior to database open")
236	}
237
238	/// Reports whether a column family with this name exists.
239	///
240	/// The lookup consults the handles currently opened by RocksDB. It does not
241	/// create a missing family.
242	#[inline]
243	#[must_use]
244	pub fn has_cf(&self, name: &str) -> bool { self.db.cf_handle(name).is_some() }
245
246	/// Returns the latest RocksDB sequence number.
247	///
248	/// RocksDB assigns sequence numbers to committed writes, so this value
249	/// marks the engine's current write position. The number is local to this
250	/// database.
251	#[inline]
252	#[must_use]
253	#[tracing::instrument(
254		name = "sequence",
255		level = "debug",
256		skip_all,
257		fields(sequence)
258	)]
259	pub fn current_sequence(&self) -> u64 {
260		let sequence = self.db.latest_sequence_number();
261
262		#[cfg(debug_assertions)]
263		tracing::Span::current().record("sequence", sequence);
264
265		sequence
266	}
267
268	/// Reports whether this engine rejects writes.
269	///
270	/// Both read-only and secondary opens reject writes through their database
271	/// handle. A writable primary open returns false.
272	#[inline]
273	#[must_use]
274	pub fn is_read_only(&self) -> bool { self.secondary || self.read_only }
275
276	/// Reports whether the database follows a primary as a secondary.
277	///
278	/// A secondary advances its view when [`Self::update`] catches up with the
279	/// primary. Writes through the secondary handle are rejected.
280	#[inline]
281	#[must_use]
282	pub fn is_secondary(&self) -> bool { self.secondary }
283}
284
285#[implement(Engine)]
286pub(crate) fn set_cf_index(&self, index: CfIndex) {
287	self.cf_index
288		.set(index)
289		.expect("cf_index initialized twice");
290}
291
292#[implement(Engine)]
293#[inline]
294pub(crate) fn map_by_cf_id(&self, cf_id: u32) -> Option<Arc<Map>> {
295	self.cf_index
296		.get()
297		.expect("cf_index initialized before writes")
298		.get(&cf_id)
299		.and_then(Weak::upgrade)
300}
301
302impl Drop for Engine {
303	#[cold]
304	fn drop(&mut self) {
305		const BLOCKING: bool = true;
306
307		debug!("Waiting for background tasks to finish...");
308		self.db.cancel_all_background_work(BLOCKING);
309
310		info!(
311			sequence = %self.current_sequence(),
312			"Closing database..."
313		);
314	}
315}