Skip to main content

tuwunel_database/engine/
env.rs

1use std::{
2	ptr::eq as ptr_eq,
3	sync::{Arc, LockResult, Mutex, MutexGuard, PoisonError, Weak},
4};
5
6use tuwunel_core::{Result, Server, debug, implement};
7
8use crate::or_else;
9
10/// The shared rocksdb environment.
11///
12/// The inner mutex guards the handle itself, which the engine needs while
13/// opening a database or a backup. Releasing the last reference shuts down and
14/// joins the environment's background threads, so an engine must hold one for
15/// as long as it is open.
16pub(super) struct Env(Mutex<rocksdb::Env>);
17
18/// The process-global rocksdb environment, held weakly so it lives exactly as
19/// long as some context needs it.
20///
21/// `Env::new` returns rocksdb's default environment singleton, so every handle
22/// addresses the same object and its thread pools are shared by every engine
23/// open in this process. Locking this slot across both acquisition and
24/// teardown is what keeps the two mutually exclusive.
25static ENV: Mutex<Weak<Env>> = Mutex::new(Weak::new());
26
27/// Take a reference to the shared environment, creating one when no context
28/// currently holds it.
29///
30/// The priority knobs apply to the environment rather than to any one context,
31/// so they are read from the config of whichever server first needs it.
32#[implement(Env)]
33pub(super) fn acquire(server: &Server) -> Result<Arc<Self>> {
34	Self::acquire_with_priorities(
35		server.config.rocksdb_compaction_prio_idle,
36		server.config.rocksdb_compaction_ioprio_idle,
37	)
38}
39
40/// Take a shared environment reference with the requested idle priorities.
41///
42/// The slot is held throughout acquisition so it cannot interleave with
43/// teardown; priorities apply only when creating the shared handle.
44#[implement(Env)]
45pub(super) fn acquire_with_priorities(cpu_idle: bool, io_idle: bool) -> Result<Arc<Self>> {
46	let mut slot = ENV.lock().expect("environment slot locked");
47
48	if let Some(env) = slot.upgrade() {
49		return Ok(env);
50	}
51
52	let mut env = rocksdb::Env::new().or_else(or_else)?;
53
54	if cpu_idle {
55		env.lower_thread_pool_cpu_priority();
56	}
57
58	if io_idle {
59		env.lower_thread_pool_io_priority();
60	}
61
62	let env = Arc::new(Self(env.into()));
63	*slot = Arc::downgrade(&env);
64
65	Ok(env)
66}
67
68#[implement(Env)]
69#[inline]
70pub(super) fn lock(&self) -> LockResult<MutexGuard<'_, rocksdb::Env>> { self.0.lock() }
71
72impl Drop for Env {
73	#[cold]
74	fn drop(&mut self) {
75		let mut slot = ENV.lock().expect("environment slot locked");
76
77		// A context which acquired after our last strong reference went away
78		// owns the same environment now, so the shutdown is its job.
79		if !ptr_eq(slot.as_ptr(), self) {
80			return;
81		}
82
83		*slot = Weak::new();
84
85		let env = self
86			.0
87			.get_mut()
88			.unwrap_or_else(PoisonError::into_inner);
89
90		debug!("Shutting down background threads");
91		env.set_high_priority_background_threads(0);
92		env.set_low_priority_background_threads(0);
93		env.set_bottom_priority_background_threads(0);
94		env.set_background_threads(0);
95
96		debug!("Joining background threads...");
97		env.join_all_threads();
98	}
99}