liminal_server/server/
runtime.rs1use std::path::Path;
2use std::sync::Arc;
3
4use crate::ServerError;
5use crate::auth_pass::PassVerifier;
6use crate::cluster::{self, ClusterHandle};
7use crate::config::file::load_config;
8use crate::config::types::{ClusterConfig, ServiceProfile};
9use crate::health::{ReadinessState, SharedReadinessState, start_health_server};
10use crate::server::connection::ConnectionSupervisor;
11use crate::server::connection::WebSocketListener;
12use crate::server::connection::services::{
13 ChannelCluster, LiminalConnectionServices, build_connection_services,
14};
15use crate::server::listener::ServerListener;
16use crate::server::shutdown::{ShutdownHandle, register_signal_handlers, run_shutdown_sequence};
17
18fn configured_authentication(
19 config: &crate::config::types::ServerConfig,
20) -> Result<(Option<Vec<u8>>, Option<PassVerifier>), ServerError> {
21 let auth_token = config
22 .auth
23 .as_ref()
24 .map(|auth| auth.token.clone().into_bytes());
25 let pass_verifier = config
26 .auth
27 .as_ref()
28 .and_then(|auth| auth.pass.as_ref())
29 .map(PassVerifier::from_config)
30 .transpose()?;
31 Ok((auth_token, pass_verifier))
32}
33
34pub fn run(config_path: &Path) -> Result<(), ServerError> {
40 if config_path.as_os_str().is_empty() {
41 return Err(ServerError::ConfigLoad {
42 message: "configuration path is empty".to_owned(),
43 });
44 }
45
46 let config = load_config(config_path)?;
47
48 crate::metrics::init();
52
53 let readiness = SharedReadinessState::new(ReadinessState::default());
54 let health_server = start_health_server(config.health_listen_address, readiness.clone())?;
55 let shutdown_handle = ShutdownHandle::new();
56 let signal_registration = register_signal_handlers(shutdown_handle.clone())?;
57
58 let (auth_token, pass_verifier) = configured_authentication(&config)?;
63
64 let (connection_supervisor, cluster_handle) = match config.services.profile()? {
72 ServiceProfile::Full => {
73 let services = Arc::new(LiminalConnectionServices::from_config(&config)?);
74 if let Some(record) = services.unloadable_conversation_record() {
84 health_server.install_unloadable_record(record);
85 }
86 if let Some(reissuer) = services.credential_reissuer() {
91 health_server.install_credential_reissuer(reissuer);
92 }
93 let channel_cluster = services.channel_cluster().clone();
94 let connection_supervisor = ConnectionSupervisor::with_fatal_shutdown(
95 services,
96 auth_token,
97 pass_verifier,
98 config.limits,
99 shutdown_handle.clone(),
100 )?;
101
102 readiness.set_cluster_configured(config.cluster.is_some());
107 let cluster_handle = match config.cluster.as_ref() {
108 Some(cluster_config) => {
109 Some(start_cluster(&channel_cluster, cluster_config, &readiness)?)
110 }
111 None => None,
112 };
113 (connection_supervisor, cluster_handle)
114 }
115 ServiceProfile::WorkerFrontDoor => {
116 let services = build_connection_services(&config)?;
117 let connection_supervisor = ConnectionSupervisor::with_fatal_shutdown(
118 services,
119 auth_token,
120 pass_verifier,
121 config.limits,
122 shutdown_handle.clone(),
123 )?;
124 readiness.set_cluster_configured(false);
125 (connection_supervisor, None)
126 }
127 };
128
129 readiness.track_admission(connection_supervisor.admission_readiness());
136
137 let mut listener = ServerListener::bind(&config, connection_supervisor)?;
138 let mut websocket_listener = match config.websocket.as_ref() {
143 Some(websocket_config) => Some(WebSocketListener::bind(
144 websocket_config,
145 listener.supervisor(),
146 )?),
147 None => None,
148 };
149 readiness.set_config_loaded(true);
150 readiness.set_listener_bound(true);
151
152 tracing::debug!(
153 config_path = %config_path.display(),
154 listen_address = %config.listen_address,
155 health_listen_address = %health_server.local_addr(),
156 "liminal server configuration validated"
157 );
158
159 tracing::info!(
160 listen_address = %listener.local_addr(),
161 health_listen_address = %health_server.local_addr(),
162 "liminal server started"
163 );
164
165 shutdown_handle.wait();
166 readiness.set_listener_bound(false);
167
168 if let Some(mut cluster_handle) = cluster_handle {
172 cluster_handle.shutdown();
173 }
174
175 let supervisor = listener.supervisor();
176 let shutdown_result = run_shutdown_sequence(
177 &mut listener,
178 websocket_listener.as_mut(),
179 &supervisor,
180 config.drain_timeout(),
181 );
182 let participant_fatal = supervisor.participant_service_fatal();
183 drop(websocket_listener);
184 drop(signal_registration);
185 health_server.shutdown()?;
186 shutdown_result?;
187 participant_fatal?.map_or(Ok(()), |fatal| {
188 Err(ServerError::ParticipantServiceFatal { fatal })
189 })
190}
191
192fn start_cluster(
204 channel_cluster: &ChannelCluster,
205 cluster_config: &ClusterConfig,
206 readiness: &SharedReadinessState,
207) -> Result<ClusterHandle, ServerError> {
208 let resolver = channel_cluster
209 .resolver()
210 .cloned()
211 .ok_or_else(|| ServerError::ClusterJoin {
212 message: "clustering configured but channel supervisor has no distribution resolver"
213 .to_owned(),
214 })?;
215 let scheduler = channel_cluster.supervisor().scheduler();
216 let supervisor = channel_cluster.supervisor().clone();
217 let readiness = readiness.clone();
218 cluster::start(
219 &scheduler,
220 resolver,
221 cluster_config,
222 move |sync| {
223 supervisor.install_observer(Arc::new(sync));
224 },
225 move || readiness.set_cluster_membership_established(true),
226 )
227}
228
229#[cfg(test)]
230mod tests {
231 use std::net::SocketAddr;
232
233 use super::{ChannelCluster, ClusterConfig, SharedReadinessState, start_cluster};
234 use crate::ServerError;
235 use crate::health::{ClusterReadiness, ReadinessCondition, ReadinessState, readiness_check};
236 use crate::server::connection::services::LiminalConnectionServices;
237
238 fn unclustered_channel_cluster() -> Result<ChannelCluster, ServerError> {
242 Ok(LiminalConnectionServices::empty()?
243 .channel_cluster()
244 .clone())
245 }
246
247 fn clustered_but_unmet_readiness() -> SharedReadinessState {
248 SharedReadinessState::new(ReadinessState::new(
249 true,
250 true,
251 ClusterReadiness::Configured {
252 membership_established: false,
253 },
254 ))
255 }
256
257 fn sample_cluster_config() -> Result<ClusterConfig, Box<dyn std::error::Error>> {
258 let listen_address: SocketAddr = "127.0.0.1:0".parse()?;
259 Ok(ClusterConfig {
260 node_name: "node-under-test@127.0.0.1".to_owned(),
261 listen_address,
262 seed_nodes: Vec::new(),
263 cookie: "runtime-test-cookie".to_owned(),
264 })
265 }
266
267 #[test]
268 fn failed_cluster_start_leaves_membership_unestablished()
269 -> Result<(), Box<dyn std::error::Error>> {
270 let readiness = clustered_but_unmet_readiness();
271 let channel_cluster = unclustered_channel_cluster()?;
272 let config = sample_cluster_config()?;
273
274 let result = start_cluster(&channel_cluster, &config, &readiness);
277 assert!(
278 result.is_err(),
279 "start_cluster must fail without a distribution resolver"
280 );
281
282 let status = readiness_check(&readiness.snapshot());
284 assert!(
285 !status.ready,
286 "readiness must remain not-ready after a failed start"
287 );
288 assert!(
289 status
290 .unmet_conditions
291 .contains(&ReadinessCondition::ClusterMembershipEstablished),
292 "cluster membership gate must stay unmet after a failed start"
293 );
294
295 Ok(())
296 }
297}