liminal_server/server/
runtime.rs1use std::path::Path;
2use std::sync::Arc;
3
4use crate::ServerError;
5use crate::cluster::{self, ClusterHandle};
6use crate::config::file::load_config;
7use crate::config::types::{ClusterConfig, ServiceProfile};
8use crate::health::{ReadinessState, SharedReadinessState, start_health_server};
9use crate::server::connection::ConnectionSupervisor;
10use crate::server::connection::WebSocketListener;
11use crate::server::connection::services::{
12 ChannelCluster, LiminalConnectionServices, build_connection_services,
13};
14use crate::server::listener::ServerListener;
15use crate::server::shutdown::{ShutdownHandle, register_signal_handlers, run_shutdown_sequence};
16
17pub fn run(config_path: &Path) -> Result<(), ServerError> {
23 if config_path.as_os_str().is_empty() {
24 return Err(ServerError::ConfigLoad {
25 message: "configuration path is empty".to_owned(),
26 });
27 }
28
29 let config = load_config(config_path)?;
30
31 crate::metrics::init();
35
36 let readiness = SharedReadinessState::new(ReadinessState::default());
37 let health_server = start_health_server(config.health_listen_address, readiness.clone())?;
38 let shutdown_handle = ShutdownHandle::new();
39 let signal_registration = register_signal_handlers(shutdown_handle.clone())?;
40
41 let auth_token = config
46 .auth
47 .as_ref()
48 .map(|auth| auth.token.clone().into_bytes());
49
50 let (connection_supervisor, cluster_handle) = match config.services.profile()? {
58 ServiceProfile::Full => {
59 let services = Arc::new(LiminalConnectionServices::from_config(&config)?);
60 if let Some(record) = services.unloadable_conversation_record() {
70 health_server.install_unloadable_record(record);
71 }
72 if let Some(reissuer) = services.credential_reissuer() {
77 health_server.install_credential_reissuer(reissuer);
78 }
79 let channel_cluster = services.channel_cluster().clone();
80 let connection_supervisor = ConnectionSupervisor::with_fatal_shutdown(
81 services,
82 auth_token,
83 config.limits,
84 shutdown_handle.clone(),
85 )?;
86
87 readiness.set_cluster_configured(config.cluster.is_some());
92 let cluster_handle = match config.cluster.as_ref() {
93 Some(cluster_config) => {
94 Some(start_cluster(&channel_cluster, cluster_config, &readiness)?)
95 }
96 None => None,
97 };
98 (connection_supervisor, cluster_handle)
99 }
100 ServiceProfile::WorkerFrontDoor => {
101 let services = build_connection_services(&config)?;
102 let connection_supervisor = ConnectionSupervisor::with_fatal_shutdown(
103 services,
104 auth_token,
105 config.limits,
106 shutdown_handle.clone(),
107 )?;
108 readiness.set_cluster_configured(false);
109 (connection_supervisor, None)
110 }
111 };
112
113 readiness.track_admission(connection_supervisor.admission_readiness());
120
121 let mut listener = ServerListener::bind(&config, connection_supervisor)?;
122 let mut websocket_listener = match config.websocket.as_ref() {
127 Some(websocket_config) => Some(WebSocketListener::bind(
128 websocket_config,
129 listener.supervisor(),
130 )?),
131 None => None,
132 };
133 readiness.set_config_loaded(true);
134 readiness.set_listener_bound(true);
135
136 tracing::debug!(
137 config_path = %config_path.display(),
138 listen_address = %config.listen_address,
139 health_listen_address = %health_server.local_addr(),
140 "liminal server configuration validated"
141 );
142
143 tracing::info!(
144 listen_address = %listener.local_addr(),
145 health_listen_address = %health_server.local_addr(),
146 "liminal server started"
147 );
148
149 shutdown_handle.wait();
150 readiness.set_listener_bound(false);
151
152 if let Some(mut cluster_handle) = cluster_handle {
156 cluster_handle.shutdown();
157 }
158
159 let supervisor = listener.supervisor();
160 let shutdown_result = run_shutdown_sequence(
161 &mut listener,
162 websocket_listener.as_mut(),
163 &supervisor,
164 config.drain_timeout(),
165 );
166 let participant_fatal = supervisor.participant_service_fatal();
167 drop(websocket_listener);
168 drop(signal_registration);
169 health_server.shutdown()?;
170 shutdown_result?;
171 participant_fatal?.map_or(Ok(()), |fatal| {
172 Err(ServerError::ParticipantServiceFatal { fatal })
173 })
174}
175
176fn start_cluster(
188 channel_cluster: &ChannelCluster,
189 cluster_config: &ClusterConfig,
190 readiness: &SharedReadinessState,
191) -> Result<ClusterHandle, ServerError> {
192 let resolver = channel_cluster
193 .resolver()
194 .cloned()
195 .ok_or_else(|| ServerError::ClusterJoin {
196 message: "clustering configured but channel supervisor has no distribution resolver"
197 .to_owned(),
198 })?;
199 let scheduler = channel_cluster.supervisor().scheduler();
200 let supervisor = channel_cluster.supervisor().clone();
201 let readiness = readiness.clone();
202 cluster::start(
203 &scheduler,
204 resolver,
205 cluster_config,
206 move |sync| {
207 supervisor.install_observer(Arc::new(sync));
208 },
209 move || readiness.set_cluster_membership_established(true),
210 )
211}
212
213#[cfg(test)]
214mod tests {
215 use std::net::SocketAddr;
216
217 use super::{ChannelCluster, ClusterConfig, SharedReadinessState, start_cluster};
218 use crate::ServerError;
219 use crate::health::{ClusterReadiness, ReadinessCondition, ReadinessState, readiness_check};
220 use crate::server::connection::services::LiminalConnectionServices;
221
222 fn unclustered_channel_cluster() -> Result<ChannelCluster, ServerError> {
226 Ok(LiminalConnectionServices::empty()?
227 .channel_cluster()
228 .clone())
229 }
230
231 fn clustered_but_unmet_readiness() -> SharedReadinessState {
232 SharedReadinessState::new(ReadinessState::new(
233 true,
234 true,
235 ClusterReadiness::Configured {
236 membership_established: false,
237 },
238 ))
239 }
240
241 fn sample_cluster_config() -> Result<ClusterConfig, Box<dyn std::error::Error>> {
242 let listen_address: SocketAddr = "127.0.0.1:0".parse()?;
243 Ok(ClusterConfig {
244 node_name: "node-under-test@127.0.0.1".to_owned(),
245 listen_address,
246 seed_nodes: Vec::new(),
247 cookie: "runtime-test-cookie".to_owned(),
248 })
249 }
250
251 #[test]
252 fn failed_cluster_start_leaves_membership_unestablished()
253 -> Result<(), Box<dyn std::error::Error>> {
254 let readiness = clustered_but_unmet_readiness();
255 let channel_cluster = unclustered_channel_cluster()?;
256 let config = sample_cluster_config()?;
257
258 let result = start_cluster(&channel_cluster, &config, &readiness);
261 assert!(
262 result.is_err(),
263 "start_cluster must fail without a distribution resolver"
264 );
265
266 let status = readiness_check(&readiness.snapshot());
268 assert!(
269 !status.ready,
270 "readiness must remain not-ready after a failed start"
271 );
272 assert!(
273 status
274 .unmet_conditions
275 .contains(&ReadinessCondition::ClusterMembershipEstablished),
276 "cluster membership gate must stay unmet after a failed start"
277 );
278
279 Ok(())
280 }
281}