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}