Skip to main content

tuwunel_core/
server.rs

1//! Tracks server lifecycle state and its runtime handle.
2//!
3//! The server coordinates reload, restart, and shutdown notifications. Shared
4//! services use its state to stop work promptly during teardown.
5
6mod progress;
7#[cfg(test)]
8mod tests;
9
10use std::{
11	sync::{
12		Arc,
13		atomic::{AtomicBool, Ordering},
14	},
15	time::SystemTime,
16};
17
18use ruma::OwnedServerName;
19use tokio::{runtime, sync::broadcast};
20
21pub use self::progress::Progress;
22use crate::{Err, Result, config, config::Config, log::Logging, metrics::Metrics};
23
24/// Server runtime state; public portion
25pub struct Server {
26	/// Configured name of server. This is the same as the one in the config
27	/// but developers can (and should) reference this string instead.
28	pub name: OwnedServerName,
29
30	/// Server-wide configuration instance
31	pub config: config::Manager,
32
33	/// Where the configuration came from, replayed on reload.
34	pub config_sources: config::Sources,
35
36	/// Timestamp server was started; used for uptime.
37	pub started: SystemTime,
38
39	/// Reload/shutdown pending indicator; server is shutting down. This is an
40	/// observable used on shutdown and should not be modified.
41	pub stopping: AtomicBool,
42
43	/// Reload/shutdown desired indicator; when false, shutdown is desired. This
44	/// is an observable used on shutdown and modifying is not recommended.
45	pub reloading: AtomicBool,
46
47	/// Restart desired; when true, restart it desired after shutdown.
48	pub restarting: AtomicBool,
49
50	/// Set when a backup restore is claimed, which is before it runs and is not
51	/// undone if it fails, so a database reopened later in the same process
52	/// does not restore a second time. Clearing this re-arms a destructive
53	/// operation and is never correct; claim it instead.
54	pub backup_restored: AtomicBool,
55
56	/// Handle to the runtime
57	pub runtime: Option<runtime::Handle>,
58
59	/// Reload/shutdown signal
60	pub signal: broadcast::Sender<&'static str>,
61
62	/// Logging subsystem state
63	pub log: Logging,
64
65	/// Metrics subsystem state
66	pub metrics: Arc<Metrics>,
67
68	/// Progress of the long startup phase in flight, when there is one.
69	pub progress: Progress,
70}
71
72impl Server {
73	#[must_use]
74	/// Creates shared server lifecycle state.
75	///
76	/// The initial configuration, source list, logging state, and metrics
77	/// become available to all services. A supplied runtime handle enables
78	/// task spawning.
79	pub fn new(
80		config: Config,
81		config_sources: config::Sources,
82		runtime: Option<&runtime::Handle>,
83		log: Logging,
84		metrics: Arc<Metrics>,
85	) -> Self {
86		Self {
87			name: config.server_name.clone(),
88			config: config::Manager::new(config),
89			config_sources,
90			started: SystemTime::now(),
91			stopping: AtomicBool::new(false),
92			reloading: AtomicBool::new(false),
93			restarting: AtomicBool::new(false),
94			backup_restored: AtomicBool::new(false),
95			runtime: runtime.cloned(),
96			signal: broadcast::channel::<&'static str>(1).0,
97			log,
98			metrics,
99			progress: Progress::default(),
100		}
101	}
102
103	/// Requests a dynamic module reload.
104	///
105	/// The request marks the server as reloading and stopping before
106	/// broadcasting `SIGINT`. Concurrent reload or shutdown requests are
107	/// rejected.
108	pub fn reload(&self) -> Result {
109		if cfg!(any(not(tuwunel_mods), not(feature = "tuwunel_mods"))) {
110			return Err!("Reloading not enabled");
111		}
112
113		if self.reloading.swap(true, Ordering::AcqRel) {
114			return Err!("Reloading already in progress");
115		}
116
117		if self.stopping.swap(true, Ordering::AcqRel) {
118			return Err!("Shutdown already in progress");
119		}
120
121		self.signal("SIGINT").inspect_err(|_| {
122			self.stopping.store(false, Ordering::Release);
123			self.reloading.store(false, Ordering::Release);
124		})
125	}
126
127	/// Requests a process restart through the normal shutdown path.
128	///
129	/// The restarting flag is claimed once before shutdown begins. A rejected
130	/// shutdown clears that flag so a later request may retry.
131	pub fn restart(&self) -> Result {
132		if self.restarting.swap(true, Ordering::AcqRel) {
133			return Err!("Restart already in progress");
134		}
135
136		self.shutdown().inspect_err(|_| {
137			self.restarting.store(false, Ordering::Release);
138		})
139	}
140
141	/// Requests an orderly server shutdown.
142	///
143	/// The stopping flag is claimed once before broadcasting `SIGTERM`. A
144	/// second request is rejected while shutdown remains in progress.
145	pub fn shutdown(&self) -> Result {
146		if self.stopping.swap(true, Ordering::AcqRel) {
147			return Err!("Shutdown already in progress");
148		}
149
150		self.signal("SIGTERM").inspect_err(|_| {
151			self.stopping.store(false, Ordering::Release);
152		})
153	}
154
155	/// Claims the one-shot backup restore, reporting whether this caller is the
156	/// one to perform it.
157	#[inline]
158	pub fn claim_backup_restore(&self) -> bool {
159		!self.backup_restored.swap(true, Ordering::AcqRel)
160	}
161
162	/// Broadcasts a process-signal name to lifecycle subscribers.
163	///
164	/// Delivery is best effort because subscribers may not yet be listening.
165	/// The method therefore succeeds even when the channel has no receivers.
166	pub fn signal(&self, sig: &'static str) -> Result {
167		self.signal.send(sig).ok();
168		Ok(())
169	}
170
171	#[inline]
172	/// Waits until the server enters its stopping state.
173	///
174	/// Lifecycle notifications wake the loop so it can recheck the shared
175	/// state. Calling it after shutdown has begun returns immediately.
176	pub async fn until_shutdown(self: &Arc<Self>) {
177		let mut signal = self.signal.subscribe();
178		while self.is_running() {
179			signal.recv().await.ok();
180		}
181	}
182
183	#[inline]
184	/// Returns the runtime handle supplied during server construction.
185	///
186	/// Services use this handle to spawn work on the embedding runtime. The
187	/// handle is borrowed for the lifetime of the server.
188	///
189	/// # Panics
190	///
191	/// Panics when the server was constructed without a runtime handle.
192	pub fn runtime(&self) -> &runtime::Handle {
193		self.runtime
194			.as_ref()
195			.expect("runtime handle available in Server")
196	}
197
198	#[inline]
199	/// Rejects new work after shutdown begins.
200	///
201	/// A running server returns success. A stopping server returns an
202	/// interrupted I/O error wrapped in the shared error type.
203	pub fn check_running(&self) -> Result {
204		use std::{io, io::ErrorKind::Interrupted};
205
206		self.is_running()
207			.then_some(())
208			.ok_or_else(|| io::Error::new(Interrupted, "Server shutting down"))
209			.map_err(Into::into)
210	}
211
212	#[inline]
213	/// Reports whether the server still accepts work.
214	///
215	/// Running is the inverse of the stopping state. Reload and restart
216	/// requests also transition the server through stopping.
217	pub fn is_running(&self) -> bool { !self.is_stopping() }
218
219	#[inline]
220	/// Reports whether shutdown has begun.
221	///
222	/// The flag is set by shutdown and reload transitions. Reads are relaxed
223	/// because callers use it as a lifecycle observation rather than a
224	/// synchronization edge.
225	pub fn is_stopping(&self) -> bool { self.stopping.load(Ordering::Relaxed) }
226
227	#[inline]
228	/// Reports whether a dynamic module reload is in progress.
229	///
230	/// Reload claims the flag before checking the stopping state.
231	/// Signal-delivery failure clears it, while rejection by an existing
232	/// shutdown can leave it set.
233	pub fn is_reloading(&self) -> bool { self.reloading.load(Ordering::Relaxed) }
234
235	#[inline]
236	/// Reports whether a process restart is in progress.
237	///
238	/// Restart claims the flag before requesting shutdown. Failed shutdown
239	/// initiation clears it for a later attempt.
240	pub fn is_restarting(&self) -> bool { self.restarting.load(Ordering::Relaxed) }
241
242	#[inline]
243	/// Reports whether a name matches the configured local server name.
244	///
245	/// The comparison uses the active configuration snapshot. It performs an
246	/// exact, case-sensitive string comparison.
247	pub fn is_ours(&self, name: &str) -> bool { name == self.config.server_name }
248}