use std::net::SocketAddr;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio::net::TcpStream;
use tokio::signal::unix::{SignalKind, signal};
use tokio::sync::Notify;
use crate::config::FrameConfig;
use crate::embedded::EmbeddedLiminal;
use crate::error::HostError;
use crate::server::{CONFIG_ROUTE, ShellConfig, ShellServer};
const MONITOR_INTERVAL: Duration = Duration::from_secs(1);
const PROBE_TIMEOUT: Duration = Duration::from_secs(2);
pub enum ServeStop {
Signal,
LiminalExited {
detail: String,
},
}
pub async fn serve(
config: &FrameConfig,
liminal: &EmbeddedLiminal,
liminal_endpoint: String,
) -> Result<ServeStop, HostError> {
let server = ShellServer::bind(ShellConfig {
bind: config.frame.bind,
asset_root: config.frame.assets.clone(),
liminal_endpoint,
auth_token: config.frame.auth_token.clone(),
channel: config.frame.channel.clone(),
channels: config
.liminal
.channels
.iter()
.map(|channel| channel.name.clone())
.collect(),
liminal_health: liminal.health_addr(),
})
.await?;
tracing::info!(
addr = %server.local_addr(),
assets = %config.frame.assets.display(),
liminal_tcp = %liminal.tcp_addr(),
liminal_ws = %liminal.websocket_addr(),
liminal_health = %liminal.health_addr(),
auth = if config.frame.auth_token.is_empty() { "open" } else { "bearer" },
channel = config.frame.channel.as_deref().unwrap_or("(SDK default)"),
channels = %config
.liminal
.channels
.iter()
.map(|channel| channel.name.as_str())
.collect::<Vec<_>>()
.join(", "),
config_route = CONFIG_ROUTE,
"frame server up: console shell serving; embedded liminal live; the page connects to liminal's WebSocket directly (D3)"
);
let mut sigterm =
signal(SignalKind::terminate()).map_err(|source| HostError::ShutdownSignal { source })?;
let mut sigint =
signal(SignalKind::interrupt()).map_err(|source| HostError::ShutdownSignal { source })?;
let shutdown = Arc::new(Notify::new());
let reason = Arc::new(Mutex::new(ServeStop::Signal));
let signal_shutdown = Arc::clone(&shutdown);
let signals = tokio::spawn(async move {
tokio::select! {
_ = sigterm.recv() => tracing::info!("SIGTERM received; draining frame server"),
_ = sigint.recv() => tracing::info!("SIGINT received; draining frame server"),
}
signal_shutdown.notify_waiters();
});
let monitor = tokio::spawn(liveness_monitor(
liminal.health_addr(),
liminal.tcp_addr(),
liminal.websocket_addr(),
Arc::clone(&shutdown),
Arc::clone(&reason),
));
let serve_result = server.serve(wait_for(Arc::clone(&shutdown))).await;
signals.abort();
monitor.abort();
serve_result?;
let stop = std::mem::replace(
&mut *reason
.lock()
.map_err(|_| HostError::SynchronizationPoisoned)?,
ServeStop::Signal,
);
Ok(stop)
}
async fn wait_for(shutdown: Arc<Notify>) {
shutdown.notified().await;
}
enum Liveness {
Alive,
Dead(String),
Transient(String),
}
async fn liveness_monitor(
health_addr: SocketAddr,
tcp_addr: SocketAddr,
websocket_addr: SocketAddr,
shutdown: Arc<Notify>,
reason: Arc<Mutex<ServeStop>>,
) {
let mut ticker = tokio::time::interval(MONITOR_INTERVAL);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
ticker.tick().await;
loop {
ticker.tick().await;
if let Some(detail) = probe_failure(health_addr, tcp_addr, websocket_addr).await {
tracing::error!(
detail,
"embedded liminal liveness probe found a dead listener; shutting down the frame server loudly"
);
match reason.lock() {
Ok(mut slot) => *slot = ServeStop::LiminalExited { detail },
Err(_) => {
tracing::error!("serve-reason slot poisoned while recording liminal exit");
}
}
shutdown.notify_waiters();
return;
}
}
}
async fn probe_failure(
health_addr: SocketAddr,
tcp_addr: SocketAddr,
websocket_addr: SocketAddr,
) -> Option<String> {
for (label, addr) in [
("health endpoint", health_addr),
("TCP wire listener", tcp_addr),
("WebSocket listener", websocket_addr),
] {
match probe_connect(addr).await {
Liveness::Alive => {}
Liveness::Dead(detail) => {
return Some(format!("liminal {label} at {addr}: {detail}"));
}
Liveness::Transient(detail) => {
tracing::debug!(
label,
detail,
"liminal liveness probe transient failure; not treated as death, retrying next tick"
);
}
}
}
None
}
async fn probe_connect(addr: SocketAddr) -> Liveness {
match tokio::time::timeout(PROBE_TIMEOUT, TcpStream::connect(addr)).await {
Ok(Ok(stream)) => {
drop(stream);
Liveness::Alive
}
Ok(Err(error)) if error.kind() == std::io::ErrorKind::ConnectionRefused => {
Liveness::Dead(format!("refused the connection ({error})"))
}
Ok(Err(error)) => Liveness::Transient(format!("connect errored ({error})")),
Err(_) => Liveness::Transient(format!("connect timed out after {PROBE_TIMEOUT:?}")),
}
}
#[cfg(test)]
mod tests {
use super::{Liveness, ServeStop, liveness_monitor, probe_connect, probe_failure};
use std::net::TcpListener;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio::sync::Notify;
#[tokio::test]
async fn probe_connect_classifies_alive_and_refused() -> Result<(), Box<dyn std::error::Error>>
{
let listener = TcpListener::bind("127.0.0.1:0")?;
let alive_addr = listener.local_addr()?;
assert!(matches!(probe_connect(alive_addr).await, Liveness::Alive));
let released = TcpListener::bind("127.0.0.1:0")?;
let dead_addr = released.local_addr()?;
drop(released);
assert!(matches!(probe_connect(dead_addr).await, Liveness::Dead(_)));
assert!(
probe_failure(alive_addr, alive_addr, alive_addr)
.await
.is_none()
);
assert!(
probe_failure(alive_addr, dead_addr, alive_addr)
.await
.is_some()
);
Ok(())
}
#[tokio::test]
async fn monitor_fires_liminal_exited_when_a_listener_dies()
-> Result<(), Box<dyn std::error::Error>> {
let health = TcpListener::bind("127.0.0.1:0")?;
let tcp = TcpListener::bind("127.0.0.1:0")?;
let websocket = TcpListener::bind("127.0.0.1:0")?;
let (health_addr, tcp_addr, websocket_addr) = (
health.local_addr()?,
tcp.local_addr()?,
websocket.local_addr()?,
);
let shutdown = Arc::new(Notify::new());
let reason = Arc::new(Mutex::new(ServeStop::Signal));
let monitor = tokio::spawn(liveness_monitor(
health_addr,
tcp_addr,
websocket_addr,
Arc::clone(&shutdown),
Arc::clone(&reason),
));
tokio::time::sleep(Duration::from_millis(2500)).await;
assert!(matches!(
*reason.lock().map_err(|_| "poisoned")?,
ServeStop::Signal
));
drop(tcp);
let notified = tokio::time::timeout(Duration::from_secs(6), shutdown.notified());
assert!(
notified.await.is_ok(),
"monitor must request shutdown on a dead listener"
);
assert!(matches!(
*reason.lock().map_err(|_| "poisoned")?,
ServeStop::LiminalExited { .. }
));
monitor.await?;
Ok(())
}
}