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::page::PageServer;
use crate::server::{CONFIG_ROUTE, DocumentAdvert, ShellConfig, ShellServer};
use crate::truth::AppTruth;
const MONITOR_INTERVAL: Duration = Duration::from_secs(1);
const PROBE_TIMEOUT: Duration = Duration::from_secs(2);
pub enum ServeStop {
Signal,
LiminalExited {
detail: String,
},
}
pub struct BoundServer {
server: ShellServer,
liminal_health_addr: SocketAddr,
liminal_tcp_addr: SocketAddr,
liminal_websocket_addr: SocketAddr,
}
pub async fn bind(
config: &FrameConfig,
liminal: &EmbeddedLiminal,
bus_endpoint: String,
truth: &Arc<AppTruth>,
page: PageServer,
) -> Result<BoundServer, HostError> {
let page_addr = page.local_addr();
let server = ShellServer::adopt(
page.into_listener(),
ShellConfig {
bind: page_addr,
asset_root: config.frame.assets.clone(),
bus_endpoint,
auth_token: config.frame.auth_token.clone(),
channel: config.frame.channel.clone(),
channels: config
.bus
.channels
.iter()
.map(|channel| channel.name.clone())
.collect(),
bus_health: liminal.health_addr(),
document: config.document.as_ref().map(|section| DocumentAdvert {
id: section.id.clone(),
component_id: section.component_id.clone(),
feed_channel: section.feed_channel.clone(),
authoring_channel: section.authoring_channel.clone(),
blink_interval_ms: section.blink_interval_ms,
}),
truth: Arc::clone(truth),
},
)?;
tracing::info!(
addr = %server.local_addr(),
assets = %config.frame.assets.display(),
bus_tcp = %liminal.tcp_addr(),
bus_ws = %liminal.websocket_addr(),
bus_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
.bus
.channels
.iter()
.map(|channel| channel.name.as_str())
.collect::<Vec<_>>()
.join(", "),
config_route = CONFIG_ROUTE,
"frame server bound: console shell reachable; embedded bus live; the page connects to \
the bus WebSocket directly (D3). The announcer pump (when [frame].channel is declared) \
starts draining its retained boot backlog only now that this surface is reachable."
);
Ok(BoundServer {
server,
liminal_health_addr: liminal.health_addr(),
liminal_tcp_addr: liminal.tcp_addr(),
liminal_websocket_addr: liminal.websocket_addr(),
})
}
impl BoundServer {
pub async fn serve(
self,
external_stop: impl Future<Output = ()> + Send + 'static,
) -> Result<ServeStop, HostError> {
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"),
() = external_stop => tracing::info!("embedder stop requested; draining frame server"),
}
signal_shutdown.notify_one();
});
let monitor = tokio::spawn(liveness_monitor(
self.liminal_health_addr,
self.liminal_tcp_addr,
self.liminal_websocket_addr,
Arc::clone(&shutdown),
Arc::clone(&reason),
));
let serve_result = self.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_one();
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(())
}
}