use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use liminal_server::config::ServerConfig;
use liminal_server::health::{
HealthServerHandle, ReadinessState, SharedReadinessState, start_health_server,
};
use liminal_server::server::connection::{LiminalConnectionServices, WebSocketListener};
use liminal_server::server::shutdown::run_shutdown_sequence;
use liminal_server::server::{ConnectionSupervisor, ServerListener};
use crate::error::HostError;
pub struct EmbeddedLiminal {
listener: ServerListener,
websocket: WebSocketListener,
health: HealthServerHandle,
readiness: SharedReadinessState,
drain_timeout: Duration,
tcp_addr: SocketAddr,
websocket_addr: SocketAddr,
health_addr: SocketAddr,
websocket_path: String,
}
impl EmbeddedLiminal {
pub fn boot(config: &ServerConfig) -> Result<Self, HostError> {
let websocket_config =
config
.websocket
.as_ref()
.ok_or_else(|| HostError::EmbeddedModeUnsupported {
detail: "internal invariant violated: EmbeddedLiminal::boot requires \
[bus.websocket]; FrameConfig::load must refuse its absence first"
.to_owned(),
})?;
liminal_server::metrics::init();
let readiness = SharedReadinessState::new(ReadinessState::default());
let health = start_health_server(config.health_listen_address, readiness.clone()).map_err(
|source| HostError::LiminalComponent {
component: "health endpoint",
source,
},
)?;
let health_addr = health.local_addr();
let auth_token = config
.auth
.as_ref()
.map(|auth| auth.token.clone().into_bytes());
let services = Arc::new(LiminalConnectionServices::from_config(config).map_err(
|source| HostError::LiminalComponent {
component: "connection services",
source,
},
)?);
let supervisor = ConnectionSupervisor::with_services_auth_and_limits(
services,
auth_token,
config.limits,
)
.map_err(|source| HostError::LiminalComponent {
component: "connection supervisor",
source,
})?;
let listener = ServerListener::bind(config, supervisor).map_err(|source| {
HostError::LiminalComponent {
component: "TCP listener",
source,
}
})?;
let tcp_addr = listener.local_addr();
let websocket =
WebSocketListener::bind(websocket_config, listener.supervisor()).map_err(|source| {
HostError::LiminalComponent {
component: "WebSocket listener",
source,
}
})?;
let websocket_addr = websocket.local_addr();
readiness.set_cluster_configured(false);
readiness.set_config_loaded(true);
readiness.set_listener_bound(true);
Ok(Self {
listener,
websocket,
health,
readiness,
drain_timeout: config.drain_timeout(),
tcp_addr,
websocket_addr,
health_addr,
websocket_path: websocket_config.path.clone(),
})
}
#[must_use]
pub fn websocket_endpoint(&self) -> String {
format!("ws://{}{}", self.websocket_addr, self.websocket_path)
}
#[must_use]
pub const fn tcp_addr(&self) -> SocketAddr {
self.tcp_addr
}
#[must_use]
pub const fn websocket_addr(&self) -> SocketAddr {
self.websocket_addr
}
#[must_use]
pub const fn health_addr(&self) -> SocketAddr {
self.health_addr
}
pub fn shutdown(mut self) -> Result<(), HostError> {
self.readiness.set_listener_bound(false);
let supervisor = self.listener.supervisor();
let drain = run_shutdown_sequence(
&mut self.listener,
Some(&mut self.websocket),
&supervisor,
self.drain_timeout,
)
.map_err(|source| HostError::LiminalShutdown { source });
let health = self
.health
.shutdown()
.map_err(|source| HostError::LiminalComponent {
component: "health endpoint",
source,
});
drain?;
health?;
Ok(())
}
}